面试知识库

M8 · AGUI 事件 + WebSocket —— 开发文档(面试向)#

这篇讲的是为什么这么做、做了哪些取舍,不讲具体代码。目标是读完能用大白话讲出来。

一句话概括#

M8 解决的是一个体验问题:Agent 跑一次购物任务要十几二十秒,这中间如果前端一片空白, 用户会以为「系统挂了」,甚至反复点发送。M8 让 Agent 每走一步——拆需求、跨平台搜、比价、 精挑——都实时往前端推一条消息,用户能看着它一步步干活。

做法是两件事:HTTP 只负责把任务踢到后台、立刻返回;过程靠一条 WebSocket 长连接持续推 事件。 再配一张「谁的事件该推给谁」的路由表(ConnectionManager)和一个统一的上报入口 (monitor)。


1. 为什么不能用「发个 HTTP 等结果回来」#

最直觉的做法是:前端发一个 HTTP 请求,服务端跑完整个任务再把结果返回。问题是这一跑就是 15-20 秒。这期间:

  • 用户什么都看不到——不知道是在跑、卡住了、还是死了。
  • 没法取消——只能干等到超时。
  • 出错了不知道卡在哪——一个黑盒,调试也难。

打个比方:这就像点了外卖,下单后 App 一直转圈、既不显示「商家接单」也不显示「骑手取餐」, 你只能盯着空屏幕猜。体验上「能看见进度」比「快一两秒」重要得多。

解法:把「启动」和「过程」拆开。HTTP 干一件最快的事——把任务丢进后台、立刻返回一个 任务编号(thread_id)。然后前端拿这个编号建一条 WebSocket 长连接,后台每走一步就往这条 连接推一个事件。HTTP 管「开始」,WebSocket 管「直播过程」。

面试怎么讲:长任务别用同步 HTTP 死等;拆成「异步启动 + 长连接推进度」,把黑盒变成直播。

注:M8 落地的是「直播」这套底层能力(事件协议 + 路由 + 上报)。真正对外的 HTTP 接口 (启动任务 / 取消 / WS 端点)是 M10 的活,M8 用一个最小示例 app 把 WS 这条路验通。


1.5 为什么选 WebSocket 而不是 SSE#

排除了同步 HTTP 之后,「服务端持续向前端推事件」有两个主流选项:SSE(Server-Sent Events)WebSocket。先说结论:两个都能做,SSE 甚至在某些维度更简单。选 WebSocket 不是因为 SSE「不行」,而是综合权衡后觉得 WebSocket 在这个项目里更合适。下面诚实地拆。

先承认两个事实#

事实一:当前 WebSocket 几乎是单向的。 看一眼代码就知道——取消任务走的是 POST /api/task/{thread_id}/cancel(独立 HTTP 端点),不走 WebSocket。客户端在 WS 连接上发的唯一消息就是 "ping" 心跳。所以现状下 WebSocket 实质上被当成了「服务端 单向推事件 + 简单 ping/pong」,和 SSE 的能力高度重叠。

事实二:行业主流是 SSE,不是 WebSocket。 ChatGPT / Claude / Gemini / Grok 的文本 对话流式输出全部用 SSE(text/event-stream)。OpenAI 唯一用 WebSocket 的是 Realtime API(语音/实时音频),那是完全不同的场景。所以如果有人说「Agent 对话产品几乎都用 WebSocket, 这是生态惯例」——恰恰相反,生态惯例是 SSE。

这两个事实摆出来,「双向通信是刚需」和「生态惯例」这两条理由都站不住。下面说真正的原因。

真正的原因:长连接 + 多事件类型 + Connect-First 协议#

ShoppingX 和上面四家的场景有一个关键区别:ChatGPT / Claude 等的 SSE 是「一问一答」模型 ——发一个请求、流式返回 token、结束断开。每轮对话是一条独立的 SSE 连接,生命周期很短。

ShoppingX 的 Agent 任务不是「流式返回一段文本」,而是:

  • 一个长任务持续 15-30 秒,中间推多种异构事件(思考中 / 工具开始 / 工具结束 / fork / 最终结果),不是同质的 token 流
  • 前端先建连接、等 ws_ready 确认、再 POST 启动任务(Connect-First 协议),需要 连接在任务启动之前就存在
  • 多轮对话期间连接保持不断,同一条连接上推多个任务的事件
  • 断线后用 last_event_id 重连补发(Redis Stream)

