Day 08 / 共 20 天 · 第 2 周 app_server 编排层
事件系统源码
事件是 OpenHands 的血液。今天读 app_server 侧的事件系统——你会发现它只做持久化和查询,不做实时推送。为什么?这背后是个漂亮的架构取舍。
📍 第 2 周 后端编排层 · 你在这里
Day6 app_server总览→
Day7 会话生命周期→
Day8 事件系统→
Day9 Sandbox→
Day10 server装配
L01
一个反直觉的事实
你可能以为 app_server 会用 WebSocket/SSE 把事件实时推给前端。但实际上:event/ 和 status/ 目录里没有任何 SSE / WebSocket / event-stream 代码。app_server 侧的事件只有两件事:
入站webhook
agent-server 推事件进来 → 持久化存成文件 → 出站REST 只读
前端查历史
agent-server 推事件进来 → 持久化存成文件 → 出站REST 只读
前端查历史
那"实时看 Agent 打字"靠什么?
靠前端直连沙箱里 agent-server 的 WebSocket(Day 05 讲过)。所以分工是:实时流走"前端↔沙箱"直连(低延迟);持久化存档走"沙箱→app_server webhook"(可靠);历史回看走"前端→app_server REST"(可查询)。app_server 不掺和实时推送这条链——它专心做"可靠的记账员"。这个分工让每条链路都只做自己最擅长的事。
图注:实时走"前端↔沙箱"直连、持久化走"沙箱→app_server webhook"、历史走"前端→app_server REST"——三链解耦,各做最擅长的事。
L02
事件读取 REST 接口
event_router.py,前缀 /conversation/{conversation_id}/events,三个只读端点:
# event_router.py
GET .../events/search # 搜事件(按 kind/时间过滤 + 分页)
GET .../events/count # 数事件总数
GET .../events # 批量取事件
读法:全是 GET(只读)——app_server 不提供"写事件"的公开 API。为什么?因为事件只能由 agent-server 产生(通过 webhook 写入,L05),前端只能"查"不能"造"。
search 支持分页,前端刷新页面重进时就靠它把历史事件拉回来重建界面。服务抽象
EventService(event_service.py:18)的方法:get_event/search_events/count_events/save_event(注释明确 save "不属于 REST API",仅内部 webhook 用)+ iter_events_for_export(导出整个会话)。读写分离得很干净。L03
存储的"模板方法"设计
event_service_base.py 是存储基类,用模板方法模式:基类写好 search/count/save 的编排逻辑,把三个"怎么读写一条"的原语留给子类实现:
# 基类实现编排,子类填原语(概念)
class EventServiceBase(EventService):
async def search_events(self, ...): # 编排:并发读所有文件→内存过滤→排序→分页
paths = self._search_paths(...) # ← 子类实现
events = await self._load_events_from_paths(paths) # 并发+信号量限流
return 过滤排序分页(events)
# 三个留给子类的原语:
async def _load_event(self, path): ... # 子类:怎么读一条
async def _store_event(self, event): ... # 子类:怎么写一条
def _search_paths(self, ...): ... # 子类:去哪找
读法:基类负责"搜索该怎么并发、怎么排序分页"这套通用流程;子类只负责"一条事件具体怎么读、怎么写、存哪"。于是新增一种存储后端(文件→S3→GCS),只要实现三个原语,编排逻辑完全复用。
模板方法模式,一句话
"骨架在父类,血肉在子类。"父类定好"做事的步骤顺序"(模板),把每一步的"具体怎么做"开成空方法让子类填。就像"泡茶模板:烧水→放茶叶→倒水→等待",具体用什么茶、水温多少由子类定。这样"流程"只写一遍,"细节"按需扩展。
L04
文件存储与"权限前缀"
FilesystemEventService(filesystem_event_service.py:17):一个事件存一个 JSON 文件,用 Pydantic 的 model_validate_json/model_dump_json 读写。路径设计是关键(event_service_base.py:66):
def get_conversation_path(self, ...):
# 路径 = {prefix}/{user_id}/{V1_CONVERSATIONS_DIR}/{conversation_id.hex}
return f'{prefix}/{user_id}/conversations/{conversation_id.hex}'
读法:存储路径里嵌入了 user_id。源码注释点明:权限控制的唯一手段,就是这个带 user_id 的严格路径前缀。查某会话事件时,路径由"当前登录用户的 user_id"拼成——你根本拼不出别人的路径,也就读不到别人的事件。
用"路径前缀"做权限,聪明在哪?
不用额外写一堆"这个用户能不能看这个会话"的检查代码——把用户隔离直接编码进存储路径。你的事件在
/data/你的id/...,别人在 /data/他的id/...,物理上就分开了。想越权?你手里的 user_id 拼出的路径永远指向你自己的目录。简单、不易出错。把安全约束"结构化"进数据布局,好过散落各处的 if 检查。(云存储版 AwsEventService/GoogleCloudEventService 同理,只是把文件换成对象存储。)L05
webhook:事件"进来"的地方
事件从哪写入?webhook_router.py(前缀 /webhooks)。沙箱里的 agent-server 每产生事件就回调这里。鉴权靠 valid_sandbox(:250)校验 X-Session-API-Key 头:
# webhook_router.py
@router.post('/webhooks/events/{conversation_id}')
async def on_event(conversation_id, events, sandbox=Depends(valid_sandbox), ...):
# valid_sandbox:校验 X-Session-API-Key → 反查 sandbox → 切换成沙箱属主用户上下文
...
读法:agent-server 带着启动时注入的
session_api_key 来敲门,valid_sandbox 验钥匙、找到对应沙箱、把请求上下文切换成"这个沙箱的属主用户"。于是后续存事件时,就能用正确的 user_id 拼路径(L04)。钥匙 + 用户上下文切换 = 安全地接收来自沙箱的事件。L06
on_event 做的四件事
on_event(webhook_router.py:468)收到一批事件后:
- 并发落库:把所有事件
save_event存下来。 - 处理特殊事件:stats 事件、执行状态变化、Agent 主动切模型(
SwitchLLMObservation)等 → 更新会话记录。 - 终态打点:任务结束时
_track_conversation_terminal做分析埋点。 - 后台跑回调:
background_tasks.add_task(_run_callbacks_in_bg_and_close, ...)。
回调处理器(
EventCallbackProcessor 抽象)里最典型的是 SetTitleCallbackProcessor——根据前几条消息自动生成会话标题(就像 ChatGPT 自动给对话起名)。Day 07 讲过"启动时总是确保挂上这个回调",现在看到它在这里被触发。还记得 Day 07 的"薄代理原则"吗?——业务副作用统一走这里的 webhook 回调,而不是塞进代理端点。L07
"回调不持有 DB 连接"的性能设计
_run_callbacks_in_bg_and_close(:583)的注释解释了一个精妙的性能考量:回调处理器运行时不持有 DB 连接。
为什么这很重要?
想象 Agent 高速运转,每秒产生几十个事件 → 每秒几十次 webhook → 如果每个 webhook 的回调都"占着一个数据库连接慢慢跑",数据库连接池瞬间被打满,整个服务卡死。所以设计成:webhook 主流程快速存完事件、释放连接,回调逻辑丢到后台不占连接地跑。面对突发流量,"尽快释放稀缺资源"是保命的设计。
还有个模块加载即执行的细节:
_import_all_tools()(:630)在模块加载时就预先 import 所有工具类型。为什么?因为 webhook 收到的事件 JSON 里含工具相关的多态类型,反序列化 Event 时必须所有类型都已注册,否则会失败。提前 import = 保证反序列化时类型齐全。另有 on_conversation_update(:348)处理会话 start/pause/resume/delete 的回调。L08
今日小结 + 动手
🧠 今天你应该能回答
- app_server 侧事件为什么"只存不推"?实时流靠什么?
- 事件 REST 为什么全是只读?谁能写事件?
- 模板方法模式在存储里怎么体现?
- 为什么用"带 user_id 的路径前缀"做权限隔离?
- 为什么回调要"不持有 DB 连接"地后台跑?
✋ 动手
grep -n 'search\|count\|save_event' openhands/app_server/event/event_service.py
grep -n 'get_conversation_path\|user_id\|_load_event\|_store_event' openhands/app_server/event/event_service_base.py
grep -n 'def on_event\|valid_sandbox\|_run_callbacks_in_bg\|_import_all_tools' openhands/app_server/event_callback/webhook_router.py
明天预告 · Day 09:Sandbox 沙箱管理——AI 的"隔离工作间"怎么用 Docker 供给、暂停、回收,双容器回调闭环怎么建立(注入 session_key + webhook 地址),以及密钥如何"按需下发"不进镜像。