基于SSE事件流的完整的前后端

架构概览:三文件分工

api_view/api/
├── chat.py      ★ 核心:SSE 流式对话 + 中断检测 + 中断恢复 + 展示消息持久化
├── history.py   ← 历史会话列表 / 消息查询 / 会话删除
└── agent_loader.py (上层) ← Agent 单例 / MongoDB 连接 / 展示消息存取

整体数据流

前端 (浏览器)                       后端 (FastAPI)                     LangGraph Agent
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

① POST /api/chat/stream
   {message, thread_id?, user_id?}
        ──────────────────────────→
                                     ② stream_chat_response()
                                        ├─ agent_loader.get_agent_for_user(user_id)
                                        ├─ agent_loader.create_config(thread_id, user_id)
                                        ├─ 构建 input:
                                        │    初始: {"messages": [{"role":"user",...}]}
                                        │    恢复: Command(resume=resume_data)
                                        │
                                     ③ agent_graph.astream(
                                          input, config,
                                          stream_mode=["messages","values"],
                                          subgraphs=True, version="v2"
                                        )
                                        ↓↓↓ 开始流式推送 ↓↓↓
  ←── SSE: token ──────────────────  ←  ④ messages 流: (token, metadata)
  ←── SSE: tool_start ─────────────  ←     token 有 tool_call_chunks → 工具开始
  ←── SSE: tool_args ──────────────  ←     工具参数增量
  ←── SSE: tool_result ────────────  ←     type=="tool" → 工具结果
  ←── SSE: tool_end ───────────────  ←     工具调用结束
  ←── SSE: interrupt ──────────────  ←  ⑤ values 流: chunk["interrupts"]
  ←── SSE: done ───────────────────  ←  ⑥ 流正常结束

SSE的事件流

chat.py 从头到尾只做一件事:把 Agent 流式输出转换成前端能消费的 SSE 事件序列。但它的实现方式体现了项目特色。

一、 双流模式 — 一条连接同时走两路

# chat.py:281-287
async for chunk in agent_loader.agent.astream(
    input=current_input,
    config=config,
    stream_mode=["messages", "values"],  # ← 双模式
    subgraphs=True,
    version="v2",
):

这是整个中断检测机制的基础。常规 SSE 只用 "messages" 模式拿 token/tool 事件,但本项目需要检测中断,所以同时开了 "values"

stream_mode=["messages", "values"]
     │              │
     │              └─→ values 流: 检测 interrupts() 是否触发
     │                   chunk_type == "values" && chunk.get("interrupts")
     │                   先于 messages 处理(代码中 if 在 messages 前面)
     │
     └─→ messages 流: token / tool_call / tool_result
          chunk_type == "messages"
          正常 SSE 事件 (token / tool_start / tool_args / tool_result / tool_end)
流 A:messages — 内容输出(AIMessage,ToolMessage)
messages 流产出的 token 类型:

┌─────────────────────────────────────────────────────────┐
│ token.tool_call_chunks 有值                             │  - AIMessage(调用工具的指令,一般没有content)
│   → token.tool_call_chunks[i]["name"] = "web_search"   │  → tool_start 事件
│   → token.tool_call_chunks[i]["args"] = '{"query":...}' │  → tool_args 事件
│                                                         │
│ token.type == "tool"                                    │  # ToolMessage (工具调用后返回的内容)
│   → token.name = "web_search"                          │  → tool_result 事件
│   → token.content = "搜索结果是..."                     │  → tool_end 事件
│                                                         │
│ token.content 是纯文本 (不是 tool, 不是 tool_call_chunks)│
│   → "根据搜索结果,博世刹车片..."                       │  → token 事件
└─────────────────────────────────────────────────────────┘
流 B:values — 中断检测
# 中断检测必须在 messages 处理之前执行 (line 296-366)
if chunk_type == "values" and chunk.get("interrupts"):
    for interrupt in chunk["interrupts"]:
        if "action_requests" in interrupt.value:
            → interrupt_type = "hitl_approval"   # 第2层: 审批中断
        elif interrupt.value.get("type") == "order_info_request":
            → interrupt_type = "order_info_supplement"  # 第1层: 数据补充
    → 保存 display_messages 到 MongoDB
    → 发送 done(interrupted=True)
    → return  # 停止本次流

为什么必须先检测中断? 因为中断时最后一条 messages 事件可能还来不及到达,先处理中断才能确保两端状态一致。

二、完整请求链路(从用户输入到 SSE 返回)

