feat(video-analysis): 群策略 v3 / 群文件投递通道 / Web 管理页

- policy.py:per-group 正交策略(自动解析 / 自动策略 / 禁用策略 / 存储 A·B·C /
  公网 / 链接 / 群文件 + 平台限定),list.json v1/v2 → v3 自动迁移,
  写入统一走 PolicyStore(加锁 + .tmp 原子替换 + 字段归一)
- 群文件并行通道 group_file.py:打包 zip(可选 pyzipper AES-256)后优先走 S3 预签名、
  本地直传兜底;设了密码但 pyzipper 不可用就放弃上传,不退化成明文
- list_proc.py 收敛到「视频策略」统一入口,权限判定改走 policy
- Web 管理页 /hub/video_analysis(群策略 + 链接解析面板)与 services/web_jobs.py
  (只复用纯函数层,Web 上下文不发消息;内存任务表 + 并发闸门 + 超时)
- 媒体命名统一到 utils.py({作者}_{作者id}/{作品名}[_短码]),cleanup 回收空目录
- 测试:policy / 命名 / 群文件 / web_jobs 四组

顺带 pyproject 的 pytest 加 testpaths=tests(避免收进 debug/ 下的调试脚本)。

Co-Authored-By: Claude Code <noreply@anthropic.com>
This commit is contained in:
2026-09-22 14:23:32 +08:00
co-authored by Claude Code
parent 51b08ccb68
commit 4badcfcf32
29 changed files with 4838 additions and 670 deletions
@@ -14,15 +14,19 @@
"""
import re
import tempfile
from datetime import datetime
from pathlib import Path
from typing import Optional, Union
from nonebot import logger
from ...models import ContentFetchError
from ...utils import get_data_dir, parse_netscape_cookies, slugify
from ...utils import (
build_author_dir,
build_work_stem,
get_data_dir,
get_temp_root,
parse_netscape_cookies,
)
DATA_DIR = get_data_dir()
@@ -81,10 +85,10 @@ async def _parse(url: str):
# 4. 文章动态 → 转 opus
if await dynamic.is_article():
return await _parse_opus(dynamic.turn_to_opus(), "文章")
return await _parse_opus(dynamic.turn_to_opus(), "文章", url)
info = await dynamic.get_info()
return await _parse_dynamic_info(info)
return await _parse_dynamic_info(info, url)
async def _parse_article(read_id: int) -> tuple[str, Union[Path, list[Path]]]:
@@ -94,10 +98,12 @@ async def _parse_article(read_id: int) -> tuple[str, Union[Path, list[Path]]]:
# 文章接口对匿名请求风控更严(-509),必须带凭证
article = Article(read_id, _build_credential())
opus = await article.turn_to_opus()
return await _parse_opus(opus, "文章")
return await _parse_opus(opus, "文章", f"cv{read_id}")
async def _parse_opus(opus, kind: str) -> tuple[str, Union[Path, list[Path]]]:
async def _parse_opus(
opus, kind: str, source: str = ""
) -> tuple[str, Union[Path, list[Path]]]:
"""图文动态/专栏解析(opus 接口返回 dict,直接访问)"""
info = await opus.get_info()
item = info.get("item") or {}
@@ -107,10 +113,12 @@ async def _parse_opus(opus, kind: str) -> tuple[str, Union[Path, list[Path]]]:
images: list[str] = []
texts: list[str] = []
author = ""
author_id = ""
for module in item.get("modules") or []:
if module.get("module_type") == "MODULE_TYPE_AUTHOR":
author_info = module.get("module_author") or {}
author = author_info.get("name", "")
author_id = str(author_info.get("mid") or "")
elif module.get("module_type") == "MODULE_TYPE_CONTENT":
content = module.get("module_content") or {}
for para in content.get("paragraphs") or []:
@@ -127,16 +135,20 @@ async def _parse_opus(opus, kind: str) -> tuple[str, Union[Path, list[Path]]]:
if not images:
return text or f"B站{kind}", []
file_name = _build_file_name(author, text or f"B站{kind}", kind)
file_paths = await _download_images(images, file_name)
rel_stem = _build_rel_stem(author, author_id, text or f"B站{kind}", source)
file_paths = await _download_images(images, rel_stem)
return text, file_paths
async def _parse_dynamic_info(info: dict) -> tuple[str, Union[Path, list[Path]]]:
async def _parse_dynamic_info(
info: dict, source: str = ""
) -> tuple[str, Union[Path, list[Path]]]:
"""动态解析(图文 / 视频 / 纯文字)"""
item = info.get("item") or {}
modules = item.get("modules") or {}
author = ((modules.get("module_author") or {}).get("name")) or "B站用户"
module_author = modules.get("module_author") or {}
author = module_author.get("name") or "B站用户"
author_id = str(module_author.get("mid") or "")
module_dynamic = modules.get("module_dynamic") or {}
major = module_dynamic.get("major") or {}
major_type = major.get("type", "")
@@ -151,7 +163,13 @@ async def _parse_dynamic_info(info: dict) -> tuple[str, Union[Path, list[Path]]]
title = archive.get("title") or desc or "B站视频动态"
logger.info(f"B站视频动态: bvid={bvid} 标题={title[:40]}")
video_path = await download_video(f"https://www.bilibili.com/video/{bvid}")
# 动态数据里已有 up 的昵称/mid,传下去才能和同一位 up 的图文
# 落在同一个作者目录(否则要赌 yt-dlp 返回的 id 对得上)
video_path, _ = await download_video(
f"https://www.bilibili.com/video/{bvid}",
author=author,
author_id=author_id,
)
if video_path:
return title, video_path
raise ContentFetchError(f"视频动态下载失败: {bvid}")
@@ -173,8 +191,8 @@ async def _parse_dynamic_info(info: dict) -> tuple[str, Union[Path, list[Path]]]
images = [u for u in images if u]
if images:
file_name = _build_file_name(author, title, "动态")
file_paths = await _download_images(images, file_name)
rel_stem = _build_rel_stem(author, author_id, title, source)
file_paths = await _download_images(images, rel_stem)
logger.info(f"B站图文动态: 作者={author}, 标题={title[:40]}, 图片={len(images)} 张")
return title, file_paths
@@ -197,14 +215,17 @@ def _extract_text(nodes: list) -> str:
return "".join(parts)
def _build_file_name(author: str, title: str, kind: str) -> str:
"""构建文件名 stem: {作者}_{标题}_{类型}_{时间}"""
slug_author = slugify(author)
slug_title = slugify(title or "", max_length=15)
if not slug_title:
slug_title = datetime.now().strftime("%H%M%S")
time_suffix = datetime.now().strftime("%H%M%S")
return f"{slug_author}_{slug_title}_{kind}_{time_suffix}"
def _build_rel_stem(
author: str, author_id: str, title: str, source: str = ""
) -> str:
"""相对平台根的路径词干:`{作者}_{mid}/{作品名}`
昵称/mid 都拿不到时用 source(动态链接)当来源码(见 utils.build_author_dir)。
"""
return (
f"{build_author_dir(author, author_id, source=source)}"
f"/{build_work_stem(title)}"
)
def _build_credential():
@@ -227,20 +248,21 @@ def _build_credential():
)
async def _download_images(image_urls: list[str], file_name: str) -> list[Path]:
"""并发下载图片(复用抖音图文的下载流程)"""
import httpx
async def _download_images(image_urls: list[str], rel_stem: str) -> list[Path]:
"""并发下载图片(复用抖音图文的下载流程)
落 hexi/data/temp/bilibili(原先落系统 temp,cleanup 扫不到、永不清理)。
"""
from .douyin_api import _process_note_with_parsed
tmp_root = Path(tempfile.gettempdir()) / "bilibili"
tmp_root.mkdir(parents=True, exist_ok=True)
tmp_root = get_temp_root("bilibili")
headers = {
"Referer": BILI_REFERER,
"User-Agent": BILI_UA,
}
return await _process_note_with_parsed(
[[u] for u in image_urls], None, tmp_root, file_name, headers
[[u] for u in image_urls], None, tmp_root, rel_stem, headers
)
@@ -10,11 +10,9 @@ from nonebot import logger
from playwright.async_api import async_playwright
from ...models import DouyinFetchError
from ...utils import ensure_unique_path, get_temp_root
from ...utils import get_temp_root, unique_media_path
from .douyin_parser import (
ParsedDouyinContent,
extract_trailing_digits,
is_animated_note,
parse_animated_note_videos,
parse_douyin_response,
parse_note_images,
@@ -194,9 +192,9 @@ async def fetch_douyin_content(
if parsed.media_type == "视频":
content = await _process_video(
api_response, tmp_root, parsed.file_name, aweme_id, headers
api_response, tmp_root, parsed.rel_stem, aweme_id, headers
)
return parsed.file_name, content
return parsed.raw_title, content
elif parsed.media_type == "图片":
# 先解析 images 列表,区分纯动图和图文/图+视频
@@ -209,15 +207,15 @@ async def fetch_douyin_content(
images_urls,
video_url,
tmp_root,
parsed.file_name,
parsed.rel_stem,
headers,
)
else:
# 纯动图(所有项都是视频)
content = await _process_animated_note(
api_response, tmp_root, parsed.file_name, headers
api_response, tmp_root, parsed.rel_stem, headers
)
return parsed.file_name, content
return parsed.raw_title, content
return None, None
@@ -228,11 +226,14 @@ async def fetch_douyin_content(
async def _process_video(
api_response: dict,
tmp_root: Path,
file_name: str,
rel_stem: str,
aweme_id: str,
headers: Dict[str, str],
) -> Path:
"""处理视频内容,返回本地文件路径"""
"""处理视频内容,返回本地文件路径
rel_stem 是相对平台根的路径词干 `{作者目录}/{作品名}`(见 ParsedDouyinContent.rel_stem)。
"""
groups = parse_video_urls(api_response)
best_group = None
@@ -251,7 +252,7 @@ async def _process_video(
best = max(full, key=lambda x: x["br"])
logger.info(f"选择码率: {best['br']} - {best['url'][:60]}...")
output_path = ensure_unique_path(tmp_root / f"{file_name}.mp4")
output_path = unique_media_path(tmp_root / f"{rel_stem}.mp4")
async with httpx.AsyncClient(headers=headers) as client:
async with client.stream("GET", best["url"]) as resp:
resp.raise_for_status()
@@ -271,9 +272,10 @@ async def _process_video(
logger.info(f"选择视频码率: {video['br']}")
logger.info(f"选择音频码率: {audio['br']}")
video_path = tmp_root / f"{file_name}_v.mp4"
audio_path = tmp_root / f"{file_name}_a.mp4"
output_path = ensure_unique_path(tmp_root / f"{file_name}.mp4")
# 先定下产物名(顺带建好作者目录),分轨中间文件与产物同目录
output_path = unique_media_path(tmp_root / f"{rel_stem}.mp4")
video_path = output_path.parent / f"{output_path.stem}_v.mp4"
audio_path = output_path.parent / f"{output_path.stem}_a.mp4"
async with httpx.AsyncClient(headers=headers) as client:
logger.info("开始下载视频...")
@@ -291,7 +293,9 @@ async def _process_video(
f.write(chunk)
logger.info("合并视频和音频...")
merge_video_audio(video_path, audio_path, output_path)
# ffmpeg 是同步子进程,直接 await 不了:不丢线程池会卡住整个事件循环
# (合并期间 Web 轮询、群消息全都停摆)
await asyncio.to_thread(merge_video_audio, video_path, audio_path, output_path)
video_path.unlink()
audio_path.unlink()
@@ -306,11 +310,14 @@ async def _process_note_with_parsed(
images_urls: List[List[str]],
video_url: Optional[str],
tmp_root: Path,
file_name: str,
rel_stem: str,
headers: Dict[str, str],
) -> List[Path]:
"""根据已解析的图片/视频 URL 列表,并行下载"""
note_dir = ensure_unique_path(tmp_root / file_name)
"""根据已解析的图片/视频 URL 列表,并行下载
一个作品一个目录:`{平台根}/{作者目录}/{作品名}[_{短码}]/001.jpg…`
"""
note_dir = unique_media_path(tmp_root / rel_stem)
note_dir.mkdir(parents=True, exist_ok=True)
logger.info(f"图文保存目录: {note_dir}")
@@ -353,7 +360,7 @@ async def _process_note(
api_response: dict,
api_response_favorite: dict,
tmp_root: Path,
file_name: str,
rel_stem: str,
aweme_id: str,
headers: Dict[str, str],
) -> List[Path]:
@@ -367,7 +374,7 @@ async def _process_note(
raise DouyinFetchError("未找到图文链接")
return await _process_note_with_parsed(
images_urls, video_url, tmp_root, file_name, headers
images_urls, video_url, tmp_root, rel_stem, headers
)
@@ -377,14 +384,14 @@ async def _process_note(
async def _process_animated_note(
api_response: dict,
tmp_root: Path,
file_name: str,
rel_stem: str,
headers: Dict[str, str],
) -> List[Path]:
"""处理动图内容(media_type=42),并行下载所有无声 mp4 视频"""
video_urls = parse_animated_note_videos(api_response)
logger.info(f"解析到的动图视频链接: {video_urls}")
note_dir = ensure_unique_path(tmp_root / file_name)
note_dir = unique_media_path(tmp_root / rel_stem)
note_dir.mkdir(parents=True, exist_ok=True)
logger.info(f"动图保存目录: {note_dir}")
@@ -4,13 +4,12 @@ import json
import re
import urllib.parse
from dataclasses import dataclass
from datetime import datetime
from typing import Dict, List, Optional
from nonebot import logger
from ...models import DouyinFetchError
from ...utils import slugify
from ...utils import build_author_dir, build_work_stem
# ============================= 数据结构 =============================
@@ -22,7 +21,13 @@ class ParsedDouyinContent:
raw_title: str
raw_nickname: str
media_type: str # "视频" | "图片"
file_name: str # 构建好的文件名 stem
file_name: str # 作品名 stem(落盘文件名/多图子目录名)
author_dir: str = "" # 作者目录 `{昵称}_{uid}`(落盘与 S3 key 的首层)
@property
def rel_stem(self) -> str:
"""相对平台根目录的路径词干:`{作者目录}/{作品名}`"""
return f"{self.author_dir}/{self.file_name}" if self.author_dir else self.file_name
# ============================= URL 工具 =============================
@@ -109,6 +114,29 @@ def extract_author_nickname(api_response: dict) -> str:
return nickname
def extract_author_id(api_response: dict) -> str:
"""从 API 响应提取作者稳定 id,优先级: uid > unique_id(抖音号) > sec_uid
SSR 路径的 `aweme_list[0].author.uid` 与 API 路径的
`aweme_detail.author.uid` 走同一套取值;都拿不到返回空串,
作者目录退化成只用昵称(见 utils.build_author_dir)。
"""
aweme_detail = api_response.get("aweme_detail") or {}
author = aweme_detail.get("author") or {}
if not author:
aweme_list = api_response.get("aweme_list") or []
if aweme_list and isinstance(aweme_list, list):
author = (aweme_list[0] or {}).get("author") or {}
for key in ("uid", "unique_id", "sec_uid"):
value = str(author.get(key) or "").strip()
if value and value != "0":
logger.info(f"RAW作者id({key}):{value}")
return value
logger.info("未获取到作者 id,作者目录只用昵称")
return ""
def detect_media_type(referer_url: str | None) -> str | None:
"""根据页面 URL 检测媒体类型(视频/图片)"""
if referer_url is None:
@@ -326,36 +354,6 @@ def parse_ssr_page(html: str) -> Optional[dict]:
return {"aweme_list": [item]}
# ============================= 文件名构建 =============================
def build_file_name(
raw_title: str,
raw_nickname: str,
media_type: str,
) -> str:
"""
构建文件名 stem,格式: {作者}_{标题}_{类型}_{时间戳}
昵称不限长,标题最多 15 字符(slugify 后),末尾 HHMMSS 防覆盖。
"""
slug_nickname = slugify(raw_nickname)
if raw_title:
slug_title = slugify(raw_title, max_length=15)
else:
slug_title = ""
if not slug_title:
slug_title = datetime.now().strftime("%H%M%S")
logger.info(f"标题为空,使用短时间戳: {slug_title}")
slug_type = slugify(media_type)
time_suffix = datetime.now().strftime("%H%M%S")
return f"{slug_nickname}_{slug_title}_{slug_type}_{time_suffix}"
# ============================= 动图检测 =============================
@@ -410,11 +408,21 @@ def parse_douyin_response(
raw_title = extract_title_from_api(api_response)
raw_nickname = extract_author_nickname(api_response)
file_name = build_file_name(raw_title, raw_nickname, media_type)
raw_author_id = extract_author_id(api_response)
# 昵称/uid 都拿不到时的来源码:优先作品链接,其次响应里的作品 id
aweme_id = str(
(api_response.get("aweme_detail") or {}).get("aweme_id")
or ((api_response.get("aweme_list") or [{}])[0] or {}).get("aweme_id")
or ""
)
author_dir = build_author_dir(
raw_nickname, raw_author_id, source=referer_url or aweme_id
)
file_name = build_work_stem(raw_title)
logger.info(
f"内容标题: {raw_title}, 作者: {raw_nickname}, "
f"类型: {media_type}, 文件名: {file_name}"
f"内容标题: {raw_title}, 作者: {raw_nickname}({raw_author_id}), "
f"类型: {media_type}, 落盘路径: {author_dir}/{file_name}"
)
return ParsedDouyinContent(
@@ -422,4 +430,5 @@ def parse_douyin_response(
raw_nickname=raw_nickname,
media_type=media_type,
file_name=file_name,
author_dir=author_dir,
)
@@ -85,13 +85,13 @@ async def fetch_douyin_note_ssr(
images_urls, video_url = parse_note_images(api_response, None, vid)
if images_urls:
file_paths = await _process_note_with_parsed(
images_urls, video_url, tmp_root, parsed.file_name, headers
images_urls, video_url, tmp_root, parsed.rel_stem, headers
)
else:
file_paths = await _process_animated_note(
api_response, tmp_root, parsed.file_name, headers
api_response, tmp_root, parsed.rel_stem, headers
)
return parsed.file_name, file_paths
return parsed.raw_title, file_paths
async def _fetch_ssr_api_response(
@@ -19,16 +19,21 @@
import asyncio
import json
import re
import tempfile
import urllib.parse
from datetime import datetime
from pathlib import Path
from typing import Optional, Union
from nonebot import logger
from ...models import ContentFetchError
from ...utils import get_data_dir, get_temp_root, parse_netscape_cookies, slugify
from ...utils import (
build_author_dir,
build_work_stem,
get_data_dir,
get_temp_root,
parse_netscape_cookies,
unique_media_path,
)
DATA_DIR = get_data_dir()
@@ -214,6 +219,12 @@ async def _build_result(note: dict) -> tuple[str, Union[Path, list[Path]]]:
desc = note.get("desc") or ""
nickname = ((note.get("user") or {}).get("nickname")) or "小红书用户"
text = title or desc or "小红书笔记"
rel_stem = _build_rel_stem(
nickname,
_extract_author_id(note),
text,
source=str(note.get("noteId") or note.get("id") or ""),
)
# 1. 视频笔记 → 无水印原片优先
if note.get("type") == "video" and note.get("video"):
@@ -223,8 +234,7 @@ async def _build_result(note: dict) -> tuple[str, Union[Path, list[Path]]]:
if okey:
video_url = f"https://sns-video-bd.xhscdn.com/{okey}"
logger.info(f"小红书视频: 无水印原片 originVideoKey={okey[:30]}...")
file_name = _build_file_name(nickname, text, "视频")
video_path = await _download_video(video_url, file_name)
video_path = await _download_video(video_url, rel_stem)
return text, video_path
# 1b. 无 originVideoKey(国内站数据)→ 从 stream 分组选无水印原片
@@ -251,8 +261,7 @@ async def _build_result(note: dict) -> tuple[str, Union[Path, list[Path]]]:
f"{best.get('width')}x{best.get('height')} {best.get('fps')}fps "
f"size={best.get('size')} duration={duration}ms"
)
file_name = _build_file_name(nickname, text, "视频")
video_path = await _download_video(video_url, file_name)
video_path = await _download_video(video_url, rel_stem)
return text, video_path
raise ContentFetchError("小红书视频流解析失败")
@@ -266,25 +275,36 @@ async def _build_result(note: dict) -> tuple[str, Union[Path, list[Path]]]:
logger.info(f"小红书文字笔记: {text[:30]}")
return text, []
file_name = _build_file_name(nickname, text, "笔记")
file_paths = await _download_images(images, file_name)
file_paths = await _download_images(images, rel_stem)
logger.info(f"小红书图文笔记: 作者={nickname}, 图片={len(images)} 张")
return text, file_paths
def _build_file_name(nickname: str, title: str, kind: str) -> str:
"""构建文件名 stem: {作者}_{标题}_{类型}_{时间}"""
slug_nickname = slugify(nickname)
slug_title = slugify(title or "", max_length=15)
if not slug_title:
slug_title = datetime.now().strftime("%H%M%S")
time_suffix = datetime.now().strftime("%H%M%S")
return f"{slug_nickname}_{slug_title}_{kind}_{time_suffix}"
def _extract_author_id(note: dict) -> str:
"""小红书作者稳定 id(页面数据字段未实测,逐个兜底;取不到返回空串)"""
user = note.get("user") or {}
for key in ("userId", "user_id", "id"):
value = user.get(key)
if isinstance(value, (str, int)) and str(value).strip() not in ("", "0"):
return str(value).strip()
return ""
async def _download_images(image_urls: list[str], file_name: str) -> list[Path]:
def _build_rel_stem(
nickname: str, author_id: str, title: str, source: str = ""
) -> str:
"""相对平台根的路径词干:`{作者}_{userId}/{作品名}`
昵称/作者 id 都拿不到时用 source(笔记 id)当来源码(见 utils.build_author_dir)。
"""
return (
f"{build_author_dir(nickname, author_id, source=source)}"
f"/{build_work_stem(title)}"
)
async def _download_images(image_urls: list[str], rel_stem: str) -> list[Path]:
"""并发下载图片(复用抖音图文的下载流程)"""
import httpx
from .douyin_api import _process_note_with_parsed
@@ -294,11 +314,11 @@ async def _download_images(image_urls: list[str], file_name: str) -> list[Path]:
"User-Agent": REDNOTE_UA,
}
return await _process_note_with_parsed(
[[u] for u in image_urls], None, tmp_root, file_name, headers
[[u] for u in image_urls], None, tmp_root, rel_stem, headers
)
async def _download_video(video_url: str, file_name: str) -> Path:
async def _download_video(video_url: str, rel_stem: str) -> Path:
"""流式下载视频
注意:sns-video-bd(无水印原片)不带 Referer 或带 xiaohongshu.com
@@ -307,7 +327,7 @@ async def _download_video(video_url: str, file_name: str) -> Path:
import httpx
tmp_root = get_temp_root("xiaohongshu")
output_path = tmp_root / f"{file_name}.mp4"
output_path = unique_media_path(tmp_root / f"{rel_stem}.mp4")
headers = {"User-Agent": REDNOTE_UA}
async with httpx.AsyncClient(headers=headers, timeout=300) as client:
async with client.stream("GET", video_url) as resp:
@@ -3,9 +3,9 @@
import asyncio
import os
import re
import shutil
import sys
import tempfile
from datetime import datetime
from pathlib import Path
from typing import Optional
@@ -14,7 +14,14 @@ from nonebot import logger
from yt_dlp import YoutubeDL
from yt_dlp.utils import DownloadError
from ...utils import get_data_dir, get_temp_root, slugify, ensure_unique_path
from ...utils import (
build_author_dir,
build_work_stem,
get_data_dir,
get_temp_root,
slugify,
unique_media_path,
)
def detect_platform(url: str) -> str:
@@ -41,6 +48,16 @@ def extract_uploader(info: dict) -> Optional[str]:
)
def extract_uploader_id(info: dict) -> str:
"""作者稳定 id:channel_id > uploader_id(@handle / B站 mid),拿不到返回空串
不用 `id`(那是作品 id,会把同一作者的作品拆到不同目录)。
"""
if not info:
return ""
return str(info.get("channel_id") or info.get("uploader_id") or "").strip()
def get_ffmpeg_path() -> str:
scripts_dir = os.path.dirname(sys.executable)
ffmpeg_path = os.path.join(scripts_dir, "ffmpeg.exe")
@@ -78,19 +95,36 @@ async def _retry_download(
logger.error(
f"yt-dlp 重试 {max_retries} 次后仍失败: {str(e)[:120]}"
)
except Exception as e:
except Exception:
# 非 DownloadError(如 OSError)不重试,直接抛出
raise
raise last_error # type: ignore[misc]
async def download_video(url: str) -> Optional[Path]:
"""下载视频,支持直链和 yt-dlp"""
async def download_video(
url: str,
*,
author: Optional[str] = None,
author_id: Optional[str] = None,
) -> tuple[Optional[Path], str]:
"""下载视频,支持直链和 yt-dlp
author / author_id 可由调用方覆盖(如 B站 视频动态已从动态数据里拿到 mid,
传进来才能和同一位 up 的图文落在同一个作者目录)。
Returns:
(本地文件, 作品标题) — 失败时 (None, "")。
落盘位置:`temp/{平台}/{作者}_{作者id}/{作品名}[_{短码}].ext`
(直链拿不到作者信息,统一进 `未知作者/`)。
"""
platform = detect_platform(url)
temp_root = get_temp_root(platform)
# ---------- 1. 直链探测 ----------
# 不含 m3u8:HLS 播放列表直下只会得到一个文本文件,交给 yt-dlp 处理
direct_media_ext = re.search(
r"\.(mp4|m3u8|ts|webm|mov|flv)(?:$|\?)", url, re.IGNORECASE
r"\.(mp4|ts|webm|mov|flv)(?:$|\?)", url, re.IGNORECASE
)
is_direct = bool(direct_media_ext)
@@ -105,39 +139,35 @@ async def download_video(url: str) -> Optional[Path]:
is_direct = False
if is_direct:
temp_dir = tempfile.mkdtemp(prefix="direct_ytcache_", dir=get_temp_root("ytcache"))
ext = "mp4"
m = re.search(r"\.([a-zA-Z0-9]{2,5})(?:$|\?)", url)
if m and len(m.group(1)) <= 5:
ext = m.group(1)
url_stem = Path(url.split("?")[0]).stem or "video"
slug_stem = slugify(url_stem, max_length=15)
if not slug_stem:
slug_stem = datetime.now().strftime("%H%M%S")
time_suffix = datetime.now().strftime("%H%M%S")
new_name = f"{slug_stem}_视频_{time_suffix}.{ext}"
filename = os.path.join(temp_dir, new_name)
slug_stem = slugify(url_stem, max_length=15) or "视频"
# 直链拿不到作者信息 → `未知作者_{来源短码}`(同一链接稳定、不同链接不撞)
author_dir = build_author_dir(author, author_id, source=url)
final_path = unique_media_path(temp_root / author_dir / f"{slug_stem}.{ext}")
try:
async with AsyncClient(follow_redirects=True, timeout=300) as client:
async with client.stream("GET", url) as resp:
resp.raise_for_status()
with open(filename, "wb") as fh:
with open(final_path, "wb") as fh:
async for chunk in resp.aiter_bytes(chunk_size=8192):
fh.write(chunk)
final_path = ensure_unique_path(Path(filename))
logger.info(f"直接下载完成: {final_path}")
return final_path
return final_path, ""
except Exception:
logger.exception("直接下载失败,回退 yt-dlp")
if os.path.exists(filename):
os.remove(filename)
final_path.unlink(missing_ok=True)
# ---------- 2. yt-dlp 下载 ----------
platform = detect_platform(url)
temp_dir = tempfile.mkdtemp(prefix="ytcache_", dir=get_temp_root("ytcache"))
# 先下到 scratch 目录(outtmpl 必须在拿到 info 之前给定),拿到 info 后再
# 按作者归位到 temp/{平台}/{作者}_{作者id}/
temp_dir = tempfile.mkdtemp(prefix="_dl_", dir=temp_root)
output_path = os.path.join(temp_dir, "%(title).80s.%(ext)s")
base_opts = {
@@ -188,7 +218,10 @@ async def download_video(url: str) -> Optional[Path]:
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8",
"Accept-Language": "zh-CN,zh;q=0.9,en;q=0.8",
}
base_opts["extractor_args"] = {"twitter": {"api": ["syndication"]}}
# 不要强制 twitter:api=syndication:该端点为未登录视角,对敏感/受限推文
# 只返回 tombstone(无 mediaDetails),且会无条件覆盖已登录 GraphQL 的结果,
# 表现为 "No video could be found in this tweet"。默认走 GraphQL + cookies,
# 遇 429 yt-dlp 会自行回退 syndication。
elif platform == "youtube":
base_opts["http_headers"] = {
"User-Agent": ua,
@@ -206,54 +239,53 @@ async def download_video(url: str) -> Optional[Path]:
info = await _retry_download(loop, url, base_opts)
except Exception:
logger.exception("yt-dlp 下载失败")
# YouTube: cookies 可能触发 bot 检测导致只返回图片无视频格式
# 回退无 cookie 模式重试
if platform == "youtube" and "cookiefile" in base_opts:
# YouTube: cookies 可能触发 bot 检测导致只返回图片无视频格式
# 回退无 cookie 模式重试
logger.info("YouTube 回退无 cookies 模式重试...")
base_opts.pop("cookiefile", None)
base_opts.pop("http_headers", None)
# 清理失败残留
for f in Path(temp_dir).glob("*.*"):
try:
f.unlink()
except Exception:
pass
try:
info = await _retry_download(loop, url, base_opts, max_retries=2)
except Exception:
logger.exception("yt-dlp 无 cookies 重试也失败")
return None
elif platform == "twitter":
# X 登录态失效(auth_token 过期)时 GraphQL 会直接拒绝请求;
# 退回未登录的 syndication 端点,公开推文仍可下载(敏感推文会失败)
logger.info("Twitter 回退 syndication 端点重试...")
base_opts["extractor_args"] = {"twitter": {"api": ["syndication"]}}
else:
return None
return None, ""
# 清理失败残留
for f in Path(temp_dir).glob("*.*"):
try:
f.unlink()
except Exception:
pass
try:
info = await _retry_download(loop, url, base_opts, max_retries=2)
except Exception:
logger.exception("yt-dlp 回退重试也失败")
return None, ""
files = list(Path(temp_dir).glob("*.*"))
if not files:
return None
return None, ""
original_file = files[0]
# 构建新文件名
uploader = extract_uploader(info or {})
# 归位到作者目录:{作者}_{作者id}/{作品名}[_{短码}].ext
uploader = author or extract_uploader(info or {})
uploader_id = author_id or extract_uploader_id(info or {})
title = ((info or {}).get("title") or "").strip()
slug_title = slugify(title, max_length=15) if title else ""
if not slug_title:
slug_title = datetime.now().strftime("%H%M%S")
time_suffix = datetime.now().strftime("%H%M%S")
if uploader:
slug_uploader = slugify(str(uploader))
new_stem = f"{slug_uploader}_{slug_title}_视频_{time_suffix}"
else:
new_stem = f"{slug_title}_视频_{time_suffix}"
new_path = ensure_unique_path(
original_file.with_name(f"{new_stem}{original_file.suffix}")
author_dir = build_author_dir(uploader, uploader_id, source=url)
new_path = unique_media_path(
temp_root / author_dir / f"{build_work_stem(title)}{original_file.suffix}"
)
original_file.rename(new_path)
shutil.move(str(original_file), str(new_path))
# 下载用的 scratch 目录已空,顺手收掉(cleanup 不删目录)
shutil.rmtree(temp_dir, ignore_errors=True)
logger.info(
f"yt-dlp 下载完成, 标题: {title}, "
f"作者: {uploader}, 重命名: {new_path}"
f"作者: {uploader}({uploader_id}), 落盘: {new_path}"
)
return new_path
return new_path, title
@@ -0,0 +1,271 @@
"""群文件上传 —— 与消息发送并行的第二条投递通道。
策略开关见 `policy.Policy.upload_group_file`:开着的时候,媒体照常发到群里,
同时另传一份到群文件(群友可随时下载、不占聊天记录)。
投递前会按全局配置打包(`video_analysis_group_file_zip`)成一个 zip:
配了解压密码(`video_analysis_group_file_password`)就用 AES-256 加密,
**密码设置了但 pyzipper 不可用时直接放弃上传,绝不退化成传明文**。
(无密码的普通 zip 走标准库,零依赖。)
失败只记日志(bot 没有群文件权限、超出群文件大小上限等都属于预期内的失败),
不影响消息发送链路;上传成功后不额外发消息,避免刷屏。
"""
from __future__ import annotations
import asyncio
import re
import zipfile
from collections.abc import Iterable
from datetime import datetime
from pathlib import Path
from typing import TYPE_CHECKING, Any
from nonebot import get_bot, logger
if TYPE_CHECKING: # 仅类型检查:运行期不导入,单测可裸加载本模块
from ..policy import Policy
try: # 缺失时只影响「加密打包」这一路,见 build_archive
import pyzipper
except ImportError: # pragma: no cover
pyzipper = None # type: ignore[assignment]
def encryption_available() -> bool:
"""加密打包是否可用(pyzipper 已安装)。"""
return pyzipper is not None
def build_archive(
files: Iterable[Path | str],
title: str = "",
*,
password: str = "",
out_dir: Path | None = None,
rel_dir: str = "",
) -> Path:
"""把文件打包成一个 zip,返回产物路径;password 非空则 AES-256 加密。
out_dir 缺省落 `hexi/data/temp/archive/{作者目录}`(rel_dir 由调用方传入,
与源媒体同一套目录结构,见 sender.media_rel_dir_of),由 cleanup 按天清理。
加密需要 pyzipper,缺失时抛 RuntimeError —— 调用方应当**放弃上传**。
"""
paths = [Path(f) for f in files]
if not paths:
raise ValueError("没有可打包的文件")
missing = [p for p in paths if not p.exists()]
if missing:
raise FileNotFoundError(f"待打包文件不存在: {missing[0]}")
if password and pyzipper is None:
raise RuntimeError("配置了解压密码,但 pyzipper 未安装,无法加密打包")
archive = _resolve_target(out_dir, title, rel_dir)
if password:
with pyzipper.AESZipFile(
archive,
"w",
compression=pyzipper.ZIP_DEFLATED,
encryption=pyzipper.WZ_AES,
) as zf:
zf.setpassword(password.encode("utf-8"))
_write_members(zf, paths)
else:
with zipfile.ZipFile(archive, "w", zipfile.ZIP_DEFLATED) as zf:
_write_members(zf, paths)
logger.info(
f"群文件打包完成: {archive.name}({len(paths)} 个文件"
f"{',已加密' if password else ''})"
)
return archive
def _default_archive_dir() -> Path:
"""默认落 temp/archive(延迟导入 utils:单测裸加载本模块时没有包上下文)。"""
from ...utils import get_temp_root
return get_temp_root("archive")
def _safe_stem(title: str, max_length: int = 15) -> str:
"""标题 → 安全的文件名片段。
原始标题(抖音文案、YouTube 标题)可能带 `/`、换行、控制字符与 `#话题`,
这些进不了文件名,统一换成 `_` 后截断。
"""
cleaned = re.sub(r"#\S+", "", title)
cleaned = re.sub(r'[\\/:*?"<>|\s]+', "_", cleaned.strip())
cleaned = re.sub(r"_+", "_", cleaned).strip("_")
return cleaned[:max_length].strip("_")
def _resolve_target(out_dir: Path | None, title: str, rel_dir: str = "") -> Path:
"""产物最终路径:默认 temp/archive/{rel_dir},重名自动加序号。
空标题 / 占位符 / 已带「群文件」前缀的一律回退成默认名。
"""
base = Path(out_dir) if out_dir is not None else _default_archive_dir()
if rel_dir:
base = base / rel_dir
base.mkdir(parents=True, exist_ok=True)
stamp = f"{datetime.now():%H%M%S}"
stem = _safe_stem(title) if title else ""
if not stem or stem in ("title", "video") or stem.startswith("群文件"):
return _unique_path(base / f"群文件_{stamp}.zip")
return _unique_path(base / f"{stem}_群文件_{stamp}.zip")
def _unique_path(path: Path) -> Path:
"""重名追加 _2/_3…(同 utils.ensure_unique_path 的语义,就地实现)。"""
if not path.exists():
return path
index = 2
while True:
candidate = path.with_name(f"{path.stem}_{index}{path.suffix}")
if not candidate.exists():
return candidate
index += 1
def _write_members(zf: Any, paths: list[Path]) -> None:
"""写入成员:同名文件自动加序号,避免在包里互相覆盖。"""
used: set[str] = set()
for path in paths:
name = path.name
if name in used:
index = 2
while f"{path.stem}_{index}{path.suffix}" in used:
index += 1
name = f"{path.stem}_{index}{path.suffix}"
used.add(name)
zf.write(path, arcname=name)
def _file_uri(path: Path) -> str:
"""本地文件 → `file:///D:/a/b.zip`(正斜杠,中文/空格不转义)。
与图片/视频消息发给 NapCat 的形式一致;直接传 Windows 反斜杠路径会被
它的 realpath 判成 ENOENT(见 upload_group_file 的降级说明)。
"""
return "file:///" + str(path.resolve()).replace("\\", "/")
async def _call_upload(
file_value: str, group_id: int, name: str, folder_id: str | None = None
) -> None:
params: dict = {"group_id": group_id, "file": file_value, "name": name}
if folder_id:
params["folder"] = folder_id
await get_bot().call_api("upload_group_file", **params)
def _s3_url(path: Path, policy: "Policy | None") -> str:
"""把文件传到局域网 S3 换预签名链接(延迟导入 s3,便于单测裸加载本模块)。"""
from .s3 import upload_with_plan
url, _ = upload_with_plan(path, policy=policy)
return url
async def upload_group_file(
file_path: Path | str,
group_id: int,
*,
folder_id: str | None = None,
policy: Policy | None = None,
name: str | None = None,
) -> bool:
"""上传单个文件到群文件,返回是否成功。
两级投递,**S3 链接优先、本地直传兜底**:
本环境实测 NapCat 读不到 bot 进程写的本地文件(`realpath ... ENOENT`,
媒体消息的本地直发 55 次全失败、换 S3 链接后次次成功),所以直传只留作
S3 不可用时的兜底。policy 决定 S3 落到哪个桶(plan)。
本地直传用 `file:///D:/…` 形式(正斜杠、不转义中文)——图片/视频消息就是
这么发的,同机部署时能work。
name 可覆盖群文件列表里显示的名字(多图作品的成员是 001.jpg,调用方会补上
作品名前缀,免得多张图在群文件里全叫 001.jpg)。
"""
path = Path(file_path)
if not path.exists():
logger.warning(f"群文件上传:文件不存在 {path}")
return False
display_name = name or path.name
# 1) 局域网 S3 预签名链接(boto3 同步,丢线程池)
url = await asyncio.to_thread(_s3_url, path, policy)
if url:
try:
await _call_upload(url, group_id, display_name, folder_id)
except Exception as e:
logger.warning(
f"群文件 S3 链接上传失败(群 {group_id} / {display_name}): "
f"{str(e)[:160]};改用本地路径重试"
)
else:
logger.info(f"群文件上传成功(群 {group_id},S3 链接): {display_name}")
return True
else:
logger.warning("群文件:S3 中转没拿到链接,改用本地路径直传")
# 2) 兜底:file:/// 直传本地文件
try:
await _call_upload(_file_uri(path), group_id, display_name, folder_id)
except Exception as e:
logger.warning(
f"群文件上传失败(群 {group_id} / {display_name}): {str(e)[:160]}"
)
return False
logger.info(f"群文件上传成功(群 {group_id},本地直传): {display_name}")
return True
async def upload_group_files(
files: Iterable[Path | str],
group_id: int,
*,
title: str = "",
zip_files: bool = True,
password: str = "",
policy: Policy | None = None,
rel_dir: str = "",
) -> bool:
"""群文件投递入口:按配置打包(可加密)后传一个包,或逐个传原文件。
rel_dir 是源媒体所在的 `{作者}_{作者id}[/{作品名}]` 子目录,打包产物落到
`temp/archive/{rel_dir}`(由调用方从文件路径推出来,见 sender)。
"""
paths = [Path(f) for f in files]
if not paths:
return False
if not zip_files:
# 多图作品逐个传时,成员名是 001.jpg…,补上作品名前缀便于在群文件里辨认
work = Path(rel_dir).name if rel_dir else ""
prefix = f"{work}_" if work and len(paths) > 1 else ""
results = [
await upload_group_file(
p, group_id, policy=policy, name=f"{prefix}{p.name}"
)
for p in paths
]
return any(results)
try:
archive = await asyncio.to_thread(
build_archive, paths, title, password=password, rel_dir=rel_dir
)
except Exception as e:
logger.error(f"群文件打包失败,跳过本次群文件上传: {e}")
return False
return await upload_group_file(archive, group_id, policy=policy)
@@ -1,7 +1,6 @@
"""统一 S3 存储模块 — 合并本地局域网 S3 和公网 MinIO"""
from pathlib import Path
from time import strftime, localtime
import boto3
from botocore.config import Config
@@ -9,6 +8,9 @@ from nonebot import logger
from hexi.web_hub.web_config import get_effective_value
from ...policy import Policy
from ...utils import media_key_of
# 该插件模块名,用于读取统一配置值库(Web 修改后生效)
_PLUGIN_ID = "hexi.plugins.nonebot_plugin_video_analysis"
@@ -171,11 +173,13 @@ def upload_to_public_s3(file_path: str | Path) -> tuple[str, str]:
"""
上传到公网 MinIO,返回 (公网URL, file_key)
私聊场景使用,生成可公网访问的链接
私聊场景使用,生成可公网访问的链接。
key 按本地落盘结构推导:`{作者}_{作者id}/{作品名}[_{短码}].ext`
(不再用 `{年月日}/` 前缀,见 utils.media_key_of)。
"""
file_path = Path(file_path)
current_day = strftime("%Y-%m-%d", localtime())
file_key = f"{current_day}/{file_path.name}"
file_key = media_key_of(file_path)
try:
url = _get_public_s3().upload_public(str(file_path), file_key)
return url, file_key
@@ -184,22 +188,19 @@ def upload_to_public_s3(file_path: str | Path) -> tuple[str, str]:
return "", ""
def upload_to_local_s3(
title: str, image_post: bool, file_path: str | Path, plan: str | None = None
) -> str:
def upload_to_local_s3(file_path: str | Path, plan: str | None = None) -> str:
"""
上传到局域网 S3,返回预签名 URL
plan="A" → PLANA 桶
plan="B" → PLANB 桶
None → PLANC 桶(默认)
key 与公网一致:`{作者}_{作者id}/{作品名}[_{短码}].ext`
(不再有 `{年月日}/` 前缀,也不再重复套一层标题目录)。
"""
file_path = Path(file_path)
current_day = strftime("%Y-%m-%d", localtime())
if image_post:
file_key = f"{current_day}/{title}/{file_path.name}"
else:
file_key = f"{current_day}/{file_path.name}"
file_key = media_key_of(file_path)
if plan == "A":
client = _get_local_s3_plana()
@@ -224,35 +225,24 @@ def delete_from_public_s3(file_key: str) -> bool:
def upload_with_plan(
file_path: str | Path,
*,
plan: str | None = None,
is_private: bool = False,
title: str = "",
image_post: bool = False,
policy: Policy | None = None,
) -> tuple[str, str | None]:
"""
统一上传入口:按 plan 路由到对应的本地桶,并按需上传公网
统一上传入口:按策略路由本地桶,并按需上传公网
plan="A" → PLANA(仅本地)
plan="B" → PLANB + 公网
私聊 → PLANB + 公网
默认 → PLANC(仅本地)
policy.plan → PLANA / PLANB / PLANC(默认 C)
policy.upload_public → 是否额外上传公网(决定能否发下载链接)
key 由文件路径推导(utils.media_key_of),调用方不再传标题。
Returns:
(local_url, public_url_or_none)
"""
# 本地上传
if plan == "A":
local_url = upload_to_local_s3(title, image_post, file_path, plan="A")
elif plan == "B":
local_url = upload_to_local_s3(title, image_post, file_path, plan="B")
elif is_private:
local_url = upload_to_local_s3(title, image_post, file_path, plan="B")
else:
local_url = upload_to_local_s3(title, image_post, file_path) # PLANC
policy = policy or Policy()
local_url = upload_to_local_s3(file_path, plan=policy.plan)
# 公网上传(仅 PLANB 或私聊)
public_url = None
if plan == "B" or is_private:
if policy.upload_public:
public_url, _ = upload_to_public_s3(file_path)
return local_url, public_url
@@ -0,0 +1,508 @@
"""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()