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 打字"靠什么?前端直连沙箱里 agent-server 的 WebSocket(Day 05 讲过)。所以分工是:实时流走"前端↔沙箱"直连(低延迟);持久化存档走"沙箱→app_server webhook"(可靠);历史回看走"前端→app_server REST"(可查询)。app_server 不掺和实时推送这条链——它专心做"可靠的记账员"。这个分工让每条链路都只做自己最擅长的事。
三条数据链各司其职:app_server 只"存"和"查",不"推" 前端看进度 / 下任务 app_server记账员(存+查) 沙箱agent-server ① 实时流:WebSocket 直连(低延迟,绕开 app_server) ② 持久化webhook 入站→存文件 ③ 历史REST 只读查询 event/ 与 status/ 目录里没有任何 SSE / WebSocket 代码——实时推送不归它管
图注:实时走"前端↔沙箱"直连、持久化走"沙箱→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 支持分页,前端刷新页面重进时就靠它把历史事件拉回来重建界面。
服务抽象 EventServiceevent_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

文件存储与"权限前缀"

FilesystemEventServicefilesystem_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_eventwebhook_router.py:468)收到一批事件后:

  1. 并发落库:把所有事件 save_event 存下来。
  2. 处理特殊事件:stats 事件、执行状态变化、Agent 主动切模型(SwitchLLMObservation)等 → 更新会话记录。
  3. 终态打点:任务结束时 _track_conversation_terminal 做分析埋点。
  4. 后台跑回调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 09Sandbox 沙箱管理——AI 的"隔离工作间"怎么用 Docker 供给、暂停、回收,双容器回调闭环怎么建立(注入 session_key + webhook 地址),以及密钥如何"按需下发"不进镜像。
← Day 07 会话 Day 09 · 沙箱管理 →