"""流水线:把 6 步串起来,每步落库。INGEST_ENABLED=False 时只组装+存,不真实发送。 依赖都可注入(fetch/extract/chat/store),便于离线测试与替换平台。 """ from __future__ import annotations import logging from typing import Callable, Optional from creation_knowledge.config import Settings from creation_knowledge.integrations.crawler import CrawlerError, fetch_post_detail from creation_knowledge.integrations.db import CkStore from creation_knowledge.integrations.extractor import ExtractorError, GeminiExtractor from creation_knowledge.integrations.llm import ChatFn, default_chat from creation_knowledge.ingest import IngestError, ingest as real_ingest from creation_knowledge.models import ExtractedContent, Post from creation_knowledge.stages import ( build_ingest_payload, deconstruct_item, screen_post, split_post, ) logger = logging.getLogger(__name__) FetchFn = Callable[[str], Post] ExtractFn = Callable[[Post], ExtractedContent] def _process_one( url: str, *, settings: Settings, store: CkStore, fetch_fn: FetchFn, extract_fn: ExtractFn, chat: ChatFn, ingest_enabled: bool, ) -> dict: # 1) 拉取 try: post = fetch_fn(url) except CrawlerError as exc: return {"url": url, "status": "fetch_failed", "error": str(exc)} # 1.5) 视频帖:抽帧补成卡片(best-effort,失败不阻塞整帖) if post.video_urls and not post.cards: try: from creation_knowledge.integrations.video_frames import extract_frames post.cards = extract_frames( post.video_urls[0], out_dir=f"{settings.frames_dir}/{post.id}", url_prefix=f"/frames/{post.id}", platform=post.platform, max_frames=settings.max_cards, ) except Exception as exc: # 抽帧失败不影响图文/正文路径 logger.warning("post %s 抽帧失败: %s", post.id, exc) store.upsert_post(post) # stage=fetched(含 cards) # 2) 多模态提取 try: content = extract_fn(post) store.set_extracted(post.id, content.model_dump()) # stage=extracted except ExtractorError as exc: store.update_stage(post.id, "failed") return {"url": url, "post_id": post.id, "status": "extract_failed", "error": str(exc)} # 3) 筛选 screening = screen_post(post, content, chat=chat) store.set_screening(post.id, screening.model_dump()) # stage=screened if not screening.passed: store.update_stage(post.id, "rejected") return {"url": url, "post_id": post.id, "status": "rejected", "score": screening.score, "reason": screening.reason} # 4) 拆分 -> 5) 解构 -> 6) 组装(+可选入库) items = split_post(post, content, chat=chat) item_ids = [] for item in items: deco = deconstruct_item(item, chat=chat) payload = build_ingest_payload(post, item, deco) item_id = store.save_item( post.id, item.model_dump(), deco.model_dump(), payload.model_dump() ) item_ids.append(item_id) if ingest_enabled: try: res = real_ingest(payload, settings=settings) store.update_item_ingest(item_id, "ingested", res.get("knowledge_id")) except IngestError: store.update_item_ingest(item_id, "failed", None) store.update_stage(post.id, "done") return {"url": url, "post_id": post.id, "status": "done", "items": len(item_ids)} def run_pipeline( urls: list[str], *, settings: Optional[Settings] = None, env_file: str = ".env", ingest_enabled: Optional[bool] = None, store: Optional[CkStore] = None, fetch_fn: Optional[FetchFn] = None, extract_fn: Optional[ExtractFn] = None, chat: Optional[ChatFn] = None, ) -> list[dict]: settings = settings or Settings.from_env(env_file) ingest_enabled = settings.ingest_enabled if ingest_enabled is None else ingest_enabled store = store or CkStore(settings.pg) fetch_fn = fetch_fn or (lambda url: fetch_post_detail(url, settings=settings)) extract_fn = extract_fn or GeminiExtractor.from_env(env_file=env_file).extract chat = chat or default_chat(env_file) results = [] for url in urls: results.append(_process_one( url, settings=settings, store=store, fetch_fn=fetch_fn, extract_fn=extract_fn, chat=chat, ingest_enabled=ingest_enabled, )) return results