Repository Wiki
nikkigallery/Whimbox

启动流程与服务生命周期

Whimbox 的启动入口由 start_rpc_server() 负责:它建立事件循环上下文,注册事件通知桥接,启动后台恢复与辅助监听任务,随后创建 WebSocket JSON-RPC 服务,并在服务退出时停止脚本监视器。

Purpose and Scope

本页聚焦 RPC 服务从启动到运行、客户端连接管理、请求调度以及退出清理的生命周期。覆盖范围包括 whimbox/rpc_server.py 中的 WebSocket 服务、事件广播、并发任务和启动期辅助组件。

本页不展开 JSON-RPC 各业务方法的具体领域逻辑;这些逻辑由 rpc_method_groups 中的处理器承载,适合在对应的能力页面中单独说明。同样,微信自动恢复、脚本监视器和 overlay 热键监听的内部实现不在本页深入展开,这里只记录它们在生命周期中的接入点。

Overview

启动流程有一个明确的“先建立运行上下文、再启动外围服务”的顺序:

  1. 从 RPC_CONFIG 读取监听地址和端口。
  2. 获取当前 asyncio 事件循环并保存到模块级 _loop,使同步事件来源也能够向该循环投递异步广播。
  3. 通过 set_notifier(notify_event) 把事件总线接到 RPC 通知出口。
  4. 创建 weixin_service.auto_restore() 后台任务,启动 overlay 热键监听和 scripts_watcher。
  5. 使用 websockets.serve(_ws_handler, host, port, max_size=10 * 1024 * 1024) 开始监听。
  6. 通过永不完成的 asyncio.Future() 保持服务运行;退出时进入 finally,停止脚本监视器。

运行期,每个 WebSocket 客户端都有独立的发送锁和待处理任务集合。收到消息后,服务为每条消息创建独立任务,因此同一连接上的请求可以异步处理;发送响应时使用锁,避免多个协程同时向同一个 WebSocket 写入。连接关闭时,所有尚未完成的请求都会被取消并等待回收。

Architecture

Loading diagram...

Source: rpc_server.py Source: rpc_server.py Source: rpc_server.py

图中的关系均对应源码中的导入、模块级生命周期状态或实际调用。start_rpc_server() 是启动编排点;_ws_handler() 是连接级入口;notify_event() 则把事件总线事件转化为广播操作。_loop 和 _client_send_locks 是跨协程协调的关键状态。

启动阶段

监听配置与事件循环绑定

start_rpc_server() 不创建新的事件循环,而是使用 asyncio.get_running_loop() 获取调用它的活动循环。这一点很重要:后续 _notify() 可能在事件循环线程内,也可能从其他线程被调用,因此必须保存真实的 loop 引用。

启动函数同时调用 set_notifier(notify_event),把系统事件的出口指向 RPC 模块。这样,事件生产者不需要直接了解 WebSocket 客户端,只需调用已注册的 notifier;RPC 层负责将事件包装成 JSON-RPC notification 并发送给当前客户端集合。

后台组件的启动顺序

启动期的三个外围动作分别承担不同职责:

  • asyncio.create_task(weixin_service.auto_restore()):以后台任务方式执行自动恢复,不阻塞 WebSocket 监听器建立。
  • _start_overlay_hotkey_listener():安装 overlay 热键监听器。
  • scripts_watcher.start():启动脚本监视器。

源码没有在本文件中等待 auto_restore() 完成,因此 RPC 服务可以先进入监听状态;自动恢复与网络入口是并行生命周期。对于需要严格等待恢复完成的调用方,本文件没有提供这样的启动屏障。

连接与请求生命周期

Loading diagram...

Source: rpc_server.py

连接建立时,WebSocket 对象被加入 _clients,并在 _client_send_locks 中获得一个 asyncio.Lock。消息循环不会直接串行等待处理结果,而是创建任务并将其放入 pending_tasks;任务完成后通过 add_done_callback(pending_tasks.discard) 自动从集合中移除。

连接断开时,finally 块取消所有仍在集合中的任务,使用 asyncio.gather(..., return_exceptions=True) 等待取消结果,然后清理客户端和锁。这种清理顺序避免了连接已经消失后仍持有客户端引用,也避免遗留任务继续尝试写入已关闭的 socket。

事件通知与跨线程投递

_notify() 根据调用位置选择两种投递策略:如果当前调用已经位于保存的 loop 上,则使用 asyncio.create_task(_broadcast(...));如果来自其他线程,则使用 asyncio.run_coroutine_threadsafe(...) 将广播协程安全地提交到目标 loop。

_broadcast() 对客户端逐个发送消息。若该客户端存在发送锁,则使用 async with lock 串行化写操作;发送抛出异常的客户端会被收集到 stale,随后从 _clients 和 _client_send_locks 中移除。广播不会因为某一个失效客户端而停止整个遍历。

python
1async def _broadcast(method: str, params: Dict[str, Any]) -> None: 2 if not _clients: 3 return 4 payload = {"jsonrpc": "2.0", "method": method, "params": params} 5 message = json.dumps(payload, ensure_ascii=False) 6 stale = [] 7 for client in _clients: 8 try: 9 lock = _client_send_locks.get(client) 10 if lock is None: 11 await client.send(message) 12 else: 13 async with lock: 14 await client.send(message) 15 except Exception: # noqa: BLE001 16 stale.append(client) 17 for client in stale: 18 _clients.discard(client) 19 _client_send_locks.pop(client, None)

Source: rpc_server.py

Sources

(1 files)