"""文件夹批量处理引擎。 本地版核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹下的 全部视频,逐个调用现有的工作流流水线(复用 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///`, 与用户的视频库目录天然隔离;暂停/失败的视频保留工作空间以便断点续跑。 暂停/恢复语义(对应前端"暂停/继续"按钮): - 暂停批量任务:把批量任务置为 PAUSED,并暂停当前正在执行的 run(写 paused.flag,whisper 按分块、OCR 按帧检查后停止),处理中的视频保持 PAUSED,后续视频不再开始。 - 继续:批量任务恢复 QUEUED,引擎从断点继续——PAUSED 视频的 run 显式 resume 后由 execute_run 从产物表断点续跑,已完成节点不重复执行。 主调度器不会抢占批量 run(next_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_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 根目录(/runs//)。 PAUSE_FLAG = "paused.flag" # 旧版批量完成标记文件名(位于视频同名文件夹根目录)。新逻辑不再写入该标记 # (产物直接放视频旁、靠旁挂字幕文件识别完成);仍保留读取能力,用于兼容 # 旧版任务在详情/下载接口中展示产物。 MARKER_NAME = "batch.done.json" # 批量处理私有工作空间根目录:位于应用存储目录下(data/storage/batch)。 # 每个视频的工作目录为 <根>///,与用户视频库目录完全隔离, # 媒体库软件只会看到最终放到视频旁的 .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: """删除批量任务在应用私有存储下的工作空间目录(/batch/)。 每个视频完成后工作空间已被逐视频清理;此处兜底清理任务级残留(删除任务 或任务异常中止时)。只作用于私有工作空间,绝不触碰用户视频目录。 """ 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///,与媒体库隔离。 "work_dir": str(BATCH_WORK_ROOT / job_id / video_id), "run_id": None, "status": status, "error": None, "created_at": now, "updated_at": now, }) 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 = len(items) self.db.update_batch_job( job_id, status="RUNNING", total=total, progress=0, current_video=None, error=None, updated_at=_now_iso(), ) for index, item in enumerate(items): # 暂停检查:批量任务被暂停后停止处理后续视频,等待用户继续。 current = self.db.get_batch_job(job_id) if current is None or current["status"] == "PAUSED": # 停下前把已完成的视频实时入账:任务可能已处理多个视频才被暂停, # 若不在暂停边界同步,前端会一直看到 0/总数 0%(回归 batch_fee668175444)。 self.db.sync_batch_job_progress(job_id) logger.info("批量任务 %s 已暂停,停止在视频 %s", job_id, item["video_path"]) return # 已完成/已跳过的视频不再处理(跳过决策在创建任务时已定)。 if item["status"] in ("COMPLETED", "SKIPPED"): # 已完成的视频同样是任务进度的一部分:continue 前实时同步汇总, # 避免长任务(大量 SKIPPED)中途汇总停留在 0。 self.db.sync_batch_job_progress(job_id) 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), progress=index / total if total else 0, 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 # 全部视频处理完成:先用明细实时对齐汇总(含已跳过),再置 COMPLETED。 self.db.sync_batch_job_progress(job_id) job = self.db.get_batch_job(job_id) 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, 中间态落在 /runs//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), 文件名稳定且含视频主名——媒体库按主名匹配字幕,下次批量扫描也会命中 "已有字幕"规则跳过该视频。同名目标直接覆盖:可能是上一次运行/旧工作流 留下的旧内容,应以本次产物为准。 """ placed: list[str] = [] video.parent.mkdir(parents=True, exist_ok=True) for alias in definition.final_outputs: artifact = self.db.get_artifact(run_id, alias) if artifact is None: continue source = Path(artifact["uri"]) if not source.is_file(): continue target = video.parent / _sidecar_product_name(video, source) # 目标名稳定 → 直接覆盖写入,避免旧同名产物被"大小一致复用"误保留。 shutil.copy2(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.flag(whisper 按分块、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())