xueyiming 1 неделя назад
Родитель
Сommit
7d285421cf

+ 17 - 13
find_agent_v2/AGENT_FLOW.md

@@ -1,7 +1,7 @@
 # 寻找 Agent v2 当前实现流程
 
 本文档描述 `find_agent_v2` 当前代码的真实执行路径、LLM 与宿主程序的职责边界、工具权限、
-数据表和状态流转。文档基于 2026-08-12 的实现。
+数据表和状态流转。文档基于 2026-08-14 的实现。
 
 ## 1. 总体架构
 
@@ -16,14 +16,15 @@ flowchart TD
     C --> D["创建 round<br/>phase=planning"]
 
     subgraph G["单轮受控自主图"]
-        D --> P["Supervisor<br/>提议 next_action / worker_count / evidence_scope"]
-        P --> Q["Policy Guard<br/>依据真实 DB 状态与预算审核"]
+        D --> P["Planner<br/>原始需求 → ExecutionPlan"]
+        P --> SP["Supervisor<br/>基于计划提议下一动作"]
+        SP --> Q["Policy Guard<br/>依据真实 DB 状态与预算审核"]
         Q --> S["Search<br/>搜索并写入候选"]
         Q --> E["Evidence Workers<br/>按建议范围补证"]
         Q --> V["Evaluator Workers<br/>评分并申请分池"]
-        S --> P
-        E --> P
-        V --> P
+        S --> SP
+        E --> SP
+        V --> SP
         Q --> X["结束本轮"]
     end
 
@@ -81,11 +82,12 @@ middleware,工具调用也接入独立重试;可通过 `FIND_AGENT_V2_FALLBA
 
 | 节点 | LLM 职责 | 可用工具 | 上下文范围 | 业务迭代提示值 |
 |---|---|---|---|---:|
-| Supervisor | 提议下一动作、搜索计划、补证范围和 Worker 数 | 无 | 全量 run 状态 | 2 |
-| Search | 执行搜索计划,不补证、不评分 | `search_videos_v2`、`query_find_agent_v2_state`、`delegate_agents_v2` | 全量状态 | 10 |
-| Evidence | 为候选补视频详情和双侧年龄画像 | `fetch_candidate_details_v2`、`fetch_candidate_portraits_v2`、`query_pending_candidates_v2`、`delegate_agents_v2` | 仅 pending 候选 | 12 |
-| Evaluator | 给出 R/E/S/V,申请 `primary` 或 `rejected` | `evaluate_candidates_v2`、`query_pending_candidates_v2`、`delegate_agents_v2` | 仅 pending 候选 | 每次 12 |
-| Report | 汇总最终结果 | `query_find_agent_v2_state` | 全量最终状态 | 4 |
+| Planner | 原始需求压缩为结构化 ExecutionPlan,并提出首次动作 | 无 | PlannerAssignment(唯一包含原始需求) | 2 |
+| Supervisor | 根据 ExecutionPlan 和执行状态提议下一动作、补证范围和 Worker 数 | 无 | SupervisorAssignment | 2 |
+| Search | 执行搜索计划,不补证、不评分 | `search_videos_v2`、`query_find_agent_v2_state`、`delegate_agents_v2` | SearchAssignment | 10 |
+| Evidence | 为候选补视频详情和双侧年龄画像 | `fetch_candidate_details_v2`、`fetch_candidate_portraits_v2`、`query_pending_candidates_v2` | EvidenceAssignment(单一分片) | 12 |
+| Evaluator | 给出 R/E/S/V,申请 `primary` 或 `rejected` | `understand_candidate_video_30s_v2`、`evaluate_candidates_v2`、`query_pending_candidates_v2` | EvaluationAssignment(评估摘要和单一分片) | 每次 12 |
+| Report | 汇总最终结果 | `query_find_agent_v2_state` | ReportAssignment | 4 |
 
 每个业务动作执行后都会回到 Supervisor。默认单轮最多 16 个业务动作、最多 3 次搜索动作;
 Worker 并发建议会被限制为 1~8。Evaluator 连续 3 次没有减少 pending 时判定技术失败。
@@ -310,7 +312,8 @@ flowchart TD
 
 每个节点记录:
 
-- 稳定的输入槽:当前轮次、原始任务、数据库快照、本轮搜索计划。
+- 单一 `stage_assignment` 输入槽;其内容与真实发送给模型的 JSON 完全一致。
+- 原始任务只出现在首次 PlannerAssignment,不进入后续阶段 Assignment。
 - 实际 system prompt、模型和工具列表。
 - 完整 LLM 消息链、usage、工具参数、工具返回、异常和迭代次数。
 - 节点最终输出。
@@ -327,7 +330,8 @@ flowchart TD
 |---|---|
 | `runner.py` | 创建 run,提供同步/异步/已准备任务入口 |
 | `agent.py` | 外层业务轮次、停止条件、终态和 Report |
-| `graph.py` | 编译真实 StateGraph、条件边和 Evaluator 消费保护 |
+| `graph.py` | 编译真实 StateGraph、构造阶段 Assignment、条件边和 Evaluator 消费保护 |
+| `context.py` | 序列化阶段 Assignment,并生成与真实模型输入一致的观测槽 |
 | `runtime.py` | LangChain create_agent、重试/fallback、usage 和受控委派 |
 | `prompts.py` | 节点提示词和公共业务规则 |
 | `tools.py` | v2 工具定义和节点 allowlist |

+ 4 - 2
find_agent_v2/README.md

@@ -39,8 +39,10 @@ observe.run(project=find_agent_v2)
   └─ report
 ```
 
-节点通过固定 `InputSlot` 声明输入,`ctx.declare()` 的返回值直接作为真实模型输入。自研 ReAct
-运行时的完整 LLM 输入/输出、usage、工具调用和消息链会汇总到各节点稳定的 `react` stage。
+Planner 把原始需求压缩成经 Pydantic 校验的 `ExecutionPlan`;Search、Evidence、Evaluator 和
+Report 分别只接收本阶段 Assignment。节点通过单一 `InputSlot` 镜像这份真实输入,Obagent
+`ctx.declare()` 只负责记录,不再覆盖发送给模型的内容。自研 ReAct 运行时的完整 LLM 输入/输出、
+usage、工具调用和消息链会汇总到各节点稳定的 `react` stage。
 `find_agent_v2_run.obagent_run_uid` 保存观测台深链 UID。
 
 配置环境变量:`OBAGENT_ENDPOINT`、`OBAGENT_API_KEY`、`OBAGENT_WAL_DIR`、

+ 20 - 15
find_agent_v2/agent.py

@@ -3,10 +3,9 @@
 from __future__ import annotations
 
 import asyncio
-import json
 from dataclasses import dataclass
 
-from find_agent_v2.context import build_node_slots
+from find_agent_v2.context import assignment_slots, render_assignment
 from find_agent_v2.graph import FindAgentRoundGraph, NodeRunner
 from find_agent_v2.observability import ObagentObserver
 from find_agent_v2.prompts import REPORT_PROMPT
@@ -17,7 +16,12 @@ from find_agent_v2.service import (
     FindAgentV2Service,
     get_find_agent_v2_service,
 )
-from find_agent_v2.state import DiscoverySnapshot, FindAgentResult, FindAgentState
+from find_agent_v2.state import (
+    DiscoverySnapshot,
+    FindAgentResult,
+    FindAgentState,
+    ReportAssignment,
+)
 from find_agent_v2.tools import REPORT_TOOLS
 from supply_agent.config import Settings
 
