1mod 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#[derive(Debug, Error)]
60#[non_exhaustive]
61pub enum PostgresWorkflowStoreError {
62 #[error("workflow table must be a portable PostgreSQL identifier of at most 48 bytes")]
64 InvalidTable,
65 #[error("PostgreSQL workflow store operation failed: {0}")]
67 Database(#[from] tokio_postgres::Error),
68}
69
70#[derive(Clone, Debug)]
76pub struct PostgresWorkflowStore {
77 client: Arc<Client>,
78 table: String,
79}
80
81impl PostgresWorkflowStore {
82 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 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}