video.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238
  1. """Native whole-video reader for creation knowledge decode."""
  2. from __future__ import annotations
  3. import base64
  4. import logging
  5. import os
  6. import re
  7. import time
  8. from pathlib import Path
  9. from typing import Any, Callable, Optional
  10. import httpx
  11. from core.config import Settings
  12. from core.jsonio import extract_json_object
  13. from core.media_download import download_media_bytes
  14. from core.models import Card, CardExtract, ExtractedContent, Post
  15. from core.prompts import load_prompt
  16. from pipeline.tracing import TraceContext, TraceWriter, hash_prompt, redact_headers, timed_ms
  17. logger = logging.getLogger(__name__)
  18. class VideoExtractError(RuntimeError):
  19. pass
  20. def _mmss_to_sec(value: Any) -> Optional[float]:
  21. if value is None:
  22. return None
  23. if isinstance(value, (int, float)):
  24. return float(value)
  25. parts = str(value).strip().split(":")
  26. try:
  27. nums = [float(part) for part in parts]
  28. except ValueError:
  29. return None
  30. sec = 0.0
  31. for number in nums:
  32. sec = sec * 60 + number
  33. return sec
  34. def _default_download(url: str, platform: str, timeout: float = 180.0) -> bytes:
  35. return download_media_bytes(url, platform, timeout=timeout)
  36. def _seg_content(segment: dict[str, Any]) -> str:
  37. parts = [segment.get("title") or ""]
  38. for label, key in (("What", "what"), ("Why", "why"), ("How", "how")):
  39. value = segment.get(key)
  40. if value and str(value).strip().lower() not in {"null", "none", ""}:
  41. parts.append(f"{label}:{value}")
  42. return "。".join(part for part in parts if part)
  43. def extract_video(
  44. post: Post,
  45. *,
  46. settings: Settings,
  47. http_post: Callable[..., Any] = httpx.post,
  48. video_path: Optional[str] = None,
  49. downloader: Optional[Callable[[str, str], bytes]] = None,
  50. timeout: float = 600.0,
  51. save_path: Optional[Path] = None,
  52. public_url: Optional[str] = None,
  53. oss_video_url: Optional[str] = None,
  54. trace_writer: TraceWriter | None = None,
  55. trace_context: TraceContext | None = None,
  56. ) -> ExtractedContent:
  57. key = settings.openrouter_api_key
  58. if not key:
  59. raise VideoExtractError("missing OPENROUTER_API_KEY")
  60. if oss_video_url:
  61. media_url = oss_video_url
  62. public_url = public_url or oss_video_url
  63. else:
  64. if post.platform == "bilibili" and not (video_path and os.path.exists(video_path)):
  65. raise VideoExtractError("B站视频暂未接入:需取 voice_data 音轨 + ffmpeg 合并")
  66. if video_path and os.path.exists(video_path):
  67. data = open(video_path, "rb").read()
  68. logger.info("video_extract using local file %s (%d bytes)", video_path, len(data))
  69. else:
  70. if not post.video_urls:
  71. raise VideoExtractError(f"post {post.id} has no video_urls and no video_path")
  72. url = post.video_urls[0]
  73. if post.platform == "douyin" and "ratio=" in url:
  74. url = re.sub(r"ratio=[^&]+", f"ratio={settings.douyin_ratio}", url)
  75. dl = downloader or _default_download
  76. try:
  77. data = dl(url, post.platform)
  78. except Exception as exc:
  79. raise VideoExtractError(f"视频下载失败: {exc}") from exc
  80. if not data:
  81. raise VideoExtractError("视频字节为空")
  82. if save_path is not None:
  83. save_path.parent.mkdir(parents=True, exist_ok=True)
  84. Path(save_path).write_bytes(data)
  85. logger.info("video_extract saved whole video %s (%d bytes)", save_path, len(data))
  86. media_url = "data:video/mp4;base64," + base64.b64encode(data).decode()
  87. prompt = load_prompt("extract_video").format()
  88. body = {
  89. "model": settings.video_model,
  90. "messages": [
  91. {
  92. "role": "user",
  93. "content": [
  94. {"type": "text", "text": prompt},
  95. {"type": "video_url", "video_url": {"url": media_url}},
  96. ],
  97. }
  98. ],
  99. }
  100. started = time.perf_counter()
  101. headers = {"Authorization": f"Bearer {key}", "Content-Type": "application/json"}
  102. try:
  103. resp = http_post(
  104. f"{settings.openrouter_base_url.rstrip('/')}/chat/completions",
  105. headers=headers,
  106. json=body,
  107. timeout=timeout,
  108. )
  109. resp.raise_for_status()
  110. response_json = resp.json()
  111. content = response_json["choices"][0]["message"]["content"]
  112. except httpx.HTTPError as exc:
  113. if trace_writer is not None:
  114. trace_writer.llm_call(
  115. context=trace_context or TraceContext(stage="decode", substage="read_video"),
  116. stage="decode",
  117. substage="read_video",
  118. provider="openrouter",
  119. model_name=settings.video_model,
  120. endpoint=f"{settings.openrouter_base_url.rstrip('/')}/chat/completions",
  121. prompt_name="extract_video",
  122. prompt_hash=hash_prompt(prompt),
  123. request_payload={"headers": redact_headers(headers), **body},
  124. status="failed",
  125. error_message=str(exc),
  126. latency_ms=timed_ms(started),
  127. )
  128. raise VideoExtractError(f"openrouter_http_error: {exc}") from exc
  129. except (KeyError, IndexError, TypeError, ValueError) as exc:
  130. if trace_writer is not None:
  131. trace_writer.llm_call(
  132. context=trace_context or TraceContext(stage="decode", substage="read_video"),
  133. stage="decode",
  134. substage="read_video",
  135. provider="openrouter",
  136. model_name=settings.video_model,
  137. endpoint=f"{settings.openrouter_base_url.rstrip('/')}/chat/completions",
  138. prompt_name="extract_video",
  139. prompt_hash=hash_prompt(prompt),
  140. request_payload={"headers": redact_headers(headers), **body},
  141. status="failed",
  142. error_message=str(exc),
  143. latency_ms=timed_ms(started),
  144. )
  145. raise VideoExtractError(f"openrouter_response_invalid: {exc}") from exc
  146. obj = extract_json_object(content)
  147. segments = obj.get("segments") or []
  148. cards: list[Card] = []
  149. card_extracts: list[CardExtract] = []
  150. for idx, segment in enumerate(segments, start=1):
  151. if not isinstance(segment, dict):
  152. continue
  153. cards.append(
  154. Card(
  155. index=idx,
  156. kind="segment",
  157. url=public_url,
  158. start=_mmss_to_sec(segment.get("start")),
  159. end=_mmss_to_sec(segment.get("end")),
  160. )
  161. )
  162. card_extracts.append(CardExtract(index=idx, content=_seg_content(segment)))
  163. post.cards = cards
  164. if trace_writer is not None:
  165. trace_writer.llm_call(
  166. context=trace_context or TraceContext(stage="decode", substage="read_video"),
  167. stage="decode",
  168. substage="read_video",
  169. provider="openrouter",
  170. model_name=settings.video_model,
  171. endpoint=f"{settings.openrouter_base_url.rstrip('/')}/chat/completions",
  172. prompt_name="extract_video",
  173. prompt_hash=hash_prompt(prompt),
  174. request_payload={"headers": redact_headers(headers), **body},
  175. response_payload=response_json,
  176. parsed_payload={
  177. "segment_count": len(card_extracts),
  178. "overall": obj.get("overall"),
  179. "video_title": obj.get("video_title"),
  180. },
  181. status="done",
  182. latency_ms=timed_ms(started),
  183. )
  184. return ExtractedContent(
  185. text=str(obj.get("overall") or obj.get("video_title") or ""),
  186. cards=card_extracts,
  187. is_empty=len(card_extracts) == 0,
  188. )
  189. def read_video(
  190. post: Post,
  191. *,
  192. settings: Settings,
  193. http_post: Callable[..., Any] | None = None,
  194. video_path: str | None = None,
  195. downloader: Callable[[str, str], bytes] | None = None,
  196. save_path: Path | None = None,
  197. public_url: str | None = None,
  198. oss_video_url: str | None = None,
  199. trace_writer: TraceWriter | None = None,
  200. trace_context: TraceContext | None = None,
  201. ) -> ExtractedContent:
  202. kwargs: dict[str, Any] = {
  203. "settings": settings,
  204. "video_path": video_path,
  205. "downloader": downloader,
  206. "save_path": save_path,
  207. "public_url": public_url or oss_video_url,
  208. "oss_video_url": oss_video_url,
  209. "trace_writer": trace_writer,
  210. "trace_context": trace_context,
  211. }
  212. if http_post is not None:
  213. kwargs["http_post"] = http_post
  214. return extract_video(post, **kwargs)
  215. __all__ = ["VideoExtractError", "extract_video", "read_video"]