| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119 |
- """流水线:把 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.integrations.video_extract import VideoExtractError, extract_video
- 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)}
- store.upsert_post(post) # 图文卡片在此落;视频段卡在提取后补落
- # 2) 多模态提取(图文=逐图;视频=原生整段视频→段卡,extract 内写 post.cards)
- try:
- content = extract_fn(post)
- except (ExtractorError, VideoExtractError) as exc:
- store.update_stage(post.id, "failed")
- return {"url": url, "post_id": post.id, "status": "extract_failed", "error": str(exc)}
- if post.cards:
- store.upsert_post(post) # 视频段卡落库(图文重复 upsert 无害)
- store.set_extracted(post.id, content.model_dump()) # stage=extracted
- # 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))
- if extract_fn is None:
- _image_client = GeminiExtractor.from_env(env_file=env_file)
- def _dispatch_extract(post: Post) -> ExtractedContent:
- # 视频帖 → 原生整段视频(OpenRouter base64);图文帖 → 逐图提取
- if (post.content_type or "").lower() == "video" or post.video_urls:
- return extract_video(post, settings=settings)
- return _image_client.extract(post)
- extract_fn = _dispatch_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
|