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 current_time_ms(&self) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>> {
292        Box::pin(async move {
293            let row = self
294                .client
295                .query_one(
296                    "SELECT (EXTRACT(EPOCH FROM clock_timestamp()) * 1000)::BIGINT",
297                    &[],
298                )
299                .await
300                .map_err(storage)?;
301            let value: i64 = row.try_get(0).map_err(storage)?;
302            u64::try_from(value).map_err(|_| {
303                WorkflowStoreError::new(
304                    WorkflowStoreErrorKind::Storage,
305                    "PostgreSQL returned a negative workflow clock",
306                )
307            })
308        })
309    }
310
311    fn set_tenant_policy(
312        &self,
313        tenant_id: WorkflowTenantId,
314        policy: WorkflowTenantPolicy,
315    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
316        Box::pin(async move {
317            let outstanding = i64::from(policy.max_outstanding_tasks());
318            let leases = i64::from(policy.max_concurrent_leases());
319            self.client
320                .execute(
321                    &format!(
322                        "INSERT INTO {table}_tenants (
323                            tenant_id, max_outstanding_tasks, max_concurrent_leases
324                         ) VALUES ($1, $2, $3)
325                         ON CONFLICT (tenant_id) DO UPDATE SET
326                            max_outstanding_tasks = EXCLUDED.max_outstanding_tasks,
327                            max_concurrent_leases = EXCLUDED.max_concurrent_leases,
328                            updated_at = clock_timestamp()",
329                        table = self.table
330                    ),
331                    &[&tenant_id.as_str(), &outstanding, &leases],
332                )
333                .await
334                .map_err(storage)?;
335            Ok(())
336        })
337    }
338
339    fn set_tenant_budget_policy(
340        &self,
341        tenant_id: WorkflowTenantId,
342        policy: WorkflowTenantBudgetPolicy,
343    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
344        self.set_tenant_budget_policy_inner(tenant_id, policy)
345    }
346
347    fn list_tenant_budgets(
348        &self,
349        after: Option<WorkflowTenantId>,
350        limit: WorkflowTenantListLimit,
351    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowTenantId>, WorkflowStoreError>> {
352        self.list_tenant_budgets_inner(after, limit)
353    }
354
355    fn inspect_tenant_budget(
356        &self,
357        tenant_id: WorkflowTenantId,
358    ) -> WorkflowStoreFuture<'_, Result<WorkflowTenantBudgetSnapshot, WorkflowStoreError>> {
359        self.inspect_tenant_budget_inner(tenant_id)
360    }
361
362    fn list_tenant_budget_audit(
363        &self,
364        tenant_id: WorkflowTenantId,
365        after: Option<WorkflowBudgetAuditCursor>,
366        limit: WorkflowBudgetAuditLimit,
367    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowBudgetAuditEvent>, WorkflowStoreError>> {
368        self.list_tenant_budget_audit_inner(tenant_id, after, limit)
369    }
370
371    fn compact_tenant_budget_audit(
372        &self,
373        tenant_id: WorkflowTenantId,
374        through: WorkflowBudgetAuditCursor,
375    ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>> {
376        self.compact_tenant_budget_audit_inner(tenant_id, through)
377    }
378
379    fn load_or_create_tenant_budget_audit_projection(
380        &self,
381        tenant_id: WorkflowTenantId,
382        projection_id: WorkflowBudgetAuditProjectionId,
383    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditCursor, WorkflowStoreError>> {
384        self.load_or_create_tenant_budget_audit_projection_inner(tenant_id, projection_id)
385    }
386
387    fn advance_tenant_budget_audit_projection(
388        &self,
389        tenant_id: WorkflowTenantId,
390        projection_id: WorkflowBudgetAuditProjectionId,
391        expected: WorkflowBudgetAuditCursor,
392        next: WorkflowBudgetAuditCursor,
393    ) -> WorkflowStoreFuture<'_, Result<bool, WorkflowStoreError>> {
394        self.advance_tenant_budget_audit_projection_inner(tenant_id, projection_id, expected, next)
395    }
396
397    fn claim_tenant_budget_audit_projection(
398        &self,
399        tenant_id: WorkflowTenantId,
400        projection_id: WorkflowBudgetAuditProjectionId,
401        owner: WorkerId,
402        lease: LeaseDuration,
403    ) -> WorkflowStoreFuture<
404        '_,
405        Result<Option<WorkflowBudgetAuditProjectionLease>, WorkflowStoreError>,
406    > {
407        self.claim_tenant_budget_audit_projection_inner(tenant_id, projection_id, owner, lease)
408    }
409
410    fn heartbeat_tenant_budget_audit_projection(
411        &self,
412        lease: WorkflowBudgetAuditProjectionLease,
413        extension: LeaseDuration,
414    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>
415    {
416        self.heartbeat_tenant_budget_audit_projection_inner(lease, extension)
417    }
418
419    fn advance_tenant_budget_audit_projection_lease(
420        &self,
421        lease: WorkflowBudgetAuditProjectionLease,
422        next: WorkflowBudgetAuditCursor,
423    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetAuditProjectionLease, WorkflowStoreError>>
424    {
425        self.advance_tenant_budget_audit_projection_lease_inner(lease, next)
426    }
427
428    fn release_tenant_budget_audit_projection(
429        &self,
430        lease: WorkflowBudgetAuditProjectionLease,
431    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
432        self.release_tenant_budget_audit_projection_inner(lease)
433    }
434
435    fn reserve_budget(
436        &self,
437        lease: WorkflowLease,
438        workflow_limit: Budget,
439        baseline: Usage,
440    ) -> WorkflowStoreFuture<'_, Result<WorkflowBudgetReservationOutcome, WorkflowStoreError>> {
441        self.reserve_budget_inner(lease, workflow_limit, baseline)
442    }
443
444    fn settle_budget(
445        &self,
446        lease: WorkflowLease,
447        cumulative: Usage,
448    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
449        self.settle_budget_inner(lease, cumulative)
450    }
451    fn enqueue(
452        &self,
453        task: WorkflowTask,
454    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
455        Box::pin(async move {
456            task.validate()?;
457            let version = i32::try_from(task.workflow_version).map_err(|_| {
458                WorkflowStoreError::new(
459                    WorkflowStoreErrorKind::InvalidInput,
460                    "workflow version exceeds PostgreSQL INTEGER",
461                )
462            })?;
463            self.client
464                .execute(
465                    &format!(
466                        "INSERT INTO {table}_tenants (
467                            tenant_id, max_outstanding_tasks, max_concurrent_leases
468                         ) VALUES ($1, 10000, 100)
469                         ON CONFLICT (tenant_id) DO NOTHING",
470                        table = self.table
471                    ),
472                    &[&task.tenant_id.as_str()],
473                )
474                .await
475                .map_err(storage)?;
476            let inserted = self
477                .client
478                .execute(
479                    &format!(
480                        r"
481                        WITH admitted AS (
482                            UPDATE {table}_tenants
483                            SET
484                                outstanding_tasks = outstanding_tasks + 1,
485                                updated_at = clock_timestamp()
486                            WHERE tenant_id = $2
487                              AND outstanding_tasks < max_outstanding_tasks
488                            RETURNING tenant_id
489                        )
490                        INSERT INTO {table} (
491                            checkpoint_id, tenant_id, workflow, workflow_version,
492                            input, priority, state
493                        )
494                        SELECT $1, $2, $3, $4, $5, $6, 'queued'
495                        FROM admitted
496                        ",
497                        table = self.table
498                    ),
499                    &[
500                        &task.checkpoint_id.as_uuid(),
501                        &task.tenant_id.as_str(),
502                        &task.workflow,
503                        &version,
504                        &task.input,
505                        &task.priority,
506                    ],
507                )
508                .await;
509            match inserted {
510                Ok(1) => Ok(()),
511                Ok(_) => Err(WorkflowStoreError::new(
512                    WorkflowStoreErrorKind::AdmissionDenied,
513                    format!(
514                        "workflow tenant `{}` reached its outstanding task limit",
515                        task.tenant_id.as_str()
516                    ),
517                )),
518                Err(error)
519                    if error
520                        .as_db_error()
521                        .is_some_and(|db| db.code() == &SqlState::UNIQUE_VIOLATION) =>
522                {
523                    Err(WorkflowStoreError::new(
524                        WorkflowStoreErrorKind::Conflict,
525                        format!("workflow task `{}` already exists", task.checkpoint_id),
526                    ))
527                }
528                Err(error) => Err(storage(error)),
529            }
530        })
531    }
532
533    fn claim(
534        &self,
535        worker: WorkerId,
536        lease: LeaseDuration,
537    ) -> WorkflowStoreFuture<'_, Result<Option<ClaimedWorkflow>, WorkflowStoreError>> {
538        Box::pin(async move { self.claim_inner(worker, lease).await })
539    }
540
541    fn heartbeat(
542        &self,
543        lease: WorkflowLease,
544        extension: LeaseDuration,
545    ) -> WorkflowStoreFuture<'_, Result<WorkflowLease, WorkflowStoreError>> {
546        Box::pin(async move { self.heartbeat_inner(lease, extension).await })
547    }
548
549    fn finish(
550        &self,
551        lease: WorkflowLease,
552        disposition: WorkflowDisposition,
553    ) -> WorkflowStoreFuture<'_, Result<(), WorkflowStoreError>> {
554        Box::pin(async move { self.finish_inner(&lease, disposition).await })
555    }
556
557    fn publish_signal(
558        &self,
559        tenant_id: WorkflowTenantId,
560        signal: WorkflowSignal,
561    ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalOutcome, WorkflowStoreError>> {
562        Box::pin(signal::publish(self, tenant_id, signal, false))
563    }
564
565    fn publish_control_signal(
566        &self,
567        tenant_id: WorkflowTenantId,
568        signal: WorkflowSignal,
569    ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalOutcome, WorkflowStoreError>> {
570        Box::pin(signal::publish(self, tenant_id, signal, true))
571    }
572
573    fn cancel(
574        &self,
575        tenant_id: WorkflowTenantId,
576        checkpoint_id: CheckpointId,
577    ) -> WorkflowStoreFuture<'_, Result<WorkflowCancelOutcome, WorkflowStoreError>> {
578        Box::pin(signal::cancel(self, tenant_id, checkpoint_id))
579    }
580
581    fn inspect_signal(
582        &self,
583        tenant_id: WorkflowTenantId,
584        signal_id: WorkflowSignalId,
585    ) -> WorkflowStoreFuture<'_, Result<WorkflowSignalSnapshot, WorkflowStoreError>> {
586        Box::pin(signal::inspect(self, tenant_id, signal_id))
587    }
588
589    fn load_signal_payload(
590        &self,
591        tenant_id: WorkflowTenantId,
592        signal_id: WorkflowSignalId,
593    ) -> WorkflowStoreFuture<'_, Result<Value, WorkflowStoreError>> {
594        Box::pin(signal::load_payload(self, tenant_id, signal_id))
595    }
596
597    fn compact_signals(
598        &self,
599        tenant_id: WorkflowTenantId,
600        retention: WorkflowSignalRetention,
601    ) -> WorkflowStoreFuture<'_, Result<u64, WorkflowStoreError>> {
602        Box::pin(signal::compact(self, tenant_id, retention))
603    }
604    fn inspect(
605        &self,
606        tenant_id: WorkflowTenantId,
607        checkpoint_id: CheckpointId,
608    ) -> WorkflowStoreFuture<'_, Result<WorkflowTaskSnapshot, WorkflowStoreError>> {
609        self.inspect_ext(tenant_id, checkpoint_id)
610    }
611
612    fn load_task_input(
613        &self,
614        tenant_id: WorkflowTenantId,
615        checkpoint_id: CheckpointId,
616    ) -> WorkflowStoreFuture<'_, Result<Value, WorkflowStoreError>> {
617        Box::pin(async move {
618            let row = self
619                .client
620                .query_opt(
621                    &format!(
622                        "SELECT input FROM {table} WHERE checkpoint_id = $1 AND tenant_id = $2",
623                        table = self.table
624                    ),
625                    &[&checkpoint_id.as_uuid(), &tenant_id.as_str()],
626                )
627                .await
628                .map_err(storage)?
629                .ok_or_else(|| {
630                    WorkflowStoreError::new(
631                        WorkflowStoreErrorKind::NotFound,
632                        "workflow task does not exist",
633                    )
634                })?;
635            row.try_get(0).map_err(storage)
636        })
637    }
638
639    fn list_checkpoint_history(
640        &self,
641        tenant_id: WorkflowTenantId,
642        checkpoint_id: CheckpointId,
643        after_revision: Option<u64>,
644        limit: WorkflowCheckpointHistoryLimit,
645    ) -> WorkflowStoreFuture<'_, Result<Vec<WorkflowCheckpointRevision>, WorkflowStoreError>> {
646        self.list_checkpoint_history_ext(tenant_id, checkpoint_id, after_revision, limit)
647    }
648
649    fn load_checkpoint_revision(
650        &self,
651        tenant_id: WorkflowTenantId,
652        checkpoint_id: CheckpointId,
653        revision: u64,
654    ) -> WorkflowStoreFuture<'_, Result<WorkflowCheckpointRevision, WorkflowStoreError>> {
655        self.load_checkpoint_revision_ext(tenant_id, checkpoint_id, revision)
656    }
657
658    fn fork_workflow(
659        &self,
660        tenant_id: WorkflowTenantId,
661        command: WorkflowForkCommand,
662    ) -> WorkflowStoreFuture<'_, Result<WorkflowForkOutcome, WorkflowStoreError>> {
663        self.fork_workflow_ext(tenant_id, command)
664    }
665
666    fn load_checkpoint(
667        &self,
668        lease: WorkflowLease,
669    ) -> WorkflowStoreFuture<'_, Result<Checkpoint, CheckpointError>> {
670        self.load_checkpoint_ext(lease)
671    }
672
673    fn compare_and_swap_checkpoint(
674        &self,
675        lease: WorkflowLease,
676        checkpoint: Checkpoint,
677        expected_revision: Option<u64>,
678    ) -> WorkflowStoreFuture<'_, Result<(), CheckpointError>> {
679        self.compare_and_swap_checkpoint_ext(lease, checkpoint, expected_revision)
680    }
681}
682
683#[cfg(test)]
684mod tests {
685    use super::{PostgresWorkflowStore, PostgresWorkflowStoreError, validate_identifier};
686
687    #[test]
688    fn table_identifiers_are_validated_before_sql_construction() {
689        assert!(validate_identifier("runifold_workflows").is_ok());
690        assert!(matches!(
691            validate_identifier("workflow; DROP TABLE users"),
692            Err(PostgresWorkflowStoreError::InvalidTable)
693        ));
694        assert!(matches!(
695            validate_identifier("9workflow"),
696            Err(PostgresWorkflowStoreError::InvalidTable)
697        ));
698    }
699
700    #[test]
701    fn lease_statements_preserve_keyword_boundaries() {
702        let claim = PostgresWorkflowStore::claim_sql_for("runifold_workflows");
703        let heartbeat = PostgresWorkflowStore::heartbeat_sql_for("runifold_workflows");
704
705        assert!(claim.contains("FROM runifold_workflows"));
706        assert!(claim.contains("FOR UPDATE OF task, tenant SKIP LOCKED"));
707        assert!(claim.contains("tenant.max_concurrent_leases"));
708        assert!(claim.contains("pg_try_advisory_xact_lock"));
709        assert!(claim.contains("nextval('runifold_workflows_claim_seq')"));
710        assert!(claim.contains("UPDATE runifold_workflows AS task"));
711        assert!(heartbeat.contains("UPDATE runifold_workflows"));
712        assert!(heartbeat.contains("FROM renewed"));
713        assert!(heartbeat.contains("_budgets AS reservation"));
714    }
715}