mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 20:23:06 +00:00
Compare commits
27 Commits
v1.0.0-rc
...
v1.0.0-rc2
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cbf22ac9f3 | ||
|
|
8f7fa6a562 | ||
|
|
7b96eb28eb | ||
|
|
14ccf96eb0 | ||
|
|
ea8557296b | ||
|
|
2a821fb62a | ||
|
|
d982009260 | ||
|
|
a3f028ee3d | ||
|
|
d34b092cdd | ||
|
|
77f4071191 | ||
|
|
653b29145d | ||
|
|
96e0a4e3cc | ||
|
|
1da92850a0 | ||
|
|
9747cf9aee | ||
|
|
7f45b3fd10 | ||
|
|
217ab47577 | ||
|
|
532f47444b | ||
|
|
24f9ee1560 | ||
|
|
d7a76c7c72 | ||
|
|
c175d66e23 | ||
|
|
05003e3e2a | ||
|
|
64c2cb1b51 | ||
|
|
7f0e4eac8a | ||
|
|
141a4d0992 | ||
|
|
645ce72e7e | ||
|
|
fca4d67a43 | ||
|
|
014cbb9676 |
@@ -419,8 +419,11 @@ luaS_clonestring(lua_State *L, TString *ts) {
|
|||||||
result = query_ptr(ts);
|
result = query_ptr(ts);
|
||||||
if (result)
|
if (result)
|
||||||
return result;
|
return result;
|
||||||
// ts is not in SSM, so recalc hash, and add it to SSM
|
|
||||||
h = luaS_hash(str, l, 0);
|
h = luaS_hash(str, l, 0);
|
||||||
|
result = query_string(h, str, l);
|
||||||
|
if (result)
|
||||||
|
return result;
|
||||||
|
// ts is not in SSM, so recalc hash, and add it to SSM
|
||||||
return add_string(h, str, l);
|
return add_string(h, str, l);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -200,18 +200,19 @@ void luaV_finishset (lua_State *L, const TValue *t, TValue *key,
|
|||||||
for (loop = 0; loop < MAXTAGLOOP; loop++) {
|
for (loop = 0; loop < MAXTAGLOOP; loop++) {
|
||||||
const TValue *tm;
|
const TValue *tm;
|
||||||
if (oldval != NULL) {
|
if (oldval != NULL) {
|
||||||
lua_assert(ttistable(t) && ttisnil(oldval));
|
Table *h = hvalue(t);
|
||||||
|
lua_assert(ttisnil(oldval));
|
||||||
/* must check the metamethod */
|
/* must check the metamethod */
|
||||||
if ((tm = fasttm(L, hvalue(t)->metatable, TM_NEWINDEX)) == NULL &&
|
if ((tm = fasttm(L, h->metatable, TM_NEWINDEX)) == NULL &&
|
||||||
/* no metamethod; is there a previous entry in the table? */
|
/* no metamethod; is there a previous entry in the table? */
|
||||||
(oldval != luaO_nilobject ||
|
(oldval != luaO_nilobject ||
|
||||||
/* no previous entry; must create one. (The next test is
|
/* no previous entry; must create one. (The next test is
|
||||||
always true; we only need the assignment.) */
|
always true; we only need the assignment.) */
|
||||||
(oldval = luaH_newkey(L, hvalue(t), key), 1))) {
|
(oldval = luaH_newkey(L, h, key), 1))) {
|
||||||
/* no metamethod and (now) there is an entry with given key */
|
/* no metamethod and (now) there is an entry with given key */
|
||||||
setobj2t(L, cast(TValue *, oldval), val);
|
setobj2t(L, cast(TValue *, oldval), val);
|
||||||
invalidateTMcache(hvalue(t));
|
invalidateTMcache(h);
|
||||||
luaC_barrierback(L, hvalue(t), val);
|
luaC_barrierback(L, h, val);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
/* else will try the metamethod */
|
/* else will try the metamethod */
|
||||||
|
|||||||
@@ -1,3 +1,12 @@
|
|||||||
|
v1.0.0-rc2 (2015-3-7)
|
||||||
|
-----------
|
||||||
|
* Fix a bug in lua 5.3.2
|
||||||
|
* Update sproto (fix bugs and add ud for package)
|
||||||
|
* Fix a bug in http
|
||||||
|
* Fix a bug in harbor
|
||||||
|
* Fix a bug in socket channel
|
||||||
|
* Enhance remote debugger
|
||||||
|
|
||||||
v1.0.0-rc (2015-12-28)
|
v1.0.0-rc (2015-12-28)
|
||||||
-----------
|
-----------
|
||||||
* Update to lua 5.3.2
|
* Update to lua 5.3.2
|
||||||
|
|||||||
@@ -885,7 +885,12 @@ int lhmac_sha1(lua_State *L);
|
|||||||
int
|
int
|
||||||
luaopen_crypt(lua_State *L) {
|
luaopen_crypt(lua_State *L) {
|
||||||
luaL_checkversion(L);
|
luaL_checkversion(L);
|
||||||
srandom(time(NULL));
|
static int init = 0;
|
||||||
|
if (!init) {
|
||||||
|
// Don't need call srandom more than once.
|
||||||
|
init = 1 ;
|
||||||
|
srandom(time(NULL));
|
||||||
|
}
|
||||||
luaL_Reg l[] = {
|
luaL_Reg l[] = {
|
||||||
{ "hashkey", lhashkey },
|
{ "hashkey", lhashkey },
|
||||||
{ "randomkey", lrandomkey },
|
{ "randomkey", lrandomkey },
|
||||||
|
|||||||
@@ -469,7 +469,7 @@ lpack(lua_State *L) {
|
|||||||
size_t sz=0;
|
size_t sz=0;
|
||||||
const void * buffer = getbuffer(L, 1, &sz);
|
const void * buffer = getbuffer(L, 1, &sz);
|
||||||
// the worst-case space overhead of packing is 2 bytes per 2 KiB of input (256 words = 2KiB).
|
// the worst-case space overhead of packing is 2 bytes per 2 KiB of input (256 words = 2KiB).
|
||||||
size_t maxsz = (sz + 2047) / 2048 * 2 + sz;
|
size_t maxsz = (sz + 2047) / 2048 * 2 + sz + 2;
|
||||||
void * output = lua_touserdata(L, lua_upvalueindex(1));
|
void * output = lua_touserdata(L, lua_upvalueindex(1));
|
||||||
int bytes;
|
int bytes;
|
||||||
int osz = lua_tointeger(L, lua_upvalueindex(2));
|
int osz = lua_tointeger(L, lua_upvalueindex(2));
|
||||||
|
|||||||
@@ -72,7 +72,8 @@ local function request(fd, method, host, url, recvheader, header, content)
|
|||||||
body = body .. padding
|
body = body .. padding
|
||||||
end
|
end
|
||||||
else
|
else
|
||||||
body = nil
|
-- no content-length, read all
|
||||||
|
body = body .. socket.readall(fd)
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|||||||
@@ -19,6 +19,8 @@ function sockethelper.readfunc(fd)
|
|||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
|
sockethelper.readall = socket.readall
|
||||||
|
|
||||||
function sockethelper.writefunc(fd)
|
function sockethelper.writefunc(fd)
|
||||||
return function(content)
|
return function(content)
|
||||||
local ok = writebytes(fd, content)
|
local ok = writebytes(fd, content)
|
||||||
|
|||||||
@@ -125,19 +125,6 @@ local function _compose_packet(self, req, size)
|
|||||||
return packet
|
return packet
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|
||||||
local function _send_packet(self, req, size)
|
|
||||||
local sock = self.sock
|
|
||||||
|
|
||||||
self.packet_no = self.packet_no + 1
|
|
||||||
|
|
||||||
|
|
||||||
local packet = _set_byte3(size) .. strchar(self.packet_no) .. req
|
|
||||||
|
|
||||||
return socket.write(self.sock,packet)
|
|
||||||
end
|
|
||||||
|
|
||||||
|
|
||||||
local function _recv_packet(self,sock)
|
local function _recv_packet(self,sock)
|
||||||
|
|
||||||
|
|
||||||
@@ -620,6 +607,10 @@ local function _query_resp(self)
|
|||||||
while err =="again" do
|
while err =="again" do
|
||||||
res, err, errno, sqlstate = read_result(self,sock)
|
res, err, errno, sqlstate = read_result(self,sock)
|
||||||
if not res then
|
if not res then
|
||||||
|
mulitresultset.badresult = true
|
||||||
|
mulitresultset.err = err
|
||||||
|
mulitresultset.errno = errno
|
||||||
|
mulitresultset.sqlstate = sqlstate
|
||||||
return true, mulitresultset
|
return true, mulitresultset
|
||||||
end
|
end
|
||||||
mulitresultset[i]=res
|
mulitresultset[i]=res
|
||||||
|
|||||||
@@ -1,6 +1,4 @@
|
|||||||
local io = io
|
|
||||||
local table = table
|
local table = table
|
||||||
local debug = debug
|
|
||||||
|
|
||||||
return function (skynet, export)
|
return function (skynet, export)
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
local skynet = require "skynet"
|
local skynet = require "skynet"
|
||||||
local debugchannel = require "debugchannel"
|
local debugchannel = require "debugchannel"
|
||||||
local socket = require "socket"
|
local socketdriver = require "socketdriver"
|
||||||
local injectrun = require "skynet.injectcode"
|
local injectrun = require "skynet.injectcode"
|
||||||
local table = table
|
local table = table
|
||||||
local debug = debug
|
local debug = debug
|
||||||
@@ -56,7 +56,7 @@ local function gen_print(fd)
|
|||||||
tmp[i] = tostring(tmp[i])
|
tmp[i] = tostring(tmp[i])
|
||||||
end
|
end
|
||||||
table.insert(tmp, "\n")
|
table.insert(tmp, "\n")
|
||||||
socket.write(fd, table.concat(tmp, "\t"))
|
socketdriver.send(fd, table.concat(tmp, "\t"))
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
@@ -140,10 +140,13 @@ end
|
|||||||
local function watch_proto(protoname, cond)
|
local function watch_proto(protoname, cond)
|
||||||
local proto = assert(replace_upvalue(skynet.register_protocol, "proto"), "Can't find proto table")
|
local proto = assert(replace_upvalue(skynet.register_protocol, "proto"), "Can't find proto table")
|
||||||
local p = proto[protoname]
|
local p = proto[protoname]
|
||||||
local dispatch = p.dispatch_origin or p.dispatch
|
if p == nil then
|
||||||
if p == nil or dispatch == nil then
|
|
||||||
return "No " .. protoname
|
return "No " .. protoname
|
||||||
end
|
end
|
||||||
|
local dispatch = p.dispatch_origin or p.dispatch
|
||||||
|
if dispatch == nil then
|
||||||
|
return "No dispatch"
|
||||||
|
end
|
||||||
p.dispatch_origin = dispatch
|
p.dispatch_origin = dispatch
|
||||||
p.dispatch = function(...)
|
p.dispatch = function(...)
|
||||||
if not cond or cond(...) then
|
if not cond or cond(...) then
|
||||||
@@ -208,7 +211,7 @@ local function hook_dispatch(dispatcher, resp, fd, channel)
|
|||||||
local function debug_hook()
|
local function debug_hook()
|
||||||
while true do
|
while true do
|
||||||
if newline then
|
if newline then
|
||||||
socket.write(fd, prompt)
|
socketdriver.send(fd, prompt)
|
||||||
newline = false
|
newline = false
|
||||||
end
|
end
|
||||||
local cmd = channel:read()
|
local cmd = channel:read()
|
||||||
@@ -244,7 +247,9 @@ local function hook_dispatch(dispatcher, resp, fd, channel)
|
|||||||
func = replace_upvalue(dispatcher, HOOK_FUNC, hook)
|
func = replace_upvalue(dispatcher, HOOK_FUNC, hook)
|
||||||
if func then
|
if func then
|
||||||
local function idle()
|
local function idle()
|
||||||
skynet.timeout(10,idle) -- idle every 0.1s
|
if raw_dispatcher then
|
||||||
|
skynet.timeout(10,idle) -- idle every 0.1s
|
||||||
|
end
|
||||||
end
|
end
|
||||||
skynet.timeout(0, idle)
|
skynet.timeout(0, idle)
|
||||||
end
|
end
|
||||||
|
|||||||
@@ -1,7 +1,4 @@
|
|||||||
local si = require "snax.interface"
|
local si = require "snax.interface"
|
||||||
local io = io
|
|
||||||
|
|
||||||
local hotfix = {}
|
|
||||||
|
|
||||||
local function envid(f)
|
local function envid(f)
|
||||||
local i = 1
|
local i = 1
|
||||||
|
|||||||
@@ -411,10 +411,10 @@ end
|
|||||||
|
|
||||||
function channel:close()
|
function channel:close()
|
||||||
if not self.__closed then
|
if not self.__closed then
|
||||||
local thread = assert(self.__dispatch_thread)
|
local thread = self.__dispatch_thread
|
||||||
self.__closed = true
|
self.__closed = true
|
||||||
close_channel_socket(self)
|
close_channel_socket(self)
|
||||||
if not self.__response and self.__dispatch_thread == thread then
|
if not self.__response and self.__dispatch_thread == thread and thread then
|
||||||
-- dispatch by order, send close signal to dispatch thread
|
-- dispatch by order, send close signal to dispatch thread
|
||||||
push_response(self, true, false) -- (true, false) is close signal
|
push_response(self, true, false) -- (true, false) is close signal
|
||||||
end
|
end
|
||||||
|
|||||||
@@ -180,9 +180,10 @@ end
|
|||||||
local header_tmp = {}
|
local header_tmp = {}
|
||||||
|
|
||||||
local function gen_response(self, response, session)
|
local function gen_response(self, response, session)
|
||||||
return function(args)
|
return function(args, ud)
|
||||||
header_tmp.type = nil
|
header_tmp.type = nil
|
||||||
header_tmp.session = session
|
header_tmp.session = session
|
||||||
|
header_tmp.ud = ud
|
||||||
local header = core.encode(self.__package, header_tmp)
|
local header = core.encode(self.__package, header_tmp)
|
||||||
if response then
|
if response then
|
||||||
local content = core.encode(response, args)
|
local content = core.encode(response, args)
|
||||||
@@ -197,6 +198,7 @@ function host:dispatch(...)
|
|||||||
local bin = core.unpack(...)
|
local bin = core.unpack(...)
|
||||||
header_tmp.type = nil
|
header_tmp.type = nil
|
||||||
header_tmp.session = nil
|
header_tmp.session = nil
|
||||||
|
header_tmp.ud = nil
|
||||||
local header, size = core.decode(self.__package, bin, header_tmp)
|
local header, size = core.decode(self.__package, bin, header_tmp)
|
||||||
local content = bin:sub(size + 1)
|
local content = bin:sub(size + 1)
|
||||||
if header.type then
|
if header.type then
|
||||||
@@ -207,9 +209,9 @@ function host:dispatch(...)
|
|||||||
result = core.decode(proto.request, content)
|
result = core.decode(proto.request, content)
|
||||||
end
|
end
|
||||||
if header_tmp.session then
|
if header_tmp.session then
|
||||||
return "REQUEST", proto.name, result, gen_response(self, proto.response, header_tmp.session)
|
return "REQUEST", proto.name, result, gen_response(self, proto.response, header_tmp.session), header.ud
|
||||||
else
|
else
|
||||||
return "REQUEST", proto.name, result
|
return "REQUEST", proto.name, result, nil, header.ud
|
||||||
end
|
end
|
||||||
else
|
else
|
||||||
-- response
|
-- response
|
||||||
@@ -217,19 +219,20 @@ function host:dispatch(...)
|
|||||||
local response = assert(self.__session[session], "Unknown session")
|
local response = assert(self.__session[session], "Unknown session")
|
||||||
self.__session[session] = nil
|
self.__session[session] = nil
|
||||||
if response == true then
|
if response == true then
|
||||||
return "RESPONSE", session
|
return "RESPONSE", session, nil, header.ud
|
||||||
else
|
else
|
||||||
local result = core.decode(response, content)
|
local result = core.decode(response, content)
|
||||||
return "RESPONSE", session, result
|
return "RESPONSE", session, result, header.ud
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
end
|
end
|
||||||
|
|
||||||
function host:attach(sp)
|
function host:attach(sp)
|
||||||
return function(name, args, session)
|
return function(name, args, session, ud)
|
||||||
local proto = queryproto(sp, name)
|
local proto = queryproto(sp, name)
|
||||||
header_tmp.type = proto.tag
|
header_tmp.type = proto.tag
|
||||||
header_tmp.session = session
|
header_tmp.session = session
|
||||||
|
header_tmp.ud = ud
|
||||||
local header = core.encode(self.__package, header_tmp)
|
local header = core.encode(self.__package, header_tmp)
|
||||||
|
|
||||||
if session then
|
if session then
|
||||||
|
|||||||
@@ -346,7 +346,6 @@ dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
|||||||
struct harbor_msg_queue * queue = node->queue;
|
struct harbor_msg_queue * queue = node->queue;
|
||||||
uint32_t handle = node->value;
|
uint32_t handle = node->value;
|
||||||
int harbor_id = handle >> HANDLE_REMOTE_SHIFT;
|
int harbor_id = handle >> HANDLE_REMOTE_SHIFT;
|
||||||
assert(harbor_id != 0);
|
|
||||||
struct skynet_context * context = h->ctx;
|
struct skynet_context * context = h->ctx;
|
||||||
struct slave *s = &h->s[harbor_id];
|
struct slave *s = &h->s[harbor_id];
|
||||||
int fd = s->fd;
|
int fd = s->fd;
|
||||||
@@ -366,6 +365,16 @@ dispatch_name_queue(struct harbor *h, struct keyvalue * node) {
|
|||||||
push_queue_msg(s->queue, m);
|
push_queue_msg(s->queue, m);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if (harbor_id == (h->slave >> HANDLE_REMOTE_SHIFT)) {
|
||||||
|
// the harbor_id is local
|
||||||
|
struct harbor_msg * m;
|
||||||
|
while ((m = pop_queue(s->queue)) != NULL) {
|
||||||
|
int type = m->header.destination >> HANDLE_REMOTE_SHIFT;
|
||||||
|
skynet_send(context, m->header.source, handle , type | PTYPE_TAG_DONTCOPY, m->header.session, m->buffer, m->size);
|
||||||
|
}
|
||||||
|
release_queue(s->queue);
|
||||||
|
s->queue = NULL;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -84,7 +84,7 @@ skynet_mq_create(uint32_t handle) {
|
|||||||
SPIN_INIT(q)
|
SPIN_INIT(q)
|
||||||
// When the queue is create (always between service create and service init) ,
|
// When the queue is create (always between service create and service init) ,
|
||||||
// set in_global flag to avoid push it to global queue .
|
// set in_global flag to avoid push it to global queue .
|
||||||
// If the service init success, skynet_context_new will call skynet_mq_force_push to push it to global queue.
|
// If the service init success, skynet_context_new will call skynet_mq_push to push it to global queue.
|
||||||
q->in_global = MQ_IN_GLOBAL;
|
q->in_global = MQ_IN_GLOBAL;
|
||||||
q->release = 0;
|
q->release = 0;
|
||||||
q->overload = 0;
|
q->overload = 0;
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ struct socket_server;
|
|||||||
struct socket_message {
|
struct socket_message {
|
||||||
int id;
|
int id;
|
||||||
uintptr_t opaque;
|
uintptr_t opaque;
|
||||||
int ud; // for accept, ud is listen id ; for data, ud is size of data
|
int ud; // for accept, ud is new connection id ; for data, ud is size of data
|
||||||
char * data;
|
char * data;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user