mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
add cluster.register and cluster.query
This commit is contained in:
@@ -1,16 +1,17 @@
|
|||||||
local skynet = require "skynet"
|
local skynet = require "skynet"
|
||||||
local cluster = require "cluster"
|
local cluster = require "cluster"
|
||||||
require "skynet.manager" -- import skynet.name
|
|
||||||
local snax = require "snax"
|
local snax = require "snax"
|
||||||
|
|
||||||
skynet.start(function()
|
skynet.start(function()
|
||||||
local sdb = skynet.newservice("simpledb")
|
local sdb = skynet.newservice("simpledb")
|
||||||
skynet.name(".simpledb", sdb)
|
-- register name "sdb" for simpledb, you can use cluster.query() later.
|
||||||
|
-- See cluster2.lua
|
||||||
|
cluster.register("sdb", sdb)
|
||||||
|
|
||||||
print(skynet.call(".simpledb", "lua", "SET", "a", "foobar"))
|
print(skynet.call(sdb, "lua", "SET", "a", "foobar"))
|
||||||
print(skynet.call(".simpledb", "lua", "SET", "b", "foobar2"))
|
print(skynet.call(sdb, "lua", "SET", "b", "foobar2"))
|
||||||
print(skynet.call(".simpledb", "lua", "GET", "a"))
|
print(skynet.call(sdb, "lua", "GET", "a"))
|
||||||
print(skynet.call(".simpledb", "lua", "GET", "b"))
|
print(skynet.call(sdb, "lua", "GET", "b"))
|
||||||
cluster.open "db"
|
cluster.open "db"
|
||||||
cluster.open "db2"
|
cluster.open "db2"
|
||||||
-- unique snax service
|
-- unique snax service
|
||||||
|
|||||||
@@ -2,15 +2,18 @@ local skynet = require "skynet"
|
|||||||
local cluster = require "cluster"
|
local cluster = require "cluster"
|
||||||
|
|
||||||
skynet.start(function()
|
skynet.start(function()
|
||||||
local proxy = cluster.proxy("db", ".simpledb")
|
-- query name "sdb" of cluster db.
|
||||||
|
local sdb = cluster.query("db", "sdb")
|
||||||
|
print("db.sbd=",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))
|
print(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)
|
||||||
|
|
||||||
print(cluster.call("db", ".simpledb", "GET", "a"))
|
print(cluster.call("db", sdb, "GET", "a"))
|
||||||
print(cluster.call("db2", ".simpledb", "GET", "b"))
|
print(cluster.call("db2", sdb, "GET", "b"))
|
||||||
|
|
||||||
-- test snax service
|
-- test snax service
|
||||||
local pingserver = cluster.snax("db", "pingserver")
|
local pingserver = cluster.snax("db", "pingserver")
|
||||||
|
|||||||
@@ -219,8 +219,8 @@ unpackreq_number(lua_State *L, const uint8_t * buf, int sz) {
|
|||||||
|
|
||||||
static int
|
static int
|
||||||
unpackmreq_number(lua_State *L, const uint8_t * buf, int sz) {
|
unpackmreq_number(lua_State *L, const uint8_t * buf, int sz) {
|
||||||
if (sz != 15) {
|
if (sz != 13) {
|
||||||
return luaL_error(L, "Invalid cluster message size %d (multi req must be 15)", sz);
|
return luaL_error(L, "Invalid cluster message size %d (multi req must be 13)", sz);
|
||||||
}
|
}
|
||||||
uint32_t address = unpack_uint32(buf+1);
|
uint32_t address = unpack_uint32(buf+1);
|
||||||
uint32_t session = unpack_uint32(buf+5);
|
uint32_t session = unpack_uint32(buf+5);
|
||||||
|
|||||||
@@ -33,6 +33,16 @@ function cluster.snax(node, name, address)
|
|||||||
return snax.bind(handle, name)
|
return snax.bind(handle, name)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
function cluster.register(name, addr)
|
||||||
|
assert(type(name) == "string")
|
||||||
|
assert(addr == nil or type(addr) == "number")
|
||||||
|
return skynet.call(clusterd, "lua", "register", name, addr)
|
||||||
|
end
|
||||||
|
|
||||||
|
function cluster.query(node, name)
|
||||||
|
return skynet.call(clusterd, "lua", "req", node, 0, skynet.pack(name))
|
||||||
|
end
|
||||||
|
|
||||||
skynet.init(function()
|
skynet.init(function()
|
||||||
clusterd = skynet.uniqueservice("clusterd")
|
clusterd = skynet.uniqueservice("clusterd")
|
||||||
end)
|
end)
|
||||||
|
|||||||
@@ -97,6 +97,21 @@ function command.proxy(source, node, name)
|
|||||||
skynet.ret(skynet.pack(proxy[fullname]))
|
skynet.ret(skynet.pack(proxy[fullname]))
|
||||||
end
|
end
|
||||||
|
|
||||||
|
local register_name = {}
|
||||||
|
|
||||||
|
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
|
||||||
|
end
|
||||||
|
register_name[addr] = name
|
||||||
|
register_name[name] = addr
|
||||||
|
skynet.ret(nil)
|
||||||
|
skynet.error(string.format("Register [%s] :%08x", name, addr))
|
||||||
|
end
|
||||||
|
|
||||||
local large_request = {}
|
local large_request = {}
|
||||||
|
|
||||||
function command.socket(source, subcmd, fd, msg)
|
function command.socket(source, subcmd, fd, msg)
|
||||||
@@ -122,8 +137,20 @@ function command.socket(source, subcmd, fd, msg)
|
|||||||
return
|
return
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
local ok , msg, sz = pcall(skynet.rawcall, addr, "lua", msg, sz)
|
local ok, response
|
||||||
local 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
|
||||||
|
else
|
||||||
|
ok , msg, sz = pcall(skynet.rawcall, addr, "lua", msg, sz)
|
||||||
|
end
|
||||||
if ok then
|
if ok then
|
||||||
response = cluster.packresponse(session, true, msg, sz)
|
response = cluster.packresponse(session, true, msg, sz)
|
||||||
if type(response) == "table" then
|
if type(response) == "table" then
|
||||||
|
|||||||
Reference in New Issue
Block a user