fix: 拒绝环形 DAG 并避免无效工作流堵塞任务队列

This commit is contained in:
2026-09-11 17:19:40 +08:00
parent 6eb65e4356
commit f99f8de171
5 changed files with 242 additions and 9 deletions
+28 -5
View File
@@ -11,11 +11,11 @@
| R01 | P0 | 普通删除接口可能删除批量源视频目录;清理逻辑从输入路径推导删除范围 | 已修复 | 普通接口拒绝单独删除批量 run;手动与自动清理只删除该上传任务的私有目录;视频、旁挂字幕及其他任务文件不受影响 | | R01 | P0 | 普通删除接口可能删除批量源视频目录;清理逻辑从输入路径推导删除范围 | 已修复 | 普通接口拒绝单独删除批量 run;手动与自动清理只删除该上传任务的私有目录;视频、旁挂字幕及其他任务文件不受影响 |
| R02 | P1 | 自适应线程池限流后可能扩容;退出标记位于积压队列尾部,缩容不及时 | 已修复 | 限流不增加并发;降低目标后不再超额提交;真实积压队列验证 | | R02 | P1 | 自适应线程池限流后可能扩容;退出标记位于积压队列尾部,缩容不及时 | 已修复 | 限流不增加并发;降低目标后不再超额提交;真实积压队列验证 |
| R03 | P1 | 最后节点执行时暂停再继续会覆盖有效下载 URI;批量缺产物仍清理并完成 | 已修复 | 收尾幂等,恢复后全部必需产物可下载;缺产物不清理工作空间 | | R03 | P1 | 最后节点执行时暂停再继续会覆盖有效下载 URI;批量缺产物仍清理并完成 | 已修复 | 收尾幂等,恢复后全部必需产物可下载;缺产物不清理工作空间 |
| R04 | P1 | 环形 DAG 可通过保存校验,执行失败后仍在 QUEUED 堵塞队列 | 待处理 | 保存/发布拒绝无效 DAG;历史无效任务进入 FAILED,不阻塞后续任务 | | R04 | P1 | 环形 DAG 可通过保存校验,执行失败后仍在 QUEUED 堵塞队列 | 已修复 | 保存/发布拒绝无效 DAG;历史无效任务进入 FAILED,不阻塞后续任务 |
| R05 | P1 | 翻译按固定四行解析 SRT;补齐行数不能保证文本与时间轴对应 | 待处理 | 合法多行 cue 正确解析;按稳定 ID 回填译文并校验缺失、重复项 | | R05 | P1 | 翻译按固定四行解析 SRT;补齐行数不能保证文本与时间轴对应 | 已修复 | 合法多行 cue 正确解析;按稳定 ID 回填译文并校验缺失、重复项 |
| R06 | P1 | OCR 跨空白帧合并相同字幕;临时请求失败被永久存为无文字 | 待处理 | 空白帧结束当前字幕段;无文字和可重试失败分开存档 | | R06 | P1 | OCR 跨空白帧合并相同字幕;临时请求失败被永久存为无文字 | 已修复 | 空白帧结束当前字幕段;无文字和可重试失败分开存档 |
| R07 | P2 | 抽帧/分块目录残留污染重跑;OCR/过滤存档无输入和参数指纹 | 待处理 | 输出减少时无旧文件混入;输入/有效参数变化时断点失效 | | R07 | P2 | 抽帧/分块目录残留污染重跑;OCR/过滤存档无输入和参数指纹 | 待处理 | 输出减少时无旧文件混入;输入/有效参数变化时断点失效 |
| R08 | P2 | 批量任务未固定工作流版本,恢复收尾可能使用新版本定义;创建任务与明细不原子 | 待处理 | 执行与收尾使用固定版本;工作线程只看到完整的已创建任务 | | R08 | P2 | 批量任务创建与明细不原子;不区分工作流版本时的执行/收尾一致性 | 待处理(方案调整) | 用户决定不区分工作流版本,不实施固定版本方案;工作线程只看到完整任务,收尾使用本次执行实际产物 |
## 优化与维护清单 ## 优化与维护清单
@@ -69,7 +69,30 @@
- 相关回归:`uv run pytest tests/test_finalization.py tests/test_scheduler.py tests/test_batch.py tests/test_apps_api.py tests/test_db.py -q --tb=short`97 passed3.22 秒。 - 相关回归:`uv run pytest tests/test_finalization.py tests/test_scheduler.py tests/test_batch.py tests/test_apps_api.py tests/test_db.py -q --tb=short`97 passed3.22 秒。
- 补充批量复制故障及安全回归:`uv run pytest tests/test_finalization.py tests/test_maintenance.py tests/test_run_deletion_safety.py -q --tb=short`22 passed,5.83 秒。真实视频/字幕验证写满故障后旧成品和工作空间完整,故障解除可成功放置并清理;API 下载验证原始文件与两个最终别名均可读取。仅有既存 Starlette/httpx 弃用警告。 - 补充批量复制故障及安全回归:`uv run pytest tests/test_finalization.py tests/test_maintenance.py tests/test_run_deletion_safety.py -q --tb=short`22 passed,5.83 秒。真实视频/字幕验证写满故障后旧成品和工作空间完整,故障解除可成功放置并清理;API 下载验证原始文件与两个最终别名均可读取。仅有既存 Starlette/httpx 弃用警告。
- 状态:修复及验证完成,纳入 `fix/review-improvements` 分支;未部署。原子性为单文件级,多文件中途失败允许部分完整成品已经更新,重试会重新放置;不承诺断电持久性。 - 状态:修复及验证完成,纳入 `fix/review-improvements` 分支;未部署。原子性为单文件级,多文件中途失败允许部分完整成品已经更新,重试会重新放置;不承诺断电持久性。
- 剩余边界:批量固定工作流版本仍由 R08 跟踪;此修复不会自动重建已经丢失的历史成品。 - 剩余边界:R08 按用户决定调整为不区分工作流版本,固定版本方案取消;版本机制调整及任务原子创建另行实施。此修复不会自动重建已经丢失的历史成品。
## R05 / R06 修复记录
- 顺序:按用户要求先完成 R05,再完成 R06;R03 已提交为 `3a61291`
- R05:新增 `nodes/srt.py` 按 cue 解析 BOM/CRLF、多行及空正文,坏字幕明确报错;翻译以全局位置 ID 的 JSON 条目请求,校验 ID 集合、类型、唯一性及正文。乱序按 ID 回填,结构错误最多尝试 3 次,耗尽 failed,取消末尾合并/补空策略。空 cue 保留时间轴且不请求模型,原有幻觉清洗继续执行。
- R05 TDD:12 个新回归用例先失败;相关节点/清洗/专名/翻译测试 93 passed、1 skipped、18 deselected(排除 Whisper)。最后调整提示词统一“条目”措辞后翻译对齐测试 13 passed(含真实短句 LLM 校准),0.76 秒。
- R06:空帧终止当前字幕段;失败/异常不写成功存档,失败帧在收紧的并发下重试一次,仍失败则节点 failed,已成功帧断点保留。新存档 status=completed 表示成功(含无文字),status=skipped 表示超长输出按既有规则跳过。
- 旧存档兼容:非空结果继续复用,无状态旧空串因无法区分超时与无文字而重新 OCR;不自动重跑已完成的历史任务。
- R06 TDD:4 个新回归先失败,覆盖跨空帧合并、failed 响应、抛异常和旧空串恢复。全量数据修正后 1942 条,旧 1666 条文件未修改;测试比较真实单线程、4/16 线程与断点结果,并逐采样点检查字幕不跨空白。基线必须走完整节点路径(包含超长过滤),不以原始 OCR 文本冒充成功存档。
- 最终组合验证:`uv run pytest tests/test_ocr_recovery.py tests/test_ocr_flow.py tests/test_subtitle_ocr_order_threading.py tests/test_integration_subtitle_ocr.py tests/test_translation_line_alignment.py tests/test_hallucination_mask.py tests/test_proper_nouns.py -q --tb=short --show-capture=no`71 passed、1 skipped58.69 秒;跳过项缺少历史专名素材,真实 LLM 与真实 OCR 集成通过。`git diff --check` 通过。
- 状态:R05/R06 工作区已修复,尚未提交、未部署。模型语义是否正确仍需内容质量评估,ID 校验保证程序不会因列表乱序或漏项贴错时间。
- R08 决策:用户计划不区分工作流版本,取消固定版本修复方案;版本机制调整、执行/收尾一致性及批量原子创建保留待办,本轮未实施数据库迁移。
## R04 修复记录
- 根因:①保存路径只做结构校验(`WorkflowDefinition.validate()`:名称/版本/节点 ID 唯一/边引用存在),环形依赖结构合法因而被写库并发布;②`execute_run` 的 DAG 解析、校验与拓扑排序在 `try` 块**之外**,环检测抛出的异常直接冒到 `_loop` 被吞掉,任务状态从未离开 QUEUED——`next_queued_run` 每轮拾起同一条队首记录,后续任务永久堵塞。
- 修复 1(保存/发布):[routers/workflows.py](../src/wov_app/routers/workflows.py) 的 `_validate_definition` 在校验后调用 `topological_sort`,环形 DAG(含自环)以 "workflow contains a cycle" 转 422;创建/校验/发布三个入口共用该函数,`publish` 额外重新校验最新版本定义,历史遗留的环形版本无法进入应用中心。校验先于任何写库,被拒请求不产生工作流或版本记录。
- 修复 2(历史无效任务):[scheduler.py](../src/wov_app/scheduler.py) 把 `from_dict` / `validate` / `topological_sort` 移入 try,任何预检异常都立刻把 run 置 FAILED 并记录错误后返回,队首随即前移,不再堵塞后续任务。
- 方案边界:环检测放在**入口校验**(保存/发布)而非 `WorkflowDefinition.validate()`,因此 `WorkflowDefinition.validate()` 保持只做结构校验、SDK 协议模型不变;批量任务在工作流级校验仍放行环形定义,由调度器把每个视频的 run 判 FAILED(保留既有"单视频失败不中断整批"语义,`test_batch_worker_cycle_fails_video_not_job` 未改动仍通过)。批量的 run 现在由调度器直接置 FAILED 并回填错误,不再依赖异常冒泡。
- TDD 红:新增 5 个回归用例先失败——创建带环工作流返回 200、校验接口返回 200、发布历史环形版本返回 200、环形任务执行后仍停在 QUEUED(队首堵塞)、缺 name 的非法定义同样卡住队列。
- TDD 绿及相关回归:`uv run pytest tests/test_workflow_api.py tests/test_scheduler.py -q`29 passed1.07 秒;`uv run pytest tests/test_batch.py tests/test_apps_api.py tests/test_seed.py tests/test_models.py tests/test_registry.py tests/test_api.py tests/test_finalization.py -q`95 passed2.59 秒。
- 全量回归:`uv run pytest -q`402 passed、5 skipped84.03 秒;仅既存 Starlette/httpx 弃用警告。`git diff --check` 通过。
- 状态:修复及验证完成,纳入 `fix/review-improvements` 分支工作区(与 R05/R06 未提交改动共存);未部署。历史遗留的环形已发布工作流不会被自动取消发布,仍可创建任务,但任务会立即 FAILED 且不堵塞队列;如需清理线上遗留数据需另行确认后执行。
## 审查基线 ## 审查基线
+15 -1
View File
@@ -13,6 +13,7 @@ from fastapi import APIRouter, Depends, HTTPException
from wov_app.db import Database from wov_app.db import Database
from wov_app.schemas import WorkflowCreate from wov_app.schemas import WorkflowCreate
from wov_app.scheduler import topological_sort
from wov_sdk.models import WorkflowDefinition from wov_sdk.models import WorkflowDefinition
router = APIRouter(prefix="/api/admin/workflows", tags=["workflows"]) router = APIRouter(prefix="/api/admin/workflows", tags=["workflows"])
@@ -32,10 +33,17 @@ def _slugify(value: str) -> str:
def _validate_definition(raw: dict) -> WorkflowDefinition: def _validate_definition(raw: dict) -> WorkflowDefinition:
"""解析并校验 DAG 定义,非法时转换为 422 HTTP 异常。""" """解析并校验 DAG 定义,非法时转换为 422 HTTP 异常。
除 `WorkflowDefinition.validate()` 的结构校验(名称/版本/节点 ID 唯一/
边引用存在)之外,还要求 DAG **可拓扑排序**:环形依赖虽然结构上合法,
但执行时无法确定节点顺序,必须拒绝保存与发布(R04)。
"""
try: try:
definition = WorkflowDefinition.from_dict(raw) definition = WorkflowDefinition.from_dict(raw)
definition.validate() definition.validate()
# 环形 DAG(含自环)在这里以 "workflow contains a cycle" 被拒绝。
topological_sort(definition)
return definition return definition
except (KeyError, TypeError, ValueError) as exc: except (KeyError, TypeError, ValueError) as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc raise HTTPException(status_code=422, detail=str(exc)) from exc
@@ -119,6 +127,12 @@ def publish_workflow(workflow_id: str, db: Database = Depends(_get_db)) -> dict:
raise HTTPException(status_code=404, detail="workflow not found") raise HTTPException(status_code=404, detail="workflow not found")
if workflow["latest_version"] == 0: if workflow["latest_version"] == 0:
raise HTTPException(status_code=422, detail="workflow has no version") raise HTTPException(status_code=422, detail="workflow has no version")
# 发布前重新校验待发布版本:历史遗留的无效定义(例如修复前保存的环形 DAG)
# 不能进入用户应用中心,否则创建出来的任务会在执行期失败(R04)。
latest = db.get_latest_workflow_version(workflow_id)
if latest is None:
raise HTTPException(status_code=422, detail="workflow has no version")
_validate_definition(latest["definition"])
# 发布只是状态切换,不修改已保存的版本数据。 # 发布只是状态切换,不修改已保存的版本数据。
db.upsert_workflow( db.upsert_workflow(
{ {
+18 -3
View File
@@ -150,9 +150,24 @@ class WorkflowScheduler:
return return
# 解析并校验 DAG,随后计算拓扑执行顺序。 # 解析并校验 DAG,随后计算拓扑执行顺序。
definition = WorkflowDefinition.from_dict(version["definition"]) #
definition.validate() # 任何预检异常(缺字段/边引用不存在/环形依赖)都必须在这里把任务标
ordered = topological_sort(definition) # FAILED:修复前这段在 try 之外,异常直接冒到 _loop 被吞掉,任务停在
# QUEUEDnext_queued_run 每轮拾起同一条队首记录,后续任务全部堵塞
# (R04:环形 DAG 卡死队列)。历史无效版本无法删除,只能就地判失败。
try:
definition = WorkflowDefinition.from_dict(version["definition"])
definition.validate()
ordered = topological_sort(definition)
except Exception as exc: # noqa: BLE001
logger.exception("任务 %s 工作流定义无效,标记失败: %s", run_id, exc)
self.db.update_run(
run_id,
status="FAILED",
error=str(exc),
updated_at=_now_iso(),
)
return
# 从已登记产物重建已完成节点的输出,支持暂停后断点续跑。 # 从已登记产物重建已完成节点的输出,支持暂停后断点续跑。
outputs_by_node = self.db.restore_run_outputs(run_id) outputs_by_node = self.db.restore_run_outputs(run_id)
run_started = time.monotonic() run_started = time.monotonic()
+96
View File
@@ -173,6 +173,102 @@ def test_execute_run_missing_version(tmp_path) -> None:
assert db.get_run("run_version")["status"] == "FAILED" assert db.get_run("run_version")["status"] == "FAILED"
def test_execute_run_invalid_dag_fails_and_does_not_block_queue(tmp_path) -> None:
"""环形 DAG 的历史任务必须立即 FAILED,且不阻塞队列后续任务(R04)。
复现场景:修复前保存校验放过环形定义,执行时拓扑排序抛错但异常发生在
execute_run 的状态翻转之前 → 任务永远停在 QUEUEDnext_queued_run 每轮
都拾起同一条队首记录,后面的任务全部被堵死。断言:无效 DAG 的任务被标记
FAILED(带环错误信息),execute_run 不向调用方抛异常,队首随即前移。
"""
db = _db(tmp_path)
_register_echo()
# 直接写库模拟修复前遗留的环形版本(保存接口现在会拒绝)。
cycle_definition = {
"name": "cycle",
"version": 1,
"nodes": [
{"id": "a", "node_type": "echo", "inputs": {"file_uri": "b.file_uri"}},
{"id": "b", "node_type": "echo", "inputs": {"file_uri": "a.file_uri"}},
],
"edges": [{"from": "a", "to": "b"}, {"from": "b", "to": "a"}],
"entry_inputs": {"video_uri": "file"},
"final_outputs": {"result": "a.file_uri"},
}
db.upsert_workflow({"id": "bad-flow", "name": "Bad", "published": 1, "latest_version": 1})
db.create_workflow_version("bad-flow", 1, cycle_definition)
# 队列里同时放入合法任务,验证它不被无效任务堵住。
db.upsert_workflow({"id": "flow", "name": "Flow", "published": 1, "latest_version": 1})
db.create_workflow_version("flow", 1, _echo_definition().to_dict())
input_file = tmp_path / "input.txt"
input_file.write_text("hello queue", encoding="utf-8")
db.create_run(
{
"id": "run_bad",
"workflow_id": "bad-flow",
"workflow_version": 1,
"status": "QUEUED",
"progress": 0,
"input_uri": str(input_file),
"created_at": "2026-01-01T00:00:00+00:00",
"updated_at": "2026-01-01T00:00:00+00:00",
}
)
db.create_run(
{
"id": "run_ok",
"workflow_id": "flow",
"workflow_version": 1,
"status": "QUEUED",
"progress": 0,
"input_uri": str(input_file),
"created_at": "2026-01-01T00:00:01+00:00",
"updated_at": "2026-01-01T00:00:01+00:00",
}
)
scheduler = WorkflowScheduler(db, tmp_path / "storage")
# 调度器每一轮的取任务 → 执行,不应把环形任务留在 QUEUED。
first = db.next_queued_run()
assert first["id"] == "run_bad"
scheduler.execute_run(first["id"])
bad = db.get_run("run_bad")
assert bad["status"] == "FAILED"
assert "cycle" in bad["error"]
# 队首已前移,后续合法任务正常执行完成。
assert db.next_queued_run()["id"] == "run_ok"
scheduler.execute_run("run_ok")
assert db.get_run("run_ok")["status"] == "COMPLETED"
assert db.next_queued_run() is None
def test_execute_run_unparsable_dag_fails_not_stuck(tmp_path) -> None:
"""缺 name 的非法定义同样标记 FAILED,而不是让队列卡死。"""
db = _db(tmp_path)
db.upsert_workflow({"id": "flow", "name": "Flow", "published": 1, "latest_version": 1})
# from_dict 解析缺 name 的定义会抛 KeyError。
db.create_workflow_version("flow", 1, {"version": 1, "nodes": [], "edges": []})
now = "2026-01-01T00:00:00+00:00"
db.create_run(
{
"id": "run_corrupt",
"workflow_id": "flow",
"workflow_version": 1,
"status": "QUEUED",
"progress": 0,
"created_at": now,
"updated_at": now,
}
)
scheduler = WorkflowScheduler(db, tmp_path / "storage")
scheduler.execute_run("run_corrupt")
run = db.get_run("run_corrupt")
assert run["status"] == "FAILED"
assert run["error"]
assert db.next_queued_run() is None
def test_execute_run_missing_node(tmp_path) -> None: def test_execute_run_missing_node(tmp_path) -> None:
"""验证未注册节点被调用时任务失败。""" """验证未注册节点被调用时任务失败。"""
db = _db(tmp_path) db = _db(tmp_path)
+85
View File
@@ -26,6 +26,91 @@ def definition() -> dict:
} }
def cyclic_definition() -> dict:
"""构造带环的 DAG 定义(a → b → a),用于验证保存/发布拒绝无效工作流。
环上的两个节点互相引用对方输出,形成了无法拓扑排序的依赖:节点 ID 唯一、
边引用的节点都存在,因此只有环检测能拦住它(R04 复现定义)。
"""
return {
"name": "cyclic",
"version": 1,
"nodes": [
{"id": "a", "node_type": "echo", "inputs": {"file_uri": "b.file_uri"}},
{"id": "b", "node_type": "echo", "inputs": {"file_uri": "a.file_uri"}},
],
"edges": [{"from": "a", "to": "b"}, {"from": "b", "to": "a"}],
"entry_inputs": {"video_uri": "file"},
"final_outputs": {"result": "a.file_uri"},
}
def test_create_workflow_rejects_cycle() -> None:
"""验证保存环形 DAG 返回 422,不把无效版本写入版本表。"""
with TestClient(app) as client:
response = client.post(
"/api/admin/workflows",
json={
"id": "cyclic-flow",
"name": "Cyclic Flow",
"definition": cyclic_definition(),
},
)
assert response.status_code == 422
assert "cycle" in response.json()["detail"]
# 拒绝保存时不得留下半成品工作流/版本记录。
db = app.state.db
assert db.get_workflow("cyclic-flow") is None
assert db.list_workflow_versions("cyclic-flow") == []
# 已有工作流重新提交带环定义:同样拒绝,不追加新版本。
assert client.post(
"/api/admin/workflows",
json={"id": "cycle-check", "name": "Cycle Check", "definition": definition()},
).status_code == 200
assert client.post(
"/api/admin/workflows",
json={"id": "cycle-check", "name": "Cycle Check", "definition": cyclic_definition()},
).status_code == 422
assert db.get_workflow("cycle-check")["latest_version"] == 1
assert len(db.list_workflow_versions("cycle-check")) == 1
db.delete_workflow("cycle-check")
def test_validate_workflow_rejects_cycle() -> None:
"""验证只校验不保存的接口同样拒绝环形 DAG。"""
with TestClient(app) as client:
db = app.state.db
# 独立工作流 ID 并显式清理,避免与其他用例(共用同一个测试数据库)互相影响。
assert client.post(
"/api/admin/workflows",
json={"id": "cycle-check", "name": "Cycle Check", "definition": definition()},
).status_code == 200
response = client.post(
"/api/admin/workflows/cycle-check/validate",
json=cyclic_definition(),
)
assert response.status_code == 422
assert "cycle" in response.json()["detail"]
db.delete_workflow("cycle-check")
def test_publish_rejects_cycle_in_latest_version() -> None:
"""验证历史遗留的环形版本不能被发布(发布前重新校验最新版本定义)。"""
with TestClient(app) as client:
db = app.state.db
# 直接写库模拟修复前已保存的无效版本(保存接口现在会拒绝,只能这样构造)。
db.upsert_workflow(
{"id": "legacy-cyclic", "name": "Legacy", "published": 0, "latest_version": 1}
)
db.create_workflow_version("legacy-cyclic", 1, cyclic_definition())
response = client.post("/api/admin/workflows/legacy-cyclic/publish")
assert response.status_code == 422
assert "cycle" in response.json()["detail"]
assert db.get_workflow("legacy-cyclic")["published"] == 0
db.delete_workflow("legacy-cyclic")
def test_workflow_crud_and_publish() -> None: def test_workflow_crud_and_publish() -> None:
"""验证工作流 CRUD、校验、发布与版本列表的完整流程。""" """验证工作流 CRUD、校验、发布与版本列表的完整流程。"""
with TestClient(app) as client: with TestClient(app) as client: