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

feat: add manual query batch entry

SamLee 2 недель назад
Родитель
Сommit
d496091b07

+ 2 - 1
app/api.py

@@ -9,7 +9,7 @@ from fastapi.middleware.cors import CORSMiddleware
 from fastapi.middleware.gzip import GZipMiddleware
 from fastapi.middleware.gzip import GZipMiddleware
 from fastapi.staticfiles import StaticFiles
 from fastapi.staticfiles import StaticFiles
 
 
-from app.routes import acquisition, decode, payloads, query_generation, runs
+from app.routes import acquisition, decode, manual_queries, payloads, query_generation, runs
 
 
 ROOT = Path(__file__).resolve().parent.parent
 ROOT = Path(__file__).resolve().parent.parent
 FRONTEND_DIST = ROOT / "app" / "frontend" / "dist"
 FRONTEND_DIST = ROOT / "app" / "frontend" / "dist"
@@ -39,6 +39,7 @@ def create_app() -> FastAPI:
     app.add_middleware(GZipMiddleware, minimum_size=1024)
     app.add_middleware(GZipMiddleware, minimum_size=1024)
     app.include_router(acquisition.router)
     app.include_router(acquisition.router)
     app.include_router(decode.router)
     app.include_router(decode.router)
+    app.include_router(manual_queries.router)
     app.include_router(payloads.router)
     app.include_router(payloads.router)
     app.include_router(query_generation.router)
     app.include_router(query_generation.router)
     app.include_router(runs.router)
     app.include_router(runs.router)

Разница между файлами не показана из-за своего большого размера
+ 0 - 8
app/frontend/dist/assets/index-D2o7sqDa.js


Разница между файлами не показана из-за своего большого размера
+ 8 - 0
app/frontend/dist/assets/index-iEdd66_J.js


Разница между файлами не показана из-за своего большого размера
+ 0 - 0
app/frontend/dist/assets/index-zUYQ6IWb.css


+ 2 - 2
app/frontend/dist/index.html

@@ -4,8 +4,8 @@
     <meta charset="UTF-8" />
     <meta charset="UTF-8" />
     <meta name="viewport" content="width=device-width, initial-scale=1.0" />
     <meta name="viewport" content="width=device-width, initial-scale=1.0" />
     <title>创作知识</title>
     <title>创作知识</title>
-    <script type="module" crossorigin src="/app/assets/index-D2o7sqDa.js"></script>
-    <link rel="stylesheet" crossorigin href="/app/assets/index-qEoc1W1z.css">
+    <script type="module" crossorigin src="/app/assets/index-iEdd66_J.js"></script>
+    <link rel="stylesheet" crossorigin href="/app/assets/index-zUYQ6IWb.css">
   </head>
   </head>
   <body>
   <body>
     <div id="root"></div>
     <div id="root"></div>

+ 19 - 2
app/frontend/src/api/client.js

@@ -7,8 +7,25 @@ export class ApiError extends Error {
   }
   }
 }
 }
 
 
