feat: 批量分块流水线、本地模型显存让渡与任务列表分工

批量引擎改为「分块流水线」:视频按 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 用例)。
This commit is contained in:
2026-09-18 10:31:52 +08:00
parent 7a7212f70c
commit dcdc5e8604
25 changed files with 1606 additions and 192 deletions
+160 -90
View File
@@ -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 根目录(<work_dir>/runs/<run_id>/)。
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 阶段结束时统一释放本地模型显存,让下一组的
# whisperASR)拿到 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,
中间态落在 <work_dir>/runs/<run_id>/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"):
# 保留产物记录只置 QUEUEDexecute_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)
# ------------------------------------------------------------------
# 收尾:产物放置与过程文件清理
+5
View File
@@ -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"
# 清理扫描周期(秒),默认每小时一次。
+29 -11
View File
@@ -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。
批量任务若停在 RUNNINGnext_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
+11 -4
View File
@@ -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}")
+59 -3
View File
@@ -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: 文件名},前端据此渲染下载链接;阶段信息来自该视频
自己的 runPENDING/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")
+29 -4
View File
@@ -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)