ShoppingX 架构与流程#
本文基于真实代码梳理(非
CLAUDE.md的建议结构),覆盖:分层架构、AgentLoop 主循环、fork 安全机制、召回管道、请求时序。与文档直觉不同、值得注意的几个真实落地点见文末「实现注记」。
1. 系统总体架构(分层)#
graph TB
subgraph FE["前端 React + Vite"]
UI["对话框 / AGUI 事件可视化 / 商品卡 / 偏好面板"]
end
subgraph API["app/api · FastAPI + asyncio"]
SRV["server.py<br/>POST /api/task · WS /ws/{tid}<br/>cancel · upload · files · prefs · history · health"]
CM["connection.py<br/>ConnectionManager<br/>thread_id → WebSocket"]
MON["monitor.py<br/>8 类 AGUI 事件"]
CTX["context.py<br/>ContextVar: thread_id / session_dir / user_id / user_profile"]
AT["active_tasks: dict[tid, asyncio.Task]"]
end
subgraph AGENT["app/agent · AgentLoop 核心"]
RA["main_agent.run_agent()<br/>Think→Act→Observe→Reflect"]
TR["tool_registry<br/>FULL_TOOL_SET (14) / TERMINAL_TOOLS"]
MW["middleware<br/>ContextCompression + ToolGuard"]
DISP["dispatch_tool / parallel_dispatch_tool<br/>同质子 Agent fork"]
FG["fork_guard · MAX_FORK_DEPTH=1"]
LLM["llm.py · get_llm / get_fast_llm / get_judge_llm"]
end
subgraph TOOLS["app/tools · 十二大工具"]
T1["planner"]
T1b["image_understand"]
T2["item_search"]
T3["category_insight"]
T4["price_compare"]
T5["shipping_calc"]
T6["item_picker"]
T7["web_search"]
T7b["ask_user"]
T7c["forget_preference"]
T8["chat_fallback ⛔终结"]
T9["shopping_summary ⛔终结"]
end
subgraph INFRA["基础设施层"]
REC["recall/<br/>TowerClient · QdrantRecall · KBClient · Reranker · fx/duty/shipping"]
MEM["memory/<br/>PreferenceStore(长期,本地JSON / Redis) · session_state(会话级 P_t) · curator(记忆管家) · injector · history"]
CMP["compress/<br/>断点压缩 + cache_control"]
EVAL["eval/<br/>Rubric · recall_metrics"]
end
subgraph EXT["外部服务(均可本地降级)"]
E1["Embedding API<br/>BGE-M3"]
E2["Qdrant"]
E3["OpenSearch Hybrid"]
E4["Reranker cross-encoder"]
E5["Tavily"]
E6["Redis"]
E7["LLM OpenAI 兼容<br/>DashScope/Qwen…"]
end
UI <-->|"HTTP + WebSocket"| SRV
SRV --> AT --> RA
SRV <--> CM
MON --> CM
RA --> CTX
RA --> MW --> TR
RA --> LLM
TR --> TOOLS
TR --> DISP --> FG
DISP -.fork.-> RA
T1b --> E7
T2 --> REC
T3 --> REC
T4 --> REC
T5 --> REC
T6 --> MEM
T9 --> MEM
RA --> MON
REC --> E1 & E2 & E3 & E4
T7 --> E5
MEM --> E6
LLM --> E7
style AGENT fill:#1e3a5f,color:#fff
style TOOLS fill:#2d4a2d,color:#fff
style INFRA fill:#4a3a2d,color:#fff
style EXT fill:#3a2d4a,color:#fff
2. AgentLoop 主循环(Think → Act → Observe → Reflect)#
flowchart TD
START(["run_agent(query, thread_id, user_id)"]) --> SETUP
subgraph SETUP_BLOCK["会话准备"]
SETUP["ensure_session_dir → output/<tid>/"]
--> SCOPE["thread_scope: 绑定 ContextVar"]
--> CAP["begin_activity_capture + fork_budget_scope(树级预算=8)"]
--> PREF["读 Store: 长期偏好 + 行为历史 + P_t 拼进当轮 human<br/>(system_prompt 纯静态·跨轮缓存 M6.1)<br/>+ user_profile(仅 like) 塞 ContextVar 供 item_search"]
end
PREF --> INVOKE["agent.ainvoke({messages: 干净历史 + 当轮 human(运行时上下文+query)})<br/>recursion_limit=61 · timeout=300s"]
INVOKE --> THINK
subgraph LOOP["Agent 内部循环"]
THINK["🧠 Think<br/>ContextCompressionMiddleware.awrap_model_call<br/>① 压缩历史视图 ② report_assistant_call ③ 调 LLM"]
--> DECIDE{"模型选择"}
DECIDE -->|"调工具"| GUARD["ToolGuard 权限闸<br/>+ 检索预算 + fork 预算"]
GUARD --> ACT["⚙️ Act 执行工具"]
ACT --> OBS["👁 Observe<br/>结果截断(≤4000tok) + LoopDetector(窗6/阈4)"]
OBS --> REFLECT{"🔁 Reflect<br/>信息够了吗?"}
REFLECT -->|"不够"| THINK
REFLECT -->|"够 → 调终结工具"| TERM
DECIDE -->|"非购物意图"| TERM
end
TERM{"TERMINAL_TOOLS?<br/>shopping_summary / chat_fallback"} -->|"是"| EXIT["循环终止"]
THINK -.超 61 步/300s/循环卡死.-> EXIT
EXIT --> POST
subgraph POST_BLOCK["收尾"]
POST["_extract_summary(messages)"]
--> WRITE["写 summary.md + result.json"]
--> HIST["append_turn → turns.json"]
--> RESULT["report_task_result(items)"]
--> CURATE["curate_turn(后处理异步):更 P_t + 提升长期偏好"]
end
RESULT --> END(["返回 + 推送前端"])
style THINK fill:#1e3a5f,color:#fff
style ACT fill:#2d4a2d,color:#fff
style OBS fill:#4a3a2d,color:#fff
style REFLECT fill:#5f1e3a,color:#fff
循环终止四重护栏:
| # | 机制 | 参数 |
|---|---|---|
| ① | 调用 TERMINAL_TOOLS(shopping_summary / chat_fallback)+ 终结硬停哨兵:主 loop 调过终结工具后,后续任何工具在执行层被拦、回哨兵逼模型直接吐收尾文案(治「收尾后还打转」的 over-loop) | — |
| ② | 迭代上限 recursion_limit | MAIN_AGENT_MAX_ITERATIONS=30 → 61 步 |
| ③ | 全局超时 asyncio.wait_for | MAIN_AGENT_TIMEOUT_SEC=300 |
| ④ | 循环检测 LoopDetector | window=6 / threshold=4 → 注入换思路提示 |
over-loop 治理当前走执行层哨兵:终结硬停、post-fork 检索收权(
_MAIN_POSTFORK_SEARCH_DENIED)等「该收敛」的策略,现在都在工具执行层回哨兵(工具表不变)。这是「摘工具(模型调用前改工具表)省解码 vs 保 prompt cache 前缀」取舍下的当前默认——项目看重跨轮缓存命中,故默认选保缓存;代价是哨兵拦执行不拦 decode。是默认不是禁令:某条链路解码浪费明确压过缓存收益时,摘工具仍是正当选项,按场景重新称量即可(见M6.1/Mperf.1)。
3. fork 同质子 Agent + 安全四层#
flowchart TD
MAIN["主 loop Think 判定: fork 三件事<br/>① 能并行 ② 要隔离 ③ 链够深"] --> CALL
CALL{"dispatch_tool(demands)<br/>parallel_dispatch_tool(list)"} --> COVER["_ensure_platform_coverage<br/>补齐模型漏列的平台"]
COVER --> L1
subgraph GUARDS["fork 安全四层"]
L1["①深度上限 enter_fork()<br/>MAX_FORK_DEPTH=1 · ContextVar"]
L1 -->|"超深 → 拒绝, 返回字符串"| REJECT["不崩, 提示主 loop 自己处理"]
L1 -->|"放行"| SPAWN
SPAWN["生成 sub-tid + report_fork(路由父连接)<br/>create_agent(get_fast_llm, FULL_TOOL_SET, 同 prompt, 同中间件)"]
SPAWN --> L2["②超时+迭代 asyncio.wait_for(timeout=90)<br/>recursion_limit=13"]
L2 --> SUBRUN["子 AgentLoop 独立跑<br/>独立 thread_id · 继承父 session_dir"]
SUBRUN --> L3["③结果截断 truncate_tool_result ≤4000tok"]
SUBRUN -.每步.-> L4["④循环检测 LoopDetector 窗6/阈4"]
end
SUBRUN -.子的事件无前端连接 → 仅记日志(隔离).-> LOG["日志"]
L3 --> RET["str 回传主 loop · 异常一律转字符串"]
REJECT --> RET
RET --> BACK["主 loop Observe 子结果"]
style GUARDS fill:#5f1e1e,color:#fff
style SPAWN fill:#1e3a5f,color:#fff
同质硬约束: 子用 同一 FULL_TOOL_SET(含 dispatch 自身 → 可递归)+ 同 system_prompt + 同中间件栈。唯一差异是模型档位:子用 get_fast_llm()(关 reasoning 降延迟)。
fork 安全四层实现位置:
| 层 | 机制 | 位置 | 参数 |
|---|---|---|---|
| ① 深度上限 | ContextVar + enter_fork() | fork_guard.py | MAX_FORK_DEPTH=1 |
| ② 超时+迭代 | asyncio.wait_for + recursion_limit | dispatch_tool.py | SUB_AGENT_TIMEOUT_SEC=90 / SUB_AGENT_RECURSION_LIMIT=13 |
| ③ 结果截断 | truncate_tool_result + ToolGuardMiddleware | middleware.py | MAX_TOOL_RESULT_TOKENS=4000 |
| ④ 循环检测 | LoopDetector 滑动窗口 | middleware.py | window=6 / threshold=4 |
4. item_search 召回管道(数据流)#
flowchart LR
Q["query 文本"] --> EQ["TowerClient.encode_query<br/>→ 语义向量"]
UP["user_profile(仅 like)<br/>来自 ContextVar"] --> EU["TowerClient.encode_user<br/>→ 个性化向量"]
EQ --> FUSE["tower.fuse<br/>α·query + β·user (0.85/0.15)"]
EU --> FUSE
FUSE --> ANN["QdrantRecall.search(top_k+1, platform)<br/>COSINE · payload filter=platform"]
ANN --> FLOOR["相关度过滤 RELEVANCE_FLOOR=0.45<br/>(只挡乱码, 挡不住品类缺货)"]
FLOOR --> OUT["top_k 候选<br/>+ total_recall + truncated"]
subgraph BACKENDS["后端可降级"]
B1["EMBED_MODEL → Embedding API(BGE-M3)<br/>否则 3-gram 哈希 256 维"]
B2["QDRANT_URL → Qdrant Server<br/>否则 on-disk / :memory:"]
end
EQ -.-> B1
ANN -.-> B2
OUT --> NOTE["⚠ item_search 不做 cross-encoder 精排<br/>(精排在 category_insight)<br/>⚠ 不做跨平台并行(由 parallel_dispatch_tool fork)"]
style FUSE fill:#1e3a5f,color:#fff
style ANN fill:#2d4a2d,color:#fff
端到端数据流(贯穿 ItemCandidate 结构逐步增补字段):
用户意图
→ planner 意图结构化 + 币种/预算确定性解析(不靠模型)
→ item_search 三塔融合 + Qdrant 召回(可选个性化向量) [+基础字段]
→ category_insight RAG: OpenSearch/本地 hybrid + cross-encoder 精排
→ price_compare fx 静态表折算 USD [+price_usd]
→ shipping_calc 运费表 + 关税表公式 → 到手价 [+shipping/duty/landed_usd]
→ item_picker 对齐论文四工具 Filter→Matcher+Attenuator→Aggregator [+pick_reason]
负硬→exclude淘汰 / 正硬must+正软→Matcher(关键词+BGE-M3语义加分)
负软→Attenuator(关键词+语义减分) / +便宜度+评分 → topK
长期dislike + 会话P_t 按 {正/负}×{硬/软} 确定性并入四路
→ shopping_summary LLM(快档)只产[开场白文案+每件reason],其余字段(标题/价格/图/链接)
按 item_id 从候选确定性组装(终结,不再产偏好;确定字段不过模型防改写/截断)
── 收尾后 ──
→ curate_turn 记忆管家(后处理异步,唯一偏好写入口):本轮约束→P_t、一贯取向→长期库
→ persist_new_preferences 一贯取向落库(本地JSON / Redis)plaintext候选走 item_id、不当工具参数重吐:工具之间只传 item_id 列表,工具内按 id 从会话级候选登记表 (
_candidates,唯一真相源)hydrate回全量候选;各阶段算完的渐进字段(price_usd/landed_usd/ pick_reason)回写登记表。整包候选不再穿过模型——省掉「模型逐字重吐候选」的解码,也杜绝改写/截断 幻觉。原list[候选]参数降级为对模型不可见的注入参数(InjectedToolArg)。详见Mperf.1。
5. 一次请求的生命周期(时序)#
sequenceDiagram
participant U as 前端
participant S as server.py
participant T as asyncio.Task<br/>run_agent
participant A as AgentLoop
participant M as monitor + ConnMgr
participant Sub as 子 Agent
U->>S: POST /api/task {query}
S->>S: 登记 active_tasks[tid]<br/>(旧任务自动 cancel)
S->>T: asyncio.create_task(_runner)
S-->>U: 返回 thread_id (立即)
U->>S: WS /ws/{tid} (connect-first)
S->>M: ConnectionManager.connect
T->>A: agent.ainvoke
A->>M: session_created / assistant_call
M-->>U: 推事件 (按 tid 路由)
loop Think→Act→Observe→Reflect
A->>M: tool_start
M-->>U: 工具开始
alt fork
A->>Sub: dispatch_tool
A->>M: fork (路由父连接)
Sub-->>A: 子结果(截断)
end
A->>M: tool_end
M-->>U: 工具完成
end
A->>A: shopping_summary (终结)
T->>T: 写产物 + persist_preferences
T->>M: task_result (items)
M-->>U: 最终清单
T->>S: finally 按身份摘除 active_tasks[tid]
opt 取消
U->>S: POST /cancel
S->>T: task.cancel()
T->>M: task_cancelled
end
AGUI 事件类型(monitor.py,经 ConnectionManager.send_to_thread 按 thread_id 路由):
session_created · assistant_call · tool_start · tool_end · fork · task_result · task_cancelled · error
6. 外部依赖与本地降级#
| 组件 | 环境变量 | 真实服务 | 本地 fallback |
|---|---|---|---|
| 编码器 | EMBED_MODEL | OpenAI 兼容 embedding(BGE-M3) | 3-gram 哈希 256 维 |
| ANN | QDRANT_URL | Qdrant Server | on-disk / :memory: |
| 知识库 | OPENSEARCH_HOST | OpenSearch Hybrid(KNN 0.7 + BM25 0.3) | 进程内 TowerClient + token 重叠 |
| 精排 | RERANKER_ENDPOINT | 远程 cross-encoder | Jaccard token 重叠 |
| 网络搜索 | TAVILY_API_KEY | Tavily Search | 返回空 + 说明 |
| 长期记忆 | STORE_BACKEND | Redis hash | JSON 本地文件 |
| LLM | LLM_MAIN / LLM_FAST / LLM_JUDGE | OpenAI 兼容(DashScope/Qwen…) | — |
7. 实现注记(与文档直觉不同的真实落地点)#
- fork 深度上限是 1(
MAX_FORK_DEPTH=1):只允许一层 fork,子 Agent 不能再 fork,递归被深度闸拦在动机层。 - 子 Agent 用
get_fast_llm()(同LLM_MAIN但关 reasoning):这是「同质 fork」的唯一现实裂口,理由是降解码延迟;工具集 / prompt / 中间件仍三同。 - 精排在
category_insight而非item_search:item_search只做三塔融合 + Qdrant 召回 + 绝对相关度阈值;cross-encoder 精排服务于 RAG 品类卡。 - 确定性优先:币种 / 预算 / 汇率 / 关税 / dislike 黑名单全走规则,不靠 LLM 自由猜测;
item_picker硬门 + 关键词打分走确定性规则、软偏好另叠一路 BGE-M3 语义打分(对齐 RecBot 四工具,不调 LLM 生成,见docs/milestones/M-检索层对齐RecBot-语义打分与记忆作用位置.md)。 - 降级不崩:所有外部依赖均有本地确定性 fallback,远程挂了 log warning 并返回降级结果,主链路不中断 —— CI / 离线可跑。
- ToolGuard 工具权限分层:
DEPTH0_ONLY_TOOLS(price_compare/item_picker/shopping_summary)、MAIN_ONLY_CONTEXT_TOOLS(planner/category_insight)、FORK_TOOLS仅主 loop 可用;RETRIEVAL_TOOLS(item_search/web_search)受树级检索预算(=8)管制。