Просмотр исходного кода

Implement V2 evidence-backed demand output

SamLee 4 недель назад
Родитель
Сommit
86ec3c6309

+ 4 - 0
examples/demand/db_manager.py

@@ -14,8 +14,10 @@ from examples.demand.pg_pattern_repository import (
     query_execution_for_evidence,
     query_itemset_evidence,
     query_itemset_items_with_categories,
+    query_seed_points_for_sources,
     query_latest_success_execution_id,
     query_seed_points_for_itemsets,
+    query_source_elements,
     query_video_ids_by_names,
 )
 
@@ -57,5 +59,7 @@ __all__ = [
     "query_itemset_items_with_categories",
     "query_latest_success_execution_id",
     "query_seed_points_for_itemsets",
+    "query_seed_points_for_sources",
+    "query_source_elements",
     "query_video_ids_by_names",
 ]

+ 63 - 143
examples/demand/demand.md

@@ -8,177 +8,97 @@ $system$
 
 # 需求选择 Agent
 
-你是一个需求产生 Agent。你的任务是基于高权重的元素,产生需求,并且根据已经选择的元素,进行拓展,发现更多的需求
+你是 DemandAgent。你的任务是基于高权重元素、高权重分类、共现关系和频繁项集,为「%merge_level2%」生成约 %count% 条真实需求。
 
-**需求 = 一个人带着某种目的或兴趣,能用一个词/短语表达出来**
-
-它的本质公式是:
-
-```
-需求 = 人的渴求 × 内容的可满足性
-```
-
-二者缺一不可:
-
-- 没有人的渴求 → 不是需求,是凭空造词
-- 内容无法满足 → 不是有效需求,是伪需求
+需求 = 人的渴求 × 内容的可满足性:
 
+- 没有人的渴求,不是需求。
+- 内容无法满足,不是有效需求。
+- 需求要能用一个词或短语表达,并且能让人理解用户想看什么。
 
 ## 背景知识
 
-### 数据来源
-
-数据来自社交媒体视频的结构化分析。每个帖子被拆解为多个"选题点"(灵感点、目的点、关键点),每个点下有三个维度的元素:
-
-- **实质**: 内容的核心主题/对象(如 "咖啡豆"、"护肤品")
-- **形式**: 内容的呈现形式(如 "测评对比"、"教程")
-- **意图**: 内容的目标/用户意图(如 "购买决策"、"学习技能")
-
-每个元素归属于一个分类树节点(如 实质 > 食品 > 饮品 > 咖啡),形成层级分类结构。
-每个元素或者分类都有自己的权重分,权重分用于评判元素或者分类受欢迎程度(核心要素)
+数据来自社交媒体视频的结构化分析。每个帖子被拆解为多个选题点:
 
-### Pattern Mining 结果
+- 灵感点
+- 目的点
+- 关键点
 
-通过 FP-Growth 算法挖掘出频繁项集 —— 在多个帖子中经常共同出现的元素组合。
+每个点下有三个维度的元素:
 
-- **频繁项集**: 一组经常共现的 items
-- **absolute_support**: 包含该项集的帖子数量
-- **combination_type**: 项集涉及的点类型组合
-- **is_cross_point**: 是否跨越多个选题点
+- 实质:内容的核心主题或对象。
+- 形式:内容的呈现形式。
+- 意图:内容的目标或用户意图。
 
-## 数据模型
+每个元素归属于分类树节点。元素和分类都有权重分;PG Pattern V2 还提供频繁项集,表示多个分类或元素在支撑帖中经常共同出现。
 
-### DemandItem(核心实体)
+## 工具
 
-需求产生过程 = ADD DemandItem。每个 DemandItem 代表一个需求。
+你可以使用:
 
-**字段:**
+- `get_weight_score_topn`:查看高权重元素/分类。
+- `get_weight_score_by_name`:查询指定元素/分类的权重。
+- `get_category_tree`:查看分类树结构。
+- `search_elements` / `search_categories`:查找真实元素或分类。
+- `get_category_co_occurrences` / `get_element_co_occurrences`:查看分类或元素共现。
+- `get_frequent_itemsets` / `get_itemset_detail`:查看 PG topic 频繁项集。
+- `get_post_elements`:查看支撑帖里的结构化元素。
+- `create_demand_item` / `create_demand_items`:写出 DemandItem。
 
-- `element_names`: 元素名称列表
-- `reason`: 产生该需求的理由
-- `desc`: 需求的描述
-- `type`: 需求的来源类型(元素/分类/关系/pattern)
-- `evidence_refs`: 候选证据引用对象。LLM 只能填写从工具返回中看到的来源,不代表最终已通过 DB 校验。
+## DemandItem 格式
 
-`evidence_refs` 推荐结构
+保留老版主结构,并增加候选证据引用:
 
 ```json
 {
-  "source_kind": "pattern_itemset",
-  "source_tool": "get_itemset_detail",
-  "itemset_ids": [123],
-  "category_ids": [456],
-  "source_post_id": "55157577",
-  "case_ids": {
-    "pattern_itemset": ["55157577"]
+  "element_names": ["需求词1", "需求词2"],
+  "reason": "为什么这个需求成立;说明它来自哪些高权重/共现/pattern 线索",
+  "desc": "用户希望看到什么内容",
+  "type": "元素/分类/关系/pattern",
+  "evidence_refs": {
+    "sources": [
+      {
+        "source_kind": "high_weight_element",
+        "source_tool": "get_weight_score_topn",
+        "element_names": ["工具返回的真实元素名"],
+        "element_type": "实质"
+      }
+    ]
   }
 }
 ```
 
-约束:
-
-- **本次 MVP 最终可写入的 DemandItem 只允许 `source_kind="pattern_itemset"`**。不要把 `element` / `category` / `co_occurrence` 作为最终 `evidence_refs.source_kind`,这些只能用于探索,不能用于创建需求。
-- Pattern 证据源是 PG Pattern V2,最终只允许使用 `pattern_itemset.scope="topic"` 的分类 Pattern;`topic_element`、script、paragraph、group scope 都不能作为最终证据。
-- 只能选择能闭合到真实点位证据的 itemset:优先使用 `dimension_mode="substance_form_only"` 的结果。不要选择 `dimension_mode="需求"`、`target_depth="all_ancestors"` 或 items 里 `dimension="需求"` 的 itemset;这类 itemset 不能稳定生成 `query_seed_points`。
-- `source_tool` 必须填写 `get_itemset_detail`,因为最终证据必须来自项集详情。
-- `itemset_ids` 必须来自 `get_frequent_itemsets` 返回的真实项集 ID,并且在创建需求前必须用 `get_itemset_detail` 查询过。
-- 每条 DemandItem 只能绑定一个 `itemset_id`。
-- `source_post_id` 必须从 `get_itemset_detail` 返回的 `post_ids` 中选择一个真实帖子 ID;没有明确帖子时不要创建该 DemandItem。
-- `case_ids` 固定按 pattern 来源表达,例如 `{"pattern_itemset": ["source_post_id"]}`。
-- 不要填写 `seed_terms`;`seed_terms`、`query_seed_points`、`source_certainty`、`validation_status` 和 `demand_scope` 都由代码从 PG DB 校验并补齐。
-
-## 工具概览
+`source_kind` 只能使用:
 
-### 查询工具(只读)
+- `high_weight_element`
+- `high_weight_category`
+- `element_co_occurrence`
+- `category_co_occurrence`
+- `pattern_itemset`
 
-- `get_category_tree` — 查看当前分类下的完整分类树(分类)
-- `get_weight_score_topn` — 元素/分类权重排行榜(元素/分类)
-- `get_weight_score_by_name` — 执行元素/分类权重查询
-- `get_frequent_itemsets` — 搜索频繁项集(pattern)
-- `get_itemset_detail` — 项集详情
-- `get_post_elements` — 帖子元素
-- `search_elements` / `search_categories` — 关键词搜索
-- `get_category_co_occurrences` / `get_element_co_occurrences` — 共现查询(关系)
+如果一条需求由多个证据支持,就在 `sources` 里放多条。不要再把“一个需求绑定几个 itemset”当成需求定义;itemset 只是可选证据来源之一。
 
-### CRUD 工具
+## 生成策略
 
-- `create_demand_item` — 创建一个新需求
-- `create_demand_items` — 批量创建新需求
+1. 先从高权重叶子元素和高权重分类节点出发。
+2. 单元素如果本身能表达清晰需求,可以直接进入候选。
+3. 分类组合必须从高权重分类作为起点,再用共现查询或频繁项集验证。
+4. 元素组合必须从高权重元素作为起点,再用元素共现、帖子数或频繁项集验证。
+5. 频繁项集用于补充和验证组合需求,不是唯一入口。
+6. 最终数量尽量接近 %count%,但不要为了凑数编造。
 
-### 输出工具
+## 硬约束
 
-- `write_execution_summary` — 写入执行总结
-
-## 硬约束(不可违反)
-
-- 执行每个操作前,必须输出自己的思考,为什么要这样做,原因是什么,目的是什么
-- category 级 item 必须来自分类树的真实节点(通过 `search_categories` 查到对应的 `category_id`),不允许凭空编造分类
-- 正确的创建顺序:先 element 后 category
-- result 中出现的每一个具体内容,都必须有对应的 DemandItem
-- 每个 DemandItem 都必须包含 `evidence_refs`,且 `evidence_refs.source_kind` 必须是 `pattern_itemset`
-- 创建 DemandItem 前,必须先调用 `get_frequent_itemsets` 找到候选项集,再调用 `get_itemset_detail` 拿到 `itemset_ids`、`items`、`post_ids`,并把一个真实 `post_id` 写入 `source_post_id`
-- 每条 DemandItem 只能绑定一个 itemset;不要把多个 itemset 合并到一条需求里。
-- 不允许用 `search_elements` / `search_categories` / 共现查询结果直接创建最终 DemandItem;这些工具只能辅助理解和筛选,最终创建时仍必须绑定到一个经过 `get_itemset_detail` 验证的 itemset
-- `evidence_refs` 是候选引用,不要写 `source_certainty=db_validated` 或 `validation_status=passed`,最终校验由代码完成
-- `search_elements` / `search_categories`只能用于查询单元素/单分类,不能用于查询完整树,完整树查询用`get_category_tree`
+- 每条 DemandItem 都必须有 `evidence_refs.sources`。
+- `element_names` 必须来自工具返回的真实元素名、分类名或 itemset item 原文;可以少量删除明显不适合表达需求的泛词/动词,但不能新增工具里没出现过的新词。
+- 不要填写 `seed_terms`、`query_seed_points`、`source_certainty`、`validation_status`,这些都由代码查 PG DB 后补齐。
+- category 级 item 必须来自真实分类树或搜索结果,不能凭空编造。
+- 共现查询的起点必须来自高权重元素或高权重分类。
+- 过于抽象、没有主体、只像制作手法、不能表达人的兴趣或目的的词,直接过滤。
+- 最终保留必须有权重分、共现帖子数或 itemset support 支撑。
 
 $user$
 
-## 前置条件
-
-`get_category_tree`工具查到到的全部分类都是属于「%merge_level2%」,`search_categories`查询的只是树种的一个或者多个分类
-
-## 任务
-
-针对「%merge_level2%」,从"高权重叶子元素","高权重分类节点"出发完成需求生成
-
-1. 单元素生成需求,需要先通过`get_weight_score_topn`工具查找高权重元素/高权重分类,判断是否能作为需求,给出理由。满足的进入需求池,不满足的给出丢弃的理由。
-2. 组合需求,分类的起点必须需要先通过`get_weight_score_topn`工具查询高权重分类,通过共现查询,找到合适的组合,计算组合权重的平均分和帖子数,综合判断保留或者移除
-3. 组合需求,分类的起点必须需要先通过`get_frequent_itemsets`工具,搜索频繁出现的分类组合,根据支持度进行移除和保留
-4. 生成的需求,必须要有实际的含义,要能体现出「%merge_level2%」,要从需求中可以了解到「%merge_level2%」,不符合要求的词,给出理由直接过滤掉
-
-## 本地 exact evidence 输出流程(必须遵守)
-
-本次输出要给下游 ContentFindAgent 做证据链回溯,所以最终需求必须能被代码校验成 `evidence_pack`。
-
-请按以下顺序工作:
-
-1. 先用 `get_weight_score_topn` / 分类树 / 搜索工具理解「%merge_level2%」的高价值元素和分类。
-2. 再用 `get_frequent_itemsets(dimension_mode="substance_form_only")` 查找能支撑这些需求的频繁项集。
-3. 对准备采用的每个项集,必须调用 `get_itemset_detail(itemset_ids=[...])`。
-4. 只从 `get_itemset_detail` 返回的详情中选择需求证据:
-   - `itemset_ids`: 当前项集 ID。
-   - `items`: 只能选择 `dimension` 属于 `实质/形式/意图` 的项集;不要选择 `需求` 维度项集。
-   - `source_post_id`: 从该项集 `post_ids` 中选择一个真实帖子。
-   - 不要填写 `seed_terms`,代码会从该项集 `items` 里的分类路径、分类名或元素名提取。
-5. 调用 `create_demand_item` 或 `create_demand_items` 时,每条必须使用如下结构:
-
-```json
-{
-  "element_names": ["需求词或短语"],
-  "reason": "为什么这个需求成立,并说明它来自哪个 itemset 的哪些 items",
-  "desc": "用户希望看到什么内容",
-  "type": "pattern",
-  "evidence_refs": {
-    "source_kind": "pattern_itemset",
-    "source_tool": "get_itemset_detail",
-    "itemset_ids": [123],
-    "source_post_id": "55157577",
-    "case_ids": {
-      "pattern_itemset": ["55157577"]
-    }
-  }
-}
-```
-
-如果某个候选需求没有可绑定的 `pattern_itemset`、`itemset_ids` 或 `source_post_id`,直接丢弃,不要为了凑数量改用 `element` / `category` / `co_occurrence` 证据。
-
-完成 `create_demand_items` 且工具返回成功后,用一句话总结完成情况即可,不要再调用 `write_execution_summary` 或输出长篇执行报告。
-
-## 要求
+请针对「%merge_level2%」生成约 %count% 条 DemandItem。
 
