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 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}