| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245 |
- """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",
- ]
|