Skip to content

异步 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使用对应的 deregisterAcceptFileDescriptorderegisterReadDataFileDescriptor
registerWriteFileDescriptor已废弃,调用会抛 RuntimeError每次发送调用 writeFileDescriptor(fd, data, onWriteComplete)
deregisterWriteFileDescriptor已废弃,调用会抛 RuntimeError不需要手动注销;该 fd 的写请求全部完成后,写侧处理器自动释放。

不要在当前 API 中混用 readiness 用法:onRead 回调已经收到完成数据,不能再对同一 fd 调用 socket.recv();写入也不能通过“监听可写事件”实现。

API 总览

API调用签名回调签名作用
registerAcceptFileDescriptorKBEngine.registerAcceptFileDescriptor(fd, callback)callback(listenerFD, clientFD, errorCode)注册监听 Socket,接收已完成的新连接。
deregisterAcceptFileDescriptorKBEngine.deregisterAcceptFileDescriptor(fd)注销监听 Socket。
registerReadDataFileDescriptorKBEngine.registerReadDataFileDescriptor(fd, callback)callback(fd, data, errorCode)注册已连接 TCP Socket,接收已完成的数据或终止状态。
deregisterReadDataFileDescriptorKBEngine.deregisterReadDataFileDescriptor(fd)注销连接的读侧并清理未消费的接收 completion。
writeFileDescriptorKBEngine.writeFileDescriptor(fd, data, callback)callback(fd, bytesWritten, errorCode)把一次 TCP 写请求加入完成式发送队列。

所有 fd 都应传入 socket.fileno() 返回的整数。引擎拒绝 0、负数、超出平台句柄范围或不是整数的值。

监听 API

registerAcceptFileDescriptor

python
KBEngine.registerAcceptFileDescriptor(listener.fileno(), on_accept)

监听 Socket 必须已经完成 bind()listen(),并建议在注册前设置为非阻塞模式:

python
listener.setblocking(False)
listener.bind(("127.0.0.1", 30040))
listener.listen(128)
KBEngine.registerAcceptFileDescriptor(listener.fileno(), on_accept)

回调:

python
def on_accept(listener_fd, client_fd, error_code):
    pass
参数类型当前语义
listener_fdint已注册的监听 fd。
client_fdint引擎已经 accept 成功并交给脚本的客户端 fd。脚本负责包装和关闭。
error_codeint当前成功 accept 的回调为 0。accept 队列溢出时引擎会关闭无法交付的客户端 fd,不会把泄漏的 fd 交给脚本。

一次引擎唤醒可能对应多个已完成连接,底层会排空当前 accept completion 队列,因此回调可能连续触发多次。回调中不要再次对监听 Socket 调用 accept()

deregisterAcceptFileDescriptor

python
KBEngine.deregisterAcceptFileDescriptor(listener.fileno())

注销只移除引擎的监听处理器,不替脚本关闭 Python Socket。正确顺序是先注销,再关闭 Socket,并确保不会再次使用该 fd。

连接读取 API

registerReadDataFileDescriptor

python
KBEngine.registerReadDataFileDescriptor(client.fileno(), on_read)

同一个 fd 的读侧只能注册一次;监听 fd 和已连接 fd 不能重复占用同一个读槽。当前 API 只支持 TCP completion 数据路径,底层不支持 completion 时会抛 RuntimeError

回调:

python
def on_read(fd, data, error_code):
    pass
条件dataerror_code处理建议
收到普通 TCP 数据非空 bytes0把 bytes 放入协议缓冲区,按应用协议拆包。
对端有序关闭bytes0关闭本地连接并注销读侧。
Socket 错误通常为空 bytes平台错误码记录错误,注销读侧并关闭连接。

注意:data 是引擎已经从 Socket 读取并交接给脚本的数据,不能再调用 client.recv()client.recv_into() 读取同一批数据。TCP 仍然是字节流,一次回调可能只有半个应用包,也可能包含多个应用包。

deregisterReadDataFileDescriptor

python
KBEngine.deregisterReadDataFileDescriptor(client.fileno())

注销后,引擎会停止该 fd 的读事件并清理尚未消费的接收 completion。注销必须先于 socket.close(),避免 fd 被操作系统复用后,迟到事件误投递给新 Socket。

完成式写 API

writeFileDescriptor

python
KBEngine.writeFileDescriptor(fd, data, on_write_complete)

参数要求:

参数类型要求
fdint有效的 TCP Socket fd。写侧不要求先调用 registerReadDataFileDescriptor
databytes必须是 Python bytesbytearraystrmemoryview 不会自动转换。引擎在入队时复制字节数据。
callbackcallable必须可调用,签名为 callback(fd, bytesWritten, errorCode)

回调:

python
def on_write_complete(fd, bytes_written, error_code):
    pass
条件bytes_writtenerror_code
非空数据已完成发送队列处理本次提交的完整字节数0
bytes00,会立即完成回调
参数有效但写入队列拒绝请求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。

推荐的生命周期

text
创建 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()

关服时:

text
deregisterAcceptFileDescriptor(listenerFD)
    -> listener.close()
    -> deregisterReadDataFileDescriptor(clientFD)
    -> client.close()

完整示例:Interfaces HTTP 回调服务器

这个示例只处理一个简单的 HTTP 请求头,重点展示当前完成式 API 的正确使用方式。生产环境需要使用成熟 HTTP 解析器、认证、请求体上限、超时和反向代理。

poller.py

python
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

python
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()

测试:

bash
curl http://127.0.0.1:30040/

常见错误

错误原因修复
调用旧 registerReadFileDescriptor 后出现 RuntimeErrorreadiness API 已废弃。改用 registerAcceptFileDescriptorregisterReadDataFileDescriptor
调用 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

注意事项

  1. Socket 应在注册前设置为非阻塞模式,避免平台适配层在边界条件下进入阻塞系统调用。
  2. 一个 fd 只能注册一个读侧消费者;监听 fd 只能用于 accept,连接 fd 只能用于 data read。
  3. 写请求会复制数据并保留回调引用,不能无限提交大块 bytes;要控制待发送队列、单连接内存和慢客户端。
  4. on_read 可能连续收到多个 completion,也可能只收到半个业务包;应用层必须处理粘包、拆包和最大消息长度。
  5. fd 是底层句柄,不是 Python Socket 对象本身。保存 fd 时要同步保存 Socket 对象,避免对象被回收导致句柄关闭。
  6. 不要让外部网络回调直接修改已销毁或已迁移的 Entity;跨异步边界后重新验证 Entity、会话和请求版本。
  7. 只需要发送 HTTP 请求时,优先使用 KBEngine.urlopen()asyncio 协程,无需自行实现 Socket 协议。

相关文档