7 种 SSE 事件
事件类型 触发条件 前端行为
token AI 文本增量 追加到当前 AI 回复气泡
tool_start Agent 开始调工具 显示工具调用卡片(含工具名)
tool_args 工具参数流式输出 更新工具卡片上的参数 JSON
tool_result 工具执行完毕 填充工具执行结果
tool_end 工具调用结束 标记工具卡片为完成状态
interrupt HITL 中断触发 显示审批面板/数据补充表单
done 流结束 停止 loading,标记对话完成
用户发送 "帮我新增采购订单"
  │
  ▼
POST /api/chat/stream {message: "帮我新增采购订单"}
  │
  ▼
chat_stream() → StreamingResponse
  │
  ▼
stream_chat_response(message="帮我新增采购订单")
  │ config = agent_loader.create_config(thread_id)
  │ current_input = {"messages": [{"role":"user", "content":"..."}]}
  │ display_messages = [user消息]
  │
  ▼
agent.astream(input, config, stream_mode=["messages","values"], subgraphs=True, version="v2")
  │
  ├─→ 主 Agent 收到消息 → 判断: 订单操作 → task(procurement-order)
  │
  ├─→ [SSE] token: "好的,我来帮您..."           (source: main)
  ├─→ [SSE] tool_start: task                    (source: main)
  ├─→ [SSE] tool_args: {...}                    (source: main)
  │
  ├─→ order 子 Agent 启动 (独立上下文)
  │     ├─→ 第1步: 提取数据
  │     ├─→ 第2步: Schema 校验 → partId 缺失!
  │     ├─→ 第3步: request_order_info(...)
  │     │     └─→ interrupt({"type": "order_info_request", ...})
  │     │
  ├─→ [stream_mode="values"] 检测到 interrupts
  │     │
  │     ├─→ [SSE] interrupt {interrupt_type: "order_info_supplement", ...}
  │     ├─→ save_display_messages(thread_id, cleaned)  ← 保存现场
  │     └─→ [SSE] done {interrupted: true}
  │         return  ← 流结束
  │
  ▼
前端 InterruptBanner 弹出: "请输入补充信息"
用户输入: "物料ID=100, 数量=50, 单价=25.5"
  │
  ▼
POST /api/chat/{thread_id}/resume {resume: {supplement: "物料ID=100..."}}
  │
  ▼
stream_chat_response(thread_id, resume_data={supplement: "..."})
  │ existing = get_display_messages()  ← 加载中断前的消息
  │ current_input = Command(resume={supplement: "..."})  ← 恢复
  │
  ▼
agent.astream(Command(resume=...), config, ...)
  │
  ├─→ order 子 Agent 恢复 → 解析补充数据 → 合并 → 校验通过
  ├─→ 构造 order_create 参数 → interrupt_on 触发
  │
  ├─→ [stream_mode="values"] 再次检测到 interrupts
  │     ├─→ [SSE] interrupt {interrupt_type: "hitl_approval", action_requests: [...]}
  │     └─→ [SSE] done {interrupted: true}
  │
  ▼
前端 InterruptBanner 切换: 审批卡片 (approve/reject)
用户点击 "approve"
  │
  ▼
POST /api/chat/{thread_id}/resume {resume: {decisions: [{type: "approve"}]}}
  │
  ▼
agent.astream(Command(resume={decisions:...}), config, ...)
  │
  ├─→ order_create 执行 → 成功
  ├─→ [SSE] tool_result {text: "订单创建成功..."}
  ├─→ [SSE] token "订单已创建,编号 PO20260513..."  (source: order)
  │
  ├─→ 主 Agent 收到子Agent 完成 → compact_conversation
  ├─→ [SSE] token "已为您完成订单创建..."           (source: main)
  │
  └─→ save_display_messages(thread_id, all_messages)
      [SSE] done {thread_id: "...", content: "..."}

亮点总结

特色 说明
双流模式 [“messages”,”values”] 同时走两路,values 优先检测中断
display_messages 独立于 checkpoint 的完整展示层,包含子 Agent 消息,逐条存 MongoDB 避免 16MB 限制
中断状态保存 中断时立即 save_display_messages,恢复时 get_display_messages 加载,前后端无缝衔接
两层中断统一处理 request_order_info 和 interrupt_on 共用同一个中断检测→保存→恢复流程
嵌套工具调用栈 显式 push/pop 管理,支持主Agent调子Agent调图表工具
content 提取管道 UUID 过滤 + 多类型兼容 + 图片提取
调试日志 每个 stream 一个独立日志文件,chunk 级别全量记录
逐条 MongoDB 存储 单条消息一个文档 + 500KB 截断兜底,彻底解决长对话存储问题
--- 本文结束 The End ---