diff --git a/docs/decisions.md b/docs/decisions.md index 6c6d64a..2887086 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -71,6 +71,20 @@ 并用半透明填充 + 半透明描边降低遮挡感。 - **调研全文**:[VR双目字幕景深与遮挡调查报告.md](./VR双目字幕景深与遮挡调查报告.md)。 +## 批量任务先写明细再排队(CREATING 状态) + +- **现象**:新建批量任务后偶发任务被标 COMPLETED、`done=0/N`,剩余视频永远不再 + 被处理(例:`batch_ac585c5458de` 唯一明细还 PENDING 而任务已 COMPLETED)。 +- **根因**:`create_job` 先把任务行以 QUEUED 入库(此时引擎就看得见),再逐条登记 + 明细(扫描媒体库时 500+ 条要数秒);引擎轮询到的快照可能没包含剩余明细,收尾时 + “无未结束明细”检查也看不到它们,于是把任务标 COMPLETED。 +- **结论**:任务行改以 `CREATING` 入库,明细全部登记完才置 QUEUED(引擎只取 + QUEUED);登记中途异常置 FAILED 并抛给路由;启动时把上一进程遗留的 CREATING + 统一置 FAILED(`fail_creating_batch_jobs`),避免静默残留。 +- **旧数据修复**:`scripts/fix_zombie_batch_jobs.py` 仍用于处理历史 + “COMPLETED 但仍有未结束明细”的脏数据。回归测试见 `test_batch.py` 的 + `test_job_hidden_from_engine_until_details_written`。 + ## 相关文档 - 当前生效的参数与协议:[node-protocol.md](./node-protocol.md) diff --git a/docs/operations.md b/docs/operations.md index a3d13b0..85ca442 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -64,7 +64,10 @@ 空列表不报 500)。只暴露目录名,不返回文件内容。 - **创建任务时一次性定位(2026-09 起)**:`POST /api/batch/jobs {folder, workflow_id, recursive}` 只扫描一次文件夹并把每个视频登记为 `batch_videos` - 明细(PENDING/RUNNING/PAUSED/COMPLETED/FAILED/SKIPPED)。**视频所在目录 + 明细(PENDING/RUNNING/PAUSED/COMPLETED/FAILED/SKIPPED)。任务行先以 + `CREATING` 入库、**明细全部登记完才置 QUEUED**:否则引擎会在明细写一半时拾起 + 任务、收尾把任务误标 COMPLETED(理由见 + [decisions.md](./decisions.md#批量任务先写明细再排队creating-状态))。**视频所在目录 (视频旁)若已存在文件名含视频名的字幕文件**(`.srt/.ass/.ssa/.vtt`, `list_sidecar_subtitles` 判定,如 `movie.CN.srt`、`movie.CN_dual_eye.ass`), 说明该视频已有字幕,创建即记 **SKIPPED**——不为它触发任何流水线。运行时 diff --git a/src/wov_app/batch.py b/src/wov_app/batch.py index b385d8e..5d731bc 100644 --- a/src/wov_app/batch.py +++ b/src/wov_app/batch.py @@ -221,7 +221,10 @@ def create_job( "folder_path": str(folder), "workflow_id": workflow_id, "recursive": int(recursive), - "status": "QUEUED", + # CREATING:明细未写完前引擎看不见本任务。逐条登记 500+ 条明细要数秒, + # 若此刻已是 QUEUED,引擎会读到半个快照并在收尾时把任务误标 COMPLETED, + # 剩下的视频就再也不会被处理(自愈分支只碰运气)。写完明细立刻置 QUEUED。 + "status": "CREATING", "progress": 0, "total": 0, "done": 0, @@ -233,27 +236,35 @@ def create_job( }) pending = 0 skipped = 0 - for video in videos: - # 视频所在目录已存在对应字幕文件 → 已处理过,直接跳过不触发流水线。 - if list_sidecar_subtitles(video): - status = "SKIPPED" - skipped += 1 - else: - status = "PENDING" - pending += 1 - video_id = f"bv_{uuid.uuid4().hex[:12]}" - db.create_batch_video({ - "id": video_id, - "job_id": job_id, - "video_path": str(video), - # 私有工作空间:storage/batch///,与媒体库隔离。 - "work_dir": str(BATCH_WORK_ROOT / job_id / video_id), - "run_id": None, - "status": status, - "error": None, - "created_at": now, - "updated_at": now, - }) + try: + for video in videos: + # 视频所在目录已存在对应字幕文件 → 已处理过,直接跳过不触发流水线。 + if list_sidecar_subtitles(video): + status = "SKIPPED" + skipped += 1 + else: + status = "PENDING" + pending += 1 + video_id = f"bv_{uuid.uuid4().hex[:12]}" + db.create_batch_video({ + "id": video_id, + "job_id": job_id, + "video_path": str(video), + # 私有工作空间:storage/batch///,与媒体库隔离。 + "work_dir": str(BATCH_WORK_ROOT / job_id / video_id), + "run_id": None, + "status": status, + "error": None, + "created_at": now, + "updated_at": now, + }) + except Exception as exc: # noqa: BLE001 + # 明细写到一半失败:记 FAILED 留可见记录(CREATING 状态没人会拾起, + # 沉默的残留任务会让用户以为什么都没发生),然后把异常交给路由层。 + db.update_batch_job( + job_id, status="FAILED", error=str(exc), updated_at=_now_iso(), + ) + raise # total = 本批真正需要处理(无字幕)的视频数;已有字幕被 SKIPPED 的 # 不计入总数也不计入完成数——进度条只反映"实际待处理"的这批。 if pending == 0: @@ -263,7 +274,10 @@ def create_job( updated_at=_now_iso(), ) else: - db.update_batch_job(job_id, total=pending, updated_at=_now_iso()) + # 明细全部就位后才排队,引擎从此拿到的快照一定是完整的。 + db.update_batch_job( + job_id, status="QUEUED", total=pending, updated_at=_now_iso(), + ) logger.info( "创建批量任务 %s: 文件夹 %s, 工作流 %s, 共 %d 个视频(%d 待处理, %d 已有字幕跳过)", job_id, folder, workflow_id, len(videos), pending, skipped, @@ -291,6 +305,11 @@ class BatchWorker: return # 独立运行时确保节点已注册;重复注册幂等。 registry.register_all() + # 上一进程中断留下的 CREATING 任务(明细登记中途被杀/热重载)没人会推进, + # 启动时统一记为 FAILED,避免用户以为任务还在创建中。 + stale = self.db.fail_creating_batch_jobs("创建明细中断(进程中断),请重新创建任务") + if stale: + logger.warning("启动清理 %d 个未登记完的批量任务(CREATING → FAILED)", stale) self._stopping = False self._thread = threading.Thread( target=self._loop, diff --git a/src/wov_app/db.py b/src/wov_app/db.py index 331c5e2..357997f 100644 --- a/src/wov_app/db.py +++ b/src/wov_app/db.py @@ -590,6 +590,20 @@ class Database: conn.execute(f"UPDATE batch_jobs SET {assignments} WHERE id = ?", values) + def fail_creating_batch_jobs(self, reason: str) -> int: + """把停留在 CREATING 的批量任务置 FAILED,返回处理条数。 + + CREATING 只存在于 create_job 逐条登记明细期间;进程被杀或热重载后没有任何 + 线程会推进它(引擎只取 QUEUED),启动时统一收尾成可见的失败记录。 + """ + with self._connect() as conn: + cursor = conn.execute( + "UPDATE batch_jobs SET status = 'FAILED', error = ?, updated_at = ? " + "WHERE status = 'CREATING'", + (reason, _now_iso()), + ) + return int(cursor.rowcount or 0) + def sync_batch_job_progress(self, job_id: str) -> None: """按视频明细实时对齐任务的 total/done/failed 汇总并落库。 diff --git a/tests/app/test_batch/test_batch.py b/tests/app/test_batch/test_batch.py index 3bbdc10..98f7ba8 100644 --- a/tests/app/test_batch/test_batch.py +++ b/tests/app/test_batch/test_batch.py @@ -748,3 +748,88 @@ def test_worker_clears_stale_keep_model_flag_before_stage(tmp_path: Path, monkey # 验证结果:prep(非 LLM)看不到残留标志;translate(LLM 阶段)才写入; # post(非 LLM)不再看到它。 assert flag_trace == [("prep", False), ("translate", True), ("post", False)] + + +def test_job_hidden_from_engine_until_details_written(tmp_path: Path, monkeypatch) -> None: + """建任务期间任务对引擎不可见:明细写完前不被拾起,避免半成品被标完成。 + + 曾因任务行先入库、524 条明细后写,引擎在明细写一半时拾起任务、收尾时快照 + 里没有剩余明细,把任务误标 COMPLETED(视频永远不再被处理)。 + """ + # 数据:3 个待处理视频的媒体库。 + folder = tmp_path / "videos" + for name in ("a.mp4", "b.mp4", "c.mp4"): + _make_video(folder / name) + db = _published_db(tmp_path) + probes: list[str | None] = [] + original = db.create_batch_video + + def _probe_after_insert(item: dict) -> None: + """每写完一条明细,立刻问一次引擎队列(模拟轮询线程的拾取时机)。""" + original(item) + picked = db.next_queued_batch_job() + probes.append(picked["id"] if picked else None) + + monkeypatch.setattr(db, "create_batch_video", _probe_after_insert) + + # 测试过程 + job_id = create_job(db, str(folder), "wf", recursive=False) + + # 验证结果:写入过程中引擎始终取不到任务;写完后才是 QUEUED 且明细完整。 + assert probes == [None, None, None] + job = db.get_batch_job(job_id) + assert job["status"] == "QUEUED" + assert job["total"] == 3 + assert len(db.list_batch_videos(job_id)) == 3 + assert db.next_queued_batch_job()["id"] == job_id + + +def test_create_job_marks_failed_when_detail_insert_breaks(tmp_path: Path, monkeypatch) -> None: + """明细写入中途失败时任务记 FAILED 并留下错误,不产生看不见的残留任务。""" + # 数据:2 个视频,第二条明细写入时抛异常。 + folder = tmp_path / "videos" + for name in ("a.mp4", "b.mp4"): + _make_video(folder / name) + db = _published_db(tmp_path) + original = db.create_batch_video + calls = {"n": 0} + + def _fail_second(item: dict) -> None: + calls["n"] += 1 + if calls["n"] == 2: + raise RuntimeError("磁盘写满") + original(item) + + monkeypatch.setattr(db, "create_batch_video", _fail_second) + + # 测试过程 + 验证结果:异常继续抛出,任务可被观察到且为 FAILED。 + with pytest.raises(RuntimeError): + create_job(db, str(folder), "wf", recursive=False) + job = db.list_batch_jobs(limit=10)[0] + assert job["status"] == "FAILED" + assert "磁盘写满" in (job["error"] or "") + assert db.next_queued_batch_job() is None + + +def test_worker_start_fails_leftover_creating_job(tmp_path: Path) -> None: + """进程中断留下的 CREATING 任务在引擎启动时记为 FAILED,不静默残留。""" + # 数据:一条只登记到一半的任务(模拟明细写入中途进程被杀/热重载)。 + db = _published_db(tmp_path) + db.create_batch_job({ + "id": "batch_halfway", "folder_path": "/videos", "workflow_id": "wf", + "recursive": 0, "status": "CREATING", "progress": 0, "total": 0, "done": 0, + "failed": 0, "current_video": None, "error": None, + "created_at": "2026-09-18T00:00:00+00:00", "updated_at": "2026-09-18T00:00:00+00:00", + }) + worker = BatchWorker(db, interval_seconds=999) + + # 测试过程 + worker.start() + try: + # 验证结果:任务变为可见的 FAILED,且不会被引擎当排队任务拾起。 + job = db.get_batch_job("batch_halfway") + assert job["status"] == "FAILED" + assert "中断" in (job["error"] or "") + assert db.next_queued_batch_job() is None + finally: + worker.stop()