| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157 |
- """M5 离线编排测试:注入 fake store/fetch/extract/chat,验证 done 与 rejected 两条路径。
- 真实落库的 e2e 见 scripts/run_batch.py(在能连 DB 的云端跑)。
- """
- from __future__ import annotations
- import json
- from pathlib import Path
- from creation_knowledge.config import PgConfig, Settings
- from creation_knowledge.integrations.crawler import parse_detail_response
- from creation_knowledge.models import ExtractedContent, Post
- from creation_knowledge.pipeline import run_pipeline
- FIXTURES = Path(__file__).parent / "fixtures"
- CID = "67e4bdf50000000006028a59"
- def _settings(data_dir: str = "") -> Settings:
- return Settings(
- pg=PgConfig(host="h", port=5432, user="u", password="p", database="d"),
- crawler_base_url="http://x", crawler_key="", crawler_timeout=30,
- video_model="m", gemini_api_key="", openrouter_base_url="http://x",
- openrouter_api_key="k", llm_model="m", knowhub_api="http://x",
- ingest_enabled=False, max_cards=12, frames_dir="runtime/frames",
- douyin_ratio="540p", data_dir=data_dir,
- )
- class FakeStore:
- def __init__(self):
- self.posts: dict = {}
- self.items: list = []
- def upsert_post(self, post):
- self.posts[post.id] = {"stage": "fetched"}
- def set_extracted(self, pid, ex):
- self.posts[pid].update(stage="extracted", extracted=ex)
- def set_screening(self, pid, sc):
- self.posts[pid].update(stage="screened", screening=sc)
- def update_stage(self, pid, stage):
- self.posts[pid]["stage"] = stage
- def save_item(self, pid, item, deco, payload):
- self.items.append({"post_id": pid, "item": item, "deco": deco, "payload": payload})
- return len(self.items)
- def update_item_ingest(self, iid, status, kid):
- self.items[iid - 1].update(ingest_status=status, knowledge_id=kid)
- def clear_items(self, pid):
- before = len(self.items)
- self.items = [it for it in self.items if it["post_id"] != pid]
- return before - len(self.items)
- def _fetch(url):
- resp = json.loads((FIXTURES / f"xhs_case_{CID}.json").read_text("utf-8"))
- return parse_detail_response(resp, fallback_content_id=CID)
- def _extract(post):
- return ExtractedContent(text=post.body_text or "内容", is_empty=False)
- def _chat_pass(system, user):
- if "筛选" in system:
- return {"passed": True, "score": 8, "reason": "ok"}
- if "拆分" in system:
- return {"items": [{"title": "脚本要素", "knowledge_types": ["how"],
- "what": "null", "why": "null", "how": "逐个填要素",
- "evidence": ["原句"]}]}
- if "解构" in system:
- return {"stages": ["脚本"], "scopes": [{"scope_type": "form", "value": "操作流程"}],
- "stage_reason": "r", "scope_reason": "r"}
- raise AssertionError(f"unexpected system: {system}")
- def _chat_reject(system, user):
- if "筛选" in system:
- return {"passed": False, "score": 2, "reason": "只是作品本身"}
- raise AssertionError("rejected post 不应进入拆分/解构")
- def test_pipeline_done_path():
- store = FakeStore()
- results = run_pipeline(
- ["https://www.xiaohongshu.com/explore/" + CID],
- settings=_settings(), ingest_enabled=False, store=store,
- fetch_fn=_fetch, extract_fn=_extract, chat=_chat_pass,
- )
- r = results[0]
- assert r["status"] == "done" and r["items"] == 1
- pid = f"xhs_{CID}"
- assert store.posts[pid]["stage"] == "done"
- assert len(store.items) == 1
- saved = store.items[0]
- assert saved["payload"]["source"]["id"] == pid
- assert saved["payload"]["dim_attributes"] == ["how"]
- def test_pipeline_rejected_path():
- store = FakeStore()
- results = run_pipeline(
- [CID], settings=_settings(), ingest_enabled=False, store=store,
- fetch_fn=_fetch, extract_fn=_extract, chat=_chat_reject,
- )
- r = results[0]
- assert r["status"] == "rejected"
- assert store.posts[f"xhs_{CID}"]["stage"] == "rejected"
- assert store.items == []
- def _extract_must_not_run(post):
- raise AssertionError("skip 的帖子不应进入提取")
- def test_pipeline_skip_xhs_video():
- """小红书视频帖:content_type=video 但无视频直链 → skip,不进提取/筛选。"""
- store = FakeStore()
- post = Post(id="xhs_vid1", platform="xiaohongshu", url="u", content_id="vid1",
- content_type="video", video_urls=[], image_urls=["http://cover.jpg"])
- results = run_pipeline(
- ["vid1"], settings=_settings(), ingest_enabled=False, store=store,
- fetch_fn=lambda url: post, extract_fn=_extract_must_not_run, chat=_chat_pass,
- )
- r = results[0]
- assert r["status"] == "skipped" and r["reason"] == "video_no_direct_url"
- assert store.posts["xhs_vid1"]["stage"] == "skipped"
- assert store.items == []
- def test_pipeline_skip_empty_media():
- """既无图也无视频 → skip(empty_media)。"""
- store = FakeStore()
- post = Post(id="xhs_empty", platform="xiaohongshu", url="u", content_id="empty",
- content_type="normal", video_urls=[], image_urls=[])
- results = run_pipeline(
- ["empty"], settings=_settings(), ingest_enabled=False, store=store,
- fetch_fn=lambda url: post, extract_fn=_extract_must_not_run, chat=_chat_pass,
- )
- assert results[0]["status"] == "skipped" and results[0]["reason"] == "empty_media"
- assert store.items == []
- def test_pipeline_idempotent_rerun():
- """同帖重跑:知识片段先清后写,不叠加。"""
- store = FakeStore()
- for _ in range(3):
- run_pipeline(
- [CID], settings=_settings(), ingest_enabled=False, store=store,
- fetch_fn=_fetch, extract_fn=_extract, chat=_chat_pass,
- )
- assert len(store.items) == 1, "重跑三次仍应只有 1 条(先清后写)"
|