Files
HeXi/tests/test_web_jobs.py
sansenhoshiandClaude Code 4badcfcf32 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>
2026-09-22 14:23:32 +08:00

572 lines
19 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Web 管理台解析任务 (services/web_jobs.py) 单元测试
模块本身只依赖标准库 + nonebot.logger,故用 importlib 按文件路径裸加载
(与 test_video_policy.py / test_video_group_file.py 同法)。fetchers 与
上传都是函数内延迟导入的注入点,这里全部塞假实现 —— 不碰网络、不碰 S3。
`platform_of` 走 `..policy` 的包内相对导入,裸加载下不可用,测试里一并替换。
"""
import asyncio
import importlib.util
import subprocess
import sys
import time
from pathlib import Path
import pytest
_MODULE_PATH = (
Path(__file__).resolve().parents[1]
/ "hexi"
/ "plugins"
/ "nonebot_plugin_video_analysis"
/ "services"
/ "web_jobs.py"
)
@pytest.fixture(scope="module")
def wj():
"""以独立模块名加载 web_jobs.py,避免与插件包 __init__ 冲突"""
spec = importlib.util.spec_from_file_location(
"video_web_jobs_under_test", _MODULE_PATH
)
module = importlib.util.module_from_spec(spec)
sys.modules[spec.name] = module
spec.loader.exec_module(module)
yield module
sys.modules.pop(spec.name, None)
@pytest.fixture
def manager(wj, monkeypatch):
"""造一个依赖全注入的 JobManager(fetch / upload / poster / policy 都传进来)"""
monkeypatch.setattr(wj, "platform_of", lambda url: "抖音")
def build(*, fetch, upload=None, poster=None, policy=None):
return wj.JobManager(
fetch=fetch,
upload=upload or _ok_upload,
poster=poster or _fake_poster,
policy=policy or (lambda: None),
)
return build
async def _fake_poster(video):
"""假抽帧:真的写一个封面文件出来(上传那一环会 stat 它)"""
out = video.with_name(f"{video.stem}_封面.jpg")
out.write_bytes(b"fake-jpeg")
return out
async def _no_poster(video):
return None
# ─────────────────────────── 假实现 ───────────────────────────
async def _ok_upload(path, policy):
return f"http://lan/{path.name}", None
def _touch(directory: Path, name: str, content: bytes = b"x") -> Path:
directory.mkdir(parents=True, exist_ok=True)
path = directory / name
path.write_bytes(content)
return path
def _ok_fetch(*paths, title="作品"):
async def fetch(url):
return title, [Path(p) for p in paths], len(paths) > 1
return fetch
async def _settle(mgr, timeout=5.0):
"""等所有任务进入终态,返回 list_jobs() 的结果。"""
deadline = time.monotonic() + timeout
while True:
jobs = mgr.list_jobs()
if jobs and all(not j.active for j in jobs):
return jobs
if time.monotonic() > deadline:
raise AssertionError("任务没有在超时内结束")
await asyncio.sleep(0.01)
# ─────────────────────── URL 提取 / 文件类型 ───────────────────────
def test_extract_urls_from_multiline_text(wj):
text = (
"看看这个 https://v.douyin.com/abc123。\n"
"还有 https://b23.tv/xyz,https://v.douyin.com/abc123(重复的只留一条)"
)
assert wj.extract_urls(text) == [
"https://v.douyin.com/abc123",
"https://b23.tv/xyz",
]
def test_extract_urls_keeps_query_string(wj):
"""小红书 xsec_token / 抖音带参数的分享链,query 不能被截掉"""
url = "https://www.xiaohongshu.com/explore/abc?xsec_token=XYZ&xsec_source=pc_feed"
assert wj.extract_urls(f"看这个 {url}") == [url]
def test_extract_urls_without_links(wj):
assert wj.extract_urls("这里没有链接,纯文本一段") == []
assert wj.extract_urls("") == []
def test_kind_of(wj):
assert wj.kind_of("001.JPG") == "image"
assert wj.kind_of(Path("作品.mp4")) == "video"
assert wj.kind_of("note.txt") == "file"
# ───────────────────────── 提交与复用 ─────────────────────────
async def test_same_url_reuses_active_job(wj, manager, tmp_path):
"""同 URL 连点两次不能起两个任务(抖音双开浏览器 + 落盘路径撞车)"""
calls: list[str] = []
async def fetch(url):
calls.append(url)
await asyncio.sleep(0.05) # 保持 running,让第二次提交能命中
return "作品", [_touch(tmp_path, "a.jpg")], False
mgr = manager(fetch=fetch)
url = "https://v.douyin.com/abc123"
first = (await mgr.submit([url]))[0]
second = (await mgr.submit([url]))[0]
assert first.id == second.id
assert first.platform == "抖音"
await _settle(mgr)
assert calls == [url]
# 跑完之后再提交同一条 → 是新任务(复用只针对在跑的)
third = (await mgr.submit([url]))[0]
assert third.id != first.id
async def test_force_creates_new_job(wj, manager, tmp_path):
calls: list[str] = []
async def fetch(url):
calls.append(url)
await asyncio.sleep(0.05)
return "作品", [_touch(tmp_path, "a.jpg")], False
mgr = manager(fetch=fetch)
url = "https://v.douyin.com/abc123"
old = (await mgr.submit([url]))[0]
new = (await mgr.submit([url], force=True))[0]
assert new.id != old.id
jobs = await _settle(mgr)
assert len(calls) == 2 and len(jobs) == 2
async def test_batch_submit_keeps_input_order(wj, manager, tmp_path):
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "a.jpg")))
jobs = await mgr.submit(["https://a/1", "https://a/2", "https://a/3"])
assert [j.url for j in jobs] == ["https://a/1", "https://a/2", "https://a/3"]
await _settle(mgr)
# ───────────────────────── 失败与超时 ─────────────────────────
async def test_no_media_marks_job_failed(wj, manager):
async def fetch(url):
return None, [], False
mgr = manager(fetch=fetch)
url = "https://v.douyin.com/gone"
(await mgr.submit([url]))[0]
job = (await _settle(mgr))[0]
assert job.status == "failed"
assert "无法解析到媒体" in job.error
assert job.title == url # 没标题时退回 URL,前端不至于只显示空白
async def test_fetch_exception_marks_job_failed(wj, manager):
async def fetch(url):
raise RuntimeError("被风控了")
mgr = manager(fetch=fetch)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.status == "failed" and job.error == "被风控了"
async def test_job_timeout(wj, manager, monkeypatch):
monkeypatch.setattr(wj, "JOB_TIMEOUT", 0.05)
async def fetch(url):
await asyncio.sleep(5)
mgr = manager(fetch=fetch)
(await mgr.submit(["https://v.douyin.com/slow"]))[0]
job = (await _settle(mgr))[0]
assert job.status == "failed" and "解析超时" in job.error
# ───────────────────────── 上传 ─────────────────────────
async def test_single_upload_failure_does_not_fail_job(wj, manager, tmp_path):
"""多图作品挂一张 → 只标那一张,其余照常出链接"""
files = [_touch(tmp_path, "001.jpg"), _touch(tmp_path, "002.jpg")]
async def upload(path, policy):
if path.name == "002.jpg":
raise RuntimeError("S3 连接超时")
return "http://lan/001.jpg?sig=1", "https://pub/001.jpg"
mgr = manager(fetch=_ok_fetch(*files), upload=upload)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.status == "done"
assert job.files[0]["local_url"].endswith("sig=1")
assert job.files[0]["public_url"] == "https://pub/001.jpg"
assert job.files[0]["error"] == ""
assert "S3 连接超时" in job.files[1]["error"]
assert job.files[1]["kind"] == "image" and job.files[1]["size"] == 1
async def test_upload_without_links_is_marked(wj, manager, tmp_path):
async def upload(path, policy):
return "", None # s3.py 失败时就是返回空串
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "a.mp4")), upload=upload)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.status == "done" and "上传失败" in job.files[0]["error"]
assert job.files[0]["kind"] == "video"
async def test_upload_gets_default_policy(wj, manager, tmp_path):
"""存储链路按默认策略走:policy 原样传给 upload_with_plan"""
seen: list = []
sentinel = object()
async def upload(path, policy):
seen.append(policy)
return "http://lan/a.jpg", None
mgr = manager(
fetch=_ok_fetch(_touch(tmp_path, "a.jpg")),
upload=upload,
policy=lambda: sentinel,
)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
await _settle(mgr)
assert seen == [sentinel]
# ───────────────────────── 刷新链接 ─────────────────────────
async def test_refresh_reruns_upload(wj, manager, tmp_path):
files = [_touch(tmp_path, "001.jpg"), _touch(tmp_path, "002.jpg")]
links = iter(
[
"http://lan/a?sig=1",
"http://lan/b?sig=1",
"http://lan/a?sig=2",
"http://lan/b?sig=2",
]
)
async def upload(path, policy):
return next(links), None
mgr = manager(fetch=_ok_fetch(*files), upload=upload)
job = (await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert [f["local_url"] for f in job.files] == [
"http://lan/a?sig=1",
"http://lan/b?sig=1",
]
refreshed = await mgr.refresh(job.id)
assert refreshed is not None
assert [f["local_url"] for f in refreshed.files] == [
"http://lan/a?sig=2",
"http://lan/b?sig=2",
]
async def test_refresh_marks_files_gone_from_disk(wj, manager, tmp_path):
"""temp 被 cleanup 清掉后再点刷新 → 只标文件没了,不炸"""
path = _touch(tmp_path, "001.jpg")
async def upload(p, policy):
return "http://lan/a.jpg", None
mgr = manager(fetch=_ok_fetch(path), upload=upload)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
path.unlink()
refreshed = await mgr.refresh(job.id)
assert refreshed is not None
assert "本地文件不存在" in refreshed.files[0]["error"]
assert refreshed.files[0]["local_url"] == ""
async def test_refresh_unknown_job_returns_none(wj, manager):
mgr = manager(fetch=_ok_fetch())
assert await mgr.refresh("nope") is None
async def fetch(url):
return None, [], False
(await mgr.submit(["https://a/nothing"]))[0]
job = (await _settle(mgr))[0]
assert job.files == []
assert await mgr.refresh(job.id) is None
# ───────────────────────── 封面 ─────────────────────────
async def test_cover_is_first_image(wj, manager, tmp_path):
"""图文作品:封面直接取第一张图,不抽帧"""
poster_calls: list = []
async def poster(video):
poster_calls.append(video)
return await _fake_poster(video)
files = [_touch(tmp_path, "001.jpg"), _touch(tmp_path, "002.jpg")]
mgr = manager(fetch=_ok_fetch(*files), poster=poster)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.cover_url == "http://lan/001.jpg"
assert poster_calls == []
async def test_cover_is_extracted_for_video_only_post(wj, manager, tmp_path):
"""纯视频作品:抽第一帧、上传、当封面"""
uploaded: list[str] = []
async def upload(path, policy):
uploaded.append(path.name)
return f"http://lan/{path.name}", None
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "作品.mp4")), upload=upload)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.cover_url == "http://lan/作品_封面.jpg"
assert uploaded == ["作品.mp4", "作品_封面.jpg"]
assert job.poster_path is not None and job.poster_path.exists()
# 封面是单独的文件,不能混进 files(否则播放器里会多出一个"文件")
assert [f["name"] for f in job.files] == ["作品.mp4"]
async def test_cover_prefers_image_over_video_in_mixed_post(wj, manager, tmp_path):
"""混排(图 + 视频):按用户定的规则取第一张图"""
files = [_touch(tmp_path, "图文.jpg"), _touch(tmp_path, "动图.mp4")]
mgr = manager(fetch=_ok_fetch(*files))
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.cover_url == "http://lan/图文.jpg"
assert job.poster_path is None
async def test_poster_failure_keeps_job_done(wj, manager, tmp_path):
"""抽帧失败 → 封面留空,任务照常完成"""
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "作品.mp4")), poster=_no_poster)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.status == "done"
assert job.cover_url == "" and job.poster_path is None
async def test_make_poster_extracts_real_frame(wj, tmp_path, monkeypatch):
"""真跑一次 ffmpeg:确认真能抽出 JPEG(假 poster 只覆盖接缝,覆盖不到命令本身)"""
ffmpeg = Path(__file__).resolve().parents[1] / ".venv" / "Scripts" / "ffmpeg.exe"
if not ffmpeg.exists():
pytest.skip("没有 .venv/Scripts/ffmpeg.exe")
monkeypatch.setattr(wj, "ffmpeg_path", lambda: str(ffmpeg))
video = tmp_path / "作品.mp4"
subprocess.run(
[str(ffmpeg), "-f", "lavfi", "-i", "testsrc=size=320x240:rate=10",
"-t", "2", "-pix_fmt", "yuv420p", "-y", str(video)],
capture_output=True, check=True,
)
poster = await wj.make_poster(video)
assert poster is not None
assert poster.name == "作品_封面.jpg"
assert poster.read_bytes()[:2] == b"\xff\xd8" # JPEG magic
async def test_make_poster_falls_back_to_first_frame(wj, tmp_path, monkeypatch):
"""不足 1 秒的视频取不到"第 1 秒"那一帧,要退回第 0 秒再试"""
ffmpeg = Path(__file__).resolve().parents[1] / ".venv" / "Scripts" / "ffmpeg.exe"
if not ffmpeg.exists():
pytest.skip("没有 .venv/Scripts/ffmpeg.exe")
monkeypatch.setattr(wj, "ffmpeg_path", lambda: str(ffmpeg))
video = tmp_path / "短片.mp4"
subprocess.run(
[str(ffmpeg), "-f", "lavfi", "-i", "testsrc=size=160x120:rate=10",
"-t", "0.5", "-pix_fmt", "yuv420p", "-y", str(video)],
capture_output=True, check=True,
)
poster = await wj.make_poster(video)
assert poster is not None and poster.stat().st_size > 0
async def test_make_poster_missing_video_returns_none(wj, tmp_path, monkeypatch):
ffmpeg = Path(__file__).resolve().parents[1] / ".venv" / "Scripts" / "ffmpeg.exe"
if not ffmpeg.exists():
pytest.skip("没有 .venv/Scripts/ffmpeg.exe")
monkeypatch.setattr(wj, "ffmpeg_path", lambda: str(ffmpeg))
assert await wj.make_poster(tmp_path / "不存在.mp4") is None
async def test_refresh_refreshes_cover(wj, manager, tmp_path):
"""刷新链接时封面一起换新(否则卡片封面 1 小时后变裂图)"""
links = iter(["http://lan/a?sig=1", "http://lan/p?sig=1", "http://lan/a?sig=2", "http://lan/p?sig=2"])
async def upload(path, policy):
return next(links), None
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "作品.mp4")), upload=upload)
(await mgr.submit(["https://v.douyin.com/x"]))[0]
job = (await _settle(mgr))[0]
assert job.cover_url == "http://lan/p?sig=1"
refreshed = await mgr.refresh(job.id)
assert refreshed is not None
assert refreshed.cover_url == "http://lan/p?sig=2"
# ───────────────────────── 清理 ─────────────────────────
async def test_remove_job(wj, manager, tmp_path):
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "a.jpg")))
job = (await mgr.submit(["https://v.douyin.com/x"]))[0]
await _settle(mgr)
assert mgr.remove(job.id) is True
assert mgr.list_jobs() == []
assert mgr.remove(job.id) is False # 再删一次:不存在
async def test_remove_running_job_cancels_it(wj, manager, tmp_path):
"""删掉在跑的任务要顺手取消,不能让它继续跑完再写回任务表"""
started = asyncio.Event()
cancelled = asyncio.Event()
async def fetch(url):
started.set()
try:
await asyncio.sleep(30)
except asyncio.CancelledError:
cancelled.set()
raise
return "作品", [_touch(tmp_path, "a.jpg")], False
mgr = manager(fetch=fetch)
job = (await mgr.submit(["https://v.douyin.com/x"]))[0]
await asyncio.wait_for(started.wait(), timeout=2)
assert mgr.remove(job.id) is True
await asyncio.wait_for(cancelled.wait(), timeout=2)
assert mgr.list_jobs() == []
async def test_clear_finished_keeps_running(wj, manager, tmp_path):
release = asyncio.Event()
async def fetch(url):
if "slow" in url:
await release.wait()
return "作品", [_touch(tmp_path, "a.jpg")], False
mgr = manager(fetch=fetch)
fast = (await mgr.submit(["https://v.douyin.com/fast"]))[0]
slow = (await mgr.submit(["https://v.douyin.com/slow"]))[0]
for _ in range(500):
if not fast.active:
break
await asyncio.sleep(0.01)
assert mgr.clear_finished() == 1
assert [j.id for j in mgr.list_jobs()] == [slow.id]
release.set()
await _settle(mgr)
# ───────────────────────── 上限与回收 ─────────────────────────
async def test_job_count_is_capped(wj, manager, tmp_path, monkeypatch):
monkeypatch.setattr(wj, "MAX_JOBS", 3)
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "a.jpg")))
await mgr.submit([f"https://v.douyin.com/{i}" for i in range(6)])
jobs = await _settle(mgr)
assert len(jobs) == 3
async def test_expired_terminal_jobs_are_pruned(wj, manager, tmp_path):
mgr = manager(fetch=_ok_fetch(_touch(tmp_path, "a.jpg")))
job = (await mgr.submit(["https://v.douyin.com/x"]))[0]
await _settle(mgr)
assert len(mgr.list_jobs()) == 1
job.updated_at -= wj.TERMINAL_TTL + 1
assert mgr.list_jobs() == []
async def test_running_jobs_survive_prune(wj, manager, monkeypatch):
"""在跑的任务不受上限/TTL 影响(TTL 设为负数,终态一律该清掉)"""
monkeypatch.setattr(wj, "TERMINAL_TTL", -1)
release = asyncio.Event()
async def fetch(url):
await release.wait()
return None, [], False
mgr = manager(fetch=fetch)
job = (await mgr.submit(["https://v.douyin.com/x"]))[0]
await asyncio.sleep(0.05)
assert [j.id for j in mgr.list_jobs()] == [job.id]
assert job.status == "running" and job.stage == "解析中"
release.set()
# 终态任务会被 TTL=-1 立即回收,所以这里盯 job 本身而不是 list_jobs()
for _ in range(500):
if not job.active:
break
await asyncio.sleep(0.01)
assert job.status == "failed"