@@ -45,7 +49,6 @@ def decide_continued_exploration(
     cumulative_rate = (
         current.valid_primary_count / current_evaluated if current_evaluated else 0.0
     )
-
     if new_candidates == 0:
         keep_going = False
         reason = "本轮没有新增候选,搜索前沿已无信息增益"
@@ -109,7 +112,11 @@ class FindAgentV2:
             run = self.service.prepare_resume(run_id)
         if str(run.get("status") or "") != "running":
             raise ValueError(f"run_id={run_id} 当前状态不可执行: {run.get('status')}")
-        state = FindAgentState(run_id=run_id, user_input=user_input)
+        load_plan = getattr(self.service, "get_latest_execution_plan", None)
+        restored_plan = load_plan(run_id) if callable(load_plan) else None
+        state = FindAgentState(
+            run_id=run_id, user_input=user_input, execution_plan=restored_plan,
+        )
         graph = FindAgentRoundGraph(
             service=self.service, runner=self.node_runner, observer=self.observer,
             max_actions=self.max_actions_per_round,
@@ -210,22 +217,20 @@ class FindAgentV2:
         )
         if not failed:
             try:
+                report_assignment = ReportAssignment(
+                    run_id=run_id,
+                    final_state=self.service.get_full_state(run_id),
+                )
                 report = await self.node_runner.run_node(
                     node="report",
                     round_index=state.round_index,
                     system_prompt=REPORT_PROMPT,
-                    user_content=json.dumps(
-                        self.service.get_full_state(run_id),
-                        ensure_ascii=False,
-                        default=str,
-                    ),
+                    user_content=render_assignment(report_assignment),
                     tools=REPORT_TOOLS,
                     max_iterations=4,
-                    slots=build_node_slots(
-                        user_input=state.user_input,
-                        full_state=self.service.get_full_state(run_id),
-                        round_index=state.round_index,
-                        plan=state.plan,
+                    slots=assignment_slots(
+                        report_assignment,
+                        source="FindAgentV2Service.get_full_state(final)",
                     ),
                 )
                 state.node_runs.append(report)

+ 13 - 26
find_agent_v2/context.py

@@ -1,16 +1,15 @@
-"""Context builders backed only by ``find_agent_v2_*`` tables."""
+"""Stage-specific input contracts backed only by ``find_agent_v2_*`` tables."""
 
 from __future__ import annotations
 
-import json
-from typing import Any
+from pydantic import BaseModel
 
 from find_agent_v2.observability import InputSlot
 from find_agent_v2.service import get_find_agent_v2_service
 from find_agent_v2.state import DiscoverySnapshot
 
 
-def load_full_state(run_id: str, *, limit: int = 100) -> dict[str, Any]:
+def load_full_state(run_id: str, *, limit: int = 100) -> dict:
     return get_find_agent_v2_service().get_full_state(run_id, limit=limit)
 
 
@@ -18,38 +17,26 @@ def snapshot_run(run_id: str) -> DiscoverySnapshot:
     return get_find_agent_v2_service().snapshot(run_id)
 
 
-def require_running_run(run_id: str) -> dict[str, Any]:
+def require_running_run(run_id: str) -> dict:
     run = get_find_agent_v2_service().require_run(run_id)
     if str(run.get("status") or "") not in {"running", "finished"}:
         raise ValueError(f"run_id={run_id} 当前状态不可执行: {run.get('status')}")
     return run
 
 
-def render_node_context(
-    *, user_input: str, full_state: dict[str, Any], round_index: int, plan: str = "",
-) -> str:
-    return (
-        f"【当前轮次】\n{round_index}\n\n"
-        f"【原始任务】\n{user_input}\n\n"
-        f"【find_agent_v2 数据库状态快照】\n"
-        f"{json.dumps(full_state, ensure_ascii=False, default=str)}\n\n"
-        f"【本轮搜索计划】\n{plan}"
-    )
+def render_assignment(assignment: BaseModel) -> str:
+    """Serialize the exact object sent to a model and shown in observability."""
+    return assignment.model_dump_json(exclude_none=True)
 
 
-def build_node_slots(
-    *, user_input: str, full_state: dict[str, Any], round_index: int, plan: str = "",
-) -> tuple[InputSlot, ...]:
-    """Stable input structure shared by all node instances and obagent declarations."""
+def assignment_slots(assignment: BaseModel, *, source: str) -> tuple[InputSlot, ...]:
+    """Use one canonical payload so observation can never change model semantics."""
     return (
-        InputSlot("当前轮次", str(round_index), "round_index", "FindAgentV2.begin_round", False),
-        InputSlot("原始任务", user_input, "task", "find_agent_v2_run.input_json", False),
         InputSlot(
-            "数据库状态快照",
-            json.dumps(full_state, ensure_ascii=False, default=str),
-            "state_snapshot",
-            "FindAgentV2Service.get_full_state",
+            "阶段任务",
+            render_assignment(assignment),
+            "stage_assignment",
+            source,
             False,
         ),
-        InputSlot("本轮搜索计划", plan, "round_plan", "supervisor output"),
     )

+ 371 - 146
find_agent_v2/graph.py

@@ -9,12 +9,29 @@ from typing import Any, Protocol
 
 from langgraph.graph import END, START, StateGraph
 
-from find_agent_v2.context import build_node_slots, render_node_context
+from find_agent_v2.context import assignment_slots, render_assignment
 from find_agent_v2.gates import evaluate_candidate_gate
 from find_agent_v2.observability import InputSlot, NullObserver, graph_spec_for
-from find_agent_v2.prompts import EVALUATOR_PROMPT, EVIDENCE_PROMPT, SEARCH_PROMPT, SUPERVISOR_PROMPT
+from find_agent_v2.prompts import (
+    EVALUATOR_PROMPT,
+    EVIDENCE_PROMPT,
+    PLANNER_PROMPT,
+    SEARCH_PROMPT,
+    SUPERVISOR_PROMPT,
+)
 from find_agent_v2.service import FindAgentV2Service
-from find_agent_v2.state import FindAgentGraphState, FindAgentState, NodeRun
+from find_agent_v2.state import (
+    EvidenceAssignment,
+    EvaluationAssignment,
+    FindAgentGraphState,
+    FindAgentState,
+    NodeRun,
+    PlannerAssignment,
+    PlanningDecision,
+    SearchAssignment,
+    SupervisorAssignment,
+    SupervisorDecision,
+)
 from find_agent_v2.tools import (
     EVALUATION_TOOLS,
     EVIDENCE_TOOLS,
@@ -45,8 +62,13 @@ class FindAgentRoundGraph:
     """Supervisor-directed graph with deterministic policy approval."""
 
     def __init__(
-        self, *, service: FindAgentV2Service, runner: NodeRunner, observer=None,
-        max_actions: int = 16, max_search_actions: int = 3,
+        self,
+        *,
+        service: FindAgentV2Service,
+        runner: NodeRunner,
+        observer=None,
+        max_actions: int = 16,
+        max_search_actions: int = 3,
     ) -> None:
         self.service = service
         self.runner = runner
@@ -58,55 +80,84 @@ class FindAgentRoundGraph:
 
     def _full_state(self, state: FindAgentGraphState, *, pending_only: bool = False):
         return self.service.get_full_state(
-            state["run_id"], pending_only=pending_only,
-        )
-
-    def _context(self, state: FindAgentGraphState, *, pending_only: bool = False) -> str:
-        return render_node_context(
-            user_input=state["user_input"],
-            full_state=self._full_state(state, pending_only=pending_only),
-            round_index=state["round_index"],
-            plan=state.get("plan", ""),
-        )
-
-    def _slots(
-        self, state: FindAgentGraphState, *, pending_only: bool = False,
-    ) -> tuple[InputSlot, ...]:
-        return build_node_slots(
-            user_input=state["user_input"],
-            full_state=self._full_state(state, pending_only=pending_only),
-            round_index=state["round_index"],
-            plan=state.get("plan", ""),
+            state["run_id"],
+            pending_only=pending_only,
         )
 
-    def _shard_state(
-        self, state: FindAgentGraphState, items: list[dict[str, Any]],
-    ) -> dict[str, Any]:
-        full_state = self._full_state(state, pending_only=True)
-        return {**full_state, "candidates": items}
-
-    def _shard_context(
-        self, state: FindAgentGraphState, items: list[dict[str, Any]],
-    ) -> str:
-        return render_node_context(
-            user_input=state["user_input"],
-            full_state=self._shard_state(state, items),
-            round_index=state["round_index"],
-            plan=state.get("plan", ""),
-        )
+    @staticmethod
+    def _assignment_input(assignment, *, source: str) -> tuple[str, tuple[InputSlot, ...]]:
+        return render_assignment(assignment), assignment_slots(assignment, source=source)
 
-    def _shard_slots(
-        self, state: FindAgentGraphState, items: list[dict[str, Any]],
-    ) -> tuple[InputSlot, ...]:
-        return build_node_slots(
-            user_input=state["user_input"],
-            full_state=self._shard_state(state, items),
-            round_index=state["round_index"],
-            plan=state.get("plan", ""),
-        )
+    def _supervisor_state(self, state: FindAgentGraphState) -> dict[str, Any]:
+        """Compact progress projection; routing never needs full candidate payloads."""
+        full_state = self._full_state(state)
+        run = full_state.get("run") or {}
+        candidates = list(full_state.get("candidates") or [])
+        pending_candidates = [
+            item
+            for item in candidates
+            if item.get("decision_bucket") == "pending_evaluation"
+        ]
+        return {
+            "run": {
+                key: run.get(key)
+                for key in (
+                    "status", "outcome_status", "current_round", "search_count",
+                    "candidate_count", "valid_primary_count",
+                )
+            },
+            "searches": full_state.get("searches") or [],
+            "candidate_progress": {
+                "pending_count": sum(
+                    1 for _item in pending_candidates
+                ),
+                "primary_count": sum(
+                    item.get("decision_bucket") == "primary" for item in candidates
+                ),
+                "rejected_count": sum(
+                    item.get("decision_bucket") == "rejected" for item in candidates
+                ),
+                "detail_pending_count": sum(
+                    item.get("detail_status") == "pending"
+                    for item in pending_candidates
+                ),
+                "detail_success_count": sum(
+                    item.get("detail_status") == "success"
+                    for item in pending_candidates
+                ),
+                "detail_failed_count": sum(
+                    item.get("detail_status") == "failed"
+                    for item in pending_candidates
+                ),
+                "portrait_pending_count": sum(
+                    item.get("portrait_status") == "pending"
+                    for item in pending_candidates
+                ),
+                "portrait_success_count": sum(
+                    item.get("portrait_status") == "success"
+                    for item in pending_candidates
+                ),
+                "portrait_failed_count": sum(
+                    item.get("portrait_status") == "failed"
+                    for item in pending_candidates
+                ),
+                "evidence_completed_count": sum(
+                    item.get("detail_status") != "pending"
+                    and item.get("portrait_status") != "pending"
+                    for item in pending_candidates
+                ),
+                "evidence_success_count": sum(
+                    item.get("detail_status") == "success"
+                    and item.get("portrait_status") == "success"
+                    for item in pending_candidates
+                ),
+            },
+        }
 
     def _video_understanding_ids(
-        self, state: FindAgentGraphState, items: list[dict[str, Any]],
+        self,
+        state: FindAgentGraphState,
+        items: list[dict[str, Any]],
     ) -> list[int]:
         """Only candidates passing every deterministic hard gate may use video understanding."""
         run = self._full_state(state).get("run") or {}
@@ -116,8 +167,7 @@ class FindAgentRoundGraph:
             gate = evaluate_candidate_gate(item, rules)
             checks = list(gate.get("checks") or [])
             hard_gate_passed = gate.get("status") == "pass" and all(
-                check.get("status") == "pass" and not check.get("compensated")
-                for check in checks
+                check.get("status") == "pass" and not check.get("compensated") for check in checks
             )
             if str(item.get("video_url") or "").strip() and hard_gate_passed:
                 selected.append(int(item["candidate_id"]))
@@ -132,7 +182,7 @@ class FindAgentRoundGraph:
         else:
             start, end = text.find("{"), text.rfind("}")
             if start >= 0 and end > start:
-                text = text[start:end + 1]
+                text = text[start : end + 1]
         try:
             value = json.loads(text)
             return value if isinstance(value, dict) else {}
@@ -140,7 +190,9 @@ class FindAgentRoundGraph:
             return {}
 
     def _approve_action(
-        self, state: FindAgentGraphState, proposal: dict[str, Any],
+        self,
+        state: FindAgentGraphState,
+        proposal: dict[str, Any],
     ) -> tuple[str, str, int, str]:
         """Turn an LLM proposal into a safe, executable transition."""
         pending = self._full_state(state, pending_only=True).get("candidates") or []
@@ -185,38 +237,135 @@ class FindAgentRoundGraph:
         return "finish", reason or "没有待处理候选,结束本轮", worker_count, scope
 
     async def _supervisor(self, state: FindAgentGraphState) -> dict[str, Any]:
-        run = await self.runner.run_node(
-            node="supervisor",
-            round_index=state["round_index"],
-            system_prompt=SUPERVISOR_PROMPT,
-            user_content=self._context(state),
-            tools=(),
-            max_iterations=2,
-            slots=self._slots(state),
-            allow_delegation=False,
-        )
-        proposal = self._parse_supervisor(run.content)
-        action, reason, workers, scope = self._approve_action(state, proposal)
-        proposed_plan = proposal.get("plan")
-        plan = (
-            json.dumps(proposed_plan, ensure_ascii=False)
-            if isinstance(proposed_plan, dict)
-            else state.get("plan", "")
-        )
-        decision = {
-            "step": int(state.get("supervisor_step") or 0) + 1,
-            "proposed_action": proposal.get("next_action"),
-            "approved_action": action,
-            "reason": reason,
-            "worker_count": workers,
-            "evidence_scope": scope,
-        }
+        execution_plan = state.get("execution_plan")
+        observed_decision: dict[str, Any] = {}
+        if execution_plan is None:
+            assignment = PlannerAssignment(
+                run_id=state["run_id"],
+                round_index=state["round_index"],
+                raw_demand=state["user_input"],
+            )
+            user_content, slots = self._assignment_input(
+                assignment,
+                source="find_agent_v2_run.input_json",
+            )
+            structured_runner = getattr(self.runner, "run_planning", None)
+            if callable(structured_runner):
+                run, planning = await structured_runner(
+                    round_index=state["round_index"],
+                    system_prompt=PLANNER_PROMPT,
+                    user_content=user_content,
+                    slots=slots,
+                )
+            else:
+                run = await self.runner.run_node(
+                    node="supervisor",
+                    round_index=state["round_index"],
+                    system_prompt=PLANNER_PROMPT,
+                    user_content=user_content,
+                    tools=(),
+                    max_iterations=2,
+                    slots=slots,
+                    allow_delegation=False,
+                )
+                planning = PlanningDecision.model_validate(self._parse_supervisor(run.content))
+            execution_plan = planning.execution_plan
+            proposal = planning.model_dump(mode="json", exclude={"execution_plan"})
+        else:
+            assignment = SupervisorAssignment(
+                run_id=state["run_id"],
+                round_index=state["round_index"],
+                execution_plan=execution_plan,
+                execution_state=self._supervisor_state(state),
+            )
+            user_content, slots = self._assignment_input(
+                assignment,
+                source="ExecutionPlan + FindAgentV2Service.get_full_state",
+            )
+            structured_runner = getattr(self.runner, "run_supervision", None)
+            def approve_for_observation(supervision: SupervisorDecision) -> dict[str, Any]:
+                proposal_payload = supervision.model_dump(mode="json")
+                approved, approved_reason, approved_workers, approved_scope = (
+                    self._approve_action(state, proposal_payload)
+                )
+                decision_payload = {
+                    "step": int(state.get("supervisor_step") or 0) + 1,
+                    "proposed_action": proposal_payload.get("next_action"),
+                    "approved_action": approved,
+                    "overridden": proposal_payload.get("next_action") != approved,
+                    "reason": approved_reason,
+                    "worker_count": approved_workers,
+                    "evidence_scope": approved_scope,
+                }
+                observed_decision.update(decision_payload)
+                return {"Supervisor决策校验": decision_payload}
+
+            if callable(structured_runner):
+                run, supervision = await structured_runner(
+                    round_index=state["round_index"],
+                    system_prompt=SUPERVISOR_PROMPT,
+                    user_content=user_content,
+                    slots=slots,
+                    output_enricher=approve_for_observation,
+                )
+            else:
+                run = await self.runner.run_node(
+                    node="supervisor",
+                    round_index=state["round_index"],
+                    system_prompt=SUPERVISOR_PROMPT,
+                    user_content=user_content,
+                    tools=(),
+                    max_iterations=2,
+                    slots=slots,
+                    allow_delegation=False,
+                )
+                supervision = SupervisorDecision.model_validate(
+                    self._parse_supervisor(run.content),
+                )
+            proposal = supervision.model_dump(mode="json")
+            if supervision.additional_search_tasks:
+                known = {item.task_id for item in execution_plan.search_tasks}
+                merged = list(execution_plan.search_tasks)
+                for item in supervision.additional_search_tasks:
+                    if item.task_id not in known and len(merged) < 24:
+                        merged.append(item)
+                        known.add(item.task_id)
+                if len(merged) != len(execution_plan.search_tasks):
+                    execution_plan = type(execution_plan).model_validate(
+                        {
+                            **execution_plan.model_dump(),
+                            "search_tasks": merged,
+                        }
+                    )
+        if observed_decision:
+            decision = observed_decision
+            action = str(decision["approved_action"])
+            reason = str(decision["reason"])
+            workers = int(decision["worker_count"])
+            scope = str(decision["evidence_scope"])
+        else:
+            action, reason, workers, scope = self._approve_action(state, proposal)
+            decision = {
+                "step": int(state.get("supervisor_step") or 0) + 1,
+                "proposed_action": proposal.get("next_action"),
+                "approved_action": action,
+                "overridden": proposal.get("next_action") != action,
+                "reason": reason,
+                "worker_count": workers,
+                "evidence_scope": scope,
+            }
         self.service.update_round(
-            state["run_id"], state["round_index"], phase="planning", plan=plan,
+            state["run_id"],
+            state["round_index"],
+            phase="planning",
+            plan=execution_plan.model_dump_json(exclude_none=True),
         )
         return {
-            "phase": "planning", "plan": plan, "approved_action": action,
-            "supervisor_step": decision["step"], "worker_count": workers,
+            "phase": "planning",
+            "execution_plan": execution_plan,
+            "approved_action": action,
+            "supervisor_step": decision["step"],
+            "worker_count": workers,
             "evidence_scope": scope,
             "decision_history": [*state.get("decision_history", []), decision],
             "node_runs": [*state.get("node_runs", []), run],
@@ -224,14 +373,49 @@ class FindAgentRoundGraph:
 
     async def _search(self, state: FindAgentGraphState) -> dict[str, Any]:
         self.service.update_round(state["run_id"], state["round_index"], phase="searching")
+        execution_plan = state.get("execution_plan")
+        if execution_plan is None:
+            raise RuntimeError("Search 缺少已校验的 ExecutionPlan")
+        existing = {
+            (
+                str(item.get("keyword") or ""),
+                str(item.get("provider") or ""),
+                str(item.get("query_reason") or ""),
+            )
+            for item in self._full_state(state).get("searches") or []
+        }
+        tasks = [
+            item
+            for item in execution_plan.search_tasks
+            if (item.keyword, item.provider, item.query_reason) not in existing
+        ][:6]
+        assignment = SearchAssignment(
+            run_id=state["run_id"],
+            round_index=state["round_index"],
+            tasks=tasks,
+        )
+        if not tasks:
+            return {
+                "phase": "searching",
+                "search_actions": int(state.get("search_actions") or 0) + 1,
+                "action_count": int(state.get("action_count") or 0) + 1,
+                "node_runs": [
+                    *state.get("node_runs", []),
+                    NodeRun("search", state["round_index"], "没有未执行的搜索任务", 0, 0),
+                ],
+            }
+        user_content, slots = self._assignment_input(
+            assignment,
+            source="ExecutionPlan.search_tasks",
+        )
         run = await self.runner.run_node(
             node="search",
             round_index=state["round_index"],
             system_prompt=SEARCH_PROMPT,
-            user_content=self._context(state),
+            user_content=user_content,
             tools=SEARCH_TOOLS,
             max_iterations=10,
-            slots=self._slots(state),
+            slots=slots,
         )
         return {
             "phase": "searching",
@@ -251,19 +435,31 @@ class FindAgentRoundGraph:
         elif scope == "portrait":
             detail_items = []
         jobs = [
-            ("detail", detail_items[index:index + 8])
-            for index in range(0, len(detail_items), 8)
+            ("detail", detail_items[index : index + 8]) for index in range(0, len(detail_items), 8)
         ] + [
-            ("portrait", portrait_items[index:index + 8])
+            ("portrait", portrait_items[index : index + 8])
             for index in range(0, len(portrait_items), 8)
         ]
 
         semaphore = asyncio.Semaphore(max(1, min(8, int(state.get("worker_count") or 4))))
 
         async def run_shard(
-            index: int, evidence_type: str, items: list[dict[str, Any]],
+            index: int,
+            evidence_type: str,
+            items: list[dict[str, Any]],
         ) -> NodeRun:
             candidate_ids = [int(item["candidate_id"]) for item in items]
+            assignment = EvidenceAssignment(
+                run_id=state["run_id"],
+                round_index=state["round_index"],
+                candidate_ids=candidate_ids,
+                evidence_type=evidence_type,
+                candidates=items,
+            )
+            user_content, slots = self._assignment_input(
+                assignment,
+                source="host evidence shard",
+            )
             selected = (
                 (EVIDENCE_TOOLS[0], EVIDENCE_TOOLS[2])
                 if evidence_type == "detail"
@@ -272,26 +468,30 @@ class FindAgentRoundGraph:
             async with semaphore:
                 return await self.runner.run_node(
                     node="evidence",
-                round_index=state["round_index"],
-                system_prompt=EVIDENCE_PROMPT,
-                user_content=(
-                    f"你只负责当前 {evidence_type} 分片 candidate_ids={candidate_ids}。"
-                    f"必须为这些候选补齐 {evidence_type},不得访问其他候选。\n\n"
-                    + self._shard_context(state, items)
-                ),
-                tools=bound_candidate_tools(
-                    selected, run_id=state["run_id"], candidate_ids=candidate_ids,
-                ),
-                max_iterations=12,
-                slots=self._shard_slots(state, items),
-                branch_key=f"{evidence_type}-shard-{index}",
-                allow_delegation=False,
+                    round_index=state["round_index"],
+                    system_prompt=EVIDENCE_PROMPT,
+                    user_content=user_content,
+                    tools=bound_candidate_tools(
+                        selected,
+                        run_id=state["run_id"],
+                        candidate_ids=candidate_ids,
+                    ),
+                    max_iterations=12,
+                    slots=slots,
+                    branch_key=f"{evidence_type}-shard-{index}",
+                    allow_delegation=False,
                 )
 
-        runs = await asyncio.gather(*(
-            run_shard(index, evidence_type, items)
-            for index, (evidence_type, items) in enumerate(jobs, start=1)
-        )) if jobs else []
+        runs = (
+            await asyncio.gather(
+                *(
+                    run_shard(index, evidence_type, items)
+                    for index, (evidence_type, items) in enumerate(jobs, start=1)
+                )
+            )
+            if jobs
+            else []
+        )
         return {
             "phase": "evidence",
             "action_count": int(state.get("action_count") or 0) + 1,
@@ -314,25 +514,38 @@ class FindAgentRoundGraph:
             else:
                 gate_failures.append((int(item["candidate_id"]), gate))
         self.service.reject_failed_gates(state["run_id"], gate_failures)
-        shards = [eligible[index:index + 8] for index in range(0, len(eligible), 8)]
+        shards = [eligible[index : index + 8] for index in range(0, len(eligible), 8)]
         semaphore = asyncio.Semaphore(max(1, min(8, int(state.get("worker_count") or 4))))
 
         async def run_shard(index: int, items: list[dict[str, Any]]) -> NodeRun:
             candidate_ids = [int(item["candidate_id"]) for item in items]
             video_understanding_ids = self._video_understanding_ids(state, items)
+            execution_plan = state.get("execution_plan")
+            if execution_plan is None:
+                raise RuntimeError("Evaluator 缺少已校验的 ExecutionPlan")
+            assignment = EvaluationAssignment(
+                run_id=state["run_id"],
+                round_index=state["round_index"],
+                candidate_ids=candidate_ids,
+                evaluation_brief=execution_plan.evaluation_brief,
+                quality_gate_rules=rules,
+                current_datetime=str(rules.get("current_datetime") or ""),
+                timezone=str(rules.get("timezone") or ""),
+                video_understanding_candidate_ids=video_understanding_ids,
+                candidates=items,
+            )
+            user_content, slots = self._assignment_input(
+                assignment,
+                source="ExecutionPlan.evaluation_brief + host candidate shard",
+            )
             video_tools = (
                 bound_candidate_tools(
                     (EVALUATION_TOOLS[0],),
                     run_id=state["run_id"],
                     candidate_ids=video_understanding_ids,
                 )
-                if video_understanding_ids else ()
-            )
-            evaluation_instruction = (
-                f"你只负责当前分片 candidate_ids={candidate_ids}。"
-                f"这些候选均已通过硬门禁;有播放地址、允许按需视频理解的 "
-                f"candidate_ids={video_understanding_ids}。"
-                "视频理解不是硬性要求,可根据已有证据决定是否调用;只允许对这个列表中的候选调用。"
+                if video_understanding_ids
+                else ()
             )
             async with semaphore:
                 structured_runner = getattr(self.runner, "run_evaluation", None)
@@ -340,17 +553,14 @@ class FindAgentRoundGraph:
                     run, proposed = await structured_runner(
                         round_index=state["round_index"],
                         system_prompt=EVALUATOR_PROMPT,
-                        user_content=(
-                            evaluation_instruction
-                            + "必须为分片内每个 candidate_id 各输出一次结构化评估。\n\n"
-                            + self._shard_context(state, items)
-                        ),
-                        slots=self._shard_slots(state, items),
+                        user_content=user_content,
+                        slots=slots,
                         branch_key=f"step-{state.get('supervisor_step', 0)}-shard-{index}",
                         tools=video_tools,
                     )
                     normalized = normalize_evaluation_items(
-                        proposed, allowed_candidates=items,
+                        proposed,
+                        allowed_candidates=items,
                     )
                     updated = self.service.evaluate(state["run_id"], normalized)
                     updated_ids = {int(item["candidate_id"]) for item in updated}
@@ -364,11 +574,7 @@ class FindAgentRoundGraph:
                     node="evaluator",
                     round_index=state["round_index"],
                     system_prompt=EVALUATOR_PROMPT,
-                    user_content=(
-                        evaluation_instruction
-                        + "必须把这些候选全部分池,不得评估其他候选。\n\n"
-                        + self._shard_context(state, items)
-                    ),
+                    user_content=user_content,
                     tools=(
                         *video_tools,
                         *bound_candidate_tools(
@@ -378,14 +584,18 @@ class FindAgentRoundGraph:
                         ),
                     ),
                     max_iterations=12,
-                    slots=self._shard_slots(state, items),
+                    slots=slots,
                     branch_key=f"step-{state.get('supervisor_step', 0)}-shard-{index}",
                     allow_delegation=False,
                 )
 
-        runs = await asyncio.gather(*(
-            run_shard(index, items) for index, items in enumerate(shards, start=1)
-        )) if shards else []
+        runs = (
+            await asyncio.gather(
+                *(run_shard(index, items) for index, items in enumerate(shards, start=1))
+            )
+            if shards
+            else []
+        )
         node_runs.extend(runs)
         self.service.recount_valid_primary(state["run_id"])
         after = self.service.snapshot(state["run_id"])
@@ -393,11 +603,11 @@ class FindAgentRoundGraph:
         stagnant = stagnant + 1 if after.pending_count >= before.pending_count else 0
         if stagnant >= 3:
             raise RuntimeError(
-                "评估节点连续 3 次未消费 pending_evaluation 候选:"
-                f"remaining={after.pending_count}"
+                f"评估节点连续 3 次未消费 pending_evaluation 候选:remaining={after.pending_count}"
             )
         return {
-            "phase": "evaluating", "node_runs": node_runs,
+            "phase": "evaluating",
+            "node_runs": node_runs,
             "action_count": int(state.get("action_count") or 0) + 1,
             "evaluator_stagnation": stagnant,
         }
@@ -413,10 +623,16 @@ class FindAgentRoundGraph:
         builder.add_node("evidence", self._evidence)
         builder.add_node("evaluator", self._evaluator)
         builder.add_edge(START, "supervisor")
-        builder.add_conditional_edges("supervisor", self._route, {
-            "search": "search", "evidence": "evidence",
-            "evaluator": "evaluator", "finish": END,
-        })
+        builder.add_conditional_edges(
+            "supervisor",
+            self._route,
+            {
+                "search": "search",
+                "evidence": "evidence",
+                "evaluator": "evaluator",
+                "finish": END,
+            },
+        )
         builder.add_edge("search", "supervisor")
         builder.add_edge("evidence", "supervisor")
         builder.add_edge("evaluator", "supervisor")
@@ -427,7 +643,7 @@ class FindAgentRoundGraph:
             "run_id": state.run_id,
             "user_input": state.user_input,
             "round_index": state.round_index,
-            "plan": state.plan,
+            "execution_plan": state.execution_plan,
             "phase": state.phase,
             "node_runs": [],
             "snapshot": state.snapshot,
@@ -440,12 +656,14 @@ class FindAgentRoundGraph:
             "evaluator_stagnation": 0,
         }
         with self.observer.round(
-            round_index=state.round_index, spec=self.obagent_spec,
+            round_index=state.round_index,
+            spec=self.obagent_spec,
         ) as round_observation:
             output = await self.app.ainvoke(
-                graph_state, config={"recursion_limit": max(64, self.max_actions * 4)},
+                graph_state,
+                config={"recursion_limit": max(64, self.max_actions * 4)},
             )
-            state.plan = str(output.get("plan") or "")
+            state.execution_plan = output.get("execution_plan")
             state.node_runs.extend(output.get("node_runs") or [])
             state.snapshot = self.service.snapshot(state.run_id)
             state.phase = "done"
@@ -456,9 +674,16 @@ class FindAgentRoundGraph:
                 status="done",
                 snapshot=state.snapshot,
             )
-            round_observation.set_output({
-                "状态快照": state.snapshot.__dict__,
-                "本轮计划": state.plan,
-                "Supervisor决策轨迹": output.get("decision_history") or [],
-            }, ok=True)
+            round_observation.set_output(
+                {
+                    "状态快照": state.snapshot.__dict__,
+                    "本轮计划": (
+                        state.execution_plan.model_dump(mode="json")
+                        if state.execution_plan is not None
+                        else {}
+                    ),
+                    "Supervisor决策轨迹": output.get("decision_history") or [],
+                },
+                ok=True,
+            )
         return state

+ 6 - 3
find_agent_v2/observability.py

@@ -160,13 +160,17 @@ class _ModuleHandle:
         model: str,
         refs: dict[str, str] | None = None,
     ) -> str:
-        del fallback, refs
-        return self.ctx.declare(
+        del refs
+        # Observation is a mirror, not an input transformer. Returning the SDK's
+        # rendered blocks here used to replace shard/delegate instructions in
+        # production while NullObserver preserved them in tests.
+        self.ctx.declare(
             system_prompt=system_prompt,
             blocks=_blocks(slots),
             tools=list(tools),
             model=model,
         )
+        return fallback
 
     def record_react(self, *, output: dict[str, Any], ok: bool) -> None:
         # One stable code stage per agent module; the payload holds the complete
@@ -233,7 +237,6 @@ class ObagentObserver:
                 source="FindAgentV2.begin_round", optional=False,
             )])
             yield _ModuleHandle(ctx)
-
     @contextmanager
     def node(self, *, node: str, branch_key: str = ""):
         from obagent_sdk import observe

+ 26 - 14
find_agent_v2/prompts.py

@@ -28,37 +28,48 @@ COMMON_RULES = """
 - 每个模块只能调用当前真实提供的工具;未提供的能力视为物理不可用。
 - 优先批量调用工具,避免同一轮重复查询或重复补证。
 """
+PLANNER_PROMPT = COMMON_RULES + """
+
+# 当前模块:Planner
+
+你是唯一可以读取原始需求的模块。把原始需求、参考视频和点位压缩为结构化 ExecutionPlan,供后续
+模块执行。DemandBrief 和 EvaluationBrief 必须保留判断相关性、人群适配、排除项和时间适用性所需
+的业务语义,但不要复制原文或整段参考视频。生成 1~6 个互补且可以独立执行的 SearchTask,每个
+task_id 唯一。首次动作通常是 search。只输出符合宿主结构化 schema 的结果。
+"""
+
 
 SUPERVISOR_PROMPT = COMMON_RULES + """
 
 # 当前模块:Supervisor
 
-你负责根据真实数据库状态决定下一步,但不直接调用业务工具。每次只能提议一个动作:
+你负责根据结构化 ExecutionPlan 和执行状态决定下一步,但不读取原始需求,也不直接调用业务工具。
+每次只能提议一个动作:
 
-- `search`:继续扩大候选池,并在 `plan.searches` 中给出最多 3 个互补搜索方向。
+- `search`:继续扩大候选池;需要新方向时通过 `additional_search_tasks` 给出结构化任务
 - `evidence`:补证,可用 `evidence_scope=detail|portrait|both` 决定顺序。
 - `evaluator`:证据已处理后评分分池。
 - `finish`:没有待处理候选且继续搜索已无信息增益时结束本轮。
 
+`candidate_progress` 中的证据数量只统计仍处于 `pending_evaluation` 的候选。仅当
+`detail_pending_count` 或 `portrait_pending_count` 大于 0 时才能提议 `evidence`;当两者都为 0 且
+`pending_count` 大于 0 时,必须提议 `evaluator`,不得再次补证。`failed_count` 表示本轮证据调用已
+结束但失败,不属于待补证项,不得据此重复调用证据;成功数量以 `detail_success_count`、
+`portrait_success_count` 和 `evidence_success_count` 为准。
+
 是否继续探索应综合判断候选新增数量、已评估样本量、本轮通过率和累计通过率:样本仍少或搜索
 仍持续产生有效候选时继续;新增枯竭,或样本已充分但通过率持续为零或很低时结束。
 
 你还可以用 `worker_count` 建议 1~8 个并行 Worker。宿主会依据真实状态、工具权限和预算审核提议;
-非法跳转会被改写为安全动作。只输出 JSON,不要 Markdown:
-{"next_action":"search|evidence|evaluator|finish","reason":"...","worker_count":4,
- "evidence_scope":"both","plan":{"intent_summary":"...","searches":[{"keyword":"...",
- "query_reason":"...","source_type":"demand|seed|point|mixed","provider":"internal_keyword|tikhub",
- "max_pages":1}]}}
+非法跳转会被改写为安全动作。只输出符合宿主结构化 schema 的结果。
 """
 
-# Backward-compatible import name for callers that customized the old prompt.
-PLANNER_PROMPT = SUPERVISOR_PROMPT
-
 SEARCH_PROMPT = COMMON_RULES + """
 
 # 当前模块:候选搜索
 
-执行本轮搜索计划。只能使用 `search_videos_v2` 和 `query_find_agent_v2_state`;不得获取详情或
+输入只包含当前 SearchAssignment。严格执行其中的 tasks,不得自行重建或扩写原始需求。只能使用
+`search_videos_v2` 和 `query_find_agent_v2_state`;不得获取详情或
 画像,不得评分或分池。优先把互不依赖的关键词放入一次批量搜索;需要隔离上下文时可用
 `delegate_agents_v2` 并发委派独立搜索任务,完成后查询状态确认写入结果。
 """
@@ -67,7 +78,8 @@ EVIDENCE_PROMPT = COMMON_RULES + """
 
 # 当前模块:证据补全
 
-宿主会给你一个互斥的候选分片,并且只注册当前分片实际缺失的详情或画像工具。使用
+输入只包含当前 EvidenceAssignment,不包含也不需要原始需求。宿主会给你一个互斥的候选分片,
+并且只注册当前分片实际缺失的详情或画像工具。使用
 `query_pending_candidates_v2` 只读当前分片,并把分片内所有候选的当前证据补齐。不得继续搜索、
 评分、分池或访问分片外候选。上游失败也必须保留为明确证据状态。
 """
@@ -78,7 +90,7 @@ EVALUATOR_PROMPT = COMMON_RULES + """
 
 ## 你拥有的能力
 
-- 宿主会提供候选详情、互动数据、内容侧/账号侧年龄画像和用户需求上下文
+- 宿主会提供候选详情、互动数据、内容侧/账号侧年龄画像和结构化 EvaluationBrief
 - 对宿主列入 `允许按需视频理解的 candidate_ids` 的候选,你可以调用
   `understand_candidate_video_30s_v2` 查看视频前 30 秒的画面、字幕、对白、主题与观点。
 - 调用时,prompt 应包含当前需求及需要核验的疑点,优先核验标题和元数据无法回答的问题;同一候选
@@ -124,6 +136,6 @@ REPORT_PROMPT = COMMON_RULES + """
 
 # 当前模块:结果报告
 
-最终状态已由宿主写入。查询数据库并报告业务结果、有效 primary 数量、每条 primary 的标题和
+输入只包含最终 ReportAssignment。最终状态已由宿主写入。报告业务结果、有效 primary 数量、每条 primary 的标题和
 关键证据,以及主要淘汰原因。不得再搜索、补证、评分或修改数据。
 """

+ 1 - 1
find_agent_v2/qwen_video_understanding_30s.py

@@ -38,7 +38,7 @@ DEFAULT_PROMPT = (
 )
 CLIP_SECONDS = 30.0
 MIN_USABLE_CLIP_SECONDS = 15.0
-REMOTE_CLIP_TIMEOUT_SECONDS = 60.0
+REMOTE_CLIP_TIMEOUT_SECONDS = 120.0
 REMOTE_CLIP_ATTEMPTS = 2
 UPLOAD_TIMEOUT_SECONDS = 60.0
 MODEL_TIMEOUT_SECONDS = 300.0

+ 102 - 5
find_agent_v2/runtime.py

@@ -9,7 +9,7 @@ from __future__ import annotations
 import asyncio
 import json
 import os
-from collections.abc import Iterable
+from collections.abc import Callable, Iterable
 from typing import Any
 
 from langchain.agents import create_agent
@@ -20,7 +20,7 @@ from langchain_openai import ChatOpenAI
 from pydantic import BaseModel, Field
 
 from find_agent_v2.observability import InputSlot, ObagentObserver
-from find_agent_v2.state import NodeRun
+from find_agent_v2.state import NodeRun, PlanningDecision, SupervisorDecision
 from find_agent_v2.tools import ToolFn
 from find_agent_v2.tools import EvaluationBatch
 from supply_agent.config import Settings, get_settings
@@ -163,7 +163,6 @@ class FindAgentNodeHost:
         system_prompt: str,
         tools: tuple[ToolFn, ...],
         max_iterations: int,
-        slots: tuple[InputSlot, ...],
     ) -> StructuredTool:
         async def delegate_agents_v2(requests: list[DelegateRequest]) -> str:
             """并发委派多个相互独立的任务给同阶段子 Agent。"""
@@ -180,7 +179,10 @@ class FindAgentNodeHost:
                     user_content=request.task,
                     tools=tools,
                     max_iterations=max_iterations,
-                    slots=slots,
+                    slots=(InputSlot(
+                        "委派任务", request.task, "delegate_assignment",
+                        "delegate_agents_v2.requests", False,
+                    ),),
                     branch_key=f"delegate-{index}",
                     allow_delegation=False,
                 )
@@ -232,7 +234,6 @@ class FindAgentNodeHost:
                 system_prompt=system_prompt,
                 tools=tool_functions,
                 max_iterations=max_iterations,
-                slots=slots,
             ))
         agent = create_agent(
             model=self._model(node),
@@ -349,6 +350,102 @@ class FindAgentNodeHost:
             return run, items
 
 
+    async def _run_supervisor_structured(
+        self,
+        *,
+        round_index: int,
+        system_prompt: str,
+        user_content: str,
+        slots: tuple[InputSlot, ...],
+        response_format: type[BaseModel],
+        output_enricher: Callable[[BaseModel], dict[str, Any]] | None = None,
+    ) -> tuple[NodeRun, BaseModel]:
+        """Run planner/supervisor with a validated host-owned output contract."""
+        node = "supervisor"
+        model_name = self.models_by_role.get(node, self.default_model)
+        agent = create_agent(
+            model=self._model(node),
+            tools=[],
+            system_prompt=system_prompt,
+            middleware=self._middleware(node),
+            response_format=response_format,
+            name="find_agent_v2_supervisor",
+        )
+        with self.observer.node(node=node) as observation:
+            actual_user_content = observation.declare(
+                fallback=user_content,
+                system_prompt=system_prompt,
+                slots=slots,
+                tools=(),
+                model=model_name,
+            )
+            result = await agent.ainvoke(
+                {"messages": [{"role": "user", "content": actual_user_content}]},
+                config={"recursion_limit": 24},
+            )
+            messages: list[BaseMessage] = list(result.get("messages") or [])
+            usage = _usage(messages)
+            for key in ("input_tokens", "output_tokens", "total_tokens"):
+                self.usage[key] = int(self.usage[key]) + int(usage[key])
+            self.usage["cost"] = round(float(self.usage["cost"]) + float(usage["cost"]), 8)
+            raw_structured = result.get("structured_response")
+            structured = (
+                raw_structured
+                if isinstance(raw_structured, response_format)
+                else response_format.model_validate(raw_structured)
+            )
+            structured_payload = structured.model_dump(mode="json", exclude_none=True)
+            enriched_output = output_enricher(structured) if output_enricher else {}
+            process = {
+                "messages": [_message_dict(message) for message in messages],
+                "events": _events(messages),
+                "usage": usage,
+                "iterations": sum(isinstance(message, AIMessage) for message in messages),
+                "tool_calls_made": 0,
+                "structured_output": structured_payload,
+                **enriched_output,
+            }
+            content = json.dumps(structured_payload, ensure_ascii=False)
+            observation.record_react(output=process, ok=True)
+            observation.set_output({"结构化输出": structured_payload, **process}, ok=True)
+            return NodeRun(node, round_index, content, int(process["iterations"]), 0), structured
+
+    async def run_planning(
+        self,
+        *,
+        round_index: int,
+        system_prompt: str,
+        user_content: str,
+        slots: tuple[InputSlot, ...] = (),
+    ) -> tuple[NodeRun, PlanningDecision]:
+        run, structured = await self._run_supervisor_structured(
+            round_index=round_index,
+            system_prompt=system_prompt,
+            user_content=user_content,
+            slots=slots,
+            response_format=PlanningDecision,
+        )
+        return run, PlanningDecision.model_validate(structured)
+
+    async def run_supervision(
+        self,
+        *,
+        round_index: int,
+        system_prompt: str,
+        user_content: str,
+        slots: tuple[InputSlot, ...] = (),
+        output_enricher: Callable[[SupervisorDecision], dict[str, Any]] | None = None,
+    ) -> tuple[NodeRun, SupervisorDecision]:
+        run, structured = await self._run_supervisor_structured(
+            round_index=round_index,
+            system_prompt=system_prompt,
+            user_content=user_content,
+            slots=slots,
+            response_format=SupervisorDecision,
+            output_enricher=output_enricher,
+        )
+        return run, SupervisorDecision.model_validate(structured)
+
 def normalize_models(
     *,
     model: str | None = None,

+ 25 - 7
find_agent_v2/service.py

@@ -17,7 +17,7 @@ from find_agent_v2.models import (
     FindAgentV2Run,
     FindAgentV2Search,
 )
-from find_agent_v2.state import DiscoverySnapshot
+from find_agent_v2.state import DiscoverySnapshot, ExecutionPlan
 from supply_infra.db.session import get_session
 from find_agent_v2.gates import (
     build_rule_snapshot,
@@ -60,12 +60,9 @@ def _ratio(value: Any) -> Decimal | None:
 def _fails_search_share_gate(
     provider: str, share_count: Any, min_share_count: int,
 ) -> bool:
-    """Only TikHub has a verified search-time share metric contract."""
-    return (
-        provider == "tikhub"
-        and share_count is not None
-        and int(share_count) < int(min_share_count)
-    )
+    """Reject any provider's explicit share metric below the configured threshold."""
+    del provider
+    return share_count is not None and int(share_count) < int(min_share_count)
 
 
 def _candidate_dict(row: FindAgentV2Candidate) -> dict[str, Any]:
@@ -229,6 +226,27 @@ class FindAgentV2Service:
                 raise ValueError(f"run_id={run_id} 缺少有效 user_input")
             return user_input
 
+    def get_latest_execution_plan(self, run_id: str) -> ExecutionPlan | None:
+        """Restore the newest validated plan for continuation/resume without replanning raw input."""
+        with get_session() as session:
+            raw = session.scalar(
+                select(FindAgentV2Round.plan_json)
+                .where(
+                    FindAgentV2Round.run_id == run_id,
+                    FindAgentV2Round.plan_json.is_not(None),
+                )
+                .order_by(FindAgentV2Round.round_index.desc(), FindAgentV2Round.id.desc())
+                .limit(1)
+            )
+        payload = _loads(raw, None)
+        if not isinstance(payload, dict):
+            return None
+        try:
+            return ExecutionPlan.model_validate(payload)
+        except ValueError:
+            # Old rounds may contain the pre-contract free-form plan shape.
+            return None
+
     def begin_round(self, run_id: str, round_index: int, snapshot: DiscoverySnapshot) -> None:
         with get_session() as session:
             run = session.scalar(select(FindAgentV2Run).where(FindAgentV2Run.run_id == run_id))

+ 112 - 2
find_agent_v2/state.py

@@ -9,6 +9,8 @@ from __future__ import annotations
 from dataclasses import dataclass, field
 from typing import Any, Literal, TypedDict
 
+from pydantic import BaseModel, Field, model_validator
+
 from supply_agent.types import AgentResult
 
 Phase = Literal["planning", "searching", "evidence", "evaluating", "done"]
@@ -16,6 +18,114 @@ SupervisorAction = Literal["search", "evidence", "evaluator", "finish"]
 EndKind = Literal["goal_met", "partial", "no_match", "failed", "stopped"]
 
 
+class DemandBrief(BaseModel):
+    """Planner-owned compact interpretation of the immutable raw demand."""
+
+    core_intent: str = Field(min_length=1)
+    target_audience: list[str] = Field(default_factory=list)
+    forwarding_audience: list[str] = Field(default_factory=list)
+    key_scenarios: list[str] = Field(default_factory=list)
+    reference_signals: list[str] = Field(default_factory=list)
+    relevance_criteria: list[str] = Field(min_length=1)
+    exclusion_criteria: list[str] = Field(default_factory=list)
+    temporal_requirements: list[str] = Field(default_factory=list)
+
+
+class SearchTask(BaseModel):
+    """One independently executable search task produced by planning."""
+
+    task_id: str = Field(min_length=1)
+    keyword: str = Field(min_length=1)
+    query_reason: str = Field(min_length=1)
+    source_type: Literal["demand", "seed", "point", "mixed"] = "mixed"
+    provider: Literal["internal_keyword", "tikhub"] = "internal_keyword"
+    max_pages: int = Field(default=1, ge=1, le=2)
+    coverage_targets: list[str] = Field(default_factory=list)
+
+
+class EvaluationBrief(BaseModel):
+    """Demand semantics required by evaluators, without the original request."""
+
+    relevance_criteria: list[str] = Field(min_length=1)
+    audience_criteria: list[str] = Field(default_factory=list)
+    temporal_requirements: list[str] = Field(default_factory=list)
+    exclusion_criteria: list[str] = Field(default_factory=list)
+
+
+class ExecutionPlan(BaseModel):
+    """Validated contract passed from the planner to later stages."""
+
+    schema_version: Literal["1.0"] = "1.0"
+    demand_brief: DemandBrief
+    search_tasks: list[SearchTask] = Field(min_length=1, max_length=24)
+    evaluation_brief: EvaluationBrief
+
+    @model_validator(mode="after")
+    def validate_unique_task_ids(self):
+        task_ids = [item.task_id for item in self.search_tasks]
+        if len(task_ids) != len(set(task_ids)):
+            raise ValueError("search_tasks.task_id 不能重复")
+        return self
+
+
+class SupervisorDecision(BaseModel):
+    """One routing decision; subsequent decisions may add search tasks."""
+
+    next_action: SupervisorAction
+    reason: str = ""
+    worker_count: int = Field(default=4, ge=1, le=8)
+    evidence_scope: Literal["detail", "portrait", "both"] = "both"
+    additional_search_tasks: list[SearchTask] = Field(default_factory=list, max_length=6)
+
+
+class PlanningDecision(SupervisorDecision):
+    execution_plan: ExecutionPlan
+
+
+class PlannerAssignment(BaseModel):
+    run_id: str
+    round_index: int
+    raw_demand: str
+
+
+class SupervisorAssignment(BaseModel):
+    run_id: str
+    round_index: int
+    execution_plan: ExecutionPlan
+    execution_state: dict[str, Any]
+
+
+class SearchAssignment(BaseModel):
+    run_id: str
+    round_index: int
+    tasks: list[SearchTask]
+
+
+class EvidenceAssignment(BaseModel):
+    run_id: str
+    round_index: int
+    candidate_ids: list[int]
+    evidence_type: Literal["detail", "portrait"]
+    candidates: list[dict[str, Any]]
+
+
+class EvaluationAssignment(BaseModel):
+    run_id: str
+    round_index: int
+    candidate_ids: list[int]
+    evaluation_brief: EvaluationBrief
+    quality_gate_rules: dict[str, Any]
+    current_datetime: str = ""
+    timezone: str = ""
+    video_understanding_candidate_ids: list[int] = Field(default_factory=list)
+    candidates: list[dict[str, Any]]
+
+
+class ReportAssignment(BaseModel):
+    run_id: str
+    final_state: dict[str, Any]
+
+
 @dataclass(frozen=True)
 class DiscoverySnapshot:
     """Small deterministic projection of the database state."""
@@ -64,7 +174,7 @@ class FindAgentState:
     user_input: str
     round_index: int = 0
     phase: Phase = "planning"
-    plan: str = ""
+    execution_plan: ExecutionPlan | None = None
     stop: bool = False
     stop_reason: str = ""
     previous_snapshot: DiscoverySnapshot | None = None
@@ -79,7 +189,7 @@ class FindAgentGraphState(TypedDict, total=False):
     run_id: str
     user_input: str
     round_index: int
-    plan: str
+    execution_plan: ExecutionPlan | None
     phase: Phase
     node_runs: list[NodeRun]
     snapshot: DiscoverySnapshot | None

+ 206 - 5
tests/supply_agent/test_find_agent_v2.py

@@ -35,10 +35,22 @@ from find_agent_v2.observability import (
     OBAGENT_PROJECT,
     OBAGENT_ROUND_ANCHOR,
     NullObserver,
+    InputSlot,
+    _ModuleHandle,
 )
 from find_agent_v2.prompts import COMMON_RULES, EVALUATOR_PROMPT
 from find_agent_v2.providers import _normalize_search_item, normalize_age_pair
-from find_agent_v2.state import DiscoverySnapshot, FindAgentState, NodeRun
+from find_agent_v2.state import (
+    DemandBrief,
+    DiscoverySnapshot,
+    EvaluationBrief,
+    ExecutionPlan,
+    FindAgentState,
+    NodeRun,
+    PlanningDecision,
+    SearchTask,
+    SupervisorDecision,
+)
 from find_agent_v2.service import _fails_search_share_gate
 from find_agent_v2.tools import (
     CandidateEvaluation,
@@ -175,12 +187,15 @@ def test_search_metrics_preserve_missing_share_count_and_real_zero() -> None:
         ("tikhub", 999, True),
         ("tikhub", 1000, False),
         ("tikhub", None, False),
-        ("internal_keyword", 0, False),
-        ("internal_keyword", 999, False),
+        ("internal_keyword", 0, True),
+        ("internal_keyword", 999, True),
+        ("internal_keyword", 1000, False),
         ("internal_keyword", None, False),
     ],
 )