这套模型的核心是:连接的生命周期远长于单次请求,是「会话级」而非「请求级」。

SSE 天然是「请求级」的——每条 SSE 流绑定一个 HTTP 请求,请求结束流就断。要做「会话级长 连接」需要在 SSE 之上再搭一层:前端主动维护一条不断开的 SSE 流、服务端对这条流做连接管理、 处理「流断了但会话还在」的状态。能做,但本质上是在用 SSE 模拟 WebSocket 的连接语义。

WebSocket 的连接模型天然就是「建一次、保持住、双方随时收发、显式关闭」。服务端有明确的 open / close 生命周期钩子,ConnectionManager 能精确知道「这条连接活着还是死了」、按 对象身份做注销防重连误删。这些在 SSE 里也能实现(EventSourceonopen / onerror), 但服务端对断连的感知不如 WebSocket 直接——SSE 断开时服务端看到的是「HTTP 响应流被中断」, 没有显式的 close frame。

一句话:ChatGPT 们用 SSE 是因为它们的流式输出是请求级的(一问一答),SSE 天然匹配。 ShoppingX 的事件推送是会话级的(长连接 + 多事件类型 + 先连后发),WebSocket 的连接语义 更直接。

SSE 在哪些场景是更好的选择#

场景推荐原因
LLM 逐 token 流式输出(一问一答)SSE一次请求、流式返回、结束即断,天然匹配
纯通知推送(日志流、告警墙)SSE只出不进,EventSource 几行搞定
需穿透不支持 WS 的企业代理/CDNSSE就是 HTTP,兼容性最好
长连接 + 多异构事件 + 先连后发WebSocket连接生命周期=会话,open/close 语义显式

诚实边界#

如果 ShoppingX 的模型也是「每轮对话开一条 SSE、推完断掉」,那 SSE 就够了,跟 ChatGPT 一样。 选 WebSocket 是因为我们选了一个不同的交互模型——Connect-First + 会话级长连接——而这个 交互模型用 WebSocket 表达更自然。不是 SSE 做不到,是用 SSE 做会话级长连接需要在上面 补一层连接管理,而 WebSocket 自带这层语义。

代价是放弃了 SSE 的「零配置断线重连」和更简单的服务端实现——但前者我们已经用 Redis Stream 自建了更强的重放机制(ENH-D),后者在 FastAPI 里 WebSocket 和 SSE 的实现复杂度差距不大。

面试怎么讲:ChatGPT、Claude 这些产品用 SSE 是因为它们的流式输出是「一问一答」模型—— 一条请求流式返回、结束断开,SSE 天然匹配。ShoppingX 不一样:Agent 任务要十几秒,推多种异构 事件,而且前端先建连接再启动任务(Connect-First),需要一条会话级的长连接。这种「先连后发、 长时间保持、多事件推送」的模式用 WebSocket 更自然,因为 WS 的 open/close 生命周期是显式 的,服务端对连接状态的管理更直接。用 SSE 也能做,但要在上面额外搭一层连接管理来模拟 WebSocket 的语义。


2. 一条事件长什么样:为什么要「统一信封」#

我定义了八种事件:会话创建、思考中、工具开始、工具结束、派发子任务(fork)、任务完成、 任务取消、出错。但不管哪种,外层信封都一模一样:类型固定写 monitor_event、带一个 event 字段说明是哪种、一个给人看的 message、一个装业务数据的 data、还有时间戳和 任务编号。

为什么要统一:前端只需要写一套「收到消息 → 看 event 字段 → 分别渲染」的逻辑,而不是 为每种事件写一套解析。新增一种事件类型时,前端加一个分支就行,老逻辑不用动。这就像所有 快递不管装什么,外面的面单格式都一样——分拣的人只看面单,不用拆箱。

一个具体取舍data 里的长文本(比如子任务的需求描述、最终清单)我做了截断 (超过 2000 字砍掉加个省略标记)。因为这些事件是给前端「看个进度」的,不是传全量数据; 不截断的话,一条几万字的最终答案塞进事件,既占带宽又可能把前端卡住。全量结果该走正经的 结果接口,不该混在进度事件里。


3. thread_id:串起整条链路的那把钥匙#

