feat: 批量处理重做——视频旁已有字幕即跳过、产物对齐 CN 命名并清理过程文件
- 创建批量任务时一次性定位视频:视频所在目录存在文件名含视频名的字幕文件 (.srt/.ass/.ssa/.vtt)直接记 SKIPPED,不触发流水线;运行时只消费已定位 的明细,不再重新扫描文件夹。 - 视频完成后把最终产物放到视频旁,命名对齐媒体库约定:中文字幕存为 <视频名>.CN.srt、双目字幕存为 <视频名>.CN_dual_eye.ass;其余扩展名产物 保留原文件名。 - 收尾删除 run 记录与过程工作空间;工作空间改到应用私有目录 storage/batch/<job_id>/<bv_id>/,与用户媒体库隔离,防止媒体库把切片数据 当视频入库。 - 批量 API 产物清单/下载改为解析视频旁字幕文件,旧版 batch.done.json 语义 别名保持兼容;删除任务时清理私有工作空间。 - 前端说明与创建提示同步;测试按新语义重写并补覆盖(278 passed,100% 行覆盖率)。
This commit is contained in:
@@ -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/<run_id>/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/<job_id>/<bv_id>/`
|
||||
(不再放视频同名文件夹),与用户视频库天然隔离;暂停/失败的视频保留工作空间
|
||||
以便断点续跑。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)。
|
||||
|
||||
|
||||
+184
-114
@@ -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/<job_id>/<bv_id>/`,
|
||||
与用户的视频库目录天然隔离;暂停/失败的视频保留工作空间以便断点续跑。
|
||||
|
||||
暂停/恢复语义(对应前端"暂停/继续"按钮):
|
||||
|
||||
@@ -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 根目录(<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 格式字符串。"""
|
||||
@@ -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:
|
||||
"""删除批量任务在应用私有存储下的工作空间目录(<storage>/batch/<job_id>)。
|
||||
|
||||
同一批量任务反复处理(暂停/续跑)时保留每个视频的状态,已完成的不重置。
|
||||
每个视频完成后工作空间已被逐视频清理;此处兜底清理任务级残留(删除任务
|
||||
或任务异常中止时)。只作用于私有工作空间,绝不触碰用户视频目录。
|
||||
"""
|
||||
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/<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,
|
||||
})
|
||||
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,中间态落在
|
||||
<work_dir>/runs/<run_id>/steps/ 下;产物表记录全部节点输出,暂停后
|
||||
续跑从产物表重建已完成节点(断点续跑)。
|
||||
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 动态组装节点执行。
|
||||
@@ -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 "(无)",
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
|
||||
@@ -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/<job_id>)的残留;
|
||||
视频旁已经放置的产物属于用户数据,保留不删。
|
||||
"""
|
||||
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)
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# 本地目录浏览(目录树选择器)
|
||||
|
||||
+2
-1
@@ -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:
|
||||
"""进程退出时清理临时测试目录。"""
|
||||
|
||||
+478
-275
File diff suppressed because it is too large
Load Diff
+4
-2
@@ -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 {
|
||||
|
||||
+3
-2
@@ -23,8 +23,9 @@
|
||||
<h1>批量处理</h1>
|
||||
<p class="muted">
|
||||
直接读取所选文件夹下的全部视频(<b>不上传副本</b>),逐个执行所选流水线;
|
||||
每个视频的中间态数据与最终产物存放在视频旁边的同名文件夹(如 movie.mp4 →
|
||||
movie/)。已处理过的视频自动跳过;暂停后重新开始时,从未完成的视频继续。
|
||||
视频旁(视频所在目录)已存在文件名含视频名的字幕文件的会自动跳过;
|
||||
处理完成的中文字幕(.srt)与双目字幕(.ass)会放到视频旁,处理过程文件
|
||||
自动清理;暂停后重新开始时,从未完成的视频继续。
|
||||
</p>
|
||||
|
||||
<!-- 创建批量任务:文件夹路径 + 工作流选择 + 是否递归。 -->
|
||||
|
||||
Reference in New Issue
Block a user