"""Image-text multimodal reader for creation knowledge decode.""" from __future__ import annotations import logging import time from typing import Any, Callable, Mapping, Optional import httpx from core.config import load_env_file from core.jsonio import extract_json_object, to_bool from core.models import Card, CardExtract, ExtractedContent, Post from core.prompts import load_prompt from pipeline.tracing import TraceContext, TraceWriter, hash_prompt, redact_headers, timed_ms logger = logging.getLogger(__name__) DEFAULT_MODEL = "qwen-vl-plus" DEFAULT_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" DEFAULT_TIMEOUT = 120.0 MAX_CARDS = 12 def _card_label(card: Card) -> str: if card.kind == "frame" and card.timestamp is not None: ts = int(card.timestamp) return f"【卡片{card.index} · {ts // 60:02d}:{ts % 60:02d}】" return f"【卡片{card.index}】" SYSTEM_PROMPT = ( "你是创作知识提取助手。从给定的小红书帖子(标题、正文、图片、视频)中," "提取真正能指导『如何创作内容』的知识;知识常在图片/视频里而非正文。" "忠实提取、不编造。只输出一个 JSON 对象,不要解释或 markdown。" ) class ExtractorError(RuntimeError): pass class BailianExtractor: def __init__( self, *, api_key: str, model: str = DEFAULT_MODEL, base_url: str = DEFAULT_BASE_URL, timeout_seconds: float = DEFAULT_TIMEOUT, http_post: Callable[..., Any] = httpx.post, max_cards: int = MAX_CARDS, ) -> None: if not api_key: raise ExtractorError("missing ALIYUN_BAILIAN_API_KEY") self.api_key = api_key self.model = model self.base_url = base_url.rstrip("/") self.timeout_seconds = timeout_seconds self.http_post = http_post self.max_cards = max_cards @classmethod def from_env(cls, env: Mapping[str, str] | None = None, env_file: str = ".env") -> "BailianExtractor": source = dict(load_env_file(env_file)) if env: source.update(env) api_key = source.get("ALIYUN_BAILIAN_API_KEY") or "" return cls( api_key=api_key, model=source.get("ALIYUN_BAILIAN_VL_MODEL") or source.get("ALIYUN_BAILIAN_MODEL") or DEFAULT_MODEL, base_url=source.get("ALIYUN_BAILIAN_BASE_URL") or DEFAULT_BASE_URL, timeout_seconds=float(source.get("ALIYUN_BAILIAN_TIMEOUT_SECONDS") or DEFAULT_TIMEOUT), max_cards=int(source.get("CK_MAX_CARDS") or MAX_CARDS), ) def _cards(self, post: Post) -> list[Card]: cards = post.cards or [ Card(index=i, kind="image", url=url) for i, url in enumerate(post.image_urls, start=1) ] if len(cards) > self.max_cards: dropped = [card.index for card in cards[self.max_cards :]] logger.warning( "post %s card count %d exceeds MAX_CARDS=%d, dropping cards %s", post.id, len(cards), self.max_cards, dropped, ) cards = cards[: self.max_cards] return cards def build_messages(self, post: Post) -> list[dict[str, Any]]: user_text = load_prompt("extract").format( title=post.title or "(无)", topics="、".join(post.topic_list) or "(无)", body=post.body_text or "(空)", ) parts: list[dict[str, Any]] = [{"type": "text", "text": user_text}] for card in self._cards(post): parts.append({"type": "text", "text": _card_label(card)}) parts.append({"type": "image_url", "image_url": {"url": card.url}}) return [ {"role": "system", "content": SYSTEM_PROMPT}, {"role": "user", "content": parts}, ] def extract( self, post: Post, *, trace_writer: TraceWriter | None = None, trace_context: TraceContext | None = None, ) -> ExtractedContent: messages = self.build_messages(post) last_exc: Optional[Exception] = None for attempt in range(2): started = time.perf_counter() headers = { "Authorization": f"Bearer {self.api_key}", "Content-Type": "application/json", } body = {"model": self.model, "messages": messages} try: resp = self.http_post( f"{self.base_url}/chat/completions", headers=headers, json=body, timeout=self.timeout_seconds, ) resp.raise_for_status() response_json = resp.json() content = response_json["choices"][0]["message"]["content"] data = extract_json_object(content) cards = [ CardExtract(index=int(card["index"]), content=str(card.get("content") or "")) for card in (data.get("cards") or []) if isinstance(card, dict) and card.get("index") is not None ] if trace_writer is not None: trace_writer.llm_call( context=trace_context or TraceContext(stage="decode", substage="read_imgtext"), stage="decode", substage="read_imgtext", provider="bailian", model_name=self.model, endpoint=f"{self.base_url}/chat/completions", prompt_name="extract_imgtext", prompt_hash=hash_prompt(SYSTEM_PROMPT), request_payload={"headers": redact_headers(headers), **body}, response_payload=response_json, parsed_payload={ "is_empty": to_bool(data.get("is_empty")), "card_count": len(cards), "text": data.get("text"), }, status="done", latency_ms=timed_ms(started), attempt_index=attempt + 1, ) return ExtractedContent( text=str(data.get("text") or ""), cards=cards, from_image=str(data.get("from_image") or ""), from_video=str(data.get("from_video") or ""), is_empty=to_bool(data.get("is_empty")), ) except httpx.HTTPError as exc: last_exc = exc if trace_writer is not None: trace_writer.llm_call( context=trace_context or TraceContext(stage="decode", substage="read_imgtext"), stage="decode", substage="read_imgtext", provider="bailian", model_name=self.model, endpoint=f"{self.base_url}/chat/completions", prompt_name="extract_imgtext", prompt_hash=hash_prompt(SYSTEM_PROMPT), request_payload={"headers": redact_headers(headers), **body}, status="failed", error_message=str(exc), latency_ms=timed_ms(started), attempt_index=attempt + 1, ) if attempt == 0: continue raise ExtractorError(f"bailian_http_error: {exc}") from exc except (KeyError, IndexError, TypeError, ValueError) as exc: last_exc = exc if trace_writer is not None: trace_writer.llm_call( context=trace_context or TraceContext(stage="decode", substage="read_imgtext"), stage="decode", substage="read_imgtext", provider="bailian", model_name=self.model, endpoint=f"{self.base_url}/chat/completions", prompt_name="extract_imgtext", prompt_hash=hash_prompt(SYSTEM_PROMPT), request_payload={"headers": redact_headers(headers), **body}, status="failed", error_message=str(exc), latency_ms=timed_ms(started), attempt_index=attempt + 1, ) if attempt == 0: continue raise ExtractorError(f"bailian_response_invalid: {exc}") from exc raise ExtractorError(f"bailian_unknown_error: {last_exc}") GeminiExtractor = BailianExtractor def extract_content( post: Post, *, client: Optional[BailianExtractor] = None, env_file: str = ".env", trace_writer: TraceWriter | None = None, trace_context: TraceContext | None = None, ) -> ExtractedContent: client = client or GeminiExtractor.from_env(env_file=env_file) return client.extract(post, trace_writer=trace_writer, trace_context=trace_context) def read_imgtext( post: Post, *, extractor: BailianExtractor | None = None, trace_writer: TraceWriter | None = None, trace_context: TraceContext | None = None, ) -> ExtractedContent: client = extractor or BailianExtractor.from_env() return client.extract(post, trace_writer=trace_writer, trace_context=trace_context) __all__ = [ "BailianExtractor", "ExtractorError", "GeminiExtractor", "extract_content", "read_imgtext", ]