diff --git a/README.md b/README.md index 6826b9d..e9b1164 100644 --- a/README.md +++ b/README.md @@ -20,6 +20,8 @@ for the ownership and compatibility policy. See [the module concept](docs/CONCEPT.md) and [BPMN interoperability contract](docs/BPMN_INTEROPERABILITY.md) for runtime semantics and adapter boundaries. +See [durable runtime recovery](docs/DURABLE_RUNTIME_RECOVERY.md) for action, +worker, timer, and unknown-provider-outcome handling. The optional `workflow_engine.service_launcher` capability lets Portal start an authorized active workflow from an exact published Service revision. It resolves @@ -37,6 +39,14 @@ human handoffs; event filters and variable mappings are bounded JSON expressions and never executable code. Cron remains an optional governed scheduler-adapter concern. +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 +**Effect confirmed** and **Effect absent** operator actions instead of a blind +retry. Instance workers, trigger deliveries, and timers are fenced across +hosts; linked Dataflow recovery states remain unresolved until Dataflow reports +a conclusive outcome. + ## Checks ```bash diff --git a/docs/CONCEPT.md b/docs/CONCEPT.md index 9e0de98..5f60613 100644 --- a/docs/CONCEPT.md +++ b/docs/CONCEPT.md @@ -121,13 +121,19 @@ The first executable slice now provides: - durable duration/deadline/event wait subscriptions and scale-out-safe claims - 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 + request evidence, provider-dispatch checkpoints, and atomic local projection + commits +- explicit effect-confirmed/effect-absent reconciliation for unknown provider + outcomes, with blind retries disabled +- distributed execution fences for instance reconciliation, trigger delivery, + and timer resumption +- fail-closed propagation of unresolved Dataflow publication outcomes The next execution depth should provide: - static workflow definition registration from configuration packages - guard hooks implemented through capability calls -- registry-driven generic module-action execution records -- action/effect previews for transitions that call other modules - explicit blocked, retryable, quarantined, manual-required, and compensation-required states - dashboard summary provider diff --git a/docs/DURABLE_RUNTIME_RECOVERY.md b/docs/DURABLE_RUNTIME_RECOVERY.md new file mode 100644 index 0000000..01854f5 --- /dev/null +++ b/docs/DURABLE_RUNTIME_RECOVERY.md @@ -0,0 +1,59 @@ +# Durable Workflow Runtime Recovery + +Workflow Engine uses Core's recovery ledger and distributed leases for every +consequential module-action invocation and every worker-owned instance, trigger +delivery, or timer transition. The ledger is evidence for an effect boundary; +it is not a replacement for the Workflow instance history. + +## Module Actions + +An action provider declares a recovery mode and concrete verification checks in +its `ActionDefinition`. Before dispatch, Workflow Engine persists the pinned +definition hash, canonical action-input and preview hashes, current access +context hash, provider idempotency-key hash, and an execution fence. Raw action +input and credentials are not copied into recovery evidence. + +A conclusive provider result commits the Workflow projection and verified +terminal checkpoint together. Pending and retryable provider results remain +visible handoffs, but each later provider poll gets a new ledger operation while +retaining the provider's stable action idempotency key. + +An exception, invalid result, or unannounced effect after a non-atomic dispatch +is not an ordinary failure. The step becomes `outcome_unknown`, continuation is +blocked, and Retry is unavailable. An operator must inspect the provider and +record durable evidence for one of two decisions: + +- **Effect confirmed** records that the provider effect occurred and advances + the Workflow without invoking the provider again. +- **Effect absent** records forward recovery and enables a deliberate retry + under a new invocation ledger key and the same provider idempotency key. + +Atomic actions are different: provider database changes and the Workflow +projection share one transaction. An exception rolls that transaction back, so +the step can become retryable without claiming an unknown external outcome. + +Snapshot-restore and irreversible actions are rejected before dispatch unless +the Workflow node pins the required backup or approval reference. + +## Workers, Triggers, And Timers + +Runtime workers acquire a distributed fence for each instance. Trigger +deliveries and wait states use their own fences. Database row locks remain a +local optimization; the fence is the authority proof across hosts. A second +runtime cannot advance an actively owned item. After lease expiry, a new +runtime receives a higher fence number and may resume from durable state. + +Linked Dataflow runs retain their own recovery boundary. If Dataflow reports an +unknown publication outcome, Workflow displays a recovery handoff and only +allows cancellation. It polls again after the Dataflow outcome is reconciled +and advances only from a conclusive run descriptor. + +## Process Failure Versus Effect Uncertainty + +A process failure before dispatch is retryable after current authority and +inputs are revalidated. A failure after non-atomic dispatch is an unresolved +effect even when no provider response was received. Operators must not infer +absence from a timeout, worker restart, HTTP error, or empty local result. + +Checkpoint evidence is hash-chained. A broken chain prevents a success commit; +it must be investigated as an integrity incident rather than bypassed. diff --git a/src/govoplan_workflow_engine/backend/instance_service.py b/src/govoplan_workflow_engine/backend/instance_service.py index d715f10..cff2c8f 100644 --- a/src/govoplan_workflow_engine/backend/instance_service.py +++ b/src/govoplan_workflow_engine/backend/instance_service.py @@ -65,6 +65,18 @@ from govoplan_workflow_engine.backend.service import ( get_definition, get_definition_revision, ) +from govoplan_workflow_engine.backend.recovery import ( + WorkflowActionRecovery, + WorkflowRecoveryBusy, + WorkflowRecoveryConflict, + WorkflowRecoveryError, + acquire_workflow_state_fence, + begin_workflow_action_recovery, + canonical_sha256, + claim_workflow_action_recovery, + reconcile_stale_workflow_action_recovery, + release_workflow_state_fence, +) INSTANCE_START_SCOPE = "workflow:instance:start" @@ -355,6 +367,54 @@ def reconcile_instance( message="The linked Dataflow run no longer exists.", ) return True + descriptor_recovery = descriptor.metadata.get("recovery") + recovery_requires_attention = bool( + isinstance(descriptor_recovery, Mapping) + and descriptor_recovery.get("requires_attention") + ) + if descriptor.status == "outcome_unknown" or recovery_requires_attention: + previous_state = str(step.handoff.get("state") or "") + step.status = "waiting" + step.error = ( + "The Dataflow output outcome is unresolved. Reconcile the run " + "before this Workflow can continue." + ) + step.handoff = { + "kind": "dataflow_recovery", + "state": "outcome_unknown", + "run_ref": descriptor.ref, + "pipeline_ref": descriptor.pipeline_ref, + "action_url": _dataflow_action_url( + descriptor.pipeline_ref, + descriptor.ref, + ), + "allowed_actions": ["cancel"], + "outcome_unknown": True, + "recovery": ( + dict(descriptor_recovery) + if isinstance(descriptor_recovery, Mapping) + else {"requires_attention": True} + ), + } + instance.status = "waiting" + instance.error = step.error + if previous_state != "outcome_unknown": + _record_event( + session, + instance, + step=step, + kind="workflow.dataflow.outcome_unknown", + actor_id=actor_id, + payload=dict(step.handoff), + ) + _notify_handoff( + session, + registry=registry, + instance=instance, + step=step, + subject="Workflow Dataflow outcome requires reconciliation", + ) + return False if descriptor.status in {"queued", "retrying", "running"}: step.handoff = { **dict(step.handoff), @@ -466,6 +526,113 @@ def resolve_step( principal=principal, registry=registry, ) + if payload.action in {"confirm_effect", "confirm_absent"}: + recovery_payload = step.handoff.get("recovery") + if not isinstance(recovery_payload, Mapping): + raise WorkflowConflictError( + "This handoff has no durable recovery operation to reconcile." + ) + operation_id = str(recovery_payload.get("operation_id") or "").strip() + if not operation_id: + raise WorkflowConflictError( + "This handoff has no durable recovery operation to reconcile." + ) + if not payload.evidence and not payload.comment: + raise WorkflowConflictError( + "Record an evidence reference or operator comment before " + "resolving an unknown provider outcome." + ) + try: + recovery = claim_workflow_action_recovery( + session, + operation_id=operation_id, + ) + except WorkflowRecoveryError as exc: + raise WorkflowConflictError(str(exc)) from exc + effect_occurred = payload.action == "confirm_effect" + evidence = { + "verified": True, + "checks": { + "operator_actor_ref": actor_id, + "evidence_refs_sha256": canonical_sha256(list(payload.evidence)), + "operator_output_sha256": canonical_sha256(dict(payload.output)), + "operator_comment_sha256": canonical_sha256(payload.comment or ""), + "effect_occurred": effect_occurred, + }, + } + if not effect_occurred: + previous_recovery = dict(recovery_payload) + call_number = max(1, int(previous_recovery.get("call_number") or 1)) + step.status = "waiting" + step.error = None + step.handoff = { + **dict(step.handoff), + "state": "retryable", + "message": ( + "The provider effect was verified absent. A deliberate retry " + "is now safe." + ), + "allowed_actions": ["retry", "cancel"], + "outcome_unknown": False, + "recovery": { + **previous_recovery, + "status": "recovered", + "requires_attention": False, + "next_call_number": call_number + 1, + }, + } + instance.status = "waiting" + instance.error = None + _record_event( + session, + instance, + step=step, + kind="workflow.action.effect_absent", + actor_id=actor_id, + payload={ + "recovery_operation_id": operation_id, + "evidence": list(payload.evidence), + }, + ) + recovery.commit_unknown_resolution( + session, + effect_occurred=False, + evidence=evidence, + summary="Operator verified that the provider effect is absent", + ) + return instance + output = { + **dict(step.output_), + **dict(payload.output), + "recovery_decision": "effect_confirmed", + "recovery_evidence": list(payload.evidence), + "recovery_comment": payload.comment, + } + next_node_id = _complete_step( + session, + instance=instance, + step=step, + graph=graph, + port="success", + output=output, + actor_id=actor_id, + ) + recovery.commit_unknown_resolution( + session, + effect_occurred=True, + evidence=evidence, + summary="Operator verified that the provider effect occurred", + ) + _drive_instance( + session, + instance=instance, + graph=graph, + next_node_id=next_node_id, + principal=principal, + registry=registry, + actor_id=actor_id, + ) + return instance if payload.action == "changes": step.handoff = { **dict(step.handoff), @@ -634,12 +801,29 @@ def reconcile_pending_instances( "skipped": 0, } for instance in instances: + try: + fence = acquire_workflow_state_fence( + session, + resource_key=f"workflow:instance:{instance.id}", + ) + except RuntimeError: + logger.warning( + "Workflow instance %s could not acquire a runtime fence", + instance.id, + exc_info=True, + ) + summary["skipped"] = int(summary["skipped"]) + 1 + continue + if fence is None: + summary["skipped"] = int(summary["skipped"]) + 1 + continue step = _current_step(session, instance) if step is None or step.node_type not in { "workflow.capability", "workflow.dataflow", }: summary["waiting"] = int(summary["waiting"]) + 1 + release_workflow_state_fence(session, fence) continue principal = _resolve_instance_principal( session, @@ -648,6 +832,7 @@ def reconcile_pending_instances( ) if principal is None: summary["skipped"] = int(summary["skipped"]) + 1 + release_workflow_state_fence(session, fence) continue changed = reconcile_instance( session, @@ -661,6 +846,7 @@ def reconcile_pending_instances( summary["advanced"] = int(summary["advanced"]) + 1 else: summary["waiting"] = int(summary["waiting"]) + 1 + release_workflow_state_fence(session, fence) session.flush() return summary @@ -1075,6 +1261,141 @@ def _execute_capability_step( preview_ref=preview.preview_ref, metadata=request.metadata, ) + capability_name = str(node.config.get("capability") or "") + revision = session.get( + WorkflowDefinitionRevision, + instance.definition_revision_id, + ) + if revision is None: + _set_action_handoff( + session, + instance=instance, + step=step, + state="blocked", + message="The pinned Workflow revision is unavailable.", + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + ) + return True + try: + action_recovery = begin_workflow_action_recovery( + session, + instance=instance, + step=step, + revision=revision, + definition=definition, + capability_name=capability_name, + request_idempotency_key=request.idempotency_key, + action_input=action_input, + preview_payload=preview_payload, + backup_reference=( + str(node.config.get("recovery_backup_reference") or "").strip() + or None + ), + approval_reference=( + str(node.config.get("recovery_approval_reference") or "").strip() + or None + ), + ) + except WorkflowRecoveryBusy as exc: + _set_action_handoff( + session, + instance=instance, + step=step, + state="running", + message=str(exc), + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + ) + return False + except WorkflowRecoveryConflict as exc: + recovery_status = exc.status + if recovery_status == "running": + try: + recovery_status = reconcile_stale_workflow_action_recovery( + session, + operation_id=exc.operation_id, + ) + except WorkflowRecoveryBusy: + recovery_status = "running" + outcome_unknown = recovery_status in { + "outcome_unknown", + "recovery_required", + "manual_intervention", + } + safe_to_retry = recovery_status in {"failed", "recovered", "rejected"} + previous_recovery = ( + dict(step.handoff.get("recovery")) + if isinstance(step.handoff.get("recovery"), Mapping) + else {} + ) + call_number = max(1, int(previous_recovery.get("call_number") or 1)) + _set_action_handoff( + session, + instance=instance, + step=step, + state=( + "outcome_unknown" + if outcome_unknown + else "retryable" + if safe_to_retry + else recovery_status + ), + message=( + "The provider outcome must be reconciled before this Workflow " + "can continue." + if outcome_unknown + else ( + "The stale action was proven not to have committed. A " + "deliberate retry is now safe." + if safe_to_retry + else str(exc) + ) + ), + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + details={ + "outcome_unknown": outcome_unknown, + "recovery": { + **previous_recovery, + "operation_id": exc.operation_id, + "status": recovery_status, + "requires_attention": outcome_unknown, + **( + {"next_call_number": call_number + 1} + if safe_to_retry + else {} + ), + }, + }, + ) + return True + except WorkflowRecoveryError as exc: + _set_action_handoff( + session, + instance=instance, + step=step, + state="blocked", + message=str(exc), + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + details={"recovery_unavailable": True}, + ) + return True + if action_recovery.replayed: + session.expire_all() + return True + action_recovery.checkpoint_dispatch( + session, + instance=instance, + step=step, + action_key=definition.action_key, + capability_name=capability_name, + ) try: result = provider.execute_action( session, @@ -1087,38 +1408,85 @@ def _execute_capability_step( instance.id, step.id, ) + session.rollback() + if action_recovery.mode.value == "atomic": + _set_action_handoff( + session, + instance=instance, + step=step, + state="retryable", + message=( + "The atomic module action failed and its database changes " + "were rolled back." + ), + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + details={ + "preview": preview_payload, + "error_type": type(exc).__name__, + "recovery": { + "operation_id": action_recovery.operation_id, + "mode": action_recovery.mode.value, + "status": "failed", + "call_number": action_recovery.call_number, + "next_call_number": action_recovery.call_number + 1, + "requires_attention": False, + }, + }, + ) + action_recovery.commit_definitive_failure( + session, + summary="The atomic provider call failed before commit", + error_type=type(exc).__name__, + ) + return True _set_action_handoff( session, instance=instance, step=step, - state="quarantined", + state="outcome_unknown", message=( "The module action outcome is unknown. Inspect the provider " - "before retrying to avoid a duplicate effect." + "and record evidence before continuing or retrying." ), action_key=definition.action_key, - capability_name=str(node.config.get("capability") or ""), + capability_name=capability_name, registry=registry, details={ "preview": preview_payload, "error_type": type(exc).__name__, "outcome_unknown": True, + "recovery": { + "operation_id": action_recovery.operation_id, + "mode": action_recovery.mode.value, + "status": "outcome_unknown", + "call_number": action_recovery.call_number, + "requires_attention": True, + }, }, ) + action_recovery.commit_unknown( + session, + error_type=type(exc).__name__, + message=( + "Inspect the provider by stable idempotency key before any retry" + ), + ) return True if not isinstance(result, ActionExecutionResult): - _set_action_handoff( + return _handle_invalid_action_result( session, instance=instance, step=step, - state="quarantined", message="The action provider returned an invalid execution result.", - action_key=definition.action_key, - capability_name=str(node.config.get("capability") or ""), + definition=definition, + capability_name=capability_name, registry=registry, - details={"preview": preview_payload}, + preview_payload=preview_payload, + action_recovery=action_recovery, + error_type="InvalidActionExecutionResult", ) - return True allowed_states = { "pending", "running", @@ -1138,29 +1506,34 @@ def _execute_capability_step( } ) if result.state not in allowed_states or unknown_effects: - _set_action_handoff( + return _handle_invalid_action_result( session, instance=instance, step=step, - state="quarantined", message=( "The action provider returned an unsupported state." if result.state not in allowed_states else "The action provider reported unannounced effects." ), - action_key=definition.action_key, - capability_name=str(node.config.get("capability") or ""), + definition=definition, + capability_name=capability_name, registry=registry, - details={ + preview_payload=preview_payload, + action_recovery=action_recovery, + error_type=( + "UnsupportedActionState" + if result.state not in allowed_states + else "UnannouncedActionEffect" + ), + extra_details={ "state": result.state, "unknown_effects": unknown_effects, }, ) - return True result_payload = _action_result_payload(result) step.output_ = { "action_key": definition.action_key, - "capability": str(node.config.get("capability") or ""), + "capability": capability_name, "idempotency_key": request.idempotency_key, "preview": preview_payload, "execution": result_payload, @@ -1177,13 +1550,29 @@ def _execute_capability_step( or f"Module action is {result.state.replace('_', ' ')}." ), action_key=definition.action_key, - capability_name=str(node.config.get("capability") or ""), + capability_name=capability_name, registry=registry, details={ "preview": preview_payload, "execution": result_payload, + "recovery": { + "operation_id": action_recovery.operation_id, + "mode": action_recovery.mode.value, + "status": "succeeded", + "call_number": action_recovery.call_number, + "next_call_number": action_recovery.call_number + 1, + "requires_attention": False, + }, }, ) + action_recovery.commit_conclusive_result( + session, + provider_state=result.state, + result_sha256=canonical_sha256(result_payload), + observed_effects_sha256=canonical_sha256( + result_payload["observed_effects"] + ), + ) return True _record_event( session, @@ -1193,7 +1582,7 @@ def _execute_capability_step( actor_id=actor_id, payload={ "action_key": definition.action_key, - "capability": str(node.config.get("capability") or ""), + "capability": capability_name, "idempotency_key": request.idempotency_key, "observed_effects": result_payload["observed_effects"], "audit_event_refs": result_payload["audit_event_refs"], @@ -1211,6 +1600,14 @@ def _execute_capability_step( output=dict(step.output_), actor_id=actor_id, ) + action_recovery.commit_conclusive_result( + session, + provider_state=result.state, + result_sha256=canonical_sha256(result_payload), + observed_effects_sha256=canonical_sha256( + result_payload["observed_effects"] + ), + ) _drive_instance( session, instance=instance, @@ -1399,6 +1796,83 @@ def _action_result_payload( } +def _handle_invalid_action_result( + session: Session, + *, + instance: WorkflowInstance, + step: WorkflowInstanceStep, + message: str, + definition: ActionDefinition, + capability_name: str, + registry: object | None, + preview_payload: Mapping[str, object], + action_recovery: WorkflowActionRecovery, + error_type: str, + extra_details: Mapping[str, object] | None = None, +) -> bool: + session.rollback() + recovery_details = { + "operation_id": action_recovery.operation_id, + "mode": action_recovery.mode.value, + "call_number": action_recovery.call_number, + } + if action_recovery.mode.value == "atomic": + _set_action_handoff( + session, + instance=instance, + step=step, + state="retryable", + message=f"{message} Atomic database changes were rolled back.", + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + details={ + "preview": dict(preview_payload), + "error_type": error_type, + **dict(extra_details or {}), + "recovery": { + **recovery_details, + "status": "failed", + "next_call_number": action_recovery.call_number + 1, + "requires_attention": False, + }, + }, + ) + action_recovery.commit_definitive_failure( + session, + summary=message, + error_type=error_type, + ) + return True + _set_action_handoff( + session, + instance=instance, + step=step, + state="outcome_unknown", + message=f"{message} Reconcile the provider before retrying.", + action_key=definition.action_key, + capability_name=capability_name, + registry=registry, + details={ + "preview": dict(preview_payload), + "error_type": error_type, + "outcome_unknown": True, + **dict(extra_details or {}), + "recovery": { + **recovery_details, + "status": "outcome_unknown", + "requires_attention": True, + }, + }, + ) + action_recovery.commit_unknown( + session, + error_type=error_type, + message="Inspect the provider by stable idempotency key before any retry", + ) + return True + + def _set_action_handoff( session: Session, *, @@ -1411,10 +1885,20 @@ def _set_action_handoff( registry: object | None, details: Mapping[str, object] | None = None, ) -> None: - allowed_actions = ( - ["cancel"] if state in {"pending", "running"} else ["retry", "reject", "cancel"] - ) previous = dict(step.handoff) + if state in {"pending", "running"}: + allowed_actions = ["cancel"] + elif state in {"outcome_unknown", "recovery_required"}: + allowed_actions = ["confirm_effect", "confirm_absent", "cancel"] + elif state == "compensation_required": + allowed_actions = ["reject", "cancel"] + else: + allowed_actions = ["retry", "reject", "cancel"] + details_payload = dict(details or {}) + if "recovery" not in details_payload and isinstance( + previous.get("recovery"), Mapping + ): + details_payload["recovery"] = dict(previous["recovery"]) step.status = "waiting" step.error = message if state not in {"pending", "running"} else None step.handoff = { @@ -1425,7 +1909,7 @@ def _set_action_handoff( "capability": capability_name, "allowed_actions": allowed_actions, "suggested_port": "failure", - **dict(details or {}), + **details_payload, } instance.status = "waiting" instance.error = step.error diff --git a/src/govoplan_workflow_engine/backend/manifest.py b/src/govoplan_workflow_engine/backend/manifest.py index 6e2d92b..715e34d 100644 --- a/src/govoplan_workflow_engine/backend/manifest.py +++ b/src/govoplan_workflow_engine/backend/manifest.py @@ -391,6 +391,30 @@ manifest = ModuleManifest( related_modules=("audit", "policy", "views"), order=77, ), + DocumentationTopic( + id="workflow.runtime-recovery", + title="Workflow runtime recovery", + summary=( + "Fenced module actions, timers, and evidence-based recovery " + "for unknown provider outcomes." + ), + body=( + "Workflow Engine records a Core recovery operation before every " + "consequential module-action dispatch and commits a conclusive " + "provider result with the local Workflow projection. A timeout or " + "lost acknowledgement after non-atomic dispatch becomes an unknown " + "outcome and disables Retry. Operators inspect the provider, record " + "evidence, and choose Effect confirmed to continue without replay or " + "Effect absent to enable a deliberate retry. Instance workers, " + "trigger deliveries, and timers use distributed fences. Linked " + "Dataflow recovery remains blocked until its result is conclusive." + ), + layer="available", + documentation_types=("admin", "user"), + audience=("operator", "module_admin", "power_user"), + related_modules=("core", "dataflow", "audit"), + order=78, + ), ), architecture=declared_module_architecture( layer="human_work_procedure", @@ -415,7 +439,7 @@ manifest = ModuleManifest( "notification", "dataflow run", ), - recovery_docs=("docs/CONCEPT.md",), + recovery_docs=("docs/CONCEPT.md", "docs/DURABLE_RUNTIME_RECOVERY.md"), security_docs=("docs/CONCEPT.md",), operations_docs=("README.md",), ), diff --git a/src/govoplan_workflow_engine/backend/recovery.py b/src/govoplan_workflow_engine/backend/recovery.py new file mode 100644 index 0000000..ecba120 --- /dev/null +++ b/src/govoplan_workflow_engine/backend/recovery.py @@ -0,0 +1,535 @@ +from __future__ import annotations + +from dataclasses import dataclass +import hashlib +import json +from typing import Mapping + +from sqlalchemy.orm import Session, sessionmaker + +from govoplan_core.core.automation import ActionDefinition +from govoplan_core.core.recovery import ( + RecoveryGuaranteeError, + RecoveryMode, + RecoveryOperation, + RecoveryPlan, + RecoveryStatus, +) +from govoplan_core.core.recovery_runtime import ( + DurableRecoveryOperation, + RecoveryOperationBusy, + RecoveryOperationStateConflict, + begin_durable_recovery_operation, + claim_durable_recovery_operation, +) +from govoplan_core.core.runtime_coordination import ( + LeaseClaim, + acquire_lease, + assert_lease_fence, + process_runtime_identity, + release_lease, +) +from govoplan_workflow_engine.backend.db.models import ( + WorkflowDefinitionRevision, + WorkflowInstance, + WorkflowInstanceStep, +) + + +class WorkflowRecoveryError(RuntimeError): + pass + + +class WorkflowRecoveryBusy(WorkflowRecoveryError): + pass + + +class WorkflowRecoveryConflict(WorkflowRecoveryError): + def __init__(self, operation_id: str, status: str) -> None: + self.operation_id = operation_id + self.status = status + super().__init__( + f"Workflow action recovery is {status}; reconcile it first" + ) + + +@dataclass(frozen=True, slots=True) +class WorkflowRecoveryDeclaration: + operation_type: str + mode: RecoveryMode | None + boundaries: tuple[Mapping[str, object], ...] + verification: tuple[str, ...] + + +WORKFLOW_RECOVERY_OPERATIONS = ( + WorkflowRecoveryDeclaration( + operation_type="instance.state-transition", + mode=RecoveryMode.ATOMIC, + boundaries=( + { + "name": "worker-claim", + "classification": "fenced", + "resume": "a current process identity and distributed fence are required", + }, + { + "name": "instance-transition", + "classification": "atomic", + "resume": "the pinned revision, step, event, and instance state commit together", + }, + ), + verification=( + "verify the pinned definition revision and current step", + "verify the process still owns the distributed fence before commit", + ), + ), + WorkflowRecoveryDeclaration( + operation_type="activity.external-effect", + mode=None, + boundaries=( + { + "name": "action-preview", + "classification": "recomputable", + "resume": "revalidate capability, authority, input, and preview", + }, + { + "name": "provider-dispatch", + "classification": "outcome-unknown", + "resume": "verify by stable provider idempotency key before retry", + }, + { + "name": "workflow-projection", + "classification": "verified", + "resume": "commit the local projection with terminal recovery evidence", + }, + ), + verification=( + "verify the pinned definition and canonical action request hashes", + "verify the provider result and every announced effect", + "reconcile an uncertain provider outcome before continuation", + ), + ), +) + + +def workflow_session_factory(session: Session) -> sessionmaker[Session]: + bind = session.get_bind() + if bind is None: + raise WorkflowRecoveryError("Workflow recovery requires a database bind") + return sessionmaker(bind=bind, expire_on_commit=False) + + +def canonical_sha256(value: object) -> str: + encoded = json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=True, + default=str, + ).encode("utf-8") + return hashlib.sha256(encoded).hexdigest() + + +@dataclass(slots=True) +class WorkflowActionRecovery: + operation: DurableRecoveryOperation | None + operation_id: str + mode: RecoveryMode + call_number: int + replayed: bool + + def checkpoint_dispatch( + self, + session: Session, + *, + instance: WorkflowInstance, + step: WorkflowInstanceStep, + action_key: str, + capability_name: str, + ) -> None: + if self.operation is None: + return + previous = dict(step.handoff or {}) + step.handoff = { + **previous, + "kind": "module_action", + "state": "running", + "action_key": action_key, + "capability": capability_name, + "allowed_actions": ["cancel"], + "recovery": { + "operation_id": self.operation_id, + "mode": self.mode.value, + "status": RecoveryStatus.RUNNING.value, + "call_number": self.call_number, + "boundary": "provider-dispatch", + "requires_attention": False, + }, + } + step.status = "running" + instance.status = "running" + try: + session.commit() + self.operation.checkpoint( + kind="provider-dispatch", + summary="Provider dispatch crossed the recoverable effect boundary", + evidence={ + "effect_started": self.mode != RecoveryMode.ATOMIC, + "action_key": action_key, + "capability": capability_name, + "call_number": self.call_number, + }, + ) + except Exception as exc: + session.rollback() + raise WorkflowRecoveryError( + "Workflow action dispatch evidence could not be persisted" + ) from exc + + def commit_conclusive_result( + self, + session: Session, + *, + provider_state: str, + result_sha256: str, + observed_effects_sha256: str, + ) -> None: + if self.operation is None: + return + evidence = { + "verified": True, + "checks": { + "provider_state": provider_state, + "result_sha256": result_sha256, + "observed_effects_sha256": observed_effects_sha256, + "call_number": self.call_number, + }, + } + if self.mode == RecoveryMode.ATOMIC: + self.operation.commit_atomic_success(session, evidence=evidence) + else: + self.operation.commit_verified_success(session, evidence=evidence) + + def commit_unknown( + self, + session: Session, + *, + error_type: str, + message: str, + ) -> None: + if self.operation is None: + raise WorkflowRecoveryError( + "A replayed action cannot acquire an unknown outcome" + ) + try: + session.commit() + self.operation.unresolved( + status=RecoveryStatus.OUTCOME_UNKNOWN, + summary="The provider acknowledgement was not conclusive", + evidence={ + "effect_started": self.mode != RecoveryMode.ATOMIC, + "error_type": error_type, + "call_number": self.call_number, + }, + failure_summary=message, + ) + except Exception as exc: + session.rollback() + raise WorkflowRecoveryError( + "Workflow action uncertainty could not be recorded" + ) from exc + + def commit_definitive_failure( + self, + session: Session, + *, + summary: str, + error_type: str, + ) -> None: + if self.operation is None: + return + evidence = { + "verified": True, + "checks": { + "effect_started": False, + "error_type": error_type, + "call_number": self.call_number, + }, + } + if self.mode == RecoveryMode.ATOMIC: + self.operation.commit_atomic_failure( + session, + summary=summary, + evidence=evidence, + ) + else: + self.operation.fail(summary=summary, evidence=evidence) + session.commit() + + +def begin_workflow_action_recovery( + session: Session, + *, + instance: WorkflowInstance, + step: WorkflowInstanceStep, + revision: WorkflowDefinitionRevision, + definition: ActionDefinition, + capability_name: str, + request_idempotency_key: str, + action_input: Mapping[str, object], + preview_payload: Mapping[str, object], + backup_reference: str | None = None, + approval_reference: str | None = None, +) -> WorkflowActionRecovery: + previous_recovery = ( + step.handoff.get("recovery") + if isinstance(step.handoff, Mapping) + else None + ) + call_number = 1 + if isinstance(previous_recovery, Mapping): + call_number = max( + 1, + int( + previous_recovery.get("next_call_number") + or previous_recovery.get("call_number") + or 1 + ), + ) + mode = RecoveryMode(definition.recovery_mode) + plan = _action_recovery_plan( + definition, + mode=mode, + backup_reference=backup_reference, + approval_reference=approval_reference, + ) + request_sha256 = canonical_sha256(dict(action_input)) + preview_sha256 = canonical_sha256(dict(preview_payload)) + access_context_sha256 = canonical_sha256(instance.authorization_) + action_contract_sha256 = canonical_sha256( + { + "action_key": definition.action_key, + "owner_module": definition.owner_module, + "contract_version": definition.contract_version, + "recovery_mode": definition.recovery_mode, + "recovery_verification": list(definition.recovery_verification), + "expected_effect_keys": list(definition.expected_effect_keys), + } + ) + session.commit() + try: + started = begin_durable_recovery_operation( + workflow_session_factory(session), + identity=process_runtime_identity(), + module_id="workflow_engine", + operation_type="activity.external-effect", + idempotency_key=( + f"workflow-step:{step.id}:action-call:{call_number}" + ), + request={ + "tenant_id": instance.tenant_id, + "instance_id": instance.id, + "step_id": step.id, + "node_id": step.node_id, + "definition_revision_id": revision.id, + "definition_hash": revision.content_hash, + "action_key": definition.action_key, + "capability": capability_name, + "action_contract_sha256": action_contract_sha256, + "action_input_sha256": request_sha256, + "preview_sha256": preview_sha256, + "provider_idempotency_sha256": canonical_sha256( + request_idempotency_key + ), + "call_number": call_number, + }, + recovery_plan=plan, + precondition_evidence={ + "definition_hash": revision.content_hash, + "action_contract_sha256": action_contract_sha256, + "action_input_sha256": request_sha256, + "preview_sha256": preview_sha256, + "access_context_sha256": access_context_sha256, + "call_number": call_number, + }, + lease_resource_key=f"workflow:step:{step.id}", + lease_ttl_seconds=300, + resource_type="workflow_instance_step", + resource_id=step.id, + metadata={ + "resources": ["postgresql", "module-provider"], + "workflow_instance_id": instance.id, + "action_key": definition.action_key, + "capability": capability_name, + }, + ) + except RecoveryOperationBusy as exc: + raise WorkflowRecoveryBusy( + "Another runtime owns this Workflow action step" + ) from exc + except RecoveryOperationStateConflict as exc: + raise WorkflowRecoveryConflict(exc.operation_id, exc.status) from exc + except (RecoveryGuaranteeError, RuntimeError, ValueError) as exc: + raise WorkflowRecoveryError( + "The recovery ledger is unavailable; the module action did not start" + ) from exc + return WorkflowActionRecovery( + operation=started.operation, + operation_id=started.operation_id, + mode=mode, + call_number=call_number, + replayed=started.replayed, + ) + + +def claim_workflow_action_recovery( + session: Session, + *, + operation_id: str, +) -> DurableRecoveryOperation: + try: + return claim_durable_recovery_operation( + workflow_session_factory(session), + identity=process_runtime_identity(), + operation_id=operation_id, + lease_ttl_seconds=300, + ) + except (RecoveryGuaranteeError, RuntimeError) as exc: + raise WorkflowRecoveryError( + "The Workflow action outcome cannot currently be reconciled" + ) from exc + + +def reconcile_stale_workflow_action_recovery( + session: Session, + *, + operation_id: str, +) -> str: + try: + handle = claim_durable_recovery_operation( + workflow_session_factory(session), + identity=process_runtime_identity(), + operation_id=operation_id, + lease_ttl_seconds=300, + ) + except RecoveryOperationBusy as exc: + raise WorkflowRecoveryBusy( + "Another runtime still owns this Workflow action step" + ) from exc + except RecoveryOperationStateConflict as exc: + return exc.status + except (RecoveryGuaranteeError, RuntimeError) as exc: + raise WorkflowRecoveryError( + "The stale Workflow action fence could not be reconciled" + ) from exc + handle.release_unresolved() + session.expire_all() + operation = session.get(RecoveryOperation, operation_id) + if operation is None: + raise WorkflowRecoveryError( + "The reconciled Workflow recovery operation is unavailable" + ) + return str(operation.status) + + +def workflow_action_recovery_state( + session: Session, + *, + operation_id: str, +) -> dict[str, object] | None: + operation = session.get(RecoveryOperation, operation_id) + if operation is None or operation.module_id != "workflow_engine": + return None + return { + "operation_id": operation.id, + "mode": operation.mode, + "status": operation.status, + "requires_attention": operation.status + in { + RecoveryStatus.OUTCOME_UNKNOWN.value, + RecoveryStatus.RECOVERY_REQUIRED.value, + RecoveryStatus.MANUAL_INTERVENTION.value, + }, + "failure_summary": operation.failure_summary, + } + + +def acquire_workflow_state_fence( + session: Session, + *, + resource_key: str, + ttl_seconds: int = 120, +) -> LeaseClaim | None: + identity = process_runtime_identity() + return acquire_lease( + session, + installation_id=identity.installation_id, + resource_key=resource_key, + holder_node_id=identity.node_id, + holder_incarnation=identity.incarnation, + ttl_seconds=ttl_seconds, + metadata={"module_id": "workflow_engine", "kind": "state-transition"}, + ) + + +def release_workflow_state_fence(session: Session, claim: LeaseClaim) -> None: + assert_lease_fence(session, claim) + release_lease(session, claim) + + +def _action_recovery_plan( + definition: ActionDefinition, + *, + mode: RecoveryMode, + backup_reference: str | None, + approval_reference: str | None, +) -> RecoveryPlan: + if mode == RecoveryMode.SNAPSHOT_RESTORE and not backup_reference: + raise WorkflowRecoveryError( + "Snapshot-restore actions require a pinned recovery backup reference" + ) + if mode == RecoveryMode.IRREVERSIBLE and not approval_reference: + raise WorkflowRecoveryError( + "Irreversible actions require a pinned approval reference" + ) + return RecoveryPlan( + mode=mode, + preconditions=( + "the pinned Workflow revision and action contract hashes are present", + "the current authority and provider capability were revalidated", + "the action step owns a distributed execution fence", + ), + compensation_steps=( + "invoke the provider-declared compensation action", + "verify the domain invariant after compensation", + ) + if mode == RecoveryMode.COMPENSATION + else (), + forward_recovery_steps=( + "inspect the provider by stable idempotency key", + "record whether the announced effect occurred", + "continue or retry only after the outcome is proven", + ) + if mode == RecoveryMode.FORWARD_RECOVERY + else (), + verification_steps=tuple(definition.recovery_verification), + backup_reference=backup_reference, + approval_reference=approval_reference, + ) + + +__all__ = [ + "WORKFLOW_RECOVERY_OPERATIONS", + "WorkflowActionRecovery", + "WorkflowRecoveryDeclaration", + "WorkflowRecoveryBusy", + "WorkflowRecoveryConflict", + "WorkflowRecoveryError", + "acquire_workflow_state_fence", + "begin_workflow_action_recovery", + "canonical_sha256", + "claim_workflow_action_recovery", + "release_workflow_state_fence", + "reconcile_stale_workflow_action_recovery", + "workflow_action_recovery_state", + "workflow_session_factory", +] diff --git a/src/govoplan_workflow_engine/backend/schemas.py b/src/govoplan_workflow_engine/backend/schemas.py index f4cb3e3..b90db85 100644 --- a/src/govoplan_workflow_engine/backend/schemas.py +++ b/src/govoplan_workflow_engine/backend/schemas.py @@ -552,6 +552,8 @@ class WorkflowStepActionRequest(BaseModel): "reject", "resume", "retry", + "confirm_effect", + "confirm_absent", "cancel", ] output: dict[str, Any] = Field(default_factory=dict) diff --git a/src/govoplan_workflow_engine/backend/triggers.py b/src/govoplan_workflow_engine/backend/triggers.py index 96eea3f..70b2c10 100644 --- a/src/govoplan_workflow_engine/backend/triggers.py +++ b/src/govoplan_workflow_engine/backend/triggers.py @@ -36,6 +36,10 @@ from govoplan_workflow_engine.backend.service import ( WorkflowConflictError, get_definition_revision, ) +from govoplan_workflow_engine.backend.recovery import ( + acquire_workflow_state_fence, + release_workflow_state_fence, +) _EVENT_TYPE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,119}$") @@ -367,11 +371,20 @@ def dispatch_due_work( ) started = blocked = failed = 0 for delivery in deliveries: - outcome = _dispatch_delivery( + fence = acquire_workflow_state_fence( session, - delivery=delivery, - registry=registry, + resource_key=f"workflow:trigger-delivery:{delivery.id}", ) + if fence is None: + continue + try: + outcome = _dispatch_delivery( + session, + delivery=delivery, + registry=registry, + ) + finally: + release_workflow_state_fence(session, fence) if outcome == "started": started += 1 elif outcome == "blocked": @@ -622,6 +635,13 @@ def _dispatch_waits( ) for state in states: + fence = acquire_workflow_state_fence( + session, + resource_key=f"workflow:wait:{state.id}", + ) + if fence is None: + skipped += 1 + continue instance = session.get(WorkflowInstance, state.instance_id) step = session.get(WorkflowInstanceStep, state.step_id) if ( @@ -635,6 +655,7 @@ def _dispatch_waits( state.error = "The waiting Workflow step is no longer current." state.revision += 1 skipped += 1 + release_workflow_state_fence(session, fence) continue principal = _resolve_instance_principal( session, @@ -644,6 +665,7 @@ def _dispatch_waits( if principal is None: state.error = "Current trigger authority could not be resolved." skipped += 1 + release_workflow_state_fence(session, fence) continue revision = session.get( WorkflowDefinitionRevision, @@ -654,6 +676,7 @@ def _dispatch_waits( state.error = "The pinned Workflow revision is unavailable." state.revision += 1 skipped += 1 + release_workflow_state_fence(session, fence) continue graph = _runtime_graph(revision) triggered = state.status == "triggered" @@ -699,6 +722,7 @@ def _dispatch_waits( resumed += 1 else: timed_out += 1 + release_workflow_state_fence(session, fence) return { "waits_selected": len(states), "waits_resumed": resumed, diff --git a/tests/test_instance_service.py b/tests/test_instance_service.py index 1cb425d..eb06900 100644 --- a/tests/test_instance_service.py +++ b/tests/test_instance_service.py @@ -4,7 +4,7 @@ from dataclasses import replace from datetime import UTC, datetime, timedelta import unittest -from sqlalchemy import create_engine +from sqlalchemy import create_engine, func, select from sqlalchemy.orm import Session, sessionmaker from govoplan_core.auth import ApiPrincipal @@ -32,7 +32,19 @@ from govoplan_core.core.institutional import ( 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.db.base import Base +from govoplan_core.db.base import utcnow from govoplan_workflow_engine.backend.db.models import ( WorkflowDefinition, WorkflowDefinitionRevision, @@ -45,12 +57,17 @@ from govoplan_workflow_engine.backend.db.models import ( ) 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, @@ -97,6 +114,17 @@ def principal() -> ApiPrincipal: ) +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=[ @@ -278,6 +306,20 @@ class FakeDataflowLifecycle: 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: @@ -337,6 +379,8 @@ class FakeActionProvider: 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, @@ -396,6 +440,9 @@ class WorkflowInstanceServiceTests(unittest.TestCase): Base.metadata.create_all( self.engine, tables=[ + DistributedLease.__table__, + RecoveryOperation.__table__, + RecoveryCheckpoint.__table__, WorkflowDefinition.__table__, WorkflowDefinitionRevision.__table__, WorkflowInstance.__table__, @@ -408,6 +455,7 @@ class WorkflowInstanceServiceTests(unittest.TestCase): ) 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( @@ -432,6 +480,7 @@ class WorkflowInstanceServiceTests(unittest.TestCase): self.session.commit() def tearDown(self) -> None: + bind_process_runtime_identity(None) self.session.close() Base.metadata.drop_all( self.engine, @@ -444,6 +493,9 @@ class WorkflowInstanceServiceTests(unittest.TestCase): WorkflowInstance.__table__, WorkflowDefinitionRevision.__table__, WorkflowDefinition.__table__, + RecoveryCheckpoint.__table__, + RecoveryOperation.__table__, + DistributedLease.__table__, ], ) self.engine.dispose() @@ -687,6 +739,51 @@ class WorkflowInstanceServiceTests(unittest.TestCase): "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") @@ -739,6 +836,383 @@ class WorkflowInstanceServiceTests(unittest.TestCase): 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: diff --git a/tests/test_triggers.py b/tests/test_triggers.py index b61cb5e..1633645 100644 --- a/tests/test_triggers.py +++ b/tests/test_triggers.py @@ -13,6 +13,12 @@ from govoplan_core.core.access import ( ) from govoplan_core.core.automation import AutomationPrincipalResolution from govoplan_core.core.events import EventTenantRef, PlatformEvent +from govoplan_core.core.recovery import RecoveryCheckpoint, RecoveryOperation +from govoplan_core.core.runtime_coordination import ( + DistributedLease, + RuntimeIdentity, + bind_process_runtime_identity, +) from govoplan_core.db.base import Base from govoplan_workflow_engine.backend.db.models import ( WorkflowDefinition, @@ -55,6 +61,17 @@ def principal() -> ApiPrincipal: ) +def runtime_identity() -> RuntimeIdentity: + return RuntimeIdentity( + installation_id="workflow-trigger-tests", + node_id="workflow-trigger-worker", + incarnation="workflow-trigger-incarnation", + role="worker", + software_version="test", + composition_hash="d" * 64, + ) + + class AutomationProvider: def resolve_automation_principal(self, _session, *, request): return AutomationPrincipalResolution( @@ -110,6 +127,9 @@ class WorkflowTriggerTests(unittest.TestCase): def setUp(self) -> None: self.engine = create_engine("sqlite:///:memory:") self.tables = [ + DistributedLease.__table__, + RecoveryOperation.__table__, + RecoveryCheckpoint.__table__, WorkflowDefinition.__table__, WorkflowDefinitionRevision.__table__, WorkflowInstance.__table__, @@ -122,9 +142,11 @@ class WorkflowTriggerTests(unittest.TestCase): Base.metadata.create_all(self.engine, tables=self.tables) self.Session = sessionmaker(bind=self.engine) self.session: Session = self.Session() + bind_process_runtime_identity(runtime_identity()) self.registry = Registry() def tearDown(self) -> None: + bind_process_runtime_identity(None) self.session.close() Base.metadata.drop_all(self.engine, tables=list(reversed(self.tables))) self.engine.dispose()