架构概览:三文件分工
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 截断兜底,彻底解决长对话存储问题 |