mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 12:20:41 +00:00
use clusteragent to dispatch cluster request
This commit is contained in:
@@ -8,7 +8,7 @@ skynet.start(function()
|
|||||||
local proxy = cluster.proxy("db", sdb)
|
local proxy = cluster.proxy("db", sdb)
|
||||||
local largekey = string.rep("X", 128*1024)
|
local largekey = string.rep("X", 128*1024)
|
||||||
local largevalue = string.rep("R", 100 * 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)
|
local v = skynet.call(proxy, "lua", "GET", largekey)
|
||||||
assert(largevalue == v)
|
assert(largevalue == v)
|
||||||
skynet.send(proxy, "lua", "PING", "proxy")
|
skynet.send(proxy, "lua", "PING", "proxy")
|
||||||
|
|||||||
@@ -312,9 +312,16 @@ unpackmreq_string(lua_State *L, const uint8_t * buf, int sz, int is_push) {
|
|||||||
|
|
||||||
static int
|
static int
|
||||||
lunpackrequest(lua_State *L) {
|
lunpackrequest(lua_State *L) {
|
||||||
size_t ssz;
|
int sz;
|
||||||
const char *msg = luaL_checklstring(L,1,&ssz);
|
const char *msg;
|
||||||
int sz = (int)ssz;
|
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]) {
|
switch (msg[0]) {
|
||||||
case 0:
|
case 0:
|
||||||
return unpackreq_number(L, (const uint8_t *)msg, sz);
|
return unpackreq_number(L, (const uint8_t *)msg, sz);
|
||||||
|
|||||||
92
service/clusteragent.lua
Normal file
92
service/clusteragent.lua
Normal file
@@ -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)
|
||||||
@@ -156,14 +156,24 @@ function command.proxy(source, node, name)
|
|||||||
skynet.ret(skynet.pack(proxy[fullname]))
|
skynet.ret(skynet.pack(proxy[fullname]))
|
||||||
end
|
end
|
||||||
|
|
||||||
|
local cluster_agent = {} -- fd:service
|
||||||
local register_name = {}
|
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)
|
function command.register(source, name, addr)
|
||||||
assert(register_name[name] == nil)
|
assert(register_name[name] == nil)
|
||||||
addr = addr or source
|
addr = addr or source
|
||||||
local old_name = register_name[addr]
|
local old_name = register_name[addr]
|
||||||
if old_name then
|
if old_name then
|
||||||
register_name[old_name] = nil
|
register_name[old_name] = nil
|
||||||
|
clearnamecache()
|
||||||
end
|
end
|
||||||
register_name[addr] = name
|
register_name[addr] = name
|
||||||
register_name[name] = addr
|
register_name[name] = addr
|
||||||
@@ -171,76 +181,35 @@ function command.register(source, name, addr)
|
|||||||
skynet.error(string.format("Register [%s] :%08x", name, addr))
|
skynet.error(string.format("Register [%s] :%08x", name, addr))
|
||||||
end
|
end
|
||||||
|
|
||||||
local large_request = {}
|
function command.queryname(source, name)
|
||||||
|
skynet.ret(skynet.pack(register_name[name]))
|
||||||
|
end
|
||||||
|
|
||||||
function command.socket(source, subcmd, fd, msg)
|
function command.socket(source, subcmd, fd, msg)
|
||||||
if subcmd == "data" then
|
if subcmd == "open" 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
|
|
||||||
skynet.error(string.format("socket accept from %s", msg))
|
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
|
else
|
||||||
large_request[fd] = nil
|
if subcmd == "close" or subcmd == "error" then
|
||||||
skynet.error(string.format("socket %s %d %s", subcmd, fd, msg or ""))
|
-- 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
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user