-1. 共现查询的地点必须来自于高权重分类,不能直接从树上寻找分类
-2. 分类的共现组合,必须来自于`get_weight_score_topn`查询到的分类作为起点
-3. 最终结果的保留,必须要有权重分或者支持度进行支持
-4. 尽量保证最终产生的需求数量在「%count%」个左右
-5. 元素的需求更有意义,尽可能多的产生元素的需求,当元素的需求不足的时候,用有意义的其他类型需求补充
+请先用高权重元素/分类理解方向,再用共现和频繁项集验证,最后只调用一次 `create_demand_items` 批量写出。

+ 70 - 38
examples/demand/demand_mysql.md

@@ -1,7 +1,7 @@
 ---
 model: deepseek-v4-flash-260425
-temperature: 0.3
-max_iterations: 6
+temperature: 0.5
+max_iterations: 200
 ---
 
 $system$
@@ -10,59 +10,91 @@ $system$
 
 你是 DemandAgent,目标是为「%merge_level2%」生成约 %count% 条可写入 `demand_content` 的需求。
 
-本入口用于云端 MySQL 单表落库,必须高效执行,不要展开整棵分类树。
+需求不是简单关键词。需求 = 人的渴求 × 内容的可满足性:
 
-## 硬约束
+- 没有人的渴求,不是需求。
+- 内容无法满足,不是有效需求。
+- 需求要能用一个词或短语表达,但必须能让人理解这个用户想看什么。
 
-- 只允许基于 PG Pattern V2 的 `scope="topic"` 频繁项集生成需求。
-- 只能选择能闭合到真实点位证据的 itemset:优先使用 `dimension_mode="substance_form_only"` 的结果。不要选择 `dimension_mode="需求"`、`target_depth="all_ancestors"` 或返回 items 里 `dimension="需求"` 的 itemset;这类 itemset 不能稳定生成 `query_seed_points`。
-- 每个 DemandItem 都必须带 `evidence_refs`。
-- `evidence_refs.source_kind` 固定为 `"pattern_itemset"`。
-- `evidence_refs.source_tool` 固定为 `"get_itemset_detail"`。
-- `itemset_ids` 必须来自 `get_frequent_itemsets` 和 `get_itemset_detail` 返回的真实 itemset。
-- 每条 DemandItem 只能绑定一个 `itemset_id`。
-- `source_post_id` 必须从 `get_itemset_detail` 返回的 `post_ids` 中选择。
-- 不要填写 `seed_terms`;`seed_terms`、`query_seed_points`、`demand_scope` 都由代码从 PG DB 派生。
-- 不要写 `source_certainty` 或 `validation_status`,这些由代码 DB 强校验后补齐。
+## 数据来源
 
-## 执行步骤
+当前事实源是 PG Pattern V2,不再使用旧 MySQL `topic_pattern_*`。
 
-1. 可先调用 `get_weight_score_topn` 查看当前品类下高权重元素/分类,用于辅助判断哪些方向值得生成需求。
-2. 调用 `get_frequent_itemsets(top_n=%count% * 3, min_support=5, sort_by="absolute_support", dimension_mode="substance_form_only")`。工具会自动按当前 `merge_leve2/platform` 过滤支撑帖。
-3. 从返回结果中挑选语义清晰、彼此尽量不重复的 itemset,数量尽量接近 %count%。
-4. 调用一次 `get_itemset_detail(itemset_ids=[...])` 查询这些 itemset 的详情。
-5. 调用一次 `create_demand_items(demand_items=[...])` 批量创建需求。
-6. 最后只用一句话总结,不要再调用其他工具。
+你能看到四类线索:
+
+1. 高权重元素:来自 `get_weight_score_topn(level="元素", dimension=...)`。
+2. 高权重分类:来自 `get_weight_score_topn(level="分类", dimension=...)`。
+3. 共现关系:来自 `get_category_co_occurrences` / `get_element_co_occurrences`。
+4. 频繁项集:来自 `get_frequent_itemsets` / `get_itemset_detail`。
+
+这些线索都只是候选。最终真实性由代码查 PG DB 生成 `evidence_pack`,你不要写 `source_certainty`、`validation_status`、`seed_terms` 或 `query_seed_points`。
 
 ## DemandItem 格式
 
+每条 DemandItem 只保留老版主结构,并增加候选证据引用:
+
 ```json
 {
   "element_names": ["需求词1", "需求词2"],
-  "reason": "说明该需求来自哪个 itemset,以及哪些分类为什么能表达用户需求",
+  "reason": "为什么这个需求成立;说明它来自哪些高权重/共现/pattern 线索",
   "desc": "用户希望看到什么内容",
-  "type": "pattern",
+  "type": "元素/分类/关系/pattern",
   "evidence_refs": {
-    "source_kind": "pattern_itemset",
-    "source_tool": "get_itemset_detail",
-    "itemset_ids": [123],
-    "source_post_id": "从 post_ids 中选择一个真实帖子",
-    "case_ids": {
-      "pattern_itemset": ["同 source_post_id"]
-    }
+    "sources": [
+      {
+        "source_kind": "high_weight_element",
+        "source_tool": "get_weight_score_topn",
+        "element_names": ["工具返回的真实元素名"],
+        "element_type": "实质"
+      }
+    ]
   }
 }
 ```
 
-## 选择标准
+`source_kind` 只能使用以下值:
+
+- `high_weight_element`:来自高权重元素或元素搜索结果。
+- `high_weight_category`:来自高权重分类或分类搜索结果。
+- `element_co_occurrence`:来自元素共现。
+- `category_co_occurrence`:来自分类共现。
+- `pattern_itemset`:来自 PG topic 频繁项集。
+
+如果一条需求同时由多个证据支持,就在 `sources` 里放多条。不要再把“一个需求绑定几个 itemset”当成需求定义;itemset 只是可选证据来源之一。
+如果使用多个 itemset 作为证据,必须把每个 itemset 拆成独立的 `pattern_itemset` source,不要把并列 itemset 塞进同一个 source。
+
+## 工作方式
 
-- 优先选择 `absolute_support` 高、分类组合语义明确的 itemset。
-- 只选择 items 里的 `dimension` 属于 `实质/形式/意图` 的 itemset;不要选择 `需求` 维度 itemset。
-- 过滤纯形式、过于抽象、难以表达用户需求的组合。
-- 如果候选不足 %count% 条,可以少于 %count%,不要编造需求。
-- 不要重复创建语义相同的需求。
-- 不要把多个 itemset 合并到同一条 DemandItem;如果一个组合值得生成需求,就为它单独创建一条。
+1. 先看高权重元素和高权重分类:
+   - 至少分别查看 `实质/形式/意图` 中你认为相关的 Top 区间。
+   - 单元素如果本身就是清晰需求,可以直接作为候选。
+
+2. 再从高权重分类或元素出发做共现:
+   - 分类组合必须从高权重分类作为起点。
+   - 元素组合必须从高权重元素作为起点。
+   - 共现结果要看帖子数和语义是否能表达真实需求。
+
+3. 再看频繁项集:
+   - `get_frequent_itemsets` 默认 TopN=20,可以按高权重分类的 `category_id` 定向查询。
+   - `get_itemset_detail` 用来查看 itemset 的真实 items 和支撑帖。
+   - 不要求每条需求都有 itemset;但如果你用了 itemset 做证据,必须把真实 `itemset_ids` 写进对应 source。
+   - 多个并列 itemset 要拆成多条 `sources[]`,让代码分别校验;每条 source 只描述一个可独立闭合的 itemset 证据。
+
+4. 最终创建需求:
+   - 数量尽量接近 %count%,但不要为了凑数编造。
+   - `element_names` 必须来自工具返回的真实元素名、分类名或 itemset item 原文;可以少量删除明显不适合表达需求的泛词/动词,但不能新增工具里没出现过的新词。
+   - 每条需求必须有 `evidence_refs.sources`;没有证据来源的候选直接丢弃。
+
+## 过滤标准
+
+- 过于抽象、没有主体、只像制作手法、不能表达人的兴趣或目的的词,丢弃。
+- 只靠“分享/讲述/展示”等动作词不能成立,除非后面有明确对象。
+- 纯形式词可以作为辅助证据,但不建议单独作为需求。
+- 同义或高度重复的需求只保留一条。
+- 最终保留必须有权重分、共现帖子数或 itemset support 支撑。
 
 $user$
 
-请针对「%merge_level2%」生成约 %count% 条 DemandItem。严格按系统指令执行,只使用 `get_weight_score_topn`、`get_weight_score_by_name`、`get_frequent_itemsets`、`get_itemset_detail`、`create_demand_items`。
+请针对「%merge_level2%」生成约 %count% 条 DemandItem。
+
+请使用高权重元素、高权重分类、共现、频繁项集一起判断,最后只调用一次 `create_demand_items` 批量写出。

+ 37 - 22
examples/demand/demand_pattern_tools.py

@@ -86,6 +86,17 @@ def _env_int(name: str, default: int) -> int:
         return default
 
 
