mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 12:20:41 +00:00
avoid issue #154
This commit is contained in:
@@ -36,17 +36,17 @@ struct remote_message_header {
|
|||||||
uint32_t session;
|
uint32_t session;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct msg {
|
struct harbor_msg {
|
||||||
struct remote_message_header header;
|
struct remote_message_header header;
|
||||||
void * buffer;
|
void * buffer;
|
||||||
size_t size;
|
size_t size;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct msg_queue {
|
struct harbor_msg_queue {
|
||||||
int size;
|
int size;
|
||||||
int head;
|
int head;
|
||||||
int tail;
|
int tail;
|
||||||
struct msg * data;
|
struct harbor_msg * data;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct keyvalue {
|
struct keyvalue {
|
||||||
@@ -54,7 +54,7 @@ struct keyvalue {
|
|||||||
char key[GLOBALNAME_LENGTH];
|
char key[GLOBALNAME_LENGTH];
|
||||||
uint32_t hash;
|
uint32_t hash;
|
||||||
uint32_t value;
|
uint32_t value;
|
||||||
struct msg_queue * queue;
|
struct harbor_msg_queue * queue;
|
||||||
};
|
};
|
||||||
|
|
||||||
struct hashmap {
|
struct hashmap {
|
||||||
@@ -69,7 +69,7 @@ struct hashmap {
|
|||||||
|
|
||||||
struct slave {
|
struct slave {
|
||||||
int fd;
|
int fd;
|
||||||
struct msg_queue *queue;
|
struct harbor_msg_queue *queue;
|
||||||
int status;
|
int status;
|
||||||
int length;
|
int length;
|
||||||
int read;
|
int read;
|
||||||
@@ -88,11 +88,11 @@ struct harbor {
|
|||||||
// hash table
|
// hash table
|
||||||
|
|
||||||
static void
|
static void
|
||||||
push_queue_msg(struct msg_queue * queue, struct msg * m) {
|
push_queue_msg(struct harbor_msg_queue * queue, struct harbor_msg * m) {
|
||||||
// If there is only 1 free slot which is reserved to distinguish full/empty
|
// If there is only 1 free slot which is reserved to distinguish full/empty
|
||||||
// of circular buffer, expand it.
|
// of circular buffer, expand it.
|
||||||
if (((queue->tail + 1) % queue->size) == queue->head) {
|
if (((queue->tail + 1) % queue->size) == queue->head) {
|
||||||
struct msg * new_buffer = skynet_malloc(queue->size * 2 * sizeof(struct msg));
|
struct harbor_msg * new_buffer = skynet_malloc(queue->size * 2 * sizeof(struct harbor_msg));
|
||||||
int i;
|
int i;
|
||||||
for (i=0;i<queue->size-1;i++) {
|
for (i=0;i<queue->size-1;i++) {
|
||||||
new_buffer[i] = queue->data[(i+queue->head) % queue->size];
|
new_buffer[i] = queue->data[(i+queue->head) % queue->size];
|
||||||
@@ -103,46 +103,46 @@ push_queue_msg(struct msg_queue * queue, struct msg * m) {
|
|||||||
queue->tail = queue->size - 1;
|
queue->tail = queue->size - 1;
|
||||||
queue->size *= 2;
|
queue->size *= 2;
|
||||||
}
|
}
|
||||||
struct msg * slot = &queue->data[queue->tail];
|
struct harbor_msg * slot = &queue->data[queue->tail];
|
||||||
*slot = *m;
|
*slot = *m;
|
||||||
queue->tail = (queue->tail + 1) % queue->size;
|
queue->tail = (queue->tail + 1) % queue->size;
|
||||||
}
|
}
|
||||||
|
|
||||||
static void
|
static void
|
||||||
push_queue(struct msg_queue * queue, void * buffer, size_t sz, struct remote_message_header * header) {
|
push_queue(struct harbor_msg_queue * queue, void * buffer, size_t sz, struct remote_message_header * header) {
|
||||||
struct msg m;
|
struct harbor_msg m;
|
||||||
m.header = *header;
|
m.header = *header;
|
||||||
m.buffer = buffer;
|
m.buffer = buffer;
|
||||||
m.size = sz;
|
m.size = sz;
|
||||||
push_queue_msg(queue, &m);
|
push_queue_msg(queue, &m);
|
||||||
}
|
}
|
||||||
|
|
||||||
static struct msg *
|
static struct harbor_msg *
|
||||||
pop_queue(struct msg_queue * queue) {
|
pop_queue(struct harbor_msg_queue * queue) {
|
||||||
if (queue->head == queue->tail) {
|
if (queue->head == queue->tail) {
|
||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
struct msg * slot = &queue->data[queue->head];
|
struct harbor_msg * slot = &queue->data[queue->head];
|
||||||
queue->head = (queue->head + 1) % queue->size;
|
queue->head = (queue->head + 1) % queue->size;
|
||||||
return slot;
|
return slot;
|
||||||
}
|
}
|
||||||
|
|
||||||
static struct msg_queue *
|
static struct harbor_msg_queue *
|
||||||
new_queue() {
|
new_queue() {
|
||||||
struct msg_queue * queue = skynet_malloc(sizeof(*queue));
|
struct harbor_msg_queue * queue = skynet_malloc(sizeof(*queue));
|
||||||
queue->size = DEFAULT_QUEUE_SIZE;
|
queue->size = DEFAULT_QUEUE_SIZE;
|
||||||
queue->head = 0;
|
queue->head = 0;
|
||||||
queue->tail = 0;
|
queue->tail = 0;
|
||||||
queue->data = skynet_malloc(DEFAULT_QUEUE_SIZE * sizeof(struct msg));
|
queue->data = skynet_malloc(DEFAULT_QUEUE_SIZE * sizeof(struct harbor_msg));
|
||||||
|
|
||||||
return queue;
|
return queue;
|
||||||
}
|
}
|
||||||
|
|
||||||
static void
|
static void
|
||||||
release_queue(struct msg_queue *queue) {
|
release_queue(struct harbor_msg_queue *queue) {
|
||||||
if (queue == NULL)
|
if (queue == NULL)
|
||||||
return;
|
return;
|
||||||
struct msg * m;
|
struct harbor_msg * m;
|
||||||
while ((m=pop_queue(queue)) != NULL) {
|
while ((m=pop_queue(queue)) != NULL) {
|
||||||
skynet_free(m->buffer);
|
skynet_free(m->buffer);
|
||||||
}
|
}
|
||||||
@@ -334,7 +334,7 @@ send_remote(struct skynet_context * ctx, int fd, const char * buffer, size_t sz,
|
|||||||
|
|
||||||
static void
|
static void
|
||||||
dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
||||||
struct msg_queue * queue = node->queue;
|
struct harbor_msg_queue * queue = node->queue;
|
||||||
uint32_t handle = node->value;
|
uint32_t handle = node->value;
|
||||||
int harbor_id = handle >> HANDLE_REMOTE_SHIFT;
|
int harbor_id = handle >> HANDLE_REMOTE_SHIFT;
|
||||||
assert(harbor_id != 0);
|
assert(harbor_id != 0);
|
||||||
@@ -352,7 +352,7 @@ dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
|||||||
s->queue = node->queue;
|
s->queue = node->queue;
|
||||||
node->queue = NULL;
|
node->queue = NULL;
|
||||||
} else {
|
} else {
|
||||||
struct msg * m;
|
struct harbor_msg * m;
|
||||||
while ((m = pop_queue(queue))!=NULL) {
|
while ((m = pop_queue(queue))!=NULL) {
|
||||||
push_queue_msg(s->queue, m);
|
push_queue_msg(s->queue, m);
|
||||||
}
|
}
|
||||||
@@ -360,7 +360,7 @@ dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
|||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
struct msg * m;
|
struct harbor_msg * m;
|
||||||
while ((m = pop_queue(queue)) != NULL) {
|
while ((m = pop_queue(queue)) != NULL) {
|
||||||
m->header.destination |= (handle & HANDLE_MASK);
|
m->header.destination |= (handle & HANDLE_MASK);
|
||||||
send_remote(context, fd, m->buffer, m->size, &m->header);
|
send_remote(context, fd, m->buffer, m->size, &m->header);
|
||||||
@@ -373,11 +373,11 @@ dispatch_queue(struct harbor *h, int id) {
|
|||||||
int fd = s->fd;
|
int fd = s->fd;
|
||||||
assert(fd != 0);
|
assert(fd != 0);
|
||||||
|
|
||||||
struct msg_queue *queue = s->queue;
|
struct harbor_msg_queue *queue = s->queue;
|
||||||
if (queue == NULL)
|
if (queue == NULL)
|
||||||
return;
|
return;
|
||||||
|
|
||||||
struct msg * m;
|
struct harbor_msg * m;
|
||||||
while ((m = pop_queue(queue)) != NULL) {
|
while ((m = pop_queue(queue)) != NULL) {
|
||||||
send_remote(h->ctx, fd, m->buffer, m->size, &m->header);
|
send_remote(h->ctx, fd, m->buffer, m->size, &m->header);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user