bugfix: mq release

This commit is contained in:
云风
2012-08-17 16:49:06 +08:00
parent c8d80d58ac
commit 897a850c49
3 changed files with 45 additions and 20 deletions

View File

@@ -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;
}

View File

@@ -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

View File

@@ -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);
}