-def test_search_share_gate_only_applies_to_tikhub(provider, share_count, expected) -> None:
+def test_search_share_gate_applies_to_explicit_metrics_from_all_providers(
+    provider, share_count, expected,
+) -> None:
     assert _fails_search_share_gate(provider, share_count, 1000) is expected
 
 
@@ -550,9 +565,12 @@ class _FakeRunner:
     def __init__(self, service: _FakeService) -> None:
         self.service = service
         self.calls: list[tuple[str, set[str]]] = []
+        self.inputs: list[tuple[str, str]] = []
+        self.supervisor_visual_outputs: list[dict] = []
 
-    async def run_node(self, *, node, round_index, tools=(), **_kwargs) -> NodeRun:
+    async def run_node(self, *, node, round_index, tools=(), user_content="", **_kwargs) -> NodeRun:
         self.calls.append((node, _names(tools)))
+        self.inputs.append((node, user_content))
         if node == "search":
             self.service.stage = "searched"
         elif node == "evidence":
@@ -578,6 +596,138 @@ class _FakeRunner:
             content = '{"searches": []}'
         return NodeRun(node, round_index, content, 1, 0)
 
+    @staticmethod
+    def execution_plan() -> ExecutionPlan:
+        return ExecutionPlan(
+            demand_brief=DemandBrief(
+                core_intent="测试需求",
+                relevance_criteria=["内容直接匹配测试需求"],
+            ),
+            search_tasks=[SearchTask(
+                task_id="search-1", keyword="测试需求", query_reason="建立候选池",
+            )],
+            evaluation_brief=EvaluationBrief(
+                relevance_criteria=["内容直接匹配测试需求"],
+            ),
+        )
+
+    async def run_planning(self, *, round_index, **kwargs):
+        run = await self.run_node(node="supervisor", round_index=round_index, **kwargs)
+        proposal = json.loads(run.content)
+        return run, PlanningDecision(
+            next_action=proposal["next_action"],
+            worker_count=proposal["worker_count"],
+            evidence_scope=proposal["evidence_scope"],
+            execution_plan=self.execution_plan(),
+        )
+
+    async def run_supervision(self, *, round_index, **kwargs):
+        output_enricher = kwargs.pop("output_enricher", None)
+        run = await self.run_node(node="supervisor", round_index=round_index, **kwargs)
+        proposal = json.loads(run.content)
+        decision = SupervisorDecision(
+            next_action=proposal["next_action"],
+            worker_count=proposal["worker_count"],
+            evidence_scope=proposal["evidence_scope"],
+        )
+        if output_enricher is not None:
+            self.supervisor_visual_outputs.append(output_enricher(decision))
+        return run, decision
+
+
+class _WrongEvidenceProposalRunner(_FakeRunner):
+    async def run_node(self, *, node, round_index, tools=(), user_content="", **kwargs):
+        run = await super().run_node(
+            node=node,
+            round_index=round_index,
+            tools=tools,
+            user_content=user_content,
+            **kwargs,
+        )
+        if node == "supervisor" and self.service.stage == "evidenced":
+            return NodeRun(
+                node,
+                round_index,
+                json.dumps({
+                    "next_action": "evidence",
+                    "worker_count": 8,
+                    "evidence_scope": "both",
+                }),
+                1,
+                0,
+            )
+        return run
+
+
+def test_supervisor_progress_counts_missing_evidence_only_for_pending_candidates() -> None:
+    class ProgressService:
+        @staticmethod
+        def get_full_state(_run_id: str, **_kwargs):
+            return {
+                "run": {},
+                "searches": [],
+                "candidates": [
+                    {
+                        "candidate_id": 1,
+                        "decision_bucket": "pending_evaluation",
+                        "detail_status": "success",
+                        "portrait_status": "pending",
+                    },
+                    {
+                        "candidate_id": 2,
+                        "decision_bucket": "rejected",
+                        "detail_status": "pending",
+                        "portrait_status": "pending",
+                    },
+                    {
+                        "candidate_id": 3,
+                        "decision_bucket": "pending_evaluation",
+                        "detail_status": "failed",
+                        "portrait_status": "success",
+                    },
+                    {
+                        "candidate_id": 4,
+                        "decision_bucket": "pending_evaluation",
+                        "detail_status": "success",
+                        "portrait_status": "success",
+                    },
+                ],
+            }
+
+    graph = FindAgentRoundGraph(service=ProgressService(), runner=object())
+    progress = graph._supervisor_state({"run_id": "run"})["candidate_progress"]
+
+    assert progress["pending_count"] == 3
+    assert progress["detail_pending_count"] == 0
+    assert progress["detail_success_count"] == 2
+    assert progress["detail_failed_count"] == 1
+    assert progress["portrait_pending_count"] == 1
+    assert progress["portrait_success_count"] == 2
+    assert progress["portrait_failed_count"] == 0
+    assert progress["evidence_completed_count"] == 2
+    assert progress["evidence_success_count"] == 1
+
+
+@pytest.mark.asyncio
+async def test_supervisor_visualization_records_host_override() -> None:
+    service = _FakeService(pending_after_search=1)
+    runner = _WrongEvidenceProposalRunner(service)
+    graph = FindAgentRoundGraph(service=service, runner=runner)
+
+    await graph.invoke(FindAgentState(run_id="override-run", user_input="task", round_index=1))
+
+    validations = [
+        item["Supervisor决策校验"] for item in runner.supervisor_visual_outputs
+    ]
+    overridden = next(
+        item
+        for item in validations
+        if item["proposed_action"] == "evidence"
+        and item["approved_action"] == "evaluator"
+    )
+    assert overridden["proposed_action"] == "evidence"
+    assert overridden["approved_action"] == "evaluator"
+
 
 @pytest.mark.asyncio
 async def test_round_graph_supervisor_routes_with_guarded_allowlists() -> None:
