multicast bugfix: don't hold context any more

This commit is contained in:
云风
2012-10-12 12:15:10 +08:00
parent 4fb7fec2fa
commit a93e6b8868
4 changed files with 51 additions and 55 deletions

View File

@@ -23,9 +23,13 @@ skynet_multicast_create(const void * msg, size_t sz, uint32_t source) {
return mc; return mc;
} }
void void
skynet_multicast_copy(struct skynet_multicast_message *mc, int copy) { skynet_multicast_copy(struct skynet_multicast_message *mc, int copy) {
__sync_fetch_and_add(&mc->ref, copy); int r = __sync_add_and_fetch(&mc->ref, copy);
if (r == 0) {
free((void *)mc->msg);
free(mc);
}
} }
void void
@@ -46,17 +50,12 @@ struct array {
uint32_t *data; uint32_t *data;
}; };
struct pair {
uint32_t handle;
struct skynet_context * ctx;
};
struct skynet_multicast_group { struct skynet_multicast_group {
struct array enter_queue; struct array enter_queue;
struct array leave_queue; struct array leave_queue;
int cap; int cap;
int number; int number;
struct pair * data; uint32_t * data;
}; };
struct skynet_multicast_group * struct skynet_multicast_group *
@@ -68,12 +67,6 @@ skynet_multicast_newgroup() {
void void
skynet_multicast_deletegroup(struct skynet_multicast_group * g) { skynet_multicast_deletegroup(struct skynet_multicast_group * g) {
int i;
for (i=0;i<g->number;i++) {
if (g->data[i].ctx) {
skynet_context_release(g->data[i].ctx);
}
}
free(g->data); free(g->data);
free(g->enter_queue.data); free(g->enter_queue.data);
free(g->leave_queue.data); free(g->leave_queue.data);
@@ -120,7 +113,7 @@ combine_queue(struct skynet_context * from, struct skynet_multicast_group * grou
int new_size = group->number + enter; int new_size = group->number + enter;
if (new_size > group->cap) { if (new_size > group->cap) {
group->data = realloc(group->data, new_size * sizeof(struct pair)); group->data = realloc(group->data, new_size * sizeof(uint32_t));
group->cap = new_size; group->cap = new_size;
} }
@@ -132,19 +125,14 @@ combine_queue(struct skynet_context * from, struct skynet_multicast_group * grou
if (handle == last) if (handle == last)
continue; continue;
last = handle; last = handle;
struct skynet_context * ctx = skynet_handle_grab(handle);
if (ctx == NULL)
continue;
if (old_index < 0) { if (old_index < 0) {
group->data[new_index].handle = handle; group->data[new_index] = handle;
group->data[new_index].ctx = ctx;
} else { } else {
struct pair * p = &group->data[old_index]; uint32_t p = group->data[old_index];
if (handle == p->handle) if (handle == p)
continue; continue;
if (handle > p->handle) { if (handle > p) {
group->data[new_index].handle = handle; group->data[new_index] = handle;
group->data[new_index].ctx = ctx;
} else { } else {
group->data[new_index] = group->data[old_index]; group->data[new_index] = group->data[old_index];
--old_index; --old_index;
@@ -174,14 +162,11 @@ combine_queue(struct skynet_context * from, struct skynet_multicast_group * grou
break; break;
} }
uint32_t handle = group->leave_queue.data[i]; uint32_t handle = group->leave_queue.data[i];
struct pair * p = &group->data[old_index]; uint32_t p = group->data[old_index];
if (handle == p->handle) { if (handle == p) {
--count; --count;
++old_index; ++old_index;
if (p->ctx) { } else if ( handle > p) {
skynet_context_release(p->ctx);
}
} else if ( handle > p->handle) {
group->data[new_index] = group->data[old_index]; group->data[new_index] = group->data[old_index];
++new_index; ++new_index;
++old_index; ++old_index;
@@ -204,28 +189,45 @@ combine_queue(struct skynet_context * from, struct skynet_multicast_group * grou
int int
skynet_multicast_castgroup(struct skynet_context * from, struct skynet_multicast_group * group, struct skynet_multicast_message *msg) { skynet_multicast_castgroup(struct skynet_context * from, struct skynet_multicast_group * group, struct skynet_multicast_message *msg) {
combine_queue(from, group); combine_queue(from, group);
if (group->number == 0) {
skynet_multicast_dispatch(msg, NULL, NULL);
return 0;
}
uint32_t source = skynet_context_handle(from);
skynet_multicast_copy(msg, group->number);
int i;
int release = 0; int release = 0;
for (i=0;i<group->number;i++) { if (group->number > 0) {
struct pair * p = &group->data[i]; uint32_t source = skynet_context_handle(from);
skynet_context_send(p->ctx, msg, 0 , source, PTYPE_MULTICAST , 0); skynet_multicast_copy(msg, group->number);
int ref = skynet_context_ref(p->ctx); int i;
if (ref == 1) { for (i=0;i<group->number;i++) {
skynet_context_release(p->ctx); uint32_t p = group->data[i];
struct skynet_context * ctx = skynet_handle_grab(p->handle); struct skynet_context * ctx = skynet_handle_grab(p);
if (ctx == NULL) { if (ctx) {
p->ctx = NULL; skynet_context_send(ctx, msg, 0 , source, PTYPE_MULTICAST , 0);
skynet_multicast_leavegroup(group, p->handle); skynet_context_release(ctx);
} else {
skynet_multicast_leavegroup(group, p);
++release; ++release;
} }
} }
} }
skynet_multicast_copy(msg, -release);
return group->number - release; return group->number - release;
} }
void
skynet_multicast_cast(struct skynet_context * from, struct skynet_multicast_message *msg, uint32_t *dests, int n) {
uint32_t source = skynet_context_handle(from);
skynet_multicast_copy(msg, n);
int i;
int release = 0;
for (i=0;i<n;i++) {
uint32_t p = dests[i];
struct skynet_context * ctx = skynet_handle_grab(p);
if (ctx) {
skynet_context_send(ctx, msg, 0 , source, PTYPE_MULTICAST , 0);
skynet_context_release(ctx);
} else {
++release;
}
}
skynet_multicast_copy(msg, -release);
}

View File

@@ -13,6 +13,7 @@ typedef void (*skynet_multicast_func)(void *ud, uint32_t source, const void * ms
struct skynet_multicast_message * skynet_multicast_create(const void * msg, size_t sz, uint32_t source); struct skynet_multicast_message * skynet_multicast_create(const void * msg, size_t sz, uint32_t source);
void skynet_multicast_copy(struct skynet_multicast_message *, int copy); void skynet_multicast_copy(struct skynet_multicast_message *, int copy);
void skynet_multicast_dispatch(struct skynet_multicast_message * msg, void * ud, skynet_multicast_func func); void skynet_multicast_dispatch(struct skynet_multicast_message * msg, void * ud, skynet_multicast_func func);
void skynet_multicast_cast(struct skynet_context * from, struct skynet_multicast_message *msg, uint32_t *dests, int n);
struct skynet_multicast_group * skynet_multicast_newgroup(); struct skynet_multicast_group * skynet_multicast_newgroup();
void skynet_multicast_deletegroup(struct skynet_multicast_group * group); void skynet_multicast_deletegroup(struct skynet_multicast_group * group);

View File

@@ -132,12 +132,6 @@ skynet_context_release(struct skynet_context *ctx) {
return ctx; return ctx;
} }
int
skynet_context_ref(struct skynet_context *ctx) {
return ctx->ref;
}
int int
skynet_context_push(uint32_t handle, struct skynet_message *message) { skynet_context_push(uint32_t handle, struct skynet_message *message) {
struct skynet_context * ctx = skynet_handle_grab(handle); struct skynet_context * ctx = skynet_handle_grab(handle);

View File

@@ -11,7 +11,6 @@ struct skynet_monitor;
struct skynet_context * skynet_context_new(const char * name, const char * parm); struct skynet_context * skynet_context_new(const char * name, const char * parm);
void skynet_context_grab(struct skynet_context *); void skynet_context_grab(struct skynet_context *);
struct skynet_context * skynet_context_release(struct skynet_context *); struct skynet_context * skynet_context_release(struct skynet_context *);
int skynet_context_ref(struct skynet_context *);
uint32_t skynet_context_handle(struct skynet_context *); uint32_t skynet_context_handle(struct skynet_context *);
void skynet_context_init(struct skynet_context *, uint32_t handle); void skynet_context_init(struct skynet_context *, uint32_t handle);
int skynet_context_push(uint32_t handle, struct skynet_message *message); int skynet_context_push(uint32_t handle, struct skynet_message *message);