1. 这不是又一个“聊天界面”,而是一套可审计的智能体执行流水线
上周五下午三点,我盯着终端里一行行滚动的curl -N http://localhost:8000/agent/stream输出发了三分钟呆——不是因为卡顿,而是因为终于看到{"step":"execute","tool":"search_web","query":"2024年Q3全球AI芯片出货量统计"}这样的结构化日志,像工厂流水线上的工单编号一样,稳稳地从 FastAPI 后端流进 Vue 前端,再被实时渲染成带时间戳、带步骤类型、带工具调用参数的卡片。那一刻我才真正意识到:我们做的不是“让大模型说话”,而是给 AI 的思考过程装上仪表盘和黑匣子。
这个项目标题里的每个词都不是装饰:CLI 是起点,它代表命令行下可复现、可脚本化的最小验证单元;浏览器是终点,但不是炫技的 UI,而是面向产品、运营、法务甚至审计人员的可视化操作台;FastAPI 是承重墙,它不只提供/chat接口,更要承载状态管理、流控、日志注入、权限校验四层责任;SSE 是血管,它比 WebSocket 更轻量、比轮询更实时,专为“单向、长时、事件驱动”的推理流设计;Vue 3 是神经末梢,用 Composition API 精准响应每一条data: {"type":"thought","content":"需要验证数据来源可靠性..."};而ReAct Agent 是灵魂——它不是把 prompt 拼得更长,而是用明确的Thought/Action/Observation/Answer四段式结构,把“AI 怎么想的”变成可拆解、可回溯、可人工干预的原子操作。
你可能刚在某篇教程里跑通过fastapi + vue的 hello world,也可能用过codex cli或zcode cli调用本地模型。但那些 demo 里,response.text是一团不可分割的字符串,console.log(data)只能看到最终答案。而本项目要解决的是真实业务场景里的三个硬需求:第一,当客户质疑“为什么推荐这款芯片”,你能立刻拉出第 7 步的search_web调用记录和返回的原始网页快照;第二,当线上 agent 卡在Observation阶段超过 15 秒,运维能通过 FastAPI 的/health接口+ Prometheus 指标,精准定位是模型响应超时还是网络抖动;第三,当合规部门要求“所有工具调用必须经审批”,你在action_router.py里加一行if action == "send_email": raise PermissionError("Email requires manual approval")就能生效——这些能力,全系于 CLI 到浏览器这条链路的每一环是否真正“可审计”。
提示:本项目不依赖任何闭源 SDK 或云服务。所有代码基于 Python 3.11+、FastAPI 0.111+、Vue 3.4+(Composition API + Pinia)、Vite 构建。核心逻辑全部开源,你可以把它嵌入现有企业内网系统,也可以作为独立服务部署在私有 Kubernetes 集群中。
2. CLI 层:用纯 Python 实现 ReAct 循环,拒绝黑盒封装
很多初学者一上来就冲着langchain或llamaindex的AgentExecutor去,结果调试时发现agent.run()抛出异常,连哪一步Thought出错都看不到。本项目的第一步,就是亲手用 200 行纯 Python 写一个最小可运行的 ReAct 循环——它不漂亮,但每一行都在你眼皮底下。
2.1 ReAct 的本质不是 Prompt 工程,而是状态机驱动
ReAct 的核心不是“让模型学会思考”,而是定义一套人类可读、机器可执行的状态转移规则。我们把整个循环抽象为四个状态:
THOUGHT: 模型输出一段自然语言,描述当前推理路径(如:“需要查证该芯片的功耗数据是否符合客户要求”)ACTION: 模型按固定格式输出工具调用指令(如:Action: search_web\nAction Input: {"query":"NVIDIA H100 250W TDP official spec"})OBSERVATION: 工具执行后返回的原始结果(如:<html><title>NVIDIA H100 Data Sheet</title>...)ANSWER: 模型综合所有 Observation 后给出最终回答(如:“H100 的典型功耗为 250W,符合客户 300W 以内要求”)
关键在于:状态切换必须由代码显式控制,而非依赖模型“自觉”输出特定 token。我们用一个while True循环 +state变量实现:
# cli/agent_core.py from typing import Dict, Any, Optional import json import re class ReActAgent: def __init__(self, llm_client): self.llm_client = llm_client # 支持 openai.Completion 或 ollama.chat self.history = [] # 存储完整的 step-by-step 日志 def run(self, user_query: str) -> str: self.history.clear() self.history.append({"role": "user", "content": user_query}) max_steps = 8 for step in range(max_steps): # 1. 生成 Thought + Action prompt = self._build_prompt() response = self.llm_client.generate(prompt) # 2. 解析模型输出,严格匹配正则 thought_match = re.search(r"Thought:\s*(.*?)(?:\n|$)", response) action_match = re.search(r"Action:\s*(\w+)\nAction Input:\s*(\{.*?\})", response, re.DOTALL) if not thought_match or not action_match: # 解析失败,强制进入 ANSWER 状态 self.history.append({"role": "assistant", "content": f"Failed to parse step {step}"}) break thought = thought_match.group(1).strip() action_name = action_match.group(1).strip() try: action_input = json.loads(action_match.group(2)) except json.JSONDecodeError: action_input = {"raw": action_match.group(2)} # 3. 记录 Thought & Action 到 history self.history.append({ "step": step, "type": "thought", "content": thought, "timestamp": time.time() }) self.history.append({ "step": step, "type": "action", "tool": action_name, "input": action_input, "timestamp": time.time() }) # 4. 执行工具并记录 Observation observation = self._execute_tool(action_name, action_input) self.history.append({ "step": step, "type": "observation", "tool": action_name, "output": observation[:500] + "..." if len(observation) > 500 else observation, "timestamp": time.time() }) # 5. 如果是 final answer,跳出循环 if action_name == "finish": return observation return "Max steps exceeded"注意:
_execute_tool方法是真正的业务胶水。本项目预置了search_web(调用 duckduckgo-search)、get_weather(调用 OpenWeatherMap API)、read_file(读取本地 JSON)三个工具。每个工具都必须返回结构化字典(如{"status": "success", "data": {...}}),而非原始 HTML 字符串。这是后续审计的关键——Observation字段必须能被前端直接解析,而不是扔给用户一堆<div>标签。
2.2 CLI 的价值:可复现、可压测、可集成到 CI/CD
写完ReActAgent类,我们立刻用argparse包装成 CLI 工具:
# cli/main.py import argparse from agent_core import ReActAgent from llm_clients import OllamaClient # 或 OpenAIClient def main(): parser = argparse.ArgumentParser(description="Run ReAct Agent from CLI") parser.add_argument("--query", type=str, required=True, help="User query") parser.add_argument("--model", type=str, default="llama3", help="LLM model name") parser.add_argument("--max-steps", type=int, default=8, help="Max ReAct steps") args = parser.parse_args() client = OllamaClient(model=args.model) agent = ReActAgent(client) result = agent.run(args.query) # 关键:输出完整 history 为 JSONL,供后续分析 for log in agent.history: print(json.dumps(log, ensure_ascii=False)) if __name__ == "__main__": main()执行python cli/main.py --query "对比 RTX 4090 和 H100 在 AI 训练场景的性价比",你会得到 10+ 行 JSONL 输出,每行是一个带type、step、timestamp的审计事件。这带来三个实操优势:
- 调试效率翻倍:当某次运行卡住,你不用重启整个 Web 服务,只需
grep '"type":"action"' output.jsonl | tail -3查看最后三次工具调用,立刻判断是search_web返回空结果,还是get_weatherAPI 密钥失效; - 压测脚本直连:用
ab或wrk对 FastAPI 接口压测时,后端实际调用的就是这个 CLI 类。你可以在pytest中写test_agent_step_by_step.py,用mock.patch替换OllamaClient,100% 覆盖所有Thought/Action/Observation分支; - CI/CD 自动化审计:在 GitHub Actions 中,每次 PR 提交后自动运行
python cli/main.py --query "测试用例",将history输出存为 artifact。如果某次type字段出现error或timeout,Pipeline 直接失败——这才是真正的“可审计”。
实操心得:我在第一次部署时,发现
search_web工具在服务器上因 DNS 解析超时导致整个 agent 卡死。解决方案不是加try/except,而是在 CLI 层增加-t/--timeout参数,并在_execute_tool中统一用requests.get(url, timeout=args.timeout)。这样,审计日志里会明确记录"type":"error","message":"Timeout on search_web",而不是让前端显示一片空白。
3. FastAPI 层:不止是 API Server,更是审计日志的中央枢纽
很多 FastAPI 教程教你写@app.post("/chat"),然后return {"response": llm.generate(...)}。但在可审计 agent 场景下,FastAPI 必须承担四项额外职责:流式响应封装、请求上下文注入、结构化日志落库、跨域与安全加固。我们逐项拆解。
3.1 SSE 流的本质:HTTP Chunked Transfer Encoding 的优雅封装
SSE(Server-Sent Events)不是魔法,它是 HTTP/1.1 的Transfer-Encoding: chunked特性 + 服务端主动推送的组合。FastAPI 用StreamingResponse实现它,但关键在于如何把 ReAct 的离散步骤,转换成符合 SSE 规范的连续数据流。
标准 SSE 格式要求:
- 每条消息以
data:开头,后跟 JSON 字符串 - 消息间用双换行分隔
- 可选
event:定义事件类型(如event: thought) - 可选
id:用于客户端断线重连
我们的stream_agent接口这样实现:
# api/main.py from fastapi import FastAPI, Request, Depends from fastapi.responses import StreamingResponse from starlette.concurrency import iterate_in_threadpool import json import time from cli.agent_core import ReActAgent from llm_clients import OllamaClient app = FastAPI() @app.post("/agent/stream") async def stream_agent(request: Request): # 1. 解析请求体,获取 query 和 session_id body = await request.json() user_query = body.get("query", "") session_id = body.get("session_id", f"sess_{int(time.time())}") # 2. 初始化 agent(注意:这里不能 new ReActAgent(),需注入 context) agent = ReActAgent(OllamaClient(model="llama3")) # 3. 定义生成器函数,yield 每个 step async def event_generator(): try: # 发送初始化事件 yield f"event: init\ndata: {json.dumps({'session_id': session_id, 'timestamp': time.time()})}\n\n" # 执行 agent.run(),但捕获每一步的 history # 注意:这里不能直接调用 agent.run(),因为它返回最终字符串 # 我们需要重写 run() 为 generator for step_log in agent.run_stream(user_query): # 新增方法 # 4. 格式化为 SSE 消息 yield f"event: {step_log['type']}\ndata: {json.dumps(step_log, ensure_ascii=False)}\n\n" # 5. 强制 flush,避免 Nginx 缓存 await asyncio.sleep(0.01) except Exception as e: yield f"event: error\ndata: {json.dumps({'error': str(e), 'timestamp': time.time()})}\n\n" return StreamingResponse( event_generator(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no" # 关键!禁用 Nginx 缓存 } )agent.run_stream()是对原run()方法的重构,它不再返回字符串,而是yield每个step_log:
# cli/agent_core.py def run_stream(self, user_query: str): self.history.clear() self.history.append({"role": "user", "content": user_query}) for step in range(self.max_steps): # ... [同前] ... # 在每次 self.history.append() 后,yield 当前 log yield self.history[-1] # 最新一条日志 # ... [继续] ...关键细节:
X-Accel-Buffering: no头是生产环境的救命稻草。没有它,Nginx 会默认缓存 8KB 数据才推送给前端,导致 SSE 消息延迟数秒甚至超时。stream disconnected before completion: idle timeout waiting for sse这个热词错误,90% 源于此。我们在nginx.conf中还额外配置了proxy_buffering off; proxy_cache off;,双重保险。
3.2 审计日志的落地:不只是 print,而是结构化存储
CLI 层的print(json.dumps(log))只适合开发。生产环境必须把每条step_log存入数据库,且满足审计要求:不可篡改、带时间戳、关联 session_id、支持 SQL 查询。
我们选用 SQLite(轻量)+ SQLAlchemy Core(非 ORM,避免性能损耗):
# api/db.py from sqlalchemy import create_engine, text from sqlalchemy.pool import StaticPool # 使用 StaticPool 避免多线程连接问题 engine = create_engine( "sqlite:///./audit.db", connect_args={"check_same_thread": False}, poolclass=StaticPool ) # 初始化表 with engine.connect() as conn: conn.execute(text(""" CREATE TABLE IF NOT EXISTS audit_logs ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL, step INTEGER NOT NULL, type TEXT NOT NULL CHECK(type IN ('thought','action','observation','answer','error')), tool TEXT, content TEXT, timestamp REAL NOT NULL, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) """)) conn.commit()在stream_agent的event_generator中,每yield一条日志,就同步插入数据库:
# api/main.py from api.db import engine async def event_generator(): # ... [初始化] ... for step_log in agent.run_stream(user_query): # 插入审计日志(同步操作,确保顺序) with engine.connect() as conn: conn.execute(text(""" INSERT INTO audit_logs (session_id, step, type, tool, content, timestamp) VALUES (:session_id, :step, :type, :tool, :content, :timestamp) """), { "session_id": session_id, "step": step_log.get("step", 0), "type": step_log["type"], "tool": step_log.get("tool", ""), "content": json.dumps(step_log, ensure_ascii=False)[:1000], # 防止超长 "timestamp": step_log["timestamp"] }) conn.commit() yield f"event: {step_log['type']}\ndata: {json.dumps(step_log, ensure_ascii=False)}\n\n"实操心得:不要用异步数据库驱动(如 asyncpg)处理审计日志。SSE 流要求事件严格按序,而异步 I/O 可能导致日志入库顺序与流输出顺序不一致。我们实测过,用
threading.Lock()包裹conn.execute()比异步方案更稳定。另外,content字段限制 1000 字符,是因为 SQLite 的TEXT类型虽无硬限制,但过长字段会显著拖慢SELECT * FROM audit_logs WHERE session_id='xxx'查询速度。
3.3 CORS 与安全加固:FastAPI 不是裸奔的玩具
fastapi cors是高频搜索词,说明很多人栽在跨域上。但真正的风险不在allow_origins=["*"],而在于未校验 session_id、未限制请求频率、未过滤恶意 query。
我们添加三层防护:
- Session ID 校验:前端必须在请求体中传
session_id,后端用uuid.uuid4()生成并返回,后续请求必须携带。防止恶意脚本批量调用; - 速率限制:用
slowapi库限制/agent/stream每 IP 每分钟 5 次; - Query 过滤:对
user_query做基础清洗,移除\x00-\x08\x0b\x0c\x0e-\x1f等控制字符,防止注入攻击。
# api/main.py from slowapi import Limiter from slowapi.util import get_remote_address limiter = Limiter(key_func=get_remote_address) @app.post("/agent/stream") @limiter.limit("5/minute") async def stream_agent(request: Request): body = await request.json() user_query = body.get("query", "") session_id = body.get("session_id", "") # 1. Session ID 校验 if not session_id or not re.match(r"^sess_\d+$", session_id): raise HTTPException(status_code=400, detail="Invalid session_id") # 2. Query 清洗 clean_query = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f]", "", user_query) if len(clean_query) < 3 or len(clean_query) > 500: raise HTTPException(status_code=400, detail="Query too short or too long") # ... [后续逻辑] ...注意:
fastapi cors错误常源于 Nginx 配置。我们nginx.conf中明确设置:location /api/ { proxy_pass http://fastapi_backend; proxy_set_header Origin $scheme://$host; add_header 'Access-Control-Allow-Origin' '$scheme://$host'; add_header 'Access-Control-Allow-Methods' 'GET, POST, OPTIONS'; add_header 'Access-Control-Allow-Headers' 'Content-Type, Authorization'; }这样,前端
fetch("http://your-domain.com/api/agent/stream")才能正确拿到 CORS 头。
4. Vue 3 层:用 Composition API 构建可追溯的 UI 状态机
Vue 3 的 Composition API 不是语法糖,它是构建复杂状态流的刚需。本项目的前端不是“展示数据”,而是精确映射 ReAct 的四个状态,并允许用户在任意步骤暂停、重试、导出。
4.1 Pinia Store:定义 agent 的单一事实源
我们创建stores/agent.ts,用defineStore管理 agent 的完整生命周期:
// stores/agent.ts import { defineStore } from 'pinia' import { ref, computed } from 'vue' export interface StepLog { step: number type: 'thought' | 'action' | 'observation' | 'answer' | 'error' tool?: string input?: Record<string, any> output?: string content?: string timestamp: number } export const useAgentStore = defineStore('agent', () => { const sessionId = ref<string>('') const isStreaming = ref<boolean>(false) const logs = ref<StepLog[]>([]) const currentStep = ref<number>(0) const status = ref<'idle' | 'thinking' | 'executing' | 'observing' | 'answering' | 'error'>('idle') // 计算属性:按 step 分组的日志 const groupedLogs = computed(() => { return logs.value.reduce((acc, log) => { if (!acc[log.step]) acc[log.step] = [] acc[log.step].push(log) return acc }, {} as Record<number, StepLog[]>) }) // 计算属性:当前 step 的最新状态 const currentStatus = computed(() => { const lastLog = logs.value[logs.value.length - 1] if (!lastLog) return 'idle' switch (lastLog.type) { case 'thought': return 'thinking' case 'action': return 'executing' case 'observation': return 'observing' case 'answer': return 'answering' case 'error': return 'error' default: return 'idle' } }) // Action:启动 stream const startStream = async (query: string) => { isStreaming.value = true sessionId.value = `sess_${Date.now()}` logs.value = [] try { const response = await fetch('/api/agent/stream', { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ query, session_id: sessionId.value }) }) if (!response.ok) throw new Error(`HTTP ${response.status}`) const reader = response.body?.getReader() if (!reader) throw new Error('No reader') while (true) { const { done, value } = await reader.read() if (done) break const decoder = new TextDecoder() const text = decoder.decode(value) const lines = text.split('\n') for (const line of lines) { if (line.startsWith('data:')) { try { const data = JSON.parse(line.substring(5).trim()) logs.value.push(data) } catch (e) { console.warn('Invalid SSE data:', line) } } } } } catch (e) { logs.value.push({ step: logs.value.length, type: 'error', content: `Stream failed: ${(e as Error).message}`, timestamp: Date.now() }) } finally { isStreaming.value = false } } // Action:导出当前 session 日志 const exportLogs = () => { const blob = new Blob([JSON.stringify(logs.value, null, 2)], { type: 'application/json' }) const url = URL.createObjectURL(blob) const a = document.createElement('a') a.href = url a.download = `agent_session_${sessionId.value}.json` a.click() URL.revokeObjectURL(url) } return { sessionId, isStreaming, logs, currentStep, status, groupedLogs, currentStatus, startStream, exportLogs } })关键设计:
groupedLogs计算属性把logs数组按step分组,这样模板里可以用<div v-for="(stepLogs, step) in groupedLogs">渲染每一步的卡片。每个卡片包含thought、action、observation三条日志,形成完整的“推理单元”。
4.2 组件化渲染:每个 step 是一个可交互的审计单元
AgentView.vue组件的核心是v-for渲染groupedLogs:
<!-- components/AgentView.vue --> <template> <div class="agent-container"> <!-- 输入区 --> <div class="input-section"> <input v-model="query" @keyup.enter="startStream" placeholder="输入你的问题..." :disabled="isStreaming" /> <button @click="startStream" :disabled="isStreaming"> {{ isStreaming ? '运行中...' : '开始推理' }} </button> </div> <!-- 日志流 --> <div class="logs-section"> <div v-for="(stepLogs, step) in groupedLogs" :key="step" class="step-card" > <div class="step-header"> <span class="step-number">Step {{ step }}</span> <span class="step-status">{{ getStatusText(stepLogs) }}</span> </div> <!-- Thought --> <div v-if="getLogByType(stepLogs, 'thought')" class="log-item thought"> <strong>Thought:</strong> {{ getLogByType(stepLogs, 'thought')!.content }} </div> <!-- Action --> <div v-if="getLogByType(stepLogs, 'action')" class="log-item action"> <strong>Action:</strong> {{ getLogByType(stepLogs, 'action')!.tool }} <span v-if="getLogByType(stepLogs, 'action')!.input"> ({{ JSON.stringify(getLogByType(stepLogs, 'action')!.input, null, 2) }}) </span> </div> <!-- Observation --> <div v-if="getLogByType(stepLogs, 'observation')" class="log-item observation"> <strong>Observation:</strong> <pre>{{ getLogByType(stepLogs, 'observation')!.output }}</pre> </div> <!-- Answer 或 Error --> <div v-if="getLogByType(stepLogs, 'answer')" class="log-item answer"> <strong>Answer:</strong> {{ getLogByType(stepLogs, 'answer')!.content }} </div> <div v-if="getLogByType(stepLogs, 'error')" class="log-item error"> <strong>Error:</strong> {{ getLogByType(stepLogs, 'error')!.content }} </div> </div> </div> <!-- 控制区 --> <div class="control-section" v-if="logs.length > 0"> <button @click="exportLogs">导出本次会话日志</button> <button @click="clearLogs">清空日志</button> </div> </div> </template> <script setup lang="ts"> import { ref, computed } from 'vue' import { useAgentStore } from '@/stores/agent' const store = useAgentStore() const query = ref('') const startStream = () => { if (!query.value.trim()) return store.startStream(query.value) query.value = '' } const clearLogs = () => { store.logs = [] store.sessionId = '' } // 辅助函数:根据 type 获取日志 const getLogByType = (logs: StepLog[], type: string) => { return logs.find(log => log.type === type) } const getStatusText = (logs: StepLog[]) => { const last = logs[logs.length - 1] if (last.type === 'answer') return '完成' if (last.type === 'error') return '错误' return last.type === 'thought' ? '思考中' : last.type === 'action' ? '执行中' : '观察中' } </script>实操心得:
<pre>标签渲染Observation是刻意为之。很多教程用v-html,但这有 XSS 风险。我们要求所有Observation输出必须是纯文本或 JSON,前端不做任何 HTML 解析。如果工具返回 HTML,_execute_tool方法必须先用BeautifulSoup提取文本,再存入output字段。这样,审计日志里存的是干净文本,前端展示也绝对安全。
4.3 Vue 3 Snippets:提升开发效率的实战技巧
vue 3 snippets是高频搜索词,说明开发者渴望开箱即用的代码片段。我们整理了三个高频场景的 snippet:
- SSE 连接重试:当网络中断,前端自动重连:
// utils/sse-reconnect.ts export function createSSEWithRetry(url: string, onMessage: (data: any) => void) { let eventSource: EventSource | null = null const connect = () => { eventSource = new EventSource(url) eventSource.onmessage = (e) => onMessage(JSON.parse(e.data)) eventSource.onerror = () => { console.warn('SSE connection lost, retrying in 3s...') setTimeout(connect, 3000) } } connect() return () => eventSource?.close() }- 日志高亮:用不同颜色区分
Thought/Action/Observation:
/* styles/agent.css */ .log-item.thought { background-color: #e6f7ff; border-left: 4px solid #1890ff; } .log-item.action { background-color: #fff0f6; border-left: 4px solid #eb2f96; } .log-item.observation { background-color: #f6ffed; border-left: 4px solid #52c418; } .log-item.answer { background-color: #f0f9ff; border-left: 4px solid #1890ff; font-weight: bold; } .log-item.error { background-color: #fff2f0; border-left: 4px solid #f5222d; }- 响应式布局适配移动端:用
@media优化小屏体验:
@media (max-width: 768px) { .step-card { padding: 12px; } .step-header { flex-direction: column; } .log-item pre { white-space: pre-wrap; word-break: break-word; } }注意:
vs code gemini cli companion或claude code cli这类工具,本质是把 LLM 的代码补全能力接入 IDE。但本项目强调:前端逻辑必须手写,不能依赖 AI 生成。因为审计要求“代码可知、行为可溯”。我们用vue-tsc --noEmit做类型检查,用vitest写单元测试,确保getLogByType等辅助函数 100% 覆盖。
5. 全链路联调与生产级避坑指南
当 CLI、FastAPI、Vue 三端各自跑通,真正的挑战才开始:如何让它们在真实网络环境下稳定协同?这里没有银弹,只有踩过的坑和验证过的方案。
5.1 “Stream disconnected before completion” 的七种根因与修复
这个错误是 SSE 最常见的报错,但原因千差万别。我们按发生位置分类:
| 位置 | 根因 | 修复方案 | 验证命令 |
|---|---|---|---|
| 客户端 | 浏览器主动关闭连接(用户切页) | 前端监听visibilitychange事件,页面隐藏时暂停 stream | document.addEventListener('visibilitychange', () => { if (document.hidden) controller.abort() }) |
| Nginx | proxy_read_timeout默认 60s | 在nginx.conf中设proxy_read_timeout 300; | curl -N http://localhost/api/agent/stream观察是否 60s 后断开 |
| FastAPI | StreamingResponse未及时 flush | 加await asyncio.sleep(0.01)强制 flush | 在event_generator中添加print("flushed")日志 |
| LLM Client | OllamaClient请求超时 | 在llm_clients.py中设requests.post(..., timeout=120) | time python cli/main.py --query "long query" |
| 数据库 | INSERT操作阻塞流 | 改为异步写入(用asyncio.to_thread包裹) | ab -n 100 -c 10 http://localhost/api/agent/stream |
| 防火墙 | 云服务商拦截长连接 | 开放 TCP keepalive,设net.ipv4.tcp_keepalive_time=600 | `ss -tn |
| 前端 | EventSource未处理error事件 | 添加eventSource.onerror = () => { /* 重试逻辑 */ } | 手动断开网络,观察前端是否重连 |
实测案例:某次上线后,AWS ALB 报错
stream disconnected before completion: idle timeout waiting for sse。排查发现 ALB 的空闲超时默认 60 秒,而我们的search_web工具在高峰时段响应达 90 秒。解决方案不是改工具,而是在 ALB 控制台将Idle timeout改为 300 秒,并在 FastAPI 中加@app.middleware("http")记录每个请求的time.time(),对比日志确认是 ALB 断连。
5.2 从 CLI 到浏览器的端到端测试脚本
自动化测试不是可选项,而是审计合规的基石。我们用playwright写了一个端到端测试:
# tests/e2e_test.py from playwright.sync_api import sync_playwright import json import time def test_react_agent_flow(): with sync_playwright() as p: browser = p.chromium.launch