Browse Source

增加寻找视频的需求数量

xueyiming 2 days ago
parent
commit
c04d84d8f8

+ 1 - 1
agents/demand_video_expand_agent/__init__.py

@@ -1,5 +1,5 @@
 """
-demand_video_expand_agent — S/A 需求视频点位拓展判断 Agent
+demand_video_expand_agent — S/A/B 需求视频点位拓展判断 Agent
 
 职责:对任务层已组装好的「需求 + 视频点位」做语义判断,
 筛选可作为拓展需求的点位并落库。不负责查库。

+ 3 - 3
agents/demand_video_expand_agent/tools/batch_save_demand_expansions.py

@@ -20,7 +20,7 @@ from supply_infra.db.session import get_session
 
 logger = logging.getLogger(__name__)
 
-_VALID_GRADES = frozenset({"S", "A"})
+_VALID_GRADES = frozenset({"S", "A", "B"})
 
 
 def _optional_str(value: Any) -> str | None:
@@ -150,7 +150,7 @@ def batch_save_demand_expansions(
         biz_dt: 业务日 YYYYMMDD。
         source_demand_grade_id: 来源 demand_grade.id(由用户消息提供,原样传入)。
         source_demand_name: 来源需求名。
-        source_grade: 来源等级 S 或 A
+        source_grade: 来源等级 S、A 或 B
         run_id: 任务批次 id;省略则自动生成。
 
     Returns:
@@ -166,7 +166,7 @@ def batch_save_demand_expansions(
 
     grade = _optional_str(source_grade)
     if grade not in _VALID_GRADES:
-        return f"source_grade 无效,只能是 S 或 A: {source_grade!r}"
+        return f"source_grade 无效,只能是 S、A 或 B: {source_grade!r}"
 
     try:
         grade_id = int(source_demand_grade_id)

+ 13 - 9
agents/find_agent/demand_run.py

@@ -21,6 +21,7 @@ from supply_infra.db.repositories.multi_demand_video_detail_repo import (
     MultiDemandVideoDetailRepository,
 )
 from supply_infra.db.session import get_session
+from supply_infra.scheduler.constants import FIND_AND_EXPAND_GRADES
 from supply_infra.services.video_discovery_service import get_video_discovery_service
 
 logger = logging.getLogger(__name__)
@@ -42,7 +43,7 @@ class FindDemandVideo:
 
 @dataclass
 class FindDemandContext:
-    """单条 find_agent 任务:一个 S/A 需求词及其下全部视频与全部拓展点位。"""
+    """单条 find_agent 任务:一个 S/A/B 需求词及其下全部视频与全部拓展点位。"""
 
     biz_dt: str
     demand_grade_id: int
@@ -134,10 +135,13 @@ def _serialize_video(video: FindDemandVideo) -> dict[str, Any]:
     }
 
 
+_GRADE_RANK = {"S": 0, "A": 1, "B": 2}
+
+
 def _grade_priority_key(row: Any) -> tuple[float, int, str, int]:
-    """S 优先于 A;同等级按 score 降序。"""
+    """score 降序;同分按 S > A > B。"""
     score = float(row.score) if row.score is not None else -1.0
-    grade_rank = 0 if str(row.grade) == "S" else 1
+    grade_rank = _GRADE_RANK.get(str(row.grade), 99)
     return (-score, grade_rank, str(row.demand_name), int(row.id))
 
 
@@ -205,10 +209,10 @@ def load_find_demand_contexts(
     session,
     biz_dt: str,
     *,
-    grades: Iterable[str] = ("S", "A"),
+    grades: Iterable[str] = FIND_AND_EXPAND_GRADES,
     top_limit: int | None = None,
 ) -> list[FindDemandContext]:
-    """加载指定业务日全部 S/A 需求(可选 top_limit 截断),每个需求词组装为一条完整上下文。"""
+    """加载指定业务日全部 S/A/B 需求(可选 top_limit 截断),每个需求词组装为一条完整上下文。"""
     grades_rows = DemandGradeRepository(session).list_by_biz_dt_and_grades(biz_dt, grades)
     if not grades_rows:
         return []
@@ -233,7 +237,7 @@ def load_find_demand_context_by_id(
 ) -> FindDemandContext | None:
     """按 demand_grade.id(需求 id)直接加载一条 find_agent 上下文。
 
