Patch with forward-port of new cpg implementation. Needs to be included on top of my previous patch (remove....)
Regards, Honza
commit 6c02ec3cc8ab24f6dda7bf3b4b2e15ba67a78279 Author: Jan Friesse <[email protected]> Date: Tue Apr 21 14:30:08 2009 +0200 CPG patch - new CPG version diff --git a/include/corosync/corotypes.h b/include/corosync/corotypes.h index 75f45a5..5253f6f 100644 --- a/include/corosync/corotypes.h +++ b/include/corosync/corotypes.h @@ -135,6 +135,7 @@ typedef enum { #define CPG_ERR_NO_MEMORY CS_ERR_NO_MEMORY #define CPG_ERR_BAD_HANDLE CS_ERR_BAD_HANDLE #define CPG_ERR_ACCESS CS_ERR_ACCESS +#define CPG_ERR_BUSY CS_ERR_BUSY #define CPG_ERR_NOT_EXIST CS_ERR_NOT_EXIST #define CPG_ERR_EXIST CS_ERR_EXIST #define CPG_ERR_NOT_SUPPORTED CS_ERR_NOT_SUPPORTED diff --git a/include/corosync/ipc_cpg.h b/include/corosync/ipc_cpg.h index 5c9adb0..a05e37c 100644 --- a/include/corosync/ipc_cpg.h +++ b/include/corosync/ipc_cpg.h @@ -4,6 +4,7 @@ * All rights reserved. * * Author: Christine Caulfield ([email protected]) + * Author: Jan Friesse ([email protected]) * * This software licensed under BSD license, the text of which follows: * @@ -44,9 +45,7 @@ enum req_cpg_types { MESSAGE_REQ_CPG_LEAVE = 1, MESSAGE_REQ_CPG_MCAST = 2, MESSAGE_REQ_CPG_MEMBERSHIP = 3, - MESSAGE_REQ_CPG_TRACKSTART = 4, - MESSAGE_REQ_CPG_TRACKSTOP = 5, - MESSAGE_REQ_CPG_LOCAL_GET = 6 + MESSAGE_REQ_CPG_LOCAL_GET = 4, }; enum res_cpg_types { @@ -56,11 +55,9 @@ enum res_cpg_types { MESSAGE_RES_CPG_MEMBERSHIP = 3, MESSAGE_RES_CPG_CONFCHG_CALLBACK = 4, MESSAGE_RES_CPG_DELIVER_CALLBACK = 5, - MESSAGE_RES_CPG_TRACKSTART = 6, - MESSAGE_RES_CPG_TRACKSTOP = 7, - MESSAGE_RES_CPG_FLOW_CONTROL_STATE_SET = 8, - MESSAGE_RES_CPG_LOCAL_GET = 9, - MESSAGE_RES_CPG_FLOWCONTROL_CALLBACK = 10 + MESSAGE_RES_CPG_FLOW_CONTROL_STATE_SET = 6, + MESSAGE_RES_CPG_LOCAL_GET = 7, + MESSAGE_RES_CPG_FLOWCONTROL_CALLBACK = 8 }; enum lib_cpg_confchg_reason { diff --git a/include/corosync/mar_cpg.h b/include/corosync/mar_cpg.h index 06b0fd4..9640ed1 100644 --- a/include/corosync/mar_cpg.h +++ b/include/corosync/mar_cpg.h @@ -86,4 +86,13 @@ static inline void marshall_to_mar_cpg_address_t ( dest->reason = src->reason; } +static inline int mar_name_compare ( + const mar_cpg_name_t *g1, + const mar_cpg_name_t *g2) +{ + return (g1->length == g2->length? + memcmp (g1->value, g2->value, g1->length): + g1->length - g2->length); +} + #endif /* MAR_CPG_H_DEFINED */ diff --git a/lib/cpg.c b/lib/cpg.c index 10238c5..6f12030 100644 --- a/lib/cpg.c +++ b/lib/cpg.c @@ -381,33 +381,12 @@ cs_error_t cpg_join ( struct iovec iov[2]; struct req_lib_cpg_join req_lib_cpg_join; struct res_lib_cpg_join res_lib_cpg_join; - struct req_lib_cpg_trackstart req_lib_cpg_trackstart; - struct res_lib_cpg_trackstart res_lib_cpg_trackstart; error = saHandleInstanceGet (&cpg_handle_t_db, handle, (void *)&cpg_inst); if (error != CS_OK) { return (error); } - pthread_mutex_lock (&cpg_inst->response_mutex); - - /* Automatically add a tracker */ - req_lib_cpg_trackstart.header.size = sizeof (struct req_lib_cpg_trackstart); - req_lib_cpg_trackstart.header.id = MESSAGE_REQ_CPG_TRACKSTART; - marshall_to_mar_cpg_name_t (&req_lib_cpg_trackstart.group_name, - group); - - iov[0].iov_base = &req_lib_cpg_trackstart; - iov[0].iov_len = sizeof (struct req_lib_cpg_trackstart); - - error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, iov, 1, - &res_lib_cpg_trackstart, sizeof (struct res_lib_cpg_trackstart)); - - if (error != CS_OK) { - pthread_mutex_unlock (&cpg_inst->response_mutex); - goto error_exit; - } - /* Now join */ req_lib_cpg_join.header.size = sizeof (struct req_lib_cpg_join); req_lib_cpg_join.header.id = MESSAGE_REQ_CPG_JOIN; @@ -418,14 +397,18 @@ cs_error_t cpg_join ( iov[0].iov_base = &req_lib_cpg_join; iov[0].iov_len = sizeof (struct req_lib_cpg_join); - error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, iov, 1, - &res_lib_cpg_join, sizeof (struct res_lib_cpg_join)); + do { + pthread_mutex_lock (&cpg_inst->response_mutex); - pthread_mutex_unlock (&cpg_inst->response_mutex); + error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, iov, 1, + &res_lib_cpg_join, sizeof (struct res_lib_cpg_join)); - if (error != CS_OK) { - goto error_exit; - } + pthread_mutex_unlock (&cpg_inst->response_mutex); + + if (error != CS_OK) { + goto error_exit; + } + } while (res_lib_cpg_join.header.error == CPG_ERR_BUSY); error = res_lib_cpg_join.header.error; @@ -459,15 +442,18 @@ cs_error_t cpg_leave ( iov[0].iov_base = &req_lib_cpg_leave; iov[0].iov_len = sizeof (struct req_lib_cpg_leave); - pthread_mutex_lock (&cpg_inst->response_mutex); + do { + pthread_mutex_lock (&cpg_inst->response_mutex); - error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, iov, 1, - &res_lib_cpg_leave, sizeof (struct res_lib_cpg_leave)); + error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, iov, 1, + &res_lib_cpg_leave, sizeof (struct res_lib_cpg_leave)); - pthread_mutex_unlock (&cpg_inst->response_mutex); - if (error != CS_OK) { - goto error_exit; - } + pthread_mutex_unlock (&cpg_inst->response_mutex); + + if (error != CS_OK) { + goto error_exit; + } + } while (res_lib_cpg_leave.header.error == CPG_ERR_BUSY); error = res_lib_cpg_leave.header.error; @@ -554,16 +540,18 @@ cs_error_t cpg_membership_get ( iov.iov_base = &req_lib_cpg_membership_get; iov.iov_len = sizeof (mar_req_header_t); - pthread_mutex_lock (&cpg_inst->response_mutex); + do { + pthread_mutex_lock (&cpg_inst->response_mutex); - error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, &iov, 1, - &res_lib_cpg_membership_get, sizeof (mar_res_header_t)); + error = coroipcc_msg_send_reply_receive (cpg_inst->ipc_ctx, &iov, 1, + &res_lib_cpg_membership_get, sizeof (mar_res_header_t)); - pthread_mutex_unlock (&cpg_inst->response_mutex); + pthread_mutex_unlock (&cpg_inst->response_mutex); - if (error != CS_OK) { - goto error_exit; - } + if (error != CS_OK) { + goto error_exit; + } + } while (res_lib_cpg_membership_get.header.error == CPG_ERR_BUSY); error = res_lib_cpg_membership_get.header.error; diff --git a/services/cpg.c b/services/cpg.c index 88a3adf..d05946c 100644 --- a/services/cpg.c +++ b/services/cpg.c @@ -4,6 +4,7 @@ * All rights reserved. * * Author: Christine Caulfield ([email protected]) + * Author: Jan Friesse ([email protected]) * * This software licensed under BSD license, the text of which follows: * @@ -69,8 +70,6 @@ LOGSYS_DECLARE_SUBSYS ("CPG"); #define GROUP_HASH_SIZE 32 -#define PI_FLAG_MEMBER 1 - enum cpg_message_req_types { MESSAGE_REQ_EXEC_CPG_PROCJOIN = 0, MESSAGE_REQ_EXEC_CPG_PROCLEAVE = 1, @@ -79,39 +78,72 @@ enum cpg_message_req_types { MESSAGE_REQ_EXEC_CPG_DOWNLIST = 4 }; -struct removed_group -{ - struct group_info *gi; - struct list_head list; /* on removed_list */ - int left_list_entries; - mar_cpg_address_t left_list[PROCESSOR_COUNT_MAX]; - int left_list_size; +/* + * state` exec deliver + * match group name, pid -> if matched deliver for YES: + * XXX indicates impossible state + * + * join leave mcast + * UNJOINED XXX XXX NO + * LEAVE_STARTED XXX YES(unjoined_enter) YES + * JOIN_STARTED YES(join_started_enter) XXX NO + * JOIN_COMPLETED XXX NO YES + * + * join_started_enter + * set JOIN_COMPLETED + * add entry to process_info list + * unjoined_enter + * set UNJOINED + * delete entry from process_info list + * + * + * library accept join error codes + * UNJOINED YES(CPG_OK) set JOIN_STARTED + * LEAVE_STARTED NO(CPG_ERR_BUSY) + * JOIN_STARTED NO(CPG_ERR_EXIST) + * JOIN_COMPlETED NO(CPG_ERR_EXIST) + * + * library accept leave error codes + * UNJOINED NO(CPG_ERR_NOT_EXIST) + * LEAVE_STARTED NO(CPG_ERR_NOT_EXIST) + * JOIN_STARTED NO(CPG_ERR_BUSY) + * JOIN_COMPLETED YES(CPG_OK) set LEAVE_STARTED + * + * library accept mcast + * UNJOINED NO(CPG_ERR_NOT_EXIST) + * LEAVE_STARTED NO(CPG_ERR_NOT_EXIST) + * JOIN_STARTED YES(CPG_OK) + * JOIN_COMPLETED YES(CPG_OK) + */ +enum cpd_state { + CPD_STATE_UNJOINED, + CPD_STATE_LEAVE_STARTED, + CPD_STATE_JOIN_STARTED, + CPD_STATE_JOIN_COMPLETED }; -struct group_info { - mar_cpg_name_t group_name; - struct list_head members; - struct list_head list; /* on hash list */ - struct removed_group *rg; /* when a node goes down */ +struct cpg_pd { + void *conn; + mar_cpg_name_t group_name; + uint32_t pid; + enum cpd_state cpd_state; + struct list_head list; }; +DECLARE_LIST_INIT(cpg_pd_list_head); struct process_info { unsigned int nodeid; uint32_t pid; - uint32_t flags; - void *conn; - void *trackerconn; - struct group_info *group; + mar_cpg_name_t group; struct list_head list; /* on the group_info members list */ }; +DECLARE_LIST_INIT(process_info_list_head); struct join_list_entry { uint32_t pid; mar_cpg_name_t group_name; }; -static struct list_head group_lists[GROUP_HASH_SIZE]; - static struct corosync_api_v1 *api = NULL; /* @@ -167,16 +199,10 @@ static void message_handler_req_lib_cpg_mcast (void *conn, const void *message); static void message_handler_req_lib_cpg_membership (void *conn, const void *message); -static void message_handler_req_lib_cpg_trackstart (void *conn, - const void *message); - -static void message_handler_req_lib_cpg_trackstop (void *conn, - const void *message); - static void message_handler_req_lib_cpg_local_get (void *conn, const void *message); -static int cpg_node_joinleave_send (struct group_info *gi, struct process_info *pi, int fn, int reason); +static int cpg_node_joinleave_send (unsigned int pid, const mar_cpg_name_t *group_name, int fn, int reason); static int cpg_exec_send_joinlist(void); @@ -214,18 +240,6 @@ static struct corosync_lib_handler cpg_lib_engine[] = .flow_control = CS_LIB_FLOW_CONTROL_NOT_REQUIRED }, { /* 4 */ - .lib_handler_fn = message_handler_req_lib_cpg_trackstart, - .response_size = sizeof (struct res_lib_cpg_trackstart), - .response_id = MESSAGE_RES_CPG_TRACKSTART, - .flow_control = CS_LIB_FLOW_CONTROL_NOT_REQUIRED - }, - { /* 5 */ - .lib_handler_fn = message_handler_req_lib_cpg_trackstop, - .response_size = sizeof (struct res_lib_cpg_trackstart), - .response_id = MESSAGE_RES_CPG_TRACKSTOP, - .flow_control = CS_LIB_FLOW_CONTROL_NOT_REQUIRED - }, - { /* 6 */ .lib_handler_fn = message_handler_req_lib_cpg_local_get, .response_size = sizeof (struct res_lib_cpg_local_get), .response_id = MESSAGE_RES_CPG_LOCAL_GET, @@ -260,7 +274,7 @@ static struct corosync_exec_handler cpg_exec_engine[] = struct corosync_service_engine cpg_service_engine = { .name = "corosync cluster closed process group service v1.01", .id = CPG_SERVICE, - .private_data_size = sizeof (struct process_info), + .private_data_size = sizeof (struct cpg_pd), .flow_control = CS_LIB_FLOW_CONTROL_REQUIRED, .allow_inquorate = CS_LIB_ALLOW_INQUORATE, .lib_init_fn = cpg_lib_init_fn, @@ -363,7 +377,7 @@ static void cpg_sync_abort (void) static int notify_lib_joinlist( - struct group_info *gi, + const mar_cpg_name_t *group_name, void *conn, int joined_list_entries, mar_cpg_address_t *joined_list, @@ -371,167 +385,139 @@ static int notify_lib_joinlist( mar_cpg_address_t *left_list, int id) { - int count = 0; + int size; char *buf; - struct res_lib_cpg_confchg_callback *res; struct list_head *iter; - struct list_head *tmp; + int count; + struct res_lib_cpg_confchg_callback *res; mar_cpg_address_t *retgi; - int size; - /* First, we need to know how many nodes are in the list. While we're - traversing this list, look for the 'us' entry so we know which - connection to send back down */ - for (iter = gi->members.next; iter != &gi->members; iter = iter->next) { - struct process_info *pi = list_entry(iter, struct process_info, list); - if (pi->pid) - count++; - } + count = 0; + + for (iter = process_info_list_head.next; iter != &process_info_list_head; iter = iter->next) { + struct process_info *pi = list_entry (iter, struct process_info, list); + if (mar_name_compare (&pi->group, group_name) == 0) { + int i; + int founded = 0; + + for (i = 0; i < left_list_entries; i++) { + if (left_list[i].nodeid == pi->nodeid && left_list[i].pid == pi->pid) { + founded++; + } + } - log_printf(LOGSYS_LEVEL_DEBUG, "Sending new joinlist (%d elements) to clients\n", count); + if (!founded) + count++; + } + } size = sizeof(struct res_lib_cpg_confchg_callback) + sizeof(mar_cpg_address_t) * (count + left_list_entries + joined_list_entries); buf = alloca(size); if (!buf) - return CS_ERR_NO_SPACE; + return CPG_ERR_LIBRARY; res = (struct res_lib_cpg_confchg_callback *)buf; res->joined_list_entries = joined_list_entries; res->left_list_entries = left_list_entries; - retgi = res->member_list; - + res->member_list_entries = count; + retgi = res->member_list; res->header.size = size; res->header.id = id; res->header.error = CS_OK; - memcpy(&res->group_name, &gi->group_name, sizeof(mar_cpg_name_t)); + memcpy(&res->group_name, group_name, sizeof(mar_cpg_name_t)); - /* Build up the message */ - count = 0; - for (iter = gi->members.next; iter != &gi->members; iter = iter->next) { - struct process_info *pi = list_entry(iter, struct process_info, list); - if (pi->pid) { - /* Processes leaving will be removed AFTER this is done (so that they get their - own leave notifications), so exclude them from the members list here */ + for (iter = process_info_list_head.next; iter != &process_info_list_head; iter = iter->next) { + struct process_info *pi=list_entry (iter, struct process_info, list); + + if (mar_name_compare (&pi->group, group_name) == 0) { int i; - for (i=0; i<left_list_entries; i++) { - if (left_list[i].pid == pi->pid && left_list[i].nodeid == pi->nodeid) - goto next_member; + int founded = 0; + + for (i = 0;i < left_list_entries; i++) { + if (left_list[i].nodeid == pi->nodeid && left_list[i].pid == pi->pid) { + founded++; + } } - retgi->nodeid = pi->nodeid; - retgi->pid = pi->pid; - retgi++; - count++; - next_member: ; + if (!founded) { + retgi->nodeid = pi->nodeid; + retgi->pid = pi->pid; + retgi++; + } } } - res->member_list_entries = count; if (left_list_entries) { - memcpy(retgi, left_list, left_list_entries * sizeof(mar_cpg_address_t)); + memcpy (retgi, left_list, left_list_entries * sizeof(mar_cpg_address_t)); retgi += left_list_entries; } if (joined_list_entries) { - memcpy(retgi, joined_list, joined_list_entries * sizeof(mar_cpg_address_t)); + memcpy (retgi, joined_list, joined_list_entries * sizeof(mar_cpg_address_t)); retgi += joined_list_entries; } if (conn) { - api->ipc_dispatch_send(conn, buf, size); - } - else { - /* Send it to all listeners */ - for (iter = gi->members.next, tmp=iter->next; iter != &gi->members; iter = tmp, tmp=iter->next) { - struct process_info *pi = list_entry(iter, struct process_info, list); - if (pi->trackerconn && (pi->flags & PI_FLAG_MEMBER)) { - if (api->ipc_dispatch_send(pi->trackerconn, buf, size) == -1) { - // Error ?? + api->ipc_dispatch_send (conn, buf, size); + } else { + for (iter = cpg_pd_list_head.next; iter != &cpg_pd_list_head; iter = iter->next) { + struct cpg_pd *cpd = list_entry (iter, struct cpg_pd, list); + if (mar_name_compare (&cpd->group_name, group_name) == 0) { + api->ipc_dispatch_send (cpd->conn, buf, size); + assert (left_list_entries <= 1); + assert (joined_list_entries <= 1); + if (left_list_entries) { + if (left_list[0].pid == cpd->pid && + left_list[0].nodeid == api->totem_nodeid_get()) { + + cpd->pid = 0; + memset (&cpd->group_name, 0, sizeof(cpd->group_name)); + cpd->cpd_state = CPD_STATE_UNJOINED; + } + } + if (joined_list_entries) { + if (joined_list[0].pid == cpd->pid && + joined_list[0].nodeid == api->totem_nodeid_get()) { + cpd->cpd_state = CPD_STATE_JOIN_COMPLETED; + } } } } } - return CS_OK; -} - -static void remove_group(struct group_info *gi) -{ - list_del(&gi->list); - free(gi); + return CPG_OK; } - static int cpg_exec_init_fn (struct corosync_api_v1 *corosync_api) { - int i; - - for (i=0; i<GROUP_HASH_SIZE; i++) { - list_init(&group_lists[i]); - } - api = corosync_api; return (0); } static int cpg_lib_exit_fn (void *conn) { - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); - struct group_info *gi = pi->group; - mar_cpg_address_t notify_info; + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); log_printf(LOGSYS_LEVEL_DEBUG, "exit_fn for conn=%p\n", conn); - if (gi) { - notify_info.pid = pi->pid; - notify_info.nodeid = api->totem_nodeid_get(); - notify_info.reason = CONFCHG_CPG_REASON_PROCDOWN; - cpg_node_joinleave_send(gi, pi, MESSAGE_REQ_EXEC_CPG_PROCLEAVE, CONFCHG_CPG_REASON_PROCDOWN); + if (cpd->group_name.length > 0) { + cpg_node_joinleave_send (cpd->pid, &cpd->group_name, + MESSAGE_REQ_EXEC_CPG_PROCLEAVE, CONFCHG_CPG_REASON_LEAVE); } - if (pi->pid) - list_del(&pi->list); + list_del (&cpd->list); api->ipc_refcnt_dec (conn); return (0); } -static struct group_info *get_group(const mar_cpg_name_t *name) -{ - struct list_head *iter; - struct group_info *gi = NULL; - struct group_info *itergi; - uint32_t hash = jhash(name->value, name->length, 0) % GROUP_HASH_SIZE; - - for (iter = group_lists[hash].next; iter != &group_lists[hash]; iter = iter->next) { - itergi = list_entry(iter, struct group_info, list); - if (memcmp(itergi->group_name.value, name->value, name->length) == 0) { - gi = itergi; - break; - } - } - - if (!gi) { - gi = malloc(sizeof(struct group_info)); - if (!gi) { - log_printf(LOGSYS_LEVEL_WARNING, "Unable to allocate group_info struct"); - return NULL; - } - memcpy(&gi->group_name, name, sizeof(mar_cpg_name_t)); - gi->rg = NULL; - list_init(&gi->members); - list_add(&gi->list, &group_lists[hash]); - } - return gi; -} - -static int cpg_node_joinleave_send (struct group_info *gi, struct process_info *pi, int fn, int reason) +static int cpg_node_joinleave_send (unsigned int pid, const mar_cpg_name_t *group_name, int fn, int reason) { struct req_exec_cpg_procjoin req_exec_cpg_procjoin; struct iovec req_exec_cpg_iovec; int result; - memcpy(&req_exec_cpg_procjoin.group_name, &gi->group_name, sizeof(mar_cpg_name_t)); - req_exec_cpg_procjoin.pid = pi->pid; + memcpy(&req_exec_cpg_procjoin.group_name, group_name, sizeof(mar_cpg_name_t)); + req_exec_cpg_procjoin.pid = pid; req_exec_cpg_procjoin.reason = reason; req_exec_cpg_procjoin.header.size = sizeof(req_exec_cpg_procjoin); @@ -545,71 +531,6 @@ static int cpg_node_joinleave_send (struct group_info *gi, struct process_info * return (result); } -static void remove_node_from_groups( - unsigned int nodeid, - struct list_head *remlist) -{ - int i; - struct list_head *iter, *iter2, *tmp; - struct process_info *pi; - struct group_info *gi; - - for (i=0; i < GROUP_HASH_SIZE; i++) { - for (iter = group_lists[i].next; iter != &group_lists[i]; iter = iter->next) { - gi = list_entry(iter, struct group_info, list); - for (iter2 = gi->members.next, tmp = iter2->next; iter2 != &gi->members; iter2 = tmp, tmp = iter2->next) { - pi = list_entry(iter2, struct process_info, list); - - if (pi->nodeid == nodeid) { - - /* Add it to the list of nodes to send notifications for */ - if (!gi->rg) { - gi->rg = malloc(sizeof(struct removed_group)); - if (gi->rg) { - list_add(&gi->rg->list, remlist); - gi->rg->gi = gi; - gi->rg->left_list_entries = 0; - gi->rg->left_list_size = PROCESSOR_COUNT_MAX; - } - else { - log_printf(LOGSYS_LEVEL_CRIT, "Unable to allocate removed group struct. CPG callbacks will be junk."); - return; - } - } - /* Do we need to increase the size ? - * Yes, I increase this exponentially. Generally, if you've got a lot of groups, - * you'll have a /lot/ of groups, and cgp_groupinfo is pretty small anyway - */ - if (gi->rg->left_list_size == gi->rg->left_list_entries) { - int newsize; - struct removed_group *newrg; - - list_del(&gi->rg->list); - newsize = gi->rg->left_list_size * 2; - newrg = realloc(gi->rg, sizeof(struct removed_group) + newsize*sizeof(mar_cpg_address_t)); - if (!newrg) { - log_printf(LOGSYS_LEVEL_CRIT, "Unable to realloc removed group struct. CPG callbacks will be junk."); - return; - } - newrg->left_list_size = newsize+PROCESSOR_COUNT_MAX; - gi->rg = newrg; - list_add(&gi->rg->list, remlist); - } - gi->rg->left_list[gi->rg->left_list_entries].pid = pi->pid; - gi->rg->left_list[gi->rg->left_list_entries].nodeid = pi->nodeid; - gi->rg->left_list[gi->rg->left_list_entries].reason = CONFCHG_CPG_REASON_NODEDOWN; - gi->rg->left_list_entries++; - - /* Remove node info for dead node */ - list_del(&pi->list); - free(pi); - } - } - } - } -} - - static void cpg_confchg_fn ( enum totem_configuration_type configuration_type, const unsigned int *member_list, size_t member_list_entries, @@ -709,55 +630,50 @@ static void exec_cpg_mcast_endian_convert (void *msg) swab_mar_message_source_t (&req_exec_cpg_mcast->source); } +static struct process_info *process_info_find(const mar_cpg_name_t *group_name, uint32_t pid, unsigned int nodeid) { + struct list_head *iter; + + for (iter = process_info_list_head.next; iter != &process_info_list_head; ) { + struct process_info *pi = list_entry (iter, struct process_info, list); + iter = iter->next; + + if (pi->pid == pid && pi->nodeid == nodeid && + mar_name_compare (&pi->group, group_name) == 0) { + return pi; + } + } + + return NULL; +} + static void do_proc_join( const mar_cpg_name_t *name, uint32_t pid, unsigned int nodeid, int reason) { - struct group_info *gi; struct process_info *pi; - struct list_head *iter; mar_cpg_address_t notify_info; - gi = get_group(name); /* this will always succeed ! */ - assert(gi); - - /* See if it already exists in this group */ - for (iter = gi->members.next; iter != &gi->members; iter = iter->next) { - pi = list_entry(iter, struct process_info, list); - if (pi->pid == pid && pi->nodeid == nodeid) { - - /* It could be a local join message */ - if ((nodeid == api->totem_nodeid_get()) && - (!pi->flags & PI_FLAG_MEMBER)) { - goto local_join; - } else { - return; - } - } - } - - pi = malloc(sizeof(struct process_info)); + if (process_info_find (name, pid, nodeid) != NULL) { + return ; + } + pi = malloc (sizeof (struct process_info)); if (!pi) { log_printf(LOGSYS_LEVEL_WARNING, "Unable to allocate process_info struct"); return; } pi->nodeid = nodeid; pi->pid = pid; - pi->group = gi; - pi->conn = NULL; - pi->trackerconn = NULL; - list_add_tail(&pi->list, &gi->members); + memcpy(&pi->group, name, sizeof(*name)); + list_init(&pi->list); + list_add(&pi->list, &process_info_list_head); -local_join: - - pi->flags = PI_FLAG_MEMBER; notify_info.pid = pi->pid; notify_info.nodeid = nodeid; notify_info.reason = reason; - notify_lib_joinlist(gi, NULL, + notify_lib_joinlist(&pi->group, NULL, 1, ¬ify_info, 0, NULL, MESSAGE_RES_CPG_CONFCHG_CALLBACK); @@ -769,29 +685,32 @@ static void message_handler_req_exec_cpg_downlist ( { const struct req_exec_cpg_downlist *req_exec_cpg_downlist = message; int i; - struct list_head removed_list; - - log_printf(LOGSYS_LEVEL_DEBUG, "downlist left_list: %d\n", req_exec_cpg_downlist->left_nodes); - - list_init(&removed_list); - - /* Remove nodes from joined groups and add removed groups to the list */ - for (i = 0; i < req_exec_cpg_downlist->left_nodes; i++) { - remove_node_from_groups( req_exec_cpg_downlist->nodeids[i], &removed_list); - } + mar_cpg_address_t left_list[1]; + struct list_head *iter; - if (!list_empty(&removed_list)) { - struct list_head *iter, *tmp; + /* + FOR OPTIMALIZATION - Make list of lists + */ - for (iter = removed_list.next, tmp=iter->next; iter != &removed_list; iter = tmp, tmp = iter->next) { - struct removed_group *rg = list_entry(iter, struct removed_group, list); + log_printf (LOGSYS_LEVEL_DEBUG, "downlist left_list: %d\n", req_exec_cpg_downlist->left_nodes); - notify_lib_joinlist(rg->gi, NULL, - 0, NULL, - rg->left_list_entries, rg->left_list, - MESSAGE_RES_CPG_CONFCHG_CALLBACK); - rg->gi->rg = NULL; - free(rg); + for (iter = process_info_list_head.next; iter != &process_info_list_head; ) { + struct process_info *pi = list_entry(iter, struct process_info, list); + iter = iter->next; + + for (i = 0; i < req_exec_cpg_downlist->left_nodes; i++) { + if (pi->nodeid == req_exec_cpg_downlist->nodeids[i]) { + left_list[0].nodeid = pi->nodeid; + left_list[0].pid = pi->pid; + left_list[0].reason = CONFCHG_CPG_REASON_NODEDOWN; + + notify_lib_joinlist(&pi->group, NULL, + 0, NULL, + 1, left_list, + MESSAGE_RES_CPG_CONFCHG_CALLBACK); + list_del (&pi->list); + free (pi); + } } } } @@ -804,7 +723,7 @@ static void message_handler_req_exec_cpg_procjoin ( log_printf(LOGSYS_LEVEL_DEBUG, "got procjoin message from cluster node %d\n", nodeid); - do_proc_join(&req_exec_cpg_procjoin->group_name, + do_proc_join (&req_exec_cpg_procjoin->group_name, req_exec_cpg_procjoin->pid, nodeid, CONFCHG_CPG_REASON_JOIN); } @@ -814,41 +733,29 @@ static void message_handler_req_exec_cpg_procleave ( unsigned int nodeid) { const struct req_exec_cpg_procjoin *req_exec_cpg_procjoin = message; - struct group_info *gi; struct process_info *pi; struct list_head *iter; mar_cpg_address_t notify_info; log_printf(LOGSYS_LEVEL_DEBUG, "got procleave message from cluster node %d\n", nodeid); - gi = get_group(&req_exec_cpg_procjoin->group_name); /* this will always succeed ! */ - assert(gi); - notify_info.pid = req_exec_cpg_procjoin->pid; notify_info.nodeid = nodeid; notify_info.reason = req_exec_cpg_procjoin->reason; - notify_lib_joinlist(gi, NULL, - 0, NULL, - 1, ¬ify_info, - MESSAGE_RES_CPG_CONFCHG_CALLBACK); + notify_lib_joinlist(&req_exec_cpg_procjoin->group_name, NULL, + 0, NULL, + 1, ¬ify_info, + MESSAGE_RES_CPG_CONFCHG_CALLBACK); - /* Find the node/PID to remove */ - for (iter = gi->members.next; iter != &gi->members; iter = iter->next) { + for (iter = process_info_list_head.next; iter != &process_info_list_head; ) { pi = list_entry(iter, struct process_info, list); - if (pi->pid == req_exec_cpg_procjoin->pid && - pi->nodeid == nodeid) { + iter = iter->next; - list_del(&pi->list); - if (!pi->conn) - free(pi); - else - pi->pid = 0; - - if (list_empty(&gi->members)) { - remove_group(gi); - } - break; + if (pi->pid == req_exec_cpg_procjoin->pid && pi->nodeid == nodeid && + mar_name_compare (&pi->group, &req_exec_cpg_procjoin->group_name)==0) { + list_del (&pi->list); + free (pi); } } } @@ -872,7 +779,7 @@ static void message_handler_req_exec_cpg_joinlist ( } while ((const char*)jle < message + res->size) { - do_proc_join(&jle->group_name, jle->pid, nodeid, + do_proc_join (&jle->group_name, jle->pid, nodeid, CONFCHG_CPG_REASON_NODEUP); jle++; } @@ -883,40 +790,33 @@ static void message_handler_req_exec_cpg_mcast ( unsigned int nodeid) { const struct req_exec_cpg_mcast *req_exec_cpg_mcast = message; - struct res_lib_cpg_deliver_callback *res_lib_cpg_mcast; + struct res_lib_cpg_deliver_callback res_lib_cpg_mcast; int msglen = req_exec_cpg_mcast->msglen; - char buf[sizeof(*res_lib_cpg_mcast) + msglen]; - struct group_info *gi; struct list_head *iter; + struct cpg_pd *cpd; + struct iovec iovec[2]; - /* - * Track local messages so that flow is controlled on the local node - */ - gi = get_group(&req_exec_cpg_mcast->group_name); /* this will always succeed ! */ - assert(gi); - - res_lib_cpg_mcast = (struct res_lib_cpg_deliver_callback *)buf; - res_lib_cpg_mcast->header.id = MESSAGE_RES_CPG_DELIVER_CALLBACK; - res_lib_cpg_mcast->header.size = sizeof(*res_lib_cpg_mcast) + msglen; - res_lib_cpg_mcast->msglen = msglen; - res_lib_cpg_mcast->pid = req_exec_cpg_mcast->pid; - res_lib_cpg_mcast->nodeid = nodeid; - if (api->ipc_source_is_local (&req_exec_cpg_mcast->source)) { - api->ipc_refcnt_dec (req_exec_cpg_mcast->source.conn); - } - memcpy(&res_lib_cpg_mcast->group_name, &gi->group_name, + res_lib_cpg_mcast.header.id = MESSAGE_RES_CPG_DELIVER_CALLBACK; + res_lib_cpg_mcast.header.size = sizeof(res_lib_cpg_mcast) + msglen; + res_lib_cpg_mcast.msglen = msglen; + res_lib_cpg_mcast.pid = req_exec_cpg_mcast->pid; + res_lib_cpg_mcast.nodeid = nodeid; + + memcpy(&res_lib_cpg_mcast.group_name, &req_exec_cpg_mcast->group_name, sizeof(mar_cpg_name_t)); - memcpy(&res_lib_cpg_mcast->message, - (const char*)message+sizeof(*req_exec_cpg_mcast), msglen); + iovec[0].iov_base = &res_lib_cpg_mcast; + iovec[0].iov_len = sizeof (res_lib_cpg_mcast); - /* Send to all interested members */ - for (iter = gi->members.next; iter != &gi->members; iter = iter->next) { - struct process_info *pi = list_entry(iter, struct process_info, list); - if (pi->trackerconn && (pi->flags & PI_FLAG_MEMBER)) { - api->ipc_dispatch_send( - pi->trackerconn, - buf, - res_lib_cpg_mcast->header.size); + iovec[1].iov_base = (char*)message+sizeof(*req_exec_cpg_mcast); + iovec[1].iov_len = msglen; + + for (iter = cpg_pd_list_head.next; iter != &cpg_pd_list_head; ) { + cpd = list_entry(iter, struct cpg_pd, list); + iter = iter->next; + + if ((cpd->cpd_state == CPD_STATE_LEAVE_STARTED || cpd->cpd_state == CPD_STATE_JOIN_COMPLETED) + && (mar_name_compare (&cpd->group_name, &req_exec_cpg_mcast->group_name) == 0)) { + api->ipc_dispatch_iov_send (cpd->conn, iovec, 2); } } } @@ -925,27 +825,17 @@ static void message_handler_req_exec_cpg_mcast ( static int cpg_exec_send_joinlist(void) { int count = 0; - char *buf; - int i; struct list_head *iter; - struct list_head *iter2; - struct group_info *gi; mar_res_header_t *res; struct join_list_entry *jle; + char *buf; struct iovec req_exec_cpg_iovec; - log_printf(LOGSYS_LEVEL_DEBUG, "sending joinlist to cluster\n"); + for (iter = process_info_list_head.next; iter != &process_info_list_head; iter = iter->next) { + struct process_info *pi = list_entry (iter, struct process_info, list); - /* Count the number of groups we are a member of */ - for (i=0; i<GROUP_HASH_SIZE; i++) { - for (iter = group_lists[i].next; iter != &group_lists[i]; iter = iter->next) { - gi = list_entry(iter, struct group_info, list); - for (iter2 = gi->members.next; iter2 != &gi->members; iter2 = iter2->next) { - struct process_info *pi = list_entry(iter2, struct process_info, list); - if (pi->pid && pi->nodeid == api->totem_nodeid_get()) { - count++; - } - } + if (pi->nodeid == api->totem_nodeid_get ()) { + count++; } } @@ -953,7 +843,7 @@ static int cpg_exec_send_joinlist(void) if (!count) return 0; - buf = alloca(sizeof(mar_res_header_t) + sizeof(struct join_list_entry) * count); + buf = alloca(sizeof (mar_res_header_t) + sizeof (struct join_list_entry) * count); if (!buf) { log_printf(LOGSYS_LEVEL_WARNING, "Unable to allocate joinlist buffer"); return -1; @@ -962,24 +852,18 @@ static int cpg_exec_send_joinlist(void) jle = (struct join_list_entry *)(buf + sizeof(mar_res_header_t)); res = (mar_res_header_t *)buf; - for (i=0; i<GROUP_HASH_SIZE; i++) { - for (iter = group_lists[i].next; iter != &group_lists[i]; iter = iter->next) { - - gi = list_entry(iter, struct group_info, list); - for (iter2 = gi->members.next; iter2 != &gi->members; iter2 = iter2->next) { + for (iter = process_info_list_head.next; iter != &process_info_list_head; iter = iter->next) { + struct process_info *pi = list_entry (iter, struct process_info, list); - struct process_info *pi = list_entry(iter2, struct process_info, list); - if (pi->pid && pi->nodeid == api->totem_nodeid_get()) { - memcpy(&jle->group_name, &gi->group_name, sizeof(mar_cpg_name_t)); - jle->pid = pi->pid; - jle++; - } - } + if (pi->nodeid == api->totem_nodeid_get ()) { + memcpy (&jle->group_name, &pi->group, sizeof (mar_cpg_name_t)); + jle->pid = pi->pid; + jle++; } } res->id = SERVICE_ID_MAKE(CPG_SERVICE, MESSAGE_REQ_EXEC_CPG_JOINLIST); - res->size = sizeof(mar_res_header_t)+sizeof(struct join_list_entry) * count; + res->size = sizeof (mar_res_header_t) + sizeof (struct join_list_entry) * count; req_exec_cpg_iovec.iov_base = buf; req_exec_cpg_iovec.iov_len = res->size; @@ -989,11 +873,13 @@ static int cpg_exec_send_joinlist(void) static int cpg_lib_init_fn (void *conn) { - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); - api->ipc_refcnt_inc (conn); - pi->conn = conn; + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); + memset (cpd, 0, sizeof(struct cpg_pd)); + cpd->conn = conn; + list_add (&cpd->list, &cpg_pd_list_head); - log_printf(LOGSYS_LEVEL_DEBUG, "lib_init_fn: conn=%p, pi=%p\n", conn, pi); + api->ipc_refcnt_inc (conn); + log_printf(LOGSYS_LEVEL_DEBUG, "lib_init_fn: conn=%p, cpd=%p\n", conn, cpd); return (0); } @@ -1001,63 +887,69 @@ static int cpg_lib_init_fn (void *conn) static void message_handler_req_lib_cpg_join (void *conn, const void *message) { const struct req_lib_cpg_join *req_lib_cpg_join = message; - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); struct res_lib_cpg_join res_lib_cpg_join; - struct group_info *gi; - cs_error_t error = CS_OK; - - log_printf(LOGSYS_LEVEL_DEBUG, "got join request on %p, pi=%p, pi->pid=%d\n", conn, pi, pi->pid); - - /* Already joined on this conn */ - if (pi->pid) { - error = CS_ERR_INVALID_PARAM; - goto join_err; - } - - gi = get_group(&req_lib_cpg_join->group_name); - if (!gi) { - error = CS_ERR_NO_SPACE; - goto join_err; + cs_error_t error = CPG_OK; + + switch (cpd->cpd_state) { + case CPD_STATE_UNJOINED: + error = CPG_OK; + cpd->cpd_state = CPD_STATE_JOIN_STARTED; + cpd->pid = req_lib_cpg_join->pid; + memcpy (&cpd->group_name, &req_lib_cpg_join->group_name, + sizeof (cpd->group_name)); + + cpg_node_joinleave_send (req_lib_cpg_join->pid, + &req_lib_cpg_join->group_name, + MESSAGE_REQ_EXEC_CPG_PROCJOIN, CONFCHG_CPG_REASON_JOIN); + break; + case CPD_STATE_LEAVE_STARTED: + error = CPG_ERR_BUSY; + break; + case CPD_STATE_JOIN_STARTED: + error = CPG_ERR_EXIST; + break; + case CPD_STATE_JOIN_COMPLETED: + error = CPG_ERR_EXIST; + break; } - /* Add a node entry for us */ - pi->nodeid = api->totem_nodeid_get(); - pi->pid = req_lib_cpg_join->pid; - pi->group = gi; - list_add(&pi->list, &gi->members); - - /* Tell the rest of the cluster */ - cpg_node_joinleave_send(gi, pi, MESSAGE_REQ_EXEC_CPG_PROCJOIN, CONFCHG_CPG_REASON_JOIN); - -join_err: res_lib_cpg_join.header.size = sizeof(res_lib_cpg_join); - res_lib_cpg_join.header.id = MESSAGE_RES_CPG_JOIN; - res_lib_cpg_join.header.error = error; - api->ipc_response_send(conn, &res_lib_cpg_join, sizeof(res_lib_cpg_join)); + res_lib_cpg_join.header.id = MESSAGE_RES_CPG_JOIN; + res_lib_cpg_join.header.error = error; + api->ipc_response_send (conn, &res_lib_cpg_join, sizeof(res_lib_cpg_join)); } /* Leave message from the library */ static void message_handler_req_lib_cpg_leave (void *conn, const void *message) { - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); struct res_lib_cpg_leave res_lib_cpg_leave; - struct group_info *gi; - cs_error_t error = CS_OK; + cs_error_t error = CPG_OK; + struct req_lib_cpg_leave *req_lib_cpg_leave = (struct req_lib_cpg_leave *)message; + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); log_printf(LOGSYS_LEVEL_DEBUG, "got leave request on %p\n", conn); - if (!pi || !pi->pid || !pi->group) { - error = CS_ERR_INVALID_PARAM; - goto leave_ret; + switch (cpd->cpd_state) { + case CPD_STATE_UNJOINED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_LEAVE_STARTED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_JOIN_STARTED: + error = CPG_ERR_BUSY; + break; + case CPD_STATE_JOIN_COMPLETED: + error = CPG_OK; + cpd->cpd_state = CPD_STATE_LEAVE_STARTED; + cpg_node_joinleave_send (req_lib_cpg_leave->pid, + &req_lib_cpg_leave->group_name, + MESSAGE_REQ_EXEC_CPG_PROCLEAVE, + CONFCHG_CPG_REASON_LEAVE); + break; } - gi = pi->group; - /* Tell other nodes we are leaving. - When we get this message back we will leave too */ - cpg_node_joinleave_send(gi, pi, MESSAGE_REQ_EXEC_CPG_PROCLEAVE, CONFCHG_CPG_REASON_LEAVE); - pi->group = NULL; - -leave_ret: /* send return */ res_lib_cpg_leave.header.size = sizeof(res_lib_cpg_leave); res_lib_cpg_leave.header.id = MESSAGE_RES_CPG_LEAVE; @@ -1069,124 +961,92 @@ leave_ret: static void message_handler_req_lib_cpg_mcast (void *conn, const void *message) { const struct req_lib_cpg_mcast *req_lib_cpg_mcast = message; - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); - struct group_info *gi = pi->group; + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); + mar_cpg_name_t group_name = cpd->group_name; + struct iovec req_exec_cpg_iovec[2]; struct req_exec_cpg_mcast req_exec_cpg_mcast; struct res_lib_cpg_mcast res_lib_cpg_mcast; int msglen = req_lib_cpg_mcast->msglen; int result; + cs_error_t error = CPG_ERR_NOT_EXIST; log_printf(LOGSYS_LEVEL_DEBUG, "got mcast request on %p\n", conn); - /* Can't send if we're not joined */ - if (!gi) { - res_lib_cpg_mcast.header.size = sizeof(res_lib_cpg_mcast); - res_lib_cpg_mcast.header.id = MESSAGE_RES_CPG_MCAST; - res_lib_cpg_mcast.header.error = CS_ERR_ACCESS; /* TODO Better error code ?? */ - api->ipc_response_send(conn, &res_lib_cpg_mcast, - sizeof(res_lib_cpg_mcast)); - return; + switch (cpd->cpd_state) { + case CPD_STATE_UNJOINED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_LEAVE_STARTED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_JOIN_STARTED: + error = CPG_OK; + break; + case CPD_STATE_JOIN_COMPLETED: + error = CPG_OK; + break; } - req_exec_cpg_mcast.header.size = sizeof(req_exec_cpg_mcast) + msglen; - req_exec_cpg_mcast.header.id = SERVICE_ID_MAKE(CPG_SERVICE, - MESSAGE_REQ_EXEC_CPG_MCAST); - req_exec_cpg_mcast.pid = pi->pid; - req_exec_cpg_mcast.msglen = msglen; - api->ipc_source_set (&req_exec_cpg_mcast.source, conn); - memcpy(&req_exec_cpg_mcast.group_name, &gi->group_name, - sizeof(mar_cpg_name_t)); - - req_exec_cpg_iovec[0].iov_base = (char *)&req_exec_cpg_mcast; - req_exec_cpg_iovec[0].iov_len = sizeof(req_exec_cpg_mcast); - req_exec_cpg_iovec[1].iov_base = (char *)&req_lib_cpg_mcast->message; - req_exec_cpg_iovec[1].iov_len = msglen; - - // TODO: guarantee type... - result = api->totem_mcast (req_exec_cpg_iovec, 2, TOTEM_AGREED); - api->ipc_refcnt_inc (conn); + if (error == CPG_OK) { + req_exec_cpg_mcast.header.size = sizeof(req_exec_cpg_mcast) + msglen; + req_exec_cpg_mcast.header.id = SERVICE_ID_MAKE(CPG_SERVICE, + MESSAGE_REQ_EXEC_CPG_MCAST); + req_exec_cpg_mcast.pid = cpd->pid; + req_exec_cpg_mcast.msglen = msglen; + api->ipc_source_set (&req_exec_cpg_mcast.source, conn); + memcpy(&req_exec_cpg_mcast.group_name, &group_name, + sizeof(mar_cpg_name_t)); + + req_exec_cpg_iovec[0].iov_base = (char *)&req_exec_cpg_mcast; + req_exec_cpg_iovec[0].iov_len = sizeof(req_exec_cpg_mcast); + req_exec_cpg_iovec[1].iov_base = (char *)&req_lib_cpg_mcast->message; + req_exec_cpg_iovec[1].iov_len = msglen; + + result = api->totem_mcast (req_exec_cpg_iovec, 2, TOTEM_AGREED); + assert(result == 0); + } res_lib_cpg_mcast.header.size = sizeof(res_lib_cpg_mcast); res_lib_cpg_mcast.header.id = MESSAGE_RES_CPG_MCAST; - res_lib_cpg_mcast.header.error = CS_OK; - api->ipc_response_send(conn, &res_lib_cpg_mcast, - sizeof(res_lib_cpg_mcast)); + res_lib_cpg_mcast.header.error = error; + api->ipc_response_send (conn, &res_lib_cpg_mcast, + sizeof (res_lib_cpg_mcast)); } static void message_handler_req_lib_cpg_membership (void *conn, const void *message) { - struct process_info *pi = (struct process_info *)api->ipc_private_data_get (conn); - - log_printf(LOGSYS_LEVEL_DEBUG, "got membership request on %p\n", conn); - if (!pi->group) { - mar_res_header_t res; - res.size = sizeof(res); - res.id = MESSAGE_RES_CPG_MEMBERSHIP; - res.error = CS_ERR_ACCESS; /* TODO Better error code */ - api->ipc_response_send(conn, &res, sizeof(res)); - return; - } - - notify_lib_joinlist(pi->group, conn, 0, NULL, 0, NULL, MESSAGE_RES_CPG_MEMBERSHIP); -} - - -static void message_handler_req_lib_cpg_trackstart (void *conn, - const void *message) -{ - const struct req_lib_cpg_trackstart *req_lib_cpg_trackstart = message; - struct res_lib_cpg_trackstart res_lib_cpg_trackstart; - struct group_info *gi; - struct process_info *otherpi; - cs_error_t error = CS_OK; - - log_printf(LOGSYS_LEVEL_DEBUG, "got trackstart request on %p\n", conn); - - gi = get_group(&req_lib_cpg_trackstart->group_name); - if (!gi) { - error = CS_ERR_NO_SPACE; - goto tstart_ret; + struct cpg_pd *cpd = (struct cpg_pd *)api->ipc_private_data_get (conn); + cs_error_t error = CPG_ERR_NOT_EXIST; + mar_res_header_t res; + + switch (cpd->cpd_state) { + case CPD_STATE_UNJOINED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_LEAVE_STARTED: + error = CPG_ERR_NOT_EXIST; + break; + case CPD_STATE_JOIN_STARTED: + error = CPG_ERR_BUSY; + break; + case CPD_STATE_JOIN_COMPLETED: + error = CPG_OK; + break; } - /* Find the partner connection and add us to it's process_info struct */ - otherpi = (struct process_info *)api->ipc_private_data_get (conn); - otherpi->trackerconn = conn; - -tstart_ret: - res_lib_cpg_trackstart.header.size = sizeof(res_lib_cpg_trackstart); - res_lib_cpg_trackstart.header.id = MESSAGE_RES_CPG_TRACKSTART; - res_lib_cpg_trackstart.header.error = CS_OK; - api->ipc_response_send(conn, &res_lib_cpg_trackstart, sizeof(res_lib_cpg_trackstart)); -} + res.size = sizeof (res); + res.id = MESSAGE_RES_CPG_MEMBERSHIP; + res.error = error; + api->ipc_response_send (conn, &res, sizeof(res)); + return; -static void message_handler_req_lib_cpg_trackstop (void *conn, - const void *message) -{ - const struct req_lib_cpg_trackstop *req_lib_cpg_trackstop = message; - struct res_lib_cpg_trackstop res_lib_cpg_trackstop; - struct process_info *otherpi; - struct group_info *gi; - cs_error_t error = CS_OK; - - log_printf(LOGSYS_LEVEL_DEBUG, "got trackstop request on %p\n", conn); - - gi = get_group(&req_lib_cpg_trackstop->group_name); - if (!gi) { - error = CS_ERR_NO_SPACE; - goto tstop_ret; + if (error == CPG_OK) { + notify_lib_joinlist (&cpd->group_name, conn, 0, NULL, 0, NULL, + MESSAGE_RES_CPG_MEMBERSHIP); } - /* Find the partner connection and add us to it's process_info struct */ - otherpi = (struct process_info *)api->ipc_private_data_get (conn); - otherpi->trackerconn = NULL; - -tstop_ret: - res_lib_cpg_trackstop.header.size = sizeof(res_lib_cpg_trackstop); - res_lib_cpg_trackstop.header.id = MESSAGE_RES_CPG_TRACKSTOP; - res_lib_cpg_trackstop.header.error = CS_OK; - api->ipc_response_send(conn, &res_lib_cpg_trackstop.header, sizeof(res_lib_cpg_trackstop)); } static void message_handler_req_lib_cpg_local_get (void *conn, @@ -1194,11 +1054,11 @@ static void message_handler_req_lib_cpg_local_get (void *conn, { struct res_lib_cpg_local_get res_lib_cpg_local_get; - res_lib_cpg_local_get.header.size = sizeof(res_lib_cpg_local_get); + res_lib_cpg_local_get.header.size = sizeof (res_lib_cpg_local_get); res_lib_cpg_local_get.header.id = MESSAGE_RES_CPG_LOCAL_GET; - res_lib_cpg_local_get.header.error = CS_OK; + res_lib_cpg_local_get.header.error = CPG_OK; res_lib_cpg_local_get.local_nodeid = api->totem_nodeid_get (); - api->ipc_response_send(conn, &res_lib_cpg_local_get, - sizeof(res_lib_cpg_local_get)); + api->ipc_response_send (conn, &res_lib_cpg_local_get, + sizeof (res_lib_cpg_local_get)); }
_______________________________________________ Openais mailing list [email protected] https://lists.linux-foundation.org/mailman/listinfo/openais
