瀏覽代碼

Refactor query generation planning pipeline

SamLee 6 天之前
父節點
當前提交
310e04860d

+ 82 - 288
acquisition/queries/builder.py

@@ -1,269 +1,23 @@
-"""Formal creation-query batch builder."""
+"""Backward-compatible façade for the query planning bounded context."""
 from __future__ import annotations
 
-import json
-import random
-from dataclasses import dataclass
-from itertools import product
-from pathlib import Path
 from typing import Any
 
-from acquisition.queries.axes import ACTIONS, STAGES, _nonleaf_d4
 from acquisition.repositories.base import AcquisitionRepository
-from core.config import Settings
-
-ROOT = Path(__file__).resolve().parents[2]
-TREES = ROOT / "scope_trees" / "trees_index.json"
-KTYPE = ["怎么做", "有哪些", "为什么"]
-MODALITY = ["视频", "图文"]
-INTENT = ["灵感", "选题", "脚本"]
-DEFAULT_ACTIVE_FAMILY_KEYS = ("f1", "f2")
-
-
-@dataclass(frozen=True)
-class QueryBuildOptions:
-    per: int = 0
-    batch_n: int = 0
-    seed: int = 7
-    active_family_keys: tuple[str, ...] = DEFAULT_ACTIVE_FAMILY_KEYS
-
-
-def _segs(path: str | None) -> list[str]:
-    return [x for x in (path or "").split("/") if x]
-
-
-def _leaves(idx: list[dict[str, Any]], source_type: str, under: str | None = None) -> list[str]:
-    paths = [(_segs(n.get("path")), n.get("name")) for n in idx if n.get("source_type") == source_type]
-    if under:
-        paths = [(s, nm) for s, nm in paths if under in s]
-    all_paths = {"/".join(s) for s, _ in paths}
-    out: list[str] = []
-    seen: set[str] = set()
-    for segs, name in paths:
-        if len(segs) < 2:
-            continue
-        full = "/".join(segs)
-        is_leaf = not any(other != full and other.startswith(full + "/") for other in all_paths)
-        value = name or segs[-1]
-        if is_leaf and value and value not in seen:
-            seen.add(value)
-            out.append(value)
-    return out
-
-
-def _nodes_at_depth(
-    idx: list[dict[str, Any]],
-    source_type: str,
-    depths: tuple[int, ...] = (3, 4),
-    under: str | None = None,
-) -> list[str]:
-    paths = [(_segs(n.get("path")), n.get("name")) for n in idx if n.get("source_type") == source_type]
-    if under:
-        paths = [(s, nm) for s, nm in paths if under in s]
-    out: list[str] = []
-    seen: set[str] = set()
-    for segs, name in paths:
-        if len(segs) not in depths:
-            continue
-        value = name or (segs[-1] if segs else "")
-        if value and value not in seen:
-            seen.add(value)
-            out.append(value)
-    return out
-
-
-def _axis_tree(idx: list[dict[str, Any]], source_type: str) -> list[dict[str, Any]]:
-    nodes: dict[str, dict[str, Any]] = {}
-    for row in idx:
-        if row.get("source_type") != source_type:
-            continue
-        segs = _segs(row.get("path"))
-        if len(segs) not in (3, 4):
-            continue
-        path = "/" + "/".join(segs)
-        nodes[path] = {
-            "name": row.get("name") or segs[-1],
-            "path": path,
-            "level": len(segs),
-            "children": [],
-        }
-
-    roots: list[dict[str, Any]] = []
-    for path, node in nodes.items():
-        if node["level"] == 3:
-            roots.append(node)
-            continue
-        parent_path = "/" + "/".join(_segs(path)[:-1])
-        parent = nodes.get(parent_path)
-        if parent:
-            parent["children"].append(node)
-
-    roots.sort(key=lambda node: node["path"])
-    for node in roots:
-        node["children"].sort(key=lambda child: child["path"])
-    return roots
-
-
-def _sample_axis(values: list[str], limit: int, rng: random.Random) -> list[str]:
-    if limit <= 0 or limit >= len(values):
-        return values
-    return rng.sample(values, limit)
-
-
-def build_creation_query_batch(
-    settings: Settings,
-    *,
-    tree_path: Path = TREES,
-    options: QueryBuildOptions | None = None,
-) -> dict[str, Any]:
-    """Build a creation-query batch without writing files or database rows."""
-    opts = options or QueryBuildOptions()
-    rng = random.Random(opts.seed)
-    idx = json.loads(tree_path.read_text("utf-8"))
-    shi = _nodes_at_depth(idx, "实质", depths=(3, 4))
-    xing = _nodes_at_depth(idx, "形式", depths=(3, 4))
-    purpose_pool = _leaves(idx, "作用") + _leaves(idx, "感受") + _leaves(idx, "意图")
-    shi_batch = _sample_axis(shi, opts.batch_n, rng)
-    xing_batch = _sample_axis(xing, opts.batch_n, rng)
-    purpose_batch = _sample_axis(purpose_pool, opts.batch_n, rng)
-    f6_zy = _nonleaf_d4("作用", 10, tree_path=tree_path)
-    f6_stage_act = [(s, a) for s in STAGES for a in ACTIONS] + [("", "")]
-
-    if not shi_batch or not xing_batch or not purpose_batch or not f6_zy:
-        raise RuntimeError("scope tree does not contain enough creation query axes")
-
-    product_pools = {
-        "实质": shi_batch,
-        "形式": xing_batch,
-        "目的": purpose_batch,
-        "模态": MODALITY,
-        "业务阶段": INTENT,
-        "知识类型": KTYPE,
-    }
-
-    def gen_product(keys: list[str]):
-        for values in product(*(product_pools[key] for key in keys)):
-            yield {"parts": dict(zip(keys, values, strict=True))}
-
-    legacy_per = opts.per if opts.per > 0 else 30
-    master: list[dict[str, str]] = []
-    for i in range(legacy_per):
-        stage, action = f6_stage_act[i % len(f6_stage_act)]
-        master.append(
-            {
-                "实质": shi_batch[i % len(shi_batch)],
-                "形式": xing_batch[i % len(xing_batch)],
-                "目的": purpose_batch[i % len(purpose_batch)],
-                "模态": MODALITY[i % len(MODALITY)],
-                "业务阶段": INTENT[i % len(INTENT)],
-                "知识类型": KTYPE[(i // 3) % len(KTYPE)],
-                "阶段": stage or "/",
-                "动作": action or "/",
-                "_st": stage,
-                "_ac": action,
-                "作用": f6_zy[i % len(f6_zy)],
-            }
-        )
-
-    def gen_old(
-        i: int,
-        *,
-        shi_axis: bool = False,
-        xing_axis: bool = False,
-        purpose_axis: bool = False,
-        effect_axis: bool = True,
-    ) -> dict[str, Any]:
-        row = master[i]
-        segment = (row["_st"] + row["_ac"]) if row["_ac"] else ""
-        head = (
-            ([row["实质"]] if shi_axis else [])
-            + ([row["形式"]] if xing_axis else [])
-            + ([row["目的"]] if purpose_axis else [])
-        )
-        tail = ([row["作用"]] if effect_axis else []) + [row["知识类型"]]
-        query = " ".join(head + ([segment] if segment else []) + tail)
-        parts: dict[str, str] = {}
-        if shi_axis:
-            parts["实质"] = row["实质"]
-        if xing_axis:
-            parts["形式"] = row["形式"]
-        if purpose_axis:
-            parts["目的"] = row["目的"]
-        parts["阶段"], parts["动作"] = row["阶段"], row["动作"]
-        if effect_axis:
-            parts["作用"] = row["作用"]
-        parts["知识类型"] = row["知识类型"]
-        return {"parts": parts, "query": query}
-
-    families = [
-        {"key": "f1", "axes": ["实质", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["实质", "模态", "业务阶段", "知识类型"])},
-        {"key": "f2", "axes": ["形式", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["形式", "模态", "业务阶段", "知识类型"])},
-        {"key": "f4", "axes": ["作用/感受/意图", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["目的", "模态", "业务阶段", "知识类型"])},
-        {"key": "f3", "axes": ["实质", "形式", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["实质", "形式", "模态", "业务阶段", "知识类型"])},
-        {"key": "f5", "axes": ["模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["模态", "业务阶段", "知识类型"])},
-        {"key": "a_shi", "axes": ["实质", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True)},
-        {"key": "a_xing", "axes": ["形式", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, xing_axis=True)},
-        {"key": "a_both", "axes": ["实质", "形式", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, xing_axis=True)},
-        {"key": "a_purpose", "axes": ["作用/感受/意图", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, purpose_axis=True)},
-        {"key": "a_tail", "axes": ["阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i)},
-        {"key": "b_shi", "axes": ["实质", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, effect_axis=False)},
-        {"key": "b_xing", "axes": ["形式", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, xing_axis=True, effect_axis=False)},
-        {"key": "b_both", "axes": ["实质", "形式", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, xing_axis=True, effect_axis=False)},
-        {"key": "b_purpose", "axes": ["作用/感受/意图", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, purpose_axis=True, effect_axis=False)},
-        {"key": "b_tail", "axes": ["阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, effect_axis=False)},
-    ]
-    order = ["实质", "形式", "目的", "模态", "业务阶段", "知识类型"]
-
-    out: dict[str, Any] = {
-        "axis_values": {
-            "实质": shi,
-            "形式": xing,
-            "目的池": purpose_pool,
-            "业务阶段": INTENT,
-            "模态": MODALITY,
-            "知识类型": KTYPE,
-            "阶段": STAGES,
-            "动作": ACTIONS,
-            "作用": f6_zy,
-        },
-        "axis_trees": {
-            "实质": _axis_tree(idx, "实质"),
-            "形式": _axis_tree(idx, "形式"),
-        },
-        "metadata": {
-            "seed": opts.seed,
-            "per": opts.per,
-            "batch_n": opts.batch_n,
-            "active_family_keys": list(opts.active_family_keys),
-        },
-        "families": [],
-    }
-    family_by_key = {family["key"]: family for family in families}
-    unknown = [key for key in opts.active_family_keys if key not in family_by_key]
-    if unknown:
-        raise ValueError(f"unknown query family key(s): {', '.join(unknown)}")
-
-    for family_key in opts.active_family_keys:
-        family = family_by_key[family_key]
-        name = " × ".join(family["axes"])
-        seen: set[str] = set()
-        items: list[dict[str, Any]] = []
-        generated_items = family["items"]() if "items" in family else (family["gen"](i) for i in range(legacy_per))
-        for generated in generated_items:
-            parts = generated["parts"]
-            query = generated.get("query") or " ".join(parts[k] for k in order if k in parts)
-            if query in seen:
-                continue
-            seen.add(query)
-            items.append({"query": query, "parts": parts})
-            if opts.per > 0 and len(items) >= opts.per:
-                break
-        for item in items:
-            item.update({"keep": True, "reason": ""})
-        out["families"].append(
-            {"key": family["key"], "name": name, "axes": family["axes"], "items": items}
-        )
-    return out
+from query_planning import (
+    GenerationRequest,
+    GeneratorKind,
+    QueryBatchWriteSpec,
+    QueryBatchWriter,
+    UnifiedQueryGenerationService,
+    planning_store_for_repository,
+)
+from query_planning.cartesian import (
+    DEFAULT_ACTIVE_FAMILY_KEYS,
+    TREES,
+    QueryBuildOptions,
+    build_creation_query_batch,
+)
 
 
 def persist_query_batch(
@@ -275,33 +29,73 @@ def persist_query_batch(
     generation_method: str = "creation_demo_v1",
     target_platforms: list[str] | None = None,
 ) -> tuple[Any, int]:
-    """Persist generated families into formal query batch/query rows."""
-    batch = repo.create_query_batch(
-        name=name,
-        source_type=source_type,
-        generation_method=generation_method,
-        target_platforms=target_platforms or ["xiaohongshu", "weixin", "douyin"],
-        status="ready",
-        metadata=generated.get("metadata") or {},
-    )
-    count = 0
-    sort_order = 0
+    """Persist legacy family JSON through the unified planning and writer flow."""
+    candidates: list[dict[str, Any]] = []
     for family in generated.get("families") or []:
         for item in family.get("items") or []:
-            sort_order += 1
-            repo.add_query(
-                batch_id=batch.id,
-                query_text=item["query"],
-                axes=item.get("parts") or {},
-                keep=bool(item.get("keep", True)),
-                filter_reason=item.get("reason") or "",
-                status="ready",
-                sort_order=sort_order,
-                metadata={
-                    "family_key": family.get("key"),
-                    "family_name": family.get("name"),
-                    "family_axes": family.get("axes") or [],
-                },
+            if not item.get("keep", True):
+                continue
+            candidates.append(
+                {
+                    "query_text": item.get("query") or "",
+                    "axes": item.get("parts") or {},
+                    "filter_reason": item.get("reason") or "",
+                    "priority": item.get("priority") or 0,
+                    "source_refs": [
+                        {
+                            "generator_kind": GeneratorKind.CARTESIAN.value,
+                            "family_key": family.get("key"),
+                            "family_name": family.get("name"),
+                            "axes": item.get("parts") or {},
+                        }
+                    ],
+                    "metadata": {
+                        "family_key": family.get("key"),
+                        "family_name": family.get("name"),
+                        "family_axes": family.get("axes") or [],
+                    },
+                }
             )
-            count += 1
-    return batch, count
+    platforms = tuple(target_platforms or ["xiaohongshu", "weixin", "douyin"])
+    result = UnifiedQueryGenerationService().generate(
+        GenerationRequest(
+            generator_kind=GeneratorKind.CARTESIAN,
+            name=name,
+            target_platforms=platforms,
+            payload={
+                "candidates": candidates,
+                "generator_config": generated.get("metadata") or {},
+                "input_snapshot": {
+                    "active_family_keys": (generated.get("metadata") or {}).get(
+                        "active_family_keys", []
+                    ),
+                    "family_count": len(generated.get("families") or []),
+                    "candidate_count": len(candidates),
+                },
+            },
+            metadata={"generation_method": generation_method},
+        )
+    )
+    write_result = QueryBatchWriter(
+        legacy_sink=repo,
+        planning_store=planning_store_for_repository(repo),
+    ).write(
+        result,
+        QueryBatchWriteSpec(
+            name=name,
+            source_type=source_type,
+            generation_method=generation_method,
+            target_platforms=platforms,
+            metadata=generated.get("metadata") or {},
+        ),
+    )
+    return write_result.batch, len(write_result.queries)
+
+
+__all__ = [
+    "DEFAULT_ACTIVE_FAMILY_KEYS",
+    "QueryBuildOptions",
+    "TREES",
+    "build_creation_query_batch",
+    "persist_query_batch",
+]

+ 63 - 33
app/routes/manual_queries.py

@@ -19,6 +19,14 @@ from core.config import CreationDbConfig
 from core.db_session import transaction
 from pipeline.postgres import PostgresPipelineRepository
 from pipeline.tracing import TraceContext, new_trace_writer
+from query_planning import (
+    GenerationRequest,
+    GeneratorKind,
+    QueryBatchWriteSpec,
+    QueryBatchWriter,
+    UnifiedQueryGenerationService,
+    planning_store_for_repository,
+)
 
 router = APIRouter(prefix="/api/query-batches", tags=["manual-query-batches"])
 
@@ -112,19 +120,9 @@ def _normalize_entries(request: ManualQueryBatchRequest) -> list[ManualQueryEntr
             if entry:
                 entries.append(entry)
 
-    deduped: list[ManualQueryEntry] = []
-    seen: set[str] = set()
-    for entry in entries:
-        key = entry.query_text
-        if key in seen:
-            continue
-        seen.add(key)
-        deduped.append(entry)
-        if len(deduped) > QUERY_MAX_COUNT:
-            raise HTTPException(status_code=422, detail=f"一次最多提交 {QUERY_MAX_COUNT} 条 query")
-    if not deduped:
+    if not entries:
         raise HTTPException(status_code=422, detail="没有可提交的 query")
-    return deduped
+    return entries
 
 
 def _normalize_platforms(platforms: list[str] | None) -> list[str]:
@@ -216,37 +214,69 @@ def create_manual_query_batch(
     timestamp = datetime.now().strftime("%Y%m%d-%H%M%S")
     name = request.name or f"manual-query-{timestamp}"
 
+    generation_result = UnifiedQueryGenerationService().generate(
+        GenerationRequest(
+            generator_kind=GeneratorKind.MANUAL,
+            name=name,
+            target_platforms=tuple(platforms),
+            payload={
+                "candidates": [
+                    {
+                        "query_text": entry.query_text,
+                        "axes": entry.axes,
+                        "filter_reason": entry.filter_reason,
+                        "metadata": entry.metadata,
+                        "source_refs": [
+                            {
+                                "generator_kind": GeneratorKind.MANUAL.value,
+                                "original_position": index,
+                                "axes": entry.axes,
+                                "source_family_key": entry.metadata.get("source_family_key"),
+                                "source_family_name": entry.metadata.get("source_family_name"),
+                            }
+                        ],
+                    }
+                    for index, entry in enumerate(entries)
+                ],
+                "generator_config": {
+                    "request_shape": "manual_query_api_v1",
+                },
+                "input_snapshot": {
+                    "submitted_count": len(entries),
+                    "target_platforms": platforms,
+                },
+            },
+            metadata={"source": "manual_api"},
+        )
+    )
+    if generation_result.stats.selected_count > QUERY_MAX_COUNT:
+        raise HTTPException(status_code=422, detail=f"一次最多提交 {QUERY_MAX_COUNT} 条 query")
+
     batch_metadata = {
         "family_key": "manual",
         "source": "manual_api",
-        "query_count": len(entries),
+        "query_count": generation_result.stats.selected_count,
         "platforms": platforms,
         **request.metadata,
     }
 
     with transaction(db_config) as conn:
         repo = PostgresAcquisitionRepository(conn)
-        batch = repo.create_query_batch(
-            name=name,
-            source_type="manual",
-            generation_method="manual_query_api_v1",
-            target_platforms=platforms,
-            status="ready",
-            metadata=batch_metadata,
+        write_result = QueryBatchWriter(
+            legacy_sink=repo,
+            planning_store=planning_store_for_repository(repo),
+        ).write(
+            generation_result,
+            QueryBatchWriteSpec(
+                name=name,
+                source_type="manual",
+                generation_method="manual_query_api_v1",
+                target_platforms=tuple(platforms),
+                metadata=batch_metadata,
+            ),
         )
-        queries = [
-            repo.add_query(
-                batch_id=batch.id,
-                query_text=entry.query_text,
-                axes=entry.axes,
-                keep=True,
-                filter_reason=entry.filter_reason,
-                status="ready",
-                sort_order=index,
-                metadata=entry.metadata,
-            )
-            for index, entry in enumerate(entries)
-        ]
+        batch = write_result.batch
+        queries = write_result.queries
         run_key = f"manual-api:{batch.id}:{timestamp}"
         log_path = RUNTIME_MANUAL_DIR / f"{run_key.replace(':', '-')}.log"
         cmd = _pipeline_command(

+ 89 - 0
db/migrations/005_query_planning_schema.sql

@@ -0,0 +1,89 @@
+-- 005: add query planning state before the frozen acquisition boundary.
+
+CREATE TABLE IF NOT EXISTS creation_knowledge.query_plans (
+    id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
+    generator_kind text NOT NULL CHECK (
+        generator_kind IN ('cartesian', 'manual', 'topic_table', 'agent_plan', 'history_gap')
+    ),
+    status text NOT NULL DEFAULT 'ready' CHECK (status IN ('planning', 'ready', 'failed')),
+    input_snapshot jsonb NOT NULL DEFAULT '{}'::jsonb,
+    generator_config jsonb NOT NULL DEFAULT '{}'::jsonb,
+    max_queries integer CHECK (max_queries IS NULL OR max_queries >= 0),
+    generated_count integer NOT NULL DEFAULT 0 CHECK (generated_count >= 0),
+    normalized_count integer NOT NULL DEFAULT 0 CHECK (normalized_count >= 0),
+    unique_count integer NOT NULL DEFAULT 0 CHECK (unique_count >= 0),
+    selected_count integer NOT NULL DEFAULT 0 CHECK (selected_count >= 0),
+    dropped_count integer NOT NULL DEFAULT 0 CHECK (dropped_count >= 0),
+    metadata jsonb NOT NULL DEFAULT '{}'::jsonb,
+    created_at timestamptz NOT NULL DEFAULT now(),
+    updated_at timestamptz NOT NULL DEFAULT now(),
+    CHECK (normalized_count <= generated_count),
+    CHECK (unique_count <= normalized_count),
+    CHECK (selected_count <= unique_count),
+    CHECK (dropped_count = generated_count - selected_count)
+);
+
+CREATE TABLE IF NOT EXISTS creation_knowledge.knowledge_needs (
+    id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
+    plan_id uuid NOT NULL REFERENCES creation_knowledge.query_plans(id) ON DELETE CASCADE,
+    need_key text NOT NULL,
+    decision_context text NOT NULL,
+    unknown_information text NOT NULL,
+    source_ref jsonb NOT NULL DEFAULT '{}'::jsonb,
+    priority integer NOT NULL DEFAULT 0,
+    metadata jsonb NOT NULL DEFAULT '{}'::jsonb,
+    created_at timestamptz NOT NULL DEFAULT now(),
+    UNIQUE (plan_id, need_key)
+);
+
+CREATE TABLE IF NOT EXISTS creation_knowledge.query_plan_batches (
+    plan_id uuid NOT NULL REFERENCES creation_knowledge.query_plans(id) ON DELETE CASCADE,
+    batch_id uuid NOT NULL REFERENCES creation_knowledge.query_batches(id) ON DELETE CASCADE,
+    created_at timestamptz NOT NULL DEFAULT now(),
+    PRIMARY KEY (plan_id, batch_id)
+);
+
+CREATE TABLE IF NOT EXISTS creation_knowledge.query_knowledge_need_links (
+    query_id uuid NOT NULL REFERENCES creation_knowledge.queries(id) ON DELETE CASCADE,
+    knowledge_need_id uuid NOT NULL REFERENCES creation_knowledge.knowledge_needs(id) ON DELETE CASCADE,
+    created_at timestamptz NOT NULL DEFAULT now(),
+    PRIMARY KEY (query_id, knowledge_need_id)
+);
+
+CREATE INDEX IF NOT EXISTS idx_query_plans_kind_created
+ON creation_knowledge.query_plans(generator_kind, created_at);
+
+CREATE INDEX IF NOT EXISTS idx_knowledge_needs_plan_priority
+ON creation_knowledge.knowledge_needs(plan_id, priority DESC);
+
+CREATE INDEX IF NOT EXISTS idx_query_plan_batches_batch
+ON creation_knowledge.query_plan_batches(batch_id);
+
+CREATE INDEX IF NOT EXISTS idx_query_knowledge_need_links_need
+ON creation_knowledge.query_knowledge_need_links(knowledge_need_id);
+
+DROP TRIGGER IF EXISTS trg_query_plans_touch_updated_at
+ON creation_knowledge.query_plans;
+
+CREATE TRIGGER trg_query_plans_touch_updated_at
+BEFORE UPDATE ON creation_knowledge.query_plans
+FOR EACH ROW EXECUTE FUNCTION creation_knowledge.touch_updated_at();
+
+DO $$
+BEGIN
+    IF EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'ck_app') THEN
+        GRANT SELECT, INSERT, UPDATE, DELETE ON
+            creation_knowledge.query_plans,
+            creation_knowledge.knowledge_needs,
+            creation_knowledge.query_plan_batches,
+            creation_knowledge.query_knowledge_need_links
+        TO ck_app;
+    END IF;
+END;
+$$;
+
+INSERT INTO creation_knowledge.schema_migrations(version, description)
+VALUES ('005_query_planning_schema', 'Add query plans, knowledge needs, and legacy query links')
+ON CONFLICT (version) DO UPDATE SET
+    description = EXCLUDED.description,
+    applied_at = now();

+ 53 - 0
query_planning/__init__.py

@@ -0,0 +1,53 @@
+"""Query planning bounded context public API."""
+
+from query_planning.domain import (
+    GenerationPlanDraft,
+    GenerationRequest,
+    GenerationResult,
+    GenerationStats,
+    GeneratorKind,
+    GeneratorOutput,
+    KnowledgeNeedDraft,
+    QueryCandidate,
+)
+from query_planning.errors import NoQueryCandidatesError, UnsupportedGeneratorError
+from query_planning.generators import (
+    CartesianQueryGenerator,
+    ManualQueryGenerator,
+    QueryGeneratorRegistry,
+    default_generator_registry,
+)
+from query_planning.ports import LegacyQueryBatchSink, QueryGenerator, QueryPlanningStore
+from query_planning.service import UnifiedQueryGenerationService, normalize_query_text
+from query_planning.writer import (
+    QueryBatchWriteResult,
+    QueryBatchWriteSpec,
+    QueryBatchWriter,
+    planning_store_for_repository,
+)
+
+__all__ = [
+    "CartesianQueryGenerator",
+    "GenerationPlanDraft",
+    "GenerationRequest",
+    "GenerationResult",
+    "GenerationStats",
+    "GeneratorKind",
+    "GeneratorOutput",
+    "KnowledgeNeedDraft",
+    "LegacyQueryBatchSink",
+    "ManualQueryGenerator",
+    "NoQueryCandidatesError",
+    "QueryBatchWriteResult",
+    "QueryBatchWriteSpec",
+    "QueryBatchWriter",
+    "QueryCandidate",
+    "QueryGenerator",
+    "QueryGeneratorRegistry",
+    "QueryPlanningStore",
+    "UnifiedQueryGenerationService",
+    "UnsupportedGeneratorError",
+    "default_generator_registry",
+    "normalize_query_text",
+    "planning_store_for_repository",
+]

+ 316 - 0
query_planning/cartesian.py

@@ -0,0 +1,316 @@
+"""Cartesian creation-query generator and legacy preview projection."""
+from __future__ import annotations
+
+import json
+import random
+from dataclasses import dataclass
+from itertools import product
+from pathlib import Path
+from typing import Any
+
+from acquisition.queries.axes import ACTIONS, STAGES, _nonleaf_d4
+from core.config import Settings
+from query_planning import (
+    GenerationRequest,
+    GeneratorKind,
+    UnifiedQueryGenerationService,
+)
+
+ROOT = Path(__file__).resolve().parents[1]
+TREES = ROOT / "scope_trees" / "trees_index.json"
+KTYPE = ["怎么做", "有哪些", "为什么"]
+MODALITY = ["视频", "图文"]
+INTENT = ["灵感", "选题", "脚本"]
+DEFAULT_ACTIVE_FAMILY_KEYS = ("f1", "f2")
+
+
+def _planning_service() -> UnifiedQueryGenerationService:
+    return UnifiedQueryGenerationService()
+
+
+@dataclass(frozen=True)
+class QueryBuildOptions:
+    per: int = 0
+    batch_n: int = 0
+    seed: int = 7
+    active_family_keys: tuple[str, ...] = DEFAULT_ACTIVE_FAMILY_KEYS
+
+
+def _segs(path: str | None) -> list[str]:
+    return [x for x in (path or "").split("/") if x]
+
+
+def _leaves(idx: list[dict[str, Any]], source_type: str, under: str | None = None) -> list[str]:
+    paths = [(_segs(n.get("path")), n.get("name")) for n in idx if n.get("source_type") == source_type]
+    if under:
+        paths = [(s, nm) for s, nm in paths if under in s]
+    all_paths = {"/".join(s) for s, _ in paths}
+    out: list[str] = []
+    seen: set[str] = set()
+    for segs, name in paths:
+        if len(segs) < 2:
+            continue
+        full = "/".join(segs)
+        is_leaf = not any(other != full and other.startswith(full + "/") for other in all_paths)
+        value = name or segs[-1]
+        if is_leaf and value and value not in seen:
+            seen.add(value)
+            out.append(value)
+    return out
+
+
+def _nodes_at_depth(
+    idx: list[dict[str, Any]],
+    source_type: str,
+    depths: tuple[int, ...] = (3, 4),
+    under: str | None = None,
+) -> list[str]:
+    paths = [(_segs(n.get("path")), n.get("name")) for n in idx if n.get("source_type") == source_type]
+    if under:
+        paths = [(s, nm) for s, nm in paths if under in s]
+    out: list[str] = []
+    seen: set[str] = set()
+    for segs, name in paths:
+        if len(segs) not in depths:
+            continue
+        value = name or (segs[-1] if segs else "")
+        if value and value not in seen:
+            seen.add(value)
+            out.append(value)
+    return out
+
+
+def _axis_tree(idx: list[dict[str, Any]], source_type: str) -> list[dict[str, Any]]:
+    nodes: dict[str, dict[str, Any]] = {}
+    for row in idx:
+        if row.get("source_type") != source_type:
+            continue
+        segs = _segs(row.get("path"))
+        if len(segs) not in (3, 4):
+            continue
+        path = "/" + "/".join(segs)
+        nodes[path] = {
+            "name": row.get("name") or segs[-1],
+            "path": path,
+            "level": len(segs),
+            "children": [],
+        }
+
+    roots: list[dict[str, Any]] = []
+    for path, node in nodes.items():
+        if node["level"] == 3:
+            roots.append(node)
+            continue
+        parent_path = "/" + "/".join(_segs(path)[:-1])
+        parent = nodes.get(parent_path)
+        if parent:
+            parent["children"].append(node)
+
+    roots.sort(key=lambda node: node["path"])
+    for node in roots:
+        node["children"].sort(key=lambda child: child["path"])
+    return roots
+
+
+def _sample_axis(values: list[str], limit: int, rng: random.Random) -> list[str]:
+    if limit <= 0 or limit >= len(values):
+        return values
+    return rng.sample(values, limit)
+
+
+def build_creation_query_batch(
+    settings: Settings,
+    *,
+    tree_path: Path = TREES,
+    options: QueryBuildOptions | None = None,
+) -> dict[str, Any]:
+    """Build a creation-query batch without writing files or database rows."""
+    opts = options or QueryBuildOptions()
+    rng = random.Random(opts.seed)
+    idx = json.loads(tree_path.read_text("utf-8"))
+    shi = _nodes_at_depth(idx, "实质", depths=(3, 4))
+    xing = _nodes_at_depth(idx, "形式", depths=(3, 4))
+    purpose_pool = _leaves(idx, "作用") + _leaves(idx, "感受") + _leaves(idx, "意图")
+    shi_batch = _sample_axis(shi, opts.batch_n, rng)
+    xing_batch = _sample_axis(xing, opts.batch_n, rng)
+    purpose_batch = _sample_axis(purpose_pool, opts.batch_n, rng)
+    f6_zy = _nonleaf_d4("作用", 10, tree_path=tree_path)
+    f6_stage_act = [(s, a) for s in STAGES for a in ACTIONS] + [("", "")]
+
+    if not shi_batch or not xing_batch or not purpose_batch or not f6_zy:
+        raise RuntimeError("scope tree does not contain enough creation query axes")
+
+    product_pools = {
+        "实质": shi_batch,
+        "形式": xing_batch,
+        "目的": purpose_batch,
+        "模态": MODALITY,
+        "业务阶段": INTENT,
+        "知识类型": KTYPE,
+    }
+
+    def gen_product(keys: list[str]):
+        for values in product(*(product_pools[key] for key in keys)):
+            yield {"parts": dict(zip(keys, values, strict=True))}
+
+    legacy_per = opts.per if opts.per > 0 else 30
+    master: list[dict[str, str]] = []
+    for i in range(legacy_per):
+        stage, action = f6_stage_act[i % len(f6_stage_act)]
+        master.append(
+            {
+                "实质": shi_batch[i % len(shi_batch)],
+                "形式": xing_batch[i % len(xing_batch)],
+                "目的": purpose_batch[i % len(purpose_batch)],
+                "模态": MODALITY[i % len(MODALITY)],
+                "业务阶段": INTENT[i % len(INTENT)],
+                "知识类型": KTYPE[(i // 3) % len(KTYPE)],
+                "阶段": stage or "/",
+                "动作": action or "/",
+                "_st": stage,
+                "_ac": action,
+                "作用": f6_zy[i % len(f6_zy)],
+            }
+        )
+
+    def gen_old(
+        i: int,
+        *,
+        shi_axis: bool = False,
+        xing_axis: bool = False,
+        purpose_axis: bool = False,
+        effect_axis: bool = True,
+    ) -> dict[str, Any]:
+        row = master[i]
+        segment = (row["_st"] + row["_ac"]) if row["_ac"] else ""
+        head = (
+            ([row["实质"]] if shi_axis else [])
+            + ([row["形式"]] if xing_axis else [])
+            + ([row["目的"]] if purpose_axis else [])
+        )
+        tail = ([row["作用"]] if effect_axis else []) + [row["知识类型"]]
+        query = " ".join(head + ([segment] if segment else []) + tail)
+        parts: dict[str, str] = {}
+        if shi_axis:
+            parts["实质"] = row["实质"]
+        if xing_axis:
+            parts["形式"] = row["形式"]
+        if purpose_axis:
+            parts["目的"] = row["目的"]
+        parts["阶段"], parts["动作"] = row["阶段"], row["动作"]
+        if effect_axis:
+            parts["作用"] = row["作用"]
+        parts["知识类型"] = row["知识类型"]
+        return {"parts": parts, "query": query}
+
+    families = [
+        {"key": "f1", "axes": ["实质", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["实质", "模态", "业务阶段", "知识类型"])},
+        {"key": "f2", "axes": ["形式", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["形式", "模态", "业务阶段", "知识类型"])},
+        {"key": "f4", "axes": ["作用/感受/意图", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["目的", "模态", "业务阶段", "知识类型"])},
+        {"key": "f3", "axes": ["实质", "形式", "模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["实质", "形式", "模态", "业务阶段", "知识类型"])},
+        {"key": "f5", "axes": ["模态", "业务阶段", "知识类型"], "items": lambda: gen_product(["模态", "业务阶段", "知识类型"])},
+        {"key": "a_shi", "axes": ["实质", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True)},
+        {"key": "a_xing", "axes": ["形式", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, xing_axis=True)},
+        {"key": "a_both", "axes": ["实质", "形式", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, xing_axis=True)},
+        {"key": "a_purpose", "axes": ["作用/感受/意图", "阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i, purpose_axis=True)},
+        {"key": "a_tail", "axes": ["阶段", "动作", "作用", "知识类型"], "gen": lambda i: gen_old(i)},
+        {"key": "b_shi", "axes": ["实质", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, effect_axis=False)},
+        {"key": "b_xing", "axes": ["形式", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, xing_axis=True, effect_axis=False)},
+        {"key": "b_both", "axes": ["实质", "形式", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, shi_axis=True, xing_axis=True, effect_axis=False)},
+        {"key": "b_purpose", "axes": ["作用/感受/意图", "阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, purpose_axis=True, effect_axis=False)},
+        {"key": "b_tail", "axes": ["阶段", "动作", "知识类型"], "gen": lambda i: gen_old(i, effect_axis=False)},
+    ]
+    order = ["实质", "形式", "目的", "模态", "业务阶段", "知识类型"]
+
+    out: dict[str, Any] = {
+        "axis_values": {
+            "实质": shi,
+            "形式": xing,
+            "目的池": purpose_pool,
+            "业务阶段": INTENT,
+            "模态": MODALITY,
+            "知识类型": KTYPE,
+            "阶段": STAGES,
+            "动作": ACTIONS,
+            "作用": f6_zy,
+        },
+        "axis_trees": {
+            "实质": _axis_tree(idx, "实质"),
+            "形式": _axis_tree(idx, "形式"),
+        },
+        "metadata": {
+            "seed": opts.seed,
+            "per": opts.per,
+            "batch_n": opts.batch_n,
+            "active_family_keys": list(opts.active_family_keys),
+        },
+        "families": [],
+    }
+    family_by_key = {family["key"]: family for family in families}
+    unknown = [key for key in opts.active_family_keys if key not in family_by_key]
+    if unknown:
+        raise ValueError(f"unknown query family key(s): {', '.join(unknown)}")
+
+    for family_key in opts.active_family_keys:
+        family = family_by_key[family_key]
+        name = " × ".join(family["axes"])
+        seen: set[str] = set()
+        items: list[dict[str, Any]] = []
+        generated_items = family["items"]() if "items" in family else (family["gen"](i) for i in range(legacy_per))
+        for generated in generated_items:
+            parts = generated["parts"]
+            query = generated.get("query") or " ".join(parts[k] for k in order if k in parts)
+            if query in seen:
+                continue
+            seen.add(query)
+            items.append({"query": query, "parts": parts})
+            if opts.per > 0 and len(items) >= opts.per:
+                break
+        for item in items:
+            item.update({"keep": True, "reason": ""})
+        planned = _planning_service().generate(
+            GenerationRequest(
+                generator_kind=GeneratorKind.CARTESIAN,
+                name=f"cartesian-preview:{family['key']}",
+                payload={
+                    "candidates": [
+                        {
+                            "query_text": item["query"],
+                            "axes": item["parts"],
+                            "filter_reason": item["reason"],
+                            "metadata": {
+                                "family_key": family["key"],
+                                "family_name": name,
+                                "family_axes": family["axes"],
+                            },
+                            "source_refs": [
+                                {
+                                    "generator_kind": GeneratorKind.CARTESIAN.value,
+                                    "family_key": family["key"],
+                                    "axes": item["parts"],
+                                }
+                            ],
+                        }
+                        for item in items
+                    ],
+                    "input_snapshot": {
+                        "family_key": family["key"],
+                        "family_axes": family["axes"],
+                        "options": out["metadata"],
+                    },
+                },
+            )
+        )
+        items = [
+            {
+                "query": candidate.query_text,
+                "parts": candidate.axes,
+                "keep": True,
+                "reason": candidate.filter_reason or "",
+            }
+            for candidate in planned.selected_queries
+        ]
+        out["families"].append(
+            {"key": family["key"], "name": name, "axes": family["axes"], "items": items}
+        )
+    return out

+ 88 - 0
query_planning/domain.py

@@ -0,0 +1,88 @@
+"""Framework-free domain values for planning search queries."""
+from __future__ import annotations
+
+from dataclasses import dataclass, field
+from enum import StrEnum
+from typing import Any
+
+
+class GeneratorKind(StrEnum):
+    CARTESIAN = "cartesian"
+    MANUAL = "manual"
+    TOPIC_TABLE = "topic_table"
+    AGENT_PLAN = "agent_plan"
+    HISTORY_GAP = "history_gap"
+
+
+@dataclass(frozen=True)
+class GenerationRequest:
+    generator_kind: GeneratorKind
+    payload: dict[str, Any]
+    name: str
+    target_platforms: tuple[str, ...] = ("xiaohongshu", "weixin", "douyin")
+    max_queries: int | None = None
+    metadata: dict[str, Any] = field(default_factory=dict)
+
+    def __post_init__(self) -> None:
+        if not self.name.strip():
+            raise ValueError("generation request name must not be empty")
+        if not self.target_platforms:
+            raise ValueError("at least one target platform is required")
+        if self.max_queries is not None and self.max_queries < 0:
+            raise ValueError("max_queries must be non-negative or None")
+
+
+@dataclass(frozen=True)
+class KnowledgeNeedDraft:
+    need_key: str
+    decision_context: str
+    unknown_information: str
+    source_ref: dict[str, Any] = field(default_factory=dict)
+    priority: int = 0
+    metadata: dict[str, Any] = field(default_factory=dict)
+
+
+@dataclass
+class QueryCandidate:
+    query_text: str
+    axes: dict[str, Any] = field(default_factory=dict)
+    priority: int = 0
+    original_position: int = 0
+    source_refs: list[dict[str, Any]] = field(default_factory=list)
+    knowledge_need_keys: list[str] = field(default_factory=list)
+    metadata: dict[str, Any] = field(default_factory=dict)
+    filter_reason: str | None = None
+    dedupe_key: str = ""
+
+
+@dataclass(frozen=True)
+class GeneratorOutput:
+    candidates: tuple[QueryCandidate, ...]
+    knowledge_needs: tuple[KnowledgeNeedDraft, ...] = ()
+    generator_config: dict[str, Any] = field(default_factory=dict)
+
+
+@dataclass(frozen=True)
+class GenerationPlanDraft:
+    generator_kind: GeneratorKind
+    input_snapshot: dict[str, Any]
+    generator_config: dict[str, Any]
+    metadata: dict[str, Any]
+
+
+@dataclass(frozen=True)
+class GenerationStats:
+    generated_count: int
+    normalized_count: int
+    unique_count: int
+    selected_count: int
+    dropped_count: int
+
+
+@dataclass(frozen=True)
+class GenerationResult:
+    request: GenerationRequest
+    plan: GenerationPlanDraft
+    knowledge_needs: tuple[KnowledgeNeedDraft, ...]
+    selected_queries: tuple[QueryCandidate, ...]
+    stats: GenerationStats

+ 9 - 0
query_planning/errors.py

@@ -0,0 +1,9 @@
+"""Query planning errors."""
+
+
+class UnsupportedGeneratorError(ValueError):
+    """Raised when a generator kind is reserved but not registered."""
+
+
+class NoQueryCandidatesError(ValueError):
+    """Raised when normalization, dedupe, or budget leaves no query."""

+ 112 - 0
query_planning/generators.py

@@ -0,0 +1,112 @@
+"""Generator strategies and registry."""
+from __future__ import annotations
+
+from collections.abc import Iterable
+from typing import Any
+
+from query_planning.domain import (
+    GenerationRequest,
+    GeneratorKind,
+    GeneratorOutput,
+    KnowledgeNeedDraft,
+    QueryCandidate,
+)
+from query_planning.errors import UnsupportedGeneratorError
+from query_planning.ports import QueryGenerator
+
+
+def _candidate(value: Any, position: int, kind: GeneratorKind) -> QueryCandidate:
+    if isinstance(value, QueryCandidate):
+        if value.original_position == 0 and position:
+            value.original_position = position
+        return value
+    if isinstance(value, str):
+        return QueryCandidate(
+            query_text=value,
+            original_position=position,
+            source_refs=[{"generator_kind": kind.value}],
+        )
+    if not isinstance(value, dict):
+        return QueryCandidate(query_text="", original_position=position)
+    source_refs = value.get("source_refs") or []
+    if not source_refs:
+        source_refs = [{"generator_kind": kind.value}]
+    filter_reason = (
+        value.get("filter_reason")
+        if "filter_reason" in value
+        else value.get("reason")
+    )
+    return QueryCandidate(
+        query_text=str(value.get("query_text") or value.get("query") or ""),
+        axes=dict(value.get("axes") or value.get("parts") or {}),
+        priority=int(value.get("priority") or 0),
+        original_position=int(value.get("original_position", position)),
+        source_refs=[dict(ref) for ref in source_refs if isinstance(ref, dict)],
+        knowledge_need_keys=[str(key) for key in value.get("knowledge_need_keys") or []],
+        metadata=dict(value.get("metadata") or {}),
+        filter_reason=filter_reason,
+    )
+
+
+def _need(value: Any) -> KnowledgeNeedDraft | None:
+    if isinstance(value, KnowledgeNeedDraft):
+        return value
+    if not isinstance(value, dict):
+        return None
+    need_key = str(value.get("need_key") or "").strip()
+    if not need_key:
+        return None
+    return KnowledgeNeedDraft(
+        need_key=need_key,
+        decision_context=str(value.get("decision_context") or ""),
+        unknown_information=str(value.get("unknown_information") or ""),
+        source_ref=dict(value.get("source_ref") or {}),
+        priority=int(value.get("priority") or 0),
+        metadata=dict(value.get("metadata") or {}),
+    )
+
+
+class _PayloadQueryGenerator:
+    kind: GeneratorKind
+
+    def generate(self, request: GenerationRequest) -> GeneratorOutput:
+        raw_candidates: Iterable[Any] = request.payload.get("candidates") or ()
+        candidates = tuple(_candidate(value, index, self.kind) for index, value in enumerate(raw_candidates))
+        needs = tuple(
+            need
+            for raw in request.payload.get("knowledge_needs") or ()
+            if (need := _need(raw)) is not None
+        )
+        return GeneratorOutput(
+            candidates=candidates,
+            knowledge_needs=needs,
+            generator_config=dict(request.payload.get("generator_config") or {}),
+        )
+
+
+class CartesianQueryGenerator(_PayloadQueryGenerator):
+    kind = GeneratorKind.CARTESIAN
+
+
+class ManualQueryGenerator(_PayloadQueryGenerator):
+    kind = GeneratorKind.MANUAL
+
+
+class QueryGeneratorRegistry:
+    def __init__(self, generators: Iterable[QueryGenerator] = ()) -> None:
+        self._generators: dict[GeneratorKind, QueryGenerator] = {}
+        for generator in generators:
+            self.register(generator)
+
+    def register(self, generator: QueryGenerator) -> None:
+        self._generators[generator.kind] = generator
+
+    def get(self, kind: GeneratorKind) -> QueryGenerator:
+        try:
+            return self._generators[kind]
+        except KeyError as exc:
+            raise UnsupportedGeneratorError(f"query generator is not registered: {kind.value}") from exc
+
+
+def default_generator_registry() -> QueryGeneratorRegistry:
+    return QueryGeneratorRegistry((CartesianQueryGenerator(), ManualQueryGenerator()))

+ 36 - 0
query_planning/ports.py

@@ -0,0 +1,36 @@
+"""Ports owned by the query planning bounded context."""
+from __future__ import annotations
+
+from typing import Any, Protocol
+from uuid import UUID
+
+from query_planning.domain import GenerationRequest, GeneratorKind, GeneratorOutput, KnowledgeNeedDraft
+
+
+class QueryGenerator(Protocol):
+    kind: GeneratorKind
+
+    def generate(self, request: GenerationRequest) -> GeneratorOutput:
+        ...
+
+
+class LegacyQueryBatchSink(Protocol):
+    def create_query_batch(self, **kwargs: Any) -> Any:
+        ...
+
+    def add_query(self, **kwargs: Any) -> Any:
+        ...
+
+
+class QueryPlanningStore(Protocol):
+    def create_plan(self, **kwargs: Any) -> UUID:
+        ...
+
+    def add_knowledge_need(self, *, plan_id: UUID, need: KnowledgeNeedDraft) -> UUID:
+        ...
+
+    def link_plan_batch(self, *, plan_id: UUID, batch_id: UUID) -> None:
+        ...
+
+    def link_query_need(self, *, query_id: UUID, knowledge_need_id: UUID) -> None:
+        ...

+ 5 - 0
query_planning/repositories/__init__.py

@@ -0,0 +1,5 @@
+"""Query planning persistence adapters."""
+
+from query_planning.repositories.postgres import PostgresQueryPlanningStore
+
+__all__ = ["PostgresQueryPlanningStore"]

+ 97 - 0
query_planning/repositories/postgres.py

@@ -0,0 +1,97 @@
+"""PostgreSQL adapter for query-planning-owned tables only."""
+from __future__ import annotations
+
+from typing import Any
+from uuid import UUID
+
+import psycopg2.extras
+
+from query_planning.domain import KnowledgeNeedDraft
+
+
+Json = psycopg2.extras.Json
+psycopg2.extras.register_uuid()
+
+
+class PostgresQueryPlanningStore:
+    def __init__(self, conn: Any) -> None:
+        self.conn = conn
+
+    def _id(self, sql: str, params: tuple[Any, ...]) -> UUID:
+        with self.conn.cursor() as cur:
+            cur.execute(sql, params)
+            row = cur.fetchone()
+        if row is None:
+            raise RuntimeError("query planning insert returned no id")
+        return row[0]
+
+    def create_plan(self, **kwargs: Any) -> UUID:
+        return self._id(
+            """
+            INSERT INTO creation_knowledge.query_plans(
+                generator_kind, status, input_snapshot, generator_config,
+                max_queries, generated_count, normalized_count, unique_count,
+                selected_count, dropped_count, metadata
+            )
+            VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
+            RETURNING id
+            """,
+            (
+                kwargs["generator_kind"],
+                kwargs["status"],
+                Json(kwargs.get("input_snapshot") or {}),
+                Json(kwargs.get("generator_config") or {}),
+                kwargs.get("max_queries"),
+                kwargs["generated_count"],
+                kwargs["normalized_count"],
+                kwargs["unique_count"],
+                kwargs["selected_count"],
+                kwargs["dropped_count"],
+                Json(kwargs.get("metadata") or {}),
+            ),
+        )
+
+    def add_knowledge_need(self, *, plan_id: UUID, need: KnowledgeNeedDraft) -> UUID:
+        return self._id(
+            """
+            INSERT INTO creation_knowledge.knowledge_needs(
+                plan_id, need_key, decision_context, unknown_information,
+                source_ref, priority, metadata
+            )
+            VALUES (%s, %s, %s, %s, %s, %s, %s)
+            RETURNING id
+            """,
+            (
+                plan_id,
+                need.need_key,
+                need.decision_context,
+                need.unknown_information,
+                Json(need.source_ref),
+                need.priority,
+                Json(need.metadata),
+            ),
+        )
+
+    def link_plan_batch(self, *, plan_id: UUID, batch_id: UUID) -> None:
+        with self.conn.cursor() as cur:
+            cur.execute(
+                """
+                INSERT INTO creation_knowledge.query_plan_batches(plan_id, batch_id)
+                VALUES (%s, %s)
+                ON CONFLICT DO NOTHING
+                """,
+                (plan_id, batch_id),
+            )
+
+    def link_query_need(self, *, query_id: UUID, knowledge_need_id: UUID) -> None:
+        with self.conn.cursor() as cur:
+            cur.execute(
+                """
+                INSERT INTO creation_knowledge.query_knowledge_need_links(
+                    query_id, knowledge_need_id
+                )
+                VALUES (%s, %s)
+                ON CONFLICT DO NOTHING
+                """,
+                (query_id, knowledge_need_id),
+            )

+ 97 - 0
query_planning/service.py

@@ -0,0 +1,97 @@
+"""Unified normalization, exact dedupe, priority, and budget service."""
+from __future__ import annotations
+
+import json
+import unicodedata
+from copy import deepcopy
+from typing import Any
+
+from query_planning.domain import (
+    GenerationPlanDraft,
+    GenerationRequest,
+    GenerationResult,
+    GenerationStats,
+    QueryCandidate,
+)
+from query_planning.errors import NoQueryCandidatesError
+from query_planning.generators import QueryGeneratorRegistry, default_generator_registry
+
+
+def normalize_query_text(value: str) -> tuple[str, str]:
+    display_text = " ".join(unicodedata.normalize("NFKC", value or "").strip().split())
+    return display_text, display_text.casefold()
+
+
+def _unique_dicts(values: list[dict[str, Any]]) -> list[dict[str, Any]]:
+    out: list[dict[str, Any]] = []
+    seen: set[str] = set()
+    for value in values:
+        key = json.dumps(value, ensure_ascii=False, sort_keys=True, default=str)
+        if key in seen:
+            continue
+        seen.add(key)
+        out.append(deepcopy(value))
+    return out
+
+
+class UnifiedQueryGenerationService:
+    def __init__(self, registry: QueryGeneratorRegistry | None = None) -> None:
+        self.registry = registry or default_generator_registry()
+
+    def generate(self, request: GenerationRequest) -> GenerationResult:
+        output = self.registry.get(request.generator_kind).generate(request)
+        generated_count = len(output.candidates)
+        normalized: list[QueryCandidate] = []
+        for index, raw in enumerate(output.candidates):
+            text, key = normalize_query_text(raw.query_text)
+            if not text:
+                continue
+            candidate = deepcopy(raw)
+            candidate.query_text = text
+            candidate.dedupe_key = key
+            candidate.original_position = raw.original_position if raw.original_position >= 0 else index
+            candidate.source_refs = _unique_dicts(candidate.source_refs)
+            normalized.append(candidate)
+
+        by_key: dict[str, QueryCandidate] = {}
+        for candidate in normalized:
+            existing = by_key.get(candidate.dedupe_key)
+            if existing is None:
+                by_key[candidate.dedupe_key] = candidate
+                continue
+            existing.priority = max(existing.priority, candidate.priority)
+            existing.source_refs = _unique_dicts(existing.source_refs + candidate.source_refs)
+            existing.knowledge_need_keys = list(
+                dict.fromkeys(existing.knowledge_need_keys + candidate.knowledge_need_keys)
+            )
+
+        unique = list(by_key.values())
+        for candidate in unique:
+            candidate.metadata = deepcopy(candidate.metadata)
+            candidate.metadata["origins"] = deepcopy(candidate.source_refs)
+            candidate.metadata["priority"] = candidate.priority
+        unique.sort(key=lambda candidate: (-candidate.priority, candidate.original_position))
+        selected = unique if request.max_queries is None else unique[: request.max_queries]
+        if not selected:
+            raise NoQueryCandidatesError("query generation produced no selectable query")
+
+        stats = GenerationStats(
+            generated_count=generated_count,
+            normalized_count=len(normalized),
+            unique_count=len(unique),
+            selected_count=len(selected),
+            dropped_count=generated_count - len(selected),
+        )
+        plan = GenerationPlanDraft(
+            generator_kind=request.generator_kind,
+            input_snapshot=deepcopy(request.payload.get("input_snapshot", request.payload)),
+            generator_config=deepcopy(output.generator_config),
+            metadata=deepcopy(request.metadata),
+        )
+        return GenerationResult(
+            request=request,
+            plan=plan,
+            knowledge_needs=output.knowledge_needs,
+            selected_queries=tuple(selected),
+            stats=stats,
+        )

+ 144 - 0
query_planning/writer.py

@@ -0,0 +1,144 @@
+"""Atomic writer coordination for planning rows and the legacy query contract."""
+from __future__ import annotations
+
+from dataclasses import dataclass, field
+from typing import Any
+from uuid import UUID, uuid4
+
+from query_planning.domain import GenerationResult, KnowledgeNeedDraft
+from query_planning.ports import LegacyQueryBatchSink, QueryPlanningStore
+
+
+@dataclass(frozen=True)
+class QueryBatchWriteSpec:
+    name: str
+    source_type: str
+    generation_method: str
+    target_platforms: tuple[str, ...]
+    metadata: dict[str, Any] = field(default_factory=dict)
+
+    def __post_init__(self) -> None:
+        if not self.name.strip():
+            raise ValueError("query batch name must not be empty")
+        if not self.target_platforms:
+            raise ValueError("query batch requires at least one target platform")
+
+
+@dataclass(frozen=True)
+class QueryBatchWriteResult:
+    plan_id: UUID
+    batch: Any
+    queries: tuple[Any, ...]
+
+
+class NullQueryPlanningStore:
+    """Compatibility adapter for legacy fakes that expose no DB connection."""
+
+    def create_plan(self, **kwargs: Any) -> UUID:
+        return uuid4()
+
+    def add_knowledge_need(self, *, plan_id: UUID, need: KnowledgeNeedDraft) -> UUID:
+        return uuid4()
+
+    def link_plan_batch(self, *, plan_id: UUID, batch_id: UUID) -> None:
+        return None
+
+    def link_query_need(self, *, query_id: UUID, knowledge_need_id: UUID) -> None:
+        return None
+
+
+class QueryBatchWriter:
+    def __init__(self, *, legacy_sink: LegacyQueryBatchSink, planning_store: QueryPlanningStore) -> None:
+        self.legacy_sink = legacy_sink
+        self.planning_store = planning_store
+
+    def write(self, result: GenerationResult, spec: QueryBatchWriteSpec) -> QueryBatchWriteResult:
+        if not result.selected_queries:
+            raise ValueError("cannot persist an empty query selection")
+        declared_need_keys = [need.need_key for need in result.knowledge_needs]
+        if len(set(declared_need_keys)) != len(declared_need_keys):
+            raise ValueError("knowledge need keys must be unique inside one plan")
+        referenced_need_keys = {
+            need_key
+            for candidate in result.selected_queries
+            for need_key in candidate.knowledge_need_keys
+        }
+        unknown_need_keys = referenced_need_keys - set(declared_need_keys)
+        if unknown_need_keys:
+            raise ValueError(
+                "query references unknown knowledge need key(s): "
+                + ", ".join(sorted(unknown_need_keys))
+            )
+        stats = result.stats
+        plan_id = self.planning_store.create_plan(
+            generator_kind=result.plan.generator_kind.value,
+            status="ready",
+            input_snapshot=result.plan.input_snapshot,
+            generator_config=result.plan.generator_config,
+            max_queries=result.request.max_queries,
+            generated_count=stats.generated_count,
+            normalized_count=stats.normalized_count,
+            unique_count=stats.unique_count,
+            selected_count=stats.selected_count,
+            dropped_count=stats.dropped_count,
+            metadata=result.plan.metadata,
+        )
+        need_ids = {
+            need.need_key: self.planning_store.add_knowledge_need(plan_id=plan_id, need=need)
+            for need in result.knowledge_needs
+        }
+        batch_metadata = dict(spec.metadata)
+        batch_metadata["query_planning"] = {
+            "plan_id": str(plan_id),
+            "generator_kind": result.plan.generator_kind.value,
+            "generated_count": stats.generated_count,
+            "unique_count": stats.unique_count,
+            "selected_count": stats.selected_count,
+            "dropped_count": stats.dropped_count,
+        }
+        batch = self.legacy_sink.create_query_batch(
+            name=spec.name,
+            source_type=spec.source_type,
+            generation_method=spec.generation_method,
+            target_platforms=list(spec.target_platforms),
+            status="ready",
+            metadata=batch_metadata,
+        )
+        if getattr(batch, "id", None) is None:
+            raise RuntimeError("legacy query batch sink returned a batch without id")
+
+        written: list[Any] = []
+        for sort_order, candidate in enumerate(result.selected_queries):
+            query = self.legacy_sink.add_query(
+                batch_id=batch.id,
+                query_text=candidate.query_text,
+                axes=candidate.axes,
+                keep=True,
+                filter_reason=candidate.filter_reason,
+                status="ready",
+                sort_order=sort_order,
+                metadata=candidate.metadata,
+            )
+            if getattr(query, "id", None) is None:
+                raise RuntimeError("legacy query batch sink returned a query without id")
+            written.append(query)
+
+        self.planning_store.link_plan_batch(plan_id=plan_id, batch_id=batch.id)
+        for query, candidate in zip(written, result.selected_queries, strict=True):
+            for need_key in candidate.knowledge_need_keys:
+                need_id = need_ids.get(need_key)
+                if need_id is not None:
+                    self.planning_store.link_query_need(
+                        query_id=query.id,
+                        knowledge_need_id=need_id,
+                    )
+        return QueryBatchWriteResult(plan_id=plan_id, batch=batch, queries=tuple(written))
+
+
+def planning_store_for_repository(repo: Any) -> QueryPlanningStore:
+    conn = getattr(repo, "conn", None)
+    if conn is None:
+        return NullQueryPlanningStore()
+    from query_planning.repositories.postgres import PostgresQueryPlanningStore
+
+    return PostgresQueryPlanningStore(conn)

+ 376 - 0
tests/test_query_planning_core.py

@@ -0,0 +1,376 @@
+from __future__ import annotations
+
+from pathlib import Path
+from uuid import uuid4
+
+import pytest
+
+from acquisition.domain import Query, QueryBatch
+from acquisition.queries.builder import persist_query_batch
+from core import db_session
+from query_planning import (
+    GenerationRequest,
+    GeneratorKind,
+    NoQueryCandidatesError,
+    QueryBatchWriteSpec,
+    QueryBatchWriter,
+    UnifiedQueryGenerationService,
+    UnsupportedGeneratorError,
+)
+
+
+def _request(kind=GeneratorKind.MANUAL, *, candidates, max_queries=None):
+    return GenerationRequest(
+        generator_kind=kind,
+        name="test-plan",
+        payload={"candidates": candidates, "input_snapshot": {"case": "test"}},
+        max_queries=max_queries,
+    )
+
+
+def test_normalizes_and_exactly_dedupes_while_preserving_near_queries():
+    result = UnifiedQueryGenerationService().generate(
+        _request(
+            candidates=[
+                {
+                    "query_text": "  A   B  ",
+                    "axes": {"first": True},
+                    "metadata": {"family_key": "first"},
+                    "source_refs": [{"family_key": "first"}],
+                },
+                {
+                    "query_text": "a b",
+                    "priority": 8,
+                    "axes": {"second": True},
+                    "metadata": {"family_key": "second"},
+                    "source_refs": [{"family_key": "second"}],
+                },
+                {
+                    "query_text": "a b 方法",
+                    "source_refs": [{"family_key": "near"}],
+                },
+            ]
+        )
+    )
+
+    assert [query.query_text for query in result.selected_queries] == ["A B", "a b 方法"]
+    merged = result.selected_queries[0]
+    assert merged.axes == {"first": True}
+    assert merged.metadata["family_key"] == "first"
+    assert merged.priority == 8
+    assert merged.metadata["origins"] == [
+        {"family_key": "first"},
+        {"family_key": "second"},
+    ]
+    assert result.stats.generated_count == 3
+    assert result.stats.unique_count == 2
+    assert result.stats.dropped_count == 1
+
+
+def test_priority_budget_is_stable_and_default_is_unlimited():
+    candidates = [
+        {"query_text": "q0", "priority": 0},
+        {"query_text": "q1", "priority": 5},
+        {"query_text": "q2", "priority": 5},
+        {"query_text": "q3", "priority": 1},
+    ]
+    service = UnifiedQueryGenerationService()
+
+    unlimited = service.generate(_request(candidates=candidates))
+    budgeted = service.generate(_request(candidates=candidates, max_queries=2))
+
+    assert [query.query_text for query in unlimited.selected_queries] == ["q1", "q2", "q3", "q0"]
+    assert [query.query_text for query in budgeted.selected_queries] == ["q1", "q2"]
+    assert budgeted.stats.selected_count == 2
+    assert budgeted.stats.dropped_count == 2
+
+
+def test_empty_selection_and_unregistered_reserved_generator_fail_explicitly():
+    service = UnifiedQueryGenerationService()
+    with pytest.raises(NoQueryCandidatesError):
+        service.generate(_request(candidates=["  "]))
+    with pytest.raises(NoQueryCandidatesError):
+        service.generate(_request(candidates=["valid"], max_queries=0))
+    with pytest.raises(UnsupportedGeneratorError, match="topic_table"):
+        service.generate(
+            _request(kind=GeneratorKind.TOPIC_TABLE, candidates=["not allowed yet"])
+        )
+
+
+class FakeLegacySink:
+    def __init__(self):
+        self.batch_kwargs = None
+        self.query_kwargs = []
+
+    def create_query_batch(self, **kwargs):
+        self.batch_kwargs = kwargs
+        return QueryBatch(id=uuid4(), **kwargs)
+
+    def add_query(self, **kwargs):
+        self.query_kwargs.append(kwargs)
+        return Query(id=uuid4(), **kwargs)
+
+
+class FakePlanningStore:
+    def __init__(self):
+        self.plan_id = uuid4()
+        self.plans = []
+        self.needs = []
+        self.plan_batch_links = []
+        self.query_need_links = []
+
+    def create_plan(self, **kwargs):
+        self.plans.append(kwargs)
+        return self.plan_id
+
+    def add_knowledge_need(self, *, plan_id, need):
+        need_id = uuid4()
+        self.needs.append((plan_id, need, need_id))
+        return need_id
+
+    def link_plan_batch(self, **kwargs):
+        self.plan_batch_links.append(kwargs)
+
+    def link_query_need(self, **kwargs):
+        self.query_need_links.append(kwargs)
+
+
+def test_writer_keeps_legacy_search_contract_and_links_plan_need_rows():
+    service = UnifiedQueryGenerationService()
+    result = service.generate(
+        GenerationRequest(
+            generator_kind=GeneratorKind.MANUAL,
+            name="manual",
+            payload={
+                "knowledge_needs": [
+                    {
+                        "need_key": "need-1",
+                        "decision_context": "选择画面主线",
+                        "unknown_information": "透明伞怎样形成柔光",
+                        "source_ref": {"topic_id": 428},
+                    }
+                ],
+                "candidates": [
+                    {
+                        "query_text": "透明伞 人像 柔光",
+                        "axes": {"道具": "透明伞"},
+                        "metadata": {"family_key": "manual"},
+                        "source_refs": [{"generator_kind": "manual"}],
+                        "knowledge_need_keys": ["need-1"],
+                    }
+                ],
+            },
+        )
+    )
+    sink = FakeLegacySink()
+    planning = FakePlanningStore()
+
+    written = QueryBatchWriter(legacy_sink=sink, planning_store=planning).write(
+        result,
+        QueryBatchWriteSpec(
+            name="manual",
+            source_type="manual",
+            generation_method="manual_query_api_v1",
+            target_platforms=("xiaohongshu",),
+            metadata={"source": "test"},
+        ),
+    )
+
+    assert written.plan_id == planning.plan_id
+    assert sink.batch_kwargs["status"] == "ready"
+    assert sink.batch_kwargs["target_platforms"] == ["xiaohongshu"]
+    assert sink.batch_kwargs["metadata"]["query_planning"]["plan_id"] == str(planning.plan_id)
+    query = sink.query_kwargs[0]
+    assert query["query_text"] == "透明伞 人像 柔光"
+    assert query["keep"] is True
+    assert query["status"] == "ready"
+    assert query["sort_order"] == 0
+    assert query["filter_reason"] is None
+    assert query["metadata"]["family_key"] == "manual"
+    assert planning.plan_batch_links == [
+        {"plan_id": planning.plan_id, "batch_id": written.batch.id}
+    ]
+    assert len(planning.query_need_links) == 1
+
+
+def test_writer_rejects_unknown_need_before_any_write():
+    result = UnifiedQueryGenerationService().generate(
+        _request(
+            candidates=[
+                {
+                    "query_text": "query",
+                    "knowledge_need_keys": ["missing"],
+                }
+            ]
+        )
+    )
+    sink = FakeLegacySink()
+    planning = FakePlanningStore()
+
+    with pytest.raises(ValueError, match="missing"):
+        QueryBatchWriter(legacy_sink=sink, planning_store=planning).write(
+            result,
+            QueryBatchWriteSpec(
+                name="manual",
+                source_type="manual",
+                generation_method="manual_query_api_v1",
+                target_platforms=("xiaohongshu",),
+            ),
+        )
+
+    assert planning.plans == []
+    assert sink.batch_kwargs is None
+
+
+def test_outer_transaction_rolls_back_partial_writer_failure(monkeypatch):
+    class FakeConnection:
+        def __init__(self):
+            self.pending = []
+            self.committed = []
+            self.rollback_called = False
+            self.closed = False
+
+        def commit(self):
+            self.committed.extend(self.pending)
+            self.pending.clear()
+
+        def rollback(self):
+            self.rollback_called = True
+            self.pending.clear()
+
+        def close(self):
+            self.closed = True
+
+    class TransactionalPlanningStore(FakePlanningStore):
+        def __init__(self, conn):
+            super().__init__()
+            self.conn = conn
+
+        def create_plan(self, **kwargs):
+            self.conn.pending.append(("plan", kwargs))
+            return self.plan_id
+
+        def link_plan_batch(self, **kwargs):
+            self.conn.pending.append(("plan_batch", kwargs))
+
+    class FailingLegacySink(FakeLegacySink):
+        def __init__(self, conn):
+            super().__init__()
+            self.conn = conn
+
+        def create_query_batch(self, **kwargs):
+            batch = super().create_query_batch(**kwargs)
+            self.conn.pending.append(("batch", batch.id))
+            return batch
+
+        def add_query(self, **kwargs):
+            if len(self.query_kwargs) == 1:
+                raise RuntimeError("second query failed")
+            query = super().add_query(**kwargs)
+            self.conn.pending.append(("query", query.id))
+            return query
+
+    result = UnifiedQueryGenerationService().generate(
+        _request(candidates=["first query", "second query"])
+    )
+    conn = FakeConnection()
+    monkeypatch.setattr(db_session, "connect", lambda _config: conn)
+
+    with pytest.raises(RuntimeError, match="second query failed"):
+        with db_session.transaction(object()) as transaction_conn:
+            QueryBatchWriter(
+                legacy_sink=FailingLegacySink(transaction_conn),
+                planning_store=TransactionalPlanningStore(transaction_conn),
+            ).write(
+                result,
+                QueryBatchWriteSpec(
+                    name="rollback",
+                    source_type="manual",
+                    generation_method="manual_query_api_v1",
+                    target_platforms=("xiaohongshu",),
+                ),
+            )
+
+    assert conn.rollback_called is True
+    assert conn.pending == []
+    assert conn.committed == []
+    assert conn.closed is True
+
+
+def test_cartesian_legacy_facade_exactly_dedupes_across_families():
+    sink = FakeLegacySink()
+    generated = {
+        "metadata": {"active_family_keys": ["f1", "f2"]},
+        "families": [
+            {
+                "key": "f1",
+                "name": "实质 × 模态",
+                "axes": ["实质", "模态"],
+                "items": [
+                    {"query": "A  B", "parts": {"实质": "A"}, "keep": True}
+                ],
+            },
+            {
+                "key": "f2",
+                "name": "形式 × 模态",
+                "axes": ["形式", "模态"],
+                "items": [
+                    {"query": "a b", "parts": {"形式": "A"}, "keep": True}
+                ],
+            },
+        ],
+    }
+
+    batch, count = persist_query_batch(sink, generated, name="cartesian")
+
+    assert batch.id is not None
+    assert count == 1
+    assert sink.query_kwargs[0]["query_text"] == "A B"
+    assert sink.query_kwargs[0]["sort_order"] == 0
+    assert sink.query_kwargs[0]["metadata"]["family_key"] == "f1"
+    assert [origin["family_key"] for origin in sink.query_kwargs[0]["metadata"]["origins"]] == [
+        "f1",
+        "f2",
+    ]
+
+
+def test_migration_is_additive_replayable_and_does_not_force_what_how_why():
+    sql = Path("db/migrations/005_query_planning_schema.sql").read_text(encoding="utf-8")
+    for table in (
+        "query_plans",
+        "knowledge_needs",
+        "query_plan_batches",
+        "query_knowledge_need_links",
+    ):
+        assert f"CREATE TABLE IF NOT EXISTS creation_knowledge.{table}" in sql
+    assert "ALTER TABLE creation_knowledge.query_batches" not in sql
+    assert "ALTER TABLE creation_knowledge.queries" not in sql
+    assert "005_query_planning_schema" in sql
+    assert "trg_query_plans_touch_updated_at" in sql
+    assert "TO ck_app" in sql
+    assert "particle_type" not in sql
+
+
+def test_frozen_search_modules_do_not_depend_on_query_planning():
+    frozen_roots = [
+        Path("acquisition/runner.py"),
+        Path("acquisition/repositories"),
+        Path("acquisition/platforms"),
+        Path("pipeline"),
+        Path("decode_content"),
+    ]
+    for root in frozen_roots:
+        files = [root] if root.is_file() else list(root.rglob("*.py"))
+        for path in files:
+            assert "query_planning" not in path.read_text(encoding="utf-8"), path
+
+    forbidden = (
+        "acquisition.runner",
+        "acquisition.platforms",
+        "acquisition.search",
+        "pipeline",
+        "decode_content",
+    )
+    for path in Path("query_planning").rglob("*.py"):
+        source = path.read_text(encoding="utf-8")
+        assert not any(f"from {module}" in source or f"import {module}" in source for module in forbidden), path