-    不限制 S/A 等级;若记录不存在、或该需求下无有效拓展点位,返回 None。
+    不限制 S/A/B 等级;若记录不存在、或该需求下无有效拓展点位,返回 None。
     biz_dt 默认取该分级记录自身的业务日。
     """
     with get_session() as session:
@@ -248,7 +252,7 @@ def pick_find_demand_context(
     *,
     index: int = 0,
     demand_grade_id: int | None = None,
-    grades: Iterable[str] = ("S", "A"),
+    grades: Iterable[str] = FIND_AND_EXPAND_GRADES,
 ) -> FindDemandContext | None:
     """从数据库选取一条待执行的 find_agent 上下文。"""
     if demand_grade_id is not None:
@@ -270,10 +274,10 @@ def pick_find_demand_context(
 def list_find_demand_contexts(
     biz_dt: str | None = None,
     *,
-    grades: Iterable[str] = ("S", "A"),
+    grades: Iterable[str] = FIND_AND_EXPAND_GRADES,
     top_limit: int | None = None,
 ) -> tuple[str, list[FindDemandContext]]:
-    """返回解析后的业务日与待执行上下文(默认当日全部 S/A,可选 top_limit)。"""
+    """返回解析后的业务日与待执行上下文(默认当日全部 S/A/B,可选 top_limit)。"""
     resolved_biz_dt = _resolve_biz_dt(biz_dt)
     with get_session() as session:
         contexts = load_find_demand_contexts(

+ 1 - 1
api/services/demand_feedback.py

@@ -20,7 +20,7 @@ from supply_infra.db.repositories.multi_demand_video_detail_repo import (
 )
 from supply_infra.db.session import get_session
 
-_SOURCE_VIDEO_GRADES = frozenset({"B", "C", "D"})
+_SOURCE_VIDEO_GRADES = frozenset({"C", "D"})
 
 
 class FeedbackTargetNotFoundError(Exception):

+ 2 - 2
api/services/video_discovery.py

@@ -19,7 +19,7 @@ from supply_infra.db.repositories.pipeline_step_run_repo import PipelineStepRunR
 from supply_infra.db.session import get_session
 
 _POINT_TYPES = {"inspiration", "purpose", "key"}
-_SOURCE_VIDEO_GRADES = frozenset({"B", "C", "D"})
+_SOURCE_VIDEO_GRADES = frozenset({"C", "D"})
 _VIDEO_DETAIL_URL_TEMPLATE = "https://admin.piaoquantv.com/cms/post-detail/{vid}/detail"
 
 
@@ -236,7 +236,7 @@ def get_video_discovery_demand(demand_grade_id: int) -> dict[str, Any] | None:
     """
     Resolve the judged video list and hit evidence from demand_video_expansion.
 
