From 7304a990cb9fe9c32d87b00a69a1a98efef59c98 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Sat, 14 Apr 2018 16:11:26 +0800 Subject: [PATCH] use clusteragent to dispatch cluster request --- examples/cluster2.lua | 2 +- lualib-src/lua-cluster.c | 13 +++-- service/clusteragent.lua | 92 +++++++++++++++++++++++++++++++++++ service/clusterd.lua | 101 ++++++++++++++------------------------- 4 files changed, 138 insertions(+), 70 deletions(-) create mode 100644 service/clusteragent.lua diff --git a/examples/cluster2.lua b/examples/cluster2.lua index 3b7a2c02..d249e7ac 100644 --- a/examples/cluster2.lua +++ b/examples/cluster2.lua @@ -8,7 +8,7 @@ skynet.start(function() local proxy = cluster.proxy("db", sdb) local largekey = string.rep("X", 128*1024) local largevalue = string.rep("R", 100 * 1024) - print(skynet.call(proxy, "lua", "SET", largekey, largevalue)) + skynet.call(proxy, "lua", "SET", largekey, largevalue) local v = skynet.call(proxy, "lua", "GET", largekey) assert(largevalue == v) skynet.send(proxy, "lua", "PING", "proxy") diff --git a/lualib-src/lua-cluster.c b/lualib-src/lua-cluster.c index 3c55ac62..e4e2f57f 100644 --- a/lualib-src/lua-cluster.c +++ b/lualib-src/lua-cluster.c @@ -312,9 +312,16 @@ unpackmreq_string(lua_State *L, const uint8_t * buf, int sz, int is_push) { static int lunpackrequest(lua_State *L) { - size_t ssz; - const char *msg = luaL_checklstring(L,1,&ssz); - int sz = (int)ssz; + int sz; + const char *msg; + if (lua_type(L, 1) == LUA_TLIGHTUSERDATA) { + msg = (const char *)lua_touserdata(L, 1); + sz = luaL_checkinteger(L, 2); + } else { + size_t ssz; + msg = luaL_checklstring(L,1,&ssz); + sz = (int)ssz; + } switch (msg[0]) { case 0: return unpackreq_number(L, (const uint8_t *)msg, sz); diff --git a/service/clusteragent.lua b/service/clusteragent.lua new file mode 100644 index 00000000..d297dbb7 --- /dev/null +++ b/service/clusteragent.lua @@ -0,0 +1,92 @@ +local skynet = require "skynet" +local sc = require "skynet.socketchannel" +local socket = require "skynet.socket" +local cluster = require "skynet.cluster.core" + +local clusterd, gate, fd = ... +clusterd = tonumber(clusterd) +gate = tonumber(gate) +fd = tonumber(fd) + +local large_request = {} +local register_name = {} + +local function dispatch_request(_,_,addr, session, msg, padding, is_push) + local sz + if padding then + local req = large_request[session] or { addr = addr , is_push = is_push } + large_request[session] = req + table.insert(req, msg) + return + else + local req = large_request[session] + if req then + large_request[session] = nil + table.insert(req, msg) + msg,sz = cluster.concat(req) + addr = req.addr + is_push = req.is_push + end + if not msg then + local response = cluster.packresponse(session, false, "Invalid large req") + socket.write(fd, response) + return + end + end + local ok, response + if addr == 0 then + local name = skynet.unpack(msg, sz) + local addr = register_name[name] + if addr == nil then + addr = skynet.call(clusterd, "lua", "queryname", name) + register_name[name] = addr + end + if addr then + ok = true + msg, sz = skynet.pack(addr) + else + ok = false + msg = "name not found" + end + elseif is_push then + skynet.rawsend(addr, "lua", msg, sz) + return -- no response + else + ok , msg, sz = pcall(skynet.rawcall, addr, "lua", msg, sz) + end + if ok then + response = cluster.packresponse(session, true, msg, sz) + if type(response) == "table" then + for _, v in ipairs(response) do + socket.lwrite(fd, v) + end + else + socket.write(fd, response) + end + else + response = cluster.packresponse(session, false, msg) + socket.write(fd, response) + end +end + +skynet.start(function() + skynet.register_protocol { + name = "client", + id = skynet.PTYPE_CLIENT, + unpack = cluster.unpackrequest, + dispatch = dispatch_request, + } + -- fd can write, but don't read fd, the data package will forward from gate though client protocol. + skynet.call(gate, "lua", "forward", fd) + + skynet.dispatch("lua", function(_,source, cmd, ...) + if cmd == "exit" then + socket.close(fd) + skynet.exit() + elseif cmd == "namechange" then + register_name = {} + else + skynet.error(string.format("Invalid command %s from %s", cmd, skynet.address(source))) + end + end) +end) diff --git a/service/clusterd.lua b/service/clusterd.lua index 3196820a..381e6ad1 100644 --- a/service/clusterd.lua +++ b/service/clusterd.lua @@ -156,14 +156,24 @@ function command.proxy(source, node, name) skynet.ret(skynet.pack(proxy[fullname])) end +local cluster_agent = {} -- fd:service local register_name = {} +local function clearnamecache() + for fd, service in pairs(cluster_agent) do + if type(service) == "number" then + skynet.send(service, "lua", "namechange") + end + end +end + function command.register(source, name, addr) assert(register_name[name] == nil) addr = addr or source local old_name = register_name[addr] if old_name then register_name[old_name] = nil + clearnamecache() end register_name[addr] = name register_name[name] = addr @@ -171,76 +181,35 @@ function command.register(source, name, addr) skynet.error(string.format("Register [%s] :%08x", name, addr)) end -local large_request = {} +function command.queryname(source, name) + skynet.ret(skynet.pack(register_name[name])) +end function command.socket(source, subcmd, fd, msg) - if subcmd == "data" then - local sz - local addr, session, msg, padding, is_push = cluster.unpackrequest(msg) - if padding then - local requests = large_request[fd] - if requests == nil then - requests = {} - large_request[fd] = requests - end - local req = requests[session] or { addr = addr , is_push = is_push } - requests[session] = req - table.insert(req, msg) - return - else - local requests = large_request[fd] - if requests then - local req = requests[session] - if req then - requests[session] = nil - table.insert(req, msg) - msg,sz = cluster.concat(req) - addr = req.addr - is_push = req.is_push - end - end - if not msg then - local response = cluster.packresponse(session, false, "Invalid large req") - socket.write(fd, response) - return - end - end - local ok, response - if addr == 0 then - local name = skynet.unpack(msg, sz) - local addr = register_name[name] - if addr then - ok = true - msg, sz = skynet.pack(addr) - else - ok = false - msg = "name not found" - end - elseif is_push then - skynet.rawsend(addr, "lua", msg, sz) - return -- no response - else - ok , msg, sz = pcall(skynet.rawcall, addr, "lua", msg, sz) - end - if ok then - response = cluster.packresponse(session, true, msg, sz) - if type(response) == "table" then - for _, v in ipairs(response) do - socket.lwrite(fd, v) - end - else - socket.write(fd, response) - end - else - response = cluster.packresponse(session, false, msg) - socket.write(fd, response) - end - elseif subcmd == "open" then + if subcmd == "open" then skynet.error(string.format("socket accept from %s", msg)) - skynet.call(source, "lua", "accept", fd) + -- new cluster agent + cluster_agent[fd] = false + local agent = skynet.newservice("clusteragent", skynet.self(), source, fd) + local closed = cluster_agent[fd] + cluster_agent[fd] = agent + if closed then + skynet.send(agent, "lua", "exit") + cluster_agent[fd] = nil + end else - large_request[fd] = nil - skynet.error(string.format("socket %s %d %s", subcmd, fd, msg or "")) + if subcmd == "close" or subcmd == "error" then + -- close cluster agent + local agent = cluster_agent[fd] + if type(agent) == "boolean" then + cluster_agent[fd] = true + else + skynet.send(agent, "lua", "exit") + cluster_agent[fd] = nil + end + else + skynet.error(string.format("socket %s %d %s", subcmd, fd, msg or "")) + end end end