diff --git a/apps/api/alembic/versions/f7a8b9c0d1e2_add_agentic_workflow_trace_columns.py b/apps/api/alembic/versions/f7a8b9c0d1e2_add_agentic_workflow_trace_columns.py new file mode 100644 index 000000000..a7a4c5513 --- /dev/null +++ b/apps/api/alembic/versions/f7a8b9c0d1e2_add_agentic_workflow_trace_columns.py @@ -0,0 +1,46 @@ +"""add agentic workflow trace columns + +Revision ID: f7a8b9c0d1e2 +Revises: f6a7b8c9d0e1 +Create Date: 2026-05-14 03:20:00.000000 + +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +revision: str = "f7a8b9c0d1e2" +down_revision: Union[str, Sequence[str], None] = "f6a7b8c9d0e1" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.add_column( + "retrieval_runs", + sa.Column("parent_run_id", sa.String(length=36), nullable=True), + ) + op.add_column( + "retrieval_runs", + sa.Column("workflow_step_id", sa.String(length=64), nullable=True), + ) + op.add_column( + "retrieval_runs", + sa.Column("workflow_plan", sa.JSON(), nullable=True), + ) + op.create_index( + "ix_retrieval_runs_parent_run_id", + "retrieval_runs", + ["parent_run_id"], + unique=False, + ) + + +def downgrade() -> None: + op.drop_index("ix_retrieval_runs_parent_run_id", table_name="retrieval_runs") + op.drop_column("retrieval_runs", "workflow_plan") + op.drop_column("retrieval_runs", "workflow_step_id") + op.drop_column("retrieval_runs", "parent_run_id") diff --git a/apps/api/tests/contract/test_agentic_answer_policy_contract.py b/apps/api/tests/contract/test_agentic_answer_policy_contract.py new file mode 100644 index 000000000..e0c0cf1d2 --- /dev/null +++ b/apps/api/tests/contract/test_agentic_answer_policy_contract.py @@ -0,0 +1,25 @@ +from __future__ import annotations + +import pytest + +from shared.services.retrieval.agentic.policy import attempt_answer +from shared.services.retrieval.agentic.types import AgentRunConfig, AgentState + + +async def _malformed_json_wrapper(_prompt: str) -> str: + return '{"status": "DONE", "answer": "truncated"' + + +@pytest.mark.asyncio +async def test_attempt_answer_should_not_expose_malformed_json_wrapper() -> None: + status, answer, reason = await attempt_answer( + _malformed_json_wrapper, + query="What changed?", + evidence_text="┈ evidence", + state=AgentState(), + config=AgentRunConfig(), + ) + + assert status == "NOT_FOUND" + assert answer == "" + assert reason == "attempt_answer returned malformed JSON" diff --git a/apps/api/tests/migrations/test_schema_contract.py b/apps/api/tests/migrations/test_schema_contract.py index 47ccfb235..7cbf27eb7 100644 --- a/apps/api/tests/migrations/test_schema_contract.py +++ b/apps/api/tests/migrations/test_schema_contract.py @@ -233,3 +233,26 @@ def test_api_standalone_mode_should_create_auth_user_table_before_migrations( "updatedAt", }.issubset(columns) assert email_unique_count == 1 + + +def test_agentic_retrieval_trace_schema_matches_orm(migrated_head_engine: Engine) -> None: + with migrated_head_engine.begin() as connection: + run_columns = set( + connection.execute( + text( + """ + SELECT column_name + FROM information_schema.columns + WHERE table_name = 'retrieval_runs' + """ + ) + ) + .scalars() + .all() + ) + + assert { + "parent_run_id", + "workflow_step_id", + "workflow_plan", + }.issubset(run_columns) diff --git a/docs/agentic-rag-audit-20260513.md b/docs/agentic-rag-audit-20260513.md deleted file mode 100644 index fe09703ce..000000000 --- a/docs/agentic-rag-audit-20260513.md +++ /dev/null @@ -1,353 +0,0 @@ -# Agentic RAG 流程审计更新 - -日期:2026-05-13 - -范围:基于 `AGENTS.md`、`.agent/skills/agentic_debug_patterns/SKILL.md`、既有 trace 目录 `/Users/wuchengke/Desktop/agentic_e2e_traces/20260513_183340`,以及额外运行的典型用例结果。 - -额外 trace 输出目录: - -- `/Users/wuchengke/Desktop/agentic_e2e_traces/20260513_190749_extra` -- `/Users/wuchengke/Desktop/agentic_e2e_traces/20260513_210646_extra_batch` - -## 结论摘要 - -这套 agentic RAG 的总体方向是成立的:外层 workflow 可以把复杂问题拆成多个 retrieve/synthesize step 并并发执行;内层 retrieval agent 能基于 KG 选文档、树导航、发现补充路径,并把 connected image/table 内嵌回证据树。尤其是全局图片/图表类问题,已有路径可以通过 `connect_to` 找回资源所属文本 section。 - -但从 harness 工程师和真实用户视角看,目前仍有几个核心逻辑风险: - -1. 图表资源的“证据渲染归属”和“返回引用归属”不一致。渲染树里通常能用 `connect_to` 找到底层 owner section,但 `referenced_chunks`/citation 仍可能显示物理路径 `Root`,这会直接破坏用户理解图表出处。 -2. discovery merge 过于积极。即使 BFS 已经 `STOP`,后置 discovery 仍会把深层或邻近年份路径并入证据,导致 outline 类问题和窄 section 问题出现噪声。 -3. 预算和状态分类混淆。无证据问题会被包装成 `budget_stop`,掩盖真实原因;同时 bootstrap/revision 小预算耗尽时,整体 wallet 仍可能很充足,用户看到的失败原因不准确。 -4. 多文件/多 step 并发在外层有效,但单 step 内的多 doc 导航仍偏串行;更重要的是 `discovery_auto` 会把弱相关文档强行并入,容易在跨年份、跨主题问题上污染预算和证据。 -5. 回答 JSON 解析不够稳健。`attempt_answer` 返回含换行的 JSON-like 文本时会解析失败,单步用户可能看到 JSON wrapper。 -6. trace DB schema 与 ORM 不一致,导致 agentic trace 入库失败,削弱 harness 可观测性。 - -## 本次补跑用例 - -### T1_Outline_Extra - -Query: - -> 民生证券这份利率专题研报的整体结构是什么?包含哪些主要章节? - -结果: - -- Router:`workflow_single_step` -- LLM calls:4 -- refs:13 -- elapsed:约 13s -- action:`kg_document_select -> navigate -> discovery_select -> attempt_answer` - -观察: - -- `navigate` 在 root 层正确选择 `STOP`,这对“整体结构/主要章节”类问题是合理的。 -- 但后续 `discovery_select` 又选入了深层路径,如 `2 阶段性调整.../2.1.1 基本面企稳` 和 `5、2024:“资产荒”的极致演绎`。 -- 最终 evidence 约 7395 chars,answer 只有约 149 chars,说明证据明显过量。 -- `referenced_chunks` 里部分 image/table 的 section 显示为 `Root`,但 evidence tree 实际把它们挂在更具体 leaf section 下。 - -判断: - -这是 discovery merge 策略的问题。对 root outline 查询,BFS 已经完成任务后,不应默认再并入深层 discovery 结果。否则用户问“目录结构”,结果引用中会混入某些深层图表,影响可信度。 - -### T2_Deep_Section_Extra - -Query: - -> 2016年债市走牛的几个阶段中,机构行为是如何推动行情演绎的?有哪些相关图表说明? - -结果: - -- Router:`workflow_single_step` -- LLM calls:4 -- refs:51 -- elapsed:约 29s -- evidence:约 22960 chars -- wallet context:`TIGHT` - -观察: - -- `navigate` 选择了 `NAVIGATE`,并带 `FIND_IMAGES`、`FIND_TABLES`,方向正确。 -- 但 root scope 的 asset tool 拉入了过多全局资源;同时 discovery 又选中父级 `1、2016:机构行为助推行情演绎`,导致 hydration 范围扩大。 -- evidence 里实际有图2、图3、图4、图5等图题和图片描述,但模型回答中仍说“未提供图表具体标题/编号”。 -- refs 达到 51,包含不少 2018、2019、2023、2024 等非目标年份资源。 -- 2016 相关图片在引用元数据中仍有 `section=Root` 的情况,虽然它们通过 `connect_to` 在 evidence 中被放回了具体 section。 - -判断: - -这是窄 section + 图表问题的典型失败形态:导航方向正确,但工具作用域过宽、discovery 过宽、证据渲染噪声大,导致模型虽然拿到了图表,却没有稳定提取图题和归属。 - -### T3_Compare_Extra - -Query: - -> 对比2024年和2025年AI安全市场规模,并结合证据给出变化原因。 - -结果: - -- Router:`workflow_decomposed` -- LLM calls:27 -- refs:1 -- elapsed:约 38s -- plan:s1 查 2024 市场规模,s2 查 2025 市场规模,s3 查变化原因,s4 synthesize - -观察: - -- 外层 workflow 确认可以并发执行多个 retrieve step,三个 retrieve step 的 KG select 和 navigate 调用是交错发生的。 -- 三个 retrieve step 最终都进入 revision,然后以 `budget_stop` 结束。 -- 总体 wallet 仍有大量剩余,但 bootstrap/revision 局部预算先被耗尽,最终对用户呈现为“预算停止”。 -- 实际语义更接近:KB 中缺少可支撑 2024/2025 AI 安全市场规模对比的证据。 -- `discovery_auto` 因年份词匹配,把债券研报等弱相关文档带入候选,造成预算消耗和路径污染。 - -判断: - -这是预算状态和无证据状态混淆。对用户来说,“知识库没有足够证据”和“预算不够”是两类完全不同的反馈;当前状态分类会误导用户,也会误导 harness 判断。 - -## 核心问题清单 - -### 1. 图表资源归属在 citation 层丢失 - -涉及核心设计: - -- 每个独立图表/图片/表格都应通过 `connect_to` 找到底层 section 归属。 -- 物理资源 chunk 的 `path` 可能是 `images/...` 或 `tables/...`,甚至 DB section 可能挂在 `Root`。 -- 逻辑归属应以 text chunk 的 `metadata.connect_to[].target` 为准,`target` 指向 image/table chunk_id。 - -当前表现: - -- evidence tree 渲染阶段多数情况下能用 owner path 把资源挂回 leaf section。 -- 但最终 `referenced_chunks`/citation 仍可能使用资源 chunk 自身的 `section_path`,因此显示 `Root`。 - -影响: - -- 用户看到图表出处为 `Root`,无法判断它属于哪个章节。 -- 对图表比较、章节归因、报告复核非常不友好。 -- 这和“每个独立图表都有 `connect_to` 找到一个底层 section 归属”的设计要求冲突。 - -建议: - -- citation/ref 组装时优先使用 `owner_section_path`,只有不存在时才回退到物理 `section_path`。 -- 返回结构中建议同时保留: - - `owner_section_path`:逻辑归属,用于用户展示和排序。 - - `physical_section_path`:数据库/资源物理挂载位置,用于调试。 - - `connect_to_source_chunk_id`:是哪一个 text chunk 证明了该资源归属。 -- 对 image/table 引用增加断言:若存在 `connect_to` owner,则展示 section 不应为 `Root`。 - -### 2. discovery merge 对 STOP 和 outline 查询缺少门控 - -当前表现: - -- T1 root outline 查询已经由 `navigate` 正确 `STOP`。 -- 后续 `discovery_select` 仍并入深层路径和资源。 - -影响: - -- 简单结构问题证据膨胀。 -- 引用混入深层内容,用户会怀疑答案是不是依据了错误章节。 -- 预算被无谓消耗。 - -建议: - -- 对 outline/structure/catalogue 类意图设置 discovery gate: - - 若 root STOP 且问题不要求“细节/图表/数据”,跳过 discovery hydration。 - - 或只允许 discovery 返回 top-level structural sections,不 hydrate leaf content/assets。 -- `discovery_select` 的 prompt 应明确区分: - - structure query:只补结构遗漏。 - - evidence query:可补 leaf 内容。 - - asset query:可补 image/table。 - -### 3. 图表工具作用域过宽 - -当前表现: - -- T2 中 `NAVIGATE + FIND_IMAGES/FIND_TABLES` 方向正确,但 root 或父级 scope asset extraction 拉入大量非目标年份图表。 -- 后续 trimming 虽然会删一部分,但已经消耗 context 和模型注意力。 - -影响: - -- 窄问题变成大范围 evidence dump。 -- 模型可能拿到正确图题却没有稳定使用,反而回答“没有具体标题”。 -- refs 过多,前端引用列表不可读。 - -建议: - -- 当 action 为 `NAVIGATE` 且有 selected leaf paths 时,asset tools 默认只对 selected paths 或其 owner-linked assets 生效。 -- 只有 action 为 root `STOP` 且 query 明确要求“列出全部图表/图片/表格”时,才允许文档级全量 asset pull。 -- 对 `FIND_IMAGES/FIND_TABLES` 的输出增加 owner filter:资源必须能通过 `connect_to` 归属到当前 selected subtree。 - -### 4. 预算分配与状态管理需要区分技术预算和语义失败 - -当前表现: - -- T3 三个 retrieve step 最终都是 `budget_stop`。 -- 但总 wallet 明显还有剩余,真正失败原因是没有足够证据。 -- bootstrap/revision 局部预算耗尽被升级成 step 级 budget stop。 - -影响: - -- 用户会以为“系统钱/上下文不够”,而不是“知识库无证据”。 -- harness 也难以判断是预算策略问题、检索召回问题还是 KB 数据缺失。 - -建议: - -- step status 拆分: - - `not_found_no_evidence` - - `not_found_low_confidence` - - `budget_exhausted_bootstrap` - - `budget_exhausted_context` - - `budget_exhausted_total` -- synthesize 时保留每个 retrieve step 的 semantic reason,不要只看 stop_reason 字符串。 -- revision loop 中,如果第一轮和第二轮文档选择高度重复且 verdict 是“KB 缺证据”,应提前停止,避免继续烧 bootstrap。 -- `BudgetWallet` 的 reclaimed budget 如果暂不重分配,snapshot 文案应避免暗示这些预算已重新可用。 - -### 5. Planner 缺少 KB inventory,导致 plan reasoning 误报 - -当前表现: - -- T4 中 planner reasoning 出现 “knowledge base is empty”。 -- 实际 trace 中 KB 并不为空。 - -判断: - -`QueryPlanner.plan()` 支持 `kb_total_docs/kb_total_chunks` 参数,但 workflow 调用路径没有传入真实 inventory,默认值为 0。 - -影响: - -- plan reasoning 不可信。 -- 对调试和用户解释都很危险。 - -建议: - -- `_load_or_plan()` 前读取当前 namespace 的 KB inventory,并传给 planner。 -- workflow plan cache key 应包含 KB version 或文档集合 fingerprint,否则 KB 更新后可能复用旧 plan。 - -### 6. `attempt_answer` JSON 解析不稳健 - -当前表现: - -- T4 中 `attempt_answer` 返回 JSON-like 内容,但因 raw newline 或不合规转义导致 parse 失败。 -- parse 失败后逻辑把原始字符串当作 DONE answer。 - -影响: - -- 单步用户可能看到 `{"status":"DONE","answer":...}` wrapper。 -- synth step 可能能“洗掉”问题,但 single-step 场景会暴露。 - -建议: - -- 增加 tolerant JSON repair,只修复回答字段中的裸换行/控制字符。 -- 如果解析失败且文本明显以 JSON object 开头,不应直接 `DONE raw`,而应降级重试或抽取 `answer` 字段。 - -### 7. trace DB schema 与 ORM 不一致 - -当前表现: - -- `retrieval_runs.parent_run_id/workflow_step_id/workflow_plan` 在 ORM 中存在。 -- alembic migration 中未创建这些列。 -- trace create_run 报 `UndefinedColumnError`。 - -影响: - -- DB trace 不可用。 -- harness 只能依赖 Markdown trace,无法做结构化聚合和回归分析。 - -建议: - -- 补 migration。 -- 增加一个轻量 schema contract test,覆盖 `RetrievalTraceRecorder.create_run()`。 - -### 8. 多文件并发导航的现状 - -已确认: - -- 外层 workflow retrieve steps 使用 topological batch 并发执行。 -- T3 中多个 retrieve step 的 KG select/navigate 调用交错,说明并发有效。 - -风险: - -- 单个 retrieve step 内 selected docs 仍偏串行。 -- `discovery_auto` 追加的弱相关文档没有足够 domain guard,T3 因年份匹配引入了债券研报。 - -建议: - -- 对 `discovery_auto` 文档追加设置最低 domain relevance: - - 文档 title/summary/keywords 至少命中主题实体。 - - 或要求 bottom chunk 与 query 的非时间词、非通用词有足够 overlap。 -- 单 step 多 doc 可考虑并发,但要先修好 doc relevance guard,否则并发只会更快地放大噪声。 - -## 遗留与冗余代码观察 - -### Legacy retrieval 路径仍和 agentic 路径混杂 - -`run_retrieval_query()` 中同时存在 agentic workflow 和 legacy 3-channel RRF 排序/graph routing。若 agentic 已是主路径,建议把 legacy 路径隔离为明确 fallback,避免后续改动时误改两套逻辑。 - -### 旧 graph/discovery helper 有疑似未使用分支 - -`agentic/orchestrator.py` 附近存在 `_grep_discover_document_ids`、`_expand_by_edges` 等老式发现逻辑痕迹。若主流程已经切到 bottom discovery + KG select,应确认这些 helper 是否仍被调用;未调用则标记删除或迁移到测试辅助。 - -### path dedup 当前依赖“一叶一文本 chunk”隐含前提 - -当前 `_hydrate_paths_to_rows` 用 path-level `seen_paths` 是安全的,因为解析模型近似保持“一 leaf section 一个 text chunk”。但如果未来 parser 把一个 leaf section 拆成多个 text chunks,path-level dedup 会丢内容。 - -建议: - -- 在注释和测试中写明该前提。 -- 或把 dedup key 改成 `(document_id, section_path, chunk_id)`,再在 render 层控制同 section 合并。 - -## 建议优先级 - -P0: - -1. 修复 image/table citation 归属:优先展示 `connect_to` owner section,不再把有 owner 的图表显示成 `Root`。 -2. 修复 trace DB migration,恢复 harness 结构化观测。 -3. 修复 `attempt_answer` JSON parse fallback,避免把 wrapper 暴露给用户。 - -P1: - -1. 对 root STOP/outline query 增加 discovery gate。 -2. 收紧 asset tool 作用域:`NAVIGATE + selected paths` 时只找 selected subtree 的 connected assets。 -3. 拆分 `budget_stop` 与 `not_found` 状态,synthesize 阶段保留真实失败原因。 - -P2: - -1. 给 planner 传真实 KB inventory,并把 KB fingerprint 纳入 plan cache key。 -2. 给 `discovery_auto` 增加 domain relevance guard。 -3. 清理 legacy helper 和未使用 discovery/graph 分支。 -4. 为 path dedup 增加未来多 chunk leaf 的保护测试。 - -## 建议回归用例 - -1. Outline STOP 不应 hydrate 深层 leaf: - - Query:`民生证券这份利率专题研报的整体结构是什么?包含哪些主要章节?` - - 断言:refs 中不应出现大量 image/table;深层 section 不应被 discovery 自动并入。 - -2. 2016 section 图表归属: - - Query:`2016年债市走牛的几个阶段中,机构行为是如何推动行情演绎的?有哪些相关图表说明?` - - 断言:所有相关 image/table citation 的展示 section 应为 2016 底层 section,而不是 `Root`。 - -3. 全量图表查询: - - Query:`列出AI安全大模型报告中所有的图表和图片,并简要描述每张图的内容。` - - 断言:允许 root/global asset pull,但每个独立图表仍应有 owner section;确实无底层 owner 的封面/前言图要显式标记为 document-level。 - -4. KB 无证据查询: - - Query:`对比2024年和2025年AI安全市场规模,并结合证据给出变化原因。` - - 断言:返回状态应是 no evidence / insufficient evidence,而不是 generic `budget_stop`。 - -5. Planner inventory: - - 构造非空 KB。 - - 断言 planner reasoning 不得出现 “knowledge base is empty”。 - -## 代码落点索引 - -- Workflow orchestration:`packages/shared-python/shared/services/retrieval/workflow/orchestrator.py` -- Planner:`packages/shared-python/shared/services/retrieval/workflow/planner.py` -- Workflow budget wallet:`packages/shared-python/shared/services/retrieval/workflow/wallet.py` -- Inner agent orchestrator:`packages/shared-python/shared/services/retrieval/agentic/orchestrator.py` -- Inner agent tools:`packages/shared-python/shared/services/retrieval/agentic/tools.py` -- Answer policy / JSON parse:`packages/shared-python/shared/services/retrieval/agentic/policy.py` -- Navigation tree render:`packages/shared-python/shared/services/retrieval/agent_navigate.py` -- Retrieval entry / hydration:`packages/shared-python/shared/services/retrieval/app_service.py` -- Retrieval trace:`packages/shared-python/shared/services/retrieval/agentic/trace.py` -- ORM retrieval tables:`packages/shared-python/shared/models/database/document.py` -- Migration:`apps/api/alembic/versions/e5f6a7b8c9d0_add_agentic_retrieval_tables.py` -- Debug harness:`apps/worker/debug_agentic_e2e.py` - diff --git a/packages/shared-python/shared/services/retrieval/agentic/budget.py b/packages/shared-python/shared/services/retrieval/agentic/budget.py index 1f4940bdd..cc8e58c7b 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/budget.py +++ b/packages/shared-python/shared/services/retrieval/agentic/budget.py @@ -3,7 +3,7 @@ import asyncio from dataclasses import dataclass -from typing import Literal +from typing import Any, Literal BudgetPoolName = Literal["bootstrap", "planning", "context"] @@ -64,7 +64,7 @@ def __init__( self.total_docs = 0 self.explored_chunks = 0 self.explored_docs = 0 - self.trimmed_paths: list[dict[str, str]] = [] + self.trimmed_paths: list[dict[str, Any]] = [] def remaining(self, pool: BudgetPoolName) -> int: return self._pools[pool].remaining diff --git a/packages/shared-python/shared/services/retrieval/agentic/orchestrator.py b/packages/shared-python/shared/services/retrieval/agentic/orchestrator.py index 77bf80bca..8aa58964d 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/orchestrator.py +++ b/packages/shared-python/shared/services/retrieval/agentic/orchestrator.py @@ -48,6 +48,43 @@ +def _with_context_prompt_projection( + snapshot: dict[str, object], + *, + prompt_tokens: int, +) -> dict[str, object]: + """Return a display snapshot that includes the upcoming answer prompt.""" + projected: dict[str, object] = dict(snapshot) + context_raw = projected.get('context') or {} + if not isinstance(context_raw, dict): + return projected + context = dict(context_raw) + used = int(context.get('used', 0) or 0) + reserved = int(context.get('reserved', 0) or 0) + capacity = int(context.get('capacity', 0) or 0) + projected_used = min(capacity, used + max(int(prompt_tokens), 0)) + projected_remaining = max(capacity - projected_used - reserved, 0) + context.update({ + 'used_projected_before_answer': projected_used, + 'answer_prompt_estimate': max(int(prompt_tokens), 0), + 'remaining': projected_remaining, + 'used_pct': 100 if capacity <= 0 else min( + 100, + int(round((projected_used + reserved) * 100 / capacity)), + ), + }) + if projected_remaining <= 0: + context['status'] = 'EXHAUSTED' + elif context['used_pct'] >= 80: + context['status'] = 'CRITICAL' + elif context['used_pct'] >= 50: + context['status'] = 'TIGHT' + else: + context['status'] = 'HEALTHY' + projected['context'] = context + return projected + + @@ -365,7 +402,7 @@ async def _trim_evidence_to_budget( ) -> str: full_text = await _render_evidence(db, doc_trees, doc_id_to_name) target = int(max(context_remaining, 0) * safety_margin) - if target <= 0 or estimate_tokens(full_text) <= target: + if estimate_tokens(full_text) <= target: return full_text candidates: list[tuple[str, str, tuple[float, float, float], int]] = [] @@ -394,7 +431,7 @@ async def _trim_evidence_to_budget( candidates.append((doc_id, path, score, _estimate_chunks_tokens(chunks))) current_estimate = estimate_tokens(full_text) - removed: list[dict[str, str]] = [] + removed: list[dict[str, Any]] = [] for doc_id, path, _score, token_estimate in sorted( candidates, key=lambda item: (item[2], -item[3]), @@ -402,7 +439,16 @@ async def _trim_evidence_to_budget( if current_estimate <= target: break if _pop_leaf_path(doc_trees[doc_id], path): - removed.append({'document_id': doc_id, 'path': path}) + confidence_score, discovery_score, importance_score = _score + removed.append({ + 'document_id': doc_id, + 'document_name': doc_id_to_name.get(doc_id, doc_id), + 'path': path, + 'confidence_score': round(confidence_score, 4), + 'discovery_score': round(discovery_score, 4), + 'importance_score': round(importance_score, 4), + 'token_estimate': token_estimate, + }) current_estimate = max(current_estimate - token_estimate, 0) if ledger is not None: @@ -536,7 +582,10 @@ async def run( is returned. """ from shared.services.retrieval.agentic import tools - from shared.services.retrieval.agentic.policy import attempt_answer + from shared.services.retrieval.agentic.policy import ( + attempt_answer, + estimate_attempt_answer_prompt_tokens, + ) from shared.services.retrieval.llm_adapter import create_retrieval_vlm_fn vlm_fn = create_retrieval_vlm_fn() @@ -724,7 +773,11 @@ async def _context_llm_call(prompt): 'chunk_id': r.get('chunk_id', ''), 'document_id': r.get('document_id', ''), 'chunk_type': r.get('chunk_type', ''), - 'section_path': r.get('section_path', ''), + 'section_path': ( + r.get('source_file_name', '') + if r.get('section_path') == 'Root' + else r.get('section_path', '') + ), 'file_path': r.get('file_path', ''), } for r in discovery_rows[:top_k] @@ -789,6 +842,7 @@ async def _context_llm_call(prompt): evidence_text = '' revision_hint: str | None = None stop_reason = 'max_revisions' + failure_reason = '' for round_idx in range(config.max_revisions + 1): if state.elapsed_ms >= config.latency_budget_ms: @@ -878,9 +932,11 @@ async def _context_llm_call(prompt): break state.step_count += 1 - # ★ Asset collection (deferred reconcile) — runs if LLM selected tools - # scope is passed directly: None (root), str, or list[str] (multi-scope). - # asset_filter_step handles all forms natively. + # ★ Asset collection (deferred reconcile) — runs if LLM selected tools. + # If this navigation call selected sections, bind asset tools to + # those selections; otherwise keep the current scope (STOP/root). + selected_asset_scopes = list(step_node.confidence.keys()) + asset_scope = selected_asset_scopes or scope for asset_tool in asset_tools: if asset_tool not in ('FIND_IMAGES', 'FIND_TABLES'): continue @@ -889,13 +945,16 @@ async def _context_llm_call(prompt): db, document_id=doc.document_id, job_result_id=job_result_id, - scope_path=scope, + scope_path=asset_scope, asset_type=asset_type, ) if asset_chunks: doc_pending_assets.extend(asset_chunks) - scope_display = scope if isinstance(scope, list) else (scope or 'root') + scope_display = ( + asset_scope if isinstance(asset_scope, list) + else (asset_scope or 'root') + ) if trace_enabled: trace.record_step( 'asset_filter_step', ToolResult( @@ -903,6 +962,7 @@ async def _context_llm_call(prompt): payload={ 'document_id': doc.document_id, 'scope': scope_display, + 'navigation_scope': scope if isinstance(scope, str) else (scope or 'root'), 'asset_type': asset_type, 'chunks_found': len(asset_chunks) if asset_chunks else 0, }, @@ -978,6 +1038,10 @@ async def _context_llm_call(prompt): # ── Post-BFS: Discovery selection step ───────────────────── doc_hints = discovery_by_doc.get(doc.document_id, []) if doc_hints and planning_llm_fn is not None and state.elapsed_ms < config.latency_budget_ms: + discovery_exclude_paths = { + key.split('::', 1)[1] + for key in root.collect_all_paths(doc.document_id) + } doc_discovery_llm_fn = self._budgeted_discovery_llm_fn( state, cast(LLMFn, llm_fn), @@ -994,6 +1058,7 @@ async def _context_llm_call(prompt): namespace=namespace, doc_name=doc_name, discovery_hints=doc_hints, + exclude_paths=discovery_exclude_paths, revision_hint=revision_hint, budget_snapshot=state.ledger.snapshot() if state.ledger else None, ) @@ -1066,15 +1131,60 @@ async def _context_llm_call(prompt): state.ledger.mark_explored(docs=1) # ── Phase 3: Render evidence + attempt_answer ──────────────── + budget_snapshot_before_answer = state.ledger.snapshot() if state.ledger else None + context_remaining = ( + state.ledger.remaining('context') if state.ledger else config.token_budget_total + ) + answer_prompt_overhead = estimate_attempt_answer_prompt_tokens( + query=query, + evidence_text='', + state=state, + config=config, + budget_snapshot=budget_snapshot_before_answer, + ) evidence_text = await _trim_evidence_to_budget( db, doc_trees=state.doc_trees, doc_id_to_name=state.doc_id_to_name, - context_remaining=state.ledger.remaining('context') if state.ledger else config.token_budget_total, + context_remaining=max(context_remaining - answer_prompt_overhead, 0), user_id=user_id, namespace=namespace, ledger=state.ledger, ) + answer_prompt_tokens = estimate_attempt_answer_prompt_tokens( + query=query, + evidence_text=evidence_text, + state=state, + config=config, + budget_snapshot=budget_snapshot_before_answer, + ) + budget_snapshot_for_answer = ( + _with_context_prompt_projection( + state.ledger.snapshot(), + prompt_tokens=answer_prompt_tokens, + ) + if state.ledger else None + ) + if budget_snapshot_for_answer is not None: + answer_prompt_tokens = estimate_attempt_answer_prompt_tokens( + query=query, + evidence_text=evidence_text, + state=state, + config=config, + budget_snapshot=budget_snapshot_for_answer, + ) + budget_snapshot_for_answer = _with_context_prompt_projection( + state.ledger.snapshot(), + prompt_tokens=answer_prompt_tokens, + ) + context_budget = (budget_snapshot_for_answer.get('context') or {}) + logger.info( + ' agentic: answer context projection ' + f'prompt_tokens={answer_prompt_tokens} ' + f'remaining={context_budget.get("remaining")}/' + f'{context_budget.get("capacity")} ' + f'status={context_budget.get("status")}' + ) if context_llm_fn is None: stop_reason = 'no_llm' @@ -1105,7 +1215,7 @@ async def vlm_context_call(prompt, _vlm_fn=vlm_fn): config=config, vlm_fn=vlm_context_call if vlm_fn else None, image_urls=evidence_image_urls or None, - budget_snapshot=state.ledger.snapshot() if state.ledger else None, + budget_snapshot=budget_snapshot_for_answer, ) except BudgetExceeded: logger.info(' agentic: context budget exhausted before attempt_answer') @@ -1136,9 +1246,11 @@ async def vlm_context_call(prompt, _vlm_fn=vlm_fn): if status == 'DONE': stop_reason = 'answer_done' + failure_reason = '' break # ── NOT_FOUND: prepare revision ────────────────────────────── + failure_reason = reason if round_idx >= config.max_revisions: stop_reason = 'max_revisions' break @@ -1211,8 +1323,9 @@ async def vlm_context_call(prompt, _vlm_fn=vlm_fn): # Collect referenced chunk IDs from all doc trees all_refs: list[dict[str, str]] = [] seen_ref_ids: set[str] = set() - for doc_tree in state.doc_trees.values(): - for ref in doc_tree.collect_referenced_ids(): + for doc_id, doc_tree in state.doc_trees.items(): + doc_name = state.doc_id_to_name.get(doc_id, doc_id) + for ref in doc_tree.collect_referenced_ids(document_name=doc_name): cid = ref.get('chunk_id', '') if cid and cid not in seen_ref_ids: seen_ref_ids.add(cid) @@ -1232,6 +1345,7 @@ async def vlm_context_call(prompt, _vlm_fn=vlm_fn): router_used=router_used, budget_snapshot=state.ledger.snapshot() if state.ledger else None, stop_reason=stop_reason, + failure_reason=failure_reason, ) logger.info( diff --git a/packages/shared-python/shared/services/retrieval/agentic/policy.py b/packages/shared-python/shared/services/retrieval/agentic/policy.py index 6d63035e7..25323580a 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/policy.py +++ b/packages/shared-python/shared/services/retrieval/agentic/policy.py @@ -21,6 +21,7 @@ from shared.services.retrieval.agentic.budget import BudgetExceeded from shared.services.retrieval.agentic.types import AgentRunConfig, AgentState from shared.services.retrieval.llm_adapter import LLMFn +from shared.utils.token_estimate import estimate_tokens def _parse_answer_response(text: str) -> dict[str, Any] | None: @@ -39,6 +40,16 @@ def _parse_answer_response(text: str) -> dict[str, Any] | None: return None +def _looks_like_json_wrapper(text: str) -> bool: + """Detect malformed JSON-ish answer wrappers without exposing them.""" + stripped = text.strip() + if not stripped: + return False + if stripped.startswith('{') or stripped.endswith('}'): + return True + return bool(re.search(r'"(?:status|answer|reason)"\s*:', stripped)) + + def _is_external_http_url(url: str) -> bool: parsed = urlparse(str(url)) if parsed.scheme not in {'http', 'https'} or not parsed.hostname: @@ -53,6 +64,30 @@ def _is_external_http_url(url: str) -> bool: return not (addr.is_private or addr.is_loopback or addr.is_link_local) +def _budget_line_parts(budget_snapshot: dict | None, pool_name: str) -> dict[str, Any]: + pool = ((budget_snapshot or {}).get(pool_name) or {}) + capacity = pool.get('capacity', 'unknown') + remaining = pool.get('remaining', 'unknown') + used_pct = pool.get('used_pct', 'unknown') + remaining_pct: int | str = 'unknown' + try: + capacity_int = int(capacity) + remaining_int = int(remaining) + remaining_pct = 0 if capacity_int <= 0 else max( + 0, + min(100, round(remaining_int * 100 / capacity_int)), + ) + except (TypeError, ValueError): + pass + return { + 'status': pool.get('status', 'HEALTHY'), + 'remaining': remaining, + 'capacity': capacity, + 'used_pct': used_pct, + 'remaining_pct': remaining_pct, + } + + async def attempt_answer( llm_fn: LLMFn, *, @@ -80,12 +115,12 @@ async def attempt_answer( - status='NOT_FOUND', answer_text='', reason= → evidence was insufficient, reason is used as revision_hint """ - prompt_text = _ATTEMPT_ANSWER_PROMPT.format( + prompt_text = build_attempt_answer_prompt( query=query, - evidence_context=evidence_text, - revision_count=state.revision_count, - max_revisions=config.max_revisions, - context_status=((budget_snapshot or {}).get('context') or {}).get('status', 'HEALTHY'), + evidence_text=evidence_text, + state=state, + config=config, + budget_snapshot=budget_snapshot, ) verbose = os.environ.get('RETRIEVAL_AGENTIC_VERBOSE', '') == 'true' @@ -146,8 +181,10 @@ async def attempt_answer( parsed = _parse_answer_response(raw_response) if not parsed: - # Parse error: treat raw text as best-effort answer - return 'DONE', raw_response.strip(), 'parse_error — treating raw response as answer' + if _looks_like_json_wrapper(raw_response): + return 'NOT_FOUND', '', 'attempt_answer returned malformed JSON' + # Keep plain-text fallback for providers that ignore JSON mode entirely. + return 'DONE', raw_response.strip(), 'parse_error — treating plain text response as answer' status = str(parsed.get('status', 'DONE')).strip().upper() answer = str(parsed.get('answer', '')).strip() @@ -162,6 +199,53 @@ async def attempt_answer( return 'DONE', answer, '' +def build_attempt_answer_prompt( + *, + query: str, + evidence_text: str, + state: AgentState, + config: AgentRunConfig, + budget_snapshot: dict | None = None, +) -> str: + """Build the final answer prompt so trimming can estimate it beforehand.""" + planning = _budget_line_parts(budget_snapshot, 'planning') + context = _budget_line_parts(budget_snapshot, 'context') + return _ATTEMPT_ANSWER_PROMPT.format( + query=query, + evidence_context=evidence_text, + revision_count=state.revision_count, + max_revisions=config.max_revisions, + planning_status=planning['status'], + planning_remaining=planning['remaining'], + planning_capacity=planning['capacity'], + planning_used_pct=planning['used_pct'], + planning_remaining_pct=planning['remaining_pct'], + context_status=context['status'], + context_remaining=context['remaining'], + context_capacity=context['capacity'], + context_used_pct=context['used_pct'], + context_remaining_pct=context['remaining_pct'], + ) + + +def estimate_attempt_answer_prompt_tokens( + *, + query: str, + evidence_text: str, + state: AgentState, + config: AgentRunConfig, + budget_snapshot: dict | None = None, +) -> int: + """Estimate the exact prompt shape that will be charged to context budget.""" + return estimate_tokens(build_attempt_answer_prompt( + query=query, + evidence_text=evidence_text, + state=state, + config=config, + budget_snapshot=budget_snapshot, + )) + + _ATTEMPT_ANSWER_PROMPT = """\ You are a knowledge retrieval assistant. Answer the user's query based STRICTLY on the provided evidence. Do NOT use any external knowledge. @@ -176,7 +260,8 @@ async def attempt_answer( {evidence_context} REVISION: {revision_count} of {max_revisions} revisions used. -Context budget remaining is {context_status}; the evidence may have been trimmed. +Planning budget: {planning_status} ({planning_used_pct}% used, {planning_remaining_pct}% remaining, {planning_remaining}/{planning_capacity} remaining). +Context budget: {context_status} ({context_used_pct}% used, {context_remaining_pct}% remaining, {context_remaining}/{context_capacity} remaining); the evidence may have been trimmed. INSTRUCTIONS: 1. If the evidence contains enough information to answer the query, diff --git a/packages/shared-python/shared/services/retrieval/agentic/tools.py b/packages/shared-python/shared/services/retrieval/agentic/tools.py index 4774d2491..10ab54757 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/tools.py +++ b/packages/shared-python/shared/services/retrieval/agentic/tools.py @@ -40,6 +40,7 @@ merge_channels_rrf, ) from shared.services.retrieval.channels import content_channel, path_channel, term_channel +from shared.services.retrieval.lexical_text import normalize_section_path from shared.services.retrieval.llm_adapter import LLMFn @@ -684,12 +685,16 @@ async def navigate_step( tools_lines = ['\nOptional asset tools (usable with NAVIGATE or STOP):\n'] if total_images > 0: tools_lines.append( - f' FIND_IMAGES — Extract all image/chart assets under this scope ({total_images} available).\n' + f' FIND_IMAGES — Extract image/chart assets under the current scope ({total_images} available).\n' ) if total_tables > 0: tools_lines.append( - f' FIND_TABLES — Extract all table/data assets under this scope ({total_tables} available).\n' + f' FIND_TABLES — Extract table/data assets under the current scope ({total_tables} available).\n' ) + tools_lines.append( + ' Note: with NAVIGATE selections, asset tools are limited to the selected sections; ' + 'with STOP or no selections, they use the current scope.\n' + ) tools_block = ''.join(tools_lines) # 4. Format tree and build prompt @@ -823,6 +828,7 @@ async def discovery_select_step( namespace: str, doc_name: str = '', discovery_hints: list[dict[str, Any]], + exclude_paths: set[str] | None = None, revision_hint: str | None = None, budget_snapshot: dict | None = None, ) -> DocTreeNode: @@ -847,17 +853,23 @@ async def discovery_select_step( t0 = time.monotonic() try: # 1. Format hints for LLM (deduplicate by section_path) + exclude_set = { + normalize_section_path(path) + for path in (exclude_paths or set()) + if path + } hint_lines: list[str] = [] hint_by_path: dict[str, dict] = {} for h in hints: - sp = h.get('section_path', '') + sp = normalize_section_path(h.get('section_path', '')) if not sp or sp == 'Root': continue + if sp in exclude_set: + continue if sp in hint_by_path: continue # skip duplicate section_path - title = sp.rsplit(' / ', 1)[-1] if ' / ' in sp else sp summary = h.get('summary', '') or '' - hint_lines.append(f'▸ path="{sp}" {title} [Leaf]') + hint_lines.append(f'▸ path="{sp}"') if summary: clipped = summary[:300] hint_lines.append(f' {clipped}') diff --git a/packages/shared-python/shared/services/retrieval/agentic/trace.py b/packages/shared-python/shared/services/retrieval/agentic/trace.py index 9c48e7a30..87df0c94c 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/trace.py +++ b/packages/shared-python/shared/services/retrieval/agentic/trace.py @@ -170,6 +170,7 @@ async def complete( 'tokens_used': step_data.get('tokens_used', 0), }, latency_ms=step_data['latency_ms'], + token_count=step_data.get('tokens_used', 0), error=step_data.get('error'), created_at=step_data['created_at'], ) @@ -201,6 +202,7 @@ async def complete( final_doc_ids=doc_ids_in_result, result_provenance=provenance, latency_ms=total_latency, + token_count=sum(step.get('tokens_used', 0) for step in self._steps), completed_at=_now_utc(), ) ) diff --git a/packages/shared-python/shared/services/retrieval/agentic/types.py b/packages/shared-python/shared/services/retrieval/agentic/types.py index a448a5f6b..92c8b7f2f 100644 --- a/packages/shared-python/shared/services/retrieval/agentic/types.py +++ b/packages/shared-python/shared/services/retrieval/agentic/types.py @@ -134,7 +134,7 @@ def reparent_leaf_content(self) -> None: child.add_leaf_chunks(leaf_path, self.leaf_content.pop(leaf_path)) child.reparent_leaf_content() - def collect_referenced_ids(self) -> list[dict[str, str]]: + def collect_referenced_ids(self, *, document_name: str = '') -> list[dict[str, str]]: """Extract minimal chunk references from all hydrated leaves. Returns deduplicated list of {chunk_id, document_id, chunk_type, @@ -146,11 +146,14 @@ def collect_referenced_ids(self) -> list[dict[str, str]]: cid = row.get('chunk_id', '') if cid and cid not in seen: seen.add(cid) + section_path = row.get('section_path', '') + if section_path == 'Root' and document_name: + section_path = document_name refs.append({ 'chunk_id': cid, 'document_id': row.get('document_id', ''), 'chunk_type': row.get('chunk_type', ''), - 'section_path': row.get('section_path', ''), + 'section_path': section_path, 'file_path': row.get('file_path', ''), 'job_id': row.get('job_id', ''), }) @@ -204,6 +207,8 @@ class AgenticResult: - ``budget_snapshot``: final budget ledger state at run completion - ``stop_reason``: why the run terminated (answer_done / max_revisions / latency_budget / context_budget / no_llm / etc.) + - ``failure_reason``: semantic reason from the answer attempt when no + answer could be produced, e.g. evidence was insufficient. """ evidence_text: str answer_text: str = '' @@ -211,6 +216,7 @@ class AgenticResult: router_used: str = 'agentic_discovery_only' budget_snapshot: dict[str, Any] | None = None stop_reason: str = '' + failure_reason: str = '' @dataclass diff --git a/packages/shared-python/shared/services/retrieval/workflow/orchestrator.py b/packages/shared-python/shared/services/retrieval/workflow/orchestrator.py index 9c5a20aff..abce68b36 100644 --- a/packages/shared-python/shared/services/retrieval/workflow/orchestrator.py +++ b/packages/shared-python/shared/services/retrieval/workflow/orchestrator.py @@ -12,7 +12,7 @@ from shared.core.database import get_db_context from shared.services.retrieval.agentic.budget import BudgetLedger -from shared.services.retrieval.agentic.orchestrator import RetrievalAgent +from shared.services.retrieval.agentic.orchestrator import RetrievalAgent, _load_budget_inventory from shared.services.retrieval.agentic.types import AgenticResult from shared.services.retrieval.cache_service import ( get_cached_workflow_plan, @@ -66,6 +66,14 @@ async def run( bootstrap=planner_budget, per_doc_min_share=0, ) + total_chunks, total_docs, _chunks_count_by_doc = await _load_budget_inventory( + db, + user_id=user_id, + namespace=namespace, + exclude_document_ids=exclude_document_ids, + ) + planner_ledger.total_chunks = total_chunks + planner_ledger.total_docs = total_docs plan = await self._load_or_plan( user_id=user_id, namespace=namespace, @@ -75,6 +83,8 @@ async def run( max_steps=max_steps, wallet_total=wallet_total, per_retrieve=per_retrieve, + kb_total_docs=total_docs, + kb_total_chunks=total_chunks, ) wallet = BudgetWallet( @@ -153,6 +163,8 @@ async def _load_or_plan( max_steps: int, wallet_total: int, per_retrieve: int, + kb_total_docs: int, + kb_total_chunks: int, ) -> QueryPlan: try: cached = await get_cached_workflow_plan(user_id=user_id, namespace=namespace, query=query) @@ -168,7 +180,11 @@ async def _load_or_plan( total_budget=wallet_total, per_step_budget=per_retrieve, ) - plan = await planner.plan(query=query) + plan = await planner.plan( + query=query, + kb_total_docs=kb_total_docs, + kb_total_chunks=kb_total_chunks, + ) try: await set_cached_workflow_plan( user_id=user_id, @@ -330,7 +346,14 @@ async def _run_synthesize_step( def _step_result_from_agentic(step: PlannedStep, result: AgenticResult) -> StepResult: - status = 'budget_stop' if 'budget' in (result.stop_reason or '') else 'done' + if result.answer_text: + status = 'done' + elif result.failure_reason: + status = 'not_found' + elif 'budget' in (result.stop_reason or ''): + status = 'budget_stop' + else: + status = 'done' return StepResult( step_id=step.id, sub_query=step.sub_query, @@ -344,6 +367,7 @@ def _step_result_from_agentic(step: PlannedStep, result: AgenticResult) -> StepR budget_snapshot=result.budget_snapshot, router_used=result.router_used, stop_reason=result.stop_reason, + failure_reason=result.failure_reason, ) diff --git a/packages/shared-python/shared/services/retrieval/workflow/synthesizer.py b/packages/shared-python/shared/services/retrieval/workflow/synthesizer.py index 0e825eed9..89d4ec031 100644 --- a/packages/shared-python/shared/services/retrieval/workflow/synthesizer.py +++ b/packages/shared-python/shared/services/retrieval/workflow/synthesizer.py @@ -87,7 +87,23 @@ def _concat_final_parts(plan: QueryPlan, results: dict[str, StepResult]) -> str: for step in plan.steps if (result := results.get(step.id)) and result.answer_text.strip() ] - return "\n\n".join(fallback_parts) + if fallback_parts: + return "\n\n".join(fallback_parts) + + missing_reasons = [ + result.failure_reason.strip() + for step in plan.steps + if (result := results.get(step.id)) + and result.status == "not_found" + and result.failure_reason.strip() + ] + if missing_reasons: + return "未能基于当前知识库证据回答该问题:" + ";".join(dict.fromkeys(missing_reasons)) + + if any((result := results.get(step.id)) and result.status == "budget_stop" for step in plan.steps): + return "Unable to return a valid answer because the retrieval budget was exhausted." + + return "" def _format_prior_outputs(depends_on: list[str], prior_results: dict[str, StepResult]) -> str: diff --git a/packages/shared-python/shared/services/retrieval/workflow/types.py b/packages/shared-python/shared/services/retrieval/workflow/types.py index e196fd656..b9147856f 100644 --- a/packages/shared-python/shared/services/retrieval/workflow/types.py +++ b/packages/shared-python/shared/services/retrieval/workflow/types.py @@ -8,7 +8,7 @@ StepKind = Literal["retrieve", "synthesize"] OutputRole = Literal["final_part", "intermediate", "consumed_by_synthesis"] FinalStrategy = Literal["concat_final_parts", "last_synthesize", "template"] -StepStatus = Literal["done", "skipped", "error", "budget_stop"] +StepStatus = Literal["done", "skipped", "error", "budget_stop", "not_found"] @dataclass @@ -175,6 +175,7 @@ class StepResult: child_run_id: str | None = None router_used: str = "" stop_reason: str = "" + failure_reason: str = "" error: str | None = None def to_api_dict(self) -> dict[str, Any]: @@ -192,6 +193,7 @@ def to_api_dict(self) -> dict[str, Any]: "child_run_id": self.child_run_id, "router_used": self.router_used, "stop_reason": self.stop_reason, + "failure_reason": self.failure_reason, "error": self.error, }