+def _env_positive_int(name: str) -> int | None:
+    value = os.getenv(name)
+    if value is None or not str(value).strip():
+        return None
+    try:
+        parsed = int(value)
+    except ValueError:
+        return None
+    return parsed if parsed > 0 else None
+
+
 def _is_mysql_demand_content_entrypoint() -> bool:
     return os.getenv("DEMAND_MYSQL_ENTRYPOINT") in {
         "run_existing_execution_mysql",
@@ -108,7 +119,9 @@ def _resolve_scope_arg(value: Any, metadata_key: str) -> Any:
 def _compact_itemset_detail_for_mysql(data: list[dict[str, Any]]) -> list[dict[str, Any]]:
     if not _is_mysql_demand_content_entrypoint():
         return data
-    max_post_ids = max(_env_int("DEMAND_ITEMSET_DETAIL_MAX_POST_IDS", 20), 1)
+    max_post_ids = _env_positive_int("DEMAND_ITEMSET_DETAIL_MAX_POST_IDS")
+    if max_post_ids is None:
+        return data
     compacted: list[dict[str, Any]] = []
     for raw_itemset in data:
         itemset = dict(raw_itemset)
@@ -207,13 +220,6 @@ def get_frequent_itemsets(
     execution_id = TopicBuildAgentContext.get_execution_id()
     merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
     platform = _resolve_scope_arg(platform, "platform")
-    if _is_mysql_demand_content_entrypoint():
-        # V2 evidence needs exact PG element bindings and query_seed_points.
-        # The PG "需求" mining configs map to category nodes but often do not
-        # have source element rows, so keep the generation funnel on topic
-        # itemsets that can close through 实质/形式/意图 elements.
-        if not dimension_mode or str(dimension_mode).strip() == "需求":
-            dimension_mode = "substance_form_only"
     params = {
         "execution_id": execution_id, "top_n": top_n,
         "category_ids": category_ids, "dimension_mode": dimension_mode,
@@ -254,8 +260,9 @@ def get_itemset_detail(itemset_ids, merge_leve2=None, platform=None) -> str:
     merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
     platform = _resolve_scope_arg(platform, "platform")
     if _is_mysql_demand_content_entrypoint():
-        max_ids = max(_env_int("DEMAND_ITEMSET_DETAIL_MAX_IDS", 12), 1)
-        itemset_ids = itemset_ids[:max_ids]
+        max_ids = _env_positive_int("DEMAND_ITEMSET_DETAIL_MAX_IDS")
+        if max_ids is not None:
+            itemset_ids = itemset_ids[:max_ids]
     params = {"itemset_ids": itemset_ids, "merge_leve2": merge_leve2, "platform": platform}
     _log_tool_input("get_itemset_detail", params)
 
@@ -315,7 +322,7 @@ def get_post_elements(post_ids: list) -> str:
     "\n- 获取元素名称后,传给 get_element_co_occurrences 查共现关系"
     "\n- 通过元素的 category_id 桥接到 get_frequent_itemsets 做分类级分析")
 def search_elements(keyword: str, element_type: str = None, limit: int = 50,
-                    account_name=None, merge_leve2=None) -> str:
+                    account_name=None, merge_leve2=None, platform=None) -> str:
     """按名称关键词搜索元素。返回去重聚合后的元素列表,每个元素附带其所属分类信息(category_id、category_path)、出现次数和帖子数。
 
     使用场景:
@@ -333,13 +340,15 @@ def search_elements(keyword: str, element_type: str = None, limit: int = 50,
         元素列表的JSON字符串,每个元素含 name、element_type、category_id、category_path、occurrence_count、post_count。
     """
     execution_id = TopicBuildAgentContext.get_execution_id()
+    merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
+    platform = _resolve_scope_arg(platform, "platform")
     params = {"execution_id": execution_id, "keyword": keyword,
               "element_type": element_type, "limit": limit,
-              "account_name": account_name, "merge_leve2": merge_leve2}
+              "account_name": account_name, "merge_leve2": merge_leve2, "platform": platform}
     _log_tool_input("search_elements", params)
 
     data = pattern_service.search_elements(execution_id, keyword, element_type=element_type, limit=limit,
-                                           account_name=account_name, merge_leve2=merge_leve2)
+                                           account_name=account_name, merge_leve2=merge_leve2, platform=platform)
     result = json.dumps({
         "keyword": keyword,
         "count": len(data),
@@ -480,7 +489,7 @@ def search_categories(keyword: str, source_type: str = None) -> str:
     "\n- 通过 point_types 了解元素在灵感点/目的点/关键点中的分布"
     "\n- 获取元素名称后,传给 get_element_co_occurrences 查元素级共现"
     "\n- 从频繁项集的分类出发,落地到可用于选题的具体元素")
-def get_category_elements(category_id: int, account_name=None, merge_leve2=None) -> str:
+def get_category_elements(category_id: int, account_name=None, merge_leve2=None, platform=None) -> str:
     """获取某个分类节点下的元素列表(按名称去重聚合),按出现次数降序。
 
     Args:
@@ -492,11 +501,13 @@ def get_category_elements(category_id: int, account_name=None, merge_leve2=None)
         元素列表的JSON字符串,每个元素含 name、element_type、occurrence_count、post_count。
     """
     execution_id = TopicBuildAgentContext.get_execution_id()
-    params = {"category_id": category_id, "account_name": account_name, "merge_leve2": merge_leve2}
+    merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
+    platform = _resolve_scope_arg(platform, "platform")
+    params = {"category_id": category_id, "account_name": account_name, "merge_leve2": merge_leve2, "platform": platform}
     _log_tool_input("get_category_elements", params)
 
     data = pattern_service.get_category_elements(category_id, execution_id=execution_id,
-                                                 account_name=account_name, merge_leve2=merge_leve2)
+                                                 account_name=account_name, merge_leve2=merge_leve2, platform=platform)
     result = json.dumps({
         "category_id": category_id,
         "element_count": len(data),
@@ -518,7 +529,7 @@ def get_category_elements(category_id: int, account_name=None, merge_leve2=None)
     "\n- 验证频繁项集:将 get_frequent_itemsets 中发现的模式用此工具做更细粒度的验证"
     "\n\n提示:需要分类节点ID,可先用 search_categories 按名称查找。")
 def get_category_co_occurrences(category_ids: list, top_n: int = 30,
-                                account_name=None, merge_leve2=None) -> str:
+                                account_name=None, merge_leve2=None, platform=None) -> str:
     """查询多个分类的共现关系。找到同时包含所有指定分类下元素的帖子,返回这些帖子中其他分类的出现频率。
 
     支持叠加多分类,传入越多分类,结果越精确(交集缩小)。
@@ -533,8 +544,10 @@ def get_category_co_occurrences(category_ids: list, top_n: int = 30,
         共现分类列表的JSON字符串,含 matched_post_count(交集帖子数)和 co_categories(共现分类排名)。
     """
     execution_id = TopicBuildAgentContext.get_execution_id()
+    merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
+    platform = _resolve_scope_arg(platform, "platform")
     params = {"execution_id": execution_id, "category_ids": category_ids, "top_n": top_n,
-              "account_name": account_name, "merge_leve2": merge_leve2}
+              "account_name": account_name, "merge_leve2": merge_leve2, "platform": platform}
     _log_tool_input("get_category_co_occurrences", params)
 
     if not category_ids:
@@ -542,7 +555,7 @@ def get_category_co_occurrences(category_ids: list, top_n: int = 30,
 
     data = pattern_service.get_category_co_occurrences(
         execution_id=execution_id, category_ids=category_ids, top_n=top_n,
-        account_name=account_name, merge_leve2=merge_leve2,
+        account_name=account_name, merge_leve2=merge_leve2, platform=platform,
     )
     result = json.dumps(data, ensure_ascii=False, indent=2)
     return _log_tool_output("get_category_co_occurrences", result)
@@ -556,7 +569,7 @@ def get_category_co_occurrences(category_ids: list, top_n: int = 30,
     "\n- 渐进聚焦:先查单个元素,发现高频共现后叠加查询缩小范围"
     "\n\n提示:element_names 需要精确匹配,可先用 search_elements 按关键词查找确切名称。")
 def get_element_co_occurrences(element_names: list, top_n: int = 30,
-                               account_name=None, merge_leve2=None) -> str:
+                               account_name=None, merge_leve2=None, platform=None) -> str:
     """查询多个元素的共现关系。找到同时包含所有指定元素的帖子,返回这些帖子中其他元素的出现频率。
 
     支持叠加多元素,传入越多元素,结果越精确(交集缩小)。
@@ -571,8 +584,10 @@ def get_element_co_occurrences(element_names: list, top_n: int = 30,
         共现元素列表的JSON字符串,含 matched_post_count(交集帖子数)和 co_elements(共现元素排名)。
     """
     execution_id = TopicBuildAgentContext.get_execution_id()
+    merge_leve2 = _resolve_scope_arg(merge_leve2, "merge_leve2")
+    platform = _resolve_scope_arg(platform, "platform")
     params = {"execution_id": execution_id, "element_names": element_names, "top_n": top_n,
-              "account_name": account_name, "merge_leve2": merge_leve2}
+              "account_name": account_name, "merge_leve2": merge_leve2, "platform": platform}
     _log_tool_input("get_element_co_occurrences", params)
 
     if not element_names:
