From dcdc5e8604875dd0ab081bcb07c07002c3d19e2f Mon Sep 17 00:00:00 2001 From: cat <1716967236@qq.com> Date: Fri, 18 Sep 2026 10:31:52 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=89=B9=E9=87=8F=E5=88=86=E5=9D=97?= =?UTF-8?q?=E6=B5=81=E6=B0=B4=E7=BA=BF=E3=80=81=E6=9C=AC=E5=9C=B0=E6=A8=A1?= =?UTF-8?q?=E5=9E=8B=E6=98=BE=E5=AD=98=E8=AE=A9=E6=B8=A1=E4=B8=8E=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=88=97=E8=A1=A8=E5=88=86=E5=B7=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 批量引擎改为「分块流水线」:视频按 WOV_BATCH_STAGE_GROUP_SIZE(默认 8)分组, 组内按 DAG 拓扑序跑完全部视频(全部 extract → 全部 ASR → 全部翻译 → 全部 ASS) 再进入下一组,本地模型每组只加载一次、卸载一次,而不是每个视频来回加载卸载; 产物仍按组增量落到视频旁。调度器新增 execute_run(run_id, stop_after=节点): 该节点完成后任务保持 RUNNING 不收尾,下一次调用从产物表跳过已完成节点继续, 用于实现阶段边界。 - nodes/llm.py:翻译节点结束释放本机 Ollama 显存(node 参数 unload_after > LLM_UNLOAD_AFTER > 本机 loopback 端点默认卸载,云端端点不卸载;卸载失败只告警), 新增 keep_model.flag 语义(阶段内保持常驻)与 release_local_model(); 新增节点内暂停(按批 20 行检查 paused.flag,抛 PauseRequested,调度器保持 PAUSED)。 - src/wov_app/batch.py:分组阶段执行与阶段末统一释放显存;失败视频只在它失败 节点的那个阶段重试(避免 LLM 已常驻时重跑 ASR 抢显存);任务没有明细时保持 QUEUED 等登记完成、仍有未完成视频时置回 QUEUED 自愈(原先留 RUNNING 会卡死: 引擎只拾取 QUEUED,任务停在“运行中但没人推进”);无失败视频时删除任务级空目录; 每个阶段开始前清理 paused.flag / keep_model.flag,避免强杀残留影响后续阶段。 - src/wov_app/config.py:新增 WOV_BATCH_STAGE_GROUP_SIZE(设为 1 即旧的每视频全链路)。 - 任务列表与批量页分工:GET /api/runs 默认排除 source=batch(一个批量任务会产生 N 条单视频 run,会把 20 条窗口占满;且任务管理页的暂停/重试/删除对批量 run 语义不成立),需要排查时用 include_batch=1;作为补偿批量页详情新增阶段列 (阶段 i/N · 中文标签,由该视频 run 的 current_node_id 在 DAG 拓扑序中的位置 推导,节点类型映射中文标签)。阶段只有节点边界粒度,句级进度不落库、只在日志。 - 顺带纳入此前未提交的批量僵尸状态恢复:recover_interrupted_batch_jobs 除 RUNNING 外也把「COMPLETED 但仍含未结束视频」的任务置回 QUEUED;fix_zombie_batch_jobs.py 改为按条件扫描并支持 --apply 预览;批量页明细只列本批真正处理过的视频。 测试新增/更新:分块流水线调用顺序(组内按节点跑完再下一组)、每组只释放一次模型、 阶段内保持常驻标志、翻译按批暂停、失败视频不跨阶段推进、任务无明细/中途登记视频时 置回 QUEUED、任务工作空间与残留信号清理、任务列表默认过滤批量 run、详情阶段字段、 前端阶段列渲染;全量 507 passed(唯一失败为既有素材缺失的 integration 用例)。 --- docs/configuration.md | 2 + docs/node-protocol.md | 2 +- docs/operations.md | 63 +++- docs/testing.md | 4 + docs/代码审查问题跟踪.md | 9 + nodes/llm.py | 166 +++++++-- scripts/fix_zombie_batch_jobs.py | 73 ++-- src/wov_app/batch.py | 250 ++++++++----- src/wov_app/config.py | 5 + src/wov_app/db.py | 40 ++- src/wov_app/routers/apps.py | 15 +- src/wov_app/routers/batch.py | 62 +++- src/wov_app/scheduler.py | 33 +- tests/app/test_batch/test_batch.py | 338 +++++++++++++++++- tests/app/test_db/test_database.py | 51 +++ tests/app/test_routers/test_apps_api.py | 21 ++ tests/app/test_routers/test_batch_api.py | 63 ++++ tests/app/test_scheduler/test_scheduler.py | 66 ++++ tests/nodes/test_llm/test_translate.py | 266 ++++++++++++++ tests/scripts/__init__.py | 0 .../test_fix_zombie_batch_jobs/__init__.py | 0 .../test_fix_zombie_batch_jobs.py | 83 +++++ tests/web/test_batch/__init__.py | 1 + tests/web/test_batch/test_batch_js.py | 154 ++++++++ web/assets/batch.js | 31 +- 25 files changed, 1606 insertions(+), 192 deletions(-) create mode 100644 tests/scripts/__init__.py create mode 100644 tests/scripts/test_fix_zombie_batch_jobs/__init__.py create mode 100644 tests/scripts/test_fix_zombie_batch_jobs/test_fix_zombie_batch_jobs.py create mode 100644 tests/web/test_batch/__init__.py create mode 100644 tests/web/test_batch/test_batch_js.py diff --git a/docs/configuration.md b/docs/configuration.md index b9032f2..1c2d022 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -33,6 +33,7 @@ http://127.0.0.1:8000/docs API 文档 | `WOV_CLEANUP_GRACE_SECONDS` | `3600` | 孤儿清理宽限期(秒) | | `WOV_BATCH_ENABLED` | `1` | 开启文件夹批量处理引擎(处理 source=batch 任务) | | `WOV_BATCH_INTERVAL_SECONDS` | `1.0` | 批量引擎轮询间隔 | +| `WOV_BATCH_STAGE_GROUP_SIZE` | `8` | 批量「分块流水线」分组大小:每组视频按节点顺序跑完全部阶段(全部 extract → 全部 ASR → 全部翻译 → 全部 ASS)再进入下一组,本地模型每组只加载一次;设为 1 等价于每个视频各跑完整链路(产物逐视频落地最及时) | | `WOV_AUTO_VAD` | `1` | 开启每视频自适应 VAD 调参(详见 [adaptive_vad.md](./adaptive_vad.md)) | | `WHISPER_MODEL_PATH` | 见 [模型权重解析](./node-protocol.md#模型权重解析本地优先) | 显式指定 whisper 模型路径 | | `WHISPER_DEVICE` | `auto` | 转写设备 | @@ -40,6 +41,7 @@ http://127.0.0.1:8000/docs API 文档 | `LLM_API_KEY` | 空(读 `.env`) | SiliconFlow Bearer Key,存于 gitignored 的 `.env` | | `LLM_MODEL` | `Qwen/Qwen3.5-35B-A3B` | LLM 模型名(默认值与例外说明见 [decisions.md](./decisions.md#翻译模型默认值切换与-subtitle-correction-例外)) | | `LLM_TIMEOUT_SECONDS` | `600` | LLM 单请求超时(llm-filter 内部默认 60) | +| `LLM_UNLOAD_AFTER` | 空(按端点自动) | llm-translate 结束后是否卸载本地模型:空=端点在本机时卸载(把显存让给后续 whisper 等节点)、`1` 强制卸载、`0` 关闭;节点参数 `unload_after` 优先级更高 | | `OLLAMA_HOST` | `http://192.168.123.70:11434` | Ollama 服务地址 | | `VLM_MODEL` | `glm-ocr:latest` | VLM OCR 模型 | | `VLM_PROMPT` | 提取图像中的文字,不要描述图片中的内容 | OCR 提示词(字幕流水线在 ocr-subtitle 工作流的 subtitle-ocr 节点参数中显式指定同一提示词) | diff --git a/docs/node-protocol.md b/docs/node-protocol.md index 3c6c340..cc94ea9 100644 --- a/docs/node-protocol.md +++ b/docs/node-protocol.md @@ -12,7 +12,7 @@ | `echo` | `text` / `file_uri` | `text`、`file_uri` | 示例节点,验证协议链路 | | `ffmpeg-extract` | `video_uri` | `audio_uri`(WAV) | 参数:`sample_rate`、`channels` | | `faster-whisper` | `audio_uri`(16kHz 单声道) | `srt_uri` | 参数:`language`、`task`、`model_path`、`device`、`compute_type`、`beam_size`、`vad_filter`(默认开)、`condition_on_previous_text`、`chunk_seconds` | -| `llm-translate` | `srt_uri` | `cn_srt_uri` | 参数:`target_language`、`model` | +| `llm-translate` | `srt_uri` | `cn_srt_uri` | 参数:`target_language`、`model`、`unload_after`。`unload_after` 控制节点结束后是否卸载本地模型释放显存(默认:端点在本机时卸载,云端不卸载)。支持**节点内暂停**:run 根目录有 `paused.flag` 时翻译在批边界(20 行/批)中止并返回 failed,调度器保持 PAUSED;批量分块流水线期间引擎写 `keep_model.flag`,此时不卸载模型(阶段结束由引擎统一释放) | | `vlm-ocr` | `image_uri` | `text`、`text_uri` | 直接调本地 Ollama 多模态模型(glm-ocr)的 `/api/chat` 做视频帧 OCR(流式 + 5s 上限),参数:`model`、`ollama_host`、`prompt`、`timeout_seconds`、`keep_alive`、`num_predict`、`temperature`、`repeat_penalty` | | `frame-extract` | `video_uri` | `frames_manifest`、`frame_count` | 按**帧间隔**抽帧(解析 fps → step=round(间隔秒×fps),ffmpeg select 按帧号精确取帧,帧时间=帧号/fps 无累计偏差)并 crop 裁切字幕区域,参数:`interval_seconds`(默认 0.5)、`crop`([x,y,w,h] 0~1,**默认画面底部 1/4** `[0,0.75,1,0.25]`——字幕很少出现在画面上半部分,2026-08 调整)。**帧文件必须按帧号数值排序读取**(`_sorted_frame_files`):ffmpeg `%04d` 编号超过 9999 帧后扩为 5 位,字典序 `sorted()` 会把 5 位编号排在 4 位之前导致时间与图像错位(真实发生于 run_339ec7ee437f 的 14236 帧任务,回归测试见 `test_frame_files_read_order_matches_frame_number`) | | `subtitle-ocr` | `frames_manifest` | `srt_uri`、`count` | 自适应线程池并发逐帧调 vlm-ocr → 垃圾过滤(无文字帧)→ 相同字幕合并(记录最后可见帧)→ 组装 SRT,消失时间=最后可见帧+采样间隔(间隔从帧清单推导),参数:`min_chars`、`min_alnum_ratio`、`garbage_tokens`、`max_result_chars`、`pool_min_workers`/`pool_max_workers`/`pool_window_seconds`/`pool_fast_threshold`/`pool_slow_threshold` | diff --git a/docs/operations.md b/docs/operations.md index f440cc5..14d0af6 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -88,11 +88,39 @@ (不再放视频同名文件夹),与用户视频库天然隔离;暂停/失败的视频保留工作空间 以便断点续跑。per-video 的 `WorkflowScheduler` 实例以该目录为 storage—— 完整复用 DAG 拓扑执行、产物表登记与**断点续跑**逻辑。 +- **任务级工作空间清理**:任务全部完成且**无失败视频**时,连任务级目录 + `storage/batch//` 一并删除(每个视频的工作空间已在收尾时各自删完, + 任务级目录只剩空壳);有失败视频时保留(它们的中间产物供断点重试)。 + `paused.flag`/`keep_model.flag` 在每个阶段开始前清理,避免进程被强杀后的 + 残留影响后续阶段。 +- **任务管理与批量页的分工**:批量 run(`source=batch`)**默认不出现在任务管理页** + (`GET /api/runs` 默认排除,排查时用 `?include_batch=1`)——一个批量任务会产生 + N 条单视频 run,混进 20 条窗口会把用户自己提交的任务挤出去,且任务管理页的 + 暂停/继续/重试/删除对批量 run 语义不成立(暂停会被引擎下一次断点续跑静默复位, + 删除被 422 拒绝)。作为补偿,批量页详情表补**阶段**列:`阶段 2/4 · 转写`,由该视频 + run 的 `current_node_id` 在 DAG 拓扑序中的位置推导、节点类型映射中文标签。 + 粒度限制:`progress` 只有**节点边界**粒度,句级进度(转写分块、翻译批次、OCR 帧) + 不落库、只在控制台日志里。 +- **分块流水线执行(本地模型只加载一次)**:批量引擎把待处理视频按 + `WOV_BATCH_STAGE_GROUP_SIZE`(默认 8)分组,**组内按节点顺序跑完全部视频** + (先全部 extract、再全部 ASR、再全部 LLM 翻译、最后 ASS)再进入下一组。每个 + 视频的 run 在阶段边界保持 RUNNING(`execute_run(stop_after=节点)`),下一阶段 + 从产物表跳过已完成节点继续,因此本地模型每组只加载一次、卸载一次,而不是每个 + 视频来回加载卸载;产物仍按组增量落地。LLM 阶段执行时引擎在 run 根目录写 + `keep_model.flag`,节点据此不在每次调用后卸载模型(`nodes/llm.py`),阶段 + 结束由引擎调 `release_local_model()` 统一释放显存,让下一组的 ASR 拿到 GPU + (否则本地模型常驻显存会让 whisper 直接 CUDA OOM)。设为 1 即回到「每个视频 + 跑完整链路」的旧行为。 +- **失败视频不跨阶段推进**:某阶段失败的视频只在**它失败节点的那个阶段**重试 + (下一次引擎循环从断点续跑),不会在后续阶段里重跑前序节点——避免本地 LLM 已 + 常驻时重跑 ASR 抢显存;视频仍按既有语义记 FAILED,任务在没有其他待处理视频时 + 以 failed>0 收尾。 - **暂停/继续**:`POST /api/batch/jobs/{id}/pause` 把任务置 PAUSED 并暂停当前 - run(写 `paused.flag`;whisper **分块间**检查、OCR 逐帧检查后中止,当前节点 - 执行完才停);`resume` 恢复 QUEUED,引擎从断点继续——PAUSED 视频的 run 显式 - resume 后从产物表续跑,未开始的视频接着处理。重启进程后 RUNNING 残留 run 由 - `recover_interrupted_runs` 恢复,暂停的继续处理。 + run(写 `paused.flag`;whisper **分块间**检查、OCR 逐帧检查、llm-translate + **按批(20 行)**检查后中止,当前节点执行完才停);`resume` 恢复 QUEUED,引擎 + 从断点继续——PAUSED 视频的 run 显式 resume 后从产物表续跑,未开始的视频接着 + 处理。重启进程后 RUNNING 残留 run 由 `recover_interrupted_runs` 恢复,暂停的 + 继续处理。 - **失败容错**:单个视频失败(节点失败/文件缺失)记为 FAILED,批量任务继续 处理后续视频,结束后统计 done/failed;DAG 解析/任务级异常把任务置 FAILED。 **重跑保留产物**(2026-08,修复 run_e2b74e89e232 实测):FAILED 视频重新处理 @@ -100,20 +128,29 @@ 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 +- **完成任务判定**:`_run_job` 置 COMPLETED 前**校验全部非 SKIPPED 视频都已 + 结束**(无 PENDING/PAUSED 残留),否则**把任务置回 QUEUED** 交引擎下一轮续跑: + 留 RUNNING 是错的——`next_queued_batch_job` 只拾取 QUEUED,任务会停在“运行中 + 但没人推进”。这条路径主要出现在**任务创建与明细写入的竞态**:任务行先于视频 + 明细写入(`create_job` 逐条插入),引擎可能在登记完成前就拾起任务,本轮只看到 + 已写入的那部分视频(实测 batch_959e510259f6:524 条明细中只看到最初 4 个非 + SKIPPED 视频),剩下的留到下一轮;已登记明细全部尚未写入时(一条明细都没有) + 同样保持 QUEUED,不能按空任务收尾。旧行为留下“N 个 PENDING 待处理却已完成” + 的假完成(batch_969fabe74b83 等 3 个任务实测),或停在运行中无人推进。 + **崩溃恢复**:重启时除 `recover_interrupted_runs` 外, + `recover_interrupted_batch_jobs` 把 RUNNING 的批量任务恢复为 QUEUED (否则停在 RUNNING 的批量任务永远不会被 `next_queued_batch_job` 再次拾起, - 未处理完的 PENDING 永久残留)。历史僵尸数据修复脚本见 - `scripts/fix_zombie_batch_jobs.py`(把误标 COMPLETED 的任务置回 QUEUED 续跑)。 + 未处理完的 PENDING 永久残留);**同一恢复也会把被提前标记 COMPLETED 但仍有 + 未结束视频的僵尸任务置回 QUEUED**(完成标记先于视频收尾写出的旧数据, + batch_351833b7d446 实测:COMPLETED/done=0 却仍有 1 个 PENDING), + 否则只靠 `fix_zombie_batch_jobs.py` 手动修数据,重启也不会自动诊好。 - **产物下载**:`GET /api/batch/jobs/{id}/videos/{vid}/download?alias=<文件名>` 解析并返回视频旁的字幕文件;旧版 `batch.done.json` 完成标记里的语义别名 (位于旧 work_dir)仍兼容可下载。详情/创建响应里每个视频的 `finals` 合并上述 两处来源。 +- **详情明细展示**:任务列表的“详情”只列出本批实际处理过的视频行,SKIPPED + (视频旁已有字幕、创建时即被跳过)不出现在明细表里;整批都已被跳过时提示 + “无待处理视频”。 - **孤儿清理保护**:`source=batch` 的运行**跳过**自动清理——其 run 位于私有 `storage/batch/...` 下,普通孤儿逻辑会误判删除,且 `_remove_run` 还会删除 `input_uri` 的父目录(用户的整个视频文件夹)。详见 diff --git a/docs/testing.md b/docs/testing.md index c3e3ae6..6468eb7 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -49,6 +49,8 @@ tests/ │ └── test_routers/ # 三组 API(apps / workflows / batch) ├── sdk/test_models/ # 对应 src/wov_sdk/(协议数据模型) ├── web/test_crop/ # 对应 web/assets/(框选几何换算) +├── web/test_batch/ # 对应 web/assets/(批量页渲染) +├── scripts/test_fix_zombie_batch_jobs/ # 对应 scripts/(僵尸批量任务修复) └── shared/ # 跨模块公共设施 ├── realdata_contract.py # 真实数据契约与对齐量化 ├── srt_entries.py # 按秒解析 SRT @@ -157,6 +159,8 @@ tests/ | `wov_app/schemas.py` | 请求模型(Pydantic) | `tests/app/test_main/`(schema 用例) | 已覆盖 | | `wov_sdk/models.py` | 协议数据模型 | `tests/sdk/test_models/` | 已覆盖 | | `web/assets/crop.js` | 框选几何换算 | `tests/web/test_crop/`(真实 node 执行) | 已覆盖 | +| `web/assets/batch.js` | 批量页明细/进度渲染 | `tests/web/test_batch/`(真实 node 执行) | 已覆盖 | +| `scripts/fix_zombie_batch_jobs.py` | 僵尸批量任务诊断与修复 | `tests/scripts/test_fix_zombie_batch_jobs/` | 已覆盖 | | `tests/shared/srt_entries.py` | 按秒解析 SRT(测试公共设施) | `tests/shared/test_srt_entries/` | 已覆盖 | | `tests/shared/realdata_contract.py` | 真实数据契约与对齐量化 | `tests/shared/test_alignment/` | 已覆盖 | | `tests/shared/env_isolation.py` | 环境/临时目录隔离 | 被 `tests/app/test_config` 等间接覆盖 | 已覆盖(间接) | diff --git a/docs/代码审查问题跟踪.md b/docs/代码审查问题跟踪.md index 415a48b..defa05b 100644 --- a/docs/代码审查问题跟踪.md +++ b/docs/代码审查问题跟踪.md @@ -99,3 +99,12 @@ - 修复前全套测试:`uv run pytest`,369 passed、6 skipped,76.75 秒。 - 隔离复现已确认:源视频目录误删、限流后 1 → 19 并发、队列积压时缩容滞后、暂停恢复后成品 URI 失效、环形 DAG 阻塞队首、多行 SRT 损坏、OCR 跨空白合并、重用抽帧目录留下旧尾帧。 - 本文件中的“已修复”只表示当前工作区实现及验证完成;部署状态需另行记录。 + +## 僵尸批量任务自动恢复(2026-09) + +- 现象:`batch_351833b7d446` 状态为 COMPLETED(done=0/failed=0),明细里仍有 1 个 PENDING 视频未处理;用户看到"已完成"却什么都没做。 +- 根因:完成标记先于视频收尾写出——旧版 `_run_job` 遍历结束后无条件把任务置 COMPLETED,而"置完成前校验无未结束明细"的修复(`leftovers` 检查)只对新记录生效;已落库的僵尸数据不会被自动纠正,因为 `next_queued_batch_job` 只拾取 QUEUED。 +- 影响面:全库仅此 1 条;其余任务的非终态明细为空。 +- 修复:[db.py](../src/wov_app/db.py) 的 `recover_interrupted_batch_jobs` 在恢复 RUNNING 任务之外,同时把"COMPLETED 且存在 PENDING/RUNNING/PAUSED 明细"的任务置回 QUEUED;启动时即执行,引擎随后从断点续跑。[fix_zombie_batch_jobs.py](../scripts/fix_zombie_batch_jobs.py) 从硬编码 job_id 列表改为动态扫描同类僵尸任务,供无需重启时手动修复。 +- 验证:README 与运维文档已同步;TDD 红为 `tests/app/test_db/test_database.py::test_recover_interrupted_batch_jobs_requeues_zombie_completed`(恢复数 0)与 `tests/app/test_batch/test_batch.py::test_worker_processes_recovered_zombie_job`(任务停在 COMPLETED 且视频未处理);绿为 `uv run pytest tests/app/test_db tests/app/test_batch tests/scripts tests/web -q`,57 passed。 +- 状态:修复及验证完成;实际数据由运行中的服务在改动落盘后重启、启动恢复时自动纠正,视频已重新进入 asr 节点处理。 diff --git a/nodes/llm.py b/nodes/llm.py index 75c5df4..d3f6d14 100755 --- a/nodes/llm.py +++ b/nodes/llm.py @@ -6,6 +6,9 @@ 本地按 cue 回填,避免模型重排断句时译文贴错时间轴。 - **提示词**:要求逐行独立翻译、碎片句按语境独立成行、禁止合并或拆分。 - **严格错误处理**:结构重试耗尽立即失败,不用补空或合并掩盖对应关系丢失。 +- **显存让渡**:LLM 节点结束前卸载本机 Ollama 模型(`unload_after` / `LLM_UNLOAD_AFTER` + 可显式控制;本机端点默认卸载),避免常驻显存与后续 whisper 转写争抢—— + Ollama 默认常驻数分钟,下一个视频的 ASR 会直接 CUDA OOM。 - 拼接 system_prompt 时用 `+` 显式连成单个字符串:括号内的隐式字符串拼接 遇到 f-string 表达式会失效,生成 tuple 后序列化成数组,API 会返回 400。 """ @@ -16,8 +19,10 @@ import json import os import time import urllib.error +import urllib.parse import urllib.request from pathlib import Path +from typing import Callable from wov_app.logging import get_logger @@ -31,9 +36,33 @@ CHUNK_SIZE = 20 # 批次翻译最大尝试次数(ID/正文结构校验失败时重发本批,不用占位恢复)。 MAX_BATCH_RETRIES = 3 +# 卸载本地模型是收尾动作,超时上限固定 30s:翻译超时(LLM_TIMEOUT_SECONDS, +# 默认 600)不适合它,否则卡住的端点会把节点拖住十分钟。 +UNLOAD_TIMEOUT_SECONDS = 30.0 + # 节点运行日志:翻译分批进度与处理速度输出到主进程控制台。 logger = get_logger("llm-translate") +# 暂停信号文件名:位于 run 根目录(/runs//paused.flag),与 +# whisper/subtitle-ocr 约定一致;翻译按批检查,暂停粒度不超过一批(20 行)。 +PAUSE_FLAG = "paused.flag" + +# 保持模型常驻信号文件名:批量分块流水线期间引擎写入 run 根目录,翻译节点据此 +# 不在每次调用后卸载模型(一组视频共用一个已加载模型),阶段结束由引擎统一释放。 +KEEP_MODEL_FLAG = "keep_model.flag" + +# 默认 LLM 端点与模型(与 docs/configuration.md 的环境变量默认值一致)。 +DEFAULT_API_BASE = "https://api.siliconflow.cn/v1/chat/completions" +DEFAULT_MODEL = "Qwen/Qwen3.5-35B-A3B" + + +class PauseRequested(Exception): + """节点内暂停信号:翻译检测到任务被暂停后抛出,由调度器保持 PAUSED。 + + 不把暂停误报为 FAILED:调度器捕获异常时若任务已是 PAUSED 则保持暂停, + 等用户 resume 后整节点重跑(已完成批次不落盘,不留半成品)。 + """ + def _system_prompt(target_language: str) -> str: """构造翻译系统提示词(返回单个字符串,不用隐式拼接避免 tuple bug)。 @@ -131,7 +160,71 @@ def _parse_translations(content: str, expected: set[int]) -> dict[int, str]: return result -def translate_lines(lines: list[str], params: dict) -> list[str]: +def _is_local_endpoint(api_base: str) -> bool: + """判断 LLM 端点是否在本机(loopback),决定是否需要默认卸载显存。""" + return (urllib.parse.urlsplit(api_base).hostname or "").lower() in ( + "localhost", "127.0.0.1", "::1", + ) + + +def _should_unload(params: dict, api_base: str, keep_model_loaded: bool = False) -> bool: + """判断节点结束时是否卸载模型:参数 > 环境变量 LLM_UNLOAD_AFTER > 本机端点默认卸载。 + + 本机(loopback)跑模型时显存是本机共用的,翻译结束后默认让出,后面还要跑 + whisper 的 ASR;云端/别的机器上的端点不占本机显存,默认不发多余请求。 + keep_model_loaded(引擎写入的保持常驻信号)为真时一律不卸载:分块流水线里 + 同一组视频共用一个已加载模型,阶段结束由引擎调用 release_local_model 释放。 + """ + if keep_model_loaded: + return False + flag = params.get("unload_after", os.getenv("LLM_UNLOAD_AFTER")) + if flag is None or str(flag).strip() == "": + return _is_local_endpoint(api_base) + return str(flag).strip().lower() in ("1", "true", "yes", "on") + + +def release_local_model(model: str | None = None) -> None: + """按当前配置卸载本机 LLM 模型释放显存(批量分阶段执行时由引擎在阶段末调用)。 + + 只对本机端点生效:云端/远端端点不占本机显存,不发无意义请求。 + """ + api_base = os.getenv("LLM_API_BASE", DEFAULT_API_BASE) + if not _is_local_endpoint(api_base): + return + _unload_local_model(api_base, str(model or os.getenv("LLM_MODEL", DEFAULT_MODEL))) + + +def _unload_local_model(api_base: str, model: str) -> None: + """请求 Ollama 卸载模型释放显存(keep_alive=0);失败只记录,不影响翻译。 + + 批量链路里 translate 是最后一个占显存的节点,之后下一个视频要跑 whisper; + Ollama 默认让模型常驻数分钟,与 ASR 抢显存会直接 CUDA OOM,所以节点结束 + 时显式释放。云端端点没有该路径,请求失败视为不支持卸载即可。 + """ + parts = urllib.parse.urlsplit(api_base) + origin = f"{parts.scheme}://{parts.netloc}" + request = urllib.request.Request( + origin + "/api/generate", + data=json.dumps({"model": model, "keep_alive": 0}).encode("utf-8"), + headers={"Content-Type": "application/json"}, + method="POST", + ) + try: + with urllib.request.urlopen(request, timeout=UNLOAD_TIMEOUT_SECONDS): + pass + except (OSError, ValueError) as exc: + # HTTPError/URLError 都是 OSError 子类;不支持卸载的端点走到这里。 + logger.warning("本地模型卸载失败(%s): %s", origin, exc) + else: + logger.info("已卸载本地模型 %s,显存让给后续节点", model) + + +def translate_lines( + lines: list[str], + params: dict, + stop_requested: Callable[[], bool] | None = None, + keep_model_loaded: bool = False, +) -> list[str]: """分批调用 LLM 翻译纯文本行,返回顺序一致的译文列表。 列表的每项是一条 cue 正文(可多行);每批按全局 ID 对齐,空 cue 原样 @@ -141,16 +234,20 @@ def translate_lines(lines: list[str], params: dict) -> list[str]: 日志:每完成一批打印总进度(已完成行数/总行数、第几批/共几批、累计 耗时与行处理速度),结束打印汇总(总耗时、累计 tokens 与 tok/s), 便于评估 LLM 处理速度。 + + stop_requested 返回 True 时抛 PauseRequested 中止(调度器保持 PAUSED); + keep_model_loaded 为真时不卸载模型(批量分阶段执行,阶段结束由引擎释放)。 """ api_base = os.getenv( "LLM_API_BASE", - "https://api.siliconflow.cn/v1/chat/completions", + DEFAULT_API_BASE, ) api_key = os.getenv("LLM_API_KEY", "") request_timeout = float(os.getenv("LLM_TIMEOUT_SECONDS", "600")) - model = str(params.get("model") or os.getenv("LLM_MODEL", "Qwen/Qwen3.5-35B-A3B")) + model = str(params.get("model") or os.getenv("LLM_MODEL", DEFAULT_MODEL)) target_language = str(params.get("target_language", "zh-CN")) system_prompt = _system_prompt(target_language) + unload_after = _should_unload(params, api_base, keep_model_loaded) total_lines = len(lines) total_batches = (total_lines + CHUNK_SIZE - 1) // CHUNK_SIZE if total_lines else 0 @@ -162,27 +259,36 @@ def translate_lines(lines: list[str], params: dict) -> list[str]: translated: list[str] = [] total_tokens = 0 all_started = time.monotonic() - for batch_index in range(1, total_batches + 1): - start = (batch_index - 1) * CHUNK_SIZE - chunk = lines[start : start + CHUNK_SIZE] - # 每批日志前缀(第几批/共几批),供单次 LLM 请求日志与批进度复用。 - log_prefix = f"第 {batch_index}/{total_batches} 批" - batch_started = time.monotonic() - batch_translated, batch_tokens = _translate_batch( - chunk, api_base, api_key, model, system_prompt, request_timeout, log_prefix, - start_id=start + 1, - ) - translated.extend(batch_translated) - total_tokens += batch_tokens - # 批进度日志:已完成行数/总行数、当前批耗时、累计耗时与行处理速度。 - done = len(translated) - elapsed_total = time.monotonic() - all_started - logger.info( - "翻译进度 %d/%d 行 (%s完成, 批耗时 %.1fs, 累计 %.1fs, %.1f 行/s)", - done, total_lines, log_prefix, - time.monotonic() - batch_started, elapsed_total, - done / elapsed_total if elapsed_total > 0 else 0.0, - ) + try: + for batch_index in range(1, total_batches + 1): + # 暂停检查:调度器置 PAUSED 并写 paused.flag 后,翻译在批边界立刻停下, + # 已完成的批保留在内存中但不落盘,恢复时整节点重跑,不留半成品。 + if stop_requested is not None and stop_requested(): + raise PauseRequested("翻译被暂停") + start = (batch_index - 1) * CHUNK_SIZE + chunk = lines[start : start + CHUNK_SIZE] + # 每批日志前缀(第几批/共几批),供单次 LLM 请求日志与批进度复用。 + log_prefix = f"第 {batch_index}/{total_batches} 批" + batch_started = time.monotonic() + batch_translated, batch_tokens = _translate_batch( + chunk, api_base, api_key, model, system_prompt, request_timeout, log_prefix, + start_id=start + 1, + ) + translated.extend(batch_translated) + total_tokens += batch_tokens + # 批进度日志:已完成行数/总行数、当前批耗时、累计耗时与行处理速度。 + done = len(translated) + elapsed_total = time.monotonic() - all_started + logger.info( + "翻译进度 %d/%d 行 (%s完成, 批耗时 %.1fs, 累计 %.1fs, %.1f 行/s)", + done, total_lines, log_prefix, + time.monotonic() - batch_started, elapsed_total, + done / elapsed_total if elapsed_total > 0 else 0.0, + ) + finally: + # 成功与失败都在此释放显存:下一个视频的 ASR 需要独占 GPU。 + if unload_after: + _unload_local_model(api_base, model) # 任务汇总日志:总耗时、累计 tokens 与 token/行处理速度。 wall = time.monotonic() - all_started tok_rate = total_tokens / wall if wall > 0 and total_tokens > 0 else 0.0 @@ -257,9 +363,19 @@ def invoke(request: InvokeRequest) -> InvokeResponse: try: entries = parse_srt(srt_path.read_text(encoding="utf-8")) - translated_lines = translate_lines([entry.text for entry in entries], request.params) + # run 根目录 = /runs//;引擎在分块流水线里会写 + # keep_model.flag(阶段内保持模型常驻),暂停接口写 paused.flag。 + run_root = Path(request.output_dir).parent.parent + translated_lines = translate_lines( + [entry.text for entry in entries], + request.params, + stop_requested=(run_root / PAUSE_FLAG).exists, + keep_model_loaded=(run_root / KEEP_MODEL_FLAG).exists(), + ) if len(translated_lines) != len(entries): raise ValueError("translation count does not match subtitle cues") + except PauseRequested as exc: + return InvokeResponse(status="failed", error=f"{exc}(run {request.run_id})") except (ValueError, TypeError, OSError) as exc: return InvokeResponse(status="failed", error=str(exc)) # 时间轴始终来自原始 cue,译文通过已校验的 ID 顺序回填。 diff --git a/scripts/fix_zombie_batch_jobs.py b/scripts/fix_zombie_batch_jobs.py index 2104645..9128803 100644 --- a/scripts/fix_zombie_batch_jobs.py +++ b/scripts/fix_zombie_batch_jobs.py @@ -1,19 +1,16 @@ -"""修复僵尸批量任务脚本:把误标 COMPLETED 但仍有 PENDING 视频的任务置回 QUEUED。 +"""修复僵尸批量任务脚本:把误标 COMPLETED 但仍有未结束视频的任务置回 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 完成) +僵尸状态:`_run_job` 在视频收尾前就把任务置 COMPLETED,留下"N 个视频待处理 +却已完成"的假完成;引擎只拾取 QUEUED,剩余视频永久无人处理。代码已修复 +(置 COMPLETED 前校验无未结束明细 + 重启恢复把这类任务放回队列),本脚本用于 +修复**历史遗留**数据(如 batch_351833b7d446:COMPLETED/done=0 但仍有 1 个 PENDING)。 -修复方式:仅把这 3 个任务置回 QUEUED 并清空误导的 total/progress,保留 -SKIPPED/COMPLETED 明细与 run 记录;批量引擎(新代码)重启后会重新拾起, -从剩余 PENDING 视频续跑,全部处理完才置 COMPLETED。 +修复方式:把 COMPLETED 且仍有非终态(PENDING/RUNNING/PAUSED)明细的任务置回 +QUEUED 并清空误导的 progress,保留 SKIPPED/COMPLETED/FAILED 明细与 run 记录; +批量引擎重启后会重新拾起,从剩余视频续跑,全部处理完才置 COMPLETED。 -安全约束:只更新明确列出的 3 个 job_id;其余任务(含正常 COMPLETED 的 -fee6681/ac8ca4f)不动。执行前打印将变更的任务与明细统计供确认。 +安全约束:只改"COMPLETED 且存在未结束明细"的任务,正常完成的任务不动。 +默认只打印预览,确认后加 --apply 实际写入。 """ from __future__ import annotations @@ -23,12 +20,8 @@ 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 -] +# 未结束的明细状态:出现任一个就不能算任务完成。 +UNFINISHED_STATUSES = ("PENDING", "RUNNING", "PAUSED") def _now_iso() -> str: @@ -36,30 +29,44 @@ def _now_iso() -> str: return datetime.now(timezone.utc).isoformat() +def find_zombie_jobs(db: Database) -> list[tuple[dict, list[dict]]]: + """返回全部僵尸任务及其未结束明细(COMPLETED 但仍有未结束视频)。""" + zombies: list[tuple[dict, list[dict]]] = [] + for job in db.list_batch_jobs(limit=1000): + if job["status"] != "COMPLETED": + continue + unfinished = [ + v for v in db.list_batch_videos(str(job["id"])) + if v["status"] in UNFINISHED_STATUSES + ] + if unfinished: + zombies.append((job, unfinished)) + return zombies + + 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 + zombies = find_zombie_jobs(db) + if not zombies: + print(" [无] 没有 COMPLETED 但仍含未结束视频的批量任务") + return + for job, unfinished in zombies: + job_id = str(job["id"]) + videos = db.list_batch_videos(job_id) + completed = sum(1 for v in videos if v["status"] == "COMPLETED") print(f" [待修] {job_id}: status={job['status']} → QUEUED, " - f"COMPLETED={completed}, PENDING={pending}") + f"COMPLETED={completed}, 未结束={len(unfinished)}" + f"({', '.join(Path(v['video_path']).name for v in unfinished[:3])}…)") if apply: db.update_batch_job( - job_id, status="QUEUED", progress=0, total=pending, + job_id, status="QUEUED", progress=0, + total=int(job["total"] or 0) or len(unfinished), done=completed, failed=int(job["failed"] or 0), current_video=None, error=None, updated_at=_now_iso(), ) - print(f" ✓ 已置回 QUEUED(引擎将续跑剩余 {pending} 个 PENDING)") + print(f" ✓ 已置回 QUEUED(引擎将续跑剩余 {len(unfinished)} 个视频)") if not apply: print("\n以上为预览。确认无误后加 --apply 实际修复。") diff --git a/src/wov_app/batch.py b/src/wov_app/batch.py index 0b792ff..b385d8e 100644 --- a/src/wov_app/batch.py +++ b/src/wov_app/batch.py @@ -11,6 +11,12 @@ 字幕文件(`.srt/.ass/.ssa/.vtt`),说明该视频已有字幕,直接记为 SKIPPED, 不为它触发任何流水线。运行时(BatchWorker)只消费已定位好的明细列表, **不再重新扫描文件夹**(运行期间新增/删除的视频不会改变本次任务的范围)。 +- **分块流水线执行**:视频按 `WOV_BATCH_STAGE_GROUP_SIZE` 分组,组内按节点 + 顺序跑完全部视频(先全部 extract、再全部 ASR、再全部 LLM 翻译、最后 ASS) + 再进入下一组——本地模型每组只加载一次、卸载一次,产物按组增量落地。 + 阶段边界用 `execute_run(stop_after=节点)` 停在节点(任务保持 RUNNING), + LLM 阶段靠 `keep_model.flag` 让节点保持模型常驻,阶段结束由引擎统一释放 + 显存(详见 docs/operations.md#文件夹批量处理)。 - **产物放在视频旁**:每个视频处理完成后,把工作流 `final_outputs` 对应的 最终产物文件(字幕流水线即中文 `.srt` 与双目 `.ass`)**复制一份到视频的 所在目录**,与 .mp4 放在一起;文件名**对齐媒体库既有约定**:中文字幕存为 @@ -44,12 +50,14 @@ from datetime import datetime, timezone from pathlib import Path from wov_app import registry -from wov_app.config import BATCH_INTERVAL_SECONDS, STORAGE_DIR +from wov_app.config import BATCH_INTERVAL_SECONDS, BATCH_STAGE_GROUP_SIZE, STORAGE_DIR from wov_app.db import Database from wov_app.logging import get_logger -from wov_app.scheduler import WorkflowScheduler +from wov_app.scheduler import WorkflowScheduler, topological_sort from wov_app.storage import atomic_copy -from wov_sdk.models import WorkflowDefinition +from wov_sdk.models import WorkflowDefinition, WorkflowNode + +from nodes.llm import release_local_model # 批量引擎运行日志:任务进度、视频逐个处理与暂停/续跑等状态变化。 logger = get_logger("batch") @@ -67,6 +75,13 @@ SUBTITLE_EXTENSIONS = {".srt", ".ass", ".ssa", ".vtt"} # 暂停信号文件名:与节点约定一致,位于 run 根目录(/runs//)。 PAUSE_FLAG = "paused.flag" +# 保持模型常驻信号文件名:LLM 阶段执行期间由引擎写入 run 根目录,节点据此不在 +# 每次调用后卸载模型(同组视频共用一份已加载模型,减少加载/卸载次数)。 +KEEP_MODEL_FLAG = "keep_model.flag" + +# LLM 节点类型前缀:这类节点加载本地大模型,阶段结束后由引擎统一释放显存。 +LLM_NODE_PREFIX = "llm" + # 兼容读取的历史完成标记文件名(旧任务用它记录产物路径)。当前逻辑不再 # 写入,产物直接放视频旁;保留读取能力以便旧任务的详情/下载仍可用。 MARKER_NAME = "batch.done.json" @@ -324,11 +339,12 @@ class BatchWorker: self.db.update_batch_job(job_id, status="FAILED", error=str(exc), updated_at=_now_iso()) def _run_job(self, job_id: str) -> None: - """批量任务主流程(内部实现,异常由 _process_job 统一处理)。 + """批量任务主流程(分块流水线,异常由 _process_job 统一处理)。 - 只消费创建任务时已定位好的 batch_videos 明细:SKIPPED/COMPLETED 直接 - 跳过,PENDING(含失败/暂停后恢复的)逐个交给 _process_video 处理, - **不再扫描文件夹**补视频。 + 视频按 WOV_BATCH_STAGE_GROUP_SIZE 分组,组内按 DAG 拓扑顺序逐节点跑完 + 全部视频(先全部 extract、再全部 ASR、再全部 LLM 翻译、最后 ASS)再 + 进入下一组:本地模型每组只加载一次、卸载一次,产物按组增量落地。 + 只消费创建任务时已定位好的 batch_videos 明细,**不再扫描文件夹**。 """ job = self.db.get_batch_job(job_id) if job is None: @@ -348,8 +364,21 @@ class BatchWorker: return definition = WorkflowDefinition.from_dict(version["definition"]) definition.validate() + order = topological_sort(definition) + node_by_id = {node.id: node for node in definition.nodes} - items = self.db.list_batch_videos(job_id) + # 任务行先于视频明细写入(create_job 逐条插入),引擎可能在登记完成前就拾起 + # 任务:此时没有任何明细,不能按“空任务”收尾,保持 QUEUED 等登记完成。 + if not self.db.list_batch_videos(job_id): + logger.info("批量任务 %s 尚无视频明细(仍在登记),保持 QUEUED 稍后重试", job_id) + self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso()) + return + + # 待处理明细:SKIPPED 创建时已定(不参与 total/done),COMPLETED 无需重跑。 + items = [ + item for item in self.db.list_batch_videos(job_id) + if item["status"] not in ("COMPLETED", "SKIPPED") + ] # total 创建时已固定为"无字幕需处理的视频数",此处不覆盖;历史任务 # 的旧口径由 sync_batch_job_progress 在读取时修正为不含 SKIPPED。 self.db.update_batch_job( @@ -359,49 +388,45 @@ class BatchWorker: # total 为进度条分母;为 0 表示整批跳过(创建即 COMPLETED)。 total = int(job["total"] or 0) - for item in items: - # 暂停检查:批量任务被暂停后停止处理后续视频,等待用户继续。 - current = self.db.get_batch_job(job_id) - if current is None or current["status"] == "PAUSED": - # 停下前先把已完成/失败项入账,让暂停中的前端看到真实进度。 - self.db.sync_batch_job_progress(job_id) - logger.info("批量任务 %s 已暂停,停止在视频 %s", job_id, item["video_path"]) - return - - # 已完成/已跳过的视频不再处理:COMPLETED 由断点续跑逻辑跳过, - # SKIPPED 在创建任务时已定(不参与 total/done,故无需同步进度)。 - if item["status"] in ("COMPLETED", "SKIPPED"): - continue - - video = Path(item["video_path"]) - if not video.is_file(): - self.db.update_batch_video(item["id"], status="FAILED", error="video file not found", updated_at=_now_iso()) - self.db.sync_batch_job_progress(job_id) - continue - - work_dir = Path(item["work_dir"]) - self.db.update_batch_job( - job_id, current_video=str(video), updated_at=_now_iso(), - ) - try: - self._process_video(job, item, version, definition, work_dir) - except Exception as exc: # noqa: BLE001 - # 单视频兜底:不中断整个批量任务,记录错误后继续下一个视频。 - logger.exception("批量任务 %s 视频 %s 处理异常", job_id, video) - self.db.update_batch_video(item["id"], status="FAILED", error=str(exc), updated_at=_now_iso()) - # 本视频处理完(成功/失败/暂停)后实时同步一次汇总,让进度尽快入账。 - self.db.sync_batch_job_progress(job_id) - - # 重新读取视频明细:_process_video 可能刚创建 run 或已收尾清理 - # (快照里 run_id 可能是旧值),必须取最新记录判断暂停状态。 - item = self.db.get_batch_video(item["id"]) - # 视频处理中被暂停:批量任务整体保持 PAUSED,等待用户继续。 - run = self.db.get_run(item["run_id"]) if item and item.get("run_id") else None - if run is not None and run["status"] == "PAUSED": - self.db.update_batch_video(item["id"], status="PAUSED", updated_at=_now_iso()) - self.db.sync_batch_job_progress(job_id) - self.db.update_batch_job(job_id, status="PAUSED", updated_at=_now_iso()) - return + group_size = max(1, int(BATCH_STAGE_GROUP_SIZE)) + for start in range(0, len(items), group_size): + group = items[start:start + group_size] + for stage_index, node_id in enumerate(order): + node_spec = node_by_id[node_id] + # 末阶段不传 stop_after:让调度器收尾(final_outputs + COMPLETED)。 + is_last_stage = stage_index == len(order) - 1 + executed = False + for item in group: + # 暂停检查:批量任务被暂停后停止处理后续视频,等待用户继续。 + current = self.db.get_batch_job(job_id) + if current is None or current["status"] == "PAUSED": + # 停下前先把已完成/失败项入账,让暂停中的前端看到真实进度。 + self.db.sync_batch_job_progress(job_id) + logger.info("批量任务 %s 已暂停,停止在视频 %s", job_id, item["video_path"]) + return + self.db.update_batch_job(job_id, current_video=str(item["video_path"]), updated_at=_now_iso()) + try: + outcome = self._run_stage( + job, item, version, definition, node_spec, order, is_last_stage, + ) + except Exception as exc: # noqa: BLE001 + # 单视频兜底:不中断整个批量任务,记录错误后继续下一个视频。 + logger.exception( + "批量任务 %s 视频 %s 阶段 %s 处理异常", job_id, item["video_path"], node_id, + ) + self.db.update_batch_video(item["id"], status="FAILED", error=str(exc), updated_at=_now_iso()) + outcome = "FAILED" + executed = executed or outcome is not None + # 每个视频每个阶段后实时同步一次汇总,让进度尽快入账。 + self.db.sync_batch_job_progress(job_id) + # 阶段内被暂停(节点内的 paused.flag):任务保持 PAUSED 等续跑。 + if outcome == "PAUSED": + self.db.update_batch_job(job_id, status="PAUSED", updated_at=_now_iso()) + return + # 阶段收尾:LLM 阶段结束时统一释放本地模型显存,让下一组的 + # whisper(ASR)拿到 GPU,否则下一个视频转写会 CUDA OOM。 + if executed and node_spec.node_type.startswith(LLM_NODE_PREFIX): + self._release_llm_model(node_spec.params) # 先按明细实时对齐汇总(done 不计 SKIPPED),再判断能否收尾。 # 仍有未结束视频时不能标 COMPLETED,否则会出现“还有待处理视频却已完成” @@ -415,11 +440,15 @@ class BatchWorker: if v["status"] not in ("SKIPPED", "COMPLETED", "FAILED") ] if leftovers: - # 有未处理完的视频:保持 RUNNING,由引擎下一轮续跑。 + # 有未处理完的视频(常见于创建任务时明细还在逐条写入,本轮快照没包含 + # 它们):置回 QUEUED 自愈,让引擎下一轮按最新明细重新分组续跑。 + # 留在 RUNNING 不会被引擎再拾起(next_queued_batch_job 只取 QUEUED), + # 任务会停在“运行中但没人推进”的状态。 logger.warning( - "批量任务 %s 仍有 %d 个视频未处理完(%s…),保持 RUNNING 待续跑,不置 COMPLETED", + "批量任务 %s 仍有 %d 个视频未处理完(%s…),置回 QUEUED 待下一轮续跑", job_id, len(leftovers), Path(leftovers[0]["video_path"]).name, ) + self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso()) return done = int(job["done"]) if job else 0 failed = int(job["failed"]) if job else 0 @@ -427,25 +456,84 @@ class BatchWorker: job_id, status="COMPLETED", progress=1.0, current_video=None, error=None, updated_at=_now_iso(), ) + # 无失败视频时每个视频的工作空间已在收尾时删除,任务目录只剩空壳; + # 有失败视频则保留(它们的工作空间供断点重试)。 + if failed == 0: + remove_job_workspace(job_id) logger.info( "批量任务 %s 完成: 待处理 %d 个视频, 完成 %d, 失败 %d", job_id, total, done, failed, ) - def _process_video( + def _run_stage( self, job: dict, item: dict, version: dict, definition: WorkflowDefinition, - work_dir: Path, - ) -> None: - """处理单个视频:建 run(复用现有调度器)执行,成功后收尾清理。 + node_spec: WorkflowNode, + order: list[str], + is_last_stage: bool, + ) -> str | None: + """执行一个视频在一个阶段节点上的工作,返回执行后的视频状态。 - per-video 的 WorkflowScheduler 以该视频的私有工作空间为 storage, - 中间态落在 /runs//steps/ 下;产物表记录全部节点 - 输出,暂停后续跑从产物表重建已完成节点(断点续跑)。视频成功后 - 把最终产物复制到视频旁并删除工作空间(见 _finalize_video)。 + 返回 None 表示本阶段无需执行(视频已完成/已跳过)。per-video 的 + WorkflowScheduler 以该视频的私有工作空间为 storage,产物表记录各节点 + 输出,所以同一视频的后续阶段直接从断点继续(不重跑已完成节点)。 + LLM 阶段会写 keep_model.flag:阶段内保持模型常驻,阶段结束由引擎统一 + 释放(见 _release_llm_model),避免每个视频重新加载/卸载模型。 + """ + fresh = self.db.get_batch_video(item["id"]) + if fresh is None or fresh["status"] in ("COMPLETED", "SKIPPED"): + return None + video = Path(fresh["video_path"]) + if not video.is_file(): + self.db.update_batch_video(fresh["id"], status="FAILED", error="video file not found", updated_at=_now_iso()) + return "FAILED" + work_dir = Path(fresh["work_dir"]) + # 失败节点在当前阶段之前:本轮不再推进(否则会在本地 LLM 已常驻时重跑 + # ASR 抢显存),留待下一次引擎循环从其失败节点重试。必须在复位 run 状态 + # 之前判断,否则 _ensure_run 已把 FAILED 改成 QUEUED、判断会失效。 + previous = self.db.get_run(fresh["run_id"]) if fresh.get("run_id") else None + if previous is not None and previous["status"] == "FAILED": + failed_node = previous["current_node_id"] + failed_index = order.index(failed_node) if failed_node in order else 0 + if failed_index != order.index(node_spec.id): + return "FAILED" + run_id = self._ensure_run(job, fresh, version, work_dir) + run_dir = work_dir / "runs" / run_id + run_dir.mkdir(parents=True, exist_ok=True) + # 清除可能残留的信号(重启/强杀/异常中断后):暂停信号会让本次执行误暂停, + # 保持常驻信号会让后续单独重跑该节点时不再卸载模型。 + (run_dir / PAUSE_FLAG).unlink(missing_ok=True) + (run_dir / KEEP_MODEL_FLAG).unlink(missing_ok=True) + if node_spec.node_type.startswith(LLM_NODE_PREFIX): + (run_dir / KEEP_MODEL_FLAG).write_text("", encoding="utf-8") + try: + scheduler = WorkflowScheduler(self.db, work_dir) + scheduler.execute_run(run_id, stop_after=None if is_last_stage else node_spec.id) + finally: + # 信号只在本阶段有效:残留会让后续单独重跑该节点时也不卸载模型。 + (run_dir / KEEP_MODEL_FLAG).unlink(missing_ok=True) + + run = self.db.get_run(run_id) + if run is None: + # execute_run 期间 run 记录被删除(极端外部操作),直接返回。 + return None + if run["status"] == "COMPLETED": + # 放置最终产物到视频旁并清理过程文件。 + self._finalize_video(fresh, run_id, video, work_dir, definition) + return "COMPLETED" + # FAILED 或 PAUSED:由调用方根据 run 状态更新视频状态与任务状态。 + self.db.update_batch_video(fresh["id"], status=run["status"], error=run.get("error"), updated_at=_now_iso()) + return run["status"] + + def _ensure_run(self, job: dict, item: dict, version: dict, work_dir: Path) -> str: + """确保视频有可执行的 run,返回 run_id。 + + run 缺失时新建(input_uri 指向本地视频,不上传副本);已存在的按断点 + 续跑语义复位:PAUSED/FAILED 显式置 QUEUED(保留产物,只重跑未完成 + 节点);RUNNING 是分阶段执行的上一个阶段或进程被杀的残留,同样置回。 """ video = Path(item["video_path"]) work_dir.mkdir(parents=True, exist_ok=True) @@ -454,8 +542,6 @@ class BatchWorker: # run 记录已不存在(收尾异常删除了 run 但状态未同步):重新新建。 run_id = None if run_id is None: - # 首次处理:创建 source=batch 的运行,input_uri 指向本地视频 - # (不上传副本),由调度器按 DAG 执行。 run_id = f"run_{uuid.uuid4().hex[:12]}" now = _now_iso() self.db.create_run({ @@ -471,38 +557,22 @@ class BatchWorker: "updated_at": now, }) self.db.update_batch_video(item["id"], run_id=run_id, updated_at=_now_iso()) - run = self.db.get_run(run_id) - # 已完成(收尾前中断):直接补做收尾。 - if run["status"] == "COMPLETED": - self._finalize_video(item, run_id, video, work_dir, definition) - return - # 暂停的 run 显式 resume 回 QUEUED,由 execute_run 从产物表断点续跑。 if run["status"] == "PAUSED": + # 暂停的 run 显式 resume 回 QUEUED,由 execute_run 从产物表断点续跑。 self.db.resume_run(run_id, _now_iso()) - elif run["status"] == "FAILED": - # 失败重跑:保留产物记录只置 QUEUED,由 execute_run 跳过已完成 - # 节点、仅重跑失败节点,避免浪费抽帧/OCR 等长耗时成果。 + elif run["status"] in ("FAILED", "RUNNING"): + # 保留产物记录只置 QUEUED:execute_run 跳过已完成节点、只重跑失败节点, + # 避免浪费抽帧/ASR 等长耗时成果。 self.db.update_run(run_id, status="QUEUED", error=None, updated_at=_now_iso()) - elif run["status"] == "RUNNING": - # 上次进程被杀残留:恢复 QUEUED(保留产物)由 execute_run 续跑。 - self.db.update_run(run_id, status="QUEUED", updated_at=_now_iso()) - # 清除可能残留的暂停信号(重启/异常中断后),避免本次执行误暂停。 - (work_dir / "runs" / run_id / PAUSE_FLAG).unlink(missing_ok=True) + return run_id - scheduler = WorkflowScheduler(self.db, work_dir) - scheduler.execute_run(run_id) - - run = self.db.get_run(run_id) - if run is None: - # execute_run 期间 run 记录被删除(极端外部操作),直接返回。 - return - if run["status"] == "COMPLETED": - # 放置最终产物到视频旁并清理过程文件。 - self._finalize_video(item, run_id, video, work_dir, definition) - else: - # FAILED 或 PAUSED:由调用方根据 run 状态更新视频状态与任务状态。 - self.db.update_batch_video(item["id"], status=run["status"], error=run.get("error"), updated_at=_now_iso()) + def _release_llm_model(self, params: dict) -> None: + """LLM 阶段结束释放本机模型显存(分块流水线里每组一次,而非每视频一次)。""" + try: + release_local_model(params.get("model")) + except Exception: # noqa: BLE001 - 释放失败不影响批次推进 + logger.warning("释放本地 LLM 模型失败", exc_info=True) # ------------------------------------------------------------------ # 收尾:产物放置与过程文件清理 diff --git a/src/wov_app/config.py b/src/wov_app/config.py index b2fea8a..4d5fe95 100644 --- a/src/wov_app/config.py +++ b/src/wov_app/config.py @@ -25,6 +25,11 @@ SCHEDULER_INTERVAL_SECONDS = float(os.getenv("WOV_SCHEDULER_INTERVAL_SECONDS", " BATCH_ENABLED = os.getenv("WOV_BATCH_ENABLED", "1") == "1" BATCH_INTERVAL_SECONDS = float(os.getenv("WOV_BATCH_INTERVAL_SECONDS", "1.0")) +# 批量"分块流水线"分组大小:每组视频按节点顺序跑完全部阶段(全部 extract → 全部 +# ASR → 全部翻译 → 全部 ASS)再处理下一组,使本地模型每组只加载一次;产物仍按 +# 组增量落地(详见 docs/operations.md#文件夹批量处理)。 +BATCH_STAGE_GROUP_SIZE = int(os.getenv("WOV_BATCH_STAGE_GROUP_SIZE", "8")) + # 孤儿数据清理器配置:定时扫描并清理无对应文件/记录的死数据。 CLEANUP_ENABLED = os.getenv("WOV_CLEANUP_ENABLED", "1") == "1" # 清理扫描周期(秒),默认每小时一次。 diff --git a/src/wov_app/db.py b/src/wov_app/db.py index ea7421c..331c5e2 100644 --- a/src/wov_app/db.py +++ b/src/wov_app/db.py @@ -300,13 +300,19 @@ class Database: result["param_overrides"] = json.loads(raw) if raw else None return result - def list_runs(self, limit: int = 20) -> list[dict[str, Any]]: - """按创建时间倒序返回最近的运行记录。""" + def list_runs(self, limit: int = 20, include_batch: bool = False) -> list[dict[str, Any]]: + """按创建时间倒序返回最近的运行记录。 + + 默认排除 source=batch:批量 run 是批量任务的单视频明细(一个任务会产生 + N 条),把 20 条窗口占满会把用户自己提交的任务挤出列表;它们由批量页 + 的 `/api/batch/jobs` 展示,需要排查时可显式 include_batch=True。 + """ + sql = "SELECT * FROM workflow_runs" + if not include_batch: + sql += " WHERE source != 'batch'" + sql += " ORDER BY created_at DESC LIMIT ?" with self._connect() as conn: - rows = conn.execute( - "SELECT * FROM workflow_runs ORDER BY created_at DESC LIMIT ?", - (limit,), - ).fetchall() + rows = conn.execute(sql, (limit,)).fetchall() return [self._parse_overrides(row) for row in rows] def list_run_ids(self) -> list[str]: @@ -385,16 +391,28 @@ class Database: return cur.rowcount def recover_interrupted_batch_jobs(self, updated_at: str) -> int: - """重启恢复:把遗留 RUNNING 的批量任务恢复为 QUEUED,返回恢复数量。 + """重启恢复:把没在运行、也永远不会被拾起的批量任务恢复为 QUEUED。 - 批量任务若停在 RUNNING,next_queued_batch_job 只拾取 QUEUED, - 永远不会重新驱动它,未处理完的 PENDING 视频会永久残留;恢复为 - QUEUED 后引擎从断点(剩余视频 + 已恢复的 run)继续。用户主动暂停的 + 两类任务需要恢复:停在 RUNNING 的(进程被杀,next_queued_batch_job + 不拾起;不恢复则剩余 PENDING 视频永久残留),以及被提前标记 COMPLETED + 但仍有未结束视频的僵尸任务(完成标记先于视频收尾写出,用户看到“已完成” + 却还有视频没处理)。恢复为 QUEUED 后引擎从断点续跑,用户主动暂停的 PAUSED 保持不变。 """ with self._connect() as conn: cur = conn.execute( - "UPDATE batch_jobs SET status = 'QUEUED', updated_at = ? WHERE status = 'RUNNING'", + """ + UPDATE batch_jobs SET status = 'QUEUED', updated_at = ? + WHERE status = 'RUNNING' + OR ( + status = 'COMPLETED' + AND EXISTS ( + SELECT 1 FROM batch_videos + WHERE batch_videos.job_id = batch_jobs.id + AND batch_videos.status NOT IN ('COMPLETED', 'FAILED', 'SKIPPED') + ) + ) + """, (updated_at,), ) return cur.rowcount diff --git a/src/wov_app/routers/apps.py b/src/wov_app/routers/apps.py index 6681209..b81ef11 100755 --- a/src/wov_app/routers/apps.py +++ b/src/wov_app/routers/apps.py @@ -12,7 +12,7 @@ import uuid from datetime import datetime, timezone from pathlib import Path -from fastapi import APIRouter, Depends, File, Form, HTTPException, UploadFile +from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, UploadFile from fastapi.responses import FileResponse from wov_app.db import Database @@ -116,9 +116,16 @@ async def create_run( @router.get("/api/runs") -def list_runs(db: Database = Depends(_get_db)) -> list[dict]: - """返回最近的运行记录,供任务管理页展示。""" - return db.list_runs() +def list_runs( + include_batch: bool = Query(False), + db: Database = Depends(_get_db), +) -> list[dict]: + """返回最近的运行记录,供任务管理页展示。 + + 默认排除批量 run(每个批量任务会产生 N 条单视频 run,属于批量页的明细, + 混进来会把列表占满);`?include_batch=1` 可包含它们供排查。 + """ + return db.list_runs(include_batch=include_batch) @router.get("/api/runs/{run_id}") diff --git a/src/wov_app/routers/batch.py b/src/wov_app/routers/batch.py index 38fe3aa..d779841 100644 --- a/src/wov_app/routers/batch.py +++ b/src/wov_app/routers/batch.py @@ -16,6 +16,24 @@ from fastapi.responses import FileResponse from wov_app import batch as batch_engine from wov_app.db import Database from wov_app.schemas import BatchJobCreate +from wov_app.scheduler import topological_sort +from wov_sdk.models import WorkflowDefinition + +# 节点类型 → 阶段中文标签(详情表展示);未登记的类型回退节点 ID。 +_STAGE_LABELS = { + "ffmpeg-extract": "提取音频", + "faster-whisper": "转写", + "llm-translate": "翻译", + "srt-to-dual-eye-ass": "合成字幕", + "frame-extract": "抽帧", + "subtitle-ocr": "OCR 识别", + "vlm-ocr": "帧 OCR", + "llm-filter": "字幕过滤", + "subtitle-correction": "字幕纠错", +} + +# 未开始的视频没有阶段信息,统一用 None 占位(前端渲染为“-”)。 +_NO_STAGE = {"stage_label": None, "stage_index": None, "stage_total": None} router = APIRouter(tags=["batch"]) @@ -50,13 +68,51 @@ def _product_finals(video: dict) -> dict[str, str]: return finals -def _enrich_videos(db: Database, videos: list[dict]) -> list[dict]: - """为每个视频补充最终产物清单(视频旁字幕 + 历史完成标记)。 +def _stage_info(db: Database, run: dict, definitions: dict) -> dict: + """由 run 的当前节点推导视频阶段:第几阶段/共几阶段 + 中文标签。 - finals 形如 {alias: 文件名},前端据此渲染下载链接;未完成的视频没有产物。 + 阶段指 DAG 拓扑序里的节点;`progress` 只是节点边界进度(句级的转写分块、 + 翻译批次进度只在日志里),所以这里能给的是"卡在哪个环节"。任务/版本或 + 节点信息缺失时返回空阶段,不影响详情展示。 """ + node_id = run.get("current_node_id") + if not node_id: + return dict(_NO_STAGE) + key = (str(run["workflow_id"]), int(run["workflow_version"])) + definition = definitions.get(key) + if definition is None: + version = db.get_workflow_version(key[0], key[1]) + if version is None: + return dict(_NO_STAGE) + definition = WorkflowDefinition.from_dict(version["definition"]) + definitions[key] = definition + order = topological_sort(definition) + if node_id not in order: + return dict(_NO_STAGE) + node = next((item for item in definition.nodes if item.id == node_id), None) + return { + "stage_label": _STAGE_LABELS.get(node.node_type if node else "", str(node_id)), + "stage_index": order.index(node_id) + 1, + "stage_total": len(order), + } + + +def _enrich_videos(db: Database, videos: list[dict]) -> list[dict]: + """为每个视频补充最终产物清单与当前阶段。 + + finals 形如 {alias: 文件名},前端据此渲染下载链接;阶段信息来自该视频 + 自己的 run(PENDING/SKIPPED/已完成的任务没有 run,保持空阶段)。 + """ + # 同一任务下的视频共用一个工作流版本,定义只解析一次。 + definitions: dict = {} for video in videos: video["finals"] = _product_finals(video) + stage = dict(_NO_STAGE) + if video.get("status") in ("RUNNING", "PAUSED") and video.get("run_id"): + run = db.get_run(video["run_id"]) + if run is not None: + stage = _stage_info(db, run, definitions) + video.update(stage) return videos @router.post("/api/batch/jobs") diff --git a/src/wov_app/scheduler.py b/src/wov_app/scheduler.py index 9c98f07..fa3dfbe 100755 --- a/src/wov_app/scheduler.py +++ b/src/wov_app/scheduler.py @@ -124,11 +124,19 @@ class WorkflowScheduler: return None return outputs_by_node.get(node_id, {}).get(key) - def execute_run(self, run_id: str) -> None: - """执行单个任务:加载 DAG、按拓扑顺序调用节点并登记产物。""" + def execute_run(self, run_id: str, stop_after: str | None = None) -> None: + """执行单个任务:加载 DAG、按拓扑顺序调用节点并登记产物。 + + stop_after 指定"只执行到该节点"(批量分块流水线的阶段执行):该节点完成 + 后任务保持 RUNNING 不收尾,下一次调用从产物表跳过已完成节点继续后面的 + 阶段;不传时执行整条 DAG 并收尾(登记 final_outputs、标 COMPLETED)。 + """ run = self.db.get_run(run_id) - # 任务不存在或不在可执行状态(排队/暂停)时直接返回,避免重复执行。 - if run is None or run["status"] not in ("QUEUED", "PAUSED"): + # 任务不存在或不在可执行状态时直接返回,避免重复执行。RUNNING 只来自 + # 分阶段执行的上一个阶段(任务保持 RUNNING 等下一阶段)或进程异常中断的 + # 残留,续跑时已完成节点由产物表跳过;调度器只拾取 QUEUED 任务、批量 + # 引擎单线程推进,不会出现两个驱动方重复执行同一任务。 + if run is None or run["status"] not in ("QUEUED", "PAUSED", "RUNNING"): return # 已暂停的任务不自动续跑:直接返回保持 PAUSED,等用户显式 resume # (resume 转回 QUEUED 后才执行);否则暂停会被立刻覆盖成 RUNNING。 @@ -153,6 +161,9 @@ class WorkflowScheduler: definition = WorkflowDefinition.from_dict(version["definition"]) definition.validate() ordered = topological_sort(definition) + # 阶段节点必须存在于 DAG:写错会让任务永远停在 RUNNING 无人推进。 + if stop_after is not None and stop_after not in ordered: + raise ValueError(f"stop_after node not in workflow: {stop_after}") except Exception as exc: # noqa: BLE001 logger.exception("任务 %s 工作流定义无效,标记失败: %s", run_id, exc) self.db.update_run( @@ -177,6 +188,9 @@ class WorkflowScheduler: return # 断点续跑:跳过已产出结果的节点(其产物已作为输入可用)。 if node_id in outputs_by_node: + # 阶段边界落在已完成的节点上:直接结束本阶段。 + if node_id == stop_after: + break continue # 当前节点进度 = 已完成节点数 / 总节点数。 node_spec = next(item for item in definition.nodes if item.id == node_id) @@ -244,6 +258,17 @@ class WorkflowScheduler: time.monotonic() - node_started, time.monotonic() - run_started, ) + # 阶段边界:本阶段节点已完成,不再执行后续节点。 + if node_id == stop_after: + break + # 分阶段执行:本阶段节点已全部完成(含本轮跳过的情况),任务保持 + # RUNNING 等下一个阶段,不做 final_outputs 与完成标记。 + if stop_after is not None: + logger.info( + "任务 %s 阶段完成: 已执行到节点 %s(分阶段执行,保持 RUNNING 等下一阶段)", + run_id, stop_after, + ) + return # 处理 final_outputs,为用户端提供简洁的下载别名。 for alias, ref in definition.final_outputs.items(): resolved = self._resolve_ref(ref, run.get("input_uri"), outputs_by_node) diff --git a/tests/app/test_batch/test_batch.py b/tests/app/test_batch/test_batch.py index ddcc441..3bbdc10 100644 --- a/tests/app/test_batch/test_batch.py +++ b/tests/app/test_batch/test_batch.py @@ -8,12 +8,14 @@ from __future__ import annotations import json +from contextlib import contextmanager from pathlib import Path import pytest from wov_app import registry from wov_app.batch import ( + KEEP_MODEL_FLAG, MARKER_NAME, SUBTITLE_EXTENSIONS, VIDEO_EXTENSIONS, @@ -26,7 +28,7 @@ from wov_app.batch import ( scan_videos, ) from wov_app.db import Database -from wov_sdk.models import WorkflowDefinition +from wov_sdk.models import InvokeRequest, InvokeResponse, WorkflowDefinition @pytest.fixture(autouse=True) @@ -70,6 +72,78 @@ def _published_db(tmp_path: Path, workflow_id: str = "wf") -> Database: return db +def _staged_definition() -> WorkflowDefinition: + """三节点分阶段链路:prep(echo) → translate(llm) → post(echo),产物为 srt。""" + return WorkflowDefinition.from_dict({ + "name": "分阶段流程", + "version": 1, + "nodes": [ + {"id": "prep", "node_type": "echo", "params": {"node_tag": "prep"}, + "inputs": {"file_uri": "input.video_uri"}}, + {"id": "translate", "node_type": "llm-translate", "params": {"node_tag": "translate"}, + "inputs": {"file_uri": "prep.file_uri"}}, + {"id": "post", "node_type": "echo", "params": {"node_tag": "post"}, + "inputs": {"file_uri": "translate.file_uri"}}, + ], + "edges": [{"from": "prep", "to": "translate"}, {"from": "translate", "to": "post"}], + "entry_inputs": {"video_uri": "file"}, + "final_outputs": {"cn_srt": "post.file_uri"}, + }) + + +def _staged_db(tmp_path: Path, workflow_id: str = "wf") -> Database: + """建好已发布的三节点分阶段工作流库。""" + db = Database(tmp_path / "wov.db") + db.upsert_workflow({ + "id": workflow_id, "name": "分阶段流程", "description": "", "published": 1, + "latest_version": 1, + }) + db.create_workflow_version(workflow_id, 1, _staged_definition().to_dict()) + return db + + +@contextmanager +def _recording_nodes( + trace: list[tuple[str, str]], + llm_flag_state: list[bool] | None = None, + fail_stage: tuple[str, str] | None = None, + flag_trace: list[tuple[str, bool]] | None = None, +): + """把 echo / llm-translate 节点换成记录调用顺序的假节点(覆盖真实注册表条目)。 + + 假节点把上游传来的视频名写成 payload.srt 透传给下一节点,因此每个阶段都 + 知道自己在处理哪个视频;记录 (节点标签, 视频名),llm 节点额外记录 + keep_model.flag 是否存在;fail_stage 指定的 (标签, 视频名) 组合返回失败。 + """ + registry.register_all() + + def handler(request: InvokeRequest) -> InvokeResponse: + tag = str(request.params.get("node_tag")) + source = str(request.inputs.get("file_uri") or "") + if source and Path(source).suffix.lower() in VIDEO_EXTENSIONS: + # 首阶段的输入就是视频文件,后续阶段拿到的是上一阶段的 payload。 + video_name = Path(source).name + else: + video_name = Path(source).read_text(encoding="utf-8").strip() if source else "" + trace.append((tag, video_name)) + run_root = Path(request.output_dir).parent.parent + if llm_flag_state is not None and tag == "translate": + llm_flag_state.append((run_root / KEEP_MODEL_FLAG).exists()) + if flag_trace is not None: + flag_trace.append((tag, (run_root / KEEP_MODEL_FLAG).exists())) + if fail_stage == (tag, video_name): + return InvokeResponse(status="failed", error="模拟阶段失败") + output_dir = Path(request.output_dir) + output_dir.mkdir(parents=True, exist_ok=True) + output = output_dir / "payload.srt" + output.write_text(video_name, encoding="utf-8") + return InvokeResponse(status="completed", outputs={"file_uri": str(output)}) + + for node_type in ("echo", "llm-translate"): + registry.register(registry.get_node(node_type), handler) + yield + + # --------------------------------------------------------------------------- # 扫描与旁挂字幕判定 # --------------------------------------------------------------------------- @@ -412,3 +486,265 @@ def test_worker_start_stop_idempotent(tmp_path: Path) -> None: # 验证结果 assert first is second assert worker._thread is None + + +def test_worker_processes_recovered_zombie_job(tmp_path: Path, monkeypatch) -> None: + """僵尸任务(COMPLETED 但明细仍 PENDING)被恢复后能真正处理完剩余视频。 + + 曾出现「任务已完成、视频仍未处理」的僵尸状态(完成标记先于视频收尾写出), + 而引擎只拾取 QUEUED:不恢复就永远不会再处理那个视频。 + """ + # 数据:已发布工作流 + 一个视频;创建任务后伪造成僵尸状态。 + folder = tmp_path / "videos" + _make_video(folder / "movie.mp4") + db = _published_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + registry.register_all() + job_id = create_job(db, str(folder), "wf") + db.update_batch_job(job_id, status="COMPLETED", progress=1.0, + updated_at="2026-09-01T00:00:00+00:00") + + # 测试过程:重启恢复把僵尸任务放回队列,引擎拾起后处理剩余视频。 + db.recover_interrupted_batch_jobs("2026-09-01T01:00:00+00:00") + worker = BatchWorker(db, interval_seconds=999) + worker._process_job(db.get_batch_job(job_id)) + + # 验证结果:视频真的处理完、产物放到视频旁、任务保持完成。 + item = db.list_batch_videos(job_id)[0] + assert item["status"] == "COMPLETED" + assert db.get_batch_job(job_id)["status"] == "COMPLETED" + assert list(folder.glob("movie.*")), "应在视频旁放置最终产物" + + +# --------------------------------------------------------------------------- +# 引擎:分块流水线(阶段化执行) +# --------------------------------------------------------------------------- + + +def test_worker_runs_grouped_stage_pipeline(tmp_path: Path, monkeypatch) -> None: + """分块流水线:组内按节点顺序跑完全部视频,而不是每个视频跑完整链路。""" + # 数据:3 个视频 + 三节点链路,分组大小 2(前两个一组、第三个一组)。 + folder = tmp_path / "videos" + for name in ("a.mp4", "b.mp4", "c.mp4"): + _make_video(folder / name) + db = _staged_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + monkeypatch.setattr("wov_app.batch.BATCH_STAGE_GROUP_SIZE", 2) + trace: list[tuple[str, str]] = [] + + # 测试过程:用记录调用顺序的假节点驱动引擎跑一轮。 + with _recording_nodes(trace): + job_id = create_job(db, str(folder), "wf") + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:组1 三个阶段各跑 a/b,再轮到组2 的 c。 + assert trace == [ + ("prep", "a.mp4"), ("prep", "b.mp4"), + ("translate", "a.mp4"), ("translate", "b.mp4"), + ("post", "a.mp4"), ("post", "b.mp4"), + ("prep", "c.mp4"), ("translate", "c.mp4"), ("post", "c.mp4"), + ] + # 三个视频都完成且产物按约定名落到视频旁。 + assert db.get_batch_job(job_id)["status"] == "COMPLETED" + assert sorted(p.name for p in folder.glob("*.srt")) == ["a.CN.srt", "b.CN.srt", "c.CN.srt"] + + +def test_worker_releases_local_llm_once_per_group(tmp_path: Path, monkeypatch) -> None: + """LLM 阶段结束由引擎统一释放显存:每组一次,而不是每个视频一次。""" + # 数据:3 个视频 + 分组 2,记录释放调用与 LLM 调用时的常驻信号状态。 + folder = tmp_path / "videos" + for name in ("a.mp4", "b.mp4", "c.mp4"): + _make_video(folder / name) + db = _staged_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + monkeypatch.setattr("wov_app.batch.BATCH_STAGE_GROUP_SIZE", 2) + releases: list[str | None] = [] + monkeypatch.setattr( + "wov_app.batch.release_local_model", + lambda model=None: releases.append(model), + ) + trace: list[tuple[str, str]] = [] + flag_state: list[bool] = [] + + # 测试过程 + with _recording_nodes(trace, llm_flag_state=flag_state): + job_id = create_job(db, str(folder), "wf") + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:两组各释放一次;每次 LLM 调用都在“保持常驻”信号下执行;信号已清理。 + assert releases == [None, None] + assert flag_state == [True, True, True] + assert not list((tmp_path / "storage").rglob(KEEP_MODEL_FLAG)) + + +def test_worker_defers_video_failed_in_earlier_stage(tmp_path: Path, monkeypatch) -> None: + """上一阶段失败的视频不在后续阶段重跑(避免 LLM 已常驻时重跑 ASR 抢显存)。""" + # 数据:2 个视频(同一组)+ 三节点链路,prep 阶段让 a 失败。 + folder = tmp_path / "videos" + for name in ("a.mp4", "b.mp4"): + _make_video(folder / name) + db = _staged_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + monkeypatch.setattr("wov_app.batch.BATCH_STAGE_GROUP_SIZE", 2) + monkeypatch.setattr("wov_app.batch.release_local_model", lambda model=None: None) + trace: list[tuple[str, str]] = [] + + # 测试过程 + with _recording_nodes(trace, fail_stage=("prep", "a.mp4")): + job_id = create_job(db, str(folder), "wf") + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:a 只在 prep 出现一次并记为 FAILED;b 三阶段跑完并落地产物。 + assert [entry for entry in trace if entry[1] == "a.mp4"] == [("prep", "a.mp4")] + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job_id)} + assert videos["a.mp4"]["status"] == "FAILED" + assert videos["b.mp4"]["status"] == "COMPLETED" + assert (folder / "b.CN.srt").is_file() + assert not (folder / "a.CN.srt").exists() + assert db.get_batch_job(job_id)["failed"] == 1 + + +# --------------------------------------------------------------------------- +# 引擎:创建期间拾起任务(明细未登记完)的自愈 +# --------------------------------------------------------------------------- + + +def test_worker_requeues_job_without_details(tmp_path: Path, monkeypatch) -> None: + """任务行先于明细写入:拾起到无明细的任务时保持 QUEUED,不按空任务收尾。""" + # 数据:只有任务行、还没写任何视频明细的批量任务。 + db = _published_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + db.create_batch_job({ + "id": "batch-registering", "folder_path": str(tmp_path), "workflow_id": "wf", + "recursive": 1, "status": "QUEUED", "progress": 0, "total": 0, "done": 0, + "failed": 0, "current_video": None, "error": None, + "created_at": "2026-09-01T00:00:00+00:00", "updated_at": "2026-09-01T00:00:00+00:00", + }) + + # 测试过程 + worker = BatchWorker(db, interval_seconds=999) + worker._process_job(db.get_batch_job("batch-registering")) + + # 验证结果:任务仍在排队等待登记完成,而不是被标成 COMPLETED。 + assert db.get_batch_job("batch-registering")["status"] == "QUEUED" + assert db.next_queued_batch_job() is not None + + +def test_worker_requeues_job_when_video_registered_mid_pass(tmp_path: Path, monkeypatch) -> None: + """明细在引擎处理中途才登记进来:本轮结束后置回 QUEUED,下一轮续跑完成。""" + # 数据:1 个视频 + 单节点工作流;处理首个阶段时登记第二个视频(模拟创建中拾起)。 + folder = tmp_path / "videos" + _make_video(folder / "first.mp4") + db = _published_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + registry.register_all() + job_id = create_job(db, str(folder), "wf") + worker = BatchWorker(db, interval_seconds=999) + original_stage = worker._run_stage + injected = {"done": False} + + def stage_with_late_video(*args, **kwargs): + # 模拟 create_job 仍在写明细:引擎快照之后新视频才出现在数据库里。 + if not injected["done"]: + injected["done"] = True + late_video = _make_video(folder / "second.mp4") + db.create_batch_video({ + "id": "bv_late", "job_id": job_id, "video_path": str(late_video), + "work_dir": str(tmp_path / "storage" / "batch" / job_id / "bv_late"), + "run_id": None, "status": "PENDING", "error": None, + "created_at": "2026-09-01T00:00:01+00:00", "updated_at": "2026-09-01T00:00:01+00:00", + }) + return original_stage(*args, **kwargs) + + monkeypatch.setattr(worker, "_run_stage", stage_with_late_video) + + # 测试过程:第一轮只看到 first.mp4。 + worker._process_job(db.get_batch_job(job_id)) + + # 验证结果:job 被置回 QUEUED 等待下一轮,新视频还没被处理。 + assert db.get_batch_job(job_id)["status"] == "QUEUED" + statuses = {Path(v["video_path"]).name: v["status"] for v in db.list_batch_videos(job_id)} + assert statuses == {"first.mp4": "COMPLETED", "second.mp4": "PENDING"} + + # 测试过程:下一轮引擎拾起后处理剩余视频并收尾。 + second = BatchWorker(db, interval_seconds=999) + second._process_job(db.get_batch_job(job_id)) + + # 验证结果:两个视频都完成、任务完成、产物都在视频旁(真实 echo 节点产物为 echo.txt)。 + statuses = {Path(v["video_path"]).name: v["status"] for v in db.list_batch_videos(job_id)} + assert statuses == {"first.mp4": "COMPLETED", "second.mp4": "COMPLETED"} + assert db.get_batch_job(job_id)["status"] == "COMPLETED" + # 真实 echo 节点的最终产物保留原扩展名(非 .srt/.ass),按视频主名放置。 + assert list(folder.glob("first.*")) and list(folder.glob("second.*")) + + +def test_worker_removes_empty_job_workspace_after_completion(tmp_path: Path, monkeypatch) -> None: + """任务全部完成后删掉任务级工作空间目录(每视频工作空间已各自清理)。""" + # 数据:一个视频的批量任务。 + folder = tmp_path / "videos" + _make_video(folder / "movie.mp4") + db = _published_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + registry.register_all() + job_id = create_job(db, str(folder), "wf") + + # 测试过程 + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:任务完成,任务级目录(收尾后只剩空壳)被删除。 + assert db.get_batch_job(job_id)["status"] == "COMPLETED" + assert not (tmp_path / "storage" / "batch" / job_id).exists() + + +def test_worker_keeps_job_workspace_when_video_failed(tmp_path: Path, monkeypatch) -> None: + """有失败视频时保留任务工作空间(失败视频的中间产物供断点重试)。""" + # 数据:三节点链路,prep 阶段让视频失败。 + folder = tmp_path / "videos" + _make_video(folder / "a.mp4") + db = _staged_db(tmp_path) + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", tmp_path / "storage" / "batch") + monkeypatch.setattr("wov_app.batch.BATCH_STAGE_GROUP_SIZE", 2) + monkeypatch.setattr("wov_app.batch.release_local_model", lambda model=None: None) + + # 测试过程 + with _recording_nodes([], fail_stage=("prep", "a.mp4")): + job_id = create_job(db, str(folder), "wf") + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:视频失败、任务工作空间仍在(可重试)。 + assert db.get_batch_job(job_id)["failed"] == 1 + assert (tmp_path / "storage" / "batch" / job_id).exists() + + +def test_worker_clears_stale_keep_model_flag_before_stage(tmp_path: Path, monkeypatch) -> None: + """强杀残留的 keep_model.flag 不会带到后续阶段:非 LLM 节点不应看到它。""" + # 数据:两节点链路 + 已存在的 run(工作空间里残留强杀时的 keep_model.flag)。 + folder = tmp_path / "videos" + _make_video(folder / "a.mp4") + db = _staged_db(tmp_path) + work_root = tmp_path / "storage" / "batch" + monkeypatch.setattr("wov_app.batch.BATCH_WORK_ROOT", work_root) + monkeypatch.setattr("wov_app.batch.BATCH_STAGE_GROUP_SIZE", 8) + monkeypatch.setattr("wov_app.batch.release_local_model", lambda model=None: None) + job_id = create_job(db, str(folder), "wf") + video = db.list_batch_videos(job_id)[0] + run_id = "run_stale_flag" + db.update_batch_video(video["id"], run_id=run_id, updated_at="2026-09-01T00:00:00+00:00") + db.create_run({ + "id": run_id, "workflow_id": "wf", "workflow_version": 1, "status": "QUEUED", + "current_node_id": None, "progress": 0.0, "error": None, + "input_uri": str(folder / "a.mp4"), "param_overrides": None, "source": "batch", + "created_at": "2026-09-01T00:00:00+00:00", "updated_at": "2026-09-01T00:00:00+00:00", + }) + run_dir = Path(video["work_dir"]) / "runs" / run_id + run_dir.mkdir(parents=True, exist_ok=True) + (run_dir / KEEP_MODEL_FLAG).write_text("", encoding="utf-8") + flag_trace: list[tuple[str, bool]] = [] + + # 测试过程 + with _recording_nodes([], flag_trace=flag_trace): + BatchWorker(db, interval_seconds=999)._process_job(db.get_batch_job(job_id)) + + # 验证结果:prep(非 LLM)看不到残留标志;translate(LLM 阶段)才写入; + # post(非 LLM)不再看到它。 + assert flag_trace == [("prep", False), ("translate", True), ("post", False)] diff --git a/tests/app/test_db/test_database.py b/tests/app/test_db/test_database.py index f865bef..280725b 100644 --- a/tests/app/test_db/test_database.py +++ b/tests/app/test_db/test_database.py @@ -164,6 +164,22 @@ def test_list_runs_orders_by_created_at_desc(db_with_workflow: Database) -> None assert ids == ["new", "mid", "old"] +def test_list_runs_excludes_batch_runs_by_default(db_with_workflow: Database) -> None: + """任务列表默认不含批量 run:它们属于批量页的任务明细,会把 20 条窗口占满。""" + # 数据:一条上传任务 + 两条批量 run。 + db_with_workflow.create_run(_run("upload", created_at="2026-09-01T00:00:00+00:00")) + db_with_workflow.create_run(_run("batch-1", source="batch", created_at="2026-09-02T00:00:00+00:00")) + db_with_workflow.create_run(_run("batch-2", source="batch", created_at="2026-09-03T00:00:00+00:00")) + + # 测试过程 + default_ids = [r["id"] for r in db_with_workflow.list_runs()] + all_ids = [r["id"] for r in db_with_workflow.list_runs(include_batch=True)] + + # 验证结果:默认只列上传任务,显式要求时才包含批量 run。 + assert default_ids == ["upload"] + assert all_ids == ["batch-2", "batch-1", "upload"] + + def test_delete_run_removes_record_and_artifacts(db_with_workflow: Database) -> None: """删除任务同时清理其产物记录。""" # 数据:任务 + 一条产物。 @@ -330,3 +346,38 @@ def test_list_run_ids(db_with_workflow: Database) -> None: # 测试过程与验证结果 assert sorted(db_with_workflow.list_run_ids()) == ["r1", "r2"] + + +def test_recover_interrupted_batch_jobs_requeues_zombie_completed(db_with_workflow: Database) -> None: + """重启恢复:COMPLETED 但仍有未结束视频的僵尸任务也要置回 QUEUED。 + + 曾出现「任务已完成、视频仍未处理」的僵尸状态:完成标记先于视频收尾写出, + 而引擎只拾取 QUEUED,剩余视频永久无人处理。 + """ + # 数据:一个仍含 PENDING 视频的 COMPLETED 任务 + 一个全部结束的 COMPLETED 任务。 + db_with_workflow.create_batch_job({ + "id": "zombie", "folder_path": "/videos", "workflow_id": "wf", "recursive": False, + "status": "COMPLETED", "created_at": "t1", "updated_at": "t1", + }) + db_with_workflow.create_batch_video({ + "id": "zombie-v1", "job_id": "zombie", "video_path": "/videos/a.mp4", + "work_dir": "/tmp/zombie", "status": "PENDING", + "created_at": "t1", "updated_at": "t1", + }) + db_with_workflow.create_batch_job({ + "id": "done", "folder_path": "/videos", "workflow_id": "wf", "recursive": False, + "status": "COMPLETED", "created_at": "t1", "updated_at": "t1", + }) + db_with_workflow.create_batch_video({ + "id": "done-v1", "job_id": "done", "video_path": "/videos/b.mp4", + "work_dir": "/tmp/done", "status": "COMPLETED", + "created_at": "t1", "updated_at": "t1", + }) + + # 测试过程 + count = db_with_workflow.recover_interrupted_batch_jobs("t2") + + # 验证结果:只有僵尸任务被置回 QUEUED,真正完成的任务不受影响。 + assert count == 1 + assert db_with_workflow.get_batch_job("zombie")["status"] == "QUEUED" + assert db_with_workflow.get_batch_job("done")["status"] == "COMPLETED" diff --git a/tests/app/test_routers/test_apps_api.py b/tests/app/test_routers/test_apps_api.py index 2a7a58a..76dd7ba 100644 --- a/tests/app/test_routers/test_apps_api.py +++ b/tests/app/test_routers/test_apps_api.py @@ -191,6 +191,27 @@ def test_get_missing_run_returns_404(client: TestClient) -> None: assert client.get("/api/runs/nope").status_code == 404 +def test_list_runs_excludes_batch_runs_by_default(client: TestClient, tmp_path: Path) -> None: + """批量 run 不进任务管理默认列表(它们是批量任务明细,由批量页展示)。""" + # 数据:一条上传任务 + 一条同工作流的批量 run。 + upload_id = _create_run(client) + db = Database(tmp_path / "wov.db") + db.create_run({ + "id": "run_batch_listed", "workflow_id": "echo-app", "workflow_version": 1, + "status": "RUNNING", "current_node_id": "step", "progress": 0.0, "error": None, + "input_uri": "/videos/movie.mp4", "param_overrides": None, "source": "batch", + "created_at": "2026-09-09T00:00:00+00:00", "updated_at": "2026-09-09T00:00:00+00:00", + }) + + # 测试过程 + default_ids = [item["id"] for item in client.get("/api/runs").json()] + all_ids = [item["id"] for item in client.get("/api/runs", params={"include_batch": 1}).json()] + + # 验证结果:默认列表只有上传任务,显式请求时包含批量 run。 + assert default_ids == [upload_id] + assert set(all_ids) == {upload_id, "run_batch_listed"} + + def test_pause_and_resume_run(client: TestClient) -> None: """暂停置 PAUSED、继续置 QUEUED,并写入/清除暂停信号文件。""" # 数据:一条任务。 diff --git a/tests/app/test_routers/test_batch_api.py b/tests/app/test_routers/test_batch_api.py index d8441a7..1ddc52e 100644 --- a/tests/app/test_routers/test_batch_api.py +++ b/tests/app/test_routers/test_batch_api.py @@ -13,6 +13,7 @@ import pytest from fastapi.testclient import TestClient from wov_app import registry +from wov_app.db import Database from wov_app.main import app as fastapi_app from wov_sdk.models import WorkflowDefinition @@ -308,3 +309,65 @@ def test_download_missing_product_returns_404(client: TestClient, tmp_path: Path # 验证结果 assert response.status_code == 404 + + +def test_job_detail_reports_video_stage(client: TestClient, tmp_path: Path) -> None: + """详情给处理中的视频补阶段信息:第几阶段/共几阶段 + 中文标签。""" + # 数据:三节点工作流(prep → translate → post),视频 run 停在第二阶段。 + definition = WorkflowDefinition.from_dict({ + "name": "分阶段流程", "version": 1, + "nodes": [ + {"id": "prep", "node_type": "echo", "inputs": {"file_uri": "input.video_uri"}}, + {"id": "translate", "node_type": "llm-translate", "inputs": {"srt_uri": "prep.file_uri"}}, + {"id": "post", "node_type": "echo", "inputs": {"file_uri": "translate.file_uri"}}, + ], + "edges": [{"from": "prep", "to": "translate"}, {"from": "translate", "to": "post"}], + "entry_inputs": {"video_uri": "file"}, + "final_outputs": {"result": "post.file_uri"}, + }).to_dict() + client.post("/api/admin/workflows", json={ + "id": "wf-staged", "name": "分阶段流程", "description": "", "definition": definition, + }) + client.post("/api/admin/workflows/wf-staged/publish") + folder = tmp_path / "videos" + _make_video(folder / "movie.mp4") + job_id = client.post("/api/batch/jobs", json={ + "folder": str(folder), "workflow_id": "wf-staged", "recursive": True, + }).json()["id"] + db = Database(tmp_path / "wov.db") + video = [v for v in db.list_batch_videos(job_id) if v["status"] != "SKIPPED"][0] + db.update_batch_video(video["id"], status="RUNNING", run_id="run_stage", updated_at="2026-09-01T00:00:00+00:00") + db.create_run({ + "id": "run_stage", "workflow_id": "wf-staged", "workflow_version": 1, + "status": "RUNNING", "current_node_id": "translate", "progress": 0.3333, + "error": None, "input_uri": str(folder / "movie.mp4"), "param_overrides": None, + "source": "batch", "created_at": "2026-09-01T00:00:00+00:00", + "updated_at": "2026-09-01T00:00:00+00:00", + }) + + # 测试过程 + body = client.get(f"/api/batch/jobs/{job_id}").json() + + # 验证结果:阶段序号/总数与节点类型对应的中文标签。 + item = [v for v in body["videos"] if v["status"] != "SKIPPED"][0] + assert item["stage_label"] == "翻译" + assert (item["stage_index"], item["stage_total"]) == (2, 3) + + +def test_job_detail_omits_stage_for_unstarted_video(client: TestClient, tmp_path: Path) -> None: + """还没开始处理的视频没有阶段信息(前端显示占位符)。""" + # 数据:一个 PENDING 视频(无 run)。 + folder = tmp_path / "videos" + _make_video(folder / "movie.mp4") + _publish_workflow(client) + job_id = client.post("/api/batch/jobs", json={ + "folder": str(folder), "workflow_id": "wf", "recursive": True, + }).json()["id"] + + # 测试过程 + body = client.get(f"/api/batch/jobs/{job_id}").json() + + # 验证结果:阶段字段为空。 + item = [v for v in body["videos"] if v["status"] != "SKIPPED"][0] + assert item["stage_label"] is None + assert item["stage_index"] is None and item["stage_total"] is None diff --git a/tests/app/test_scheduler/test_scheduler.py b/tests/app/test_scheduler/test_scheduler.py index 6848ead..4b31124 100644 --- a/tests/app/test_scheduler/test_scheduler.py +++ b/tests/app/test_scheduler/test_scheduler.py @@ -116,6 +116,72 @@ def test_topological_sort_diamond() -> None: assert set(order[1:3]) == {"b", "c"} +def test_execute_run_stop_after_leaves_run_running_and_resumes(tmp_path: Path) -> None: + """分阶段执行:stop_after 指定阶段节点后停下(保持 RUNNING、不收尾),再次调用续跑完成。""" + # 数据:a → b → c 三段 echo 链 + 最终别名。 + definition = _definition( + nodes=[ + {"id": "a", "node_type": "echo", "inputs": {"file_uri": "input.video_uri"}}, + {"id": "b", "node_type": "echo", "inputs": {"file_uri": "a.file_uri"}}, + {"id": "c", "node_type": "echo", "inputs": {"file_uri": "b.file_uri"}}, + ], + edges=[{"from": "a", "to": "b"}, {"from": "b", "to": "c"}], + final_outputs={"result": "c.file_uri"}, + ) + storage = tmp_path / "storage" + db = _db_with_workflow(tmp_path, definition) + source = tmp_path / "input.txt" + source.write_text("分阶段内容", encoding="utf-8") + db.create_run(_run("run-stage", str(source))) + registry.register_all() + scheduler = _scheduler(db, storage) + + # 测试过程:第一阶段只执行到 b 为止。 + scheduler.execute_run("run-stage", stop_after="b") + + # 验证结果:任务保持 RUNNING 未收尾,a/b 产物已登记,c 与最终别名都没有。 + staged = db.get_run("run-stage") + assert staged["status"] == "RUNNING" + assert staged["current_node_id"] == "b" + names = {artifact["name"] for artifact in db.list_artifacts("run-stage")} + assert {"a.file_uri", "b.file_uri"} <= names + assert "c.file_uri" not in names + assert "result" not in names + + # 测试过程:不传 stop_after 时整条 DAG 跑完并收尾。 + scheduler.execute_run("run-stage") + + # 验证结果:完成、进度 1.0、最终别名登记。 + finished = db.get_run("run-stage") + assert finished["status"] == "COMPLETED" + assert finished["progress"] == 1.0 + assert "result" in {artifact["name"] for artifact in db.list_artifacts("run-stage")} + + +def test_execute_run_stop_after_unknown_node_marks_failed(tmp_path: Path) -> None: + """stop_after 指向不存在的节点时任务标 FAILED(不留下永远 RUNNING 的任务)。""" + # 数据:单 echo 节点任务。 + definition = _definition( + nodes=[{"id": "step", "node_type": "echo", "inputs": {"file_uri": "input.video_uri"}}], + edges=[], + final_outputs={"result": "step.file_uri"}, + ) + storage = tmp_path / "storage" + db = _db_with_workflow(tmp_path, definition) + source = tmp_path / "input.txt" + source.write_text("内容", encoding="utf-8") + db.create_run(_run("run-bad-stage", str(source))) + registry.register_all() + + # 测试过程 + _scheduler(db, storage).execute_run("run-bad-stage", stop_after="nope") + + # 验证结果:FAILED 且错误说明阶段节点不存在。 + stored = db.get_run("run-bad-stage") + assert stored["status"] == "FAILED" + assert "stop_after" in (stored["error"] or "") + + # --------------------------------------------------------------------------- # 执行:成功路径 # --------------------------------------------------------------------------- diff --git a/tests/nodes/test_llm/test_translate.py b/tests/nodes/test_llm/test_translate.py index b0efe64..d163e62 100644 --- a/tests/nodes/test_llm/test_translate.py +++ b/tests/nodes/test_llm/test_translate.py @@ -17,9 +17,11 @@ import pytest from nodes.llm import ( CHUNK_SIZE, MAX_BATCH_RETRIES, + PauseRequested, _parse_translations, _system_prompt, invoke, + release_local_model, translate_lines, ) from wov_sdk.models import InvokeRequest @@ -28,6 +30,17 @@ from wov_sdk.models import InvokeRequest DATA_DIR = Path(__file__).resolve().parent / "data" +@pytest.fixture(autouse=True) +def _isolate_llm_env(monkeypatch): + """清掉外部泄漏的 LLM 路由变量,保证用例只受自己显式设置的环境变量影响。 + + 全量跑时其它模块 import 应用会触发 load_dotenv(),把开发者 .env 里的 + LLM_API_BASE(可能指向本机 Ollama)带进来,从而改变端点判定与请求数量。 + """ + for name in ("LLM_API_BASE", "LLM_MODEL", "LLM_UNLOAD_AFTER"): + monkeypatch.delenv(name, raising=False) + + class _FakeHTTPResponse: """假的 HTTP 响应:返回预置 JSON 体(供 urlopen mock 使用)。""" @@ -70,6 +83,35 @@ def _capture_urlopen(calls: list[dict], responses: list[_FakeHTTPResponse]): return fake_urlopen +def _capture_urlopen_routing_unload( + calls: list[dict], + responses: list[_FakeHTTPResponse], + unload_error: Exception | None = None, +): + """按 URL 分流的假 urlopen:卸载请求只记录,翻译请求按序返回预置响应。""" + + def fake_urlopen(http_request, timeout=None): + calls.append({ + "url": http_request.full_url, + "body": json.loads(http_request.data.decode("utf-8")), + "timeout": timeout, + }) + if http_request.full_url.endswith("/api/generate"): + if unload_error is not None: + raise unload_error + return _FakeHTTPResponse({"done": True}) + return responses.pop(0) if responses else _llm_reply([]) + + return fake_urlopen + + +def _local_llm_env(monkeypatch) -> None: + """把 LLM 端点指向本地 Ollama(qwen3:30b-a3b)。""" + monkeypatch.setenv("LLM_API_KEY", "") + monkeypatch.setenv("LLM_API_BASE", "http://localhost:11434/v1/chat/completions") + monkeypatch.setenv("LLM_MODEL", "qwen3:30b-a3b") + + # --------------------------------------------------------------------------- # 提示词与响应解析(纯函数) # --------------------------------------------------------------------------- @@ -363,6 +405,200 @@ def test_translate_lines_sends_bearer_key(monkeypatch) -> None: assert headers.get("authorization") == "Bearer sk-abc" +def test_translate_lines_unloads_local_model_when_env_enabled(monkeypatch) -> None: + """开启 LLM_UNLOAD_AFTER 时翻译结束请求 Ollama 卸载模型,把显存让给 whisper。""" + # 数据:本地端点 + 开启卸载。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.setenv("LLM_UNLOAD_AFTER", "1") + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {}) + + # 验证结果:翻译后向 Ollama 原生端点发 keep_alive=0 的卸载请求。 + assert [call["url"] for call in calls] == [ + "http://localhost:11434/v1/chat/completions", + "http://localhost:11434/api/generate", + ] + assert calls[1]["body"] == {"model": "qwen3:30b-a3b", "keep_alive": 0} + + +def test_translate_lines_auto_unloads_loopback_endpoint(monkeypatch) -> None: + """端点在本机(loopback)时默认卸载:无需开关,默认就让出显存。""" + # 数据:本地端点 + 不设置任何开关。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.delenv("LLM_UNLOAD_AFTER", raising=False) + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {}) + + # 验证结果:本机端点默认发出卸载请求。 + assert calls[-1]["url"] == "http://localhost:11434/api/generate" + + +def test_translate_lines_does_not_unload_remote_endpoint(monkeypatch) -> None: + """云端端点默认不卸载:显存不由本机持有,多发请求只是噪声。""" + # 数据:远程端点 + 不设置任何开关。 + calls: list[dict] = [] + monkeypatch.setenv("LLM_API_KEY", "sk-test") + monkeypatch.setenv("LLM_API_BASE", "https://api.siliconflow.cn/v1/chat/completions") + monkeypatch.delenv("LLM_UNLOAD_AFTER", raising=False) + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {}) + + # 验证结果:只有翻译请求。 + assert [call["url"] for call in calls] == ["https://api.siliconflow.cn/v1/chat/completions"] + + +def test_translate_lines_param_unload_after_enables_unload(monkeypatch) -> None: + """节点参数 unload_after=True 等效于环境变量开关。""" + # 数据:只给节点参数。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.delenv("LLM_UNLOAD_AFTER", raising=False) + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {"unload_after": True}) + + # 验证结果 + assert calls[-1]["url"] == "http://localhost:11434/api/generate" + + +def test_translate_lines_param_unload_after_false_overrides_default(monkeypatch) -> None: + """节点参数 unload_after=False 可显式关闭本机端点的默认卸载。""" + # 数据:本机端点 + 节点参数显式关闭。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.setenv("LLM_UNLOAD_AFTER", "1") + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {"unload_after": False}) + + # 验证结果 + assert [call["url"] for call in calls] == ["http://localhost:11434/v1/chat/completions"] + + +def test_translate_lines_keeps_translation_when_unload_fails(monkeypatch) -> None: + """卸载请求失败不影响译文(端点不支持卸载时只是跳过释放)。""" + # 数据:远程端点显式开启卸载 + 卸载请求返回 404。 + calls: list[dict] = [] + monkeypatch.setenv("LLM_API_KEY", "sk-test") + monkeypatch.setenv("LLM_UNLOAD_AFTER", "1") + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload( + calls, + [_llm_reply([(1, "译文")])], + unload_error=urllib.error.HTTPError( + "https://api.example.com/api/generate", 404, "not found", {}, None + ), + ), + ) + + # 测试过程 + translated = translate_lines(["一"], {}) + + # 验证结果:译文正常返回,卸载失败只记录。 + assert translated == ["译文"] + assert calls[-1]["url"].endswith("/api/generate") + + +def test_translate_lines_keeps_model_loaded_in_staged_batch(monkeypatch) -> None: + """引擎标记阶段内保持常驻时不卸载模型:整批只加载一次,阶段结束统一释放。""" + # 数据:本机端点(默认会卸载)+ keep_model_loaded=True。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + + # 测试过程 + translate_lines(["一"], {}, keep_model_loaded=True) + + # 验证结果:只发翻译请求,不发卸载请求。 + assert [call["url"] for call in calls] == ["http://localhost:11434/v1/chat/completions"] + + +def test_translate_lines_aborts_between_batches_when_stop_requested(monkeypatch) -> None: + """暂停信号在两批之间生效:已完成的批保留,后续批不再发请求。""" + # 数据:CHUNK_SIZE + 1 行(两批),第二次检查返回“应停止”。 + lines = [f"行{i}" for i in range(1, CHUNK_SIZE + 2)] + calls: list[dict] = [] + checks = {"count": 0} + + def stop_requested() -> bool: + checks["count"] += 1 + return checks["count"] > 1 + + monkeypatch.setenv("LLM_API_KEY", "sk-test") + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload( + calls, + [_llm_reply([(i, f"t{i}") for i in range(1, CHUNK_SIZE + 1)])], + ), + ) + + # 测试过程与验证结果:抛暂停异常,且只发出第一批的请求。 + with pytest.raises(PauseRequested): + translate_lines(lines, {}, stop_requested=stop_requested) + assert [call["url"] for call in calls] == ["https://api.siliconflow.cn/v1/chat/completions"] + + +def test_release_local_model_unloads_loopback_endpoint(monkeypatch) -> None: + """阶段收尾释放模型:本机端点发 keep_alive=0 卸载请求,模型名参数优先。""" + # 数据:本机端点 + 显式指定的模型名。 + calls: list[dict] = [] + _local_llm_env(monkeypatch) + monkeypatch.setattr(urllib.request, "urlopen", _capture_urlopen_routing_unload(calls, [])) + + # 测试过程 + release_local_model("local/替换模型") + + # 验证结果:命中 Ollama 原生卸载端点,使用传入的模型名。 + assert [(call["url"], call["body"]) for call in calls] == [ + ("http://localhost:11434/api/generate", {"model": "local/替换模型", "keep_alive": 0}), + ] + + +def test_release_local_model_skips_remote_endpoint(monkeypatch) -> None: + """云端端点不占本机显存,阶段收尾不发卸载请求。""" + # 数据:云端端点。 + calls: list[dict] = [] + monkeypatch.setenv("LLM_API_KEY", "sk-test") + monkeypatch.setenv("LLM_API_BASE", "https://api.siliconflow.cn/v1/chat/completions") + monkeypatch.setattr(urllib.request, "urlopen", _capture_urlopen_routing_unload(calls, [])) + + # 测试过程 + release_local_model() + + # 验证结果:没有任何请求。 + assert calls == [] + + # --------------------------------------------------------------------------- # invoke 全流程 # --------------------------------------------------------------------------- @@ -396,6 +632,36 @@ def test_invoke_translates_srt_and_writes_artifact(monkeypatch, tmp_path: Path) assert content.index("你好") < content.index("再见") +def test_invoke_stops_on_pause_flag(monkeypatch, tmp_path: Path) -> None: + """run 根目录有暂停信号时节点中止且不写产物(调度器保持任务 PAUSED)。""" + # 数据:一条真实 SRT + run 根目录下的 paused.flag。 + srt_path = tmp_path / "in.srt" + srt_path.write_text("1\n00:00:01,000 --> 00:00:02,000\nこんにちは\n", encoding="utf-8") + run_root = tmp_path / "runs" / "run-paused" + run_root.mkdir(parents=True, exist_ok=True) + (run_root / "paused.flag").write_text("", encoding="utf-8") + output_dir = run_root / "steps" / "translate" + calls: list[dict] = [] + monkeypatch.setenv("LLM_API_KEY", "sk-test") + monkeypatch.setattr( + urllib.request, "urlopen", + _capture_urlopen_routing_unload(calls, [_llm_reply([(1, "译文")])]), + ) + request = InvokeRequest( + run_id="run-paused", node_instance_id="n", params={}, + inputs={"srt_uri": str(srt_path)}, output_dir=str(output_dir), + ) + + # 测试过程 + response = invoke(request) + + # 验证结果:failed 且原因为暂停;未调用 LLM、未写 cn.srt。 + assert response.status == "failed" + assert "暂停" in (response.error or "") + assert calls == [] + assert not (output_dir / "cn.srt").exists() + + def test_invoke_fails_without_input(tmp_path: Path) -> None: """缺少 srt_uri 时失败。""" # 数据:空输入。 diff --git a/tests/scripts/__init__.py b/tests/scripts/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/scripts/test_fix_zombie_batch_jobs/__init__.py b/tests/scripts/test_fix_zombie_batch_jobs/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/scripts/test_fix_zombie_batch_jobs/test_fix_zombie_batch_jobs.py b/tests/scripts/test_fix_zombie_batch_jobs/test_fix_zombie_batch_jobs.py new file mode 100644 index 0000000..0484f5c --- /dev/null +++ b/tests/scripts/test_fix_zombie_batch_jobs/test_fix_zombie_batch_jobs.py @@ -0,0 +1,83 @@ +"""scripts/fix_zombie_batch_jobs.py 的模块级测试(数据 → 测试过程 → 验证结果)。 + +被测模块:僵尸批量任务修复脚本(把「COMPLETED 但仍有未结束视频」的历史脏数据 +置回 QUEUED),可独立调用。用例在临时目录构造真实 SQLite 库与真实明细记录, +调用脚本的真实函数而不是重写扫描逻辑。 +""" + +from __future__ import annotations + +from pathlib import Path + +from scripts.fix_zombie_batch_jobs import find_zombie_jobs, main +from wov_app.db import Database + + +def _db_with_workflow(tmp_path: Path) -> Database: + """建好已发布工作流(版本)的临时库:批量明细无外键约束,任务表需要它。""" + db = Database(tmp_path / "wov.db") + db.upsert_workflow({"id": "wf", "name": "流程", "description": ""}) + db.create_workflow_version("wf", 1, {"nodes": []}) + return db + + +def _job(db: Database, job_id: str, status: str) -> None: + """登记一条批量任务。""" + db.create_batch_job({ + "id": job_id, "folder_path": "/videos", "workflow_id": "wf", "recursive": 0, + "status": status, "created_at": "t1", "updated_at": "t1", + }) + + +def _video(db: Database, job_id: str, video_id: str, status: str) -> None: + """登记一条批量视频明细。""" + db.create_batch_video({ + "id": video_id, "job_id": job_id, "video_path": f"/videos/{video_id}.mp4", + "work_dir": f"/tmp/{video_id}", "status": status, + "created_at": "t1", "updated_at": "t1", + }) + + +def test_find_zombie_jobs_selects_completed_with_unfinished_videos(tmp_path: Path) -> None: + """只挑出 COMPLETED 但仍有未结束视频的任务,已真正完成的跳过。""" + # 数据:僵尸任务(COMPLETED + PENDING)与真正完成的任务(COMPLETED + SKIPPED)。 + db = _db_with_workflow(tmp_path) + _job(db, "zombie", "COMPLETED") + _video(db, "zombie", "z1", "PENDING") + _video(db, "zombie", "z2", "SKIPPED") + _job(db, "done", "COMPLETED") + _video(db, "done", "d1", "COMPLETED") + _video(db, "done", "d2", "SKIPPED") + + # 测试过程 + zombies = find_zombie_jobs(db) + + # 验证结果:只有僵尸任务入选,且带出未结束明细。 + assert [str(job["id"]) for job, _ in zombies] == ["zombie"] + assert [v["id"] for v in zombies[0][1]] == ["z1"] + + +def test_main_apply_requeues_zombie_and_keeps_done(tmp_path: Path) -> None: + """--apply 把僵尸任务置回 QUEUED 并清 progress,正常完成的任务不动。""" + # 数据:僵尸任务(1 完成 + 1 PENDING)与正常完成任务。 + db_path = tmp_path / "wov.db" + db = Database(db_path) + db.upsert_workflow({"id": "wf", "name": "流程", "description": ""}) + db.create_workflow_version("wf", 1, {"nodes": []}) + _job(db, "zombie", "COMPLETED") + db.update_batch_job("zombie", progress=1.0, total=1, done=0, updated_at="t1") + _video(db, "zombie", "z1", "COMPLETED") + _video(db, "zombie", "z2", "PENDING") + _job(db, "done", "COMPLETED") + _video(db, "done", "d1", "COMPLETED") + + # 测试过程 + main(str(db_path), apply=True) + + # 验证结果:僵尸任务可被引擎拾起,completed 明细与任务保持原状。 + zombie = db.get_batch_job("zombie") + assert zombie["status"] == "QUEUED" + assert zombie["progress"] == 0 + assert zombie["done"] == 1 + assert db.get_batch_job("done")["status"] == "COMPLETED" + assert db.next_queued_batch_job()["id"] == "zombie" diff --git a/tests/web/test_batch/__init__.py b/tests/web/test_batch/__init__.py new file mode 100644 index 0000000..9fc575f --- /dev/null +++ b/tests/web/test_batch/__init__.py @@ -0,0 +1 @@ +"""web/assets/batch.js 的模块级测试。""" diff --git a/tests/web/test_batch/test_batch_js.py b/tests/web/test_batch/test_batch_js.py new file mode 100644 index 0000000..6bc5ec1 --- /dev/null +++ b/tests/web/test_batch/test_batch_js.py @@ -0,0 +1,154 @@ +"""web/assets/batch.js 的模块级测试(数据 → 测试过程 → 验证结果)。 + +被测模块:`web/assets/batch.js`(批量处理页的渲染逻辑:任务进度、操作按钮、 +视频明细表 HTML)。 + +实现方式:在 Node.js 子进程中加载真实 JS 文件并调用真实函数(不重写逻辑), +验证渲染出的 HTML 内容;环境无 node 时跳过(保持跨平台可运行)。 +""" + +from __future__ import annotations + +import json +import shutil +import subprocess +from pathlib import Path + +import pytest + +# 仓库根与被测脚本。 +WORKSPACE = Path(__file__).resolve().parents[3] +APP_JS = WORKSPACE / "web" / "assets" / "app.js" +BATCH_JS = WORKSPACE / "web" / "assets" / "batch.js" + +# 用真实 JS 引擎执行调用的封装:按页面加载顺序执行 app.js(提供 +# escapeHtml/badge 等公共函数)与 batch.js,再调用后者的真实函数。 +# 渲染不需要真实 DOM,只提供脚本顶层引用到的 document 桩。 +_CALL_SCRIPT = """ +const fs = require("fs"); +const vm = require("vm"); +const sandbox = { + module: { exports: {} }, + document: { addEventListener: () => {}, getElementById: () => null }, + setInterval: () => 0, + console, +}; +vm.createContext(sandbox); +for (const path of [process.argv[1], process.argv[2]]) { + vm.runInContext(fs.readFileSync(path, "utf8"), sandbox); +} +const payload = JSON.parse(process.argv[3]); +const fn = sandbox.module.exports[payload.fn]; +process.stdout.write(JSON.stringify(fn(...payload.args))); +""" + + +def _run(fn: str, *args) -> object: + """在 node 中调用 batch.js 的真实函数并返回解析后的结果。""" + node = shutil.which("node") + if node is None: + pytest.skip("环境没有 node,跳过 JS 模块测试") + completed = subprocess.run( + [ + node, + "-e", + _CALL_SCRIPT, + str(APP_JS), + str(BATCH_JS), + json.dumps({"fn": fn, "args": list(args)}), + ], + capture_output=True, + text=True, + ) + assert completed.returncode == 0, completed.stderr + return json.loads(completed.stdout) + + +# --------------------------------------------------------------------------- +# videoDetailTable:详情明细表 +# --------------------------------------------------------------------------- + + +def _video(name: str, status: str) -> dict: + """构造一条明细数据(字段与 GET /api/batch/jobs/{id} 返回的一致)。""" + return { + "id": f"bv_{name}", + "video_path": f"/videos/{name}", + "status": status, + "error": None, + "finals": {}, + } + + +def test_detail_table_lists_processing_and_completed_videos() -> None: + """正常明细:待处理与已完成的视频都出现在表格行里。""" + # 数据:一个待处理、一个已完成。 + videos = [_video("a.mp4", "PENDING"), _video("c.mp4", "COMPLETED")] + + # 测试过程 + html = _run("videoDetailTable", "batch_1", videos) + + # 验证结果 + assert "a.mp4" in html + assert "c.mp4" in html + assert "PENDING" in html and "COMPLETED" in html + + +def test_detail_table_hides_skipped_videos() -> None: + """详情列表不展示 SKIPPED(视频旁已有字幕、本次未处理)的视频行。""" + # 数据:待处理、跳过、完成各一个。 + videos = [ + _video("a.mp4", "PENDING"), + _video("b.mp4", "SKIPPED"), + _video("c.mp4", "COMPLETED"), + ] + + # 测试过程 + html = _run("videoDetailTable", "batch_1", videos) + + # 验证结果:跳过的那行完全不出现。 + assert "a.mp4" in html + assert "c.mp4" in html + assert "b.mp4" not in html + assert "SKIPPED" not in html + + +def test_detail_table_with_only_skipped_shows_hint() -> None: + """整批视频都已被跳过时给出提示,而不是渲染空表格。""" + # 数据:两个 SKIPPED。 + videos = [_video("a.mp4", "SKIPPED"), _video("b.mp4", "SKIPPED")] + + # 测试过程 + html = _run("videoDetailTable", "batch_1", videos) + + # 验证结果:无表格行,只提示无待处理视频。 + assert "a.mp4" not in html and "b.mp4" not in html + assert " None: + """处理中的视频显示阶段(第几阶段/共几阶段 + 中文标签)。""" + # 数据:一个处理到第二阶段的视频(字段与 GET /api/batch/jobs/{id} 一致)。 + video = _video("a.mp4", "RUNNING") + video.update({"stage_label": "转写", "stage_index": 2, "stage_total": 4}) + + # 测试过程 + html = _run("videoDetailTable", "batch_1", [video]) + + # 验证结果:表头与单元格都带阶段信息。 + assert "阶段" in html + assert "阶段 2/4 · 转写" in html + + +def test_detail_table_shows_placeholder_without_stage() -> None: + """未开始的视频阶段列显示占位符(不报错)。""" + # 数据:一个待处理视频(没有阶段字段)。 + video = _video("a.mp4", "PENDING") + + # 测试过程 + html = _run("videoDetailTable", "batch_1", [video]) + + # 验证结果:阶段列为占位符。 + assert "阶段" in html + assert "阶段 " not in html diff --git a/web/assets/batch.js b/web/assets/batch.js index 5e6cdfd..487e143 100644 --- a/web/assets/batch.js +++ b/web/assets/batch.js @@ -132,12 +132,13 @@ function PathBase(path) { return String(path).split(/[\\/]/).pop() || path; } -// 拉取任务详情并渲染视频明细表:文件、状态、错误与产物下载链接。 -async function videoDetailHtml(jobId) { - const job = await api(`/api/batch/jobs/${encodeURIComponent(jobId)}`); - const videos = job.videos || []; - if (!videos.length) return "(暂无视频)"; - const rows = videos +// 渲染视频明细表:文件、状态、阶段、错误与产物下载链接。 +// 只列本批真正处理过的视频:SKIPPED 表示视频旁已有字幕、本次未被处理, +// 展示出来会让用户误以为它被处理过。 +function videoDetailTable(jobId, videos) { + const pending = (videos || []).filter((video) => video.status !== "SKIPPED"); + if (!pending.length) return '无待处理视频'; + const rows = pending .map((video) => { const finals = video.finals || {}; const links = Object.keys(finals) @@ -146,16 +147,27 @@ async function videoDetailHtml(jobId) { `${escapeHtml(alias)}`, ) .join(" "); + // 阶段来自该视频 run 的当前节点(例:阶段 2/4 · 转写);未开始显示 -。 + const stage = video.stage_label + ? `阶段 ${video.stage_index}/${video.stage_total} · ${video.stage_label}` + : "-"; return ` ${escapeHtml(PathBase(video.video_path))} ${badge(video.status)} + ${escapeHtml(stage)} ${escapeHtml(video.error || "-")} ${links || "-"} `; }) .join(""); - return `${rows}
视频状态错误产物
`; + return `${rows}
视频状态阶段错误产物
`; +} + +// 拉取任务详情并渲染视频明细表。 +async function videoDetailHtml(jobId) { + const job = await api(`/api/batch/jobs/${encodeURIComponent(jobId)}`); + return videoDetailTable(jobId, job.videos || []); } // 暂停批量任务:当前 run 在分块/帧边界停下,后续视频不再开始。 @@ -351,3 +363,8 @@ document.addEventListener("DOMContentLoaded", async () => { await loadBatchJobs(); setInterval(loadBatchJobs, 2000); }); + +// 供测试在 Node 中调用纯渲染函数(浏览器下无 module,不执行)。 +if (typeof module !== "undefined" && module.exports) { + module.exports = { videoDetailTable }; +}