-export async function request(path, { signal } = {}) {
-  const response = await fetch(path, { signal })
+export async function request(path, {
+  signal,
+  method = 'GET',
+  body,
+  json,
+  headers = {},
+} = {}) {
+  const finalHeaders = { ...headers }
+  let finalBody = body
+  if (json !== undefined) {
+    finalHeaders['Content-Type'] = finalHeaders['Content-Type'] || 'application/json'
+    finalBody = JSON.stringify(json)
+  }
+  const response = await fetch(path, {
+    signal,
+    method,
+    headers: finalHeaders,
+    body: finalBody,
+  })
   const contentType = response.headers.get('content-type') || ''
   const contentType = response.headers.get('content-type') || ''
   const payload = contentType.includes('application/json')
   const payload = contentType.includes('application/json')
     ? await response.json().catch(() => null)
     ? await response.json().catch(() => null)

+ 8 - 0
app/frontend/src/api/workbench.js

@@ -11,3 +11,11 @@ export function getLatestBoard({ signal } = {}) {
 export function getLatestQueryDetail(queryId, { signal } = {}) {
 export function getLatestQueryDetail(queryId, { signal } = {}) {
   return request(`/api/query-generation/latest/queries/${encodeURIComponent(queryId)}`, { signal })
   return request(`/api/query-generation/latest/queries/${encodeURIComponent(queryId)}`, { signal })
 }
 }
+
+export function submitManualQueryBatch(payload, { signal } = {}) {
+  return request('/api/query-batches/manual', {
+    method: 'POST',
+    json: payload,
+    signal,
+  })
+}

+ 127 - 0
app/frontend/src/features/query-board/ManualQueryModal.jsx

@@ -0,0 +1,127 @@
+import { useMemo, useState } from 'react'
+import { submitManualQueryBatch } from '../../api/workbench.js'
+
+function normalizeJsonQueries(value) {
+  if (Array.isArray(value)) return value
+  if (!value || typeof value !== 'object') return []
+  if (Array.isArray(value.queries)) return value.queries
+  if (Array.isArray(value.families)) {
+    return value.families.flatMap((family) => (
+      Array.isArray(family?.items)
+        ? family.items.map((item) => ({
+          ...item,
+          metadata: {
+            ...(item.metadata || {}),
+            source_family_key: family.key || family.family_key,
+            source_family_name: family.name || family.title,
+          },
+        }))
+        : []
+    ))
+  }
+  if (value.query || value.query_text) return [value]
+  return []
+}
+
+function queryLabel(value) {
+  if (typeof value === 'string') return value.trim()
+  if (value && typeof value === 'object') return String(value.query_text || value.query || '').trim()
+  return ''
+}
+
+export default function ManualQueryModal({ mode, onClose, onSubmitted }) {
+  const [text, setText] = useState('')
+  const [jsonQueries, setJsonQueries] = useState([])
+  const [fileName, setFileName] = useState('')
+  const [error, setError] = useState('')
+  const [submitting, setSubmitting] = useState(false)
+
+  const queries = useMemo(() => {
+    if (mode === 'json') return jsonQueries
+    return text.split('\n').map((line) => line.trim()).filter(Boolean)
+  }, [jsonQueries, mode, text])
+  const previews = queries.map(queryLabel).filter(Boolean).slice(0, 5)
+
+  const onFileChange = async (event) => {
+    const file = event.target.files?.[0]
+    if (!file) return
+    setError('')
+    setFileName(file.name)
+    try {
+      const raw = await file.text()
+      const parsed = JSON.parse(raw)
+      const nextQueries = normalizeJsonQueries(parsed).filter((item) => queryLabel(item))
+      setJsonQueries(nextQueries)
+      if (!nextQueries.length) setError('这个 JSON 里没有识别到 query')
+    } catch {
+      setJsonQueries([])
+      setError('JSON 解析失败')
+    }
+  }
+
+  const onSubmit = async () => {
+    const clean = queries.filter((item) => queryLabel(item))
+    if (!clean.length) {
+      setError('请先输入或上传 query')
+      return
+    }
+    setSubmitting(true)
+    setError('')
+    try {
+      const result = await submitManualQueryBatch({ queries: clean })
+      await onSubmitted?.(result)
+    } catch (e) {
+      setError(e.message || '提交失败')
+    } finally {
+      setSubmitting(false)
+    }
+  }
+
+  return (
+    <div className="manual-modal-mask" onClick={onClose}>
+      <div className="manual-modal-panel" onClick={(event) => event.stopPropagation()}>
+        <div className="manual-modal-head">
+          <div>
+            <h2>{mode === 'json' ? '上传 JSON' : '手动输入'}</h2>
+            <p>提交后会创建手动 query 批次,并在后台跑搜索、初筛、解构和 payload dry-run。</p>
+          </div>
+          <button className="manual-close" type="button" onClick={onClose}>×</button>
+        </div>
+
+        {mode === 'json' ? (
+          <label className="manual-upload">
+            <input accept=".json,application/json" type="file" onChange={onFileChange} />
+            <span>{fileName || '选择 JSON 文件'}</span>
+          </label>
+        ) : (
+          <textarea
+            className="manual-textarea"
+            placeholder="一行一个 query,例如:&#10;符号 视频 灵感 怎么做&#10;叙事体裁 视频 灵感 有哪些"
+            value={text}
+            onChange={(event) => setText(event.target.value)}
+          />
+        )}
+
+        <div className="manual-preview">
+          <strong>识别到 {queries.length} 条</strong>
+          {previews.length ? (
+            <ul>
+              {previews.map((item, index) => <li key={`${item}-${index}`}>{item}</li>)}
+            </ul>
+          ) : (
+            <span>暂无预览</span>
+          )}
+        </div>
+
+        {error && <div className="manual-error">{error}</div>}
+
+        <div className="manual-actions">
+          <button className="manual-secondary" type="button" onClick={onClose}>取消</button>
+          <button className="manual-primary" disabled={submitting} type="button" onClick={onSubmit}>
+            {submitting ? '提交中...' : '提交并后台运行'}
+          </button>
+        </div>
+      </div>
+    </div>
+  )
+}

+ 40 - 2
app/frontend/src/features/query-board/QueryBoardPage.jsx

@@ -3,6 +3,7 @@ import ErrorState from '../../components/ErrorState.jsx'
 import LoadingState from '../../components/LoadingState.jsx'
 import LoadingState from '../../components/LoadingState.jsx'
 import DecodeKnowledgeModal from '../decode/DecodeKnowledgeModal.jsx'
 import DecodeKnowledgeModal from '../decode/DecodeKnowledgeModal.jsx'
 import AxisColumn from './AxisColumn.jsx'
 import AxisColumn from './AxisColumn.jsx'
+import ManualQueryModal from './ManualQueryModal.jsx'
 import QueryColumn from './QueryColumn.jsx'
 import QueryColumn from './QueryColumn.jsx'
 import SearchResultColumn from './SearchResultColumn.jsx'
 import SearchResultColumn from './SearchResultColumn.jsx'
 import { useLatestQueryDetail, useQueryBoard } from './useQueryBoard.js'
 import { useLatestQueryDetail, useQueryBoard } from './useQueryBoard.js'
@@ -22,6 +23,22 @@ function selectedQueryPayload(queryText, latest) {
   }
   }
 }
 }
 
 
