Files
govoplan-workflow-engine/tests/test_instance_service.py
T

1634 lines
55 KiB
Python

from __future__ import annotations
from dataclasses import replace
from datetime import UTC, datetime, timedelta
import unittest
from sqlalchemy import create_engine, func, select
from sqlalchemy.orm import Session, sessionmaker
from govoplan_core.auth import ApiPrincipal
from govoplan_core.core.access import (
CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER,
PrincipalRef,
)
from govoplan_core.core.automation import (
ActionDefinition,
ActionExecutionResult,
ActionPreview,
AutomationPrincipalResolution,
EffectDefinition,
EffectPreview,
ObservedEffect,
)
from govoplan_core.core.dataflows import (
CAPABILITY_DATAFLOW_RUN_LIFECYCLE,
DataflowRunDescriptor,
)
from govoplan_core.core.institutional import (
InstitutionalReference,
ServiceBinding,
ServiceDefinition,
ServiceLaunchRequest,
TemporalRevision,
)
from govoplan_core.core.recovery import (
RecoveryCheckpoint,
RecoveryOperation,
RecoveryStatus,
verify_recovery_evidence_chain,
)
from govoplan_core.core.runtime_coordination import (
DistributedLease,
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 (
WorkflowDefinition,
WorkflowDefinitionRevision,
WorkflowInstance,
WorkflowInstanceEvent,
WorkflowInstanceStep,
WorkflowTrigger,
WorkflowTriggerDelivery,
WorkflowWaitState,
)
from govoplan_workflow_engine.backend.instance_service import (
SqlWorkflowRuntimeWorker,
_action_preview_payload,
cancel_instance,
instance_response,
reconcile_instance,
resolve_step,
start_instance,
)
from govoplan_workflow_engine.backend.recovery import (
acquire_workflow_state_fence,
begin_workflow_action_recovery,
)
from govoplan_workflow_engine.backend.schemas import (
BpmnRevisionInput,
WorkflowDefinitionCreateRequest,
WorkflowEdge,
WorkflowGraph,
WorkflowInstanceStartRequest,
WorkflowNode,
WorkflowStepActionRequest,
)
from govoplan_workflow_engine.backend.bpmn_adapters import NATIVE_LINEAR_ADAPTER_ID
from govoplan_workflow_engine.backend.service import (
WorkflowConflictError,
activate_definition,
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
except ModuleNotFoundError as exc:
if exc.name != "test_bpmn":
raise
from tests.test_bpmn import NATIVE_BPMN
def principal() -> ApiPrincipal:
return ApiPrincipal(
principal=PrincipalRef(
account_id="account-1",
membership_id="membership-1",
tenant_id="tenant-1",
scopes=frozenset(
{
"workflow:definition:read",
"workflow:instance:read",
"workflow:instance:start",
"workflow:instance:transition",
"dataflow:pipeline:run",
}
),
),
account=object(),
user=object(),
)
def runtime_identity() -> RuntimeIdentity:
return RuntimeIdentity(
installation_id="workflow-engine-tests",
node_id="workflow-worker",
incarnation="workflow-worker-incarnation",
role="worker",
software_version="test",
composition_hash="c" * 64,
)
def runtime_graph() -> WorkflowGraph:
return WorkflowGraph(
nodes=[
WorkflowNode(
id="start",
type="workflow.start.manual",
label="Start",
config={"input_schema_ref": ""},
),
WorkflowNode(
id="flow",
type="workflow.dataflow",
label="Prepare evidence",
config={
"pipeline_ref": "pipeline:pipeline-1",
"revision": 3,
"environment": "development",
"row_limit": 250,
"publication_target_ref": "",
"warning_policy": "review",
"input_mapping": {},
"view_surface_ids": [
"dataflow.module",
"dataflow.route.pipelines",
],
},
),
WorkflowNode(
id="complete",
type="workflow.end.completed",
label="Complete",
config={"output_mapping": {}},
),
WorkflowNode(
id="cancelled",
type="workflow.end.cancelled",
label="Rejected",
config={"reason": "Rejected during review"},
),
],
edges=[
WorkflowEdge(
id="start-flow",
source="start",
target="flow",
),
WorkflowEdge(
id="flow-complete",
source="flow",
source_port="success",
target="complete",
),
WorkflowEdge(
id="flow-warning",
source="flow",
source_port="warning",
target="complete",
),
WorkflowEdge(
id="flow-review",
source="flow",
source_port="review_required",
target="complete",
),
WorkflowEdge(
id="flow-failure",
source="flow",
source_port="failure",
target="cancelled",
),
],
)
def action_graph() -> WorkflowGraph:
return WorkflowGraph(
nodes=[
WorkflowNode(
id="start",
type="workflow.start.manual",
config={"input_schema_ref": ""},
),
WorkflowNode(
id="action",
type="workflow.capability",
config={
"capability": "test.actions",
"operation": "test.case.record",
"input_mapping": {"case_id": "$input.case_id"},
"idempotency_key": "$input.case_id",
"failure_policy": "manual",
},
),
WorkflowNode(
id="complete",
type="workflow.end.completed",
config={"output_mapping": {}},
),
WorkflowNode(
id="failed",
type="workflow.end.cancelled",
config={"reason": "Action rejected"},
),
],
edges=[
WorkflowEdge(
id="start-action",
source="start",
target="action",
),
WorkflowEdge(
id="action-complete",
source="action",
source_port="success",
target="complete",
),
WorkflowEdge(
id="action-failed",
source="action",
source_port="failure",
target="failed",
),
],
)
class FakeDataflowLifecycle:
def __init__(self) -> None:
self.runs: dict[str, DataflowRunDescriptor] = {}
self.requests = []
self.cancelled: list[str] = []
def start_run(self, _session, _principal, *, request):
self.requests.append(request)
run_ref = f"run:{len(self.requests)}"
descriptor = DataflowRunDescriptor(
ref=run_ref,
pipeline_ref=request.pipeline_ref,
revision=request.revision,
status="queued",
definition_hash="definition-hash",
executor_version="test",
metadata={"progress_percent": 0, "progress_phase": "queued"},
)
self.runs[run_ref] = descriptor
return descriptor
def get_run(self, _session, _principal, *, run_ref):
return self.runs.get(run_ref)
def cancel_run(self, _session, _principal, *, run_ref):
descriptor = self.runs[run_ref]
descriptor = replace(descriptor, status="cancelled")
self.runs[run_ref] = descriptor
self.cancelled.append(run_ref)
return descriptor
def finish(
self,
run_ref: str,
*,
diagnostics: list[dict[str, object]] | None = None,
) -> None:
self.runs[run_ref] = replace(
self.runs[run_ref],
status="succeeded",
output_publication_ref="publication:1",
output_datasource_ref="datasource:1",
output_materialization_ref="materialization:1",
input_row_count=12,
output_row_count=10,
metadata={"diagnostics": diagnostics or []},
)
def fail(self, run_ref: str) -> None:
self.runs[run_ref] = replace(
self.runs[run_ref],
status="failed",
error="Data quality gate failed.",
)
def mark_outcome_unknown(self, run_ref: str) -> None:
self.runs[run_ref] = replace(
self.runs[run_ref],
status="outcome_unknown",
error="Output acknowledgement was lost.",
metadata={
"recovery": {
"operation_id": "recovery:dataflow:1",
"status": "outcome_unknown",
"requires_attention": True,
}
},
)
class FakeAutomationProvider:
def __init__(self) -> None:
self.requests = []
def resolve_automation_principal(self, _session, *, request):
self.requests.append(request)
return AutomationPrincipalResolution(
allowed=True,
principal=principal(),
granted_scopes=request.grant_scopes,
provenance={"status": "rechecked"},
)
class FakeActionProvider:
action = ActionDefinition(
action_key="test.case.record",
owner_module="test",
description="Record a test case.",
input_schema_ref="schema:test.case.record@1",
expected_effect_keys=("test.case.recorded",),
)
effect = EffectDefinition(
effect_key="test.case.recorded",
owner_module="test",
operation="created",
description="A test case was recorded.",
)
def __init__(self, *states: str) -> None:
self.states = list(states or ("completed",))
self.requests = []
def action_definitions(self):
return (self.action,)
def effect_definitions(self):
return (self.effect,)
def preview_action(self, _session, _principal, *, request):
return ActionPreview(
action_key=request.action_key,
allowed=True,
summary="Record one case.",
risk_level=self.action.risk_level,
reversibility=self.action.reversibility,
effects=(
EffectPreview(
effect_key=self.effect.effect_key,
summary="Record case.",
),
),
preview_ref="preview:test",
)
def execute_action(self, _session, _principal, *, request):
self.requests.append(request)
state = self.states.pop(0) if self.states else "completed"
if state == "exception":
raise RuntimeError("Provider acknowledgement was lost.")
if state != "completed":
return ActionExecutionResult(
state=state,
error="Temporary action failure.",
)
return ActionExecutionResult(
state="completed",
output={"case_ref": "case:1"},
observed_effects=(
ObservedEffect(
effect_key=self.effect.effect_key,
operation="created",
resource_ref="case:1",
),
),
audit_event_refs=("audit:1",),
)
class Registry:
def __init__(
self,
dataflow: FakeDataflowLifecycle,
automation: FakeAutomationProvider | None = None,
action: FakeActionProvider | None = None,
) -> None:
self.dataflow = dataflow
self.automation = automation
self.action = action
def has_capability(self, name: str) -> bool:
return (
name == CAPABILITY_DATAFLOW_RUN_LIFECYCLE
or (
name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER
and self.automation is not None
)
or (name == "test.actions" and self.action is not None)
)
def capability(self, name: str):
if name == CAPABILITY_DATAFLOW_RUN_LIFECYCLE:
return self.dataflow
if (
name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER
and self.automation is not None
):
return self.automation
if name == "test.actions" and self.action is not None:
return self.action
raise KeyError(name)
class WorkflowInstanceServiceTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite:///:memory:")
Base.metadata.create_all(
self.engine,
tables=[
DistributedLease.__table__,
RecoveryOperation.__table__,
RecoveryCheckpoint.__table__,
WorkflowDefinition.__table__,
WorkflowDefinitionRevision.__table__,
WorkflowInstance.__table__,
WorkflowInstanceStep.__table__,
WorkflowInstanceEvent.__table__,
WorkflowTrigger.__table__,
WorkflowTriggerDelivery.__table__,
WorkflowWaitState.__table__,
],
)
self.Session = sessionmaker(bind=self.engine)
self.session: Session = self.Session()
bind_process_runtime_identity(runtime_identity())
self.dataflow = FakeDataflowLifecycle()
self.registry = Registry(self.dataflow)
self.definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Monthly governed processing",
graph=runtime_graph(),
execution_mode="hybrid",
view_id="view-1",
view_revision_id="view-revision-1",
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=self.definition.id,
actor_id="account-1",
revision=1,
)
self.session.commit()
def tearDown(self) -> None:
bind_process_runtime_identity(None)
self.session.close()
Base.metadata.drop_all(
self.engine,
tables=[
WorkflowWaitState.__table__,
WorkflowTriggerDelivery.__table__,
WorkflowTrigger.__table__,
WorkflowInstanceEvent.__table__,
WorkflowInstanceStep.__table__,
WorkflowInstance.__table__,
WorkflowDefinitionRevision.__table__,
WorkflowDefinition.__table__,
RecoveryCheckpoint.__table__,
RecoveryOperation.__table__,
DistributedLease.__table__,
],
)
self.engine.dispose()
def _start(self, key: str = "request-1") -> WorkflowInstance:
instance, replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=self.definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key=key,
input={"case_id": "case-1"},
correlation_id="correlation-1",
),
)
self.assertFalse(replayed)
return instance
def test_start_pins_revision_and_replays_idempotently(self) -> None:
instance = self._start()
replayed, was_replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=self.definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="request-1",
input={"case_id": "case-1"},
correlation_id="correlation-1",
),
)
self.assertTrue(was_replayed)
self.assertEqual(instance.id, replayed.id)
self.assertEqual("waiting", instance.status)
self.assertEqual("user", instance.start_origin)
self.assertEqual(1, len(self.dataflow.requests))
response = instance_response(self.session, instance)
self.assertEqual("hybrid", response.execution_mode)
self.assertEqual("user", response.start_origin)
self.assertIsNotNone(response.view_context)
self.assertEqual(
[
"dataflow.module",
"dataflow.route.pipelines",
],
response.view_context.visible_surface_ids,
)
self.assertEqual([1, 2], [step.sequence for step in response.steps])
self.assertEqual("run:1", response.steps[-1].external_ref)
self.assertGreaterEqual(len(response.events), 4)
with self.assertRaises(WorkflowConflictError):
start_instance(
self.session,
tenant_id="tenant-1",
definition_id=self.definition.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="request-1",
input={"case_id": "another-case"},
correlation_id="correlation-1",
),
)
def test_service_launcher_starts_and_replays_exact_active_workflow(self) -> None:
now = datetime.now(tz=UTC)
service_ref = InstitutionalReference(
kind="service",
owner_module="services",
object_id="monthly-service",
tenant_id="tenant-1",
version="3",
valid_at=now,
)
binding = ServiceBinding("workflow", self.definition.id)
definition = ServiceDefinition(
reference=service_ref,
key="monthly.processing",
temporal=TemporalRevision(
revision="3",
valid_from=now - timedelta(days=1),
recorded_at=now - timedelta(days=2),
change_reason="Published workflow service.",
),
title="Monthly processing",
audience=("authenticated",),
bindings=(binding,),
publication_state="published",
)
request = ServiceLaunchRequest(
service_ref=service_ref,
binding=binding,
idempotency_key="service-workflow-1",
requested_at=now,
parameters={"case_id": "case-1"},
)
launcher = WorkflowServiceLauncher(self.registry)
first = launcher.launch_service(
self.session,
principal(),
definition=definition,
request=request,
)
second = launcher.launch_service(
self.session,
principal(),
definition=definition,
request=request,
)
self.assertFalse(first.replayed)
self.assertTrue(second.replayed)
self.assertEqual(first.target_ref, second.target_ref)
def test_guided_workflow_rejects_automated_start_origin(self) -> None:
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Guided review",
graph=WorkflowGraph(
nodes=[
WorkflowNode(
id="start",
type="workflow.start.api",
config={
"input_schema_ref": "schema:input",
"authorization_policy_ref": "policy:start",
},
),
WorkflowNode(
id="activity",
type="workflow.activity",
config={"title": "Review"},
),
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",
allow_automation=True,
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
with self.assertRaisesRegex(
WorkflowConflictError,
"Guided workflows must be started by a user",
):
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="automated-guided",
),
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:
action = FakeActionProvider()
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Action workflow",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="action-instance",
input={"case_id": "case-1"},
),
)
response = instance_response(self.session, instance)
self.assertFalse(replayed)
self.assertEqual("completed", response.status)
self.assertEqual(1, len(action.requests))
self.assertEqual(
"test.actions:test.case.record:case-1",
action.requests[0].idempotency_key,
)
self.assertEqual(
"case:1",
response.steps[1].output["execution"]["observed_effects"][0][
"resource_ref"
],
)
self.assertIn(
"workflow.action.completed",
{event.kind for event in response.events},
)
operation = self.session.scalar(
select(RecoveryOperation).where(
RecoveryOperation.resource_id == response.steps[1].id
)
)
assert operation is not None
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
self.assertTrue(verify_recovery_evidence_chain(self.session, operation.id))
def test_missing_optional_action_provider_blocks_without_recovery_effect(
self,
) -> None:
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Unavailable action provider",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
with self.assertRaisesRegex(WorkflowConflictError, "not available"):
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="missing-action-provider",
input={"case_id": "case-missing"},
),
)
self.assertEqual(
0,
self.session.scalar(select(func.count(RecoveryOperation.id))),
)
def test_retryable_module_action_reuses_the_idempotency_key(self) -> None:
action = FakeActionProvider("retryable", "completed")
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Retry action",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="retry-action-instance",
input={"case_id": "case-2"},
),
)
waiting = instance_response(self.session, instance)
current = waiting.steps[-1]
self.assertEqual("retryable", current.handoff["state"])
resolved = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(action="retry"),
)
self.assertEqual("completed", resolved.status)
self.assertEqual(
action.requests[0].idempotency_key,
action.requests[1].idempotency_key,
)
def test_unknown_module_action_requires_evidence_before_retry(self) -> None:
action = FakeActionProvider("exception", "completed")
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Uncertain action",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="uncertain-action",
input={"case_id": "case-unknown"},
),
)
waiting = instance_response(self.session, instance)
current = waiting.steps[-1]
self.assertEqual("outcome_unknown", current.handoff["state"])
self.assertNotIn("retry", current.handoff["allowed_actions"])
self.assertEqual(
{"confirm_effect", "confirm_absent", "cancel"},
set(current.handoff["allowed_actions"]),
)
operation = self.session.get(
RecoveryOperation,
current.handoff["recovery"]["operation_id"],
)
assert operation is not None
self.assertEqual(RecoveryStatus.OUTCOME_UNKNOWN.value, operation.status)
with self.assertRaisesRegex(WorkflowConflictError, "evidence reference"):
resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(action="confirm_absent"),
)
resolved = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(
action="confirm_absent",
evidence=["provider-search:case-unknown"],
),
)
current = next(step for step in resolved.steps if step.id == current.id)
self.assertEqual("retryable", current.handoff["state"])
self.session.refresh(operation)
self.assertEqual(RecoveryStatus.RECOVERED.value, operation.status)
completed = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(action="retry"),
)
self.assertEqual("completed", completed.status)
self.assertEqual(2, len(action.requests))
self.assertEqual(
action.requests[0].idempotency_key,
action.requests[1].idempotency_key,
)
def test_confirmed_unknown_effect_advances_without_provider_replay(self) -> None:
action = FakeActionProvider("exception")
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Confirmed uncertain action",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="confirmed-uncertain-action",
input={"case_id": "case-confirmed"},
),
)
current = instance_response(self.session, instance).steps[-1]
operation_id = current.handoff["recovery"]["operation_id"]
resolved = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(
action="confirm_effect",
output={"case_ref": "case:confirmed"},
evidence=["provider-receipt:case-confirmed"],
),
)
self.assertEqual("completed", resolved.status)
self.assertEqual(1, len(action.requests))
operation = self.session.get(RecoveryOperation, operation_id)
assert operation is not None
self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status)
def test_tampered_action_evidence_prevents_workflow_completion(self) -> None:
class TamperingProvider(FakeActionProvider):
def execute_action(self, session, principal, *, request):
result = super().execute_action(
session,
principal,
request=request,
)
checkpoint = session.scalar(
select(RecoveryCheckpoint).order_by(RecoveryCheckpoint.sequence)
)
assert checkpoint is not None
checkpoint.summary = "tampered provider evidence"
session.commit()
return result
action = TamperingProvider()
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Tampered action",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
with self.assertRaisesRegex(ValueError, "chain verification failed"):
start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="tampered-action",
input={"case_id": "case-tampered"},
),
)
operation = self.session.scalar(
select(RecoveryOperation).where(
RecoveryOperation.module_id == "workflow_engine"
)
)
assert operation is not None
self.assertEqual(RecoveryStatus.RUNNING.value, operation.status)
self.assertFalse(verify_recovery_evidence_chain(self.session, operation.id))
def test_dataflow_unknown_outcome_blocks_workflow_until_resolved(self) -> None:
instance = self._start("dataflow-unknown")
self.dataflow.mark_outcome_unknown("run:1")
changed = reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
)
self.assertFalse(changed)
self.assertEqual("waiting", instance.status)
current = next(
step for step in instance.steps if step.id == instance.current_step_id
)
self.assertEqual("dataflow_recovery", current.handoff["kind"])
self.assertEqual(["cancel"], current.handoff["allowed_actions"])
self.assertTrue(current.handoff["recovery"]["requires_attention"])
self.dataflow.finish("run:1")
changed = reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
)
self.assertTrue(changed)
self.assertEqual("completed", instance.status)
def test_worker_does_not_advance_an_instance_owned_by_another_runtime(
self,
) -> None:
instance = self._start("fenced-instance")
self.session.commit()
self.dataflow.finish("run:1")
fence = acquire_workflow_state_fence(
self.session,
resource_key=f"workflow:instance:{instance.id}",
)
assert fence is not None
self.session.commit()
bind_process_runtime_identity(
RuntimeIdentity(
installation_id="workflow-engine-tests",
node_id="other-worker",
incarnation="other-worker-incarnation",
role="worker",
software_version="test",
composition_hash="e" * 64,
)
)
worker = SqlWorkflowRuntimeWorker(
registry=Registry(self.dataflow, FakeAutomationProvider()),
)
summary = worker.reconcile_pending(self.session)
self.assertEqual(0, summary["advanced"])
self.assertEqual(1, summary["skipped"])
self.session.refresh(instance)
self.assertEqual("waiting", instance.status)
lease = self.session.scalar(
select(DistributedLease).where(
DistributedLease.resource_key == f"workflow:instance:{instance.id}"
)
)
assert lease is not None
lease.expires_at = utcnow() - timedelta(seconds=1)
self.session.commit()
summary = worker.reconcile_pending(self.session)
self.assertEqual(1, summary["advanced"])
self.session.refresh(instance)
self.assertEqual("completed", instance.status)
bind_process_runtime_identity(runtime_identity())
def test_stale_action_attempt_becomes_unknown_instead_of_replaying(
self,
) -> None:
action = FakeActionProvider("retryable", "completed")
registry = Registry(self.dataflow, action=action)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Stale action attempt",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="stale-action-attempt",
input={"case_id": "case-stale"},
),
)
current = instance_response(self.session, instance).steps[-1]
revision = self.session.get(
WorkflowDefinitionRevision,
instance.definition_revision_id,
)
assert revision is not None
request = action.requests[0]
preview = action.preview_action(
self.session,
principal(),
request=request,
)
started = begin_workflow_action_recovery(
self.session,
instance=instance,
step=current,
revision=revision,
definition=action.action,
capability_name="test.actions",
request_idempotency_key=request.idempotency_key,
action_input={"case_id": "case-stale"},
preview_payload=_action_preview_payload(preview),
)
self.assertFalse(started.replayed)
lease = self.session.scalar(
select(DistributedLease).where(
DistributedLease.resource_key == f"workflow:step:{current.id}"
)
)
assert lease is not None
lease.expires_at = utcnow() - timedelta(seconds=1)
self.session.commit()
bind_process_runtime_identity(
RuntimeIdentity(
installation_id="workflow-engine-tests",
node_id="takeover-worker",
incarnation="takeover-worker-incarnation",
role="worker",
software_version="test",
composition_hash="f" * 64,
)
)
resolved = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=current.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowStepActionRequest(action="retry"),
)
current = next(step for step in resolved.steps if step.id == current.id)
self.assertEqual("outcome_unknown", current.handoff["state"])
self.assertNotIn("retry", current.handoff["allowed_actions"])
self.assertEqual(1, len(action.requests))
operation = self.session.get(RecoveryOperation, started.operation_id)
assert operation is not None
self.assertEqual(RecoveryStatus.OUTCOME_UNKNOWN.value, operation.status)
bind_process_runtime_identity(runtime_identity())
def test_automated_dataflow_failure_policy_fails_without_handoff(
self,
) -> None:
graph = runtime_graph()
flow = next(node for node in graph.nodes if node.id == "flow")
flow.config["warning_policy"] = "continue"
flow.config["failure_policy"] = "fail"
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Automated Dataflow",
graph=graph,
execution_mode="automated",
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
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="automated-dataflow",
),
)
self.dataflow.fail("run:1")
changed = reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
)
self.assertTrue(changed)
self.assertEqual("failed", instance.status)
self.assertIsNone(instance.current_step_id)
self.assertFalse(
any(
step.handoff.get("kind") == "dataflow_failure"
for step in instance.steps
)
)
def test_native_bpmn_profile_runs_through_canonical_instance_state(self) -> None:
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="BPMN governed review",
graph=runtime_graph(),
bpmn=BpmnRevisionInput(
xml=NATIVE_BPMN,
adapter_id=NATIVE_LINEAR_ADAPTER_ID,
),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
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="bpmn-request-1",
input={"case_id": "case-bpmn"},
),
)
self.assertFalse(replayed)
self.assertEqual("waiting", instance.status)
waiting = instance_response(self.session, instance).steps[-1]
self.assertEqual("workflow.activity", waiting.node_type)
self.assertEqual("Review request", waiting.handoff["title"])
replay, was_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="bpmn-request-1",
input={"case_id": "case-bpmn"},
),
)
self.assertTrue(was_replayed)
self.assertEqual(instance.id, replay.id)
resolved = resolve_step(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
step_id=waiting.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
payload=WorkflowStepActionRequest(
action="complete",
comment="Reviewed.",
evidence=["case:case-bpmn"],
),
)
response = instance_response(self.session, resolved)
self.assertEqual("completed", response.status)
self.assertEqual(
"workflow.instance.completed",
response.events[-1].kind,
)
def test_reconcile_completes_with_stable_dataflow_output_refs(self) -> None:
instance = self._start()
self.dataflow.finish("run:1")
changed = reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
actor_id="account-1",
)
response = instance_response(self.session, instance)
self.assertTrue(changed)
self.assertEqual("completed", response.status)
flow_output = response.context["steps"]["flow"]
self.assertEqual("publication:1", flow_output["output_publication_ref"])
self.assertEqual("datasource:1", flow_output["output_datasource_ref"])
self.assertEqual(
"materialization:1",
flow_output["output_materialization_ref"],
)
self.assertEqual("workflow.instance.completed", response.events[-1].kind)
def test_warning_requires_review_and_approve_resumes(self) -> None:
instance = self._start()
self.dataflow.finish(
"run:1",
diagnostics=[
{
"severity": "warning",
"code": "review.required",
"message": "Verify unmatched records.",
}
],
)
reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
)
step_id = str(instance.current_step_id)
self.assertEqual(
"review_required",
instance_response(self.session, instance).steps[-1].handoff["state"],
)
resolved = 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="approve",
comment="Evidence verified.",
evidence=["publication:1"],
),
)
self.assertEqual("completed", resolved.status)
def test_failure_can_retry_and_reject_invalid_actions(self) -> None:
instance = self._start()
self.dataflow.fail("run:1")
reconcile_instance(
self.session,
instance=instance,
principal=principal(),
registry=self.registry,
)
step_id = str(instance.current_step_id)
with self.assertRaises(WorkflowConflictError):
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="approve"),
)
retried = 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="retry"),
)
self.assertEqual("waiting", retried.status)
self.assertEqual(2, len(self.dataflow.requests))
self.assertEqual("run:2", retried.steps[-1].external_ref)
self.assertEqual("superseded", retried.steps[-2].status)
def test_cancel_propagates_to_linked_dataflow(self) -> None:
instance = self._start()
cancelled = cancel_instance(
self.session,
tenant_id="tenant-1",
instance_id=instance.id,
actor_id="account-1",
principal=principal(),
registry=self.registry,
)
self.assertEqual("cancelled", cancelled.status)
self.assertEqual(["run:1"], self.dataflow.cancelled)
def test_worker_rechecks_authorization_before_reconciling(self) -> None:
instance = self._start()
self.session.commit()
self.dataflow.finish("run:1")
automation = FakeAutomationProvider()
worker = SqlWorkflowRuntimeWorker(
registry=Registry(self.dataflow, automation),
)
summary = worker.reconcile_pending(self.session)
self.assertEqual(1, summary["advanced"])
self.assertEqual("completed", instance.status)
self.assertEqual(1, len(automation.requests))
self.assertEqual(
"rechecked",
instance.authorization_["last_resolution"]["status"],
)
def test_worker_reconciles_pending_module_action(self) -> None:
action = FakeActionProvider("pending", "completed")
automation = FakeAutomationProvider()
registry = Registry(
self.dataflow,
automation=automation,
action=action,
)
definition = create_definition(
self.session,
tenant_id="tenant-1",
actor_id="account-1",
payload=WorkflowDefinitionCreateRequest(
name="Asynchronous action",
graph=action_graph(),
),
)
activate_definition(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
)
instance, _replayed = start_instance(
self.session,
tenant_id="tenant-1",
definition_id=definition.id,
actor_id="account-1",
principal=principal(),
registry=registry,
payload=WorkflowInstanceStartRequest(
idempotency_key="pending-action",
input={"case_id": "case-3"},
),
)
self.session.commit()
worker = SqlWorkflowRuntimeWorker(registry=registry)
summary = worker.reconcile_pending(self.session)
self.assertEqual(1, summary["advanced"])
self.assertEqual("completed", instance.status)
self.assertEqual(2, len(action.requests))
self.assertEqual(
action.requests[0].idempotency_key,
action.requests[1].idempotency_key,
)
self.assertEqual(1, len(automation.requests))
if __name__ == "__main__":
unittest.main()