deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/persistence/workflows/base.py
2026-09-07 18:24:55 +08:00

113 lines
3.6 KiB
Python

"""Abstract workflow definition / version store."""
from __future__ import annotations
from abc import ABC, abstractmethod
from typing import Any
class WorkflowDraftConflictError(Exception):
"""Raised when ``expected_revision`` does not match the stored draft revision."""
def __init__(self, workflow_id: str, current_revision: int) -> None:
super().__init__(f"draft conflict for {workflow_id}: current={current_revision}")
self.workflow_id = workflow_id
self.current_revision = current_revision
class WorkflowNotFoundError(Exception):
def __init__(self, workflow_id: str) -> None:
super().__init__(f"workflow not found: {workflow_id}")
self.workflow_id = workflow_id
class WorkflowVersionNotFoundError(Exception):
def __init__(self, version_id: str) -> None:
super().__init__(f"workflow version not found: {version_id}")
self.version_id = version_id
class WorkflowStore(ABC):
@abstractmethod
async def list_definitions(
self,
*,
owner_id: str | None = None,
status: str | None = None,
include_archived: bool = False,
) -> list[dict[str, Any]]:
raise NotImplementedError
@abstractmethod
async def get_definition(self, workflow_id: str, *, include_draft: bool = True) -> dict[str, Any] | None:
raise NotImplementedError
@abstractmethod
async def create_definition(self, data: dict[str, Any]) -> dict[str, Any]:
raise NotImplementedError
@abstractmethod
async def update_definition(self, workflow_id: str, data: dict[str, Any]) -> dict[str, Any] | None:
raise NotImplementedError
@abstractmethod
async def archive_definition(self, workflow_id: str) -> dict[str, Any] | None:
raise NotImplementedError
@abstractmethod
async def save_draft(
self,
workflow_id: str,
*,
expected_revision: int,
graph: dict[str, Any],
updated_by: str | None = None,
) -> dict[str, Any]:
"""CAS update of ``draft_graph_json``. Raises ``WorkflowDraftConflictError``."""
raise NotImplementedError
@abstractmethod
async def save_studio_draft(
self,
workflow_id: str,
*,
expected_revision: int,
graph: dict[str, Any],
canvas_schema: str = "{}",
updated_by: str | None = None,
) -> dict[str, Any]:
"""Atomically CAS-update BOTH the execution graph and the canvas document.
A single conditional UPDATE bumps ``draft_revision`` and writes
``draft_graph_json`` + ``draft_canvas_schema_json`` together, so the two
representations can never diverge across revisions. Raises
``WorkflowDraftConflictError`` on a stale ``expected_revision``.
"""
raise NotImplementedError
@abstractmethod
async def publish_version(
self,
workflow_id: str,
*,
graph: dict[str, Any],
graph_hash: str,
published_by: str | None = None,
change_note: str = "",
canvas_schema_hash: str | None = None,
) -> dict[str, Any]:
raise NotImplementedError
@abstractmethod
async def list_versions(self, workflow_id: str) -> list[dict[str, Any]]:
raise NotImplementedError
@abstractmethod
async def get_version(self, version_id: str) -> dict[str, Any] | None:
raise NotImplementedError
@abstractmethod
async def list_published_summaries(self, *, exclude_workflow_id: str | None = None) -> list[dict[str, Any]]:
"""Lightweight list for subworkflow resource catalog."""
raise NotImplementedError