整个系统里有个反复出现的编号——thread_id(任务编号)。它一个人串起了好几件事:

  • 前端的 WebSocket 连接按它建立;
  • 后台任务表按它登记(方便取消);
  • 会话产物目录按它隔离;
  • 子 Agent 的执行上下文按它区分。

为什么这把钥匙这么关键:它是「别串台」的根。如果处理不好,会出现 A 用户的进度推给了 B 用户、或者 A 的子任务把文件写进了 B 的目录。多用户同时在线时,这种串台是最难查的 bug。

怎么做到不用层层传参:我没有把 thread_id 当参数一路传到每个工具里——那样既啰嗦又容易 漏。而是用了 Python 的「上下文变量」(ContextVar,M0 就铺好了):在任务入口写一次,之后 无论调用链多深,工具内部要上报事件时「就近」就能读到当前是哪个任务。工具因此完全不必关心 thread_id 是什么、连接在哪——它只管喊一句「我开始干 X 了」,剩下的路由是底层的事。

面试怎么讲:用一个贯穿全链路的 ID 把「连接 / 任务 / 目录 / 上下文」对齐;用上下文变量 做隐式传递,避免把这个 ID 当参数污染每一层函数签名。


4. ConnectionManager:一张「谁的事件推给谁」的路由表#

它本质就是一个字典:任务编号 → 那条 WebSocket 连接。上报事件时拿编号查表,找到连接推过去。 简单,但有两个坑必须处理:

坑一:刷新页面导致的「误删」。 用户刷新页面会建一条连接,而连接的「断开」 回调往往晚一点才触发。如果断开时只认编号、盲目地把这个编号从表里删掉,就会把刚建好的 新连接一起删了——结果用户刷新完反而收不到事件了。

我的解法:断开时不只看编号,还要核对是不是同一个连接对象。只有「表里登记的就是正在断开 的这一条」时才删;如果表里已经是新连接了,旧连接的断开就什么都不做。这就像门禁换了卡,旧卡 来注销时系统发现「现在登记的是新卡」,就不动它。

坑二:并发改表 + 慢连接拖累全局。 服务是单线程协程并发,多个任务会交替读写这张表,所以 用了把锁串行化增删,避免一边遍历一边被改。但有个讲究:真正往连接里写数据(网络 IO)这一步 放在锁外面做。如果握着锁去等一条慢连接发送完,会把所有其他任务的上报都堵住——一条卡顿的 连接拖垮全场。所以是「握锁拿到连接对象 → 放锁 → 再慢慢发」。

另外,发送失败(连接已经死了)时,我会顺手把这条死连接从表里摘掉,免得越攒越多。

面试怎么讲:路由表两个细节——注销要按对象身份核对防重连误删;写 IO 放锁外防慢连接 阻塞全局,锁只保护「改表」这一瞬。


5. 监控绝不能拖垮主任务:上报失败要「咽下去」#

这是一条贯穿 monitor 的原则:上报事件失败,绝不能让主任务跟着崩。

场景:用户把浏览器标签关了,连接断了,但后台任务还在跑。这时工具上报事件会发送失败。如果这个 失败往上抛,就会把正在好好干活的 Agent 给搞崩——本末倒置,监控反而害了主角。

所以上报这条链路从头到尾都是「失败就记个日志、安静降级」:没有任务编号(比如离线脚本/评测在 跑)就只打日志不推送;有编号但没有活跃连接就直接跳过;推送过程中报错就摘掉死连接、返回失败但 不抛异常。监控是配角,配角出问题不能影响主角。

一个有意的设计:即使没有任何前端连着,事件也照样写进调试日志。这样离线评测时翻日志就能复盘 「这次 Agent 到底依次调了哪些工具」,不依赖前端也能观测。


6. 子 Agent 的 fork 事件:推给「父任务」而不是子任务#

主 Agent 会按需 fork 出同质子 Agent 去并行干活(比如同时搜六个平台,M2 做的)。这时该不该 上报、推给谁,有个讲究。

做法:在子 Agent 真正跑起来之前就上报一条 fork 事件,而且这条事件路由到父任务的 连接。原因是:用户的浏览器连的是父任务那条线,子任务有自己独立的编号、根本没有前端连着它。 所以「分叉出一个子任务」这件事,要在还处于父任务上下文的那一刻喊出来,才能落到用户眼前。