@@ -603,6 +753,57 @@ async def test_round_graph_supervisor_routes_with_guarded_allowlists() -> None:
     assert result.phase == "done"
     assert result.snapshot is not None and result.snapshot.pending_count == 0
 
+@pytest.mark.asyncio
+async def test_raw_demand_is_only_sent_to_initial_planner() -> None:
+    service = _FakeService(pending_after_search=1)
+    runner = _FakeRunner(service)
+    graph = FindAgentRoundGraph(service=service, runner=runner)
+    marker = "RAW-DEMAND-MUST-NOT-LEAK"
+
+    await graph.invoke(FindAgentState(
+        run_id="assignment-run", user_input=marker, round_index=1,
+    ))
+
+    assert marker in runner.inputs[0][1]
+    assert all(marker not in content for _node, content in runner.inputs[1:])
+    search_payload = json.loads(next(
+        content for node, content in runner.inputs if node == "search"
+    ))
+    evidence_payload = json.loads(next(
+        content for node, content in runner.inputs if node == "evidence"
+    ))
+    supervisor_payload = json.loads([
+        content for node, content in runner.inputs if node == "supervisor"
+    ][1])
+    assert set(search_payload) == {"run_id", "round_index", "tasks"}
+    assert set(evidence_payload) == {
+        "run_id", "round_index", "candidate_ids", "evidence_type", "candidates",
+    }
+    assert "candidates" not in supervisor_payload["execution_state"]
+    assert "candidate_progress" in supervisor_payload["execution_state"]
+
+
+def test_obagent_declaration_cannot_override_actual_model_input() -> None:
+    class Context:
+        declared = False
+
+        def declare(self, **_kwargs):
+            self.declared = True
+            return "SDK-rendered-input-that-must-not-reach-model"
+
+    context = Context()
+    handle = _ModuleHandle(context)
+    actual = handle.declare(
+        fallback='{"run_id":"r","tasks":[]}',
+        system_prompt="prompt",
+        slots=(InputSlot("阶段任务", "{}", "stage_assignment", "test", False),),
+        tools=(),
+        model="model",
+    )
+
+    assert context.declared is True
+    assert actual == '{"run_id":"r","tasks":[]}'
+
 
 @pytest.mark.asyncio
 async def test_round_graph_skips_evidence_and_evaluation_without_candidates() -> None:

