diff --git a/lualib/mqueue.lua b/lualib/mqueue.lua new file mode 100644 index 00000000..bc866b63 --- /dev/null +++ b/lualib/mqueue.lua @@ -0,0 +1,60 @@ +local skynet = require "skynet" +local c = require "skynet.c" + +local mqueue = {} +local init_once +local thread_id +local message_queue = {} + +skynet.register_protocol { + name = "queue", + id = 10, + pack = skynet.pack, + unpack = skynet.unpack, + dispatch = function(session, from, ...) + table.insert(message_queue, {..., session = session, addr = from}) + if thread_id then + skynet.wakeup(thread_id) + thread_id = nil + end + end +} + +local function message_dispatch(f) + while true do + if #message_queue==0 then + thread_id = coroutine.running() + skynet.wait() + else + local msg = table.remove(message_queue,1) + local session = msg.session + if session == 0 then + assert(f(table.unpack(msg))==nil, "message in queue returns value") + else + local data, size = skynet.pack(f(table.unpack(msg))) + -- 1 means response + c.send(msg.addr, 1, session, data, size) + end + end + end +end + +function mqueue.register(f) + assert(init_once == nil) + init_once = true + skynet.fork(message_dispatch,f) +end + +function mqueue.call(addr, ...) + return skynet.call(addr, "queue", ...) +end + +function mqueue.send(addr, ...) + return skynet.send(addr, "queue", ...) +end + +function mqueue.size() + return #message_queue +end + +return mqueue diff --git a/service/pingqueue.lua b/service/pingqueue.lua new file mode 100644 index 00000000..eca240f8 --- /dev/null +++ b/service/pingqueue.lua @@ -0,0 +1,12 @@ +local skynet = require "skynet" +local mqueue = require "mqueue" + +skynet.start(function() + local id = 0 + local pingserver = skynet.newservice "pingserver" + mqueue.register(function(str) + id = id + 1 + str = string.format("id = %d , %s",id, str) + return skynet.call(pingserver, "lua", "PING", str) + end) +end) diff --git a/service/testqueue.lua b/service/testqueue.lua new file mode 100644 index 00000000..47a86e0f --- /dev/null +++ b/service/testqueue.lua @@ -0,0 +1,10 @@ +local skynet = require "skynet" +local mqueue = require "mqueue" + +skynet.start(function() + local pingqueue = skynet.newservice "pingqueue" + print(mqueue.call(pingqueue, "A")) + print(mqueue.call(pingqueue, "B")) + print(mqueue.call(pingqueue, "C")) + skynet.exit() +end) diff --git a/skynet-src/skynet.h b/skynet-src/skynet.h index 83831f90..7be4d873 100644 --- a/skynet-src/skynet.h +++ b/skynet-src/skynet.h @@ -13,8 +13,8 @@ #define PTYPE_SOCKET 6 // don't use these id #define PTYPE_RESERVED_0 7 -#define PTYPE_RESERVED_1 8 -// read lualib/skynet.lua +// read lualib/skynet.lua lualib/mqueue.lua +#define PTYPE_RESERVED_QUEUE 8 #define PTYPE_RESERVED_DEBUG 9 #define PTYPE_RESERVED_LUA 10