Skip to main content

runifold_store_postgres/
workflow.rs

1//! `PostgreSQL` distributed workflow task-control adapter for Runifold.
2
3mod budget;
4mod budget_store;
5mod checkpoint;
6mod codec;
7mod lifecycle;
8mod retention;
9mod signal;
10mod sql;
11mod support;
12mod tombstone;
13
14use budget::{
15    StoredBudgetLimit, budget_decoding, budget_encoding, decode_budget_reservation_status,
16    decode_budget_settlement_status, postgres_budget_request,
17};
18use checkpoint::CheckpointStoreExt;
19use codec::{
20    ForkSource, PreparedFork, decode_budget_audit_event, decode_budget_audit_projection_lease,
21    decode_budget_snapshot, decode_snapshot, fork_storage_fields,
22};
23use support::{
24    checkpoint_decoding, checkpoint_domain_error, checkpoint_encoding, checkpoint_i64,
25    checkpoint_lease_lost, checkpoint_storage, database_i64, decode_u64, lease_lost,
26    projection_lease_lost, storage, tenant_mismatch, validate_identifier,
27};
28
29use std::sync::Arc;
30
31use runifold_core::{
32    Budget, Checkpoint, CheckpointError, CheckpointErrorKind, CheckpointId, Usage,
33};
34use runifold_workflow::{
35    ClaimedWorkflow, LeaseDuration, WorkerId, WorkflowBudgetAuditCursor, WorkflowBudgetAuditEvent,
36    WorkflowBudgetAuditLimit, WorkflowBudgetAuditProjectionId, WorkflowBudgetAuditProjectionLease,
37    WorkflowBudgetReservationOutcome, WorkflowCancelOutcome, WorkflowCheckpointHistoryLimit,
38    WorkflowCheckpointRevision, WorkflowDisposition, WorkflowForkCommand, WorkflowForkOutcome,
39    WorkflowLease, WorkflowLineage, WorkflowSignal, WorkflowSignalId, WorkflowSignalOutcome,
40    WorkflowSignalRetention, WorkflowSignalSnapshot, WorkflowStore, WorkflowStoreError,
41    WorkflowStoreErrorKind, WorkflowStoreFuture, WorkflowTask, WorkflowTaskCleanupLease,
42    WorkflowTaskCleanupLimit, WorkflowTaskLegalHold, WorkflowTaskLegalHoldReason,
43    WorkflowTaskRetention, WorkflowTaskRetentionStore, WorkflowTaskSnapshot, WorkflowTaskTombstone,
44    WorkflowTaskTombstoneApprovalInboxItem, WorkflowTaskTombstoneApprovalInboxLimit,
45    WorkflowTaskTombstoneApprovalLease, WorkflowTaskTombstoneApprovalWindow,
46    WorkflowTaskTombstoneCursor, WorkflowTaskTombstoneExport, WorkflowTaskTombstoneExportReceipt,
47    WorkflowTaskTombstoneGovernanceStore, WorkflowTaskTombstoneLimit,
48    WorkflowTaskTombstonePurgeEvidence, WorkflowTaskTombstonePurgeId,
49    WorkflowTaskTombstonePurgeIntent, WorkflowTaskTombstonePurgeLimit,
50    WorkflowTaskTombstoneRejectionReason, WorkflowTaskTombstoneRetention,
51    WorkflowTenantBudgetPolicy, WorkflowTenantBudgetSnapshot, WorkflowTenantId,
52    WorkflowTenantListLimit, WorkflowTenantPolicy,
53};
54use serde_json::Value;
55use thiserror::Error;
56use tokio_postgres::{Client, NoTls, error::SqlState};
57
58/// `PostgreSQL` workflow-store configuration or connection failure.
59#[derive(Debug, Error)]
60#[non_exhaustive]
61pub enum PostgresWorkflowStoreError {
62    /// The configured table name is unsafe for SQL interpolation.
63    #[error("workflow table must be a portable PostgreSQL identifier of at most 48 bytes")]
64    InvalidTable,
65    /// `PostgreSQL` connection or schema setup failed.
66    #[error("PostgreSQL workflow store operation failed: {0}")]
67    Database(#[from] tokio_postgres::Error),
68}
69
70/// PostgreSQL-backed distributed workflow task store.
71///
72/// Claim, heartbeat, and finish operations use the database clock. Every
73/// ownership mutation compares worker identity, fencing token, and lease
74/// expiration.
75#[derive(Clone, Debug)]
76pub struct PostgresWorkflowStore {
77    client: Arc<Client>,
78    table: String,
79}
80
81impl PostgresWorkflowStore {
82    /// Connects without creating or changing schema.
83    ///
84    /// # Errors
85    ///
86    /// Rejects unsafe table identifiers and propagates connection failures.
87    pub async fn connect(
88        connection: &str,
89        table: &str,
90    ) -> Result<Self, PostgresWorkflowStoreError> {
91        validate_identifier(table)?;
92        let (client, connection) = tokio_postgres::connect(connection, NoTls).await?;
93        tokio::spawn(async move {
94            let _ = connection.await;
95        });
96        Ok(Self {
97            client: Arc::new(client),
98            table: table.into(),
99        })
100    }
101
102    /// Explicitly creates the workflow task table and claim index.
103    ///
104    /// Runtime queue operations never perform hidden migrations.
105    ///
106    /// # Errors
107    ///
108    /// Propagates `PostgreSQL` DDL failures.
109    pub async fn ensure_schema(&self) -> Result<(), PostgresWorkflowStoreError> {
110        self.client
111            .batch_execute(&Self::task_schema_sql(&self.table))
112            .await?;
113        self.client
114            .batch_execute(&Self::checkpoint_history_schema_sql(&self.table))
115            .await?;
116        self.client
117            .batch_execute(&Self::signal_schema_sql(&self.table))
118            .await?;
119        self.client
120            .batch_execute(&Self::budget_schema_sql(&self.table))
121            .await?;
122        self.client
123            .batch_execute(&Self::task_retention_schema_sql(&self.table))
124            .await?;
125        Ok(())
126    }
127}
128
129impl WorkflowTaskRetentionStore for PostgresWorkflowStore {
130    fn list_task_cleanup_tenants(
131        &self,
132        after: Option<WorkflowTenantId>,
133        limit: WorkflowTenantListLimit,
134    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTenantId>, WorkflowStoreError>> {
135        retention::list_tenants(self, after, limit)
136    }
137
138    fn claim_task_cleanup(
139        &self,
140        tenant_id: WorkflowTenantId,
141        owner: WorkerId,
142        lease: LeaseDuration,
143    ) -> WorkflowStoreFuture<'_, Result<Option<WorkflowTaskCleanupLease>, WorkflowStoreError>> {
144        retention::claim(self, tenant_id, owner, lease)
145    }
146
147    fn compact_terminal_tasks(
148        &self,
149        lease: WorkflowTaskCleanupLease,
150        retention: WorkflowTaskRetention,
151        limit: WorkflowTaskCleanupLimit,
152    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTaskTombstone>, WorkflowStoreError>> {
153        retention::compact(self, lease, retention, limit)
154    }
155
156    fn heartbeat_task_cleanup(
157        &self,
158        lease: WorkflowTaskCleanupLease,
159        extension: LeaseDuration,
160    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskCleanupLease, WorkflowStoreError>> {
161        retention::heartbeat(self, lease, extension)
162    }
163
164    fn list_task_tombstones(
165        &self,
166        tenant_id: WorkflowTenantId,
167        after: Option<WorkflowTaskTombstoneCursor>,
168        limit: WorkflowTaskTombstoneLimit,
169    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTaskTombstone>, WorkflowStoreError>> {
170        retention::list(self, tenant_id, after, limit)
171    }
172
173    fn release_task_cleanup(
174        &self,
175        lease: WorkflowTaskCleanupLease,
176    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
177        retention::release(self, lease)
178    }
179}
180
181impl WorkflowTaskTombstoneGovernanceStore for PostgresWorkflowStore {
182    fn place_task_tombstone_hold(
183        &self,
184        tenant_id: WorkflowTenantId,
185        checkpoint_id: CheckpointId,
186        actor: WorkerId,
187        reason: WorkflowTaskLegalHoldReason,
188    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>> {
189        tombstone::place_hold(self, tenant_id, checkpoint_id, actor, reason)
190    }
191
192    fn release_task_tombstone_hold(
193        &self,
194        tenant_id: WorkflowTenantId,
195        checkpoint_id: CheckpointId,
196        actor: WorkerId,
197    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskLegalHold, WorkflowStoreError>> {
198        tombstone::release_hold(self, tenant_id, checkpoint_id, actor)
199    }
200
201    fn confirm_task_tombstone_export(
202        &self,
203        tenant_id: WorkflowTenantId,
204        through: WorkflowTaskTombstoneCursor,
205        receipt: WorkflowTaskTombstoneExportReceipt,
206        actor: WorkerId,
207    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneExport, WorkflowStoreError>> {
208        tombstone::confirm_export(self, tenant_id, through, receipt, actor)
209    }
210
211    fn prepare_task_tombstone_purge(
212        &self,
213        lease: WorkflowTaskCleanupLease,
214        retention: WorkflowTaskTombstoneRetention,
215        limit: WorkflowTaskTombstonePurgeLimit,
216        approval_window: WorkflowTaskTombstoneApprovalWindow,
217    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>> {
218        tombstone::prepare_purge(self, lease, retention, limit, approval_window)
219    }
220
221    fn approve_task_tombstone_purge(
222        &self,
223        tenant_id: WorkflowTenantId,
224        purge_id: WorkflowTaskTombstonePurgeId,
225        approver: WorkerId,
226    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>> {
227        tombstone::approve_purge(self, tenant_id, purge_id, approver)
228    }
229
230    fn list_task_tombstone_purge_approvals(
231        &self,
232        tenant_id: WorkflowTenantId,
233        limit: WorkflowTaskTombstoneApprovalInboxLimit,
234    ) -> WorkflowStoreFuture<
235        '_,
236        Result<Vec<WorkflowTaskTombstoneApprovalInboxItem>, WorkflowStoreError>,
237    > {
238        tombstone::list_approvals(self, tenant_id, limit)
239    }
240
241    fn claim_task_tombstone_purge_approval(
242        &self,
243        tenant_id: WorkflowTenantId,
244        reviewer: WorkerId,
245        lease: LeaseDuration,
246    ) -> WorkflowStoreFuture<
247        '_,
248        Result<Option<WorkflowTaskTombstoneApprovalLease>, WorkflowStoreError>,
249    > {
250        tombstone::claim_approval(self, tenant_id, reviewer, lease)
251    }
252
253    fn approve_claimed_task_tombstone_purge(
254        &self,
255        lease: WorkflowTaskTombstoneApprovalLease,
256    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeIntent, WorkflowStoreError>> {
257        tombstone::approve_claimed(self, lease)
258    }
259
260    fn reject_claimed_task_tombstone_purge(
261        &self,
262        lease: WorkflowTaskTombstoneApprovalLease,
263        reason: WorkflowTaskTombstoneRejectionReason,
264    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstoneApprovalInboxItem, WorkflowStoreError>>
265    {
266        tombstone::reject_claimed(self, lease, reason)
267    }
268
269    fn execute_task_tombstone_purge(
270        &self,
271        lease: WorkflowTaskCleanupLease,
272        purge_id: WorkflowTaskTombstonePurgeId,
273    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskTombstonePurgeEvidence, WorkflowStoreError>>
274    {
275        tombstone::execute_purge(self, lease, purge_id)
276    }
277
278    fn get_task_tombstone_purge_evidence(
279        &self,
280        tenant_id: WorkflowTenantId,
281        purge_id: WorkflowTaskTombstonePurgeId,
282    ) -> WorkflowStoreFuture<
283        '_,
284        Result<Option<WorkflowTaskTombstonePurgeEvidence>, WorkflowStoreError>,
285    > {
286        tombstone::get_evidence(self, tenant_id, purge_id)
287    }
288}
289
290impl WorkflowStore for PostgresWorkflowStore {
291    fn set_tenant_policy(
292        &self,
293        tenant_id: WorkflowTenantId,
294        policy: WorkflowTenantPolicy,
295    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
296        Box::pin(async move {
297            let outstanding = i64::from(policy.max_outstanding_tasks());
298            let leases = i64::from(policy.max_concurrent_leases());
299            self.client
300                .execute(
301                    &format!(
302                        "INSERT INTO {table}_tenants (
303                            tenant_id, max_outstanding_tasks, max_concurrent_leases
304                         ) VALUES ($1, $2, $3)
305                         ON CONFLICT (tenant_id) DO UPDATE SET
306                            max_outstanding_tasks = EXCLUDED.max_outstanding_tasks,
307                            max_concurrent_leases = EXCLUDED.max_concurrent_leases,
308                            updated_at = clock_timestamp()",
309                        table = self.table
310                    ),
311                    &[&tenant_id.as_str(), &outstanding, &leases],
312                )
313                .await
314                .map_err(storage)?;
315            Ok(())
316        })
317    }
318
319    fn set_tenant_budget_policy(
320        &self,
321        tenant_id: WorkflowTenantId,
322        policy: WorkflowTenantBudgetPolicy,
323    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
324        self.set_tenant_budget_policy_inner(tenant_id, policy)
325    }
326
327    fn list_tenant_budgets(
328        &self,
329        after: Option<WorkflowTenantId>,
330        limit: WorkflowTenantListLimit,
331    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTenantId>, WorkflowStoreError>> {
332        self.list_tenant_budgets_inner(after, limit)
333    }
334
335    fn inspect_tenant_budget(
336        &self,
337        tenant_id: WorkflowTenantId,
338    ) -> WorkflowStoreFuture<'_, Result<WorkflowTenantBudgetSnapshot, WorkflowStoreError>> {
339        self.inspect_tenant_budget_inner(tenant_id)
340    }
341
342    fn list_tenant_budget_audit(
343        &self,
344        tenant_id: WorkflowTenantId,
345        after: Option<WorkflowBudgetAuditCursor>,
346        limit: WorkflowBudgetAuditLimit,
347    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowBudgetAuditEvent>, WorkflowStoreError>> {
348        self.list_tenant_budget_audit_inner(tenant_id, after, limit)
349    }
350
351    fn compact_tenant_budget_audit(
352        &self,
353        tenant_id: WorkflowTenantId,
354        through: WorkflowBudgetAuditCursor,
355    ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>> {
356        self.compact_tenant_budget_audit_inner(tenant_id, through)
357    }
358
359    fn load_or_create_tenant_budget_audit_projection(
360        &self,
361        tenant_id: WorkflowTenantId,
362        projection_id: WorkflowBudgetAuditProjectionId,
363    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditCursor, WorkflowStoreError>> {
364        self.load_or_create_tenant_budget_audit_projection_inner(tenant_id, projection_id)
365    }
366
367    fn advance_tenant_budget_audit_projection(
368        &self,
369        tenant_id: WorkflowTenantId,
370        projection_id: WorkflowBudgetAuditProjectionId,
371        expected: WorkflowBudgetAuditCursor,
372        next: WorkflowBudgetAuditCursor,
373    ) -> WorkflowStoreFuture<'_, Result<bool, WorkflowStoreError>> {
374        self.advance_tenant_budget_audit_projection_inner(tenant_id, projection_id, expected, next)
375    }
376
377    fn claim_tenant_budget_audit_projection(
378        &self,
379        tenant_id: WorkflowTenantId,
380        projection_id: WorkflowBudgetAuditProjectionId,
381        owner: WorkerId,
382        lease: LeaseDuration,
383    ) -> WorkflowStoreFuture<
384        '_,
385        Result<Option<WorkflowBudgetAuditProjectionLease>, WorkflowStoreError>,
386    > {
387        self.claim_tenant_budget_audit_projection_inner(tenant_id, projection_id, owner, lease)
388    }
389
390    fn heartbeat_tenant_budget_audit_projection(
391        &self,
392        lease: WorkflowBudgetAuditProjectionLease,
393        extension: LeaseDuration,
394    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>
395    {
396        self.heartbeat_tenant_budget_audit_projection_inner(lease, extension)
397    }
398
399    fn advance_tenant_budget_audit_projection_lease(
400        &self,
401        lease: WorkflowBudgetAuditProjectionLease,
402        next: WorkflowBudgetAuditCursor,
403    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>
404    {
405        self.advance_tenant_budget_audit_projection_lease_inner(lease, next)
406    }
407
408    fn release_tenant_budget_audit_projection(
409        &self,
410        lease: WorkflowBudgetAuditProjectionLease,
411    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
412        self.release_tenant_budget_audit_projection_inner(lease)
413    }
414
415    fn reserve_budget(
416        &self,
417        lease: WorkflowLease,
418        workflow_limit: Budget,
419        baseline: Usage,
420    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetReservationOutcome, WorkflowStoreError>> {
421        self.reserve_budget_inner(lease, workflow_limit, baseline)
422    }
423
424    fn settle_budget(
425        &self,
426        lease: WorkflowLease,
427        cumulative: Usage,
428    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
429        self.settle_budget_inner(lease, cumulative)
430    }
431    fn enqueue(
432        &self,
433        task: WorkflowTask,
434    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
435        Box::pin(async move {
436            task.validate()?;
437            let version = i32::try_from(task.workflow_version).map_err(|_| {
438                WorkflowStoreError::new(
439                    WorkflowStoreErrorKind::InvalidInput,
440                    "workflow version exceeds PostgreSQL INTEGER",
441                )
442            })?;
443            self.client
444                .execute(
445                    &format!(
446                        "INSERT INTO {table}_tenants (
447                            tenant_id, max_outstanding_tasks, max_concurrent_leases
448                         ) VALUES ($1, 10000, 100)
449                         ON CONFLICT (tenant_id) DO NOTHING",
450                        table = self.table
451                    ),
452                    &[&task.tenant_id.as_str()],
453                )
454                .await
455                .map_err(storage)?;
456            let inserted = self
457                .client
458                .execute(
459                    &format!(
460                        r"
461                        WITH admitted AS (
462                            UPDATE {table}_tenants
463                            SET
464                                outstanding_tasks = outstanding_tasks + 1,
465                                updated_at = clock_timestamp()
466                            WHERE tenant_id = $2
467                              AND outstanding_tasks < max_outstanding_tasks
468                            RETURNING tenant_id
469                        )
470                        INSERT INTO {table} (
471                            checkpoint_id, tenant_id, workflow, workflow_version,
472                            input, priority, state
473                        )
474                        SELECT $1, $2, $3, $4, $5, $6, 'queued'
475                        FROM admitted
476                        ",
477                        table = self.table
478                    ),
479                    &[
480                        &task.checkpoint_id.as_uuid(),
481                        &task.tenant_id.as_str(),
482                        &task.workflow,
483                        &version,
484                        &task.input,
485                        &task.priority,
486                    ],
487                )
488                .await;
489            match inserted {
490                Ok(1) => Ok(()),
491                Ok(_) => Err(WorkflowStoreError::new(
492                    WorkflowStoreErrorKind::AdmissionDenied,
493                    format!(
494                        "workflow tenant `{}` reached its outstanding task limit",
495                        task.tenant_id.as_str()
496                    ),
497                )),
498                Err(error)
499                    if error
500                        .as_db_error()
501                        .is_some_and(|db| db.code() == &SqlState::UNIQUE_VIOLATION) =>
502                {
503                    Err(WorkflowStoreError::new(
504                        WorkflowStoreErrorKind::Conflict,
505                        format!("workflow task `{}` already exists", task.checkpoint_id),
506                    ))
507                }
508                Err(error) => Err(storage(error)),
509            }
510        })
511    }
512
513    fn claim(
514        &self,
515        worker: WorkerId,
516        lease: LeaseDuration,
517    ) -> WorkflowStoreFuture<'_, Result<Option<ClaimedWorkflow>, WorkflowStoreError>> {
518        Box::pin(async move { self.claim_inner(worker, lease).await })
519    }
520
521    fn heartbeat(
522        &self,
523        lease: WorkflowLease,
524        extension: LeaseDuration,
525    ) -> WorkflowStoreFuture<'_, Result<WorkflowLease, WorkflowStoreError>> {
526        Box::pin(async move { self.heartbeat_inner(lease, extension).await })
527    }
528
529    fn finish(
530        &self,
531        lease: WorkflowLease,
532        disposition: WorkflowDisposition,
533    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
534        Box::pin(async move { self.finish_inner(&lease, disposition).await })
535    }
536
537    fn publish_signal(
538        &self,
539        tenant_id: WorkflowTenantId,
540        signal: WorkflowSignal,
541    ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalOutcome, WorkflowStoreError>> {
542        Box::pin(signal::publish(self, tenant_id, signal))
543    }
544
545    fn cancel(
546        &self,
547        tenant_id: WorkflowTenantId,
548        checkpoint_id: CheckpointId,
549    ) -> WorkflowStoreFuture<'_, Result<WorkflowCancelOutcome, WorkflowStoreError>> {
550        Box::pin(signal::cancel(self, tenant_id, checkpoint_id))
551    }
552
553    fn inspect_signal(
554        &self,
555        tenant_id: WorkflowTenantId,
556        signal_id: WorkflowSignalId,
557    ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalSnapshot, WorkflowStoreError>> {
558        Box::pin(signal::inspect(self, tenant_id, signal_id))
559    }
560
561    fn compact_signals(
562        &self,
563        tenant_id: WorkflowTenantId,
564        retention: WorkflowSignalRetention,
565    ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>> {
566        Box::pin(signal::compact(self, tenant_id, retention))
567    }
568    fn inspect(
569        &self,
570        tenant_id: WorkflowTenantId,
571        checkpoint_id: CheckpointId,
572    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskSnapshot, WorkflowStoreError>> {
573        self.inspect_ext(tenant_id, checkpoint_id)
574    }
575
576    fn list_checkpoint_history(
577        &self,
578        tenant_id: WorkflowTenantId,
579        checkpoint_id: CheckpointId,
580        after_revision: Option<u64>,
581        limit: WorkflowCheckpointHistoryLimit,
582    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowCheckpointRevision>, WorkflowStoreError>> {
583        self.list_checkpoint_history_ext(tenant_id, checkpoint_id, after_revision, limit)
584    }
585
586    fn load_checkpoint_revision(
587        &self,
588        tenant_id: WorkflowTenantId,
589        checkpoint_id: CheckpointId,
590        revision: u64,
591    ) -> WorkflowStoreFuture<'_, Result<WorkflowCheckpointRevision, WorkflowStoreError>> {
592        self.load_checkpoint_revision_ext(tenant_id, checkpoint_id, revision)
593    }
594
595    fn fork_workflow(
596        &self,
597        tenant_id: WorkflowTenantId,
598        command: WorkflowForkCommand,
599    ) -> WorkflowStoreFuture<'_, Result<WorkflowForkOutcome, WorkflowStoreError>> {
600        self.fork_workflow_ext(tenant_id, command)
601    }
602
603    fn load_checkpoint(
604        &self,
605        lease: WorkflowLease,
606    ) -> WorkflowStoreFuture<'_, Result<Checkpoint, CheckpointError>> {
607        self.load_checkpoint_ext(lease)
608    }
609
610    fn compare_and_swap_checkpoint(
611        &self,
612        lease: WorkflowLease,
613        checkpoint: Checkpoint,
614        expected_revision: Option<u64>,
615    ) -> WorkflowStoreFuture<'_, Result<(), CheckpointError>> {
616        self.compare_and_swap_checkpoint_ext(lease, checkpoint, expected_revision)
617    }
618}
619
620#[cfg(test)]
621mod tests {
622    use super::{PostgresWorkflowStore, PostgresWorkflowStoreError, validate_identifier};
623
624    #[test]
625    fn table_identifiers_are_validated_before_sql_construction() {
626        assert!(validate_identifier("runifold_workflows").is_ok());
627        assert!(matches!(
628            validate_identifier("workflow; DROP TABLE users"),
629            Err(PostgresWorkflowStoreError::InvalidTable)
630        ));
631        assert!(matches!(
632            validate_identifier("9workflow"),
633            Err(PostgresWorkflowStoreError::InvalidTable)
634        ));
635    }
636
637    #[test]
638    fn lease_statements_preserve_keyword_boundaries() {
639        let claim = PostgresWorkflowStore::claim_sql_for("runifold_workflows");
640        let heartbeat = PostgresWorkflowStore::heartbeat_sql_for("runifold_workflows");
641
642        assert!(claim.contains("FROM runifold_workflows"));
643        assert!(claim.contains("FOR UPDATE OF task, tenant SKIP LOCKED"));
644        assert!(claim.contains("tenant.max_concurrent_leases"));
645        assert!(claim.contains("pg_try_advisory_xact_lock"));
646        assert!(claim.contains("nextval('runifold_workflows_claim_seq')"));
647        assert!(claim.contains("UPDATE runifold_workflows AS task"));
648        assert!(heartbeat.contains("UPDATE runifold_workflows"));
649        assert!(heartbeat.contains("FROM renewed"));
650        assert!(heartbeat.contains("_budgets AS reservation"));
651    }
652}