3 Commits
Author SHA1 Message Date
zemion b0a9dc9739 feat: add resumable external work hand-offs
Module Package Release / publish-packages (push) Successful in 10s
2026-08-22 02:14:35 +02:00
zemion ceb61b5867 feat(workflow-engine): add governed DSAR coverage 2026-08-21 03:33:06 +02:00
zemion 9174e07118 Project durable workflow handoffs as work items 2026-08-06 16:06:17 +02:00
17 changed files with 3333 additions and 61 deletions
+24
View File
@@ -1,5 +1,22 @@
# GovOPlaN Workflow Engine
## Data-subject requests
Workflow Engine publishes `privacy.dsar.workflow_engine` for exact definition,
revision, instance, step, event, trigger, delivery, and wait references and for
structured instance authorization, work assignments, automation authority, and
minimized operator attribution. It never exports graphs/BPMN, inputs, context,
outputs, handoffs, events/configuration, authorization snapshots, errors,
replay keys, external references, hashes, or credentials. Owning service, case,
form, and other source modules locate and correct facts inside arbitrary
runtime payloads.
Subject-linked automation can be disabled and revoked, while exact terminal
delivery detail can be minimized idempotently without changing its replay key.
Definitions, instances, live work, assignments, transition and decision
events, waits, and institutional attribution require authorized review or
retention.
<!-- govoplan-repository-type:start -->
**Repository type:** module (platform).
<!-- govoplan-repository-type:end -->
@@ -39,6 +56,13 @@ human handoffs; event filters and variable mappings are bounded JSON
expressions and never executable code. Cron remains an optional governed
scheduler-adapter concern.
External hand-off waits bind an immutable, revision-bearing domain reference
to a focused action URL and resume only from declared terminal platform events.
The runtime rechecks the provider's current authorization and resource revision
before continuing. Observational events are duplicate-safe, navigation alone
never completes work, and a missing optional projection remains visible without
changing the domain outcome.
Consequential module actions are staged in Core's durable recovery ledger
before provider dispatch. Conclusive effects commit with the Workflow
projection. Lost acknowledgements block continuation and expose evidence-based
+3
View File
@@ -119,6 +119,9 @@ The first executable slice now provides:
- activation-bound API, one-time/interval schedule, platform-event, and
parent-workflow trigger registrations with exact-revision dispatch
- durable duration/deadline/event wait subscriptions and scale-out-safe claims
- resumable external hand-offs with immutable references, focused action URLs,
duplicate-safe observed events, terminal outcome ports, timeout paths, and a
provider authorization/revision recheck before continuation
- a separate transactional platform-event consumer with bounded JSON filters
and variable mappings, idempotent delivery, and current-authority rechecks
- Core recovery-ledger operations for module actions, including canonical
+18
View File
@@ -135,3 +135,21 @@ baseline revision inactive when an older revision is active, and fails closed
when required capabilities or interfaces are absent. Baselines are immutable;
editing derives a pinned local override, and reset archives that override
without removing revision or instance history.
## Common Work Inbox
Workflow Engine also projects the current human handoff of a waiting instance
through Core's versioned work-item provider contract. Tasks may aggregate that
projection when it is enabled, while Workflow Engine remains the sole owner of
the instance, step, allowed actions, and completion state.
Human activities and reviews accept account, group, role, function, and
function-assignment responsibility references. A missing assignee defaults to
the account that started the instance. `due_after` uses the same bounded
duration syntax as Workflow timers (`5m`, `2h`, `1d`, or ISO 8601). The Engine
materializes assignment and due-date columns on the current step so inbox reads
remain tenant-scoped and indexed; historical JSON handoffs are interpreted
conservatively for upgrade compatibility. Automated waits and in-flight
provider calls are not presented as human work. Recovery-required, unknown,
dependency, and failed handoffs appear as blocked work and always link back to
the pinned Workflow instance.
+2 -2
View File
@@ -4,13 +4,13 @@ build-backend = "setuptools.build_meta"
[project]
name = "govoplan-workflow-engine"
version = "0.1.18"
version = "0.1.19"
description = "Headless, versioned workflow definition and execution engine for GovOPlaN."
readme = "README.md"
requires-python = ">=3.12"
license = "AGPL-3.0-or-later"
authors = [{ name = "GovOPlaN" }]
dependencies = ["defusedxml>=0.7,<1", "govoplan-core>=0.1.18"]
dependencies = ["defusedxml>=0.7,<1", "govoplan-core>=0.1.28"]
[tool.setuptools.packages.find]
where = ["src"]
@@ -228,6 +228,10 @@ def legacy_graph_to_bpmn(graph: WorkflowGraph) -> WorkflowGraph:
config["wait_mode"] = config.pop("mode", "manual")
config.setdefault("event_definition", "none")
bpmn_type = "bpmn.intermediateCatchEvent"
elif node_type == "workflow.external_handoff":
config["wait_mode"] = "external_handoff"
config.setdefault("event_definition", "message")
bpmn_type = "bpmn.intermediateCatchEvent"
elif node_type == "workflow.capability":
config["implementation"] = "capability"
bpmn_type = "bpmn.serviceTask"
@@ -1002,8 +1006,13 @@ def _runtime_node(node: WorkflowNode) -> WorkflowNode:
elif node.type in {"bpmn.task", "bpmn.userTask", "bpmn.manualTask"}:
node_type = "workflow.activity"
elif node.type in {"bpmn.receiveTask", "bpmn.intermediateCatchEvent"}:
node_type = "workflow.wait"
config["mode"] = config.get("wait_mode") or "event"
wait_mode = str(config.get("wait_mode") or "event")
node_type = (
"workflow.external_handoff"
if wait_mode == "external_handoff"
else "workflow.wait"
)
config["mode"] = wait_mode
elif node.type in {"bpmn.serviceTask", "bpmn.sendTask"}:
implementation = str(config.get("implementation") or "capability")
node_type = (
@@ -408,6 +408,13 @@ class WorkflowInstanceStep(Base, TimestampMixin):
"tenant_id",
"status",
),
Index(
"ix_workflow_instance_steps_work_assignment",
"tenant_id",
"status",
"work_assignment_kind",
"work_assignment_id",
),
)
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid)
@@ -465,6 +472,25 @@ class WorkflowInstanceStep(Base, TimestampMixin):
String(255),
nullable=True,
)
work_assignment_kind: Mapped[str | None] = mapped_column(
String(40),
nullable=True,
index=True,
)
work_assignment_id: Mapped[str | None] = mapped_column(
String(255),
nullable=True,
index=True,
)
work_assignment_label: Mapped[str | None] = mapped_column(
String(500),
nullable=True,
)
work_due_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True),
nullable=True,
index=True,
)
instance: Mapped[WorkflowInstance] = relationship(back_populates="steps")
@@ -0,0 +1,879 @@
from __future__ import annotations
from collections.abc import Sequence
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Any
from sqlalchemy import or_
from sqlalchemy.orm import Session
from govoplan_core.core.dsar import (
DsarErasureActionRef,
DsarExecutionResultRef,
DsarRecordRef,
DsarSubjectRef,
dsar_capability_name,
)
from govoplan_workflow_engine.backend.db.models import (
WorkflowDefinition,
WorkflowDefinitionRevision,
WorkflowInstance,
WorkflowInstanceEvent,
WorkflowInstanceStep,
WorkflowTrigger,
WorkflowTriggerDelivery,
WorkflowWaitState,
)
WORKFLOW_ENGINE_DSAR_CAPABILITY = dsar_capability_name("workflow_engine")
_MAX_RECORDS = 5_000
_CONFLICT = object()
_DIRECT_ALIASES = {
"definition_id": ("workflow_engine.definition", "workflow.definition"),
"revision_id": (
"workflow_engine.definition_revision",
"workflow_engine.revision",
),
"instance_id": ("workflow_engine.instance", "workflow.instance"),
"step_id": ("workflow_engine.step", "workflow.step"),
"event_id": ("workflow_engine.event", "workflow.event"),
"trigger_id": ("workflow_engine.trigger", "workflow.trigger"),
"delivery_id": (
"workflow_engine.trigger_delivery",
"workflow_engine.delivery",
),
"wait_id": ("workflow_engine.wait_state", "workflow_engine.wait"),
}
_RESOURCE_MODELS = {
"workflow_definition": WorkflowDefinition,
"workflow_definition_revision": WorkflowDefinitionRevision,
"workflow_instance": WorkflowInstance,
"workflow_instance_step": WorkflowInstanceStep,
"workflow_instance_event": WorkflowInstanceEvent,
"workflow_trigger": WorkflowTrigger,
"workflow_trigger_delivery": WorkflowTriggerDelivery,
"workflow_wait_state": WorkflowWaitState,
}
@dataclass(frozen=True, slots=True)
class _Selectors:
account_id: str | None
identity_id: str | None
membership_id: str | None
direct: dict[str, str]
@property
def subject_ids(self) -> tuple[str, ...]:
return tuple(
value
for value in (self.account_id, self.identity_id, self.membership_id)
if value
)
@dataclass(frozen=True, slots=True)
class _Match:
resource_type: str
row: Any
category: str
class WorkflowEngineDsarProvider:
provider_id = "workflow_engine"
module_id = "workflow_engine"
def search_subject(
self,
session: object,
*,
tenant_id: str,
subject: DsarSubjectRef,
) -> Sequence[DsarRecordRef]:
db = _session(session)
selectors = _selectors(subject)
if selectors is None or not (selectors.subject_ids or selectors.direct):
return ()
direct = _direct_matches(db, tenant_id=tenant_id, selectors=selectors)
if direct is None:
return ()
if direct:
if selectors.subject_ids and not all(
_correlates(
db,
tenant_id=tenant_id,
match=match,
subject_ids=selectors.subject_ids,
)
for match in direct
):
return ()
matches = direct
else:
matches = _canonical_matches(
db,
tenant_id=tenant_id,
selectors=selectors,
)
records: list[DsarRecordRef] = []
seen: set[tuple[str, str]] = set()
for match in matches:
key = (match.resource_type, str(match.row.id))
if key in seen:
continue
if len(records) >= _MAX_RECORDS:
raise ValueError(
"Workflow Engine DSAR result limit exceeded; narrow the selectors."
)
seen.add(key)
records.append(_record(match))
return tuple(records)
def plan_erasure(
self,
session: object,
*,
tenant_id: str,
subject: DsarSubjectRef,
records: Sequence[DsarRecordRef],
) -> Sequence[DsarErasureActionRef]:
del tenant_id
_session(session)
if _selectors(subject) is None:
raise ValueError("Workflow Engine DSAR subject selectors conflict.")
actions: list[DsarErasureActionRef] = []
for record in records:
_validate_record(record)
kind = _planned_kind(record)
executable = kind in {"anonymize", "revoke"}
actions.append(
DsarErasureActionRef(
action_id=(
f"workflow_engine:{kind}:{record.resource_type}:"
f"{record.resource_id}"
),
provider_id=self.provider_id,
module_id=self.module_id,
kind=kind,
resource_type=record.resource_type,
resource_id=record.resource_id,
title=(
f"Minimize {record.title}"
if kind == "anonymize"
else f"{kind.replace('_', ' ').title()} {record.title}"
),
rationale=_rationale(record, kind=kind),
executable=executable,
irreversible=kind == "anonymize",
metadata={"record_category": record.category},
)
)
return tuple(actions)
def execute_erasure(
self,
session: object,
*,
tenant_id: str,
subject: DsarSubjectRef,
actions: Sequence[DsarErasureActionRef],
request_id: str,
) -> Sequence[DsarExecutionResultRef]:
db = _session(session)
selectors = _selectors(subject)
if selectors is None:
raise ValueError("Workflow Engine DSAR subject selectors conflict.")
results: list[DsarExecutionResultRef] = []
for action in actions:
_validate_action(action)
if not action.executable:
results.append(
DsarExecutionResultRef(
action_id=action.action_id,
status="blocked",
summary=(
"Review live work, third parties, decision evidence, "
"retention, and downstream effects before changing state."
),
evidence={"request_id": request_id},
)
)
continue
model = _RESOURCE_MODELS[action.resource_type]
row = (
db.query(model)
.filter(model.tenant_id == tenant_id, model.id == action.resource_id)
.with_for_update()
.one_or_none()
)
if row is None:
status = "unchanged"
summary = "Workflow row was already absent or minimized."
else:
match = _Match(action.resource_type, row, "execution")
if not (
_directly_targets(selectors, match)
or _correlates(
db,
tenant_id=tenant_id,
match=match,
subject_ids=selectors.subject_ids,
)
or _already_revoked(action.resource_type, row)
):
raise ValueError(
"Workflow Engine DSAR action is not corroborated by the subject."
)
status, summary = _execute_action(
db,
resource_type=action.resource_type,
row=row,
kind=action.kind,
)
results.append(
DsarExecutionResultRef(
action_id=action.action_id,
status=status,
summary=summary,
evidence={"request_id": request_id},
)
)
return tuple(results)
def _direct_matches(
session: Session,
*,
tenant_id: str,
selectors: _Selectors,
) -> list[_Match] | None:
matches: list[_Match] = []
for selector, raw_value in selectors.direct.items():
value = _strip_prefix(raw_value)
if selector == "definition_id":
definition = _one(session, WorkflowDefinition, tenant_id, value)
if definition is None:
return None
current = _definition_package(
session,
tenant_id=tenant_id,
definition=definition,
)
elif selector == "instance_id":
instance = _one(session, WorkflowInstance, tenant_id, value)
if instance is None:
return None
current = _instance_package(
session,
tenant_id=tenant_id,
instance=instance,
)
else:
model, resource_type = {
"revision_id": (
WorkflowDefinitionRevision,
"workflow_definition_revision",
),
"step_id": (WorkflowInstanceStep, "workflow_instance_step"),
"event_id": (WorkflowInstanceEvent, "workflow_instance_event"),
"trigger_id": (WorkflowTrigger, "workflow_trigger"),
"delivery_id": (
WorkflowTriggerDelivery,
"workflow_trigger_delivery",
),
"wait_id": (WorkflowWaitState, "workflow_wait_state"),
}[selector]
row = _one(session, model, tenant_id, value)
if row is None:
return None
current = [_Match(resource_type, row, _direct_category(resource_type, row))]
matches.extend(current)
if len(matches) > _MAX_RECORDS:
raise ValueError(
"Workflow Engine DSAR result limit exceeded; narrow the selectors."
)
return matches
def _definition_package(
session: Session,
*,
tenant_id: str,
definition: WorkflowDefinition,
) -> list[_Match]:
matches = [_Match("workflow_definition", definition, "workflow_definition_package")]
specs = (
(WorkflowDefinitionRevision, "workflow_definition_revision"),
(WorkflowInstance, "workflow_instance"),
(WorkflowTrigger, "workflow_trigger"),
(WorkflowTriggerDelivery, "workflow_trigger_delivery"),
)
instances: list[WorkflowInstance] = []
for model, resource_type in specs:
rows = (
session.query(model)
.filter(model.tenant_id == tenant_id, model.definition_id == definition.id)
.order_by(model.id)
.limit(_MAX_RECORDS + 1)
.all()
)
if model is WorkflowInstance:
instances = rows
matches.extend(
_Match(resource_type, row, "workflow_definition_package") for row in rows
)
instance_ids = {row.id for row in instances}
if instance_ids:
for model, resource_type in (
(WorkflowInstanceStep, "workflow_instance_step"),
(WorkflowInstanceEvent, "workflow_instance_event"),
(WorkflowWaitState, "workflow_wait_state"),
):
rows = (
session.query(model)
.filter(
model.tenant_id == tenant_id,
model.instance_id.in_(instance_ids),
)
.order_by(model.id)
.limit(_MAX_RECORDS + 1)
.all()
)
matches.extend(
_Match(resource_type, row, "workflow_definition_package")
for row in rows
)
return matches
def _instance_package(
session: Session,
*,
tenant_id: str,
instance: WorkflowInstance,
) -> list[_Match]:
matches = [_Match("workflow_instance", instance, "workflow_instance_package")]
specs = (
(WorkflowInstanceStep, "workflow_instance_step"),
(WorkflowInstanceEvent, "workflow_instance_event"),
(WorkflowWaitState, "workflow_wait_state"),
)
for model, resource_type in specs:
rows = (
session.query(model)
.filter(model.tenant_id == tenant_id, model.instance_id == instance.id)
.order_by(model.id)
.limit(_MAX_RECORDS + 1)
.all()
)
matches.extend(
_Match(resource_type, row, "workflow_instance_package") for row in rows
)
return matches
def _canonical_matches(
session: Session,
*,
tenant_id: str,
selectors: _Selectors,
) -> list[_Match]:
subject_ids = selectors.subject_ids
if not subject_ids:
return []
triggers = (
session.query(WorkflowTrigger)
.filter(
WorkflowTrigger.tenant_id == tenant_id,
or_(
WorkflowTrigger.authorization_account_id.in_(subject_ids),
WorkflowTrigger.authorization_membership_id.in_(subject_ids),
),
)
.order_by(WorkflowTrigger.id)
.limit(_MAX_RECORDS + 1)
.all()
)
matches = [
_Match("workflow_trigger", row, "workflow_automation_authority")
for row in triggers
]
instances = (
session.query(WorkflowInstance)
.filter(WorkflowInstance.tenant_id == tenant_id)
.order_by(WorkflowInstance.id)
.limit(_MAX_RECORDS + 1)
.all()
)
if len(instances) > _MAX_RECORDS:
raise ValueError(
"Workflow Engine DSAR instance scan limit exceeded; use an exact instance reference."
)
for row in instances:
authorization = dict(row.authorization_ or {})
if any(
str(authorization.get(key) or "") == value
for key, value in (
("account_id", selectors.account_id),
("identity_id", selectors.identity_id),
("membership_id", selectors.membership_id),
)
if value
):
matches.append(
_Match(
"workflow_instance",
row,
"workflow_subject_authorization",
)
)
steps = (
session.query(WorkflowInstanceStep)
.filter(
WorkflowInstanceStep.tenant_id == tenant_id,
WorkflowInstanceStep.work_assignment_id.in_(subject_ids),
)
.order_by(WorkflowInstanceStep.id)
.limit(_MAX_RECORDS + 1)
.all()
)
matches.extend(
_Match("workflow_instance_step", row, "workflow_subject_work_assignment")
for row in steps
)
if len(matches) > _MAX_RECORDS:
raise ValueError(
"Workflow Engine DSAR result limit exceeded; narrow the selectors."
)
specs = (
(
WorkflowDefinition,
or_(
WorkflowDefinition.created_by.in_(subject_ids),
WorkflowDefinition.updated_by.in_(subject_ids),
),
"workflow_definition",
),
(
WorkflowDefinitionRevision,
WorkflowDefinitionRevision.created_by.in_(subject_ids),
"workflow_definition_revision",
),
(
WorkflowInstance,
WorkflowInstance.created_by.in_(subject_ids),
"workflow_instance",
),
(
WorkflowInstanceStep,
WorkflowInstanceStep.completed_by.in_(subject_ids),
"workflow_instance_step",
),
(
WorkflowInstanceEvent,
WorkflowInstanceEvent.actor_id.in_(subject_ids),
"workflow_instance_event",
),
(
WorkflowTrigger,
or_(
WorkflowTrigger.created_by.in_(subject_ids),
WorkflowTrigger.updated_by.in_(subject_ids),
),
"workflow_trigger",
),
)
for model, condition, resource_type in specs:
rows = (
session.query(model)
.filter(model.tenant_id == tenant_id, condition)
.order_by(model.id)
.limit(_MAX_RECORDS + 1)
.all()
)
matches.extend(
_Match(resource_type, row, "workflow_operator_attribution") for row in rows
)
if len(matches) > _MAX_RECORDS:
raise ValueError(
"Workflow Engine DSAR result limit exceeded; narrow the selectors."
)
return matches
def _one(session: Session, model: Any, tenant_id: str, row_id: str) -> Any | None:
return (
session.query(model)
.filter(model.tenant_id == tenant_id, model.id == row_id)
.one_or_none()
)
def _direct_category(resource_type: str, row: Any) -> str:
if resource_type == "workflow_trigger_delivery":
return (
"terminal_workflow_delivery"
if row.status in {"succeeded", "failed", "skipped", "cancelled"}
else "active_workflow_delivery"
)
return "workflow_governed_state"
def _correlates(
session: Session,
*,
tenant_id: str,
match: _Match,
subject_ids: tuple[str, ...],
) -> bool:
row = match.row
if any(
str(getattr(row, field, "") or "") in subject_ids
for field in (
"created_by",
"updated_by",
"completed_by",
"actor_id",
"work_assignment_id",
"authorization_account_id",
"authorization_membership_id",
)
):
return True
if match.resource_type == "workflow_instance":
authorization = dict(row.authorization_ or {})
if any(
str(authorization.get(key) or "") in subject_ids
for key in ("account_id", "identity_id", "membership_id")
):
return True
instance_id = getattr(row, "instance_id", None)
if instance_id:
instance = session.get(WorkflowInstance, instance_id)
if instance and instance.tenant_id == tenant_id:
return _correlates(
session,
tenant_id=tenant_id,
match=_Match("workflow_instance", instance, "related"),
subject_ids=subject_ids,
)
trigger_id = getattr(row, "trigger_id", None)
if trigger_id:
trigger = session.get(WorkflowTrigger, trigger_id)
if trigger and trigger.tenant_id == tenant_id:
return _correlates(
session,
tenant_id=tenant_id,
match=_Match("workflow_trigger", trigger, "related"),
subject_ids=subject_ids,
)
definition_id = getattr(row, "definition_id", None)
if definition_id:
definition = session.get(WorkflowDefinition, definition_id)
return bool(
definition
and definition.tenant_id == tenant_id
and (
definition.created_by in subject_ids
or definition.updated_by in subject_ids
)
)
return False
def _directly_targets(selectors: _Selectors, match: _Match) -> bool:
selector = {
"workflow_definition": "definition_id",
"workflow_definition_revision": "revision_id",
"workflow_instance": "instance_id",
"workflow_instance_step": "step_id",
"workflow_instance_event": "event_id",
"workflow_trigger": "trigger_id",
"workflow_trigger_delivery": "delivery_id",
"workflow_wait_state": "wait_id",
}[match.resource_type]
return _strip_prefix(selectors.direct.get(selector, "")) == str(match.row.id)
def _record(match: _Match) -> DsarRecordRef:
row = match.row
immutable = match.category == "workflow_operator_attribution"
return DsarRecordRef(
provider_id="workflow_engine",
module_id="workflow_engine",
resource_type=match.resource_type,
resource_id=str(row.id),
category=match.category,
title=match.resource_type.removeprefix("workflow_").replace("_", " ").title(),
data={
key: value
for key, value in _record_data(match.resource_type, row).items()
if value is not None
},
observed_at=_observed_at(row),
immutable_evidence=immutable,
retention_reason=(
"Institutional Workflow authorship, assignments, transitions, and "
"decision events remain attributable for governance and audit."
if immutable
else None
),
source_path="/workflow",
)
def _record_data(resource_type: str, row: Any) -> dict[str, object]:
if resource_type == "workflow_definition":
return {
"scope_type": row.scope_type,
"definition_kind": row.definition_kind,
"status": row.status,
"current_revision": row.current_revision,
"active_revision": row.active_revision,
"allow_start": row.allow_start,
"allow_reuse": row.allow_reuse,
"allow_automation": row.allow_automation,
"deleted_at": _iso(row.deleted_at),
}
if resource_type == "workflow_definition_revision":
return {
"revision": row.revision,
"schema_version": row.schema_version,
"execution_mode": row.execution_mode,
"bpmn_runtime_kind": row.bpmn_runtime_kind,
"bpmn_executable": row.bpmn_executable,
"created_at": _iso(row.created_at),
}
if resource_type == "workflow_instance":
return {
"status": row.status,
"start_origin": row.start_origin,
"started_at": _iso(row.started_at),
"finished_at": _iso(row.finished_at),
"cancellation_requested_at": _iso(row.cancellation_requested_at),
}
if resource_type == "workflow_instance_step":
return {
"sequence": row.sequence,
"node_type": row.node_type,
"status": row.status,
"attempt": row.attempt,
"work_assignment_kind": row.work_assignment_kind,
"work_due_at": _iso(row.work_due_at),
"started_at": _iso(row.started_at),
"finished_at": _iso(row.finished_at),
}
if resource_type == "workflow_instance_event":
return {
"sequence": row.sequence,
"kind": row.kind,
"created_at": _iso(row.created_at),
}
if resource_type == "workflow_trigger":
return {
"kind": row.kind,
"status": row.status,
"event_type": row.event_type,
"authorization_subject_kind": row.authorization_subject_kind,
"next_fire_at": _iso(row.next_fire_at),
"last_fire_at": _iso(row.last_fire_at),
}
if resource_type == "workflow_trigger_delivery":
return {
"invocation_kind": row.invocation_kind,
"status": row.status,
"attempts": row.attempts,
"scheduled_for": _iso(row.scheduled_for),
"created_at": _iso(row.created_at),
}
return {
"mode": row.mode,
"status": row.status,
"due_at": _iso(row.due_at),
"event_type": row.event_type,
"revision": row.revision,
"created_at": _iso(row.created_at),
"updated_at": _iso(row.updated_at),
}
def _planned_kind(record: DsarRecordRef) -> str:
if record.category == "terminal_workflow_delivery":
return "anonymize"
if record.category == "workflow_automation_authority":
return "revoke"
if record.category == "workflow_operator_attribution":
return "retain"
return "manual_review"
def _rationale(record: DsarRecordRef, *, kind: str) -> str:
if kind == "anonymize":
return (
"Clear terminal trigger event and error detail while preserving the "
"replay identity and minimal delivery evidence."
)
if kind == "revoke":
return (
"Disable automation and remove subject-linked delegated authority "
"without deleting historical workflow evidence."
)
if kind == "retain":
return record.retention_reason or "Retain institutional attribution evidence."
return (
"A Workflow owner must review live work, assignments, third parties, "
"decision evidence, source authority, and retention."
)
def _execute_action(
session: Session,
*,
resource_type: str,
row: Any,
kind: str,
) -> tuple[str, str]:
if resource_type == "workflow_trigger_delivery" and kind == "anonymize":
if row.status not in {"succeeded", "failed", "skipped", "cancelled"}:
raise ValueError("Active Workflow deliveries require manual review.")
changed = _replace_fields(row, {"event_": None, "error": None})
elif resource_type == "workflow_trigger" and kind == "revoke":
changed = _replace_fields(
row,
{
"status": "disabled",
"config_": {},
"last_error": None,
"authorization_account_id": None,
"authorization_membership_id": None,
"authorization_service_account_id": None,
"authorization_ref": f"revoked:{row.id}",
"grant_scopes": [],
},
)
else:
raise ValueError("Workflow Engine DSAR executable action is unsupported.")
if changed:
session.flush()
return "executed", "Personal Workflow detail minimized."
return "unchanged", "Personal Workflow detail was already minimized."
def _replace_fields(row: Any, values: dict[str, object]) -> bool:
changed = False
for field, value in values.items():
if getattr(row, field) != value:
setattr(row, field, value)
changed = True
return changed
def _already_revoked(resource_type: str, row: Any) -> bool:
return bool(
resource_type == "workflow_trigger"
and row.status == "disabled"
and row.authorization_account_id is None
and row.authorization_membership_id is None
and row.authorization_service_account_id is None
and row.authorization_ref == f"revoked:{row.id}"
)
def _selectors(subject: DsarSubjectRef) -> _Selectors | None:
refs = subject.external_references
canonical = (
_coalesce(
subject.account_id,
refs.get("workflow_engine.account"),
refs.get("access.account"),
),
_coalesce(
subject.identity_id,
refs.get("workflow_engine.identity"),
refs.get("identity.id"),
),
_coalesce(
subject.membership_id,
refs.get("workflow_engine.membership"),
refs.get("tenancy.membership"),
),
)
if any(value is _CONFLICT for value in canonical):
return None
direct: dict[str, str] = {}
for selector, aliases in _DIRECT_ALIASES.items():
value = _coalesce(*(refs.get(alias) for alias in aliases))
if value is _CONFLICT:
return None
if value:
direct[selector] = str(value)
return _Selectors(
account_id=_optional(canonical[0]),
identity_id=_optional(canonical[1]),
membership_id=_optional(canonical[2]),
direct=direct,
)
def _coalesce(*values: str | None) -> str | None | object:
normalized = {str(value).strip() for value in values if str(value or "").strip()}
if len(normalized) > 1:
return _CONFLICT
return next(iter(normalized), None)
def _optional(value: object) -> str | None:
return value if isinstance(value, str) and value else None
def _strip_prefix(value: str) -> str:
return value.partition(":")[2] if ":" in value else value
def _observed_at(row: Any) -> datetime | None:
for field in ("finished_at", "created_at", "updated_at", "started_at"):
value = getattr(row, field, None)
if isinstance(value, datetime):
return _aware(value)
return None
def _iso(value: datetime | None) -> str | None:
aware = _aware(value)
return aware.isoformat() if aware else None
def _aware(value: datetime | None) -> datetime | None:
if value is None or value.tzinfo is not None:
return value
return value.replace(tzinfo=timezone.utc)
def _session(value: object) -> Session:
if not isinstance(value, Session):
raise TypeError("Workflow Engine DSAR requires a SQLAlchemy Session.")
return value
def _validate_record(record: DsarRecordRef) -> None:
if record.provider_id != "workflow_engine" or record.module_id != "workflow_engine":
raise ValueError("Workflow Engine DSAR cannot plan a foreign provider record.")
if record.resource_type not in _RESOURCE_MODELS or not record.resource_id:
raise ValueError("Workflow Engine DSAR record identity is invalid.")
def _validate_action(action: DsarErasureActionRef) -> None:
if action.provider_id != "workflow_engine" or action.module_id != "workflow_engine":
raise ValueError(
"Workflow Engine DSAR cannot execute a foreign provider action."
)
if action.resource_type not in _RESOURCE_MODELS or not action.action_id.startswith(
"workflow_engine:"
):
raise ValueError("Workflow Engine DSAR action identity is invalid.")
__all__ = ["WORKFLOW_ENGINE_DSAR_CAPABILITY", "WorkflowEngineDsarProvider"]
@@ -1,7 +1,7 @@
from __future__ import annotations
from collections.abc import Mapping
from datetime import datetime
from datetime import datetime, timedelta
import hashlib
import logging
@@ -751,7 +751,7 @@ def cancel_instance(
instance.id,
exc,
)
if step.node_type == "workflow.wait":
if step.node_type in {"workflow.wait", "workflow.external_handoff"}:
from govoplan_workflow_engine.backend.triggers import (
resolve_wait_state,
)
@@ -1010,6 +1010,25 @@ def _drive_instance(
registry=registry,
)
return
if node.type == "workflow.external_handoff":
from govoplan_workflow_engine.backend.triggers import (
register_external_handoff_state,
)
handoff_state = register_external_handoff_state(
session,
instance=instance,
step=step,
node=node,
)
_set_external_handoff(
session,
instance=instance,
step=step,
node=node,
wait_state=handoff_state,
)
return
if node.type == "workflow.wait":
from govoplan_workflow_engine.backend.triggers import (
register_wait_state,
@@ -1294,8 +1313,7 @@ def _execute_capability_step(
action_input=action_input,
preview_payload=preview_payload,
backup_reference=(
str(node.config.get("recovery_backup_reference") or "").strip()
or None
str(node.config.get("recovery_backup_reference") or "").strip() or None
),
approval_reference=(
str(node.config.get("recovery_approval_reference") or "").strip()
@@ -1368,11 +1386,7 @@ def _execute_capability_step(
"operation_id": exc.operation_id,
"status": recovery_status,
"requires_attention": outcome_unknown,
**(
{"next_call_number": call_number + 1}
if safe_to_retry
else {}
),
**({"next_call_number": call_number + 1} if safe_to_retry else {}),
},
},
)
@@ -1473,9 +1487,7 @@ def _execute_capability_step(
action_recovery.commit_unknown(
session,
error_type=type(exc).__name__,
message=(
"Inspect the provider by stable idempotency key before any retry"
),
message=("Inspect the provider by stable idempotency key before any retry"),
)
return True
if not isinstance(result, ActionExecutionResult):
@@ -1608,9 +1620,7 @@ def _execute_capability_step(
session,
provider_state=result.state,
result_sha256=canonical_sha256(result_payload),
observed_effects_sha256=canonical_sha256(
result_payload["observed_effects"]
),
observed_effects_sha256=canonical_sha256(result_payload["observed_effects"]),
)
_drive_instance(
session,
@@ -1915,6 +1925,10 @@ def _set_action_handoff(
"suggested_port": "failure",
**details_payload,
}
if state in {"pending", "running"}:
_clear_work_projection(step)
else:
_apply_work_projection(instance, step)
instance.status = "waiting"
instance.error = step.error
if previous.get("state") != state or previous.get("message") != message:
@@ -2109,6 +2123,7 @@ def _handle_dataflow_success(
"warnings": warnings,
"output": output,
}
_apply_work_projection(instance, step)
instance.status = "waiting"
_record_event(
session,
@@ -2165,6 +2180,7 @@ def _complete_step(
step.finished_at = utcnow()
step.completed_by = actor_id
step.handoff = {}
_clear_work_projection(step)
context = dict(instance.context_)
step_values = dict(context.get("steps") or {})
step_values[step.node_id] = dict(output)
@@ -2183,6 +2199,80 @@ def _complete_step(
return _next_node_id(graph, step.node_id, port)
_WORK_ASSIGNMENT_KINDS = {
"account",
"group",
"role",
"function",
"function_assignment",
"anyone",
}
def _work_assignment(
instance: WorkflowInstance,
configured: object | None = None,
) -> dict[str, str | None] | None:
if isinstance(configured, Mapping):
kind = str(configured.get("kind") or "").strip()
assignment_id = str(configured.get("id") or "").strip()
label = str(configured.get("label") or "").strip() or None
if kind in _WORK_ASSIGNMENT_KINDS and assignment_id:
if kind == "anyone" and assignment_id != "*":
return None
return {"kind": kind, "id": assignment_id, "label": label}
return None
value = str(configured or "").strip()
if value:
prefix, separator, remainder = value.partition(":")
if separator and prefix in _WORK_ASSIGNMENT_KINDS and remainder.strip():
assignment_id = remainder.strip()
if prefix == "anyone" and assignment_id != "*":
return None
return {"kind": prefix, "id": assignment_id, "label": None}
return {"kind": "account", "id": value, "label": None}
account_id = str(instance.authorization_.get("account_id") or "").strip()
if not account_id:
return None
return {"kind": "account", "id": account_id, "label": None}
def _work_due_at(configured: object | None) -> datetime | None:
if not str(configured or "").strip():
return None
from govoplan_workflow_engine.backend.triggers import duration_seconds
return utcnow() + timedelta(seconds=duration_seconds(configured))
def _apply_work_projection(
instance: WorkflowInstance,
step: WorkflowInstanceStep,
*,
assignment: Mapping[str, object] | None = None,
due_at: datetime | None = None,
) -> None:
normalized = (
dict(assignment) if assignment is not None else _work_assignment(instance)
)
if normalized is None:
_clear_work_projection(step)
return
step.work_assignment_kind = str(normalized.get("kind") or "") or None
step.work_assignment_id = str(normalized.get("id") or "") or None
step.work_assignment_label = str(normalized.get("label") or "") or None
step.work_due_at = due_at
def _clear_work_projection(step: WorkflowInstanceStep) -> None:
step.work_assignment_kind = None
step.work_assignment_id = None
step.work_assignment_label = None
step.work_due_at = None
def _set_human_handoff(
session: Session,
*,
@@ -2200,6 +2290,11 @@ def _set_human_handoff(
else:
actions = ["complete", "cancel"]
kind = "activity"
assignment = _work_assignment(
instance,
node.config.get("reviewer") or node.config.get("assignee"),
)
due_at = _work_due_at(node.config.get("due_after"))
step.status = "waiting"
step.handoff = {
"kind": kind,
@@ -2207,9 +2302,12 @@ def _set_human_handoff(
"title": str(node.config.get("title") or node.label or node.type),
"instructions": str(node.config.get("instructions") or ""),
"assignee": node.config.get("reviewer") or node.config.get("assignee"),
"assignment": assignment,
"due_at": due_at.isoformat() if due_at else None,
"required_evidence": list(node.config.get("required_evidence") or []),
"allowed_actions": actions,
}
_apply_work_projection(instance, step, assignment=assignment, due_at=due_at)
instance.status = "waiting"
_record_event(
session,
@@ -2249,6 +2347,7 @@ def _set_automated_wait(
"event_type": event_type,
"allowed_actions": ["cancel"],
}
_clear_work_projection(step)
instance.status = "waiting"
_record_event(
session,
@@ -2260,6 +2359,51 @@ def _set_automated_wait(
)
def _set_external_handoff(
session: Session,
*,
instance: WorkflowInstance,
step: WorkflowInstanceStep,
node: WorkflowNode,
wait_state: object,
) -> None:
config = dict(getattr(wait_state, "config_"))
due_at = getattr(wait_state, "due_at")
optional_capabilities = {
str(key): bool(value)
for key, value in dict(config.get("optional_capabilities") or {}).items()
}
unavailable = sorted(
name for name, available in optional_capabilities.items() if not available
)
step.status = "waiting"
step.external_ref = str(config.get("immutable_ref") or "") or None
step.handoff = {
"kind": "external_handoff",
"state": "assigned",
"title": str(node.config.get("title") or node.label or "External hand-off"),
"instructions": str(node.config.get("instructions") or ""),
"event_type": getattr(wait_state, "event_type"),
"external_id": config.get("external_id"),
"immutable_ref": config.get("immutable_ref"),
"action_url": config.get("action_url"),
"due_at": due_at.isoformat() if due_at else None,
"optional_capabilities": optional_capabilities,
"unavailable_optional_capabilities": unavailable,
"allowed_actions": [],
}
_clear_work_projection(step)
instance.status = "waiting"
_record_event(
session,
instance,
step=step,
kind="workflow.external_handoff.registered",
actor_id=instance.created_by,
payload=dict(step.handoff),
)
def _set_dependency_handoff(
session: Session,
*,
@@ -2275,6 +2419,7 @@ def _set_dependency_handoff(
"message": message,
"allowed_actions": ["retry", "cancel"],
}
_apply_work_projection(instance, step)
instance.status = "waiting"
instance.error = message
_record_event(
@@ -2382,6 +2527,7 @@ def _set_failure_handoff(
"allowed_actions": ["retry", "reject", "cancel"],
"suggested_port": "failure",
}
_apply_work_projection(instance, step)
instance.status = "waiting"
instance.error = message
_record_event(
@@ -2725,6 +2871,26 @@ def _require_runtime_dependencies(
principal=principal,
registry=registry,
)
if node.type == "workflow.external_handoff":
from govoplan_core.core.campaigns import (
CampaignWorkOrchestrationProvider,
)
capability_name = str(
node.config.get("provider_capability") or ""
).strip()
capability = (
registry.capability(capability_name)
if capability_name
and registry is not None
and hasattr(registry, "has_capability")
and registry.has_capability(capability_name)
else None
)
if not isinstance(capability, CampaignWorkOrchestrationProvider):
raise WorkflowConflictError(
f"External hand-off capability {capability_name!r} is not available."
)
def _authorization_payload(
@@ -9,12 +9,15 @@ from govoplan_core.core.access import (
CAPABILITY_AUTH_PRINCIPAL_RESOLVER,
)
from govoplan_core.core.dataflows import CAPABILITY_DATAFLOW_RUN_LIFECYCLE
from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION
from govoplan_core.core.idm import CAPABILITY_IDM_DIRECTORY
from govoplan_core.core.module_guards import (
drop_table_retirement_provider,
persistent_table_uninstall_guard,
)
from govoplan_core.core.modules import (
CapabilityDocumentation,
DocumentationCondition,
DocumentationTopic,
MigrationSpec,
ModuleContext,
@@ -33,6 +36,7 @@ from govoplan_core.core.notifications import (
)
from govoplan_core.core.references import CAPABILITY_ACCESS_REFERENCE_OPTIONS
from govoplan_core.core.views import CAPABILITY_VIEWS_RESOLVER
from govoplan_core.core.tasks import WorkItemProviderRegistration
from govoplan_core.core.workflows import (
CAPABILITY_WORKFLOW_DEFINITION_CONTRIBUTIONS,
CAPABILITY_WORKFLOW_ORCHESTRATION,
@@ -41,6 +45,10 @@ from govoplan_core.core.workflows import (
)
from govoplan_core.db.base import Base
from govoplan_workflow_engine.backend.db import models as workflow_models
from govoplan_workflow_engine.backend.dsar_provider import (
WORKFLOW_ENGINE_DSAR_CAPABILITY,
WorkflowEngineDsarProvider,
)
from govoplan_workflow_engine.backend.configuration_provider import (
WORKFLOW_CONFIGURATION_CAPABILITY,
)
@@ -52,7 +60,7 @@ from govoplan_workflow_engine.backend.service_launcher import (
MODULE_ID = "workflow_engine"
MODULE_NAME = "Workflow Engine"
MODULE_VERSION = "0.1.18"
MODULE_VERSION = "0.1.19"
DEFINITION_READ_SCOPE = "workflow:definition:read"
DEFINITION_WRITE_SCOPE = "workflow:definition:write"
@@ -187,6 +195,17 @@ def _service_launcher(context: ModuleContext) -> WorkflowServiceLauncher:
return WorkflowServiceLauncher(registry=context.registry)
def _dsar_provider(context: ModuleContext) -> WorkflowEngineDsarProvider:
del context
return WorkflowEngineDsarProvider()
def _work_items(context: ModuleContext):
from govoplan_workflow_engine.backend.work_items import WorkflowWorkItemProvider
return WorkflowWorkItemProvider(registry=context.registry)
manifest = ModuleManifest(
id=MODULE_ID,
name=MODULE_NAME,
@@ -196,8 +215,10 @@ manifest = ModuleManifest(
optional_dependencies=(
"access",
"audit",
"campaigns",
"dataflow",
"datasources",
"idm",
"notifications",
"policy",
"tasks",
@@ -209,7 +230,9 @@ manifest = ModuleManifest(
CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER,
CAPABILITY_AUTH_PRINCIPAL_RESOLVER,
CAPABILITY_AUTH_PERMISSION_EVALUATOR,
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
CAPABILITY_DATAFLOW_RUN_LIFECYCLE,
CAPABILITY_IDM_DIRECTORY,
CAPABILITY_NOTIFICATIONS_DISPATCH,
CAPABILITY_POLICY_DEFINITION_GOVERNANCE,
CAPABILITY_VIEWS_RESOLVER,
@@ -244,8 +267,18 @@ manifest = ModuleManifest(
name="workflow.bpmn_execution_adapters",
version="1.0.0",
),
ModuleInterfaceProvider(
name=WORKFLOW_ENGINE_DSAR_CAPABILITY,
version="0.1.0",
),
),
requires_interfaces=(
ModuleInterfaceRequirement(
name="campaigns.work_orchestration",
version_min="1.0.0",
version_max_exclusive="2.0.0",
optional=True,
),
ModuleInterfaceRequirement(
name=CAPABILITY_ACCESS_REFERENCE_OPTIONS,
version_min="0.1.0",
@@ -295,6 +328,7 @@ manifest = ModuleManifest(
CAPABILITY_WORKFLOW_ORCHESTRATION: _orchestration_provider,
WORKFLOW_CONFIGURATION_CAPABILITY: _configuration_provider,
CAPABILITY_WORKFLOW_SERVICE_LAUNCHER: _service_launcher,
WORKFLOW_ENGINE_DSAR_CAPABILITY: _dsar_provider,
},
capability_documentation={
CAPABILITY_WORKFLOW_SERVICE_LAUNCHER: CapabilityDocumentation(
@@ -302,7 +336,21 @@ manifest = ModuleManifest(
summary="Starts an authorized active Workflow from an exact available Service revision.",
contract_version="0.1.0",
),
WORKFLOW_ENGINE_DSAR_CAPABILITY: CapabilityDocumentation(
label="Workflow Engine data-subject request provider",
summary="Finds subject-linked workflow state while preserving live work and decision evidence.",
contract_version="0.1.0",
documentation_types=("admin", "user"),
audience=("privacy_officer", "workflow_operator", "user"),
),
},
work_item_providers=(
WorkItemProviderRegistration(
id="workflow_engine.handoffs",
factory=_work_items,
order=20,
),
),
migration_spec=MigrationSpec(
module_id=MODULE_ID,
metadata=Base.metadata,
@@ -338,6 +386,71 @@ manifest = ModuleManifest(
),
),
documentation=(
DocumentationTopic(
id="workflow.external-campaign-handoffs",
title="Resume Workflow from accountable Campaign work",
summary=(
"Open revision-bearing Campaign work, focus the relevant UI, "
"and continue from declared lifecycle events without polling."
),
body=(
"The external hand-off primitive follows a module action that "
"creates or references a Campaign and opens one accountable work "
"assignment. Workflow stores the exact Campaign version and assignment "
"revision, a safe action link, correlation and idempotency provenance, "
"and a durable event subscription. Assigned, accepted, started, and "
"reassigned events update the waiting state; completed, rejected, "
"cancelled, and timed-out outcomes take separate graph paths. Duplicate "
"events are ignored. Before a terminal event resumes execution, Workflow "
"re-resolves its authority and Campaign rechecks both access and the event "
"revision. An assignment never grants Campaign access. Missing Campaign or "
"Views capabilities leave the instance inspectable and blocked; missing "
"Tasks or Notifications is reported as reduced optional integration and "
"does not change the Campaign outcome."
),
layer="available",
documentation_types=("admin", "user"),
audience=("operator", "workflow_designer", "campaign_manager", "module_admin"),
related_modules=("campaigns", "views", "tasks", "notifications", "audit"),
conditions=(
DocumentationCondition(
any_scopes=(INSTANCE_READ_SCOPE, ADMIN_SCOPE),
),
),
order=75,
metadata={
"kind": "workflow",
"help_contexts": [
"workflow.external-handoff",
"workflow.instances",
],
"limitations": [
"The Campaign provider must be installed to start or resume a Campaign hand-off.",
"Focused View projection is optional; the safe Campaign action link remains available without it.",
],
},
),
DocumentationTopic(
id="workflow.data-subject-requests",
title="Workflow data-subject requests",
summary="Review structured subject links without exposing or guessing arbitrary process payload content.",
body=(
"Workflow Engine matches exact tenant-scoped definition, revision, instance, step, event, trigger, delivery, and wait identifiers. It also matches structured instance authorization, work assignment, automation authority, and minimized staff attribution for account, identity, and membership selectors. DSAR results never copy graphs or BPMN, inputs, context, outputs, handoffs, event/configuration payloads, authorization snapshots, errors, replay keys, external references, hashes, or credentials. Arbitrary runtime payloads are not scanned for identifiers because the owning service, case, form, or other source module remains authoritative. "
"Subject-linked automation can be disabled and revoked, and explicitly selected terminal trigger delivery detail can be minimized idempotently while its replay key remains. Definitions, instances, live work, assignments, immutable transition/decision events, waits, and operator attribution require workflow-owner, retention, third-party, and source-authority review."
),
layer="configured",
documentation_types=("admin", "user"),
audience=("user", "operator", "module_admin", "auditor"),
related_modules=("core", "cases", "forms_runtime", "services", "tasks"),
order=75,
metadata={
"help_contexts": [
"workflow.data-subject-requests",
"workflow.instances",
"workflow.triggers",
],
},
),
DocumentationTopic(
id="workflow.definition-graphs",
title="Workflow definition graphs",
@@ -351,6 +464,8 @@ manifest = ModuleManifest(
"revisions; activation pins the exact revision used by future instances. "
"The optional service-launch capability starts an authorized active "
"revision from an exact Portal Service binding and records that provenance. "
"When Tasks is enabled, current human handoffs are projected into the common "
"work inbox with typed responsibility, due date, and a resumable source link."
),
layer="available",
documentation_types=("admin", "user"),
@@ -431,6 +546,7 @@ manifest = ModuleManifest(
"workflow revision",
"workflow instance",
"work transition",
"external event hand-off",
"execution adapter binding",
),
non_owned_concepts=(
@@ -0,0 +1,65 @@
"""v0.1.18 Workflow work projections.
Revision ID: 8d5a2f7c1b4e
Revises: e4a1f8c2d7b6
"""
from __future__ import annotations
from alembic import op
import sqlalchemy as sa
revision = "8d5a2f7c1b4e"
down_revision = "e4a1f8c2d7b6"
branch_labels = None
depends_on = None
def upgrade() -> None:
with op.batch_alter_table("workflow_instance_steps") as batch_op:
batch_op.add_column(
sa.Column("work_assignment_kind", sa.String(length=40), nullable=True)
)
batch_op.add_column(
sa.Column("work_assignment_id", sa.String(length=255), nullable=True)
)
batch_op.add_column(
sa.Column("work_assignment_label", sa.String(length=500), nullable=True)
)
batch_op.add_column(
sa.Column("work_due_at", sa.DateTime(timezone=True), nullable=True)
)
batch_op.create_index(
"ix_workflow_instance_steps_work_assignment_kind",
["work_assignment_kind"],
)
batch_op.create_index(
"ix_workflow_instance_steps_work_assignment_id",
["work_assignment_id"],
)
batch_op.create_index(
"ix_workflow_instance_steps_work_due_at",
["work_due_at"],
)
batch_op.create_index(
"ix_workflow_instance_steps_work_assignment",
[
"tenant_id",
"status",
"work_assignment_kind",
"work_assignment_id",
],
)
def downgrade() -> None:
with op.batch_alter_table("workflow_instance_steps") as batch_op:
batch_op.drop_index("ix_workflow_instance_steps_work_assignment")
batch_op.drop_index("ix_workflow_instance_steps_work_due_at")
batch_op.drop_index("ix_workflow_instance_steps_work_assignment_id")
batch_op.drop_index("ix_workflow_instance_steps_work_assignment_kind")
batch_op.drop_column("work_due_at")
batch_op.drop_column("work_assignment_label")
batch_op.drop_column("work_assignment_id")
batch_op.drop_column("work_assignment_kind")
@@ -84,8 +84,7 @@ LEGACY_WORKFLOW_NODE_TYPES = (
category="trigger",
label="Parent workflow",
description=(
"Start as a pinned child or dependency of another Workflow "
"instance."
"Start as a pinned child or dependency of another Workflow instance."
),
icon="git-branch",
default_config={
@@ -165,8 +164,12 @@ LEGACY_WORKFLOW_NODE_TYPES = (
icon="square-check-big",
input_ports=(DefinitionPort(id="input", label="Input"),),
config_fields=(
DefinitionConfigField(id="title", label="Title", kind="text", required=True),
DefinitionConfigField(id="instructions", label="Instructions", kind="textarea"),
DefinitionConfigField(
id="title", label="Title", kind="text", required=True
),
DefinitionConfigField(
id="instructions", label="Instructions", kind="textarea"
),
DefinitionConfigField(id="assignee", label="Assignee", kind="subject"),
DefinitionConfigField(id="due_after", label="Due after", kind="duration"),
FOCUSED_VIEW_SURFACES_FIELD,
@@ -192,8 +195,11 @@ LEGACY_WORKFLOW_NODE_TYPES = (
DefinitionPort(id="rejected", label="Rejected", required=False),
),
config_fields=(
DefinitionConfigField(id="title", label="Title", kind="text", required=True),
DefinitionConfigField(
id="title", label="Title", kind="text", required=True
),
DefinitionConfigField(id="reviewer", label="Reviewer", kind="subject"),
DefinitionConfigField(id="due_after", label="Due after", kind="duration"),
DefinitionConfigField(
id="required_evidence",
label="Required evidence",
@@ -204,6 +210,7 @@ LEGACY_WORKFLOW_NODE_TYPES = (
default_config={
"title": "",
"reviewer": "",
"due_after": "",
"required_evidence": [],
"view_surface_ids": [],
},
@@ -262,6 +269,86 @@ LEGACY_WORKFLOW_NODE_TYPES = (
"view_surface_ids": [],
},
),
DefinitionNodeType(
type="workflow.external_handoff",
category="integration",
label="External hand-off",
description=(
"Wait for revision-bearing module events without polling, while "
"preserving a focused action link and explicit terminal outcomes."
),
icon="arrow-left-right",
input_ports=(DefinitionPort(id="input", label="Input"),),
output_ports=(
DefinitionPort(id="completed", label="Completed", required=False),
DefinitionPort(id="rejected", label="Rejected", required=False),
DefinitionPort(id="cancelled", label="Cancelled", required=False),
DefinitionPort(id="timed_out", label="Timed out", required=False),
),
config_fields=(
DefinitionConfigField(
id="provider_capability",
label="Hand-off provider",
kind="capability",
required=True,
),
DefinitionConfigField(
id="event_type",
label="Event type",
kind="text",
required=True,
),
DefinitionConfigField(
id="event_filter",
label="Event filter",
kind="mapping",
required=True,
),
DefinitionConfigField(
id="external_id",
label="External reference path",
kind="expression",
required=True,
),
DefinitionConfigField(
id="action_url",
label="Action URL path",
kind="expression",
required=True,
),
DefinitionConfigField(
id="immutable_ref",
label="Immutable reference path",
kind="expression",
required=True,
),
DefinitionConfigField(
id="timeout_after",
label="Timeout after",
kind="duration",
),
FOCUSED_VIEW_SURFACES_FIELD,
),
default_config={
"provider_capability": "",
"event_type": "",
"event_filter": {},
"outcome_path": "payload.outcome",
"terminal_outcomes": {
"completed": "completed",
"rejected": "rejected",
"cancelled": "cancelled",
},
"observed_outcomes": ["assigned", "accepted", "started", "reassigned"],
"external_id": "",
"expected_revision": "",
"action_url": "",
"immutable_ref": "",
"optional_capabilities": {},
"timeout_after": "",
"view_surface_ids": [],
},
),
DefinitionNodeType(
type="workflow.capability",
category="integration",
@@ -292,7 +379,9 @@ LEGACY_WORKFLOW_NODE_TYPES = (
kind="text",
required=True,
),
DefinitionConfigField(id="input_mapping", label="Input mapping", kind="mapping"),
DefinitionConfigField(
id="input_mapping", label="Input mapping", kind="mapping"
),
DefinitionConfigField(
id="idempotency_key",
label="Idempotency key",
@@ -373,8 +462,7 @@ LEGACY_WORKFLOW_NODE_TYPES = (
label="Publication datasource",
kind="text",
description=(
"Optional stable Datasource target for materialized "
"output."
"Optional stable Datasource target for materialized output."
),
),
DefinitionConfigField(
@@ -397,7 +485,9 @@ LEGACY_WORKFLOW_NODE_TYPES = (
("continue", "Follow failure path"),
),
),
DefinitionConfigField(id="input_mapping", label="Input mapping", kind="mapping"),
DefinitionConfigField(
id="input_mapping", label="Input mapping", kind="mapping"
),
FOCUSED_VIEW_SURFACES_FIELD,
),
default_config={
@@ -428,7 +518,9 @@ LEGACY_WORKFLOW_NODE_TYPES = (
),
output_ports=(),
config_fields=(
DefinitionConfigField(id="output_mapping", label="Output mapping", kind="mapping"),
DefinitionConfigField(
id="output_mapping", label="Output mapping", kind="mapping"
),
),
default_config={"output_mapping": {}},
),
@@ -870,7 +962,9 @@ BPMN_NODE_TYPES = (
shape="activity",
config_fields=(
*_TASK_FIELDS,
DefinitionConfigField(id="message_ref", label="Message reference", kind="text"),
DefinitionConfigField(
id="message_ref", label="Message reference", kind="text"
),
FOCUSED_VIEW_SURFACES_FIELD,
),
default_config={
@@ -898,7 +992,13 @@ BPMN_NODE_TYPES = (
"Script task",
"A BPMN script task retained as notation; arbitrary scripts are not executed.",
"file-code-2",
(*_TASK_FIELDS, DefinitionConfigField(id="script_format", label="Script format", kind="text"), DefinitionConfigField(id="script", label="Script", kind="textarea")),
(
*_TASK_FIELDS,
DefinitionConfigField(
id="script_format", label="Script format", kind="text"
),
DefinitionConfigField(id="script", label="Script", kind="textarea"),
),
{"title": "", "instructions": "", "script_format": "", "script": ""},
),
(
@@ -906,7 +1006,14 @@ BPMN_NODE_TYPES = (
"Business rule task",
"Evaluate a governed business-rule implementation.",
"scale",
(*_TASK_FIELDS, DefinitionConfigField(id="implementation_ref", label="Implementation reference", kind="text")),
(
*_TASK_FIELDS,
DefinitionConfigField(
id="implementation_ref",
label="Implementation reference",
kind="text",
),
),
{"title": "", "instructions": "", "implementation_ref": ""},
),
(
@@ -914,7 +1021,15 @@ BPMN_NODE_TYPES = (
"Call activity",
"Call another reusable BPMN process or GovOPlaN workflow.",
"external-link",
(*_TASK_FIELDS, DefinitionConfigField(id="called_element", label="Called element", kind="text", required=True)),
(
*_TASK_FIELDS,
DefinitionConfigField(
id="called_element",
label="Called element",
kind="text",
required=True,
),
),
{"title": "", "instructions": "", "called_element": ""},
),
(
@@ -956,11 +1071,41 @@ BPMN_NODE_TYPES = (
runtime_support=runtime_support,
)
for type_name, label, description, icon, runtime_support in (
("exclusiveGateway", "Exclusive gateway", "Choose exactly one matching sequence flow.", "diamond", "native"),
("parallelGateway", "Parallel gateway", "Split or join concurrent sequence flows.", "plus", "model_only"),
("inclusiveGateway", "Inclusive gateway", "Choose one or more matching sequence flows.", "circle-plus", "model_only"),
("eventBasedGateway", "Event-based gateway", "Choose a path according to the first caught event.", "radio-tower", "model_only"),
("complexGateway", "Complex gateway", "Apply a complex activation condition.", "asterisk", "model_only"),
(
"exclusiveGateway",
"Exclusive gateway",
"Choose exactly one matching sequence flow.",
"diamond",
"native",
),
(
"parallelGateway",
"Parallel gateway",
"Split or join concurrent sequence flows.",
"plus",
"model_only",
),
(
"inclusiveGateway",
"Inclusive gateway",
"Choose one or more matching sequence flows.",
"circle-plus",
"model_only",
),
(
"eventBasedGateway",
"Event-based gateway",
"Choose a path according to the first caught event.",
"radio-tower",
"model_only",
),
(
"complexGateway",
"Complex gateway",
"Apply a complex activation condition.",
"asterisk",
"model_only",
),
)
),
_bpmn_node(
@@ -972,8 +1117,12 @@ BPMN_NODE_TYPES = (
shape="data-object",
input_ports=_OPTIONAL_INCOMING,
config_fields=(
DefinitionConfigField(id="data_object_ref", label="Data object reference", kind="text"),
DefinitionConfigField(id="item_subject_ref", label="Item definition", kind="text"),
DefinitionConfigField(
id="data_object_ref", label="Data object reference", kind="text"
),
DefinitionConfigField(
id="item_subject_ref", label="Item definition", kind="text"
),
),
default_config={"data_object_ref": "", "item_subject_ref": ""},
),
@@ -986,8 +1135,12 @@ BPMN_NODE_TYPES = (
shape="data-store",
input_ports=_OPTIONAL_INCOMING,
config_fields=(
DefinitionConfigField(id="data_store_ref", label="Data store reference", kind="text"),
DefinitionConfigField(id="item_subject_ref", label="Item definition", kind="text"),
DefinitionConfigField(
id="data_store_ref", label="Data store reference", kind="text"
),
DefinitionConfigField(
id="item_subject_ref", label="Item definition", kind="text"
),
),
default_config={"data_store_ref": "", "item_subject_ref": ""},
),
@@ -1000,7 +1153,9 @@ BPMN_NODE_TYPES = (
shape="participant",
input_ports=_OPTIONAL_INCOMING,
config_fields=(
DefinitionConfigField(id="process_ref", label="Process reference", kind="text"),
DefinitionConfigField(
id="process_ref", label="Process reference", kind="text"
),
),
default_config={"process_ref": ""},
),
@@ -1013,7 +1168,9 @@ BPMN_NODE_TYPES = (
shape="lane",
input_ports=_OPTIONAL_INCOMING,
config_fields=(
DefinitionConfigField(id="flow_node_refs", label="Flow node references", kind="string_list"),
DefinitionConfigField(
id="flow_node_refs", label="Flow node references", kind="string_list"
),
),
default_config={"flow_node_refs": []},
),
@@ -1040,7 +1197,9 @@ BPMN_NODE_TYPES = (
shape="group",
input_ports=_OPTIONAL_INCOMING,
config_fields=(
DefinitionConfigField(id="category_value_ref", label="Category value", kind="text"),
DefinitionConfigField(
id="category_value_ref", label="Category value", kind="text"
),
),
default_config={"category_value_ref": ""},
),
@@ -264,6 +264,151 @@ def register_wait_state(
return state
def register_external_handoff_state(
session: Session,
*,
instance: WorkflowInstance,
step: WorkflowInstanceStep,
node: WorkflowNode,
now: datetime | None = None,
) -> WorkflowWaitState:
"""Persist one event-driven external hand-off with a bounded timeout."""
current = _as_utc(now or utcnow())
event_type = _validated_event_type(node.config.get("event_type"))
event_filter = _resolved_context_value(
_event_filter(node.config.get("event_filter")),
instance.context_,
)
if not isinstance(event_filter, Mapping) or not event_filter:
raise WorkflowConflictError(
"External hand-offs require a non-empty event filter."
)
terminal = node.config.get("terminal_outcomes")
if not isinstance(terminal, Mapping) or not terminal:
raise WorkflowConflictError(
"External hand-offs require terminal outcome mappings."
)
terminal_outcomes = {
str(key).strip(): str(value).strip()
for key, value in terminal.items()
if str(key).strip() and str(value).strip()
}
allowed_ports = {"completed", "rejected", "cancelled"}
if (
not terminal_outcomes
or any(port not in allowed_ports for port in terminal_outcomes.values())
):
raise WorkflowConflictError(
"External hand-off outcomes must map to completed, rejected, or cancelled."
)
observed = tuple(
dict.fromkeys(
str(value).strip()
for value in node.config.get("observed_outcomes") or ()
if str(value).strip()
)
)
overlap = set(observed) & set(terminal_outcomes)
if overlap:
raise WorkflowConflictError(
"External hand-off outcomes cannot be both observed and terminal: "
+ ", ".join(sorted(overlap))
)
timeout_after = _resolved_context_value(
node.config.get("timeout_after"),
instance.context_,
)
due_at = (
current + timedelta(seconds=_duration_seconds(timeout_after, minimum=1))
if str(timeout_after or "").strip()
else None
)
provider_capability = str(
_resolved_context_value(
node.config.get("provider_capability"),
instance.context_,
)
or ""
).strip()
external_id = str(
_resolved_context_value(
node.config.get("external_id"),
instance.context_,
)
or ""
).strip()
if not provider_capability or not external_id:
raise WorkflowConflictError(
"External hand-offs require a provider capability and external reference."
)
expected_revision_value = _resolved_context_value(
node.config.get("expected_revision"),
instance.context_,
)
expected_revision = (
int(expected_revision_value)
if expected_revision_value not in (None, "")
else None
)
if expected_revision is not None and expected_revision < 1:
raise WorkflowConflictError(
"External hand-off revisions start at one."
)
action_url = _safe_action_url(
_resolved_context_value(
node.config.get("action_url"),
instance.context_,
)
)
immutable_ref = str(
_resolved_context_value(
node.config.get("immutable_ref"),
instance.context_,
)
or ""
).strip()[:1_000]
if not immutable_ref:
raise WorkflowConflictError(
"External hand-offs require an immutable reference."
)
optional_capabilities = _resolved_context_value(
node.config.get("optional_capabilities") or {},
instance.context_,
)
if not isinstance(optional_capabilities, Mapping):
raise WorkflowConflictError(
"External hand-off capability availability must be an object."
)
state = WorkflowWaitState(
tenant_id=instance.tenant_id,
instance_id=instance.id,
step_id=step.id,
mode="external_handoff",
status="waiting",
due_at=due_at,
event_type=event_type,
config_={
"filter": dict(event_filter),
"outcome_path": str(
node.config.get("outcome_path") or "payload.outcome"
).strip(),
"terminal_outcomes": terminal_outcomes,
"observed_outcomes": list(observed),
"observed_event_ids": [],
"provider_capability": provider_capability,
"external_id": external_id,
"expected_revision": expected_revision,
"action_url": action_url,
"immutable_ref": immutable_ref,
"optional_capabilities": dict(optional_capabilities),
},
)
session.add(state)
session.flush()
return state
def resolve_wait_state(
session: Session,
*,
@@ -328,7 +473,7 @@ def ingest_platform_event(
.where(
WorkflowWaitState.tenant_id == tenant_id,
WorkflowWaitState.status == "waiting",
WorkflowWaitState.mode == "event",
WorkflowWaitState.mode.in_(("event", "external_handoff")),
WorkflowWaitState.event_type == event.type,
)
.with_for_update(skip_locked=True)
@@ -338,6 +483,71 @@ def ingest_platform_event(
for state in waits:
if not _matches_filter(envelope, state.config_.get("filter")):
continue
if state.mode == "external_handoff":
observed_ids = [
str(value)
for value in state.config_.get("observed_event_ids") or ()
]
if event.event_id in observed_ids:
continue
outcome = str(
_value_at_path(
envelope,
state.config_.get("outcome_path") or "payload.outcome",
)
or ""
).strip()
terminal = state.config_.get("terminal_outcomes")
terminal_port = (
str(terminal.get(outcome) or "").strip()
if isinstance(terminal, Mapping)
else ""
)
observed_outcomes = {
str(value).strip()
for value in state.config_.get("observed_outcomes") or ()
if str(value).strip()
}
if not terminal_port and outcome not in observed_outcomes:
continue
state.config_ = {
**dict(state.config_),
"observed_event_ids": [*observed_ids[-99:], event.event_id],
"last_outcome": outcome,
**(
{"selected_port": terminal_port}
if terminal_port
else {}
),
}
if not terminal_port:
step = session.get(WorkflowInstanceStep, state.step_id)
instance = session.get(WorkflowInstance, state.instance_id)
if step is not None and instance is not None:
step.handoff = {
**dict(step.handoff),
"state": outcome,
"last_event_id": event.event_id,
}
from govoplan_workflow_engine.backend.instance_service import (
_record_event,
)
_record_event(
session,
instance,
step=step,
kind="workflow.external_handoff.observed",
actor_id=event.actor.id if event.actor else None,
payload={
"event_id": event.event_id,
"event_type": event.type,
"outcome": outcome,
"external_id": state.config_.get("external_id"),
},
)
state.revision += 1
continue
state.status = "triggered"
state.source_event_id = event.event_id
state.event_ = envelope
@@ -700,11 +910,45 @@ def _dispatch_waits(
continue
graph = _runtime_graph(revision)
triggered = state.status == "triggered"
external_inspection: dict[str, object] | None = None
if triggered and state.mode == "external_handoff":
external_inspection = _inspect_external_handoff(
session,
state=state,
step=step,
principal=principal,
registry=registry,
)
if external_inspection.get("allowed") is not True:
reason = str(
external_inspection.get("reason")
or "External hand-off authorization is unavailable."
)
state.error = reason
step.handoff = {
**dict(step.handoff),
"state": "blocked",
"message": reason,
"inspection": external_inspection,
}
skipped += 1
release_workflow_state_fence(session, fence)
continue
state.status = "resumed" if triggered else "timed_out"
state.error = None
state.revision += 1
output = (
{"event": dict(state.event_ or {})}
{
"event": dict(state.event_ or {}),
**(
{
"external_handoff": external_inspection,
"immutable_ref": state.config_.get("immutable_ref"),
}
if external_inspection is not None
else {}
),
}
if triggered
else {"due_at": state.due_at.isoformat() if state.due_at else None}
)
@@ -713,8 +957,12 @@ def _dispatch_waits(
instance,
step=step,
kind=(
"workflow.wait.event_received"
"workflow.external_handoff.completed"
if triggered and state.mode == "external_handoff"
else "workflow.wait.event_received"
if triggered
else "workflow.external_handoff.timed_out"
if state.mode == "external_handoff"
else "workflow.wait.timed_out"
),
actor_id=None,
@@ -725,7 +973,13 @@ def _dispatch_waits(
instance=instance,
step=step,
graph=graph,
port="resumed" if triggered else "timed_out",
port=(
str(state.config_.get("selected_port") or "completed")
if triggered and state.mode == "external_handoff"
else "resumed"
if triggered
else "timed_out"
),
output=output,
actor_id=None,
)
@@ -751,6 +1005,92 @@ def _dispatch_waits(
}
def _inspect_external_handoff(
session: Session,
*,
state: WorkflowWaitState,
step: WorkflowInstanceStep,
principal: ApiPrincipal,
registry: object | None,
) -> dict[str, object]:
from govoplan_core.core.campaigns import CampaignWorkOrchestrationProvider
capability_name = str(
state.config_.get("provider_capability") or ""
).strip()
if (
not capability_name
or registry is None
or not hasattr(registry, "has_capability")
or not registry.has_capability(capability_name)
):
return {
"allowed": False,
"reason": (
f"External hand-off capability {capability_name!r} is unavailable."
),
"code": "workflow_external_handoff_provider_unavailable",
}
provider = registry.capability(capability_name)
if not isinstance(provider, CampaignWorkOrchestrationProvider):
return {
"allowed": False,
"reason": "External hand-off provider has an incompatible contract.",
"code": "workflow_external_handoff_provider_invalid",
}
event_revision = _value_at_path(
state.event_ or {},
"payload.assignment_revision",
)
try:
expected_revision = (
int(event_revision)
if event_revision not in (None, "")
else None
)
except (TypeError, ValueError):
return {
"allowed": False,
"reason": "External hand-off event has no valid resource revision.",
"code": "workflow_external_handoff_event_revision_invalid",
}
inspection = provider.inspect_handoff(
session,
principal,
tenant_id=state.tenant_id,
assignment_id=str(state.config_.get("external_id") or ""),
expected_revision=expected_revision,
)
payload: dict[str, object] = {
"allowed": inspection.allowed,
"status": inspection.status,
"assignment_revision": inspection.assignment_revision,
"action_url": inspection.action_url,
"assignment_ref": inspection.assignment_ref,
"reason": inspection.reason,
"provenance": dict(inspection.provenance),
}
selected_port = str(state.config_.get("selected_port") or "")
if inspection.allowed and inspection.status != selected_port:
payload.update(
{
"allowed": False,
"reason": (
"External hand-off state does not match the terminal event; "
"reload and reconcile the provider."
),
"code": "workflow_external_handoff_state_mismatch",
}
)
if inspection.allowed and inspection.assignment_ref:
state.config_ = {
**dict(state.config_),
"immutable_ref": inspection.assignment_ref,
}
step.external_ref = inspection.assignment_ref
return payload
def _resolve_trigger_principal(
session: Session,
*,
@@ -938,6 +1278,12 @@ def _duration_seconds(value: object, *, minimum: int) -> int:
return seconds
def duration_seconds(value: object, *, minimum: int = 1) -> int:
"""Parse the duration syntax shared by timers and human-work due dates."""
return _duration_seconds(value, minimum=minimum)
def _parse_instant(value: object, *, timezone_name: str) -> datetime:
text = str(value or "").strip()
if not text:
@@ -972,6 +1318,70 @@ def _event_filter(value: object) -> dict[str, object]:
return parsed
def _resolved_context_value(
value: object,
context: Mapping[str, object],
*,
depth: int = 0,
) -> object:
if depth > 10:
raise WorkflowConflictError(
"External hand-off configuration is nested too deeply."
)
if isinstance(value, str) and value.startswith("$"):
path = value[1:].lstrip(".")
current: object = context
if not path:
return dict(context)
for segment in path.split("."):
if not isinstance(current, Mapping) or segment not in current:
raise WorkflowConflictError(
f"External hand-off input path {value!r} is unavailable."
)
current = current[segment]
return current
if isinstance(value, Mapping):
return {
str(key): _resolved_context_value(
item,
context,
depth=depth + 1,
)
for key, item in value.items()
}
if isinstance(value, list):
return [
_resolved_context_value(item, context, depth=depth + 1)
for item in value
]
return value
def _value_at_path(value: object, path: object) -> object | None:
current = value
for segment in str(path or "").strip().lstrip("$").lstrip(".").split("."):
if not segment:
continue
if not isinstance(current, Mapping) or segment not in current:
return None
current = current[segment]
return current
def _safe_action_url(value: object) -> str:
candidate = str(value or "").strip()
if (
not candidate.startswith("/")
or candidate.startswith("//")
or "\\" in candidate
or any(ord(character) < 32 or ord(character) == 127 for character in candidate)
):
raise WorkflowConflictError(
"External hand-off action URLs must be safe application-relative paths."
)
return candidate[:1_500]
def _validated_event_type(value: object) -> str:
event_type = str(value or "").strip()
if not _EVENT_TYPE.fullmatch(event_type):
@@ -0,0 +1,367 @@
from __future__ import annotations
from collections.abc import Mapping
from datetime import UTC, datetime
from sqlalchemy import and_, or_, select
from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, has_scope
from govoplan_core.core.idm import CAPABILITY_IDM_DIRECTORY, IdmDirectory
from govoplan_core.core.tasks import (
WorkAssignmentRef,
WorkItem,
WorkItemPage,
WorkItemQuery,
WorkSourceRef,
)
from govoplan_workflow_engine.backend.db.models import (
WorkflowDefinition,
WorkflowInstance,
WorkflowInstanceStep,
)
from govoplan_workflow_engine.backend.governance import definition_decision
PROVIDER_ID = "workflow_engine.handoffs"
INSTANCE_READ_SCOPE = "workflow:instance:read"
ADMIN_SCOPE = "workflow:instance:admin"
_NON_HUMAN_KINDS = {"timer", "event_wait", "dataflow_run"}
_BLOCKED_STATES = {
"blocked",
"failed",
"outcome_unknown",
"recovery_required",
"compensation_required",
}
_ASSIGNMENT_KINDS = {
"account",
"group",
"role",
"function",
"function_assignment",
"anyone",
}
class WorkflowWorkItemProvider:
def __init__(self, *, registry: object | None = None) -> None:
self.registry = registry
def list_items(
self,
session: object,
principal: object,
*,
query: WorkItemQuery,
) -> WorkItemPage:
if not isinstance(session, Session):
raise TypeError("Workflow work aggregation requires a SQLAlchemy Session.")
if not isinstance(principal, ApiPrincipal):
return WorkItemPage(items=(), total=0)
if principal.tenant_id != query.tenant_id:
return WorkItemPage(items=(), total=0)
administrative = has_scope(principal, ADMIN_SCOPE)
if not administrative and not has_scope(principal, INSTANCE_READ_SCOPE):
return WorkItemPage(items=(), total=0)
statement = (
select(WorkflowInstance, WorkflowInstanceStep, WorkflowDefinition)
.join(
WorkflowInstanceStep,
WorkflowInstanceStep.id == WorkflowInstance.current_step_id,
)
.join(
WorkflowDefinition,
WorkflowDefinition.id == WorkflowInstance.definition_id,
)
.where(
WorkflowInstance.tenant_id == query.tenant_id,
WorkflowInstance.status == "waiting",
WorkflowInstanceStep.status == "waiting",
)
.order_by(
WorkflowInstanceStep.work_due_at.is_(None),
WorkflowInstanceStep.work_due_at.asc(),
WorkflowInstanceStep.updated_at.desc(),
WorkflowInstanceStep.id.desc(),
)
)
if not administrative:
targets = self._targets(principal, query.tenant_id)
conditions = [
and_(
WorkflowInstanceStep.work_assignment_kind == kind,
WorkflowInstanceStep.work_assignment_id.in_(tuple(values)),
)
for kind, values in targets.items()
if values
]
conditions.append(
and_(
WorkflowInstanceStep.work_assignment_kind.is_(None),
WorkflowInstanceStep.work_assignment_id.is_(None),
)
)
statement = statement.where(or_(*conditions))
items: list[WorkItem] = []
total = 0
decisions: dict[str, bool] = {}
targets = self._targets(principal, query.tenant_id)
for instance, step, definition in session.execute(statement).yield_per(250):
assignment = _step_assignment(instance, step)
if not administrative and not _assignment_matches(assignment, targets):
continue
if not _is_actionable_handoff(step.handoff):
continue
allowed = decisions.get(definition.id)
if allowed is None:
allowed = definition_decision(
definition,
principal=principal,
registry=self.registry,
action="view",
).allowed
decisions[definition.id] = allowed
if not allowed:
continue
item = _work_item(instance, step, definition, assignment)
if query.statuses and item.status not in query.statuses:
continue
if query.priorities and item.priority not in query.priorities:
continue
if query.due_before is not None and (
item.due_at is None or _aware(item.due_at) > _aware(query.due_before)
):
continue
if query.text and query.text.casefold() not in _search_text(item):
continue
total += 1
if len(items) < query.limit:
items.append(item)
return WorkItemPage(
items=tuple(items),
total=total,
truncated=total > len(items),
)
def _targets(self, principal: ApiPrincipal, tenant_id: str) -> dict[str, set[str]]:
result = {
"account": {principal.account_id} if principal.account_id else set(),
"group": set(principal.group_ids),
"role": set(principal.role_ids),
"function_assignment": set(principal.function_assignment_ids),
"function": set(),
"anyone": {"*"},
}
directory = self._idm_directory()
if directory is not None and principal.account_id:
result["function"].update(
item.function_id
for item in directory.organization_function_assignments_for_account(
principal.account_id,
tenant_id=tenant_id,
)
if item.tenant_id == tenant_id and item.status == "active"
)
return result
def _idm_directory(self) -> IdmDirectory | None:
registry = self.registry
if (
registry is None
or not hasattr(registry, "has_capability")
or not registry.has_capability(CAPABILITY_IDM_DIRECTORY)
):
return None
provider = registry.capability(CAPABILITY_IDM_DIRECTORY)
return provider if isinstance(provider, IdmDirectory) else None
def _step_assignment(
instance: WorkflowInstance,
step: WorkflowInstanceStep,
) -> WorkAssignmentRef | None:
if step.work_assignment_kind and step.work_assignment_id:
return _assignment_ref(
step.work_assignment_kind,
step.work_assignment_id,
step.work_assignment_label,
)
handoff_assignment = step.handoff.get("assignment")
if isinstance(handoff_assignment, Mapping):
kind = str(handoff_assignment.get("kind") or "").strip()
assignment_id = str(handoff_assignment.get("id") or "").strip()
if kind and assignment_id:
return _assignment_ref(
kind,
assignment_id,
str(handoff_assignment.get("label") or "").strip() or None,
)
account_id = str(instance.authorization_.get("account_id") or "").strip()
return WorkAssignmentRef(kind="account", id=account_id) if account_id else None
def _assignment_matches(
assignment: WorkAssignmentRef | None,
targets: Mapping[str, set[str]],
) -> bool:
return bool(
assignment is not None and assignment.id in targets.get(assignment.kind, set())
)
def _is_actionable_handoff(handoff: Mapping[str, object]) -> bool:
kind = str(handoff.get("kind") or "")
state = str(handoff.get("state") or "waiting")
return kind not in _NON_HUMAN_KINDS and state not in {"pending", "running"}
def _work_item(
instance: WorkflowInstance,
step: WorkflowInstanceStep,
definition: WorkflowDefinition,
assignment: WorkAssignmentRef | None,
) -> WorkItem:
handoff = dict(step.handoff or {})
state = str(handoff.get("state") or "waiting")
status = "blocked" if state in _BLOCKED_STATES else "open"
title = str(
handoff.get("title") or handoff.get("message") or f"Continue {definition.name}"
)
required_action = str(
handoff.get("instructions")
or handoff.get("message")
or "Continue the current workflow handoff."
).strip()
action_url = _action_url(
handoff.get("action_url")
or f"/workflow?definition={definition.id}&run={instance.id}"
)
due_at = step.work_due_at or _date(handoff.get("due_at"))
priority = str(handoff.get("priority") or "normal").casefold()
if priority not in {"low", "normal", "high", "urgent"}:
priority = "normal"
updated_at = step.updated_at or instance.updated_at
revision = f"{step.attempt}:{updated_at.isoformat() if updated_at else '1'}"
return WorkItem(
id=step.id,
provider_id=PROVIDER_ID,
owner_module="workflow_engine",
tenant_id=instance.tenant_id,
title=title,
summary=f"{definition.name} · {step.node_type}",
status=status, # type: ignore[arg-type]
priority=priority, # type: ignore[arg-type]
required_action=required_action or None,
action_url=action_url,
due_at=due_at,
assignments=(assignment,) if assignment else (),
sources=(
WorkSourceRef(
module_id="workflow_engine",
resource_type="workflow_instance",
resource_id=instance.id,
revision=instance.definition_revision_id,
url=f"/workflow?definition={definition.id}&run={instance.id}",
label=definition.name,
),
WorkSourceRef(
module_id="workflow_engine",
resource_type="workflow_step",
resource_id=step.id,
revision=str(step.attempt),
),
),
provenance={
"definition_id": definition.id,
"definition_revision_id": instance.definition_revision_id,
"workflow_instance_id": instance.id,
"workflow_step_id": step.id,
},
metadata={
"handoff_kind": handoff.get("kind"),
"handoff_state": state,
"allowed_actions": _allowed_actions(handoff.get("allowed_actions")),
},
revision=revision,
created_at=step.created_at,
updated_at=updated_at,
)
def _assignment_ref(
kind: object,
assignment_id: object,
label: object = None,
) -> WorkAssignmentRef | None:
normalized_kind = str(kind or "").strip()
normalized_id = str(assignment_id or "").strip()
if normalized_kind not in _ASSIGNMENT_KINDS or not normalized_id:
return None
if normalized_kind == "anyone":
normalized_id = "*"
try:
normalized_label = str(label).strip()[:500] if label is not None else ""
return WorkAssignmentRef(
kind=normalized_kind, # type: ignore[arg-type]
id=normalized_id,
label=normalized_label or None,
)
except ValueError:
return None
def _action_url(value: object) -> str:
candidate = str(value or "").strip()
if (
candidate.startswith("/")
and not candidate.startswith("//")
and "\\" not in candidate
and all(
ord(character) >= 32 and ord(character) != 127 for character in candidate
)
):
return candidate[:1_500]
return "/workflow"
def _allowed_actions(value: object) -> list[str]:
if not isinstance(value, (list, tuple, set, frozenset)):
return []
return [normalized for item in value if (normalized := str(item or "").strip())][
:100
]
def _date(value: object) -> datetime | None:
if isinstance(value, datetime):
return value
text = str(value or "").strip()
if not text:
return None
try:
return datetime.fromisoformat(text.replace("Z", "+00:00"))
except ValueError:
return None
def _aware(value: datetime) -> datetime:
return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
def _search_text(item: WorkItem) -> str:
return " ".join(
value
for value in (
item.title,
item.summary,
item.required_action,
item.owner_module,
)
if value
).casefold()
__all__ = ["PROVIDER_ID", "WorkflowWorkItemProvider"]
+578
View File
@@ -0,0 +1,578 @@
from __future__ import annotations
import json
import unittest
from datetime import UTC, datetime, timedelta
from sqlalchemy import create_engine
from sqlalchemy.orm import Session
from govoplan_core.core.dsar import (
DsarErasureActionRef,
DsarProvider,
DsarRecordRef,
DsarSubjectRef,
)
from govoplan_core.db.base import Base
from govoplan_core.privacy.dsar_workflow import (
create_data_subject_request,
search_data_subject_request,
)
from govoplan_workflow_engine.backend.db.models import (
WorkflowDefinition,
WorkflowDefinitionRevision,
WorkflowInstance,
WorkflowInstanceEvent,
WorkflowInstanceStep,
WorkflowTrigger,
WorkflowTriggerDelivery,
WorkflowWaitState,
)
from govoplan_workflow_engine.backend.dsar_provider import (
WORKFLOW_ENGINE_DSAR_CAPABILITY,
WorkflowEngineDsarProvider,
)
from govoplan_workflow_engine.backend.manifest import manifest
NOW = datetime(2026, 8, 21, 21, 0, tzinfo=UTC)
SECRET = "private-workflow-detail-do-not-export"
class _Registry:
def __init__(
self,
provider: WorkflowEngineDsarProvider,
*,
active: bool = True,
) -> None:
self.provider = provider
self.active = active
def capability_names(self):
return (WORKFLOW_ENGINE_DSAR_CAPABILITY,)
def capability_owner(self, name):
self._assert_capability(name)
return "workflow_engine"
def tenant_entitlement_resolver(self):
active = self.active
class _Resolver:
@staticmethod
def resolve(session, tenant_id):
del session, tenant_id
return type(
"State",
(),
{"effective_modules": ("workflow_engine",) if active else ()},
)()
return _Resolver()
def require_tenant_capability(self, name, session, **kwargs):
del session, kwargs
self._assert_capability(name)
return self.provider
def manifests(self):
return (type("Manifest", (), {"id": "workflow_engine"})(),)
@staticmethod
def _assert_capability(name: str) -> None:
if name != WORKFLOW_ENGINE_DSAR_CAPABILITY:
raise KeyError(name)
class WorkflowEngineDsarProviderTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite+pysqlite:///:memory:")
Base.metadata.create_all(self.engine)
self.session = Session(self.engine)
self.provider = WorkflowEngineDsarProvider()
self.assertIsInstance(self.provider, DsarProvider)
self._seed()
self.session.commit()
def tearDown(self) -> None:
self.session.close()
self.engine.dispose()
def _seed(self) -> None:
definition = WorkflowDefinition(
id="definition-1",
tenant_id="tenant-1",
scope_type="tenant",
scope_id="tenant-1",
scope_key="tenant:tenant-1",
definition_kind="flow",
definition_key="resident-permit",
name=SECRET,
description=SECRET,
status="active",
current_revision=1,
active_revision=1,
metadata_={"secret": SECRET},
derivation_provenance={"secret": SECRET},
created_by="account-1",
updated_by="account-1",
)
other = WorkflowDefinition(
id="definition-other",
tenant_id="tenant-2",
scope_type="tenant",
scope_id="tenant-2",
scope_key="tenant:tenant-2",
definition_kind="flow",
definition_key="other",
name=SECRET,
status="active",
current_revision=1,
created_by="account-1",
updated_by="account-1",
)
self.session.add_all((definition, other))
self.session.flush()
revision = WorkflowDefinitionRevision(
id="revision-1",
tenant_id="tenant-1",
definition_id=definition.id,
revision=1,
schema_version=1,
graph={"secret": SECRET},
content_hash="a" * 64,
library_id="bpmn",
library_version="1.0.0",
execution_mode="hybrid",
bpmn_xml=SECRET,
bpmn_hash="b" * 64,
bpmn_runtime_kind="native_graph",
bpmn_executable=True,
contribution_metadata={"secret": SECRET},
created_by="account-1",
)
self.session.add(revision)
self.session.flush()
instance = WorkflowInstance(
id="instance-1",
tenant_id="tenant-1",
definition_id=definition.id,
definition_revision_id=revision.id,
status="completed",
start_origin="service",
idempotency_key=SECRET,
correlation_id=SECRET,
input_={"secret": SECRET},
context_={"secret": SECRET},
output_={"secret": SECRET},
authorization_={
"account_id": "account-1",
"identity_id": "identity-1",
"membership_id": "membership-1",
"secret": SECRET,
},
started_at=NOW,
finished_at=NOW,
error=SECRET,
created_by="account-1",
)
self.session.add(instance)
self.session.flush()
step = WorkflowInstanceStep(
id="step-1",
tenant_id="tenant-1",
instance_id=instance.id,
sequence=1,
node_id="review",
node_type="human_task",
status="completed",
attempt=1,
idempotency_key=SECRET,
input_={"secret": SECRET},
output_={"secret": SECRET},
handoff={"secret": SECRET},
external_ref=SECRET,
started_at=NOW,
finished_at=NOW,
error=SECRET,
completed_by="account-1",
work_assignment_kind="membership",
work_assignment_id="membership-1",
work_assignment_label=SECRET,
work_due_at=NOW,
)
trigger = WorkflowTrigger(
id="trigger-1",
tenant_id="tenant-1",
definition_id=definition.id,
definition_revision_id=revision.id,
node_id="start-event",
kind="event",
status="active",
config_={"secret": SECRET},
event_type="permit.submitted",
next_fire_at=NOW + timedelta(hours=1),
last_error=SECRET,
authorization_subject_kind="delegated_user",
authorization_account_id="account-1",
authorization_membership_id="membership-1",
authorization_ref=SECRET,
grant_scopes=[SECRET],
created_by="account-1",
updated_by="account-1",
)
self.session.add_all((step, trigger))
self.session.flush()
self.session.add_all(
(
WorkflowInstanceEvent(
id="event-1",
tenant_id="tenant-1",
instance_id=instance.id,
step_id=step.id,
sequence=1,
kind="step.completed",
actor_id="account-1",
payload={"secret": SECRET},
created_at=NOW,
),
WorkflowTriggerDelivery(
id="delivery-1",
tenant_id="tenant-1",
trigger_id=trigger.id,
definition_id=definition.id,
definition_revision_id=revision.id,
source_key=SECRET,
invocation_kind="event",
status="succeeded",
scheduled_for=NOW,
event_={"secret": SECRET},
instance_id=instance.id,
attempts=1,
error=SECRET,
),
WorkflowWaitState(
id="wait-1",
tenant_id="tenant-1",
instance_id=instance.id,
step_id=step.id,
mode="event",
status="completed",
due_at=NOW,
event_type="permit.approved",
config_={"secret": SECRET},
source_event_id=SECRET,
event_={"secret": SECRET},
error=SECRET,
revision=1,
),
)
)
def test_canonical_selector_uses_structured_links_and_minimized_output(
self,
) -> None:
records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
account_id="account-1",
identity_id="identity-1",
membership_id="membership-1",
),
)
self.assertEqual(6, len(records))
categories = {record.resource_type: record.category for record in records}
self.assertEqual(
"workflow_automation_authority",
categories["workflow_trigger"],
)
self.assertEqual(
"workflow_subject_authorization",
categories["workflow_instance"],
)
self.assertEqual(
"workflow_subject_work_assignment",
categories["workflow_instance_step"],
)
exported = json.dumps([record.to_dict() for record in records])
self.assertNotIn(SECRET, exported)
self.assertNotIn("account-1", exported)
self.assertNotIn("identity-1", exported)
self.assertNotIn("membership-1", exported)
self.assertNotIn("definition-other", exported)
def test_exact_instance_and_definition_packages_are_review_only(self) -> None:
instance_records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
external_references={"workflow.instance": "instance:instance-1"}
),
)
definition_records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
external_references={"workflow.definition": "definition-1"}
),
)
self.assertEqual(4, len(instance_records))
self.assertEqual(8, len(definition_records))
self.assertNotIn(
SECRET,
json.dumps([record.to_dict() for record in definition_records]),
)
actions = self.provider.plan_erasure(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
external_references={"workflow.definition": "definition-1"}
),
records=definition_records,
)
self.assertEqual({"manual_review"}, {action.kind for action in actions})
def test_every_exact_reference_and_conflicts_fail_closed(self) -> None:
references = {
"workflow_engine.definition_revision": "revision-1",
"workflow_engine.step": "step-1",
"workflow_engine.event": "event-1",
"workflow_engine.trigger": "trigger-1",
"workflow_engine.trigger_delivery": "delivery-1",
"workflow_engine.wait_state": "wait-1",
}
for key, value in references.items():
with self.subTest(key=key):
records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(external_references={key: value}),
)
self.assertEqual(1, len(records))
mismatch = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
account_id="account-2",
external_references={"workflow.instance": "instance-1"},
),
)
wrong_tenant = self.provider.search_subject(
self.session,
tenant_id="tenant-2",
subject=DsarSubjectRef(
external_references={"workflow.instance": "instance-1"}
),
)
conflict = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=DsarSubjectRef(
external_references={
"workflow_engine.revision": "revision-1",
"workflow_engine.definition_revision": "different",
}
),
)
self.assertEqual((), mismatch)
self.assertEqual((), wrong_tenant)
self.assertEqual((), conflict)
def test_terminal_delivery_minimization_preserves_replay_identity(self) -> None:
subject = DsarSubjectRef(
external_references={"workflow_engine.delivery": "delivery-1"}
)
records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=subject,
)
actions = self.provider.plan_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
records=records,
)
self.assertEqual("anonymize", actions[0].kind)
first = self.provider.execute_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
actions=actions,
request_id="dsar-1",
)
second = self.provider.execute_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
actions=actions,
request_id="dsar-1-retry",
)
self.assertEqual("executed", first[0].status)
self.assertEqual("unchanged", second[0].status)
delivery = self.session.get(WorkflowTriggerDelivery, "delivery-1")
self.assertIsNone(delivery.event_)
self.assertIsNone(delivery.error)
self.assertEqual(SECRET, delivery.source_key)
def test_automation_authority_revocation_is_idempotent(self) -> None:
subject = DsarSubjectRef(
account_id="account-1",
membership_id="membership-1",
)
records = self.provider.search_subject(
self.session,
tenant_id="tenant-1",
subject=subject,
)
actions = self.provider.plan_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
records=records,
)
trigger_action = next(
action for action in actions if action.resource_type == "workflow_trigger"
)
self.assertEqual("revoke", trigger_action.kind)
first = self.provider.execute_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
actions=(trigger_action,),
request_id="dsar-2",
)
second = self.provider.execute_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
actions=(trigger_action,),
request_id="dsar-2-retry",
)
self.assertEqual("executed", first[0].status)
self.assertEqual("unchanged", second[0].status)
trigger = self.session.get(WorkflowTrigger, "trigger-1")
self.assertEqual("disabled", trigger.status)
self.assertIsNone(trigger.authorization_account_id)
self.assertIsNone(trigger.authorization_membership_id)
self.assertEqual([], trigger.grant_scopes)
self.assertEqual({}, trigger.config_)
def test_foreign_records_and_actions_are_rejected(self) -> None:
subject = DsarSubjectRef(account_id="account-1")
with self.assertRaisesRegex(ValueError, "foreign provider record"):
self.provider.plan_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
records=(
DsarRecordRef(
provider_id="cases",
module_id="cases",
resource_type="case",
resource_id="case-1",
category="case",
title="Case",
),
),
)
with self.assertRaisesRegex(ValueError, "foreign provider action"):
self.provider.execute_erasure(
self.session,
tenant_id="tenant-1",
subject=subject,
actions=(
DsarErasureActionRef(
action_id="cases:delete:case:case-1",
provider_id="cases",
module_id="cases",
kind="delete",
resource_type="case",
resource_id="case-1",
title="Delete case",
rationale="Foreign",
executable=True,
),
),
request_id="dsar-3",
)
def test_core_workflow_reports_active_and_inactive_provider(self) -> None:
row = create_data_subject_request(
self.session,
tenant_id="tenant-1",
reference="DSAR-WORKFLOW-1",
request_kind="access_and_erasure",
subject=DsarSubjectRef(account_id="account-1"),
purpose="Respond to a verified request.",
legal_basis="Article 15 and 17 GDPR",
due_at=None,
requested_by_account_id="privacy-officer",
)
self.session.commit()
search_data_subject_request(
self.session,
registry=_Registry(self.provider),
row=row,
expected_revision=1,
)
self.assertEqual(
[WORKFLOW_ENGINE_DSAR_CAPABILITY],
row.coverage["provider_capabilities"],
)
self.assertEqual(6, row.search_result["record_count"])
inactive = create_data_subject_request(
self.session,
tenant_id="tenant-1",
reference="DSAR-WORKFLOW-2",
request_kind="access",
subject=DsarSubjectRef(account_id="account-1"),
purpose="Respond to a verified request.",
legal_basis="Article 15 GDPR",
due_at=None,
requested_by_account_id="privacy-officer",
)
self.session.commit()
search_data_subject_request(
self.session,
registry=_Registry(self.provider, active=False),
row=inactive,
expected_revision=1,
)
self.assertEqual([], inactive.coverage["provider_capabilities"])
self.assertEqual(
[WORKFLOW_ENGINE_DSAR_CAPABILITY],
inactive.coverage["inactive_provider_capabilities"],
)
self.assertEqual(0, inactive.search_result["record_count"])
def test_manifest_registers_and_documents_capability(self) -> None:
self.assertIn(
WORKFLOW_ENGINE_DSAR_CAPABILITY,
manifest.capability_factories,
)
self.assertIn(
WORKFLOW_ENGINE_DSAR_CAPABILITY,
manifest.capability_documentation,
)
self.assertIn(
WORKFLOW_ENGINE_DSAR_CAPABILITY,
{item.name for item in manifest.provides_interfaces},
)
self.assertTrue(
any(
topic.id == "workflow.data-subject-requests"
and {"admin", "user"}.issubset(topic.documentation_types)
for topic in manifest.documentation
)
)
if __name__ == "__main__":
unittest.main()
+104 -3
View File
@@ -43,6 +43,7 @@ from govoplan_core.core.runtime_coordination import (
RuntimeIdentity,
bind_process_runtime_identity,
)
from govoplan_core.core.tasks import WorkItemQuery
from govoplan_core.db.base import Base
from govoplan_core.db.base import utcnow
from govoplan_workflow_engine.backend.db.models import (
@@ -84,6 +85,7 @@ from govoplan_workflow_engine.backend.service import (
create_definition,
)
from govoplan_workflow_engine.backend.service_launcher import WorkflowServiceLauncher
from govoplan_workflow_engine.backend.work_items import WorkflowWorkItemProvider
try:
from test_bpmn import NATIVE_BPMN
@@ -687,6 +689,107 @@ class WorkflowInstanceServiceTests(unittest.TestCase):
start_origin="api",
)
def test_human_handoff_projects_typed_due_work_and_disappears_on_completion(
self,
) -> None:
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Guided case review",
graph=WorkflowGraph(
nodes=[
WorkflowNode(
id="start",
type="workflow.start.manual",
config={"input_schema_ref": ""},
),
WorkflowNode(
id="activity",
type="workflow.activity",
config={
"title": "Assess the application",
"instructions": "Record the assessment evidence.",
"assignee": "account:account-1",
"due_after": "2h",
},
),
WorkflowNode(
id="done",
type="workflow.end.completed",
),
],
edges=[
WorkflowEdge(
id="start-activity", source="start", target="activity"
),
WorkflowEdge(
id="activity-done", source="activity", target="done"
),
],
),
execution_mode="guided",
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
before = utcnow()
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(idempotency_key="guided-work-1"),
)
step = self.session.get(WorkflowInstanceStep, instance.current_step_id)
self.assertIsNotNone(step)
assert step is not None
self.assertEqual("account", step.work_assignment_kind)
self.assertEqual("account-1", step.work_assignment_id)
self.assertIsNotNone(step.work_due_at)
assert step.work_due_at is not None
due_at = (
step.work_due_at.replace(tzinfo=UTC)
if step.work_due_at.tzinfo is None
else step.work_due_at
)
self.assertGreaterEqual(due_at, before + timedelta(hours=1, minutes=59))
provider = WorkflowWorkItemProvider(registry=self.registry)
page = provider.list_items(
self.session,
principal(),
query=WorkItemQuery(tenant_id="tenant-1"),
)
self.assertEqual(1, page.total)
self.assertEqual("Assess the application", page.items[0].title)
self.assertEqual("account-1", page.items[0].assignments[0].id)
self.assertEqual(step.id, page.items[0].id)
resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=step.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowStepActionRequest(action="complete"),
)
closed_page = provider.list_items(
self.session,
principal(),
query=WorkItemQuery(tenant_id="tenant-1"),
)
self.assertEqual(0, closed_page.total)
def test_module_action_records_effects_and_completes_idempotently(
self,
) -> None:
@@ -990,9 +1093,7 @@ class WorkflowInstanceServiceTests(unittest.TestCase):
request=request,
)
checkpoint = session.scalar(
select(RecoveryCheckpoint).order_by(
RecoveryCheckpoint.sequence
)
select(RecoveryCheckpoint).order_by(RecoveryCheckpoint.sequence)
)
assert checkpoint is not None
checkpoint.summary = "tampered provider evidence"
+43 -5
View File
@@ -29,7 +29,7 @@ class WorkflowMigrationTests(unittest.TestCase):
try:
with engine.connect() as connection:
self.assertIn(
"e4a1f8c2d7b6",
"8d5a2f7c1b4e",
set(MigrationContext.configure(connection).get_current_heads()),
)
self.assertEqual(
@@ -102,7 +102,12 @@ class WorkflowMigrationTests(unittest.TestCase):
)
for path in current_revisions.glob("*.py"):
if path.name.startswith(
("0b4e7c9a2d6f_", "b2e4f6a8c0d1_", "e4a1f8c2d7b6_")
(
"0b4e7c9a2d6f_",
"b2e4f6a8c0d1_",
"e4a1f8c2d7b6_",
"8d5a2f7c1b4e_",
)
):
continue
shutil.copy2(path, legacy_revisions / path.name)
@@ -146,7 +151,7 @@ class WorkflowMigrationTests(unittest.TestCase):
manifest_factories=(get_manifest,),
)
self.assertIn("e4a1f8c2d7b6", result.current_revision or "")
self.assertIn("8d5a2f7c1b4e", result.current_revision or "")
engine = create_engine(url)
try:
upgraded_tables = set(inspect(engine).get_table_names())
@@ -181,6 +186,25 @@ class WorkflowMigrationTests(unittest.TestCase):
engine = create_engine(url)
try:
with engine.begin() as connection:
for index_name in (
"ix_workflow_instance_steps_work_assignment",
"ix_workflow_instance_steps_work_due_at",
"ix_workflow_instance_steps_work_assignment_id",
"ix_workflow_instance_steps_work_assignment_kind",
):
connection.execute(text(f"DROP INDEX {index_name}"))
for column_name in (
"work_due_at",
"work_assignment_label",
"work_assignment_id",
"work_assignment_kind",
):
connection.execute(
text(
"ALTER TABLE workflow_instance_steps "
f"DROP COLUMN {column_name}"
)
)
connection.execute(
text(
"ALTER TABLE workflow_definition_revisions "
@@ -196,7 +220,7 @@ class WorkflowMigrationTests(unittest.TestCase):
connection.execute(
text(
"UPDATE alembic_version SET version_num = "
"'b2e4f6a8c0d1' WHERE version_num = 'e4a1f8c2d7b6'"
"'b2e4f6a8c0d1' WHERE version_num = '8d5a2f7c1b4e'"
)
)
finally:
@@ -218,9 +242,23 @@ class WorkflowMigrationTests(unittest.TestCase):
}
self.assertIn("bpmn_runtime_kind", columns)
self.assertIn("bpmn_executable", columns)
step_columns = {
item["name"]
for item in inspect(engine).get_columns(
"workflow_instance_steps"
)
}
self.assertTrue(
{
"work_assignment_kind",
"work_assignment_id",
"work_assignment_label",
"work_due_at",
}.issubset(step_columns)
)
with engine.connect() as connection:
self.assertIn(
"e4a1f8c2d7b6",
"8d5a2f7c1b4e",
set(MigrationContext.configure(connection).get_current_heads()),
)
finally:
+314 -1
View File
@@ -12,6 +12,10 @@ from govoplan_core.core.access import (
PrincipalRef,
)
from govoplan_core.core.automation import AutomationPrincipalResolution
from govoplan_core.core.campaigns import (
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
CampaignWorkHandoffInspection,
)
from govoplan_core.core.events import EventTenantRef, PlatformEvent
from govoplan_core.core.recovery import RecoveryCheckpoint, RecoveryOperation
from govoplan_core.core.runtime_coordination import (
@@ -39,6 +43,7 @@ from govoplan_workflow_engine.backend.schemas import (
WorkflowNode,
)
from govoplan_workflow_engine.backend.service import (
WorkflowConflictError,
activate_definition,
create_definition,
)
@@ -82,17 +87,58 @@ class AutomationProvider:
)
class CampaignHandoffProvider:
def __init__(self) -> None:
self.allowed = True
self.status = "completed"
self.revision = 2
self.inspections: list[tuple[str, int | None]] = []
def prepare_handoff(self, _session, _principal, *, request):
raise AssertionError("The external wait must not create Campaign work.")
def inspect_handoff(
self,
_session,
_principal,
*,
tenant_id,
assignment_id,
expected_revision=None,
):
assert tenant_id == "tenant-1"
self.inspections.append((assignment_id, expected_revision))
return CampaignWorkHandoffInspection(
allowed=self.allowed,
status=self.status,
assignment_revision=self.revision,
action_url="/campaigns/campaign-1/work?assignment=assignment-1",
assignment_ref=f"campaign-work-assignment:assignment-1:r{self.revision}",
reason=None if self.allowed else "Campaign access was revoked.",
provenance={"access_rechecked": True},
)
class Registry:
def __init__(self) -> None:
self.provider = AutomationProvider()
self.campaign = CampaignHandoffProvider()
def has_capability(self, name: str) -> bool:
return name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER
return (
name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER
or (
name == CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION
and self.campaign is not None
)
)
def capability(self, name: str):
if not self.has_capability(name):
raise KeyError(name)
if name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER:
return self.provider
return self.campaign
def graph(start_type: str, *, wait: WorkflowNode | None = None) -> WorkflowGraph:
@@ -123,6 +169,83 @@ def graph(start_type: str, *, wait: WorkflowNode | None = None) -> WorkflowGraph
return WorkflowGraph(nodes=nodes, edges=edges)
def external_handoff_graph(*, timeout_after: str = "1h") -> WorkflowGraph:
return WorkflowGraph(
nodes=[
WorkflowNode(id="start", type="workflow.start.manual"),
WorkflowNode(
id="campaign_work",
type="workflow.external_handoff",
label="Complete Campaign review",
config={
"provider_capability": CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
"event_type": "campaign.work.changed",
"event_filter": {
"payload": {"assignment_id": "$input.assignment_id"}
},
"outcome_path": "payload.outcome",
"terminal_outcomes": {
"completed": "completed",
"rejected": "rejected",
"cancelled": "cancelled",
},
"observed_outcomes": ["assigned", "accepted", "reassigned"],
"external_id": "$input.assignment_id",
"expected_revision": "$input.assignment_revision",
"action_url": "$input.action_url",
"immutable_ref": "$input.assignment_ref",
"optional_capabilities": "$input.optional_capabilities",
"timeout_after": timeout_after,
"view_surface_ids": ["campaigns.page.work"],
},
),
WorkflowNode(id="completed", type="workflow.end.completed"),
WorkflowNode(
id="rejected",
type="workflow.end.cancelled",
config={"reason": "Campaign work rejected"},
),
WorkflowNode(
id="cancelled",
type="workflow.end.cancelled",
config={"reason": "Campaign work cancelled"},
),
WorkflowNode(
id="timed_out",
type="workflow.end.cancelled",
config={"reason": "Campaign work timed out"},
),
],
edges=[
WorkflowEdge(id="start-work", source="start", target="campaign_work"),
WorkflowEdge(
id="work-completed",
source="campaign_work",
source_port="completed",
target="completed",
),
WorkflowEdge(
id="work-rejected",
source="campaign_work",
source_port="rejected",
target="rejected",
),
WorkflowEdge(
id="work-cancelled",
source="campaign_work",
source_port="cancelled",
target="cancelled",
),
WorkflowEdge(
id="work-timeout",
source="campaign_work",
source_port="timed_out",
target="timed_out",
),
],
)
class WorkflowTriggerTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite:///:memory:")
@@ -175,6 +298,18 @@ class WorkflowTriggerTests(unittest.TestCase):
)
return definition
def _external_input(self) -> dict[str, object]:
return {
"assignment_id": "assignment-1",
"assignment_revision": 1,
"action_url": "/campaigns/campaign-1/work?assignment=assignment-1",
"assignment_ref": "campaign-work-assignment:assignment-1:r1",
"optional_capabilities": {
"tasks": False,
"notifications": False,
},
}
def test_schedule_registration_dispatch_and_replay_are_durable(self) -> None:
definition = self._definition(
graph("workflow.start.schedule"),
@@ -280,6 +415,184 @@ class WorkflowTriggerTests(unittest.TestCase):
self.assertEqual(1, result["waits_timed_out"])
self.assertEqual("completed", instance.status)
def test_external_handoff_observes_duplicate_safe_events_and_resumes(self) -> None:
definition = self._definition(
external_handoff_graph(),
automation=False,
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="campaign-handoff-1",
input=self._external_input(),
),
)
step = self.session.get(WorkflowInstanceStep, instance.current_step_id)
state = self.session.scalar(select(WorkflowWaitState))
assert step is not None and state is not None
self.assertEqual("external_handoff", state.mode)
self.assertEqual("assigned", step.handoff["state"])
self.assertEqual(
["notifications", "tasks"],
step.handoff["unavailable_optional_capabilities"],
)
dispatcher = SqlWorkflowTriggerDispatcher(registry=self.registry)
accepted = PlatformEvent(
type="campaign.work.changed",
module_id="campaigns",
event_id="campaign-event-accepted",
tenant=EventTenantRef(id="tenant-1"),
payload={
"assignment_id": "assignment-1",
"assignment_revision": 2,
"outcome": "accepted",
},
)
observed = dispatcher.ingest_event(self.session, event=accepted)
duplicate = dispatcher.ingest_event(self.session, event=accepted)
self.assertEqual(0, observed["waits_triggered"])
self.assertEqual(0, duplicate["waits_triggered"])
self.assertEqual("accepted", step.handoff["state"])
self.assertEqual(
1,
self.session.query(WorkflowInstanceEvent)
.filter(
WorkflowInstanceEvent.kind
== "workflow.external_handoff.observed"
)
.count(),
)
completed = dispatcher.ingest_event(
self.session,
event=PlatformEvent(
type="campaign.work.changed",
module_id="campaigns",
event_id="campaign-event-completed",
tenant=EventTenantRef(id="tenant-1"),
payload={
"assignment_id": "assignment-1",
"assignment_revision": 2,
"outcome": "completed",
},
),
)
result = dispatcher.dispatch_due(self.session)
self.assertEqual(1, completed["waits_triggered"])
self.assertEqual(1, result["waits_resumed"])
self.assertEqual("completed", instance.status)
self.assertEqual(
[("assignment-1", 2)],
self.registry.campaign.inspections,
)
self.assertEqual(
"campaign-work-assignment:assignment-1:r2",
step.external_ref,
)
def test_external_handoff_revoked_access_blocks_until_rechecked(self) -> None:
definition = self._definition(
external_handoff_graph(),
automation=False,
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="campaign-handoff-revoked",
input=self._external_input(),
),
)
dispatcher = SqlWorkflowTriggerDispatcher(registry=self.registry)
dispatcher.ingest_event(
self.session,
event=PlatformEvent(
type="campaign.work.changed",
module_id="campaigns",
tenant=EventTenantRef(id="tenant-1"),
payload={
"assignment_id": "assignment-1",
"assignment_revision": 2,
"outcome": "completed",
},
),
)
self.registry.campaign.allowed = False
blocked = dispatcher.dispatch_due(self.session)
step = self.session.get(WorkflowInstanceStep, instance.current_step_id)
assert step is not None
self.assertEqual(1, blocked["waits_skipped"])
self.assertEqual("waiting", instance.status)
self.assertEqual("blocked", step.handoff["state"])
self.assertIn("revoked", str(step.handoff["message"]))
self.registry.campaign.allowed = True
resumed = dispatcher.dispatch_due(self.session)
self.assertEqual(1, resumed["waits_resumed"])
self.assertEqual("completed", instance.status)
def test_external_handoff_timeout_and_optional_provider_absence(self) -> None:
definition = self._definition(
external_handoff_graph(timeout_after="1s"),
automation=False,
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="campaign-handoff-timeout",
input=self._external_input(),
),
)
result = SqlWorkflowTriggerDispatcher(
registry=self.registry
).dispatch_due(
self.session,
now=datetime.now(tz=UTC) + timedelta(seconds=2),
)
self.assertEqual(1, result["waits_timed_out"])
self.assertEqual("cancelled", instance.status)
unavailable = self._definition(
external_handoff_graph(),
automation=False,
)
self.registry.campaign = None # type: ignore[assignment]
with self.assertRaisesRegex(WorkflowConflictError, "is not available"):
start_instance(
self.session,
tenant_id="tenant-1",
definition_id=unavailable.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="campaign-handoff-unavailable",
input=self._external_input(),
),
)
def test_parent_workflow_outcome_starts_pinned_child(self) -> None:
parent = self._definition(
graph("workflow.start.manual"),