This patch adds the calling process' pid to saMsgMessageGet and
saMsgMessageCancel operatioin in the message service. The reason is
simple: the cancel operation will cancel all pending get operations
for that process, so we need to identify the process along with nodeid
and ipc connection.

Ryan

Index: include/ipc_msg.h
===================================================================
--- include/ipc_msg.h   (revision 1875)
+++ include/ipc_msg.h   (working copy)
@@ -263,6 +263,7 @@
        coroipc_request_header_t header;
        SaNameT queue_name;
        SaUint32T queue_id;
+       SaUint32T pid;
        SaTimeT timeout;
 };
 
@@ -285,6 +286,7 @@
        coroipc_request_header_t header;
        SaNameT queue_name;
        SaUint32T queue_id;
+       SaUint32T pid;
 };
 
 struct res_lib_msg_messagecancel {
Index: lib/msg.c
===================================================================
--- lib/msg.c   (revision 1875)
+++ lib/msg.c   (working copy)
@@ -1391,14 +1391,14 @@
                sizeof (struct req_lib_msg_messageget);
        req_lib_msg_messageget.header.id =
                MESSAGE_REQ_MSG_MESSAGEGET;
-       req_lib_msg_messageget.queue_id =
-               queueInstance->queue_id;
 
+       req_lib_msg_messageget.queue_id = queueInstance->queue_id;
+       req_lib_msg_messageget.pid = (SaUint32T)(getpid());
+       req_lib_msg_messageget.timeout = timeout;
+
        memcpy (&req_lib_msg_messageget.queue_name,
                &queueInstance->queue_name, sizeof (SaNameT));
 
-       req_lib_msg_messageget.timeout = timeout;
-
        iov.iov_base = &req_lib_msg_messageget;
        iov.iov_len = sizeof (struct req_lib_msg_messageget);
 
@@ -1499,9 +1499,10 @@
                sizeof (struct req_lib_msg_messagecancel);
        req_lib_msg_messagecancel.header.id =
                MESSAGE_REQ_MSG_MESSAGECANCEL;
-       req_lib_msg_messagecancel.queue_id =
-               queueInstance->queue_id;
 
+       req_lib_msg_messagecancel.queue_id = queueInstance->queue_id;
+       req_lib_msg_messagecancel.pid = (SaUint32T)(getpid());
+
        memcpy (&req_lib_msg_messagecancel.queue_name,
                &queueInstance->queue_name, sizeof (SaNameT));
 
Index: services/msg.c
===================================================================
--- services/msg.c      (revision 1875)
+++ services/msg.c      (working copy)
@@ -157,6 +157,7 @@
        mar_message_source_t source;
        corosync_timer_handle_t timer_handle;
        SaNameT queue_name;
+       SaUint32T pid;
        struct list_head list;
 };
 
@@ -867,6 +868,7 @@
        mar_message_source_t source;
        SaNameT queue_name;
        SaUint32T queue_id;
+       SaUint32T pid;
        SaTimeT timeout;
 };
 
@@ -880,6 +882,7 @@
        mar_message_source_t source;
        SaNameT queue_name;
        SaUint32T queue_id;
+       SaUint32T pid;
 };
 
 struct req_exec_msg_messagesendreceive {
@@ -1727,8 +1730,13 @@
        struct pending_entry *pending = (struct pending_entry *)data;
 
        /* DEBUG */
-       log_printf (LOGSYS_LEVEL_DEBUG, "[DEBUG]: msg_expire_pending { %p }\n", 
data);
-       log_printf (LOGSYS_LEVEL_DEBUG, "[DEBUG]:\t queue = %s\n", (char 
*)(pending->queue_name.value));
+       log_printf (LOGSYS_LEVEL_DEBUG, "[DEBUG]: msg_expire_pending\n");
+       log_printf (LOGSYS_LEVEL_DEBUG, "[DEBUG]:\t queue = %s\n",
+                   (char *)(pending->queue_name.value));
+       log_printf (LOGSYS_LEVEL_DEBUG, "[DEBUG]:\t pending = { nodeid=%x 
pid=%u conn=%p }\n",
+                   (unsigned int)(pending->source.nodeid),
+                   (unsigned int)(pending->pid),
+                   (void *)(pending->source.conn));
 
        req_exec_msg_pending_timeout.header.size =
                sizeof (struct req_exec_msg_pending_timeout);
@@ -3759,6 +3767,14 @@
                        &req_exec_msg_messageget->queue_name,
                        sizeof (SaNameT));
 
+               get->pid = req_exec_msg_messageget->pid;
+
+               /* DEBUG */
+               log_printf (LOGSYS_LEVEL_DEBUG, "\t pending = { nodeid=%x 
pid=%u conn=%p }\n",
+                           (unsigned int)(get->source.nodeid),
+                           (unsigned int)(get->pid),
+                           (void *)(get->source.conn));
+
                list_add_tail (&get->list, &queue->pending_head);
 
                if (api->ipc_source_is_local 
(&req_exec_msg_messageget->source)) {
@@ -5168,6 +5184,8 @@
 
        req_exec_msg_messageget.queue_id =
                req_lib_msg_messageget->queue_id;
+       req_exec_msg_messageget.pid =
+               req_lib_msg_messageget->pid;
        req_exec_msg_messageget.timeout =
                req_lib_msg_messageget->timeout;
 
@@ -5221,6 +5239,8 @@
 
        req_exec_msg_messagecancel.queue_id =
                req_lib_msg_messagecancel->queue_id;
+       req_exec_msg_messagecancel.pid =
+               req_lib_msg_messagecancel->pid;
 
        iovec.iov_base = (char *)&req_exec_msg_messagecancel;
        iovec.iov_len = sizeof (req_exec_msg_messagecancel);
_______________________________________________
Openais mailing list
[email protected]
https://lists.linux-foundation.org/mailman/listinfo/openais

Reply via email to