Files
vrsub/src/wov_app/batch.py
T

602 lines
28 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""文件夹批量处理引擎。
本地版核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹下的
全部视频,逐个调用现有的工作流流水线(复用 WorkflowScheduler 的 DAG 执行与
断点续跑逻辑)。
处理约定(2026-09 起):
- **创建任务时一次性定位**:`create_job` 扫描文件夹并把每个视频登记为
batch_videos 明细;视频所在目录(视频旁)若已存在**文件名包含视频名**的
字幕文件(`.srt/.ass/.ssa/.vtt`),说明该视频已有字幕,直接记为 SKIPPED,
不为它触发任何流水线。运行时(BatchWorker)只消费已定位好的明细列表,
**不再重新扫描文件夹**(运行期间新增/删除的视频不会改变本次任务的范围)。
- **产物放在视频旁**:每个视频处理完成后,把工作流 `final_outputs` 对应的
最终产物文件(字幕流水线即中文 `.srt` 与双目 `.ass`**复制一份到视频的
所在目录**,与 .mp4 放在一起;文件名**对齐媒体库既有约定**:中文字幕存为
`<视频名>.CN.srt`、双目字幕存为 `<视频名>.CN_dual_eye.ass`(文件名稳定且
含视频主名,媒体库可自动匹配,下次批量扫描也会命中"已有字幕"规则跳过)。
- **过程文件清理**:视频收尾完成后删除该视频的整个工作空间与 run 记录,
中间产物(音频/分块/帧图/节点产物)不残留在媒体库,也不会被影视库软件
当作视频载入。工作空间位于应用私有目录 `storage/batch/<job_id>/<bv_id>/`
与用户的视频库目录天然隔离;暂停/失败的视频保留工作空间以便断点续跑。
暂停/恢复语义(对应前端"暂停/继续"按钮):
- 暂停批量任务:把批量任务置为 PAUSED,并暂停当前正在执行的 run(写
paused.flagwhisper 按分块、OCR 按帧检查后停止),处理中的视频保持
PAUSED,后续视频不再开始。
- 继续:批量任务恢复 QUEUED,引擎从断点继续——PAUSED 视频的 run 显式
resume 后由 execute_run 从产物表断点续跑,已完成节点不重复执行。
主调度器不会抢占批量 runnext_queued_run 排除 source=batch),批量引擎
使用 per-video 的 WorkflowScheduler 实例,storage 指向该视频的私有工作空间。
"""
from __future__ import annotations
import json
import shutil
import threading
import time
import uuid
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.db import Database
from wov_app.logging import get_logger
from wov_app.scheduler import WorkflowScheduler
from wov_app.storage import atomic_copy
from wov_sdk.models import WorkflowDefinition
# 批量引擎运行日志:任务进度、视频逐个处理与暂停/续跑等状态变化。
logger = get_logger("batch")
# 识别为视频文件的扩展名(大小写不敏感,扫描时统一转小写比较)。
VIDEO_EXTENSIONS = {
".mp4", ".mkv", ".avi", ".mov", ".webm", ".flv", ".ts",
".m4v", ".wmv", ".mpg", ".mpeg", ".3gp",
}
# 字幕文件扩展名:批量扫描时按它识别"视频旁已有字幕";处理后放回视频旁的
# 最终产物(.srt/.ass)也在该集合内,保证下次扫描能命中同一规则直接跳过。
SUBTITLE_EXTENSIONS = {".srt", ".ass", ".ssa", ".vtt"}
# 暂停信号文件名:与节点约定一致,位于 run 根目录(<work_dir>/runs/<run_id>/)。
PAUSE_FLAG = "paused.flag"
# 旧版批量完成标记文件名(位于视频同名文件夹根目录)。新逻辑不再写入该标记
# (产物直接放视频旁、靠旁挂字幕文件识别完成);仍保留读取能力,用于兼容
# 旧版任务在详情/下载接口中展示产物。
MARKER_NAME = "batch.done.json"
# 批量处理私有工作空间根目录:位于应用存储目录下(data/storage/batch)。
# 每个视频的工作目录为 <根>/<job_id>/<bv_id>/,与用户视频库目录完全隔离,
# 媒体库软件只会看到最终放到视频旁的 .srt/.ass 字幕成品。
BATCH_WORK_ROOT = STORAGE_DIR / "batch"
def _now_iso() -> str:
"""返回当前 UTC 时间的 ISO 格式字符串。"""
return datetime.now(timezone.utc).isoformat()
def scan_videos(folder: Path, recursive: bool = True) -> list[Path]:
"""扫描文件夹下的全部视频文件,按路径排序保证处理顺序确定。
recursive=True 时递归扫描子文件夹;recursive=False 只扫描顶层。
"""
if recursive:
paths = [
p for p in folder.rglob("*")
if p.is_file() and p.suffix.lower() in VIDEO_EXTENSIONS
]
else:
paths = [
p for p in folder.glob("*")
if p.is_file() and p.suffix.lower() in VIDEO_EXTENSIONS
]
return sorted(paths)
def list_sidecar_subtitles(video: Path) -> list[Path]:
"""列出视频所在目录(视频旁)与视频"对应"的字幕文件。
判定规则:与视频同一目录、扩展名为字幕格式、且文件名包含视频主名
(大小写不敏感)的文件都视为该视频已带的字幕。典型命中如
`movie.srt`、`movie.CN.srt`、`movie.CN_dual_eye.ass`,以及本引擎处理
完成后放到视频旁的 `movie.CN.srt` / `movie.CN_dual_eye.ass`。
视频主名过短(单个字符)时只接受"主名."前缀,避免 a.mp4 误配 apple.srt。
目录不可读时保守返回空列表,不影响批量任务创建。
"""
stem = video.stem.lower()
try:
siblings = list(video.parent.iterdir())
except OSError:
return []
found: list[Path] = []
for item in siblings:
try:
if not item.is_file() or item.suffix.lower() not in SUBTITLE_EXTENSIONS:
continue
except OSError:
# 单个子项不可读(权限不足)时跳过,不拖垮整个目录。
continue
item_name = item.name.lower()
# 视频本身不是字幕扩展名,此处无需再排除同名文件;直接按主名匹配。
if len(stem) <= 1:
matched = item_name.startswith(stem + ".")
else:
matched = stem in item_name
if matched:
found.append(item)
return sorted(found)
def _sidecar_product_name(video: Path, source: Path) -> str:
"""把最终产物映射为放在视频旁时的标准字幕文件名。
对齐媒体库既有约定(文件名稳定、无时间戳,媒体库可按视频主名自动匹配):
- 中文 `.srt` 产物 → `<视频名>.CN.srt`
- 双目 `.ass` 产物 → `<视频名>.CN_dual_eye.ass`
- 其余扩展名的最终产物保留原文件名(含时间戳),避免误改语义。
注:当前内置字幕工作流(demo / zh-direct)在视频旁放置的 `.srt/.ass` 即
中文字幕与双目字幕;若未来出现"非中文 .srt"类产物需在此按产物来源区分。
"""
suffix = source.suffix.lower()
if suffix == ".srt":
return f"{video.stem}.CN.srt"
if suffix == ".ass":
return f"{video.stem}.CN_dual_eye.ass"
return source.name
def load_marker(work_dir: Path) -> dict | None:
"""读取同名文件夹里的旧版完成标记;不存在或损坏时返回 None。"""
path = work_dir / MARKER_NAME
if not path.is_file():
return None
try:
data = json.loads(path.read_text(encoding="utf-8"))
return data if isinstance(data, dict) else None
except (json.JSONDecodeError, OSError):
# 半行写入或权限异常时保守视为无标记,走旁挂字幕判定。
return None
def remove_job_workspace(job_id: str) -> None:
"""删除批量任务在应用私有存储下的工作空间目录(<storage>/batch/<job_id>)。
每个视频完成后工作空间已被逐视频清理;此处兜底清理任务级残留(删除任务
或任务异常中止时)。只作用于私有工作空间,绝不触碰用户视频目录。
"""
shutil.rmtree(BATCH_WORK_ROOT / job_id, ignore_errors=True)
def create_job(
db: Database,
folder_path: str,
workflow_id: str,
recursive: bool = True,
) -> str:
"""创建批量任务:校验文件夹与工作流、**一次性定位**视频并登记明细。
扫描到的每个视频都会登记为 batch_videos 明细:视频旁已有对应字幕文件
的直接记 SKIPPED(不触发流水线),否则记 PENDING(等待引擎处理)。
引擎运行时只消费这批已定位的明细,不再重新扫描文件夹。
校验失败抛出 ValueError(由路由层转为 422 响应)。
"""
folder = Path(folder_path).expanduser()
if not folder.is_dir():
raise ValueError("folder not found")
workflow = db.get_workflow(workflow_id)
if workflow is None or not workflow["published"]:
raise ValueError("published workflow not found")
if db.get_latest_workflow_version(workflow_id) is None:
raise ValueError("workflow has no version")
videos = scan_videos(folder, recursive)
if not videos:
raise ValueError("no videos found in folder")
job_id = f"batch_{uuid.uuid4().hex[:12]}"
now = _now_iso()
db.create_batch_job({
"id": job_id,
"folder_path": str(folder),
"workflow_id": workflow_id,
"recursive": int(recursive),
"status": "QUEUED",
"progress": 0,
"total": 0,
"done": 0,
"failed": 0,
"current_video": None,
"error": None,
"created_at": now,
"updated_at": now,
})
pending = 0
skipped = 0
for video in videos:
# 视频所在目录已存在对应字幕文件 → 已处理过,直接跳过不触发流水线。
if list_sidecar_subtitles(video):
status = "SKIPPED"
skipped += 1
else:
status = "PENDING"
pending += 1
video_id = f"bv_{uuid.uuid4().hex[:12]}"
db.create_batch_video({
"id": video_id,
"job_id": job_id,
"video_path": str(video),
# 私有工作空间:storage/batch/<job_id>/<bv_id>/,与媒体库隔离。
"work_dir": str(BATCH_WORK_ROOT / job_id / video_id),
"run_id": None,
"status": status,
"error": None,
"created_at": now,
"updated_at": now,
})
# total = 本批真正需要处理(无字幕)的视频数;已有字幕被 SKIPPED 的
# 不计入总数也不计入完成数——进度条只反映"实际待处理"的这批。
if pending == 0:
# 整批都已有字幕、无任何可处理项:直接视为完成,不排队。
db.update_batch_job(
job_id, status="COMPLETED", total=0, progress=1.0,
updated_at=_now_iso(),
)
else:
db.update_batch_job(job_id, total=pending, updated_at=_now_iso())
logger.info(
"创建批量任务 %s: 文件夹 %s, 工作流 %s, 共 %d 个视频(%d 待处理, %d 已有字幕跳过)",
job_id, folder, workflow_id, len(videos), pending, skipped,
)
return job_id
class BatchWorker:
"""批量处理引擎:单线程轮询 QUEUED 批量任务,逐视频调用现有调度器执行。"""
def __init__(
self,
db: Database,
interval_seconds: float | None = None,
) -> None:
"""保存数据库依赖并初始化轮询线程控制字段。"""
self.db = db
self.interval_seconds = interval_seconds or BATCH_INTERVAL_SECONDS
self._thread: threading.Thread | None = None
self._stopping = False
def start(self) -> None:
"""启动批量处理线程;重复调用无副作用。"""
if self._thread is not None:
return
# 独立运行时确保节点已注册;重复注册幂等。
registry.register_all()
self._stopping = False
self._thread = threading.Thread(
target=self._loop,
name="wov-batch-worker",
daemon=True,
)
self._thread.start()
def stop(self) -> None:
"""请求停止并等待轮询线程退出。"""
self._stopping = True
if self._thread is not None:
self._thread.join(timeout=5)
self._thread = None
def _loop(self) -> None:
"""轮询循环:有排队中的批量任务就处理,否则休眠一个间隔。"""
while not self._stopping:
try:
job = self.db.next_queued_batch_job()
if job is not None:
self._process_job(job)
else:
time.sleep(self.interval_seconds)
except Exception: # noqa: BLE001
# 单次轮询异常不杀死线程,记录后跳过本轮(与主调度器一致)。
logger.exception("批量引擎轮询异常,跳过本轮")
time.sleep(self.interval_seconds)
# ------------------------------------------------------------------
# 任务执行
# ------------------------------------------------------------------
def _process_job(self, job: dict) -> None:
"""处理一个批量任务:校验、逐个消费已定位的视频并放置产物。
job 以 QUEUED 状态进入,处理期间置 RUNNING;全部视频处理完置
COMPLETED;被暂停时保持 PAUSED;校验失败置 FAILED。
"""
job_id = str(job["id"])
try:
self._run_job(job_id)
except Exception as exc: # noqa: BLE001
# 任务级兜底:任何未捕获异常都记录到任务而不是卡死在 RUNNING。
logger.exception("批量任务 %s 处理异常", job_id)
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 统一处理)。
只消费创建任务时已定位好的 batch_videos 明细:SKIPPED/COMPLETED 直接
跳过,PENDING(含失败/暂停后恢复的)逐个交给 _process_video 处理,
**不再扫描文件夹**补视频。
"""
job = self.db.get_batch_job(job_id)
if job is None:
return
folder = Path(job["folder_path"])
if not folder.is_dir():
self.db.update_batch_job(job_id, status="FAILED", error="folder not found", updated_at=_now_iso())
return
workflow = self.db.get_workflow(job["workflow_id"])
if workflow is None or not workflow["published"]:
self.db.update_batch_job(job_id, status="FAILED", error="workflow not found or unpublished", updated_at=_now_iso())
return
version = self.db.get_latest_workflow_version(job["workflow_id"])
if version is None:
self.db.update_batch_job(job_id, status="FAILED", error="workflow has no version", updated_at=_now_iso())
return
definition = WorkflowDefinition.from_dict(version["definition"])
definition.validate()
items = self.db.list_batch_videos(job_id)
# total 在创建任务时已固定为"无字幕需处理的视频数",这里不覆盖;
# 老任务(历史口径 total=全部视频数)由 sync_batch_job_progress 在读取
# 时自我修正为不含 SKIPPED 的口径。
self.db.update_batch_job(
job_id, status="RUNNING", progress=0,
current_video=None, error=None, updated_at=_now_iso(),
)
# 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":
# 停下前把已完成/失败的视频实时入账,让暂停中的前端也能看到
# 真实进度(回归 batch_fee668175444)。
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
# 全部视频处理完成:先用明细实时对齐汇总(done 只计实际完成的,
# 不含 SKIPPED),再置 COMPLETED。
#
# 置 COMPLETED 前必须校验**没有未处理完的视频残留**:若本轮循环因
# 视频处理中断/异常(_process_video 返回但视频仍 PENDING,等同进程
# 在处理中被杀)而没有真正处理完所有 PENDING,就**不能**标完成——
# 否则会出现"明细还有 N 个待处理、任务却已完成"的僵尸状态
# batch_969fabe74b83 等 3 个任务真实发生:引擎串行处理到 9.9GB
# 大视频时中断,10 个视频留 PENDING 却被无条件置 COMPLETED)。
# 此时保持 RUNNING,让引擎下一轮(重启后重新拾起 RUNNING 任务)
# 继续处理剩余 PENDING,全部结束才真正置 COMPLETED。
self.db.sync_batch_job_progress(job_id)
job = self.db.get_batch_job(job_id)
if job is None:
return
leftovers = [
v for v in self.db.list_batch_videos(job_id)
if v["status"] not in ("SKIPPED", "COMPLETED", "FAILED")
]
if leftovers:
# 有未处理完的视频(PENDING/PAUSED/QUEUED 等):保持 RUNNING
# 由引擎下一轮续跑;记录日志便于排查中断位置。
logger.warning(
"批量任务 %s 仍有 %d 个视频未处理完(%s…),保持 RUNNING 待续跑,不置 COMPLETED",
job_id, len(leftovers), Path(leftovers[0]["video_path"]).name,
)
return
done = int(job["done"]) if job else 0
failed = int(job["failed"]) if job else 0
self.db.update_batch_job(
job_id, status="COMPLETED", progress=1.0,
current_video=None, error=None, updated_at=_now_iso(),
)
logger.info(
"批量任务 %s 完成: 待处理 %d 个视频, 完成 %d, 失败 %d",
job_id, total, done, failed,
)
def _process_video(
self,
job: dict,
item: dict,
version: dict,
definition: WorkflowDefinition,
work_dir: Path,
) -> None:
"""处理单个视频:建 run(复用现有调度器)执行,成功后收尾清理。
per-video 的 WorkflowScheduler 以该视频的私有工作空间为 storage,
中间态落在 <work_dir>/runs/<run_id>/steps/ 下;产物表记录全部节点
输出,暂停后续跑从产物表重建已完成节点(断点续跑)。视频成功后
把最终产物复制到视频旁并删除工作空间(见 _finalize_video)。
"""
video = Path(item["video_path"])
work_dir.mkdir(parents=True, exist_ok=True)
run_id = item.get("run_id")
if run_id is not None and self.db.get_run(run_id) is None:
# 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({
"id": run_id,
"workflow_id": job["workflow_id"],
"workflow_version": int(version["version"]),
"status": "QUEUED",
"progress": 0,
"input_uri": str(video),
"param_overrides": None,
"source": "batch",
"created_at": now,
"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":
self.db.resume_run(run_id, _now_iso())
elif run["status"] == "FAILED":
# 失败重跑:**保留**已完成节点的产物记录,只恢复 QUEUED——
# execute_run 从产物表重建已完成节点并跳过,只重跑失败节点。
# 不再 reset_run 清空产物:extract/ocr 等长耗时节点的成果会被
# 白白丢弃重做(run_e2b74e89e232 实测 22222 帧 OCR)。
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)
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 _place_products(self, run_id: str, video: Path, definition: WorkflowDefinition) -> list[str]:
"""把最终产物文件复制到视频所在目录(视频旁),返回放置的文件名。
只为 `final_outputs` 声明的最终产物放置副本:字幕流水线的产物即中文
`.srt` 与双目 `.ass`,按库内约定命名(见 _sidecar_product_name),
文件名稳定且含视频主名——媒体库按主名匹配字幕,下次批量扫描也会命中
"已有字幕"规则跳过该视频。同名目标直接覆盖:可能是上一次运行/旧工作流
留下的旧内容,应以本次产物为准。
"""
# 先校验全部必需输出,缺文件时不覆盖任何视频旁成品,更不能继续清理。
products: list[tuple[Path, Path]] = []
for alias in definition.final_outputs:
artifact = self.db.get_artifact(run_id, alias)
if artifact is None:
raise ValueError(f"missing final artifact: {alias}")
source = Path(artifact["uri"])
if not source.is_file():
raise ValueError(f"missing final artifact file: {alias} ({source})")
target = video.parent / _sidecar_product_name(video, source)
products.append((source, target))
placed: list[str] = []
for source, target in products:
# 稳定目标名 + 原子替换:复制失败保留旧成品,异常交调用方记录 FAILED。
# 多文件中途失败仍保留完整工作空间,下次可以幂等地重新放置。
atomic_copy(source, target)
placed.append(target.name)
return placed
def _finalize_video(
self,
item: dict,
run_id: str,
video: Path,
work_dir: Path,
definition: WorkflowDefinition,
) -> None:
"""视频成功处理后的收尾:产物放视频旁、清理 run 记录与过程文件。
顺序:先复制最终产物(失败则保持现状可重试),再删除 run 与产物
记录(产物已复制到视频旁不再依赖原文件),最后删除整个工作空间
(音频/分块/帧图等过程文件一并清除,防止影视库把切片数据当视频)。
"""
placed = self._place_products(run_id, video, definition)
self.db.delete_run(run_id)
self.db.update_batch_video(item["id"], status="COMPLETED", error=None, run_id=None, updated_at=_now_iso())
shutil.rmtree(work_dir, ignore_errors=True)
logger.info(
"视频 %s 处理完成,产物已放视频旁: %s,过程文件已清理",
video.name, ", ".join(placed) if placed else "(无)",
)
# ------------------------------------------------------------------
# 暂停/继续
# ------------------------------------------------------------------
def pause_job(self, job_id: str) -> None:
"""暂停批量任务:停止当前 run 与后续视频处理。
先把任务置为 PAUSED(引擎在视频间检查后停下),再暂停所有排队/运行
中的 run 并写 paused.flagwhisper 按分块、OCR 按帧检查后中止)。
"""
self.db.update_batch_job(job_id, status="PAUSED", updated_at=_now_iso())
for item in self.db.list_batch_videos(job_id):
if not item.get("run_id"):
continue
run = self.db.get_run(item["run_id"])
if run is None or run["status"] not in ("QUEUED", "RUNNING"):
continue
self.db.pause_run(item["run_id"], _now_iso())
run_dir = Path(item["work_dir"]) / "runs" / item["run_id"]
run_dir.mkdir(parents=True, exist_ok=True)
(run_dir / PAUSE_FLAG).write_text("", encoding="utf-8")
def resume_job(self, job_id: str) -> None:
"""继续批量任务:置回 QUEUED,引擎从断点续跑(PAUSED 视频逐个 resume)。"""
self.db.update_batch_job(job_id, status="QUEUED", updated_at=_now_iso())