diff --git a/skynet-src/skynet_mq.c b/skynet-src/skynet_mq.c index c324d302..f36b30c1 100644 --- a/skynet-src/skynet_mq.c +++ b/skynet-src/skynet_mq.c @@ -12,8 +12,14 @@ #define DEFAULT_QUEUE_SIZE 64; #define MAX_GLOBAL_MQ 0x10000 +// 0 means mq is not in global mq. +// 1 means mq is in global mq , or the message is dispatching. +// 2 means message is dispatching with locked session set. +// 3 means mq is not in global mq, and locked session has been set. + #define MQ_IN_GLOBAL 1 -#define MQ_LOCKED 2 +#define MQ_DISPATCHING 2 +#define MQ_LOCKED 3 struct message_queue { uint32_t handle; @@ -158,11 +164,10 @@ _pushhead(struct message_queue *q, struct skynet_message *message) { // this api use in push a unlock message, so the in_global flags must not be 0 , // but the q is not exist in global queue. - if (q->in_global == MQ_IN_GLOBAL) { + if (q->in_global == MQ_LOCKED) { skynet_globalmq_push(q); - } else { - assert(q->in_global == MQ_LOCKED); - } + q->in_global = MQ_IN_GLOBAL; + } q->lock_session = 0; } @@ -222,21 +227,17 @@ skynet_mq_force_push(struct message_queue * queue) { void skynet_mq_pushglobal(struct message_queue *queue) { + LOCK(queue) assert(queue->in_global); - if (queue->in_global == MQ_LOCKED) { - // lock message queue just now, unlock it. - LOCK(queue) - queue->in_global = MQ_IN_GLOBAL; - if (queue->lock_session == 0) { - skynet_globalmq_push(queue); - } - UNLOCK(queue) - } else { - queue->in_global = MQ_IN_GLOBAL; - if (queue->lock_session == 0) { - skynet_globalmq_push(queue); - } + if (queue->in_global == MQ_DISPATCHING) { + // lock message queue just now. + queue->in_global = MQ_LOCKED; } + if (queue->lock_session == 0) { + skynet_globalmq_push(queue); + queue->in_global = MQ_IN_GLOBAL; + } + UNLOCK(queue) } void