Compare commits

...

48 Commits

Author SHA1 Message Date
Cloud Wu
8214c53891 release alpha7 2015-06-08 10:31:00 +08:00
Cloud Wu
5e642b15a4 Merge branch 'master' of github.com:cloudwu/skynet 2015-06-07 12:39:33 +08:00
Cloud Wu
ac093c7142 fix memory leak, see Issue #290 2015-06-07 12:38:55 +08:00
云风
01dcc45827 Merge pull request #289 from harrywong/typo-patch
typo in type of variable
2015-06-05 23:02:16 +08:00
Harry
d3ac522dc6 typo in type of variable 2015-06-05 19:04:43 +08:00
Cloud Wu
55d0d57d5c skynet.fork returns the coroutine, see Issue #287 2015-06-04 16:06:21 +08:00
Cloud Wu
1fec2e6063 sep size may greater than buffer node, See Issue #286 2015-06-03 09:42:49 +08:00
Cloud Wu
cc4756de35 fix Issue #285 2015-06-02 11:52:55 +08:00
Cloud Wu
7fb109dbb0 See Issue #284 2015-06-01 10:57:30 +08:00
Cloud Wu
7f578be649 bugfix: Issue #283 2015-06-01 10:42:29 +08:00
Cloud Wu
69946d75c5 add config.logservice for user defined log service 2015-05-30 22:17:53 +08:00
Cloud Wu
24f6994b50 close uncomplete when socket disconnect, see Issue #280 2015-05-29 21:11:02 +08:00
Cloud Wu
07dbfd8651 add underscore 2015-05-28 16:18:56 +08:00
Cloud Wu
ca50a5f518 dns support underscore 2015-05-28 16:10:19 +08:00
Cloud Wu
b8c54cbac2 skynet.kill move into skynet.manager 2015-05-27 21:17:43 +08:00
Cloud Wu
3a2c43e9e0 set nodelay in clusterd, Issue #278 2015-05-27 14:54:26 +08:00
Cloud Wu
205824ab12 use LUA_API instead of extern 2015-05-26 16:05:31 +08:00
云风
4ac047b274 Merge pull request #277 from xiyanxiyan10/master
调整格式
2015-05-23 17:02:37 +08:00
xiyanxiyan10
987a90af8b 调整格式 2015-05-23 16:35:15 +08:00
Cloud Wu
250531c9a7 bugfix: sproto.default 2015-05-22 20:51:00 +08:00
Cloud Wu
7b471ae27d snax service support cluster 2015-05-20 16:13:11 +08:00
云风
4f93054f3a Merge pull request #275 from ximenpo/master
console服务增加启动snax服务功能
2015-05-20 15:43:55 +08:00
ximenpo
65fad89794 console服务增加启动snax服务功能
输入:
svc  -> 用snlua启动svc服务
snax svc -> 用snax启动svc服务
2015-05-20 15:14:28 +08:00
云风
286c0f2c23 Merge pull request #273 from flashjay/patch-4
Update lua-crypt.c
2015-05-19 11:41:14 +08:00
HuaYang Huang
de57064b96 Update lua-crypt.c
添加 static 避免可能导致的链接冲突
2015-05-18 23:05:24 +08:00
Cloud Wu
6c43e90708 release alpha6 2015-05-18 12:02:08 +08:00
云风
1f25b79722 Merge pull request #272 from antsmallant/patch-1
Update inject.lua
2015-05-15 17:04:27 +08:00
Cloud Wu
6a9080157f need require skynet.manager 2015-05-15 16:24:01 +08:00
antsmallant
5cc2ef3ac8 Update inject.lua
免得变量i污染全局空间
2015-05-15 15:59:47 +08:00
Cloud Wu
3baeb62b0b move some api from skynet.lua to skynet/manager.lua 2015-05-13 11:04:25 +08:00
Cloud Wu
9e27f59033 remove task overload warning 2015-05-13 09:50:53 +08:00
Cloud Wu
82fa2f979c skynet.exit will call every unresponse handle 2015-05-12 23:15:03 +08:00
Cloud Wu
07d7324332 Merge branch 'master' of github.com:cloudwu/skynet 2015-05-12 22:53:16 +08:00
Cloud Wu
d92adda2c9 If skynet.wakeup a skynet.call, raise an error 2015-05-12 22:53:03 +08:00
Cloud Wu
0587400e2d update sproto: handle empty protocol 2015-05-12 10:20:16 +08:00
Cloud Wu
35c8c63793 update sproto, bugfix sproto_dump 2015-05-11 21:39:56 +08:00
Cloud Wu
c7f5145e9e update sproto , add new api sproto:default 2015-05-11 18:20:17 +08:00
Cloud Wu
9452169f0a remove unused comment 2015-05-07 18:17:46 +08:00
云风
a15f43e0f8 Merge pull request #267 from gaopan461/master
fix:udp and udp_send
2015-05-07 18:18:36 +08:00
gaopan
0eb4754b31 @fix udp_address, in big-endian machine maybe error 2015-05-07 16:52:48 +08:00
gaopan
f987ff8199 fix:udp and udp_send 2015-05-07 15:25:47 +08:00
云风
64a5d1ca24 Merge pull request #266 from xiyanxiyan10/master
错误注释更正
2015-05-06 23:27:13 +08:00
xiyanxiyan10
fa17081012 错误注释更正 2015-05-06 23:13:06 +08:00
Cloud Wu
856cb0737d Merge branch 'master' of github.com:cloudwu/skynet 2015-05-02 20:47:29 +08:00
Cloud Wu
f61e3f46e8 fix issue #265 2015-05-02 20:47:15 +08:00
云风
6b6b943b73 Merge pull request #262 from flashjay/patch-3
fix: httpc.get
2015-04-29 22:19:31 +08:00
HuaYang Huang
75d9b07158 fix: httpc.get
:-(
此错误会导致httpc.get无法找到正确的虚拟主机
2015-04-29 20:18:13 +08:00
Cloud Wu
8682f7f82f response error when return/response package is too large 2015-04-27 16:23:12 +08:00
49 changed files with 395 additions and 155 deletions

View File

@@ -460,7 +460,7 @@ struct lua_Debug {
/* Add by skynet */ /* Add by skynet */
extern lua_State * skynet_sig_L; LUA_API lua_State * skynet_sig_L;
LUA_API void (lua_checksig_)(lua_State *L); LUA_API void (lua_checksig_)(lua_State *L);
#define lua_checksig(L) if (skynet_sig_L) { lua_checksig_(L); } #define lua_checksig(L) if (skynet_sig_L) { lua_checksig_(L); }

View File

@@ -1,3 +1,28 @@
v1.0.0-alpha7 (2015-6-8)
-----------
* console support launch snax service
* Add cluster.snax
* Add nodelay in clusterd
* Merge sproto bugfix patch
* Move some skynet api into skynet.manager
* DNS support underscore
* Add logservice in config file for user defined log service
* skynet.fork returns coroutine
* Fix a few of bugs , see the commits log
v1.0.0-alpha6 (2015-5-18)
-----------
* bugfix: httpc.get
* bugfix: seri lib stack overflow
* bugfix: udp send
* bugfix: udp address
* bugfix: sproto dump
* add: sproto default
* improve: skynet.wakeup (can wakeup skynet.call by raise an error)
* improve: skynet.exit (raise error when uncall response)
* remove: task overload warning
* move: some skynet api move into skynet.manager
v1.0.0-alpha5 (2015-4-27) v1.0.0-alpha5 (2015-4-27)
----------- -----------
* merge lua 5.3 offical bugfix * merge lua 5.3 offical bugfix

View File

@@ -1,3 +1,4 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager" -- import skynet.abort
skynet.abort() skynet.abort()

View File

@@ -4,6 +4,7 @@ local socket = require "socket"
local sproto = require "sproto" local sproto = require "sproto"
local sprotoloader = require "sprotoloader" local sprotoloader = require "sprotoloader"
local WATCHDOG
local host local host
local send_request local send_request
@@ -26,6 +27,10 @@ function REQUEST:handshake()
return { msg = "Welcome to skynet, I will send heartbeat every 5 sec." } return { msg = "Welcome to skynet, I will send heartbeat every 5 sec." }
end end
function REQUEST:quit()
skynet.call(WATCHDOG, "lua", "close", client_fd)
end
local function request(name, args, response) local function request(name, args, response)
local f = assert(REQUEST[name]) local f = assert(REQUEST[name])
local r = f(args) local r = f(args)
@@ -62,7 +67,10 @@ skynet.register_protocol {
end end
} }
function CMD.start(gate, fd) function CMD.start(conf)
local fd = conf.client
local gate = conf.gate
WATCHDOG = conf.watchdog
-- slot 1,2 set at main.lua -- slot 1,2 set at main.lua
host = sprotoloader.load(1):host "package" host = sprotoloader.load(1):host "package"
send_request = host:attach(sprotoloader.load(2)) send_request = host:attach(sprotoloader.load(2))
@@ -77,6 +85,11 @@ function CMD.start(gate, fd)
skynet.call(gate, "lua", "forward", fd) skynet.call(gate, "lua", "forward", fd)
end end
function CMD.disconnect()
-- todo: do something before exit
skynet.exit()
end
skynet.start(function() skynet.start(function()
skynet.dispatch("lua", function(_,_, command, ...) skynet.dispatch("lua", function(_,_, command, ...)
local f = CMD[command] local f = CMD[command]

View File

@@ -104,7 +104,11 @@ while true do
dispatch_package() dispatch_package()
local cmd = socket.readstdin() local cmd = socket.readstdin()
if cmd then if cmd then
send_request("get", { what = cmd }) if cmd == "quit" then
send_request("quit")
else
send_request("get", { what = cmd })
end
else else
socket.usleep(100) socket.usleep(100)
end end

View File

@@ -1,5 +1,7 @@
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"
skynet.start(function() skynet.start(function()
local sdb = skynet.newservice("simpledb") local sdb = skynet.newservice("simpledb")
@@ -10,4 +12,6 @@ skynet.start(function()
print(skynet.call(".simpledb", "lua", "GET", "b")) print(skynet.call(".simpledb", "lua", "GET", "b"))
cluster.open "db" cluster.open "db"
cluster.open "db2" cluster.open "db2"
-- unique snax service
snax.uniqueservice "pingserver"
end) end)

View File

@@ -6,4 +6,8 @@ skynet.start(function()
print(skynet.call(proxy, "lua", "GET", "a")) print(skynet.call(proxy, "lua", "GET", "a"))
print(cluster.call("db", ".simpledb", "GET", "a")) print(cluster.call("db", ".simpledb", "GET", "a"))
print(cluster.call("db2", ".simpledb", "GET", "b")) print(cluster.call("db2", ".simpledb", "GET", "b"))
-- test snax service
local pingserver = cluster.snax("db", "pingserver")
print(pingserver.req.ping "hello")
end) end)

View File

@@ -7,3 +7,4 @@ luaservice = "./service/?.lua;./test/?.lua;./examples/?.lua"
lualoader = "lualib/loader.lua" lualoader = "lualib/loader.lua"
cpath = "./cservice/?.so" cpath = "./cservice/?.so"
cluster = "./examples/clustername.lua" cluster = "./examples/clustername.lua"
snax = "./test/?.lua"

View File

@@ -7,3 +7,4 @@ luaservice = "./service/?.lua;./test/?.lua;./examples/?.lua"
lualoader = "lualib/loader.lua" lualoader = "lualib/loader.lua"
cpath = "./cservice/?.so" cpath = "./cservice/?.so"
cluster = "./examples/clustername.lua" cluster = "./examples/clustername.lua"
snax = "./test/?.lua"

View File

@@ -1,4 +1,5 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager" -- import skynet.register
skynet.start(function() skynet.start(function()
skynet.dispatch("lua", function(session, address, ...) skynet.dispatch("lua", function(session, address, ...)

View File

@@ -1,5 +1,6 @@
local skynet = require "skynet" local skynet = require "skynet"
local harbor = require "skynet.harbor" local harbor = require "skynet.harbor"
require "skynet.manager" -- import skynet.monitor
local function monitor_master() local function monitor_master()
harbor.linkmaster() harbor.linkmaster()

View File

@@ -30,6 +30,8 @@ set 3 {
} }
} }
quit 4 {}
]] ]]
proto.s2c = sprotoparser.parse [[ proto.s2c = sprotoparser.parse [[

View File

@@ -1,4 +1,5 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager" -- import skynet.register
local db = {} local db = {}
local command = {} local command = {}

View File

@@ -9,14 +9,16 @@ local agent = {}
function SOCKET.open(fd, addr) function SOCKET.open(fd, addr)
skynet.error("New client from : " .. addr) skynet.error("New client from : " .. addr)
agent[fd] = skynet.newservice("agent") agent[fd] = skynet.newservice("agent")
skynet.call(agent[fd], "lua", "start", gate, fd) skynet.call(agent[fd], "lua", "start", { gate = gate, client = fd, watchdog = skynet.self() })
end end
local function close_agent(fd) local function close_agent(fd)
local a = agent[fd] local a = agent[fd]
agent[fd] = nil
if a then if a then
skynet.kill(a) skynet.call(gate, "lua", "kick", fd)
agent[fd] = nil -- disconnect never return
skynet.send(a, "lua", "disconnect")
end end
end end
@@ -37,6 +39,10 @@ function CMD.start(conf)
skynet.call(gate, "lua", "open" , conf) skynet.call(gate, "lua", "open" , conf)
end end
function CMD.close(fd)
close_agent(fd)
end
skynet.start(function() skynet.start(function()
skynet.dispatch("lua", function(session, source, cmd, subcmd, ...) skynet.dispatch("lua", function(session, source, cmd, subcmd, ...)
if cmd == "socket" then if cmd == "socket" then

View File

@@ -10,7 +10,7 @@
/* the eight DES S-boxes */ /* the eight DES S-boxes */
uint32_t SB1[64] = { static uint32_t SB1[64] = {
0x01010400, 0x00000000, 0x00010000, 0x01010404, 0x01010400, 0x00000000, 0x00010000, 0x01010404,
0x01010004, 0x00010404, 0x00000004, 0x00010000, 0x01010004, 0x00010404, 0x00000004, 0x00010000,
0x00000400, 0x01010400, 0x01010404, 0x00000400, 0x00000400, 0x01010400, 0x01010404, 0x00000400,

View File

@@ -59,6 +59,7 @@ channel_release(struct channel *c) {
free(p); free(p);
p = next; p = next;
} }
free(c);
return NULL; return NULL;
} }

View File

@@ -210,6 +210,16 @@ push_more(lua_State *L, int fd, uint8_t *buffer, int size) {
} }
} }
static void
close_uncomplete(lua_State *L, int fd) {
struct queue *q = lua_touserdata(L,1);
struct uncomplete * uc = find_uncomplete(q, fd);
if (uc) {
skynet_free(uc->pack.buffer);
skynet_free(uc);
}
}
static int static int
filter_data_(lua_State *L, int fd, uint8_t * buffer, int size) { filter_data_(lua_State *L, int fd, uint8_t * buffer, int size) {
struct queue *q = lua_touserdata(L,1); struct queue *q = lua_touserdata(L,1);
@@ -343,6 +353,8 @@ lfilter(lua_State *L) {
// ignore listen fd connect // ignore listen fd connect
return 1; return 1;
case SKYNET_SOCKET_TYPE_CLOSE: case SKYNET_SOCKET_TYPE_CLOSE:
// no more data in fd (message->id)
close_uncomplete(L, message->id);
lua_pushvalue(L, lua_upvalueindex(TYPE_CLOSE)); lua_pushvalue(L, lua_upvalueindex(TYPE_CLOSE));
lua_pushinteger(L, message->id); lua_pushinteger(L, message->id);
return 3; return 3;
@@ -353,6 +365,8 @@ lfilter(lua_State *L) {
pushstring(L, buffer, size); pushstring(L, buffer, size);
return 4; return 4;
case SKYNET_SOCKET_TYPE_ERROR: case SKYNET_SOCKET_TYPE_ERROR:
// no more data in fd (message->id)
close_uncomplete(L, message->id);
lua_pushvalue(L, lua_upvalueindex(TYPE_ERROR)); lua_pushvalue(L, lua_upvalueindex(TYPE_ERROR));
lua_pushinteger(L, message->id); lua_pushinteger(L, message->id);
pushstring(L, buffer, size); pushstring(L, buffer, size);

View File

@@ -256,6 +256,7 @@ wb_table_hash(lua_State *L, struct write_block * wb, int index, int depth, int a
static void static void
wb_table(lua_State *L, struct write_block *wb, int index, int depth) { wb_table(lua_State *L, struct write_block *wb, int index, int depth) {
luaL_checkstack(L, LUA_MINSTACK, NULL);
if (index < 0) { if (index < 0) {
index = lua_gettop(L) + index + 1; index = lua_gettop(L) + index + 1;
} }
@@ -416,6 +417,7 @@ unpack_table(lua_State *L, struct read_block *rb, int array_size) {
} }
array_size = get_integer(L,rb,cookie); array_size = get_integer(L,rb,cookie);
} }
luaL_checkstack(L,LUA_MINSTACK,NULL);
lua_createtable(L,array_size,0); lua_createtable(L,array_size,0);
int i; int i;
for (i=1;i<=array_size;i++) { for (i=1;i<=array_size;i++) {
@@ -549,8 +551,8 @@ _luaseri_unpack(lua_State *L) {
int i; int i;
for (i=0;;i++) { for (i=0;;i++) {
if (i%16==15) { if (i%8==7) {
lua_checkstack(L,i); luaL_checkstack(L,LUA_MINSTACK,NULL);
} }
uint8_t type = 0; uint8_t type = 0;
uint8_t *t = rb_read(&rb, sizeof(type)); uint8_t *t = rb_read(&rb, sizeof(type));

View File

@@ -190,7 +190,10 @@ pop_lstring(lua_State *L, struct socket_buffer *sb, int sz, int skip) {
} }
break; break;
} }
luaL_addlstring(&b, current->msg + sb->offset, (sz - skip < bytes) ? sz - skip : bytes); int real_sz = sz - skip;
if (real_sz > 0) {
luaL_addlstring(&b, current->msg + sb->offset, (real_sz < bytes) ? real_sz : bytes);
}
return_free_node(L,2,sb); return_free_node(L,2,sb);
sz-=bytes; sz-=bytes;
if (sz==0) if (sz==0)
@@ -534,12 +537,6 @@ lnodelay(lua_State *L) {
skynet_socket_nodelay(ctx,id); skynet_socket_nodelay(ctx,id);
return 0; return 0;
} }
/*
int skynet_socket_udp(struct skynet_context *ctx, const char * addr, int port);
int skynet_socket_udp_connect(struct skynet_context *ctx, int id, const char * addr, int port);
int skynet_socket_udp_send(struct skynet_context *ctx, int id, const char * address, const void *buffer, int sz);
const char * skynet_socket_udp_address(struct skynet_context *ctx, struct skynet_socket_message *, int *addrsz);
*/
static int static int
ludp(lua_State *L) { ludp(lua_State *L) {
@@ -599,7 +596,9 @@ static int
ludp_address(lua_State *L) { ludp_address(lua_State *L) {
size_t sz = 0; size_t sz = 0;
const uint8_t * addr = (const uint8_t *)luaL_checklstring(L, 1, &sz); const uint8_t * addr = (const uint8_t *)luaL_checklstring(L, 1, &sz);
int port = addr[1] * 256 + addr[2]; uint16_t port = 0;
memcpy(&port, addr+1, sizeof(uint16_t));
port = ntohs(port);
const void * src = addr+3; const void * src = addr+3;
char tmp[256]; char tmp[256];
int family; int family;

View File

@@ -573,6 +573,66 @@ lloadproto(lua_State *L) {
return 1; return 1;
} }
static int
encode_default(const struct sproto_arg *args) {
lua_State *L = args->ud;
lua_pushstring(L, args->tagname);
if (args->index > 0) {
lua_newtable(L);
} else {
switch(args->type) {
case SPROTO_TINTEGER:
lua_pushinteger(L, 0);
break;
case SPROTO_TBOOLEAN:
lua_pushboolean(L, 0);
break;
case SPROTO_TSTRING:
lua_pushliteral(L, "");
break;
case SPROTO_TSTRUCT:
lua_createtable(L, 0, 1);
lua_pushstring(L, sproto_name(args->subtype));
lua_setfield(L, -2, "__type");
break;
}
}
lua_rawset(L, -3);
return 0;
}
/*
lightuserdata sproto_type
return default table
*/
static int
ldefault(lua_State *L) {
int ret;
// 64 is always enough for dummy buffer, except the type has many fields ( > 27).
char dummy[64];
struct sproto_type * st = lua_touserdata(L, 1);
if (st == NULL) {
return luaL_argerror(L, 1, "Need a sproto_type object");
}
lua_newtable(L);
ret = sproto_encode(st, dummy, sizeof(dummy), encode_default, L);
if (ret<0) {
// try again
int sz = sizeof(dummy) * 2;
void * tmp = lua_newuserdata(L, sz);
lua_insert(L, -2);
for (;;) {
ret = sproto_encode(st, tmp, sz, encode_default, L);
if (ret >= 0)
break;
sz *= 2;
tmp = lua_newuserdata(L, sz);
lua_replace(L, -3);
}
}
return 1;
}
int int
luaopen_sproto_core(lua_State *L) { luaopen_sproto_core(lua_State *L) {
#ifdef luaL_checkversion #ifdef luaL_checkversion
@@ -587,6 +647,7 @@ luaopen_sproto_core(lua_State *L) {
{ "protocol", lprotocol }, { "protocol", lprotocol },
{ "loadproto", lloadproto }, { "loadproto", lloadproto },
{ "saveproto", lsaveproto }, { "saveproto", lsaveproto },
{ "default", ldefault },
{ NULL, NULL }, { NULL, NULL },
}; };
luaL_newlib(L,l); luaL_newlib(L,l);

View File

@@ -500,7 +500,11 @@ sproto_dump(struct sproto *s) {
printf("=== %d protocol ===\n", s->protocol_n); printf("=== %d protocol ===\n", s->protocol_n);
for (i=0;i<s->protocol_n;i++) { for (i=0;i<s->protocol_n;i++) {
struct protocol *p = &s->proto[i]; struct protocol *p = &s->proto[i];
printf("\t%s (%d) request:%s", p->name, p->tag, p->p[SPROTO_REQUEST]->name); if (p->p[SPROTO_REQUEST]) {
printf("\t%s (%d) request:%s", p->name, p->tag, p->p[SPROTO_REQUEST]->name);
} else {
printf("\t%s (%d) request:(null)", p->name, p->tag);
}
if (p->p[SPROTO_RESPONSE]) { if (p->p[SPROTO_RESPONSE]) {
printf(" response:%s", p->p[SPROTO_RESPONSE]->name); printf(" response:%s", p->p[SPROTO_RESPONSE]->name);
} }

View File

@@ -24,6 +24,15 @@ function cluster.proxy(node, name)
return skynet.call(clusterd, "lua", "proxy", node, name) return skynet.call(clusterd, "lua", "proxy", node, name)
end end
function cluster.snax(node, name, address)
local snax = require "snax"
if not address then
address = cluster.call(node, ".service", "QUERY", "snaxd" , name)
end
local handle = skynet.call(clusterd, "lua", "proxy", node, address)
return snax.bind(handle, name)
end
skynet.init(function() skynet.init(function()
clusterd = skynet.uniqueservice("clusterd") clusterd = skynet.uniqueservice("clusterd")
end) end)

View File

@@ -88,10 +88,10 @@ local function verify_domain_name(name)
if #name > MAX_DOMAIN_LEN then if #name > MAX_DOMAIN_LEN then
return false return false
end end
if not name:match("^[%l%d-%.]+$") then if not name:match("^[_%l%d%-%.]+$") then
return false return false
end end
for w in name:gmatch("([%w-]+)%.?") do for w in name:gmatch("([_%w%-]+)%.?") do
if #w > MAX_LABEL_LEN then if #w > MAX_LABEL_LEN then
return false return false
end end
@@ -113,7 +113,7 @@ end
local function pack_question(name, qtype, qclass) local function pack_question(name, qtype, qclass)
local labels = {} local labels = {}
for w in name:gmatch("([%w-]+)%.?") do for w in name:gmatch("([_%w%-]+)%.?") do
table.insert(labels, string.pack("s1",w)) table.insert(labels, string.pack("s1",w))
end end
table.insert(labels, '\0') table.insert(labels, '\0')
@@ -282,7 +282,7 @@ function dns.resolve(name, ipv6)
qdcount = 1, qdcount = 1,
} }
local req = pack_header(question_header) .. pack_question(name, qtype, QCLASS.IN) local req = pack_header(question_header) .. pack_question(name, qtype, QCLASS.IN)
assert(dns_server, "Call dns.server fist") assert(dns_server, "Call dns.server first")
socket.write(dns_server, req) socket.write(dns_server, req)
return suspend(question_header.tid, name, qtype) return suspend(question_header.tid, name, qtype)
end end

View File

@@ -11,23 +11,21 @@ local function request(fd, method, host, url, recvheader, header, content)
local write = socket.writefunc(fd) local write = socket.writefunc(fd)
local header_content = "" local header_content = ""
if header then if header then
if not header.host then
header.host = host
end
for k,v in pairs(header) do for k,v in pairs(header) do
header_content = string.format("%s%s:%s\r\n", header_content, k, v) header_content = string.format("%s%s:%s\r\n", header_content, k, v)
end end
if header.host then
host = ""
else
host = string.format("host:%s\r\n", host)
end
else else
host = string.format("host:%s\r\n",host) header_content = string.format("host:%s\r\n",host)
end end
if content then if content then
local data = string.format("%s %s HTTP/1.1\r\n%scontent-length:%d\r\n%s\r\n%s", method, url, host, #content, header_content, content) local data = string.format("%s %s HTTP/1.1\r\n%scontent-length:%d\r\n\r\n%s", method, url, header_content, #content, content)
write(data) write(data)
else else
local request_header = string.format("%s %s HTTP/1.1\r\nhost:%s\r\ncontent-length:0\r\n%s\r\n", method, url, host, header_content) local request_header = string.format("%s %s HTTP/1.1\r\n%scontent-length:0\r\n\r\n", method, url, header_content)
write(request_header) write(request_header)
end end

View File

@@ -44,6 +44,7 @@ local session_id_coroutine = {}
local session_coroutine_id = {} local session_coroutine_id = {}
local session_coroutine_address = {} local session_coroutine_address = {}
local session_response = {} local session_response = {}
local unresponse = {}
local wakeup_session = {} local wakeup_session = {}
local sleep_session = {} local sleep_session = {}
@@ -96,7 +97,6 @@ end
local coroutine_pool = {} local coroutine_pool = {}
local coroutine_yield = coroutine.yield local coroutine_yield = coroutine.yield
local coroutine_count = 0
local function co_create(f) local function co_create(f)
local co = table.remove(coroutine_pool) local co = table.remove(coroutine_pool)
@@ -110,11 +110,6 @@ local function co_create(f)
f(coroutine_yield()) f(coroutine_yield())
end end
end) end)
coroutine_count = coroutine_count + 1
if coroutine_count > 1024 then
skynet.error("May overload, create 1024 task")
coroutine_count = 0
end
else else
coroutine.resume(co, f) coroutine.resume(co, f)
end end
@@ -128,7 +123,7 @@ local function dispatch_wakeup()
local session = sleep_session[co] local session = sleep_session[co]
if session then if session then
session_id_coroutine[session] = "BREAK" session_id_coroutine[session] = "BREAK"
return suspend(co, coroutine.resume(co, true, "BREAK")) return suspend(co, coroutine.resume(co, false, "BREAK"))
end end
end end
end end
@@ -175,7 +170,11 @@ function suspend(co, result, command, param, size)
local ret local ret
if not dead_service[co_address] then if not dead_service[co_address] then
ret = c.send(co_address, skynet.PTYPE_RESPONSE, co_session, param, size) ~= nil ret = c.send(co_address, skynet.PTYPE_RESPONSE, co_session, param, size) ~= nil
elseif size == nil then if not ret then
-- If the package is too large, returns nil. so we should report error back
c.send(co_address, skynet.PTYPE_ERROR, co_session, "")
end
elseif size ~= nil then
c.trash(param, size) c.trash(param, size)
ret = false ret = false
end end
@@ -209,6 +208,10 @@ function suspend(co, result, command, param, size)
if not dead_service[co_address] then if not dead_service[co_address] then
if ok then if ok then
ret = c.send(co_address, skynet.PTYPE_RESPONSE, co_session, f(...)) ~= nil ret = c.send(co_address, skynet.PTYPE_RESPONSE, co_session, f(...)) ~= nil
if not ret then
-- If the package is too large, returns false. so we should report error back
c.send(co_address, skynet.PTYPE_ERROR, co_session, "")
end
else else
ret = c.send(co_address, skynet.PTYPE_ERROR, co_session, "") ~= nil ret = c.send(co_address, skynet.PTYPE_ERROR, co_session, "") ~= nil
end end
@@ -216,11 +219,13 @@ function suspend(co, result, command, param, size)
ret = false ret = false
end end
release_watching(co_address) release_watching(co_address)
unresponse[response] = nil
f = nil f = nil
return ret return ret
end end
watching_service[co_address] = watching_service[co_address] + 1 watching_service[co_address] = watching_service[co_address] + 1
session_response[co] = response session_response[co] = response
unresponse[response] = true
return suspend(co, coroutine.resume(co, response)) return suspend(co, coroutine.resume(co, response))
elseif command == "EXIT" then elseif command == "EXIT" then
-- coroutine exit -- coroutine exit
@@ -257,9 +262,13 @@ function skynet.sleep(ti)
session = tonumber(session) session = tonumber(session)
local succ, ret = coroutine_yield("SLEEP", session) local succ, ret = coroutine_yield("SLEEP", session)
sleep_session[coroutine.running()] = nil sleep_session[coroutine.running()] = nil
assert(succ, ret) if succ then
return
end
if ret == "BREAK" then if ret == "BREAK" then
return "BREAK" return "BREAK"
else
error(ret)
end end
end end
@@ -269,41 +278,12 @@ end
function skynet.wait() function skynet.wait()
local session = c.genid() local session = c.genid()
coroutine_yield("SLEEP", session) local ret, msg = coroutine_yield("SLEEP", session)
local co = coroutine.running() local co = coroutine.running()
sleep_session[co] = nil sleep_session[co] = nil
session_id_coroutine[session] = nil session_id_coroutine[session] = nil
end 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)
if not globalname(name) then
c.command("REG", name)
end
end
function skynet.name(name, handle)
if not globalname(name, handle) then
c.command("NAME", name .. " " .. skynet.address(handle))
end
end
local self_handle local self_handle
function skynet.self() function skynet.self()
if self_handle then if self_handle then
@@ -320,13 +300,6 @@ function skynet.localname(name)
end end
end end
function skynet.launch(...)
local addr = c.command("LAUNCH", table.concat({...}," "))
if addr then
return string_to_handle(addr)
end
end
function skynet.now() function skynet.now()
return tonumber(c.command("NOW")) return tonumber(c.command("NOW"))
end end
@@ -349,6 +322,9 @@ function skynet.exit()
c.redirect(address, 0, skynet.PTYPE_ERROR, session, "") c.redirect(address, 0, skynet.PTYPE_ERROR, session, "")
end end
end end
for resp in pairs(unresponse) do
resp(false)
end
-- report the sources I call but haven't return -- report the sources I call but haven't return
local tmp = {} local tmp = {}
for session, address in pairs(watching_session) do for session, address in pairs(watching_session) do
@@ -362,14 +338,6 @@ function skynet.exit()
coroutine_yield "QUIT" coroutine_yield "QUIT"
end end
function skynet.kill(name)
if type(name) == "number" then
skynet.send(".launcher","lua","REMOVE",name, true)
name = skynet.address(name)
end
c.command("KILL",name)
end
function skynet.getenv(key) function skynet.getenv(key)
local ret = c.command("GETENV",key) local ret = c.command("GETENV",key)
if ret == "" then if ret == "" then
@@ -405,7 +373,7 @@ local function yield_call(service, session)
local succ, msg, sz = coroutine_yield("CALL", session) local succ, msg, sz = coroutine_yield("CALL", session)
watching_session[session] = nil watching_session[session] = nil
if not succ then if not succ then
error(debug.traceback()) error "call failed"
end end
return msg,sz return msg,sz
end end
@@ -464,7 +432,7 @@ function skynet.dispatch_unknown_request(unknown)
end end
local function unknown_response(session, address, msg, sz) local function unknown_response(session, address, msg, sz)
skynet.error(string.format("Response message :" , c.tostring(msg,sz))) skynet.error(string.format("Response message : %s" , c.tostring(msg,sz)))
error(string.format("Unknown session : %d from %x", session, address)) error(string.format("Unknown session : %d from %x", session, address))
end end
@@ -482,6 +450,7 @@ function skynet.fork(func,...)
func(tunpack(args)) func(tunpack(args))
end) end)
table.insert(fork_queue, co) table.insert(fork_queue, co)
return co
end end
local function raw_dispatch_message(prototype, msg, sz, session, source, ...) local function raw_dispatch_message(prototype, msg, sz, session, source, ...)
@@ -511,12 +480,12 @@ local function raw_dispatch_message(prototype, msg, sz, session, source, ...)
session_coroutine_address[co] = source session_coroutine_address[co] = source
suspend(co, coroutine.resume(co, session,source, p.unpack(msg,sz, ...))) suspend(co, coroutine.resume(co, session,source, p.unpack(msg,sz, ...)))
else else
unknown_request(session, source, msg, sz, proto[prototype]) unknown_request(session, source, msg, sz, proto[prototype].name)
end end
end end
end end
local function dispatch_message(...) function skynet.dispatch_message(...)
local succ, err = pcall(raw_dispatch_message,...) local succ, err = pcall(raw_dispatch_message,...)
while true do while true do
local key,co = next(fork_queue) local key,co = next(fork_queue)
@@ -638,7 +607,7 @@ function skynet.pcall(start)
return xpcall(init_template, debug.traceback, start) return xpcall(init_template, debug.traceback, start)
end end
local function init_service(start) function skynet.init_service(start)
local ok, err = skynet.pcall(start) local ok, err = skynet.pcall(start)
if not ok then if not ok then
skynet.error("init service failed: " .. tostring(err)) skynet.error("init service failed: " .. tostring(err))
@@ -650,33 +619,9 @@ local function init_service(start)
end end
function skynet.start(start_func) function skynet.start(start_func)
c.callback(dispatch_message) c.callback(skynet.dispatch_message)
skynet.timeout(0, function() skynet.timeout(0, function()
init_service(start_func) skynet.init_service(start_func)
end)
end
function skynet.filter(f ,start_func)
c.callback(function(...)
dispatch_message(f(...))
end)
skynet.timeout(0, function()
init_service(start_func)
end)
end
function skynet.forward_type(map, start_func)
c.callback(function(ptype, msg, sz, ...)
local prototype = map[ptype]
if prototype then
dispatch_message(prototype, msg, sz, ...)
else
dispatch_message(ptype, msg, sz, ...)
c.trash(msg, sz)
end
end, true)
skynet.timeout(0, function()
init_service(start_func)
end) end)
end end
@@ -684,22 +629,6 @@ function skynet.endless()
return c.command("ENDLESS")~=nil return c.command("ENDLESS")~=nil
end end
function skynet.abort()
c.command("ABORT")
end
function skynet.monitor(service, query)
local monitor
if query then
monitor = skynet.queryservice(true, service)
else
monitor = skynet.uniqueservice(true, service)
end
assert(monitor, "Monitor launch failed")
c.command("MONITOR", string.format(":%08x", monitor))
return monitor
end
function skynet.mqlen() function skynet.mqlen()
return tonumber(c.command "MQLEN") return tonumber(c.command "MQLEN")
end end
@@ -726,7 +655,7 @@ end
-- Inject internal debug framework -- Inject internal debug framework
local debug = require "skynet.debug" local debug = require "skynet.debug"
debug(skynet, { debug(skynet, {
dispatch = dispatch_message, dispatch = skynet.dispatch_message,
clear = clear_pool, clear = clear_pool,
suspend = suspend, suspend = suspend,
}) })

View File

@@ -1,5 +1,5 @@
local function getupvaluetable(u, func, unique) local function getupvaluetable(u, func, unique)
i = 1 local i = 1
while true do while true do
local name, value = debug.getupvalue(func, i) local name, value = debug.getupvalue(func, i)
if name == nil then if name == nil then
@@ -44,7 +44,7 @@ return function(skynet, source, filename , ...)
if proto then if proto then
for k,v in pairs(proto) do for k,v in pairs(proto) do
local name, dispatch = v.name, v.dispatch local name, dispatch = v.name, v.dispatch
if name and dispatch then if name and dispatch and not p[name] then
local pp = {} local pp = {}
p[name] = pp p[name] = pp
getupvaluetable(pp, dispatch, unique) getupvaluetable(pp, dispatch, unique)

90
lualib/skynet/manager.lua Normal file
View File

@@ -0,0 +1,90 @@
local skynet = require "skynet"
local c = require "skynet.core"
function skynet.launch(...)
local addr = c.command("LAUNCH", table.concat({...}," "))
if addr then
return tonumber("0x" .. string.sub(addr , 2))
end
end
function skynet.kill(name)
if type(name) == "number" then
skynet.send(".launcher","lua","REMOVE",name, true)
name = skynet.address(name)
end
c.command("KILL",name)
end
function skynet.abort()
c.command("ABORT")
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)
if not globalname(name) then
c.command("REG", name)
end
end
function skynet.name(name, handle)
if not globalname(name, handle) then
c.command("NAME", name .. " " .. skynet.address(handle))
end
end
local dispatch_message = skynet.dispatch_message
function skynet.forward_type(map, start_func)
c.callback(function(ptype, msg, sz, ...)
local prototype = map[ptype]
if prototype then
dispatch_message(prototype, msg, sz, ...)
else
dispatch_message(ptype, msg, sz, ...)
c.trash(msg, sz)
end
end, true)
skynet.timeout(0, function()
skynet.init_service(start_func)
end)
end
function skynet.filter(f ,start_func)
c.callback(function(...)
dispatch_message(f(...))
end)
skynet.timeout(0, function()
skynet.init_service(start_func)
end)
end
function skynet.monitor(service, query)
local monitor
if query then
monitor = skynet.queryservice(true, service)
else
monitor = skynet.uniqueservice(true, service)
end
assert(monitor, "Monitor launch failed")
c.command("MONITOR", string.format(":%08x", monitor))
return monitor
end
return skynet

View File

@@ -64,7 +64,7 @@ return function (name , G, loader)
local pattern local pattern
do do
local path = skynet.getenv "snax" local path = assert(skynet.getenv "snax" , "please set snax in config file")
local errlist = {} local errlist = {}

View File

@@ -1,4 +1,5 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager"
local socket = require "socket" local socket = require "socket"
local crypt = require "crypt" local crypt = require "crypt"
local table = table local table = table

View File

@@ -101,27 +101,64 @@ end
function sproto:request_encode(protoname, tbl) function sproto:request_encode(protoname, tbl)
local p = queryproto(self, protoname) local p = queryproto(self, protoname)
return core.encode(p.request,tbl) , p.tag local request = p.request
if request then
return core.encode(request,tbl) , p.tag
else
return "" , p.tag
end
end end
function sproto:response_encode(protoname, tbl) function sproto:response_encode(protoname, tbl)
local p = queryproto(self, protoname) local p = queryproto(self, protoname)
return core.encode(p.response,tbl) local response = p.response
if response then
return core.encode(response,tbl)
else
return ""
end
end end
function sproto:request_decode(protoname, ...) function sproto:request_decode(protoname, ...)
local p = queryproto(self, protoname) local p = queryproto(self, protoname)
return core.decode(p.request,...) , p.name local request = p.request
if request then
return core.decode(request,...) , p.name
else
return nil, p.name
end
end end
function sproto:response_decode(protoname, ...) function sproto:response_decode(protoname, ...)
local p = queryproto(self, protoname) local p = queryproto(self, protoname)
return core.decode(p.response,...) local response = p.response
if response then
return core.decode(response,...)
end
end end
sproto.pack = core.pack sproto.pack = core.pack
sproto.unpack = core.unpack sproto.unpack = core.unpack
function sproto:default(typename, type)
if type == nil then
return core.default(querytype(self, typename))
else
local p = queryproto(self, typename)
if type == "REQUEST" then
if p.request then
return core.default(p.request)
end
elseif type == "RESPONSE" then
if p.response then
return core.default(p.response)
end
else
error "Invalid type"
end
end
end
local header_tmp = {} local header_tmp = {}
local function gen_response(self, response, session) local function gen_response(self, response, session)

View File

@@ -340,9 +340,6 @@ local function packtype(name, t, alltypes)
end end
local function packproto(name, p, alltypes) local function packproto(name, p, alltypes)
-- if p.request == nil then
-- error(string.format("Protocol %s need request", name))
-- end
if p.request then if p.request then
local request = alltypes[p.request] local request = alltypes[p.request]
if request == nil then if request == nil then

View File

@@ -135,7 +135,7 @@ _ctrl(struct gate * g, const void * msg, int sz) {
skynet_socket_start(ctx, g->listen_id); skynet_socket_start(ctx, g->listen_id);
return; return;
} }
if (memcmp(command, "close", i) == 0) { if (memcmp(command, "close", i) == 0) {
if (g->listen_id >= 0) { if (g->listen_id >= 0) {
skynet_socket_close(ctx, g->listen_id); skynet_socket_close(ctx, g->listen_id);
g->listen_id = -1; g->listen_id = -1;

View File

@@ -365,6 +365,7 @@ dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
while ((m = pop_queue(queue)) != NULL) { while ((m = pop_queue(queue)) != NULL) {
m->header.destination |= (handle & HANDLE_MASK); m->header.destination |= (handle & HANDLE_MASK);
send_remote(context, fd, m->buffer, m->size, &m->header); send_remote(context, fd, m->buffer, m->size, &m->header);
skynet_free(m->buffer);
} }
} }
@@ -381,6 +382,7 @@ dispatch_queue(struct harbor *h, int id) {
struct harbor_msg * m; struct harbor_msg * m;
while ((m = pop_queue(queue)) != NULL) { while ((m = pop_queue(queue)) != NULL) {
send_remote(h->ctx, fd, m->buffer, m->size, &m->header); send_remote(h->ctx, fd, m->buffer, m->size, &m->header);
skynet_free(m->buffer);
} }
release_queue(queue); release_queue(queue);
s->queue = NULL; s->queue = NULL;

View File

@@ -1,5 +1,6 @@
local skynet = require "skynet" local skynet = require "skynet"
local harbor = require "skynet.harbor" local harbor = require "skynet.harbor"
require "skynet.manager" -- import skynet.launch, ...
skynet.start(function() skynet.start(function()
local standalone = skynet.getenv "standalone" local standalone = skynet.getenv "standalone"

View File

@@ -1,4 +1,5 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager" -- import skynet.launch, ...
local globalname = {} local globalname = {}
local queryname = {} local queryname = {}

View File

@@ -20,6 +20,7 @@ local function open_channel(t, key)
host = host, host = host,
port = tonumber(port), port = tonumber(port),
response = read_response, response = read_response,
nodelay = true,
} }
assert(c:connect(true)) assert(c:connect(true))
t[key] = c t[key] = c

View File

@@ -1,5 +1,6 @@
local skynet = require "skynet" local skynet = require "skynet"
local cluster = require "cluster" local cluster = require "cluster"
require "skynet.manager" -- inject skynet.forward_type
local node, address = ... local node, address = ...
@@ -10,6 +11,7 @@ skynet.register_protocol {
} }
local forward_map = { local forward_map = {
[skynet.PTYPE_SNAX] = skynet.PTYPE_SYSTEM,
[skynet.PTYPE_LUA] = skynet.PTYPE_SYSTEM, [skynet.PTYPE_LUA] = skynet.PTYPE_SYSTEM,
[skynet.PTYPE_RESPONSE] = skynet.PTYPE_RESPONSE, -- don't free response message [skynet.PTYPE_RESPONSE] = skynet.PTYPE_RESPONSE, -- don't free response message
} }

View File

@@ -1,13 +1,26 @@
local skynet = require "skynet" local skynet = require "skynet"
local snax = require "snax"
local socket = require "socket" local socket = require "socket"
local function split_cmdline(cmdline)
local split = {}
for i in string.gmatch(cmdline, "%S+") do
table.insert(split,i)
end
return split
end
local function console_main_loop() local function console_main_loop()
local stdin = socket.stdin() local stdin = socket.stdin()
socket.lock(stdin) socket.lock(stdin)
while true do while true do
local cmdline = socket.readline(stdin, "\n") local cmdline = socket.readline(stdin, "\n")
if cmdline ~= "" then local split = split_cmdline(cmdline)
pcall(skynet.newservice,cmdline) local command = split[1]
if command == "snax" then
pcall(snax.newservice, select(2, table.unpack(split)))
elseif cmdline ~= "" then
pcall(skynet.newservice, cmdline)
end end
end end
socket.unlock(stdin) socket.unlock(stdin)

View File

@@ -1,5 +1,6 @@
local skynet = require "skynet" local skynet = require "skynet"
local socket = require "socket" local socket = require "socket"
require "skynet.manager" -- import skynet.launch, ...
local table = table local table = table
local slaves = {} local slaves = {}

View File

@@ -171,7 +171,7 @@ end
local function adjust_address(address) local function adjust_address(address)
if address:sub(1,1) ~= ":" then if address:sub(1,1) ~= ":" then
address = tonumber("0x" .. address) | (skynet.harbor(skynet.self()) << 24) address = assert(tonumber("0x" .. address), "Need an address") | (skynet.harbor(skynet.self()) << 24)
end end
return address return address
end end

View File

@@ -1,5 +1,6 @@
local skynet = require "skynet" local skynet = require "skynet"
local core = require "skynet.core" local core = require "skynet.core"
require "skynet.manager" -- import manager apis
local string = string local string = string
local services = {} local services = {}

View File

@@ -1,4 +1,5 @@
local skynet = require "skynet" local skynet = require "skynet"
require "skynet.manager" -- import skynet.register
local snax = require "snax" local snax = require "snax"
local cmd = {} local cmd = {}

View File

@@ -54,7 +54,7 @@ end
skynet.start(function() skynet.start(function()
local init = false local init = false
skynet.dispatch("snax", function ( session , source , id, ...) local function dispatcher( session , source , id, ...)
local method = func[id] local method = func[id]
if method[2] == "system" then if method[2] == "system" then
@@ -84,5 +84,11 @@ skynet.start(function()
assert(init, "Init first") assert(init, "Init first")
timing(method, ...) timing(method, ...)
end end
end) end
skynet.dispatch("snax", dispatcher)
-- set lua dispatcher
function snax.enablecluster()
skynet.dispatch("lua", dispatcher)
end
end) end)

View File

@@ -8,6 +8,7 @@ struct skynet_config {
const char * module_path; const char * module_path;
const char * bootstrap; const char * bootstrap;
const char * logger; const char * logger;
const char * logservice;
}; };
#define THREAD_WORKER 0 #define THREAD_WORKER 0

View File

@@ -133,6 +133,7 @@ main(int argc, char *argv[]) {
config.bootstrap = optstring("bootstrap","snlua bootstrap"); config.bootstrap = optstring("bootstrap","snlua bootstrap");
config.daemon = optstring("daemon", NULL); config.daemon = optstring("daemon", NULL);
config.logger = optstring("logger", NULL); config.logger = optstring("logger", NULL);
config.logservice = optstring("logservice", "logger");
lua_close(L); lua_close(L);

View File

@@ -82,7 +82,7 @@ skynet_current_handle(void) {
void * handle = pthread_getspecific(G_NODE.handle_key); void * handle = pthread_getspecific(G_NODE.handle_key);
return (uint32_t)(uintptr_t)handle; return (uint32_t)(uintptr_t)handle;
} else { } else {
uintptr_t v = (uint32_t)(-THREAD_MAIN); uint32_t v = (uint32_t)(-THREAD_MAIN);
return v; return v;
} }
} }

View File

@@ -223,9 +223,9 @@ skynet_start(struct skynet_config * config) {
skynet_timer_init(); skynet_timer_init();
skynet_socket_init(); skynet_socket_init();
struct skynet_context *ctx = skynet_context_new("logger", config->logger); struct skynet_context *ctx = skynet_context_new(config->logservice, config->logger);
if (ctx == NULL) { if (ctx == NULL) {
fprintf(stderr, "Can't launch logger service\n"); fprintf(stderr, "Can't launch %s service\n", config->logservice);
exit(1); exit(1);
} }

View File

@@ -508,7 +508,7 @@ udp_socket_address(struct socket *s, const uint8_t udp_address[UDP_ADDRESS_SIZE]
memset(&sa->v6, 0, sizeof(sa->v6)); memset(&sa->v6, 0, sizeof(sa->v6));
sa->s.sa_family = AF_INET6; sa->s.sa_family = AF_INET6;
sa->v6.sin6_port = port; sa->v6.sin6_port = port;
memcpy(&sa->v6.sin6_addr, udp_address + 1 + sizeof(uint16_t), sizeof(sa->v6.sin6_addr)); // ipv4 address is 128 bits memcpy(&sa->v6.sin6_addr, udp_address + 1 + sizeof(uint16_t), sizeof(sa->v6.sin6_addr)); // ipv6 address is 128 bits
return sizeof(sa->v6); return sizeof(sa->v6);
} }
return 0; return 0;
@@ -720,6 +720,7 @@ send_socket(struct socket_server *ss, struct request_send * request, struct sock
append_sendbuffer_udp(ss,s,priority,request,udp_address); append_sendbuffer_udp(ss,s,priority,request,udp_address);
} else { } else {
so.free_func(request->buffer); so.free_func(request->buffer);
return -1;
} }
} }
sp_write(ss->event_fd, s->fd, s, true); sp_write(ss->event_fd, s->fd, s, true);
@@ -887,6 +888,7 @@ add_udp_socket(struct socket_server *ss, struct request_udp *udp) {
if (ns == NULL) { if (ns == NULL) {
close(udp->fd); close(udp->fd);
ss->slot[HASH_ID(id)].type = SOCKET_TYPE_INVALID; ss->slot[HASH_ID(id)].type = SOCKET_TYPE_INVALID;
return;
} }
ns->type = SOCKET_TYPE_CONNECTED; ns->type = SOCKET_TYPE_CONNECTED;
memset(ns->p.udp_address, 0, sizeof(ns->p.udp_address)); memset(ns->p.udp_address, 0, sizeof(ns->p.udp_address));

View File

@@ -43,6 +43,7 @@ end
function init( ... ) function init( ... )
print ("ping server start:", ...) print ("ping server start:", ...)
snax.enablecluster() -- enable cluster call
-- init queue -- init queue
lock = queue() lock = queue()
end end