M6 · WebSocket 实时通信与异步任务#
简历 Bullet Point: 设计 connect-first WebSocket 协议零事件丢失(先建连再起任务),Redis Stream 持久化 + 断线重连 last_event_id 补发;双池优先级队列按任务重量分流(normal/heavy 独立槽位),长任务堵不死短任务;幂等三层 + 覆盖重发 + 冷启动续看
开场钩子#
第一版是”先起任务再连 WS”——POST /api/task 返回 thread_id,前端拿到后才建 WebSocket。问题:任务跑得快的时候,session_created 和前几个 tool_start 事件在 WS 建连之前就推了出去,前端连上后这些事件已经消失。用户看到的第一个事件是某个工具的 tool_end,前面发生了什么一无所知。
改成 connect-first:前端先生成 thread_id → 连 WS → 收 ws_ready 确认登记 → 才 POST 起任务。零缓冲零竞态。
一、模块运作流程#
1.1 connect-first 协议#
前端 后端
│ 生成 thread_id │
│ ──── WS /ws/{tid} ──────→ │ ConnectionManager.connect(tid, ws)
│ ←──── ws_ready ─────────── │
│ ──── POST /api/task ─────→ │ create_task(run_agent, tid)
│ ←──── session_created ──── │ ← 不丢!WS 已注册
│ ←──── tool_start ───────── │
│ ←──── tool_end ────────── │
│ ... │
│ ←──── task_result ──────── │plaintext1.2 AGUI 事件体系#
十余类事件:session_created / assistant_call / tool_start / tool_end / summary_delta / items_preview / fork / queue_status / clarification_request / memory_updated / memory_applied / task_result / task_cancelled / error。统一信封:{type, event, message, data, thread_id, timestamp}。
fork 事件在进子 thread_scope 之前上报、路由到父 thread(前端连的是父任务)。
1.3 双池优先级队列#
| 池 | 槽位 | 适用 | 分流依据 |
|---|---|---|---|
| normal | 5 | 新对话/短任务 | 该 thread 历史 ≤3 轮 |
| heavy | 3 | 续聊/长链路 | 历史 >3 轮 |
有界等待队列 + 溢出 429(保住”无界排队”的反对论点)。动态再平衡:normal 积压时压缩 heavy 容量但不抢占已在跑的任务。queue_status AGUI 事件让前端能显示排队进度。
1.4 幂等三层#
- 同 thread 同 query →
already_running - 同 thread 换 query → 覆盖重发(取消旧任务,新任务接管)
- 不带 thread_id → 指纹去重
1.5 事件回放与冷启动续看#
Redis Stream 持久化每会话事件。断线重连带 last_event_id 补发缺口。GET /api/task/{tid}/inflight 探口 + replay_current_run 按轮回放 + 前端 resumeIfRunning 自动重建。切对话不再隐式取消任务。
1.6 收尾流式#
items_preview 事件先推商品卡(先出货),summary_delta 事件流式推文案(后出文案)。定稿前 3-8 秒开始逐字渲染,用户感知延迟大幅降低。
二、踩坑实录#
坑 1:先起任务再连 WS 丢早期事件#
- 改 connect-first(先 WS 再 POST)。
坑 2:active_tasks 按 key 盲删的重发竞态#
- 同 tid 快速重发,新任务的 TaskHandle 被旧任务 finally 块删掉。改成按对象身份校验防误删。
坑 3:前端 WS 生命周期三个坑#
- 非 JSON 帧炸回调 / 断连无兜底卡 running / 卸载不关连接泄漏。逐一修复。
坑 4:ConnectionManager 死连接残留#
- WS 仅正常断开才注销。异常断开时连接泄漏。改成
finally注销。
三、面试问答#
Q1: 为什么 connect-first 而不是缓冲早期事件?#
缓冲需要知道”缓冲到什么时候为止”——WS 什么时候连上是不确定的。connect-first 在协议层关死竞态,比应用层缓冲更可靠。
Q2: 为什么双池不是单池加优先级?#
单池里长任务把所有槽位占满,新的简单查询(品类行情 2 步出结果)也要排队。双池保证 normal 池不被 heavy 任务堵死。
Q3: 事件回放怎么保证不重复?#
每个事件有单调递增的 event_id。断线重连带 last_event_id,服务端只推 > last_event_id 的事件。
四、诚实边界#
| 维度 | 做了 | 没做 |
|---|---|---|
| 实时通信 | WS + connect-first | SSE(不需要双向) |
| 事件持久化 | Redis Stream | 多副本路由 |
| 排队 | 双池有界 + 429 | 分布式任务队列 |
| 幂等 | 同 thread 去重 | 跨 thread 指纹(仅脚本生效) |
| 冷启动 | 续看在跑任务 | 历史任务完整回放 |