Add governed reusable workflow definitions
This commit is contained in:
@@ -1,21 +1,28 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping
|
||||
import hashlib
|
||||
import json
|
||||
import re
|
||||
import unicodedata
|
||||
|
||||
from sqlalchemy import select
|
||||
from sqlalchemy import or_, select
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from govoplan_core.auth import ApiPrincipal
|
||||
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.governance import (
|
||||
definition_governance_payload,
|
||||
require_definition_action,
|
||||
)
|
||||
from govoplan_workflow.backend.schemas import (
|
||||
WorkflowDefinitionCreateRequest,
|
||||
WorkflowDefinitionDeriveRequest,
|
||||
WorkflowDefinitionResponse,
|
||||
WorkflowDefinitionRevisionResponse,
|
||||
WorkflowDefinitionUpdateRequest,
|
||||
@@ -54,7 +61,10 @@ def list_definitions(
|
||||
session.scalars(
|
||||
select(WorkflowDefinition)
|
||||
.where(
|
||||
WorkflowDefinition.tenant_id == tenant_id,
|
||||
or_(
|
||||
WorkflowDefinition.tenant_id == tenant_id,
|
||||
WorkflowDefinition.tenant_id.is_(None),
|
||||
),
|
||||
WorkflowDefinition.deleted_at.is_(None),
|
||||
)
|
||||
.order_by(
|
||||
@@ -74,7 +84,10 @@ def get_definition(
|
||||
definition = session.scalar(
|
||||
select(WorkflowDefinition).where(
|
||||
WorkflowDefinition.id == definition_id,
|
||||
WorkflowDefinition.tenant_id == tenant_id,
|
||||
or_(
|
||||
WorkflowDefinition.tenant_id == tenant_id,
|
||||
WorkflowDefinition.tenant_id.is_(None),
|
||||
),
|
||||
WorkflowDefinition.deleted_at.is_(None),
|
||||
)
|
||||
)
|
||||
@@ -127,11 +140,34 @@ def create_definition(
|
||||
payload: WorkflowDefinitionCreateRequest,
|
||||
) -> WorkflowDefinition:
|
||||
graph = _validated_graph(payload.graph)
|
||||
stored_tenant_id = (
|
||||
None if payload.scope_type == "system" else tenant_id
|
||||
)
|
||||
scope_id = (
|
||||
None
|
||||
if payload.scope_type == "system"
|
||||
else tenant_id
|
||||
if payload.scope_type == "tenant"
|
||||
else payload.scope_id
|
||||
)
|
||||
scope_key = (
|
||||
"system"
|
||||
if payload.scope_type == "system"
|
||||
else f"{payload.scope_type}:{scope_id}"
|
||||
)
|
||||
definition = WorkflowDefinition(
|
||||
tenant_id=tenant_id,
|
||||
tenant_id=stored_tenant_id,
|
||||
scope_type=payload.scope_type,
|
||||
scope_id=scope_id,
|
||||
scope_key=scope_key,
|
||||
definition_kind=payload.definition_kind,
|
||||
inherit_to_lower_scopes=payload.inherit_to_lower_scopes,
|
||||
allow_start=payload.allow_start,
|
||||
allow_reuse=payload.allow_reuse,
|
||||
allow_automation=payload.allow_automation,
|
||||
definition_key=_available_key(
|
||||
session,
|
||||
tenant_id=tenant_id,
|
||||
scope_key=scope_key,
|
||||
requested=payload.key,
|
||||
name=payload.name,
|
||||
),
|
||||
@@ -146,7 +182,7 @@ def create_definition(
|
||||
)
|
||||
definition.revisions.append(
|
||||
_new_revision(
|
||||
tenant_id=tenant_id,
|
||||
tenant_id=stored_tenant_id,
|
||||
revision=1,
|
||||
graph=graph,
|
||||
actor_id=actor_id,
|
||||
@@ -176,18 +212,49 @@ def update_definition(
|
||||
f"expected revision {payload.expected_revision}, "
|
||||
f"current revision is {definition.current_revision}."
|
||||
)
|
||||
if (
|
||||
payload.scope_type != definition.scope_type
|
||||
or payload.scope_id != definition.scope_id
|
||||
and not (
|
||||
definition.scope_type == "tenant"
|
||||
and payload.scope_id in {None, definition.scope_id}
|
||||
)
|
||||
):
|
||||
raise WorkflowConflictError(
|
||||
"Definition scope is immutable; derive a scoped copy instead."
|
||||
)
|
||||
if payload.definition_kind != definition.definition_kind:
|
||||
raise WorkflowConflictError(
|
||||
"Definition kind is immutable; derive a flow or template instead."
|
||||
)
|
||||
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)
|
||||
ancestor_limits = _ancestor_governance_limits(
|
||||
definition.derivation_provenance
|
||||
)
|
||||
definition.inherit_to_lower_scopes = (
|
||||
payload.inherit_to_lower_scopes
|
||||
and ancestor_limits["inherit_to_lower_scopes"]
|
||||
)
|
||||
definition.allow_start = (
|
||||
payload.allow_start and ancestor_limits["allow_start"]
|
||||
)
|
||||
definition.allow_reuse = (
|
||||
payload.allow_reuse and ancestor_limits["allow_reuse"]
|
||||
)
|
||||
definition.allow_automation = (
|
||||
payload.allow_automation and ancestor_limits["allow_automation"]
|
||||
)
|
||||
definition.updated_by = actor_id
|
||||
if current.content_hash != graph_hash:
|
||||
definition.current_revision += 1
|
||||
definition.revisions.append(
|
||||
_new_revision(
|
||||
tenant_id=tenant_id,
|
||||
tenant_id=definition.tenant_id,
|
||||
revision=definition.current_revision,
|
||||
graph=graph,
|
||||
actor_id=actor_id,
|
||||
@@ -199,6 +266,127 @@ def update_definition(
|
||||
return definition
|
||||
|
||||
|
||||
def derive_definition(
|
||||
session: Session,
|
||||
*,
|
||||
tenant_id: str,
|
||||
actor_id: str | None,
|
||||
principal: ApiPrincipal,
|
||||
registry: object | None,
|
||||
source_definition_id: str,
|
||||
payload: WorkflowDefinitionDeriveRequest,
|
||||
) -> WorkflowDefinition:
|
||||
source = get_definition(
|
||||
session,
|
||||
tenant_id=tenant_id,
|
||||
definition_id=source_definition_id,
|
||||
)
|
||||
decision = require_definition_action(
|
||||
source,
|
||||
principal=principal,
|
||||
registry=registry,
|
||||
action="derive",
|
||||
)
|
||||
source_revision = get_definition_revision(
|
||||
session,
|
||||
definition=source,
|
||||
revision=payload.source_revision,
|
||||
)
|
||||
stored_tenant_id = (
|
||||
None if payload.scope_type == "system" else tenant_id
|
||||
)
|
||||
scope_id = (
|
||||
None
|
||||
if payload.scope_type == "system"
|
||||
else tenant_id
|
||||
if payload.scope_type == "tenant"
|
||||
else payload.scope_id
|
||||
)
|
||||
scope_key = (
|
||||
"system"
|
||||
if payload.scope_type == "system"
|
||||
else f"{payload.scope_type}:{scope_id}"
|
||||
)
|
||||
source_limits = _effective_governance_limits(
|
||||
source,
|
||||
decision_details=decision.details,
|
||||
)
|
||||
limits = {
|
||||
"inherit_to_lower_scopes": (
|
||||
source_limits["inherit_to_lower_scopes"]
|
||||
and payload.inherit_to_lower_scopes
|
||||
),
|
||||
"allow_start": (
|
||||
source_limits["allow_start"] and payload.allow_start
|
||||
),
|
||||
"allow_reuse": (
|
||||
source_limits["allow_reuse"] and payload.allow_reuse
|
||||
),
|
||||
"allow_automation": (
|
||||
source_limits["allow_automation"]
|
||||
and payload.allow_automation
|
||||
),
|
||||
}
|
||||
provenance = {
|
||||
"source_ref": f"workflow-definition:{source.id}",
|
||||
"source_scope": {
|
||||
"scope_type": source.scope_type,
|
||||
"scope_id": source.scope_id,
|
||||
},
|
||||
"source_definition_kind": source.definition_kind,
|
||||
"source_revision": source_revision.revision,
|
||||
"source_hash": source_revision.content_hash,
|
||||
"source_effective_limits": limits,
|
||||
"policy_decision": decision.to_dict(),
|
||||
"derived_by": actor_id,
|
||||
"derived_at": utcnow().isoformat(),
|
||||
}
|
||||
definition = WorkflowDefinition(
|
||||
tenant_id=stored_tenant_id,
|
||||
scope_type=payload.scope_type,
|
||||
scope_id=scope_id,
|
||||
scope_key=scope_key,
|
||||
definition_kind=payload.definition_kind,
|
||||
inherit_to_lower_scopes=limits["inherit_to_lower_scopes"],
|
||||
allow_start=limits["allow_start"],
|
||||
allow_reuse=limits["allow_reuse"],
|
||||
allow_automation=limits["allow_automation"],
|
||||
derived_from_definition_id=source.id,
|
||||
derived_from_revision=source_revision.revision,
|
||||
derived_from_hash=source_revision.content_hash,
|
||||
derivation_provenance=provenance,
|
||||
definition_key=_available_key(
|
||||
session,
|
||||
scope_key=scope_key,
|
||||
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(
|
||||
WorkflowDefinitionRevision(
|
||||
tenant_id=stored_tenant_id,
|
||||
revision=1,
|
||||
schema_version=source_revision.schema_version,
|
||||
graph=dict(source_revision.graph),
|
||||
content_hash=source_revision.content_hash,
|
||||
library_id=source_revision.library_id,
|
||||
library_version=source_revision.library_version,
|
||||
created_by=actor_id,
|
||||
)
|
||||
)
|
||||
session.add(definition)
|
||||
session.flush()
|
||||
return definition
|
||||
|
||||
|
||||
def activate_definition(
|
||||
session: Session,
|
||||
*,
|
||||
@@ -217,6 +405,10 @@ def activate_definition(
|
||||
definition=definition,
|
||||
revision=revision,
|
||||
)
|
||||
if definition.definition_kind == "template":
|
||||
raise WorkflowConflictError(
|
||||
"Workflow templates cannot be activated or started."
|
||||
)
|
||||
_validated_graph(WorkflowGraph.model_validate(selected.graph))
|
||||
definition.active_revision = selected.revision
|
||||
definition.status = "active"
|
||||
@@ -265,6 +457,8 @@ def definition_response(
|
||||
session: Session,
|
||||
definition: WorkflowDefinition,
|
||||
*,
|
||||
principal: ApiPrincipal,
|
||||
registry: object | None,
|
||||
revision: int | None = None,
|
||||
) -> WorkflowDefinitionResponse:
|
||||
selected = get_definition_revision(
|
||||
@@ -287,6 +481,11 @@ def definition_response(
|
||||
created_at=definition.created_at,
|
||||
updated_at=definition.updated_at,
|
||||
revision=revision_response(selected),
|
||||
governance=definition_governance_payload(
|
||||
definition,
|
||||
principal=principal,
|
||||
registry=registry,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
@@ -308,7 +507,7 @@ def revision_response(
|
||||
|
||||
def _new_revision(
|
||||
*,
|
||||
tenant_id: str,
|
||||
tenant_id: str | None,
|
||||
revision: int,
|
||||
graph: WorkflowGraph,
|
||||
actor_id: str | None,
|
||||
@@ -348,7 +547,7 @@ def _content_hash(graph: WorkflowGraph) -> str:
|
||||
def _available_key(
|
||||
session: Session,
|
||||
*,
|
||||
tenant_id: str,
|
||||
scope_key: str,
|
||||
requested: str | None,
|
||||
name: str,
|
||||
) -> str:
|
||||
@@ -357,7 +556,7 @@ def _available_key(
|
||||
suffix = 2
|
||||
while session.scalar(
|
||||
select(WorkflowDefinition.id).where(
|
||||
WorkflowDefinition.tenant_id == tenant_id,
|
||||
WorkflowDefinition.scope_key == scope_key,
|
||||
WorkflowDefinition.definition_key == candidate,
|
||||
)
|
||||
):
|
||||
@@ -380,6 +579,65 @@ def _clean_optional(value: str | None) -> str | None:
|
||||
return cleaned or None
|
||||
|
||||
|
||||
def _ancestor_governance_limits(
|
||||
provenance: Mapping[str, object],
|
||||
) -> dict[str, bool]:
|
||||
raw = provenance.get("source_effective_limits")
|
||||
limits = raw if isinstance(raw, Mapping) else {}
|
||||
return {
|
||||
key: value if isinstance((value := limits.get(key)), bool) else True
|
||||
for key in (
|
||||
"inherit_to_lower_scopes",
|
||||
"allow_start",
|
||||
"allow_reuse",
|
||||
"allow_automation",
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
def _effective_governance_limits(
|
||||
definition: WorkflowDefinition,
|
||||
*,
|
||||
decision_details: Mapping[str, object] | None = None,
|
||||
) -> dict[str, bool]:
|
||||
ancestor = _ancestor_governance_limits(
|
||||
definition.derivation_provenance
|
||||
)
|
||||
effective = {
|
||||
"inherit_to_lower_scopes": (
|
||||
definition.inherit_to_lower_scopes
|
||||
and ancestor["inherit_to_lower_scopes"]
|
||||
),
|
||||
"allow_start": (
|
||||
definition.allow_start and ancestor["allow_start"]
|
||||
),
|
||||
"allow_reuse": (
|
||||
definition.allow_reuse and ancestor["allow_reuse"]
|
||||
),
|
||||
"allow_automation": (
|
||||
definition.allow_automation
|
||||
and ancestor["allow_automation"]
|
||||
),
|
||||
}
|
||||
policy_limits = (
|
||||
decision_details.get("effective_limits")
|
||||
if decision_details is not None
|
||||
else None
|
||||
)
|
||||
if isinstance(policy_limits, Mapping):
|
||||
policy_key_by_local_key = {
|
||||
"inherit_to_lower_scopes": "inherit_to_lower_scopes",
|
||||
"allow_start": "allow_run",
|
||||
"allow_reuse": "allow_reuse",
|
||||
"allow_automation": "allow_automation",
|
||||
}
|
||||
for local_key, policy_key in policy_key_by_local_key.items():
|
||||
value = policy_limits.get(policy_key)
|
||||
if isinstance(value, bool):
|
||||
effective[local_key] = effective[local_key] and value
|
||||
return effective
|
||||
|
||||
|
||||
__all__ = [
|
||||
"WorkflowConflictError",
|
||||
"WorkflowError",
|
||||
@@ -388,6 +646,7 @@ __all__ = [
|
||||
"activate_definition",
|
||||
"archive_definition",
|
||||
"create_definition",
|
||||
"derive_definition",
|
||||
"definition_response",
|
||||
"delete_definition",
|
||||
"get_definition",
|
||||
|
||||
Reference in New Issue
Block a user