Files
govoplan-workflow/src/govoplan_workflow/backend/service.py

400 lines
11 KiB
Python

from __future__ import annotations
import hashlib
import json
import re
import unicodedata
from sqlalchemy import select
from sqlalchemy.orm import Session
from govoplan_core.db.base import utcnow
from govoplan_workflow.backend.db.models import (
WorkflowDefinition,
WorkflowDefinitionRevision,
)
from govoplan_workflow.backend.node_library import WORKFLOW_GRAPH_LIBRARY
from govoplan_workflow.backend.schemas import (
WorkflowDefinitionCreateRequest,
WorkflowDefinitionResponse,
WorkflowDefinitionRevisionResponse,
WorkflowDefinitionUpdateRequest,
WorkflowGraph,
)
from govoplan_workflow.backend.validation import validate_workflow_graph
class WorkflowError(RuntimeError):
pass
class WorkflowNotFoundError(WorkflowError):
pass
class WorkflowConflictError(WorkflowError):
pass
class WorkflowValidationError(WorkflowError):
def __init__(self, diagnostics: tuple[object, ...]) -> None:
first = diagnostics[0] if diagnostics else None
super().__init__(
str(getattr(first, "message", "Workflow definition validation failed."))
)
self.diagnostics = diagnostics
def list_definitions(
session: Session,
*,
tenant_id: str,
) -> list[WorkflowDefinition]:
return list(
session.scalars(
select(WorkflowDefinition)
.where(
WorkflowDefinition.tenant_id == tenant_id,
WorkflowDefinition.deleted_at.is_(None),
)
.order_by(
WorkflowDefinition.updated_at.desc(),
WorkflowDefinition.name.asc(),
)
)
)
def get_definition(
session: Session,
*,
tenant_id: str,
definition_id: str,
) -> WorkflowDefinition:
definition = session.scalar(
select(WorkflowDefinition).where(
WorkflowDefinition.id == definition_id,
WorkflowDefinition.tenant_id == tenant_id,
WorkflowDefinition.deleted_at.is_(None),
)
)
if definition is None:
raise WorkflowNotFoundError("Workflow definition not found.")
return definition
def get_definition_revision(
session: Session,
*,
definition: WorkflowDefinition,
revision: int | None = None,
) -> WorkflowDefinitionRevision:
revision_number = revision or definition.current_revision
item = session.scalar(
select(WorkflowDefinitionRevision).where(
WorkflowDefinitionRevision.definition_id == definition.id,
WorkflowDefinitionRevision.tenant_id == definition.tenant_id,
WorkflowDefinitionRevision.revision == revision_number,
)
)
if item is None:
raise WorkflowNotFoundError("Workflow definition revision not found.")
return item
def list_definition_revisions(
session: Session,
*,
definition: WorkflowDefinition,
) -> list[WorkflowDefinitionRevision]:
return list(
session.scalars(
select(WorkflowDefinitionRevision)
.where(
WorkflowDefinitionRevision.definition_id == definition.id,
WorkflowDefinitionRevision.tenant_id == definition.tenant_id,
)
.order_by(WorkflowDefinitionRevision.revision.desc())
)
)
def create_definition(
session: Session,
*,
tenant_id: str,
actor_id: str | None,
payload: WorkflowDefinitionCreateRequest,
) -> WorkflowDefinition:
graph = _validated_graph(payload.graph)
definition = WorkflowDefinition(
tenant_id=tenant_id,
definition_key=_available_key(
session,
tenant_id=tenant_id,
requested=payload.key,
name=payload.name,
),
name=payload.name.strip(),
description=_clean_optional(payload.description),
status="draft",
current_revision=1,
active_revision=None,
metadata_=dict(payload.metadata),
created_by=actor_id,
updated_by=actor_id,
)
definition.revisions.append(
_new_revision(
tenant_id=tenant_id,
revision=1,
graph=graph,
actor_id=actor_id,
)
)
session.add(definition)
session.flush()
return definition
def update_definition(
session: Session,
*,
tenant_id: str,
definition_id: str,
actor_id: str | None,
payload: WorkflowDefinitionUpdateRequest,
) -> WorkflowDefinition:
definition = get_definition(
session,
tenant_id=tenant_id,
definition_id=definition_id,
)
if payload.expected_revision != definition.current_revision:
raise WorkflowConflictError(
"Workflow definition changed on the server; "
f"expected revision {payload.expected_revision}, "
f"current revision is {definition.current_revision}."
)
graph = _validated_graph(payload.graph)
current = get_definition_revision(session, definition=definition)
graph_hash = _content_hash(graph)
definition.name = payload.name.strip()
definition.description = _clean_optional(payload.description)
definition.metadata_ = dict(payload.metadata)
definition.updated_by = actor_id
if current.content_hash != graph_hash:
definition.current_revision += 1
definition.revisions.append(
_new_revision(
tenant_id=tenant_id,
revision=definition.current_revision,
graph=graph,
actor_id=actor_id,
)
)
if definition.status == "active":
definition.status = "draft"
session.flush()
return definition
def activate_definition(
session: Session,
*,
tenant_id: str,
definition_id: str,
actor_id: str | None,
revision: int | None = None,
) -> WorkflowDefinition:
definition = get_definition(
session,
tenant_id=tenant_id,
definition_id=definition_id,
)
selected = get_definition_revision(
session,
definition=definition,
revision=revision,
)
_validated_graph(WorkflowGraph.model_validate(selected.graph))
definition.active_revision = selected.revision
definition.status = "active"
definition.updated_by = actor_id
session.flush()
return definition
def archive_definition(
session: Session,
*,
tenant_id: str,
definition_id: str,
actor_id: str | None,
) -> WorkflowDefinition:
definition = get_definition(
session,
tenant_id=tenant_id,
definition_id=definition_id,
)
definition.status = "archived"
definition.updated_by = actor_id
session.flush()
return definition
def delete_definition(
session: Session,
*,
tenant_id: str,
definition_id: str,
actor_id: str | None,
) -> WorkflowDefinition:
definition = get_definition(
session,
tenant_id=tenant_id,
definition_id=definition_id,
)
definition.deleted_at = utcnow()
definition.updated_by = actor_id
session.flush()
return definition
def definition_response(
session: Session,
definition: WorkflowDefinition,
*,
revision: int | None = None,
) -> WorkflowDefinitionResponse:
selected = get_definition_revision(
session,
definition=definition,
revision=revision,
)
return WorkflowDefinitionResponse(
id=definition.id,
tenant_id=definition.tenant_id,
key=definition.definition_key,
name=definition.name,
description=definition.description,
status=definition.status,
current_revision=definition.current_revision,
active_revision=definition.active_revision,
metadata=dict(definition.metadata_),
created_by=definition.created_by,
updated_by=definition.updated_by,
created_at=definition.created_at,
updated_at=definition.updated_at,
revision=revision_response(selected),
)
def revision_response(
revision: WorkflowDefinitionRevision,
) -> WorkflowDefinitionRevisionResponse:
return WorkflowDefinitionRevisionResponse(
id=revision.id,
revision=revision.revision,
schema_version=revision.schema_version,
graph=WorkflowGraph.model_validate(revision.graph),
content_hash=revision.content_hash,
library_id=revision.library_id,
library_version=revision.library_version,
created_by=revision.created_by,
created_at=revision.created_at,
)
def _new_revision(
*,
tenant_id: str,
revision: int,
graph: WorkflowGraph,
actor_id: str | None,
) -> WorkflowDefinitionRevision:
return WorkflowDefinitionRevision(
tenant_id=tenant_id,
revision=revision,
schema_version=graph.schema_version,
graph=_canonical_graph(graph),
content_hash=_content_hash(graph),
library_id=WORKFLOW_GRAPH_LIBRARY.id,
library_version=WORKFLOW_GRAPH_LIBRARY.version,
created_by=actor_id,
)
def _validated_graph(graph: WorkflowGraph) -> WorkflowGraph:
diagnostics = validate_workflow_graph(graph)
if any(item.severity == "error" for item in diagnostics):
raise WorkflowValidationError(diagnostics)
return graph
def _canonical_graph(graph: WorkflowGraph) -> dict[str, object]:
return graph.model_dump(mode="json")
def _content_hash(graph: WorkflowGraph) -> str:
encoded = json.dumps(
_canonical_graph(graph),
sort_keys=True,
separators=(",", ":"),
)
return hashlib.sha256(encoded.encode("utf-8")).hexdigest()
def _available_key(
session: Session,
*,
tenant_id: str,
requested: str | None,
name: str,
) -> str:
base = _slug(requested or name)
candidate = base
suffix = 2
while session.scalar(
select(WorkflowDefinition.id).where(
WorkflowDefinition.tenant_id == tenant_id,
WorkflowDefinition.definition_key == candidate,
)
):
candidate = f"{base[: max(1, 120 - len(str(suffix)) - 1)]}-{suffix}"
suffix += 1
return candidate
def _slug(value: str) -> str:
normalized = unicodedata.normalize("NFKD", value)
ascii_value = normalized.encode("ascii", "ignore").decode("ascii").lower()
cleaned = re.sub(r"[^a-z0-9]+", "-", ascii_value).strip("-")
return (cleaned or "workflow")[:120]
def _clean_optional(value: str | None) -> str | None:
if value is None:
return None
cleaned = value.strip()
return cleaned or None
__all__ = [
"WorkflowConflictError",
"WorkflowError",
"WorkflowNotFoundError",
"WorkflowValidationError",
"activate_definition",
"archive_definition",
"create_definition",
"definition_response",
"delete_definition",
"get_definition",
"get_definition_revision",
"list_definition_revisions",
"list_definitions",
"revision_response",
"update_definition",
]