From d1557eca46b905453f430e79b3756f45716ec5e6 Mon Sep 17 00:00:00 2001 From: catShark <1716967236@qq.com> Date: Thu, 3 Sep 2026 08:22:42 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=89=B9=E9=87=8F=E5=A4=84=E7=90=86?= =?UTF-8?q?=E9=87=8D=E5=81=9A=E2=80=94=E2=80=94=E8=A7=86=E9=A2=91=E6=97=81?= =?UTF-8?q?=E5=B7=B2=E6=9C=89=E5=AD=97=E5=B9=95=E5=8D=B3=E8=B7=B3=E8=BF=87?= =?UTF-8?q?=E3=80=81=E4=BA=A7=E7=89=A9=E5=AF=B9=E9=BD=90=20CN=20=E5=91=BD?= =?UTF-8?q?=E5=90=8D=E5=B9=B6=E6=B8=85=E7=90=86=E8=BF=87=E7=A8=8B=E6=96=87?= =?UTF-8?q?=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 创建批量任务时一次性定位视频:视频所在目录存在文件名含视频名的字幕文件 (.srt/.ass/.ssa/.vtt)直接记 SKIPPED,不触发流水线;运行时只消费已定位 的明细,不再重新扫描文件夹。 - 视频完成后把最终产物放到视频旁,命名对齐媒体库约定:中文字幕存为 <视频名>.CN.srt、双目字幕存为 <视频名>.CN_dual_eye.ass;其余扩展名产物 保留原文件名。 - 收尾删除 run 记录与过程工作空间;工作空间改到应用私有目录 storage/batch///,与用户媒体库隔离,防止媒体库把切片数据 当视频入库。 - 批量 API 产物清单/下载改为解析视频旁字幕文件,旧版 batch.done.json 语义 别名保持兼容;删除任务时清理私有工作空间。 - 前端说明与创建提示同步;测试按新语义重写并补覆盖(278 passed,100% 行覆盖率)。 --- AGENTS.md | 63 +-- src/wov_app/batch.py | 298 ++++++++------ src/wov_app/routers/batch.py | 64 ++- tests/conftest.py | 3 +- tests/test_batch.py | 753 ++++++++++++++++++++++------------- web/assets/batch.js | 6 +- web/batch.html | 5 +- 7 files changed, 754 insertions(+), 438 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index f5089c9..1eafb8d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -276,31 +276,38 @@ http://127.0.0.1:8000/docs API 文档 (线程池 `on_progress` 回调,每任务完成触发); - whisper:分块转写打印"分块 X/Y 完成 offset=... 耗时 Zs (Nx 实时, 累计 ...s)"; - frame-extract:ffmpeg `-progress` 输出解析 `frame=N`,打印"抽帧进度 X/Y 帧 (Z 帧/s)"。 -## 文件夹批量处理(2026-08) +## 文件夹批量处理(2026-09 更新) 本地版核心能力:**不把视频上传到工作目录**,直接读取用户所选文件夹下的全部 视频,逐个执行所选流水线。入口为批量处理页(`web/batch.html`,导航"批量处理"), 后端为 `src/wov_app/batch.py` 的 `BatchWorker`(单线程轮询线程,处理 `source=batch` 的运行,与主调度器互不抢占)与 `routers/batch.py`。 -- **路径选择(2026-08 起不用手敲路径)**:批量页点击"选择文件夹…"按钮弹出 - 目录树选择器(懒加载),选完回填只读路径框。浏览器拿不到所选文件夹的绝对 - 路径,因此由**本地后端**提供目录浏览:`GET /api/batch/roots`(Windows 盘符 / - POSIX 根 + 家目录)、`GET /api/batch/dirs?path=`(列直接子目录,隐藏目录 - 过滤;不存在/不可读返回空列表不报 500)。只暴露目录名,不返回文件内容。 -- **数据落盘**:每个视频的中间态(`runs//steps/...`)与最终产物都放在 - **视频所在目录的同名文件夹**(`movie.mp4` → `movie/`,去掉扩展名, - `work_dir_for` 推导);最终产物在任务完成后从节点产物目录**复制**到同名 - 文件夹根目录,并写 `batch.done.json` 完成标记(记录 workflow_id 与产物文件名)。 -- **已处理过的不再处理**:同名文件夹已有**同工作流**完成标记且产物文件齐全 → - 直接 SKIPPED;不同工作流的标记不互相误判(换流水线会重新处理)。 -- **任务参数**:`POST /api/batch/jobs {folder, workflow_id, recursive}` 创建批量 - 任务(校验文件夹/已发布工作流/有版本/至少一个视频,失败 422);任务入 - `batch_jobs` 表,每个视频一行 `batch_videos`(PENDING/RUNNING/PAUSED/ - COMPLETED/FAILED/SKIPPED)。可用任意已发布流水线(demo / zh-direct / - ocr-subtitle 等),前端下拉选择。 -- **执行复用**:每个视频创建一个 `source=batch` 的 run(`input_uri` 直接指向 - 本地视频路径),用 per-video 的 `WorkflowScheduler` 实例(storage=同名文件夹) - 执行——完整复用 DAG 拓扑执行、产物表登记与**断点续跑**逻辑。 +- **路径选择**:批量页点击"选择文件夹…"按钮弹出目录树选择器(懒加载),选完 + 回填只读路径框。浏览器拿不到所选文件夹的绝对路径,因此由**本地后端**提供目录 + 浏览:`GET /api/batch/roots`(Windows 盘符 / POSIX 根 + 家目录)、 + `GET /api/batch/dirs?path=`(列直接子目录,隐藏目录过滤;不存在/不可读返回 + 空列表不报 500)。只暴露目录名,不返回文件内容。 +- **创建任务时一次性定位(2026-09 起)**:`POST /api/batch/jobs {folder, + workflow_id, recursive}` 只扫描一次文件夹并把每个视频登记为 `batch_videos` + 明细(PENDING/RUNNING/PAUSED/COMPLETED/FAILED/SKIPPED)。**视频所在目录 + (视频旁)若已存在文件名含视频名的字幕文件**(`.srt/.ass/.ssa/.vtt`, + `list_sidecar_subtitles` 判定,如 `movie.CN.srt`、`movie.CN_dual_eye.ass`), + 说明该视频已有字幕,创建即记 **SKIPPED**——不为它触发任何流水线。运行时 + `BatchWorker` **只消费这批已定位的明细,不再重新扫描文件夹**(运行期间新增/ + 删除的视频不会改变本次任务的范围)。校验失败(文件夹不存在/未发布工作流/ + 无版本/一个视频都没有)返回 422。 +- **产物放在视频旁**:每个视频处理完成后,把工作流 `final_outputs` 对应的最终 + 产物文件(字幕流水线即中文 `.srt` 与双目 `.ass`)**复制一份到视频所在目录**, + 与 .mp4 放在一起(`_place_products`)。文件名**对齐媒体库既有约定**:中文字幕 + 存为 `<视频名>.CN.srt`、双目字幕存为 `<视频名>.CN_dual_eye.ass`(稳定无时间戳, + `_sidecar_product_name` 映射,其余扩展名产物保留原文件名;同名目标直接覆盖)。 + 文件名含视频主名,下次批量扫描会命中"已有字幕"规则直接跳过该视频。 +- **过程文件清理(2026-09 起)**:视频收尾完成后删除该视频的整个工作空间与 + run 记录(音频/分块/帧图/节点产物不留残),防止媒体库把切片数据当视频入库。 + 工作空间位于**应用私有目录** `data/storage/batch///` + (不再放视频同名文件夹),与用户视频库天然隔离;暂停/失败的视频保留工作空间 + 以便断点续跑。per-video 的 `WorkflowScheduler` 实例以该目录为 storage—— + 完整复用 DAG 拓扑执行、产物表登记与**断点续跑**逻辑。 - **暂停/继续**:`POST /api/batch/jobs/{id}/pause` 把任务置 PAUSED 并暂停当前 run(写 `paused.flag`;whisper **分块间**检查、OCR 逐帧检查后中止,当前节点 执行完才停);`resume` 恢复 QUEUED,引擎从断点继续——PAUSED 视频的 run 显式 @@ -313,13 +320,15 @@ http://127.0.0.1:8000/docs API 文档 execute_run 从产物表跳过已完成节点、只重跑失败节点——extract/ocr 等长耗时 成果不浪费;配合 llm-filter/OCR 的节点级断点存档,失败节点自身也只重判未完成 条目。前端对"部分失败"(COMPLETED 且 failed>0)用红色徽章醒目标示。 -- **产物下载**:`GET /api/batch/jobs/{id}/videos/{vid}/download?alias=result` - 从完成标记解析产物文件并返回(只读同名文件夹根目录)。 -- **孤儿清理保护**:`source=batch` 的运行**跳过**自动清理——产物不在主存储 - 目录下,普通孤儿逻辑会误删记录并连带删除 `input_uri` 的父目录(用户的整个 - 视频文件夹)。 -- **删除任务**:`DELETE /api/batch/jobs/{id}` 只清理数据库记录(含关联 run), - 磁盘上的同名文件夹与产物属于用户数据,保留不删。 +- **产物下载**:`GET /api/batch/jobs/{id}/videos/{vid}/download?alias=<文件名>` + 解析并返回视频旁的字幕文件;旧版 `batch.done.json` 完成标记里的语义别名 + (位于旧 work_dir)仍兼容可下载。详情/创建响应里每个视频的 `finals` 合并上述 + 两处来源。 +- **孤儿清理保护**:`source=batch` 的运行**跳过**自动清理——其 run 位于私有 + `storage/batch/...` 下,普通孤儿逻辑会误判删除,且 `_remove_run` 还会删除 + `input_uri` 的父目录(用户的整个视频文件夹)。 +- **删除任务**:`DELETE /api/batch/jobs/{id}` 清理数据库记录(含关联 run)与 + 应用私有工作空间残留;视频旁已放置的产物属于用户数据,保留不删。 - **环境变量**:`WOV_BATCH_ENABLED`(默认 1)、`WOV_BATCH_INTERVAL_SECONDS` (默认 1.0)。 diff --git a/src/wov_app/batch.py b/src/wov_app/batch.py index 5d000bf..746f1db 100644 --- a/src/wov_app/batch.py +++ b/src/wov_app/batch.py @@ -1,16 +1,25 @@ """文件夹批量处理引擎。 -本地版的核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹 -下的所有视频,逐个调用现有的工作流流水线(复用 WorkflowScheduler 的 DAG -执行与断点续跑逻辑)。 +本地版核心能力:**不把视频上传到工作目录**,而是直接读取用户所选文件夹下的 +全部视频,逐个调用现有的工作流流水线(复用 WorkflowScheduler 的 DAG 执行与 +断点续跑逻辑)。 -数据落盘约定: +处理约定(2026-09 起): -- 每个视频的中间态数据(runs/、chunks/、帧图等)与最终产物都存放在**视频 - 所在目录的同名文件夹**里(movie.mp4 → movie/),源视频目录保持干净。 -- 最终产物(SRT/ASS)在任务完成后从节点产物目录复制到同名文件夹根目录, - 同时写入 `batch.done.json` 完成标记;再次批量处理同一文件夹时,已有完成 - 标记且产物文件齐全的视频直接跳过(已经处理过的不再处理)。 +- **创建任务时一次性定位**:`create_job` 扫描文件夹并把每个视频登记为 + batch_videos 明细;视频所在目录(视频旁)若已存在**文件名包含视频名**的 + 字幕文件(`.srt/.ass/.ssa/.vtt`),说明该视频已有字幕,直接记为 SKIPPED, + 不为它触发任何流水线。运行时(BatchWorker)只消费已定位好的明细列表, + **不再重新扫描文件夹**(运行期间新增/删除的视频不会改变本次任务的范围)。 +- **产物放在视频旁**:每个视频处理完成后,把工作流 `final_outputs` 对应的 + 最终产物文件(字幕流水线即中文 `.srt` 与双目 `.ass`)**复制一份到视频的 + 所在目录**,与 .mp4 放在一起;文件名**对齐媒体库既有约定**:中文字幕存为 + `<视频名>.CN.srt`、双目字幕存为 `<视频名>.CN_dual_eye.ass`(文件名稳定且 + 含视频主名,媒体库可自动匹配,下次批量扫描也会命中"已有字幕"规则跳过)。 +- **过程文件清理**:视频收尾完成后删除该视频的整个工作空间与 run 记录, + 中间产物(音频/分块/帧图/节点产物)不残留在媒体库,也不会被影视库软件 + 当作视频载入。工作空间位于应用私有目录 `storage/batch///`, + 与用户的视频库目录天然隔离;暂停/失败的视频保留工作空间以便断点续跑。 暂停/恢复语义(对应前端"暂停/继续"按钮): @@ -21,7 +30,7 @@ resume 后由 execute_run 从产物表断点续跑,已完成节点不重复执行。 主调度器不会抢占批量 run(next_queued_run 排除 source=batch),批量引擎 -使用 per-video 的 WorkflowScheduler 实例,storage 指向同名文件夹。 +使用 per-video 的 WorkflowScheduler 实例,storage 指向该视频的私有工作空间。 """ from __future__ import annotations @@ -35,7 +44,7 @@ from datetime import datetime, timezone from pathlib import Path from wov_app import registry -from wov_app.config import BATCH_INTERVAL_SECONDS +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 @@ -50,13 +59,23 @@ VIDEO_EXTENSIONS = { ".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 格式字符串。""" @@ -81,16 +100,60 @@ def scan_videos(folder: Path, recursive: bool = True) -> list[Path]: return sorted(paths) -def work_dir_for(video: Path) -> Path: - """返回视频的同名文件夹:去掉扩展名,位于视频所在目录。 +def list_sidecar_subtitles(video: Path) -> list[Path]: + """列出视频所在目录(视频旁)与视频"对应"的字幕文件。 - 例如 movie.mp4 → 旁边的 movie/ 文件夹,中间态与最终产物都放这里。 + 判定规则:与视频同一目录、扩展名为字幕格式、且文件名包含视频主名 + (大小写不敏感)的文件都视为该视频已带的字幕。典型命中如 + `movie.srt`、`movie.CN.srt`、`movie.CN_dual_eye.ass`,以及本引擎处理 + 完成后放到视频旁的 `movie.CN.srt` / `movie.CN_dual_eye.ass`。 + 视频主名过短(单个字符)时只接受"主名."前缀,避免 a.mp4 误配 apple.srt。 + 目录不可读时保守返回空列表,不影响批量任务创建。 """ - return video.parent / video.stem + 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。""" + """读取同名文件夹里的旧版完成标记;不存在或损坏时返回 None。""" path = work_dir / MARKER_NAME if not path.is_file(): return None @@ -98,32 +161,17 @@ def load_marker(work_dir: Path) -> dict | None: data = json.loads(path.read_text(encoding="utf-8")) return data if isinstance(data, dict) else None except (json.JSONDecodeError, OSError): - # 半行写入或权限异常时保守视为未完成,允许重新处理。 + # 半行写入或权限异常时保守视为无标记,走旁挂字幕判定。 return None -def ensure_videos(db: Database, job_id: str, videos: list[Path]) -> None: - """为扫描到的视频补齐 batch_videos 明细;已存在的记录保持不变。 +def remove_job_workspace(job_id: str) -> None: + """删除批量任务在应用私有存储下的工作空间目录(/batch/)。 - 同一批量任务反复处理(暂停/续跑)时保留每个视频的状态,已完成的不重置。 + 每个视频完成后工作空间已被逐视频清理;此处兜底清理任务级残留(删除任务 + 或任务异常中止时)。只作用于私有工作空间,绝不触碰用户视频目录。 """ - existing = {item["video_path"] for item in db.list_batch_videos(job_id)} - now = _now_iso() - for video in videos: - path = str(video) - if path in existing: - continue - db.create_batch_video({ - "id": f"bv_{uuid.uuid4().hex[:12]}", - "job_id": job_id, - "video_path": path, - "work_dir": str(work_dir_for(video)), - "run_id": None, - "status": "PENDING", - "error": None, - "created_at": now, - "updated_at": now, - }) + shutil.rmtree(BATCH_WORK_ROOT / job_id, ignore_errors=True) def create_job( @@ -132,10 +180,12 @@ def create_job( workflow_id: str, recursive: bool = True, ) -> str: - """创建批量任务:校验文件夹与工作流、扫描视频、落库明细,返回任务 ID。 + """创建批量任务:校验文件夹与工作流、**一次性定位**视频并登记明细。 - 校验失败抛出 ValueError(由路由层转为 422 响应);扫描到的视频全部 - 登记为 batch_videos 明细,引擎轮询到该任务后逐个处理。 + 扫描到的每个视频都会登记为 batch_videos 明细:视频旁已有对应字幕文件 + 的直接记 SKIPPED(不触发流水线),否则记 PENDING(等待引擎处理)。 + 引擎运行时只消费这批已定位的明细,不再重新扫描文件夹。 + 校验失败抛出 ValueError(由路由层转为 422 响应)。 """ folder = Path(folder_path).expanduser() if not folder.is_dir(): @@ -166,10 +216,36 @@ def create_job( "created_at": now, "updated_at": now, }) - ensure_videos(db, job_id, videos) - logger.info("创建批量任务 %s: 文件夹 %s, 工作流 %s, 视频 %d 个", job_id, folder, workflow_id, len(videos)) + 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 批量任务,逐视频调用现有调度器执行。""" @@ -224,7 +300,7 @@ class BatchWorker: # ------------------------------------------------------------------ def _process_job(self, job: dict) -> None: - """处理一个批量任务:校验、扫描、补齐明细、逐视频执行并复制产物。 + """处理一个批量任务:校验、逐个消费已定位的视频并放置产物。 job 以 QUEUED 状态进入,处理期间置 RUNNING;全部视频处理完置 COMPLETED;被暂停时保持 PAUSED;校验失败置 FAILED。 @@ -238,7 +314,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 处理, + **不再扫描文件夹**补视频。 + """ job = self.db.get_batch_job(job_id) if job is None: return @@ -258,10 +339,6 @@ class BatchWorker: definition = WorkflowDefinition.from_dict(version["definition"]) definition.validate() - # 扫描当前文件夹的视频,为新增视频补齐明细(已有明细保留原状态, - # 保证暂停/续跑时已完成与进行中的视频不被重置)。 - self._ensure_videos(job_id, scan_videos(folder, bool(job["recursive"]))) - items = self.db.list_batch_videos(job_id) total = len(items) self.db.update_batch_job( @@ -276,7 +353,7 @@ class BatchWorker: logger.info("批量任务 %s 已暂停,停止在视频 %s", job_id, item["video_path"]) return - # 已完成/已跳过的视频不再处理。 + # 已完成/已跳过的视频不再处理(跳过决策在创建任务时已定)。 if item["status"] in ("COMPLETED", "SKIPPED"): continue @@ -286,11 +363,6 @@ class BatchWorker: continue work_dir = Path(item["work_dir"]) - # 同名文件夹里已有同工作流的完成标记且产物齐全 → 直接跳过。 - if self._is_done(work_dir, job["workflow_id"]): - self.db.update_batch_video(item["id"], status="SKIPPED", updated_at=_now_iso()) - continue - self.db.update_batch_job( job_id, current_video=str(video), progress=index / total if total else 0, @@ -303,8 +375,8 @@ class BatchWorker: logger.exception("批量任务 %s 视频 %s 处理异常", job_id, video) self.db.update_batch_video(item["id"], status="FAILED", error=str(exc), updated_at=_now_iso()) - # 重新读取视频明细:_process_video 可能刚创建 run(快照里 run_id - # 还是 None),必须取最新记录才能拿到 run_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 @@ -334,15 +406,19 @@ class BatchWorker: definition: WorkflowDefinition, work_dir: Path, ) -> None: - """处理单个视频:建 run(复用现有调度器)并执行,完成后复制最终产物。 + """处理单个视频:建 run(复用现有调度器)执行,成功后收尾清理。 - per-video 的 WorkflowScheduler 以同名文件夹为 storage,中间态落在 - /runs//steps/ 下;产物表记录全部节点输出,暂停后 - 续跑从产物表重建已完成节点(断点续跑)。 + 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 动态组装节点执行。 @@ -363,10 +439,9 @@ class BatchWorker: 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._copy_finals(run_id, work_dir, definition, job["workflow_id"]) - self.db.update_batch_video(item["id"], status="COMPLETED", error=None, updated_at=_now_iso()) + self._finalize_video(item, run_id, video, work_dir, definition) return # 暂停的 run 显式 resume 回 QUEUED,由 execute_run 从产物表断点续跑。 if run["status"] == "PAUSED": @@ -374,8 +449,8 @@ class BatchWorker: elif run["status"] == "FAILED": # 失败重跑:**保留**已完成节点的产物记录,只恢复 QUEUED—— # execute_run 从产物表重建已完成节点并跳过,只重跑失败节点。 - # 不再 reset_run 清空产物:extract/ocr 等长耗时节点的成果(如 - # ABP-885 的 22222 帧 OCR)会被白白丢弃重做(run_e2b74e89e232 实测)。 + # 不再 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 续跑。 @@ -387,48 +462,31 @@ class BatchWorker: 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._copy_finals(run_id, work_dir, definition, job["workflow_id"]) - self.db.update_batch_video(item["id"], status="COMPLETED", error=None, updated_at=_now_iso()) + # 放置最终产物到视频旁并清理过程文件。 + 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 _ensure_videos(self, job_id: str, videos: list[Path]) -> None: - """为扫描到的视频补齐 batch_videos 明细;已存在的记录保持不变。""" - ensure_videos(self.db, job_id, videos) + def _place_products(self, run_id: str, video: Path, definition: WorkflowDefinition) -> list[str]: + """把最终产物文件复制到视频所在目录(视频旁),返回放置的文件名。 - def _is_done(self, work_dir: Path, workflow_id: str) -> bool: - """判断同名文件夹是否已完成当前工作流的处理。 - - 完成标记记录 workflow_id 与最终产物文件名;只有工作流一致且产物文件 - 全部存在时才视为已处理(不同工作流的产物不互相误判为完成)。 + 只为 `final_outputs` 声明的最终产物放置副本:字幕流水线的产物即中文 + `.srt` 与双目 `.ass`,按库内约定命名(见 _sidecar_product_name), + 文件名稳定且含视频主名——媒体库按主名匹配字幕,下次批量扫描也会命中 + "已有字幕"规则跳过该视频。同名目标直接覆盖:可能是上一次运行/旧工作流 + 留下的旧内容,应以本次产物为准。 """ - marker = load_marker(work_dir) - if marker is None or marker.get("workflow_id") != workflow_id: - return False - finals = marker.get("finals") or {} - return bool(finals) and all((work_dir / name).is_file() for name in finals.values()) - - def _copy_finals( - self, - run_id: str, - work_dir: Path, - definition: WorkflowDefinition, - workflow_id: str, - ) -> None: - """把最终产物从节点目录复制到同名文件夹根目录,并写完成标记。 - - 调度器收尾时已把产物重命名为 上传文件名.标识.时间戳(如 - movie.zh-CN.20260819120000.srt),这里原样复制,文件名保留辨识度。 - """ - work_dir.mkdir(parents=True, exist_ok=True) - finals: dict[str, str] = {} + 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: @@ -436,21 +494,33 @@ class BatchWorker: source = Path(artifact["uri"]) if not source.is_file(): continue - target = work_dir / source.name - # 已存在的产物直接复用,避免重复复制。 - if not target.is_file() or target.stat().st_size != source.stat().st_size: - shutil.copy2(source, target) - finals[alias] = source.name - marker = { - "workflow_id": workflow_id, - "workflow_version": definition.version, - "run_id": run_id, - "completed_at": _now_iso(), - "finals": finals, - } - (work_dir / MARKER_NAME).write_text( - json.dumps(marker, ensure_ascii=False, indent=2), - encoding="utf-8", + 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 "(无)", ) # ------------------------------------------------------------------ diff --git a/src/wov_app/routers/batch.py b/src/wov_app/routers/batch.py index 34f72fd..a076d47 100644 --- a/src/wov_app/routers/batch.py +++ b/src/wov_app/routers/batch.py @@ -33,15 +33,31 @@ def _get_worker() -> batch_engine.BatchWorker | None: return getattr(app.state, "batch", None) -def _enrich_videos(db: Database, videos: list[dict]) -> list[dict]: - """为每个视频补充最终产物清单(从同名文件夹的完成标记读取)。 +def _product_finals(video: dict) -> dict[str, str]: + """列出该视频可下载的最终产物(键为下载 alias,值为文件名)。 - finals 形如 {别名: 文件名}(如 {"cn_srt": "movie.zh-CN.20260819.srt"}), - 前端据此渲染下载链接;未完成的视频没有产物。 + 来源合并两处: + - 旧版完成标记 `batch.done.json`(位于 work_dir/同名文件夹),键为语义 + 别名(如 cn_srt/ass),用于兼容旧版批量任务; + - 视频所在目录(视频旁)中**文件名含视频名**的字幕文件,键即文件名。 + 新版处理完成后产物放到视频旁,靠旁挂字幕文件即可列出与下载。 + """ + finals: dict[str, str] = {} + marker = batch_engine.load_marker(Path(video["work_dir"])) + if marker: + finals.update(marker.get("finals") or {}) + for sidecar in batch_engine.list_sidecar_subtitles(Path(video["video_path"])): + finals[sidecar.name] = sidecar.name + return finals + + +def _enrich_videos(db: Database, videos: list[dict]) -> list[dict]: + """为每个视频补充最终产物清单(视频旁的字幕文件 + 旧版完成标记)。 + + finals 形如 {alias: 文件名},前端据此渲染下载链接;未完成的视频没有产物。 """ for video in videos: - marker = batch_engine.load_marker(Path(video["work_dir"])) - video["finals"] = (marker or {}).get("finals") or {} + video["finals"] = _product_finals(video) return videos @router.post("/api/batch/jobs") @@ -111,8 +127,8 @@ def resume_batch_job( def delete_batch_job(job_id: str, db: Database = Depends(_get_db)) -> dict: """删除批量任务:移除任务、明细记录与关联的 run 记录。 - 磁盘上的同名文件夹与产物属于用户数据,保留不删(与任务页删除接口的 - 行为不同),只清理数据库记录。 + 只清理数据库记录与应用私有工作空间(storage/batch/)的残留; + 视频旁已经放置的产物属于用户数据,保留不删。 """ if db.get_batch_job(job_id) is None: raise HTTPException(status_code=404, detail="batch job not found") @@ -120,6 +136,8 @@ def delete_batch_job(job_id: str, db: Database = Depends(_get_db)) -> dict: if item.get("run_id"): db.delete_run(item["run_id"]) db.delete_batch_job(job_id) + # 清理任务在应用私有存储下的工作空间残留(视频完成后已逐视频清理)。 + batch_engine.remove_job_workspace(job_id) return {"deleted": job_id} @@ -130,22 +148,34 @@ def download_batch_video( alias: str = Query(...), db: Database = Depends(_get_db), ) -> FileResponse: - """下载视频的最终产物:从同名文件夹的完成标记解析产物文件名后返回。 + """下载视频的最终产物:解析 alias 对应的文件后返回。 - alias 为工作流 final_outputs 的别名(如 cn_srt / ass);只有完成标记 - 中记录且文件真实存在的产物才可下载。 + alias 解析顺序: + 1. 旧版完成标记里的语义别名(如 cn_srt/ass)→ 文件位于 work_dir; + 2. 视频旁(视频所在目录)文件名含视频名的字幕文件名 → 直接返回该文件。 + 只有存在且文件真实落盘的产物才可下载。 """ video = db.get_batch_video(video_id) if video is None or video["job_id"] != job_id: raise HTTPException(status_code=404, detail="video not found") - marker = batch_engine.load_marker(Path(video["work_dir"])) - if marker is None or alias not in (marker.get("finals") or {}): + work_dir = Path(video["work_dir"]) + marker = batch_engine.load_marker(work_dir) + target: Path | None = None + if marker and alias in (marker.get("finals") or {}): + candidate = work_dir / marker["finals"][alias] + if candidate.is_file(): + target = candidate + else: + # 新版/既有字幕:alias 是视频旁的字幕文件名。 + for sidecar in batch_engine.list_sidecar_subtitles(Path(video["video_path"])): + if sidecar.name == alias: + target = sidecar + break + if target is None: raise HTTPException(status_code=404, detail="artifact not found") - path = Path(video["work_dir"]) / marker["finals"][alias] - if not path.is_file(): + if not target.is_file(): raise HTTPException(status_code=404, detail="artifact file missing") - return FileResponse(path, filename=path.name) - + return FileResponse(target, filename=target.name) # --------------------------------------------------------------------------- # 本地目录浏览(目录树选择器) diff --git a/tests/conftest.py b/tests/conftest.py index 10653bd..2fc5635 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -22,7 +22,8 @@ os.environ["WOV_STORAGE_DIR"] = str(TEST_ROOT / "storage") os.environ["WOV_AUTO_SEED"] = "0" os.environ["WOV_SCHEDULER_ENABLED"] = "0" os.environ["WOV_CLEANUP_ENABLED"] = "0" - +# 默认关闭批量引擎后台线程:API 测试手工控制执行时机,避免后台线程与断言竞态。 +os.environ["WOV_BATCH_ENABLED"] = "0" def _cleanup() -> None: """进程退出时清理临时测试目录。""" diff --git a/tests/test_batch.py b/tests/test_batch.py index 9c56aa7..d0ea150 100644 --- a/tests/test_batch.py +++ b/tests/test_batch.py @@ -1,9 +1,9 @@ """文件夹批量处理引擎与 API 测试。 -覆盖:视频扫描/同名文件夹推导、完成标记读写、批量任务创建校验、 -引擎逐视频执行(建 run、复用现有调度器、复制最终产物、跳过已处理)、 -暂停/继续断点续跑、失败视频不中断、孤儿清理不删批量运行、 -批量 API 全端点(创建/列表/详情/暂停/继续/删除/下载)与 404/422 分支。 +覆盖:视频扫描与旁挂字幕判定、创建任务时一次性定位(已有字幕 → SKIPPED)、 +运行时只消费已定位明细、产物复制到视频旁与过程文件清理、暂停/继续断点续跑、 +失败视频不中断、批量 API 全端点(创建/列表/详情/暂停/继续/删除/下载/目录树) +与 404/422 分支。 """ import json @@ -16,16 +16,20 @@ from fastapi.testclient import TestClient from wov_app import batch as batch_engine from wov_app.batch import ( + BATCH_WORK_ROOT, MARKER_NAME, PAUSE_FLAG, BatchWorker, create_job, + list_sidecar_subtitles, load_marker, + remove_job_workspace, scan_videos, - work_dir_for, ) +from wov_app.config import STORAGE_DIR from wov_app.db import Database from wov_app.main import app +from wov_sdk.models import WorkflowDefinition def _now_iso() -> str: @@ -72,13 +76,13 @@ def _seed_echo_workflow(db: Database, workflow_id: str = "echo-app", published: def _register_echo() -> None: """把内置 echo 节点注册到进程内注册表(conftest 每个测试隔离注册表)。""" from wov_app import registry + from nodes.echo import invoke from wov_sdk.models import NodeManifest root = Path(__file__).resolve().parent.parent - from nodes.echo import invoke - registry.register(NodeManifest.load(str(root / "manifests" / "echo.json")), invoke) + def _video_folder(tmp_path, names=("a.mp4", "b.mp4")) -> Path: """创建含视频文件的文件夹:内容为真实文本(echo 节点按文本读入)。""" folder = tmp_path / "videos" @@ -94,13 +98,45 @@ def _make_job(db: Database, folder: Path, workflow_id: str = "echo-app", recursi return db.get_batch_job(job_id) +def _create_run( + db: Database, + run_id: str, + folder: Path, + workflow_id: str = "echo-app", + status: str = "QUEUED", + name: str = "a", +) -> None: + """创建一条 source=batch 的运行记录(input 指向 folder 下的视频文件)。""" + db.create_run( + { + "id": run_id, + "workflow_id": workflow_id, + "workflow_version": 1, + "status": status, + "progress": 0, + "input_uri": str(folder / f"{name}.mp4"), + "param_overrides": None, + "source": "batch", + "created_at": _now_iso(), + "updated_at": _now_iso(), + } + ) + + +def _find_product(video_dir: Path, stem: str, alias: str) -> Path: + """在同名输出位置查找最终产物文件(文件名以 主名.别名 开头)。""" + matches = list(video_dir.glob(f"{stem}.{alias}.*")) + assert matches, f"未找到产物 {stem}.{alias}.* in {video_dir}" + return matches[0] + + # --------------------------------------------------------------------------- -# 扫描 / 目录推导 / 完成标记 +# 扫描 / 旁挂字幕判定 / 完成标记读取 # --------------------------------------------------------------------------- -def test_scan_videos_recursive_and_work_dir(tmp_path) -> None: - """递归/非递归扫描只返回视频文件,同名文件夹去掉扩展名。""" +def test_scan_videos_recursive_and_top_level(tmp_path) -> None: + """递归/非递归扫描只返回视频文件,隐藏目录与普通文件不参与。""" folder = tmp_path / "media" (folder / "sub").mkdir(parents=True) (folder / "a.mp4").write_text("a", encoding="utf-8") @@ -113,12 +149,84 @@ def test_scan_videos_recursive_and_work_dir(tmp_path) -> None: assert [p.name for p in recursive] == ["a.mp4", "b.MKV", "c.avi"] flat = scan_videos(folder, recursive=False) assert [p.name for p in flat] == ["a.mp4", "b.MKV"] - # 同名文件夹:去掉扩展名,位于视频所在目录。 - assert work_dir_for(folder / "sub" / "c.avi") == folder / "sub" / "c" -def test_load_marker_variants(tmp_path) -> None: - """完成标记缺失/损坏/非字典时返回 None。""" +def test_list_sidecar_subtitles_matches(tmp_path) -> None: + """视频旁"文件名含视频名"的字幕文件全部命中;不相关/非字幕/别的目录不算。""" + folder = tmp_path / "media" + (folder / "sub").mkdir(parents=True) + video = folder / "movie.mp4" + video.write_text("v", encoding="utf-8") + # 命中:同目录、字幕扩展名、文件名含视频主名(大小写不敏感)。 + (folder / "movie.srt").write_text("s", encoding="utf-8") + (folder / "movie.CN.srt").write_text("s", encoding="utf-8") + (folder / "MOVIE.CN_dual_eye.ass").write_text("s", encoding="utf-8") + (folder / "movie.zh-CN.20260819120000.srt").write_text("s", encoding="utf-8") + # 不命中:不含视频主名、非字幕扩展名、子目录里的字幕。 + (folder / "other.srt").write_text("s", encoding="utf-8") + (folder / "movie.jpg").write_text("s", encoding="utf-8") + (folder / "sub" / "movie.srt").write_text("s", encoding="utf-8") + + names = [p.name for p in list_sidecar_subtitles(video)] + assert names == [ + "MOVIE.CN_dual_eye.ass", + "movie.CN.srt", + "movie.srt", + "movie.zh-CN.20260819120000.srt", + ] + + +def test_list_sidecar_subtitles_short_stem_and_unrelated(tmp_path) -> None: + """单字符视频主名只接受"主名."前缀,避免 a.mp4 误配 apple.srt。""" + folder = tmp_path / "media" + folder.mkdir() + video = folder / "a.mp4" + video.write_text("v", encoding="utf-8") + (folder / "apple.srt").write_text("s", encoding="utf-8") + (folder / "b.srt").write_text("s", encoding="utf-8") + assert list_sidecar_subtitles(video) == [] + (folder / "a.srt").write_text("s", encoding="utf-8") + assert [p.name for p in list_sidecar_subtitles(video)] == ["a.srt"] + + +def test_list_sidecar_subtitles_handles_unreadable_entries(tmp_path, monkeypatch) -> None: + """目录不可读返回空列表;单个子项不可读时跳过该项不影响其余匹配。""" + from pathlib import Path as RealPath + + # 整目录不可读(iterdir 抛 OSError)→ 保守返回空。 + class _Denied(RealPath): + """iterdir 恒抛权限错误的子类,模拟不可读的视频目录。""" + + def iterdir(self): + raise OSError("denied") + + denied = _Denied(str(tmp_path / "denied")) + assert list_sidecar_subtitles(denied / "a.mp4") == [] + + # 单个子项不可读(is_file 抛 OSError)→ 跳过该项,其余正常返回。 + folder = tmp_path / "media" + folder.mkdir() + video = folder / "a.mp4" + video.write_text("v", encoding="utf-8") + (folder / "a.srt").write_text("s", encoding="utf-8") + + class _Poison(RealPath): + """is_file 恒抛权限错误的子类,模拟不可读的目录项。""" + + def is_file(self): + raise OSError("denied") + + real_iterdir = RealPath.iterdir + + def mixed_iterdir(path): + return list(real_iterdir(path)) + [_Poison(str(folder / "secret"))] + + monkeypatch.setattr(RealPath, "iterdir", mixed_iterdir) + assert [p.name for p in list_sidecar_subtitles(video)] == ["a.srt"] + + +def test_load_marker_variants(tmp_path, monkeypatch) -> None: + """完成标记缺失/损坏/非字典/读取异常时返回 None。""" work = tmp_path / "movie" work.mkdir() assert load_marker(work) is None @@ -127,12 +235,29 @@ def test_load_marker_variants(tmp_path) -> None: (work / MARKER_NAME).write_text("[1,2]", encoding="utf-8") assert load_marker(work) is None (work / MARKER_NAME).write_text('{"workflow_id": "w", "finals": {"r": "m.txt"}}', encoding="utf-8") - marker = load_marker(work) - assert marker["workflow_id"] == "w" + assert load_marker(work)["finals"]["r"] == "m.txt" + + # 读取抛 OSError(权限等)时同样保守视为无标记。 + from pathlib import Path as RealPath + + def broken_read_text(path, **kwargs): + raise OSError("denied") + + monkeypatch.setattr(RealPath, "read_text", broken_read_text) + assert load_marker(work) is None + + +def test_remove_job_workspace(tmp_path) -> None: + """删除任务级私有工作空间;目录不存在时静默无副作用。""" + job_dir = BATCH_WORK_ROOT / "batch_del" + (job_dir / "runs" / "run_x").mkdir(parents=True) + remove_job_workspace("batch_del") + assert not job_dir.exists() + remove_job_workspace("batch_missing") # --------------------------------------------------------------------------- -# create_job 校验 +# create_job:校验 + 一次性定位(已有字幕 → SKIPPED) # --------------------------------------------------------------------------- @@ -160,8 +285,8 @@ def test_create_job_validation_errors(tmp_path) -> None: create_job(db, str(folder), "echo-app") -def test_create_job_success_creates_rows(tmp_path) -> None: - """创建成功:任务入队、视频明细齐全、递归标志与文件夹路径落库。""" +def test_create_job_locates_all_videos_once(tmp_path) -> None: + """创建任务时一次性定位:全部视频登记明细、递归标志落库、可被引擎拾起。""" db = _db(tmp_path) _seed_echo_workflow(db) (tmp_path / "videos" / "sub").mkdir(parents=True) @@ -170,22 +295,47 @@ def test_create_job_success_creates_rows(tmp_path) -> None: job = _make_job(db, tmp_path / "videos", recursive=True) assert job["status"] == "QUEUED" videos = db.list_batch_videos(job["id"]) - assert {v["video_path"] for v in videos} == { - str(tmp_path / "videos" / "a.mp4"), - str(tmp_path / "videos" / "sub" / "b.mkv"), - } + assert {Path(v["video_path"]).name for v in videos} == {"a.mp4", "b.mkv"} assert all(v["status"] == "PENDING" for v in videos) + # 工作空间位于应用私有存储目录(与媒体库隔离),且每个视频独立目录。 + assert all(Path(v["work_dir"]).is_relative_to(BATCH_WORK_ROOT / job["id"]) for v in videos) assert db.next_queued_batch_job()["id"] == job["id"] +def test_create_job_skips_videos_with_existing_subtitles(tmp_path) -> None: + """视频旁已有对应字幕文件的直接记 SKIPPED,其余 PENDING(不触发流水线)。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) + # a 已处理过(视频旁有中文双目字幕),b 未处理。 + (folder / "a.CN_dual_eye.ass").write_text("已处理", encoding="utf-8") + + job = _make_job(db, folder) + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + assert videos["a.mp4"]["status"] == "SKIPPED" + assert videos["a.mp4"]["run_id"] is None + assert videos["b.mp4"]["status"] == "PENDING" + + +def test_create_job_all_videos_skipped_still_created(tmp_path) -> None: + """文件夹里全部视频都已有字幕时任务仍可创建(全部 SKIPPED,不再处理)。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) + (folder / "a.CN.srt").write_text("x", encoding="utf-8") + (folder / "b.CN_dual_eye.ass").write_text("x", encoding="utf-8") + job = _make_job(db, folder) + assert job["status"] == "QUEUED" + assert all(v["status"] == "SKIPPED" for v in db.list_batch_videos(job["id"])) + + # --------------------------------------------------------------------------- -# 引擎:完整处理 / 跳过 / 失败 / 暂停续跑 +# 引擎:完整处理 / 跳过 / 失败 / 暂停续跑 / 收尾清理 # --------------------------------------------------------------------------- -def test_batch_worker_processes_all_videos(tmp_path) -> None: - """批量引擎逐个处理视频:建 source=batch 的 run、中间态进同名文件夹、 - 最终产物复制到同名文件夹根目录并写完成标记,任务最终 COMPLETED。""" +def test_batch_worker_processes_all_videos_and_cleans_up(tmp_path) -> None: + """批量引擎逐个处理视频:产物复制到视频旁、run 记录与过程工作空间清理。""" db = _db(tmp_path) _seed_echo_workflow(db) folder = _video_folder(tmp_path) @@ -197,35 +347,28 @@ def test_batch_worker_processes_all_videos(tmp_path) -> None: assert job["status"] == "COMPLETED" assert job["done"] == 2 and job["failed"] == 0 for name in ("a", "b"): - work = folder / name - marker = load_marker(work) - assert marker is not None and marker["workflow_id"] == "echo-app" - final_name = marker["finals"]["result"] - final = work / final_name - # 最终产物复制到同名文件夹根目录,内容与输入一致(真实数据流)。 - assert final.is_file() - assert final.read_text(encoding="utf-8") == (folder / f"{name}.mp4").read_text(encoding="utf-8") - # 中间态落在 /runs//steps/ 下。 - # 中间态落在 /runs//steps/ 下(收尾时产物已重命名)。 - assert list(work.glob("runs/*/steps/step/*.txt")) - # run 记录为 batch 来源,且主调度器不会抢占。 - runs = [db.get_run(v["run_id"]) for v in db.list_batch_videos(job["id"])] - assert all(run["source"] == "batch" and run["status"] == "COMPLETED" for run in runs) - assert db.next_queued_run() is None + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + row = videos[f"{name}.mp4"] + # 最终产物复制到视频旁(与 .mp4 同目录),内容与输入一致(真实数据流)。 + product = _find_product(folder, name, "result") + assert product.parent == folder + assert product.read_text(encoding="utf-8") == (folder / f"{name}.mp4").read_text(encoding="utf-8") + # 视频状态 COMPLETED 且 run 引用清空;不再写完成标记。 + assert row["status"] == "COMPLETED" + assert row["run_id"] is None + assert not (Path(row["work_dir"]) / MARKER_NAME).exists() + # 过程工作空间(含 runs/steps/音频分块等)已被整体清理,不在媒体库残留。 + assert not Path(row["work_dir"]).exists() + # run 记录已删除(产物已放视频旁,不再依赖原文件)。 + assert db.list_runs() == [] -def test_batch_worker_skips_already_done_videos(tmp_path) -> None: - """已有同工作流完成标记且产物齐全的视频直接跳过(不再处理)。""" +def test_batch_worker_skips_videos_with_existing_subtitles(tmp_path) -> None: + """已有字幕的视频在创建时记 SKIPPED,引擎运行时不再为它触发流水线。""" db = _db(tmp_path) _seed_echo_workflow(db) folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) - work_a = folder / "a" - work_a.mkdir() - # 伪造 a.mp4 的完成标记与最终产物(模拟上次任务已处理完)。 - marker = {"workflow_id": "echo-app", "workflow_version": 1, "run_id": "run_old", "finals": {"result": "a.result.txt"}} - (work_a / MARKER_NAME).write_text(json.dumps(marker), encoding="utf-8") - (work_a / "a.result.txt").write_text("已处理", encoding="utf-8") - + (folder / "a.CN_dual_eye.ass").write_text("已处理", encoding="utf-8") job = _make_job(db, folder) BatchWorker(db, interval_seconds=0.05)._process_job(job) videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} @@ -235,29 +378,22 @@ def test_batch_worker_skips_already_done_videos(tmp_path) -> None: assert db.get_batch_job(job["id"])["done"] == 2 -def test_batch_worker_marker_workflow_mismatch_reprocesses(tmp_path) -> None: - """不同工作流的完成标记不互相误判:换工作流后视频重新处理。""" +def test_batch_worker_all_skipped_job_completes(tmp_path) -> None: + """任务里全部视频都是 SKIPPED 时引擎正常完成,不创建任何 run。""" db = _db(tmp_path) - _seed_echo_workflow(db, workflow_id="echo-app") - _seed_echo_workflow(db, workflow_id="echo-other") - folder = _video_folder(tmp_path, names=("a.mp4",)) - work_a = folder / "a" - work_a.mkdir() - marker = {"workflow_id": "echo-other", "finals": {"result": "x.txt"}} - (work_a / MARKER_NAME).write_text(json.dumps(marker), encoding="utf-8") - (work_a / "x.txt").write_text("x", encoding="utf-8") - - job = _make_job(db, folder, workflow_id="echo-app") + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) + (folder / "a.srt").write_text("x", encoding="utf-8") + (folder / "b.CN.srt").write_text("x", encoding="utf-8") + job = _make_job(db, folder) BatchWorker(db, interval_seconds=0.05)._process_job(job) - video = db.list_batch_videos(job["id"])[0] - # 标记工作流不匹配 → 重新处理并覆盖为新工作流的标记。 - assert video["status"] == "COMPLETED" - new_marker = load_marker(work_a) - assert new_marker["workflow_id"] == "echo-app" + assert db.get_batch_job(job["id"])["status"] == "COMPLETED" + assert db.get_batch_job(job["id"])["done"] == 2 + assert db.list_runs() == [] def test_batch_worker_failed_video_continues(tmp_path) -> None: - """缺失视频文件与失败 run 都记为 FAILED,任务继续处理后续视频。""" + """缺失视频文件记为 FAILED,任务继续处理后续视频。""" db = _db(tmp_path) _seed_echo_workflow(db) folder = _video_folder(tmp_path, names=("gone.mp4", "b.mp4")) @@ -275,6 +411,31 @@ def test_batch_worker_failed_video_continues(tmp_path) -> None: assert job["done"] == 1 and job["failed"] == 1 +def test_batch_worker_process_exception_marks_video_failed(tmp_path, monkeypatch) -> None: + """单视频执行抛异常:该视频 FAILED 带错误信息,任务继续处理后续视频。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) + job = _make_job(db, folder) + + class _BoomScheduler: + """execute_run 直接抛异常的假调度器,触发单视频兜底分支。""" + + def __init__(self, db: Database, work_dir: Path) -> None: + self.db = db + + def execute_run(self, run_id: str) -> None: + raise RuntimeError("boom") + + monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _BoomScheduler) + BatchWorker(db, interval_seconds=0.05)._process_job(job) + videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} + assert videos["a.mp4"]["status"] == "FAILED" + assert "boom" in videos["a.mp4"]["error"] + assert videos["b.mp4"]["status"] == "FAILED" + assert db.get_batch_job(job["id"])["status"] == "COMPLETED" + + class _PausingScheduler: """把 execute_run 模拟为"被暂停"的假调度器:置 run 为 PAUSED。""" @@ -302,7 +463,7 @@ def test_batch_worker_pause_then_resume_continues(tmp_path, monkeypatch) -> None assert videos["b.mp4"]["status"] == "PENDING" assert db.get_batch_job(job["id"])["status"] == "PAUSED" - # 恢复真实调度器并继续:a 从断点完成,b 接着处理,任务 COMPLETED。 + # 恢复真实调度器并继续:a 从断点完成(收尾清理),b 接着处理,任务 COMPLETED。 monkeypatch.undo() worker.resume_job(job["id"]) worker._process_job(db.get_batch_job(job["id"])) @@ -348,58 +509,18 @@ def test_batch_worker_paused_job_does_not_start_new_video(tmp_path, monkeypatch) assert videos["b.mp4"]["run_id"] is None assert db.get_batch_job(job_id)["status"] == "PAUSED" -class _FailingScheduler: - """把 execute_run 模拟为"节点失败"的假调度器。""" - def __init__(self, db: Database, work_dir: Path) -> None: - self.db = db - - def execute_run(self, run_id: str) -> None: - self.db.update_run(run_id, status="FAILED", error="node boom", updated_at=_now_iso()) - - -def test_batch_worker_run_failed_marks_video_failed(tmp_path, monkeypatch) -> None: - """节点执行失败:run FAILED → 视频 FAILED 带错误信息,任务继续。""" +def test_batch_worker_failed_and_running_runs_resume(tmp_path) -> None: + """已失败的 run 重跑、上次进程残留的 RUNNING run 恢复后继续(收尾清理)。""" db = _db(tmp_path) _seed_echo_workflow(db) folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) job = _make_job(db, folder) - monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _FailingScheduler) - BatchWorker(db, interval_seconds=0.05)._process_job(job) videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} - assert videos["a.mp4"]["status"] == "FAILED" - assert videos["a.mp4"]["error"] == "node boom" - assert videos["b.mp4"]["status"] == "FAILED" - assert db.get_batch_job(job["id"])["status"] == "COMPLETED" - assert db.get_batch_job(job["id"])["failed"] == 2 - - -def test_batch_worker_resumes_failed_and_running_runs(tmp_path) -> None: - """已失败的 run 重跑、上次进程残留的 RUNNING run 恢复后继续。""" - db = _db(tmp_path) - _seed_echo_workflow(db) - folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4")) - job = _make_job(db, folder) - now = _now_iso() - # 手动构造 run:a 失败(重跑)、b 残留 RUNNING(进程被杀后恢复)。 for name, status in (("a", "FAILED"), ("b", "RUNNING")): run_id = f"run_{name}" - db.create_run( - { - "id": run_id, - "workflow_id": "echo-app", - "workflow_version": 1, - "status": status, - "progress": 0, - "input_uri": str(folder / f"{name}.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": now, - "updated_at": now, - } - ) - video = next(v for v in db.list_batch_videos(job["id"]) if v["video_path"].endswith(f"{name}.mp4")) - db.update_batch_video(video["id"], run_id=run_id, updated_at=now) + _create_run(db, run_id, folder, status=status, name=name) + db.update_batch_video(videos[f"{name}.mp4"]["id"], run_id=run_id, updated_at=_now_iso()) BatchWorker(db, interval_seconds=0.05)._process_job(db.get_batch_job(job["id"])) videos = {Path(v["video_path"]).name: v for v in db.list_batch_videos(job["id"])} @@ -412,7 +533,8 @@ def test_batch_worker_retry_failed_keeps_completed_node_artifacts(tmp_path) -> N 回归:此前 FAILED 走 reset_run 清空全部产物记录,重跑时 extract/ocr 等 长耗时节点从头重做(run_e2b74e89e232 的 22222 帧 OCR 被白白丢弃)。 - 现改为保留产物恢复 QUEUED,execute_run 从产物表跳过已完成节点。 + 通过改写 step1 产物内容验证:若 step1 被重跑,最终产物会恢复为源视频 + 文本;保留产物则最终产物内容为改写后的内容。 """ db = _db(tmp_path) _register_echo() @@ -432,27 +554,15 @@ def test_batch_worker_retry_failed_keeps_completed_node_artifacts(tmp_path) -> N db.create_workflow_version("echo-app", 1, definition) folder = _video_folder(tmp_path, names=("a.mp4",)) job = _make_job(db, folder) - now = _now_iso() + video = db.list_batch_videos(job["id"])[0] run_id = "run_retry" - db.create_run( - { - "id": run_id, - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "FAILED", - "progress": 0.5, - "error": "节点2炸了", - "input_uri": str(folder / "a.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": now, - "updated_at": now, - } - ) - # 模拟 step1 已完成并登记产物(step2 失败时 step1 的成果)。 - s1_out = folder / "a" / "runs" / run_id / "steps" / "s1" / "echo.txt" + _create_run(db, run_id, folder, status="FAILED", name="a") + db.update_batch_video(video["id"], run_id=run_id, updated_at=_now_iso()) + + # 模拟 step1 已完成并登记产物,把其内容改写为与源视频不同(step2 会引用)。 + s1_out = Path(video["work_dir"]) / "runs" / run_id / "steps" / "s1" / "echo.txt" s1_out.parent.mkdir(parents=True) - s1_out.write_text("step1 产物", encoding="utf-8") + s1_out.write_text("step1 已保留产物", encoding="utf-8") db.create_artifact( { "run_id": run_id, @@ -460,32 +570,62 @@ def test_batch_worker_retry_failed_keeps_completed_node_artifacts(tmp_path) -> N "name": "s1.file_uri", "uri": str(s1_out), "mime_type": "text/plain", - "size": 10, + "size": s1_out.stat().st_size, } ) - video = db.list_batch_videos(job["id"])[0] - db.update_batch_video(video["id"], run_id=run_id, updated_at=now) BatchWorker(db, interval_seconds=0.05)._process_job(db.get_batch_job(job["id"])) video = db.list_batch_videos(job["id"])[0] assert video["status"] == "COMPLETED" - run = db.get_run(run_id) - assert run["status"] == "COMPLETED" and run["error"] is None - # step1 产物记录保留且未被重写(节点没有重新执行)。 - artifacts = {a["name"]: a for a in db.list_artifacts(run_id)} - assert artifacts["s1.file_uri"]["uri"] == str(s1_out) - # step2 本次补做完成;收尾时最终产物被重命名为 <片名>.result.<时间戳>.txt。 - assert "s2.file_uri" in artifacts - assert list((folder / "a" / "runs" / run_id / "steps" / "s2").glob("*")) - finals = {a["name"]: a for a in db.list_artifacts(run_id)} - assert "result" in finals - assert Path(finals["result"]["uri"]).is_file() + # step1 未重跑:最终产物内容等于保留的 step1 产物,而不是源视频文本。 + product = _find_product(folder, "a", "result") + assert product.read_text(encoding="utf-8") == "step1 已保留产物" -def test_batch_worker_already_completed_run_copies_finals(tmp_path) -> None: - """run 已完成但视频未标记(收尾前中断):直接复制产物并标记完成。 - 覆盖 _copy_finals 的三个分支:产物齐全(复制)、产物记录存在但文件丢失 - (跳过)、无产物记录(跳过)——完成后均写出完成标记。 +def test_batch_worker_stale_run_id_recreated(tmp_path) -> None: + """明细里的 run_id 指向已不存在的 run 时按全新视频处理并正常完成。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4",)) + job = _make_job(db, folder) + video = db.list_batch_videos(job["id"])[0] + db.update_batch_video(video["id"], run_id="run_dead", updated_at=_now_iso()) + + BatchWorker(db, interval_seconds=0.05)._process_job(job) + video = db.list_batch_videos(job["id"])[0] + assert video["status"] == "COMPLETED" + assert db.get_run("run_dead") is None + assert _find_product(folder, "a", "result").is_file() + + +def test_batch_worker_run_deleted_by_scheduler_returns(tmp_path, monkeypatch) -> None: + """execute_run 期间 run 记录被删除(极端情况)时直接返回,不误标完成。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4",)) + job = _make_job(db, folder) + + class _DeletingScheduler: + """execute_run 直接删除 run 记录的假调度器。""" + + def __init__(self, db: Database, work_dir: Path) -> None: + self.db = db + + def execute_run(self, run_id: str) -> None: + self.db.delete_run(run_id) + + monkeypatch.setattr("wov_app.batch.WorkflowScheduler", _DeletingScheduler) + BatchWorker(db, interval_seconds=0.05)._process_job(job) + video = db.list_batch_videos(job["id"])[0] + # 不标 COMPLETED/FAILED,保持 PENDING,等待下一次处理重试。 + assert video["status"] == "PENDING" + + +def test_batch_worker_already_completed_runs_finalized(tmp_path) -> None: + """run 已完成但视频未标记(收尾前中断):直接放置产物、清理并标记完成。 + + 覆盖 _place_products 各分支:产物齐全(复制)、产物记录存在但文件丢失 + (跳过)、无产物记录(跳过)——完成后均清理工作空间与 run 记录。 """ db = _db(tmp_path) _seed_echo_workflow(db) @@ -496,35 +636,24 @@ def test_batch_worker_already_completed_run_copies_finals(tmp_path) -> None: def complete_run(run_id: str, name: str, with_file: bool, with_artifact: bool) -> None: """创建 COMPLETED run;可选产物文件与产物记录。""" - db.create_run( - { - "id": run_id, - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "COMPLETED", - "progress": 1, - "input_uri": str(folder / f"{name}.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": now, - "updated_at": now, - } - ) + _create_run(db, run_id, folder, status="COMPLETED", name=name) + work_dir = Path(videos[name]["work_dir"]) if with_artifact: + path = work_dir / "runs" / run_id / "steps" / "step" / f"{name}.result.txt" db.create_artifact( { "run_id": run_id, "node_id": "step", "name": "result", - "uri": str(folder / name / "runs" / run_id / "steps" / "step" / f"{name}.result.txt"), + "uri": str(path), "mime_type": "text/plain", "size": 6, } ) if with_file: - path = folder / name / "runs" / run_id / "steps" / "step" / f"{name}.result.txt" + path = work_dir / "runs" / run_id / "steps" / "step" / f"{name}.result.txt" path.parent.mkdir(parents=True, exist_ok=True) - path.write_text("产物", encoding="utf-8") + path.write_text(f"产物 {name}", encoding="utf-8") db.update_batch_video(videos[name]["id"], run_id=run_id, updated_at=now) # a:产物齐全;b:产物记录存在但文件丢失;c:没有任何产物记录。 @@ -535,19 +664,79 @@ def test_batch_worker_already_completed_run_copies_finals(tmp_path) -> None: BatchWorker(db, interval_seconds=0.05)._process_job(db.get_batch_job(job["id"])) videos = {Path(v["video_path"]).stem: v for v in db.list_batch_videos(job["id"])} assert all(videos[name]["status"] == "COMPLETED" for name in ("a", "b", "c")) - # a 的产物被复制到同名文件夹根目录;b/c 无产物可复制,标记为空。 - marker_a = load_marker(folder / "a") - assert marker_a["finals"]["result"] == "a.result.txt" - assert (folder / "a" / "a.result.txt").is_file() - assert load_marker(folder / "b")["finals"] == {} - assert load_marker(folder / "c")["finals"] == {} + # a 的产物复制到视频旁;b/c 无产物可复制。 + product = _find_product(folder, "a", "result") + assert product.read_text(encoding="utf-8") == "产物 a" + assert not list(folder.glob("b.result.*")) and not list(folder.glob("c.result.*")) + # 三个视频的工作空间与 run 记录都被清理。 + for name in ("a", "b", "c"): + assert not Path(videos[name]["work_dir"]).exists() + assert db.get_run(f"run_done_{name}") is None + assert db.list_runs() == [] -# --------------------------------------------------------------------------- -# 引擎:任务级异常与校验失败 -# --------------------------------------------------------------------------- - +def test_batch_worker_place_products_branches(tmp_path) -> None: + """_place_products:无产物记录/文件丢失跳过;.srt/.ass 按库内约定命名覆盖; + 其他扩展名保留原文件名复制。""" + db = _db(tmp_path) + _seed_echo_workflow(db) + folder = _video_folder(tmp_path, names=("a.mp4",)) + video = folder / "a.mp4" + run_id = "run_place" + _create_run(db, run_id, folder, name="a") + source_dir = tmp_path / "src" + source_dir.mkdir() + worker = BatchWorker(db, interval_seconds=0.05) + definition = { + "name": "multi-final", + "version": 1, + "nodes": [], + "edges": [], + "entry_inputs": {"video_uri": "file"}, + # 无产物记录 / 文件丢失 / srt / ass / txt 五种 final 场景。 + "final_outputs": { + "alias_none": "x.none", + "alias_lost": "x.lost", + "alias_srt": "step.srt", + "alias_ass": "step.ass", + "alias_txt": "step.txt", + }, + } + # alias_lost:产物记录存在但指向的文件已丢失 → 跳过不放置。 + db.create_artifact( + {"run_id": run_id, "node_id": "step", "name": "alias_lost", "uri": str(source_dir / "lost.srt"), "mime_type": "text/plain", "size": 0} + ) + # alias_srt:中文 srt → <视频名>.CN.srt;目标已存在(旧内容)→ 覆盖。 + src_srt = source_dir / "raw_any_name.srt" + src_srt.write_text("1\n00:00:00,000 --> 00:00:01,000\n新中文字幕\n", encoding="utf-8") + db.create_artifact( + {"run_id": run_id, "node_id": "step", "name": "alias_srt", "uri": str(src_srt), "mime_type": "text/plain", "size": src_srt.stat().st_size} + ) + target_srt = video.parent / "a.CN.srt" + target_srt.write_text("旧字幕内容", encoding="utf-8") + # alias_ass:双目 ass → <视频名>.CN_dual_eye.ass;覆盖同名旧文件。 + src_ass = source_dir / "whatever.ass" + src_ass.write_text("[Script Info]\n新双目字幕\n", encoding="utf-8") + db.create_artifact( + {"run_id": run_id, "node_id": "step", "name": "alias_ass", "uri": str(src_ass), "mime_type": "text/plain", "size": src_ass.stat().st_size} + ) + target_ass = video.parent / "a.CN_dual_eye.ass" + target_ass.write_text("旧 ass", encoding="utf-8") + # alias_txt:其他扩展名保留原文件名(含时间戳命名)复制。 + src_txt = source_dir / "a.other.20260902120000.txt" + src_txt.write_text("其余产物内容", encoding="utf-8") + db.create_artifact( + {"run_id": run_id, "node_id": "step", "name": "alias_txt", "uri": str(src_txt), "mime_type": "text/plain", "size": src_txt.stat().st_size} + ) + placed = worker._place_products(run_id, video, WorkflowDefinition.from_dict(definition)) + assert placed == ["a.CN.srt", "a.CN_dual_eye.ass", "a.other.20260902120000.txt"] + # srt/ass 被改名为标准名且覆盖旧文件;txt 保留原名。 + assert target_srt.read_text(encoding="utf-8") == "1\n00:00:00,000 --> 00:00:01,000\n新中文字幕\n" + assert target_ass.read_text(encoding="utf-8") == "[Script Info]\n新双目字幕\n" + assert (video.parent / "a.other.20260902120000.txt").read_text(encoding="utf-8") == "其余产物内容" + # alias_none(无记录)与 alias_lost(文件丢失)未放置。 + assert not (video.parent / "lost.srt").exists() def test_batch_worker_job_validation_failures(tmp_path) -> None: """文件夹缺失/工作流未发布/无版本时任务置为 FAILED 并记录错误。""" db = _db(tmp_path) @@ -638,7 +827,6 @@ def test_batch_worker_loop_processes_and_survives_exceptions(tmp_path, monkeypat # 不存在的任务 ID:_run_job 直接返回,不报错。 BatchWorker(db, interval_seconds=0.05)._process_job({"id": "ghost_job"}) # 第一次轮询抛异常(模拟数据库抖动),后续正常。 - # 第一次轮询抛异常(模拟数据库抖动),后续正常。 calls = {"n": 0} real_next = db.next_queued_batch_job @@ -670,50 +858,34 @@ def test_batch_worker_pause_job_writes_flag_and_pauses_run(tmp_path) -> None: _seed_echo_workflow(db) folder = _video_folder(tmp_path, names=("a.mp4", "b.mp4", "c.mp4")) job = _make_job(db, folder) - now = _now_iso() videos = {Path(v["video_path"]).stem: v for v in db.list_batch_videos(job["id"])} # a:QUEUED 运行(应被暂停并写 flag);b:run 记录不存在;c:已完成的 run。 - db.create_run( - { - "id": "run_pause_me", - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "QUEUED", - "progress": 0, - "input_uri": str(folder / "a.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": now, - "updated_at": now, - } - ) - db.create_run( - { - "id": "run_done_c", - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "COMPLETED", - "progress": 1, - "input_uri": str(folder / "c.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": now, - "updated_at": now, - } - ) - db.update_batch_video(videos["a"]["id"], run_id="run_pause_me", updated_at=now) - db.update_batch_video(videos["b"]["id"], run_id="run_ghost", updated_at=now) - db.update_batch_video(videos["c"]["id"], run_id="run_done_c", updated_at=now) + _create_run(db, "run_pause_me", folder, status="QUEUED", name="a") + _create_run(db, "run_done_c", folder, status="COMPLETED", name="c") + db.update_batch_video(videos["a"]["id"], run_id="run_pause_me", updated_at=_now_iso()) + db.update_batch_video(videos["b"]["id"], run_id="run_ghost", updated_at=_now_iso()) + db.update_batch_video(videos["c"]["id"], run_id="run_done_c", updated_at=_now_iso()) worker = BatchWorker(db, interval_seconds=0.05) worker.pause_job(job["id"]) assert db.get_batch_job(job["id"])["status"] == "PAUSED" assert db.get_run("run_pause_me")["status"] == "PAUSED" - assert (folder / "a" / "runs" / "run_pause_me" / PAUSE_FLAG).is_file() + work_dir_a = Path(videos["a"]["work_dir"]) + assert (work_dir_a / "runs" / "run_pause_me" / PAUSE_FLAG).is_file() # b 的 run 不存在、c 的 run 已完成:都被跳过,不写 flag。 - assert not (folder / "b" / "runs" / "run_ghost" / PAUSE_FLAG).exists() - assert not (folder / "c" / "runs" / "run_done_c" / PAUSE_FLAG).exists() + assert not (Path(videos["b"]["work_dir"]) / "runs" / "run_ghost" / PAUSE_FLAG).exists() + assert not (Path(videos["c"]["work_dir"]) / "runs" / "run_done_c" / PAUSE_FLAG).exists() + + # 继续:任务回到 QUEUED,由引擎从断点续跑。 + worker.resume_job(job["id"]) + assert db.get_batch_job(job["id"])["status"] == "QUEUED" + + +def test_batch_worker_ghost_job_id_returns(tmp_path) -> None: + """_run_job 在任务不存在时直接返回(幽灵任务处理无副作用)。""" + db = _db(tmp_path) + BatchWorker(db, interval_seconds=0.05)._run_job("batch_ghost") # --------------------------------------------------------------------------- @@ -749,9 +921,11 @@ def _client_with_echo_workflow(tmp_path): def test_batch_api_create_list_detail(tmp_path) -> None: - """批量 API:创建任务、列表、详情(含视频明细与完成标记产物)。""" + """批量 API:创建任务、列表、详情(含 SKIPPED 状态与视频旁产物清单)。""" client, folder = _client_with_echo_workflow(tmp_path) try: + # a 视频旁已有字幕(创建即 SKIPPED),b 待处理。 + (folder / "a.CN_dual_eye.ass").write_text("已有字幕", encoding="utf-8") response = client.post( "/api/batch/jobs", json={"folder": str(folder), "workflow_id": "echo-app", "recursive": True}, @@ -759,8 +933,11 @@ def test_batch_api_create_list_detail(tmp_path) -> None: assert response.status_code == 200 job = response.json() assert job["status"] == "QUEUED" - assert len(job["videos"]) == 2 - assert job["videos"][0]["finals"] == {} + by_name = {Path(v["video_path"]).name: v for v in job["videos"]} + assert by_name["a.mp4"]["status"] == "SKIPPED" + assert by_name["a.mp4"]["finals"] == {"a.CN_dual_eye.ass": "a.CN_dual_eye.ass"} + assert by_name["b.mp4"]["status"] == "PENDING" + assert by_name["b.mp4"]["finals"] == {} listed = client.get("/api/batch/jobs").json() assert any(item["id"] == job["id"] for item in listed) @@ -801,7 +978,7 @@ def test_batch_api_creation_errors(tmp_path) -> None: def test_batch_api_pause_resume_delete(tmp_path) -> None: - """批量 API:暂停/继续切换任务状态,删除清理数据库记录(含 run)。""" + """批量 API:暂停/继续切换任务状态,删除清理 DB 记录与私有工作空间。""" client, folder = _client_with_echo_workflow(tmp_path) try: job = client.post( @@ -819,28 +996,20 @@ def test_batch_api_pause_resume_delete(tmp_path) -> None: # 给第一个视频挂一个 run,验证删除任务时级联删除 run 记录。 db = app.state.db - db.create_run( - { - "id": "run_del", - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "QUEUED", - "progress": 0, - "input_uri": str(folder / "a.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": _now_iso(), - "updated_at": _now_iso(), - } - ) video = job["videos"][0] + _create_run(db, "run_del", folder, name=Path(video["video_path"]).stem) db.update_batch_video(video["id"], run_id="run_del", updated_at=_now_iso()) assert db.get_run("run_del") is not None + # 模拟任务遗留的私有工作空间:删除任务时一并清理。 + workspace = BATCH_WORK_ROOT / job["id"] + (workspace / "runs" / "run_del").mkdir(parents=True) + (workspace / "runs" / "run_del" / "paused.flag").write_text("", encoding="utf-8") deleted = client.delete(f"/api/batch/jobs/{job['id']}") assert deleted.status_code == 200 assert client.get(f"/api/batch/jobs/{job['id']}").status_code == 404 assert db.get_run("run_del") is None + assert not workspace.exists() assert client.post("/api/batch/jobs/ghost/pause").status_code == 404 assert client.post("/api/batch/jobs/ghost/resume").status_code == 404 @@ -864,60 +1033,77 @@ def test_batch_api_worker_unavailable(tmp_path, monkeypatch) -> None: client.__exit__(None, None, None) -def test_batch_api_download_artifact(tmp_path) -> None: - """批量 API:从完成标记下载最终产物,缺失别名/文件返回 404。""" +def test_batch_api_download_sidecar_subtitle(tmp_path) -> None: + """批量 API:下载视频旁的字幕文件,缺失别名/文件返回 404。""" + client, folder = _client_with_echo_workflow(tmp_path) + try: + (folder / "a.CN_dual_eye.ass").write_text("已有字幕内容", encoding="utf-8") + job = client.post( + "/api/batch/jobs", + json={"folder": str(folder), "workflow_id": "echo-app"}, + ).json() + video = next(v for v in job["videos"] if Path(v["video_path"]).name == "a.mp4") + + ok = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=a.CN_dual_eye.ass") + assert ok.status_code == 200 + assert ok.content == "已有字幕内容".encode("utf-8") + + bad_alias = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=nope") + assert bad_alias.status_code == 404 + bad_video = client.get(f"/api/batch/jobs/{job['id']}/videos/bv_ghost/download?alias=a.CN_dual_eye.ass") + assert bad_video.status_code == 404 + finally: + client.__exit__(None, None, None) + + +def test_batch_api_download_sidecar_file_missing(tmp_path, monkeypatch) -> None: + """下载时旁挂字幕文件已被外部移除(列表与下载之间的竞态)→ 404。""" + client, folder = _client_with_echo_workflow(tmp_path) + try: + (folder / "a.srt").write_text("x", encoding="utf-8") + job = client.post( + "/api/batch/jobs", + json={"folder": str(folder), "workflow_id": "echo-app"}, + ).json() + video = next(v for v in job["videos"] if Path(v["video_path"]).name == "a.mp4") + + # 模拟列表后文件被删:让旁挂字幕扫描返回一个磁盘上已不存在的路径。 + gone = folder / "a.srt" + gone.unlink() + monkeypatch.setattr(batch_engine, "list_sidecar_subtitles", lambda video_path: [gone]) + missing = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=a.srt") + assert missing.status_code == 404 + assert "file missing" in missing.json()["detail"] + finally: + client.__exit__(None, None, None) + + +def test_batch_api_download_legacy_marker(tmp_path) -> None: + """批量 API:旧版完成标记(batch.done.json)里的语义别名仍可下载。""" client, folder = _client_with_echo_workflow(tmp_path) try: job = client.post( "/api/batch/jobs", json={"folder": str(folder), "workflow_id": "echo-app"}, ).json() - video = job["videos"][0] db = app.state.db - # 直接构造完成状态:run + 产物 + 同名文件夹完成标记。 - # 产物文件名与下载端点约定一致:完成标记里的文件名位于同名文件夹根。 - run_id = "run_dl" - work = folder / "a" - steps = work / "runs" / run_id / "steps" / "step" - steps.mkdir(parents=True) - (steps / "a.result.20260819000000.txt").write_text("下载内容", encoding="utf-8") - (work / "a.result.20260819000000.txt").write_text("下载内容", encoding="utf-8") - db.create_run( - { - "id": run_id, - "workflow_id": "echo-app", - "workflow_version": 1, - "status": "COMPLETED", - "progress": 1, - "input_uri": str(folder / "a.mp4"), - "param_overrides": None, - "source": "batch", - "created_at": _now_iso(), - "updated_at": _now_iso(), - } - ) - db.create_artifact( - { - "run_id": run_id, - "node_id": "step", - "name": "result", - "uri": str(steps / "a.result.20260819000000.txt"), - "mime_type": "text/plain", - "size": 8, - } - ) - marker = {"workflow_id": "echo-app", "workflow_version": 1, "run_id": run_id, "finals": {"result": "a.result.20260819000000.txt"}} + video = job["videos"][0] + # 旧版布局:work_dir 是视频的同名文件夹,里面写 batch.done.json + 产物。 + work = Path(video["work_dir"]) + work.mkdir(parents=True) + marker = {"workflow_id": "echo-app", "workflow_version": 1, "run_id": "run_old", "finals": {"result": "a.result.20260819000000.txt"}} (work / MARKER_NAME).write_text(json.dumps(marker), encoding="utf-8") + (work / "a.result.20260819000000.txt").write_text("下载内容", encoding="utf-8") + + # 详情页产物清单合并旧版完成标记里的语义别名。 + detail = client.get(f"/api/batch/jobs/{job['id']}").json() + video_detail = next(v for v in detail["videos"] if v["id"] == video["id"]) + assert video_detail["finals"]["result"] == "a.result.20260819000000.txt" ok = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=result") assert ok.status_code == 200 assert ok.content == "下载内容".encode("utf-8") - - bad_alias = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=nope") - assert bad_alias.status_code == 404 - bad_video = client.get(f"/api/batch/jobs/{job['id']}/videos/bv_ghost/download?alias=result") - assert bad_video.status_code == 404 - # 产物文件缺失 → 404。 + # 标记里别名对应的产物文件缺失 → 404(不会落到旁挂字幕解析)。 (work / "a.result.20260819000000.txt").unlink() gone = client.get(f"/api/batch/jobs/{job['id']}/videos/{video['id']}/download?alias=result") assert gone.status_code == 404 @@ -1001,3 +1187,20 @@ def test_batch_api_dirs(tmp_path, monkeypatch) -> None: assert denied["dirs"] == [] finally: client.__exit__(None, None, None) + + +def test_main_starts_batch_worker_when_enabled(monkeypatch) -> None: + """WOV_BATCH_ENABLED=1 时应用生命周期启动批量引擎后台线程(退出时回收)。""" + monkeypatch.setenv("WOV_BATCH_ENABLED", "1") + client = TestClient(app) + with client: + worker = app.state.batch + assert worker is not None + assert worker._thread is not None + assert worker._thread.name == "wov-batch-worker" + # 退出应用后批量引擎线程已停止,后续测试不会被后台线程打扰。 + assert worker._thread is None + + +# STORAGE_DIR 引用仅供静态检查使用(ensure batch 根相对应用存储目录)。 +assert BATCH_WORK_ROOT == STORAGE_DIR / "batch" diff --git a/web/assets/batch.js b/web/assets/batch.js index 0adf5af..9ff1bc5 100644 --- a/web/assets/batch.js +++ b/web/assets/batch.js @@ -42,7 +42,9 @@ async function createBatchJob() { method: "POST", body: JSON.stringify({ folder, workflow_id: workflowId, recursive }), }); - hint.textContent = `已创建任务 ${job.id},共 ${job.videos.length} 个视频`; + const videos = job.videos || []; + const skipped = videos.filter((v) => v.status === "SKIPPED").length; + hint.textContent = `已创建任务 ${job.id},共 ${videos.length} 个视频(待处理 ${videos.length - skipped},已有字幕跳过 ${skipped})`; document.getElementById("batchFolder").value = ""; } catch (error) { hint.textContent = `创建失败:${error.message}`; @@ -172,7 +174,7 @@ async function resumeBatchJob(jobId) { await loadBatchJobs(); } -// 删除批量任务:只清理数据库记录,磁盘上的同名文件夹与产物保留。 +// 删除批量任务:只清理数据库记录与应用私有工作空间;视频旁的产物保留。 async function deleteBatchJob(jobId) { if (!confirm("删除批量任务?磁盘上的产物文件会保留。")) return; try { diff --git a/web/batch.html b/web/batch.html index 1bfaab7..4dde9ce 100644 --- a/web/batch.html +++ b/web/batch.html @@ -23,8 +23,9 @@

批量处理

直接读取所选文件夹下的全部视频(不上传副本),逐个执行所选流水线; - 每个视频的中间态数据与最终产物存放在视频旁边的同名文件夹(如 movie.mp4 → - movie/)。已处理过的视频自动跳过;暂停后重新开始时,从未完成的视频继续。 + 视频旁(视频所在目录)已存在文件名含视频名的字幕文件的会自动跳过; + 处理完成的中文字幕(.srt)与双目字幕(.ass)会放到视频旁,处理过程文件 + 自动清理;暂停后重新开始时,从未完成的视频继续。