-    S/A 使用拓展命中视频;B/C/D 使用 demand_grade.video_list 中的源视频。
+    S/A/B 使用拓展命中视频;C/D 使用 demand_grade.video_list 中的源视频。
     multi_demand_video_detail 补充 title 与解构 url1/url2。
     """
     with get_session() as session:

+ 3 - 1
find_agent_v2/AGENT_FLOW.md

@@ -105,8 +105,10 @@ Search、Evidence、Evaluator 还可以通过 `delegate_agents_v2` 将最多 8 
 ```mermaid
 flowchart LR
     SP["Supervisor 搜索计划"] --> ST["search_videos_v2"]
+    ST -->|未指定 provider| IK["内部关键词搜索接口"]
+    ST -->|provider=internal_keyword| IK
     ST -->|provider=tikhub| TK["TikHub 搜索接口"]
-    ST -->|其他 provider| IK["内部关键词搜索接口"]
+    IK -->|未指定且无候选| TK
     TK --> SR["find_agent_v2_search"]
     IK --> SR
     SR --> C["find_agent_v2_candidate<br/>pending_evaluation"]

+ 13 - 1
find_agent_v2/graph.py

@@ -30,6 +30,7 @@ from find_agent_v2.state import (
     PlannerAssignment,
     PlanningDecision,
     SearchAssignment,
+    SearchTask,
     SupervisorAssignment,
     SupervisorDecision,
 )
@@ -388,6 +389,17 @@ class FindAgentRoundGraph:
             "node_runs": [*state.get("node_runs", []), run],
         }
 
+    @staticmethod
+    def _search_task_already_executed(
+        task: SearchTask,
+        existing: set[tuple[str, str, str]],
+    ) -> bool:
+        keyword = task.keyword
+        reason = task.query_reason
+        if not task.provider:
+            return (keyword, "internal_keyword", reason) in existing
+        return (keyword, task.provider, reason) in existing
+
     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")
@@ -404,7 +416,7 @@ class FindAgentRoundGraph:
         tasks = [
             item
             for item in execution_plan.search_tasks
-            if (item.keyword, item.provider, item.query_reason) not in existing
+            if not self._search_task_already_executed(item, existing)
         ][:6]
         assignment = SearchAssignment(
             run_id=state["run_id"],

+ 4 - 1
find_agent_v2/prompts.py

@@ -35,7 +35,8 @@ PLANNER_PROMPT = COMMON_RULES + """
 你是唯一可以读取原始需求的模块。把原始需求、参考视频和点位压缩为结构化 ExecutionPlan,供后续
 模块执行。DemandBrief 和 EvaluationBrief 必须保留判断相关性、人群适配、排除项和时间适用性所需
 的业务语义,但不要复制原文或整段参考视频。生成 1~6 个互补且可以独立执行的 SearchTask,每个
-task_id 唯一。首次动作通常是 search。只输出符合宿主结构化 schema 的结果。
+task_id 唯一。`provider` 不填则先内部搜索、无结果再 TikHub;填 `internal_keyword` 只用内部搜索;
+填 `tikhub` 只用 TikHub。首次动作通常是 search。只输出符合宿主结构化 schema 的结果。
 """
 
 
@@ -47,6 +48,8 @@ SUPERVISOR_PROMPT = COMMON_RULES + """
 每次只能提议一个动作:
 
 - `search`:继续扩大候选池;需要新方向时通过 `additional_search_tasks` 给出结构化任务。
+  `provider` 不填则先内部搜索、无结果再 TikHub;填 `internal_keyword` 只用内部搜索;填 `tikhub`
+  只用 TikHub。
 - `evidence`:补证,可用 `evidence_scope=detail|portrait|both` 决定顺序。
 - `evaluator`:证据已处理后评分分池。
 - `finish`:没有待处理候选且继续搜索已无信息增益时结束本轮。

+ 4 - 1
find_agent_v2/state.py

@@ -38,7 +38,10 @@ class SearchTask(BaseModel):
     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"
+    provider: Literal["internal_keyword", "tikhub"] | None = Field(
+        default=None,
+        description="不填则先内部搜索、无结果再 TikHub;internal_keyword 只用内部搜索;tikhub 只用 TikHub。",
+    )
     max_pages: int = Field(default=1, ge=1, le=2)
     coverage_targets: list[str] = Field(default_factory=list)
 

+ 109 - 45
find_agent_v2/tools.py