@@ -580,7 +595,7 @@ def get_element_co_occurrences(element_names: list, top_n: int = 30,
 
     data = pattern_service.get_element_co_occurrences(
         execution_id=execution_id, element_names=element_names, top_n=top_n,
-        account_name=account_name, merge_leve2=merge_leve2,
+        account_name=account_name, merge_leve2=merge_leve2, platform=platform,
     )
     result = json.dumps(data, ensure_ascii=False, indent=2)
     return _log_tool_output("get_element_co_occurrences", result)

+ 626 - 94
examples/demand/evidence_pack_builder.py

@@ -15,11 +15,24 @@ from examples.demand.db_manager import (
     query_execution_for_evidence,
     query_itemset_evidence,
     query_itemset_items_with_categories,
-    query_seed_points_for_itemsets,
+    query_seed_points_for_sources,
+    query_source_elements,
 )
 
 
 SOURCE_KIND_PATTERN_ITEMSET = "pattern_itemset"
+SOURCE_KIND_HIGH_WEIGHT_ELEMENT = "high_weight_element"
+SOURCE_KIND_HIGH_WEIGHT_CATEGORY = "high_weight_category"
+SOURCE_KIND_ELEMENT_CO_OCCURRENCE = "element_co_occurrence"
+SOURCE_KIND_CATEGORY_CO_OCCURRENCE = "category_co_occurrence"
+SOURCE_KIND_MULTI_SOURCE = "multi_source"
+SUPPORTED_SOURCE_KINDS = {
+    SOURCE_KIND_PATTERN_ITEMSET,
+    SOURCE_KIND_HIGH_WEIGHT_ELEMENT,
+    SOURCE_KIND_HIGH_WEIGHT_CATEGORY,
+    SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+    SOURCE_KIND_CATEGORY_CO_OCCURRENCE,
+}
 PATTERN_SOURCE_SYSTEM = "pg_pattern_v2"
 PATTERN_ITEMSET_SCOPE = "topic"
 
@@ -77,8 +90,184 @@ def _build_evidence_pack(
     )
     scope_merge_leve2 = _clean_str(resolved_scope.get("merge_leve2"))
     scope_platform = _clean_str(resolved_scope.get("platform"))
+
+    execution = query_execution_for_evidence(int(execution_id))
+    if not execution:
+        return _reject(f"execution_id {execution_id} not found")
+    execution_status = _clean_str(execution.get("status")).lower()
+    if execution_status != "success":
+        return _reject(
+            f"execution_id {execution_id} status is {execution.get('status')}, not success"
+        )
+
+    source_candidates = _extract_evidence_sources(evidence_refs, demand_item)
+    if not source_candidates:
+        return _reject("missing evidence_refs.sources or resolvable evidence source")
+
+    resolved_sources: list[dict[str, Any]] = []
+    reject_reasons: list[str] = []
+    for source in source_candidates:
+        source_kind = _resolve_source_kind(source, demand_item, [])
+        if source_kind == SOURCE_KIND_PATTERN_ITEMSET:
+            resolved, reason = _resolve_pattern_itemset_source(
+                execution_id=int(execution_id),
+                source=source,
+                demand_item=demand_item,
+                evidence_refs=evidence_refs,
+                scope_merge_leve2=scope_merge_leve2 or None,
+                scope_platform=scope_platform or None,
+            )
+        else:
+            resolved, reason = _resolve_elementish_source(
+                execution_id=int(execution_id),
+                source=source,
+                demand_item=demand_item,
+                source_kind=source_kind,
+                scope_merge_leve2=scope_merge_leve2 or None,
+                scope_platform=scope_platform or None,
+            )
+        if resolved:
+            resolved_sources.append(resolved)
+        elif reason:
+            reject_reasons.append(reason)
+
+    if not resolved_sources:
+        return _reject("no evidence source passed DB validation: " + "; ".join(reject_reasons[:5]))
+
+    matched_post_ids = _merge_post_ids(source.get("matched_post_ids") for source in resolved_sources)
+    if not matched_post_ids:
+        return _reject("validated evidence sources produced no matched_post_ids")
+
+    validated_source_post_ids = _unique_strings(source.get("source_post_id") for source in resolved_sources)
+    source_post_id = next((post_id for post_id in validated_source_post_ids if post_id in set(matched_post_ids)), "")
+    if not source_post_id:
+        source_post_id = _choose_source_post_id(evidence_refs, demand_item, matched_post_ids)
+    if not source_post_id:
+        source_post_id = matched_post_ids[0]
+    if source_post_id not in set(matched_post_ids):
+        return _reject(f"source_post_id {source_post_id} is not in validated matched_post_ids")
+
+    itemset_ids = _unique_ints(
+        itemset_id
+        for source in resolved_sources
+        for itemset_id in source.get("itemset_ids", [])
+    )
+    itemset_items = _dedupe_itemset_items(
+        item
+        for source in resolved_sources
+        for item in source.get("raw_itemset_items", [])
+    )
+    category_bindings = _dedupe_bindings(
+        binding
+        for source in resolved_sources
+        for binding in source.get("category_bindings", [])
+    )
+    element_bindings = _dedupe_bindings(
+        binding
+        for source in resolved_sources
+        for binding in source.get("element_bindings", [])
+    )
+    if not category_bindings:
+        return _reject("validated evidence sources produced no category_bindings")
+    if not element_bindings:
+        return _reject("validated evidence sources produced no element_bindings")
+
+    seed_terms = _build_seed_terms_from_sources(demand_item, resolved_sources, itemset_items, element_bindings)
+    if not seed_terms:
+        return _reject("seed_terms cannot be resolved from DB-validated evidence sources")
+
+    case_rows = query_case_ids_by_post_ids(matched_post_ids)
+    decode_case_ids = _unique_strings(row.get("case_id") for row in case_rows)
+    query_seed_points = query_seed_points_for_sources(
+        execution_id=int(execution_id),
+        matched_post_ids=matched_post_ids,
+        category_ids=_unique_ints(binding.get("category_id") for binding in category_bindings),
+        top_k=_env_int("DEMAND_QUERY_SEED_POINTS_TOP_K", 30),
+    )
+    mining_config_ids = _unique_ints(
+        config_id
+        for source in resolved_sources
+        for config_id in source.get("mining_config_ids", [])
+    )
+    absolute_support = _resolve_absolute_support(execution, resolved_sources, matched_post_ids)
+    source_kinds = _unique_strings(source.get("source_kind") for source in resolved_sources)
+    source_kind = source_kinds[0] if len(source_kinds) == 1 else SOURCE_KIND_MULTI_SOURCE
+
+    evidence_pack = {
+        "pattern_source_system": PATTERN_SOURCE_SYSTEM,
+        "pattern_execution_id": int(execution_id),
+        "mining_config_id": mining_config_ids[0] if mining_config_ids else None,
+        "source_kind": source_kind,
+        "case_id_type": "post_id",
+        "source_post_id": source_post_id,
+        "evidence_sources": _build_evidence_sources_payload(resolved_sources),
+        "category_bindings": category_bindings,
+        "element_bindings": element_bindings,
+        "itemset_ids": itemset_ids,
+        "itemset_items": _build_itemset_items(itemset_items),
+        "support": _resolve_support(execution, resolved_sources, matched_post_ids),
+        "absolute_support": absolute_support,
+        "filtered_absolute_support": len(matched_post_ids),
+        "scoped_post_count": len(matched_post_ids),
+        "matched_post_ids": matched_post_ids,
+        "video_ids": matched_post_ids,
+        "case_ids": matched_post_ids,
+        "decode_case_ids": decode_case_ids,
+        "seed_terms": seed_terms,
+        "query_seed_points": query_seed_points,
+        "demand_scope": resolved_scope,
+        "trace_id": trace_id,
+        "demand_task_id": demand_task_id,
+        "demand_content_id": demand_content_id,
+        "source_certainty": "db_validated",
+        "validation_status": "passed",
+    }
+    return {"success": True, "evidence_pack": evidence_pack}
+
+
+def _reject(reason: str) -> dict[str, Any]:
+    return {"success": False, "reject_reason": reason}
+
+
+def _extract_evidence_sources(evidence_refs: Mapping[str, Any], demand_item: Any) -> list[dict[str, Any]]:
+    raw_sources = evidence_refs.get("sources")
+    sources: list[dict[str, Any]] = []
+    if isinstance(raw_sources, list):
+        for raw_source in raw_sources:
+            if isinstance(raw_source, Mapping):
+                sources.append(dict(raw_source))
+    elif isinstance(raw_sources, Mapping):
+        sources.append(dict(raw_sources))
+
+    if not sources and evidence_refs:
+        sources.append(dict(evidence_refs))
+
+    if not sources:
+        names = _normalize_string_list(_read_field(demand_item, "element_names"))
+        if names:
+            sources.append(
+                {
+                    "source_kind": SOURCE_KIND_HIGH_WEIGHT_ELEMENT,
+                    "source_tool": "create_demand_items.element_names_fallback",
+                    "element_names": names,
+                }
+            )
+    return sources
+
+
+def _resolve_pattern_itemset_source(
+        *,
+        execution_id: int,
+        source: Mapping[str, Any],
+        demand_item: Any,
+        evidence_refs: Mapping[str, Any],
+        scope_merge_leve2: str | None,
+        scope_platform: str | None,
+) -> tuple[dict[str, Any] | None, str | None]:
     itemset_ids, invalid_itemset_ids = _normalize_int_list(
         _first_present(
+            source.get("itemset_ids"),
+            source.get("itemset_id"),
             evidence_refs.get("itemset_ids"),
             evidence_refs.get("itemset_id"),
             _read_field(demand_item, "itemset_ids"),
@@ -86,145 +275,188 @@ def _build_evidence_pack(
         )
     )
     if invalid_itemset_ids:
-        return _reject(f"invalid itemset_ids: {invalid_itemset_ids}")
-
-    source_kind = _resolve_source_kind(evidence_refs, demand_item, itemset_ids)
-    if source_kind and source_kind != SOURCE_KIND_PATTERN_ITEMSET:
-        return _reject(f"unsupported source_kind={source_kind}; only pattern_itemset is supported")
-    if not source_kind:
-        return _reject("missing source_kind=pattern_itemset")
+        return None, f"invalid itemset_ids: {invalid_itemset_ids}"
     if not itemset_ids:
-        return _reject("missing itemset_ids for source_kind=pattern_itemset")
-    if len(itemset_ids) != 1:
-        return _reject("V2 exact evidence requires exactly one itemset_id per DemandItem")
-
-    execution = query_execution_for_evidence(int(execution_id))
-    if not execution:
-        return _reject(f"execution_id {execution_id} not found")
-    execution_status = _clean_str(execution.get("status")).lower()
-    if execution_status != "success":
-        return _reject(
-            f"execution_id {execution_id} status is {execution.get('status')}, not success"
-        )
+        return None, "pattern_itemset source missing itemset_ids"
 
     itemsets = query_itemset_evidence(
-        execution_id=int(execution_id),
+        execution_id=execution_id,
         itemset_ids=itemset_ids,
-        merge_leve2=scope_merge_leve2 or None,
-        platform=scope_platform or None,
+        merge_leve2=scope_merge_leve2,
+        platform=scope_platform,
     )
     found_itemset_ids = {int(itemset["id"]) for itemset in itemsets}
     missing_itemset_ids = [itemset_id for itemset_id in itemset_ids if itemset_id not in found_itemset_ids]
     if missing_itemset_ids:
-        return _reject(
-            f"itemset_id {missing_itemset_ids[0]} not found under execution_id {execution_id}"
-        )
+        return None, f"itemset_id {missing_itemset_ids[0]} not found under execution_id {execution_id}"
 
     invalid_itemset_reason = _validate_itemset_facts(itemsets)
     if invalid_itemset_reason:
-        return _reject(invalid_itemset_reason)
-
-    mining_config_ids = _unique_ints(itemset.get("mining_config_id") for itemset in itemsets)
-    if len(mining_config_ids) != 1:
-        return _reject(
-            "itemset_ids must resolve to exactly one mining_config_id; "
-            f"got {mining_config_ids}"
-        )
+        return None, invalid_itemset_reason
 
     itemset_items = query_itemset_items_with_categories(
-        execution_id=int(execution_id),
+        execution_id=execution_id,
         itemset_ids=itemset_ids,
     )
     items_by_itemset: dict[int, list[dict[str, Any]]] = defaultdict(list)
     for item in itemset_items:
         items_by_itemset[int(item["itemset_id"])].append(item)
-
     item_reason = _validate_itemset_items(
-        execution_id=int(execution_id),
+        execution_id=execution_id,
         itemsets=itemsets,
         items_by_itemset=items_by_itemset,
     )
     if item_reason:
-        return _reject(item_reason)
+        return None, item_reason
 
     matched_post_ids = _merge_post_ids(itemset.get("matched_post_ids") for itemset in itemsets)
-    requested_source_post_id = _choose_source_post_id(evidence_refs, demand_item, matched_post_ids)
+    requested_source_post_id = _choose_source_post_id(source, demand_item, matched_post_ids)
     if not requested_source_post_id:
-        return _reject("source_post_id cannot be resolved from DB-validated matched_post_ids")
+        return None, "source_post_id cannot be resolved from DB-validated itemset posts"
     if requested_source_post_id not in set(matched_post_ids):
-        return _reject(
-            f"source_post_id {requested_source_post_id} is not in matched_post_ids for itemset_ids={itemset_ids}"
+        return None, (
+            f"source_post_id {requested_source_post_id} is not in matched_post_ids "
+            f"for itemset_ids={itemset_ids}"
         )
 
     element_bindings = query_element_bindings_for_items(
-        execution_id=int(execution_id),
+        execution_id=execution_id,
         itemset_items=itemset_items,
         post_ids=[requested_source_post_id],
     )
     binding_reason = _validate_element_bindings(itemset_items, element_bindings)
     if binding_reason:
-        resolved_source_post_id, element_bindings, binding_reason = _resolve_source_post_with_bindings(
-            execution_id=int(execution_id),
+        source_post_id, element_bindings, binding_reason = _resolve_source_post_with_bindings(
+            execution_id=execution_id,
             itemset_items=itemset_items,
             matched_post_ids=matched_post_ids,
             preferred_post_id=requested_source_post_id,
         )
         if binding_reason:
-            return _reject(binding_reason)
-        source_post_id = resolved_source_post_id
+            return None, binding_reason
     else:
         source_post_id = requested_source_post_id
 
-    seed_terms = _build_seed_terms_from_itemset_items(itemset_items)
-    seed_reason = _validate_seed_terms(seed_terms, itemset_items, [])
-    if seed_reason:
-        return _reject(seed_reason)
-
-    case_rows = query_case_ids_by_post_ids(matched_post_ids)
-    decode_case_ids = _unique_strings(row.get("case_id") for row in case_rows)
-    query_seed_points = query_seed_points_for_itemsets(
-        execution_id=int(execution_id),
-        itemset_ids=itemset_ids,
-        matched_post_ids=matched_post_ids,
-        top_k=_env_int("DEMAND_QUERY_SEED_POINTS_TOP_K", 30),
-    )
-    support_itemset = itemsets[0]
-    scoped_post_count = support_itemset.get("scoped_post_count")
-    filtered_absolute_support = support_itemset.get("filtered_absolute_support")
-
-    evidence_pack = {
-        "pattern_source_system": PATTERN_SOURCE_SYSTEM,
-        "pattern_execution_id": int(execution_id),
-        "mining_config_id": mining_config_ids[0],
+    return {
         "source_kind": SOURCE_KIND_PATTERN_ITEMSET,
-        "case_id_type": "post_id",
-        "source_post_id": source_post_id,
-        "category_bindings": _build_category_bindings(itemset_items),
-        "element_bindings": _build_element_bindings(element_bindings),
+        "source_tool": _clean_str(source.get("source_tool")) or "get_itemset_detail",
         "itemset_ids": itemset_ids,
-        "itemset_items": _build_itemset_items(itemset_items),
+        "mining_config_ids": _unique_ints(itemset.get("mining_config_id") for itemset in itemsets),
         "support": min(float(itemset["support"]) for itemset in itemsets),
         "absolute_support": min(int(itemset["absolute_support"]) for itemset in itemsets),
-        "filtered_absolute_support": int(filtered_absolute_support) if filtered_absolute_support is not None else None,
-        "scoped_post_count": int(scoped_post_count) if scoped_post_count is not None else None,
         "matched_post_ids": matched_post_ids,
-        "video_ids": matched_post_ids,
-        "case_ids": matched_post_ids,
-        "decode_case_ids": decode_case_ids,
-        "seed_terms": seed_terms,
-        "query_seed_points": query_seed_points,
-        "demand_scope": resolved_scope,
-        "trace_id": trace_id,
-        "demand_task_id": demand_task_id,
-        "demand_content_id": demand_content_id,
-        "source_certainty": "db_validated",
-        "validation_status": "passed",
-    }
-    return {"success": True, "evidence_pack": evidence_pack}
+        "source_post_id": source_post_id,
+        "raw_itemset_items": itemset_items,
+        "category_bindings": _build_category_bindings(itemset_items),
+        "element_bindings": _build_element_bindings(element_bindings),
+        "seed_terms": _build_seed_terms_from_itemset_items(itemset_items),
+    }, None
 
 
-def _reject(reason: str) -> dict[str, Any]:
-    return {"success": False, "reject_reason": reason}
+def _resolve_elementish_source(
+        *,
+        execution_id: int,
+        source: Mapping[str, Any],
+        demand_item: Any,
+        source_kind: str,
+        scope_merge_leve2: str | None,
+        scope_platform: str | None,
+) -> tuple[dict[str, Any] | None, str | None]:
+    if source_kind not in SUPPORTED_SOURCE_KINDS or source_kind == SOURCE_KIND_PATTERN_ITEMSET:
+        return None, f"unsupported source_kind={source_kind or '<empty>'}"
+
+    element_names = _source_element_names(source, demand_item)
+    category_ids, invalid_category_ids = _normalize_int_list(
+        _first_present(source.get("category_ids"), source.get("category_id"))
+    )
+    if invalid_category_ids:
+        return None, f"invalid category_ids: {invalid_category_ids}"
+    category_names = _normalize_string_list(
+        _first_present(source.get("category_names"), source.get("category_name"), source.get("categories"))
+    )
+    category_paths = _normalize_string_list(
+        _first_present(source.get("category_paths"), source.get("category_path"))
+    )
+    element_types = _normalize_string_list(
+        _first_present(source.get("element_types"), source.get("element_type"), source.get("dimension"))
+    )
+
+    groups: list[dict[str, Any]] = []
+    if source_kind in {SOURCE_KIND_HIGH_WEIGHT_ELEMENT, SOURCE_KIND_ELEMENT_CO_OCCURRENCE} and element_names:
+        groups.extend({"element_names": [name]} for name in element_names)
+    elif source_kind in {SOURCE_KIND_HIGH_WEIGHT_CATEGORY, SOURCE_KIND_CATEGORY_CO_OCCURRENCE}:
+        if category_ids:
+            groups.extend({"category_ids": [category_id]} for category_id in category_ids)
+        elif category_paths:
+            groups.extend({"category_paths": [path]} for path in category_paths)
+        elif category_names:
+            groups.extend({"category_names": [name]} for name in category_names)
+
+    if not groups:
+        groups.append(
+            {
+                "element_names": element_names,
+                "category_ids": category_ids,
+                "category_names": category_names,
+                "category_paths": category_paths,
+            }
+        )
+
+    if not any(_source_group_has_criteria(group) for group in groups):
+        return None, f"{source_kind} source has no resolvable names or categories"
+
+    rows_by_group: list[list[dict[str, Any]]] = []
+    for group in groups:
+        rows = query_source_elements(
+            execution_id=execution_id,
+            element_names=group.get("element_names"),
+            category_ids=group.get("category_ids"),
+            category_names=group.get("category_names"),
+            category_paths=group.get("category_paths"),
+            element_types=element_types,
+            merge_leve2=scope_merge_leve2,
+            platform=scope_platform,
+            limit=_env_int("DEMAND_EVIDENCE_SOURCE_ELEMENT_LIMIT", 20000),
+        )
+        if not rows:
+            return None, f"{source_kind} source cannot resolve DB rows for {group}"
+        rows_by_group.append(rows)
+
+    matched_sets = [
+        {str(row["post_id"]) for row in rows if row.get("post_id") is not None}
+        for rows in rows_by_group
+    ]
+    matched_post_set = set.intersection(*matched_sets) if matched_sets else set()
+    if not matched_post_set:
+        return None, f"{source_kind} source has no common scoped support posts"
+
+    matched_post_ids = sorted(matched_post_set)
+    all_rows = [
+        dict(row)
+        for rows in rows_by_group
+        for row in rows
+        if str(row.get("post_id") or "") in matched_post_set
+    ]
+    category_bindings = _build_source_category_bindings(all_rows, source_kind=source_kind)
+    element_bindings = _build_source_element_bindings(all_rows, source_kind=source_kind)
+    if not category_bindings or not element_bindings:
+        return None, f"{source_kind} source produced no category/element bindings"
+
+    return {
+        "source_kind": source_kind,
+        "source_tool": _clean_str(source.get("source_tool")),
+        "source_terms": _source_terms_from_bindings(category_bindings, element_bindings),
+        "itemset_ids": [],
+        "mining_config_ids": [],
+        "support": None,
+        "absolute_support": len(matched_post_ids),
+        "matched_post_ids": matched_post_ids,
+        "source_post_id": _choose_source_post_id(source, demand_item, matched_post_ids),
+        "raw_itemset_items": [],
+        "category_bindings": category_bindings,
+        "element_bindings": element_bindings,
+        "seed_terms": _source_terms_from_bindings(category_bindings, element_bindings),
+    }, None
 
 
 def _resolve_source_kind(
@@ -239,11 +471,39 @@ def _resolve_source_kind(
         )
     )
     if explicit:
-        return explicit
+        aliases = {
+            "element": SOURCE_KIND_HIGH_WEIGHT_ELEMENT,
+            "元素": SOURCE_KIND_HIGH_WEIGHT_ELEMENT,
+            "high_weight": SOURCE_KIND_HIGH_WEIGHT_ELEMENT,
+            "category": SOURCE_KIND_HIGH_WEIGHT_CATEGORY,
+            "分类": SOURCE_KIND_HIGH_WEIGHT_CATEGORY,
+            "co_occurrence": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "co-occurrence": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "co occurrence": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "element_co-occurrence": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "element co occurrence": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "category_co-occurrence": SOURCE_KIND_CATEGORY_CO_OCCURRENCE,
+            "category co occurrence": SOURCE_KIND_CATEGORY_CO_OCCURRENCE,
+            "关系": SOURCE_KIND_ELEMENT_CO_OCCURRENCE,
+            "pattern": SOURCE_KIND_PATTERN_ITEMSET,
+        }
+        return aliases.get(explicit, explicit)
 
     source_tool = _clean_str(evidence_refs.get("source_tool"))
     if itemset_ids or source_tool == "get_itemset_detail":
         return SOURCE_KIND_PATTERN_ITEMSET
+    if source_tool == "get_element_co_occurrences":
+        return SOURCE_KIND_ELEMENT_CO_OCCURRENCE
+    if source_tool == "get_category_co_occurrences":
+        return SOURCE_KIND_CATEGORY_CO_OCCURRENCE
+    if evidence_refs.get("category_id") or evidence_refs.get("category_ids"):
+        return SOURCE_KIND_HIGH_WEIGHT_CATEGORY
+    if evidence_refs.get("category_path") or evidence_refs.get("category_paths"):
+        return SOURCE_KIND_HIGH_WEIGHT_CATEGORY
+    if evidence_refs.get("category_name") or evidence_refs.get("category_names"):
+        return SOURCE_KIND_HIGH_WEIGHT_CATEGORY
+    if evidence_refs.get("element_names") or evidence_refs.get("element_name") or evidence_refs.get("name"):
+        return SOURCE_KIND_HIGH_WEIGHT_ELEMENT
     return ""
 
 
@@ -394,13 +654,44 @@ def _normalize_int_list(value: Any) -> tuple[list[int], list[Any]]:
     return result, invalid
 
 
+def _normalize_string_list(value: Any) -> list[str]:
+    result: list[str] = []
+    seen: set[str] = set()
+    for raw in _normalize_list(value):
+        text = _clean_str(raw)
+        if text and text not in seen:
+            seen.add(text)
+            result.append(text)
+    return result
+
+
+def _source_group_has_criteria(group: Mapping[str, Any]) -> bool:
+    return any(_normalize_list(value) for value in group.values())
+
+
+def _source_element_names(source: Mapping[str, Any], demand_item: Any) -> list[str]:
+    return _normalize_string_list(
+        _first_present(
+            source.get("element_names"),
+            source.get("element_name"),
+            source.get("names"),
+            source.get("name"),
+            source.get("seed_terms"),
+            _read_field(demand_item, "element_names"),
+        )
+    )
+
+
 def _unique_ints(values: Iterable[Any]) -> list[int]:
     result: list[int] = []
     seen: set[int] = set()
     for value in values:
         if value is None:
             continue
-        parsed = int(value)
+        try:
+            parsed = int(value)
+        except (TypeError, ValueError):
+            continue
         if parsed not in seen:
             seen.add(parsed)
             result.append(parsed)
@@ -625,8 +916,7 @@ def _build_seed_terms_from_itemset_items(itemset_items: list[dict[str, Any]]) ->
         if _normalize_term(term) not in {"其他", "其它", "未知", "内容", "视频"}
     ]
     source = filtered or candidates
-    max_terms = max(_env_int("DEMAND_SEED_TERMS_MAX", 10), 1)
-    return source[:max_terms]
+    return _maybe_limit_seed_terms(source)
 
 
 def _validate_seed_terms(
@@ -712,6 +1002,36 @@ def _build_category_bindings(itemset_items: list[dict[str, Any]]) -> list[dict[s
     return bindings
 
 
+def _build_source_category_bindings(rows: list[dict[str, Any]], *, source_kind: str) -> list[dict[str, Any]]:
+    bindings: list[dict[str, Any]] = []
+    seen: set[tuple[Any, ...]] = set()
+    for row in rows:
+        category_id = row.get("category_id")
+        category_path = row.get("category_path") or row.get("category_full_path")
+        if category_id is None and not category_path:
+            continue
+        key = (category_id, category_path, row.get("element_type"), source_kind)
+        if key in seen:
+            continue
+        seen.add(key)
+        bindings.append(
+            {
+                "source_kind": source_kind,
+                "category_id": int(category_id) if category_id is not None else None,
+                "category_name": row.get("category_name") or _last_path_part(category_path),
+                "category_path": category_path,
+                "category_full_path": row.get("category_full_path") or category_path,
+                "category_source_type": row.get("category_source_type"),
+                "category_level": row.get("category_level"),
+                "category_nature": row.get("category_nature"),
+                "point_type": row.get("point_type"),
+                "dimension": row.get("element_type"),
+                "element_name": row.get("name"),
+            }
+        )
+    return bindings
+
+
 def _build_element_bindings(element_bindings: list[dict[str, Any]]) -> list[dict[str, Any]]:
     return [
         {
@@ -730,6 +1050,66 @@ def _build_element_bindings(element_bindings: list[dict[str, Any]]) -> list[dict
     ]
 
 
+def _build_source_element_bindings(rows: list[dict[str, Any]], *, source_kind: str) -> list[dict[str, Any]]:
+    grouped: dict[tuple[Any, ...], dict[str, Any]] = {}
+    for row in rows:
+        key = (
+            source_kind,
+            row.get("category_id"),
+            row.get("point_type"),
+            row.get("element_type"),
+            row.get("name"),
+        )
+        group = grouped.setdefault(
+            key,
+            {
+                "source_kind": source_kind,
+                "category_id": int(row["category_id"]) if row.get("category_id") is not None else None,
+                "point_type": row.get("point_type"),
+                "dimension": row.get("element_type"),
+                "element_name": row.get("name"),
+                "matched_element_count": 0,
+                "matched_post_ids": [],
+                "sample_elements": [],
+            },
+        )
+        group["matched_element_count"] += 1
+        post_id = _clean_str(row.get("post_id"))
+        if post_id and post_id not in group["matched_post_ids"]:
+            group["matched_post_ids"].append(post_id)
+        if len(group["sample_elements"]) < 20:
+            group["sample_elements"].append(_source_element_sample(row))
+
+    result: list[dict[str, Any]] = []
+    for group in grouped.values():
+        group["matched_post_count"] = len(group["matched_post_ids"])
+        result.append(group)
+    result.sort(
+        key=lambda item: (
+            -int(item.get("matched_post_count") or 0),
+            str(item.get("element_name") or ""),
+            str(item.get("category_id") or ""),
+        )
+    )
+    return result
+
+
+def _source_element_sample(row: Mapping[str, Any]) -> dict[str, Any]:
+    return {
+        "id": row.get("id"),
+        "name": row.get("name"),
+        "post_id": _clean_str(row.get("post_id")),
+        "point_text": row.get("point_text"),
+        "point_type": row.get("point_type"),
+        "category_id": row.get("category_id"),
+        "element_type": row.get("element_type"),
+        "source_table": row.get("source_table"),
+        "category_path": row.get("category_path") or row.get("category_full_path"),
+        "topic_point_id": row.get("topic_point_id"),
+        "source_element_id": row.get("source_element_id"),
+    }
+
+
 def _build_itemset_items(itemset_items: list[dict[str, Any]]) -> list[dict[str, Any]]:
     return [
         {
@@ -742,3 +1122,155 @@ def _build_itemset_items(itemset_items: list[dict[str, Any]]) -> list[dict[str,
         }
         for item in itemset_items
     ]
+
+
+def _dedupe_itemset_items(items: Iterable[dict[str, Any]]) -> list[dict[str, Any]]:
+    result: list[dict[str, Any]] = []
+    seen: set[tuple[Any, ...]] = set()
+    for item in items:
+        key = (
+            item.get("itemset_item_id"),
+            item.get("itemset_id"),
+            item.get("category_id"),
+            item.get("point_type"),
+            item.get("dimension"),
+            item.get("element_name"),
+        )
+        if key in seen:
+            continue
+        seen.add(key)
+        result.append(dict(item))
+    return result
+
+
+def _dedupe_bindings(bindings: Iterable[dict[str, Any]]) -> list[dict[str, Any]]:
+    result: list[dict[str, Any]] = []
+    seen: set[tuple[Any, ...]] = set()
+    for binding in bindings:
+        key = (
+            binding.get("source_kind"),
+            binding.get("itemset_id"),
+            binding.get("itemset_item_id"),
+            binding.get("category_id"),
+            binding.get("category_path"),
+            binding.get("point_type"),
+            binding.get("dimension"),
+            binding.get("element_name"),
+        )
+        if key in seen:
+            continue
+        seen.add(key)
+        result.append(dict(binding))
+    return result
+
+
+def _source_terms_from_bindings(
+        category_bindings: list[dict[str, Any]],
+        element_bindings: list[dict[str, Any]],
+) -> list[str]:
+    terms: list[str] = []
+    for binding in element_bindings:
+        for value in (binding.get("element_name"),):
+            text = _clean_str(value)
+            if text and text not in terms:
+                terms.append(text)
+    for binding in category_bindings:
+        for value in (binding.get("category_name"), _last_path_part(binding.get("category_path"))):
+            text = _clean_str(value)
+            if text and text not in terms:
+                terms.append(text)
+    return terms
+
+
+def _build_seed_terms_from_sources(
+        demand_item: Any,
+        resolved_sources: list[dict[str, Any]],
+        itemset_items: list[dict[str, Any]],
+        element_bindings: list[dict[str, Any]],
+) -> list[str]:
+    del demand_item
+    candidates: list[str] = []
+    for source in resolved_sources:
+        for value in source.get("seed_terms", []) or []:
+            text = _clean_str(value)
+            if text and text not in candidates:
+                candidates.append(text)
+        for value in source.get("source_terms", []) or []:
+            text = _clean_str(value)
+            if text and text not in candidates:
+                candidates.append(text)
+
+    if not candidates and itemset_items:
+        candidates.extend(_build_seed_terms_from_itemset_items(itemset_items))
+    if not candidates:
+        candidates.extend(_source_terms_from_bindings([], element_bindings))
+
+    filtered = [
+        term for term in candidates
+        if _normalize_term(term) not in {"其他", "其它", "未知", "内容", "视频"}
+    ]
+    source = filtered or candidates
+    return _maybe_limit_seed_terms(source)
+
+
+def _maybe_limit_seed_terms(seed_terms: list[str]) -> list[str]:
+    raw_limit = os.getenv("DEMAND_SEED_TERMS_MAX")
+    if raw_limit is None or not str(raw_limit).strip():
+        return seed_terms
+    try:
+        max_terms = int(str(raw_limit).strip())
+    except ValueError:
+        return seed_terms
+    if max_terms <= 0:
+        return seed_terms
+    return seed_terms[:max_terms]
+
+
+def _resolve_support(
+        execution: Mapping[str, Any],
+        resolved_sources: list[dict[str, Any]],
+        matched_post_ids: list[str],
+) -> float:
+    itemset_supports = [
+        float(source["support"])
+        for source in resolved_sources
+        if source.get("support") is not None
+    ]
+    if itemset_supports:
+        return min(itemset_supports)
+    post_count = int(execution.get("post_count") or 0)
+    return float(len(matched_post_ids)) / float(post_count) if post_count > 0 else 0.0
+
+
+def _resolve_absolute_support(
+        execution: Mapping[str, Any],
+        resolved_sources: list[dict[str, Any]],
+        matched_post_ids: list[str],
+) -> int:
+    del execution
+    itemset_supports = [
+        int(source["absolute_support"])
+        for source in resolved_sources
+        if source.get("absolute_support") is not None and source.get("itemset_ids")
+    ]
+    return max([len(matched_post_ids), *itemset_supports])
+
+
+def _build_evidence_sources_payload(resolved_sources: list[dict[str, Any]]) -> list[dict[str, Any]]:
+    payload: list[dict[str, Any]] = []
+    for source in resolved_sources:
+        payload.append(
+            {
+                "source_kind": source.get("source_kind"),
+                "source_tool": source.get("source_tool"),
+                "itemset_ids": source.get("itemset_ids") or [],
+                "mining_config_ids": source.get("mining_config_ids") or [],
+                "source_terms": source.get("source_terms") or source.get("seed_terms") or [],
+                "source_post_id": source.get("source_post_id"),
+                "matched_post_count": len(source.get("matched_post_ids") or []),
+                "matched_post_ids": source.get("matched_post_ids") or [],
+                "support": source.get("support"),
+                "absolute_support": source.get("absolute_support"),
+            }
+        )
+    return payload

+ 17 - 6
examples/demand/mysql_demand_content_sink.py

@@ -119,7 +119,6 @@ def _validate_row(row: dict[str, Any], ext_data: dict[str, Any]) -> None:
         raise ValueError("missing ext_data.evidence_pack")
     expected = {
         "pattern_source_system": "pg_pattern_v2",
-        "source_kind": "pattern_itemset",
         "case_id_type": "post_id",
         "source_certainty": "db_validated",
         "validation_status": "passed",
@@ -131,9 +130,8 @@ def _validate_row(row: dict[str, Any], ext_data: dict[str, Any]) -> None:
     required = [
         "source_post_id",
         "pattern_execution_id",
-        "mining_config_id",
-        "itemset_ids",
-        "itemset_items",
+        "source_kind",
+        "evidence_sources",
         "category_bindings",
         "element_bindings",
         "support",
@@ -154,8 +152,21 @@ def _validate_row(row: dict[str, Any], ext_data: dict[str, Any]) -> None:
             continue
         if value is None or value == "" or value == []:
             raise ValueError(f"missing evidence_pack.{field}")
-    if not isinstance(evidence_pack.get("itemset_ids"), list) or len(evidence_pack["itemset_ids"]) != 1:
-        raise ValueError("evidence_pack.itemset_ids must contain exactly one itemset_id")
+    if evidence_pack.get("source_kind") not in {
+        "high_weight_element",
+        "high_weight_category",
+        "element_co_occurrence",
+        "category_co_occurrence",
+        "pattern_itemset",
+        "multi_source",
+    }:
+        raise ValueError(f"invalid evidence_pack.source_kind: {evidence_pack.get('source_kind')!r}")
+    if not isinstance(evidence_pack.get("evidence_sources"), list) or not evidence_pack["evidence_sources"]:
+        raise ValueError("missing evidence_pack.evidence_sources")
+    if not isinstance(evidence_pack.get("itemset_ids", []), list):
+        raise ValueError("evidence_pack.itemset_ids must be an array")
+    if not isinstance(evidence_pack.get("itemset_items", []), list):
+        raise ValueError("evidence_pack.itemset_items must be an array")
 
     matched_post_ids = [str(post_id) for post_id in evidence_pack.get("matched_post_ids") or []]
     if str(evidence_pack["source_post_id"]) not in set(matched_post_ids):

+ 55 - 10
examples/demand/pattern_builds/pg_pattern_service.py

@@ -248,9 +248,17 @@ def search_elements(
         limit: int = 50,
         account_name=None,
         merge_leve2=None,
+        platform=None,
 ) -> list[dict[str, Any]]:
-    del account_name, merge_leve2
-    rows = repo.query_elements(execution_id, keyword=keyword, element_type=element_type, limit=5000)
+    del account_name
+    rows = repo.query_elements(
+        execution_id,
+        keyword=keyword,
+        element_type=element_type,
+        merge_leve2=merge_leve2,
+        platform=platform,
+        limit=5000,
+    )
     grouped: dict[tuple[str, str, int | None, str | None], list[dict[str, Any]]] = defaultdict(list)
     for row in rows:
         key = (
@@ -364,9 +372,16 @@ def get_category_elements(
         execution_id: int,
         account_name=None,
         merge_leve2=None,
+        platform=None,
 ) -> list[dict[str, Any]]:
-    del account_name, merge_leve2
-    rows = repo.query_elements(execution_id, category_ids=[category_id], limit=5000)
+    del account_name
+    rows = repo.query_elements(
+        execution_id,
+        category_ids=[category_id],
+        merge_leve2=merge_leve2,
+        platform=platform,
+        limit=5000,
+    )
     grouped: dict[tuple[str, str], list[dict[str, Any]]] = defaultdict(list)
     for row in rows:
         name = str(row.get("name") or "").strip()
@@ -394,12 +409,19 @@ def get_category_co_occurrences(
         top_n: int = 30,
         account_name=None,
         merge_leve2=None,
+        platform=None,
 ) -> dict[str, Any]:
-    del account_name, merge_leve2
+    del account_name
     clean_ids = _clean_ints(category_ids)
     if not clean_ids:
         return {"matched_post_count": 0, "co_categories": []}
-    rows = repo.query_elements(execution_id, category_ids=clean_ids, limit=10000)
+    rows = repo.query_elements(
+        execution_id,
+        category_ids=clean_ids,
+        merge_leve2=merge_leve2,
+        platform=platform,
+        limit=10000,
+    )
     posts_by_category: dict[int, set[str]] = defaultdict(set)
     for row in rows:
         if row.get("category_id") is not None and row.get("post_id") is not None:
@@ -409,7 +431,15 @@ def get_category_co_occurrences(
         current = posts_by_category.get(category_id, set())
         matched_posts = current if matched_posts is None else matched_posts & current
     matched_posts = matched_posts or set()
-    co_rows = repo.query_elements(execution_id, post_ids=list(matched_posts)[:5000], limit=20000)
+    if not matched_posts:
+        return {"matched_post_count": 0, "matched_post_ids": [], "co_categories": []}
+    co_rows = repo.query_elements(
+        execution_id,
+        post_ids=list(matched_posts)[:5000],
+        merge_leve2=merge_leve2,
+        platform=platform,
+        limit=20000,
+    )
     counts: dict[tuple[int, str, str], set[str]] = defaultdict(set)
     for row in co_rows:
         cat_id = row.get("category_id")
@@ -440,14 +470,21 @@ def get_element_co_occurrences(
         top_n: int = 30,
         account_name=None,
         merge_leve2=None,
+        platform=None,
 ) -> dict[str, Any]:
-    del account_name, merge_leve2
+    del account_name
     names = [str(name).strip() for name in element_names or [] if str(name).strip()]
     if not names:
         return {"matched_post_count": 0, "co_elements": []}
     posts_by_name: dict[str, set[str]] = {}
     for name in names:
-        rows = repo.query_elements(execution_id, keyword=name, limit=10000)
+        rows = repo.query_elements(
+            execution_id,
+            keyword=name,
+            merge_leve2=merge_leve2,
+            platform=platform,
+            limit=10000,
+        )
         posts_by_name[name] = {
             str(row.get("post_id"))
             for row in rows
@@ -457,7 +494,15 @@ def get_element_co_occurrences(
     for posts in posts_by_name.values():
         matched_posts = posts if matched_posts is None else matched_posts & posts
     matched_posts = matched_posts or set()
-    co_rows = repo.query_elements(execution_id, post_ids=list(matched_posts)[:5000], limit=20000)
+    if not matched_posts:
+        return {"matched_post_count": 0, "matched_post_ids": [], "co_elements": []}
+    co_rows = repo.query_elements(
+        execution_id,
+        post_ids=list(matched_posts)[:5000],
+        merge_leve2=merge_leve2,
+        platform=platform,
+        limit=20000,
+    )
     grouped: dict[tuple[str, str, int | None, str | None], list[dict[str, Any]]] = defaultdict(list)
     for row in co_rows:
         name = str(row.get("name") or "")

+ 209 - 23
examples/demand/pg_pattern_repository.py

@@ -166,6 +166,22 @@ def _to_str_filter_list(value: Any) -> list[str]:
     return result
 
 
+def _expand_platform_filters(values: list[str]) -> list[str]:
+    """Expand known platform aliases used by Hive and PG post rows."""
+    aliases = {
+        "piaoquan": ["piaoquan", "票圈"],
+        "票圈": ["票圈", "piaoquan"],
+    }
+    result: list[str] = []
+    seen: set[str] = set()
+    for value in values:
+        for candidate in aliases.get(value, [value]):
+            if candidate and candidate not in seen:
+                seen.add(candidate)
+                result.append(candidate)
+    return result
+
+
 def _to_int_list(values: Any) -> list[int]:
     result: list[int] = []
     seen: set[int] = set()
@@ -185,7 +201,7 @@ def _to_int_list(values: Any) -> list[int]:
 def _scope_filter_payload(merge_leve2: Any = None, platform: Any = None) -> dict[str, list[str]]:
     return {
         "merge_leve2": _to_str_filter_list(merge_leve2),
-        "platform": _to_str_filter_list(platform),
+        "platform": _expand_platform_filters(_to_str_filter_list(platform)),
     }
 
 
@@ -666,57 +682,173 @@ def query_elements(
         element_type: str | None = None,
         post_ids: Iterable[Any] | None = None,
         category_ids: Iterable[Any] | None = None,
+        merge_leve2: Any = None,
+        platform: Any = None,
         limit: int | None = None,
 ) -> list[dict[str, Any]]:
     clauses = [
-        "execution_id = %s",
-        "source_table = %s",
+        "e.execution_id = %s",
+        "e.source_table = %s",
     ]
     params: list[Any] = [int(execution_id), TOPIC_ELEMENT_SOURCE_TABLE]
     if keyword:
-        clauses.append("name ILIKE %s")
+        clauses.append("e.name ILIKE %s")
         params.append(f"%{keyword}%")
     if element_type:
-        clauses.append("element_type = %s")
+        clauses.append("e.element_type = %s")
         params.append(element_type)
     clean_posts = _to_post_id_list(post_ids)
+    if post_ids is not None and not clean_posts:
+        return []
     if clean_posts:
-        clauses.append("post_id = ANY(%s)")
+        clauses.append("e.post_id = ANY(%s)")
         params.append(clean_posts)
     clean_categories = _to_int_list(category_ids or [])
     if clean_categories:
-        clauses.append("category_id = ANY(%s)")
+        clauses.append("e.category_id = ANY(%s)")
         params.append(clean_categories)
 
+    scope_filter = _scope_filter_payload(merge_leve2=merge_leve2, platform=platform)
+    scope_active = _scope_is_active(scope_filter)
+    join_post_sql = ""
+    if scope_active:
+        join_post_sql = "JOIN post post_scope ON post_scope.post_id = e.post_id"
+        if scope_filter["merge_leve2"]:
+            clauses.append("post_scope.merge_leve2 = ANY(%s)")
+            params.append(scope_filter["merge_leve2"])
+        if scope_filter["platform"]:
+            clauses.append("post_scope.platform = ANY(%s)")
+            params.append(scope_filter["platform"])
+
     limit_sql = "LIMIT %s" if limit else ""
     if limit:
         params.append(int(limit))
     return _fetch_all(
         f"""
         SELECT
-            id,
-            execution_id,
-            post_id,
-            source_table,
-            source_element_id,
-            element_type,
-            element_sub_type,
-            name,
-            description,
-            category_id,
-            category_path,
-            point_type,
-            point_text,
-            topic_point_id
-        FROM pattern_mining_element
+            e.id,
+            e.execution_id,
+            e.post_id,
+            e.source_table,
+            e.source_element_id,
+            e.element_type,
+            e.element_sub_type,
+            e.name,
+            e.description,
+            e.category_id,
+            e.category_path,
+            e.point_type,
+            e.point_text,
+            e.topic_point_id
+        FROM pattern_mining_element e
+        {join_post_sql}
         WHERE {" AND ".join(clauses)}
-        ORDER BY post_id, id
+        ORDER BY e.post_id, e.id
         {limit_sql}
         """,
         tuple(params),
     )
 
 
+def query_source_elements(
+        execution_id: int,
+        *,
+        element_names: Iterable[Any] | None = None,
+        category_ids: Iterable[Any] | None = None,
+        category_names: Iterable[Any] | None = None,
+        category_paths: Iterable[Any] | None = None,
+        element_types: Iterable[Any] | None = None,
+        post_ids: Iterable[Any] | None = None,
+        merge_leve2: Any = None,
+        platform: Any = None,
+        limit: int = 20000,
+) -> list[dict[str, Any]]:
+    """Read source elements for DB-validating non-itemset evidence sources."""
+    clauses = [
+        "e.execution_id = %s",
+        "e.source_table = %s",
+    ]
+    params: list[Any] = [int(execution_id), TOPIC_ELEMENT_SOURCE_TABLE]
+
+    clean_names = _clean_names(str(name) for name in _jsonish_to_list(element_names))
+    if clean_names:
+        clauses.append("e.name = ANY(%s)")
+        params.append(clean_names)
+
+    clean_category_ids = _to_int_list(category_ids or [])
+    if clean_category_ids:
+        clauses.append("e.category_id = ANY(%s)")
+        params.append(clean_category_ids)
+
+    clean_category_names = _clean_names(str(name) for name in _jsonish_to_list(category_names))
+    if clean_category_names:
+        clauses.append("(c.name = ANY(%s) OR regexp_replace(e.category_path, '^.*/', '') = ANY(%s))")
+        params.extend([clean_category_names, clean_category_names])
+
+    clean_category_paths = _clean_names(str(path) for path in _jsonish_to_list(category_paths))
+    if clean_category_paths:
+        clauses.append("(c.path = ANY(%s) OR e.category_path = ANY(%s))")
+        params.extend([clean_category_paths, clean_category_paths])
+
+    clean_element_types = _clean_names(str(kind) for kind in _jsonish_to_list(element_types))
+    if clean_element_types:
+        clauses.append("e.element_type = ANY(%s)")
+        params.append(clean_element_types)
+
+    clean_post_ids = _to_post_id_list(post_ids)
+    if clean_post_ids:
+        clauses.append("e.post_id = ANY(%s)")
+        params.append(clean_post_ids)
+
+    scope_filter = _scope_filter_payload(merge_leve2=merge_leve2, platform=platform)
+    scope_active = _scope_is_active(scope_filter)
+    join_post_sql = ""
+    if scope_active:
+        join_post_sql = "JOIN post post_scope ON post_scope.post_id = e.post_id"
+        if scope_filter["merge_leve2"]:
+            clauses.append("post_scope.merge_leve2 = ANY(%s)")
+            params.append(scope_filter["merge_leve2"])
+        if scope_filter["platform"]:
+            clauses.append("post_scope.platform = ANY(%s)")
+            params.append(scope_filter["platform"])
+
+    params.append(int(limit))
+    return _fetch_all(
+        f"""
+        SELECT
+            e.id,
+            e.execution_id,
+            e.post_id,
+            e.source_table,
+            e.source_element_id,
+            e.element_type,
+            e.element_sub_type,
+            e.name,
+            e.description,
+            e.category_id,
+            e.category_path,
+            e.point_type,
+            e.point_text,
+            e.topic_point_id,
+            c.name AS category_name,
+            c.path AS category_full_path,
+            c.source_type AS category_source_type,
+            c.level AS category_level,
+            c.parent_id AS category_parent_id,
+            c.classified_as AS category_nature
+        FROM pattern_mining_element e
+        LEFT JOIN pattern_mining_category c
+          ON c.id = e.category_id
+         AND c.execution_id = e.execution_id
+        {join_post_sql}
+        WHERE {" AND ".join(clauses)}
+        ORDER BY e.post_id, e.id
+        LIMIT %s
+        """,
+        tuple(params),
+    )
+
+
 def query_topic_itemsets(
         execution_id: int,
         *,
@@ -971,6 +1103,60 @@ def query_seed_points_for_itemsets(
         ),
     )
 
+    return _rank_query_seed_point_rows(rows, top_k=top_k)
+
+
+def query_seed_points_for_sources(
+        execution_id: int,
+        matched_post_ids: Iterable[Any],
+        *,
+        category_ids: Iterable[Any] | None = None,
+        top_k: int = 30,
+) -> list[dict[str, Any]]:
+    """Rank searchable source points for any DB-validated evidence source."""
+    clean_post_ids = _to_post_id_list(list(matched_post_ids or []))
+    if not clean_post_ids:
+        return []
+
+    clauses = [
+        "e.execution_id = %s",
+        "e.source_table = %s",
+        "e.post_id = ANY(%s)",
+        "e.point_type IN ('灵感点', '目的点')",
+        "e.point_text IS NOT NULL",
+        "e.point_text <> ''",
+    ]
+    params: list[Any] = [int(execution_id), TOPIC_ELEMENT_SOURCE_TABLE, clean_post_ids]
+    clean_category_ids = _to_int_list(category_ids or [])
+    if clean_category_ids:
+        clauses.append("e.category_id = ANY(%s)")
+        params.append(clean_category_ids)
+
+    rows = _fetch_all(
+        f"""
+        SELECT
+            e.id,
+            e.post_id,
+            e.source_table,
+            e.source_element_id,
+            e.point_type,
+            e.point_text,
+            e.element_type,
+            e.name,
+            e.category_id,
+            e.category_path,
+            e.topic_point_id,
+            e.category_id AS matched_category_id
+        FROM pattern_mining_element e
+        WHERE {" AND ".join(clauses)}
+        ORDER BY e.point_text, e.point_type, e.post_id, e.id
+        """,
+        tuple(params),
+    )
+    return _rank_query_seed_point_rows(rows, top_k=top_k)
+
+
+def _rank_query_seed_point_rows(rows: list[dict[str, Any]], *, top_k: int = 30) -> list[dict[str, Any]]:
     grouped: dict[tuple[str, str], dict[str, Any]] = {}
     for row in rows:
         point_text = str(row.get("point_text") or "").strip()

+ 10 - 0
examples/demand/run.py

@@ -215,8 +215,18 @@ def _enabled_tools_for_run(configured_tools: list[str]) -> list[str]:
     tools = configured_tools.copy()
     if _is_mysql_demand_content_mode():
         allowed = {
+            "think_and_plan",
+            "get_category_tree",
             "get_frequent_itemsets",
             "get_itemset_detail",
+            "get_post_elements",
+            "search_elements",
+            "get_element_category_chain",
+            "get_category_detail",
+            "search_categories",
+            "get_category_elements",
+            "get_category_co_occurrences",
+            "get_element_co_occurrences",
             "get_weight_score_topn",
             "get_weight_score_by_name",
             "create_demand_item",

+ 2 - 1
examples/demand/run_existing_execution_local.py

@@ -57,7 +57,8 @@ def _configure_local_env(execution_id: int, output_root: Path) -> None:
     os.environ["DEMAND_WEIGHT_DATA_DIR"] = str(output_root / "intermediate" / "data")
     os.environ["DEMAND_TRACE_STORE_PATH"] = str(output_root / ".trace")
     os.environ["DEMAND_OUTPUT_BASE_DIR"] = str(output_root / "output")
-    os.environ.setdefault("DEMAND_LOCAL_TASK_ID", "1")
+    if not os.environ.get("DEMAND_LOCAL_TASK_ID", "").strip():
+        os.environ["DEMAND_LOCAL_TASK_ID"] = "1"
 
 
 def _validate_execution_success(execution_id: int) -> None:

+ 0 - 2
examples/demand/run_existing_execution_mysql.py

@@ -48,8 +48,6 @@ def _configure_mysql_env(run_label: str) -> None:
     os.environ["DEMAND_OUTPUT_MODE"] = "mysql_demand_content"
     os.environ["DEMAND_MYSQL_ENTRYPOINT"] = "run_existing_execution_mysql"
     os.environ["DEMAND_RUN_LABEL"] = run_label
-    os.environ.setdefault("DEMAND_ITEMSET_DETAIL_MAX_IDS", "12")
-    os.environ.setdefault("DEMAND_ITEMSET_DETAIL_MAX_POST_IDS", "20")
 
 
 def _validate_execution_success(execution_id: int) -> None:

+ 0 - 2
examples/demand/run_hive_gap_mysql.py

@@ -48,8 +48,6 @@ def _configure_mysql_env(run_label: str) -> None:
     os.environ["DEMAND_OUTPUT_MODE"] = "mysql_demand_content"
     os.environ["DEMAND_MYSQL_ENTRYPOINT"] = "run_hive_gap_mysql"
     os.environ["DEMAND_RUN_LABEL"] = run_label
-    os.environ.setdefault("DEMAND_ITEMSET_DETAIL_MAX_IDS", "12")
-    os.environ.setdefault("DEMAND_ITEMSET_DETAIL_MAX_POST_IDS", "20")
 
 
 def _validate_execution_success(execution_id: int) -> None:

+ 18 - 3
examples/demand/tests/test_mysql_demand_content_sink.py

@@ -61,6 +61,20 @@ def _valid_row():
         "source_post_id": "p1",
         "pattern_execution_id": 581,
         "mining_config_id": 7,
+        "evidence_sources": [
+            {
+                "source_kind": "pattern_itemset",
+                "source_tool": "get_itemset_detail",
+                "itemset_ids": [11],
+                "mining_config_ids": [7],
+                "source_terms": ["x"],
+                "source_post_id": "p1",
+                "matched_post_count": 2,
+                "matched_post_ids": ["p1", "p2"],
+                "support": 0.5,
+                "absolute_support": 2,
+            }
+        ],
         "itemset_ids": [11],
         "itemset_items": [{"itemset_id": 11, "category_id": 22, "category_path": "A>B"}],
         "category_bindings": [{"category_id": 22, "category_path": "A>B"}],
@@ -186,15 +200,16 @@ class MySQLDemandContentSinkTest(unittest.TestCase):
             with self.assertRaisesRegex(ValueError, "demand_scope"):
                 sink.write_demand_content_rows([row], run_label="batch01")
 
-    def test_rejects_multiple_itemset_ids(self):
+    def test_accepts_multiple_itemset_ids(self):
         from examples.demand import mysql_demand_content_sink as sink
 
         row = _valid_row()
         row["ext_data"]["evidence_pack"]["itemset_ids"] = [11, 12]
+        row["ext_data"]["evidence_pack"]["evidence_sources"][0]["itemset_ids"] = [11, 12]
 
         with patch.object(sink, "_connect", return_value=_FakeConnection()):
-            with self.assertRaisesRegex(ValueError, "exactly one"):
-                sink.write_demand_content_rows([row], run_label="batch01")
+            result = sink.write_demand_content_rows([row], run_label="batch01")
+        self.assertEqual(result.inserted_count, 1)
 
 
 if __name__ == "__main__":

+ 209 - 10
examples/demand/tests/test_pg_evidence_builder.py

@@ -25,8 +25,8 @@ def _itemset(**overrides):
     return data
 
 
-def _item():
-    return {
+def _item(**overrides):
+    data = {
         "itemset_item_id": 1,
         "itemset_id": 1607313,
         "category_id": 10,
@@ -40,10 +40,12 @@ def _item():
         "point_type": "关键点",
         "element_name": None,
     }
+    data.update(overrides)
+    return data
 
 
-def _binding():
-    return {
+def _binding(**overrides):
+    data = {
         "itemset_id": 1607313,
         "itemset_item_id": 1,
         "category_id": 10,
@@ -55,6 +57,8 @@ def _binding():
         "matched_post_ids": ["p1"],
         "sample_elements": [{"name": "综合性腐败", "category_path": "/事件行为/违纪违法/综合性腐败"}],
     }
+    data.update(overrides)
+    return data
 
 
 def _empty_binding(post_ids=None):
@@ -83,7 +87,7 @@ class PgEvidenceBuilderTest(unittest.TestCase):
             patch("examples.demand.evidence_pack_builder.query_itemset_items_with_categories", return_value=[_item()]),
             patch("examples.demand.evidence_pack_builder.query_element_bindings_for_items", return_value=[_binding()]),
             patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
-            patch("examples.demand.evidence_pack_builder.query_seed_points_for_itemsets", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
         ]
         for patcher in patches:
             patcher.start()
@@ -149,7 +153,7 @@ class PgEvidenceBuilderTest(unittest.TestCase):
             patch("examples.demand.evidence_pack_builder.query_itemset_items_with_categories", return_value=[_item()]),
             patch("examples.demand.evidence_pack_builder.query_element_bindings_for_items", side_effect=bindings),
             patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
-            patch("examples.demand.evidence_pack_builder.query_seed_points_for_itemsets", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
         ):
             result = self._build(
                 {
@@ -162,8 +166,36 @@ class PgEvidenceBuilderTest(unittest.TestCase):
         self.assertTrue(result["success"])
         self.assertEqual(result["evidence_pack"]["source_post_id"], "p2")
 
-    def test_multiple_itemsets_are_rejected(self):
-        self._patch_success()
+    def test_multiple_itemsets_are_supported(self):
+        patches = [
+            patch("examples.demand.evidence_pack_builder.query_execution_for_evidence", return_value=_execution()),
+            patch(
+                "examples.demand.evidence_pack_builder.query_itemset_evidence",
+                return_value=[
+                    _itemset(),
+                    _itemset(id=1607314, mining_config_id=2082, absolute_support=2),
+                ],
+            ),
+            patch(
+                "examples.demand.evidence_pack_builder.query_itemset_items_with_categories",
+                return_value=[
+                    _item(),
+                    _item(itemset_item_id=2, itemset_id=1607314, category_id=11, bound_category_id=11),
+                ],
+            ),
+            patch(
+                "examples.demand.evidence_pack_builder.query_element_bindings_for_items",
+                return_value=[
+                    _binding(),
+                    _binding(itemset_id=1607314, itemset_item_id=2, category_id=11),
+                ],
+            ),
+            patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
+        ]
+        for patcher in patches:
+            patcher.start()
+            self.addCleanup(patcher.stop)
         result = self._build(
             {
                 "source_kind": "pattern_itemset",
@@ -172,8 +204,175 @@ class PgEvidenceBuilderTest(unittest.TestCase):
             }
         )
 
-        self.assertFalse(result["success"])
-        self.assertIn("exactly one itemset_id", result["reject_reason"])
+        self.assertTrue(result["success"])
+        self.assertEqual(result["evidence_pack"]["itemset_ids"], [1607313, 1607314])
+        self.assertEqual(result["evidence_pack"]["source_kind"], "pattern_itemset")
+
+    def test_high_weight_element_source_builds_pg_evidence_pack(self):
+        rows = [
+            {
+                "id": 1,
+                "post_id": "p1",
+                "source_table": "post_decode_topic_point_element",
+                "source_element_id": 11,
+                "point_type": "灵感点",
+                "point_text": "反腐故事",
+                "element_type": "实质",
+                "name": "综合性腐败",
+                "category_id": 10,
+                "category_name": "综合性腐败",
+                "category_path": "/事件行为/违纪违法/综合性腐败",
+                "category_full_path": "/事件行为/违纪违法/综合性腐败",
+                "topic_point_id": 101,
+            }
+        ]
+        with (
+            patch("examples.demand.evidence_pack_builder.query_execution_for_evidence", return_value={**_execution(), "post_count": 10}),
+            patch("examples.demand.evidence_pack_builder.query_source_elements", return_value=rows),
+            patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
+        ):
+            result = self._build(
+                {
+                    "sources": [
+                        {
+                            "source_kind": "high_weight_element",
+                            "source_tool": "get_weight_score_topn",
+                            "element_names": ["综合性腐败"],
+                            "element_type": "实质",
+                        }
+                    ]
+                }
+            )
+
+        self.assertTrue(result["success"])
+        pack = result["evidence_pack"]
+        self.assertEqual(pack["source_kind"], "high_weight_element")
+        self.assertEqual(pack["itemset_ids"], [])
+        self.assertEqual(pack["matched_post_ids"], ["p1"])
+        self.assertEqual(pack["seed_terms"], ["综合性腐败"])
+        self.assertEqual(pack["evidence_sources"][0]["source_kind"], "high_weight_element")
+
+    def test_co_occurrence_alias_is_normalized(self):
+        rows = [
+            {
+                "id": 1,
+                "post_id": "p1",
+                "source_table": "post_decode_topic_point_element",
+                "source_element_id": 11,
+                "point_type": "灵感点",
+                "point_text": "反腐故事",
+                "element_type": "实质",
+                "name": "综合性腐败",
+                "category_id": 10,
+                "category_name": "综合性腐败",
+                "category_path": "/事件行为/违纪违法/综合性腐败",
+                "category_full_path": "/事件行为/违纪违法/综合性腐败",
+                "topic_point_id": 101,
+            }
+        ]
+        with (
+            patch("examples.demand.evidence_pack_builder.query_execution_for_evidence", return_value={**_execution(), "post_count": 10}),
+            patch("examples.demand.evidence_pack_builder.query_source_elements", return_value=rows),
+            patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
+        ):
+            result = self._build(
+                {
+                    "sources": [
+                        {
+                            "source_kind": "co-occurrence",
+                            "source_tool": "get_element_co_occurrences",
+                            "element_names": ["综合性腐败"],
+                        }
+                    ]
+                }
+            )
+
+        self.assertTrue(result["success"])
+        self.assertEqual(result["evidence_pack"]["source_kind"], "element_co_occurrence")
+
+    def test_high_weight_category_source_builds_pg_evidence_pack(self):
+        rows = [
+            {
+                "id": 1,
+                "post_id": "p1",
+                "source_table": "post_decode_topic_point_element",
+                "source_element_id": 11,
+                "point_type": "灵感点",
+                "point_text": "反腐故事",
+                "element_type": "实质",
+                "name": "综合性腐败",
+                "category_id": 10,
+                "category_name": "综合性腐败",
+                "category_path": "/事件行为/违纪违法/综合性腐败",
+                "category_full_path": "/事件行为/违纪违法/综合性腐败",
+                "topic_point_id": 101,
+            }
+        ]
+        with (
+            patch("examples.demand.evidence_pack_builder.query_execution_for_evidence", return_value={**_execution(), "post_count": 10}),
+            patch("examples.demand.evidence_pack_builder.query_source_elements", return_value=rows),
+            patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
+        ):
+            result = self._build(
+                {
+                    "sources": [
+                        {
+                            "source_kind": "high_weight_category",
+                            "source_tool": "get_weight_score_topn",
+                            "category_ids": [10],
+                        }
+                    ]
+                }
+            )
+
+        self.assertTrue(result["success"])
+        pack = result["evidence_pack"]
+        self.assertEqual(pack["source_kind"], "high_weight_category")
+        self.assertEqual(pack["matched_post_ids"], ["p1"])
+        self.assertEqual(pack["seed_terms"], ["综合性腐败"])
+        self.assertEqual(pack["evidence_sources"][0]["source_kind"], "high_weight_category")
+
+    def test_category_co_occurrence_source_builds_pg_evidence_pack(self):
+        rows = [
+            {
+                "id": 1,
+                "post_id": "p1",
+                "source_table": "post_decode_topic_point_element",
+                "source_element_id": 11,
+                "point_type": "灵感点",
+                "point_text": "反腐故事",
+                "element_type": "实质",
+                "name": "综合性腐败",
+                "category_id": 10,
+                "category_name": "综合性腐败",
+                "category_path": "/事件行为/违纪违法/综合性腐败",
+                "category_full_path": "/事件行为/违纪违法/综合性腐败",
+                "topic_point_id": 101,
+            }
+        ]
+        with (
+            patch("examples.demand.evidence_pack_builder.query_execution_for_evidence", return_value={**_execution(), "post_count": 10}),
+            patch("examples.demand.evidence_pack_builder.query_source_elements", return_value=rows),
+            patch("examples.demand.evidence_pack_builder.query_case_ids_by_post_ids", return_value=[]),
+            patch("examples.demand.evidence_pack_builder.query_seed_points_for_sources", return_value=[]),
+        ):
+            result = self._build(
+                {
+                    "sources": [
+                        {
+                            "source_kind": "category_co_occurrence",
+                            "source_tool": "get_category_co_occurrences",
+                            "category_names": ["综合性腐败"],
+                        }
+                    ]
+                }
+            )
+
+        self.assertTrue(result["success"])
+        self.assertEqual(result["evidence_pack"]["source_kind"], "category_co_occurrence")
 
 
 if __name__ == "__main__":

+ 12 - 1
examples/demand/tests/test_pg_pattern_repository.py

@@ -1,7 +1,11 @@
 import unittest
 from unittest.mock import patch
 
-from examples.demand.pg_pattern_repository import _is_element_type_dimension, query_seed_points_for_itemsets
+from examples.demand.pg_pattern_repository import (
+    _is_element_type_dimension,
+    query_elements,
+    query_seed_points_for_itemsets,
+)
 
 
 class PgPatternRepositoryTest(unittest.TestCase):
@@ -75,6 +79,13 @@ class PgPatternRepositoryTest(unittest.TestCase):
         self.assertEqual(result[0]["rank"], 1)
         self.assertEqual(result[1]["rank"], 2)
 
+    def test_query_elements_empty_post_ids_does_not_fallback_to_global(self):
+        with patch("examples.demand.pg_pattern_repository._fetch_all") as fetch_all:
+            result = query_elements(581, post_ids=[], merge_leve2="历史名人", platform="piaoquan")
+
+        self.assertEqual(result, [])
+        fetch_all.assert_not_called()
+
 
 if __name__ == "__main__":
     unittest.main()

+ 22 - 7
examples/demand/validate_v2_demand_content.py

@@ -60,7 +60,6 @@ def validate_row(row: dict[str, Any], *, require_scope: bool = True) -> list[str
 
     expected = {
         "pattern_source_system": "pg_pattern_v2",
-        "source_kind": "pattern_itemset",
         "case_id_type": "post_id",
         "source_certainty": "db_validated",
         "validation_status": "passed",
@@ -72,9 +71,8 @@ def validate_row(row: dict[str, Any], *, require_scope: bool = True) -> list[str
     required_non_empty = [
         "source_post_id",
         "pattern_execution_id",
-        "mining_config_id",
-        "itemset_ids",
-        "itemset_items",
+        "source_kind",
+        "evidence_sources",
         "category_bindings",
         "element_bindings",
         "support",
@@ -89,8 +87,21 @@ def validate_row(row: dict[str, Any], *, require_scope: bool = True) -> list[str
         value = evidence_pack.get(field)
         if value is None or value == "" or value == []:
             errors.append(f"missing evidence_pack.{field}")
-    if not isinstance(evidence_pack.get("itemset_ids"), list) or len(evidence_pack["itemset_ids"]) != 1:
-        errors.append("itemset_ids must contain exactly one itemset_id")
+    if evidence_pack.get("source_kind") not in {
+        "high_weight_element",
+        "high_weight_category",
+        "element_co_occurrence",
+        "category_co_occurrence",
+        "pattern_itemset",
+        "multi_source",
+    }:
+        errors.append(f"invalid source_kind: {evidence_pack.get('source_kind')!r}")
+    if not isinstance(evidence_pack.get("evidence_sources"), list) or not evidence_pack["evidence_sources"]:
+        errors.append("missing evidence_sources")
+    if not isinstance(evidence_pack.get("itemset_ids", []), list):
+        errors.append("itemset_ids must be an array")
+    if not isinstance(evidence_pack.get("itemset_items", []), list):
+        errors.append("itemset_items must be an array")
 
     matched_post_ids = [str(post_id) for post_id in evidence_pack.get("matched_post_ids") or []]
     if str(evidence_pack.get("source_post_id") or "") not in set(matched_post_ids):
@@ -122,9 +133,13 @@ def validate_row(row: dict[str, Any], *, require_scope: bool = True) -> list[str
 
     itemset_items = evidence_pack.get("itemset_items") or []
     seed_candidates = _derive_seed_candidates(itemset_items)
+    for source in evidence_pack.get("evidence_sources") or []:
+        if isinstance(source, dict):
+            for value in source.get("source_terms") or []:
+                seed_candidates.add(str(value).replace(" ", ""))
     for term in evidence_pack.get("seed_terms") or []:
         if str(term).replace(" ", "") not in seed_candidates:
-            errors.append(f"seed_term not derived from itemset_items: {term}")
+            errors.append(f"seed_term not derived from DB evidence sources: {term}")
 
     demand_scope = evidence_pack.get("demand_scope")
     if require_scope and not isinstance(demand_scope, dict):