一个诚实的边界:子 Agent 内部再调工具产生的那些事件,因为挂在子任务编号上、而子任务没有 前端连接,所以对用户是静默的。这其实和「子 Agent 要做上下文隔离、不污染主 loop」是一致 的——用户看到的是「分出去一个并行子任务」这个粗粒度信号,而不是子任务内部的每一步细节。要不 要把子任务内部也透传给前端,是后续可以加的(比如给并行批次一个共同标记让前端分组显示),M8 先把粗粒度做对。

面试怎么讲:fork 事件在进子上下文之前发、落在父连接上,因为前端连的是父任务;子任务内部 静默,正好契合上下文隔离。


7. 踩过的坑:测试不该依赖一个「本地才有」的配置文件#

做 M8 时跑测试,一个 M2 就写好的、跟本次改动八竿子打不着的测试突然红了。查下来发现:它间接 要构造大模型客户端,而构造时要读几个配置(模型名、密钥、地址)。这些配置在开发者本地的 .env 文件里有,但 .env 是不进仓库的,到了全新的工作区或 CI 就没有,于是构造直接报「找不到配置」。

这其实是个早就埋着的隐患——之前能过,只是因为开发者本地正好有 .env。测试依赖一个不在 仓库里、换台机器就没有的东西,本身就不健康。

解法:给测试加了一组「哑配置」(假的模型名/密钥/地址),让客户端能构造出来——注意构造客户端 本身不发网络请求,只有真去调模型才会,而这些测试要么用假模型替身、要么根本不会真调。而且这组 哑值用的是「没有才填」的方式,本地有真 .env 就不覆盖。这样测试自己就能跑,不再挂靠某台机器 的私有文件。

面试怎么讲:单元测试要「自包含」,不能依赖某台机器上的私有配置;缺什么就在测试里造一个不发真 请求的替身,让任何环境都能一键跑绿。


8. 这一步做了什么、没做什么(诚实边界)#

做了

  • 八种 AGUI 事件 + 统一信封结构;模块级上报入口(工具侧从 M4 起一行没改,桩直接接真)。
  • 路由表 ConnectionManager:按对象身份注销防重连误删、锁内改表 / 锁外发送、死连接自动摘除。
  • fork 事件接入主链路,路由到父任务。
  • 一个端到端示例:真起一个 WebSocket 服务、真用客户端订阅,验证十个事件按序推出。

后续增强(已落地)

  • 工具思考结果摘要tool_end 事件现在带一个可选的 result 字段——每件工具在执行完成后把关键发现拼成一段人读摘要(planner 输出拆解后的结构化字段、item_search 输出 Top 候选标题、price_compare 输出归一价排序、shipping_calc 输出到手价、item_picker 输出精选理由、category_insight 输出爆款/价位/维度分布、web_search 输出 Tavily 摘要)。前端 ActivityFeed 优先展示 result(这一步查到了什么),无 result 时退回入参/元信息的 k=v 串。
  • token 用量事件task_result 事件新增 tokens 字段(input / output / total / cost_usd),是全树(主 + 各 fork 子 Agent)记账口径。前端在每轮右下角与「用时」并排显示「token 消耗」,hover 看输入/输出/成本拆分。
  • 事件回放(ENH-D):补上了断线重连后的事件补发(Redis Stream),包括后续的冷启动续看路径。
  • WebSocket 双向通信 + 用户澄清:见下文 §9。

M8 原始边界(已由后续收口)

  • 对外正式 HTTP 接口 → M10 已接上 FastAPI + React 前端。
  • 断线重连补发 → ENH-D(Redis Stream + last_event_id + 去重)已覆盖。
  • WebSocket 只有单向推送 → §9 已接上双向通信(用户澄清)。
  • 子任务内部事件透传给前端 / 并行批次分组——仍只发粗粒度的 fork 信号(够用、不加噪)。

一句话:M8 把「Agent 在干什么」这件事变得实时可见,并且把可见性这套底层能力(协议 + 路由 + 上报)做扎实、做到「监控崩了也不连累主任务」,为 M10 接上真前端铺好路。后续增强(工具 result 摘要、token 用量、事件回放、双向通信)都是在这套事件协议上叠加字段和能力,信封格式不变。


9. WebSocket 双向通信与用户澄清#

§1.5 诚实承认了 WebSocket「几乎是单向的」这个事实。这一节补上双向的那一半,并落地第一个真实用例:Agent 在跑任务的过程中暂停、向用户提问、拿到回复后继续

