异步 Socket 通信
适用范围
KBEngine 为 Python 逻辑层提供完成式文件描述符 API。接口把监听、接收和发送交给引擎的 Dispatcher/EventPoller,完成结果再通过 Python 回调返回。
这套 API 适合:
| 场景 | 建议 |
|---|---|
| Interfaces 接收第三方 HTTP/TCP 回调 | 推荐。Interfaces 通常是单线程,不能使用同步 recv() 或阻塞 HTTP 服务。 |
| BaseApp/CellApp 接收少量外部连接 | 可以使用,但要评估主线程 Tick、连接数和业务处理时间。 |
| 只向外部服务发起 HTTP 请求 | 优先使用 KBEngine.urlopen() 或 asyncio HTTP 客户端,不要重复实现 Socket 协议。 |
| 大型公网网关、TLS 终止、连接数很高 | 更适合独立网关或反向代理,避免把大量连接和协议解析放进逻辑组件。 |
当前受支持平台已经统一为完成式上层模型:Windows 使用 IOCP,Linux 使用 io_uring,macOS 使用基于 kqueue 的 completion adapter。逻辑层使用相同的 accept、read-data 和 write-completion API,不需要判断底层事件后端。
与旧 API 的区别
旧版 readiness API 只通知“fd 当前可读/可写”,脚本还需要自行调用 accept() 或 recv()。当前引擎已经切换为完成式语义:底层先完成 accept/recv/send,再把结果交给脚本。
| 旧 API | 当前状态 | 替代方式 |
|---|---|---|
registerReadFileDescriptor | 已废弃,调用会抛 RuntimeError | 监听 fd 使用 registerAcceptFileDescriptor;连接 fd 使用 registerReadDataFileDescriptor。 |
deregisterReadFileDescriptor | 已废弃,调用会抛 RuntimeError | 使用对应的 deregisterAcceptFileDescriptor 或 deregisterReadDataFileDescriptor。 |
registerWriteFileDescriptor | 已废弃,调用会抛 RuntimeError | 每次发送调用 writeFileDescriptor(fd, data, onWriteComplete)。 |
deregisterWriteFileDescriptor | 已废弃,调用会抛 RuntimeError | 不需要手动注销;该 fd 的写请求全部完成后,写侧处理器自动释放。 |
不要在当前 API 中混用 readiness 用法:onRead 回调已经收到完成数据,不能再对同一 fd 调用 socket.recv();写入也不能通过“监听可写事件”实现。
API 总览
| API | 调用签名 | 回调签名 | 作用 |
|---|---|---|---|
registerAcceptFileDescriptor | KBEngine.registerAcceptFileDescriptor(fd, callback) | callback(listenerFD, clientFD, errorCode) | 注册监听 Socket,接收已完成的新连接。 |
deregisterAcceptFileDescriptor | KBEngine.deregisterAcceptFileDescriptor(fd) | 无 | 注销监听 Socket。 |
registerReadDataFileDescriptor | KBEngine.registerReadDataFileDescriptor(fd, callback) | callback(fd, data, errorCode) | 注册已连接 TCP Socket,接收已完成的数据或终止状态。 |
deregisterReadDataFileDescriptor | KBEngine.deregisterReadDataFileDescriptor(fd) | 无 | 注销连接的读侧并清理未消费的接收 completion。 |
writeFileDescriptor | KBEngine.writeFileDescriptor(fd, data, callback) | callback(fd, bytesWritten, errorCode) | 把一次 TCP 写请求加入完成式发送队列。 |
所有 fd 都应传入 socket.fileno() 返回的整数。引擎拒绝 0、负数、超出平台句柄范围或不是整数的值。
监听 API
registerAcceptFileDescriptor
KBEngine.registerAcceptFileDescriptor(listener.fileno(), on_accept)监听 Socket 必须已经完成 bind() 和 listen(),并建议在注册前设置为非阻塞模式:
listener.setblocking(False)
listener.bind(("127.0.0.1", 30040))
listener.listen(128)
KBEngine.registerAcceptFileDescriptor(listener.fileno(), on_accept)回调:
def on_accept(listener_fd, client_fd, error_code):
pass| 参数 | 类型 | 当前语义 |
|---|---|---|
listener_fd | int | 已注册的监听 fd。 |
client_fd | int | 引擎已经 accept 成功并交给脚本的客户端 fd。脚本负责包装和关闭。 |
error_code | int | 当前成功 accept 的回调为 0。accept 队列溢出时引擎会关闭无法交付的客户端 fd,不会把泄漏的 fd 交给脚本。 |
一次引擎唤醒可能对应多个已完成连接,底层会排空当前 accept completion 队列,因此回调可能连续触发多次。回调中不要再次对监听 Socket 调用 accept()。
deregisterAcceptFileDescriptor
KBEngine.deregisterAcceptFileDescriptor(listener.fileno())注销只移除引擎的监听处理器,不替脚本关闭 Python Socket。正确顺序是先注销,再关闭 Socket,并确保不会再次使用该 fd。
连接读取 API
registerReadDataFileDescriptor
KBEngine.registerReadDataFileDescriptor(client.fileno(), on_read)同一个 fd 的读侧只能注册一次;监听 fd 和已连接 fd 不能重复占用同一个读槽。当前 API 只支持 TCP completion 数据路径,底层不支持 completion 时会抛 RuntimeError。
回调:
def on_read(fd, data, error_code):
pass| 条件 | data | error_code | 处理建议 |
|---|---|---|---|
| 收到普通 TCP 数据 | 非空 bytes | 0 | 把 bytes 放入协议缓冲区,按应用协议拆包。 |
| 对端有序关闭 | 空 bytes | 0 | 关闭本地连接并注销读侧。 |
| Socket 错误 | 通常为空 bytes | 平台错误码 | 记录错误,注销读侧并关闭连接。 |
注意:data 是引擎已经从 Socket 读取并交接给脚本的数据,不能再调用 client.recv() 或 client.recv_into() 读取同一批数据。TCP 仍然是字节流,一次回调可能只有半个应用包,也可能包含多个应用包。
deregisterReadDataFileDescriptor
KBEngine.deregisterReadDataFileDescriptor(client.fileno())注销后,引擎会停止该 fd 的读事件并清理尚未消费的接收 completion。注销必须先于 socket.close(),避免 fd 被操作系统复用后,迟到事件误投递给新 Socket。
完成式写 API
writeFileDescriptor
KBEngine.writeFileDescriptor(fd, data, on_write_complete)参数要求:
| 参数 | 类型 | 要求 |
|---|---|---|
fd | int | 有效的 TCP Socket fd。写侧不要求先调用 registerReadDataFileDescriptor。 |
data | bytes | 必须是 Python bytes,bytearray、str 和 memoryview 不会自动转换。引擎在入队时复制字节数据。 |
callback | callable | 必须可调用,签名为 callback(fd, bytesWritten, errorCode)。 |
回调:
def on_write_complete(fd, bytes_written, error_code):
pass| 条件 | bytes_written | error_code |
|---|---|---|
| 非空数据已完成发送队列处理 | 本次提交的完整字节数 | 0 |
空 bytes | 0 | 0,会立即完成回调 |
| 参数有效但写入队列拒绝请求 | 0 | 平台错误码,回调可能在 writeFileDescriptor() 返回前同步发生 |
同一 fd 可以连续提交多个写请求。引擎按提交顺序排队,并为每个请求调用一次完成回调。写回调中再次提交的请求会留到下一次底层完成阶段处理,不会破坏当前请求快照。
bytes_written 表示本次 API 请求对应的数据长度,不是某一次底层 send() 系统调用返回的分段长度。当前每个 fd 的 completion TCP 待发送数据上限为 1 MiB;超过上限或单次数据本身过大时,写请求会以平台的 would-block 或 message-too-large 错误返回。业务仍应设置更小、更符合协议的单连接上限。
底层已经入队后的异步 Socket 错误还可能通过 on_read(fd, b"", error_code) 进入读侧终止路径,因此不能把写完成回调当作唯一的断线检测。读回调和写回调都应使用同一个幂等关闭函数。
写侧没有单独的注销 API。所有排队请求完成后,写处理器自动释放。若连接需要关闭,应先停止提交新的写请求,再注销读侧并关闭 Socket;不要在写完成回调中无条件重复关闭同一个已被其他错误路径关闭的 fd。
推荐的生命周期
创建 listener
-> setblocking(False)
-> bind/listen
-> registerAcceptFileDescriptor
-> on_accept(listenerFD, clientFD, 0)
-> 用 socket.socket(fileno=clientFD) 包装客户端 Socket
-> setblocking(False)
-> registerReadDataFileDescriptor
-> on_read(fd, data, errorCode)
-> writeFileDescriptor(fd, response, on_write_complete)
-> on_write_complete(fd, bytesWritten, errorCode)
-> deregisterReadDataFileDescriptor
-> client.close()关服时:
deregisterAcceptFileDescriptor(listenerFD)
-> listener.close()
-> deregisterReadDataFileDescriptor(clientFD)
-> client.close()完整示例:Interfaces HTTP 回调服务器
这个示例只处理一个简单的 HTTP 请求头,重点展示当前完成式 API 的正确使用方式。生产环境需要使用成熟 HTTP 解析器、认证、请求体上限、超时和反向代理。
poller.py
import socket
import KBEngine
from KBEDebug import DEBUG_MSG, ERROR_MSG, INFO_MSG
class Poller:
def __init__(self):
self._listener = None
self._clients = {}
def start(self, host, port):
if self._listener is not None:
return
listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
listener.setblocking(False)
listener.bind((host, port))
listener.listen(128)
self._listener = listener
try:
KBEngine.registerAcceptFileDescriptor(listener.fileno(), self.on_accept)
except Exception:
self._listener = None
listener.close()
raise
INFO_MSG("Poller::start: listen %s:%s" % (host, port))
def stop(self):
listener = self._listener
self._listener = None
if listener is not None:
KBEngine.deregisterAcceptFileDescriptor(listener.fileno())
listener.close()
for fd in list(self._clients):
self.close_client(fd)
def on_accept(self, listener_fd, client_fd, error_code):
if error_code != 0:
ERROR_MSG("Poller::on_accept: listenerFD=%i error=%i" %
(listener_fd, error_code))
return
try:
client = socket.socket(fileno=client_fd)
client.setblocking(False)
client.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
except Exception:
ERROR_MSG("Poller::on_accept: wrap client fd failed: %i" % client_fd)
# client_fd 已由引擎交给脚本;包装失败时必须关闭它。
try:
socket.close(client_fd)
except Exception:
pass
return
self._clients[client_fd] = {
"socket": client,
"buffer": bytearray(),
"responded": False,
}
try:
KBEngine.registerReadDataFileDescriptor(client_fd, self.on_read)
except Exception:
self.close_client(client_fd)
raise
DEBUG_MSG("Poller::on_accept: new clientFD=%i" % client_fd)
def on_read(self, fd, data, error_code):
client = self._clients.get(fd)
if client is None:
return
if error_code != 0 or not data:
self.close_client(fd)
return
if client["responded"]:
return
client["buffer"].extend(data)
if b"\r\n\r\n" not in client["buffer"]:
return
client["responded"] = True
self.send_response(fd)
def send_response(self, fd):
body = b"Hello KBEngine completion API\n"
response = (
b"HTTP/1.1 200 OK\r\n"
b"Content-Type: text/plain; charset=utf-8\r\n"
b"Content-Length: " + str(len(body)).encode("ascii") + b"\r\n"
b"Connection: close\r\n"
b"\r\n" + body
)
try:
KBEngine.writeFileDescriptor(fd, response, self.on_write_complete)
except Exception:
self.close_client(fd)
def on_write_complete(self, fd, bytes_written, error_code):
if error_code != 0:
ERROR_MSG("Poller::on_write_complete: fd=%i error=%i" %
(fd, error_code))
else:
DEBUG_MSG("Poller::on_write_complete: fd=%i bytes=%i" %
(fd, bytes_written))
self.close_client(fd)
def close_client(self, fd):
client = self._clients.pop(fd, None)
if client is None:
return
try:
KBEngine.deregisterReadDataFileDescriptor(fd)
except Exception:
pass
try:
client["socket"].close()
except Exception:
pass
DEBUG_MSG("Poller::close_client: fd=%i" % fd)kbemain.py
import KBEngine
from KBEDebug import INFO_MSG
from poller import Poller
_poller = Poller()
def onInterfaceAppReady():
INFO_MSG("onInterfaceAppReady")
_poller.start("127.0.0.1", 30040)
def onInterfaceAppShutDown():
INFO_MSG("onInterfaceAppShutDown")
_poller.stop()测试:
curl http://127.0.0.1:30040/常见错误
| 错误 | 原因 | 修复 |
|---|---|---|
调用旧 registerReadFileDescriptor 后出现 RuntimeError | readiness API 已废弃。 | 改用 registerAcceptFileDescriptor 或 registerReadDataFileDescriptor。 |
调用 registerWriteFileDescriptor 后出现 RuntimeError | 当前写入使用完成回调,不再注册“可写”通知。 | 改用 writeFileDescriptor(fd, data, callback)。 |
on_read 中再次 recv() | completion 后端已经把数据读入交接队列。 | 只消费回调参数 data。 |
| 写入前强制注册读侧 | 误把旧文档前置条件带入新 API。 | 写侧可以独立提交;但 fd 必须有效且由当前组件拥有。 |
| 关闭 Socket 后偶发 fd 错误 | 先 close 后注销,或 fd 已被系统复用。 | 先注销对应读/监听 API,再关闭 Socket。 |
| 收到半个 HTTP 请求就返回 | TCP 没有消息边界。 | 使用缓冲区累积并实现拆包,不能假设一次回调就是一个请求。 |
| 大量连接导致 Tick 变慢 | 回调中做了同步 IO、复杂解析或一次处理过多数据。 | 限制单次处理量,拆批,或迁移到独立网关。 |
| 读回调异常后连接状态混乱 | 回调异常没有统一清理。 | 捕获业务异常,记录日志,并走幂等的 close_client。 |
注意事项
- Socket 应在注册前设置为非阻塞模式,避免平台适配层在边界条件下进入阻塞系统调用。
- 一个 fd 只能注册一个读侧消费者;监听 fd 只能用于 accept,连接 fd 只能用于 data read。
- 写请求会复制数据并保留回调引用,不能无限提交大块 bytes;要控制待发送队列、单连接内存和慢客户端。
on_read可能连续收到多个 completion,也可能只收到半个业务包;应用层必须处理粘包、拆包和最大消息长度。- fd 是底层句柄,不是 Python Socket 对象本身。保存 fd 时要同步保存 Socket 对象,避免对象被回收导致句柄关闭。
- 不要让外部网络回调直接修改已销毁或已迁移的 Entity;跨异步边界后重新验证 Entity、会话和请求版本。
- 只需要发送 HTTP 请求时,优先使用
KBEngine.urlopen()或 asyncio 协程,无需自行实现 Socket 协议。
