mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 03:53:09 +00:00
Use 2 level queues for local message
This commit is contained in:
@@ -38,9 +38,8 @@ skynet_error(struct skynet_context * context, const char *msg, ...) {
|
||||
} else {
|
||||
smsg.source = skynet_context_handle(context);
|
||||
}
|
||||
smsg.destination = logger;
|
||||
smsg.data = strdup(tmp);
|
||||
smsg.sz = len;
|
||||
skynet_mq_push(&smsg);
|
||||
skynet_context_push(logger, &smsg);
|
||||
}
|
||||
|
||||
|
||||
@@ -102,7 +102,7 @@ skynet_handle_grab(uint32_t handle) {
|
||||
|
||||
uint32_t hash = handle & (s->slot_size-1);
|
||||
struct skynet_context * ctx = s->slot[hash];
|
||||
if (skynet_context_handle(ctx) == handle) {
|
||||
if (ctx && skynet_context_handle(ctx) == handle) {
|
||||
result = ctx;
|
||||
skynet_context_grab(result);
|
||||
}
|
||||
|
||||
122
skynet_harbor.c
122
skynet_harbor.c
@@ -2,6 +2,7 @@
|
||||
#include "skynet_mq.h"
|
||||
#include "skynet_handle.h"
|
||||
#include "skynet_system.h"
|
||||
#include "skynet_server.h"
|
||||
#include "skynet.h"
|
||||
|
||||
#include <zmq.h>
|
||||
@@ -35,7 +36,7 @@ struct remote_header {
|
||||
|
||||
struct remote {
|
||||
void *socket;
|
||||
struct message_queue *queue;
|
||||
struct message_remote_queue *queue;
|
||||
};
|
||||
|
||||
struct harbor {
|
||||
@@ -46,7 +47,7 @@ struct harbor {
|
||||
int notice_event;
|
||||
struct hashmap *map;
|
||||
struct remote remote[REMOTE_MAX];
|
||||
struct message_queue *queue;
|
||||
struct message_remote_queue *queue;
|
||||
int harbor;
|
||||
|
||||
int lock;
|
||||
@@ -150,46 +151,49 @@ send_notice() {
|
||||
|
||||
// thread safe function
|
||||
void
|
||||
skynet_harbor_send(const char *name, struct skynet_message * message) {
|
||||
skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * msg) {
|
||||
if (name == NULL) {
|
||||
int remote_id = message->destination >> HANDLE_REMOTE_SHIFT;
|
||||
assert(destination!=0);
|
||||
int remote_id = destination >> HANDLE_REMOTE_SHIFT;
|
||||
assert(remote_id > 0 && remote_id <= REMOTE_MAX);
|
||||
--remote_id;
|
||||
struct remote * r = &Z->remote[remote_id];
|
||||
struct skynet_remote_message message;
|
||||
message.destination = destination;
|
||||
message.message = *msg;
|
||||
if (r->socket) {
|
||||
skynet_mq_enter(Z->queue, message);
|
||||
skynet_remotemq_push(Z->queue, &message);
|
||||
send_notice();
|
||||
} else {
|
||||
skynet_mq_enter(r->queue, message);
|
||||
skynet_remotemq_push(r->queue, &message);
|
||||
}
|
||||
} else {
|
||||
_lock();
|
||||
struct keyvalue * node = _hash_search(Z->map, name);
|
||||
printf("send to %s %p\n",name,node);
|
||||
if (node) {
|
||||
uint32_t dest = node->value;
|
||||
_unlock();
|
||||
if (dest == 0) {
|
||||
printf("queue a unknown address %s\n",name);
|
||||
// push message to unknown name service queue
|
||||
skynet_mq_enter(node->queue, message);
|
||||
skynet_mq_push(node->queue, msg);
|
||||
} else {
|
||||
message->destination = dest;
|
||||
if (!skynet_harbor_message_isremote(dest)) {
|
||||
// local message
|
||||
printf("%s is a local adress\n",name);
|
||||
skynet_mq_push(message);
|
||||
if (skynet_context_push(dest, msg)) {
|
||||
skynet_error(NULL, "Drop local message from %u to %s",msg->source, name);
|
||||
}
|
||||
return;
|
||||
}
|
||||
printf("queue remote message %s to global queue\n",name);
|
||||
skynet_mq_enter(Z->queue,message);
|
||||
struct skynet_remote_message message;
|
||||
message.destination = dest;
|
||||
message.message = *msg;
|
||||
skynet_remotemq_push(Z->queue,&message);
|
||||
send_notice();
|
||||
}
|
||||
} else {
|
||||
printf("Create new queue %s\n",name);
|
||||
// never seen name before
|
||||
struct message_queue * queue = skynet_mq_create(DEFAULT_QUEUE_SIZE);
|
||||
skynet_mq_enter(queue, message);
|
||||
struct message_queue * queue = skynet_mq_create(0);
|
||||
skynet_mq_push(queue, msg);
|
||||
_hash_insert(Z->map, name, 0, queue);
|
||||
_unlock();
|
||||
// 0 for query
|
||||
@@ -202,12 +206,13 @@ skynet_harbor_send(const char *name, struct skynet_message * message) {
|
||||
//queue a register message (destination = 0)
|
||||
void
|
||||
skynet_harbor_register(const char *name, uint32_t handle) {
|
||||
struct skynet_message msg;
|
||||
msg.source = handle;
|
||||
struct skynet_remote_message msg;
|
||||
msg.destination = SKYNET_SYSTEM_NAME;
|
||||
msg.data = strdup(name);
|
||||
msg.sz = 0;
|
||||
skynet_mq_enter(Z->queue,&msg);
|
||||
msg.message.source = handle;
|
||||
msg.message.data = strdup(name);
|
||||
|
||||
msg.message.sz = 0;
|
||||
skynet_remotemq_push(Z->queue,&msg);
|
||||
send_notice();
|
||||
}
|
||||
|
||||
@@ -227,24 +232,30 @@ _register_name(const char *name, uint32_t addr) {
|
||||
_hash_insert(Z->map, name, addr, NULL);
|
||||
}
|
||||
|
||||
if (addr == 0) {
|
||||
_unlock();
|
||||
return;
|
||||
}
|
||||
|
||||
struct skynet_message msg;
|
||||
struct message_queue * queue = node ? node->queue : NULL;
|
||||
|
||||
if (queue) {
|
||||
if (skynet_harbor_message_isremote(addr)) {
|
||||
while (skynet_mq_leave(queue, &msg)) {
|
||||
msg.destination = addr;
|
||||
skynet_mq_enter(Z->queue, &msg);
|
||||
while (!skynet_mq_pop(queue, &msg)) {
|
||||
struct skynet_remote_message message;
|
||||
message.destination = addr;
|
||||
message.message = msg;
|
||||
skynet_remotemq_push(Z->queue, &message);
|
||||
}
|
||||
} else {
|
||||
while (skynet_mq_leave(queue, &msg)) {
|
||||
msg.destination = addr;
|
||||
skynet_mq_push(&msg);
|
||||
while (!skynet_mq_pop(queue, &msg)) {
|
||||
if (skynet_context_push(addr,&msg)) {
|
||||
skynet_error(NULL,"Drop local message from %u to %s",msg.source,name);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
skynet_mq_release(queue);
|
||||
|
||||
node->queue = NULL;
|
||||
}
|
||||
|
||||
@@ -272,14 +283,14 @@ _remote_harbor_update(int harbor_id, const char * addr) {
|
||||
if (old_socket) {
|
||||
zmq_close(old_socket);
|
||||
}
|
||||
struct message_queue * queue = r->queue;
|
||||
struct message_remote_queue * queue = r->queue;
|
||||
|
||||
if (queue) {
|
||||
struct skynet_message msg;
|
||||
while (skynet_mq_leave(queue, &msg)) {
|
||||
skynet_mq_enter(Z->queue, &msg);
|
||||
struct skynet_remote_message msg;
|
||||
while (!skynet_remotemq_pop(queue, &msg)) {
|
||||
skynet_remotemq_push(Z->queue, &msg);
|
||||
}
|
||||
skynet_mq_release(queue);
|
||||
skynet_remotemq_release(queue);
|
||||
r->queue = NULL;
|
||||
}
|
||||
|
||||
@@ -403,13 +414,14 @@ _remote_query_name(const char *name) {
|
||||
|
||||
// Always in main harbor thread
|
||||
static void
|
||||
free_message (void *data, void *hint) {
|
||||
free_message(void *data, void *hint) {
|
||||
free(data);
|
||||
}
|
||||
|
||||
static void
|
||||
remote_socket_send(void * socket, struct skynet_message *msg) {
|
||||
remote_socket_send(void * socket, struct skynet_remote_message *msg) {
|
||||
struct remote_header rh;
|
||||
rh.source = msg->source;
|
||||
rh.source = msg->message.source;
|
||||
rh.destination = msg->destination;
|
||||
zmq_msg_t part;
|
||||
zmq_msg_init_size(&part,8);
|
||||
@@ -418,7 +430,7 @@ remote_socket_send(void * socket, struct skynet_message *msg) {
|
||||
zmq_send(socket, &part, ZMQ_SNDMORE);
|
||||
zmq_msg_close(&part);
|
||||
|
||||
zmq_msg_init_data(&part,msg->data,msg->sz,free_message,NULL);
|
||||
zmq_msg_init_data(&part,msg->message.data,msg->message.sz,free_message,NULL);
|
||||
zmq_send(socket, &part, 0);
|
||||
zmq_msg_close(&part);
|
||||
}
|
||||
@@ -455,43 +467,48 @@ _remote_recv() {
|
||||
zmq_msg_init(data);
|
||||
rc = zmq_recv(Z->zmq_local,data,0);
|
||||
_report_zmq_error(rc);
|
||||
|
||||
struct skynet_message msg;
|
||||
msg.source = rh.source;
|
||||
msg.destination = rh.destination;
|
||||
msg.data = data;
|
||||
msg.sz = zmq_msg_size(data);
|
||||
|
||||
// push remote message to local message queue
|
||||
skynet_mq_push(&msg);
|
||||
if (skynet_context_push(rh.destination, &msg)) {
|
||||
zmq_msg_close(data);
|
||||
free(data);
|
||||
skynet_error(NULL, "Drop remote message from %u to %u",rh.source, rh.destination);
|
||||
}
|
||||
}
|
||||
|
||||
// Always in main harbor thread
|
||||
static void
|
||||
_remote_send() {
|
||||
struct skynet_message msg;
|
||||
while (skynet_mq_leave(Z->queue,&msg)) {
|
||||
struct skynet_remote_message msg;
|
||||
while (!skynet_remotemq_pop(Z->queue,&msg)) {
|
||||
_goback:
|
||||
if (msg.destination == SKYNET_SYSTEM_NAME) {
|
||||
// register name
|
||||
const char * name = msg.data;
|
||||
if (msg.source) {
|
||||
_remote_register_name(name, msg.source);
|
||||
char * name = msg.message.data;
|
||||
|
||||
if (msg.message.source) {
|
||||
_remote_register_name(name, msg.message.source);
|
||||
} else {
|
||||
_remote_query_name(name);
|
||||
}
|
||||
|
||||
free(msg.data);
|
||||
free(name);
|
||||
} else {
|
||||
int harbor_id = (msg.destination >> HANDLE_REMOTE_SHIFT);
|
||||
assert(harbor_id > 0);
|
||||
struct remote * r = &Z->remote[harbor_id-1];
|
||||
if (r->socket == NULL) {
|
||||
if (r->queue == NULL) {
|
||||
r->queue = skynet_mq_create(DEFAULT_QUEUE_SIZE);
|
||||
skynet_mq_enter(r->queue, &msg);
|
||||
r->queue = skynet_remotemq_create();
|
||||
skynet_remotemq_push(r->queue, &msg);
|
||||
remote_query_harbor(harbor_id);
|
||||
} else {
|
||||
skynet_mq_enter(r->queue, &msg);
|
||||
skynet_remotemq_push(r->queue, &msg);
|
||||
}
|
||||
} else {
|
||||
remote_socket_send(r->socket, &msg);
|
||||
@@ -500,7 +517,8 @@ _goback:
|
||||
}
|
||||
__sync_lock_release(&Z->notice_event);
|
||||
// double check
|
||||
if (skynet_mq_leave(Z->queue,&msg)) {
|
||||
if (!skynet_remotemq_pop(Z->queue,&msg)) {
|
||||
printf("goback %x\n",msg.destination);
|
||||
goto _goback;
|
||||
}
|
||||
}
|
||||
@@ -600,7 +618,7 @@ skynet_harbor_init(const char * master, const char *local, int harbor) {
|
||||
h->zmq_local = harbor_socket;
|
||||
h->map = _hash_new();
|
||||
h->harbor = harbor;
|
||||
h->queue = skynet_mq_create(DEFAULT_QUEUE_SIZE);
|
||||
h->queue = skynet_remotemq_create();
|
||||
h->zmq_queue_notice = zmq_socket(context, ZMQ_PULL);
|
||||
r = zmq_bind(h->zmq_queue_notice, "inproc://notice");
|
||||
assert(r==0);
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
|
||||
struct skynet_message;
|
||||
|
||||
void skynet_harbor_send(const char *name, struct skynet_message * message);
|
||||
void skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * message);
|
||||
void skynet_harbor_register(const char *name, uint32_t handle);
|
||||
|
||||
// remote message is diffrent from local message.
|
||||
|
||||
193
skynet_mq.c
193
skynet_mq.c
@@ -2,9 +2,13 @@
|
||||
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <assert.h>
|
||||
|
||||
#define DEFAULT_QUEUE_SIZE 64;
|
||||
|
||||
struct message_queue {
|
||||
uint32_t handle;
|
||||
int cap;
|
||||
int head;
|
||||
int tail;
|
||||
@@ -12,16 +16,86 @@ struct message_queue {
|
||||
struct skynet_message *queue;
|
||||
};
|
||||
|
||||
static struct message_queue *Q = NULL;
|
||||
struct global_queue {
|
||||
int cap;
|
||||
int head;
|
||||
int tail;
|
||||
int lock;
|
||||
struct message_queue ** queue;
|
||||
};
|
||||
|
||||
struct message_remote_queue {
|
||||
int cap;
|
||||
int head;
|
||||
int tail;
|
||||
int lock;
|
||||
struct skynet_remote_message *queue;
|
||||
};
|
||||
|
||||
static struct global_queue *Q = NULL;
|
||||
|
||||
static inline void
|
||||
_lock_global_queue() {
|
||||
while (__sync_lock_test_and_set(&Q->lock,1)) {}
|
||||
}
|
||||
|
||||
#define LOCK(q) while (__sync_lock_test_and_set(&(q)->lock,1)) {}
|
||||
#define UNLOCK(q) __sync_lock_release(&(q)->lock);
|
||||
|
||||
void
|
||||
skynet_globalmq_push(struct message_queue * queue) {
|
||||
struct global_queue *q= Q;
|
||||
LOCK(q)
|
||||
|
||||
q->queue[q->tail] = queue;
|
||||
if (++ q->tail >= q->cap) {
|
||||
q->tail = 0;
|
||||
}
|
||||
|
||||
if (q->head == q->tail) {
|
||||
struct message_queue **new_queue = malloc(sizeof(struct message_queue *) * q->cap * 2);
|
||||
int i;
|
||||
for (i=0;i<q->cap;i++) {
|
||||
new_queue[i] = q->queue[(q->head + i) % q->cap];
|
||||
}
|
||||
q->head = 0;
|
||||
q->tail = q->cap;
|
||||
q->cap *= 2;
|
||||
|
||||
free(q->queue);
|
||||
q->queue = new_queue;
|
||||
}
|
||||
|
||||
UNLOCK(q)
|
||||
}
|
||||
|
||||
struct message_queue *
|
||||
skynet_mq_create(int cap) {
|
||||
skynet_globalmq_pop() {
|
||||
struct global_queue *q = Q;
|
||||
struct message_queue * ret = NULL;
|
||||
LOCK(q)
|
||||
|
||||
if (q->head != q->tail) {
|
||||
ret = q->queue[q->head];
|
||||
if ( ++ q->head >= q->cap) {
|
||||
q->head = 0;
|
||||
}
|
||||
}
|
||||
|
||||
UNLOCK(q)
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
struct message_queue *
|
||||
skynet_mq_create(uint32_t handle) {
|
||||
struct message_queue *q = malloc(sizeof(*q));
|
||||
q->cap = cap;
|
||||
q->handle = handle;
|
||||
q->cap = DEFAULT_QUEUE_SIZE;
|
||||
q->head = 0;
|
||||
q->tail = 0;
|
||||
q->lock = 0;
|
||||
q->queue = malloc(sizeof(struct skynet_message) * cap);
|
||||
q->queue = malloc(sizeof(struct skynet_message) * q->cap);
|
||||
|
||||
return q;
|
||||
}
|
||||
@@ -32,37 +106,33 @@ skynet_mq_release(struct message_queue *q) {
|
||||
free(q);
|
||||
}
|
||||
|
||||
static inline void
|
||||
_lock_queue(struct message_queue *q) {
|
||||
while (__sync_lock_test_and_set(&q->lock,1)) {}
|
||||
uint32_t
|
||||
skynet_mq_handle(struct message_queue *q) {
|
||||
return q->handle;
|
||||
}
|
||||
|
||||
static inline void
|
||||
_unlock_queue(struct message_queue *q) {
|
||||
__sync_lock_release(&q->lock);
|
||||
}
|
||||
|
||||
uint32_t
|
||||
skynet_mq_leave(struct message_queue *q, struct skynet_message *message) {
|
||||
uint32_t ret = 0;
|
||||
_lock_queue(q);
|
||||
int
|
||||
skynet_mq_pop(struct message_queue *q, struct skynet_message *message) {
|
||||
int ret = -1;
|
||||
LOCK(q)
|
||||
|
||||
if (q->head != q->tail) {
|
||||
*message = q->queue[q->head];
|
||||
ret = message->destination;
|
||||
ret = 0;
|
||||
if ( ++ q->head >= q->cap) {
|
||||
q->head = 0;
|
||||
}
|
||||
}
|
||||
|
||||
_unlock_queue(q);
|
||||
UNLOCK(q)
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
void
|
||||
skynet_mq_enter(struct message_queue *q, struct skynet_message *message) {
|
||||
_lock_queue(q);
|
||||
skynet_mq_push(struct message_queue *q, struct skynet_message *message) {
|
||||
LOCK(q)
|
||||
|
||||
q->queue[q->tail] = *message;
|
||||
if (++ q->tail >= q->cap) {
|
||||
@@ -83,20 +153,85 @@ skynet_mq_enter(struct message_queue *q, struct skynet_message *message) {
|
||||
q->queue = new_queue;
|
||||
}
|
||||
|
||||
_unlock_queue(q);
|
||||
}
|
||||
|
||||
uint32_t
|
||||
skynet_mq_pop(struct skynet_message *message) {
|
||||
return skynet_mq_leave(Q,message);
|
||||
UNLOCK(q)
|
||||
}
|
||||
|
||||
void
|
||||
skynet_mq_push(struct skynet_message *message) {
|
||||
skynet_mq_enter(Q,message);
|
||||
skynet_mq_init(int n) {
|
||||
struct global_queue *q = malloc(sizeof(*q));
|
||||
memset(q,0,sizeof(*q));
|
||||
int cap = 2;
|
||||
while (cap < n) {
|
||||
cap *=2;
|
||||
}
|
||||
|
||||
q->cap = cap;
|
||||
q->queue = malloc(cap * sizeof(struct skynet_message*));
|
||||
Q=q;
|
||||
}
|
||||
|
||||
// remote message queue
|
||||
|
||||
struct message_remote_queue *
|
||||
skynet_remotemq_create(void) {
|
||||
struct message_remote_queue *q = malloc(sizeof(*q));
|
||||
q->cap = DEFAULT_QUEUE_SIZE;
|
||||
q->head = 0;
|
||||
q->tail = 0;
|
||||
q->lock = 0;
|
||||
q->queue = malloc(sizeof(struct skynet_remote_message) * q->cap);
|
||||
|
||||
return q;
|
||||
}
|
||||
|
||||
void
|
||||
skynet_mq_init(int cap) {
|
||||
Q = skynet_mq_create(cap);
|
||||
skynet_remotemq_release(struct message_remote_queue *q) {
|
||||
free(q->queue);
|
||||
free(q);
|
||||
}
|
||||
|
||||
int
|
||||
skynet_remotemq_pop(struct message_remote_queue *q, struct skynet_remote_message *message) {
|
||||
int ret = -1;
|
||||
LOCK(q)
|
||||
|
||||
if (q->head != q->tail) {
|
||||
*message = q->queue[q->head];
|
||||
ret = 0;
|
||||
if ( ++ q->head >= q->cap) {
|
||||
q->head = 0;
|
||||
}
|
||||
}
|
||||
|
||||
UNLOCK(q)
|
||||
|
||||
return ret;
|
||||
}
|
||||
|
||||
void
|
||||
skynet_remotemq_push(struct message_remote_queue *q, struct skynet_remote_message *message) {
|
||||
assert(message->destination != 0);
|
||||
LOCK(q)
|
||||
|
||||
q->queue[q->tail] = *message;
|
||||
|
||||
if (++ q->tail >= q->cap) {
|
||||
q->tail = 0;
|
||||
}
|
||||
|
||||
if (q->head == q->tail) {
|
||||
struct skynet_remote_message *new_queue = malloc(sizeof(struct skynet_remote_message) * q->cap * 2);
|
||||
int i;
|
||||
for (i=0;i<q->cap;i++) {
|
||||
new_queue[i] = q->queue[(q->head + i) % q->cap];
|
||||
}
|
||||
q->head = 0;
|
||||
q->tail = q->cap;
|
||||
q->cap *= 2;
|
||||
|
||||
free(q->queue);
|
||||
q->queue = new_queue;
|
||||
}
|
||||
|
||||
UNLOCK(q)
|
||||
}
|
||||
|
||||
29
skynet_mq.h
29
skynet_mq.h
@@ -6,20 +6,35 @@
|
||||
|
||||
struct skynet_message {
|
||||
uint32_t source;
|
||||
uint32_t destination;
|
||||
void * data;
|
||||
size_t sz;
|
||||
};
|
||||
|
||||
struct message_queue;
|
||||
|
||||
uint32_t skynet_mq_pop(struct skynet_message *message);
|
||||
void skynet_mq_push(struct skynet_message *message);
|
||||
void skynet_globalmq_push(struct message_queue *);
|
||||
struct message_queue * skynet_globalmq_pop(void);
|
||||
|
||||
struct message_queue * skynet_mq_create(int cap);
|
||||
void skynet_mq_release(struct message_queue *q);
|
||||
uint32_t skynet_mq_leave(struct message_queue *q, struct skynet_message *message);
|
||||
void skynet_mq_enter(struct message_queue *q, struct skynet_message *message);
|
||||
struct message_queue * skynet_mq_create(uint32_t handle);
|
||||
void skynet_mq_release(struct message_queue *);
|
||||
uint32_t skynet_mq_handle(struct message_queue *);
|
||||
|
||||
// 0 for success
|
||||
int skynet_mq_pop(struct message_queue *q, struct skynet_message *message);
|
||||
void skynet_mq_push(struct message_queue *q, struct skynet_message *message);
|
||||
|
||||
struct skynet_remote_message {
|
||||
uint32_t destination;
|
||||
struct skynet_message message;
|
||||
};
|
||||
|
||||
struct message_remote_queue;
|
||||
|
||||
struct message_remote_queue * skynet_remotemq_create(void);
|
||||
void skynet_remotemq_release(struct message_remote_queue *);
|
||||
|
||||
int skynet_remotemq_pop(struct message_remote_queue *q, struct skynet_remote_message *message);
|
||||
void skynet_remotemq_push(struct message_remote_queue *q, struct skynet_remote_message *message);
|
||||
|
||||
void skynet_mq_init(int cap);
|
||||
|
||||
|
||||
103
skynet_server.c
103
skynet_server.c
@@ -12,18 +12,18 @@
|
||||
#include <stdio.h>
|
||||
|
||||
#define BLACKHOLE "blackhole"
|
||||
#define DEFAULT_MESSAGE_QUEUE 16
|
||||
#define DEFAULT_MESSAGE_QUEUE 16
|
||||
|
||||
struct skynet_context {
|
||||
void * instance;
|
||||
struct skynet_module * mod;
|
||||
uint32_t handle;
|
||||
int calling;
|
||||
int ref;
|
||||
char handle_name[10];
|
||||
char result[32];
|
||||
void * cb_ud;
|
||||
skynet_cb cb;
|
||||
int in_global_queue;
|
||||
struct message_queue *queue;
|
||||
};
|
||||
|
||||
@@ -53,19 +53,17 @@ skynet_context_new(const char * name, const char *parm) {
|
||||
ctx->ref = 2;
|
||||
ctx->cb = NULL;
|
||||
ctx->cb_ud = NULL;
|
||||
ctx->in_global_queue = 0;
|
||||
char * uid = ctx->handle_name;
|
||||
uid[0] = ':';
|
||||
_id_to_hex(uid+1, ctx->handle);
|
||||
|
||||
ctx->handle = skynet_handle_register(ctx);
|
||||
ctx->queue = skynet_mq_create(DEFAULT_MESSAGE_QUEUE);
|
||||
ctx->calling = 1;
|
||||
ctx->queue = skynet_mq_create(ctx->handle);
|
||||
// init function maybe use ctx->handle, so it must init at last
|
||||
|
||||
int r = skynet_module_instance_init(mod, inst, ctx, parm);
|
||||
if (r == 0) {
|
||||
__sync_synchronize();
|
||||
ctx->calling = 0;
|
||||
return skynet_context_release(ctx);
|
||||
} else {
|
||||
skynet_context_release(ctx);
|
||||
@@ -82,7 +80,9 @@ 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_release(ctx->queue);
|
||||
if (!ctx->in_global_queue) {
|
||||
skynet_mq_release(ctx->queue);
|
||||
}
|
||||
free(ctx);
|
||||
}
|
||||
|
||||
@@ -115,40 +115,58 @@ _dispatch_message(struct skynet_context *ctx, struct skynet_message *msg) {
|
||||
}
|
||||
}
|
||||
|
||||
static void
|
||||
_drop_queue(struct message_queue *q) {
|
||||
// todo: send message back to message source
|
||||
struct skynet_message msg;
|
||||
while(!skynet_mq_pop(q, &msg)) {
|
||||
if (skynet_harbor_message_isremote(msg.source)) {
|
||||
skynet_harbor_message_close(&msg);
|
||||
}
|
||||
free(msg.data);
|
||||
}
|
||||
skynet_mq_release(q);
|
||||
}
|
||||
|
||||
int
|
||||
skynet_context_message_dispatch(void) {
|
||||
struct skynet_message msg;
|
||||
uint32_t handle = skynet_mq_pop(&msg);
|
||||
if (handle == 0) {
|
||||
struct message_queue * q = skynet_globalmq_pop();
|
||||
if (q==NULL)
|
||||
return 1;
|
||||
}
|
||||
|
||||
uint32_t handle = skynet_mq_handle(q);
|
||||
|
||||
struct skynet_context * ctx = skynet_handle_grab(handle);
|
||||
if (ctx == NULL) {
|
||||
free(msg.data);
|
||||
skynet_error(NULL, "Drop message from %u to %u , size = %d",msg.source, msg.destination, (int)msg.sz);
|
||||
skynet_error(NULL, "Drop message queue %u ", handle);
|
||||
_drop_queue(q);
|
||||
return 0;
|
||||
}
|
||||
if (__sync_lock_test_and_set(&ctx->calling, 1)) {
|
||||
// When calling, push to context's message queue
|
||||
skynet_mq_enter(ctx->queue, &msg);
|
||||
} else {
|
||||
if (ctx->cb == NULL) {
|
||||
if (skynet_harbor_message_isremote(msg.source)) {
|
||||
skynet_harbor_message_close(&msg);
|
||||
}
|
||||
free(msg.data);
|
||||
skynet_error(NULL, "Drop message from %u to %u without callback , size = %d",msg.source, msg.destination, (int)msg.sz);
|
||||
} else {
|
||||
_dispatch_message(ctx, &msg);
|
||||
while(skynet_mq_leave(ctx->queue,&msg)) {
|
||||
_dispatch_message(ctx,&msg);
|
||||
}
|
||||
|
||||
assert(ctx->in_global_queue);
|
||||
|
||||
struct skynet_message msg;
|
||||
if (skynet_mq_pop(q,&msg)) {
|
||||
// empty queue
|
||||
__sync_lock_release(&ctx->in_global_queue);
|
||||
skynet_context_release(ctx);
|
||||
return 0;
|
||||
}
|
||||
|
||||
if (ctx->cb == NULL) {
|
||||
if (skynet_harbor_message_isremote(msg.source)) {
|
||||
skynet_harbor_message_close(&msg);
|
||||
}
|
||||
__sync_lock_release(&ctx->calling);
|
||||
free(msg.data);
|
||||
skynet_error(NULL, "Drop message from %u to %u without callback , size = %d",msg.source, handle, (int)msg.sz);
|
||||
} else {
|
||||
_dispatch_message(ctx, &msg);
|
||||
}
|
||||
|
||||
skynet_context_release(ctx);
|
||||
|
||||
skynet_globalmq_push(q);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
@@ -222,10 +240,9 @@ skynet_send(struct skynet_context * context, const char * addr , void * msg, siz
|
||||
} else {
|
||||
struct skynet_message smsg;
|
||||
smsg.source = context->handle;
|
||||
smsg.destination = 0;
|
||||
smsg.data = msg;
|
||||
smsg.sz = sz;
|
||||
skynet_harbor_send(addr, &smsg);
|
||||
skynet_harbor_send(addr, 0, &smsg);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -233,13 +250,15 @@ skynet_send(struct skynet_context * context, const char * addr , void * msg, siz
|
||||
|
||||
struct skynet_message smsg;
|
||||
smsg.source = context->handle;
|
||||
smsg.destination = des;
|
||||
smsg.data = msg;
|
||||
smsg.sz = sz;
|
||||
|
||||
if (skynet_harbor_message_isremote(des)) {
|
||||
skynet_harbor_send(NULL, &smsg);
|
||||
} else {
|
||||
skynet_mq_push(&smsg);
|
||||
skynet_harbor_send(NULL, des, &smsg);
|
||||
} else if (skynet_context_push(des, &smsg)) {
|
||||
free(msg);
|
||||
skynet_error(NULL, "Drop message from %u to %s (size=%d)", smsg.source, addr, (int)sz);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -260,3 +279,17 @@ skynet_callback(struct skynet_context * context, void *ud, skynet_cb cb) {
|
||||
context->cb_ud = ud;
|
||||
}
|
||||
|
||||
int
|
||||
skynet_context_push(uint32_t handle, struct skynet_message *message) {
|
||||
struct skynet_context * ctx = skynet_handle_grab(handle);
|
||||
if (ctx == NULL) {
|
||||
return -1;
|
||||
}
|
||||
skynet_mq_push(ctx->queue, message);
|
||||
if (__sync_lock_test_and_set(&ctx->in_global_queue,1) == 0) {
|
||||
skynet_globalmq_push(ctx->queue);
|
||||
}
|
||||
skynet_context_release(ctx);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
@@ -11,8 +11,7 @@ void skynet_context_grab(struct skynet_context *);
|
||||
struct skynet_context * skynet_context_release(struct skynet_context *);
|
||||
uint32_t skynet_context_handle(struct skynet_context *);
|
||||
void skynet_context_init(struct skynet_context *, uint32_t handle);
|
||||
void skynet_context_push(struct skynet_context *, struct skynet_message *message);
|
||||
int skynet_context_pop(struct skynet_context *, struct skynet_message *message);
|
||||
int skynet_context_push(uint32_t handle, struct skynet_message *message);
|
||||
int skynet_context_message_dispatch(void); // return 1 when block
|
||||
|
||||
#endif
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
#include "skynet_timer.h"
|
||||
#include "skynet_timer.h"
|
||||
#include "skynet_mq.h"
|
||||
#include "skynet_server.h"
|
||||
|
||||
#include <time.h>
|
||||
#include <assert.h>
|
||||
@@ -15,6 +16,11 @@ typedef void (*timer_execute_func)(void *ud,void *arg);
|
||||
#define TIME_NEAR_MASK (TIME_NEAR-1)
|
||||
#define TIME_LEVEL_MASK (TIME_LEVEL-1)
|
||||
|
||||
struct timer_event {
|
||||
uint32_t handle;
|
||||
int session;
|
||||
};
|
||||
|
||||
struct timer_node {
|
||||
struct timer_node *next;
|
||||
int expire;
|
||||
@@ -101,8 +107,15 @@ timer_execute(struct timer *T)
|
||||
current=link_clear(&T->near[idx]);
|
||||
|
||||
do {
|
||||
struct timer_node *temp=current;
|
||||
skynet_mq_push((struct skynet_message *)(temp+1));
|
||||
struct timer_event * event = (struct timer_event *)(current+1);
|
||||
struct skynet_message message;
|
||||
message.source = SKYNET_SYSTEM_TIMER;
|
||||
message.data = NULL;
|
||||
message.sz = (size_t) event->session;
|
||||
|
||||
skynet_context_push(event->handle, &message);
|
||||
|
||||
struct timer_node * temp = current;
|
||||
current=current->next;
|
||||
free(temp);
|
||||
} while (current);
|
||||
@@ -158,16 +171,21 @@ timer_create_timer()
|
||||
}
|
||||
|
||||
void
|
||||
skynet_timeout(int handle, int time, int session) {
|
||||
struct skynet_message message;
|
||||
message.source = SKYNET_SYSTEM_TIMER;
|
||||
message.destination = handle;
|
||||
message.data = NULL;
|
||||
message.sz = (size_t) session;
|
||||
skynet_timeout(uint32_t handle, int time, int session) {
|
||||
if (time == 0) {
|
||||
skynet_mq_push(&message);
|
||||
struct skynet_message message;
|
||||
message.source = SKYNET_SYSTEM_TIMER;
|
||||
message.data = NULL;
|
||||
message.sz = (size_t) session;
|
||||
|
||||
if (skynet_context_push(handle, &message)) {
|
||||
return;
|
||||
}
|
||||
} else {
|
||||
timer_add(TI, &message, sizeof(message), time);
|
||||
struct timer_event event;
|
||||
event.handle = handle;
|
||||
event.session = session;
|
||||
timer_add(TI, &event, sizeof(event), time);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2,9 +2,10 @@
|
||||
#define SKYNET_TIMER_H
|
||||
|
||||
#include "skynet_system.h"
|
||||
|
||||
#include <stdint.h>
|
||||
|
||||
void skynet_timeout(int handle, int time, int session);
|
||||
void skynet_timeout(uint32_t handle, int time, int session);
|
||||
void skynet_updatetime(void);
|
||||
uint32_t skynet_gettime(void);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user