+ 1 - 1
tests/supply_agent/test_find_agent_v2_video_understanding.py

@@ -27,7 +27,7 @@ async def test_remote_clip_retries_once_and_prefers_complete_clip(monkeypatch, t
     assert path.name == "clipped_attempt_2.mp4"
     assert duration == 30.0
     assert complete is True
-    assert video_tool.REMOTE_CLIP_TIMEOUT_SECONDS == 60.0
+    assert video_tool.REMOTE_CLIP_TIMEOUT_SECONDS == 120.0
 
 
 @pytest.mark.asyncio

+ 14 - 1
web/src/views/FindAgentV2View.vue

@@ -40,9 +40,22 @@ const obagentUrl = computed(() => run.value?.obagent_run_uid
   ? `http://ob.aiddit.com/#client_uid=${run.value.obagent_run_uid}` : '')
 const filteredCandidates = computed(() => {
   const items = candidates.value
-  return candidateFilter.value === 'all'
+  const filtered = candidateFilter.value === 'all'
     ? items
     : items.filter((item) => item.decision_bucket === candidateFilter.value)
+  const bucketOrder: Record<string, number> = {
+    primary: 0,
+    pending_evaluation: 1,
+    rejected: 2,
+  }
+  return filtered
+    .map((item, index) => ({ item, index }))
+    .sort((left, right) => (
+      (bucketOrder[left.item.decision_bucket] ?? 1)
+      - (bucketOrder[right.item.decision_bucket] ?? 1)
+      || left.index - right.index
+    ))
+    .map(({ item }) => item)
 })
 const taskPayload = computed<Row>(() => {
   const stored = (detail.value?.run.input as Row | undefined)?.user_input