mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 20:23:06 +00:00
@@ -1076,9 +1076,15 @@ pack_seg(const uint8_t *src, uint8_t * buffer, int sz, int n) {
|
||||
|
||||
static inline void
|
||||
write_ff(const uint8_t * src, uint8_t * des, int n) {
|
||||
int i;
|
||||
int align8_n = (n+7)&(~7);
|
||||
|
||||
des[0] = 0xff;
|
||||
des[1] = n-1;
|
||||
memcpy(des+2, src, n * 8);
|
||||
des[1] = align8_n/8 - 1;
|
||||
memcpy(des+2, src, n);
|
||||
for(i=0; i< align8_n-n; i++){
|
||||
des[n+2+i] = 0;
|
||||
}
|
||||
}
|
||||
|
||||
int
|
||||
@@ -1092,6 +1098,7 @@ sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) {
|
||||
const uint8_t * src = srcv;
|
||||
uint8_t * buffer = bufferv;
|
||||
for (i=0;i<srcsz;i+=8) {
|
||||
int n;
|
||||
int padding = i+8 - srcsz;
|
||||
if (padding > 0) {
|
||||
int j;
|
||||
@@ -1101,7 +1108,7 @@ sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) {
|
||||
}
|
||||
src = tmp;
|
||||
}
|
||||
int n = pack_seg(src, buffer, bufsz, ff_n);
|
||||
n = pack_seg(src, buffer, bufsz, ff_n);
|
||||
bufsz -= n;
|
||||
if (n == 10) {
|
||||
// first FF
|
||||
@@ -1112,14 +1119,14 @@ sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) {
|
||||
++ff_n;
|
||||
if (ff_n == 256) {
|
||||
if (bufsz >= 0) {
|
||||
write_ff(ff_srcstart, ff_desstart, 256);
|
||||
write_ff(ff_srcstart, ff_desstart, 256*8);
|
||||
}
|
||||
ff_n = 0;
|
||||
}
|
||||
} else {
|
||||
if (ff_n > 0) {
|
||||
if (bufsz >= 0) {
|
||||
write_ff(ff_srcstart, ff_desstart, ff_n);
|
||||
write_ff(ff_srcstart, ff_desstart, ff_n*8);
|
||||
}
|
||||
ff_n = 0;
|
||||
}
|
||||
@@ -1128,8 +1135,11 @@ sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) {
|
||||
buffer += n;
|
||||
size += n;
|
||||
}
|
||||
if (ff_n > 0 && bufsz >= 0) {
|
||||
write_ff(ff_srcstart, ff_desstart, ff_n);
|
||||
if(bufsz >= 0){
|
||||
if(ff_n == 1)
|
||||
write_ff(ff_srcstart, ff_desstart, 8);
|
||||
else if (ff_n > 1)
|
||||
write_ff(ff_srcstart, ff_desstart, srcsz - (intptr_t)(ff_srcstart - (const uint8_t*)srcv));
|
||||
}
|
||||
return size;
|
||||
}
|
||||
@@ -1144,10 +1154,11 @@ sproto_unpack(const void * srcv, int srcsz, void * bufferv, int bufsz) {
|
||||
--srcsz;
|
||||
++src;
|
||||
if (header == 0xff) {
|
||||
int n;
|
||||
if (srcsz < 0) {
|
||||
return -1;
|
||||
}
|
||||
int n = (src[0] + 1) * 8;
|
||||
n = (src[0] + 1) * 8;
|
||||
if (srcsz < n + 1)
|
||||
return -1;
|
||||
srcsz -= n + 1;
|
||||
|
||||
@@ -129,6 +129,7 @@ function mongo.client( conf )
|
||||
response = dispatch_reply,
|
||||
auth = mongo_auth(obj),
|
||||
backup = backup,
|
||||
nodelay = true,
|
||||
}
|
||||
setmetatable(obj, client_meta)
|
||||
obj.__sock:connect(true) -- try connect only once
|
||||
|
||||
@@ -84,6 +84,7 @@ function redis.connect(db_conf)
|
||||
host = db_conf.host,
|
||||
port = db_conf.port or 6379,
|
||||
auth = redis_login(db_conf.auth, db_conf.db),
|
||||
nodelay = true,
|
||||
}
|
||||
-- try connect first only once
|
||||
channel:connect(true)
|
||||
@@ -199,6 +200,7 @@ function redis.watch(db_conf)
|
||||
host = db_conf.host,
|
||||
port = db_conf.port or 6379,
|
||||
auth = watch_login(obj, db_conf.auth),
|
||||
nodelay = true,
|
||||
}
|
||||
obj.__sock = channel
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
local skynet = require "skynet"
|
||||
local socket = require "socket"
|
||||
local socketdriver = require "socketdriver"
|
||||
|
||||
-- channel support auto reconnect , and capture socket error in request/response transaction
|
||||
-- { host = "", port = , auth = function(so) , response = function(so) session, data }
|
||||
@@ -37,6 +38,7 @@ function socket_channel.channel(desc)
|
||||
__sock = false,
|
||||
__closed = false,
|
||||
__authcoroutine = false,
|
||||
__nodelay = desc.nodelay,
|
||||
}
|
||||
|
||||
return setmetatable(c, channel_meta)
|
||||
@@ -186,6 +188,9 @@ local function connect_once(self)
|
||||
return false
|
||||
end
|
||||
end
|
||||
if self.__nodelay then
|
||||
socketdriver.nodelay(fd)
|
||||
end
|
||||
|
||||
self.__sock = setmetatable( {fd} , channel_socket_meta )
|
||||
skynet.fork(dispatch_function(self), self)
|
||||
|
||||
@@ -47,7 +47,7 @@ forward_message(int type, bool padding, struct socket_message * result) {
|
||||
sm->ud = result->ud;
|
||||
if (padding) {
|
||||
sm->buffer = NULL;
|
||||
strcpy((char*)(sm+1), result->data);
|
||||
memcpy(sm+1, result->data, sz - sizeof(*sm));
|
||||
} else {
|
||||
sm->buffer = result->data;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user