Skip to main content

PostgresWorkflowStore

Struct PostgresWorkflowStore 

Source
pub struct PostgresWorkflowStore { /* private fields */ }
Expand description

PostgreSQL-backed distributed workflow task store.

Claim, heartbeat, and finish operations use the database clock. Every ownership mutation compares worker identity, fencing token, and lease expiration.

Implementations§

Source§

impl PostgresWorkflowStore

Source

pub async fn connect( connection: &str, table: &str, ) -> Result<Self, PostgresWorkflowStoreError>

Connects without creating or changing schema.

§Errors

Rejects unsafe table identifiers and propagates connection failures.

Source

pub async fn ensure_schema(&self) -> Result<(), PostgresWorkflowStoreError>

Explicitly creates the workflow task table and claim index.

Runtime queue operations never perform hidden migrations.

§Errors

Propagates PostgreSQL DDL failures.

Trait Implementations§

Source§

impl Clone for PostgresWorkflowStore

Source§

fn clone(&self) -> PostgresWorkflowStore

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for PostgresWorkflowStore

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl WorkflowStore for PostgresWorkflowStore

Source§

fn current_time_ms( &self, ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>>

Returns the store-authoritative Unix time in milliseconds.
Source§

fn set_tenant_policy( &self, tenant_id: WorkflowTenantId, policy: WorkflowTenantPolicy, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Creates or replaces one tenant’s admission policy.
Source§

fn set_tenant_budget_policy( &self, tenant_id: WorkflowTenantId, policy: WorkflowTenantBudgetPolicy, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Creates or replaces one tenant’s persistent aggregate budget policy.
Source§

fn list_tenant_budgets( &self, after: Option<WorkflowTenantId>, limit: WorkflowTenantListLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTenantId>, WorkflowStoreError>>

Discovers budget-enabled tenants in stable identity order.
Source§

fn inspect_tenant_budget( &self, tenant_id: WorkflowTenantId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTenantBudgetSnapshot, WorkflowStoreError>>

Reads a tenant budget after reclaiming expired reservations.
Source§

fn list_tenant_budget_audit( &self, tenant_id: WorkflowTenantId, after: Option<WorkflowBudgetAuditCursor>, limit: WorkflowBudgetAuditLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowBudgetAuditEvent>, WorkflowStoreError>>

Reads durable budget decisions strictly after an optional cursor.
Source§

fn compact_tenant_budget_audit( &self, tenant_id: WorkflowTenantId, through: WorkflowBudgetAuditCursor, ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>>

Deletes tenant audit facts at or before an explicitly acknowledged cursor.
Source§

fn load_or_create_tenant_budget_audit_projection( &self, tenant_id: WorkflowTenantId, projection_id: WorkflowBudgetAuditProjectionId, ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditCursor, WorkflowStoreError>>

Loads or atomically registers one named consumer at cursor zero.
Source§

fn advance_tenant_budget_audit_projection( &self, tenant_id: WorkflowTenantId, projection_id: WorkflowBudgetAuditProjectionId, expected: WorkflowBudgetAuditCursor, next: WorkflowBudgetAuditCursor, ) -> WorkflowStoreFuture<'_, Result<bool, WorkflowStoreError>>

Monotonically advances a projection cursor using compare-and-set. Read more
Source§

fn claim_tenant_budget_audit_projection( &self, tenant_id: WorkflowTenantId, projection_id: WorkflowBudgetAuditProjectionId, owner: WorkerId, lease: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<Option<WorkflowBudgetAuditProjectionLease>, WorkflowStoreError>>

Exclusively claims an idle or expired named audit projection.
Source§

fn heartbeat_tenant_budget_audit_projection( &self, lease: WorkflowBudgetAuditProjectionLease, extension: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>

Extends an active projection lease under its current fencing token.
Source§

fn advance_tenant_budget_audit_projection_lease( &self, lease: WorkflowBudgetAuditProjectionLease, next: WorkflowBudgetAuditCursor, ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>

Advances a projection cursor only while its fenced lease remains active.
Source§

fn release_tenant_budget_audit_projection( &self, lease: WorkflowBudgetAuditProjectionLease, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Releases a currently fenced projection without changing its cursor.
Source§

fn reserve_budget( &self, lease: WorkflowLease, workflow_limit: Budget, baseline: Usage, ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetReservationOutcome, WorkflowStoreError>>

Idempotently reserves the remaining workflow envelope under a lease.
Source§

fn settle_budget( &self, lease: WorkflowLease, cumulative: Usage, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Commits observed cumulative usage and releases unused reservation.
Source§

fn enqueue( &self, task: WorkflowTask, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Enqueues a task exactly once.
Source§

fn claim( &self, worker: WorkerId, lease: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<Option<ClaimedWorkflow>, WorkflowStoreError>>

Atomically claims the highest-priority eligible task.
Source§

fn heartbeat( &self, lease: WorkflowLease, extension: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<WorkflowLease, WorkflowStoreError>>

Extends a currently owned, unexpired lease.
Source§

fn finish( &self, lease: WorkflowLease, disposition: WorkflowDisposition, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Applies a terminal or retry disposition under the current lease.
Source§

fn publish_signal( &self, tenant_id: WorkflowTenantId, signal: WorkflowSignal, ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalOutcome, WorkflowStoreError>>

Idempotently publishes an external signal, buffering it when necessary.
Source§

fn publish_control_signal( &self, tenant_id: WorkflowTenantId, signal: WorkflowSignal, ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalOutcome, WorkflowStoreError>>

Publishes durable coordination metadata excluded from signal retention. Read more
Source§

fn cancel( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, ) -> WorkflowStoreFuture<'_, Result<WorkflowCancelOutcome, WorkflowStoreError>>

Idempotently cancels queued, waiting, or currently leased work.
Source§

fn inspect_signal( &self, tenant_id: WorkflowTenantId, signal_id: WorkflowSignalId, ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalSnapshot, WorkflowStoreError>>

Loads safe signal lifecycle metadata without exposing its payload.
Source§

fn load_signal_payload( &self, tenant_id: WorkflowTenantId, signal_id: WorkflowSignalId, ) -> WorkflowStoreFuture<'_, Result<Value, WorkflowStoreError>>

Loads one accepted signal payload under tenant authorization. Read more
Source§

fn compact_signals( &self, tenant_id: WorkflowTenantId, retention: WorkflowSignalRetention, ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>>

Deletes only consumed or dead-letter signals older than retention.
Source§

fn inspect( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskSnapshot, WorkflowStoreError>>

Loads safe control-plane state for inspection.
Source§

fn load_task_input( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, ) -> WorkflowStoreFuture<'_, Result<Value, WorkflowStoreError>>

Loads the immutable original task input under tenant authorization. Read more
Source§

fn list_checkpoint_history( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, after_revision: Option<u64>, limit: WorkflowCheckpointHistoryLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowCheckpointRevision>, WorkflowStoreError>>

Lists immutable checkpoint revisions after an optional revision cursor.
Source§

fn load_checkpoint_revision( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, revision: u64, ) -> WorkflowStoreFuture<'_, Result<WorkflowCheckpointRevision, WorkflowStoreError>>

Loads one exact immutable checkpoint revision for state inspection.
Source§

fn fork_workflow( &self, tenant_id: WorkflowTenantId, command: WorkflowForkCommand, ) -> WorkflowStoreFuture<'_, Result<WorkflowForkOutcome, WorkflowStoreError>>

Idempotently creates a new execution branch from immutable history.
Source§

fn load_checkpoint( &self, lease: WorkflowLease, ) -> WorkflowStoreFuture<'_, Result<Checkpoint, CheckpointError>>

Loads a checkpoint under a current worker lease.
Source§

fn compare_and_swap_checkpoint( &self, lease: WorkflowLease, checkpoint: Checkpoint, expected_revision: Option<u64>, ) -> WorkflowStoreFuture<'_, Result<(), CheckpointError>>

Creates or compare-and-swaps a checkpoint under a current worker lease.
Source§

fn decide_interrupt( &self, tenant_id: WorkflowTenantId, command: WorkflowInterruptCommand, ) -> Pin<Box<dyn Future<Output = Result<WorkflowInterruptDecisionOutcome, WorkflowStoreError>> + Send + '_>>

Idempotently applies a typed human decision to a durable interrupt.
Source§

impl WorkflowTaskRetentionStore for PostgresWorkflowStore

Source§

fn list_task_cleanup_tenants( &self, after: Option<WorkflowTenantId>, limit: WorkflowTenantListLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTenantId>, WorkflowStoreError>>

Discovers tenants that currently own terminal Tasks. Read more
Source§

fn claim_task_cleanup( &self, tenant_id: WorkflowTenantId, owner: WorkerId, lease: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<Option<WorkflowTaskCleanupLease>, WorkflowStoreError>>

Claims one tenant’s cleanup partition if it is idle or expired.
Source§

fn compact_terminal_tasks( &self, lease: WorkflowTaskCleanupLease, retention: WorkflowTaskRetention, limit: WorkflowTaskCleanupLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTaskTombstone>, WorkflowStoreError>>

Atomically tombstones and removes one bounded terminal Task batch.
Source§

fn heartbeat_task_cleanup( &self, lease: WorkflowTaskCleanupLease, extension: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskCleanupLease, WorkflowStoreError>>

Extends an exact current cleanup lease using store-authoritative time.
Source§

fn list_task_tombstones( &self, tenant_id: WorkflowTenantId, after: Option<WorkflowTaskTombstoneCursor>, limit: WorkflowTaskTombstoneLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTaskTombstone>, WorkflowStoreError>>

Lists immutable tombstones strictly after an optional cursor.
Source§

fn release_task_cleanup( &self, lease: WorkflowTaskCleanupLease, ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>>

Releases a current unexpired cleanup lease.
Source§

impl WorkflowTaskTombstoneGovernanceStore for PostgresWorkflowStore

Source§

fn place_task_tombstone_hold( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, actor: WorkerId, reason: WorkflowTaskLegalHoldReason, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>

Places or idempotently reads a legal hold on an existing tombstone.
Source§

fn release_task_tombstone_hold( &self, tenant_id: WorkflowTenantId, checkpoint_id: CheckpointId, actor: WorkerId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>>

Releases an exact active legal hold while retaining its audit row.
Source§

fn confirm_task_tombstone_export( &self, tenant_id: WorkflowTenantId, through: WorkflowTaskTombstoneCursor, receipt: WorkflowTaskTombstoneExportReceipt, actor: WorkerId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneExport, WorkflowStoreError>>

Monotonically confirms an externally archived tenant cursor prefix.
Source§

fn prepare_task_tombstone_purge( &self, lease: WorkflowTaskCleanupLease, retention: WorkflowTaskTombstoneRetention, limit: WorkflowTaskTombstonePurgeLimit, approval_window: WorkflowTaskTombstoneApprovalWindow, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>

Freezes one bounded, exported, unheld, old-enough purge candidate set.
Source§

fn approve_task_tombstone_purge( &self, tenant_id: WorkflowTenantId, purge_id: WorkflowTaskTombstonePurgeId, approver: WorkerId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>

Approves a pending intent using a different principal from its preparer.
Source§

fn list_task_tombstone_purge_approvals( &self, tenant_id: WorkflowTenantId, limit: WorkflowTaskTombstoneApprovalInboxLimit, ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTaskTombstoneApprovalInboxItem>, WorkflowStoreError>>

Lists a bounded tenant approval inbox with expired claims normalized.
Source§

fn claim_task_tombstone_purge_approval( &self, tenant_id: WorkflowTenantId, reviewer: WorkerId, lease: LeaseDuration, ) -> WorkflowStoreFuture<'_, Result<Option<WorkflowTaskTombstoneApprovalLease>, WorkflowStoreError>>

Atomically claims the oldest eligible request for an independent reviewer.
Source§

fn approve_claimed_task_tombstone_purge( &self, lease: WorkflowTaskTombstoneApprovalLease, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>>

Approves under an exact, unexpired, fenced reviewer lease.
Source§

fn reject_claimed_task_tombstone_purge( &self, lease: WorkflowTaskTombstoneApprovalLease, reason: WorkflowTaskTombstoneRejectionReason, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneApprovalInboxItem, WorkflowStoreError>>

Rejects under an exact reviewer lease and preserves the reason.
Source§

fn execute_task_tombstone_purge( &self, lease: WorkflowTaskCleanupLease, purge_id: WorkflowTaskTombstonePurgeId, ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeEvidence, WorkflowStoreError>>

Executes an approved intent under a current fenced cleanup lease. Read more
Source§

fn get_task_tombstone_purge_evidence( &self, tenant_id: WorkflowTenantId, purge_id: WorkflowTaskTombstonePurgeId, ) -> WorkflowStoreFuture<'_, Result<Option<WorkflowTaskTombstonePurgeEvidence>, WorkflowStoreError>>

Reads immutable evidence for one executed purge.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.