diff --git a/AGENTS.md b/AGENTS.md index 30f9d15..68b0302 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -339,6 +339,16 @@ http://127.0.0.1:8000/docs API 文档 execute_run 从产物表跳过已完成节点、只重跑失败节点——extract/ocr 等长耗时 成果不浪费;配合 llm-filter/OCR 的节点级断点存档,失败节点自身也只重判未完成 条目。前端对"部分失败"(COMPLETED 且 failed>0)用红色徽章醒目标示。 +- **完成任务判定(2026-09 修复)**:`_run_job` 置 COMPLETED 前**校验全部非 + SKIPPED 视频都已结束**(无 PENDING/PAUSED 残留),否则保持 RUNNING 交引擎 + 下一轮续跑——修复僵尸状态:引擎串行处理到 9.9GB 大视频时中断,`_run_job` + 无条件收尾把任务置 COMPLETED,留下"N 个 PENDING 待处理却已完成"的假完成 + (batch_969fabe74b83 等 3 个任务实测:遗留的 10 个 PENDING 完全相同且卡在 + kiwvr-887 大文件前)。**崩溃恢复**:重启时除 `recover_interrupted_runs` 外, + 新增 `recover_interrupted_batch_jobs` 把 RUNNING 的批量任务恢复为 QUEUED + (否则停在 RUNNING 的批量任务永远不会被 `next_queued_batch_job` 再次拾起, + 未处理完的 PENDING 永久残留)。历史僵尸数据修复脚本见 + `scripts/fix_zombie_batch_jobs.py`(把误标 COMPLETED 的任务置回 QUEUED 续跑)。 - **产物下载**:`GET /api/batch/jobs/{id}/videos/{vid}/download?alias=<文件名>` 解析并返回视频旁的字幕文件;旧版 `batch.done.json` 完成标记里的语义别名 (位于旧 work_dir)仍兼容可下载。详情/创建响应里每个视频的 `finals` 合并上述 diff --git a/scripts/fix_zombie_batch_jobs.py b/scripts/fix_zombie_batch_jobs.py new file mode 100644 index 0000000..2104645 --- /dev/null +++ b/scripts/fix_zombie_batch_jobs.py @@ -0,0 +1,68 @@ +"""修复僵尸批量任务脚本:把误标 COMPLETED 但仍有 PENDING 视频的任务置回 QUEUED。 + +背景(batch_969fabe74b83 事故):批量引擎处理大视频时中断,_run_job 无条件 +收尾把任务置 COMPLETED,留下"N 个 PENDING 待处理却已完成"的僵尸状态。 +代码已修复(置 COMPLETED 前校验无 PENDING 残留 + 重启恢复 RUNNING 批量任务), +本脚本用于修复**历史遗留**的 3 个僵尸任务数据: +- batch_969fabe74b83(10 个 PENDING,0 完成) +- batch_1febe532a7cd(10 个 PENDING,11 完成) +- batch_73c2b723456a(10 个 PENDING,6 完成) + +修复方式:仅把这 3 个任务置回 QUEUED 并清空误导的 total/progress,保留 +SKIPPED/COMPLETED 明细与 run 记录;批量引擎(新代码)重启后会重新拾起, +从剩余 PENDING 视频续跑,全部处理完才置 COMPLETED。 + +安全约束:只更新明确列出的 3 个 job_id;其余任务(含正常 COMPLETED 的 +fee6681/ac8ca4f)不动。执行前打印将变更的任务与明细统计供确认。 +""" +from __future__ import annotations + +import sys +from datetime import datetime, timezone +from pathlib import Path + +from wov_app.db import Database + +# 待修复的僵尸任务(经诊断确认:COMPLETED 但仍有 PENDING 残留)。 +ZOMBIE_JOBS = [ + "batch_969fabe74b83", # 0 完成 / 10 PENDING(最严重,从未真正处理) + "batch_1febe532a7cd", # 11 完成 / 10 PENDING + "batch_73c2b723456a", # 6 完成 / 10 PENDING +] + + +def _now_iso() -> str: + """返回当前 UTC 时间的 ISO 格式字符串。""" + return datetime.now(timezone.utc).isoformat() + + +def main(db_path: str, apply: bool = False) -> None: + """诊断(默认)或修复(--apply)僵尸批量任务。""" + db = Database(Path(db_path)) + print(f"数据库: {db_path}\n") + for job_id in ZOMBIE_JOBS: + job = db.get_batch_job(job_id) + if job is None: + print(f" [跳过] {job_id}: 任务不存在") + continue + pending = sum(1 for v in db.list_batch_videos(job_id) if v["status"] == "PENDING") + completed = sum(1 for v in db.list_batch_videos(job_id) if v["status"] == "COMPLETED") + # 安全校验:只修复"COMPLETED 但仍有 PENDING"的僵尸状态;已正常完成的跳过。 + if job["status"] != "COMPLETED" or pending == 0: + print(f" [跳过] {job_id}: status={job['status']}, PENDING={pending},非僵尸状态") + continue + print(f" [待修] {job_id}: status={job['status']} → QUEUED, " + f"COMPLETED={completed}, PENDING={pending}") + if apply: + db.update_batch_job( + job_id, status="QUEUED", progress=0, total=pending, + done=completed, failed=int(job["failed"] or 0), + current_video=None, error=None, updated_at=_now_iso(), + ) + print(f" ✓ 已置回 QUEUED(引擎将续跑剩余 {pending} 个 PENDING)") + if not apply: + print("\n以上为预览。确认无误后加 --apply 实际修复。") + + +if __name__ == "__main__": + main(sys.argv[1] if len(sys.argv) > 1 else "data/wov.db", "--apply" in sys.argv) diff --git a/src/wov_app/batch.py b/src/wov_app/batch.py index 290ed98..6c01803 100644 --- a/src/wov_app/batch.py +++ b/src/wov_app/batch.py @@ -408,8 +408,31 @@ class BatchWorker: # 全部视频处理完成:先用明细实时对齐汇总(done 只计实际完成的, # 不含 SKIPPED),再置 COMPLETED。 + # + # 置 COMPLETED 前必须校验**没有未处理完的视频残留**:若本轮循环因 + # 视频处理中断/异常(_process_video 返回但视频仍 PENDING,等同进程 + # 在处理中被杀)而没有真正处理完所有 PENDING,就**不能**标完成—— + # 否则会出现"明细还有 N 个待处理、任务却已完成"的僵尸状态 + # (batch_969fabe74b83 等 3 个任务真实发生:引擎串行处理到 9.9GB + # 大视频时中断,10 个视频留 PENDING 却被无条件置 COMPLETED)。 + # 此时保持 RUNNING,让引擎下一轮(重启后重新拾起 RUNNING 任务) + # 继续处理剩余 PENDING,全部结束才真正置 COMPLETED。 self.db.sync_batch_job_progress(job_id) job = self.db.get_batch_job(job_id) + if job is None: + return + leftovers = [ + v for v in self.db.list_batch_videos(job_id) + if v["status"] not in ("SKIPPED", "COMPLETED", "FAILED") + ] + if leftovers: + # 有未处理完的视频(PENDING/PAUSED/QUEUED 等):保持 RUNNING, + # 由引擎下一轮续跑;记录日志便于排查中断位置。 + logger.warning( + "批量任务 %s 仍有 %d 个视频未处理完(%s…),保持 RUNNING 待续跑,不置 COMPLETED", + job_id, len(leftovers), Path(leftovers[0]["video_path"]).name, + ) + return done = int(job["done"]) if job else 0 failed = int(job["failed"]) if job else 0 self.db.update_batch_job( diff --git a/src/wov_app/db.py b/src/wov_app/db.py index 7512cc9..197b5b4 100644 --- a/src/wov_app/db.py +++ b/src/wov_app/db.py @@ -385,6 +385,24 @@ class Database: ) return cur.rowcount + def recover_interrupted_batch_jobs(self, updated_at: str) -> int: + """重启恢复:把遗留 RUNNING 的批量任务恢复为 QUEUED,返回恢复数量。 + + 批量引擎处理视频(尤其大文件)时进程被杀/重启,批量任务会停在 + RUNNING:其关联 run 由 recover_interrupted_runs 恢复为 QUEUED,但 + 批量任务本身若保持 RUNNING,next_queued_batch_job 只拾取 QUEUED, + 永远不会重新驱动它 → 未处理完的 PENDING 视频永久残留 + (batch_969fabe74b83 事故链路之一)。恢复为 QUEUED 后引擎重新拾起, + 从断点(剩余 PENDING 视频 + 已恢复的 run)继续处理。用户主动暂停的 + PAUSED 批量任务保持不变,等待显式 resume。 + """ + with self._connect() as conn: + cur = conn.execute( + "UPDATE batch_jobs SET status = 'QUEUED', updated_at = ? WHERE status = 'RUNNING'", + (updated_at,), + ) + return cur.rowcount + def pause_run(self, run_id: str, updated_at: str) -> None: """暂停任务:置为 PAUSED;调度器会在节点边界检查并停止推进。""" with self._connect() as conn: diff --git a/src/wov_app/main.py b/src/wov_app/main.py index 7c3cb87..e379c48 100644 --- a/src/wov_app/main.py +++ b/src/wov_app/main.py @@ -42,7 +42,11 @@ async def lifespan(app: FastAPI): seed_default_workflows(db) # 重启恢复:上次进程被杀时遗留的 RUNNING 任务恢复为 QUEUED, # 调度器会从产物表断点续跑(不重复已完成节点);PAUSED 保持等待显式 resume。 + # 批量任务同样恢复:RUNNING 的批量 job 置回 QUEUED,其关联 run 由 + # recover_interrupted_runs 恢复,批量引擎重新拾起后从断点续跑剩余 PENDING + # (不恢复则批量任务停在 RUNNING 永远不会被再次驱动,未处理视频永久残留)。 db.recover_interrupted_runs(datetime.now(timezone.utc).isoformat()) + db.recover_interrupted_batch_jobs(datetime.now(timezone.utc).isoformat()) scheduler = WorkflowScheduler(db, STORAGE_DIR) # 调度器默认开启,处理排队中的任务;测试可关闭后手动执行。 if os.getenv("WOV_SCHEDULER_ENABLED", "1") == "1": diff --git a/tests/test_batch.py b/tests/test_batch.py index f9e0aaf..2eb8d39 100644 --- a/tests/test_batch.py +++ b/tests/test_batch.py @@ -634,6 +634,69 @@ def test_batch_worker_run_deleted_by_scheduler_returns(tmp_path, monkeypatch) -> assert video["status"] == "PENDING" +def test_batch_worker_pending_leftover_never_marks_completed(tmp_path, monkeypatch) -> None: + """残留未处理 PENDING 视频时,任务**不得**置 COMPLETED(状态机缺陷回归)。 + + 真实事故(batch_969fabe74b83 等 3 个任务):引擎串行处理到 9.9GB 大视频时 + execute_run 异常中断,_process_video 返回但该视频未被标 FAILED/COMPLETED + (保持 PENDING),_run_job 循环照常走完剩余 SKIPPED 后**无条件**收尾置 + COMPLETED——留下"10 个待处理却已完成"的僵尸状态。 + + 本测试用假调度器复现:第一个视频 execute_run 抛异常且不标任何状态 + (run 被删除的极端情形同路径),第二个视频正常;断言任务必须保持 + RUNNING(而非 COMPLETED),等待引擎下一次拾起重跑剩余 PENDING。 + """ + db = _db(tmp_path) + _seed_echo_workflow(db) + # a.mp4 待处理(处理中崩溃);b.mp4 旁放好字幕 → 创建即 SKIPPED,不参与处理。 + folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) + (folder / "b.CN.srt").write_text("1\n00:00:00,000 --> 00:00:01,000\nb\n", encoding="utf-8") + job = _make_job(db, folder) + videos0 = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + assert videos0["b.mp4"]["status"] == "SKIPPED" # 前置:旁挂字幕命中跳过 + + class _CrashScheduler: + """模拟进程在第一个(唯一待处理)视频处理中被杀:run 消失、视频保持 PENDING。""" + + def __init__(self, db: Database, work_dir: Path) -> None: + self.db = db + + def execute_run(self, run_id: str) -> None: + # 等同进程崩溃后 run 记录残留被清/丢失,视频仍是 PENDING。 + self.db.delete_run(run_id) + + monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _CrashScheduler) + worker = BatchWorker(db, interval_seconds=0.05) + # 第一轮:a 崩溃残留 PENDING → 任务必须仍 RUNNING(允许续跑),不得 COMPLETED。 + worker._process_job(job) + job = db.get_batch_job(job["id"]) + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + # 缺陷行为:job 被置 COMPLETED(无条件收尾);期望保持 RUNNING。 + assert videos["a.mp4"]["status"] == "PENDING" + assert videos["b.mp4"]["status"] == "SKIPPED" + assert job["status"] == "RUNNING", ( + f"残留 PENDING 时任务被误置为 {job['status']}(缺陷:无条件收尾置 COMPLETED)" + ) + + # 第二轮(引擎重启后再次拾起 RUNNING 任务):a 这次正常完成 → 全部完成才 COMPLETED。 + monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _CrashScheduler) + # 放开假调度器的崩溃:第二次调用时不再删 run,让真实调度器跑通。 + class _RecoveringScheduler: + """第二轮:不再崩溃,把 run 置 COMPLETED(由引擎收尾放产物)。""" + + def __init__(self, db: Database, work_dir: Path) -> None: + self.db = db + + def execute_run(self, run_id: str) -> None: + self.db.update_run(run_id, status="COMPLETED", updated_at=_now_iso()) + + monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _RecoveringScheduler) + worker._process_job(job) + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + assert videos["a.mp4"]["status"] == "COMPLETED" + assert db.get_batch_job(job["id"])["status"] == "COMPLETED" + + def test_batch_worker_already_completed_runs_finalized(tmp_path) -> None: """run 已完成但视频未标记(收尾前中断):直接放置产物、清理并标记完成。 diff --git a/tests/test_db.py b/tests/test_db.py index 75f8276..44cdf2a 100644 --- a/tests/test_db.py +++ b/tests/test_db.py @@ -299,6 +299,39 @@ def test_recover_interrupted_runs(tmp_path) -> None: # 恢复后调度器可拾起并断点续跑。 assert db.next_queued_run()["id"] == "run_orphan" +def test_recover_interrupted_batch_jobs(tmp_path) -> None: + """重启恢复:遗留 RUNNING 的批量任务恢复为 QUEUED,供引擎重新拾起。 + + 批量引擎处理视频(尤其大文件)时进程被杀/重启,批量任务停在 RUNNING。 + 若保持 RUNNING,next_queued_batch_job 只拾取 QUEUED,任务永远不会被再次 + 驱动 → 未处理完的 PENDING 视频永久残留(batch_969fabe74b83 事故链路之一)。 + 恢复为 QUEUED 后引擎重新拾起续跑;PAUSED 批量任务保持不变等待显式 resume。 + """ + db = Database(tmp_path / "wov.db") + now = "2026-01-01T00:00:00+00:00" + + def add_job(job_id: str, status: str) -> None: + """创建指定状态的批量任务记录。""" + db.create_batch_job( + { + "id": job_id, "folder_path": "/tmp/videos", "workflow_id": "demo", + "recursive": 1, "status": status, "progress": 0.5, "total": 3, + "done": 1, "failed": 0, "error": None, + "created_at": now, "updated_at": now, + } + ) + + add_job("job_orphan", "RUNNING") + add_job("job_paused", "PAUSED") + add_job("job_done", "COMPLETED") + recovered = db.recover_interrupted_batch_jobs("2026-01-02T00:00:00+00:00") + assert recovered == 1 # 只有 RUNNING 被恢复。 + assert db.get_batch_job("job_orphan")["status"] == "QUEUED" + assert db.get_batch_job("job_paused")["status"] == "PAUSED" + assert db.get_batch_job("job_done")["status"] == "COMPLETED" + # 恢复后批量引擎可拾起并从断点续跑。 + assert db.next_queued_batch_job()["id"] == "job_orphan" + def test_restore_run_outputs(tmp_path) -> None: """验证从产物重建节点输出(断点续跑的依据)。""" db = Database(tmp_path / "wov.db")