lua层处理
上述代码省略了部分冗余的判断逻辑,我们聚焦于核心执行流程。可以清晰地看到,socket.read 方法首先尝试读取指定字节数的数据。这一过程调用了 lua-socket.c 中的 lpopbuffer 方法,该部分将在后续深入分析。若成功读取到足够数据,则直接返回结果;否则,当前协程将主动让出执行权,进入挂起状态。系统会等待对端发送更多数据,直至满足预设的字节数要求,随后重新恢复该协程的执行。
在此机制中,s 对象通过 read_required 和 co 字段分别保存了期望读取的字节数以及当前运行的协程引用。这种设计确保了当接收到新数据时,系统能够准确判断数据量是否充足,并精准定位到应当被唤醒的特定协程,从而保证数据处理的正确性与顺序性。
关于 skynet.wait 方法的具体实现细节,此处暂不作展开分析。若读者对此感兴趣,可参考 zhuanlan.zhihu.com/p/84653538 一文,其中提供了更为详尽的解读。
接下来,我们将目光转向接收到新数据后的处理逻辑,这是理解整个异步读取机制的关键环节。
socket.read(id, sz)
我们首先看一下官方文档是怎么写的: 从一个 socket 上读 sz 指定的字节数。如果读到了指定长度的字符串,它把这个字符串返回。如果连接断开导致字节数不够,将返回一个 false 加上读到的字符串。如果 sz 为 nil ,则返回尽可能多的字节数,但至少读一个字节(若无新数据,会阻塞)。 所谓阻塞模式,实际上是利用了 lua 的 coroutine 机制。当你调用 socket api 时,服务有可能被挂起(时间片被让给其他业务处理),待结果通过 socket 消息返回,coroutine 将延续执行。
分析read方法有两个重点,一是同步调用的原理,二是数据的存取方式。
这里主要涉及到socket.lua 和 lua-socket.c 两个文件
我们首先来看一下lua层
lua层处理
function socket.read(id, sz)
local s = socket_pool[id]
local ret = driver.pop(s.buffer, s.pool, sz) --调用lpopbuffer
if ret then --如果返回的不是nil,表示读取到了需要的数据,直接返回
return ret
end
--否则就只有暂停当前协程,直到从对端发来足够的数据后再恢复执行
--在暂停之前会判断对端是不是已经断开连接了,如果是的话,就不暂停当前协程了,直接返回目前已经读到的所有数据
if not s.connected then
return false, driver.readall(s.buffer, s.pool)
end
s.read_required = sz --保存需要读取的size,这样当对端发送过来更多的数据时我们才能判断够不够
suspend(s) --挂起当前协程,当数据够或者连接断开时恢复协程
ret = driver.pop(s.buffer, s.pool, sz)
if ret then
return ret
else
return false, driver.readall(s.buffer, s.pool)
end
end
local function suspend(s)
s.co = coroutine.running() --保存协程
skynet.wait(s.co) --让出当前协程执行权,再次恢复时也是从这里开始执行
end
上面的代码删除了一些额外判断,我们只看一下主要的流程,可以看到socket.read 首先尝试读取指定字节数的数据(这里调用了 lua-socket.c 的 lpopbuffer 方法,我们后面具体分析),如果读取到了,就直接返回,否则就让出当前协程的执行权,直到对端发来了更多的数据满足了我们的要求,才会重新恢复当前协程的执行。我们看到这里s用read_required和co字段分别保存了要求的字节数和当前运行的协程,这样当收到数据时我们才能够判断够不够以及应该恢复哪个协程。
这里skynet.wait 方法 就不分析了,如果感兴趣的话可以参考 zhuanlan.zhihu.com/p/84653538 这篇文章
下面我们看一下接收到新数据后的处理
-- read skynet_socket.h for these macro
-- SKYNET_SOCKET_TYPE_DATA = 1
socket_message[1] = function(id, size, data)
local s = socket_pool[id]
local sz = driver.push(s.buffer, s.pool, data, size) --调用lpushbuffer,返回未读取的数据长度
local rr = s.read_required --取出要读取的数据长度
local rrt = type(rr)
if rrt == "number" then
if sz >= rr then --如果有足够的数据可读
s.read_required = nil
wakeup(s) --唤起调用read的协程
end
else
end
end
local function wakeup(s)
local co = s.co
if co then
s.co = nil
skynet.wakeup(co)
end
end
c层处理
在具体查看读取及写入方法前,我们先来深入了解一下存储网络数据包的核心数据结构。这一结构的设计直接决定了数据处理的效率与逻辑的清晰度。
//存放所有未读取的数据
struct socket_buffer {
int size; // 还未读取的网络数据总长度
int offset; // head 已读数据的偏移
struct buffer_node *head; // 数据buff_node链表的头部指针
struct buffer_node *tail;
};
struct buffer_node {
char *msg;
int sz; //该buffer_node存储的数据size
struct buffer_node *next;
};
从上述代码结构可以看出,socket_buffer本质上是一个链表结构。当系统接收到新的网络数据时,会动态构造一个新的节点,并将其插入到链表的尾部。而在执行读取操作时,系统始终从头节点开始读取。如果单次读取无法获取全部数据,系统会在 offset 变量中记录当前头节点剩余未读的数据量。一旦头节点的数据被完全读取完毕,系统便会移除该头节点,并将链表中的下一个节点提升为新的头节点,从而保证读取指针的正确推进。
接下来,我们将详细分析具体的读取方法实现。
lpopbuffer 读取数据
/*
userdata send_buffer
table pool
integer sz
*/
static int
lpopbuffer(lua_State *L) {
struct socket_buffer *sb = lua_touserdata(L, 1);
if (sb == NULL) {
return luaL_error(L, "Need buffer object at param 1");
}
luaL_checktype(L, 2, LUA_TTABLE);
int sz = luaL_checkinteger(L, 3);
if (sb->size < sz || sz == 0) { //如果此次没有足够的数据可读
lua_pushnil(L); //直接返回一个nil
} else {
pop_lstring(L, sb, sz, 0); //从sb中取出sz字节的数据
sb->size -= sz; //更新size
}
lua_pushinteger(L, sb->size); //同时返回剩下可读的字节数
return 2;
}
观察 lpopbuffer 函数的实现逻辑,它首先会判断当前缓冲区中可读的字节数是否满足请求长度。如果数据量不足,函数将直接返回 nil;反之,若数据充足,则提取指定长度的数据。除了返回实际读取的数据外,该函数还会额外返回剩余可读的字节数,以便调用方了解缓冲区的状态。随后,我们进入 pop_lstring 方法,深入探究从 sb 中提取数据的具体实现细节。
static void
pop_lstring(lua_State *L, struct socket_buffer *sb, int sz, int skip) {
struct buffer_node *current = sb->head; //取出头节点
if (sz < current->sz - sb->offset) { //如果有足够数据,并且还有多的
lua_pushlstring(L, current->msg + sb->offset, sz - skip);
sb->offset += sz;
return;
}
if (sz == current->sz - sb->offset) { //刚好够(此时就需要移除头结点了)
lua_pushlstring(L, current->msg + sb->offset, sz - skip);
return_free_node(L, 2, sb); //移除头结点
return;
}
luaL_Buffer b;
luaL_buffinitsize(L, &b, sz);
for (;;) {
int bytes = current->sz - sb->offset; //sb当前节点可读的数据
if (bytes >= sz) { //如果当前节点的数据够了,就读取需要的,然后退出循环返回
if (sz > skip) {
luaL_addlstring(&b, current->msg + sb->offset, sz - skip);
}
sb->offset += sz;
if (bytes == sz) {
return_free_node(L, 2, sb);
}
break;
}
//如果当前节点的数据不够,就全部读完,更新当前节点为下一个节点,继续循环
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);
sz -= bytes; //更新需要读取的数据
if (sz == 0) //需要读取的数据为0,表示已经读够了,退出循环
break;
current = sb->head; //更新current节点为新的头结点
assert(current);
}
luaL_pushresult(&b);
}
在 pop_lstring 方法中,读取操作从头节点开始,持续进行直到获取到所需的全部字节数后返回结果。如果头节点剩余的数据量恰好满足请求,则只需读取该节点即可;否则,系统需要从头节点开始,依次遍历并读取链表中的后续节点。回顾上一步的调用栈,传入的 skip 参数值为 0,因此在当前的应用场景中,我们可以直接忽略该参数的影响。值得注意的是,每次头节点数据读取完毕后,系统都必须执行 sb 链表的更新操作,即将头节点的下一个节点设置为新的头节点。在此过程中,return_free_node 方法涉及到了读写缓冲区的复杂交互,我们将在后续章节中对其进行分析。
最后,我们来看一下数据的写入方法。
lpushbuffer 写入数据
/*
userdata send_buffer
table pool
lightuserdata msg
*/
int size
static int
lpushbuffer(lua_State *L) {
struct socket_buffer *sb = lua_touserdata(L, 1);
if (sb == NULL) {
return luaL_error(L, "need buffer object at param 1");
}
char *msg = lua_touserdata(L, 3);
if (msg == NULL) {
return luaL_error(L, "need message block at param 3");
}
int pool_index = 2;
luaL_checktype(L, pool_index, LUA_TTABLE);
int sz = luaL_checkinteger(L, 4);
lua_rawgeti(L, pool_index, 1); //把pool[1]压栈
struct buffer_node *free_node = lua_touserdata(L, -1); // 拿到pool[1]处存的 free_node
//将网络数据的指针和大小保存在这个空闲的 buff_node上
free_node->msg = msg;
free_node->sz = sz;
free_node->next = NULL;
//把存储了数据的 buff_node 插入到sb的尾部
if (sb->head == NULL) {
assert(sb->tail == NULL);
sb->head = sb->tail = free_node;
} else {
sb->tail->next = free_node;
sb->tail = free_node;
}
sb->size += sz;
lua_pushinteger(L, sb->size);
return 1;
}
观察 lpushbuffer 函数的实现,它首先会从缓冲池中获取一个空闲节点 free_node(具体的获取机制暂不展开讨论)。接着,将待写入的数据保存到该 free_node 中,最后将该节点插入到 sb 链表的尾部。通过这一系列操作,数据便成功完成了存放过程,等待后续的读取处理。
缓冲池的数据结构设计
在深入剖析具体实现之前,我们需要明确缓冲池的核心数据结构。该池本质上是一个存储 buffer_node 的集合,其中既包含已填充数据的节点,也包含作为空闲资源存在的节点。为了高效地管理这些空闲节点,我们引入了一个关键的指针 free_node。这个指针直接指向当前可用的第一个空闲节点,从而避免了遍历整个池子来寻找可用资源的开销。
为了实现高效的增删操作,每个空闲节点内部都包含一个 next 指针,指向链表中的下一个空闲节点。这种链式结构使得我们可以以常数时间复杂度将新产生的空闲节点插入到链表头部,或者将已使用的节点从链表中移除并重新分配。通过维护 free_node 和节点间的 next 链接,缓冲池能够迅速响应数据的读写请求,确保数据流转的流畅性。
动态扩容策略与实现
缓冲池并非静态不变,当接收的数据量超过当前池容量时,必须执行扩容操作。扩容的核心问题在于确定每次新增 buffer_node 的数量。如果每次增加的数量过少,会导致频繁的扩容操作,产生额外的内存分配开销;反之,如果每次增加过多,则可能造成内存浪费。参考业界通用的扩容策略,采用“指数级增长”是最为合理的方案,即每次新增的节点数量是上一次的两倍。这种策略在内存利用率与扩容频率之间取得了良好的平衡。
然而,这里存在一个实现上的挑战:扩容策略所需的“上次增量”信息需要被持久化。由于 pool 对象主要保存在 Lua 层,而具体的节点管理在 C 层进行,如果将增量信息单独存储在 Lua 变量中,每次调用 C 函数时都需要在两者之间传递该参数,增加了调用的复杂性和开销。为此,Skynet 采用了一种巧妙的优化手段:利用 Lua Table 自身的容量特性来隐式存储扩容信息。
具体而言,每次执行扩容操作时,系统会在 Lua Table 中新增一个表项。这个新表项的索引值(index)被用来编码此次需要新增的 buffer_node 数量。通过读取该表项的索引,C 代码可以直接计算出当前应分配的节点数,而无需额外的参数传递。这种设计不仅简化了 C 与 Lua 之间的交互接口,还巧妙地利用了 Lua 语言的特性,实现了状态信息的紧凑存储与高效访问。
缓冲池扩容机制
最终结构如前所述,初始状态下相关表项均不存在。接着分析 lpushbuffer 方法,此处仅关注与缓冲池相关的逻辑。如代码所示,空闲节点使用完毕后触发扩容,扩容后将新增 buffer_node 的首个节点赋值给 free_node。此外,在扩容过程中若检测到为首次使用,系统会重新设置 __gc 元方法,以确保垃圾回收机制的正确运行。
节点回收处理
随后考察 sb 头结点数据读取完毕后调用的 return_free_node 函数。该函数核心功能是将已读取完毕的 head 节点重新归还至缓冲池,并将其更新为新的 free_node,即更新空闲链表的头结点,从而完成节点的循环利用。
参考资料
pool = {
[1] = free_node, -- 指向下面31个表项中某个buff_node,作为空闲链表的头结点
[2] = buffer_node_pool, -- 存放16个buff_node
[3] = buffer_node_pool, -- 存放32个buff_node
...
[32] = buffer_node_pool, -- 存放4096个buff_node
}
static int
lpushbuffer(lua_State *L) {
int pool_index = 2;
luaL_checktype(L, pool_index, LUA_TTABLE);
lua_rawgeti(L, pool_index, 1); //把pool[1]压栈
struct buffer_node *free_node = lua_touserdata(L, -1); // 拿到pool[1]处存的 free_node
lua_pop(L, 1); //弹出pool[1]
if (free_node == NULL) { //为 null,则表示已经没有空闲的 buff_node
int tsz = lua_rawlen(L, pool_index); //拿到作为pool的table的元素个数
if (tsz == 0)
tsz++;
int size = 8;
//每次*2,达到4096就不再增了
if (tsz <= LARGE_PAGE_NODE - 3) {
size <<= tsz;
} else {
size <<= LARGE_PAGE_NODE - 3;
}
lnewpool(L, size); //初始化这个表项存储的所有buffer_node
free_node = lua_touserdata(L, -1); //free_node指向这个表项的第一个元素
lua_rawseti(L, pool_index, tsz + 1); //弹出栈顶值并赋值给pool[tsz+1]
}
lua_pushlightuserdata(L, free_node->next); // 将free_node 指向的下一个空闲 buff_node 指针赋值到 pool[1]
lua_rawseti(L, pool_index, 1); // pool[1]重新赋值到pool里
}
static int
lnewpool(lua_State *L, int sz) {
struct buffer_node *pool = lua_newuserdatauv(L, sizeof(struct buffer_node) * sz, 0); //新建pool的一个表项,并放到栈顶
int i;
for (i = 0; i < sz; i++) {
pool[i].msg = NULL;
pool[i].sz = 0;
pool[i].next = &pool[i + 1];
}
pool[sz - 1].next = NULL;
if (luaL_newmetatable(L, "buffer_pool")) { //只有第一次执行时,表不存在才会返回1,不过存不存在都会把表压栈
lua_pushcfunction(L, lfreepool);
lua_setfield(L, -2, "__gc");
}
lua_setmetatable(L, -2); //把一张表弹出栈,并将其设为给定索引处的值的元表。
return 1;
}
//取出sb头结点,清空,并将它塞入到pool[1]里作为新的free_node结点(重新放回到缓冲池里),原来的free_node节点作为它的下一节点
static void
return_free_node(lua_State *L, int pool, struct socket_buffer *sb) {
struct buffer_node *free_node = sb->head; //把head设为free_node
//以下是调整sb链表
sb->offset = 0;
sb->head = free_node->next;
if (sb->head == NULL) {
sb->tail = NULL;
}
lua_rawgeti(L, pool, 1); //pool[1]压栈
free_node->next = lua_touserdata(L, -1); //原来的free_node成为新的free_node的next节点
lua_pop(L, 1); //弹出pool[1]
skynet_free(free_node->msg);
free_node->msg = NULL;
free_node->sz = 0;
lua_pushlightuserdata(L, free_node); //把新free_node压栈
lua_rawseti(L, pool, 1); //弹出刚刚压栈的free_node,并设置为pool[1]
}










