- 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>
572 lines
19 KiB
Python
572 lines
19 KiB
Python
"""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"
|