diff --git a/README.md b/README.md index 5e61a60..6826b9d 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,15 @@ records the Service and binding in trusted instance context, and safely replays the same launch after an ambiguous response. Portal never accesses Workflow Engine tables. +Active definitions reconcile durable API, one-time/interval schedule, +platform-event, and parent-workflow trigger registrations. A shared worker +claims due deliveries with scale-out-safe locking, rechecks the exact active +revision and automation authority, and starts instances idempotently. Duration, +deadline, and platform-event waits are persistent runtime state rather than +human handoffs; event filters and variable mappings are bounded JSON +expressions and never executable code. Cron remains an optional governed +scheduler-adapter concern. + ## Checks ```bash diff --git a/docs/CONCEPT.md b/docs/CONCEPT.md index 798f768..9e0de98 100644 --- a/docs/CONCEPT.md +++ b/docs/CONCEPT.md @@ -116,18 +116,22 @@ The first executable slice now provides: - an operator dialog for starting, inspecting, and advancing instances - an owner-side Service launcher that starts an authorized active definition, retains the exact Service/binding provenance, and safely replays Portal calls +- activation-bound API, one-time/interval schedule, platform-event, and + parent-workflow trigger registrations with exact-revision dispatch +- durable duration/deadline/event wait subscriptions and scale-out-safe claims +- a separate transactional platform-event consumer with bounded JSON filters + and variable mappings, idempotent delivery, and current-authority rechecks -The next execution slices should provide: +The next execution depth should provide: - static workflow definition registration from configuration packages -- event, API, schedule, and parent-workflow start dispatchers - 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 -- event emission and audit integration +- cron/calendar scheduling through a governed scheduler adapter ## Permissions @@ -165,11 +169,13 @@ Current tables: - `workflow_instances` - `workflow_instance_steps` - `workflow_instance_events` +- `workflow_triggers` +- `workflow_trigger_deliveries` +- `workflow_wait_states` -Future generic action execution and timers may add: +Future generic action execution may add: - `workflow_command_records` -- `workflow_timers` Definitions should be immutable by version after activation. Instances should reference the exact version used at start. @@ -196,8 +202,11 @@ Minimum tests: - events are emitted for start/transition/completion - configuration package can install a simple workflow definition -## Open Decisions +## Bounded Decisions -- Whether long-running timers use Celery beat, a module scheduler, or an ops - scheduler abstraction. -- How workflow variables are redacted and retained. +- Native scheduling deliberately supports one-time and minimum-60-second + interval triggers. Cron/calendar semantics belong to a future governed + scheduler adapter rather than an unbounded expression evaluator in Engine. +- Platform events contain sanitized lifecycle context, not raw instance + variables. Definition mappings can select only bounded event fields; broader + variable retention/redaction remains policy-controlled product depth. diff --git a/src/govoplan_workflow_engine/backend/db/models.py b/src/govoplan_workflow_engine/backend/db/models.py index 6b6e996..ce03b32 100644 --- a/src/govoplan_workflow_engine/backend/db/models.py +++ b/src/govoplan_workflow_engine/backend/db/models.py @@ -382,9 +382,7 @@ class WorkflowInstance(Base, TimestampMixin): index=True, ) - definition: Mapped[WorkflowDefinition] = relationship( - back_populates="instances" - ) + definition: Mapped[WorkflowDefinition] = relationship(back_populates="instances") steps: Mapped[list["WorkflowInstanceStep"]] = relationship( back_populates="instance", cascade="all, delete-orphan", @@ -518,11 +516,185 @@ class WorkflowInstanceEvent(Base): instance: Mapped[WorkflowInstance] = relationship(back_populates="events") +class WorkflowTrigger(Base, TimestampMixin): + __tablename__ = "workflow_triggers" + __table_args__ = ( + UniqueConstraint( + "definition_id", + "node_id", + name="uq_workflow_trigger_definition_node", + ), + Index("ix_workflow_triggers_due", "status", "next_fire_at"), + Index( + "ix_workflow_triggers_event", + "tenant_id", + "status", + "event_type", + ), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + definition_id: Mapped[str] = mapped_column( + ForeignKey("workflow_definitions.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + definition_revision_id: Mapped[str] = mapped_column( + ForeignKey("workflow_definition_revisions.id", ondelete="RESTRICT"), + nullable=False, + index=True, + ) + node_id: Mapped[str] = mapped_column(String(120), nullable=False) + kind: Mapped[str] = mapped_column(String(30), nullable=False, index=True) + status: Mapped[str] = mapped_column( + String(30), default="disabled", nullable=False, index=True + ) + config_: Mapped[dict[str, Any]] = mapped_column( + "config", JSON, default=dict, nullable=False + ) + event_type: Mapped[str | None] = mapped_column( + String(120), nullable=True, index=True + ) + next_fire_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + last_fire_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + last_status: Mapped[str | None] = mapped_column(String(30), nullable=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + authorization_subject_kind: Mapped[str] = mapped_column(String(30), nullable=False) + authorization_account_id: Mapped[str | None] = mapped_column( + String(36), nullable=True + ) + authorization_membership_id: Mapped[str | None] = mapped_column( + String(36), nullable=True + ) + authorization_service_account_id: Mapped[str | None] = mapped_column( + String(36), nullable=True + ) + authorization_ref: Mapped[str] = mapped_column( + String(255), nullable=False, unique=True + ) + grant_scopes: Mapped[list[str]] = mapped_column(JSON, default=list, nullable=False) + created_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + updated_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + + +class WorkflowTriggerDelivery(Base, TimestampMixin): + __tablename__ = "workflow_trigger_deliveries" + __table_args__ = ( + UniqueConstraint( + "trigger_id", + "source_key", + name="uq_workflow_trigger_delivery_source", + ), + Index( + "ix_workflow_trigger_deliveries_queue", + "status", + "created_at", + ), + Index( + "ix_workflow_trigger_deliveries_tenant_trigger", + "tenant_id", + "trigger_id", + ), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + trigger_id: Mapped[str] = mapped_column( + ForeignKey("workflow_triggers.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + definition_id: Mapped[str] = mapped_column( + ForeignKey("workflow_definitions.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + definition_revision_id: Mapped[str] = mapped_column( + ForeignKey("workflow_definition_revisions.id", ondelete="RESTRICT"), + nullable=False, + index=True, + ) + source_key: Mapped[str] = mapped_column(String(255), nullable=False) + invocation_kind: Mapped[str] = mapped_column(String(30), nullable=False) + status: Mapped[str] = mapped_column( + String(30), default="queued", nullable=False, index=True + ) + scheduled_for: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + event_: Mapped[dict[str, Any] | None] = mapped_column("event", JSON, nullable=True) + instance_id: Mapped[str | None] = mapped_column( + ForeignKey("workflow_instances.id", ondelete="SET NULL"), + nullable=True, + index=True, + ) + attempts: Mapped[int] = mapped_column(Integer, default=0, nullable=False) + error: Mapped[str | None] = mapped_column(Text, nullable=True) + + +class WorkflowWaitState(Base, TimestampMixin): + __tablename__ = "workflow_wait_states" + __table_args__ = ( + UniqueConstraint("step_id", name="uq_workflow_wait_state_step"), + Index("ix_workflow_wait_states_due", "status", "due_at"), + Index( + "ix_workflow_wait_states_event", + "tenant_id", + "status", + "event_type", + ), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + instance_id: Mapped[str] = mapped_column( + ForeignKey("workflow_instances.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + step_id: Mapped[str] = mapped_column( + ForeignKey("workflow_instance_steps.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + mode: Mapped[str] = mapped_column(String(30), nullable=False, index=True) + status: Mapped[str] = mapped_column( + String(30), default="waiting", nullable=False, index=True + ) + due_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + event_type: Mapped[str | None] = mapped_column( + String(120), nullable=True, index=True + ) + config_: Mapped[dict[str, Any]] = mapped_column( + "config", JSON, default=dict, nullable=False + ) + source_event_id: Mapped[str | None] = mapped_column( + String(128), nullable=True, index=True + ) + event_: Mapped[dict[str, Any] | None] = mapped_column("event", JSON, nullable=True) + error: Mapped[str | None] = mapped_column(Text, nullable=True) + revision: Mapped[int] = mapped_column(Integer, default=1, nullable=False) + + __all__ = [ "WorkflowDefinition", "WorkflowDefinitionRevision", "WorkflowInstance", "WorkflowInstanceEvent", "WorkflowInstanceStep", + "WorkflowTrigger", + "WorkflowTriggerDelivery", + "WorkflowWaitState", "new_uuid", ] diff --git a/src/govoplan_workflow_engine/backend/instance_service.py b/src/govoplan_workflow_engine/backend/instance_service.py index 0f4ae13..d715f10 100644 --- a/src/govoplan_workflow_engine/backend/instance_service.py +++ b/src/govoplan_workflow_engine/backend/instance_service.py @@ -29,6 +29,13 @@ from govoplan_core.core.notifications import ( NotificationDispatchRequest, notification_dispatch_provider, ) +from govoplan_core.core.events import ( + EventActorRef, + EventObjectRef, + EventTenantRef, + PlatformEvent, + emit_platform_event, +) from govoplan_core.db.base import utcnow from govoplan_workflow_engine.backend.db.models import ( WorkflowDefinition, @@ -89,9 +96,7 @@ def list_instances( .limit(max(1, min(int(limit), 200))) ) if definition_id: - statement = statement.where( - WorkflowInstance.definition_id == definition_id - ) + statement = statement.where(WorkflowInstance.definition_id == definition_id) return list(session.scalars(statement)) @@ -160,9 +165,7 @@ def start_instance( ) normalized_origin = _normalize_start_origin(start_origin) if normalized_origin != "user" and not definition.allow_automation: - raise WorkflowConflictError( - "This Workflow does not allow automated starts." - ) + raise WorkflowConflictError("This Workflow does not allow automated starts.") if revision.execution_mode == "guided" and normalized_origin != "user": raise WorkflowConflictError( "Guided workflows must be started by a user; use hybrid mode " @@ -189,8 +192,7 @@ def start_instance( or existing.start_origin != normalized_origin ): raise WorkflowConflictError( - "The Workflow idempotency key was already used with " - "different input." + "The Workflow idempotency key was already used with different input." ) return get_instance( session, @@ -357,9 +359,7 @@ def reconcile_instance( step.handoff = { **dict(step.handoff), "state": descriptor.status, - "progress_percent": int( - descriptor.metadata.get("progress_percent") or 0 - ), + "progress_percent": int(descriptor.metadata.get("progress_percent") or 0), "progress_phase": str( descriptor.metadata.get("progress_phase") or descriptor.status ), @@ -434,14 +434,11 @@ def resolve_step( instance.definition_revision_id, ) if revision is None: - raise WorkflowConflictError( - "Pinned Workflow revision no longer exists." - ) + raise WorkflowConflictError("Pinned Workflow revision no longer exists.") graph = _runtime_graph(revision) node = _node(graph, step.node_id) allowed_actions = { - str(action) - for action in step.handoff.get("allowed_actions") or () + str(action) for action in step.handoff.get("allowed_actions") or () } if payload.action not in allowed_actions: raise WorkflowConflictError( @@ -525,6 +522,10 @@ def resolve_step( "comment": payload.comment, "evidence": list(payload.evidence), } + if step.node_type == "workflow.wait": + from govoplan_workflow_engine.backend.triggers import resolve_wait_state + + resolve_wait_state(session, step_id=step.id, status="resumed") next_node_id = _complete_step( session, instance=instance, @@ -562,28 +563,33 @@ def cancel_instance( for_update=True, ) if instance.status in {"completed", "failed", "cancelled"}: - raise WorkflowConflictError( - f"Workflow instance is already {instance.status}." - ) + raise WorkflowConflictError(f"Workflow instance is already {instance.status}.") now = utcnow() instance.cancellation_requested_at = now step = _current_step(session, instance) - if step is not None and step.external_ref: - provider = dataflow_run_lifecycle(registry) - if provider is not None: - try: - provider.cancel_run( - session, - principal, - run_ref=step.external_ref, - ) - except ValueError as exc: - logger.info( - "Linked Dataflow run could not be cancelled for " - "Workflow instance %s: %s", - instance.id, - exc, - ) + if step is not None: + if step.external_ref: + provider = dataflow_run_lifecycle(registry) + if provider is not None: + try: + provider.cancel_run( + session, + principal, + run_ref=step.external_ref, + ) + except ValueError as exc: + logger.info( + "Linked Dataflow run could not be cancelled for " + "Workflow instance %s: %s", + instance.id, + exc, + ) + if step.node_type == "workflow.wait": + from govoplan_workflow_engine.backend.triggers import ( + resolve_wait_state, + ) + + resolve_wait_state(session, step_id=step.id, status="cancelled") step.status = "cancelled" step.finished_at = now step.completed_by = actor_id @@ -805,7 +811,6 @@ def _drive_instance( if node.type in { "workflow.activity", "workflow.review", - "workflow.wait", }: _set_human_handoff( session, @@ -815,6 +820,34 @@ def _drive_instance( registry=registry, ) return + if node.type == "workflow.wait": + from govoplan_workflow_engine.backend.triggers import ( + register_wait_state, + ) + + wait_state = register_wait_state( + session, + instance=instance, + step=step, + node=node, + ) + if wait_state is None: + _set_human_handoff( + session, + instance=instance, + step=step, + node=node, + registry=registry, + ) + else: + _set_automated_wait( + session, + instance=instance, + step=step, + node=node, + wait_state=wait_state, + ) + return if node.type == "workflow.end.completed": _complete_step( session, @@ -1096,9 +1129,7 @@ def _execute_capability_step( "manual_required", "compensation_required", } - announced_effects = { - item.effect_key for item in provider.effect_definitions() - } + announced_effects = {item.effect_key for item in provider.effect_definitions()} unknown_effects = sorted( { effect.effect_key @@ -1210,25 +1241,19 @@ def _capability_action_definition( f"Action capability {capability_name!r} is not available." ) definitions = [ - item - for item in provider.action_definitions() - if item.action_key == action_key + item for item in provider.action_definitions() if item.action_key == action_key ] if len(definitions) != 1: raise WorkflowConflictError( - f"Action {action_key!r} is not uniquely announced by " - f"{capability_name!r}." + f"Action {action_key!r} is not uniquely announced by {capability_name!r}." ) definition = definitions[0] missing_scopes = [ - scope - for scope in definition.required_scopes - if not has_scope(principal, scope) + scope for scope in definition.required_scopes if not has_scope(principal, scope) ] if missing_scopes: raise WorkflowConflictError( - "Module action requires scopes: " - + ", ".join(sorted(missing_scopes)) + "Module action requires scopes: " + ", ".join(sorted(missing_scopes)) ) missing_capabilities = [ capability @@ -1244,12 +1269,8 @@ def _capability_action_definition( "Module action requires capabilities: " + ", ".join(sorted(missing_capabilities)) ) - effect_keys = { - item.effect_key for item in provider.effect_definitions() - } - missing_effects = sorted( - set(definition.expected_effect_keys) - effect_keys - ) + effect_keys = {item.effect_key for item in provider.effect_definitions()} + missing_effects = sorted(set(definition.expected_effect_keys) - effect_keys) if missing_effects: raise WorkflowConflictError( "Action provider does not define its expected effects: " @@ -1265,9 +1286,7 @@ def _mapped_action_input( if raw_mapping is None or raw_mapping == "": return dict(context) if not isinstance(raw_mapping, Mapping): - raise WorkflowConflictError( - "Module-action input mapping must be an object." - ) + raise WorkflowConflictError("Module-action input mapping must be an object.") return { str(key): _resolve_action_value(value, context, depth=0) for key, value in raw_mapping.items() @@ -1282,9 +1301,7 @@ def _resolve_action_value( depth: int, ) -> object: if depth > 10: - raise WorkflowConflictError( - "Module-action input mapping is nested too deeply." - ) + raise WorkflowConflictError("Module-action input mapping is nested too deeply.") if isinstance(value, str) and value.startswith("$"): path = value[1:].lstrip(".") current: object = context @@ -1307,10 +1324,7 @@ def _resolve_action_value( for key, nested in value.items() } if isinstance(value, list): - return [ - _resolve_action_value(item, context, depth=depth + 1) - for item in value - ] + return [_resolve_action_value(item, context, depth=depth + 1) for item in value] return value @@ -1322,9 +1336,7 @@ def _action_idempotency_key( action_key: str, context: Mapping[str, object], ) -> str: - expression = str( - node.config.get("idempotency_key") or "workflow-step" - ).strip() + expression = str(node.config.get("idempotency_key") or "workflow-step").strip() if expression == "workflow-step": return step.idempotency_key resolved = _resolve_action_value(expression, context, depth=0) @@ -1345,8 +1357,7 @@ def _action_preview_payload(preview: object) -> dict[str, object]: "preview_ref": getattr(preview, "preview_ref", None), "blockers": list(getattr(preview, "blockers", ()) or ()), "policy_provenance": [ - dict(item) - for item in getattr(preview, "policy_provenance", ()) or () + dict(item) for item in getattr(preview, "policy_provenance", ()) or () ], "effects": [ { @@ -1380,9 +1391,7 @@ def _action_result_payload( ], "error": result.error, "retry_after": ( - result.retry_after.isoformat() - if result.retry_after is not None - else None + result.retry_after.isoformat() if result.retry_after is not None else None ), "manual_instructions": result.manual_instructions, "compensation_action_key": result.compensation_action_key, @@ -1403,9 +1412,7 @@ def _set_action_handoff( details: Mapping[str, object] | None = None, ) -> None: allowed_actions = ( - ["cancel"] - if state in {"pending", "running"} - else ["retry", "reject", "cancel"] + ["cancel"] if state in {"pending", "running"} else ["retry", "reject", "cancel"] ) previous = dict(step.handoff) step.status = "waiting" @@ -1422,10 +1429,7 @@ def _set_action_handoff( } instance.status = "waiting" instance.error = step.error - if ( - previous.get("state") != state - or previous.get("message") != message - ): + if previous.get("state") != state or previous.get("message") != message: _record_event( session, instance, @@ -1497,9 +1501,7 @@ def _start_dataflow_step( message="Dataflow steps require a pipeline and pinned revision.", ) return - target_ref = str( - node.config.get("publication_target_ref") or "" - ).strip() + target_ref = str(node.config.get("publication_target_ref") or "").strip() try: run = provider.start_run( session, @@ -1509,13 +1511,9 @@ def _start_dataflow_step( revision=revision, idempotency_key=step.idempotency_key, row_limit=row_limit, - environment=str( - node.config.get("environment") or "development" - ), + environment=str(node.config.get("environment") or "development"), publication=( - DataflowPublicationTarget( - target_datasource_ref=target_ref - ) + DataflowPublicationTarget(target_datasource_ref=target_ref) if target_ref else None ), @@ -1525,9 +1523,7 @@ def _start_dataflow_step( causation_id=f"workflow-step:{step.id}", requested_by=instance.created_by, metadata={ - "workflow_instance_ref": ( - f"workflow-instance:{instance.id}" - ), + "workflow_instance_ref": (f"workflow-instance:{instance.id}"), "workflow_step_ref": f"workflow-step:{step.id}", }, ), @@ -1557,12 +1553,8 @@ def _start_dataflow_step( "pipeline_revision": revision, "action_url": _dataflow_action_url(pipeline_ref, run.ref), "allowed_actions": ["cancel"], - "progress_percent": int( - run.metadata.get("progress_percent") or 0 - ), - "progress_phase": str( - run.metadata.get("progress_phase") or run.status - ), + "progress_percent": int(run.metadata.get("progress_percent") or 0), + "progress_phase": str(run.metadata.get("progress_phase") or run.status), } instance.status = "waiting" _record_event( @@ -1592,11 +1584,11 @@ def _handle_dataflow_success( warnings = [ dict(item) for item in items - if isinstance(item, Mapping) - and str(item.get("severity") or "") == "warning" + if isinstance(item, Mapping) and str(item.get("severity") or "") == "warning" ] explicit_review = any( - str(item.get("code") or "") in { + str(item.get("code") or "") + in { "review.required", "reconciliation.review_required", } @@ -1606,8 +1598,7 @@ def _handle_dataflow_success( output = _dataflow_output(descriptor) step.output_ = output if explicit_review or ( - warnings - and str(node.config.get("warning_policy") or "review") == "review" + warnings and str(node.config.get("warning_policy") or "review") == "review" ): step.status = "waiting" step.handoff = { @@ -1626,9 +1617,7 @@ def _handle_dataflow_success( "retry", "cancel", ], - "suggested_port": ( - "review_required" if explicit_review else "warning" - ), + "suggested_port": ("review_required" if explicit_review else "warning"), "warnings": warnings, "output": output, } @@ -1729,11 +1718,8 @@ def _set_human_handoff( "state": "waiting", "title": str(node.config.get("title") or node.label or node.type), "instructions": str(node.config.get("instructions") or ""), - "assignee": node.config.get("reviewer") - or node.config.get("assignee"), - "required_evidence": list( - node.config.get("required_evidence") or [] - ), + "assignee": node.config.get("reviewer") or node.config.get("assignee"), + "required_evidence": list(node.config.get("required_evidence") or []), "allowed_actions": actions, } instance.status = "waiting" @@ -1754,6 +1740,38 @@ def _set_human_handoff( ) +def _set_automated_wait( + session: Session, + *, + instance: WorkflowInstance, + step: WorkflowInstanceStep, + node: WorkflowNode, + wait_state: object, +) -> None: + mode = str(getattr(wait_state, "mode")) + due_at = getattr(wait_state, "due_at") + event_type = getattr(wait_state, "event_type") + step.status = "waiting" + step.handoff = { + "kind": "event_wait" if mode == "event" else "timer", + "state": "waiting", + "title": str(node.config.get("title") or node.label or "Wait"), + "mode": mode, + "due_at": due_at.isoformat() if due_at else None, + "event_type": event_type, + "allowed_actions": ["cancel"], + } + instance.status = "waiting" + _record_event( + session, + instance, + step=step, + kind="workflow.wait.registered", + actor_id=instance.created_by, + payload=dict(step.handoff), + ) + + def _set_dependency_handoff( session: Session, *, @@ -1913,25 +1931,31 @@ def _new_step( instance: WorkflowInstance, node: WorkflowNode, ) -> WorkflowInstanceStep: - sequence = int( - session.scalar( - select(func.max(WorkflowInstanceStep.sequence)).where( - WorkflowInstanceStep.instance_id == instance.id + sequence = ( + int( + session.scalar( + select(func.max(WorkflowInstanceStep.sequence)).where( + WorkflowInstanceStep.instance_id == instance.id + ) ) + or 0 ) - or 0 - ) + 1 - attempt = int( - session.scalar( - select(func.count()) - .select_from(WorkflowInstanceStep) - .where( - WorkflowInstanceStep.instance_id == instance.id, - WorkflowInstanceStep.node_id == node.id, + + 1 + ) + attempt = ( + int( + session.scalar( + select(func.count()) + .select_from(WorkflowInstanceStep) + .where( + WorkflowInstanceStep.instance_id == instance.id, + WorkflowInstanceStep.node_id == node.id, + ) ) + or 0 ) - or 0 - ) + 1 + + 1 + ) step = WorkflowInstanceStep( tenant_id=instance.tenant_id, instance=instance, @@ -1940,9 +1964,7 @@ def _new_step( node_type=node.type, status="running", attempt=attempt, - idempotency_key=( - f"workflow:{instance.id}:node:{node.id}:attempt:{attempt}" - ), + idempotency_key=(f"workflow:{instance.id}:node:{node.id}:attempt:{attempt}"), input_=dict(instance.context_), output_={}, handoff={}, @@ -1962,14 +1984,17 @@ def _record_event( payload: Mapping[str, object], step: WorkflowInstanceStep | None = None, ) -> None: - sequence = int( - session.scalar( - select(func.max(WorkflowInstanceEvent.sequence)).where( - WorkflowInstanceEvent.instance_id == instance.id + sequence = ( + int( + session.scalar( + select(func.max(WorkflowInstanceEvent.sequence)).where( + WorkflowInstanceEvent.instance_id == instance.id + ) ) + or 0 ) - or 0 - ) + 1 + + 1 + ) event = WorkflowInstanceEvent( tenant_id=instance.tenant_id, instance=instance, @@ -1982,6 +2007,46 @@ def _record_event( ) session.add(event) session.flush() + if kind.startswith("workflow.instance."): + from govoplan_workflow_engine.backend.runtime import get_registry + + emit_platform_event( + session, + PlatformEvent( + type=kind, + module_id="workflow_engine", + event_id=event.id, + occurred_at=event.created_at, + correlation_id=instance.correlation_id, + causation_id=( + str(payload.get("event_id")) if payload.get("event_id") else None + ), + actor=( + EventActorRef(type="account", id=actor_id) + if actor_id + else EventActorRef(type="system_actor") + ), + tenant=EventTenantRef(id=instance.tenant_id), + subject=EventObjectRef( + type="workflow_instance", + id=instance.id, + ), + resource=EventObjectRef( + type="workflow_definition", + id=instance.definition_id, + ), + classification="internal", + payload={ + "instance_id": instance.id, + "definition_id": instance.definition_id, + "definition_revision_id": instance.definition_revision_id, + "step_id": step.id if step else None, + "status": instance.status, + "start_origin": instance.start_origin, + }, + ), + registry=get_registry(), + ) def _current_step( @@ -1998,15 +2063,12 @@ def _start_node(graph: WorkflowGraph, *, kind: str) -> WorkflowNode: node = next((item for item in graph.nodes if item.type == expected), None) if node is None: starts = [ - item for item in graph.nodes - if item.type.startswith("workflow.start.") + item for item in graph.nodes if item.type.startswith("workflow.start.") ] if len(starts) == 1: node = starts[0] if node is None: - raise WorkflowConflictError( - f"Workflow has no {kind} start node." - ) + raise WorkflowConflictError(f"Workflow has no {kind} start node.") return node @@ -2024,9 +2086,7 @@ def _normalize_start_origin(value: str) -> str: "backfill", } if normalized not in allowed: - raise WorkflowConflictError( - f"Unsupported Workflow start origin {value!r}." - ) + raise WorkflowConflictError(f"Unsupported Workflow start origin {value!r}.") return normalized @@ -2048,17 +2108,10 @@ def _instance_view_context( instance: WorkflowInstance, revision: WorkflowDefinitionRevision, ) -> WorkflowViewContextResponse | None: - if ( - not revision.view_id - or instance.status not in {"running", "waiting"} - ): + if not revision.view_id or instance.status not in {"running", "waiting"}: return None step = next( - ( - item - for item in instance.steps - if item.id == instance.current_step_id - ), + (item for item in instance.steps if item.id == instance.current_step_id), None, ) node = None @@ -2089,17 +2142,13 @@ def _instance_view_context( def _node(graph: WorkflowGraph, node_id: str) -> WorkflowNode: node = next((item for item in graph.nodes if item.id == node_id), None) if node is None: - raise WorkflowConflictError( - f"Workflow node {node_id!r} no longer exists." - ) + raise WorkflowConflictError(f"Workflow node {node_id!r} no longer exists.") return node def _runtime_graph(revision: WorkflowDefinitionRevision) -> WorkflowGraph: try: - return materialize_runtime_graph( - WorkflowGraph.model_validate(revision.graph) - ) + return materialize_runtime_graph(WorkflowGraph.model_validate(revision.graph)) except BpmnGraphError as exc: raise WorkflowConflictError(str(exc)) from exc @@ -2157,9 +2206,7 @@ def _dataflow_output( "output_materialization_ref": descriptor.output_materialization_ref, "input_row_count": descriptor.input_row_count, "output_row_count": descriptor.output_row_count, - "diagnostics": list( - descriptor.metadata.get("diagnostics") or [] - ), + "diagnostics": list(descriptor.metadata.get("diagnostics") or []), } @@ -2199,6 +2246,32 @@ def _authorization_payload( registry: object | None, ) -> dict[str, object]: principal_ref = principal.to_platform_principal() + scopes = required_instance_scopes( + graph, + principal=principal, + registry=registry, + ) + return { + "contract_version": "1", + "subject_kind": ( + "service_account" if principal_ref.service_account_id else "delegated_user" + ), + "account_id": principal_ref.account_id, + "membership_id": principal_ref.membership_id, + "service_account_id": principal_ref.service_account_id, + "grant_scopes": list(scopes), + "authorization_ref": None, + } + + +def required_instance_scopes( + graph: WorkflowGraph, + *, + principal: ApiPrincipal, + registry: object | None, +) -> tuple[str, ...]: + """Return scopes pinned into an instance or trigger authorization artifact.""" + scopes = {INSTANCE_START_SCOPE} if any(node.type == "workflow.dataflow" for node in graph.nodes): scopes.add(DATAFLOW_RUN_SCOPE) @@ -2211,19 +2284,7 @@ def _authorization_payload( registry=registry, ) scopes.update(definition.required_scopes) - return { - "contract_version": "1", - "subject_kind": ( - "service_account" - if principal_ref.service_account_id - else "delegated_user" - ), - "account_id": principal_ref.account_id, - "membership_id": principal_ref.membership_id, - "service_account_id": principal_ref.service_account_id, - "grant_scopes": sorted(scopes), - "authorization_ref": None, - } + return tuple(sorted(scopes)) def _resolve_instance_principal( @@ -2239,25 +2300,18 @@ def _resolve_instance_principal( common = { "tenant_id": instance.tenant_id, "authorization_ref": str( - value.get("authorization_ref") - or f"workflow-instance:{instance.id}" - ), - "grant_scopes": tuple( - str(scope) for scope in value.get("grant_scopes") or () + value.get("authorization_ref") or f"workflow-instance:{instance.id}" ), + "grant_scopes": tuple(str(scope) for scope in value.get("grant_scopes") or ()), "context": { "workflow_instance_ref": f"workflow-instance:{instance.id}", - "definition_ref": ( - f"workflow-definition:{instance.definition_id}" - ), + "definition_ref": (f"workflow-definition:{instance.definition_id}"), }, } try: if value.get("subject_kind") == "service_account": request = AutomationPrincipalRequest.service_account( - service_account_id=str( - value.get("service_account_id") or "" - ), + service_account_id=str(value.get("service_account_id") or ""), **common, ) else: @@ -2287,8 +2341,7 @@ def _resolve_instance_principal( } return ( resolution.principal - if resolution.allowed - and isinstance(resolution.principal, ApiPrincipal) + if resolution.allowed and isinstance(resolution.principal, ApiPrincipal) else None ) @@ -2302,9 +2355,7 @@ def _notify_handoff( subject: str, ) -> None: provider = notification_dispatch_provider(registry) - account_id = str( - instance.authorization_.get("account_id") or "" - ).strip() + account_id = str(instance.authorization_.get("account_id") or "").strip() if provider is None or not account_id: return try: @@ -2320,9 +2371,7 @@ def _notify_handoff( recipient_id=account_id, subject=subject, action_url=( - "/workflow?" - f"definition={instance.definition_id}" - f"&run={instance.id}" + f"/workflow?definition={instance.definition_id}&run={instance.id}" ), payload={ "instance_id": instance.id, @@ -2351,7 +2400,6 @@ class SqlWorkflowRuntimeWorker: now: datetime | None = None, limit: int = 50, ) -> Mapping[str, object]: - del now if not isinstance(session, Session): raise TypeError("Workflow reconciliation requires a Session.") standards: Mapping[str, object] | None = None @@ -2370,8 +2418,17 @@ class SqlWorkflowRuntimeWorker: registry=self._registry, limit=limit, ) + from govoplan_workflow_engine.backend.triggers import dispatch_due_work + + triggers = dispatch_due_work( + session, + registry=self._registry, + now=now, + limit=limit, + ) return { **runtime, + "triggers": triggers, **({"standards": standards} if standards is not None else {}), } @@ -2384,6 +2441,7 @@ __all__ = [ "list_instances", "reconcile_instance", "reconcile_pending_instances", + "required_instance_scopes", "resolve_step", "start_instance", ] diff --git a/src/govoplan_workflow_engine/backend/manifest.py b/src/govoplan_workflow_engine/backend/manifest.py index 49acf8b..6e2d92b 100644 --- a/src/govoplan_workflow_engine/backend/manifest.py +++ b/src/govoplan_workflow_engine/backend/manifest.py @@ -37,6 +37,7 @@ from govoplan_core.core.workflows import ( CAPABILITY_WORKFLOW_DEFINITION_CONTRIBUTIONS, CAPABILITY_WORKFLOW_ORCHESTRATION, CAPABILITY_WORKFLOW_RUNTIME_WORKER, + CAPABILITY_WORKFLOW_TRIGGER_DISPATCHER, ) from govoplan_core.db.base import Base from govoplan_workflow_engine.backend.db import models as workflow_models @@ -113,7 +114,11 @@ ROLE_TEMPLATES = ( slug="workflow_designer", name="Workflow designer", description="Design, validate, and publish workflow definitions.", - permissions=(DEFINITION_READ_SCOPE, DEFINITION_WRITE_SCOPE, INSTANCE_READ_SCOPE), + permissions=( + DEFINITION_READ_SCOPE, + DEFINITION_WRITE_SCOPE, + INSTANCE_READ_SCOPE, + ), ), RoleTemplate( slug="workflow_operator", @@ -146,6 +151,14 @@ def _runtime_worker(context: ModuleContext): return SqlWorkflowRuntimeWorker(registry=context.registry) +def _trigger_dispatcher(context: ModuleContext): + from govoplan_workflow_engine.backend.triggers import ( + SqlWorkflowTriggerDispatcher, + ) + + return SqlWorkflowTriggerDispatcher(registry=context.registry) + + def _definition_contribution_provider(context: ModuleContext): from govoplan_workflow_engine.backend.contributions import ( SqlWorkflowDefinitionContributionProvider, @@ -218,6 +231,10 @@ manifest = ModuleManifest( version="1.0.0", ), ModuleInterfaceProvider(name="workflow.runtime_worker", version=MODULE_VERSION), + ModuleInterfaceProvider( + name="workflow.trigger_dispatcher", + version="1.0.0", + ), ModuleInterfaceProvider(name="workflow.bpmn_interchange", version="1.0.0"), ModuleInterfaceProvider( name=CAPABILITY_WORKFLOW_SERVICE_LAUNCHER, @@ -274,6 +291,7 @@ manifest = ModuleManifest( _definition_contribution_provider ), CAPABILITY_WORKFLOW_RUNTIME_WORKER: _runtime_worker, + CAPABILITY_WORKFLOW_TRIGGER_DISPATCHER: _trigger_dispatcher, CAPABILITY_WORKFLOW_ORCHESTRATION: _orchestration_provider, WORKFLOW_CONFIGURATION_CAPABILITY: _configuration_provider, CAPABILITY_WORKFLOW_SERVICE_LAUNCHER: _service_launcher, @@ -291,6 +309,9 @@ manifest = ModuleManifest( script_location=str(Path(__file__).with_name("migrations") / "versions"), retirement_supported=True, retirement_provider=drop_table_retirement_provider( + workflow_models.WorkflowWaitState, + workflow_models.WorkflowTriggerDelivery, + workflow_models.WorkflowTrigger, workflow_models.WorkflowInstanceEvent, workflow_models.WorkflowInstanceStep, workflow_models.WorkflowInstance, @@ -310,6 +331,9 @@ manifest = ModuleManifest( workflow_models.WorkflowInstance, workflow_models.WorkflowInstanceStep, workflow_models.WorkflowInstanceEvent, + workflow_models.WorkflowTrigger, + workflow_models.WorkflowTriggerDelivery, + workflow_models.WorkflowWaitState, label="Workflow", ), ), @@ -331,7 +355,13 @@ manifest = ModuleManifest( layer="available", documentation_types=("admin", "user"), audience=("operator", "module_admin", "power_user", "product_owner"), - related_modules=("dataflow", "datasources", "tasks", "notifications", "audit"), + related_modules=( + "dataflow", + "datasources", + "tasks", + "notifications", + "audit", + ), order=76, ), DocumentationTopic( @@ -368,9 +398,23 @@ manifest = ModuleManifest( maturity="vertical_slice", documentation_ref="docs/ENGINE_EDITOR_SPLIT.md", test_ref="tests/test_instance_service.py", - known_limits=("Execution adapters support declared conformance profiles but do not cover every editable BPMN semantic.",), - owned_concepts=("workflow definition", "workflow revision", "workflow instance", "work transition", "execution adapter binding"), - non_owned_concepts=("visual editor", "domain action", "notification", "dataflow run"), + known_limits=( + "Execution adapters support declared conformance profiles but do not cover every editable BPMN semantic.", + "Native schedules intentionally support one-time and bounded interval starts; cron requires a future governed scheduler adapter.", + ), + owned_concepts=( + "workflow definition", + "workflow revision", + "workflow instance", + "work transition", + "execution adapter binding", + ), + non_owned_concepts=( + "visual editor", + "domain action", + "notification", + "dataflow run", + ), recovery_docs=("docs/CONCEPT.md",), security_docs=("docs/CONCEPT.md",), operations_docs=("README.md",), diff --git a/src/govoplan_workflow_engine/backend/migrations/versions/b2e4f6a8c0d1_v0114_workflow_triggers_and_waits.py b/src/govoplan_workflow_engine/backend/migrations/versions/b2e4f6a8c0d1_v0114_workflow_triggers_and_waits.py new file mode 100644 index 0000000..220bb75 --- /dev/null +++ b/src/govoplan_workflow_engine/backend/migrations/versions/b2e4f6a8c0d1_v0114_workflow_triggers_and_waits.py @@ -0,0 +1,211 @@ +"""v0.1.14 durable Workflow triggers and wait states + +Revision ID: b2e4f6a8c0d1 +Revises: 0b4e7c9a2d6f +Create Date: 2026-08-01 00:00:00.000000 +""" + +from __future__ import annotations + +from alembic import op +import sqlalchemy as sa + + +revision = "b2e4f6a8c0d1" +down_revision = "0b4e7c9a2d6f" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "workflow_triggers", + sa.Column("id", sa.String(36), primary_key=True), + sa.Column("tenant_id", sa.String(36), nullable=False), + sa.Column("definition_id", sa.String(36), nullable=False), + sa.Column("definition_revision_id", sa.String(36), nullable=False), + sa.Column("node_id", sa.String(120), nullable=False), + sa.Column("kind", sa.String(30), nullable=False), + sa.Column("status", sa.String(30), nullable=False), + sa.Column("config", sa.JSON(), nullable=False), + sa.Column("event_type", sa.String(120), nullable=True), + sa.Column("next_fire_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_fire_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("last_status", sa.String(30), nullable=True), + sa.Column("last_error", sa.Text(), nullable=True), + sa.Column("authorization_subject_kind", sa.String(30), nullable=False), + sa.Column("authorization_account_id", sa.String(36), nullable=True), + sa.Column("authorization_membership_id", sa.String(36), nullable=True), + sa.Column("authorization_service_account_id", sa.String(36), nullable=True), + sa.Column("authorization_ref", sa.String(255), nullable=False), + sa.Column("grant_scopes", sa.JSON(), nullable=False), + sa.Column("created_by", sa.String(255), nullable=True), + sa.Column("updated_by", sa.String(255), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["definition_id"], + ["workflow_definitions.id"], + ondelete="CASCADE", + ), + sa.ForeignKeyConstraint( + ["definition_revision_id"], + ["workflow_definition_revisions.id"], + ondelete="RESTRICT", + ), + sa.UniqueConstraint( + "definition_id", + "node_id", + name="uq_workflow_trigger_definition_node", + ), + sa.UniqueConstraint( + "authorization_ref", + name="uq_workflow_triggers_authorization_ref", + ), + ) + op.create_index( + "ix_workflow_triggers_due", + "workflow_triggers", + ["status", "next_fire_at"], + ) + op.create_index( + "ix_workflow_triggers_event", + "workflow_triggers", + ["tenant_id", "status", "event_type"], + ) + for name in ( + "tenant_id", + "definition_id", + "definition_revision_id", + "kind", + "status", + "event_type", + "next_fire_at", + "created_by", + "updated_by", + ): + op.create_index( + f"ix_workflow_triggers_{name}", + "workflow_triggers", + [name], + ) + + op.create_table( + "workflow_trigger_deliveries", + sa.Column("id", sa.String(36), primary_key=True), + sa.Column("tenant_id", sa.String(36), nullable=False), + sa.Column("trigger_id", sa.String(36), nullable=False), + sa.Column("definition_id", sa.String(36), nullable=False), + sa.Column("definition_revision_id", sa.String(36), nullable=False), + sa.Column("source_key", sa.String(255), nullable=False), + sa.Column("invocation_kind", sa.String(30), nullable=False), + sa.Column("status", sa.String(30), nullable=False), + sa.Column("scheduled_for", sa.DateTime(timezone=True), nullable=True), + sa.Column("event", sa.JSON(), nullable=True), + sa.Column("instance_id", sa.String(36), nullable=True), + sa.Column("attempts", sa.Integer(), nullable=False), + sa.Column("error", sa.Text(), nullable=True), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["trigger_id"], ["workflow_triggers.id"], ondelete="CASCADE" + ), + sa.ForeignKeyConstraint( + ["definition_id"], + ["workflow_definitions.id"], + ondelete="CASCADE", + ), + sa.ForeignKeyConstraint( + ["definition_revision_id"], + ["workflow_definition_revisions.id"], + ondelete="RESTRICT", + ), + sa.ForeignKeyConstraint( + ["instance_id"], ["workflow_instances.id"], ondelete="SET NULL" + ), + sa.UniqueConstraint( + "trigger_id", + "source_key", + name="uq_workflow_trigger_delivery_source", + ), + ) + op.create_index( + "ix_workflow_trigger_deliveries_queue", + "workflow_trigger_deliveries", + ["status", "created_at"], + ) + op.create_index( + "ix_workflow_trigger_deliveries_tenant_trigger", + "workflow_trigger_deliveries", + ["tenant_id", "trigger_id"], + ) + for name in ( + "tenant_id", + "trigger_id", + "definition_id", + "definition_revision_id", + "status", + "instance_id", + ): + op.create_index( + f"ix_workflow_trigger_deliveries_{name}", + "workflow_trigger_deliveries", + [name], + ) + + op.create_table( + "workflow_wait_states", + sa.Column("id", sa.String(36), primary_key=True), + sa.Column("tenant_id", sa.String(36), nullable=False), + sa.Column("instance_id", sa.String(36), nullable=False), + sa.Column("step_id", sa.String(36), nullable=False), + sa.Column("mode", sa.String(30), nullable=False), + sa.Column("status", sa.String(30), nullable=False), + sa.Column("due_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("event_type", sa.String(120), nullable=True), + sa.Column("config", sa.JSON(), nullable=False), + sa.Column("source_event_id", sa.String(128), nullable=True), + sa.Column("event", sa.JSON(), nullable=True), + sa.Column("error", sa.Text(), nullable=True), + sa.Column("revision", sa.Integer(), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["instance_id"], ["workflow_instances.id"], ondelete="CASCADE" + ), + sa.ForeignKeyConstraint( + ["step_id"], ["workflow_instance_steps.id"], ondelete="CASCADE" + ), + sa.UniqueConstraint("step_id", name="uq_workflow_wait_state_step"), + ) + op.create_index( + "ix_workflow_wait_states_due", + "workflow_wait_states", + ["status", "due_at"], + ) + op.create_index( + "ix_workflow_wait_states_event", + "workflow_wait_states", + ["tenant_id", "status", "event_type"], + ) + for name in ( + "tenant_id", + "instance_id", + "step_id", + "mode", + "status", + "due_at", + "event_type", + "source_event_id", + ): + op.create_index( + f"ix_workflow_wait_states_{name}", + "workflow_wait_states", + [name], + ) + + +def downgrade() -> None: + op.drop_table("workflow_wait_states") + op.drop_table("workflow_trigger_deliveries") + op.drop_table("workflow_triggers") diff --git a/src/govoplan_workflow_engine/backend/router.py b/src/govoplan_workflow_engine/backend/router.py index 80c8b8e..3124222 100644 --- a/src/govoplan_workflow_engine/backend/router.py +++ b/src/govoplan_workflow_engine/backend/router.py @@ -90,6 +90,8 @@ from govoplan_workflow_engine.backend.schemas import ( WorkflowPortResponse, WorkflowStandardDiffResponse, WorkflowStepActionRequest, + WorkflowTriggerListResponse, + WorkflowTriggerResponse, ) from govoplan_workflow_engine.backend.instance_service import ( cancel_instance, @@ -122,6 +124,11 @@ from govoplan_workflow_engine.backend.service import ( update_definition, ) from govoplan_workflow_engine.backend.validation import validate_workflow_graph +from govoplan_workflow_engine.backend.triggers import ( + disable_definition_triggers, + list_definition_triggers, + reconcile_definition_triggers, +) router = APIRouter(prefix="/workflow", tags=["workflow"]) @@ -303,8 +310,7 @@ def _bpmn_inspection_response( ) ) deduplicated = { - (item.code, item.element_id, item.message): item - for item in diagnostics + (item.code, item.element_id, item.message): item for item in diagnostics }.values() diagnostic_items = list(deduplicated) return BpmnInspectionResponse( @@ -338,9 +344,7 @@ def _bpmn_inspection_response( for item in diagnostic_items ], adapter_id=adapter.profile.id if adapter else adapter_id, - adapter_version=( - adapter.profile.version if adapter else adapter_version - ), + adapter_version=(adapter.profile.version if adapter else adapter_version), runtime_kind=adapter.profile.runtime_kind if adapter else None, executable=bool(adapter and adapter.profile.executable), activatable=bool( @@ -514,7 +518,9 @@ def api_node_types( WorkflowNodeTypeResponse( type=definition.type, category=definition.category, - category_label=WORKFLOW_GRAPH_LIBRARY.category_labels[definition.category], + category_label=WORKFLOW_GRAPH_LIBRARY.category_labels[ + definition.category + ], label=definition.label, description=definition.description, icon=definition.icon, @@ -687,10 +693,7 @@ def api_list_instances( ).allowed ] return WorkflowInstanceListResponse( - instances=[ - instance_response(session, instance) - for instance in instances - ] + instances=[instance_response(session, instance) for instance in instances] ) except WorkflowError as exc: raise _http_error(exc) from exc @@ -715,11 +718,7 @@ def api_reconcile_standards( action="workflow.standards.reconciled", object_type="workflow_standard_catalogue", object_id="module-contributions", - details={ - key: value - for key, value in result.items() - if key != "items" - }, + details={key: value for key, value in result.items() if key != "items"}, ) session.commit() return result @@ -746,9 +745,7 @@ def api_start_instance( principal=principal, registry=get_registry(), payload=payload, - start_origin=( - "user" if principal.auth_method == "session" else "api" - ), + start_origin=("user" if principal.auth_method == "session" else "api"), ) except WorkflowError as exc: raise _http_error(exc) from exc @@ -756,9 +753,7 @@ def api_start_instance( session, principal, action=( - "workflow.instance.replayed" - if replayed - else "workflow.instance.started" + "workflow.instance.replayed" if replayed else "workflow.instance.started" ), instance_id=instance.id, details={ @@ -1026,13 +1021,11 @@ def api_create_definition( ) -> WorkflowDefinitionResponse: _require_any_scope(principal, DEFINITION_WRITE_SCOPE, ADMIN_SCOPE) try: - tenant_id, scope_type, scope_id, _scope_key = ( - normalize_definition_scope( - principal, - scope_type=payload.scope_type, - scope_id=payload.scope_id, - administrative=has_scope(principal, ADMIN_SCOPE), - ) + tenant_id, scope_type, scope_id, _scope_key = normalize_definition_scope( + principal, + scope_type=payload.scope_type, + scope_id=payload.scope_id, + administrative=has_scope(principal, ADMIN_SCOPE), ) payload = payload.model_copy( update={"scope_type": scope_type, "scope_id": scope_id} @@ -1107,6 +1100,58 @@ def api_get_definition( raise _http_error(exc) from exc +@router.get( + "/definitions/{definition_id}/triggers", + response_model=WorkflowTriggerListResponse, +) +def api_list_definition_triggers( + definition_id: str, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), +) -> WorkflowTriggerListResponse: + _require_any_scope(principal, DEFINITION_READ_SCOPE, ADMIN_SCOPE) + try: + definition = get_definition( + session, + tenant_id=principal.tenant_id, + definition_id=definition_id, + ) + require_definition_action( + definition, + principal=principal, + registry=get_registry(), + action="view", + ) + except PermissionError as exc: + raise _governance_http_error(exc) from exc + except WorkflowError as exc: + raise _http_error(exc) from exc + return WorkflowTriggerListResponse( + triggers=[ + WorkflowTriggerResponse( + id=item.id, + definition_id=item.definition_id, + definition_revision_id=item.definition_revision_id, + node_id=item.node_id, + kind=item.kind, + status=item.status, + event_type=item.event_type, + next_fire_at=item.next_fire_at, + last_fire_at=item.last_fire_at, + last_status=item.last_status, + last_error=item.last_error, + authorization_subject_kind=item.authorization_subject_kind, + grant_scopes=list(item.grant_scopes), + ) + for item in list_definition_triggers( + session, + tenant_id=principal.tenant_id, + definition_id=definition_id, + ) + ] + ) + + @router.put( "/definitions/{definition_id}", response_model=WorkflowDefinitionResponse, @@ -1142,9 +1187,7 @@ def api_update_definition( scope_type=scope_type, scope_id=scope_id, preserve_existing=( - existing.scope_id - if existing.scope_type == scope_type - else None + existing.scope_id if existing.scope_type == scope_type else None ), ) payload = payload.model_copy( @@ -1189,13 +1232,11 @@ def api_derive_definition( ) -> WorkflowDefinitionResponse: _require_any_scope(principal, DEFINITION_WRITE_SCOPE, ADMIN_SCOPE) try: - tenant_id, scope_type, scope_id, _scope_key = ( - normalize_definition_scope( - principal, - scope_type=payload.scope_type, - scope_id=payload.scope_id, - administrative=has_scope(principal, ADMIN_SCOPE), - ) + tenant_id, scope_type, scope_id, _scope_key = normalize_definition_scope( + principal, + scope_type=payload.scope_type, + scope_id=payload.scope_id, + administrative=has_scope(principal, ADMIN_SCOPE), ) payload = payload.model_copy( update={"scope_type": scope_type, "scope_id": scope_id} @@ -1226,9 +1267,7 @@ def api_derive_definition( action="workflow.definition.derived", definition_id=definition.id, details={ - "source_definition_id": ( - definition.derived_from_definition_id - ), + "source_definition_id": (definition.derived_from_definition_id), "source_revision": definition.derived_from_revision, "source_hash": definition.derived_from_hash, "scope_type": definition.scope_type, @@ -1396,6 +1435,12 @@ def api_activate_definition( actor_id=_actor_id(principal), revision=payload.revision, ) + trigger_summary = reconcile_definition_triggers( + session, + definition=definition, + principal=principal, + registry=get_registry(), + ) except PermissionError as exc: raise _governance_http_error(exc) from exc except WorkflowError as exc: @@ -1405,7 +1450,10 @@ def api_activate_definition( principal, action="workflow.definition.activated", definition_id=definition.id, - details={"active_revision": definition.active_revision}, + details={ + "active_revision": definition.active_revision, + "triggers": trigger_summary, + }, ) response = _definition_response(session, definition, principal) session.commit() @@ -1440,6 +1488,10 @@ def api_archive_definition( definition_id=definition_id, actor_id=_actor_id(principal), ) + disabled_triggers = disable_definition_triggers( + session, + definition_id=definition.id, + ) except PermissionError as exc: raise _governance_http_error(exc) from exc except WorkflowError as exc: @@ -1449,7 +1501,10 @@ def api_archive_definition( principal, action="workflow.definition.archived", definition_id=definition.id, - details={"active_revision": definition.active_revision}, + details={ + "active_revision": definition.active_revision, + "disabled_triggers": disabled_triggers, + }, ) response = _definition_response(session, definition, principal) session.commit() diff --git a/src/govoplan_workflow_engine/backend/schemas.py b/src/govoplan_workflow_engine/backend/schemas.py index cd92fe2..f4cb3e3 100644 --- a/src/govoplan_workflow_engine/backend/schemas.py +++ b/src/govoplan_workflow_engine/backend/schemas.py @@ -501,6 +501,26 @@ class WorkflowDefinitionDeleteResponse(BaseModel): definition_id: str +class WorkflowTriggerResponse(BaseModel): + id: str + definition_id: str + definition_revision_id: str + node_id: str + kind: Literal["schedule", "event"] + status: str + event_type: str | None = None + next_fire_at: datetime | None = None + last_fire_at: datetime | None = None + last_status: str | None = None + last_error: str | None = None + authorization_subject_kind: Literal["delegated_user", "service_account"] + grant_scopes: list[str] = Field(default_factory=list) + + +class WorkflowTriggerListResponse(BaseModel): + triggers: list[WorkflowTriggerResponse] + + WorkflowInstanceStatus = Literal[ "running", "waiting", diff --git a/src/govoplan_workflow_engine/backend/triggers.py b/src/govoplan_workflow_engine/backend/triggers.py new file mode 100644 index 0000000..96eea3f --- /dev/null +++ b/src/govoplan_workflow_engine/backend/triggers.py @@ -0,0 +1,1030 @@ +from __future__ import annotations + +from collections.abc import Mapping +from datetime import datetime, timedelta, timezone +import hashlib +import json +import re +from zoneinfo import ZoneInfo, ZoneInfoNotFoundError + +from sqlalchemy import and_, or_, select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.orm import Session + +from govoplan_core.auth import ApiPrincipal, has_scope +from govoplan_core.core.automation import ( + AutomationPrincipalRequest, + automation_principal_provider, +) +from govoplan_core.core.events import PlatformEvent +from govoplan_core.db.base import utcnow +from govoplan_workflow_engine.backend.db.models import ( + WorkflowDefinition, + WorkflowDefinitionRevision, + WorkflowInstance, + WorkflowInstanceStep, + WorkflowTrigger, + WorkflowTriggerDelivery, + WorkflowWaitState, +) +from govoplan_workflow_engine.backend.schemas import ( + WorkflowGraph, + WorkflowInstanceStartRequest, + WorkflowNode, +) +from govoplan_workflow_engine.backend.service import ( + WorkflowConflictError, + get_definition_revision, +) + + +_EVENT_TYPE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,119}$") +_DURATION = re.compile( + r"^(?P[1-9][0-9]*)(?Ps|m|h|d)?$", + re.IGNORECASE, +) +_ISO_DURATION = re.compile( + r"^P(?:(?P[0-9]+)D)?(?:T(?:(?P[0-9]+)H)?" + r"(?:(?P[0-9]+)M)?(?:(?P[0-9]+)S)?)?$", + re.IGNORECASE, +) +MINIMUM_INTERVAL_SECONDS = 60 + + +def reconcile_definition_triggers( + session: Session, + *, + definition: WorkflowDefinition, + principal: ApiPrincipal, + registry: object | None, +) -> dict[str, int]: + """Replace active graph trigger registrations in the activation transaction.""" + + if definition.active_revision is None: + raise WorkflowConflictError("An active Workflow revision is required.") + revision = get_definition_revision( + session, + definition=definition, + revision=definition.active_revision, + ) + graph = _runtime_graph(revision) + nodes = [ + node + for node in graph.nodes + if node.type + in { + "workflow.start.schedule", + "workflow.start.event", + "workflow.start.workflow", + } + ] + existing = { + item.node_id: item + for item in session.scalars( + select(WorkflowTrigger).where( + WorkflowTrigger.definition_id == definition.id + ) + ) + } + if nodes and not definition.allow_automation: + raise WorkflowConflictError( + "Scheduled and event starts require automation to be enabled." + ) + if nodes and automation_principal_provider(registry) is None: + raise WorkflowConflictError( + "Scheduled and event starts require the automation-principal provider." + ) + + from govoplan_workflow_engine.backend.instance_service import ( + required_instance_scopes, + ) + + grant_scopes = required_instance_scopes( + graph, + principal=principal, + registry=registry, + ) + missing = [scope for scope in grant_scopes if not has_scope(principal, scope)] + if nodes and missing: + raise WorkflowConflictError( + "The activation principal cannot delegate required scopes: " + + ", ".join(missing) + ) + authorization = _authorization_owner(principal) + now = _as_utc(utcnow()) + retained: set[str] = set() + created = 0 + updated = 0 + for node in nodes: + retained.add(node.id) + kind = node.type.removeprefix("workflow.start.") + config, event_type, next_fire_at = _trigger_config( + node, + now=now, + ) + trigger = existing.get(node.id) + if trigger is None: + trigger = WorkflowTrigger( + tenant_id=str(definition.tenant_id or principal.tenant_id), + definition_id=definition.id, + definition_revision_id=revision.id, + node_id=node.id, + kind=kind, + status="enabled", + config_=config, + event_type=event_type, + next_fire_at=next_fire_at, + authorization_ref=f"workflow-trigger:{definition.id}:{node.id}", + grant_scopes=list(grant_scopes), + created_by=principal.account_id, + updated_by=principal.account_id, + **authorization, + ) + session.add(trigger) + created += 1 + else: + _cancel_queued_deliveries( + session, + trigger=trigger, + reason="The Workflow trigger was replaced by a new active revision.", + ) + trigger.definition_revision_id = revision.id + trigger.kind = kind + trigger.status = "enabled" + trigger.config_ = config + trigger.event_type = event_type + trigger.next_fire_at = next_fire_at + trigger.last_error = None + trigger.grant_scopes = list(grant_scopes) + trigger.updated_by = principal.account_id + for key, value in authorization.items(): + setattr(trigger, key, value) + updated += 1 + disabled = 0 + for node_id, trigger in existing.items(): + if node_id in retained: + continue + _disable_trigger( + session, + trigger=trigger, + reason="The active Workflow revision no longer contains this trigger.", + ) + disabled += 1 + session.flush() + return {"created": created, "updated": updated, "disabled": disabled} + + +def disable_definition_triggers( + session: Session, + *, + definition_id: str, + reason: str = "The Workflow definition was archived.", +) -> int: + triggers = list( + session.scalars( + select(WorkflowTrigger).where( + WorkflowTrigger.definition_id == definition_id, + WorkflowTrigger.status == "enabled", + ) + ) + ) + for trigger in triggers: + _disable_trigger(session, trigger=trigger, reason=reason) + session.flush() + return len(triggers) + + +def list_definition_triggers( + session: Session, + *, + tenant_id: str, + definition_id: str, +) -> list[WorkflowTrigger]: + return list( + session.scalars( + select(WorkflowTrigger) + .where( + WorkflowTrigger.tenant_id == tenant_id, + WorkflowTrigger.definition_id == definition_id, + ) + .order_by(WorkflowTrigger.node_id) + ) + ) + + +def register_wait_state( + session: Session, + *, + instance: WorkflowInstance, + step: WorkflowInstanceStep, + node: WorkflowNode, + now: datetime | None = None, +) -> WorkflowWaitState | None: + mode = str(node.config.get("mode") or "manual").strip().lower() + if mode == "manual": + return None + current = _as_utc(now or utcnow()) + due_at: datetime | None = None + event_type: str | None = None + config: dict[str, object] = {"value": node.config.get("value")} + if mode == "duration": + seconds = _duration_seconds(node.config.get("value"), minimum=1) + due_at = current + timedelta(seconds=seconds) + config["duration_seconds"] = seconds + elif mode == "deadline": + due_at = _parse_instant( + node.config.get("value"), + timezone_name=str(node.config.get("timezone") or "UTC"), + ) + if due_at < current: + due_at = current + elif mode == "event": + event_type = _validated_event_type( + node.config.get("event_type") or node.config.get("value") + ) + config["filter"] = _event_filter(node.config.get("filter")) + else: + raise WorkflowConflictError(f"Unsupported wait mode {mode!r}.") + state = WorkflowWaitState( + tenant_id=instance.tenant_id, + instance_id=instance.id, + step_id=step.id, + mode=mode, + status="waiting", + due_at=due_at, + event_type=event_type, + config_=config, + ) + session.add(state) + session.flush() + return state + + +def resolve_wait_state( + session: Session, + *, + step_id: str, + status: str, +) -> None: + state = session.scalar( + select(WorkflowWaitState).where(WorkflowWaitState.step_id == step_id) + ) + if state is None or state.status in {"resumed", "timed_out", "cancelled"}: + return + state.status = status + state.revision += 1 + + +def ingest_platform_event( + session: Session, + *, + event: PlatformEvent, +) -> dict[str, int]: + if event.classification not in {"public", "internal"} or event.tenant is None: + return {"trigger_deliveries": 0, "waits_triggered": 0} + tenant_id = event.tenant.id + envelope = event.to_dict() + triggers = list( + session.scalars( + select(WorkflowTrigger).where( + WorkflowTrigger.tenant_id == tenant_id, + WorkflowTrigger.kind.in_(("event", "workflow")), + WorkflowTrigger.status == "enabled", + WorkflowTrigger.event_type == event.type, + ) + ) + ) + queued = 0 + for trigger in triggers: + if not _matches_filter(envelope, trigger.config_.get("filter")): + continue + delivery = WorkflowTriggerDelivery( + tenant_id=tenant_id, + trigger_id=trigger.id, + definition_id=trigger.definition_id, + definition_revision_id=trigger.definition_revision_id, + source_key=f"event:{event.event_id}", + invocation_kind=( + "parent_workflow" if trigger.kind == "workflow" else "event" + ), + status="queued", + event_=envelope, + ) + try: + with session.begin_nested(): + session.add(delivery) + session.flush() + except IntegrityError: + continue + queued += 1 + + waits = list( + session.scalars( + select(WorkflowWaitState) + .where( + WorkflowWaitState.tenant_id == tenant_id, + WorkflowWaitState.status == "waiting", + WorkflowWaitState.mode == "event", + WorkflowWaitState.event_type == event.type, + ) + .with_for_update(skip_locked=True) + ) + ) + triggered = 0 + for state in waits: + if not _matches_filter(envelope, state.config_.get("filter")): + continue + state.status = "triggered" + state.source_event_id = event.event_id + state.event_ = envelope + state.revision += 1 + triggered += 1 + session.flush() + return {"trigger_deliveries": queued, "waits_triggered": triggered} + + +def dispatch_due_work( + session: Session, + *, + registry: object | None, + now: datetime | None = None, + limit: int = 50, +) -> dict[str, int]: + current = _as_utc(now or utcnow()) + bounded = max(1, min(int(limit), 200)) + scheduled = _queue_due_schedules(session, now=current, limit=bounded) + deliveries = list( + session.scalars( + select(WorkflowTriggerDelivery) + .where(WorkflowTriggerDelivery.status == "queued") + .order_by( + WorkflowTriggerDelivery.created_at, + WorkflowTriggerDelivery.id, + ) + .limit(bounded) + .with_for_update(skip_locked=True) + ) + ) + started = blocked = failed = 0 + for delivery in deliveries: + outcome = _dispatch_delivery( + session, + delivery=delivery, + registry=registry, + ) + if outcome == "started": + started += 1 + elif outcome == "blocked": + blocked += 1 + else: + failed += 1 + waits = _dispatch_waits( + session, + registry=registry, + now=current, + limit=bounded, + ) + session.flush() + return { + "scheduled": scheduled, + "deliveries": len(deliveries), + "started": started, + "blocked": blocked, + "failed": failed, + **waits, + } + + +class SqlWorkflowTriggerDispatcher: + def __init__(self, *, registry: object | None = None) -> None: + self._registry = registry + + def dispatch_due( + self, + session: object, + *, + now: datetime | None = None, + limit: int = 50, + ) -> Mapping[str, object]: + if not isinstance(session, Session): + raise TypeError("Workflow trigger dispatch requires a Session.") + return dispatch_due_work( + session, + registry=self._registry, + now=now, + limit=limit, + ) + + def ingest_event( + self, + session: object, + *, + event: PlatformEvent, + ) -> Mapping[str, object]: + if not isinstance(session, Session): + raise TypeError("Workflow event ingestion requires a Session.") + return ingest_platform_event(session, event=event) + + +def _queue_due_schedules( + session: Session, + *, + now: datetime, + limit: int, +) -> int: + triggers = list( + session.scalars( + select(WorkflowTrigger) + .where( + WorkflowTrigger.kind == "schedule", + WorkflowTrigger.status == "enabled", + WorkflowTrigger.next_fire_at.is_not(None), + WorkflowTrigger.next_fire_at <= now, + ) + .order_by(WorkflowTrigger.next_fire_at, WorkflowTrigger.id) + .limit(limit) + .with_for_update(skip_locked=True) + ) + ) + queued = 0 + for trigger in triggers: + scheduled_for = _as_utc(trigger.next_fire_at or now) + source = f"schedule:{scheduled_for.isoformat()}" + delivery = WorkflowTriggerDelivery( + tenant_id=trigger.tenant_id, + trigger_id=trigger.id, + definition_id=trigger.definition_id, + definition_revision_id=trigger.definition_revision_id, + source_key=source, + invocation_kind="schedule", + status="queued", + scheduled_for=scheduled_for, + ) + try: + with session.begin_nested(): + session.add(delivery) + session.flush() + except IntegrityError: + pass + else: + queued += 1 + interval = int(trigger.config_.get("interval_seconds") or 0) + if interval: + elapsed = max(0, int((now - scheduled_for).total_seconds())) + periods = (elapsed // interval) + 1 + trigger.next_fire_at = scheduled_for + timedelta(seconds=periods * interval) + else: + trigger.status = "disabled" + trigger.next_fire_at = None + trigger.last_fire_at = scheduled_for + trigger.last_status = "queued" + trigger.last_error = None + return queued + + +def _dispatch_delivery( + session: Session, + *, + delivery: WorkflowTriggerDelivery, + registry: object | None, +) -> str: + trigger = session.get(WorkflowTrigger, delivery.trigger_id) + definition = session.get(WorkflowDefinition, delivery.definition_id) + revision = session.get( + WorkflowDefinitionRevision, + delivery.definition_revision_id, + ) + if ( + trigger is None + or definition is None + or revision is None + or definition.status != "active" + or definition.active_revision != revision.revision + or trigger.definition_revision_id != revision.id + ): + delivery.status = "blocked" + delivery.error = "The exact active Workflow trigger revision is unavailable." + return "blocked" + principal = _resolve_trigger_principal( + session, + trigger=trigger, + registry=registry, + ) + if principal is None: + delivery.status = "blocked" + delivery.error = trigger.last_error or "Trigger authority is unavailable." + trigger.last_status = "blocked" + return "blocked" + from govoplan_workflow_engine.backend.instance_service import start_instance + + delivery.status = "running" + delivery.attempts += 1 + event = dict(delivery.event_ or {}) + input_value: dict[str, object] + correlation_id: str | None = None + if event: + input_value = _mapped_event_input( + trigger.config_.get("input_mapping"), + event, + ) + if not input_value: + input_value = { + ( + "parent_workflow" + if delivery.invocation_kind == "parent_workflow" + else "event" + ): event + } + correlation_id = ( + str(event.get("correlation_id") or event.get("event_id") or "") or None + ) + else: + input_value = { + "schedule": { + "trigger_id": trigger.id, + "scheduled_for": ( + delivery.scheduled_for.isoformat() + if delivery.scheduled_for + else None + ), + } + } + key_hash = hashlib.sha256( + f"{trigger.id}:{delivery.source_key}".encode("utf-8") + ).hexdigest() + try: + instance, _replayed = start_instance( + session, + tenant_id=delivery.tenant_id, + definition_id=delivery.definition_id, + actor_id=trigger.created_by, + principal=principal, + registry=registry, + payload=WorkflowInstanceStartRequest( + idempotency_key=f"trigger:{key_hash}", + input=input_value, + correlation_id=correlation_id, + ), + start_origin=delivery.invocation_kind, + ) + except WorkflowConflictError as exc: + delivery.status = "blocked" + delivery.error = str(exc) + trigger.last_status = "blocked" + trigger.last_error = str(exc) + return "blocked" + except (TypeError, ValueError) as exc: + delivery.status = "failed" + delivery.error = str(exc) + trigger.last_status = "failed" + trigger.last_error = str(exc) + return "failed" + delivery.instance_id = instance.id + delivery.status = "succeeded" + delivery.error = None + trigger.last_status = "succeeded" + trigger.last_error = None + return "started" + + +def _dispatch_waits( + session: Session, + *, + registry: object | None, + now: datetime, + limit: int, +) -> dict[str, int]: + states = list( + session.scalars( + select(WorkflowWaitState) + .where( + or_( + WorkflowWaitState.status == "triggered", + and_( + WorkflowWaitState.status == "waiting", + WorkflowWaitState.due_at.is_not(None), + WorkflowWaitState.due_at <= now, + ), + ) + ) + .order_by(WorkflowWaitState.updated_at, WorkflowWaitState.id) + .limit(limit) + .with_for_update(skip_locked=True) + ) + ) + resumed = timed_out = skipped = 0 + from govoplan_workflow_engine.backend.instance_service import ( + _complete_step, + _drive_instance, + _record_event, + _resolve_instance_principal, + _runtime_graph, + ) + + for state in states: + instance = session.get(WorkflowInstance, state.instance_id) + step = session.get(WorkflowInstanceStep, state.step_id) + if ( + instance is None + or step is None + or instance.current_step_id != step.id + or instance.status != "waiting" + or step.status != "waiting" + ): + state.status = "cancelled" + state.error = "The waiting Workflow step is no longer current." + state.revision += 1 + skipped += 1 + continue + principal = _resolve_instance_principal( + session, + instance=instance, + registry=registry, + ) + if principal is None: + state.error = "Current trigger authority could not be resolved." + skipped += 1 + continue + revision = session.get( + WorkflowDefinitionRevision, + instance.definition_revision_id, + ) + if revision is None: + state.status = "cancelled" + state.error = "The pinned Workflow revision is unavailable." + state.revision += 1 + skipped += 1 + continue + graph = _runtime_graph(revision) + triggered = state.status == "triggered" + state.status = "resumed" if triggered else "timed_out" + state.error = None + state.revision += 1 + output = ( + {"event": dict(state.event_ or {})} + if triggered + else {"due_at": state.due_at.isoformat() if state.due_at else None} + ) + _record_event( + session, + instance, + step=step, + kind=( + "workflow.wait.event_received" + if triggered + else "workflow.wait.timed_out" + ), + actor_id=None, + payload=output, + ) + next_node_id = _complete_step( + session, + instance=instance, + step=step, + graph=graph, + port="resumed" if triggered else "timed_out", + output=output, + actor_id=None, + ) + _drive_instance( + session, + instance=instance, + graph=graph, + next_node_id=next_node_id, + principal=principal, + registry=registry, + actor_id=None, + ) + if triggered: + resumed += 1 + else: + timed_out += 1 + return { + "waits_selected": len(states), + "waits_resumed": resumed, + "waits_timed_out": timed_out, + "waits_skipped": skipped, + } + + +def _resolve_trigger_principal( + session: Session, + *, + trigger: WorkflowTrigger, + registry: object | None, +) -> ApiPrincipal | None: + provider = automation_principal_provider(registry) + if provider is None: + trigger.last_error = "The automation-principal provider is unavailable." + return None + common = { + "tenant_id": trigger.tenant_id, + "authorization_ref": trigger.authorization_ref, + "grant_scopes": tuple(trigger.grant_scopes), + "context": { + "workflow_trigger_ref": f"workflow-trigger:{trigger.id}", + "definition_ref": f"workflow-definition:{trigger.definition_id}", + }, + } + try: + if trigger.authorization_subject_kind == "service_account": + request = AutomationPrincipalRequest.service_account( + service_account_id=str(trigger.authorization_service_account_id or ""), + **common, + ) + else: + request = AutomationPrincipalRequest.delegated_user( + account_id=str(trigger.authorization_account_id or ""), + membership_id=str(trigger.authorization_membership_id or ""), + **common, + ) + except ValueError as exc: + trigger.last_error = str(exc) + return None + resolution = provider.resolve_automation_principal(session, request=request) + if not resolution.allowed or not isinstance(resolution.principal, ApiPrincipal): + trigger.last_error = str( + resolution.provenance.get("reason") + or "Current trigger authority was denied." + ) + return None + return resolution.principal + + +def _authorization_owner(principal: ApiPrincipal) -> dict[str, str | None]: + ref = principal.to_platform_principal() + if ref.service_account_id: + return { + "authorization_subject_kind": "service_account", + "authorization_account_id": None, + "authorization_membership_id": None, + "authorization_service_account_id": ref.service_account_id, + } + if not ref.account_id or not ref.membership_id: + raise WorkflowConflictError( + "Automated Workflow activation requires a tenant membership." + ) + return { + "authorization_subject_kind": "delegated_user", + "authorization_account_id": ref.account_id, + "authorization_membership_id": ref.membership_id, + "authorization_service_account_id": None, + } + + +def _trigger_config( + node: WorkflowNode, + *, + now: datetime, +) -> tuple[dict[str, object], str | None, datetime | None]: + if node.type == "workflow.start.event": + event_type = _validated_event_type(node.config.get("event_type")) + return ( + { + "event_type": event_type, + "filter": _event_filter(node.config.get("filter")), + }, + event_type, + None, + ) + if node.type == "workflow.start.workflow": + parent_ref = str(node.config.get("parent_definition_ref") or "").strip() + parent_id = parent_ref.removeprefix("workflow-definition:") + if not parent_id: + raise WorkflowConflictError( + "A parent Workflow definition reference is required." + ) + outcome = str(node.config.get("parent_outcome") or "completed").strip().lower() + if outcome not in {"completed", "cancelled", "failed"}: + raise WorkflowConflictError( + "Parent Workflow outcomes are completed, cancelled, or failed." + ) + event_type = f"workflow.instance.{outcome}" + return ( + { + "event_type": event_type, + "filter": {"payload": {"definition_id": parent_id}}, + "input_mapping": dict(node.config.get("input_mapping") or {}), + "parent_definition_id": parent_id, + "parent_outcome": outcome, + }, + event_type, + None, + ) + schedule = node.config.get("schedule") + timezone_name = str(node.config.get("timezone") or "UTC") + config, next_fire_at = _schedule_config( + schedule, + timezone_name=timezone_name, + now=now, + ) + return config, None, next_fire_at + + +def _schedule_config( + value: object, + *, + timezone_name: str, + now: datetime, +) -> tuple[dict[str, object], datetime]: + if isinstance(value, Mapping): + kind = str(value.get("kind") or "").strip().lower() + if kind == "interval": + seconds = _duration_seconds( + value.get("seconds") or value.get("interval"), + minimum=MINIMUM_INTERVAL_SECONDS, + ) + start = value.get("start_at") + first = ( + _parse_instant(start, timezone_name=timezone_name) + if start + else now + timedelta(seconds=seconds) + ) + return { + "kind": "interval", + "interval_seconds": seconds, + "timezone": timezone_name, + }, max(first, now) + if kind == "once": + instant = _parse_instant(value.get("at"), timezone_name=timezone_name) + return {"kind": "once", "timezone": timezone_name}, max(instant, now) + raise WorkflowConflictError("Schedules must use kind 'once' or 'interval'.") + text = str(value or "").strip() + if not text: + raise WorkflowConflictError("A schedule value is required.") + lower = text.lower() + if lower.startswith("interval:"): + seconds = _duration_seconds( + text.split(":", 1)[1], + minimum=MINIMUM_INTERVAL_SECONDS, + ) + return { + "kind": "interval", + "interval_seconds": seconds, + "timezone": timezone_name, + }, now + timedelta(seconds=seconds) + if lower.startswith("once:"): + text = text.split(":", 1)[1].strip() + instant = _parse_instant(text, timezone_name=timezone_name) + return {"kind": "once", "timezone": timezone_name}, max(instant, now) + + +def _duration_seconds(value: object, *, minimum: int) -> int: + text = str(value or "").strip() + match = _DURATION.fullmatch(text) + if match: + multiplier = {None: 1, "s": 1, "m": 60, "h": 3600, "d": 86400}[ + (match.group("unit") or "").lower() or None + ] + seconds = int(match.group("value")) * multiplier + else: + iso = _ISO_DURATION.fullmatch(text) + if not iso: + raise WorkflowConflictError( + "Durations use seconds, 5m/2h/1d, or ISO 8601 such as PT15M." + ) + seconds = ( + int(iso.group("days") or 0) * 86400 + + int(iso.group("hours") or 0) * 3600 + + int(iso.group("minutes") or 0) * 60 + + int(iso.group("seconds") or 0) + ) + if seconds < minimum: + raise WorkflowConflictError(f"The interval must be at least {minimum} seconds.") + return seconds + + +def _parse_instant(value: object, *, timezone_name: str) -> datetime: + text = str(value or "").strip() + if not text: + raise WorkflowConflictError("A date and time are required.") + try: + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + except ValueError as exc: + raise WorkflowConflictError("Date/time values must use ISO 8601.") from exc + if parsed.tzinfo is None: + try: + parsed = parsed.replace(tzinfo=ZoneInfo(timezone_name)) + except ZoneInfoNotFoundError as exc: + raise WorkflowConflictError( + f"Unknown time zone {timezone_name!r}." + ) from exc + return _as_utc(parsed) + + +def _event_filter(value: object) -> dict[str, object]: + if value in (None, "", {}): + return {} + if isinstance(value, Mapping): + return {str(key): item for key, item in value.items()} + try: + parsed = json.loads(str(value)) + except json.JSONDecodeError as exc: + raise WorkflowConflictError( + "Event filters must be JSON objects; executable expressions are not allowed." + ) from exc + if not isinstance(parsed, dict): + raise WorkflowConflictError("Event filters must be JSON objects.") + return parsed + + +def _validated_event_type(value: object) -> str: + event_type = str(value or "").strip() + if not _EVENT_TYPE.fullmatch(event_type): + raise WorkflowConflictError("A valid platform event type is required.") + return event_type + + +def _matches_filter(value: object, expected: object) -> bool: + if not expected: + return True + if not isinstance(expected, Mapping) or not isinstance(value, Mapping): + return value == expected + return all( + key in value and _matches_filter(value[key], item) + for key, item in expected.items() + ) + + +def _mapped_event_input( + value: object, + event: Mapping[str, object], +) -> dict[str, object]: + if not isinstance(value, Mapping) or not value: + return {} + result: dict[str, object] = {} + for key, raw in value.items(): + output_key = str(key).strip() + if not output_key: + raise WorkflowConflictError( + "Parent Workflow input mappings require non-empty keys." + ) + if not isinstance(raw, str) or not raw.startswith("$event"): + result[output_key] = raw + continue + current: object = event + path = raw.removeprefix("$event").removeprefix(".") + for segment in path.split(".") if path else (): + if not isinstance(current, Mapping) or segment not in current: + current = None + break + current = current[segment] + result[output_key] = current + return result + + +def _runtime_graph(revision: WorkflowDefinitionRevision) -> WorkflowGraph: + from govoplan_workflow_engine.backend.bpmn_graph import materialize_runtime_graph + + graph = WorkflowGraph.model_validate(revision.graph) + try: + return materialize_runtime_graph(graph) + except Exception as exc: + raise WorkflowConflictError(str(exc)) from exc + + +def _disable_trigger( + session: Session, + *, + trigger: WorkflowTrigger, + reason: str, +) -> None: + trigger.status = "disabled" + trigger.next_fire_at = None + trigger.last_status = "disabled" + trigger.last_error = reason + _cancel_queued_deliveries(session, trigger=trigger, reason=reason) + + +def _cancel_queued_deliveries( + session: Session, + *, + trigger: WorkflowTrigger, + reason: str, +) -> None: + for delivery in session.scalars( + select(WorkflowTriggerDelivery).where( + WorkflowTriggerDelivery.trigger_id == trigger.id, + WorkflowTriggerDelivery.status == "queued", + ) + ): + delivery.status = "cancelled" + delivery.error = reason + + +def _as_utc(value: datetime) -> datetime: + if value.tzinfo is None: + return value.replace(tzinfo=timezone.utc) + return value.astimezone(timezone.utc) + + +__all__ = [ + "SqlWorkflowTriggerDispatcher", + "disable_definition_triggers", + "dispatch_due_work", + "ingest_platform_event", + "list_definition_triggers", + "reconcile_definition_triggers", + "register_wait_state", + "resolve_wait_state", +] diff --git a/tests/test_instance_service.py b/tests/test_instance_service.py index e4119c2..1cb425d 100644 --- a/tests/test_instance_service.py +++ b/tests/test_instance_service.py @@ -39,6 +39,9 @@ from govoplan_workflow_engine.backend.db.models import ( WorkflowInstance, WorkflowInstanceEvent, WorkflowInstanceStep, + WorkflowTrigger, + WorkflowTriggerDelivery, + WorkflowWaitState, ) from govoplan_workflow_engine.backend.instance_service import ( SqlWorkflowRuntimeWorker, @@ -64,6 +67,7 @@ from govoplan_workflow_engine.backend.service import ( create_definition, ) from govoplan_workflow_engine.backend.service_launcher import WorkflowServiceLauncher + try: from test_bpmn import NATIVE_BPMN except ModuleNotFoundError as exc: @@ -364,12 +368,13 @@ class Registry: 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 + 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): @@ -396,6 +401,9 @@ class WorkflowInstanceServiceTests(unittest.TestCase): WorkflowInstance.__table__, WorkflowInstanceStep.__table__, WorkflowInstanceEvent.__table__, + WorkflowTrigger.__table__, + WorkflowTriggerDelivery.__table__, + WorkflowWaitState.__table__, ], ) self.Session = sessionmaker(bind=self.engine) @@ -428,6 +436,9 @@ class WorkflowInstanceServiceTests(unittest.TestCase): Base.metadata.drop_all( self.engine, tables=[ + WorkflowWaitState.__table__, + WorkflowTriggerDelivery.__table__, + WorkflowTrigger.__table__, WorkflowInstanceEvent.__table__, WorkflowInstanceStep.__table__, WorkflowInstance.__table__, diff --git a/tests/test_migrations.py b/tests/test_migrations.py index 6647103..3814fbc 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -29,7 +29,7 @@ class WorkflowMigrationTests(unittest.TestCase): try: with engine.connect() as connection: self.assertIn( - "0b4e7c9a2d6f", + "b2e4f6a8c0d1", set(MigrationContext.configure(connection).get_current_heads()), ) self.assertEqual( @@ -39,6 +39,9 @@ class WorkflowMigrationTests(unittest.TestCase): "workflow_instance_events", "workflow_instance_steps", "workflow_instances", + "workflow_triggers", + "workflow_trigger_deliveries", + "workflow_wait_states", }, { name @@ -98,7 +101,7 @@ class WorkflowMigrationTests(unittest.TestCase): engine_manifest.migration_spec.script_location or "" ) for path in current_revisions.glob("*.py"): - if path.name.startswith("0b4e7c9a2d6f_"): + if path.name.startswith(("0b4e7c9a2d6f_", "b2e4f6a8c0d1_")): continue shutil.copy2(path, legacy_revisions / path.name) legacy_manifest = replace( @@ -123,19 +126,13 @@ class WorkflowMigrationTests(unittest.TestCase): with engine.connect() as connection: self.assertIn( "f1b7d3e5a9c2", - set( - MigrationContext.configure( - connection - ).get_current_heads() - ), + set(MigrationContext.configure(connection).get_current_heads()), ) self.assertNotIn( "standard_origin_module_id", { item["name"] - for item in inspect(engine).get_columns( - "workflow_definitions" - ) + for item in inspect(engine).get_columns("workflow_definitions") }, ) finally: @@ -147,17 +144,23 @@ class WorkflowMigrationTests(unittest.TestCase): manifest_factories=(get_manifest,), ) - self.assertIn("0b4e7c9a2d6f", result.current_revision or "") + self.assertIn("b2e4f6a8c0d1", result.current_revision or "") engine = create_engine(url) try: - self.assertEqual(first_tables, set(inspect(engine).get_table_names())) + upgraded_tables = set(inspect(engine).get_table_names()) + self.assertTrue(first_tables.issubset(upgraded_tables)) + self.assertTrue( + { + "workflow_triggers", + "workflow_trigger_deliveries", + "workflow_wait_states", + }.issubset(upgraded_tables) + ) self.assertIn( "standard_origin_module_id", { item["name"] - for item in inspect(engine).get_columns( - "workflow_definitions" - ) + for item in inspect(engine).get_columns("workflow_definitions") }, ) finally: diff --git a/tests/test_triggers.py b/tests/test_triggers.py new file mode 100644 index 0000000..b61cb5e --- /dev/null +++ b/tests/test_triggers.py @@ -0,0 +1,316 @@ +from __future__ import annotations + +from datetime import UTC, datetime, timedelta +import unittest + +from sqlalchemy import create_engine, 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 AutomationPrincipalResolution +from govoplan_core.core.events import EventTenantRef, PlatformEvent +from govoplan_core.db.base import Base +from govoplan_workflow_engine.backend.db.models import ( + WorkflowDefinition, + WorkflowDefinitionRevision, + WorkflowInstance, + WorkflowInstanceEvent, + WorkflowInstanceStep, + WorkflowTrigger, + WorkflowTriggerDelivery, + WorkflowWaitState, +) +from govoplan_workflow_engine.backend.instance_service import start_instance +from govoplan_workflow_engine.backend.schemas import ( + WorkflowDefinitionCreateRequest, + WorkflowEdge, + WorkflowGraph, + WorkflowInstanceStartRequest, + WorkflowNode, +) +from govoplan_workflow_engine.backend.service import ( + activate_definition, + create_definition, +) +from govoplan_workflow_engine.backend.triggers import ( + SqlWorkflowTriggerDispatcher, + reconcile_definition_triggers, +) + + +def principal() -> ApiPrincipal: + return ApiPrincipal( + principal=PrincipalRef( + account_id="account-1", + membership_id="membership-1", + tenant_id="tenant-1", + scopes=frozenset({"workflow:instance:start"}), + ), + account=object(), + user=object(), + ) + + +class AutomationProvider: + def resolve_automation_principal(self, _session, *, request): + return AutomationPrincipalResolution( + allowed=True, + principal=principal(), + granted_scopes=request.grant_scopes, + provenance={"status": "rechecked"}, + ) + + +class Registry: + def __init__(self) -> None: + self.provider = AutomationProvider() + + def has_capability(self, name: str) -> bool: + return name == CAPABILITY_AUTH_AUTOMATION_PRINCIPAL_PROVIDER + + def capability(self, name: str): + if not self.has_capability(name): + raise KeyError(name) + return self.provider + + +def graph(start_type: str, *, wait: WorkflowNode | None = None) -> WorkflowGraph: + start_config: dict[str, object] = {} + if start_type == "workflow.start.schedule": + start_config = {"schedule": "interval:60", "timezone": "UTC"} + elif start_type == "workflow.start.event": + start_config = { + "event_type": "case.updated", + "filter": {"payload": {"state": "ready"}}, + } + nodes = [WorkflowNode(id="start", type=start_type, config=start_config)] + edges: list[WorkflowEdge] = [] + previous = "start" + if wait is not None: + nodes.append(wait) + edges.append(WorkflowEdge(id="start-wait", source="start", target=wait.id)) + previous = wait.id + nodes.append(WorkflowNode(id="end", type="workflow.end.completed")) + edges.append( + WorkflowEdge( + id="to-end", + source=previous, + source_port="timed_out" if wait is not None else "output", + target="end", + ) + ) + return WorkflowGraph(nodes=nodes, edges=edges) + + +class WorkflowTriggerTests(unittest.TestCase): + def setUp(self) -> None: + self.engine = create_engine("sqlite:///:memory:") + self.tables = [ + WorkflowDefinition.__table__, + WorkflowDefinitionRevision.__table__, + WorkflowInstance.__table__, + WorkflowInstanceStep.__table__, + WorkflowInstanceEvent.__table__, + WorkflowTrigger.__table__, + WorkflowTriggerDelivery.__table__, + WorkflowWaitState.__table__, + ] + Base.metadata.create_all(self.engine, tables=self.tables) + self.Session = sessionmaker(bind=self.engine) + self.session: Session = self.Session() + self.registry = Registry() + + def tearDown(self) -> None: + self.session.close() + Base.metadata.drop_all(self.engine, tables=list(reversed(self.tables))) + self.engine.dispose() + + def _definition( + self, + definition_graph: WorkflowGraph, + *, + automation: bool, + ) -> WorkflowDefinition: + definition = create_definition( + self.session, + tenant_id="tenant-1", + actor_id="account-1", + payload=WorkflowDefinitionCreateRequest( + name="Trigger test", + graph=definition_graph, + allow_automation=automation, + ), + ) + activate_definition( + self.session, + tenant_id="tenant-1", + definition_id=definition.id, + actor_id="account-1", + ) + return definition + + def test_schedule_registration_dispatch_and_replay_are_durable(self) -> None: + definition = self._definition( + graph("workflow.start.schedule"), + automation=True, + ) + summary = reconcile_definition_triggers( + self.session, + definition=definition, + principal=principal(), + registry=self.registry, + ) + trigger = self.session.scalar(select(WorkflowTrigger)) + assert trigger is not None + trigger.next_fire_at = datetime.now(tz=UTC) - timedelta(seconds=1) + + result = SqlWorkflowTriggerDispatcher(registry=self.registry).dispatch_due( + self.session, now=datetime.now(tz=UTC) + ) + + self.assertEqual({"created": 1, "updated": 0, "disabled": 0}, summary) + self.assertEqual(1, result["started"]) + delivery = self.session.scalar(select(WorkflowTriggerDelivery)) + assert delivery is not None + self.assertEqual("succeeded", delivery.status) + instance = self.session.get(WorkflowInstance, delivery.instance_id) + assert instance is not None + self.assertEqual("schedule", instance.start_origin) + self.assertEqual("completed", instance.status) + + replay = SqlWorkflowTriggerDispatcher(registry=self.registry).dispatch_due( + self.session, now=datetime.now(tz=UTC) + ) + self.assertEqual(0, replay["started"]) + self.assertEqual(1, self.session.query(WorkflowInstance).count()) + + def test_event_filter_queues_only_matching_event(self) -> None: + definition = self._definition( + graph("workflow.start.event"), + automation=True, + ) + reconcile_definition_triggers( + self.session, + definition=definition, + principal=principal(), + registry=self.registry, + ) + dispatcher = SqlWorkflowTriggerDispatcher(registry=self.registry) + ignored = dispatcher.ingest_event( + self.session, + event=PlatformEvent( + type="case.updated", + module_id="cases", + tenant=EventTenantRef(id="tenant-1"), + payload={"state": "draft"}, + ), + ) + accepted = dispatcher.ingest_event( + self.session, + event=PlatformEvent( + type="case.updated", + module_id="cases", + tenant=EventTenantRef(id="tenant-1"), + payload={"state": "ready"}, + ), + ) + result = dispatcher.dispatch_due(self.session) + + self.assertEqual(0, ignored["trigger_deliveries"]) + self.assertEqual(1, accepted["trigger_deliveries"]) + self.assertEqual(1, result["started"]) + + def test_duration_wait_resumes_through_persisted_timer(self) -> None: + definition = self._definition( + graph( + "workflow.start.manual", + wait=WorkflowNode( + id="wait", + type="workflow.wait", + config={"mode": "duration", "value": "1"}, + ), + ), + automation=False, + ) + instance, replayed = start_instance( + self.session, + tenant_id="tenant-1", + definition_id=definition.id, + actor_id="account-1", + principal=principal(), + registry=self.registry, + payload=WorkflowInstanceStartRequest(idempotency_key="wait-1"), + ) + state = self.session.scalar(select(WorkflowWaitState)) + assert state is not None + + result = SqlWorkflowTriggerDispatcher(registry=self.registry).dispatch_due( + self.session, + now=datetime.now(tz=UTC) + timedelta(seconds=2), + ) + + self.assertFalse(replayed) + self.assertEqual("timed_out", state.status) + self.assertEqual(1, result["waits_timed_out"]) + self.assertEqual("completed", instance.status) + + def test_parent_workflow_outcome_starts_pinned_child(self) -> None: + parent = self._definition( + graph("workflow.start.manual"), + automation=False, + ) + child_graph = WorkflowGraph( + nodes=[ + WorkflowNode( + id="start", + type="workflow.start.workflow", + config={ + "parent_definition_ref": (f"workflow-definition:{parent.id}"), + "parent_outcome": "completed", + "input_mapping": {"parent_id": "$event.payload.instance_id"}, + }, + ), + WorkflowNode(id="end", type="workflow.end.completed"), + ], + edges=[WorkflowEdge(id="finish", source="start", target="end")], + ) + child = self._definition(child_graph, automation=True) + reconcile_definition_triggers( + self.session, + definition=child, + principal=principal(), + registry=self.registry, + ) + dispatcher = SqlWorkflowTriggerDispatcher(registry=self.registry) + event = PlatformEvent( + type="workflow.instance.completed", + module_id="workflow_engine", + tenant=EventTenantRef(id="tenant-1"), + payload={ + "instance_id": "parent-instance-1", + "definition_id": parent.id, + }, + ) + + queued = dispatcher.ingest_event(self.session, event=event) + result = dispatcher.dispatch_due(self.session) + + self.assertEqual(1, queued["trigger_deliveries"]) + self.assertEqual(1, result["started"]) + instance = self.session.scalar( + select(WorkflowInstance).where(WorkflowInstance.definition_id == child.id) + ) + assert instance is not None + self.assertEqual("parent_workflow", instance.start_origin) + self.assertEqual( + "parent-instance-1", + instance.input_["parent_id"], + ) + + +if __name__ == "__main__": + unittest.main()