From ad9898a209e5ee38629e9c6b5f991eca9118f1c7 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Thu, 19 Jun 2014 15:43:41 +0800 Subject: [PATCH] add module skynet.harbor --- .gitignore | 8 +- HISTORY.md | 5 + lualib-src/lua-skynet.c | 2 +- lualib/skynet.lua | 25 +++- lualib/skynet/harbor.lua | 40 +++++ service-src/service_dummy.c | 35 ++++- service-src/service_harbor.c | 137 +++++++++++++++++- service/bootstrap.lua | 7 +- skynet-src/skynet_harbor.c | 17 --- skynet-src/skynet_harbor.h | 1 - skynet-src/skynet_server.c | 11 +- test/testharborlink.lua | 11 ++ examples/redistest.lua => test/testredis2.lua | 0 13 files changed, 253 insertions(+), 46 deletions(-) create mode 100644 lualib/skynet/harbor.lua create mode 100644 test/testharborlink.lua rename examples/redistest.lua => test/testredis2.lua (100%) 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 9d39efaf..addc8955 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -1,3 +1,8 @@ +Dev version +----------- +* Optimize redis driver `compose_message`. +* Add module skynet.harbor for monitor harbor connect/disconnect, see test/testharborlink.lua . + v0.3.1 (2014-6-16) ----------- * Bugfix: lua mongo driver . Hold reply string before decode bson data. diff --git a/lualib-src/lua-skynet.c b/lualib-src/lua-skynet.c index 0025dd52..db49e997 100644 --- a/lualib-src/lua-skynet.c +++ b/lualib-src/lua-skynet.c @@ -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); diff --git a/lualib/skynet.lua b/lualib/skynet.lua index 430881fe..93b67d65 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 diff --git a/lualib/skynet/harbor.lua b/lualib/skynet/harbor.lua new file mode 100644 index 00000000..06f47bb6 --- /dev/null +++ b/lualib/skynet/harbor.lua @@ -0,0 +1,40 @@ +local skynet = require "skynet" + +local harbor = {} + +local HARBOR = skynet.getenv "harbor_address" + +if HARBOR then + HARBOR = tonumber("0x" .. string.sub(HARBOR , 2)) +end + +function harbor.globalname(name, handle) + assert(HARBOR) + handle = handle or skynet.self() + skynet.redirect(HARBOR, handle, "system", 0, "R " .. name) +end + +function harbor.init(h) + assert(HARBOR == nil) + HARBOR = h + skynet.setenv("harbor_address", skynet.address(h)) +end + +function harbor.link(id) + assert(HARBOR) + skynet.call(HARBOR, "system", "M " .. tostring(id)) +end + +function harbor.connect(id) + assert(HARBOR) + skynet.call(HARBOR, "system", "C " .. tostring(id)) +end + +skynet.register_protocol { + name = "system", + id = skynet.PTYPE_SYSTEM, + pack = function(...) return ... end, + unpack = skynet.tostring, +} + +return harbor diff --git a/service-src/service_dummy.c b/service-src/service_dummy.c index 006a276f..a1f19d4c 100644 --- a/service-src/service_dummy.c +++ b/service-src/service_dummy.c @@ -257,15 +257,42 @@ _send_name(struct dummy *h, uint32_t source, const char name[GLOBALNAME_LENGTH], } } +static void +dummy_command(struct dummy * h, const char * msg, size_t sz, int session, uint32_t source) { + switch(msg[0]) { + case 'R' : { + // register global name + const char * name = msg + 2; + int s = (int)sz; + s -= 2; + 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 'C' : + case 'M' : + skynet_error(h->ctx, "Don't support harbor monitor in cluster dummy mode"); + skynet_send(h->ctx, 0, source, PTYPE_ERROR, session, NULL, 0); + break; + default: + skynet_error(h->ctx, "Unknown command %s", msg); + return; + } +} + 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); + dummy_command(h, msg, sz, session, source); return 0; } default: { diff --git a/service-src/service_harbor.c b/service-src/service_harbor.c index 926ed46f..aa6763a7 100644 --- a/service-src/service_harbor.c +++ b/service-src/service_harbor.c @@ -47,6 +47,17 @@ struct remote_message_header { uint32_t session; }; +struct monitor_response { + uint32_t addr; + int session; +}; + +struct monitor_set { + int cap; + int n; + struct monitor_response * resp; +}; + // 12 is sizeof(struct remote_message_header) #define HEADER_COOKIE_LENGTH 12 @@ -60,8 +71,61 @@ struct harbor { int remote_fd[REMOTE_MAX]; bool connected[REMOTE_MAX]; char * remote_addr[REMOTE_MAX]; + struct monitor_set * monitor[REMOTE_MAX]; }; +static void +monitor_free(struct harbor *h) { + int i; + for (i=0;imonitor[i]; + if (m) { + skynet_free(m->resp); + skynet_free(m); + h->monitor[i] = NULL; + } + } +} + +static void +monitor_add(struct harbor *h, int id, uint32_t addr, int session) { + struct monitor_set * m = h->monitor[id]; + if (m == NULL) { + m = skynet_malloc(sizeof(*m)); + m->cap = 4; + m->n = 0; + m->resp = skynet_malloc(m->cap * sizeof(struct monitor_response)); + h->monitor[id] = m; + } + if (m->n >= m->cap) { + assert(m->n == m->cap); + struct monitor_response * resp = skynet_malloc(m->cap * 2 * sizeof(struct monitor_response)); + int i; + for (i=0;in;i++) { + resp[i] = m->resp[i]; + } + m->cap *= 2; + skynet_free(m->resp); + m->resp = resp; + } + struct monitor_response * resp = &m->resp[m->n++]; + resp->addr = addr; + resp->session = session; +} + +static void +monitor_clear(struct harbor *h, int id) { + struct monitor_set * m = h->monitor[id]; + if (m) { + int i; + for (i=0;in;i++) { + struct monitor_response * resp = &m->resp[i]; + skynet_send(h->ctx, 0, resp->addr, PTYPE_RESPONSE, resp->session, NULL, 0); + } + m->n = 0; + } +} + // hash table static void @@ -211,6 +275,7 @@ harbor_create(void) { h->remote_fd[i] = -1; h->connected[i] = false; h->remote_addr[i] = NULL; + h->monitor[i] = NULL; } h->map = _hash_new(); return h; @@ -232,6 +297,7 @@ harbor_release(struct harbor *h) { } } _hash_delete(h->map); + monitor_free(h); skynet_free(h); } @@ -313,6 +379,14 @@ _send_remote(struct skynet_context * ctx, int fd, const char * buffer, size_t sz } } +static void +response_close(struct harbor *h, int id) { + if (h->connected[id]) { + monitor_clear(h, id); + } + h->connected[id] = false; +} + static void _update_remote_address(struct harbor *h, int harbor_id, const char * ipaddr) { if (harbor_id == h->id) { @@ -326,7 +400,7 @@ _update_remote_address(struct harbor *h, int harbor_id, const char * ipaddr) { h->remote_addr[harbor_id] = NULL; } h->remote_fd[harbor_id] = _connect_to(h, ipaddr, false); - h->connected[harbor_id] = false; + response_close(h, harbor_id); } static void @@ -466,7 +540,7 @@ close_harbor(struct harbor *h, int fd) { skynet_error(h->ctx, "Harbor %d closed",id); skynet_socket_close(h->ctx, fd); h->remote_fd[id] = -1; - h->connected[id] = false; + response_close(h,id); } static void @@ -475,9 +549,63 @@ open_harbor(struct harbor *h, int fd) { if (id == 0) return; assert(h->connected[id] == false); + monitor_clear(h, id); h->connected[id] = true; } +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 'R' : { + // register global name + 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; + _remote_register_name(h, rn.name, rn.handle); + break; + } + case 'C' : + case 'M' : { + if (s <= 0) { + skynet_error(h->ctx, "Invalid harbor montior"); + skynet_send(h->ctx, 0, source, PTYPE_ERROR, session, NULL, 0); + return; + } + int hid = strtol(name, NULL, 10); + if (hid <= 0 || hid >= REMOTE_MAX) { + skynet_error(h->ctx, "Invalid harbor montior id : %s", name); + skynet_send(h->ctx, 0, source, PTYPE_ERROR, session, NULL, 0); + return; + } + if (msg[0] == 'M') { + if (!h->connected[hid]) { + skynet_send(h->ctx, 0, source, PTYPE_RESPONSE, session, NULL, 0); + return; + } + } else { + assert(msg[0] == 'C'); + if (h->connected[hid]) { + skynet_send(h->ctx, 0, source, PTYPE_RESPONSE, session, NULL, 0); + return; + } + } + monitor_add(h, hid, source, session); + break; + } + default: + skynet_error(h->ctx, "Unknown command %s", msg); + return; + } +} + static int _mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { struct harbor * h = ud; @@ -536,10 +664,7 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin 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: { diff --git a/service/bootstrap.lua b/service/bootstrap.lua index 7c92b419..3d66115b 100644 --- a/service/bootstrap.lua +++ b/service/bootstrap.lua @@ -1,4 +1,5 @@ local skynet = require "skynet" +local harbor = require "skynet.harbor" skynet.start(function() assert(skynet.launch("logger", skynet.getenv "logger")) @@ -9,7 +10,8 @@ skynet.start(function() assert(standalone == nil) standalone = true skynet.setenv("standalone", "true") - assert(skynet.launch("dummy")) + local dummy = assert(skynet.launch("dummy")) + harbor.init(dummy) else local master_addr = skynet.getenv "master" @@ -19,7 +21,8 @@ skynet.start(function() local local_addr = skynet.getenv "address" - assert(skynet.launch("harbor",master_addr, local_addr, harbor_id)) + local h = assert(skynet.launch("harbor",master_addr, local_addr, harbor_id)) + harbor.init(h) end local launcher = assert(skynet.launch("snlua","launcher")) 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_server.c b/skynet-src/skynet_server.c index 98fcfcd5..e20f8102 100644 --- a/skynet-src/skynet_server.c +++ b/skynet-src/skynet_server.c @@ -342,11 +342,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 +373,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; } 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/examples/redistest.lua b/test/testredis2.lua similarity index 100% rename from examples/redistest.lua rename to test/testredis2.lua