352 lines
13 KiB
Python
352 lines
13 KiB
Python
"""Studio facade: camelCase DTOs for the standalone Coze canvas page.
|
|
|
|
Four endpoints in front of the standard workflow store:
|
|
|
|
GET /api/workflows/{id}/studio-document
|
|
PUT /api/workflows/{id}/studio-draft
|
|
POST /api/workflows/{id}/studio-validate
|
|
POST /api/workflows/{id}/studio-publish
|
|
|
|
Design contract (docs/WORKFLOW_STUDIO_COZE_FRONTEND_ADAPTATION_ZH.md §3-§4):
|
|
|
|
- One workflow, two documents. ``canvasSchemaJson`` is an opaque Coze canvas
|
|
saved for editing; ``executionGraph`` is the WorkflowGraph v1.0 the runtime
|
|
executes. Both land in the SAME CAS update, so they can never disagree on
|
|
revision.
|
|
- The facade only maps DTO naming and delegates to the standard store /
|
|
validator / publish service — it never forks versioning logic.
|
|
- Publishing a workflow whose draft has no valid execution graph is a 400, not
|
|
a silent empty-graph publish.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
from hashlib import sha256
|
|
from typing import Any
|
|
|
|
from fastapi import APIRouter, HTTPException, Request
|
|
from pydantic import BaseModel, Field
|
|
|
|
from app.gateway.deps import get_current_user
|
|
from app.gateway.workflow_audit import audit_workflow
|
|
from app.gateway.workflow_resource_checks import build_resource_catalog, make_resource_checker
|
|
from deerflow.config.app_config import get_app_config
|
|
from deerflow.persistence.workflows import (
|
|
WorkflowDraftConflictError,
|
|
WorkflowNotFoundError,
|
|
WorkflowStore,
|
|
)
|
|
from deerflow.workflows.schemas import WorkflowGraph
|
|
from deerflow.workflows.validator import validate_workflow_graph
|
|
|
|
router = APIRouter(prefix="/api/workflows", tags=["workflow-studio"])
|
|
|
|
# Canvas document guards: it is stored and replayed as an opaque document, so
|
|
# the only checks are size / JSON shape / depth / obvious credential keys.
|
|
_CANVAS_MAX_CHARS = 2_000_000
|
|
_CANVAS_MAX_DEPTH = 64
|
|
_CANVAS_SECRET_KEYS = frozenset(
|
|
{
|
|
"authorization",
|
|
"proxy-authorization",
|
|
"cookie",
|
|
"x-api-key",
|
|
"apikey",
|
|
"api_key",
|
|
"secret",
|
|
"client_secret",
|
|
"password",
|
|
"dsn",
|
|
"access_token",
|
|
"refresh_token",
|
|
"resumetoken",
|
|
"resume_token",
|
|
}
|
|
)
|
|
|
|
|
|
class StudioDraftBody(BaseModel):
|
|
expected_revision: int = Field(alias="expectedRevision", ge=0)
|
|
# Canvas may arrive as a JSON string or an embedded object; both are
|
|
# canonicalised to compact JSON text before storage.
|
|
canvas_schema_json: Any = Field(default=None, alias="canvasSchemaJson")
|
|
execution_graph: dict[str, Any] = Field(alias="executionGraph")
|
|
|
|
model_config = {"populate_by_name": True}
|
|
|
|
|
|
class StudioValidateBody(BaseModel):
|
|
execution_graph: dict[str, Any] | None = Field(default=None, alias="executionGraph")
|
|
|
|
model_config = {"populate_by_name": True}
|
|
|
|
|
|
class StudioPublishBody(BaseModel):
|
|
expected_revision: int | None = Field(default=None, alias="expectedRevision", ge=0)
|
|
change_note: str = Field(default="", alias="changeNote", max_length=2000)
|
|
|
|
model_config = {"populate_by_name": True}
|
|
|
|
|
|
def _require_enabled() -> None:
|
|
if not get_app_config().workflows.enabled:
|
|
raise HTTPException(status_code=503, detail="Workflow Studio is disabled")
|
|
|
|
|
|
def _get_store(request: Request) -> WorkflowStore:
|
|
store = getattr(request.app.state, "workflow_store", None)
|
|
if store is None:
|
|
raise HTTPException(status_code=503, detail="Workflow store not available")
|
|
return store
|
|
|
|
|
|
async def _load_owned(request: Request, workflow_id: str) -> dict[str, Any]:
|
|
"""Owner or admin only — the studio edits the draft, so read == write."""
|
|
from app.gateway.routers.workflows import _load_owned as _standard_load_owned
|
|
|
|
return await _standard_load_owned(request, workflow_id, write=True)
|
|
|
|
|
|
def _limits() -> dict[str, int]:
|
|
cfg = get_app_config().workflows
|
|
return {
|
|
"system_max_steps": cfg.max_steps,
|
|
"system_max_loop_iterations": cfg.max_loop_iterations,
|
|
"system_max_parallelism": cfg.max_parallelism,
|
|
}
|
|
|
|
|
|
def _parse_canvas(raw: Any) -> str:
|
|
"""Validate + canonicalise the canvas document into compact JSON text."""
|
|
if raw is None or raw == "":
|
|
return "{}"
|
|
if isinstance(raw, str):
|
|
text = raw
|
|
try:
|
|
document = json.loads(text)
|
|
except json.JSONDecodeError as exc:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_CANVAS_INVALID", "message": "画布 JSON 无法解析"},
|
|
) from exc
|
|
elif isinstance(raw, dict):
|
|
document = raw
|
|
else:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_CANVAS_INVALID", "message": "画布必须是 JSON 对象或 JSON 字符串"},
|
|
)
|
|
text = json.dumps(document, ensure_ascii=False, separators=(",", ":"))
|
|
if len(text) > _CANVAS_MAX_CHARS:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_CANVAS_INVALID", "message": f"画布文档超过 {_CANVAS_MAX_CHARS} 字符上限"},
|
|
)
|
|
offending = _find_secret_key(document, 0)
|
|
if offending:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={
|
|
"code": "WORKFLOW_CANVAS_INVALID",
|
|
"message": f"画布中出现凭据类字段 {offending},请改用凭据引用",
|
|
"field": offending,
|
|
},
|
|
)
|
|
return text
|
|
|
|
|
|
def _find_secret_key(node: Any, depth: int) -> str | None:
|
|
if depth > _CANVAS_MAX_DEPTH:
|
|
return "<max-depth-exceeded>"
|
|
if isinstance(node, dict):
|
|
for key, value in node.items():
|
|
name = str(key).strip().lower()
|
|
if name in _CANVAS_SECRET_KEYS and value not in (None, "", {}, []):
|
|
return str(key)
|
|
found = _find_secret_key(value, depth + 1)
|
|
if found:
|
|
return found
|
|
elif isinstance(node, list):
|
|
for item in node:
|
|
found = _find_secret_key(item, depth + 1)
|
|
if found:
|
|
return found
|
|
return None
|
|
|
|
|
|
def _studio_issues(issues: list[Any]) -> list[dict[str, Any]]:
|
|
out: list[dict[str, Any]] = []
|
|
for issue in issues:
|
|
body = issue.model_dump(by_alias=True)
|
|
details = body.get("details") or {}
|
|
out.append(
|
|
{
|
|
"code": body.get("code"),
|
|
"nodeId": body.get("nodeId"),
|
|
"field": details.get("field") if isinstance(details, dict) else None,
|
|
"message": body.get("message"),
|
|
}
|
|
)
|
|
return out
|
|
|
|
|
|
async def _validate_with_resources(request: Request, workflow_id: str, raw_graph: dict[str, Any], user_id: str) -> list[Any]:
|
|
catalog = await build_resource_catalog(request, user_id=user_id)
|
|
checker = make_resource_checker(catalog, config=get_app_config().workflows, current_workflow_id=workflow_id)
|
|
return validate_workflow_graph(raw_graph, resource_checker=checker, **_limits())
|
|
|
|
|
|
@router.get("/{workflow_id}/studio-document")
|
|
async def get_studio_document(request: Request, workflow_id: str) -> dict[str, Any]:
|
|
_require_enabled()
|
|
definition = await _load_owned(request, workflow_id)
|
|
store = _get_store(request)
|
|
graph = definition.get("draft_graph") or {}
|
|
versions = await store.list_versions(workflow_id)
|
|
return {
|
|
"workflowId": definition["id"],
|
|
"name": definition.get("name"),
|
|
"description": definition.get("description") or "",
|
|
"status": definition.get("status"),
|
|
"draftRevision": int(definition.get("draft_revision") or 0),
|
|
"publishedVersionId": versions[0]["id"] if versions else None,
|
|
# Canvas stays an opaque JSON document (string); never null.
|
|
"canvasSchemaJson": definition.get("draft_canvas_schema") or "{}",
|
|
"executionGraph": graph,
|
|
"inputSchema": graph.get("inputSchema") or graph.get("input_schema") or {"type": "object", "properties": {}},
|
|
"outputSchema": graph.get("outputSchema") or graph.get("output_schema") or {"type": "object", "properties": {}},
|
|
"updatedAt": definition.get("updated_at"),
|
|
}
|
|
|
|
|
|
@router.put("/{workflow_id}/studio-draft")
|
|
async def put_studio_draft(request: Request, workflow_id: str, body: StudioDraftBody) -> dict[str, Any]:
|
|
_require_enabled()
|
|
await _load_owned(request, workflow_id)
|
|
store = _get_store(request)
|
|
user_id = await get_current_user(request)
|
|
|
|
try:
|
|
WorkflowGraph.model_validate(body.execution_graph)
|
|
except Exception as exc: # noqa: BLE001 - surface pydantic as 400
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_SCHEMA_INVALID", "message": f"执行图不符合 WorkflowGraph v1.0: {exc}"},
|
|
) from exc
|
|
|
|
canvas_text = _parse_canvas(body.canvas_schema_json)
|
|
try:
|
|
saved = await store.save_studio_draft(
|
|
workflow_id,
|
|
expected_revision=body.expected_revision,
|
|
graph=body.execution_graph,
|
|
canvas_schema=canvas_text,
|
|
updated_by=user_id,
|
|
)
|
|
except WorkflowDraftConflictError as exc:
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail={
|
|
"code": "WORKFLOW_DRAFT_CONFLICT",
|
|
"message": "草稿版本冲突,请刷新后重试或另存副本",
|
|
"currentRevision": exc.current_revision,
|
|
},
|
|
) from exc
|
|
except WorkflowNotFoundError as exc:
|
|
raise HTTPException(status_code=404, detail="工作流不存在") from exc
|
|
audit_workflow("workflow.studio_draft", user_id=user_id, workflowId=workflow_id, revision=saved.get("draft_revision"))
|
|
return {"revision": int(saved.get("draft_revision") or 0), "updatedAt": saved.get("updated_at")}
|
|
|
|
|
|
@router.post("/{workflow_id}/studio-validate")
|
|
async def studio_validate(request: Request, workflow_id: str, body: StudioValidateBody | None = None) -> dict[str, Any]:
|
|
_require_enabled()
|
|
definition = await _load_owned(request, workflow_id)
|
|
user_id = await get_current_user(request)
|
|
graph = (body.execution_graph if body and body.execution_graph is not None else definition.get("draft_graph")) or {}
|
|
try:
|
|
issues = await _validate_with_resources(request, workflow_id, graph, user_id)
|
|
except Exception as exc: # noqa: BLE001 - unparsable graph is a validation result, not a 500
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_SCHEMA_INVALID", "message": f"图结构无效: {exc}"},
|
|
) from exc
|
|
valid = not issues
|
|
graph_hash = None
|
|
if valid:
|
|
try:
|
|
graph_hash = WorkflowGraph.model_validate(graph).graph_hash()
|
|
except Exception: # noqa: BLE001
|
|
valid = False
|
|
return {"valid": valid, "issues": _studio_issues(issues), "graphHash": graph_hash}
|
|
|
|
|
|
@router.post("/{workflow_id}/studio-publish")
|
|
async def studio_publish(request: Request, workflow_id: str, body: StudioPublishBody | None = None) -> dict[str, Any]:
|
|
_require_enabled()
|
|
definition = await _load_owned(request, workflow_id)
|
|
store = _get_store(request)
|
|
user_id = await get_current_user(request)
|
|
body = body or StudioPublishBody()
|
|
|
|
graph_raw = definition.get("draft_graph") or {}
|
|
if not graph_raw or not (graph_raw.get("nodes") or []):
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={
|
|
"code": "WORKFLOW_SCHEMA_INVALID",
|
|
"message": "草稿没有可发布的执行图(只有画布文档不能发布)",
|
|
},
|
|
)
|
|
try:
|
|
issues = await _validate_with_resources(request, workflow_id, graph_raw, user_id)
|
|
except Exception as exc: # noqa: BLE001
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={"code": "WORKFLOW_SCHEMA_INVALID", "message": f"图结构无效: {exc}"},
|
|
) from exc
|
|
if issues:
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail={
|
|
"code": "WORKFLOW_SCHEMA_INVALID",
|
|
"message": "发布校验未通过",
|
|
"issues": _studio_issues(issues),
|
|
},
|
|
)
|
|
|
|
if body.expected_revision is not None and body.expected_revision != definition.get("draft_revision"):
|
|
raise HTTPException(
|
|
status_code=409,
|
|
detail={
|
|
"code": "WORKFLOW_DRAFT_CONFLICT",
|
|
"message": "发布前草稿版本已变化",
|
|
"currentRevision": definition.get("draft_revision"),
|
|
},
|
|
)
|
|
|
|
graph = WorkflowGraph.model_validate(graph_raw)
|
|
canvas_text = str(definition.get("draft_canvas_schema") or "{}")
|
|
canvas_hash = sha256(canvas_text.encode("utf-8")).hexdigest()
|
|
version = await store.publish_version(
|
|
workflow_id,
|
|
graph=graph.model_dump(by_alias=True),
|
|
graph_hash=graph.graph_hash(),
|
|
published_by=user_id,
|
|
change_note=body.change_note,
|
|
canvas_schema_hash=canvas_hash,
|
|
)
|
|
audit_workflow(
|
|
"workflow.studio_publish",
|
|
user_id=user_id,
|
|
workflowId=workflow_id,
|
|
versionId=version.get("id"),
|
|
graphHash=graph.graph_hash(),
|
|
)
|
|
return {
|
|
"versionId": version.get("id"),
|
|
"versionNumber": version.get("version_number"),
|
|
"publishedAt": version.get("published_at"),
|
|
}
|