feat(workflow-engine): add governed DSAR coverage

This commit is contained in:
2026-08-21 03:33:06 +02:00
parent 9174e07118
commit ceb61b5867
4 changed files with 1516 additions and 0 deletions
@@ -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"]