mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
bugfix : connection
This commit is contained in:
@@ -65,6 +65,7 @@ _write(lua_State *L) {
|
|||||||
case EINTR:
|
case EINTR:
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
return 0;
|
||||||
}
|
}
|
||||||
assert(err == sz);
|
assert(err == sz);
|
||||||
return 0;
|
return 0;
|
||||||
@@ -104,6 +105,7 @@ _writeblock(lua_State *L) {
|
|||||||
case EINTR:
|
case EINTR:
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
return 0;
|
||||||
}
|
}
|
||||||
assert(err == sz +2);
|
assert(err == sz +2);
|
||||||
return 0;
|
return 0;
|
||||||
|
|||||||
@@ -95,6 +95,7 @@ _del(struct connection_server * server, int fd) {
|
|||||||
static void
|
static void
|
||||||
_poll(struct connection_server * server) {
|
_poll(struct connection_server * server) {
|
||||||
int timeout = 100;
|
int timeout = 100;
|
||||||
|
void * buffer = NULL;
|
||||||
for (;;) {
|
for (;;) {
|
||||||
struct connection * c = connection_poll(server->pool, timeout);
|
struct connection * c = connection_poll(server->pool, timeout);
|
||||||
if (c==NULL) {
|
if (c==NULL) {
|
||||||
@@ -103,7 +104,9 @@ _poll(struct connection_server * server) {
|
|||||||
}
|
}
|
||||||
timeout = 0;
|
timeout = 0;
|
||||||
|
|
||||||
void * buffer = malloc(DEFAULT_BUFFER_SIZE);
|
if (buffer == NULL) {
|
||||||
|
buffer = malloc(DEFAULT_BUFFER_SIZE);
|
||||||
|
}
|
||||||
|
|
||||||
int size = recv(c->fd, buffer, DEFAULT_BUFFER_SIZE, MSG_DONTWAIT);
|
int size = recv(c->fd, buffer, DEFAULT_BUFFER_SIZE, MSG_DONTWAIT);
|
||||||
if (size < 0) {
|
if (size < 0) {
|
||||||
@@ -112,9 +115,11 @@ _poll(struct connection_server * server) {
|
|||||||
if (size == 0) {
|
if (size == 0) {
|
||||||
connection_del(server->pool, c->fd);
|
connection_del(server->pool, c->fd);
|
||||||
free(buffer);
|
free(buffer);
|
||||||
|
buffer = NULL;
|
||||||
skynet_send(server->ctx, 0, c->address, SESSION_CLIENT, NULL, 0, DONTCOPY);
|
skynet_send(server->ctx, 0, c->address, SESSION_CLIENT, NULL, 0, DONTCOPY);
|
||||||
} else {
|
} else {
|
||||||
skynet_send(server->ctx, 0, c->address, SESSION_CLIENT, buffer, size, DONTCOPY);
|
skynet_send(server->ctx, 0, c->address, SESSION_CLIENT, buffer, size, DONTCOPY);
|
||||||
|
buffer = NULL;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -137,8 +142,8 @@ _main(struct skynet_context * ctx, void * ud, int session, uint32_t source, cons
|
|||||||
char addr [addr_sz];
|
char addr [addr_sz];
|
||||||
memcpy(addr, endptr+1, addr_sz-1);
|
memcpy(addr, endptr+1, addr_sz-1);
|
||||||
addr[addr_sz-1] = '\0';
|
addr[addr_sz-1] = '\0';
|
||||||
uint32_t address = strtoul(addr, NULL, 16);
|
uint32_t address = strtoul(addr+1, NULL, 16);
|
||||||
if (address != 0) {
|
if (address == 0) {
|
||||||
skynet_error(ctx, "[connection] Invalid ADD command from %x (session = %d)", source, session);
|
skynet_error(ctx, "[connection] Invalid ADD command from %x (session = %d)", source, session);
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -30,8 +30,8 @@ local meta = {
|
|||||||
}
|
}
|
||||||
|
|
||||||
function redis.connect(dbname)
|
function redis.connect(dbname)
|
||||||
local handle = skynet.call(".redis-manager",dbname)
|
local handle = skynet.call(".redis-manager",skynet.unpack, skynet.pack(dbname))
|
||||||
assert(handle ~= "")
|
assert(handle ~= nil)
|
||||||
return setmetatable({ __handle = handle } , meta)
|
return setmetatable({ __handle = handle } , meta)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|||||||
@@ -12,7 +12,9 @@ function socket.connect(addr)
|
|||||||
if fd == nil then
|
if fd == nil then
|
||||||
return true
|
return true
|
||||||
end
|
end
|
||||||
skynet.send(".connection","ADD "..fd.." "..skynet.self())
|
print("connect to " .. addr)
|
||||||
|
local command = "ADD "..fd.." ".. skynet.address(skynet.self())
|
||||||
|
skynet.send(".connection", command )
|
||||||
object = c.new()
|
object = c.new()
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|||||||
@@ -56,10 +56,9 @@ _cb(struct skynet_context * context, void * ud, int session, uint32_t source, co
|
|||||||
if (b->init < DEFAULT_NUMBER) {
|
if (b->init < DEFAULT_NUMBER) {
|
||||||
if (source != b->launcher)
|
if (source != b->launcher)
|
||||||
return 0;
|
return 0;
|
||||||
assert(sz == 9);
|
char addr[sz+1];
|
||||||
char addr[10];
|
memcpy(addr, msg, sz);
|
||||||
memcpy(addr, msg, 9);
|
addr[sz] = '\0';
|
||||||
addr[9] = '\0';
|
|
||||||
uint32_t address = strtoul(addr+1, NULL, 16);
|
uint32_t address = strtoul(addr+1, NULL, 16);
|
||||||
assert(address != 0);
|
assert(address != 0);
|
||||||
_init(b, session, address);
|
_init(b, session, address);
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ skynet.dispatch(function(msg, sz , session, address)
|
|||||||
-- init notice
|
-- init notice
|
||||||
local reply = instance[address]
|
local reply = instance[address]
|
||||||
if reply then
|
if reply then
|
||||||
skynet.send(reply[2] , reply[1], address)
|
skynet.send(reply[2] , reply[1], skynet.address(address))
|
||||||
instance[address] = nil
|
instance[address] = nil
|
||||||
end
|
end
|
||||||
else
|
else
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
local skynet = require "skynet"
|
local skynet = require "skynet"
|
||||||
local socket = require "socket"
|
local socket = require "socket"
|
||||||
|
local int64 = require "int64"
|
||||||
local string = string
|
local string = string
|
||||||
local table = table
|
local table = table
|
||||||
local tonumber = tonumber
|
local tonumber = tonumber
|
||||||
@@ -10,6 +11,12 @@ local redis_server, redis_db = ...
|
|||||||
local function compose_message(msg)
|
local function compose_message(msg)
|
||||||
local lines = { "*" .. #msg }
|
local lines = { "*" .. #msg }
|
||||||
for _,v in ipairs(msg) do
|
for _,v in ipairs(msg) do
|
||||||
|
local t = type(v)
|
||||||
|
if t == "number" then
|
||||||
|
v = tostring(v)
|
||||||
|
elseif t == "userdata" then
|
||||||
|
v = int64.tostring(int64.new(v),10)
|
||||||
|
end
|
||||||
table.insert(lines,"$"..#v)
|
table.insert(lines,"$"..#v)
|
||||||
table.insert(lines,v)
|
table.insert(lines,v)
|
||||||
end
|
end
|
||||||
@@ -101,6 +108,7 @@ redcmd[45] = function(data) -- '-'
|
|||||||
end
|
end
|
||||||
|
|
||||||
redcmd[58] = function(data) -- ':'
|
redcmd[58] = function(data) -- ':'
|
||||||
|
-- todo: return string later
|
||||||
response(true, tonumber(data))
|
response(true, tonumber(data))
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -117,7 +125,6 @@ end
|
|||||||
|
|
||||||
local function init()
|
local function init()
|
||||||
while socket.connect(redis_server) do
|
while socket.connect(redis_server) do
|
||||||
print("Connect failed : "..redis_server)
|
|
||||||
skynet.sleep(1000)
|
skynet.sleep(1000)
|
||||||
end
|
end
|
||||||
if redis_db then
|
if redis_db then
|
||||||
|
|||||||
@@ -1,27 +1,26 @@
|
|||||||
local skynet = require "skynet"
|
local skynet = require "skynet"
|
||||||
local log = require "log"
|
local log = require "log"
|
||||||
|
local config = require "config"
|
||||||
|
|
||||||
local name = {
|
local redis_conf = skynet.getenv "redis"
|
||||||
main = "127.0.0.1:6379",
|
local name = config (redis_conf)
|
||||||
}
|
|
||||||
|
|
||||||
local connection = {}
|
local connection = {}
|
||||||
|
|
||||||
skynet.dispatch(function(msg, sz , session, from)
|
skynet.dispatch(function(msg, sz , session, from)
|
||||||
local dbname = skynet.tostring(msg,sz)
|
local dbname = skynet.unpack(msg,sz)
|
||||||
if connection[dbname] then
|
if connection[dbname] then
|
||||||
skynet.ret(connection[dbname])
|
skynet.ret(skynet.pack(connection[dbname]))
|
||||||
return
|
return
|
||||||
end
|
end
|
||||||
if name[dbname] == nil then
|
if name[dbname] == nil then
|
||||||
log.Error("Invalid db name : "..dbname)
|
log.Error("Invalid db name : "..dbname)
|
||||||
skynet.ret("")
|
skynet.ret(skynet.pack(nil))
|
||||||
return
|
return
|
||||||
end
|
end
|
||||||
|
|
||||||
local redis_cli = skynet.launch("snlua", "redis-cli", name[dbname])
|
local redis_cli = skynet.launch("snlua", "redis-cli", name[dbname])
|
||||||
connection[dbname] = redis_cli
|
connection[dbname] = redis_cli
|
||||||
skynet.ret(redis_cli)
|
skynet.ret(skynet.pack(redis_cli))
|
||||||
end)
|
end)
|
||||||
|
|
||||||
skynet.register ".redis-manager"
|
skynet.register ".redis-manager"
|
||||||
|
|||||||
Reference in New Issue
Block a user