"""工作流管理路由。 提供工作流的创建、查询、校验、发布和删除能力。工作流以版本化 DAG 数据保存, 不写死在业务代码中。 """ from __future__ import annotations import re import uuid from fastapi import APIRouter, Depends, HTTPException from wov_app.db import Database from wov_app.schemas import WorkflowCreate from wov_app.scheduler import topological_sort from wov_sdk.models import WorkflowDefinition router = APIRouter(prefix="/api/admin/workflows", tags=["workflows"]) def _get_db() -> Database: """从应用状态延迟获取数据库实例。""" from wov_app.main import app return app.state.db def _slugify(value: str) -> str: """把工作流名称转换为小写连字符 ID;无有效字符时生成随机 ID。""" slug = re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") return slug or uuid.uuid4().hex[:8] def _validate_definition(raw: dict) -> WorkflowDefinition: """解析并校验 DAG 定义,非法时转换为 422 HTTP 异常。 除 `WorkflowDefinition.validate()` 的结构校验(名称/版本/节点 ID 唯一/ 边引用存在)之外,还要求 DAG **可拓扑排序**:环形依赖结构上合法但无法 确定执行顺序,必须拒绝保存与发布。 """ try: definition = WorkflowDefinition.from_dict(raw) definition.validate() # 环形 DAG(含自环)在这里以 "workflow contains a cycle" 被拒绝。 topological_sort(definition) return definition except (KeyError, TypeError, ValueError) as exc: raise HTTPException(status_code=422, detail=str(exc)) from exc @router.get("") def list_workflows(db: Database = Depends(_get_db)) -> list[dict]: """返回全部工作流概要。""" return db.list_workflows() @router.post("") def create_workflow( payload: WorkflowCreate, db: Database = Depends(_get_db), ) -> dict: """创建新工作流或为已有工作流追加一个版本。""" definition = _validate_definition(payload.definition) # 未显式指定 ID 时由名称生成;已有工作流则版本号递增。 workflow_id = payload.id or _slugify(payload.name) existing = db.get_workflow(workflow_id) version = (existing or {}).get("latest_version", 0) + 1 # 每次创建都保存新版本,发布操作只切换 published 标记。 db.upsert_workflow( { "id": workflow_id, "name": payload.name, "description": payload.description, "published": 0, "latest_version": version, } ) db.create_workflow_version(workflow_id, version, definition.to_dict()) return { "id": workflow_id, "name": payload.name, "description": payload.description, "published": False, "latest_version": version, } @router.get("/{workflow_id}") def get_workflow(workflow_id: str, db: Database = Depends(_get_db)) -> dict: """返回工作流概要及最新版本定义。""" workflow = db.get_workflow(workflow_id) if workflow is None: raise HTTPException(status_code=404, detail="workflow not found") latest = db.get_latest_workflow_version(workflow_id) workflow["latest_version_data"] = latest return workflow @router.delete("/{workflow_id}") def delete_workflow(workflow_id: str, db: Database = Depends(_get_db)) -> dict: """删除工作流及其版本、任务和产物记录。""" if db.get_workflow(workflow_id) is None: raise HTTPException(status_code=404, detail="workflow not found") db.delete_workflow(workflow_id) return {"deleted": workflow_id} @router.post("/{workflow_id}/validate") def validate_workflow( workflow_id: str, definition: dict, db: Database = Depends(_get_db), ) -> dict: """在不保存的情况下校验一份 DAG 定义。""" if db.get_workflow(workflow_id) is None: raise HTTPException(status_code=404, detail="workflow not found") parsed = _validate_definition(definition) return {"valid": True, "node_ids": [node.id for node in parsed.nodes]} @router.post("/{workflow_id}/publish") def publish_workflow(workflow_id: str, db: Database = Depends(_get_db)) -> dict: """把工作流标记为已发布,使其出现在用户应用中心。""" workflow = db.get_workflow(workflow_id) if workflow is None: raise HTTPException(status_code=404, detail="workflow not found") if workflow["latest_version"] == 0: raise HTTPException(status_code=422, detail="workflow has no version") # 发布前重新校验:无效定义(如环形 DAG)不能进入应用中心,否则任务会在 # 执行期失败。 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( { "id": workflow_id, "name": workflow["name"], "description": workflow["description"], "published": 1, "latest_version": workflow["latest_version"], } ) return {"published": workflow_id} @router.get("/{workflow_id}/versions") def list_versions(workflow_id: str, db: Database = Depends(_get_db)) -> list[dict]: """返回工作流全部版本定义。""" if db.get_workflow(workflow_id) is None: raise HTTPException(status_code=404, detail="workflow not found") return db.list_workflow_versions(workflow_id)