@@ -197,9 +197,78 @@ def bound_candidate_tools(
     return tuple(output)
 
 
+def _search_has_results(payload: dict[str, Any]) -> bool:
+    return bool(payload.get("search_results"))
+
+
+async def _search_provider_pages(
+    *,
+    service: Any,
+    run_id: str,
+    round_index: int,
+    keyword: str,
+    query_reason: str,
+    source_type: str,
+    provider: str,
+    raw_task: dict[str, Any],
+    common: dict[str, Any],
+    max_pages: int,
+    cursor: str | int,
+    provider_search_id: str = "",
+    backtrace: str = "",
+    extra_output: dict[str, Any] | None = None,
+) -> tuple[list[dict[str, Any]], bool]:
+    outputs: list[dict[str, Any]] = []
+    has_results = False
+    for page_no in range(1, max_pages + 1):
+        if provider == "tikhub":
+            payload = await search_tikhub(
+                **common,
+                cursor=int(cursor or 0),
+                filter_duration=str(raw_task.get("filter_duration") or "不限"),
+                search_id=provider_search_id,
+                backtrace=backtrace,
+            )
+        else:
+            payload = await search_internal(**common, cursor=str(cursor or "0"))
+        saved = service.save_search(
+            run_id=run_id,
+            round_index=int(round_index),
+            keyword=keyword,
+            query_reason=query_reason,
+            source_type=source_type,
+            provider=provider,
+            cursor=str(cursor),
+            page_no=page_no,
+            payload=payload,
+        )
+        if _search_has_results(payload):
+            has_results = True
+        outputs.append({
+            "keyword": keyword,
+            "provider": provider,
+            "page_no": page_no,
+            "error": payload.get("error"),
+            "has_more": bool(payload.get("has_more")),
+            "next_cursor": payload.get("next_cursor"),
+            **(extra_output or {}),
+            **saved,
+        })
+        if payload.get("error") or not payload.get("has_more"):
+            break
+        cursor = payload.get("next_cursor") or cursor
+        provider_search_id = str(payload.get("search_id") or provider_search_id)
+        backtrace = str(payload.get("backtrace") or backtrace)
+    return outputs, has_results
+
+
 @tool
 async def search_videos_v2(run_id: str, round_index: int, searches: list[dict[str, Any]]) -> str:
-    """批量搜索并仅写入 find_agent_v2_search/candidate;searches 最多 6 项。"""
+    """批量搜索并仅写入 find_agent_v2_search/candidate;searches 最多 6 项。
+
+    未指定 provider 时先打内部搜索,无候选再回退 TikHub。
+    显式指定 internal_keyword 只用内部搜索;显式指定 tikhub 只用 TikHub。
+    """
     if not searches or len(searches) > 6:
         return json.dumps({"error": "searches 必须为 1~6 项"}, ensure_ascii=False)
     service = get_find_agent_v2_service()
@@ -207,58 +276,53 @@ async def search_videos_v2(run_id: str, round_index: int, searches: list[dict[st
     for raw_task in searches:
         keyword = str(raw_task.get("keyword") or "").strip()
         reason = str(raw_task.get("query_reason") or "").strip()
-        provider = str(raw_task.get("provider") or "internal_keyword")
+        requested_provider = str(raw_task.get("provider") or "").strip()
+        allow_tikhub_fallback = not requested_provider
+        provider = "tikhub" if requested_provider == "tikhub" else "internal_keyword"
         if not keyword or not reason:
             outputs.append({"error": "keyword/query_reason 不能为空"})
             continue
         max_pages = max(1, min(int(raw_task.get("max_pages") or 1), 2))
-        cursor: str | int = raw_task.get("cursor") or 0
-        provider_search_id = str(raw_task.get("search_id") or "")
-        backtrace = str(raw_task.get("backtrace") or "")
-        for page_no in range(1, max_pages + 1):
-            common = {
-                "keyword": keyword,
-                "content_type": str(raw_task.get("content_type") or "视频"),
-                "sort_type": str(raw_task.get("sort_type") or "综合排序"),
-                "publish_time": str(raw_task.get("publish_time") or "不限"),
-                "min_duration_seconds": int(raw_task.get("min_duration_seconds") or 30),
-            }
-            if provider == "tikhub":
-                payload = await search_tikhub(
-                    **common,
-                    cursor=int(cursor or 0),
-                    filter_duration=str(raw_task.get("filter_duration") or "不限"),
-                    search_id=provider_search_id,
-                    backtrace=backtrace,
-                )
-            else:
-                provider = "internal_keyword"
-                payload = await search_internal(**common, cursor=str(cursor or "0"))
-            saved = service.save_search(
+        source_type = str(raw_task.get("source_type") or "mixed")
+        common = {
+            "keyword": keyword,
+            "content_type": str(raw_task.get("content_type") or "视频"),
+            "sort_type": str(raw_task.get("sort_type") or "综合排序"),
+            "publish_time": str(raw_task.get("publish_time") or "不限"),
+            "min_duration_seconds": int(raw_task.get("min_duration_seconds") or 30),
+        }
+        pages, has_results = await _search_provider_pages(
+            service=service,
+            run_id=run_id,
+            round_index=round_index,
+            keyword=keyword,
+            query_reason=reason,
+            source_type=source_type,
+            provider=provider,
+            raw_task=raw_task,
+            common=common,
+            max_pages=max_pages,
+            cursor=raw_task.get("cursor") or 0,
+            provider_search_id=str(raw_task.get("search_id") or ""),
+            backtrace=str(raw_task.get("backtrace") or ""),
+        )
+        outputs.extend(pages)
+        if allow_tikhub_fallback and not has_results:
+            fallback, _ = await _search_provider_pages(
+                service=service,
                 run_id=run_id,
-                round_index=int(round_index),
+                round_index=round_index,
                 keyword=keyword,
                 query_reason=reason,
-                source_type=str(raw_task.get("source_type") or "mixed"),
-                provider=provider,
-                cursor=str(cursor),
-                page_no=page_no,
-                payload=payload,
+                source_type=source_type,
+                provider="tikhub",
+                raw_task=raw_task,
+                common=common,
+                max_pages=max_pages,
+                cursor=0,
+                extra_output={"fallback_from": "internal_keyword"},
             )
-            outputs.append({
-                "keyword": keyword,
-                "provider": provider,
-                "page_no": page_no,
-                "error": payload.get("error"),
-                "has_more": bool(payload.get("has_more")),
-                "next_cursor": payload.get("next_cursor"),
-                **saved,
-            })
-            if payload.get("error") or not payload.get("has_more"):
-                break
-            cursor = payload.get("next_cursor") or cursor
-            provider_search_id = str(payload.get("search_id") or provider_search_id)
-            backtrace = str(payload.get("backtrace") or backtrace)
+            outputs.extend(fallback)
     return json.dumps({"run_id": run_id, "searches": outputs}, ensure_ascii=False)
 
 

+ 2 - 2
supply_infra/db/models/demand_video_expansion.py

@@ -9,7 +9,7 @@ from supply_infra.db.base import Base
 
 
 class DemandVideoExpansion(Base):
-    """S/A 需求关联视频点位拓展结果。"""
+    """S/A/B 需求关联视频点位拓展结果。"""
 
     __tablename__ = "demand_video_expansion"
     __table_args__ = (
@@ -33,7 +33,7 @@ class DemandVideoExpansion(Base):
     source_demand_name: Mapped[str] = mapped_column(
         String(256), nullable=False, comment="来源需求名"
     )
-    source_grade: Mapped[str] = mapped_column(String(4), nullable=False, comment="来源等级 S/A")
+    source_grade: Mapped[str] = mapped_column(String(4), nullable=False, comment="来源等级 S/A/B")
     expanded_text: Mapped[str] = mapped_column(
         String(512), nullable=False, comment="拓展需求文本(来自 point_data)"
     )

+ 1 - 1
supply_infra/db/repositories/demand_grade_repo.py

@@ -57,7 +57,7 @@ class DemandGradeRepository(BaseRepository[DemandGrade]):
     def list_by_biz_dt_and_grades(
         self,
         biz_dt: str,
-        grades: Iterable[str] = ("S", "A"),
+        grades: Iterable[str] = ("S", "A", "B"),
     ) -> list[DemandGrade]:
         """返回指定业务日、指定等级的分级结果,按等级、需求名排序。"""
         grade_list = [g for g in grades if g]

+ 2 - 1
supply_infra/scheduler/constants.py

@@ -7,5 +7,6 @@ SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 AIGC_PUBLISH_JOB_ID = "publish_videos_from_discovery"
 AIGC_PUBLISH_JOB_NAME = "AIGC候选分发"
 
-# find_agent:当日全部 S/A 需求(有拓展点位),任务队列 + 固定 worker 并发找视频
+# find_agent:当日 S/A/B 需求(有拓展点位),任务队列 + 固定 worker 并发找视频
 PIPELINE_FIND_AGENT_WORKERS = 5
+FIND_AND_EXPAND_GRADES = ("S", "A", "B")

+ 2 - 2
supply_infra/scheduler/jobs/discover_videos_from_demands.py

@@ -1,5 +1,5 @@
 """
-从全部 S/A 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 1000 提前结束。
+从全部 S/A/B 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 1000 提前结束。
 
 任务层负责查库与组装上下文;Agent 负责搜索、画像与分池落库。
 
@@ -478,7 +478,7 @@ def discover_videos_from_demands(
     force: bool = False,
 ) -> dict[str, Any]:
     """
-    对指定业务日全部 S/A 级需求(有拓展点位)以任务队列方式并发调用 find_agent。
+    对指定业务日全部 S/A/B 级需求(有拓展点位)以任务队列方式并发调用 find_agent。
 
     - 先把当天待执行需求全部放入内存队列(此时不写 video_discovery_run)
     - 固定 worker 线程循环取任务;单任务开始执行时再创建 run 记录

+ 7 - 5
supply_infra/scheduler/jobs/expand_demand_from_video_points.py

@@ -1,5 +1,5 @@
 """
-从 S/A 级需求关联视频中挖掘拓展需求。
+从 S/A/B 级需求关联视频中挖掘拓展需求。
 
 任务层负责查库与组装上下文;Agent 仅做语义判断与落库。
 """
@@ -29,10 +29,12 @@ from supply_infra.db.repositories.multi_demand_video_point_repo import (
     MultiDemandVideoPointRepository,
 )
 from supply_infra.db.session import get_session
+from supply_infra.scheduler.constants import FIND_AND_EXPAND_GRADES
 
 logger = logging.getLogger(__name__)
 
 _DEFAULT_WORKERS = 5
+EXPAND_GRADES = FIND_AND_EXPAND_GRADES
 
 
 def _resolve_biz_dt(biz_dt: str | None) -> str:
@@ -94,14 +96,14 @@ def load_expand_contexts(
 ) -> tuple[list[DemandExpandContext], dict[str, int]]:
     """加载待拓展判断的需求上下文,并返回跳过统计。"""
     stats = {
-        "total_sa": 0,
+        "total": 0,
         "skipped_no_video": 0,
         "skipped_no_points": 0,
         "skipped_already_done": 0,
     }
 
-    grades = DemandGradeRepository(session).list_by_biz_dt_and_grades(biz_dt, ("S", "A"))
-    stats["total_sa"] = len(grades)
+    grades = DemandGradeRepository(session).list_by_biz_dt_and_grades(biz_dt, EXPAND_GRADES)
+    stats["total"] = len(grades)
 
     finished_ids: set[int] = set()
     if skip_finished:
@@ -241,7 +243,7 @@ def expand_demand_from_video_points(
     workers: int = _DEFAULT_WORKERS,
 ) -> dict[str, Any]:
     """
-    对指定业务日的 S/A 需求执行视频点位拓展判断。
+    对指定业务日的 S/A/B 需求执行视频点位拓展判断。
 
     程序负责查需求与点位;无视频或无点位则跳过;有数据则并发调用 Agent。
     """

+ 177 - 1
tests/supply_agent/test_find_agent_v2.py

@@ -155,7 +155,7 @@ async def test_host_direct_search_and_evidence_paths_use_no_llm(monkeypatch) ->
             "keyword": "日本间谍",
             "query_reason": "测试",
             "source_type": "mixed",
-            "provider": "internal_keyword",
+            "provider": None,
             "max_pages": 1,
             "coverage_targets": [],
         }]),
@@ -509,6 +509,182 @@ async def test_langchain_runtime_runs_without_network(monkeypatch) -> None:
     assert result.tool_calls_made == 0
 
 
+class _FakeSearchService:
+    def __init__(self) -> None:
+        self.saved: list[dict] = []
+
+    def save_search(self, **kwargs):
+        self.saved.append(kwargs)
+        results = list((kwargs.get("payload") or {}).get("search_results") or [])
+        return {
+            "search_id": len(self.saved),
+            "new_candidate_count": len(results),
+            "result_count": len(results),
+            "share_gate_rejected_count": 0,
+        }
+
+
+def _empty_search_payload(*, provider: str, error: str | None = None) -> dict:
+    return {
+        "provider": provider,
+        "search_results": [],
+        "has_more": False,
+        "next_cursor": "",
+        "error": error,
+    }
+
+
+def _hit_search_payload(*, provider: str, aweme_id: str) -> dict:
+    return {
+        "provider": provider,
+        "search_results": [{"aweme_id": aweme_id, "desc": "命中"}],
+        "has_more": False,
+        "next_cursor": "",
+    }
+
+
+@pytest.mark.asyncio
+async def test_default_search_falls_back_to_tikhub_when_internal_is_empty(
+    monkeypatch,
+) -> None:
+    service = _FakeSearchService()
+    providers: list[str] = []
+
+    async def fake_internal(**_kwargs):
+        providers.append("internal_keyword")
+        return _empty_search_payload(provider="internal_keyword")
+
+    async def fake_tikhub(**_kwargs):
+        providers.append("tikhub")
+        return _hit_search_payload(provider="tikhub", aweme_id="1")
+
+    monkeypatch.setattr(find_agent_tools, "get_find_agent_v2_service", lambda: service)
+    monkeypatch.setattr(find_agent_tools, "search_internal", fake_internal)
+    monkeypatch.setattr(find_agent_tools, "search_tikhub", fake_tikhub)
+
+    raw = await find_agent_tools.search_videos_v2(
+        run_id="run",
+        round_index=1,
+        searches=[{
+            "keyword": "甲午战争 民族觉醒 历史真相",
+            "query_reason": "扩大候选池",
+        }],
+    )
+    payload = json.loads(raw)
+
+    assert providers == ["internal_keyword", "tikhub"]
+    assert [item["provider"] for item in payload["searches"]] == [
+        "internal_keyword",
+        "tikhub",
+    ]
+    assert payload["searches"][1]["fallback_from"] == "internal_keyword"
+    assert [item["provider"] for item in service.saved] == [
+        "internal_keyword",
+        "tikhub",
+    ]
+
+
+@pytest.mark.asyncio
+async def test_default_search_skips_tikhub_when_internal_has_results(
+    monkeypatch,
+) -> None:
+    service = _FakeSearchService()
+    providers: list[str] = []
+
+    async def fake_internal(**_kwargs):
+        providers.append("internal_keyword")
+        return _hit_search_payload(provider="internal_keyword", aweme_id="2")
+
+    async def fake_tikhub(**_kwargs):
+        providers.append("tikhub")
+        return _hit_search_payload(provider="tikhub", aweme_id="3")
+
+    monkeypatch.setattr(find_agent_tools, "get_find_agent_v2_service", lambda: service)
+    monkeypatch.setattr(find_agent_tools, "search_internal", fake_internal)
+    monkeypatch.setattr(find_agent_tools, "search_tikhub", fake_tikhub)
+
+    raw = await find_agent_tools.search_videos_v2(
+        run_id="run",
+        round_index=1,
+        searches=[{
+            "keyword": "国家安全 反间谍",
+            "query_reason": "未指定 provider 且内部有结果",
+        }],
+    )
+    payload = json.loads(raw)
+
+    assert providers == ["internal_keyword"]
+    assert [item["provider"] for item in payload["searches"]] == ["internal_keyword"]
+    assert "fallback_from" not in payload["searches"][0]
+
+
+@pytest.mark.asyncio
+async def test_explicit_internal_search_does_not_fallback_when_empty(
+    monkeypatch,
+) -> None:
+    service = _FakeSearchService()
+    providers: list[str] = []
+
+    async def fake_internal(**_kwargs):
+        providers.append("internal_keyword")
+        return _empty_search_payload(provider="internal_keyword")
+
+    async def fake_tikhub(**_kwargs):
+        providers.append("tikhub")
+        return _hit_search_payload(provider="tikhub", aweme_id="3")
+
+    monkeypatch.setattr(find_agent_tools, "get_find_agent_v2_service", lambda: service)
+    monkeypatch.setattr(find_agent_tools, "search_internal", fake_internal)
+    monkeypatch.setattr(find_agent_tools, "search_tikhub", fake_tikhub)
+
+    raw = await find_agent_tools.search_videos_v2(
+        run_id="run",
+        round_index=1,
+        searches=[{
+            "keyword": "国家安全 反间谍",
+            "query_reason": "指定只用内部搜索",
+            "provider": "internal_keyword",
+        }],
+    )
+    payload = json.loads(raw)
+
+    assert providers == ["internal_keyword"]
+    assert [item["provider"] for item in payload["searches"]] == ["internal_keyword"]
+    assert "fallback_from" not in payload["searches"][0]
+
+
+@pytest.mark.asyncio
+async def test_explicit_tikhub_search_does_not_call_internal(monkeypatch) -> None:
+    service = _FakeSearchService()
+    providers: list[str] = []
+
+    async def fake_internal(**_kwargs):
+        providers.append("internal_keyword")
+        return _hit_search_payload(provider="internal_keyword", aweme_id="4")
+
+    async def fake_tikhub(**_kwargs):
+        providers.append("tikhub")
+        return _hit_search_payload(provider="tikhub", aweme_id="5")
+
+    monkeypatch.setattr(find_agent_tools, "get_find_agent_v2_service", lambda: service)
+    monkeypatch.setattr(find_agent_tools, "search_internal", fake_internal)
+    monkeypatch.setattr(find_agent_tools, "search_tikhub", fake_tikhub)
+
+    raw = await find_agent_tools.search_videos_v2(
+        run_id="run",
+        round_index=1,
+        searches=[{
+            "keyword": "抗日战争 历史档案",
+            "query_reason": "指定 TikHub",
+            "provider": "tikhub",
+        }],
+    )
+    payload = json.loads(raw)
+
+    assert providers == ["tikhub"]
+    assert [item["provider"] for item in payload["searches"]] == ["tikhub"]
+
+
 def test_stage_tool_allowlists_are_physical_and_isolated() -> None:
     assert _names(SEARCH_TOOLS) == {"search_videos_v2", "query_find_agent_v2_state"}
     assert _names(EVIDENCE_TOOLS) == {

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

@@ -16,6 +16,7 @@ const totalPages = ref(1)
 const hasPreviousPage = ref(false)
 const hasNextPage = ref(false)
 const pageSize = 30
+const POLL_INTERVAL_MS = 5 * 60 * 1000
 const detail = ref<FindAgentV2Detail | null>(null)
 const selectedId = ref('')
 const demandWord = ref('')
@@ -107,7 +108,7 @@ onMounted(async () => {
   await loadRuns()
   poll = window.setInterval(() => {
     if (run.value?.status === 'running') void refresh(false)
-  }, 5000)
+  }, POLL_INTERVAL_MS)
 })
 onBeforeUnmount(() => poll && window.clearInterval(poll))