+function manualFamilyFromBoard(board) {
+  const rows = (board?.queries || []).filter((row) => row.family_key === 'manual')
+  if (!rows.length) return null
+  return {
+    key: 'manual',
+    axes: ['手动 Query'],
+    items: rows.map((row, index) => ({
+      query: row.query_text,
+      parts: { '手动 Query': '手动 Query' },
+      keep: true,
+      reason: row.filter_reason || '',
+      sort_order: row.sort_order ?? index,
+    })),
+  }
+}
+
 export default function QueryBoardPage() {
 export default function QueryBoardPage() {
   const {
   const {
     preview,
     preview,
@@ -30,12 +47,18 @@ export default function QueryBoardPage() {
     loading,
     loading,
     error,
     error,
     lastUpdatedAt,
     lastUpdatedAt,
+    refreshBoard,
   } = useQueryBoard()
   } = useQueryBoard()
   const [activeKey, setActiveKey] = useState('f1')
   const [activeKey, setActiveKey] = useState('f1')
   const [selectedQuery, setSelectedQuery] = useState(null)
   const [selectedQuery, setSelectedQuery] = useState(null)
   const [decodeItemId, setDecodeItemId] = useState('')
   const [decodeItemId, setDecodeItemId] = useState('')
+  const [manualMode, setManualMode] = useState('')
 
 
-  const families = preview?.families || []
+  const families = useMemo(() => {
+    const base = preview?.families || []
+    const manual = manualFamilyFromBoard(board)
+    return manual ? [...base, manual] : base
+  }, [board, preview])
   const family = families.find((row) => row.key === activeKey) || families[0] || { axes: [], items: [] }
   const family = families.find((row) => row.key === activeKey) || families[0] || { axes: [], items: [] }
   const familyItems = family.items || []
   const familyItems = family.items || []
   const familyRunCount = familyItems.filter((item) => isQuerySearched(latestByQuery.get(item.query))).length
   const familyRunCount = familyItems.filter((item) => isQuerySearched(latestByQuery.get(item.query))).length
@@ -96,9 +119,13 @@ export default function QueryBoardPage() {
         <span className="refresh-state">
         <span className="refresh-state">
           {lastUpdatedText ? `上次刷新 ${lastUpdatedText}` : '等待后台写入数据'}
           {lastUpdatedText ? `上次刷新 ${lastUpdatedText}` : '等待后台写入数据'}
         </span>
         </span>
+        <div className="manual-entry-actions">
+          <button type="button" onClick={() => setManualMode('text')}>手动输入</button>
+          <button type="button" onClick={() => setManualMode('json')}>上传 JSON</button>
+        </div>
       </div>
       </div>
 
 
-      <div className="cdemo">
+      <div className={`cdemo ${family.key === 'manual' ? 'manual-cdemo' : ''}`}>
         {(family.axes || []).map((axis) => (
         {(family.axes || []).map((axis) => (
           <AxisColumn
           <AxisColumn
             key={axis}
             key={axis}
@@ -133,6 +160,17 @@ export default function QueryBoardPage() {
           </div>
           </div>
         </div>
         </div>
       )}
       )}
+      {manualMode && (
+        <ManualQueryModal
+          mode={manualMode}
+          onClose={() => setManualMode('')}
+          onSubmitted={async () => {
+            setManualMode('')
+            await refreshBoard()
+            setActiveKey('manual')
+          }}
+        />
+      )}
     </div>
     </div>
   )
   )
 }
 }

+ 177 - 0
app/frontend/src/styles/query-board.css

@@ -5,6 +5,183 @@
   white-space: nowrap;
   white-space: nowrap;
 }
 }
 
 
+.manual-entry-actions {
+  display: flex;
+  align-items: center;
+  gap: 8px;
+}
+
+.manual-entry-actions button {
+  height: 30px;
+  padding: 0 10px;
+  border: 1px solid #d8dee9;
+  border-radius: 6px;
+  background: #fff;
+  color: #165dff;
+  font-size: 12px;
+  font-weight: 600;
+  cursor: pointer;
+}
+
+.manual-entry-actions button:hover {
+  border-color: #94bfff;
+  background: #f2f7ff;
+}
+
+.dashboard-page .cdemo.manual-cdemo {
+  grid-template-columns: 180px minmax(420px, 560px) minmax(520px, 1fr);
+}
+
+.manual-modal-mask {
+  position: fixed;
+  inset: 0;
+  z-index: 1200;
+  display: flex;
+  align-items: center;
+  justify-content: center;
+  padding: 24px;
+  background: rgba(29, 33, 41, .42);
+}
+
+.manual-modal-panel {
+  width: min(680px, 100%);
+  max-height: calc(100vh - 72px);
+  display: flex;
+  flex-direction: column;
+  gap: 14px;
+  padding: 18px;
+  overflow: auto;
+  border-radius: 8px;
+  background: #fff;
+  box-shadow: 0 18px 48px rgba(29, 33, 41, .22);
+}
+
+.manual-modal-head {
+  display: flex;
+  align-items: flex-start;
+  justify-content: space-between;
+  gap: 18px;
+}
+
+.manual-modal-head h2 {
+  margin: 0 0 6px;
+  color: #1d2129;
+  font-size: 18px;
+  line-height: 1.25;
+}
+
+.manual-modal-head p {
+  margin: 0;
+  color: #86909c;
+  font-size: 13px;
+  line-height: 1.55;
+}
+
+.manual-close {
+  width: 30px;
+  height: 30px;
+  border: 0;
+  border-radius: 6px;
+  background: #f2f3f5;
+  color: #4e5969;
+  font-size: 20px;
+  line-height: 1;
+  cursor: pointer;
+}
+
+.manual-textarea {
+  min-height: 180px;
+  padding: 12px;
+  resize: vertical;
+  border: 1px solid #d8dee9;
+  border-radius: 6px;
+  color: #1d2129;
+  font: 13px/1.6 ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas, monospace;
+}
+
+.manual-upload {
+  display: flex;
+  align-items: center;
+  justify-content: center;
+  min-height: 96px;
+  border: 1px dashed #94bfff;
+  border-radius: 8px;
+  background: #f7fbff;
+  color: #165dff;
+  font-size: 13px;
+  font-weight: 700;
+  cursor: pointer;
+}
+
+.manual-upload input {
+  display: none;
+}
+
+.manual-preview {
+  padding: 10px 12px;
+  border: 1px solid #e5e6eb;
+  border-radius: 6px;
+  background: #f7f8fa;
+  color: #4e5969;
+  font-size: 12px;
+}
+
+.manual-preview strong {
+  display: block;
+  margin-bottom: 6px;
+  color: #1d2129;
+}
+
+.manual-preview ul {
+  margin: 0;
+  padding-left: 18px;
+}
+
+.manual-preview li {
+  margin: 3px 0;
+}
+
+.manual-error {
+  padding: 8px 10px;
+  border-radius: 6px;
+  background: #fff1f0;
+  color: #f53f3f;
+  font-size: 13px;
+}
+
+.manual-actions {
+  display: flex;
+  justify-content: flex-end;
+  gap: 10px;
+}
+
+.manual-actions button {
+  height: 32px;
+  padding: 0 14px;
+  border-radius: 6px;
+  font-size: 13px;
+  font-weight: 700;
+  cursor: pointer;
+}
+
+.manual-secondary {
+  border: 1px solid #d8dee9;
+  background: #fff;
+  color: #4e5969;
+}
+
+.manual-primary {
+  border: 1px solid #165dff;
+  background: #165dff;
+  color: #fff;
+}
+
+.manual-primary:disabled {
+  border-color: #c9cdd4;
+  background: #c9cdd4;
+  cursor: default;
+}
+
 .virtual-query-list {
 .virtual-query-list {
   padding: 0;
   padding: 0;
 }
 }

+ 277 - 0
app/routes/manual_queries.py

@@ -0,0 +1,277 @@
+"""Manual query batch submission API."""
+from __future__ import annotations
+
+import os
+import subprocess
+import sys
+from datetime import datetime
+from pathlib import Path
+from typing import Any
+
+from fastapi import APIRouter, Body, Depends, HTTPException
+from pydantic import BaseModel, Field
+
+from acquisition.repositories.postgres import PostgresAcquisitionRepository
+from app.dependencies import _env_file, get_creation_db_config
+from app.routes.acquisition import _model_dump
+from core.config import CreationDbConfig
+from core.db_session import transaction
+
+router = APIRouter(prefix="/api/query-batches", tags=["manual-query-batches"])
+
+ROOT = Path(__file__).resolve().parents[2]
+RUNTIME_MANUAL_DIR = ROOT / "runtime" / "manual_query_runs"
+DEFAULT_PLATFORMS = ("xiaohongshu", "weixin", "douyin")
+QUERY_MAX_COUNT = 500
+QUERY_MAX_CHARS = 500
+
+
+class ManualQueryEntry(BaseModel):
+    query_text: str
+    axes: dict[str, Any] = Field(default_factory=dict)
+    metadata: dict[str, Any] = Field(default_factory=dict)
+    filter_reason: str | None = None
+
+
+class ManualQueryBatchRequest(BaseModel):
+    name: str | None = None
+    queries: list[Any] | None = None
+    families: list[dict[str, Any]] | None = None
+    target_platforms: list[str] | None = None
+    metadata: dict[str, Any] = Field(default_factory=dict)
+    search_limit: int = Field(default=10, ge=1, le=50)
+    display_limit: int = Field(default=5, ge=1, le=50)
+    decode_limit: int = Field(default=100, ge=0, le=500)
+
+
+def _as_request(payload: Any) -> ManualQueryBatchRequest:
+    if isinstance(payload, list):
+        payload = {"queries": payload}
+    if not isinstance(payload, dict):
+        raise HTTPException(status_code=422, detail="请求体必须是对象或 query 数组")
+    return ManualQueryBatchRequest.model_validate(payload)
+
+
+def _query_text(value: Any) -> str:
+    if value is None:
+        return ""
+    return str(value).strip()
+
+
+def _entry_from_value(value: Any, *, family: dict[str, Any] | None = None) -> ManualQueryEntry | None:
+    if isinstance(value, str):
+        query_text = _query_text(value)
+        axes: dict[str, Any] = {}
+        metadata: dict[str, Any] = {}
+        filter_reason = None
+    elif isinstance(value, dict):
+        query_text = _query_text(value.get("query_text") or value.get("query"))
+        axes = value.get("axes") or value.get("parts") or {}
+        if not isinstance(axes, dict):
+            axes = {}
+        metadata = value.get("metadata") or {}
+        if not isinstance(metadata, dict):
+            metadata = {}
+        filter_reason = value.get("filter_reason") or value.get("reason")
+    else:
+        return None
+
+    if not query_text:
+        return None
+    if len(query_text) > QUERY_MAX_CHARS:
+        raise HTTPException(status_code=422, detail=f"query 过长,最多 {QUERY_MAX_CHARS} 字符")
+
+    merged_metadata = dict(metadata)
+    merged_metadata["family_key"] = "manual"
+    if family:
+        merged_metadata["source_family_key"] = family.get("key") or family.get("family_key")
+        merged_metadata["source_family_name"] = family.get("name") or family.get("title")
+    return ManualQueryEntry(
+        query_text=query_text,
+        axes=axes,
+        metadata=merged_metadata,
+        filter_reason=_query_text(filter_reason) or None,
+    )
+
+
+def _normalize_entries(request: ManualQueryBatchRequest) -> list[ManualQueryEntry]:
+    entries: list[ManualQueryEntry] = []
+    for value in request.queries or []:
+        entry = _entry_from_value(value)
+        if entry:
+            entries.append(entry)
+
+    for family in request.families or []:
+        if not isinstance(family, dict):
+            continue
+        for value in family.get("items") or []:
+            entry = _entry_from_value(value, family=family)
+            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:
+        raise HTTPException(status_code=422, detail="没有可提交的 query")
+    return deduped
+
+
+def _normalize_platforms(platforms: list[str] | None) -> list[str]:
+    values = platforms or list(DEFAULT_PLATFORMS)
+    if not values:
+        raise HTTPException(status_code=422, detail="至少选择一个平台")
+    seen: set[str] = set()
+    out: list[str] = []
+    for platform in values:
+        key = str(platform).strip()
+        if key not in DEFAULT_PLATFORMS:
+            raise HTTPException(status_code=422, detail=f"不支持的平台:{key}")
+        if key not in seen:
+            seen.add(key)
+            out.append(key)
+    return out
+
+
+def _pipeline_command(
+    *,
+    batch_id: str,
+    run_key: str,
+    platforms: list[str],
+    search_limit: int,
+    display_limit: int,
+    decode_limit: int,
+) -> list[str]:
+    cmd = [
+        sys.executable,
+        str(ROOT / "scripts" / "run_creation_pipeline.py"),
+        "--batch-id",
+        batch_id,
+        "--search-limit",
+        str(search_limit),
+        "--display-limit",
+        str(display_limit),
+        "--decode-limit",
+        str(decode_limit),
+        "--run-key",
+        run_key,
+        "--env-file",
+        _env_file(),
+    ]
+    for platform in platforms:
+        cmd.extend(["--platform", platform])
+    return cmd
+
+
+@router.post("/manual")
+def create_manual_query_batch(
+    payload: Any = Body(...),
+    db_config: CreationDbConfig = Depends(get_creation_db_config),
+) -> dict[str, Any]:
+    request = _as_request(payload)
+    entries = _normalize_entries(request)
+    platforms = _normalize_platforms(request.target_platforms)
+    timestamp = datetime.now().strftime("%Y%m%d-%H%M%S")
+    name = request.name or f"manual-query-{timestamp}"
+
+    batch_metadata = {
+        "family_key": "manual",
+        "source": "manual_api",
+        "query_count": len(entries),
+        "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,
+        )
+        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)
+        ]
+        run_key = f"manual-api:{batch.id}:{timestamp}"
+        log_path = RUNTIME_MANUAL_DIR / f"{run_key.replace(':', '-')}.log"
+        cmd = _pipeline_command(
+            batch_id=str(batch.id),
+            run_key=run_key,
+            platforms=platforms,
+            search_limit=request.search_limit,
+            display_limit=request.display_limit,
+            decode_limit=request.decode_limit,
+        )
+        run_metadata = {
+            "source": "manual_api",
+            "batch_id": str(batch.id),
+            "platforms": platforms,
+            "search_limit": request.search_limit,
+            "display_limit": request.display_limit,
+            "decode_limit": request.decode_limit,
+            "dry_ingest_record": True,
+            "log_path": str(log_path),
+            "command": cmd,
+        }
+        run = repo.create_acquisition_run(
+            batch_id=batch.id,
+            run_key=run_key,
+            status="pending",
+            note="manual query API queued",
+            metadata=run_metadata,
+        )
+
+    RUNTIME_MANUAL_DIR.mkdir(parents=True, exist_ok=True)
+    env = os.environ.copy()
+    env["PYTHONPATH"] = f"{ROOT}:{env.get('PYTHONPATH', '')}".rstrip(":")
+    try:
+        with log_path.open("ab") as stream:
+            process = subprocess.Popen(
+                cmd,
+                cwd=str(ROOT),
+                env=env,
+                stdout=stream,
+                stderr=subprocess.STDOUT,
+                start_new_session=True,
+            )
+    except OSError as exc:
+        with transaction(db_config) as conn:
+            PostgresAcquisitionRepository(conn).update_acquisition_run(
+                run.id,
+                status="failed",
+                error_message=f"启动手动 query pipeline 失败:{exc}",
+                metadata={"log_path": str(log_path), "command": cmd},
+        )
+        raise HTTPException(status_code=500, detail="启动后台 pipeline 失败") from exc
+
+    return {
+        "status": "queued",
+        "batch_id": str(batch.id),
+        "run_id": str(run.id),
+        "run_key": run_key,
+        "pid": process.pid,
+        "query_count": len(queries),
+        "queries": [_model_dump(query) for query in queries],
+        "summary_url": f"/api/acquisition/runs/{run.id}/summary",
+        "log_path": str(log_path),
+    }

