From 4f6826de4059adb7f15994cb3818f8a874bf3d4a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BA=91=E9=A3=8E?= Date: Fri, 17 Aug 2012 11:50:45 +0800 Subject: [PATCH] new harbor --- Makefile | 76 ++-- README.md | 2 +- config | 6 +- config_log | 4 +- connection/main.c | 10 +- gate/main.c | 45 ++- gate/mread.c | 6 +- gate/mread.h | 4 +- lualib-src/lua-skynet.c | 33 +- service-src/service_broker.c | 4 +- service-src/service_client.c | 4 +- service/agent.lua | 4 +- service/main_log.lua | 2 +- skynet-src/skynet.h | 6 +- skynet-src/skynet_handle.h | 5 +- skynet-src/skynet_harbor.c | 691 ++--------------------------------- skynet-src/skynet_harbor.h | 32 +- skynet-src/skynet_imp.h | 2 +- skynet-src/skynet_logger.c | 4 +- skynet-src/skynet_main.c | 9 +- skynet-src/skynet_mq.c | 121 ++---- skynet-src/skynet_mq.h | 15 +- skynet-src/skynet_server.c | 196 +++++----- skynet-src/skynet_server.h | 2 + skynet-src/skynet_start.c | 58 ++- 25 files changed, 331 insertions(+), 1010 deletions(-) diff --git a/Makefile b/Makefile index cb63c9a8..2f93ad57 100644 --- a/Makefile +++ b/Makefile @@ -1,4 +1,7 @@ -.PHONY : all clean install +.PHONY : all clean + +CFLAGS = -g -Wall +SHARED = -fPIC --shared all : \ skynet \ @@ -14,7 +17,8 @@ lualib/socket.so \ lualib/lpeg.so \ lualib/protobuf.so \ lualib/int64.so \ -skynet-master +service/master.so \ +service/harbor.so skynet : \ skynet-src/skynet_main.c \ @@ -26,36 +30,41 @@ skynet-src/skynet_start.c \ skynet-src/skynet_timer.c \ skynet-src/skynet_error.c \ skynet-src/skynet_harbor.c \ -skynet-src/skynet_env.c \ -master/master.c - gcc -Wall -g -Wl,-E -o $@ $^ -Iskynet-src -Imaster -lpthread -ldl -lrt -Wl,-E -llua -lm -lzmq +skynet-src/skynet_env.c + gcc $(CFLAGS) -Wl,-E -o $@ $^ -Iskynet-src -lpthread -ldl -lrt -Wl,-E -llua -lm + +service/master.so : service-src/service_master.c + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src + +service/harbor.so : service-src/service_harbor.c + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src service/logger.so : skynet-src/skynet_logger.c - gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src service/snlua.so : service-src/service_lua.c - gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src service/gate.so : gate/mread.c gate/ringbuffer.c gate/main.c - gcc -Wall -g -fPIC --shared -o $@ $^ -Igate -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Igate -Iskynet-src lualib/skynet.so : lualib-src/lua-skynet.c lua-serialize/serialize.c - gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src -Ilua-serialize + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Ilua-serialize service/client.so : service-src/service_client.c - gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src service/connection.so : connection/connection.c connection/main.c - gcc -Wall -g -fPIC --shared -o $@ $^ -Iconnection -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Iconnection lualib/socket.so : connection/lua-socket.c - gcc -Wall -g -fPIC --shared -o $@ $^ -Iconnection -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Iconnection service/broker.so : service-src/service_broker.c - gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src + gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src lualib/lpeg.so : lpeg/lpeg.c - gcc -Wall -O2 -fPIC --shared $^ -o $@ -Ilpeg + gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@ -Ilpeg PROTOBUFSRC = \ lua-protobuf/context.c \ @@ -73,45 +82,14 @@ PROTOBUFSRC = \ lua-protobuf/decode.c lualib/protobuf.so : $(PROTOBUFSRC) lua-protobuf/pbc-lua.c - gcc -O2 -Wall --shared -fPIC -o $@ $^ + gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@ lualib/int64.so : lua-int64/int64.c - gcc -O2 -Wall --shared -fPIC -o $@ $^ + gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@ client : client-src/client.c - gcc -Wall -g $^ -o $@ -lpthread - -skynet-master : master/master.c master/main.c - gcc -g -Wall -Imaster -o $@ $^ -lzmq + gcc $(CFLAGS) $^ -o $@ -lpthread clean : - rm skynet client skynet-master lualib/*.so service/*.so - - -$(SKYNET_PATH) : - mkdir $@ - -install-libs : | $(SKYNET_PATH) - cp -R lualib $(SKYNET_PATH) - -install-bins : | $(SKYNET_PATH) - cp skynet $(SKYNET_PATH) - -$(SKYNET_PATH)/service : | $(SKYNET_PATH) - mkdir $@ - -install-services : | $(SKYNET_PATH)/service - cp service/connection.so $(SKYNET_PATH)/service/ - cp service/logger.so $(SKYNET_PATH)/service/ - cp service/gate.so $(SKYNET_PATH)/service/ - cp service/broker.so $(SKYNET_PATH)/service/ - cp service/snlua.so $(SKYNET_PATH)/service/ - cp service/client.so $(SKYNET_PATH)/service/ - cp service/launcher.lua $(SKYNET_PATH)/service/ - cp service/redis-cli.lua $(SKYNET_PATH)/service/ - - -install : install-libs install-bins install-services - - + rm skynet client lualib/*.so service/*.so \ No newline at end of file diff --git a/README.md b/README.md index 8e447449..cce2abce 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ ## Build -Install zeromq 2.2 and lua 5.2 first. +Install lua 5.2 first. ``` make diff --git a/config b/config index 54bdc917..1a28330f 100644 --- a/config +++ b/config @@ -3,10 +3,10 @@ thread = 8 mqueue = 256 logger = nil harbor = 1 -address = "tcp://127.0.0.1:2525" -master = "tcp://127.0.0.1:2012" +address = "127.0.0.1:2525" +master = "127.0.0.1:2012" start = "main" -standalone = true +standalone = "0.0.0.0:2012" luaservice = root.."service/?.lua;"..root.."service/?/init.lua" cpath = root.."service/?.so" protopath = root.."proto" diff --git a/config_log b/config_log index 74358d05..f4e4e0f4 100644 --- a/config_log +++ b/config_log @@ -3,6 +3,6 @@ mqueue = 256 cpath = "./service/?.so" logger = nil harbor = 2 -address = "tcp://127.0.0.1:2526" -master = "tcp://127.0.0.1:2012" +address = "127.0.0.1:2526" +master = "127.0.0.1:2012" start = "main_log" diff --git a/connection/main.c b/connection/main.c index d7046ff6..a62f3880 100644 --- a/connection/main.c +++ b/connection/main.c @@ -124,11 +124,11 @@ _poll(struct connection_server * server) { } } -static void +static int _main(struct skynet_context * ctx, void * ud, int session, const char * uid, const void * msg, size_t sz) { if (msg == NULL) { _poll(ud); - return; + return 0; } const char * param = (const char *)msg + 4; if (memcmp(msg, "ADD ", 4)==0) { @@ -136,7 +136,7 @@ _main(struct skynet_context * ctx, void * ud, int session, const char * uid, con int fd = strtol(param, &endptr, 10); if (endptr == NULL) { skynet_error(ctx, "[connection] Invalid ADD command from %s (session = %d)", uid, session); - return; + return 0; } int addr_sz = sz - (endptr - (char *)msg); char * addr = malloc(addr_sz); @@ -148,12 +148,14 @@ _main(struct skynet_context * ctx, void * ud, int session, const char * uid, con int fd = strtol(param, &endptr, 10); if (endptr == NULL) { skynet_error(ctx, "[connection] Invalid DEL command from %s (session = %d)", uid, session); - return; + return 0; } _del(ud, fd); } else { skynet_error(ctx, "[connection] Invalid command from %s (session = %d)", uid, session); } + + return 0; } int diff --git a/gate/main.c b/gate/main.c index 63009f5b..39f9b92f 100644 --- a/gate/main.c +++ b/gate/main.c @@ -122,6 +122,9 @@ _ctrl(struct skynet_context * ctx, struct gate * g, const void * msg, int sz) { static void _report(struct gate *g, struct skynet_context * ctx, const char * data, ...) { + if (g->watchdog == NULL) { + return; + } va_list ap; va_start(ap, data); char tmp[1024]; @@ -134,14 +137,14 @@ _report(struct gate *g, struct skynet_context * ctx, const char * data, ...) { static void _forward(struct skynet_context * ctx,struct gate *g, int uid, void * data, size_t len) { if (g->broker) { - skynet_send(ctx, NULL, g->broker, 0x7fffffff, data, len, 0); + skynet_send(ctx, NULL, g->broker, SESSION_CLIENT, data, len, 0); return; } struct connection * agent = _id_to_agent(g,uid); if (agent->agent) { // todo: client package has not session , send 0x7fffffff - skynet_send(ctx, agent->client, agent->agent, 0x7fffffff, data, len, 0); - } else { + skynet_send(ctx, agent->client, agent->agent, SESSION_CLIENT, data, len, 0); + } else if (g->watchdog) { char * tmp = malloc(len + 32); int n = snprintf(tmp,len+32,"%d data ",uid); memcpy(tmp+n,data,len); @@ -178,12 +181,12 @@ _remove_id(struct gate *g, int uid) { } } -static void +static int _cb(struct skynet_context * ctx, void * ud, int session, const char * uid, const void * msg, size_t sz) { struct gate *g = ud; if (msg) { _ctrl(ctx, g , msg , (int)sz); - return; + return 0; } struct mread_pool * m = g->pool; int connection_id = mread_poll(m,100); // timeout : 100ms @@ -223,6 +226,7 @@ _cb(struct skynet_context * ctx, void * ud, int session, const char * uid, const _break: skynet_command(ctx, "TIMEOUT", "0"); } + return 0; } int @@ -230,18 +234,41 @@ gate_init(struct gate *g , struct skynet_context * ctx, char * parm) { int port = 0; int max = 0; int buffer = 0; - char watchdog[strlen(parm)+1]; - int n = sscanf(parm, "%s %d %d %d",watchdog, &port,&max,&buffer); + int sz = strlen(parm)+1; + char watchdog[sz]; + char binding[sz]; + int n = sscanf(parm, "%s %s %d %d",watchdog, binding,&max,&buffer); if (n!=4) { skynet_error(ctx, "Invalid gate parm %s",parm); return 1; } - struct mread_pool * pool = mread_create(port, max, buffer); + char * portstr = strchr(binding,':'); + uint32_t addr = INADDR_ANY; + if (portstr == NULL) { + port = strtol(binding, NULL, 10); + if (port <= 0) { + skynet_error(ctx, "Invalid gate address %s",parm); + return 1; + } + } else { + port = strtol(portstr + 1, NULL, 10); + if (port <= 0) { + skynet_error(ctx, "Invalid gate address %s",parm); + return 1; + } + portstr[0] = '\0'; + addr=inet_addr(binding); + } + struct mread_pool * pool = mread_create(addr, port, max, buffer); if (pool == NULL) { skynet_error(ctx, "Create gate %s failed",parm); return 1; } - g->watchdog = strdup(watchdog); + if (watchdog[0] == '!') { + g->watchdog = NULL; + } else { + g->watchdog = strdup(watchdog); + } g->pool = pool; int cap = 1; while (cap < max) { diff --git a/gate/mread.c b/gate/mread.c index b610f560..be70c23d 100644 --- a/gate/mread.c +++ b/gate/mread.c @@ -93,7 +93,7 @@ _set_nonblocking(int fd) } struct mread_pool * -mread_create(int port , int max , int buffer_size) { +mread_create(uint32_t addr, int port , int max , int buffer_size) { int listen_fd = socket(AF_INET, SOCK_STREAM, 0); if (listen_fd == -1) { return NULL; @@ -109,8 +109,8 @@ mread_create(int port , int max , int buffer_size) { memset(&my_addr, 0, sizeof(struct sockaddr_in)); my_addr.sin_family = AF_INET; my_addr.sin_port = htons(port); - my_addr.sin_addr.s_addr = htonl(INADDR_ANY); // INADDR_LOOPBACK -// printf("MREAD bind %s:%u\n",inet_ntoa(my_addr.sin_addr),ntohs(my_addr.sin_port)); + my_addr.sin_addr.s_addr = addr; + printf("MREAD bind %s:%u\n",inet_ntoa(my_addr.sin_addr),ntohs(my_addr.sin_port)); if (bind(listen_fd, (struct sockaddr *)&my_addr, sizeof(struct sockaddr)) == -1) { close(listen_fd); return NULL; diff --git a/gate/mread.h b/gate/mread.h index c1e271b0..e1f27837 100644 --- a/gate/mread.h +++ b/gate/mread.h @@ -1,9 +1,11 @@ #ifndef MREAD_H #define MREAD_H +#include + struct mread_pool; -struct mread_pool * mread_create(int port , int max , int buffer); +struct mread_pool * mread_create(uint32_t addr, int port , int max , int buffer); void mread_close(struct mread_pool *m); int mread_poll(struct mread_pool *m , int timeout); diff --git a/lualib-src/lua-skynet.c b/lualib-src/lua-skynet.c index 0454c5a6..80b1374c 100644 --- a/lualib-src/lua-skynet.c +++ b/lualib-src/lua-skynet.c @@ -5,24 +5,20 @@ #include #include #include +#include -static int -traceback (lua_State *L) { - const char *msg = lua_tostring(L, 1); - if (msg) { - luaL_traceback(L, L, msg, 1); - } else { - lua_pushliteral(L, "(no error message)"); - } - return 1; -} - -static void +static int _cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) { lua_State *L = ud; - lua_rawgetp(L, LUA_REGISTRYINDEX, traceback); - int trace = lua_gettop(L); - lua_rawgetp(L, LUA_REGISTRYINDEX, _cb); + int trace = 1; + int top = lua_gettop(L); + if (top == 1) { + lua_rawgetp(L, LUA_REGISTRYINDEX, _cb); + } else { + assert(top == 2); + lua_pushvalue(L,2); + } + int r; if (msg == NULL) { if (addr == NULL) { @@ -41,7 +37,7 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr, r = lua_pcall(L, 4, 0 , trace); } if (r == LUA_OK) - return; + return 0; const char * self = skynet_command(context, "REG", NULL); switch (r) { case LUA_ERRRUN: @@ -59,6 +55,8 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr, }; lua_pop(L,1); + + return 0; } static int @@ -201,9 +199,6 @@ luaopen_skynet_c(lua_State *L) { { "unpack", _luaseri_unpack }, { NULL, NULL }, }; - lua_pushcfunction(L, traceback); - lua_rawsetp(L, LUA_REGISTRYINDEX, traceback); - luaL_newlibtable(L,l); lua_getfield(L, LUA_REGISTRYINDEX, "skynet_context"); diff --git a/service-src/service_broker.c b/service-src/service_broker.c index bcac111b..56bbffe4 100644 --- a/service-src/service_broker.c +++ b/service-src/service_broker.c @@ -56,7 +56,7 @@ _forward(struct broker *b, struct skynet_context * context) { b->id = (b->id + 1) % DEFAULT_NUMBER; } -static void +static int _cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) { struct broker * b = ud; if (b->init < DEFAULT_NUMBER) { @@ -68,6 +68,8 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr, } else { _forward(b, context); } + + return 0; } diff --git a/service-src/service_client.c b/service-src/service_client.c index 6412183b..ebddac7e 100644 --- a/service-src/service_client.c +++ b/service-src/service_client.c @@ -8,7 +8,7 @@ #include #include -static void +static int _cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) { assert(sz <= 65535); int fd = (int)(intptr_t)ud; @@ -31,7 +31,7 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr, } } assert(err == sz +2); - return; + return 0; } } diff --git a/service/agent.lua b/service/agent.lua index aa20af13..ba5bdfab 100644 --- a/service/agent.lua +++ b/service/agent.lua @@ -4,12 +4,12 @@ local client = ... local session_id = 0 skynet.filter(function (session, address , msg, sz) if session == 0x7fffffff then - print("client message",skynet.tostring(msg,sz)) + skynet.send("LOG", "client message :" .. skynet.tostring(msg,sz)) -- It's client, there is no session session_id = session_id + 1 session = - session_id else - print("skynet message",msg,sz) + skynet.send("LOG", "skynet message") end return session, address , msg, sz end, function (msg,sz) diff --git a/service/main_log.lua b/service/main_log.lua index 9e613c51..976e0f49 100644 --- a/service/main_log.lua +++ b/service/main_log.lua @@ -2,7 +2,7 @@ local skynet = require "skynet" print("Log server start") -local log = skynet.launch("snlua","globallog.lua") +local log = skynet.launch("snlua","globallog") print("log",log) skynet.exit() diff --git a/skynet-src/skynet.h b/skynet-src/skynet.h index 4d46766a..70e671de 100644 --- a/skynet-src/skynet.h +++ b/skynet-src/skynet.h @@ -5,6 +5,7 @@ #include #define DONTCOPY 1 +#define SESSION_CLIENT 0x7fffffff struct skynet_context; @@ -12,8 +13,9 @@ void skynet_error(struct skynet_context * context, const char *msg, ...); const char * skynet_command(struct skynet_context * context, const char * cmd , const char * parm); int skynet_send(struct skynet_context * context, const char * source, const char * addr , int session, void * msg, size_t sz, int flags); -void skynet_forward(struct skynet_context *, const char * addr); -typedef void (*skynet_cb)(struct skynet_context * context, void *ud, int session, const char * addr , const void * msg, size_t sz); +void skynet_forward(struct skynet_context *, const char * destination); + +typedef int (*skynet_cb)(struct skynet_context * context, void *ud, int session, const char * addr , const void * msg, size_t sz); void skynet_callback(struct skynet_context * context, void *ud, skynet_cb cb); #endif diff --git a/skynet-src/skynet_handle.h b/skynet-src/skynet_handle.h index 28bd548c..ce5d4976 100644 --- a/skynet-src/skynet_handle.h +++ b/skynet-src/skynet_handle.h @@ -3,10 +3,7 @@ #include -// reserve high 8 bits for remote id -// see skynet_harbor.c REMOTE_MAX -#define HANDLE_MASK 0xffffff -#define HANDLE_REMOTE_SHIFT 24 +#include "skynet_harbor.h" struct skynet_context; diff --git a/skynet-src/skynet_harbor.c b/skynet-src/skynet_harbor.c index 16e6397f..6c4cf80e 100644 --- a/skynet-src/skynet_harbor.c +++ b/skynet-src/skynet_harbor.c @@ -1,674 +1,53 @@ #include "skynet_harbor.h" -#include "skynet_mq.h" -#include "skynet_handle.h" -#include "skynet_system.h" #include "skynet_server.h" -#include "skynet.h" -#include #include -#include #include #include -#define HASH_SIZE 4096 -#define DEFAULT_QUEUE_SIZE 1024 +static struct skynet_context * REMOTE = 0; +static int HARBOR = 0; -// see skynet_handle.h for HANDLE_REMOTE_SHIFT -#define REMOTE_MAX 255 - -struct keyvalue { - struct keyvalue * next; - uint32_t hash; - char * key; - uint32_t value; - struct message_queue * queue; -}; - -struct hashmap { - struct keyvalue *node[HASH_SIZE]; -}; - -struct remote_header { - uint32_t source; - uint32_t destination; - uint32_t session; -}; - -struct remote { - void *socket; - struct message_remote_queue *queue; -}; - -struct harbor { - void * zmq_context; - void * zmq_master_request; - void * zmq_local; - void * zmq_queue_notice; - int notice_event; - struct hashmap *map; - struct remote remote[REMOTE_MAX]; - struct message_remote_queue *queue; - int harbor; - - int lock; -}; - -static struct harbor *Z = NULL; - -// todo: optimize for little endian system - -static inline void -buffer_to_remote_header(uint8_t *buffer, struct remote_header *header) { - header->source = buffer[0] | buffer[1] << 8 | buffer[2] << 16 | buffer[3] << 24; - header->destination = buffer[4] | buffer[5] << 8 | buffer[6] << 16 | buffer[7] << 24; - header->session = buffer[8] | buffer[9] << 8 | buffer[10] << 16 | buffer[11] << 24; +void +skynet_harbor_send(struct remote_message *rmsg, uint32_t source, int session) { + skynet_context_send(REMOTE, rmsg, sizeof(*rmsg), source, session); } -static inline void -remote_header_to_buffer(struct remote_header *header, uint8_t *buffer) { - buffer[0] = header->source & 0xff; - buffer[1] = (header->source >> 8) & 0xff; - buffer[2] = (header->source >>16) & 0xff; - buffer[3] = (header->source >>24)& 0xff; - buffer[4] = (header->destination) & 0xff; - buffer[5] = (header->destination >>8) & 0xff; - buffer[6] = (header->destination >>16) & 0xff; - buffer[7] = (header->destination >>24) & 0xff; - buffer[8] = (header->session) & 0xff; - buffer[9] = (header->session >>8) & 0xff; - buffer[10] = (header->session >>16) & 0xff; - buffer[11] = (header->session >>24) & 0xff; -} - -static uint32_t -calc_hash(const char *name) { +void +skynet_harbor_register(struct remote_name *rname) { int i; - uint32_t h = 0; - for (i=0;name[i];i++) { - h = h ^ ((h<<5)+(h>>2)+(uint8_t)name[i]); - } - h ^= i; - return h; -} - -static inline void -_lock() { - while (__sync_lock_test_and_set(&Z->lock,1)) {} -} - -static inline void -_unlock() { - __sync_lock_release(&Z->lock); -} - -static struct hashmap * -_hash_new(void) { - struct hashmap * hash = malloc(sizeof(*hash)); - memset(hash, 0, sizeof(*hash)); - return hash; -} - -static struct keyvalue * -_hash_search(struct hashmap * hash, const char * key) { - uint32_t h = calc_hash(key); - struct keyvalue * n = hash->node[h & (HASH_SIZE-1)]; - while (n) { - if (n->hash == h && strcmp(n->key, key) == 0) { - return n; - } - n = n->next; - } - return NULL; -} - -static void -_hash_insert(struct hashmap * hash, const char * key, uint32_t handle, struct message_queue *queue) { - uint32_t h = calc_hash(key); - struct keyvalue * node = malloc(sizeof(*node)); - node->next = hash->node[h & (HASH_SIZE-1)]; - node->hash = h; - node->key = strdup(key); - node->value = handle; - node->queue = queue; - - hash->node[h & (HASH_SIZE-1)] = node; -} - -// thread safe function -static void -send_notice() { - if (__sync_lock_test_and_set(&Z->notice_event,1)) { - // already send notice - return; - } - static __thread void * queue_notice = NULL; - if (queue_notice == NULL) { - void * pub = zmq_socket(Z->zmq_context, ZMQ_PUSH); - int r = zmq_connect(pub , "inproc://notice"); - assert(r==0); - queue_notice = pub; - } - zmq_msg_t dummy; - zmq_msg_init(&dummy); - zmq_send(queue_notice,&dummy,0); - zmq_msg_close(&dummy); -} - -// thread safe function -void -skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * msg) { - if (name == NULL) { - assert(destination!=0); - int remote_id = destination >> HANDLE_REMOTE_SHIFT; - assert(remote_id > 0 && remote_id <= REMOTE_MAX); - struct skynet_remote_message message; - message.destination = destination; - message.message = *msg; - skynet_remotemq_push(Z->queue, &message); - send_notice(); - } else { - _lock(); - struct keyvalue * node = _hash_search(Z->map, name); - if (node) { - uint32_t dest = node->value; - _unlock(); - if (dest == 0) { - // push message to unknown name service queue - skynet_mq_push(node->queue, msg); - } else { - if (!skynet_harbor_message_isremote(dest)) { - // local message - if (skynet_context_push(dest, msg)) { - skynet_error(NULL, "Drop local message from %u to %s",msg->source, name); - } - return; - } - struct skynet_remote_message message; - message.destination = dest; - message.message = *msg; - skynet_remotemq_push(Z->queue,&message); - send_notice(); - } - } else { - // never seen name before - struct message_queue * queue = skynet_mq_create(0); - skynet_mq_push(queue, msg); - _hash_insert(Z->map, name, 0, queue); - _unlock(); - // 0 for query - skynet_harbor_register(name,0); + int number = 1; + for (i=0;iname[i]; + if (!(c >= '0' && c <='9')) { + number = 0; + break; } } -} - -// thread safe function -//queue a register message (destination = 0) -void -skynet_harbor_register(const char *name, uint32_t handle) { - struct skynet_remote_message msg; - msg.destination = SKYNET_SYSTEM_NAME; - msg.message.source = handle; - msg.message.data = strdup(name); - - msg.message.sz = 0; - skynet_remotemq_push(Z->queue,&msg); - send_notice(); -} - -// Always in main harbor thread -static void -_register_name(const char *name, uint32_t addr) { - _lock(); - struct keyvalue * node = _hash_search(Z->map, name); - if (node) { - if (node->value) { - node->value = addr; - assert(node->queue == NULL); - } else { - node->value = addr; - } - } else { - _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_pop(queue, &msg)) { - struct skynet_remote_message message; - message.destination = addr; - message.message = msg; - skynet_remotemq_push(Z->queue, &message); - send_notice(); - } - } else { - 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); - } - } - } - - node->queue = NULL; - } - - _unlock(); - - if (queue) { - skynet_mq_release(queue); - } -} - -// Always in main harbor thread - -static void -_remote_harbor_update(int harbor_id, const char * addr) { - struct remote * r = &Z->remote[harbor_id-1]; - void *socket = zmq_socket( Z->zmq_context, ZMQ_PUSH); - int rc = zmq_connect(socket, addr); - if (rc<0) { - skynet_error(NULL, "Can't connect to harbor %d %s",harbor_id,addr); - zmq_close(socket); - socket = NULL; - } - if (socket) { - void *old_socket = r->socket; - if (old_socket) { - zmq_close(old_socket); - } - struct message_remote_queue * queue = r->queue; - - if (queue) { - struct skynet_remote_message msg; - while (!skynet_remotemq_pop(queue, &msg)) { - skynet_remotemq_push(Z->queue, &msg); - } - skynet_remotemq_release(queue); - r->queue = NULL; - } - - r->socket = socket; - } -} - -static void -_report_zmq_error(int rc) { - if (rc) { - fprintf(stderr, "zmq error : %s\n",zmq_strerror(errno)); - exit(1); - } -} - -static int -_isdecimal(int c) { - return c>='0' && c<='9'; -} - -// Name-updating protocols: -// -// 1) harbor_id=harbor_address -// 2) context_name=context_handle -static int -_split_name(uint8_t *buf, int len, int *np) { - uint8_t *sep; - if (len > 0 && _isdecimal(buf[0])) { - int i=0; - int n=0; - do { - n = n*10 + (buf[i]-'0'); - } while(++izmq_local,&content,0); - _report_zmq_error(rc); - int sz = zmq_msg_size(&content); - uint8_t * buffer = zmq_msg_data(&content); - - int n = 0; - int i = _split_name(buffer, sz, &n); - if (i == -1) { - char tmp[sz+1]; - memcpy(tmp,buffer,sz); - tmp[sz] = '\0'; - skynet_error(NULL, "Invalid master update [%s]",tmp); - zmq_msg_close(&content); - return; - } - - char tmp[sz-i]; - memcpy(tmp,buffer+i+1,sz-i-1); - tmp[sz-i-1]='\0'; - - if (n>0 && n <= REMOTE_MAX) { - _remote_harbor_update(n, tmp); - } else { - uint32_t source = strtoul(tmp,NULL,16); - if (source == 0) { - skynet_error(NULL, "Invalid master update [%s=%s]",(const char *)buffer,tmp); - } else { - _register_name((const char *)buffer, source); - } - } - - zmq_msg_close(&content); -} - -// Always in main harbor thread -static void -remote_query_harbor(int harbor_id) { - char tmp[32]; - int sz = sprintf(tmp,"%d",harbor_id); - zmq_msg_t request; - zmq_msg_init_size(&request,sz); - memcpy(zmq_msg_data(&request),tmp,sz); - zmq_send(Z->zmq_master_request, &request, 0); - zmq_msg_close(&request); - zmq_msg_t reply; - zmq_msg_init(&reply); - int rc = zmq_recv(Z->zmq_master_request, &reply, 0); - _report_zmq_error(rc); - sz = zmq_msg_size(&reply); - char tmp2[sz+1]; - memcpy(tmp2,zmq_msg_data(&reply),sz); - tmp2[sz] = '\0'; - if (sz == 0) { - skynet_error(NULL,"Request harbor %d failed",harbor_id); - } else { - _remote_harbor_update(harbor_id, tmp2); - } - zmq_msg_close(&reply); -} - -// Always in main harbor thread -static void -_remote_register_name(const char *name, uint32_t source) { - char tmp[strlen(name) + 20]; - int sz = sprintf(tmp,"%s=%X",name,source); - zmq_msg_t msg; - zmq_msg_init_size(&msg,sz); - memcpy(zmq_msg_data(&msg), tmp , sz); - zmq_send(Z->zmq_master_request, &msg,0); - zmq_msg_close(&msg); - zmq_msg_init(&msg); - int rc = zmq_recv(Z->zmq_master_request, &msg,0); - _report_zmq_error(rc); - zmq_msg_close(&msg); -} - -// Always in main harbor thread -static void -_remote_query_name(const char *name) { - int sz = strlen(name); - zmq_msg_t msg; - zmq_msg_init_size(&msg,sz); - memcpy(zmq_msg_data(&msg), name , sz); - zmq_send(Z->zmq_master_request, &msg,0); - zmq_msg_close(&msg); - zmq_msg_init(&msg); - int rc = zmq_recv(Z->zmq_master_request, &msg,0); - _report_zmq_error(rc); - sz = zmq_msg_size(&msg); - char tmp[sz+1]; - memcpy(tmp, zmq_msg_data(&msg),sz); - tmp[sz] = '\0'; - - uint32_t addr = strtoul(tmp,NULL,16); - _register_name(name,addr); - - zmq_msg_close(&msg); -} - -// Always in main harbor thread -static void -free_message(void *data, void *hint) { - free(data); -} - -static void -remote_socket_send(void * socket, struct skynet_remote_message *msg) { - struct remote_header rh; - rh.source = msg->message.source; - rh.destination = msg->destination; - rh.session = msg->message.session; - zmq_msg_t part; - zmq_msg_init_size(&part,sizeof(struct remote_header)); - uint8_t * buffer = zmq_msg_data(&part); - remote_header_to_buffer(&rh,buffer); - zmq_send(socket, &part, ZMQ_SNDMORE); - zmq_msg_close(&part); - - zmq_msg_init_data(&part,msg->message.data,msg->message.sz,free_message,NULL); - zmq_send(socket, &part, 0); - zmq_msg_close(&part); -} - -// Always in main harbor thread - -// remote message has two part -// when part one is nil (size == 0), part two is name update -// Or part one is source:destination (8 bytes little endian), part two is a binary block for message -static void -_remote_recv() { - zmq_msg_t header; - zmq_msg_init(&header); - int rc = zmq_recv(Z->zmq_local,&header,0); - _report_zmq_error(rc); - size_t s = zmq_msg_size(&header); - if (s!=sizeof(struct remote_header)) { - // s should be 0 - if (s>0) { - char tmp[s+1]; - memcpy(tmp, zmq_msg_data(&header),s); - tmp[s] = '\0'; - skynet_error(NULL,"Invalid master header [%s]",tmp); - } - _name_update(); - return; - } - uint8_t * buffer = zmq_msg_data(&header); - struct remote_header rh; - buffer_to_remote_header(buffer, &rh); - zmq_close(&header); - - zmq_msg_t * data = malloc(sizeof(zmq_msg_t)); - zmq_msg_init(data); - rc = zmq_recv(Z->zmq_local,data,0); - _report_zmq_error(rc); - - struct skynet_message msg; - msg.session = (int)rh.session; - msg.source = rh.source; - msg.data = data; - msg.sz = zmq_msg_size(data); - - // push remote message to local message queue - 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_remote_message msg; - while (!skynet_remotemq_pop(Z->queue,&msg)) { -_goback: - if (msg.destination == SKYNET_SYSTEM_NAME) { - // register name - char * name = msg.message.data; - - if (msg.message.source) { - _remote_register_name(name, msg.message.source); - } else { - _remote_query_name(name); - } - - 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_remotemq_create(); - skynet_remotemq_push(r->queue, &msg); - remote_query_harbor(harbor_id); - } else { - skynet_remotemq_push(r->queue, &msg); - } - } else { - remote_socket_send(r->socket, &msg); - } - } - } - __sync_lock_release(&Z->notice_event); - // double check - if (!skynet_remotemq_pop(Z->queue,&msg)) { - __sync_lock_test_and_set(&Z->notice_event, 1); - goto _goback; - } -} - -// Main harbor thread -void * -skynet_harbor_dispatch_thread(void *ud) { - zmq_pollitem_t items[2]; - - items[0].socket = Z->zmq_queue_notice; - items[0].events = ZMQ_POLLIN; - items[1].socket = Z->zmq_local; - items[1].events = ZMQ_POLLIN; - - for (;;) { - zmq_poll(items,2,-1); - if (items[0].revents) { - zmq_msg_t msg; - zmq_msg_init(&msg); - int rc = zmq_recv(Z->zmq_queue_notice,&msg,0); - _report_zmq_error(rc); - zmq_msg_close(&msg); - _remote_send(); - } - if (items[1].revents) { - _remote_recv(); - } - } -} - -// Call only at init -static void -register_harbor(void *request, const char *local, int harbor) { - char tmp[1024]; - sprintf(tmp,"%d",harbor); - size_t sz = strlen(tmp); - - zmq_msg_t req; - zmq_msg_init_size (&req , sz); - memcpy(zmq_msg_data(&req),tmp,sz); - zmq_send (request, &req, 0); - zmq_msg_close (&req); - - zmq_msg_t reply; - zmq_msg_init (&reply); - int rc = zmq_recv(request, &reply, 0); - _report_zmq_error(rc); - - sz = zmq_msg_size (&reply); - if (sz > 0) { - memcpy(tmp,zmq_msg_data(&reply),sz); - tmp[sz] = '\0'; - fprintf(stderr, "Harbor %d is already registered by %s\n", harbor, tmp); - exit(1); - } - zmq_msg_close (&reply); - - sprintf(tmp,"%d=%s",harbor,local); - sz = strlen(tmp); - zmq_msg_init_size (&req , sz); - memcpy(zmq_msg_data(&req),tmp,sz); - zmq_send (request, &req, 0); - zmq_msg_close (&req); - - zmq_msg_init (&reply); - rc = zmq_recv (request, &reply, 0); - _report_zmq_error(rc); - zmq_msg_close (&reply); -} - -void -skynet_harbor_init(void * context, const char * master, const char *local, int harbor) { - if (harbor <=0 || harbor>255 || strlen(local) > 512) { - fprintf(stderr,"Invalid harbor id\n"); - exit(1); - } - void *request = zmq_socket (context, ZMQ_REQ); - int r = zmq_connect(request, master); - if (r<0) { - fprintf(stderr, "Can't connect to master: %s\n",master); - exit(1); - } - void *harbor_socket = zmq_socket(context, ZMQ_PULL); - r = zmq_bind(harbor_socket, local); - if (r<0) { - fprintf(stderr, "Can't bind to local : %s\n",local); - exit(1); - } - register_harbor(request,local, harbor); - printf("Start harbor on : %s\n",local); - - struct harbor * h = malloc(sizeof(*h)); - memset(h, 0, sizeof(*h)); - h->zmq_context = context; - h->zmq_master_request = request; - h->zmq_local = harbor_socket; - h->map = _hash_new(); - h->harbor = harbor; - 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); - - Z = h; -} - -// thread safe api -void * -skynet_harbor_message_open(struct skynet_message * message) { - return zmq_msg_data(message->data); -} - -void -skynet_harbor_message_close(struct skynet_message * message) { - zmq_msg_close(message->data); + assert(number == 0); + skynet_context_send(REMOTE, rname, sizeof(*rname), 0, 0); } int skynet_harbor_message_isremote(uint32_t handle) { - int harbor_id = handle >> HANDLE_REMOTE_SHIFT; - return !(harbor_id == 0 || harbor_id == Z->harbor); + return !(handle & HARBOR); +} + +void +skynet_harbor_init(int harbor) { + HARBOR = harbor << HANDLE_REMOTE_SHIFT; +} + +int +skynet_harbor_start(const char * master, const char *local) { + size_t sz = strlen(master) + strlen(local) + 32; + char args[sz]; + sprintf(args, "%s %s %d",master,local,HARBOR >> HANDLE_REMOTE_SHIFT); + struct skynet_context * inst = skynet_context_new("harbor",args); + if (inst == NULL) { + return 1; + } + REMOTE = inst; + + return 0; } diff --git a/skynet-src/skynet_harbor.h b/skynet-src/skynet_harbor.h index 235420d9..6ad6023a 100644 --- a/skynet-src/skynet_harbor.h +++ b/skynet-src/skynet_harbor.h @@ -2,20 +2,30 @@ #define SKYNET_HARBOR_H #include +#include -struct skynet_message; +#define GLOBALNAME_LENGTH 16 +#define REMOTE_MAX 256 -void skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * message); -void skynet_harbor_register(const char *name, uint32_t handle); +// reserve high 8 bits for remote id +#define HANDLE_MASK 0xffffff +#define HANDLE_REMOTE_SHIFT 24 -// remote message is diffrent from local message. -// We must use these api to open and close message , see skynet_server.c +struct remote_name { + char name[GLOBALNAME_LENGTH]; + uint32_t handle; +}; + +struct remote_message { + struct remote_name destination; + const void * message; + size_t sz; +}; + +void skynet_harbor_send(struct remote_message *rmsg, uint32_t source, int session); +void skynet_harbor_register(struct remote_name *rname); int skynet_harbor_message_isremote(uint32_t handle); -void * skynet_harbor_message_open(struct skynet_message * message); -void skynet_harbor_message_close(struct skynet_message * message); - -// harbor worker thread -void * skynet_harbor_dispatch_thread(void *ud); -void skynet_harbor_init(void * context, const char * master, const char *local, int harbor); +void skynet_harbor_init(int harbor); +int skynet_harbor_start(const char * master, const char *local); #endif diff --git a/skynet-src/skynet_imp.h b/skynet-src/skynet_imp.h index a9eedf11..bd76b1bd 100644 --- a/skynet-src/skynet_imp.h +++ b/skynet-src/skynet_imp.h @@ -10,7 +10,7 @@ struct skynet_config { const char * master; const char * local; const char * start; - int standalone; + const char * standalone; }; void skynet_start(struct skynet_config * config); diff --git a/skynet-src/skynet_logger.c b/skynet-src/skynet_logger.c index 64a655be..4df320be 100644 --- a/skynet-src/skynet_logger.c +++ b/skynet-src/skynet_logger.c @@ -24,12 +24,14 @@ logger_release(struct logger * inst) { free(inst); } -static void +static int _logger(struct skynet_context * context, void *ud, int session, const char * uid, const void * msg, size_t sz) { struct logger * inst = ud; fprintf(inst->handle, "[%s] ",uid); fwrite(msg, sz , 1, inst->handle); fprintf(inst->handle, "\n"); + + return 0; } int diff --git a/skynet-src/skynet_main.c b/skynet-src/skynet_main.c index 8d4655c9..607dc7bb 100644 --- a/skynet-src/skynet_main.c +++ b/skynet-src/skynet_main.c @@ -20,6 +20,7 @@ optint(const char *key, int opt) { return strtol(str, NULL, 10); } +/* static int optboolean(const char *key, int opt) { const char * str = skynet_getenv(key); @@ -29,7 +30,7 @@ optboolean(const char *key, int opt) { } return strcmp(str,"true")==0; } - +*/ static const char * optstring(const char *key,const char * opt) { const char * str = skynet_getenv(key); @@ -104,10 +105,10 @@ main(int argc, char *argv[]) { config.module_path = optstring("cpath","./service/?.so"); config.logger = optstring("logger",NULL); config.harbor = optint("harbor", 1); - config.master = optstring("master","tcp://127.0.0.1:2012"); + config.master = optstring("master","127.0.0.1:2012"); config.start = optstring("start","main.lua"); - config.local = optstring("address","tcp://127.0.0.1:2525"); - config.standalone = optboolean("standalone",0); + config.local = optstring("address","127.0.0.1:2525"); + config.standalone = optstring("standalone",NULL); lua_close(L); diff --git a/skynet-src/skynet_mq.c b/skynet-src/skynet_mq.c index d83ef831..5d96543c 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 in_global; struct skynet_message *queue; }; @@ -24,14 +25,6 @@ struct global_queue { 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 @@ -42,7 +35,7 @@ _lock_global_queue() { #define LOCK(q) while (__sync_lock_test_and_set(&(q)->lock,1)) {} #define UNLOCK(q) __sync_lock_release(&(q)->lock); -void +static void skynet_globalmq_push(struct message_queue * queue) { struct global_queue *q= Q; LOCK(q) @@ -95,6 +88,7 @@ skynet_mq_create(uint32_t handle) { q->head = 0; q->tail = 0; q->lock = 0; + q->in_global = 1; q->queue = malloc(sizeof(struct skynet_message) * q->cap); return q; @@ -114,7 +108,7 @@ skynet_mq_handle(struct message_queue *q) { int skynet_mq_pop(struct message_queue *q, struct skynet_message *message) { - int ret = -1; + int ret = 1; LOCK(q) if (q->head != q->tail) { @@ -124,6 +118,10 @@ skynet_mq_pop(struct message_queue *q, struct skynet_message *message) { q->head = 0; } } + + if (ret) { + q->in_global = 0; + } UNLOCK(q) @@ -134,23 +132,30 @@ void skynet_mq_push(struct message_queue *q, struct skynet_message *message) { LOCK(q) - q->queue[q->tail] = *message; - if (++ q->tail >= q->cap) { - q->tail = 0; + if (message) { + q->queue[q->tail] = *message; + if (++ q->tail >= q->cap) { + q->tail = 0; + } + + if (q->head == q->tail) { + struct skynet_message *new_queue = malloc(sizeof(struct skynet_message) * q->cap * 2); + int i; + for (i=0;icap;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; + } } - if (q->head == q->tail) { - struct skynet_message *new_queue = malloc(sizeof(struct skynet_message) * q->cap * 2); - int i; - for (i=0;icap;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; + if (q->in_global == 0) { + skynet_globalmq_push(q); + q->in_global = 1; } UNLOCK(q) @@ -170,68 +175,8 @@ skynet_mq_init(int n) { 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_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;icap;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) +skynet_mq_force_push(struct message_queue * queue) { + assert(queue->in_global); + skynet_globalmq_push(queue); } diff --git a/skynet-src/skynet_mq.h b/skynet-src/skynet_mq.h index acd222f1..375700e1 100644 --- a/skynet-src/skynet_mq.h +++ b/skynet-src/skynet_mq.h @@ -13,7 +13,6 @@ struct skynet_message { struct message_queue; -void skynet_globalmq_push(struct message_queue *); struct message_queue * skynet_globalmq_pop(void); struct message_queue * skynet_mq_create(uint32_t handle); @@ -23,19 +22,7 @@ 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_force_push(struct message_queue *q); void skynet_mq_init(int cap); diff --git a/skynet-src/skynet_server.c b/skynet-src/skynet_server.c index ee46e6f0..c2e092ca 100644 --- a/skynet-src/skynet_server.c +++ b/skynet-src/skynet_server.c @@ -25,10 +25,9 @@ struct skynet_context { void * cb_ud; skynet_cb cb; int session_id; - int in_global_queue; int init; uint32_t forward; - char * forward_address; + char forward_address[GLOBALNAME_LENGTH]; struct message_queue *queue; }; @@ -36,10 +35,11 @@ static void _id_to_hex(char * str, uint32_t id) { int i; static char hex[16] = { '0','1','2','3','4','5','6','7','8','9','A','B','C','D','E','F' }; + str[0] = ':'; for (i=0;i<8;i++) { - str[i] = hex[(id >> ((7-i) * 4))&0xf]; + str[i+1] = hex[(id >> ((7-i) * 4))&0xf]; } - str[8] = '\0'; + str[9] = '\0'; } struct skynet_context * @@ -58,17 +58,13 @@ skynet_context_new(const char * name, const char *param) { ctx->ref = 2; ctx->cb = NULL; ctx->cb_ud = NULL; - // a trick, global loop can't be dispatch in init process - ctx->in_global_queue = 1; ctx->forward = 0; - ctx->forward_address = NULL; - ctx->session_id = 0; + ctx->forward_address[0] = '\0'; ctx->init = 0; ctx->handle = skynet_handle_register(ctx); char * uid = ctx->handle_name; - uid[0] = ':'; - _id_to_hex(uid+1, ctx->handle); + _id_to_hex(uid, ctx->handle); struct message_queue * queue = ctx->queue = skynet_mq_create(ctx->handle); // init function maybe use ctx->handle, so it must init at last @@ -77,15 +73,10 @@ skynet_context_new(const char * name, const char *param) { struct skynet_context * ret = skynet_context_release(ctx); if (ret) { ctx->init = 1; - skynet_globalmq_push(ctx->queue); - return ret; - } else { - // because of ctx->in_global_queue == 1 , so we should release queue here. - skynet_mq_release(queue); } - return NULL; + skynet_mq_force_push(queue); + return ret; } else { - ctx->in_global_queue = 0; skynet_context_release(ctx); skynet_handle_retire(ctx->handle); return NULL; @@ -95,7 +86,7 @@ skynet_context_new(const char * name, const char *param) { int skynet_context_newsession(struct skynet_context *ctx) { int session = ++ctx->session_id; - if (session >= 0x7fffffff) { + if (session >= SESSION_CLIENT) { ctx->session_id = 1; return 1; } @@ -111,9 +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); - if (!ctx->in_global_queue) { - skynet_mq_release(ctx->queue); - } + skynet_mq_push(ctx->queue , NULL); free(ctx); } @@ -126,13 +115,29 @@ skynet_context_release(struct skynet_context *ctx) { return ctx; } +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); + skynet_context_release(ctx); + + return 0; +} + static int _forwarding(struct skynet_context *ctx, struct skynet_message *msg) { if (ctx->forward) { uint32_t des = ctx->forward; ctx->forward = 0; if (skynet_harbor_message_isremote(des)) { - skynet_harbor_send(NULL, des, msg); + struct remote_message * rmsg = malloc(sizeof(*rmsg)); + rmsg->destination.handle = des; + rmsg->message = msg->data; + rmsg->sz = msg->sz; + skynet_harbor_send(rmsg, msg->source, msg->session); } else { if (skynet_context_push(des, msg)) { free(msg->data); @@ -142,10 +147,14 @@ _forwarding(struct skynet_context *ctx, struct skynet_message *msg) { } return 1; } - if (ctx->forward_address) { - skynet_harbor_send(ctx->forward_address, 0, msg); - free(ctx->forward_address); - ctx->forward_address = NULL; + if (ctx->forward_address[0]) { + struct remote_message * rmsg = malloc(sizeof(*rmsg)); + memcpy(rmsg->destination.name, ctx->forward_address,GLOBALNAME_LENGTH); + rmsg->destination.handle = 0; + rmsg->message = msg->data; + rmsg->sz = msg->sz; + skynet_harbor_send(rmsg, msg->source, msg->session); + ctx->forward_address[0] = '\0'; return 1; } return 0; @@ -158,21 +167,10 @@ _dispatch_message(struct skynet_context *ctx, struct skynet_message *msg) { ctx->cb(ctx, ctx->cb_ud, msg->session, NULL, msg->data, msg->sz); } else { char tmp[10]; - tmp[0] = ':'; - int not_delete; - _id_to_hex(tmp+1, msg->source); - if (skynet_harbor_message_isremote(msg->source)) { - void * data = skynet_harbor_message_open(msg); - ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, data, msg->sz); - not_delete = _forwarding(ctx, msg); - if (!not_delete) { - skynet_harbor_message_close(msg); - } - } else { - ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, msg->data, msg->sz); - not_delete = _forwarding(ctx, msg); - } - if (!not_delete) { + _id_to_hex(tmp, msg->source); + int reserve = ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, msg->data, msg->sz); + reserve |= _forwarding(ctx, msg); + if (!reserve) { free(msg->data); } } @@ -185,9 +183,6 @@ _drop_queue(struct message_queue *q) { int s = 0; while(!skynet_mq_pop(q, &msg)) { ++s; - if (skynet_harbor_message_isremote(msg.source)) { - skynet_harbor_message_close(&msg); - } free(msg.data); } skynet_mq_release(q); @@ -211,19 +206,13 @@ skynet_context_message_dispatch(void) { return 0; } - assert(ctx->in_global_queue); - struct skynet_message msg; if (skynet_mq_pop(q,&msg)) { - __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); - } free(msg.data); skynet_error(NULL, "Drop message from %x to %x without callback , size = %d",msg.source, handle, (int)msg.sz); } else { @@ -231,12 +220,23 @@ skynet_context_message_dispatch(void) { } assert(q == ctx->queue); - skynet_globalmq_push(q); + skynet_mq_force_push(q); skynet_context_release(ctx); return 0; } +static void +_copy_name(char name[GLOBALNAME_LENGTH], const char * addr) { + int i; + for (i=0;ihandle, param + 1); } else { assert(context->handle!=0); - int i; - for (i=0;i= '0' && param[i] <= '9')) { - break; - } - } - assert(param[i]); - skynet_harbor_register(param, context->handle); + struct remote_name *rname = malloc(sizeof(*rname)); + _copy_name(rname->name, param); + rname->handle = context->handle; + skynet_harbor_register(rname); return NULL; } } @@ -290,7 +286,10 @@ skynet_command(struct skynet_context * context, const char * cmd , const char * if (name[0] == '.') { return skynet_handle_namehandle(handle_id, name + 1); } else { - skynet_harbor_register(name, handle_id); + struct remote_name *rname = malloc(sizeof(*rname)); + _copy_name(rname->name, name); + rname->handle = handle_id; + skynet_harbor_register(rname); } return NULL; } @@ -327,8 +326,7 @@ skynet_command(struct skynet_context * context, const char * cmd , const char * fprintf(stderr, "Launch %s %s failed\n",mod,args); return NULL; } else { - context->result[0] = ':'; - _id_to_hex(context->result+1, inst->handle); + _id_to_hex(context->result, inst->handle); return context->result; } } @@ -366,7 +364,7 @@ skynet_command(struct skynet_context * context, const char * cmd , const char * void skynet_forward(struct skynet_context * context, const char * addr) { uint32_t des = 0; - assert(context->forward == 0 && context->forward_address == NULL); + assert(context->forward == 0 && context->forward_address[0] == '\0'); if (addr[0] == ':') { des = strtol(addr+1, NULL, 16); assert(des != 0); @@ -377,7 +375,7 @@ skynet_forward(struct skynet_context * context, const char * addr) { return; } } else { - context->forward_address = strdup(addr); + _copy_name(context->forward_address,addr); return; } context->forward = des; @@ -385,9 +383,14 @@ skynet_forward(struct skynet_context * context, const char * addr) { int skynet_send(struct skynet_context * context, const char * source, const char * addr , int session, void * data, size_t sz, int flags) { + int session_id = session; uint32_t source_handle; if (source == NULL) { source_handle = context->handle; + if (session < 0) { + session = skynet_context_newsession(context); + session_id = - session; + } } else { assert (source[0] == ':'); source_handle = strtoul(source+1, NULL, 16); @@ -401,11 +404,6 @@ skynet_send(struct skynet_context * context, const char * source, const char * a memcpy(msg, data, sz); msg[sz] = '\0'; } - int session_id = session; - if (session < 0) { - session = skynet_context_newsession(context); - session_id = - session; - } if (addr == NULL) { return session; } @@ -420,29 +418,35 @@ skynet_send(struct skynet_context * context, const char * source, const char * a return session; } } else { - struct skynet_message smsg; - smsg.source = source_handle; - smsg.session = session_id; - smsg.data = msg; - smsg.sz = sz; - skynet_harbor_send(addr, 0, &smsg); + struct remote_message * rmsg = malloc(sizeof(*rmsg)); + _copy_name(rmsg->destination.name, addr); + rmsg->destination.handle = 0; + rmsg->message = msg; + rmsg->sz = sz; + skynet_harbor_send(rmsg, source_handle, session_id); return session; } assert(des > 0); - struct skynet_message smsg; - smsg.source = source_handle; - smsg.session = session_id; - smsg.data = msg; - smsg.sz = sz; - if (skynet_harbor_message_isremote(des)) { - skynet_harbor_send(NULL, des, &smsg); - } else if (skynet_context_push(des, &smsg)) { - free(msg); - skynet_error(NULL, "Drop message from %x to %s (size=%d)", smsg.source, addr, (int)sz); - return -1; + struct remote_message * rmsg = malloc(sizeof(*rmsg)); + rmsg->destination.handle = des; + rmsg->message = msg; + rmsg->sz = sz; + skynet_harbor_send(rmsg, source_handle, session_id); + } else { + struct skynet_message smsg; + smsg.source = source_handle; + smsg.session = session_id; + smsg.data = msg; + smsg.sz = sz; + + if (skynet_context_push(des, &smsg)) { + free(msg); + skynet_error(NULL, "Drop message from %x to %s (size=%d)", smsg.source, addr, (int)sz); + return -1; + } } return session; } @@ -464,17 +468,13 @@ 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); +void +skynet_context_send(struct skynet_context * ctx, void * msg, size_t sz, uint32_t source, int session) { + struct skynet_message smsg; + smsg.source = source; + smsg.session = session; + smsg.data = msg; + smsg.sz = sz; - return 0; + skynet_mq_push(ctx->queue, &smsg); } diff --git a/skynet-src/skynet_server.h b/skynet-src/skynet_server.h index 95ad8a86..992ee1d2 100644 --- a/skynet-src/skynet_server.h +++ b/skynet-src/skynet_server.h @@ -2,6 +2,7 @@ #define SKYNET_SERVER_H #include +#include struct skynet_context; struct skynet_message; @@ -12,6 +13,7 @@ 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); int skynet_context_push(uint32_t handle, struct skynet_message *message); +void skynet_context_send(struct skynet_context * context, void * msg, size_t sz, uint32_t source, int session); int skynet_context_newsession(struct skynet_context *); int skynet_context_message_dispatch(void); // return 1 when block diff --git a/skynet-src/skynet_start.c b/skynet-src/skynet_start.c index e7e4f321..ef863095 100644 --- a/skynet-src/skynet_start.c +++ b/skynet-src/skynet_start.c @@ -5,12 +5,10 @@ #include "skynet_module.h" #include "skynet_timer.h" #include "skynet_harbor.h" -#include "skynet_master.h" #include #include #include -#include static void * _timer(void *p) { @@ -33,59 +31,51 @@ _worker(void *p) { static void _start(int thread) { - pthread_t pid[thread+2]; + pthread_t pid[thread+1]; pthread_create(&pid[0], NULL, _timer, NULL); - pthread_create(&pid[1], NULL, skynet_harbor_dispatch_thread, NULL); int i; - for (i=2;icontext, args->port); - free(args); - return NULL; -} - -static void -_start_master(void * context, const char * port) { - pthread_t pid; - struct master_arg * args = malloc(sizeof(*args)); - args->context = context; - args->port = port; - pthread_create(&pid, NULL, _master_thread, args); +static int +_start_master(const char * master) { + struct skynet_context *ctx = skynet_context_new("master", master); + if (ctx == NULL) + return 1; + return 0; } void skynet_start(struct skynet_config * config) { - void *context = zmq_init (1); - assert(context); - if (config->standalone) { - _start_master(context, config->master); - } - // harbor must be init first - skynet_harbor_init(context, config->master , config->local, config->harbor); + skynet_harbor_init(config->harbor); skynet_handle_init(config->harbor); skynet_mq_init(config->mqueue_size); skynet_module_init(config->module_path); skynet_timer_init(); + + if (config->standalone) { + if (_start_master(config->standalone)) { + return; + } + } + // harbor must be init first + if (skynet_harbor_start(config->master , config->local)) { + return; + } + struct skynet_context *ctx; ctx = skynet_context_new("logger", config->logger); - assert(ctx); + if (ctx == NULL) { + return; + } ctx = skynet_context_new("snlua", config->start); _start(config->thread);