فهرست منبع

feat(mode_workflow): 新增工具知识导入全套工具并优化页面展示

- 新增`tool_url_ingest.json`配置与`ingest_urls.py`脚本,支持跳过搜索环节直接通过指定渠道+链接批量导入工具知识
- 新增维度提取表导入脚本`import_knowledge.py`,支持结构化维度数据批量导入知识接口
- 新增品类聚焦维度配置`category_focus_dimensions.json`及对应导入脚本
- 新增工具解构知识导入脚本`import_destruction_knowledge.py`
- 在`db.py`中新增解构知识防重台账与相关操作函数
- 优化`index.html`工具表格布局,调整列宽与表头结构提升展示效果
刘文武 2 هفته پیش
والد
کامیت
38e7733b29

+ 5 - 1
examples/mode_workflow/.gitignore

@@ -13,4 +13,8 @@ import_process_knowledge.处理方式.md
 工序接口文档.md
 工序接口文档.md
 流程执行手册.md
 流程执行手册.md
 拆表_3-II_设计.md
 拆表_3-II_设计.md
-docs/
+docs/
+stages/output/
+stages/input/
+流程执行手册.md
+工序接口文档.md

+ 71 - 0
examples/mode_workflow/db.py

@@ -219,6 +219,24 @@ CREATE TABLE IF NOT EXISTS tools_ingest_log (
 ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工具知识已导入台账(防重复上传)';
 ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='工具知识已导入台账(防重复上传)';
 """
 """
 
 
+# 「解构知识(工具·仅实质/形式作用域 → what-制作还原)」已导入台账。
+# 与 tools_ingest_log 分表:两者同源 mode_tools、同键 (case_id, tool_index),但上传到不同维度
+# (本表 = what-制作还原,tools_ingest_log = how工具),共表会在 (case_id, tool_index) 撞键,
+# 使一方上传后另一方被误判「已传过」而跳过。故独立成表(stages/import_destruction_knowledge.py 用)。
+DDL_DESTRUCTION_INGEST_LOG = """
+CREATE TABLE IF NOT EXISTS destruction_ingest_log (
+  id            BIGINT AUTO_INCREMENT PRIMARY KEY,
+  case_id       VARCHAR(128)  NOT NULL,
+  tool_index    INT           NOT NULL COMMENT '工具序号(1-based),对齐导入脚本枚举',
+  version       VARCHAR(32)   NULL COMMENT '导入时 mode_tools 版本;变了应重导',
+  knowledge_id  VARCHAR(128)  NULL COMMENT '接口返回的 knowledge_id',
+  api_url       VARCHAR(255)  NULL,
+  ingested_at   TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
+  updated_at    TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
+  UNIQUE KEY uk_case_tool (case_id, tool_index)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='解构知识(what-制作还原)已导入台账(防重复上传)';
+"""
+
 
 
 # step 维度归类结果独立表(3-I,见 docs/step_classification_3-I_设计.md):一行 = 某 case 某版本
 # step 维度归类结果独立表(3-I,见 docs/step_classification_3-I_设计.md):一行 = 某 case 某版本
 # 某工序某 step 某维度某子项的命中分类。与 mode_process 解耦,归类重跑只覆盖本表、不重写 steps blob;
 # 某工序某 step 某维度某子项的命中分类。与 mode_process 解耦,归类重跑只覆盖本表、不重写 steps blob;
@@ -344,6 +362,7 @@ def init_tables():
             cur.execute(DDL_TOOLS)
             cur.execute(DDL_TOOLS)
             cur.execute(DDL_INGEST_LOG)
             cur.execute(DDL_INGEST_LOG)
             cur.execute(DDL_TOOLS_INGEST_LOG)
             cur.execute(DDL_TOOLS_INGEST_LOG)
+            cur.execute(DDL_DESTRUCTION_INGEST_LOG)
             cur.execute(DDL_STEP_CLASSIFICATION)
             cur.execute(DDL_STEP_CLASSIFICATION)
             cur.execute(DDL_PROCESS_RUN)
             cur.execute(DDL_PROCESS_RUN)
             cur.execute(DDL_PROCESS_PROCEDURE)
             cur.execute(DDL_PROCESS_PROCEDURE)
@@ -1715,6 +1734,30 @@ def fetch_adopted_process_cases(query_id=None):
     return sorted(set(cases))
     return sorted(set(cases))
 
 
 
 
+def fetch_destructed_tools_cases(query_ids=None):
+    """返回「有工具解构」的 case_id 列表(**不看采纳**),按 mode_tools.query_id 过滤。
+
+    与 fetch_adopted_tools_cases 的区别:不 JOIN search_tools、不做 is_adopted_rel 过滤,
+    直接取 mode_tools 里存在解构行的帖子。query_ids=None 取全部;给列表则 IN 过滤
+    (可跨多个 query;同一 case 在多 query 下解构时按 case_id 去重)。去重、排序返回。
+    供「所有已解构(含未采纳)」上传口径用。
+    """
+    sql = "SELECT DISTINCT case_id FROM mode_tools"
+    params = ()
+    if query_ids is not None:
+        if not query_ids:
+            return []   # 显式空列表:直接空结果
+        sql += " WHERE query_id IN (" + ",".join(["%s"] * len(query_ids)) + ")"
+        params = tuple(query_ids)
+    conn = _conn()
+    try:
+        with conn.cursor() as cur:
+            cur.execute(sql, params)
+            return sorted({r["case_id"] for r in cur.fetchall()})
+    finally:
+        conn.close()
+
+
 def fetch_adopted_tools_cases(query_id=None):
 def fetch_adopted_tools_cases(query_id=None):
     """返回「已采纳且有工具解构」的 case_id 列表(供工具知识上传脚本用)。
     """返回「已采纳且有工具解构」的 case_id 列表(供工具知识上传脚本用)。
 
 
@@ -1865,6 +1908,34 @@ def mark_tools_ingested(case_id, tool_index, version, knowledge_id=None, api_url
         conn.close()
         conn.close()
 
 
 
 
+def fetch_destruction_ingested_map(case_id):
+    """返回 {tool_index: version} —— 该 case 各工具的解构知识(what-制作还原)已导入的版本。
+    独立台账(destruction_ingest_log),与工具方向 tools_ingest_log 互不干扰(见 DDL 注释)。"""
+    conn = _conn()
+    try:
+        with conn.cursor() as cur:
+            cur.execute("SELECT tool_index, version FROM destruction_ingest_log WHERE case_id=%s",
+                        (case_id,))
+            return {r["tool_index"]: r["version"] for r in cur.fetchall()}
+    finally:
+        conn.close()
+
+
+def mark_destruction_ingested(case_id, tool_index, version, knowledge_id=None, api_url=None):
+    """记一条解构知识「已导入」台账(case_id+tool_index 唯一,重导同序号则更新版本/knowledge_id)。"""
+    conn = _conn()
+    try:
+        with conn.cursor() as cur:
+            cur.execute("""INSERT INTO destruction_ingest_log
+                             (case_id, tool_index, version, knowledge_id, api_url)
+                           VALUES (%s,%s,%s,%s,%s)
+                           ON DUPLICATE KEY UPDATE version=VALUES(version),
+                             knowledge_id=VALUES(knowledge_id), api_url=VALUES(api_url)""",
+                        (case_id, tool_index, version, knowledge_id, api_url))
+    finally:
+        conn.close()
+
+
 def fetch_dashboard_rows():
 def fetch_dashboard_rows():
     """拉 Dashboard 计算所需的轻量行。数据量级:百~千行,Python 聚合足够。
     """拉 Dashboard 计算所需的轻量行。数据量级:百~千行,Python 聚合足够。
     优化:① 不传 llm_evaluation 整块,SQL 只取采纳判定要的相关性得分;
     优化:① 不传 llm_evaluation 整块,SQL 只取采纳判定要的相关性得分;

+ 13 - 24
examples/mode_workflow/index.html

@@ -1141,7 +1141,8 @@
         border-collapse: separate;
         border-collapse: separate;
         border-spacing: 0;
         border-spacing: 0;
         width: 100%;
         width: 100%;
-        min-width: 1180px;
+        min-width: 640px;
+        table-layout: fixed;
         background: #fff;
         background: #fff;
         font-size: 12.5px;
         font-size: 12.5px;
       }
       }
@@ -3881,7 +3882,8 @@
           inner = _toolCell(t[c]);
           inner = _toolCell(t[c]);
         }
         }
         if (["输入", "输出", "用法", "缺点"].includes(c)) style = "max-width:240px;";
         if (["输入", "输出", "用法", "缺点"].includes(c)) style = "max-width:240px;";
-        else if (c === "实质作用域" || c === "形式作用域") style = "max-width:170px;";
+        else if (c === "工具名称") style = "width:150px;white-space:normal;";
+        else if (c === "创作层级") style = "width:96px;white-space:nowrap;";
         else if (!clampable) style = "white-space:nowrap;";
         else if (!clampable) style = "white-space:nowrap;";
         return { inner, cls, clampable, style };
         return { inner, cls, clampable, style };
       }
       }