+ 133 - 1
tests/test_app_api.py

@@ -1,17 +1,20 @@
 from __future__ import annotations
 from __future__ import annotations
 
 
+from contextlib import nullcontext
 from pathlib import Path
 from pathlib import Path
 from uuid import uuid4
 from uuid import uuid4
 
 
 from fastapi.testclient import TestClient
 from fastapi.testclient import TestClient
 
 
-from acquisition.domain import Query, QueryBatch
+from acquisition.domain import AcquisitionRun, Query, QueryBatch
 from app.api import app
 from app.api import app
 from app.dependencies import (
 from app.dependencies import (
     get_acquisition_repository,
     get_acquisition_repository,
+    get_creation_db_config,
     get_decode_repository,
     get_decode_repository,
     get_pipeline_repository,
     get_pipeline_repository,
 )
 )
+from app.routes import manual_queries
 from decode_content.models import (
 from decode_content.models import (
     DecodeJob,
     DecodeJob,
     DecodeResult,
     DecodeResult,
@@ -290,11 +293,140 @@ class FakePipelineRepo:
         )
         )
 
 
 
 
+class FakeManualRepo:
+    def __init__(self):
+        self.batch_id = uuid4()
+        self.run_id = uuid4()
+        self.queries = []
+        self.runs = []
+        self.updates = []
+
+    def create_query_batch(self, **kwargs):
+        self.batch_kwargs = kwargs
+        return QueryBatch(id=self.batch_id, **kwargs)
+
+    def add_query(self, **kwargs):
+        query = Query(id=uuid4(), **kwargs)
+        self.queries.append(query)
+        return query
+
+    def create_acquisition_run(self, **kwargs):
+        run = AcquisitionRun(id=self.run_id, **kwargs)
+        self.runs.append(run)
+        return run
+
+    def update_acquisition_run(self, run_id, **kwargs):
+        assert run_id == self.run_id
+        self.updates.append(kwargs)
+        base = self.runs[-1].model_dump()
+        base.update(kwargs)
+        return AcquisitionRun(**base)
+
+
 def _client(repo: FakeAcquisitionRepo):
 def _client(repo: FakeAcquisitionRepo):
     app.dependency_overrides[get_acquisition_repository] = lambda: repo
     app.dependency_overrides[get_acquisition_repository] = lambda: repo
     return TestClient(app)
     return TestClient(app)
 
 
 
 
