diff --git a/.gitignore b/.gitignore index 2cbccc98..9ee6cde3 100644 --- a/.gitignore +++ b/.gitignore @@ -1,10 +1,10 @@ *.o *.a -skynet -skynet.pid +./skynet +./skynet.pid 3rd/lua/lua 3rd/lua/luac -cservice -luaclib +./cservice +./luaclib *.so *.dSYM diff --git a/HISTORY.md b/HISTORY.md index 89dd918e..d4066410 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -1,6 +1,16 @@ +Dev version +----------- +* Optimize redis driver `compose_message`. +* Add module skynet.harbor for monitor harbor connect/disconnect, see test/testharborlink.lua . +* cluster.open support cluster name. +* Add new api skynet.packstring , and skynet.unpack support lua string +* socket.listen support put port into address. (address:port) +* Redesign harbor/master/dummy, remove lots of C code and rewite in lua. +* Remove block connect api, queue sending message during connecting now. +* Add skynet.time() + v0.3.2 (2014-6-23) ---------- -* Bugifx : mongo driver and lua-bson . (objectid encoding, and gc problem when mongdo driver reply bson object). * Bugfix : cluster (double free). * Add socket.header() to decode big-endian package header (and fix the bug in cluster). diff --git a/Makefile b/Makefile index 3c5fd446..45631905 100644 --- a/Makefile +++ b/Makefile @@ -40,7 +40,7 @@ jemalloc : $(MALLOC_STATICLIB) # skynet -CSERVICE = snlua logger gate master harbor dummy +CSERVICE = snlua logger gate harbor LUA_CLIB = skynet socketdriver int64 bson mongo md5 netpack \ cjson clientsocket memory profile multicast \ cluster diff --git a/examples/cluster1.lua b/examples/cluster1.lua index 73fad761..551b0f38 100644 --- a/examples/cluster1.lua +++ b/examples/cluster1.lua @@ -5,5 +5,5 @@ skynet.start(function() skynet.newservice("simpledb") print(skynet.call("SIMPLEDB", "lua", "SET", "a", "foobar")) print(skynet.call("SIMPLEDB", "lua", "GET", "a")) - cluster.open(2528) + cluster.open "db" end) diff --git a/examples/globallog.lua b/examples/globallog.lua index ec431cce..ffd0a0f3 100644 --- a/examples/globallog.lua +++ b/examples/globallog.lua @@ -1,8 +1,8 @@ local skynet = require "skynet" skynet.start(function() - skynet.dispatch("text", function(session, address, text) - print("[GLOBALLOG]", skynet.address(address),text) + skynet.dispatch("lua", function(session, address, ...) + print("[GLOBALLOG]", skynet.address(address), ...) end) skynet.register "LOG" end) diff --git a/examples/main_log.lua b/examples/main_log.lua index 8a5194a0..f74b563f 100644 --- a/examples/main_log.lua +++ b/examples/main_log.lua @@ -2,7 +2,6 @@ local skynet = require "skynet" skynet.start(function() print("Log server start") - local service = skynet.newservice("service_mgr") skynet.monitor "simplemonitor" local log = skynet.newservice("globallog") skynet.exit() diff --git a/lualib-src/lua-seri.c b/lualib-src/lua-seri.c index 69e1c2af..c6544462 100644 --- a/lualib-src/lua-seri.c +++ b/lualib-src/lua-seri.c @@ -581,8 +581,16 @@ _luaseri_unpack(lua_State *L) { if (lua_isnoneornil(L,1)) { return 0; } - void * buffer = lua_touserdata(L,1); - int len = luaL_checkinteger(L,2); + void * buffer; + int len; + if (lua_type(L,1) == LUA_TSTRING) { + size_t sz; + buffer = (void *)lua_tolstring(L,1,&sz); + len = (int)sz; + } else { + buffer = lua_touserdata(L,1); + len = luaL_checkinteger(L,2); + } if (len == 0) { return 0; } diff --git a/lualib-src/lua-skynet.c b/lualib-src/lua-skynet.c index 0025dd52..0189b61b 100644 --- a/lualib-src/lua-skynet.c +++ b/lualib-src/lua-skynet.c @@ -149,7 +149,7 @@ _sendname(lua_State *L, struct skynet_context * context, const char * dest) { luaL_error(L, "skynet.send invalid param %s", lua_typename(L,lua_type(L,4))); } if (session < 0) { - luaL_error(L, "skynet.send session (%d) < 0", session); + return 0; } lua_pushinteger(L,session); return 1; @@ -221,7 +221,7 @@ _send(lua_State *L) { } if (session < 0) { // send to invalid address - // todo: maybe throw error is better + // todo: maybe throw error whould be better return 0; } lua_pushinteger(L,session); @@ -289,6 +289,16 @@ _harbor(lua_State *L) { return 2; } +static int +lpackstring(lua_State *L) { + _luaseri_pack(L); + char * str = (char *)lua_touserdata(L, -2); + int sz = lua_tointeger(L, -1); + lua_pushlstring(L, str, sz); + skynet_free(str); + return 1; +} + int luaopen_skynet_c(lua_State *L) { luaL_checkversion(L); @@ -303,6 +313,7 @@ luaopen_skynet_c(lua_State *L) { { "harbor", _harbor }, { "pack", _luaseri_pack }, { "unpack", _luaseri_unpack }, + { "packstring", lpackstring }, { "callback", _callback }, { NULL, NULL }, }; diff --git a/lualib/cluster.lua b/lualib/cluster.lua index e4f6cfa2..bc4312b3 100644 --- a/lualib/cluster.lua +++ b/lualib/cluster.lua @@ -9,7 +9,11 @@ function cluster.call(node, address, ...) end function cluster.open(port) - skynet.call(clusterd, "lua", "listen", "0.0.0.0", port) + if type(port) == "string" then + skynet.call(clusterd, "lua", "listen", port) + else + skynet.call(clusterd, "lua", "listen", "0.0.0.0", port) + end end skynet.init(function() diff --git a/lualib/redis.lua b/lualib/redis.lua index d28278dc..cd9f43f4 100644 --- a/lualib/redis.lua +++ b/lualib/redis.lua @@ -1,6 +1,7 @@ local skynet = require "skynet" local socket = require "socket" local socketchannel = require "socketchannel" +local int64 = require "int64" local table = table local string = string @@ -94,31 +95,56 @@ function command:disconnect() setmetatable(self, nil) end -local function compose_message(msg) - if #msg == 1 then - return msg[1] .. "\r\n" +-- msg could be any type of value +local function pack_value(lines, v) + if v == nil then + return end - local lines = { "*" .. #msg } - for _,v in ipairs(msg) do - local t = type(v) - if t == "number" then - v = tostring(v) - elseif t == "userdata" then - v = int64.tostring(int64.new(v),10) - end - table.insert(lines,"$"..#v) - table.insert(lines,v) - end - table.insert(lines,"") - local cmd = table.concat(lines,"\r\n") - return cmd + local t = type(v) + if t == "number" then + v = tostring(v) + elseif t == "userdata" then + v = int64.tostring(int64.new(v),10) + end + table.insert(lines,"$"..#v) + table.insert(lines,v) +end + +local function compose_message(cmd, msg) + local len = 1 + local t = type(msg) + + if t == "table" then + len = len + #msg + elseif t ~= nil then + len = len + 1 + end + + local lines = {"*" .. len} + pack_value(lines, cmd) + + if t == "table" then + for _,v in ipairs(msg) do + pack_value(lines, v) + end + else + pack_value(lines, msg) + end + table.insert(lines, "") + + local chunk = table.concat(lines,"\r\n") + return chunk end setmetatable(command, { __index = function(t,k) local cmd = string.upper(k) - local f = function (self, ...) - return self[1]:request(compose_message { cmd, ... }, read_response) + local f = function (self, v, ...) + if type(v) == "table" then + return self[1]:request(compose_message(cmd, v), read_response) + else + return self[1]:request(compose_message(cmd, {v, ...}), read_response) + end end t[k] = f return f @@ -131,12 +157,12 @@ end function command:exists(key) local fd = self[1] - return fd:request(compose_message { "EXISTS", key }, read_boolean) + return fd:request(compose_message ("EXISTS", key), read_boolean) end function command:sismember(key, value) local fd = self[1] - return fd:request(compose_message { "SISMEMBER", key, value }, read_boolean) + return fd:request(compose_message ("SISMEMBER", {key, value}), read_boolean) end --- watch mode @@ -156,10 +182,10 @@ local function watch_login(obj, auth) so:request("AUTH "..auth.."\r\n", read_response) end for k in pairs(obj.__psubscribe) do - so:request(compose_message { "PSUBSCRIBE", k }) + so:request(compose_message ("PSUBSCRIBE", k)) end for k in pairs(obj.__subscribe) do - so:request(compose_message { "SUBSCRIBE", k }) + so:request(compose_message("SUBSCRIBE", k)) end end end @@ -192,7 +218,7 @@ local function watch_func( name ) local so = self.__sock for i = 1, select("#", ...) do local v = select(i, ...) - so:request(compose_message { NAME, v }) + so:request(compose_message(NAME, v)) end end end diff --git a/lualib/skynet.lua b/lualib/skynet.lua index 1d936052..2019f530 100644 --- a/lualib/skynet.lua +++ b/lualib/skynet.lua @@ -190,12 +190,33 @@ function skynet.wait() session_id_coroutine[session] = nil end +local function globalname(name, handle) + local c = string.sub(name,1,1) + assert(c ~= ':') + if c == '.' then + return false + end + + assert(#name <= 16) -- GLOBALNAME_LENGTH is 16, defined in skynet_harbor.h + assert(tonumber(name) == nil) -- global name can't be number + + local harbor = require "skynet.harbor" + + harbor.globalname(name, handle) + + return true +end + function skynet.register(name) - c.command("REG", name) + if not globalname(name) then + c.command("REG", name) + end end function skynet.name(name, handle) - c.command("NAME", name .. " " .. skynet.address(handle)) + if not globalname(name, handle) then + c.command("NAME", name .. " " .. skynet.address(handle)) + end end local self_handle @@ -226,6 +247,10 @@ function skynet.starttime() return tonumber(c.command("STARTTIME")) end +function skynet.time() + return skynet.now()/100 + skynet.starttime() -- get now first would be better +end + function skynet.exit() skynet.send(".launcher","lua","REMOVE",skynet.self()) c.command("EXIT") @@ -267,6 +292,7 @@ skynet.redirect = function(dest,source,typename,...) end skynet.pack = assert(c.pack) +skynet.packstring = assert(c.packstring) skynet.unpack = assert(c.unpack) skynet.tostring = assert(c.tostring) @@ -318,8 +344,8 @@ function skynet.dispatch(typename, func) p.dispatch = func end -local function unknown_request(session, address, msg, sz) - print("Unknown request :" , c.tostring(msg,sz)) +local function unknown_request(session, address, msg, sz, prototype) + skynet.error(string.format("Unknown request (%s): %s", prototype, c.tostring(msg,sz))) error(string.format("Unknown session : %d from %x", session, address)) end @@ -373,7 +399,7 @@ local function raw_dispatch_message(prototype, msg, sz, session, source, ...) session_coroutine_address[co] = source suspend(co, coroutine.resume(co, session,source, p.unpack(msg,sz, ...))) else - unknown_request(session, source, msg, sz) + unknown_request(session, source, msg, sz, proto[prototype]) end end end @@ -502,7 +528,7 @@ end local function init_service(start) local ok, err = xpcall(init_template, debug.traceback, start) if not ok then - print("init service failed:", err) + skynet.error("init service failed: " .. tostring(err)) skynet.send(".launcher","lua", "ERROR") skynet.exit() else diff --git a/lualib/skynet/harbor.lua b/lualib/skynet/harbor.lua new file mode 100644 index 00000000..5f470be8 --- /dev/null +++ b/lualib/skynet/harbor.lua @@ -0,0 +1,18 @@ +local skynet = require "skynet" + +local harbor = {} + +function harbor.globalname(name, handle) + handle = handle or skynet.self() + skynet.send(".slave", "lua", "REGISTER", name, handle) +end + +function harbor.link(id) + skynet.call(".slave", "lua", "LINK", id) +end + +function harbor.connect(id) + skynet.call(".slave", "lua", "CONNECT", id) +end + +return harbor diff --git a/lualib/socket.lua b/lualib/socket.lua index f5319ccf..d7bc49d5 100644 --- a/lualib/socket.lua +++ b/lualib/socket.lua @@ -272,7 +272,13 @@ function socket.invalid(id) return socket_pool[id] == nil end -socket.listen = assert(driver.listen) +function socket.listen(host, port, backlog) + if port == nil then + host, port = string.match(host, "([^:]+):(.+)$") + port = tonumber(port) + end + return driver.listen(host, port, backlog) +end function socket.lock(id) local s = socket_pool[id] diff --git a/platform.mk b/platform.mk index 345acabe..a112f0f2 100644 --- a/platform.mk +++ b/platform.mk @@ -31,7 +31,7 @@ macosx : EXPORT := macosx linux : SKYNET_LIBS += -ldl linux freebsd : SKYNET_LIBS += -lrt -# Turn off jemalloc and malloc hook on macosx and freebsd +# Turn off jemalloc and malloc hook on macosx macosx : MALLOC_STATICLIB := macosx : SKYNET_DEFINES :=-DNOUSE_JEMALLOC diff --git a/service-src/service_dummy.c b/service-src/service_dummy.c deleted file mode 100644 index 006a276f..00000000 --- a/service-src/service_dummy.c +++ /dev/null @@ -1,292 +0,0 @@ -#include "skynet.h" -#include "skynet_harbor.h" - -#include -#include - -#define HASH_SIZE 4096 -#define DEFAULT_QUEUE_SIZE 1024 - -struct msg { - uint8_t * buffer; - size_t size; -}; - -struct msg_queue { - int size; - int head; - int tail; - struct msg * data; -}; - -struct keyvalue { - struct keyvalue * next; - char key[GLOBALNAME_LENGTH]; - uint32_t hash; - uint32_t value; - struct msg_queue * queue; -}; - -struct hashmap { - struct keyvalue *node[HASH_SIZE]; -}; - -/* - message type (8bits) is in destination high 8bits - harbor id (8bits) is also in that place , but remote message doesn't need harbor id. - */ -struct remote_message_header { - uint32_t source; - uint32_t destination; - uint32_t session; -}; - -// 12 is sizeof(struct remote_message_header) -#define HEADER_COOKIE_LENGTH 12 - -struct dummy { - struct skynet_context *ctx; - struct hashmap * map; -}; - -// hash table - -static void -_push_queue(struct msg_queue * queue, const void * buffer, size_t sz, struct remote_message_header * header) { - // If there is only 1 free slot which is reserved to distinguish full/empty - // of circular buffer, expand it. - if (((queue->tail + 1) % queue->size) == queue->head) { - struct msg * new_buffer = skynet_malloc(queue->size * 2 * sizeof(struct msg)); - int i; - for (i=0;isize-1;i++) { - new_buffer[i] = queue->data[(i+queue->head) % queue->size]; - } - skynet_free(queue->data); - queue->data = new_buffer; - queue->head = 0; - queue->tail = queue->size - 1; - queue->size *= 2; - } - struct msg * slot = &queue->data[queue->tail]; - queue->tail = (queue->tail + 1) % queue->size; - - slot->buffer = skynet_malloc(sz + sizeof(*header)); - memcpy(slot->buffer, buffer, sz); - memcpy(slot->buffer + sz, header, sizeof(*header)); - slot->size = sz + sizeof(*header); -} - -static struct msg * -_pop_queue(struct msg_queue * queue) { - if (queue->head == queue->tail) { - return NULL; - } - struct msg * slot = &queue->data[queue->head]; - queue->head = (queue->head + 1) % queue->size; - return slot; -} - -static struct msg_queue * -_new_queue() { - struct msg_queue * queue = skynet_malloc(sizeof(*queue)); - queue->size = DEFAULT_QUEUE_SIZE; - queue->head = 0; - queue->tail = 0; - queue->data = skynet_malloc(DEFAULT_QUEUE_SIZE * sizeof(struct msg)); - - return queue; -} - -static void -_release_queue(struct msg_queue *queue) { - if (queue == NULL) - return; - struct msg * m = _pop_queue(queue); - while (m) { - skynet_free(m->buffer); - m = _pop_queue(queue); - } - skynet_free(queue->data); - skynet_free(queue); -} - -static struct keyvalue * -_hash_search(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { - uint32_t *ptr = (uint32_t*) name; - uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; - struct keyvalue * node = hash->node[h % HASH_SIZE]; - while (node) { - if (node->hash == h && strncmp(node->key, name, GLOBALNAME_LENGTH) == 0) { - return node; - } - node = node->next; - } - return NULL; -} - -static struct keyvalue * -_hash_insert(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { - uint32_t *ptr = (uint32_t *)name; - uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; - struct keyvalue ** pkv = &hash->node[h % HASH_SIZE]; - struct keyvalue * node = skynet_malloc(sizeof(*node)); - memcpy(node->key, name, GLOBALNAME_LENGTH); - node->next = *pkv; - node->queue = NULL; - node->hash = h; - node->value = 0; - *pkv = node; - - return node; -} - -static struct hashmap * -_hash_new() { - struct hashmap * h = skynet_malloc(sizeof(struct hashmap)); - memset(h,0,sizeof(*h)); - return h; -} - -static void -_hash_delete(struct hashmap *hash) { - int i; - for (i=0;inode[i]; - while (node) { - struct keyvalue * next = node->next; - _release_queue(node->queue); - skynet_free(node); - node = next; - } - } - skynet_free(hash); -} - -/////////////// - -struct dummy * -dummy_create(void) { - struct dummy * d = skynet_malloc(sizeof(*d)); - d->map = _hash_new(); - return d; -} - -void -dummy_release(struct dummy *d) { - _hash_delete(d->map); - skynet_free(d); -} - -static inline void -to_bigendian(uint8_t *buffer, uint32_t n) { - buffer[0] = (n >> 24) & 0xff; - buffer[1] = (n >> 16) & 0xff; - buffer[2] = (n >> 8) & 0xff; - buffer[3] = n & 0xff; -} - -static inline void -_header_to_message(const struct remote_message_header * header, uint8_t * message) { - to_bigendian(message , header->source); - to_bigendian(message+4 , header->destination); - to_bigendian(message+8 , header->session); -} - -static inline uint32_t -from_bigendian(uint32_t n) { - union { - uint32_t big; - uint8_t bytes[4]; - } u; - u.big = n; - return u.bytes[0] << 24 | u.bytes[1] << 16 | u.bytes[2] << 8 | u.bytes[3]; -} - -static inline void -_message_to_header(const uint32_t *message, struct remote_message_header *header) { - header->source = from_bigendian(message[0]); - header->destination = from_bigendian(message[1]); - header->session = from_bigendian(message[2]); -} - -static void -_dispatch_queue(struct dummy *h, struct msg_queue * queue, uint32_t handle, const char name[GLOBALNAME_LENGTH] ) { - struct msg * m = _pop_queue(queue); - while (m) { - struct remote_message_header cookie; - uint8_t *ptr = m->buffer + m->size - sizeof(cookie); - memcpy(&cookie, ptr, sizeof(cookie)); - int type = cookie.destination >> HANDLE_REMOTE_SHIFT; - skynet_send(h->ctx, cookie.source, handle , type | PTYPE_TAG_DONTCOPY, cookie.session, m->buffer, m->size - sizeof(cookie)); - m = _pop_queue(queue); - } -} - -static void -_update_name(struct dummy *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { - struct keyvalue * node = _hash_search(h->map, name); - if (node == NULL) { - node = _hash_insert(h->map, name); - } - node->value = handle; - if (node->queue) { - _dispatch_queue(h, node->queue, handle, name); - _release_queue(node->queue); - node->queue = NULL; - } -} - -static void -_send_name(struct dummy *h, uint32_t source, const char name[GLOBALNAME_LENGTH], int type, int session, const char * msg, size_t sz) { - struct keyvalue * node = _hash_search(h->map, name); - if (node == NULL) { - node = _hash_insert(h->map, name); - } - if (node->value == 0) { - if (node->queue == NULL) { - node->queue = _new_queue(); - } - struct remote_message_header header; - header.source = source; - header.destination = type << HANDLE_REMOTE_SHIFT; - header.session = (uint32_t)session; - _push_queue(node->queue, msg, sz, &header); - } else { - // local message - skynet_send(h->ctx, source, node->value , type | PTYPE_TAG_DONTCOPY, session, (void *)msg, sz); - } -} - -static int -_mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { - struct dummy * h = ud; - switch (type) { - case PTYPE_SYSTEM: { - // register name message - const struct remote_message *rmsg = msg; - assert (sz == sizeof(rmsg->destination)); - _update_name(h, rmsg->destination.name, rmsg->destination.handle); - return 0; - } - default: { - // remote message out - const struct remote_message *rmsg = msg; - if (rmsg->destination.handle == 0) { - _send_name(h, source , rmsg->destination.name, type, session, rmsg->message, rmsg->sz); - } else { - // local message - skynet_send(context, source, rmsg->destination.handle , type | PTYPE_TAG_DONTCOPY, session, (void *)rmsg->message, rmsg->sz); - } - return 0; - } - } -} - -int -dummy_init(struct dummy *d, struct skynet_context *ctx, const char * args) { - d->ctx = ctx; - skynet_harbor_start(ctx); - skynet_callback(ctx, d, _mainloop); - - return 0; -} diff --git a/service-src/service_harbor.c b/service-src/service_harbor.c index 926ed46f..f7a88f2b 100644 --- a/service-src/service_harbor.c +++ b/service-src/service_harbor.c @@ -2,6 +2,16 @@ #include "skynet_harbor.h" #include "skynet_socket.h" +/* + harbor listen the PTYPE_HARBOR (in text) + N name : update the global name + S fd id: connect to new harbor , we should send self_id to fd first , and then recv a id (check it), and at last send queue. + A fd id: accept new harbor , we should send self_id to fd , and then send queue. + + If the fd is disconnected, send message to slave in PTYPE_TEXT. D id + If we don't known a globalname, send message to slave in PTYPE_TEXT. Q name + */ + #include #include #include @@ -13,8 +23,22 @@ #define HASH_SIZE 4096 #define DEFAULT_QUEUE_SIZE 1024 +// 12 is sizeof(struct remote_message_header) +#define HEADER_COOKIE_LENGTH 12 + +/* + message type (8bits) is in destination high 8bits + harbor id (8bits) is also in that place , but remote message doesn't need harbor id. + */ +struct remote_message_header { + uint32_t source; + uint32_t destination; + uint32_t session; +}; + struct msg { - uint8_t * buffer; + struct remote_message_header header; + void * buffer; size_t size; }; @@ -37,35 +61,34 @@ struct hashmap { struct keyvalue *node[HASH_SIZE]; }; -/* - message type (8bits) is in destination high 8bits - harbor id (8bits) is also in that place , but remote message doesn't need harbor id. - */ -struct remote_message_header { - uint32_t source; - uint32_t destination; - uint32_t session; -}; +#define STATUS_WAIT 0 +#define STATUS_HANDSHAKE 1 +#define STATUS_HEADER 2 +#define STATUS_CONTENT 3 +#define STATUS_DOWN 4 -// 12 is sizeof(struct remote_message_header) -#define HEADER_COOKIE_LENGTH 12 +struct slave { + int fd; + struct msg_queue *queue; + int status; + int length; + int read; + uint8_t size[4]; + char * recv_buffer; +}; struct harbor { struct skynet_context *ctx; - char * local_addr; int id; + uint32_t slave; struct hashmap * map; - int master_fd; - char * master_addr; - int remote_fd[REMOTE_MAX]; - bool connected[REMOTE_MAX]; - char * remote_addr[REMOTE_MAX]; + struct slave s[REMOTE_MAX]; }; // hash table static void -_push_queue(struct msg_queue * queue, const void * buffer, size_t sz, struct remote_message_header * header) { +push_queue_msg(struct msg_queue * queue, struct msg * m) { // If there is only 1 free slot which is reserved to distinguish full/empty // of circular buffer, expand it. if (((queue->tail + 1) % queue->size) == queue->head) { @@ -81,16 +104,21 @@ _push_queue(struct msg_queue * queue, const void * buffer, size_t sz, struct rem queue->size *= 2; } struct msg * slot = &queue->data[queue->tail]; + *slot = *m; queue->tail = (queue->tail + 1) % queue->size; +} - slot->buffer = skynet_malloc(sz + sizeof(*header)); - memcpy(slot->buffer, buffer, sz); - memcpy(slot->buffer + sz, header, sizeof(*header)); - slot->size = sz + sizeof(*header); +static void +push_queue(struct msg_queue * queue, void * buffer, size_t sz, struct remote_message_header * header) { + struct msg m; + m.header = *header; + m.buffer = buffer; + m.size = sz; + push_queue_msg(queue, &m); } static struct msg * -_pop_queue(struct msg_queue * queue) { +pop_queue(struct msg_queue * queue) { if (queue->head == queue->tail) { return NULL; } @@ -100,7 +128,7 @@ _pop_queue(struct msg_queue * queue) { } static struct msg_queue * -_new_queue() { +new_queue() { struct msg_queue * queue = skynet_malloc(sizeof(*queue)); queue->size = DEFAULT_QUEUE_SIZE; queue->head = 0; @@ -111,20 +139,19 @@ _new_queue() { } static void -_release_queue(struct msg_queue *queue) { +release_queue(struct msg_queue *queue) { if (queue == NULL) return; - struct msg * m = _pop_queue(queue); - while (m) { + struct msg * m; + while ((m=pop_queue(queue)) != NULL) { skynet_free(m->buffer); - m = _pop_queue(queue); } skynet_free(queue->data); skynet_free(queue); } static struct keyvalue * -_hash_search(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { +hash_search(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { uint32_t *ptr = (uint32_t*) name; uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; struct keyvalue * node = hash->node[h % HASH_SIZE]; @@ -142,7 +169,7 @@ _hash_search(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { // Don't support erase name yet static struct void -_hash_erase(struct hashmap * hash, char name[GLOBALNAME_LENGTH) { +hash_erase(struct hashmap * hash, char name[GLOBALNAME_LENGTH) { uint32_t *ptr = name; uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; struct keyvalue ** ptr = &hash->node[h % HASH_SIZE]; @@ -160,7 +187,7 @@ _hash_erase(struct hashmap * hash, char name[GLOBALNAME_LENGTH) { */ static struct keyvalue * -_hash_insert(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { +hash_insert(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { uint32_t *ptr = (uint32_t *)name; uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; struct keyvalue ** pkv = &hash->node[h % HASH_SIZE]; @@ -176,20 +203,20 @@ _hash_insert(struct hashmap * hash, const char name[GLOBALNAME_LENGTH]) { } static struct hashmap * -_hash_new() { +hash_new() { struct hashmap * h = skynet_malloc(sizeof(struct hashmap)); memset(h,0,sizeof(*h)); return h; } static void -_hash_delete(struct hashmap *hash) { +hash_delete(struct hashmap *hash) { int i; for (i=0;inode[i]; while (node) { struct keyvalue * next = node->next; - _release_queue(node->queue); + release_queue(node->queue); skynet_free(node); node = next; } @@ -199,64 +226,49 @@ _hash_delete(struct hashmap *hash) { /////////////// +static void +close_harbor(struct harbor *h, int id) { + struct slave *s = &h->s[id]; + s->status = STATUS_DOWN; + if (s->fd) { + skynet_socket_close(h->ctx, s->fd); + } + if (s->queue) { + release_queue(s->queue); + s->queue = NULL; + } +} + +static void +report_harbor_down(struct harbor *h, int id) { + char down[64]; + int n = sprintf(down, "D %d",id); + + skynet_send(h->ctx, 0, h->slave, PTYPE_TEXT, 0, down, n); +} + struct harbor * harbor_create(void) { struct harbor * h = skynet_malloc(sizeof(*h)); - h->ctx = NULL; - h->id = 0; - h->master_fd = -1; - h->master_addr = NULL; - int i; - for (i=0;iremote_fd[i] = -1; - h->connected[i] = false; - h->remote_addr[i] = NULL; - } - h->map = _hash_new(); + memset(h,0,sizeof(*h)); + h->map = hash_new(); return h; } void harbor_release(struct harbor *h) { - struct skynet_context *ctx = h->ctx; - if (h->master_fd >= 0) { - skynet_socket_close(ctx, h->master_fd); - } - skynet_free(h->master_addr); - skynet_free(h->local_addr); int i; - for (i=0;iremote_fd[i] >= 0) { - skynet_socket_close(ctx, h->remote_fd[i]); - skynet_free(h->remote_addr[i]); + for (i=1;is[i]; + if (s->fd && s->status != STATUS_DOWN) { + close_harbor(h,i); + report_harbor_down(h,i); } } - _hash_delete(h->map); + hash_delete(h->map); skynet_free(h); } -static int -_connect_to(struct harbor *h, const char *ipaddress, bool blocking) { - char * port = strchr(ipaddress,':'); - if (port==NULL) { - return -1; - } - int sz = port - ipaddress; - char tmp[sz + 1]; - memcpy(tmp,ipaddress,sz); - tmp[sz] = '\0'; - - int portid = strtol(port+1, NULL,10); - - skynet_error(h->ctx, "Harbor(%d) connect to %s:%d", h->id, tmp, portid); - - if (blocking) { - return skynet_socket_block_connect(h->ctx, tmp, portid); - } else { - return skynet_socket_connect(h->ctx, tmp, portid); - } -} - static inline void to_bigendian(uint8_t *buffer, uint32_t n) { buffer[0] = (n >> 24) & 0xff; @@ -266,7 +278,7 @@ to_bigendian(uint8_t *buffer, uint32_t n) { } static inline void -_header_to_message(const struct remote_message_header * header, uint8_t * message) { +header_to_message(const struct remote_message_header * header, uint8_t * message) { to_bigendian(message , header->source); to_bigendian(message+4 , header->destination); to_bigendian(message+8 , header->session); @@ -283,112 +295,205 @@ from_bigendian(uint32_t n) { } static inline void -_message_to_header(const uint32_t *message, struct remote_message_header *header) { +message_to_header(const uint32_t *message, struct remote_message_header *header) { header->source = from_bigendian(message[0]); header->destination = from_bigendian(message[1]); header->session = from_bigendian(message[2]); } -static void -_send_package(struct skynet_context *ctx, int fd, const void * buffer, size_t sz) { - uint8_t * sendbuf = skynet_malloc(sz+4); - to_bigendian(sendbuf, sz); - memcpy(sendbuf+4, buffer, sz); +// socket package - if (skynet_socket_send(ctx, fd, sendbuf, sz+4)) { - skynet_error(ctx, "Send to %d error", fd); - } +static void +forward_local_messsage(struct harbor *h, void *msg, int sz) { + const char * cookie = msg; + cookie += sz - HEADER_COOKIE_LENGTH; + struct remote_message_header header; + message_to_header((const uint32_t *)cookie, &header); + + uint32_t destination = header.destination; + int type = (destination >> HANDLE_REMOTE_SHIFT) | PTYPE_TAG_DONTCOPY; + destination = (destination & HANDLE_MASK) | ((uint32_t)h->id << HANDLE_REMOTE_SHIFT); + + skynet_send(h->ctx, header.source, destination, type, (int)header.session, (void *)msg, sz-HEADER_COOKIE_LENGTH); } static void -_send_remote(struct skynet_context * ctx, int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) { +send_remote(struct skynet_context * ctx, int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) { uint32_t sz_header = sz+sizeof(*cookie); uint8_t * sendbuf = skynet_malloc(sz_header+4); to_bigendian(sendbuf, sz_header); memcpy(sendbuf+4, buffer, sz); - _header_to_message(cookie, sendbuf+4+sz); + header_to_message(cookie, sendbuf+4+sz); - if (skynet_socket_send(ctx, fd, sendbuf, sz_header+4)) { - skynet_error(ctx, "Remote send to %d error", fd); - } + // ignore send error, because if the connection is broken, the mainloop will recv a message. + skynet_socket_send(ctx, fd, sendbuf, sz_header+4); } static void -_update_remote_address(struct harbor *h, int harbor_id, const char * ipaddr) { - if (harbor_id == h->id) { - return; - } - assert(harbor_id > 0 && harbor_id< REMOTE_MAX); - struct skynet_context * context = h->ctx; - if (h->remote_fd[harbor_id] >=0) { - skynet_socket_close(context, h->remote_fd[harbor_id]); - skynet_free(h->remote_addr[harbor_id]); - h->remote_addr[harbor_id] = NULL; - } - h->remote_fd[harbor_id] = _connect_to(h, ipaddr, false); - h->connected[harbor_id] = false; -} - -static void -_dispatch_queue(struct harbor *h, struct msg_queue * queue, uint32_t handle, const char name[GLOBALNAME_LENGTH] ) { +dispatch_name_queue(struct harbor *h, struct keyvalue * node) { + struct msg_queue * queue = node->queue; + uint32_t handle = node->value; int harbor_id = handle >> HANDLE_REMOTE_SHIFT; assert(harbor_id != 0); struct skynet_context * context = h->ctx; - int fd = h->remote_fd[harbor_id]; - if (fd < 0) { - char tmp [GLOBALNAME_LENGTH+1]; - memcpy(tmp, name , GLOBALNAME_LENGTH); - tmp[GLOBALNAME_LENGTH] = '\0'; - skynet_error(context, "Drop message to %s (in harbor %d)",tmp,harbor_id); + struct slave *s = &h->s[harbor_id]; + int fd = s->fd; + if (fd == 0) { + if (s->status == STATUS_DOWN) { + char tmp [GLOBALNAME_LENGTH+1]; + memcpy(tmp, node->key, GLOBALNAME_LENGTH); + tmp[GLOBALNAME_LENGTH] = '\0'; + skynet_error(context, "Drop message to %s (in harbor %d)",tmp,harbor_id); + } else { + if (s->queue == NULL) { + s->queue = node->queue; + node->queue = NULL; + } else { + struct msg * m; + while ((m = pop_queue(queue))!=NULL) { + push_queue_msg(s->queue, m); + } + } + } return; } - struct msg * m = _pop_queue(queue); - while (m) { - struct remote_message_header cookie; - uint8_t *ptr = m->buffer + m->size - sizeof(cookie); - memcpy(&cookie, ptr, sizeof(cookie)); - cookie.destination |= (handle & HANDLE_MASK); - _header_to_message(&cookie, ptr); - _send_package(context, fd, m->buffer, m->size); - m = _pop_queue(queue); + struct msg * m; + while ((m = pop_queue(queue)) != NULL) { + m->header.destination |= (handle & HANDLE_MASK); + send_remote(context, fd, m->buffer, m->size, &m->header); } } static void -_update_remote_name(struct harbor *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { - struct keyvalue * node = _hash_search(h->map, name); +dispatch_queue(struct harbor *h, int id) { + struct slave *s = &h->s[id]; + int fd = s->fd; + assert(fd != 0); + + struct msg_queue *queue = s->queue; + if (queue == NULL) + return; + + struct msg * m; + while ((m = pop_queue(queue)) != NULL) { + send_remote(h->ctx, fd, m->buffer, m->size, &m->header); + } + release_queue(queue); + s->queue = NULL; +} + +static void +push_socket_data(struct harbor *h, const struct skynet_socket_message * message) { + assert(message->type == SKYNET_SOCKET_TYPE_DATA); + int fd = message->id; + int i; + int id = 0; + struct slave * s = NULL; + for (i=1;is[i].fd == fd) { + s = &h->s[i]; + id = i; + break; + } + } + if (s == NULL) { + skynet_free(message->buffer); + skynet_error(h->ctx, "Invalid socket fd (%d) data", fd); + return; + } + uint8_t * buffer = (uint8_t *)message->buffer; + int size = message->ud; + + for (;;) { + switch(s->status) { + case STATUS_HANDSHAKE: { + // check id + uint8_t remote_id = buffer[0]; + if (remote_id != id) { + skynet_error(h->ctx, "Invalid shakehand id (%d) from fd = %d , harbor = %d", id, fd, remote_id); + close_harbor(h,id); + return; + } + ++buffer; + --size; + s->status = STATUS_HEADER; + + dispatch_queue(h, id); + + if (size == 0) { + break; + } + // go though + } + case STATUS_HEADER: { + // big endian 4 bytes length, the first one must be 0. + int need = 4 - s->read; + if (size < need) { + memcpy(s->size + s->read, buffer, size); + s->read += size; + return; + } else { + memcpy(s->size + s->read, buffer, need); + buffer += need; + size -= need; + + if (s->size[0] != 0) { + skynet_error(h->ctx, "Message is too long from harbor %d", id); + close_harbor(h,id); + return; + } + s->length = s->size[1] << 16 | s->size[2] << 8 | s->size[3]; + s->read = 0; + s->recv_buffer = skynet_malloc(s->length); + s->status = STATUS_CONTENT; + if (size == 0) { + return; + } + } + } + // go though + case STATUS_CONTENT: { + int need = s->length - s->read; + if (size < need) { + memcpy(s->recv_buffer + s->read, buffer, size); + s->read += size; + return; + } + memcpy(s->recv_buffer + s->read, buffer, need); + forward_local_messsage(h, s->recv_buffer, s->length); + s->length = 0; + s->read = 0; + s->recv_buffer = NULL; + size -= need; + buffer += need; + s->status = STATUS_HEADER; + if (size == 0) + return; + break; + } + default: + return; + } + } +} + +static void +update_name(struct harbor *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { + struct keyvalue * node = hash_search(h->map, name); if (node == NULL) { - node = _hash_insert(h->map, name); + node = hash_insert(h->map, name); } node->value = handle; if (node->queue) { - _dispatch_queue(h, node->queue, handle, name); - _release_queue(node->queue); + dispatch_name_queue(h, node); + release_queue(node->queue); node->queue = NULL; } } -static void -_request_master(struct harbor *h, const char name[GLOBALNAME_LENGTH], size_t i, uint32_t handle) { - uint8_t buffer[4+i]; - to_bigendian(buffer, handle); - memcpy(buffer+4,name,i); - - _send_package(h->ctx, h->master_fd, buffer, 4+i); -} - -/* - update global name to master - - 2 bytes (size) - 4 bytes (handle) (handle == 0 for request) - n bytes string (name) - */ - static int -_remote_send_handle(struct harbor *h, uint32_t source, uint32_t destination, int type, int session, const char * msg, size_t sz) { +remote_send_handle(struct harbor *h, uint32_t source, uint32_t destination, int type, int session, const char * msg, size_t sz) { int harbor_id = destination >> HANDLE_REMOTE_SHIFT; - assert(harbor_id != 0); struct skynet_context * context = h->ctx; if (harbor_id == h->id) { // local message @@ -396,55 +501,114 @@ _remote_send_handle(struct harbor *h, uint32_t source, uint32_t destination, int return 1; } - int fd = h->remote_fd[harbor_id]; - if (fd >= 0 && h->connected[harbor_id]) { + struct slave * s = &h->s[harbor_id]; + if (s->fd == 0 || s->status == STATUS_HANDSHAKE) { + if (s->status == STATUS_DOWN) { + // throw an error return to source + // report the destination is dead + skynet_send(context, destination, source, PTYPE_ERROR, 0 , NULL, 0); + skynet_error(context, "Drop message to harbor %d from %x to %x (session = %d, msgsz = %d)",harbor_id, source, destination,session,(int)sz); + } else { + struct remote_message_header header; + header.source = source; + header.destination = type << HANDLE_REMOTE_SHIFT; + header.session = (uint32_t)session; + push_queue(s->queue, (void *)msg, sz, &header); + return 1; + } + } else { struct remote_message_header cookie; cookie.source = source; cookie.destination = (destination & HANDLE_MASK) | ((uint32_t)type << HANDLE_REMOTE_SHIFT); cookie.session = (uint32_t)session; - _send_remote(context, fd, msg,sz,&cookie); - } else { - // throw an error return to source - // report the destination is dead - skynet_send(context, destination, source, PTYPE_ERROR, 0 , NULL, 0); - skynet_error(context, "Drop message to harbor %d from %x to %x (session = %d, msgsz = %d)",harbor_id, source, destination,session,(int)sz); + send_remote(context, s->fd, msg,sz,&cookie); } + return 0; } -static void -_remote_register_name(struct harbor *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { - int i; - for (i=0;imap, name); +remote_send_name(struct harbor *h, uint32_t source, const char name[GLOBALNAME_LENGTH], int type, int session, const char * msg, size_t sz) { + struct keyvalue * node = hash_search(h->map, name); if (node == NULL) { - node = _hash_insert(h->map, name); + node = hash_insert(h->map, name); } if (node->value == 0) { if (node->queue == NULL) { - node->queue = _new_queue(); + node->queue = new_queue(); } struct remote_message_header header; header.source = source; header.destination = type << HANDLE_REMOTE_SHIFT; header.session = (uint32_t)session; - _push_queue(node->queue, msg, sz, &header); - // 0 for request - _remote_register_name(h, name, 0); + push_queue(node->queue, (void *)msg, sz, &header); + char query[2+GLOBALNAME_LENGTH+1] = "Q "; + query[2+GLOBALNAME_LENGTH] = 0; + memcpy(query+2, name, GLOBALNAME_LENGTH); + skynet_send(h->ctx, 0, h->slave, PTYPE_TEXT, 0, query, strlen(query)); return 1; } else { - return _remote_send_handle(h, source, node->value, type, session, msg, sz); + return remote_send_handle(h, source, node->value, type, session, msg, sz); + } +} + +static void +handshake(struct harbor *h, int id) { + struct slave *s = &h->s[id]; + uint8_t * handshake = skynet_malloc(1); + handshake[0] = (uint8_t)h->id; + skynet_socket_send(h->ctx, s->fd, handshake, 1); +} + +static void +harbor_command(struct harbor * h, const char * msg, size_t sz, int session, uint32_t source) { + const char * name = msg + 2; + int s = (int)sz; + s -= 2; + switch(msg[0]) { + case 'N' : { + if (s <=0 || s>= GLOBALNAME_LENGTH) { + skynet_error(h->ctx, "Invalid global name %s", name); + return; + } + struct remote_name rn; + memset(&rn, 0, sizeof(rn)); + memcpy(rn.name, name, s); + rn.handle = source; + update_name(h, rn.name, rn.handle); + break; + } + case 'S' : + case 'A' : { + char buffer[s+1]; + memcpy(buffer, name, s); + buffer[s] = 0; + int fd=0, id=0; + sscanf(buffer, "%d %d",&fd,&id); + if (fd == 0 || id <= 0 || id>=REMOTE_MAX) { + skynet_error(h->ctx, "Invalid command %c %s", msg[0], buffer); + return; + } + struct slave * slave = &h->s[id]; + if (slave->fd != 0) { + skynet_error(h->ctx, "Harbor %d alreay exist", id); + return; + } + slave->fd = fd; + + skynet_socket_start(h->ctx, fd); + handshake(h, id); + if (msg[0] == 'S') { + slave->status = STATUS_HANDSHAKE; + } else { + slave->status = STATUS_HEADER; + dispatch_queue(h,id); + } + break; + } + default: + skynet_error(h->ctx, "Unknown command %s", msg); + return; } } @@ -452,105 +616,57 @@ static int harbor_id(struct harbor *h, int fd) { int i; for (i=1;iremote_fd[i] == fd) + struct slave *s = &h->s[i]; + if (s->fd == fd) { return i; + } } return 0; } -static void -close_harbor(struct harbor *h, int fd) { - int id = harbor_id(h,fd); - if (id == 0) - return; - skynet_error(h->ctx, "Harbor %d closed",id); - skynet_socket_close(h->ctx, fd); - h->remote_fd[id] = -1; - h->connected[id] = false; -} - -static void -open_harbor(struct harbor *h, int fd) { - int id = harbor_id(h,fd); - if (id == 0) - return; - assert(h->connected[id] == false); - h->connected[id] = true; -} - static int -_mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { +mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { struct harbor * h = ud; switch (type) { case PTYPE_SOCKET: { const struct skynet_socket_message * message = msg; switch(message->type) { case SKYNET_SOCKET_TYPE_DATA: + push_socket_data(h, message); skynet_free(message->buffer); - skynet_error(context, "recv invalid socket message (size=%d)", message->ud); - break; - case SKYNET_SOCKET_TYPE_ACCEPT: - skynet_error(context, "recv invalid socket accept message"); break; case SKYNET_SOCKET_TYPE_ERROR: - case SKYNET_SOCKET_TYPE_CLOSE: - close_harbor(h, message->id); + case SKYNET_SOCKET_TYPE_CLOSE: { + int id = harbor_id(h, message->id); + if (id) { + report_harbor_down(h,id); + } else { + skynet_error(context, "Unkown fd (%d) closed", message->id); + } break; + } case SKYNET_SOCKET_TYPE_CONNECT: - open_harbor(h, message->id); + // fd forward to this service + break; + default: + skynet_error(context, "recv invalid socket message type %d", type); break; } return 0; } case PTYPE_HARBOR: { - // remote message in - const char * cookie = msg; - cookie += sz - HEADER_COOKIE_LENGTH; - struct remote_message_header header; - _message_to_header((const uint32_t *)cookie, &header); - if (header.source == 0) { - if (header.destination < REMOTE_MAX) { - // 1 byte harbor id (0~255) - // update remote harbor address - char ip [sz - HEADER_COOKIE_LENGTH + 1]; - memcpy(ip, msg, sz-HEADER_COOKIE_LENGTH); - ip[sz-HEADER_COOKIE_LENGTH] = '\0'; - _update_remote_address(h, header.destination, ip); - } else { - // update global name - if (sz - HEADER_COOKIE_LENGTH > GLOBALNAME_LENGTH) { - char name[sz-HEADER_COOKIE_LENGTH+1]; - memcpy(name, msg, sz-HEADER_COOKIE_LENGTH); - name[sz-HEADER_COOKIE_LENGTH] = '\0'; - skynet_error(context, "Global name is too long %s", name); - } - _update_remote_name(h, msg, header.destination); - } - } else { - uint32_t destination = header.destination; - int type = (destination >> HANDLE_REMOTE_SHIFT) | PTYPE_TAG_DONTCOPY; - destination = (destination & HANDLE_MASK) | ((uint32_t)h->id << HANDLE_REMOTE_SHIFT); - skynet_send(context, header.source, destination, type, (int)header.session, (void *)msg, sz-HEADER_COOKIE_LENGTH); - return 1; - } - return 0; - } - case PTYPE_SYSTEM: { - // register name message - const struct remote_message *rmsg = msg; - assert (sz == sizeof(rmsg->destination)); - _remote_register_name(h, rmsg->destination.name, rmsg->destination.handle); + harbor_command(h, msg,sz,session,source); return 0; } default: { // remote message out const struct remote_message *rmsg = msg; if (rmsg->destination.handle == 0) { - if (_remote_send_name(h, source , rmsg->destination.name, type, session, rmsg->message, rmsg->sz)) { + if (remote_send_name(h, source , rmsg->destination.name, type, session, rmsg->message, rmsg->sz)) { return 0; } } else { - if (_remote_send_handle(h, source , rmsg->destination.handle, type, session, rmsg->message, rmsg->sz)) { + if (remote_send_handle(h, source , rmsg->destination.handle, type, session, rmsg->message, rmsg->sz)) { return 0; } } @@ -560,48 +676,19 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin } } -static void -_launch_gate(struct skynet_context * ctx, const char * local_addr) { - char tmp[128]; - sprintf(tmp,"gate L ! %s %d %d 0",local_addr, PTYPE_HARBOR, REMOTE_MAX); - const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp); - if (gate_addr == NULL) { - fprintf(stderr, "Harbor : launch gate failed\n"); - exit(1); - } - uint32_t gate = strtoul(gate_addr+1 , NULL, 16); - if (gate == 0) { - fprintf(stderr, "Harbor : launch gate invalid %s", gate_addr); - exit(1); - } - const char * self_addr = skynet_command(ctx, "REG", NULL); - int n = sprintf(tmp,"broker %s",self_addr); - skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, tmp, n); - skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, "start", 5); -} - int harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) { h->ctx = ctx; - int sz = strlen(args)+1; - char master_addr[sz]; - char local_addr[sz]; int harbor_id = 0; - sscanf(args,"%s %s %d",master_addr, local_addr, &harbor_id); - h->master_addr = skynet_strdup(master_addr); - h->id = harbor_id; - h->master_fd = _connect_to(h, master_addr, true); - if (h->master_fd == -1) { - fprintf(stderr, "Harbor: Connect to master failed\n"); - exit(1); + uint32_t slave = 0; + sscanf(args,"%d %u", &harbor_id, &slave); + if (slave == 0) { + return 1; } + h->id = harbor_id; + h->slave = slave; + skynet_callback(ctx, h, mainloop); skynet_harbor_start(ctx); - h->local_addr = skynet_strdup(local_addr); - - _launch_gate(ctx, local_addr); - skynet_callback(ctx, h, _mainloop); - _request_master(h, local_addr, strlen(local_addr), harbor_id); - return 0; } diff --git a/service-src/service_master.c b/service-src/service_master.c deleted file mode 100644 index b5d8a560..00000000 --- a/service-src/service_master.c +++ /dev/null @@ -1,313 +0,0 @@ -#include "skynet.h" -#include "skynet_harbor.h" -#include "skynet_socket.h" - -#include -#include -#include -#include -#include -#include - -#define HASH_SIZE 4096 - -struct name { - struct name * next; - char key[GLOBALNAME_LENGTH]; - uint32_t hash; - uint32_t value; -}; - -struct namemap { - struct name *node[HASH_SIZE]; -}; - -struct master { - struct skynet_context *ctx; - int remote_fd[REMOTE_MAX]; - bool connected[REMOTE_MAX]; - char * remote_addr[REMOTE_MAX]; - struct namemap map; -}; - -struct master * -master_create() { - struct master *m = skynet_malloc(sizeof(*m)); - int i; - for (i=0;iremote_fd[i] = -1; - m->remote_addr[i] = NULL; - m->connected[i] = false; - } - memset(&m->map, 0, sizeof(m->map)); - return m; -} - -void -master_release(struct master * m) { - int i; - struct skynet_context *ctx = m->ctx; - for (i=0;iremote_fd[i]; - if (fd >= 0) { - assert(ctx); - skynet_socket_close(ctx, fd); - } - skynet_free(m->remote_addr[i]); - } - for (i=0;imap.node[i]; - while (node) { - struct name * next = node->next; - skynet_free(node); - node = next; - } - } - skynet_free(m); -} - -static struct name * -_search_name(struct master *m, char name[GLOBALNAME_LENGTH]) { - uint32_t *ptr = (uint32_t *) name; - uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; - struct name * node = m->map.node[h % HASH_SIZE]; - while (node) { - if (node->hash == h && strncmp(node->key, name, GLOBALNAME_LENGTH) == 0) { - return node; - } - node = node->next; - } - return NULL; -} - -static struct name * -_insert_name(struct master *m, char name[GLOBALNAME_LENGTH]) { - uint32_t *ptr = (uint32_t *)name; - uint32_t h = ptr[0] ^ ptr[1] ^ ptr[2] ^ ptr[3]; - struct name **pname = &m->map.node[h % HASH_SIZE]; - struct name * node = skynet_malloc(sizeof(*node)); - memcpy(node->key, name, GLOBALNAME_LENGTH); - node->next = *pname; - node->hash = h; - node->value = 0; - *pname = node; - return node; -} - -static void -_copy_name(char *name, const char * buffer, size_t sz) { - if (sz < GLOBALNAME_LENGTH) { - memcpy(name, buffer, sz); - memset(name+sz, 0 , GLOBALNAME_LENGTH - sz); - } else { - memcpy(name, buffer, GLOBALNAME_LENGTH); - } -} - -static void -_connect_to(struct master *m, int id) { - assert(m->connected[id] == false); - struct skynet_context * ctx = m->ctx; - const char *ipaddress = m->remote_addr[id]; - char * portstr = strchr(ipaddress,':'); - if (portstr==NULL) { - skynet_error(ctx, "Harbor %d : address invalid (%s)",id, ipaddress); - return; - } - int sz = portstr - ipaddress; - char tmp[sz + 1]; - memcpy(tmp,ipaddress,sz); - tmp[sz] = '\0'; - int port = strtol(portstr+1,NULL,10); - skynet_error(ctx, "Master connect to harbor(%d) %s:%d", id, tmp, port); - m->remote_fd[id] = skynet_socket_connect(ctx, tmp, port); -} - -static inline void -to_bigendian(uint8_t *buffer, uint32_t n) { - buffer[0] = (n >> 24) & 0xff; - buffer[1] = (n >> 16) & 0xff; - buffer[2] = (n >> 8) & 0xff; - buffer[3] = n & 0xff; -} - -static void -_send_to(struct master *m, int id, const void * buf, int sz, uint32_t handle) { - uint8_t * buffer= (uint8_t *)skynet_malloc(4 + sz + 12); - to_bigendian(buffer, sz+12); - memcpy(buffer+4, buf, sz); - to_bigendian(buffer+4+sz, 0); - to_bigendian(buffer+4+sz+4, handle); - to_bigendian(buffer+4+sz+8, 0); - - sz += 4 + 12; - - if (skynet_socket_send(m->ctx, m->remote_fd[id], buffer, sz)) { - skynet_error(m->ctx, "Harbor %d : send error", id); - } -} - -static void -_broadcast(struct master *m, const char *name, size_t sz, uint32_t handle) { - int i; - for (i=1;iremote_fd[i]; - if (fd < 0 || m->connected[i]==false) - continue; - _send_to(m, i , name, sz, handle); - } -} - -static void -_request_name(struct master *m, const char * buffer, size_t sz) { - char name[GLOBALNAME_LENGTH]; - _copy_name(name, buffer, sz); - struct name * n = _search_name(m, name); - if (n == NULL) { - return; - } - _broadcast(m, name, GLOBALNAME_LENGTH, n->value); -} - -static void -_update_name(struct master *m, uint32_t handle, const char * buffer, size_t sz) { - char name[GLOBALNAME_LENGTH]; - _copy_name(name, buffer, sz); - struct name * n = _search_name(m, name); - if (n==NULL) { - n = _insert_name(m,name); - } - n->value = handle; - _broadcast(m,name,GLOBALNAME_LENGTH, handle); -} - -static void -close_harbor(struct master *m, int harbor_id) { - if (m->connected[harbor_id]) { - struct skynet_context * context = m->ctx; - skynet_socket_close(context, m->remote_fd[harbor_id]); - m->remote_fd[harbor_id] = -1; - m->connected[harbor_id] = false; - } -} - -static void -_update_address(struct master *m, int harbor_id, const char * buffer, size_t sz) { - if (m->remote_fd[harbor_id] >= 0) { - close_harbor(m, harbor_id); - } - skynet_free(m->remote_addr[harbor_id]); - char * addr = skynet_malloc(sz+1); - memcpy(addr, buffer, sz); - addr[sz] = '\0'; - m->remote_addr[harbor_id] = addr; - _connect_to(m, harbor_id); -} - -static int -socket_id(struct master *m, int id) { - int i; - for (i=1;iremote_fd[i] == id) - return i; - } - return 0; -} - -static void -on_connected(struct master *m, int id) { - _broadcast(m, m->remote_addr[id], strlen(m->remote_addr[id]), id); - m->connected[id] = true; - int i; - for (i=1;iremote_addr[i]; - if (addr == NULL || m->connected[i] == false) { - continue; - } - _send_to(m, id , addr, strlen(addr), i); - } -} - -static void -dispatch_socket(struct master *m, const struct skynet_socket_message *msg, int sz) { - int id = socket_id(m, msg->id); - switch(msg->type) { - case SKYNET_SOCKET_TYPE_CONNECT: - assert(id); - on_connected(m, id); - break; - case SKYNET_SOCKET_TYPE_ERROR: - skynet_error(m->ctx, "socket error on harbor %d", id); - // go though, close socket - case SKYNET_SOCKET_TYPE_CLOSE: - close_harbor(m, id); - break; - default: - skynet_error(m->ctx, "Invalid socket message type %d", msg->type); - break; - } -} - - -/* - update global name to master - - 4 bytes (handle) (handle == 0 for request) - n bytes string (name) - */ - -static int -_mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { - if (type == PTYPE_SOCKET) { - dispatch_socket(ud, msg, (int)sz); - return 0; - } - if (type != PTYPE_HARBOR) { - skynet_error(context, "None harbor message recv from %x (type = %d)", source, type); - return 0; - } - assert(sz >= 4); - struct master *m = ud; - const uint8_t *handlen = msg; - uint32_t handle = handlen[0]<<24 | handlen[1]<<16 | handlen[2]<<8 | handlen[3]; - sz -= 4; - const char * name = msg; - name += 4; - - if (handle == 0) { - _request_name(m , name, sz); - } else if (handle < REMOTE_MAX) { - _update_address(m , handle, name, sz); - } else { - _update_name(m , handle, name, sz); - } - - return 0; -} - -int -master_init(struct master *m, struct skynet_context *ctx, const char * args) { - char tmp[strlen(args) + 32]; - sprintf(tmp,"gate L ! %s %d %d 0",args,PTYPE_HARBOR,REMOTE_MAX); - const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp); - if (gate_addr == NULL) { - skynet_error(ctx, "Master : launch gate failed"); - return 1; - } - uint32_t gate = strtoul(gate_addr+1, NULL, 16); - if (gate == 0) { - skynet_error(ctx, "Master : launch gate invalid %s", gate_addr); - return 1; - } - const char * self_addr = skynet_command(ctx, "REG", NULL); - int n = sprintf(tmp,"broker %s",self_addr); - skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, tmp, n); - skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, "start", 5); - - skynet_callback(ctx, m, _mainloop); - - m->ctx = ctx; - return 0; -} diff --git a/service/bootstrap.lua b/service/bootstrap.lua index 7c92b419..11303577 100644 --- a/service/bootstrap.lua +++ b/service/bootstrap.lua @@ -1,30 +1,38 @@ local skynet = require "skynet" +local harbor = require "skynet.harbor" skynet.start(function() - assert(skynet.launch("logger", skynet.getenv "logger")) - local standalone = skynet.getenv "standalone" + + local launcher = assert(skynet.launch("snlua","launcher")) + skynet.name(".launcher", launcher) + local harbor_id = tonumber(skynet.getenv "harbor") if harbor_id == 0 then assert(standalone == nil) standalone = true skynet.setenv("standalone", "true") - assert(skynet.launch("dummy")) - else - local master_addr = skynet.getenv "master" + local slave = skynet.newservice "cdummy" + if slave == nil then + skynet.abort() + end + skynet.name(".slave", slave) + + else if standalone then - assert(skynet.launch("master", master_addr)) + if not skynet.newservice "cmaster" then + skynet.abort() + end end - local local_addr = skynet.getenv "address" - - assert(skynet.launch("harbor",master_addr, local_addr, harbor_id)) + local slave = skynet.newservice "cslave" + if slave == nil then + skynet.abort() + end + skynet.name(".slave", slave) end - local launcher = assert(skynet.launch("snlua","launcher")) - skynet.name(".launcher", launcher) - if standalone then local datacenter = assert(skynet.newservice "datacenterd") skynet.name("DATACENTER", datacenter) diff --git a/service/cdummy.lua b/service/cdummy.lua new file mode 100644 index 00000000..fddb2c01 --- /dev/null +++ b/service/cdummy.lua @@ -0,0 +1,47 @@ +local skynet = require "skynet" + +local globalname = {} +local harbor = {} + +skynet.register_protocol { + name = "harbor", + id = skynet.PTYPE_HARBOR, + pack = function(...) return ... end, + unpack = skynet.tostring, +} + +skynet.register_protocol { + name = "text", + id = skynet.PTYPE_TEXT, + pack = function(...) return ... end, + unpack = skynet.tostring, +} + +function harbor.REGISTER(name, handle) + assert(globalname[name] == nil) + globalname[name] = handle + skynet.redirect(harbor_service, handle, "harbor", 0, "N " .. name) +end + +function harbor.LINK(id) + skynet.ret() +end + +function harbor.CONNECT(id) + skynet.error("Can't connect to other harbor in single node mode") +end + +skynet.start(function() + local harbor_id = tonumber(skynet.getenv "harbor") + assert(harbor_id == 0) + + skynet.dispatch("lua", function (session,source,command,...) + local f = assert(harbor[command]) + f(...) + end) + skynet.dispatch("text", function(session,source,command) + -- ignore all the command + end) + + harbor_service = assert(skynet.launch("harbor", harbor_id, skynet.self())) +end) diff --git a/service/clusterd.lua b/service/clusterd.lua index 563869ac..405e28ab 100644 --- a/service/clusterd.lua +++ b/service/clusterd.lua @@ -32,6 +32,9 @@ local node_channel = setmetatable({}, { __index = open_channel }) function command.listen(source, addr, port) local gate = skynet.newservice("gate") + if port == nil then + addr, port = string.match(node_address[addr], "([^:]+):(.*)$") + end skynet.call(gate, "lua", "open", { address = addr, port = port }) skynet.ret(skynet.pack(nil)) end diff --git a/service/cmaster.lua b/service/cmaster.lua new file mode 100644 index 00000000..59bc61f1 --- /dev/null +++ b/service/cmaster.lua @@ -0,0 +1,125 @@ +local skynet = require "skynet" +local socket = require "socket" + +--[[ + master manage data : + 1. all the slaves address : id -> ipaddr:port + 2. all the global names : name -> address + + master hold connections from slaves . + + protocol slave->master : + package size 1 byte + type 1 byte : + 'H' : HANDSHAKE, report slave id, and address. + 'R' : REGISTER name address + 'Q' : QUERY name + + + protocol master->slave: + package size 1 byte + type 1 byte : + 'W' : WAIT n + 'C' : CONNECT slave_id slave_address + 'N' : NAME globalname address + 'D' : DISCONNECT slave_id +]] + +local slave_node = {} +local global_name = {} + +local function read_package(fd) + local sz = socket.read(fd, 1) + assert(sz, "closed") + sz = string.byte(sz) + local content = socket.read(fd, sz) + return skynet.unpack(content) +end + +local function pack_package(...) + local message = skynet.packstring(...) + local size = #message + assert(size <= 255 , "too long") + return string.char(size) .. message +end + +local function report_slave(fd, slave_id, slave_addr) + local message = pack_package("C", slave_id, slave_addr) + local n = 0 + for k,v in pairs(slave_node) do + if v.fd ~= 0 then + socket.write(v.fd, message) + n = n + 1 + end + end + socket.write(fd, pack_package("W", n)) +end + +local function handshake(fd) + local t, slave_id, slave_addr = read_package(fd) + assert(t=='H', "Invalid handshake type " .. t) + assert(slave_id ~= 0 , "Invalid slave id 0") + if slave_node[slave_id] then + error(string.format("Slave %d already register on %s", slave_id, slave_node[slave_id].addr)) + end + report_slave(fd, slave_id, slave_addr) + slave_node[slave_id] = { + fd = fd, + id = slave_id, + addr = slave_addr, + } + return slave_id , slave_addr +end + +local function dispatch_slave(fd) + local t, name, address = read_package(fd) + if t == 'R' then + -- register name + assert(type(address)=="number", "Invalid request") + if not global_name[name] then + global_name[name] = address + end + local message = pack_package("N", name, address) + for k,v in pairs(slave_node) do + socket.write(v.fd, message) + end + elseif t == 'Q' then + -- query name + local address = global_name[name] + if address then + socket.write(fd, pack_package("N", name, address)) + end + else + skynet.error("Invalid slave message type " .. t) + end +end + +local function monitor_slave(slave_id, slave_address) + local fd = slave_node[slave_id].fd + skynet.error(string.format("Harbor %d (fd=%d) report %s", slave_id, fd, slave_address)) + while pcall(dispatch_slave, fd) do end + skynet.error("slave " ..slave_id .. " is down") + local message = pack_package("D", slave_id) + slave_node[slave_id].fd = 0 + for k,v in pairs(slave_node) do + socket.write(v.fd, message) + end + socket.close(fd) +end + +skynet.start(function() + local master_addr = skynet.getenv "standalone" + skynet.error("master listen socket " .. tostring(master_addr)) + local fd = socket.listen(master_addr) + socket.start(fd , function(id, addr) + skynet.error("connect from " .. addr .. " " .. id) + socket.start(id) + local ok, slave, slave_addr = pcall(handshake, id) + if ok then + skynet.fork(monitor_slave, slave, slave_addr) + else + skynet.error(string.format("disconnect fd = %d, error = %s", id, slave)) + socket.close(id) + end + end) +end) diff --git a/service/cslave.lua b/service/cslave.lua new file mode 100644 index 00000000..6f005b2f --- /dev/null +++ b/service/cslave.lua @@ -0,0 +1,222 @@ +local skynet = require "skynet" +local socket = require "socket" + +local slaves = {} +local connect_queue = {} +local globalname = {} +local harbor = {} +local harbor_service +local monitor = {} + +local function read_package(fd) + local sz = socket.read(fd, 1) + assert(sz, "closed") + sz = string.byte(sz) + local content = socket.read(fd, sz) + return skynet.unpack(content) +end + +local function pack_package(...) + local message = skynet.packstring(...) + local size = #message + assert(size <= 255 , "too long") + return string.char(size) .. message +end + +local function monitor_clear(id) + local v = monitor[id] + if v then + monitor[id] = nil + for _, v in ipairs(v) do + skynet.redirect(v.address, 0, "response", v.session, "") + end + end +end + +local function connect_slave(slave_id, address) + local ok, err = pcall(function() + if slaves[slave_id] == nil then + local fd = assert(socket.open(address), "Can't connect to "..address) + skynet.error(string.format("Connect to harbor %d (fd=%d), %s", slave_id, fd, address)) + slaves[slave_id] = fd + monitor_clear(slave_id) + socket.abandon(fd) + skynet.send(harbor_service, "harbor", string.format("S %d %d",fd,slave_id)) + end + end) + if not ok then + skynet.error(err) + end +end + +local function ready() + local queue = connect_queue + connect_queue = nil + for k,v in pairs(queue) do + connect_slave(k,v) + end + for name,address in pairs(globalname) do + skynet.redirect(harbor_service, address, "harbor", "N " .. name) + end +end + +local function monitor_master(master_fd) + while true do + local ok, t, id_name, address = pcall(read_package,master_fd) + if ok then + if t == 'C' then + if connect_queue then + connect_queue[id_name] = address + else + connect_slave(id_name, address) + end + elseif t == 'N' then + globalname[id_name] = address + if connect_queue == nil then + skynet.redirect(harbor_service, address, "harbor", 0, "N " .. id_name) + end + elseif t == 'D' then + local fd = slaves[id_name] + slaves[id_name] = false + if fd then + socket.close(fd) + end + end + else + skynet.error("Master disconnect") + socket.close(master_fd) + break + end + end +end + +local function accept_slave(fd) + socket.start(fd) + local id = socket.read(fd, 1) + if not id then + skynet.error(string.format("Connection (fd =%d) closed", fd)) + socket.close(fd) + return + end + id = string.byte(id) + if slaves[id] ~= nil then + skynet.error(string.format("Slave %d exist (fd =%d)", id, fd)) + socket.close(fd) + return + end + slaves[id] = fd + monitor_clear(id) + socket.abandon(fd) + skynet.error(string.format("Harbor %d connected (fd = %d)", id, fd)) + skynet.send(harbor_service, "harbor", string.format("A %d %d", fd, id)) +end + +skynet.register_protocol { + name = "harbor", + id = skynet.PTYPE_HARBOR, + pack = function(...) return ... end, + unpack = skynet.tostring, +} + +skynet.register_protocol { + name = "text", + id = skynet.PTYPE_TEXT, + pack = function(...) return ... end, + unpack = skynet.tostring, +} + +local function monitor_harbor(master_fd) + return function(session, source, command) + local t = string.sub(command, 1, 1) + local arg = string.sub(command, 3) + if t == "Q" then + -- query name + if globalname[arg] then + skynet.redirect(harbor_service, globalname[arg], "harbor", 0, "N " .. arg) + else + socket.write(master_fd, pack_package("Q", arg)) + end + elseif t == "D" then + -- harbor down + local id = tonumber(arg) + if slaves[id] then + monitor_clear(id) + end + slaves[id] = false + else + skynet.error("Unknown command ", command) + end + end +end + +function harbor.REGISTER(_,_, fd, name, handle) + assert(globalname[name] == nil) + globalname[name] = handle + socket.write(fd, pack_package("R", name, handle)) + skynet.redirect(harbor_service, handle, "harbor", 0, "N " .. name) +end + +function harbor.LINK(session, source, fd, id) + if slaves[id] then + if monitor[id] == nil then + monitor[id] = {} + end + table.insert(monitor[id], { address = source, session = session }) + else + skynet.ret() + end +end + +function harbor.CONNECT(session, source, fd, id) + if not slaves[id] then + if monitor[id] == nil then + monitor[id] = {} + end + table.insert(monitor[id], { address = source, session = session }) + else + skynet.ret() + end +end + +skynet.start(function() + local master_addr = skynet.getenv "master" + local harbor_id = tonumber(skynet.getenv "harbor") + local slave_address = assert(skynet.getenv "address") + local slave_fd = socket.listen(slave_address) + skynet.error("slave connect to master " .. tostring(master_addr)) + local master_fd = socket.open(master_addr) + + skynet.dispatch("lua", function (session,source,command,...) + local f = assert(harbor[command]) + f(session, source, master_fd, ...) + end) + skynet.dispatch("text", monitor_harbor(master_fd)) + + harbor_service = assert(skynet.launch("harbor", harbor_id, skynet.self())) + + local hs_message = pack_package("H", harbor_id, slave_address) + socket.write(master_fd, hs_message) + local t, n = read_package(master_fd) + assert(t == "W" and type(n) == "number", "slave shakehand failed") + skynet.error(string.format("Waiting for %d harbors", n)) + skynet.fork(monitor_master, master_fd) + if n > 0 then + local co = coroutine.running() + socket.start(slave_fd, function(fd, addr) + skynet.error(string.format("New connection (fd = %d, %s)",fd, addr)) + if pcall(accept_slave,fd) then + local s = 0 + for k,v in pairs(slaves) do + s = s + 1 + end + if s >= n then + skynet.wakeup(co) + end + end + end) + skynet.wait() + end + socket.close(slave_fd) + skynet.error("Shakehand ready") + skynet.fork(ready) +end) diff --git a/service/debug_console.lua b/service/debug_console.lua index d49dd809..b15c3889 100644 --- a/service/debug_console.lua +++ b/service/debug_console.lua @@ -88,7 +88,7 @@ end skynet.start(function() local listen_socket = socket.listen ("127.0.0.1", port) - print("Start debug console at 127.0.0.1",port) + skynet.error("Start debug console at 127.0.0.1 " .. port) socket.start(listen_socket , function(id, addr) local function print(...) local t = { ... } diff --git a/skynet-src/skynet_harbor.c b/skynet-src/skynet_harbor.c index 160b1994..48814c83 100644 --- a/skynet-src/skynet_harbor.c +++ b/skynet-src/skynet_harbor.c @@ -17,23 +17,6 @@ skynet_harbor_send(struct remote_message *rmsg, uint32_t source, int session) { skynet_context_send(REMOTE, rmsg, sizeof(*rmsg) , source, type , session); } -void -skynet_harbor_register(struct remote_name *rname) { - if (REMOTE == NULL) - return; - int i; - int number = 1; - for (i=0;iname[i]; - if (!(c >= '0' && c <='9')) { - number = 0; - break; - } - } - assert(number == 0); - skynet_context_send(REMOTE, rname, sizeof(*rname), 0, PTYPE_SYSTEM , 0); -} - int skynet_harbor_message_isremote(uint32_t handle) { assert(HARBOR != ~0); diff --git a/skynet-src/skynet_harbor.h b/skynet-src/skynet_harbor.h index 0116a91a..b699f625 100644 --- a/skynet-src/skynet_harbor.h +++ b/skynet-src/skynet_harbor.h @@ -23,7 +23,6 @@ struct remote_message { }; 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_init(int harbor); void skynet_harbor_start(void * ctx); diff --git a/skynet-src/skynet_imp.h b/skynet-src/skynet_imp.h index 11b12ee1..c5a144a4 100644 --- a/skynet-src/skynet_imp.h +++ b/skynet-src/skynet_imp.h @@ -7,6 +7,7 @@ struct skynet_config { const char * daemon; const char * module_path; const char * bootstrap; + const char * logger; }; #define THREAD_WORKER 0 diff --git a/skynet-src/skynet_main.c b/skynet-src/skynet_main.c index 8596c778..9c38fd25 100644 --- a/skynet-src/skynet_main.c +++ b/skynet-src/skynet_main.c @@ -115,6 +115,7 @@ main(int argc, char *argv[]) { config.harbor = optint("harbor", 1); config.bootstrap = optstring("bootstrap","snlua bootstrap"); config.daemon = optstring("daemon", NULL); + config.logger = optstring("logger", NULL); lua_close(L); diff --git a/skynet-src/skynet_server.c b/skynet-src/skynet_server.c index 98fcfcd5..017d7206 100644 --- a/skynet-src/skynet_server.c +++ b/skynet-src/skynet_server.c @@ -234,6 +234,16 @@ _dispatch_message(struct skynet_context *ctx, struct skynet_message *msg) { CHECKCALLING_END(ctx) } +void +skynet_context_dispatchall(struct skynet_context * ctx) { + // for skynet_error + struct skynet_message msg; + struct message_queue *q = ctx->queue; + while (!skynet_mq_pop(q,&msg)) { + _dispatch_message(ctx, &msg); + } +} + struct message_queue * skynet_context_message_dispatch(struct skynet_monitor *sm, struct message_queue *q) { if (q == NULL) { @@ -342,11 +352,7 @@ cmd_reg(struct skynet_context * context, const char * param) { } else if (param[0] == '.') { return skynet_handle_namehandle(context->handle, param + 1); } else { - assert(context->handle!=0); - struct remote_name *rname = skynet_malloc(sizeof(*rname)); - copy_name(rname->name, param); - rname->handle = context->handle; - skynet_harbor_register(rname); + skynet_error(context, "Can't register global name %s in C", param); return NULL; } } @@ -377,10 +383,7 @@ cmd_name(struct skynet_context * context, const char * param) { if (name[0] == '.') { return skynet_handle_namehandle(handle_id, name + 1); } else { - struct remote_name *rname = skynet_malloc(sizeof(*rname)); - copy_name(rname->name, name); - rname->handle = handle_id; - skynet_harbor_register(rname); + skynet_error(context, "Can't set global name %s in C", name); } return NULL; } @@ -559,12 +562,16 @@ _filter_args(struct skynet_context * context, int type, int *session, void ** da *data = msg; } - assert((*sz & HANDLE_MASK) == *sz); *sz |= type << HANDLE_REMOTE_SHIFT; } int skynet_send(struct skynet_context * context, uint32_t source, uint32_t destination , int type, int session, void * data, size_t sz) { + if ((sz & HANDLE_MASK) != sz) { + skynet_error(context, "The message to %x is too large (sz = %lu)", destination, sz); + skynet_free(data); + return -1; + } _filter_args(context, type, &session, (void **)&data, &sz); if (source == 0) { diff --git a/skynet-src/skynet_server.h b/skynet-src/skynet_server.h index 51b8d1b2..be4e301f 100644 --- a/skynet-src/skynet_server.h +++ b/skynet-src/skynet_server.h @@ -18,6 +18,7 @@ void skynet_context_send(struct skynet_context * context, void * msg, size_t sz, int skynet_context_newsession(struct skynet_context *); struct message_queue * skynet_context_message_dispatch(struct skynet_monitor *, struct message_queue *); // return next queue int skynet_context_total(); +void skynet_context_dispatchall(struct skynet_context * context); // for skynet_error output before exit void skynet_context_endless(uint32_t handle); // for monitor diff --git a/skynet-src/skynet_socket.c b/skynet-src/skynet_socket.c index 3654d591..610bf2fc 100644 --- a/skynet-src/skynet_socket.c +++ b/skynet-src/skynet_socket.c @@ -133,12 +133,6 @@ skynet_socket_connect(struct skynet_context *ctx, const char *host, int port) { return socket_server_connect(SOCKET_SERVER, source, host, port); } -int -skynet_socket_block_connect(struct skynet_context *ctx, const char *host, int port) { - uint32_t source = skynet_context_handle(ctx); - return socket_server_block_connect(SOCKET_SERVER, source, host, port); -} - int skynet_socket_bind(struct skynet_context *ctx, int fd) { uint32_t source = skynet_context_handle(ctx); diff --git a/skynet-src/skynet_socket.h b/skynet-src/skynet_socket.h index 3275521d..2c0eb93b 100644 --- a/skynet-src/skynet_socket.h +++ b/skynet-src/skynet_socket.h @@ -25,7 +25,6 @@ int skynet_socket_send(struct skynet_context *ctx, int id, void *buffer, int sz) void skynet_socket_send_lowpriority(struct skynet_context *ctx, int id, void *buffer, int sz); int skynet_socket_listen(struct skynet_context *ctx, const char *host, int port, int backlog); int skynet_socket_connect(struct skynet_context *ctx, const char *host, int port); -int skynet_socket_block_connect(struct skynet_context *ctx, const char *host, int port); int skynet_socket_bind(struct skynet_context *ctx, int fd); void skynet_socket_close(struct skynet_context *ctx, int id); void skynet_socket_start(struct skynet_context *ctx, int id); diff --git a/skynet-src/skynet_start.c b/skynet-src/skynet_start.c index d45ded49..c101da30 100644 --- a/skynet-src/skynet_start.c +++ b/skynet-src/skynet_start.c @@ -184,7 +184,7 @@ _start(int thread) { } static void -bootstrap(const char * cmdline) { +bootstrap(struct skynet_context * logger, const char * cmdline) { int sz = strlen(cmdline); char name[sz+1]; char args[sz+1]; @@ -192,6 +192,7 @@ bootstrap(const char * cmdline) { struct skynet_context *ctx = skynet_context_new(name, args); if (ctx == NULL) { skynet_error(NULL, "Bootstrap error : %s\n", cmdline); + skynet_context_dispatchall(logger); exit(1); } } @@ -210,7 +211,13 @@ skynet_start(struct skynet_config * config) { skynet_timer_init(); skynet_socket_init(); - bootstrap(config->bootstrap); + struct skynet_context *ctx = skynet_context_new("logger", config->logger); + if (ctx == NULL) { + fprintf(stderr, "Can't launch logger service\n"); + exit(1); + } + + bootstrap(ctx, config->bootstrap); _start(config->thread); skynet_socket_free(); diff --git a/skynet-src/skynet_timer.c b/skynet-src/skynet_timer.c index 32b71f77..3cd0c51d 100644 --- a/skynet-src/skynet_timer.c +++ b/skynet-src/skynet_timer.c @@ -48,6 +48,8 @@ struct timer { int time; uint32_t current; uint32_t starttime; + uint64_t current_point; + uint64_t origin_point; }; static struct timer * TI = NULL; @@ -218,18 +220,41 @@ skynet_timeout(uint32_t handle, int time, int session) { return session; } -static uint32_t -_gettime(void) { - uint32_t t; +// centisecond: 1/100 second +static void +systime(uint32_t *sec, uint32_t *cs) { #if !defined(__APPLE__) struct timespec ti; - clock_gettime(CLOCK_MONOTONIC, &ti); - t = (uint32_t)(ti.tv_sec & 0xffffff) * 100; + clock_gettime(CLOCK_REALTIME, &ti); + *sec = (uint32_t)ti.tv_sec; + *cs = (uint32_t)(ti.tv_nsec / 10000000); +#else + struct timeval tv; + gettimeofday(&tv, NULL); + *sec = tv.tv_sec; + *cs = tv.tv_usec / 10000; +#endif +} + +static uint64_t +gettime() { + uint64_t t; +#if !defined(__APPLE__) + +#ifdef CLOCK_MONOTONIC_RAW +#define CLOCK_TIMER CLOCK_MONOTONIC_RAW +#else +#define CLOCK_TIMER CLOCK_MONOTONIC +#endif + + struct timespec ti; + clock_gettime(CLOCK_TIMER, &ti); + t = (uint64_t)ti.tv_sec * 100; t += ti.tv_nsec / 10000000; #else struct timeval tv; gettimeofday(&tv, NULL); - t = (uint32_t)(tv.tv_sec & 0xffffff) * 100; + t = (uint64_t)tv.tv_sec * 100; t += tv.tv_usec / 10000; #endif return t; @@ -237,10 +262,19 @@ _gettime(void) { void skynet_updatetime(void) { - uint32_t ct = _gettime(); - if (ct != TI->current) { - int diff = ct>=TI->current?ct-TI->current:(0xffffff+1)*100-TI->current+ct; - TI->current = ct; + uint64_t cp = gettime(); + if(cp < TI->current_point) { + skynet_error(NULL, "time diff error: change from %lld to %lld", cp, TI->current_point); + } else if (cp != TI->current_point) { + uint32_t diff = (uint32_t)(cp - TI->current_point); + TI->current_point = cp; + + uint32_t oc = TI->current; + TI->current += diff; + if (TI->current < oc) { + // when cs > 0xffffffff(about 497 days), time rewind + TI->starttime += 0xffffffff / 100; + } int i; for (i=0;icurrent = _gettime(); - -#if !defined(__APPLE__) - struct timespec ti; - clock_gettime(CLOCK_REALTIME, &ti); - uint32_t sec = (uint32_t)ti.tv_sec; -#else - struct timeval tv; - gettimeofday(&tv, NULL); - uint32_t sec = (uint32_t)tv.tv_sec; -#endif - uint32_t mono = _gettime() / 100; - - TI->starttime = sec - mono; + systime(&TI->starttime, &TI->current); + uint64_t point = gettime(); + TI->current_point = point; + TI->origin_point = point; } + diff --git a/skynet-src/socket_server.c b/skynet-src/socket_server.c index f013fd1b..aa105a16 100644 --- a/skynet-src/socket_server.c +++ b/skynet-src/socket_server.c @@ -147,6 +147,8 @@ reserve_id(struct socket_server *ss) { struct socket *s = &ss->slot[id % MAX_SOCKET]; if (s->type == SOCKET_TYPE_INVALID) { if (__sync_bool_compare_and_swap(&s->type, SOCKET_TYPE_INVALID, SOCKET_TYPE_RESERVE)) { + s->id = id; + s->fd = -1; return id; } else { // retry @@ -287,7 +289,7 @@ new_fd(struct socket_server *ss, int id, int fd, uintptr_t opaque, bool add) { // return -1 when connecting static int -open_socket(struct socket_server *ss, struct request_open * request, struct socket_message *result, bool blocking) { +open_socket(struct socket_server *ss, struct request_open * request, struct socket_message *result) { int id = request->id; result->opaque = request->opaque; result->id = id; @@ -316,18 +318,14 @@ open_socket(struct socket_server *ss, struct request_open * request, struct sock continue; } socket_keepalive(sock); - if (!blocking) { - sp_nonblocking(sock); - } + sp_nonblocking(sock); status = connect( sock, ai_ptr->ai_addr, ai_ptr->ai_addrlen); if ( status != 0 && errno != EINPROGRESS) { close(sock); sock = -1; continue; } - if (blocking) { - sp_nonblocking(sock); - } + sp_nonblocking(sock); break; } @@ -512,7 +510,7 @@ send_socket(struct socket_server *ss, struct request_send * request, struct sock return -1; } assert(s->type != SOCKET_TYPE_PLISTEN && s->type != SOCKET_TYPE_LISTEN); - if (send_buffer_empty(s)) { + if (send_buffer_empty(s) && s->type == SOCKET_TYPE_CONNECTED) { int n = write(s->fd, request->buffer, request->sz); if (n<0) { switch(errno) { @@ -687,7 +685,7 @@ ctrl_cmd(struct socket_server *ss, struct socket_message *result) { case 'K': return close_socket(ss,(struct request_close *)buffer, result); case 'O': - return open_socket(ss, (struct request_open *)buffer, result, false); + return open_socket(ss, (struct request_open *)buffer, result); case 'X': result->opaque = 0; result->id = 0; @@ -765,7 +763,9 @@ report_connect(struct socket_server *ss, struct socket *s, struct socket_message result->opaque = s->opaque; result->id = s->id; result->ud = 0; - sp_write(ss->event_fd, s->fd, s, false); + if (send_buffer_empty(s)) { + sp_write(ss->event_fd, s->fd, s, false); + } union sockaddr_all u; socklen_t slen = sizeof(u); if (getpeername(s->fd, &u.s, &slen) == 0) { @@ -939,19 +939,6 @@ socket_server_connect(struct socket_server *ss, uintptr_t opaque, const char * a return request.u.open.id; } -int -socket_server_block_connect(struct socket_server *ss, uintptr_t opaque, const char * addr, int port) { - struct request_package request; - struct socket_message result; - open_request(ss, &request, opaque, addr, port); - int ret = open_socket(ss, &request.u.open, &result, true); - if (ret == SOCKET_OPEN) { - return result.id; - } else { - return -1; - } -} - // return -1 when error int64_t socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz) { @@ -959,7 +946,6 @@ socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz if (s->id != id || s->type == SOCKET_TYPE_INVALID) { return -1; } - assert(s->type != SOCKET_TYPE_RESERVE); struct request_package request; request.u.send.id = id; @@ -976,7 +962,6 @@ socket_server_send_lowpriority(struct socket_server *ss, int id, const void * bu if (s->id != id || s->type == SOCKET_TYPE_INVALID) { return; } - assert(s->type != SOCKET_TYPE_RESERVE); struct request_package request; request.u.send.id = id; diff --git a/skynet-src/socket_server.h b/skynet-src/socket_server.h index 2fbb6939..ea15c0e1 100644 --- a/skynet-src/socket_server.h +++ b/skynet-src/socket_server.h @@ -36,6 +36,4 @@ int socket_server_listen(struct socket_server *, uintptr_t opaque, const char * int socket_server_connect(struct socket_server *, uintptr_t opaque, const char * addr, int port); int socket_server_bind(struct socket_server *, uintptr_t opaque, int fd); -int socket_server_block_connect(struct socket_server *, uintptr_t opaque, const char * addr, int port); - #endif diff --git a/test/testharborlink.lua b/test/testharborlink.lua new file mode 100644 index 00000000..031e38ec --- /dev/null +++ b/test/testharborlink.lua @@ -0,0 +1,11 @@ +local skynet = require "skynet" +local harbor = require "skynet.harbor" + +skynet.start(function() + print("wait for harbor 2") + print("run skynet examples/config_log please") + harbor.connect(2) + print("harbor 2 connected") + harbor.link(2) + print("disconnected") +end) diff --git a/test/testredis2.lua b/test/testredis2.lua new file mode 100644 index 00000000..f12d8100 --- /dev/null +++ b/test/testredis2.lua @@ -0,0 +1,53 @@ +local skynet = require "skynet" +local redis = require "redis" + +local db + +function add1(key, count) + local t = {} + for i = 1, count do + t[2*i -1] = "key" ..i + t[2*i] = "value" .. i + end + db:hmset(key, table.unpack(t)) +end + +function add2(key, count) + local t = {} + for i = 1, count do + t[2*i -1] = "key" ..i + t[2*i] = "value" .. i + end + table.insert(t, 1, key) + db:hmset(t) +end + +function __init__() + db = redis.connect { + host = "127.0.0.1", + port = 6300, + db = 0, + auth = "foobared" + } + print("dbsize:", db:dbsize()) + local ok, msg = xpcall(add1, debug.traceback, "test1", 250000) + if not ok then + print("add1 failed", msg) + else + print("add1 succeed") + + end + + local ok, msg = xpcall(add2, debug.traceback, "test2", 250000) + if not ok then + print("add2 failed", msg) + else + print("add2 succeed") + end + print("dbsize:", db:dbsize()) + + print("redistest launched") +end + +skynet.start(__init__) + diff --git a/test/time.lua b/test/time.lua new file mode 100644 index 00000000..be02805c --- /dev/null +++ b/test/time.lua @@ -0,0 +1,22 @@ +local skynet = require "skynet" +skynet.start(function() + print(skynet.starttime()) + print(skynet.now()) + + skynet.timeout(1, function() + print("in 1", skynet.now()) + end) + skynet.timeout(2, function() + print("in 2", skynet.now()) + end) + skynet.timeout(3, function() + print("in 3", skynet.now()) + end) + + skynet.timeout(4, function() + print("in 4", skynet.now()) + end) + skynet.timeout(100, function() + print("in 100", skynet.now()) + end) +end)