mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 20:23:06 +00:00
add module skynet.harbor
This commit is contained in:
8
.gitignore
vendored
8
.gitignore
vendored
@@ -1,10 +1,10 @@
|
||||
*.o
|
||||
*.a
|
||||
skynet
|
||||
skynet.pid
|
||||
./skynet
|
||||
./skynet.pid
|
||||
3rd/lua/lua
|
||||
3rd/lua/luac
|
||||
cservice
|
||||
luaclib
|
||||
./cservice
|
||||
./luaclib
|
||||
*.so
|
||||
*.dSYM
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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
|
||||
|
||||
40
lualib/skynet/harbor.lua
Normal file
40
lualib/skynet/harbor.lua
Normal file
@@ -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
|
||||
@@ -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: {
|
||||
|
||||
@@ -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;i<REMOTE_MAX;i++) {
|
||||
struct monitor_set * m = h->monitor[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;i<m->n;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;i<m->n;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: {
|
||||
|
||||
@@ -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"))
|
||||
|
||||
@@ -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;i<GLOBALNAME_LENGTH;i++) {
|
||||
char c = rname->name[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);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
11
test/testharborlink.lua
Normal file
11
test/testharborlink.lua
Normal file
@@ -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)
|
||||
Reference in New Issue
Block a user