"""Web 管理台的「链接解析 + 预览」任务层(供 web_hub.py 调用)。 群消息链路(`handlers/entry.py::dispatch_url`)从解析到投递都绑在 event/bot 上, Web 上下文里第一次 `UniMessage.send()` 就会抛 SerializeFailed,所以这里**只复用 纯函数层**:fetchers 的解析/下载 + `storage/s3.py::upload_with_plan`, 全程不发消息、不碰 event、不读群策略。 存储按**默认策略**走(`STORE.default_policy()` 的 plan / upload_public): 产出什么链接就返回什么链接 —— 开了 `upload_public` 才有公网链接,否则只有 局域网 S3 的预签名链接(1 小时过期,靠 `refresh()` 重传换新)。 任务表在内存里(进程重启即清空,前端按"查不到就算了"处理): submit(urls, force=False) 新建任务;同 URL 已有 queued/running 任务时**复用** list_jobs() 全部任务,新的在前 refresh(job_id) 对已落盘文件重跑上传,换一批新链接 抖音每次都新起一个 Chrome,所以并发闸门固定 `Semaphore(2)`;单任务 240s 超时兜底。 单文件上传失败只标该文件(`files[i].error`),不整体判失败 —— 多图作品挂一张 不该让整条任务失败,缺的那个文件点「刷新链接」还能补回来。 进度只有两档(解析中 → 上传中):解析和下载都在 fetcher 内部完成,中间没有可挂的钩子。 """ from __future__ import annotations import asyncio import re import subprocess import time import uuid from collections.abc import Awaitable, Callable, Iterable from dataclasses import dataclass, field from pathlib import Path from typing import TYPE_CHECKING, Any, Optional from nonebot import logger if TYPE_CHECKING: # 仅类型检查:运行期不导入,单测可裸加载本模块 from ..policy import Policy # ─────────────────────────── 常量 ─────────────────────────── #: URL 的终止符:空白 + 中文标点(全角括号/引号/破折号也要断) URL_STOP = ",。!?、;:()【】《》「」『』…—“”‘’" #: 从整段文本里挑 http(s) 链接。形制同 handlers/entry.py::URL_PATTERN, #: 但把中文标点也当成终止符 —— `链接A,链接B` 中间没有空格也能各自成条 #: (`\S+` 会把它们连同后面的中文连成一个词,整段粘进来时更容易踩到)。 URL_PATTERN = re.compile(r"https?://[^" + re.escape(URL_STOP) + r"\s]+") #: 兜底再洗一遍收尾标点(entry.py::dispatch_url 同款) URL_TRAILING = URL_STOP #: 任务表上限与终态任务的保留时长(秒) MAX_JOBS = 100 TERMINAL_TTL = 30 * 60 #: 单任务超时(含下载,不含排队等闸门的时间) JOB_TIMEOUT = 240 #: 并发闸门:抖音每次都新起 Chrome,必须限流 MAX_CONCURRENCY = 2 VIDEO_SUFFIXES = {".mp4", ".webm", ".mov", ".flv", ".mkv", ".ts", ".m4v"} IMAGE_SUFFIXES = {".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp", ".avif"} #: 视频封面(关键帧):宽度压到 480 够卡片用,别把 4K 原帧传上云 POSTER_WIDTH = 480 POSTER_SUFFIX = "_封面.jpg" POSTER_TIMEOUT = 20 #: 状态机:queued → running → done / failed ACTIVE_STATUSES = ("queued", "running") # ─────────────────────────── 纯函数 ─────────────────────────── def extract_urls(text: str) -> list[str]: """整段文本 → 去重后的 URL 列表(保序,去掉尾随中文标点)。 多行粘贴时一行一条;URL 后面跟的"。"这类标点不能进链接。 """ urls: list[str] = [] for raw in URL_PATTERN.findall(text or ""): url = raw.rstrip(URL_TRAILING) if url and url not in urls: urls.append(url) return urls def kind_of(path: Path | str) -> str: """文件类型(前端据此决定内联预览方式)。""" suffix = Path(path).suffix.lower() if suffix in VIDEO_SUFFIXES: return "video" if suffix in IMAGE_SUFFIXES: return "image" return "file" # ─────────────────────────── 依赖注入点 ─────────────────────────── #: (标题, 落盘文件列表, 是否图文作品) FetchResult = tuple[Optional[str], list[Path], bool] Fetcher = Callable[[str], Awaitable[FetchResult]] #: (局域网链接, 公网链接) —— 拿不到时返回 "" / None Uploader = Callable[[Path, "Optional[Policy]"], Awaitable[tuple[str, Optional[str]]]] #: 视频 → 封面图(抽关键帧);抽不出来返回 None PosterMaker = Callable[[Path], Awaitable[Optional[Path]]] def ffmpeg_path() -> str: """ffmpeg 可执行文件(解析逻辑复用 video_downloader;单测可替换)""" from .fetchers.video_downloader import get_ffmpeg_path return get_ffmpeg_path() async def make_poster(video: Path) -> Optional[Path]: """抽一帧当视频封面,落同目录 `{作品名}_封面.jpg`;失败返回 None(不影响任务)。 先取第 1 秒(躲开黑场/淡入),整段不到 1 秒的视频退回第 0 秒再试一次。 """ out = video.with_name(f"{video.stem}{POSTER_SUFFIX}") for seek in ("1", "0"): cmd = [ ffmpeg_path(), "-ss", seek, "-i", str(video), "-frames:v", "1", "-vf", f"scale={POSTER_WIDTH}:-2", "-q:v", "4", "-y", str(out), ] try: # 同步子进程必须丢线程池,否则抽帧期间整个事件循环都停摆 await asyncio.to_thread( subprocess.run, cmd, capture_output=True, timeout=POSTER_TIMEOUT ) except Exception: # noqa: BLE001 — 封面失败不该拖累任务 logger.warning(f"抽关键帧失败:{video.name}") return None if out.exists() and out.stat().st_size > 0: return out return None def platform_of(url: str) -> Optional[str]: """URL → 平台规范标签(延迟导入 policy,便于单测裸加载本模块)。""" from ..policy import match_platform return match_platform(url) def default_policy() -> "Optional[Policy]": """Web 解析统一按默认策略走(不是某个群的策略)。""" from ..policy import STORE return STORE.default_policy() async def fetch_media(url: str) -> FetchResult: """按平台解析并下载媒体(纯函数层:不发消息、不碰 event、不上传)。 分派规则与 `handlers/entry.py::dispatch_url` 一致:抖音 → parse_douyin; B站动态/专栏 → fetch_bilibili_content;其余 B站链接 → download_video(yt-dlp); 小红书 → fetch_rednote_content;其它平台 → download_video。 b23 / xhslink / 抖音短链的重定向在各 fetcher 内部自己处理。 """ from ..handlers.douyin import parse_douyin from ..policy import BILIBILI, DOUYIN, XHS from .fetchers.bilibili_content import fetch_bilibili_content from .fetchers.rednote_content import fetch_rednote_content from .fetchers.video_downloader import download_video platform = platform_of(url) title: Optional[str] = None parsed: Path | list[Path] | None = None image_post = False if platform == DOUYIN: title, parsed, image_post = await parse_douyin(url) elif platform == XHS: title, parsed = await fetch_rednote_content(url) image_post = isinstance(parsed, list) elif platform == BILIBILI and any( kw in url for kw in ( "bilibili.com/opus", "bilibili.com/dynamic", "t.bilibili.com", "bilibili.com/read", ) ): title, parsed = await fetch_bilibili_content(url) image_post = isinstance(parsed, list) else: video_file, title = await download_video(url) parsed = video_file if isinstance(parsed, list): files = [Path(p) for p in parsed] else: files = [Path(parsed)] if parsed else [] return title, files, image_post async def upload_file( path: Path, policy: "Optional[Policy]" ) -> tuple[str, Optional[str]]: """上传单个文件 → (局域网链接, 公网链接)。 boto3 是同步的,必须丢线程池:直接在协程里跑会卡死整个事件循环 (Web 轮询、群消息全都停摆)。 """ from .storage.s3 import upload_with_plan return await asyncio.to_thread(upload_with_plan, path, policy=policy) # ─────────────────────────── 任务模型 ─────────────────────────── @dataclass class WebJob: """一条解析任务;`to_dict()` 就是 API 的返回体。""" id: str url: str platform: Optional[str] = None title: str = "" status: str = "queued" #: 运行中的进度文案:解析中 / 上传中(见模块文档:解析与下载不细分) stage: str = "" #: 每个文件:{name, kind, size, local_url, public_url, error} files: list[dict[str, Any]] = field(default_factory=list) #: 卡片封面:有图取第一张图,全是视频则抽第一帧(见 `_resolve_cover`) cover_url: str = "" error: str = "" created_at: float = field(default_factory=time.time) updated_at: float = field(default_factory=time.time) #: 落盘的源文件(不出 API;`refresh()` 重传要用) paths: list[Path] = field(default_factory=list, repr=False) #: 抽出来的封面文件(不出 API;`refresh()` 重传要用) poster_path: Optional[Path] = field(default=None, repr=False) @property def active(self) -> bool: """还在排队或运行(终态任务才会被 prune 掉)。""" return self.status in ACTIVE_STATUSES def to_dict(self) -> dict[str, Any]: return { "id": self.id, "url": self.url, "platform": self.platform, "title": self.title, "status": self.status, "stage": self.stage, "files": self.files, "cover_url": self.cover_url, "error": self.error, "created_at": self.created_at, "updated_at": self.updated_at, } class JobManager: """内存任务表:上限 MAX_JOBS 条,终态任务 TTL 30 分钟后自动清掉。 fetch / upload / poster / policy 四个依赖可注入(单测塞假实现,不碰网络与磁盘)。 """ def __init__( self, *, fetch: Optional[Fetcher] = None, upload: Optional[Uploader] = None, poster: Optional[PosterMaker] = None, policy: Optional[Callable[[], "Optional[Policy]"]] = None, ) -> None: self._jobs: dict[str, WebJob] = {} self._tasks: dict[str, asyncio.Task[None]] = {} self._sem: Optional[asyncio.Semaphore] = None self._fetch: Fetcher = fetch or fetch_media self._upload: Uploader = upload or upload_file self._poster: PosterMaker = poster or make_poster self._policy: Callable[[], "Optional[Policy]"] = policy or default_policy # ── 查询 ────────────────────────────────────────────── def list_jobs(self) -> list[WebJob]: """全部任务,新的在前。""" self._prune() return sorted(self._jobs.values(), key=lambda j: j.created_at, reverse=True) def get(self, job_id: str) -> Optional[WebJob]: return self._jobs.get(job_id) def _find_active(self, url: str) -> Optional[WebJob]: for job in self._jobs.values(): if job.url == url and job.active: return job return None # ── 提交与执行 ──────────────────────────────────────── async def submit(self, urls: Iterable[str], *, force: bool = False) -> list[WebJob]: """提交一批 URL;同 URL 已有在跑的任务时直接复用(force=True 强行新建)。 复用不只是省一次解析:抖音双开浏览器毫无意义,而且两个任务同时往 `unique_media_path` 的同一个路径写会撞车(它是 exists → 改名的写法)。 """ self._prune() jobs: list[WebJob] = [] for url in urls: job = None if force else self._find_active(url) if job is None: job = WebJob( id=uuid.uuid4().hex[:12], url=url, platform=platform_of(url) ) self._jobs[job.id] = job self._spawn(job) jobs.append(job) return jobs async def refresh(self, job_id: str) -> Optional[WebJob]: """对已落盘文件重跑上传,换一批新的预签名链接(旧的 1 小时过期)。 任务不存在 / 还没有文件时返回 None,由调用方给提示。 """ job = self._jobs.get(job_id) if job is None or not job.files: return None policy = self._policy() job.files = [await self._upload_one(path, policy) for path in job.paths] job.cover_url = await self._resolve_cover(job, policy) job.updated_at = time.time() return job # ── 清理 ────────────────────────────────────────────── def remove(self, job_id: str) -> bool: """删掉一条任务;还在跑的一并取消(Chrome 那头由 playwright 自己收尾)。""" job = self._jobs.pop(job_id, None) if job is None: return False task = self._tasks.pop(job_id, None) if task is not None and not task.done(): task.cancel() return True def clear_finished(self) -> int: """清掉所有终态任务(排队/运行中的不动),返回清掉几条。""" finished = [jid for jid, job in self._jobs.items() if not job.active] for job_id in finished: self._jobs.pop(job_id, None) return len(finished) def _spawn(self, job: WebJob) -> None: task = asyncio.create_task(self._run(job)) self._tasks[job.id] = task task.add_done_callback(lambda _t: self._tasks.pop(job.id, None)) async def _run(self, job: WebJob) -> None: """排队等闸门 → 执行;任何异常都落到任务状态里,不外抛。""" try: async with self._gate(): await asyncio.wait_for(self._execute(job), JOB_TIMEOUT) except asyncio.CancelledError: raise except asyncio.TimeoutError: # 3.11+ 就是内置 TimeoutError,wait_for 抛的 self._update(job, status="failed", error=f"解析超时(超过 {JOB_TIMEOUT}s)") except Exception as e: # noqa: BLE001 — 任务边界,失败即任务状态 logger.exception(f"Web 解析任务失败:{job.url}") self._update(job, status="failed", error=str(e) or type(e).__name__) async def _execute(self, job: WebJob) -> None: self._update(job, status="running", stage="解析中") policy = self._policy() title, files, _ = await self._fetch(job.url) job.title = title or job.url if not files: raise RuntimeError("无法解析到媒体(链接失效 / 风控 / cookies 过期)") job.paths = list(files) self._update(job, stage="上传中") entries: list[dict[str, Any]] = [] for path in files: entries.append(await self._upload_one(path, policy)) job.files = entries # 落一个刷一个,前端能看着进度出图 job.cover_url = await self._resolve_cover(job, policy) self._update(job, status="done", stage="") async def _resolve_cover(self, job: WebJob, policy: "Optional[Policy]") -> str: """任务封面:有图就用第一张图;全是视频就抽第一帧上传当封面。 抽帧/上传失败都只返回空串 —— 卡片那边退化成占位块,不影响任务本身。 """ for entry in job.files: if entry["kind"] == "image": url = entry["public_url"] or entry["local_url"] if url: return url if job.poster_path and job.poster_path.exists(): entry = await self._upload_one(job.poster_path, policy) return entry["public_url"] or entry["local_url"] for path in job.paths: if kind_of(path) != "video" or not path.exists(): continue poster = await self._poster(path) if poster is None: return "" job.poster_path = poster entry = await self._upload_one(poster, policy) return entry["public_url"] or entry["local_url"] return "" async def _upload_one( self, path: Path, policy: "Optional[Policy]" ) -> dict[str, Any]: """上传单个文件 → files 里的一项;失败只标这项。""" if not path.exists(): return _file_entry(path, error="本地文件不存在(可能已被 temp 清理)") try: local_url, public_url = await self._upload(path, policy) except Exception as e: # noqa: BLE001 — 单文件失败不拖累整条任务 logger.exception(f"Web 任务上传失败:{path}") return _file_entry(path, error=f"上传失败:{str(e)[:120]}") entry = _file_entry(path, local_url=local_url, public_url=public_url) if not entry["local_url"] and not entry["public_url"]: entry["error"] = "上传失败(S3 没返回链接)" return entry # ── 内部工具 ────────────────────────────────────────── def _gate(self) -> asyncio.Semaphore: """惰性建闸门:构造必须发生在跑着的事件循环里。""" if self._sem is None: self._sem = asyncio.Semaphore(MAX_CONCURRENCY) return self._sem def _update( self, job: WebJob, *, status: Optional[str] = None, stage: Optional[str] = None, error: Optional[str] = None, ) -> None: if status is not None: job.status = status if stage is not None: job.stage = stage if error is not None: job.error = error job.updated_at = time.time() def _prune(self) -> None: """先清超龄的终态任务,再按上限砍掉最旧的终态任务(在跑的不动)。""" now = time.time() for job_id, job in list(self._jobs.items()): if not job.active and now - job.updated_at > TERMINAL_TTL: self._jobs.pop(job_id, None) overflow = len(self._jobs) - MAX_JOBS if overflow <= 0: return finished = sorted( (j for j in self._jobs.values() if not j.active), key=lambda j: j.created_at, ) for job in finished[:overflow]: self._jobs.pop(job.id, None) def _file_entry( path: Path, *, local_url: str = "", public_url: Optional[str] = None, error: str = "", ) -> dict[str, Any]: try: size = path.stat().st_size except OSError: size = 0 return { "name": path.name, "kind": kind_of(path), "size": size, "local_url": local_url or "", "public_url": public_url or "", "error": error, } #: 全局单例(web_hub.py 的路由共用) JOBS = JobManager()