@@ -3897,31 +3899,18 @@
       function renderTools(data) {
       function renderTools(data) {
         const tools = data.tools || [];
         const tools = data.tools || [];
         if (!tools.length) return renderSourceBlock() + '<div class="empty">本版本无工具</div>';
         if (!tools.length) return renderSourceBlock() + '<div class="empty">本版本无工具</div>';
-        /* 案例 group(输入/输出/效果)放在 用法 后、缺点 前;用 colspan/rowspan 做两层表头 */
-        const before = ["工具名称", "创作层级", "实质作用域", "形式作用域", "输入", "输出", "用法"];
-        const after = ["缺点", "来源链接", "最新更新时间"];
-        const thead = `<thead>
-    <tr>
-      ${before.map((c) => `<th rowspan="2">${c}</th>`).join("")}
-      <th colspan="3" class="th-group">案例</th>
-      ${after.map((c) => `<th rowspan="2">${c}</th>`).join("")}
-    </tr>
-    <tr>${["输入", "输出", "效果"].map((c) => `<th class="th-sub">${c}</th>`).join("")}</tr>
-  </thead>`;
+        /* 只展示 4 个核心列(其余字段仍在后台提取/存储,此处不渲染) */
+        const cols = ["工具名称", "创作层级", "实质作用域", "形式作用域"];
+        /* 只给前两列固定宽度;实质/形式作用域不设宽度,由它们平分剩余宽度(表格 width:100%)。
+           不再补空白填充列——否则右侧会留出一大块空的表头/单元格。 */
+        const _thStyle = { 工具名称: "width:150px;", 创作层级: "width:96px;" };
+        const thead = `<thead><tr>${cols
+          .map((c) => `<th${_thStyle[c] ? ` style="${_thStyle[c]}"` : ""}>${c}</th>`)
+          .join("")}</tr></thead>`;
         const rows = tools
         const rows = tools
           .map((t, ti) => {
           .map((t, ti) => {
-            const cases = Array.isArray(t["案例"]) && t["案例"].length ? t["案例"] : [null];
-            const K = cases.length;
             const par = ti % 2 ? "tr-b" : "tr-a";
             const par = ti % 2 ? "tr-b" : "tr-a";
-            return cases
-              .map((cse, i) => {
-                const caseTds = `${_caseTd(cse, "输入")}${_caseTd(cse, "输出")}${_caseTd(cse, "效果")}`;
-                if (i === 0) {
-                  return `<tr class="${par}">${before.map((c) => _td(c, t, K)).join("")}${caseTds}${after.map((c) => _td(c, t, K)).join("")}</tr>`;
-                }
-                return `<tr class="${par}">${caseTds}</tr>`;
-              })
-              .join("");
+            return `<tr class="${par}">${cols.map((c) => _td(c, t)).join("")}</tr>`;
           })
           })
           .join("");
           .join("");
         return (
         return (

+ 449 - 0
examples/mode_workflow/stages/import_destruction_knowledge.py

@@ -0,0 +1,449 @@
+"""
+把数据库里「已采纳」的工具解构(mode_tools)按「解构知识」口径批量导入知识接口。
+
+与 stages/import_tools_knowledge.py 同源(都从 mode_tools 取数),区别只在上传口径:
+  import_tools_knowledge —— dim=how工具,scopes(实质/形式) + custom_ext(工具名/输入/输出)。
+  本脚本(解构知识)       —— dim=what-制作解构,**只上传实质作用域 + 形式作用域**两类 scope,
+                            不带 custom_ext。即只把「这类图里有什么 / 什么风格」沉淀为还原维度。
+
+字段映射:
+  source.id          ← case_id(DB 主键)
+  source.title       ← mode_tools.source 块里的帖子标题(老数据回退 post_title)
+  每个 tool → 一条知识:
+    title          ← tool.工具名称;为空回退 "来源标题 — 工具N"
+    content        ← 整个 tool 对象的 JSON 串(保留上下文)
+    dim_attributes ← 固定 ["what-制作解构"]   dim_creations ← 固定 ["制作"]
+    scopes         ← 实质作用域(substance)/形式作用域(form)按顿号拆分去重(**仅此两类**)
+
+采纳口径:db.is_adopted_rel(相关性<4 / 实现完整性<4 / 发布超两年 / 综合分<6 任一命中即不采纳)。
+
+去重台账:独立表 destruction_ingest_log(case_id, tool_index),与 tools_ingest_log 分表——
+两者同源 mode_tools、同键,但上传到不同维度(what-制作还原 vs how工具),共表会互相误判「已传过」。
+
+用法:
+    python stages/import_destruction_knowledge.py                      # 真实导入(采纳工具全量)
+    python stages/import_destruction_knowledge.py --dry-run            # 只组装,不调接口
+    python stages/import_destruction_knowledge.py --dry-run --verbose  # 打印完整 payload JSON
+    python stages/import_destruction_knowledge.py --dump               # 组装并把入参落盘,不调接口 ★
+    python stages/import_destruction_knowledge.py --dump --limit 10    # 只落盘前 10 个 case 便于核对
+    python stages/import_destruction_knowledge.py --query-id q0001         # 只传某搜索任务下的采纳 case
+    python stages/import_destruction_knowledge.py --query-id q3300,q3299   # 多个 query(逗号分隔,按 case 去重)
+    python stages/import_destruction_knowledge.py --query-id q3300 --all-destructed  # 该 query 下所有已解构(含未采纳)★
+    python stages/import_destruction_knowledge.py --case-ids xhs_a,gzh_b   # 只传指定 case(当前口径集内精确补传)
+    python stages/import_destruction_knowledge.py --limit 5            # 只处理前 5 个 case(调试)
+    python stages/import_destruction_knowledge.py --api-url http://... # 指定后端地址
+    python stages/import_destruction_knowledge.py --delay 200          # 每次调用间隔 200ms
+    python stages/import_destruction_knowledge.py --force              # 忽略去重台账,强制重导
+
+--dump 落盘目录:stages/output/destruction_knowledge/<YYYYMMDD_HHMMSS>/
+    ├─ <case_id>__tool{N}.json   每个工具一条 payload(即上传入参)
+    └─ _index.json               本次汇总(case/工具/scope 数/版本/是否会被去重跳过)
+"""
+
+import argparse
+import json
+import logging
+import sys
+import time
+from datetime import datetime
+from pathlib import Path
+
+import requests
+
+# 本脚本归档在 stages/ 子目录,补 mode_workflow/ 到 sys.path 以裸 import db
+sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
+import db
+
+# ── 配置 ──────────────────────────────────────────────────────────────────────
+
+DEFAULT_API_URL = "http://47.236.83.130:8001"
+INGEST_ENDPOINT = "/api/v1/knowledge/ingest"
+
+DIM_ATTRIBUTES = ["what-制作解构"]
+DIM_CREATIONS = ["制作"]
+
+# --dump 落盘根目录(相对本脚本:stages/output/destruction_knowledge/)
+DUMP_ROOT = Path(__file__).resolve().parent / "output" / "destruction_knowledge"
+
+# ── 日志 ──────────────────────────────────────────────────────────────────────
+
+logging.basicConfig(
+    level=logging.INFO,
+    format="%(asctime)s [%(levelname)s] %(message)s",
+    datefmt="%H:%M:%S",
+)
+logger = logging.getLogger(__name__)
+
+
+# ── 数据来源:从 DB 取采纳工具(同 import_tools_knowledge)──────────────────────────
+
+def _build_source_data(payload):
+    """从 fetch_tools 的 payload 提炼上传所需的来源字段。
+    优先用解构时存下的 source 块(tool_extract._row_to_source 产出);老数据无 source
+    时回退 mode_tools 自带的 post_title/platform(此时点赞/正文/时间取不到,留空)。
+    """
+    src = payload.get("source") or {}
+    post = src.get("post") or {}
+    return {
+        "title": (post.get("title") or payload.get("title") or "").strip() or None,
+        "platform": src.get("platform") or payload.get("platform") or "",
+        "url": src.get("source_url") or post.get("link") or None,
+        "like_count": post.get("like_count"),
+        "publish_timestamp": post.get("publish_timestamp") or "",
+        "excerpt": (post.get("body_text") or "")[:500],
+    }
+
+
+def iter_cases_from_db(query_ids=None, limit=None, case_ids=None, all_destructed=False):
+    """产出 (case_id, source_data, tools, version)。
+
+    候选 case 集合(base)按口径二选一:
+      all_destructed=True  → fetch_destructed_tools_cases(query_ids):所有已解构(**不看采纳**)。
+      all_destructed=False → 各 query 的 fetch_adopted_tools_cases 并集:已解构且已采纳。
+    query_ids 为多值列表(跨 query 时按 case_id 去重);None=全量。
+    再逐个 fetch_tools 重建解构详情(取最新版本)。version 用于上传去重。
+    case_ids 不为空时,在 base 集内精确筛选(点名 case 不在集内则告警跳过,不绕过口径)。
+    """
+    if all_destructed:
+        base = db.fetch_destructed_tools_cases(query_ids)
+    elif query_ids:
+        base = sorted({c for q in query_ids for c in db.fetch_adopted_tools_cases(q)})
+    else:
+        base = db.fetch_adopted_tools_cases(None)
+    if case_ids:
+        base_set = set(base)
+        missing = [c for c in case_ids if c not in base_set]
+        if missing:
+            gate = "已解构集" if all_destructed else "采纳集"
+            logger.warning("--case-ids 中 %d 个不在%s(跳过):%s",
+                           len(missing), gate, ",".join(missing))
+        base = [c for c in case_ids if c in base_set]
+    if limit:
+        base = base[:limit]
+    for case_id in base:
+        payload = db.fetch_tools(case_id)   # 最新版本
+        if not payload:
+            continue
+        yield (case_id, _build_source_data(payload),
+               (payload.get("tools") or []), payload.get("version"))
+
+
+# ── 作用域提取(仅实质/形式,复用 import_tools_knowledge 口径)──────────────────────
+
+def _split_values(raw):
+    """按顿号分割,括号内的顿号不作为分隔符,结果去重保序。"""
+    parts, current, depth = [], [], 0
+    for ch in raw:
+        if ch in ("(", "("):
+            depth += 1
+            current.append(ch)
+        elif ch in (")", ")"):
+            depth -= 1
+            current.append(ch)
+        elif ch == "、" and depth == 0:
+            part = "".join(current).strip()
+            if part:
+                parts.append(part)
+            current = []
+        else:
+            current.append(ch)
+    part = "".join(current).strip()
+    if part:
+        parts.append(part)
+
+    seen, result = set(), []
+    for p in parts:
+        if p not in seen:
+            seen.add(p)
+            result.append(p)
+    return result
+
+
+def _iter_scope_items(raw):
+    """实质/形式作用域在 mode_tools 里存为 JSON 数组(也兼容旧的顿号串):
+    先归一为字符串列表,再对每个元素按顿号拆分。逐个产出拆分后的标量值。"""
+    if raw is None:
+        return
+    items = raw if isinstance(raw, list) else [raw]
+    for item in items:
+        text = str(item).strip()
+        if not text:
+            continue
+        for value in _split_values(text):
+            yield value
+
+
+def build_scopes(tool):
+    """只从工具的实质作用域(substance)/形式作用域(form)各自去重,返回 scope 列表。
+    这是本脚本与 import_tools_knowledge 的关键差异:不采 effect、不带 custom_ext。"""
+    seen_sub, seen_form, scopes = set(), set(), []
+    for sub in _iter_scope_items(tool.get("实质作用域")):
+        if sub not in seen_sub:
+            scopes.append({"scope_type": "substance", "value": sub})
+            seen_sub.add(sub)
+    for form in _iter_scope_items(tool.get("形式作用域")):
+        if form not in seen_form:
+            scopes.append({"scope_type": "form", "value": form})
+            seen_form.add(form)
+    return scopes
+
+
+# ── 单条 payload 组装(source_id 直接用 case_id;只带 scopes,不带 custom_ext)──────────
+
+def build_payload(source_id, source_data, tool, tool_index):
+    source_title = (source_data.get("title") or "").strip()
+    tool_name = (tool.get("工具名称") or "").strip()
+
+    if tool_name:
+        knowledge_title = tool_name
+    elif source_title:
+        knowledge_title = f"{source_title} — 工具{tool_index}"
+    else:
+        knowledge_title = f"工具{tool_index}"
+
+    content = json.dumps(tool, ensure_ascii=False)
+
+    source_metadata = {
+        "platform": source_data.get("platform") or "",
+        # 链接优先取帖子来源,缺失时回退工具自身的来源链接(如 GitHub)。
+        "url": source_data.get("url") or tool.get("来源链接") or None,
+        "like_count": source_data.get("like_count"),
+        "publish_timestamp": source_data.get("publish_timestamp") or "",
+        "excerpt": source_data.get("excerpt") or "",
+        "creation_layer": tool.get("创作层级") or "",
+        "updated_time": tool.get("最新更新时间") or "",
+    }
+
+    payload = {
+        # 工具无独立 id,以 tool_index 兜底(对齐 import_process_knowledge 的
+        # "{source_id}-{procedure.id or proc_index}" 形式)。后端按 knowledge_id 幂等 upsert。
+        "knowledge_id": f"{source_id}-{tool_index}",
+        "source": {
+            "id": source_id,
+            "source_type": "post",
+            "title": source_title or None,
+            "author": None,
+            "source_metadata": source_metadata,
+        },
+        "title": knowledge_title[:512],
+        "content": content,
+        "dim_attributes": DIM_ATTRIBUTES,
+        "dim_creations": DIM_CREATIONS,
+    }
+
+    scopes = build_scopes(tool)
+    if scopes:
+        payload["scopes"] = scopes
+    # 解构知识口径:只上传实质/形式作用域,不带 custom_ext(与 import_tools_knowledge 的差异)。
+    return payload
+
+
+# ── 单条写入 ────────────────────────────────────────────────────────────────────
+
+def ingest_one(api_url, payload, dry_run):
+    """调用导入接口写入一条知识,返回 (success, info_message, knowledge_id)。"""
+    if dry_run:
+        return True, "(dry-run, skipped)", None
+
+    url = api_url.rstrip("/") + INGEST_ENDPOINT
+    try:
+        resp = requests.post(url, json=payload, timeout=30)
+        if resp.status_code == 201:
+            kid = resp.json().get("knowledge_id", "?")
+            return True, f"knowledge_id={kid}", kid
+        try:
+            detail = resp.json().get("detail", resp.text[:300])
+        except Exception:
+            detail = resp.text[:300]
+        return False, f"HTTP {resp.status_code}: {detail}", None
+    except requests.Timeout:
+        return False, "超时(30s)", None
+    except requests.RequestException as exc:
+        return False, str(exc), None
+
+
+# ── --dump 落盘(把上传入参写到 stages/output/destruction_knowledge/<时间>/)──────────
+
+def _open_dump_dir():
+    """按时间新建落盘目录并返回 Path。datetime 取本地时间,秒级足够区分多次运行。"""
+    stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
+    dump_dir = DUMP_ROOT / stamp
+    dump_dir.mkdir(parents=True, exist_ok=True)
+    return dump_dir
+
+
+def _dump_payload(dump_dir, case_id, tool_index, payload):
+    """单条 payload 落盘为 <case_id>__tool{N}.json,返回文件名。"""
+    fname = f"{case_id}__tool{tool_index}.json"
+    (dump_dir / fname).write_text(
+        json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
+    return fname
+
+
+# ── 主循环 ────────────────────────────────────────────────────────────────────
+
+def run(api_url, dry_run, verbose, delay_ms, query_ids, limit, force,
+        dump=False, case_ids=None, all_destructed=False):
+    cases = list(iter_cases_from_db(query_ids, limit, case_ids, all_destructed))
+    if not cases:
+        parts = []
+        if query_ids:
+            parts.append(f"query_id={','.join(query_ids)}")
+        if case_ids:
+            parts.append(f"case_ids={','.join(case_ids)}")
+        scope = f"({'; '.join(parts)})" if parts else ""
+        gate = "有工具解构" if all_destructed else "已采纳且有工具解构"
+        logger.error("DB 中未发现任何「%s」的 case %s", gate, scope)
+        sys.exit(1)
+
+    # --dump 隐含 dry-run(只组装+落盘,绝不打接口),便于先核对入参。
+    if dump:
+        dry_run = True
+    dump_dir = _open_dump_dir() if dump else None
+    dump_index = []
+
+    tags = "".join(t for t in ("  [DRY-RUN]" if dry_run and not dump else "",
+                               "  [DUMP]" if dump else "",
+                               "  [FORCE]" if force else ""))
+    gate = "已解构" if all_destructed else "采纳"
+    logger.info("发现 %d 个%s case。目标接口:%s%s", len(cases), gate, api_url, tags)
+    if dump:
+        logger.info("落盘目录:%s", dump_dir)
+
+    ok_count = fail_count = skip_count = dup_count = 0
+
+    for case_id, source_data, tools, version in cases:
+        logger.info("── %s", case_id)
+        if not tools:
+            logger.warning("  无 tools,跳过")
+            skip_count += 1
+            continue
+        # 去重台账:该 case 各工具已导入的版本(解构知识独立表)。force 时清空(强制重导)。
+        # dump 只核对入参,同样读台账,把「是否会被跳过」标进 _index.json(但不真跳过、全量落盘)。
+        ingested = {} if force else db.fetch_destruction_ingested_map(case_id)
+        logger.info("  source_id=%-45s  tools=%d  version=%s", case_id, len(tools), version)
+
+        for idx, tool in enumerate(tools, 1):
+            would_skip = (not force) and ingested.get(idx) == version
+
+            # 真实/普通 dry-run:已导入且版本未变 → 跳过。dump 例外:全量落盘,只在 index 标记。
+            if would_skip and not dump:
+                logger.info("  ♻️ [%d/%d] 已导入(版本 %s),跳过", idx, len(tools), version)
+                dup_count += 1
+                continue
+
+            payload = build_payload(case_id, source_data, tool, idx)
+            title = payload["title"]
+            n_scopes = len(payload.get("scopes", []))
+
+            if dump:
+                fname = _dump_payload(dump_dir, case_id, idx, payload)
+                dump_index.append({
+                    "case_id": case_id, "tool_index": idx, "version": version,
+                    "title": title, "n_scopes": n_scopes,
+                    "would_skip_by_dedup": would_skip, "file": fname,
+                })
+                if would_skip:
+                    dup_count += 1
+                logger.info("  📝 [%d/%d] title=%r  scopes=%d%s  → %s",
+                            idx, len(tools), title[:40], n_scopes,
+                            "  (去重会跳过)" if would_skip else "", fname)
+                ok_count += 1
+                continue
+
+            if dry_run and verbose:
+                print(f"\n{'=' * 60}")
+                print(f"[{case_id}] 工具 {idx}/{len(tools)}")
+                print(json.dumps(payload, ensure_ascii=False, indent=2))
+
+            ok, msg, kid = ingest_one(api_url, payload, dry_run)
+            status_icon = "✓" if ok else "✗"
+            level = logging.INFO if ok else logging.WARNING
+            logger.log(level, "  %s [%d/%d] title=%r  scopes=%d  %s",
+                       status_icon, idx, len(tools), title[:40], n_scopes, msg)
+
+            if ok:
+                ok_count += 1
+                # 仅真实导入成功才记台账(dry-run 不写,免污染去重状态)
+                if not dry_run:
+                    try:
+                        db.mark_destruction_ingested(case_id, idx, version, kid, api_url)
+                    except Exception as exc:
+                        logger.warning("  ⚠️ 台账写入失败(不影响本次导入):%s", exc)
+            else:
+                fail_count += 1
+
+            if delay_ms > 0 and not dry_run:
+                time.sleep(delay_ms / 1000)
+
+    if dump:
+        (dump_dir / "_index.json").write_text(
+            json.dumps({
+                "generated_at": datetime.now().isoformat(timespec="seconds"),
+                "api_url": api_url, "query_ids": query_ids, "case_ids": case_ids,
+                "all_destructed": all_destructed,
+                "dim_attributes": DIM_ATTRIBUTES, "dim_creations": DIM_CREATIONS,
+                "case_count": len(cases), "payload_count": len(dump_index),
+                "would_skip_by_dedup": dup_count, "items": dump_index,
+            }, ensure_ascii=False, indent=2), encoding="utf-8")
+        logger.info("落盘完成。共 %d 条入参 → %s(其中 %d 条按去重台账本会跳过)",
+                    len(dump_index), dump_dir, dup_count)
+        return
+
+    logger.info(
+        "完成。成功=%d  失败=%d  无工具跳过=%d  已传过跳过=%d  合计导入=%d",
+        ok_count, fail_count, skip_count, dup_count, ok_count,
+    )
+    if fail_count:
+        sys.exit(1)
+
+
+# ── CLI ───────────────────────────────────────────────────────────────────────
+
+def main():
+    parser = argparse.ArgumentParser(
+        description="把 DB 中已采纳的工具解构(mode_tools)按「解构知识」口径(仅实质/形式→what-制作解构)批量导入",
+        formatter_class=argparse.RawDescriptionHelpFormatter,
+    )
+    parser.add_argument("--api-url", default=DEFAULT_API_URL, metavar="URL",
+                        help=f"后端 API 根地址(默认:{DEFAULT_API_URL})")
+    parser.add_argument("--dry-run", action="store_true",
+                        help="仅从 DB 取数并组装 payload,不实际调用接口")
+    parser.add_argument("--dump", action="store_true",
+                        help="组装并把上传入参落盘到 stages/output/destruction_knowledge/<时间>/(隐含 dry-run,不调接口)")
+    parser.add_argument("--verbose", "-v", action="store_true",
+                        help="dry-run 时打印完整 payload JSON")
+    parser.add_argument("--delay", type=int, default=100, metavar="MS",
+                        help="两次 API 调用之间的间隔毫秒数(默认:100)")
+    parser.add_argument("--query-id", "--query-ids", dest="query_id", default=None, metavar="Q1,Q2",
+                        help="只导入这些搜索任务(query_id)下的 case,逗号分隔可多值")
+    parser.add_argument("--all-destructed", action="store_true",
+                        help="旁路采纳门槛:取 query 下**所有已解构**的 case(含未采纳),按 mode_tools.query_id 过滤")
+    parser.add_argument("--case-ids", default=None, metavar="C1,C2",
+                        help="只导入指定 case_id(逗号分隔);仍受当前口径约束(采纳/已解构),不在集内的跳过")
+    parser.add_argument("--limit", type=int, default=None, metavar="N",
+                        help="只处理前 N 个 case(调试用)")
+    parser.add_argument("--force", action="store_true",
+                        help="忽略去重台账,强制重导(换 prompt/模型、需覆盖时用)")
+
+    args = parser.parse_args()
+    case_ids = ([c.strip() for c in args.case_ids.split(",") if c.strip()]
+                if args.case_ids else None)
+    query_ids = ([q.strip() for q in args.query_id.split(",") if q.strip()]
+                 if args.query_id else None)
+    run(
+        api_url=args.api_url,
+        dry_run=args.dry_run,
+        verbose=args.verbose,
+        delay_ms=args.delay,
+        query_ids=query_ids,
+        limit=args.limit,
+        force=args.force,
+        dump=args.dump,
+        case_ids=case_ids,
+        all_destructed=args.all_destructed,
+    )
+
+
+if __name__ == "__main__":
+    main()

+ 453 - 0
examples/mode_workflow/stages/import_knowledge.py

@@ -0,0 +1,453 @@
+"""
+把「维度提取表」源数据(input/dimension_extraction_forms.json)拆解成知识导入接口的
+入参并上传。每个维度(姿势 / 身材体型 / 服装单品 / …)对应一条知识。
+
+字段映射(对齐示例 curl):
+  source.id          ← 随机数(src_ + 13 位随机串,每次运行现生成)
+  source.source_type ← 固定 "other"
+  source.title/author← 固定 null
+  title              ← "{维度}怎么解构"          (如 姿势 → "姿势怎么解构")
+  content            ← 把该维度全部字段(别名/推荐结果形态/提取方式/提取指令要点/通过标准/
+                        容忍/否决项/禁止/说明/视频提取)用大模型「变成人话」,重点突出
+                        【推荐结果形态】与【提取方式】;别名等关键词与推荐形态/方式原样保留。
+  dim_attributes     ← 固定 ["what-制作解构"]   dim_creations ← 固定 ["解构"]
+  custom_ext         ← 由「别名」逐个生成 {key, type:"str", value}(key=value=别名)
+
+LLM 润色(默认开启):借助大模型把 title 与 content 写成自然、专业、好读的人话。
+  长字段(指令要点/通过标准/否决项…)是「知识提取说明」,允许改写成人话;但别名关键词、
+  推荐结果形态、提取方式必须原样出现(校验器兜底,缺则重试),不丢关键信息。
+  关掉润色用 --no-enhance(只用结构化基线正文)。
+
+输出:每次运行在 stages/output/ 下按时间戳建一个子文件夹(如 output/20260709_144530/),
+      单维度一个 JSON + 汇总 _all.json,方便按时间回看每次执行的拆解结果。
+
+三条执行指令:
+  # 1) 拆解入参 + 大模型润色,结果写 output/时间戳/,并打印完整 payload
+  python stages/import_knowledge.py --dry-run -v
+  # 1') 不想看打印就去掉 -v:一样润色 + 落时间戳文件夹,只是不刷屏
+  python stages/import_knowledge.py --dry-run
+
+  # 2) 单个上传(--only 传维度名,可逗号分隔多个)
+  python stages/import_knowledge.py --only 姿势
+
+  # 3) 批量上传(默认全量,input 里所有维度)
+  python stages/import_knowledge.py
+
+  其它:--input <源文件>  --output <目录>  --api-url <后端根地址>
+        --from-output(直接传 output 里已审阅的 payload,不重建/不润色)
+        --no-enhance(关闭润色)  --model <模型>  --delay <毫秒,调用间隔>
+"""
+
+import argparse
+import asyncio
+import json
+import logging
+import random
+import string
+import sys
+import time
+from datetime import datetime
+from pathlib import Path
+
+import requests
+
+# stages/ 归档在子目录;补 Agent 根到 sys.path 以 import agent.llm / examples.*
+PROJECT_ROOT = Path(__file__).resolve().parents[3]   # …/Agent
+sys.path.insert(0, str(PROJECT_ROOT))
+
+from dotenv import load_dotenv
+load_dotenv()
+
+# ── 配置(对齐示例 curl)──────────────────────────────────────────────────────
+
+HERE = Path(__file__).resolve().parent
+DEFAULT_INPUT = HERE / "input" / "dimension_extraction_forms.json"
+DEFAULT_OUTPUT = HERE / "output"
+
+DEFAULT_API_URL = "http://47.236.83.130:8001"
+INGEST_ENDPOINT = "/api/v1/knowledge/ingest"
+
+DIM_ATTRIBUTES = ["what-制作解构"]
+DIM_CREATIONS = ["解构"]
+
+KEY_ALIAS = "别名"
+KEY_FORM = "推荐结果形态"      # 图 #3 圈中之一:重点突出
+KEY_METHOD = "提取方式"        # 图 #3 圈中之一:重点突出
+KEY_VIDEO = "视频提取"
+
+# content 里按此顺序铺陈的字段(视频提取单独嵌套处理)
+CONTENT_FIELDS = [
+    KEY_ALIAS, KEY_FORM, KEY_METHOD,
+    "提取指令要点", "通过标准", "容忍", "否决项", "禁止", "说明",
+]
+
+DEFAULT_MODEL = "anthropic/claude-sonnet-4-6"
+
+# ── 日志 ──────────────────────────────────────────────────────────────────────
+
+logging.basicConfig(
+    level=logging.INFO,
+    format="%(asctime)s [%(levelname)s] %(message)s",
+    datefmt="%H:%M:%S",
+)
+logger = logging.getLogger(__name__)
+
+
+# ── 组装 ──────────────────────────────────────────────────────────────────────
+
+def gen_source_id():
+    """随机 source.id,形如示例的 src_mrd37w05ppua7(src_ + 13 位小写数字串)。"""
+    rand = "".join(random.choices(string.ascii_lowercase + string.digits, k=13))
+    return f"src_{rand}"
+
+
+def _rand_num(k=8):
+    """一个随机数(k 位数字串),用于拼 knowledge_id。"""
+    return "".join(random.choices(string.digits, k=k))
+
+
+def _fmt_value(v):
+    if isinstance(v, list):
+        return "、".join(str(x) for x in v)
+    return str(v).strip()
+
+
+def build_content(name, data):
+    """结构化基线正文——列出该维度全部字段(润色前的草稿,也是 --no-enhance 的产物)。"""
+    lines = [f"「{name}」的知识提取维度说明:"]
+    for field in CONTENT_FIELDS:
+        if field in data and data[field] not in (None, "", []):
+            lines.append(f"- {field}:{_fmt_value(data[field])}")
+    video = data.get(KEY_VIDEO)
+    if isinstance(video, dict) and video:
+        lines.append("- 视频提取:")
+        for k, v in video.items():
+            if v not in (None, "", []):
+                lines.append(f"    · {k}:{_fmt_value(v)}")
+    return "\n".join(lines)
+
+
+def build_custom_ext(data):
+    """custom_ext ← 别名 列表,逐个 {key, type:"str", value}(key=value=别名,保序去重)。"""
+    ext, seen = [], set()
+    for alias in (data.get(KEY_ALIAS) or []):
+        alias = (alias or "").strip()
+        if alias and alias not in seen:
+            seen.add(alias)
+            ext.append({"key": alias, "type": "str", "value": alias})
+    return ext
+
+
+def build_payload(name, data):
+    sid = gen_source_id()
+    return {
+        "knowledge_id": f"{sid}-{_rand_num()}",   # 外层,与 source 同级;= source.id + 一个随机数
+        "source": {
+            "id": sid,
+            "source_type": "other",
+            "title": None,
+            "author": None,
+        },
+        "title": f"{name}怎么解构",
+        "content": build_content(name, data),
+        "dim_attributes": DIM_ATTRIBUTES,
+        "dim_creations": DIM_CREATIONS,
+        "custom_ext": build_custom_ext(data),
+    }
+
+
+# ── LLM 润色(长字段可改写成人话;关键词/推荐形态/提取方式原样保留)────────────────
+
+ENHANCE_SYSTEM = (
+    "你是知识库文案润色器。给你一个「内容制作」里某个解构维度的结构化提取说明(维度名、别名、"
+    "推荐结果形态、提取方式,以及提取指令要点/通过标准/容忍/否决项/禁止/说明/视频提取等长字段)与一版基线草稿,"
+    "请把它整理成一条自然、专业、好读的『该维度怎么解构』知识条目,输出 title 和 content。\n\n"
+    "写作要求:\n"
+    "1. 把冗长的指令/标准/否决项等用人话讲清楚(可改写措辞、合并同类、分点),让人一看就懂"
+    "『这个维度是什么、该提取成什么、怎么提取、什么算通过、什么会被否决』。\n"
+    "2. 重点突出【推荐结果形态】和【提取方式】——这是该维度最核心的两点,要在正文里清楚点出。\n"
+    "3. 必须保留的关键信息(原样出现,不得改写/翻译/漏掉):全部【别名】词、全部【推荐结果形态】、【提取方式】。\n"
+    "4. 若有【视频提取】,也用一小段说清楚(载体/工具/结果形态)。\n"
+    "5. title 用一句话:『{维度}怎么解构』或语义等价的简洁专业说法(需含维度名与「解构」)。\n"
+    '只输出一个 JSON 对象: {"title": "...", "content": "..."},不要任何额外文字或解释。'
+)
+
+
+def _build_enhance_user(name, data, base):
+    return (
+        f"【维度】{name}\n"
+        f"【结构化数据(原始)】\n{json.dumps(data, ensure_ascii=False, indent=2)}\n\n"
+        f"【基线草稿 title】{base['title']}\n"
+        f"【基线草稿 content】\n{base['content']}\n\n"
+        "请输出润色后的 JSON。"
+    )
+
+
+# 关键词包含判断前的规范化:抹平全角/半角括号等标点差异(LLM 常把 (线框) 写成(线框)),
+# 只用于「是否保留了该词」的校验,不改动实际落盘正文。
+_PUNCT_NORM = str.maketrans({
+    "(": "(", ")": ")", "[": "[", "]": "]", "{": "{", "}": "}",
+    "/": "/", " ": " ",
+})
+
+
+def _norm(s):
+    return s.translate(_PUNCT_NORM)
+
+
+def _make_validator(name, data):
+    aliases = [t.strip() for t in (data.get(KEY_ALIAS) or []) if t and t.strip()]
+    forms = [t.strip() for t in (data.get(KEY_FORM) or []) if t and t.strip()]
+    method = (data.get(KEY_METHOD) or "").strip()
+
+    def _v(d):
+        if not isinstance(d, dict):
+            return "需 JSON 对象"
+        title = (d.get("title") or "").strip()
+        content = (d.get("content") or "").strip()
+        if not title:
+            return "title 缺失"
+        if name not in title or "解构" not in title:
+            return f"title 需含维度名「{name}」与「解构」"
+        if not content:
+            return "content 缺失"
+        nc = _norm(content)
+        miss_alias = [t for t in aliases if _norm(t) not in nc]
+        if miss_alias:
+            return "content 丢了别名(须原样保留):" + "、".join(miss_alias[:10])
+        miss_form = [t for t in forms if _norm(t) not in nc]
+        if miss_form:
+            return "content 缺推荐结果形态(须原样保留):" + "、".join(miss_form)
+        if method and _norm(method) not in nc:
+            return f"content 缺提取方式「{method}」(须原样保留)"
+        return None
+    return _v
+
+
+async def _enhance_one(name, data, payload, llm_call, model, sem):
+    from examples.process_pipeline.script.llm_helper import call_llm_with_retry
+    base = {"title": payload["title"], "content": payload["content"]}
+    messages = [{"role": "system", "content": ENHANCE_SYSTEM},
+                {"role": "user", "content": _build_enhance_user(name, data, base)}]
+    async with sem:
+        d, cost = await call_llm_with_retry(
+            llm_call=llm_call, messages=messages, model=model,
+            temperature=0.3, max_tokens=2500,
+            validate_fn=_make_validator(name, data), task_name=f"Enhance[{name}]")
+    if d:   # 校验通过才覆盖;失败则沿用基线草稿(不丢数据)
+        payload["title"] = d["title"].strip()
+        payload["content"] = d["content"].strip()
+        return cost, True
+    logger.warning("  ⚠️ [%s] 润色失败,沿用基线草稿", name)
+    return cost, False
+
+
+async def enhance_all(items, payloads, model, concurrency=5):
+    from agent.llm.openrouter import create_openrouter_llm_call
+    llm_call = create_openrouter_llm_call(model=model)
+    sem = asyncio.Semaphore(concurrency)
+    results = await asyncio.gather(*[
+        _enhance_one(name, data, payloads[i], llm_call, model, sem)
+        for i, (name, data) in enumerate(items)])
+    total_cost = sum(c for c, _ in results)
+    ok = sum(1 for _, good in results if good)
+    logger.info("🪄 LLM 润色完成:%d/%d 成功  成本 $%.4f", ok, len(items), total_cost)
+
+
+# ── 输出落盘(按时间戳建子文件夹)──────────────────────────────────────────────
+
+def write_outputs(output_base, items, payloads):
+    ts = datetime.now().strftime("%Y%m%d_%H%M%S")
+    run_dir = Path(output_base) / ts
+    run_dir.mkdir(parents=True, exist_ok=True)
+    combined = []
+    for (name, _), payload in zip(items, payloads):
+        (run_dir / f"{name}.json").write_text(
+            json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
+        combined.append({"category": name, "payload": payload})
+    (run_dir / "_all.json").write_text(
+        json.dumps(combined, ensure_ascii=False, indent=2), encoding="utf-8")
+    logger.info("📁 已写出 %d 个 payload → %s", len(payloads), run_dir)
+    return run_dir
+
+
+def resolve_read_dir(output_base):
+    """--from-output 读取目录:优先取 output_base 下最新的时间戳子文件夹;
+    若 output_base 直接就是含 {name}.json 的目录(用户显式指定某次运行)则用它本身。"""
+    base = Path(output_base)
+    if not base.is_dir():
+        logger.error("output 目录不存在:%s(先跑一次 --dry-run 生成)", base)
+        sys.exit(1)
+    subs = [d for d in base.iterdir()
+            if d.is_dir() and any(p.stem != "_all" for p in d.glob("*.json"))]
+    if subs:
+        return max(subs, key=lambda d: d.name)   # 时间戳名可按字典序取最新
+    if any(p.stem != "_all" for p in base.glob("*.json")):
+        return base
+    logger.error("output 下未找到任何 payload:%s", base)
+    sys.exit(1)
+
+
+def load_payloads_from_output(output_base, only):
+    """从 output 读回已审阅的 payload(每维度一个 {name}.json)。返回 [(name, payload)]。"""
+    read_dir = resolve_read_dir(output_base)
+    logger.info("📂 从已审阅目录读取:%s", read_dir)
+    if only:
+        names = [c.strip() for c in only if c.strip()]
+    else:
+        names = sorted(p.stem for p in read_dir.glob("*.json") if p.stem != "_all")
+    items = []
+    for name in names:
+        pf = read_dir / f"{name}.json"
+        if not pf.is_file():
+            logger.warning("  ⚠️ 该目录缺 %s.json,跳过", name)
+            continue
+        items.append((name, json.loads(pf.read_text(encoding="utf-8"))))
+    return items
+
+
+# ── 单条写入 ──────────────────────────────────────────────────────────────────
+
+def ingest_one(api_url, payload, dry_run):
+    """调用导入接口写入一条知识,返回 (success, info_message, knowledge_id)。"""
+    if dry_run:
+        return True, "(dry-run, skipped)", None
+
+    url = api_url.rstrip("/") + INGEST_ENDPOINT
+    try:
+        resp = requests.post(url, json=payload, timeout=30)
+        if resp.status_code in (200, 201):
+            try:
+                kid = resp.json().get("knowledge_id", "?")
+            except Exception:
+                kid = "?"
+            return True, f"knowledge_id={kid}", kid
+        try:
+            detail = resp.json().get("detail", resp.text[:300])
+        except Exception:
+            detail = resp.text[:300]
+        return False, f"HTTP {resp.status_code}: {detail}", None
+    except requests.Timeout:
+        return False, "超时(30s)", None
+    except requests.RequestException as exc:
+        return False, str(exc), None
+
+
+# ── 主循环 ────────────────────────────────────────────────────────────────────
+
+def load_categories(input_path, only=None):
+    """读源 JSON,返回有序 [(name, data)]。only 非空时按维度名精确筛选。"""
+    with open(input_path, "r", encoding="utf-8") as f:
+        raw = json.load(f)
+    items = list(raw.items())
+    if only:
+        want = [c.strip() for c in only if c.strip()]
+        missing = [c for c in want if c not in raw]
+        if missing:
+            logger.warning("--only 中 %d 个维度不存在(跳过):%s", len(missing), ",".join(missing))
+        items = [(n, d) for n, d in items if n in want]
+    return items
+
+
+def run(input_path, output_dir, api_url, dry_run, verbose, delay_ms, only, enhance, model,
+        from_output=False):
+    mode_tag = "  [DRY-RUN]" if dry_run else ""
+
+    if from_output:
+        # 直接上传 output 里已审阅的 payload,不重建/不润色(保证上传=所见)
+        loaded = load_payloads_from_output(output_dir, only)
+        if not loaded:
+            logger.error("output 中未匹配到任何 payload"); sys.exit(1)
+        items = [(name, None) for name, _ in loaded]
+        payloads = [p for _, p in loaded]
+        logger.info("发现 %d 个已审阅 payload。目标接口:%s%s", len(payloads), api_url, mode_tag)
+    else:
+        items = load_categories(input_path, only)
+        if not items:
+            logger.error("源文件中未匹配到任何维度:%s%s", input_path,
+                          f"(--only {','.join(only)})" if only else "")
+            sys.exit(1)
+        enh_tag = "  [ENHANCE]" if enhance else "  [NO-ENHANCE]"
+        logger.info("发现 %d 个维度。源:%s  目标接口:%s%s%s",
+                    len(items), input_path, api_url, mode_tag, enh_tag)
+        payloads = [build_payload(name, data) for name, data in items]
+        if enhance:
+            asyncio.run(enhance_all(items, payloads, model))
+        # 落盘(dry-run 与真实上传都留一份记录,按时间戳分文件夹便于回看)
+        write_outputs(output_dir, items, payloads)
+
+    ok_count = fail_count = 0
+    for (name, _), payload in zip(items, payloads):
+        n_ext = len(payload["custom_ext"])
+        if dry_run and verbose:
+            print(f"\n{'=' * 60}")
+            print(f"[{name}] source.id={payload['source']['id']}")
+            print(json.dumps(payload, ensure_ascii=False, indent=2))
+
+        ok, msg, _ = ingest_one(api_url, payload, dry_run)
+        icon = "✓" if ok else "✗"
+        level = logging.INFO if ok else logging.WARNING
+        logger.log(level, "  %s %-8s title=%r  ext=%d  %s",
+                   icon, name, payload["title"], n_ext, msg)
+
+        if ok:
+            ok_count += 1
+        else:
+            fail_count += 1
+        if delay_ms > 0 and not dry_run:
+            time.sleep(delay_ms / 1000)
+
+    logger.info("完成。成功=%d  失败=%d  合计=%d", ok_count, fail_count, ok_count)
+    if fail_count:
+        sys.exit(1)
+
+
+# ── CLI ───────────────────────────────────────────────────────────────────────
+
+def main():
+    parser = argparse.ArgumentParser(
+        description="把维度提取表源数据拆解成知识接口入参(默认大模型润色)并上传",
+        formatter_class=argparse.RawDescriptionHelpFormatter,
+    )
+    parser.add_argument("--input", default=str(DEFAULT_INPUT), metavar="PATH",
+                        help=f"源 JSON 路径(默认:{DEFAULT_INPUT})")
+    parser.add_argument("--output", default=str(DEFAULT_OUTPUT), metavar="DIR",
+                        help=f"payload 落盘根目录,每次运行下建时间戳子文件夹(默认:{DEFAULT_OUTPUT})")
+    parser.add_argument("--api-url", default=DEFAULT_API_URL, metavar="URL",
+                        help=f"后端 API 根地址(默认:{DEFAULT_API_URL})")
+    parser.add_argument("--dry-run", action="store_true",
+                        help="仅拆解 + 组装(+ 润色 + 落盘),不实际调用接口")
+    parser.add_argument("--verbose", "-v", action="store_true",
+                        help="打印完整 payload JSON(不影响是否润色)")
+    parser.add_argument("--only", default=None, metavar="C1,C2",
+                        help="只处理指定维度(逗号分隔维度名,如 姿势,表情)")
+    parser.add_argument("--from-output", action="store_true",
+                        help="直接上传 output 最新时间戳目录里已审阅的 payload(不重建/不润色)")
+    parser.add_argument("--enhance", dest="enhance", action="store_true", default=True,
+                        help="用大模型润色 title/content(默认开启)")
+    parser.add_argument("--no-enhance", dest="enhance", action="store_false",
+                        help="关闭润色,只用结构化基线正文")
+    parser.add_argument("--model", default=DEFAULT_MODEL, metavar="MODEL",
+                        help=f"润色用模型(默认:{DEFAULT_MODEL})")
+    parser.add_argument("--delay", type=int, default=100, metavar="MS",
+                        help="两次 API 调用之间的间隔毫秒数(默认:100)")
+
+    args = parser.parse_args()
+    only = ([c.strip() for c in args.only.split(",") if c.strip()]
+            if args.only else None)
+    run(
+        input_path=args.input,
+        output_dir=args.output,
+        api_url=args.api_url,
+        dry_run=args.dry_run,
+        verbose=args.verbose,
+        delay_ms=args.delay,
+        only=only,
+        enhance=args.enhance,
+        model=args.model,
+        from_output=args.from_output,
+    )
+
+
+if __name__ == "__main__":
+    main()

+ 413 - 0
examples/mode_workflow/stages/import_knowledge_category.py

@@ -0,0 +1,413 @@
+"""
+把「品类聚焦维度」源数据(input/category_focus_dimensions.json)拆解成知识导入接口的
+入参并上传。每个品类(美食 / 穿搭服饰 / …)对应一条知识。
+
+—— 与 import_knowledge.py 的分工:
+     import_knowledge.py           处理 dimension_extraction_forms.json(维度提取表,title「怎么解构」)
+     import_knowledge_category.py  处理 category_focus_dimensions.json(品类聚焦维度,title「怎么分类」) ← 本文件
+
+字段映射(对齐示例 curl):
+  source.id          ← 随机数(src_ + 13 位随机串,每次运行现生成)
+  source.knowledge_id← source.id + 一个随机数
+  source.source_type ← 固定 "other"
+  source.title/author← 固定 null
+  title              ← "{品类}怎么分类"        (如 美食 → "美食怎么分类")
+  content            ← 把 别名/重点维度/提示 用大模型「变成人话」(数据词原样保留,只美化排版)
+  dim_attributes     ← 固定 ["what-制作解构"]   dim_creations ← 固定 ["解构"]
+  custom_ext         ← 由「别名」逐个生成 {key, type:"str", value}(key=value=别名)
+
+LLM 润色(默认开启):把 title/content 写成自然、专业、好读的人话。硬约束:别名、重点维度里
+  每个词、提示原文都是【数据】,必须原样保留(校验器逐词兜底,丢词即重试)。--no-enhance 关闭。
+
+输出:每次运行在 stages/output_category/ 下按时间戳建子文件夹,单品类一个 JSON + 汇总 _all.json。
+
+三条执行指令:
+  python stages/import_knowledge_category.py --dry-run -v   # 拆解 + 润色 + 落时间戳文件夹 + 打印
+  python stages/import_knowledge_category.py --dry-run      # 同上但不刷屏
+  python stages/import_knowledge_category.py --only 美食     # 单个上传
+  python stages/import_knowledge_category.py                # 批量上传(全量)
+  python stages/import_knowledge_category.py --from-output  # 上传最新时间戳目录里已审阅的 payload
+"""
+
+import argparse
+import asyncio
+import json
+import logging
+import random
+import string
+import sys
+import time
+from datetime import datetime
+from pathlib import Path
+
+import requests
+
+PROJECT_ROOT = Path(__file__).resolve().parents[3]   # …/Agent
+sys.path.insert(0, str(PROJECT_ROOT))
+
+from dotenv import load_dotenv
+load_dotenv()
+
+# ── 配置 ──────────────────────────────────────────────────────────────────────
+
+HERE = Path(__file__).resolve().parent
+DEFAULT_INPUT = HERE / "input" / "category_focus_dimensions.json"
+DEFAULT_OUTPUT = HERE / "output_category"
+
+DEFAULT_API_URL = "http://47.236.83.130:8001"
+INGEST_ENDPOINT = "/api/v1/knowledge/ingest"
+
+DIM_ATTRIBUTES = ["what-制作解构"]
+DIM_CREATIONS = ["解构"]
+
+KEY_ALIAS = "别名"
+KEY_DIMENSIONS = "重点维度"
+KEY_TIP = "提示"
+
+DEFAULT_MODEL = "anthropic/claude-sonnet-4-6"
+
+# ── 日志 ──────────────────────────────────────────────────────────────────────
+
+logging.basicConfig(
+    level=logging.INFO,
+    format="%(asctime)s [%(levelname)s] %(message)s",
+    datefmt="%H:%M:%S",
+)
+logger = logging.getLogger(__name__)
+
+
+# ── 组装 ──────────────────────────────────────────────────────────────────────
+
+def gen_source_id():
+    rand = "".join(random.choices(string.ascii_lowercase + string.digits, k=13))
+    return f"src_{rand}"
+
+
+def _rand_num(k=8):
+    return "".join(random.choices(string.digits, k=k))
+
+
+def build_content(name, data):
+    """把 别名/重点维度/提示 美化成人话正文(润色前的基线,也是 --no-enhance 的产物)。"""
+    aliases = data.get(KEY_ALIAS) or []
+    dimensions = data.get(KEY_DIMENSIONS) or []
+    tip = (data.get(KEY_TIP) or "").strip()
+
+    lines = []
+    if aliases:
+        lines.append(f"1. {name}可以通过这些别名去检索:{'、'.join(aliases)}。")
+    if dimensions:
+        lines.append(f"2. {name}可以从这些重点维度去查找:{'、'.join(dimensions)}。")
+    if tip:
+        lines.append(f"3. {tip}")
+    return "\n".join(lines)
+
+
+def build_custom_ext(data):
+    """custom_ext ← 别名 列表,逐个 {key, type:"str", value}(保序去重)。"""
+    ext, seen = [], set()
+    for alias in (data.get(KEY_ALIAS) or []):
+        alias = (alias or "").strip()
+        if alias and alias not in seen:
+            seen.add(alias)
+            ext.append({"key": alias, "type": "str", "value": alias})
+    return ext
+
+
+def build_payload(name, data):
+    sid = gen_source_id()
+    return {
+        "knowledge_id": f"{sid}-{_rand_num()}",   # 外层,与 source 同级;= source.id + 一个随机数
+        "source": {
+            "id": sid,
+            "source_type": "other",
+            "title": None,
+            "author": None,
+        },
+        "title": f"{name}怎么分类",
+        "content": build_content(name, data),
+        "dim_attributes": DIM_ATTRIBUTES,
+        "dim_creations": DIM_CREATIONS,
+        "custom_ext": build_custom_ext(data),
+    }
+
+
+# ── LLM 润色(别名/重点维度/提示 逐词原样保留)──────────────────────────────────
+
+ENHANCE_SYSTEM = (
+    "你是知识库文案润色器。给你一个内容品类的结构化信息(品类名、别名列表、重点维度列表、提示)与一版基线草稿,"
+    "请把它整理成一条更自然、更专业、更好读的知识条目。\n\n"
+    "严格约束:\n"
+    "1. 别名、重点维度里的每一个词,以及【提示】原文,都是【数据】——必须原样保留,"
+    "不得改写/翻译/增删/替换其中任何一个词,也不得漏掉任何一个。\n"
+    "2. 你只能优化句子结构与连接措辞,让正文读起来通顺、像人话、有条理。\n"
+    "3. content 必须包含:全部别名(讲清可用于检索)、全部重点维度(讲清可用于查找/解构)、以及提示原文。\n"
+    "4. title 用一句话概括『如何对该品类做内容分类』,简洁专业,不堆砌。\n"
+    '只输出一个 JSON 对象: {"title": "...", "content": "..."},不要任何额外文字或解释。'
+)
+
+_PUNCT_NORM = str.maketrans({
+    "(": "(", ")": ")", "[": "[", "]": "]", "{": "{", "}": "}", "/": "/", " ": " ",
+})
+
+
+def _norm(s):
+    return s.translate(_PUNCT_NORM)
+
+
+def _build_enhance_user(name, data, base):
+    return (
+        f"【品类】{name}\n"
+        f"【别名】{json.dumps(data.get(KEY_ALIAS) or [], ensure_ascii=False)}\n"
+        f"【重点维度】{json.dumps(data.get(KEY_DIMENSIONS) or [], ensure_ascii=False)}\n"
+        f"【提示】{(data.get(KEY_TIP) or '').strip()}\n\n"
+        f"【基线 title】{base['title']}\n"
+        f"【基线 content】\n{base['content']}\n\n"
+        "请输出润色后的 JSON。"
+    )
+
+
+def _make_validator(data):
+    terms = [t.strip() for t in (data.get(KEY_ALIAS) or []) + (data.get(KEY_DIMENSIONS) or [])
+             if t and t.strip()]
+    tip = (data.get(KEY_TIP) or "").strip()
+
+    def _v(d):
+        if not isinstance(d, dict):
+            return "需 JSON 对象"
+        title = (d.get("title") or "").strip()
+        content = (d.get("content") or "").strip()
+        if not title:
+            return "title 缺失"
+        if not content:
+            return "content 缺失"
+        nc = _norm(content)
+        missing = [t for t in terms if _norm(t) not in nc]
+        if missing:
+            return "content 丢了数据词(必须原样保留):" + "、".join(missing[:10])
+        if tip and _norm(tip) not in nc:
+            return "content 缺【提示】原文(必须原样保留)"
+        return None
+    return _v
+
+
+async def _enhance_one(name, data, payload, llm_call, model, sem):
+    from examples.process_pipeline.script.llm_helper import call_llm_with_retry
+    base = {"title": payload["title"], "content": payload["content"]}
+    messages = [{"role": "system", "content": ENHANCE_SYSTEM},
+                {"role": "user", "content": _build_enhance_user(name, data, base)}]
+    async with sem:
+        d, cost = await call_llm_with_retry(
+            llm_call=llm_call, messages=messages, model=model,
+            temperature=0.3, max_tokens=1500,
+            validate_fn=_make_validator(data), task_name=f"Enhance[{name}]")
+    if d:
+        payload["title"] = d["title"].strip()
+        payload["content"] = d["content"].strip()
+        return cost, True
+    logger.warning("  ⚠️ [%s] 润色失败,沿用基线草稿", name)
+    return cost, False
+
+
+async def enhance_all(items, payloads, model, concurrency=5):
+    from agent.llm.openrouter import create_openrouter_llm_call
+    llm_call = create_openrouter_llm_call(model=model)
+    sem = asyncio.Semaphore(concurrency)
+    results = await asyncio.gather(*[
+        _enhance_one(name, data, payloads[i], llm_call, model, sem)
+        for i, (name, data) in enumerate(items)])
+    total_cost = sum(c for c, _ in results)
+    ok = sum(1 for _, good in results if good)
+    logger.info("🪄 LLM 润色完成:%d/%d 成功  成本 $%.4f", ok, len(items), total_cost)
+
+
+# ── 输出落盘(按时间戳建子文件夹)──────────────────────────────────────────────
+
+def write_outputs(output_base, items, payloads):
+    ts = datetime.now().strftime("%Y%m%d_%H%M%S")
+    run_dir = Path(output_base) / ts
+    run_dir.mkdir(parents=True, exist_ok=True)
+    combined = []
+    for (name, _), payload in zip(items, payloads):
+        (run_dir / f"{name}.json").write_text(
+            json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
+        combined.append({"category": name, "payload": payload})
+    (run_dir / "_all.json").write_text(
+        json.dumps(combined, ensure_ascii=False, indent=2), encoding="utf-8")
+    logger.info("📁 已写出 %d 个 payload → %s", len(payloads), run_dir)
+    return run_dir
+
+
+def resolve_read_dir(output_base):
+    base = Path(output_base)
+    if not base.is_dir():
+        logger.error("output 目录不存在:%s(先跑一次 --dry-run 生成)", base)
+        sys.exit(1)
+    subs = [d for d in base.iterdir()
+            if d.is_dir() and any(p.stem != "_all" for p in d.glob("*.json"))]
+    if subs:
+        return max(subs, key=lambda d: d.name)
+    if any(p.stem != "_all" for p in base.glob("*.json")):
+        return base
+    logger.error("output 下未找到任何 payload:%s", base)
+    sys.exit(1)
+
+
+def load_payloads_from_output(output_base, only):
+    read_dir = resolve_read_dir(output_base)
+    logger.info("📂 从已审阅目录读取:%s", read_dir)
+    if only:
+        names = [c.strip() for c in only if c.strip()]
+    else:
+        names = sorted(p.stem for p in read_dir.glob("*.json") if p.stem != "_all")
+    items = []
+    for name in names:
+        pf = read_dir / f"{name}.json"
+        if not pf.is_file():
+            logger.warning("  ⚠️ 该目录缺 %s.json,跳过", name)
+            continue
+        items.append((name, json.loads(pf.read_text(encoding="utf-8"))))
+    return items
+
+
+# ── 单条写入 ──────────────────────────────────────────────────────────────────
+
+def ingest_one(api_url, payload, dry_run):
+    if dry_run:
+        return True, "(dry-run, skipped)", None
+    url = api_url.rstrip("/") + INGEST_ENDPOINT
+    try:
+        resp = requests.post(url, json=payload, timeout=30)
+        if resp.status_code in (200, 201):
+            try:
+                kid = resp.json().get("knowledge_id", "?")
+            except Exception:
+                kid = "?"
+            return True, f"knowledge_id={kid}", kid
+        try:
+            detail = resp.json().get("detail", resp.text[:300])
+        except Exception:
+            detail = resp.text[:300]
+        return False, f"HTTP {resp.status_code}: {detail}", None
+    except requests.Timeout:
+        return False, "超时(30s)", None
+    except requests.RequestException as exc:
+        return False, str(exc), None
+
+
+# ── 主循环 ────────────────────────────────────────────────────────────────────
+
+def load_categories(input_path, only=None):
+    with open(input_path, "r", encoding="utf-8") as f:
+        raw = json.load(f)
+    items = list(raw.items())
+    if only:
+        want = [c.strip() for c in only if c.strip()]
+        missing = [c for c in want if c not in raw]
+        if missing:
+            logger.warning("--only 中 %d 个品类不存在(跳过):%s", len(missing), ",".join(missing))
+        items = [(n, d) for n, d in items if n in want]
+    return items
+
+
+def run(input_path, output_dir, api_url, dry_run, verbose, delay_ms, only, enhance, model,
+        from_output=False):
+    mode_tag = "  [DRY-RUN]" if dry_run else ""
+
+    if from_output:
+        loaded = load_payloads_from_output(output_dir, only)
+        if not loaded:
+            logger.error("output 中未匹配到任何 payload"); sys.exit(1)
+        items = [(name, None) for name, _ in loaded]
+        payloads = [p for _, p in loaded]
+        logger.info("发现 %d 个已审阅 payload。目标接口:%s%s", len(payloads), api_url, mode_tag)
+    else:
+        items = load_categories(input_path, only)
+        if not items:
+            logger.error("源文件中未匹配到任何品类:%s%s", input_path,
+                          f"(--only {','.join(only)})" if only else "")
+            sys.exit(1)
+        enh_tag = "  [ENHANCE]" if enhance else "  [NO-ENHANCE]"
+        logger.info("发现 %d 个品类。源:%s  目标接口:%s%s%s",
+                    len(items), input_path, api_url, mode_tag, enh_tag)
+        payloads = [build_payload(name, data) for name, data in items]
+        if enhance:
+            asyncio.run(enhance_all(items, payloads, model))
+        write_outputs(output_dir, items, payloads)
+
+    ok_count = fail_count = 0
+    for (name, _), payload in zip(items, payloads):
+        n_ext = len(payload["custom_ext"])
+        if dry_run and verbose:
+            print(f"\n{'=' * 60}")
+            print(f"[{name}] knowledge_id={payload['knowledge_id']}  "
+                  f"source.id={payload['source']['id']}")
+            print(json.dumps(payload, ensure_ascii=False, indent=2))
+
+        ok, msg, _ = ingest_one(api_url, payload, dry_run)
+        icon = "✓" if ok else "✗"
+        level = logging.INFO if ok else logging.WARNING
+        logger.log(level, "  %s %-8s title=%r  ext=%d  %s",
+                   icon, name, payload["title"], n_ext, msg)
+
+        if ok:
+            ok_count += 1
+        else:
+            fail_count += 1
+        if delay_ms > 0 and not dry_run:
+            time.sleep(delay_ms / 1000)
+
+    logger.info("完成。成功=%d  失败=%d  合计=%d", ok_count, fail_count, ok_count)
+    if fail_count:
+        sys.exit(1)
+
+
+# ── CLI ───────────────────────────────────────────────────────────────────────
+
+def main():
+    parser = argparse.ArgumentParser(
+        description="把品类聚焦维度源数据拆解成知识接口入参(默认大模型润色)并上传",
+        formatter_class=argparse.RawDescriptionHelpFormatter,
+    )
+    parser.add_argument("--input", default=str(DEFAULT_INPUT), metavar="PATH",
+                        help=f"源 JSON 路径(默认:{DEFAULT_INPUT})")
+    parser.add_argument("--output", default=str(DEFAULT_OUTPUT), metavar="DIR",
+                        help=f"payload 落盘根目录,每次运行下建时间戳子文件夹(默认:{DEFAULT_OUTPUT})")
+    parser.add_argument("--api-url", default=DEFAULT_API_URL, metavar="URL",
+                        help=f"后端 API 根地址(默认:{DEFAULT_API_URL})")
+    parser.add_argument("--dry-run", action="store_true",
+                        help="仅拆解 + 组装(+ 润色 + 落盘),不实际调用接口")
+    parser.add_argument("--verbose", "-v", action="store_true",
+                        help="打印完整 payload JSON(不影响是否润色)")
+    parser.add_argument("--only", default=None, metavar="C1,C2",
+                        help="只处理指定品类(逗号分隔品类名,如 美食,穿搭服饰)")
+    parser.add_argument("--from-output", action="store_true",
+                        help="直接上传 output 最新时间戳目录里已审阅的 payload(不重建/不润色)")
+    parser.add_argument("--enhance", dest="enhance", action="store_true", default=True,
+                        help="用大模型润色 title/content(默认开启)")
+    parser.add_argument("--no-enhance", dest="enhance", action="store_false",
+                        help="关闭润色,只用结构化基线正文")
+    parser.add_argument("--model", default=DEFAULT_MODEL, metavar="MODEL",
+                        help=f"润色用模型(默认:{DEFAULT_MODEL})")
+    parser.add_argument("--delay", type=int, default=100, metavar="MS",
+                        help="两次 API 调用之间的间隔毫秒数(默认:100)")
+
+    args = parser.parse_args()
+    only = ([c.strip() for c in args.only.split(",") if c.strip()]
+            if args.only else None)
+    run(
+        input_path=args.input,
+        output_dir=args.output,
+        api_url=args.api_url,
+        dry_run=args.dry_run,
+        verbose=args.verbose,
+        delay_ms=args.delay,
+        only=only,
+        enhance=args.enhance,
+        model=args.model,
+        from_output=args.from_output,
+    )
+
+
+if __name__ == "__main__":
+    main()

+ 282 - 0
examples/mode_workflow/stages/ingest_urls.py

@@ -0,0 +1,282 @@
+# -*- coding: utf-8 -*-
+"""按「渠道 + 帖子链接」直接注入 → 评估 → search_tools(工具方向)
+================================================================================
+正常流程是 stages/search_eval.py 按 query 词搜索;本脚本处理**特殊情况**:
+已知一批好帖的直链(跳过搜索),直接把它们塞进工具方向的评估管线。
+
+设计:只替换最前面的「搜索」环节 —— 按 URL 直取帖子、拼成与 search_all 输出
+**同构的 source dict**,后面完全复用 search_eval 的下游(转写 → 评估 → 入库),
+核心管线一行不改。落表固定 search_tools(手工指定即工具方向,不再按评估标签路由)。
+
+按-URL 取帖的两条免鉴权通路(平台后端本身不支持「按 id 取单帖」,故绕行):
+  - gzh: aigc-channel/detail  {"type":"gzh","content_link":<完整文章URL>}
+         → 直接返回 channel_content_id/title/body_text/images/videos/link
+  - x  : cdn.syndication.twimg.com/tweet-result?id=<status_id>
+         → text/mediaDetails(视频取 mp4、图片取 photos)/favorite_count/created_at
+
+清单文件(默认 stages/input/tool_url_ingest.json)格式见该文件 _comment。
+
+用法:
+  python stages/ingest_urls.py --dry-run              # 只取帖+拼 source,不评估不入库(先验证)
+  python stages/ingest_urls.py --dry-run -v           # 同上,打印每帖正文/媒体概览
+  python stages/ingest_urls.py                        # 真实:取帖 → 转写 → 评估 → 写 search_tools
+  python stages/ingest_urls.py --query-id q3302       # 只处理清单里的某个 query
+  python stages/ingest_urls.py --no-transcribe --no-images   # 省钱:不转写、纯文本评估
+"""
+import argparse
+import asyncio
+import copy
+import json
+import re
+import sys
+from pathlib import Path
+
+PROJECT_ROOT = Path(__file__).resolve().parents[3]   # …/Agent
+sys.path.insert(0, str(PROJECT_ROOT))
+
+from dotenv import load_dotenv
+load_dotenv()
+
+import httpx
+
+from examples.process_pipeline.script.search_eval.search_and_evaluate import (
+    evaluate_posts, transcribe_video_posts,
+)
+from examples.process_pipeline.script.llm_evaluate_sources import (
+    build_eval_llm_call, EVAL_MODELS, DEFAULT_EVAL_MODEL,
+)
+
+HERE = Path(__file__).resolve().parent           # stages/
+MW = HERE.parent                                  # mode_workflow/
+sys.path.insert(0, str(MW))
+import db
+# 评估去重的轻量 query 相关重算,复用 search_eval 的实现(该模块 __main__ 有 guard,import 安全)
+from search_eval import _rescore_query_relevance
+
+DEFAULT_MANIFEST = HERE / "input" / "tool_url_ingest.json"
+AIGC_DETAIL = "http://aigc-channel.aiddit.com/aigc/channel/detail"
+X_SYNDICATION = "https://cdn.syndication.twimg.com/tweet-result"
+UA = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36"
+
+
+# ── 按 URL 取帖 → source dict(与 search_and_evaluate.search_all 输出同构) ──────────
+
+def _source(platform, cid, url, post):
+    return {
+        "case_id": f"{platform}_{cid}", "platform": platform,
+        "channel_content_id": cid, "source_url": url,
+        "post": post, "comments": [], "found_by_queries": [],
+    }
+
+
+async def fetch_gzh(url, retries=3):
+    """公众号文章 → 完整 post(aigc-channel 详情接口,免鉴权)。
+
+    后端偶发「未知错误」(瞬时),故带退避重试。"""
+    last = ""
+    for i in range(retries):
+        async with httpx.AsyncClient(timeout=45) as c:
+            r = await c.post(AIGC_DETAIL, json={"type": "gzh", "content_link": url})
+            r.raise_for_status()
+            j = r.json()
+        if j.get("code") == 0 and j.get("data"):
+            break
+        last = j.get("message") or "无返回"
+        if i < retries - 1:
+            await asyncio.sleep(2 + 2 * i)
+    else:
+        raise RuntimeError(f"gzh 详情失败(重试 {retries} 次): {last}")
+    post = j["data"]
+    post.setdefault("link", url)
+    post["channel"] = "gzh"
+    cid = post.get("channel_content_id")
+    if not cid:
+        raise RuntimeError("gzh 详情缺 channel_content_id")
+    return _source("gzh", cid, url, post)
+
+
+async def fetch_x(url):
+    """推文 → post(X syndication CDN,免鉴权)。视频取最高码率 mp4,图片取 photos。"""
+    m = re.search(r"/status/(\d+)", url)
+    if not m:
+        raise RuntimeError(f"无法从 URL 解析 status id: {url}")
+    sid = m.group(1)
+    async with httpx.AsyncClient(timeout=30, follow_redirects=True) as c:
+        r = await c.get(X_SYNDICATION, params={"id": sid, "lang": "en", "token": "a"},
+                        headers={"User-Agent": UA})
+        r.raise_for_status()
+        d = r.json()
+    text = d.get("text", "") or ""
+    media = d.get("mediaDetails") or []
+    images, videos = [], []
+    for it in media:
+        t = it.get("type")
+        if t == "photo":
+            u = it.get("media_url_https")
+            if u:
+                images.append(u)
+        elif t in ("video", "animated_gif"):
+            variants = [v for v in (it.get("video_info") or {}).get("variants", [])
+                        if v.get("content_type") == "video/mp4" and v.get("url")]
+            if variants:
+                videos.append(max(variants, key=lambda v: v.get("bitrate", 0))["url"])
+    post = {
+        "channel_content_id": sid,
+        "title": text[:80],
+        "body_text": text,
+        "images": images,
+        "videos": videos,
+        # extract_video_url(platform="x") 读 video_url_list → 触发 Deepgram 转写
+        "video_url_list": [{"video_url": u} for u in videos],
+        "link": url,
+        "like_count": d.get("favorite_count"),
+        "publish_timestamp": d.get("created_at"),
+        "channel": "x",
+        "channel_account_name": (d.get("user") or {}).get("name"),
+        "content_type": "video" if videos else "normal",
+    }
+    return _source("x", sid, url, post)
+
+
+async def fetch_one(platform, url):
+    if platform == "gzh":
+        return await fetch_gzh(url)
+    if platform == "x":
+        return await fetch_x(url)
+    raise RuntimeError(f"暂不支持的渠道: {platform}(仅 gzh / x)")
+
+
+# ── 单个 query 的注入(取帖 → 转写 → 评估 → 写 search_tools) ──────────────────────
+
+async def ingest_query(q, args, eval_llm, eval_model_id):
+    qid, qtext, url_specs = q["query_id"], q["query"], q["urls"]
+    print(f"\n▶ {qid}  {qtext!r}  ({len(url_specs)} 链接)")
+
+    sources = []
+    for spec in url_specs:
+        platform, url = spec["platform"], spec["url"]
+        try:
+            s = await fetch_one(platform, url)
+            s["found_by_queries"] = [qtext]
+            post = s["post"]
+            print(f"   ✓ {platform:3s} {s['case_id']}  "
+                  f"图{len(post.get('images') or [])} 视{len(post.get('videos') or [])}  "
+                  f"「{(post.get('title') or post.get('body_text') or '')[:30]}」")
+            if args.verbose:
+                print(f"        body[:120]: {(post.get('body_text') or '')[:120]!r}")
+            sources.append(s)
+        except Exception as e:
+            print(f"   ✗ {platform:3s} 取帖失败 {url}\n       {type(e).__name__}: {e}")
+    if not sources:
+        print("   ❌ 无可用帖子,跳过"); return 0, 0.0
+
+    if args.dry_run:
+        print(f"   (dry-run)取到 {len(sources)} 帖,不评估不入库")
+        return 0, 0.0
+
+    # 统一时间戳(尽量,失败无妨)
+    try:
+        from examples.process_pipeline.script.extract_sources import _convert_timestamps
+        _convert_timestamps(sources)
+    except Exception:
+        pass
+
+    if not args.no_transcribe:
+        n = await transcribe_video_posts(sources, concurrency=args.max_concurrent)
+        if n:
+            print(f"   🎙️  视频转写 {n} 条")
+
+    # ── 评估(去重口径同 search_eval:已评过的帖复用通用分、只重算 query 相关分) ──
+    cost = 0.0
+    if not args.no_eval:
+        prior = {}
+        if not args.force_eval:
+            for s in sources:
+                e = db.fetch_existing_eval_any(s["case_id"])
+                if e:
+                    prior[s["case_id"]] = e
+        fresh = [s for s in sources if s["case_id"] not in prior]
+        reused = [s for s in sources if s["case_id"] in prior]
+        if reused:
+            print(f"   ♻️ 评估去重:{len(reused)} 帖已评过 → 复用+重算 query 相关;"
+                  f"{len(fresh)} 帖走完整评估")
+        if fresh:
+            _, c = await evaluate_posts(
+                fresh, "", eval_llm, eval_model_id, args.max_concurrent,
+                include_images=not args.no_images, max_images=args.max_images,
+                image_mode=args.image_mode, query=qtext,
+            )
+            cost += c
+        if reused:
+            sem = asyncio.Semaphore(args.max_concurrent)
+            rr = await asyncio.gather(*[
+                _rescore_query_relevance(s, qtext, eval_llm, eval_model_id, sem)
+                for s in reused])
+            for s, (qr, c) in zip(reused, rr):
+                cost += c
+                blob = copy.deepcopy(prior[s["case_id"]])
+                if qr is not None:
+                    blob.setdefault("相关性", {})["和 query 相关"] = {
+                        "得分": qr.get("得分"), "理由": qr.get("理由", "")}
+                s["llm_evaluation"] = blob
+    for s in sources:
+        s.pop("_image_data_urls", None)
+
+    n = db.upsert_search_posts(qid, qtext, sources, table="search_tools")
+    print(f"   🗄️  search_tools 入库 {n} 行 · 评估成本 ${cost:.4f}")
+
+    # 运行副本
+    out_dir = MW / "runs" / "search"
+    out_dir.mkdir(parents=True, exist_ok=True)
+    (out_dir / f"{qid}.json").write_text(json.dumps({
+        "query_id": qid, "query": qtext, "via": "ingest_urls",
+        "total": len(sources), "results": sources,
+    }, ensure_ascii=False, indent=2), encoding="utf-8")
+    return n, cost
+
+
+async def run(args):
+    manifest = json.loads(Path(args.manifest).read_text(encoding="utf-8"))
+    queries = manifest["queries"]
+    if args.query_id:
+        queries = [q for q in queries if q["query_id"] == args.query_id]
+        if not queries:
+            print(f"❌ 清单里没有 query_id={args.query_id}"); return 1
+
+    eval_llm, eval_model_id = build_eval_llm_call(args.eval_model)
+    print(f"清单 {args.manifest} · {len(queries)} 个 query · 评估模型 {eval_model_id}"
+          + ("  [dry-run]" if args.dry_run else ""))
+
+    total_rows = total_cost = 0.0
+    for q in queries:
+        n, c = await ingest_query(q, args, eval_llm, eval_model_id)
+        total_rows += n; total_cost += c
+    print(f"\n📊 合计写 {int(total_rows)} 行 · 评估成本 ${total_cost:.4f}")
+    if not args.dry_run:
+        print("下一步(工具解构):对每个 query 用显式 --case-ids 解构手挑帖(绕过采纳门槛):")
+        for q in queries:
+            print(f"  python stages/tool_extract.py --query-id {q['query_id']} "
+                  f"--case-ids <该 query 全部 case_id>")
+    return 0
+
+
+def main():
+    p = argparse.ArgumentParser(description="按渠道+链接直接注入工具方向(search_tools)")
+    p.add_argument("--manifest", default=str(DEFAULT_MANIFEST))
+    p.add_argument("--query-id", default=None, help="只处理清单里的某个 query")
+    p.add_argument("--eval-model", default=DEFAULT_EVAL_MODEL, choices=list(EVAL_MODELS))
+    p.add_argument("--max-concurrent", type=int, default=3)
+    p.add_argument("--max-images", type=int, default=4)
+    p.add_argument("--image-mode", choices=["url", "base64"], default="url")
+    p.add_argument("--no-transcribe", action="store_true")
+    p.add_argument("--no-eval", action="store_true")
+    p.add_argument("--no-images", action="store_true")
+    p.add_argument("--force-eval", action="store_true", help="跳过评估去重,全部走完整评估")
+    p.add_argument("--dry-run", action="store_true", help="只取帖+拼 source,不评估不入库")
+    p.add_argument("-v", "--verbose", action="store_true")
+    args = p.parse_args()
+    raise SystemExit(asyncio.run(run(args)))
+
+
+if __name__ == "__main__":
+    main()