背景:原来的澄清方式为什么不行#

Agent 遇到关键信息缺失(比如”帮我找同款”但从没说过任何具体商品)时,原来只有两条路:

  1. 猜一个品类硬搜——大概率不是用户要的,白费一圈检索。
  2. chat_fallback 反问——但 chat_fallback 是终结性工具,调了就结束循环。用户得重新发一轮任务才能继续,中间辛苦做的 planner 拆解、品类洞察全丢了

缺的是一个非终结性的澄清:问完了 Agent 继续跑,已有上下文不丢。

方案:三个零件#

零件一:ask_user 工具。 对 Agent 来说就是另一个工具——和 planneritem_search 平级,在 FULL_TOOL_SET 里。Agent Think 阶段觉得信息不够时调它,传入问题,拿到回复文本,继续 Observe → Reflect。不是终结性工具:调完循环继续。

零件二:asyncio.Future 做阻塞桥梁。 工具协程和 WS handler 协程在同一个事件循环但不同调用栈里(不共享 ContextVar),用模块级 dict[thread_id → Future] 桥接——工具侧创建 Future 并 await;WS 侧收到回复后按 thread_id resolve。单进程 asyncio 下 Future 就是最轻量的跨协程信号,不需要消息队列或 Redis。

零件三:WS 消息循环扩展。 原来只认 "ping" 纯文本。改成先尝试 JSON 解析,认到 {"type": "clarification_response", "text": "..."} 就 resolve 对应 Future。

全部复用现有基础设施——没有新的进程、服务或中间件层。

几个关键取舍#

深度守卫:只有主 Agent 能问用户。 子 Agent 调 ask_user 直接返回硬挡文案——子是被派去搜单个平台的,五个子同时问用户岂不乱套。和 DEPTH0_ONLY_TOOLS 是同一套”能力同质但授权不同质”的思路。

超时: 120 秒不回复返回兜底文案,Agent 自己判断继续还是收尾。不无限等——用户可能已经关了浏览器。

清理: 任务结束(正常/超时/取消)时 cancel_pending 收掉 Future,防泄漏。

新事件类型而非复用 tool_start 普通 tool_start 在 ActivityFeed 加一行”运行中”就够了;clarification_request改变交互模态——激活输入框、显示问题 banner。语义差别大到不该塞同一个事件类型。

不改 prompt 的主动性纪律: System prompt 保留”品类偏宽时直接搜,不要停下来反问”。ask_user 的 docstring 和 few-shot 示例都强调”只在硬信息确实缺失时用”——两条线夹出一个窄走廊:“真推不了了才问,问了还能继续”。

前端:一个新状态#

加了 "waiting" 状态(介于 runningdone 之间):收到 clarification_request 事件时切入,用户回复后切回 running。InputBar 在 waiting 态激活输入(非 disabled),回复走 WS 发 JSON(不是 POST 新任务)。琥珀色调(区别于正常绿色)统一标识”需要你回复”。

断线重连天然兼容clarification_request_emit 标准管线 → 进 Redis Stream → 重连补发 → 前端重新进入 waiting 态。D 块的基础设施白嫖一次,不需要额外处理。

诚实边界#

  • 做了文本澄清的阻塞-等待-恢复。不做多选/表单/审批。
  • LLM 是否可靠调用 ask_user 取决于模型判断,few-shot 缓解但不消除不确定性。
  • 中途追加约束(“别看 shein 了”)和人工审批工具调用是 WS 双向的另外两个用例,留后续。

§9 面试可讲点#

  1. 「原来的澄清为什么不行」——chat_fallback 终结循环丢上下文 vs 瞎猜搜出来不对。
  2. 「为什么用 Future 而不是 ContextVar」——两个协程不共享 CV context;单进程 asyncio 下 Future 是最轻量的信号。
  3. 「§1.5 的伏笔收上了」——当时诚实说 WS 只用了单向,这里把双向补齐,而且是一个不靠新依赖的三零件组合。
  4. 「D 块白嫖」——clarification_request 走标准 _emit 管线,断线重连补发天然覆盖。
  5. 「深度守卫 + prompt 走廊」——子 Agent 硬挡 + 主动性纪律保留 + few-shot 教范式,三管齐下控制”什么时候该问”。