fix: 批量任务残留PENDING仍误标COMPLETED——置完成前校验无残留+崩溃恢复批量任务
事故(batch_969fabe74b83 等 3 个任务实测):批量引擎串行处理到 9.9GB 大视频时中断,_run_job 无条件收尾把任务置 COMPLETED,留下"N 个 PENDING 待处理却已完成"的僵尸状态——3 个任务遗留的 PENDING 完全相同(kiwvr-887/ mdvr-413/savr-1077 共 10 个大视频),且均卡在大文件前从未被真正处理。 - batch.py: _run_job 置 COMPLETED 前校验全部非 SKIPPED 视频已结束 (无 PENDING/PAUSED 残留),否则保持 RUNNING 交引擎下一轮续跑 - db.py/main.py: 新增 recover_interrupted_batch_jobs,重启时把 RUNNING 批量任务恢复 QUEUED(否则停在 RUNNING 永远不会被再次拾起) - tests: 新增 2 条 TDD 回归(残留 PENDING 不置 COMPLETED、批量崩溃恢复) - scripts/fix_zombie_batch_jobs.py: 历史僵尸数据修复脚本(置回 QUEUED 续跑) - AGENTS.md: 补充完成任务判定与崩溃恢复约定
This commit is contained in:
@@ -339,6 +339,16 @@ http://127.0.0.1:8000/docs API 文档
|
|||||||
execute_run 从产物表跳过已完成节点、只重跑失败节点——extract/ocr 等长耗时
|
execute_run 从产物表跳过已完成节点、只重跑失败节点——extract/ocr 等长耗时
|
||||||
成果不浪费;配合 llm-filter/OCR 的节点级断点存档,失败节点自身也只重判未完成
|
成果不浪费;配合 llm-filter/OCR 的节点级断点存档,失败节点自身也只重判未完成
|
||||||
条目。前端对"部分失败"(COMPLETED 且 failed>0)用红色徽章醒目标示。
|
条目。前端对"部分失败"(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=<文件名>`
|
- **产物下载**:`GET /api/batch/jobs/{id}/videos/{vid}/download?alias=<文件名>`
|
||||||
解析并返回视频旁的字幕文件;旧版 `batch.done.json` 完成标记里的语义别名
|
解析并返回视频旁的字幕文件;旧版 `batch.done.json` 完成标记里的语义别名
|
||||||
(位于旧 work_dir)仍兼容可下载。详情/创建响应里每个视频的 `finals` 合并上述
|
(位于旧 work_dir)仍兼容可下载。详情/创建响应里每个视频的 `finals` 合并上述
|
||||||
|
|||||||
@@ -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)
|
||||||
@@ -408,8 +408,31 @@ class BatchWorker:
|
|||||||
|
|
||||||
# 全部视频处理完成:先用明细实时对齐汇总(done 只计实际完成的,
|
# 全部视频处理完成:先用明细实时对齐汇总(done 只计实际完成的,
|
||||||
# 不含 SKIPPED),再置 COMPLETED。
|
# 不含 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)
|
self.db.sync_batch_job_progress(job_id)
|
||||||
job = self.db.get_batch_job(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
|
done = int(job["done"]) if job else 0
|
||||||
failed = int(job["failed"]) if job else 0
|
failed = int(job["failed"]) if job else 0
|
||||||
self.db.update_batch_job(
|
self.db.update_batch_job(
|
||||||
|
|||||||
@@ -385,6 +385,24 @@ class Database:
|
|||||||
)
|
)
|
||||||
return cur.rowcount
|
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:
|
def pause_run(self, run_id: str, updated_at: str) -> None:
|
||||||
"""暂停任务:置为 PAUSED;调度器会在节点边界检查并停止推进。"""
|
"""暂停任务:置为 PAUSED;调度器会在节点边界检查并停止推进。"""
|
||||||
with self._connect() as conn:
|
with self._connect() as conn:
|
||||||
|
|||||||
@@ -42,7 +42,11 @@ async def lifespan(app: FastAPI):
|
|||||||
seed_default_workflows(db)
|
seed_default_workflows(db)
|
||||||
# 重启恢复:上次进程被杀时遗留的 RUNNING 任务恢复为 QUEUED,
|
# 重启恢复:上次进程被杀时遗留的 RUNNING 任务恢复为 QUEUED,
|
||||||
# 调度器会从产物表断点续跑(不重复已完成节点);PAUSED 保持等待显式 resume。
|
# 调度器会从产物表断点续跑(不重复已完成节点);PAUSED 保持等待显式 resume。
|
||||||
|
# 批量任务同样恢复:RUNNING 的批量 job 置回 QUEUED,其关联 run 由
|
||||||
|
# recover_interrupted_runs 恢复,批量引擎重新拾起后从断点续跑剩余 PENDING
|
||||||
|
# (不恢复则批量任务停在 RUNNING 永远不会被再次驱动,未处理视频永久残留)。
|
||||||
db.recover_interrupted_runs(datetime.now(timezone.utc).isoformat())
|
db.recover_interrupted_runs(datetime.now(timezone.utc).isoformat())
|
||||||
|
db.recover_interrupted_batch_jobs(datetime.now(timezone.utc).isoformat())
|
||||||
scheduler = WorkflowScheduler(db, STORAGE_DIR)
|
scheduler = WorkflowScheduler(db, STORAGE_DIR)
|
||||||
# 调度器默认开启,处理排队中的任务;测试可关闭后手动执行。
|
# 调度器默认开启,处理排队中的任务;测试可关闭后手动执行。
|
||||||
if os.getenv("WOV_SCHEDULER_ENABLED", "1") == "1":
|
if os.getenv("WOV_SCHEDULER_ENABLED", "1") == "1":
|
||||||
|
|||||||
@@ -634,6 +634,69 @@ def test_batch_worker_run_deleted_by_scheduler_returns(tmp_path, monkeypatch) ->
|
|||||||
assert video["status"] == "PENDING"
|
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:
|
def test_batch_worker_already_completed_runs_finalized(tmp_path) -> None:
|
||||||
"""run 已完成但视频未标记(收尾前中断):直接放置产物、清理并标记完成。
|
"""run 已完成但视频未标记(收尾前中断):直接放置产物、清理并标记完成。
|
||||||
|
|
||||||
|
|||||||
@@ -299,6 +299,39 @@ def test_recover_interrupted_runs(tmp_path) -> None:
|
|||||||
# 恢复后调度器可拾起并断点续跑。
|
# 恢复后调度器可拾起并断点续跑。
|
||||||
assert db.next_queued_run()["id"] == "run_orphan"
|
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:
|
def test_restore_run_outputs(tmp_path) -> None:
|
||||||
"""验证从产物重建节点输出(断点续跑的依据)。"""
|
"""验证从产物重建节点输出(断点续跑的依据)。"""
|
||||||
db = Database(tmp_path / "wov.db")
|
db = Database(tmp_path / "wov.db")
|
||||||
|
|||||||
Reference in New Issue
Block a user