From 897a850c490a5aa78ea75c0df9a46f6d7b7747fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BA=91=E9=A3=8E?= Date: Fri, 17 Aug 2012 16:49:06 +0800 Subject: [PATCH] bugfix: mq release --- skynet-src/skynet_mq.c | 45 ++++++++++++++++++++++++++++++++++---- skynet-src/skynet_mq.h | 3 ++- skynet-src/skynet_server.c | 17 ++------------ 3 files changed, 45 insertions(+), 20 deletions(-) diff --git a/skynet-src/skynet_mq.c b/skynet-src/skynet_mq.c index 5d96543c..bdd3eae5 100644 --- a/skynet-src/skynet_mq.c +++ b/skynet-src/skynet_mq.c @@ -13,6 +13,7 @@ struct message_queue { int head; int tail; int lock; + int release; int in_global; struct skynet_message *queue; }; @@ -76,7 +77,7 @@ skynet_globalmq_pop() { } UNLOCK(q) - + return ret; } @@ -89,13 +90,14 @@ skynet_mq_create(uint32_t handle) { q->tail = 0; q->lock = 0; q->in_global = 1; + q->release = 0; q->queue = malloc(sizeof(struct skynet_message) * q->cap); return q; } -void -skynet_mq_release(struct message_queue *q) { +static void +_release(struct message_queue *q) { free(q->queue); free(q); } @@ -154,8 +156,8 @@ skynet_mq_push(struct message_queue *q, struct skynet_message *message) { } if (q->in_global == 0) { - skynet_globalmq_push(q); q->in_global = 1; + skynet_globalmq_push(q); } UNLOCK(q) @@ -180,3 +182,38 @@ skynet_mq_force_push(struct message_queue * queue) { assert(queue->in_global); skynet_globalmq_push(queue); } + +void +skynet_mq_mark_release(struct message_queue *q) { + assert(q->release == 0); + q->release = 1; +} + +static int +_drop_queue(struct message_queue *q) { + // todo: send message back to message source + struct skynet_message msg; + int s = 0; + while(!skynet_mq_pop(q, &msg)) { + ++s; + free(msg.data); + } + _release(q); + return s; +} + +int +skynet_mq_release(struct message_queue *q) { + int ret = 0; + LOCK(q) + + if (q->release) { + ret = _drop_queue(q); + } else { + skynet_mq_force_push(q); + } + + UNLOCK(q) + + return ret; +} diff --git a/skynet-src/skynet_mq.h b/skynet-src/skynet_mq.h index 375700e1..c024250f 100644 --- a/skynet-src/skynet_mq.h +++ b/skynet-src/skynet_mq.h @@ -16,7 +16,8 @@ struct message_queue; struct message_queue * skynet_globalmq_pop(void); struct message_queue * skynet_mq_create(uint32_t handle); -void skynet_mq_release(struct message_queue *); +void skynet_mq_mark_release(struct message_queue *q); +int skynet_mq_release(struct message_queue *q); uint32_t skynet_mq_handle(struct message_queue *); // 0 for success diff --git a/skynet-src/skynet_server.c b/skynet-src/skynet_server.c index c2e092ca..1786c255 100644 --- a/skynet-src/skynet_server.c +++ b/skynet-src/skynet_server.c @@ -102,7 +102,7 @@ skynet_context_grab(struct skynet_context *ctx) { static void _delete_context(struct skynet_context *ctx) { skynet_module_instance_release(ctx->mod, ctx->instance); - skynet_mq_push(ctx->queue , NULL); + skynet_mq_mark_release(ctx->queue); free(ctx); } @@ -176,19 +176,6 @@ _dispatch_message(struct skynet_context *ctx, struct skynet_message *msg) { } } -static int -_drop_queue(struct message_queue *q) { - // todo: send message back to message source - struct skynet_message msg; - int s = 0; - while(!skynet_mq_pop(q, &msg)) { - ++s; - free(msg.data); - } - skynet_mq_release(q); - return s; -} - int skynet_context_message_dispatch(void) { struct message_queue * q = skynet_globalmq_pop(); @@ -199,7 +186,7 @@ skynet_context_message_dispatch(void) { struct skynet_context * ctx = skynet_handle_grab(handle); if (ctx == NULL) { - int s = _drop_queue(q); + int s = skynet_mq_release(q); if (s>0) { skynet_error(NULL, "Drop message queue %x (%d messages)", handle,s); }