启动流程与服务生命周期
Whimbox 的启动入口由 start_rpc_server() 负责:它建立事件循环上下文,注册事件通知桥接,启动后台恢复与辅助监听任务,随后创建 WebSocket JSON-RPC 服务,并在服务退出时停止脚本监视器。
Purpose and Scope
本页聚焦 RPC 服务从启动到运行、客户端连接管理、请求调度以及退出清理的生命周期。覆盖范围包括 whimbox/rpc_server.py 中的 WebSocket 服务、事件广播、并发任务和启动期辅助组件。
本页不展开 JSON-RPC 各业务方法的具体领域逻辑;这些逻辑由 rpc_method_groups 中的处理器承载,适合在对应的能力页面中单独说明。同样,微信自动恢复、脚本监视器和 overlay 热键监听的内部实现不在本页深入展开,这里只记录它们在生命周期中的接入点。
Overview
启动流程有一个明确的“先建立运行上下文、再启动外围服务”的顺序:
- 从
RPC_CONFIG读取监听地址和端口。 - 获取当前
asyncio事件循环并保存到模块级_loop,使同步事件来源也能够向该循环投递异步广播。 - 通过
set_notifier(notify_event)把事件总线接到 RPC 通知出口。 - 创建
weixin_service.auto_restore()后台任务,启动 overlay 热键监听和scripts_watcher。 - 使用
websockets.serve(_ws_handler, host, port, max_size=10 * 1024 * 1024)开始监听。 - 通过永不完成的
asyncio.Future()保持服务运行;退出时进入finally,停止脚本监视器。
运行期,每个 WebSocket 客户端都有独立的发送锁和待处理任务集合。收到消息后,服务为每条消息创建独立任务,因此同一连接上的请求可以异步处理;发送响应时使用锁,避免多个协程同时向同一个 WebSocket 写入。连接关闭时,所有尚未完成的请求都会被取消并等待回收。
Architecture
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 服务可以先进入监听状态;自动恢复与网络入口是并行生命周期。对于需要严格等待恢复完成的调用方,本文件没有提供这样的启动屏障。
连接与请求生命周期
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 中移除。广播不会因为某一个失效客户端而停止整个遍历。
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