Files

572 lines
19 KiB
Python
Raw Permalink Normal View History

"""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"