+def test_manual_query_batch_api_creates_ready_batch_and_starts_pipeline(monkeypatch, tmp_path):
+    repo = FakeManualRepo()
+    popen_calls = []
+
+    class FakeProcess:
+        pid = 12345
+
+    def fake_popen(cmd, **kwargs):
+        popen_calls.append((cmd, kwargs))
+        return FakeProcess()
+
+    monkeypatch.setattr(manual_queries, "transaction", lambda _config: nullcontext(object()))
+    monkeypatch.setattr(manual_queries, "PostgresAcquisitionRepository", lambda _conn: repo)
+    monkeypatch.setattr(manual_queries, "RUNTIME_MANUAL_DIR", tmp_path)
+    monkeypatch.setattr(manual_queries.subprocess, "Popen", fake_popen)
+    app.dependency_overrides[get_creation_db_config] = lambda: object()
+    client = TestClient(app)
+
+    try:
+        response = client.post(
+            "/api/query-batches/manual",
+            json={
+                "queries": [
+                    "符号 视频 灵感 怎么做",
+                    "符号 视频 灵感 怎么做",
+                    {"query_text": "叙事体裁 视频 灵感 有哪些", "axes": {"实质": "叙事体裁"}},
+                ]
+            },
+        )
+        assert response.status_code == 200
+        data = response.json()
+        assert data["status"] == "queued"
+        assert data["batch_id"] == str(repo.batch_id)
+        assert data["run_id"] == str(repo.run_id)
+        assert data["pid"] == 12345
+        assert data["query_count"] == 2
+        assert data["summary_url"] == f"/api/acquisition/runs/{repo.run_id}/summary"
+        assert repo.batch_kwargs["generation_method"] == "manual_query_api_v1"
+        assert repo.batch_kwargs["target_platforms"] == ["xiaohongshu", "weixin", "douyin"]
+        assert [query.query_text for query in repo.queries] == [
+            "符号 视频 灵感 怎么做",
+            "叙事体裁 视频 灵感 有哪些",
+        ]
+        assert all(query.keep is True and query.status == "ready" for query in repo.queries)
+        assert all(query.metadata["family_key"] == "manual" for query in repo.queries)
+        assert repo.runs[0].status == "pending"
+        assert repo.runs[0].metadata["log_path"].endswith(".log")
+        assert repo.runs[0].metadata["command"]
+        assert popen_calls
+        cmd, kwargs = popen_calls[0]
+        assert "scripts/run_creation_pipeline.py" in cmd[1]
+        assert ["--platform", "xiaohongshu"] == cmd[-6:-4]
+        assert "--no-dry-ingest-record" not in cmd
+        assert kwargs["cwd"] == str(manual_queries.ROOT)
+        assert str(manual_queries.ROOT) in kwargs["env"]["PYTHONPATH"]
+        assert repo.updates == []
+    finally:
+        app.dependency_overrides.clear()
+
+
+def test_manual_query_batch_api_extracts_family_json(monkeypatch, tmp_path):
+    repo = FakeManualRepo()
+    monkeypatch.setattr(manual_queries, "transaction", lambda _config: nullcontext(object()))
+    monkeypatch.setattr(manual_queries, "PostgresAcquisitionRepository", lambda _conn: repo)
+    monkeypatch.setattr(manual_queries, "RUNTIME_MANUAL_DIR", tmp_path)
+    monkeypatch.setattr(manual_queries.subprocess, "Popen", lambda *args, **kwargs: type("P", (), {"pid": 1})())
+    app.dependency_overrides[get_creation_db_config] = lambda: object()
+    client = TestClient(app)
+
+    try:
+        response = client.post(
+            "/api/query-batches/manual",
+            json={
+                "families": [
+                    {
+                        "key": "f1",
+                        "name": "实质正交",
+                        "items": [{"query": "公共安全 视频 灵感 怎么做"}],
+                    }
+                ]
+            },
+        )
+        assert response.status_code == 200
+        assert repo.queries[0].query_text == "公共安全 视频 灵感 怎么做"
+        assert repo.queries[0].metadata["family_key"] == "manual"
+        assert repo.queries[0].metadata["source_family_key"] == "f1"
+    finally:
+        app.dependency_overrides.clear()
+
+
+def test_manual_query_batch_api_rejects_empty_and_invalid_platform():
+    client = TestClient(app)
+    assert client.post("/api/query-batches/manual", json={"queries": ["  "]}).status_code == 422
+    assert client.post(
+        "/api/query-batches/manual",
+        json={"queries": ["符号 视频 灵感 怎么做"], "target_platforms": ["bilibili"]},
+    ).status_code == 422
+
+
 def test_app_health_and_formal_acquisition_routes():
 def test_app_health_and_formal_acquisition_routes():
     repo = FakeAcquisitionRepo()
     repo = FakeAcquisitionRepo()
     client = _client(repo)
     client = _client(repo)

Некоторые файлы не были показаны из-за большого количества измененных файлов