1use anyhow::Result;
2use sqlx::PgPool;
3
4use crate::store::{RetryEvent, WorkflowStore, retry_denial};
5use crate::types::*;
6
7const RETRY_ACTIVITY_SELECT: &str = "SELECT id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat FROM workflow.activities WHERE workflow_id = $1 AND status = 'FAILED' ORDER BY seq DESC LIMIT 1 FOR UPDATE";
8const RETRY_ACTIVITY_UPDATE: &str = "UPDATE workflow.activities SET status = 'PENDING', result = NULL, error = NULL, attempt = 1, claimed_by = NULL, scheduled_at = $1, started_at = NULL, completed_at = NULL, last_heartbeat = NULL WHERE id = $2 RETURNING id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat";
9
10const SCHEMA: &str = r#"
16CREATE SCHEMA IF NOT EXISTS workflow;
17CREATE SCHEMA IF NOT EXISTS engine;
18
19CREATE TABLE IF NOT EXISTS workflow.namespaces (
20 name TEXT PRIMARY KEY,
21 created_at DOUBLE PRECISION NOT NULL
22);
23INSERT INTO workflow.namespaces (name, created_at)
24 VALUES ('main', EXTRACT(EPOCH FROM NOW()))
25 ON CONFLICT DO NOTHING;
26
27CREATE TABLE IF NOT EXISTS workflow.workflows (
28 id TEXT PRIMARY KEY,
29 namespace TEXT NOT NULL DEFAULT 'main',
30 run_id TEXT NOT NULL,
31 workflow_type TEXT NOT NULL,
32 task_queue TEXT NOT NULL DEFAULT 'main',
33 status TEXT NOT NULL DEFAULT 'PENDING',
34 input TEXT,
35 result TEXT,
36 error TEXT,
37 parent_id TEXT,
38 claimed_by TEXT,
39 search_attributes TEXT,
40 archived_at DOUBLE PRECISION,
41 archive_uri TEXT,
42 -- Workflow-task dispatch (Phase 9): see sqlite.rs for the full comment.
43 needs_dispatch BOOLEAN NOT NULL DEFAULT FALSE,
44 dispatch_claimed_by TEXT,
45 dispatch_last_heartbeat DOUBLE PRECISION,
46 created_at DOUBLE PRECISION NOT NULL,
47 updated_at DOUBLE PRECISION NOT NULL,
48 completed_at DOUBLE PRECISION
49);
50CREATE INDEX IF NOT EXISTS idx_wf_status_queue ON workflow.workflows(status, task_queue);
51CREATE INDEX IF NOT EXISTS idx_wf_namespace ON workflow.workflows(namespace);
52CREATE INDEX IF NOT EXISTS idx_wf_dispatch ON workflow.workflows(task_queue, needs_dispatch, dispatch_claimed_by);
53
54CREATE TABLE IF NOT EXISTS workflow.events (
55 id BIGSERIAL PRIMARY KEY,
56 workflow_id TEXT NOT NULL REFERENCES workflow.workflows(id),
57 seq INTEGER NOT NULL,
58 event_type TEXT NOT NULL,
59 payload TEXT,
60 timestamp DOUBLE PRECISION NOT NULL
61);
62CREATE INDEX IF NOT EXISTS idx_wf_events_lookup ON workflow.events(workflow_id, seq);
63
64CREATE TABLE IF NOT EXISTS workflow.activities (
65 id BIGSERIAL PRIMARY KEY,
66 workflow_id TEXT NOT NULL REFERENCES workflow.workflows(id),
67 seq INTEGER NOT NULL,
68 name TEXT NOT NULL,
69 task_queue TEXT NOT NULL DEFAULT 'main',
70 input TEXT,
71 status TEXT NOT NULL DEFAULT 'PENDING',
72 result TEXT,
73 error TEXT,
74 attempt INTEGER NOT NULL DEFAULT 1,
75 max_attempts INTEGER NOT NULL DEFAULT 3,
76 initial_interval_secs DOUBLE PRECISION NOT NULL DEFAULT 1,
77 backoff_coefficient DOUBLE PRECISION NOT NULL DEFAULT 2,
78 start_to_close_secs DOUBLE PRECISION NOT NULL DEFAULT 300,
79 heartbeat_timeout_secs DOUBLE PRECISION,
80 claimed_by TEXT,
81 scheduled_at DOUBLE PRECISION NOT NULL,
82 started_at DOUBLE PRECISION,
83 completed_at DOUBLE PRECISION,
84 last_heartbeat DOUBLE PRECISION,
85 UNIQUE (workflow_id, seq)
86);
87CREATE INDEX IF NOT EXISTS idx_wf_act_pending ON workflow.activities(task_queue, status, scheduled_at);
88
89CREATE TABLE IF NOT EXISTS workflow.timers (
90 id BIGSERIAL PRIMARY KEY,
91 workflow_id TEXT NOT NULL REFERENCES workflow.workflows(id),
92 seq INTEGER NOT NULL,
93 fire_at DOUBLE PRECISION NOT NULL,
94 fired BOOLEAN NOT NULL DEFAULT FALSE,
95 UNIQUE (workflow_id, seq)
96);
97CREATE INDEX IF NOT EXISTS idx_wf_timers_due ON workflow.timers(fire_at) WHERE fired = FALSE;
98
99CREATE TABLE IF NOT EXISTS workflow.signals (
100 id BIGSERIAL PRIMARY KEY,
101 workflow_id TEXT NOT NULL REFERENCES workflow.workflows(id),
102 name TEXT NOT NULL,
103 payload TEXT,
104 consumed BOOLEAN NOT NULL DEFAULT FALSE,
105 received_at DOUBLE PRECISION NOT NULL
106);
107CREATE INDEX IF NOT EXISTS idx_wf_signals_lookup ON workflow.signals(workflow_id, name, consumed);
108
109CREATE TABLE IF NOT EXISTS workflow.schedules (
110 namespace TEXT NOT NULL DEFAULT 'main',
111 name TEXT NOT NULL,
112 workflow_type TEXT NOT NULL,
113 cron_expr TEXT NOT NULL,
114 timezone TEXT NOT NULL DEFAULT 'UTC',
115 input TEXT,
116 task_queue TEXT NOT NULL DEFAULT 'main',
117 overlap_policy TEXT NOT NULL DEFAULT 'skip',
118 paused BOOLEAN NOT NULL DEFAULT FALSE,
119 last_run_at DOUBLE PRECISION,
120 next_run_at DOUBLE PRECISION,
121 last_workflow_id TEXT,
122 created_at DOUBLE PRECISION NOT NULL,
123 PRIMARY KEY (namespace, name)
124);
125
126CREATE TABLE IF NOT EXISTS workflow.workers (
127 id TEXT PRIMARY KEY,
128 namespace TEXT NOT NULL DEFAULT 'main',
129 identity TEXT NOT NULL,
130 task_queue TEXT NOT NULL,
131 workflows TEXT,
132 activities TEXT,
133 max_concurrent_workflows INTEGER NOT NULL DEFAULT 10,
134 max_concurrent_activities INTEGER NOT NULL DEFAULT 10,
135 active_tasks INTEGER NOT NULL DEFAULT 0,
136 last_heartbeat DOUBLE PRECISION NOT NULL,
137 registered_at DOUBLE PRECISION NOT NULL
138);
139
140CREATE TABLE IF NOT EXISTS workflow.snapshots (
141 workflow_id TEXT NOT NULL REFERENCES workflow.workflows(id),
142 event_seq INTEGER NOT NULL,
143 state_json TEXT NOT NULL,
144 created_at DOUBLE PRECISION NOT NULL,
145 PRIMARY KEY (workflow_id, event_seq)
146);
147
148-- Plan-15 slice 3: workflow.api_keys retired in favour of the auth
149-- module (sessions / JWT / Zanzibar tuples). Table is dropped on
150-- migration; nothing here re-creates it.
151DROP TABLE IF EXISTS workflow.api_keys CASCADE;
152
153CREATE TABLE IF NOT EXISTS engine.events (
154 id BIGSERIAL PRIMARY KEY,
155 ts DOUBLE PRECISION NOT NULL DEFAULT EXTRACT(EPOCH FROM NOW()),
156 namespace TEXT NOT NULL,
157 subsystem TEXT NOT NULL,
158 kind TEXT NOT NULL,
159 payload JSONB NOT NULL DEFAULT '{}'::jsonb
160);
161CREATE INDEX IF NOT EXISTS idx_engine_events_ns_id ON engine.events(namespace, id);
162CREATE INDEX IF NOT EXISTS idx_engine_events_ts_prune ON engine.events(ts);
163
164"#;
165
166const V0_13_2_RELOCATION_SQL: &str = r#"
177DO $$
178DECLARE
179 has_old BOOLEAN;
180BEGIN
181 -- Each table: if the legacy public.<old> exists, drop the empty
182 -- schema-qualified twin (created above by SCHEMA) and move the
183 -- legacy table into its new home.
184
185 -- workflows
186 SELECT to_regclass('public.workflows') IS NOT NULL INTO has_old;
187 IF has_old THEN
188 DROP TABLE IF EXISTS workflow.workflows CASCADE;
189 ALTER TABLE public.workflows SET SCHEMA workflow;
190 END IF;
191
192 -- workflow_events → workflow.events
193 SELECT to_regclass('public.workflow_events') IS NOT NULL INTO has_old;
194 IF has_old THEN
195 DROP TABLE IF EXISTS workflow.events CASCADE;
196 ALTER TABLE public.workflow_events SET SCHEMA workflow;
197 ALTER TABLE workflow.workflow_events RENAME TO events;
198 END IF;
199
200 -- workflow_activities → workflow.activities
201 SELECT to_regclass('public.workflow_activities') IS NOT NULL INTO has_old;
202 IF has_old THEN
203 DROP TABLE IF EXISTS workflow.activities CASCADE;
204 ALTER TABLE public.workflow_activities SET SCHEMA workflow;
205 ALTER TABLE workflow.workflow_activities RENAME TO activities;
206 END IF;
207
208 -- workflow_timers → workflow.timers
209 SELECT to_regclass('public.workflow_timers') IS NOT NULL INTO has_old;
210 IF has_old THEN
211 DROP TABLE IF EXISTS workflow.timers CASCADE;
212 ALTER TABLE public.workflow_timers SET SCHEMA workflow;
213 ALTER TABLE workflow.workflow_timers RENAME TO timers;
214 END IF;
215
216 -- workflow_signals → workflow.signals
217 SELECT to_regclass('public.workflow_signals') IS NOT NULL INTO has_old;
218 IF has_old THEN
219 DROP TABLE IF EXISTS workflow.signals CASCADE;
220 ALTER TABLE public.workflow_signals SET SCHEMA workflow;
221 ALTER TABLE workflow.workflow_signals RENAME TO signals;
222 END IF;
223
224 -- workflow_snapshots → workflow.snapshots
225 SELECT to_regclass('public.workflow_snapshots') IS NOT NULL INTO has_old;
226 IF has_old THEN
227 DROP TABLE IF EXISTS workflow.snapshots CASCADE;
228 ALTER TABLE public.workflow_snapshots SET SCHEMA workflow;
229 ALTER TABLE workflow.workflow_snapshots RENAME TO snapshots;
230 END IF;
231
232 -- workflow_schedules → workflow.schedules
233 SELECT to_regclass('public.workflow_schedules') IS NOT NULL INTO has_old;
234 IF has_old THEN
235 DROP TABLE IF EXISTS workflow.schedules CASCADE;
236 ALTER TABLE public.workflow_schedules SET SCHEMA workflow;
237 ALTER TABLE workflow.workflow_schedules RENAME TO schedules;
238 END IF;
239
240 -- workflow_workers → workflow.workers
241 SELECT to_regclass('public.workflow_workers') IS NOT NULL INTO has_old;
242 IF has_old THEN
243 DROP TABLE IF EXISTS workflow.workers CASCADE;
244 ALTER TABLE public.workflow_workers SET SCHEMA workflow;
245 ALTER TABLE workflow.workflow_workers RENAME TO workers;
246 END IF;
247
248 -- namespaces → workflow.namespaces
249 SELECT to_regclass('public.namespaces') IS NOT NULL INTO has_old;
250 IF has_old THEN
251 DROP TABLE IF EXISTS workflow.namespaces CASCADE;
252 ALTER TABLE public.namespaces SET SCHEMA workflow;
253 END IF;
254
255 -- public.api_keys: retired in plan-15 slice 3 (workflow REST API
256 -- auth moved to the auth module — see CHANGELOG). Drop any
257 -- orphaned legacy table so an upgraded v0.13.1 install doesn't
258 -- carry it forward.
259 SELECT to_regclass('public.api_keys') IS NOT NULL INTO has_old;
260 IF has_old THEN
261 DROP TABLE public.api_keys CASCADE;
262 END IF;
263
264 -- engine_events → engine.events (notification outbox; preserves the
265 -- v0.13.1 publish-on-commit guarantee since the new INSERT into
266 -- engine.events sits in the same transaction as the pg_notify call).
267 SELECT to_regclass('public.engine_events') IS NOT NULL INTO has_old;
268 IF has_old THEN
269 DROP TABLE IF EXISTS engine.events CASCADE;
270 ALTER TABLE public.engine_events SET SCHEMA engine;
271 ALTER TABLE engine.engine_events RENAME TO events;
272 END IF;
273END
274$$;
275"#;
276
277fn sanitise_schema(schema: &str) -> Vec<String> {
289 let without_comments: String = schema
290 .lines()
291 .filter(|line| !line.trim_start().starts_with("--"))
292 .collect::<Vec<_>>()
293 .join("\n");
294
295 without_comments
296 .split(';')
297 .map(|s| s.trim().to_string())
298 .filter(|s| !s.is_empty())
299 .collect()
300}
301
302#[derive(Clone)]
308pub struct PostgresStore {
309 pool: PgPool,
310}
311
312impl PostgresStore {
313 pub async fn new(url: &str) -> Result<Self> {
314 let pool = PgPool::connect(url).await?;
315 Self::from_pool(pool).await
316 }
317
318 pub async fn from_pool(pool: PgPool) -> Result<Self> {
323 let store = Self { pool };
324 store.migrate().await?;
325 Ok(store)
326 }
327
328 pub fn pool(&self) -> &PgPool {
331 &self.pool
332 }
333
334 async fn migrate(&self) -> Result<()> {
335 for statement in sanitise_schema(SCHEMA) {
342 sqlx::query(&statement).execute(&self.pool).await?;
343 }
344 sqlx::raw_sql(V0_13_2_RELOCATION_SQL)
347 .execute(&self.pool)
348 .await?;
349 sqlx::raw_sql(
357 r#"
358 DROP TRIGGER IF EXISTS workflow_runnable_notify ON workflow.workflows;
359 DROP TRIGGER IF EXISTS workflow_task_notify ON workflow.activities;
360 DROP FUNCTION IF EXISTS assay_notify_runnable();
361 DROP FUNCTION IF EXISTS assay_notify_task();
362 "#,
363 )
364 .execute(&self.pool)
365 .await?;
366 Ok(())
367 }
368
369 pub async fn try_acquire_leader_lock(&self) -> Result<bool> {
372 let row: (bool,) = sqlx::query_as("SELECT pg_try_advisory_lock(1)")
373 .fetch_one(&self.pool)
374 .await?;
375 Ok(row.0)
376 }
377}
378
379impl WorkflowStore for PostgresStore {
380 async fn create_namespace(&self, name: &str) -> Result<()> {
383 sqlx::query("INSERT INTO workflow.namespaces (name, created_at) VALUES ($1, EXTRACT(EPOCH FROM NOW()))")
384 .bind(name)
385 .execute(&self.pool)
386 .await?;
387 Ok(())
388 }
389
390 async fn list_namespaces(&self) -> Result<Vec<crate::store::NamespaceRecord>> {
391 let rows = sqlx::query_as::<_, (String, f64)>(
392 "SELECT name, created_at FROM workflow.namespaces ORDER BY name",
393 )
394 .fetch_all(&self.pool)
395 .await?;
396 Ok(rows
397 .into_iter()
398 .map(|(name, created_at)| crate::store::NamespaceRecord { name, created_at })
399 .collect())
400 }
401
402 async fn delete_namespace(&self, name: &str) -> Result<bool> {
403 let res = sqlx::query("DELETE FROM workflow.namespaces WHERE name = $1 AND name != 'main'")
404 .bind(name)
405 .execute(&self.pool)
406 .await?;
407 Ok(res.rows_affected() > 0)
408 }
409
410 async fn get_namespace_stats(&self, namespace: &str) -> Result<crate::store::NamespaceStats> {
411 let total: (i64,) =
412 sqlx::query_as("SELECT COUNT(*) FROM workflow.workflows WHERE namespace = $1")
413 .bind(namespace)
414 .fetch_one(&self.pool)
415 .await?;
416 let running: (i64,) = sqlx::query_as(
417 "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = $1 AND status = 'RUNNING'",
418 )
419 .bind(namespace)
420 .fetch_one(&self.pool)
421 .await?;
422 let pending: (i64,) = sqlx::query_as(
423 "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = $1 AND status = 'PENDING'",
424 )
425 .bind(namespace)
426 .fetch_one(&self.pool)
427 .await?;
428 let completed: (i64,) = sqlx::query_as(
429 "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = $1 AND status = 'COMPLETED'",
430 )
431 .bind(namespace)
432 .fetch_one(&self.pool)
433 .await?;
434 let failed: (i64,) = sqlx::query_as(
435 "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = $1 AND status = 'FAILED'",
436 )
437 .bind(namespace)
438 .fetch_one(&self.pool)
439 .await?;
440 let schedules: (i64,) =
441 sqlx::query_as("SELECT COUNT(*) FROM workflow.schedules WHERE namespace = $1")
442 .bind(namespace)
443 .fetch_one(&self.pool)
444 .await?;
445 let workers: (i64,) =
446 sqlx::query_as("SELECT COUNT(*) FROM workflow.workers WHERE namespace = $1")
447 .bind(namespace)
448 .fetch_one(&self.pool)
449 .await?;
450
451 Ok(crate::store::NamespaceStats {
452 namespace: namespace.to_string(),
453 total_workflows: total.0,
454 running: running.0,
455 pending: pending.0,
456 completed: completed.0,
457 failed: failed.0,
458 schedules: schedules.0,
459 workers: workers.0,
460 })
461 }
462
463 async fn create_workflow(&self, wf: &WorkflowRecord) -> Result<()> {
466 sqlx::query(
467 "INSERT INTO workflow.workflows (id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at)
468 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17)",
469 )
470 .bind(&wf.id)
471 .bind(&wf.namespace)
472 .bind(&wf.run_id)
473 .bind(&wf.workflow_type)
474 .bind(&wf.task_queue)
475 .bind(&wf.status)
476 .bind(&wf.input)
477 .bind(&wf.result)
478 .bind(&wf.error)
479 .bind(&wf.parent_id)
480 .bind(&wf.claimed_by)
481 .bind(&wf.search_attributes)
482 .bind(wf.archived_at)
483 .bind(&wf.archive_uri)
484 .bind(wf.created_at)
485 .bind(wf.updated_at)
486 .bind(wf.completed_at)
487 .execute(&self.pool)
488 .await?;
489 Ok(())
490 }
491
492 async fn get_workflow(&self, id: &str) -> Result<Option<WorkflowRecord>> {
493 let row = sqlx::query_as::<_, PgWorkflowRow>(
494 "SELECT id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at FROM workflow.workflows WHERE id = $1",
495 )
496 .bind(id)
497 .fetch_optional(&self.pool)
498 .await?;
499 Ok(row.map(Into::into))
500 }
501
502 async fn list_workflows(
503 &self,
504 namespace: &str,
505 status: Option<WorkflowStatus>,
506 workflow_type: Option<&str>,
507 search_attrs_filter: Option<&str>,
508 limit: i64,
509 offset: i64,
510 ) -> Result<Vec<WorkflowRecord>> {
511 let status_str = status.map(|s| s.to_string());
512
513 let filter_pairs: Vec<(String, serde_json::Value)> = search_attrs_filter
514 .and_then(|s| serde_json::from_str::<serde_json::Value>(s).ok())
515 .and_then(|v| v.as_object().cloned())
516 .map(|m| m.into_iter().collect())
517 .unwrap_or_default();
518
519 let mut sql = String::from(
520 "SELECT id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at
521 FROM workflow.workflows
522 WHERE namespace = $1
523 AND ($2::TEXT IS NULL OR status = $2)
524 AND ($3::TEXT IS NULL OR workflow_type = $3)",
525 );
526 let mut idx = 4usize;
528 for _ in &filter_pairs {
529 sql.push_str(&format!(
530 " AND (search_attributes::jsonb)->>${} = ${}",
531 idx,
532 idx + 1
533 ));
534 idx += 2;
535 }
536 sql.push_str(&format!(
537 " ORDER BY created_at DESC LIMIT ${} OFFSET ${}",
538 idx,
539 idx + 1
540 ));
541
542 let mut q = sqlx::query_as::<_, PgWorkflowRow>(&sql)
543 .bind(namespace)
544 .bind(&status_str)
545 .bind(workflow_type);
546 for (key, value) in &filter_pairs {
547 q = q.bind(key.clone());
548 let as_text = match value {
550 serde_json::Value::String(s) => s.clone(),
551 other => other.to_string(),
552 };
553 q = q.bind(as_text);
554 }
555 let rows = q.bind(limit).bind(offset).fetch_all(&self.pool).await?;
556 Ok(rows.into_iter().map(Into::into).collect())
557 }
558
559 async fn update_workflow_status(
560 &self,
561 id: &str,
562 status: WorkflowStatus,
563 result: Option<&str>,
564 error: Option<&str>,
565 ) -> Result<()> {
566 let now = timestamp_now();
567 let completed_at = if status.is_terminal() {
568 Some(now)
569 } else {
570 None
571 };
572 sqlx::query(
573 "UPDATE workflow.workflows SET status = $1, result = COALESCE($2, result), error = COALESCE($3, error), updated_at = $4, completed_at = COALESCE($5, completed_at) WHERE id = $6",
574 )
575 .bind(status.to_string())
576 .bind(result)
577 .bind(error)
578 .bind(now)
579 .bind(completed_at)
580 .bind(id)
581 .execute(&self.pool)
582 .await?;
583 Ok(())
584 }
585
586 async fn claim_workflow(&self, id: &str, worker_id: &str) -> Result<bool> {
587 let res = sqlx::query(
588 "UPDATE workflow.workflows SET claimed_by = $1, status = 'RUNNING', updated_at = $2 WHERE id = $3 AND claimed_by IS NULL",
589 )
590 .bind(worker_id)
591 .bind(timestamp_now())
592 .bind(id)
593 .execute(&self.pool)
594 .await?;
595 Ok(res.rows_affected() > 0)
596 }
597
598 async fn mark_workflow_dispatchable(&self, workflow_id: &str) -> Result<()> {
599 sqlx::query("UPDATE workflow.workflows SET needs_dispatch = TRUE WHERE id = $1")
600 .bind(workflow_id)
601 .execute(&self.pool)
602 .await?;
603 Ok(())
604 }
605
606 async fn claim_workflow_task(
607 &self,
608 task_queue: &str,
609 worker_id: &str,
610 ) -> Result<Option<WorkflowRecord>> {
611 let now = timestamp_now();
612 let row = sqlx::query_as::<_, PgWorkflowRow>(
615 "UPDATE workflow.workflows
616 SET dispatch_claimed_by = $1, dispatch_last_heartbeat = $2, needs_dispatch = FALSE
617 WHERE id = (
618 SELECT id FROM workflow.workflows
619 WHERE task_queue = $3
620 AND needs_dispatch = TRUE
621 AND dispatch_claimed_by IS NULL
622 AND status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')
623 ORDER BY updated_at ASC
624 FOR UPDATE SKIP LOCKED
625 LIMIT 1
626 )
627 RETURNING id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at",
628 )
629 .bind(worker_id)
630 .bind(now)
631 .bind(task_queue)
632 .fetch_optional(&self.pool)
633 .await?;
634 Ok(row.map(Into::into))
635 }
636
637 async fn release_workflow_task(&self, workflow_id: &str, worker_id: &str) -> Result<()> {
638 sqlx::query(
639 "UPDATE workflow.workflows
640 SET dispatch_claimed_by = NULL, dispatch_last_heartbeat = NULL
641 WHERE id = $1 AND dispatch_claimed_by = $2",
642 )
643 .bind(workflow_id)
644 .bind(worker_id)
645 .execute(&self.pool)
646 .await?;
647 Ok(())
648 }
649
650 async fn release_stale_dispatch_leases(&self, now: f64, timeout_secs: f64) -> Result<u64> {
651 let res = sqlx::query(
652 "UPDATE workflow.workflows
653 SET dispatch_claimed_by = NULL,
654 dispatch_last_heartbeat = NULL,
655 needs_dispatch = TRUE
656 WHERE dispatch_claimed_by IS NOT NULL
657 AND ($1 - dispatch_last_heartbeat) > $2
658 AND status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')",
659 )
660 .bind(now)
661 .bind(timeout_secs)
662 .execute(&self.pool)
663 .await?;
664 Ok(res.rows_affected())
665 }
666
667 async fn append_event(&self, ev: &WorkflowEvent) -> Result<i64> {
670 let row: (i64,) = sqlx::query_as(
671 "INSERT INTO workflow.events (workflow_id, seq, event_type, payload, timestamp) VALUES ($1, $2, $3, $4, $5) RETURNING id",
672 )
673 .bind(&ev.workflow_id)
674 .bind(ev.seq)
675 .bind(&ev.event_type)
676 .bind(&ev.payload)
677 .bind(ev.timestamp)
678 .fetch_one(&self.pool)
679 .await?;
680 Ok(row.0)
681 }
682
683 async fn list_events(&self, workflow_id: &str) -> Result<Vec<WorkflowEvent>> {
684 let rows = sqlx::query_as::<_, PgEventRow>(
685 "SELECT id, workflow_id, seq, event_type, payload, timestamp FROM workflow.events WHERE workflow_id = $1 ORDER BY seq ASC",
686 )
687 .bind(workflow_id)
688 .fetch_all(&self.pool)
689 .await?;
690 Ok(rows.into_iter().map(Into::into).collect())
691 }
692
693 async fn list_events_page(
694 &self,
695 workflow_id: &str,
696 cursor: Option<i32>,
697 limit: i64,
698 descending: bool,
699 ) -> Result<Vec<WorkflowEvent>> {
700 let limit = limit.clamp(0, 1_000);
701 if limit == 0 {
702 return Ok(Vec::new());
703 }
704 let rows = if descending {
705 sqlx::query_as::<_, PgEventRow>(
706 "SELECT id, workflow_id, seq, event_type, payload, timestamp
707 FROM workflow.events
708 WHERE workflow_id = $1 AND ($2::INTEGER IS NULL OR seq < $2)
709 ORDER BY seq DESC LIMIT $3",
710 )
711 .bind(workflow_id)
712 .bind(cursor)
713 .bind(limit)
714 .fetch_all(&self.pool)
715 .await?
716 } else {
717 sqlx::query_as::<_, PgEventRow>(
718 "SELECT id, workflow_id, seq, event_type, payload, timestamp
719 FROM workflow.events
720 WHERE workflow_id = $1 AND ($2::INTEGER IS NULL OR seq > $2)
721 ORDER BY seq ASC LIMIT $3",
722 )
723 .bind(workflow_id)
724 .bind(cursor)
725 .bind(limit)
726 .fetch_all(&self.pool)
727 .await?
728 };
729 Ok(rows.into_iter().map(Into::into).collect())
730 }
731
732 async fn get_event_count(&self, workflow_id: &str) -> Result<i64> {
733 let row: (i64,) =
734 sqlx::query_as("SELECT COUNT(*) FROM workflow.events WHERE workflow_id = $1")
735 .bind(workflow_id)
736 .fetch_one(&self.pool)
737 .await?;
738 Ok(row.0)
739 }
740
741 async fn create_activity(&self, act: &WorkflowActivity) -> Result<i64> {
744 let row: (i64,) = sqlx::query_as(
745 "INSERT INTO workflow.activities (workflow_id, seq, name, task_queue, input, status, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, scheduled_at)
746 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13) RETURNING id",
747 )
748 .bind(&act.workflow_id)
749 .bind(act.seq)
750 .bind(&act.name)
751 .bind(&act.task_queue)
752 .bind(&act.input)
753 .bind(&act.status)
754 .bind(act.attempt)
755 .bind(act.max_attempts)
756 .bind(act.initial_interval_secs)
757 .bind(act.backoff_coefficient)
758 .bind(act.start_to_close_secs)
759 .bind(act.heartbeat_timeout_secs)
760 .bind(act.scheduled_at)
761 .fetch_one(&self.pool)
762 .await?;
763 Ok(row.0)
764 }
765
766 async fn get_activity(&self, id: i64) -> Result<Option<WorkflowActivity>> {
767 let row = sqlx::query_as::<_, PgActivityRow>(
768 "SELECT id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat
769 FROM workflow.activities WHERE id = $1",
770 )
771 .bind(id)
772 .fetch_optional(&self.pool)
773 .await?;
774 Ok(row.map(Into::into))
775 }
776
777 async fn get_activity_by_workflow_seq(
778 &self,
779 workflow_id: &str,
780 seq: i32,
781 ) -> Result<Option<WorkflowActivity>> {
782 let row = sqlx::query_as::<_, PgActivityRow>(
783 "SELECT id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat
784 FROM workflow.activities WHERE workflow_id = $1 AND seq = $2",
785 )
786 .bind(workflow_id)
787 .bind(seq)
788 .fetch_optional(&self.pool)
789 .await?;
790 Ok(row.map(Into::into))
791 }
792
793 async fn claim_activity(
794 &self,
795 task_queue: &str,
796 worker_id: &str,
797 ) -> Result<Option<WorkflowActivity>> {
798 let now = timestamp_now();
799 let row = sqlx::query_as::<_, PgActivityRow>(
802 "UPDATE workflow.activities SET status = 'RUNNING', claimed_by = $1, started_at = $2
803 WHERE id = (
804 SELECT id FROM workflow.activities
805 WHERE task_queue = $3 AND status = 'PENDING'
806 ORDER BY scheduled_at ASC
807 FOR UPDATE SKIP LOCKED
808 LIMIT 1
809 )
810 RETURNING id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat",
811 )
812 .bind(worker_id)
813 .bind(now)
814 .bind(task_queue)
815 .fetch_optional(&self.pool)
816 .await?;
817 Ok(row.map(Into::into))
818 }
819
820 async fn requeue_activity_for_retry(
821 &self,
822 id: i64,
823 next_attempt: i32,
824 next_scheduled_at: f64,
825 ) -> Result<()> {
826 sqlx::query(
827 "UPDATE workflow.activities
828 SET status = 'PENDING', attempt = $1, scheduled_at = $2,
829 claimed_by = NULL, started_at = NULL, last_heartbeat = NULL,
830 error = NULL
831 WHERE id = $3",
832 )
833 .bind(next_attempt)
834 .bind(next_scheduled_at)
835 .bind(id)
836 .execute(&self.pool)
837 .await?;
838 Ok(())
839 }
840
841 async fn retry_failed_activity(
842 &self,
843 workflow_id: &str,
844 requested_by: &str,
845 reason: &str,
846 requested_at: f64,
847 ) -> Result<RetryFailedActivityResult> {
848 let mut tx = self.pool.begin().await?;
849 let workflow: Option<(String, Option<String>, Option<f64>)> = sqlx::query_as(
850 "SELECT status, parent_id, archived_at FROM workflow.workflows WHERE id = $1 FOR UPDATE",
851 )
852 .bind(workflow_id)
853 .fetch_optional(&mut *tx)
854 .await?;
855 let Some((status, parent_id, archived_at)) = workflow else {
856 return Ok(RetryFailedActivityResult::NotFound);
857 };
858 if let Some(denial) = retry_denial(status, parent_id, archived_at) {
859 return Ok(denial);
860 }
861
862 let failed = sqlx::query_as::<_, PgActivityRow>(RETRY_ACTIVITY_SELECT)
863 .bind(workflow_id)
864 .fetch_optional(&mut *tx)
865 .await?;
866 let Some(failed) = failed else {
867 return Ok(RetryFailedActivityResult::NoFailedActivity);
868 };
869 let failed_event_seq: (i32,) = sqlx::query_as(
870 "SELECT seq FROM workflow.events
871 WHERE workflow_id = $1 AND event_type = 'ActivityFailed'
872 ORDER BY seq DESC LIMIT 1",
873 )
874 .bind(workflow_id)
875 .fetch_one(&mut *tx)
876 .await?;
877 let invalidated =
878 sqlx::query("DELETE FROM workflow.activities WHERE workflow_id = $1 AND seq > $2")
879 .bind(workflow_id)
880 .bind(failed.seq)
881 .execute(&mut *tx)
882 .await?
883 .rows_affected();
884 let activity = sqlx::query_as::<_, PgActivityRow>(RETRY_ACTIVITY_UPDATE)
885 .bind(requested_at)
886 .bind(failed.id)
887 .fetch_one(&mut *tx)
888 .await?;
889 sqlx::query(
890 "UPDATE workflow.workflows
891 SET status = 'WAITING', result = NULL, error = NULL, completed_at = NULL,
892 updated_at = $1, needs_dispatch = FALSE, dispatch_claimed_by = NULL,
893 dispatch_last_heartbeat = NULL
894 WHERE id = $2",
895 )
896 .bind(requested_at)
897 .bind(workflow_id)
898 .execute(&mut *tx)
899 .await?;
900 let event_seq: (i32,) = sqlx::query_as(
901 "SELECT COALESCE(MAX(seq), 0) + 1 FROM workflow.events WHERE workflow_id = $1",
902 )
903 .bind(workflow_id)
904 .fetch_one(&mut *tx)
905 .await?;
906 let payload = RetryEvent {
907 activity_id: failed.id,
908 activity_seq: failed.seq,
909 activity_name: &failed.name,
910 failed_event_seq: failed_event_seq.0,
911 requested_by,
912 reason,
913 invalidated_activities: invalidated,
914 }
915 .payload();
916 sqlx::query(
917 "INSERT INTO workflow.events (workflow_id, seq, event_type, payload, timestamp)
918 VALUES ($1, $2, 'ActivityRetryRequested', $3, $4)",
919 )
920 .bind(workflow_id)
921 .bind(event_seq.0)
922 .bind(payload.to_string())
923 .bind(requested_at)
924 .execute(&mut *tx)
925 .await?;
926 tx.commit().await?;
927 Ok(RetryFailedActivityResult::Retried(Box::new(
928 RetriedActivity {
929 activity: activity.into(),
930 invalidated_activities: invalidated,
931 },
932 )))
933 }
934
935 async fn complete_activity(
936 &self,
937 id: i64,
938 result: Option<&str>,
939 error: Option<&str>,
940 failed: bool,
941 ) -> Result<()> {
942 let status = if failed { "FAILED" } else { "COMPLETED" };
943 sqlx::query(
944 "UPDATE workflow.activities SET status = $1, result = $2, error = $3, completed_at = $4 WHERE id = $5",
945 )
946 .bind(status)
947 .bind(result)
948 .bind(error)
949 .bind(timestamp_now())
950 .bind(id)
951 .execute(&self.pool)
952 .await?;
953 Ok(())
954 }
955
956 async fn heartbeat_activity(&self, id: i64, _details: Option<&str>) -> Result<()> {
957 sqlx::query("UPDATE workflow.activities SET last_heartbeat = $1 WHERE id = $2")
958 .bind(timestamp_now())
959 .bind(id)
960 .execute(&self.pool)
961 .await?;
962 Ok(())
963 }
964
965 async fn get_timed_out_activities(&self, now: f64) -> Result<Vec<WorkflowActivity>> {
966 let rows = sqlx::query_as::<_, PgActivityRow>(
967 "SELECT id, workflow_id, seq, name, task_queue, input, status, result, error, attempt, max_attempts, initial_interval_secs, backoff_coefficient, start_to_close_secs, heartbeat_timeout_secs, claimed_by, scheduled_at, started_at, completed_at, last_heartbeat
968 FROM workflow.activities
969 WHERE status = 'RUNNING'
970 AND heartbeat_timeout_secs IS NOT NULL
971 AND ($1 - COALESCE(last_heartbeat, started_at)) > heartbeat_timeout_secs",
972 )
973 .bind(now)
974 .fetch_all(&self.pool)
975 .await?;
976 Ok(rows.into_iter().map(Into::into).collect())
977 }
978
979 async fn create_timer(&self, timer: &WorkflowTimer) -> Result<i64> {
982 let inserted: Option<(i64,)> = sqlx::query_as(
985 "INSERT INTO workflow.timers (workflow_id, seq, fire_at, fired)
986 VALUES ($1, $2, $3, FALSE)
987 ON CONFLICT (workflow_id, seq) DO NOTHING
988 RETURNING id",
989 )
990 .bind(&timer.workflow_id)
991 .bind(timer.seq)
992 .bind(timer.fire_at)
993 .fetch_optional(&self.pool)
994 .await?;
995
996 if let Some((id,)) = inserted {
997 return Ok(id);
998 }
999
1000 let (id,): (i64,) =
1002 sqlx::query_as("SELECT id FROM workflow.timers WHERE workflow_id = $1 AND seq = $2")
1003 .bind(&timer.workflow_id)
1004 .bind(timer.seq)
1005 .fetch_one(&self.pool)
1006 .await?;
1007 Ok(id)
1008 }
1009
1010 async fn cancel_pending_activities(&self, workflow_id: &str) -> Result<u64> {
1011 let res = sqlx::query(
1012 "UPDATE workflow.activities SET status = 'CANCELLED', completed_at = $1
1013 WHERE workflow_id = $2 AND status = 'PENDING'",
1014 )
1015 .bind(timestamp_now())
1016 .bind(workflow_id)
1017 .execute(&self.pool)
1018 .await?;
1019 Ok(res.rows_affected())
1020 }
1021
1022 async fn cancel_pending_timers(&self, workflow_id: &str) -> Result<u64> {
1023 let res = sqlx::query(
1024 "UPDATE workflow.timers SET fired = TRUE
1025 WHERE workflow_id = $1 AND fired = FALSE",
1026 )
1027 .bind(workflow_id)
1028 .execute(&self.pool)
1029 .await?;
1030 Ok(res.rows_affected())
1031 }
1032
1033 async fn get_timer_by_workflow_seq(
1034 &self,
1035 workflow_id: &str,
1036 seq: i32,
1037 ) -> Result<Option<WorkflowTimer>> {
1038 let row = sqlx::query_as::<_, PgTimerRow>(
1039 "SELECT id, workflow_id, seq, fire_at, fired
1040 FROM workflow.timers WHERE workflow_id = $1 AND seq = $2",
1041 )
1042 .bind(workflow_id)
1043 .bind(seq)
1044 .fetch_optional(&self.pool)
1045 .await?;
1046 Ok(row.map(Into::into))
1047 }
1048
1049 async fn fire_due_timers(&self, now: f64) -> Result<Vec<WorkflowTimer>> {
1050 let rows = sqlx::query_as::<_, PgTimerRow>(
1051 "UPDATE workflow.timers SET fired = TRUE
1052 WHERE fired = FALSE AND fire_at <= $1
1053 RETURNING id, workflow_id, seq, fire_at, fired",
1054 )
1055 .bind(now)
1056 .fetch_all(&self.pool)
1057 .await?;
1058 Ok(rows.into_iter().map(Into::into).collect())
1059 }
1060
1061 async fn send_signal(&self, sig: &WorkflowSignal) -> Result<i64> {
1064 let row: (i64,) = sqlx::query_as(
1065 "INSERT INTO workflow.signals (workflow_id, name, payload, consumed, received_at) VALUES ($1, $2, $3, FALSE, $4) RETURNING id",
1066 )
1067 .bind(&sig.workflow_id)
1068 .bind(&sig.name)
1069 .bind(&sig.payload)
1070 .bind(sig.received_at)
1071 .fetch_one(&self.pool)
1072 .await?;
1073 Ok(row.0)
1074 }
1075
1076 async fn consume_signals(&self, workflow_id: &str, name: &str) -> Result<Vec<WorkflowSignal>> {
1077 let rows = sqlx::query_as::<_, PgSignalRow>(
1078 "UPDATE workflow.signals SET consumed = TRUE
1079 WHERE workflow_id = $1 AND name = $2 AND consumed = FALSE
1080 RETURNING id, workflow_id, name, payload, consumed, received_at",
1081 )
1082 .bind(workflow_id)
1083 .bind(name)
1084 .fetch_all(&self.pool)
1085 .await?;
1086 Ok(rows.into_iter().map(Into::into).collect())
1087 }
1088
1089 async fn create_schedule(&self, sched: &WorkflowSchedule) -> Result<()> {
1092 sqlx::query(
1093 "INSERT INTO workflow.schedules (namespace, name, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at)
1094 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13)",
1095 )
1096 .bind(&sched.namespace)
1097 .bind(&sched.name)
1098 .bind(&sched.workflow_type)
1099 .bind(&sched.cron_expr)
1100 .bind(&sched.timezone)
1101 .bind(&sched.input)
1102 .bind(&sched.task_queue)
1103 .bind(&sched.overlap_policy)
1104 .bind(sched.paused)
1105 .bind(sched.last_run_at)
1106 .bind(sched.next_run_at)
1107 .bind(&sched.last_workflow_id)
1108 .bind(sched.created_at)
1109 .execute(&self.pool)
1110 .await?;
1111 Ok(())
1112 }
1113
1114 async fn get_schedule(&self, namespace: &str, name: &str) -> Result<Option<WorkflowSchedule>> {
1115 let row = sqlx::query_as::<_, PgScheduleRow>(
1116 "SELECT namespace, name, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at FROM workflow.schedules WHERE namespace = $1 AND name = $2",
1117 )
1118 .bind(namespace)
1119 .bind(name)
1120 .fetch_optional(&self.pool)
1121 .await?;
1122 Ok(row.map(Into::into))
1123 }
1124
1125 async fn list_schedules(&self, namespace: &str) -> Result<Vec<WorkflowSchedule>> {
1126 let rows = sqlx::query_as::<_, PgScheduleRow>(
1127 "SELECT namespace, name, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at FROM workflow.schedules WHERE namespace = $1 ORDER BY name",
1128 )
1129 .bind(namespace)
1130 .fetch_all(&self.pool)
1131 .await?;
1132 Ok(rows.into_iter().map(Into::into).collect())
1133 }
1134
1135 async fn update_schedule_last_run(
1136 &self,
1137 namespace: &str,
1138 name: &str,
1139 last_run_at: f64,
1140 next_run_at: f64,
1141 workflow_id: &str,
1142 ) -> Result<()> {
1143 sqlx::query(
1144 "UPDATE workflow.schedules SET last_run_at = $1, next_run_at = $2, last_workflow_id = $3 WHERE namespace = $4 AND name = $5",
1145 )
1146 .bind(last_run_at)
1147 .bind(next_run_at)
1148 .bind(workflow_id)
1149 .bind(namespace)
1150 .bind(name)
1151 .execute(&self.pool)
1152 .await?;
1153 Ok(())
1154 }
1155
1156 async fn delete_schedule(&self, namespace: &str, name: &str) -> Result<bool> {
1157 let res = sqlx::query("DELETE FROM workflow.schedules WHERE namespace = $1 AND name = $2")
1158 .bind(namespace)
1159 .bind(name)
1160 .execute(&self.pool)
1161 .await?;
1162 Ok(res.rows_affected() > 0)
1163 }
1164
1165 async fn list_archivable_workflows(
1166 &self,
1167 cutoff: f64,
1168 limit: i64,
1169 ) -> Result<Vec<WorkflowRecord>> {
1170 let rows = sqlx::query_as::<_, PgWorkflowRow>(
1171 "SELECT id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at
1172 FROM workflow.workflows
1173 WHERE status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')
1174 AND completed_at IS NOT NULL
1175 AND completed_at < $1
1176 AND archived_at IS NULL
1177 ORDER BY completed_at ASC
1178 LIMIT $2",
1179 )
1180 .bind(cutoff)
1181 .bind(limit)
1182 .fetch_all(&self.pool)
1183 .await?;
1184 Ok(rows.into_iter().map(Into::into).collect())
1185 }
1186
1187 async fn mark_archived_and_purge(
1188 &self,
1189 workflow_id: &str,
1190 archive_uri: &str,
1191 archived_at: f64,
1192 ) -> Result<()> {
1193 let mut tx = self.pool.begin().await?;
1194 sqlx::query("DELETE FROM workflow.events WHERE workflow_id = $1")
1195 .bind(workflow_id)
1196 .execute(&mut *tx)
1197 .await?;
1198 sqlx::query("DELETE FROM workflow.activities WHERE workflow_id = $1")
1199 .bind(workflow_id)
1200 .execute(&mut *tx)
1201 .await?;
1202 sqlx::query("DELETE FROM workflow.timers WHERE workflow_id = $1")
1203 .bind(workflow_id)
1204 .execute(&mut *tx)
1205 .await?;
1206 sqlx::query("DELETE FROM workflow.signals WHERE workflow_id = $1")
1207 .bind(workflow_id)
1208 .execute(&mut *tx)
1209 .await?;
1210 sqlx::query("DELETE FROM workflow.snapshots WHERE workflow_id = $1")
1211 .bind(workflow_id)
1212 .execute(&mut *tx)
1213 .await?;
1214 sqlx::query(
1215 "UPDATE workflow.workflows SET archived_at = $1, archive_uri = $2 WHERE id = $3",
1216 )
1217 .bind(archived_at)
1218 .bind(archive_uri)
1219 .bind(workflow_id)
1220 .execute(&mut *tx)
1221 .await?;
1222 tx.commit().await?;
1223 Ok(())
1224 }
1225
1226 async fn upsert_search_attributes(&self, workflow_id: &str, patch_json: &str) -> Result<()> {
1227 let current: Option<(Option<String>,)> =
1228 sqlx::query_as("SELECT search_attributes FROM workflow.workflows WHERE id = $1")
1229 .bind(workflow_id)
1230 .fetch_optional(&self.pool)
1231 .await?;
1232 let merged = crate::store::sqlite::merge_search_attrs(
1233 current.and_then(|(s,)| s).as_deref(),
1234 patch_json,
1235 )?;
1236 sqlx::query("UPDATE workflow.workflows SET search_attributes = $1 WHERE id = $2")
1237 .bind(merged)
1238 .bind(workflow_id)
1239 .execute(&self.pool)
1240 .await?;
1241 Ok(())
1242 }
1243
1244 async fn update_schedule(
1245 &self,
1246 namespace: &str,
1247 name: &str,
1248 patch: &SchedulePatch,
1249 ) -> Result<Option<WorkflowSchedule>> {
1250 let mut sets: Vec<String> = Vec::new();
1251 let mut idx = 1usize;
1252 if patch.cron_expr.is_some() {
1253 sets.push(format!("cron_expr = ${idx}"));
1254 idx += 1;
1255 }
1256 if patch.timezone.is_some() {
1257 sets.push(format!("timezone = ${idx}"));
1258 idx += 1;
1259 }
1260 if patch.input.is_some() {
1261 sets.push(format!("input = ${idx}"));
1262 idx += 1;
1263 }
1264 if patch.task_queue.is_some() {
1265 sets.push(format!("task_queue = ${idx}"));
1266 idx += 1;
1267 }
1268 if patch.overlap_policy.is_some() {
1269 sets.push(format!("overlap_policy = ${idx}"));
1270 idx += 1;
1271 }
1272 if sets.is_empty() {
1273 return self.get_schedule(namespace, name).await;
1274 }
1275 let sql = format!(
1276 "UPDATE workflow.schedules SET {} WHERE namespace = ${} AND name = ${}",
1277 sets.join(", "),
1278 idx,
1279 idx + 1
1280 );
1281 let mut q = sqlx::query(&sql);
1282 if let Some(ref v) = patch.cron_expr {
1283 q = q.bind(v);
1284 }
1285 if let Some(ref v) = patch.timezone {
1286 q = q.bind(v);
1287 }
1288 if let Some(ref v) = patch.input {
1289 q = q.bind(v.to_string());
1290 }
1291 if let Some(ref v) = patch.task_queue {
1292 q = q.bind(v);
1293 }
1294 if let Some(ref v) = patch.overlap_policy {
1295 q = q.bind(v);
1296 }
1297 let res = q.bind(namespace).bind(name).execute(&self.pool).await?;
1298 if res.rows_affected() == 0 {
1299 return Ok(None);
1300 }
1301 self.get_schedule(namespace, name).await
1302 }
1303
1304 async fn set_schedule_paused(
1305 &self,
1306 namespace: &str,
1307 name: &str,
1308 paused: bool,
1309 ) -> Result<Option<WorkflowSchedule>> {
1310 let res = sqlx::query(
1311 "UPDATE workflow.schedules SET paused = $1 WHERE namespace = $2 AND name = $3",
1312 )
1313 .bind(paused)
1314 .bind(namespace)
1315 .bind(name)
1316 .execute(&self.pool)
1317 .await?;
1318 if res.rows_affected() == 0 {
1319 return Ok(None);
1320 }
1321 self.get_schedule(namespace, name).await
1322 }
1323
1324 async fn register_worker(&self, w: &WorkflowWorker) -> Result<()> {
1327 sqlx::query(
1328 "INSERT INTO workflow.workers (id, namespace, identity, task_queue, workflows, activities, max_concurrent_workflows, max_concurrent_activities, active_tasks, last_heartbeat, registered_at)
1329 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
1330 ON CONFLICT (id) DO UPDATE SET last_heartbeat = EXCLUDED.last_heartbeat, identity = EXCLUDED.identity",
1331 )
1332 .bind(&w.id)
1333 .bind(&w.namespace)
1334 .bind(&w.identity)
1335 .bind(&w.task_queue)
1336 .bind(&w.workflows)
1337 .bind(&w.activities)
1338 .bind(w.max_concurrent_workflows)
1339 .bind(w.max_concurrent_activities)
1340 .bind(w.active_tasks)
1341 .bind(w.last_heartbeat)
1342 .bind(w.registered_at)
1343 .execute(&self.pool)
1344 .await?;
1345 Ok(())
1346 }
1347
1348 async fn heartbeat_worker(&self, id: &str, now: f64) -> Result<()> {
1349 sqlx::query("UPDATE workflow.workers SET last_heartbeat = $1 WHERE id = $2")
1350 .bind(now)
1351 .bind(id)
1352 .execute(&self.pool)
1353 .await?;
1354 Ok(())
1355 }
1356
1357 async fn list_workers(&self, namespace: &str) -> Result<Vec<WorkflowWorker>> {
1358 let rows = sqlx::query_as::<_, PgWorkerRow>(
1359 "SELECT id, namespace, identity, task_queue, workflows, activities, max_concurrent_workflows, max_concurrent_activities, active_tasks, last_heartbeat, registered_at FROM workflow.workers WHERE namespace = $1 ORDER BY registered_at",
1360 )
1361 .bind(namespace)
1362 .fetch_all(&self.pool)
1363 .await?;
1364 Ok(rows.into_iter().map(Into::into).collect())
1365 }
1366
1367 async fn remove_dead_workers(&self, cutoff: f64) -> Result<Vec<String>> {
1368 let rows: Vec<(String,)> =
1369 sqlx::query_as("SELECT id FROM workflow.workers WHERE last_heartbeat < $1")
1370 .bind(cutoff)
1371 .fetch_all(&self.pool)
1372 .await?;
1373 let ids: Vec<String> = rows.into_iter().map(|r| r.0).collect();
1374 if !ids.is_empty() {
1375 sqlx::query("DELETE FROM workflow.workers WHERE last_heartbeat < $1")
1376 .bind(cutoff)
1377 .execute(&self.pool)
1378 .await?;
1379 }
1380 Ok(ids)
1381 }
1382
1383 async fn list_child_workflows(&self, parent_id: &str) -> Result<Vec<WorkflowRecord>> {
1386 let rows = sqlx::query_as::<_, PgWorkflowRow>(
1387 "SELECT id, namespace, run_id, workflow_type, task_queue, status, input, result, error, parent_id, claimed_by, search_attributes, archived_at, archive_uri, created_at, updated_at, completed_at
1388 FROM workflow.workflows WHERE parent_id = $1 ORDER BY created_at ASC",
1389 )
1390 .bind(parent_id)
1391 .fetch_all(&self.pool)
1392 .await?;
1393 Ok(rows.into_iter().map(Into::into).collect())
1394 }
1395
1396 async fn create_snapshot(
1399 &self,
1400 workflow_id: &str,
1401 event_seq: i32,
1402 state_json: &str,
1403 ) -> Result<()> {
1404 sqlx::query(
1405 "INSERT INTO workflow.snapshots (workflow_id, event_seq, state_json, created_at)
1406 VALUES ($1, $2, $3, $4)
1407 ON CONFLICT (workflow_id, event_seq) DO UPDATE SET state_json = EXCLUDED.state_json, created_at = EXCLUDED.created_at",
1408 )
1409 .bind(workflow_id)
1410 .bind(event_seq)
1411 .bind(state_json)
1412 .bind(timestamp_now())
1413 .execute(&self.pool)
1414 .await?;
1415 Ok(())
1416 }
1417
1418 async fn get_latest_snapshot(&self, workflow_id: &str) -> Result<Option<WorkflowSnapshot>> {
1419 let row = sqlx::query_as::<_, (String, i32, String, f64)>(
1420 "SELECT workflow_id, event_seq, state_json, created_at
1421 FROM workflow.snapshots WHERE workflow_id = $1
1422 ORDER BY event_seq DESC LIMIT 1",
1423 )
1424 .bind(workflow_id)
1425 .fetch_optional(&self.pool)
1426 .await?;
1427
1428 Ok(row.map(
1429 |(workflow_id, event_seq, state_json, created_at)| WorkflowSnapshot {
1430 workflow_id,
1431 event_seq,
1432 state_json,
1433 created_at,
1434 },
1435 ))
1436 }
1437
1438 async fn get_queue_stats(&self, namespace: &str) -> Result<Vec<crate::store::QueueStats>> {
1441 let rows = sqlx::query_as::<_, (String, i64, i64, i64)>(
1442 "SELECT
1443 a.task_queue AS queue,
1444 SUM(CASE WHEN a.status = 'PENDING' THEN 1 ELSE 0 END) AS pending,
1445 SUM(CASE WHEN a.status = 'RUNNING' THEN 1 ELSE 0 END) AS running,
1446 (SELECT COUNT(*) FROM workflow.workers w WHERE w.task_queue = a.task_queue AND w.namespace = $1) AS workers
1447 FROM workflow.activities a
1448 JOIN workflow.workflows wf ON a.workflow_id = wf.id AND wf.namespace = $1
1449 GROUP BY a.task_queue",
1450 )
1451 .bind(namespace)
1452 .fetch_all(&self.pool)
1453 .await?;
1454
1455 Ok(rows
1456 .into_iter()
1457 .map(
1458 |(queue, pending, running, workers)| crate::store::QueueStats {
1459 queue,
1460 pending_activities: pending,
1461 running_activities: running,
1462 workers,
1463 },
1464 )
1465 .collect())
1466 }
1467
1468 async fn try_acquire_scheduler_lock(&self) -> Result<bool> {
1471 let row: (bool,) = sqlx::query_as("SELECT pg_try_advisory_lock(42)")
1475 .fetch_one(&self.pool)
1476 .await?;
1477 Ok(row.0)
1478 }
1479}
1480
1481fn timestamp_now() -> f64 {
1482 std::time::SystemTime::now()
1483 .duration_since(std::time::UNIX_EPOCH)
1484 .unwrap()
1485 .as_secs_f64()
1486}
1487
1488#[derive(sqlx::FromRow)]
1491struct PgWorkflowRow {
1492 id: String,
1493 namespace: String,
1494 run_id: String,
1495 workflow_type: String,
1496 task_queue: String,
1497 status: String,
1498 input: Option<String>,
1499 result: Option<String>,
1500 error: Option<String>,
1501 parent_id: Option<String>,
1502 claimed_by: Option<String>,
1503 search_attributes: Option<String>,
1504 archived_at: Option<f64>,
1505 archive_uri: Option<String>,
1506 created_at: f64,
1507 updated_at: f64,
1508 completed_at: Option<f64>,
1509}
1510
1511impl From<PgWorkflowRow> for WorkflowRecord {
1512 fn from(r: PgWorkflowRow) -> Self {
1513 Self {
1514 id: r.id,
1515 namespace: r.namespace,
1516 run_id: r.run_id,
1517 workflow_type: r.workflow_type,
1518 task_queue: r.task_queue,
1519 status: r.status,
1520 input: r.input,
1521 result: r.result,
1522 error: r.error,
1523 parent_id: r.parent_id,
1524 claimed_by: r.claimed_by,
1525 search_attributes: r.search_attributes,
1526 archived_at: r.archived_at,
1527 archive_uri: r.archive_uri,
1528 created_at: r.created_at,
1529 updated_at: r.updated_at,
1530 completed_at: r.completed_at,
1531 }
1532 }
1533}
1534
1535#[derive(sqlx::FromRow)]
1536struct PgEventRow {
1537 id: i64,
1538 workflow_id: String,
1539 seq: i32,
1540 event_type: String,
1541 payload: Option<String>,
1542 timestamp: f64,
1543}
1544
1545impl From<PgEventRow> for WorkflowEvent {
1546 fn from(r: PgEventRow) -> Self {
1547 Self {
1548 id: Some(r.id),
1549 workflow_id: r.workflow_id,
1550 seq: r.seq,
1551 event_type: r.event_type,
1552 payload: r.payload,
1553 timestamp: r.timestamp,
1554 }
1555 }
1556}
1557
1558#[derive(sqlx::FromRow)]
1559struct PgActivityRow {
1560 id: i64,
1561 workflow_id: String,
1562 seq: i32,
1563 name: String,
1564 task_queue: String,
1565 input: Option<String>,
1566 status: String,
1567 result: Option<String>,
1568 error: Option<String>,
1569 attempt: i32,
1570 max_attempts: i32,
1571 initial_interval_secs: f64,
1572 backoff_coefficient: f64,
1573 start_to_close_secs: f64,
1574 heartbeat_timeout_secs: Option<f64>,
1575 claimed_by: Option<String>,
1576 scheduled_at: f64,
1577 started_at: Option<f64>,
1578 completed_at: Option<f64>,
1579 last_heartbeat: Option<f64>,
1580}
1581
1582impl From<PgActivityRow> for WorkflowActivity {
1583 fn from(r: PgActivityRow) -> Self {
1584 Self {
1585 id: Some(r.id),
1586 workflow_id: r.workflow_id,
1587 seq: r.seq,
1588 name: r.name,
1589 task_queue: r.task_queue,
1590 input: r.input,
1591 status: r.status,
1592 result: r.result,
1593 error: r.error,
1594 attempt: r.attempt,
1595 max_attempts: r.max_attempts,
1596 initial_interval_secs: r.initial_interval_secs,
1597 backoff_coefficient: r.backoff_coefficient,
1598 start_to_close_secs: r.start_to_close_secs,
1599 heartbeat_timeout_secs: r.heartbeat_timeout_secs,
1600 claimed_by: r.claimed_by,
1601 scheduled_at: r.scheduled_at,
1602 started_at: r.started_at,
1603 completed_at: r.completed_at,
1604 last_heartbeat: r.last_heartbeat,
1605 }
1606 }
1607}
1608
1609#[derive(sqlx::FromRow)]
1610struct PgTimerRow {
1611 id: i64,
1612 workflow_id: String,
1613 seq: i32,
1614 fire_at: f64,
1615 fired: bool,
1616}
1617
1618impl From<PgTimerRow> for WorkflowTimer {
1619 fn from(r: PgTimerRow) -> Self {
1620 Self {
1621 id: Some(r.id),
1622 workflow_id: r.workflow_id,
1623 seq: r.seq,
1624 fire_at: r.fire_at,
1625 fired: r.fired,
1626 }
1627 }
1628}
1629
1630#[derive(sqlx::FromRow)]
1631struct PgSignalRow {
1632 id: i64,
1633 workflow_id: String,
1634 name: String,
1635 payload: Option<String>,
1636 consumed: bool,
1637 received_at: f64,
1638}
1639
1640impl From<PgSignalRow> for WorkflowSignal {
1641 fn from(r: PgSignalRow) -> Self {
1642 Self {
1643 id: Some(r.id),
1644 workflow_id: r.workflow_id,
1645 name: r.name,
1646 payload: r.payload,
1647 consumed: r.consumed,
1648 received_at: r.received_at,
1649 }
1650 }
1651}
1652
1653#[derive(sqlx::FromRow)]
1654struct PgScheduleRow {
1655 namespace: String,
1656 name: String,
1657 workflow_type: String,
1658 cron_expr: String,
1659 timezone: String,
1660 input: Option<String>,
1661 task_queue: String,
1662 overlap_policy: String,
1663 paused: bool,
1664 last_run_at: Option<f64>,
1665 next_run_at: Option<f64>,
1666 last_workflow_id: Option<String>,
1667 created_at: f64,
1668}
1669
1670impl From<PgScheduleRow> for WorkflowSchedule {
1671 fn from(r: PgScheduleRow) -> Self {
1672 Self {
1673 namespace: r.namespace,
1674 name: r.name,
1675 workflow_type: r.workflow_type,
1676 cron_expr: r.cron_expr,
1677 timezone: r.timezone,
1678 input: r.input,
1679 task_queue: r.task_queue,
1680 overlap_policy: r.overlap_policy,
1681 paused: r.paused,
1682 last_run_at: r.last_run_at,
1683 next_run_at: r.next_run_at,
1684 last_workflow_id: r.last_workflow_id,
1685 created_at: r.created_at,
1686 }
1687 }
1688}
1689
1690#[derive(sqlx::FromRow)]
1691struct PgWorkerRow {
1692 id: String,
1693 namespace: String,
1694 identity: String,
1695 task_queue: String,
1696 workflows: Option<String>,
1697 activities: Option<String>,
1698 max_concurrent_workflows: i32,
1699 max_concurrent_activities: i32,
1700 active_tasks: i32,
1701 last_heartbeat: f64,
1702 registered_at: f64,
1703}
1704
1705impl From<PgWorkerRow> for WorkflowWorker {
1706 fn from(r: PgWorkerRow) -> Self {
1707 Self {
1708 id: r.id,
1709 namespace: r.namespace,
1710 identity: r.identity,
1711 task_queue: r.task_queue,
1712 workflows: r.workflows,
1713 activities: r.activities,
1714 max_concurrent_workflows: r.max_concurrent_workflows,
1715 max_concurrent_activities: r.max_concurrent_activities,
1716 active_tasks: r.active_tasks,
1717 last_heartbeat: r.last_heartbeat,
1718 registered_at: r.registered_at,
1719 }
1720 }
1721}
1722
1723#[cfg(test)]
1724mod tests {
1725 use super::*;
1726
1727 #[test]
1728 fn sanitise_schema_keeps_statements_intact() {
1729 let input = "CREATE TABLE foo (x INT);\nCREATE INDEX idx_foo ON foo(x);\n";
1730 let out = sanitise_schema(input);
1731 assert_eq!(out.len(), 2);
1732 assert!(out[0].starts_with("CREATE TABLE foo"));
1733 assert!(out[1].starts_with("CREATE INDEX idx_foo"));
1734 }
1735
1736 #[test]
1737 fn sanitise_schema_drops_pure_comment_lines() {
1738 let input = "-- header comment\nCREATE TABLE foo (x INT);\n-- trailing comment\n";
1739 let out = sanitise_schema(input);
1740 assert_eq!(out.len(), 1);
1741 assert!(out[0].starts_with("CREATE TABLE foo"));
1742 }
1743
1744 #[test]
1745 fn sanitise_schema_ignores_semicolons_inside_comment_prose() {
1746 let input = "\
1750CREATE TABLE foo (x INT);
1751-- Idempotent across startups; fresh installs pick the column up from the
1752-- CREATE TABLE above so the ADD is a no-op.
1753";
1754 let out = sanitise_schema(input);
1755 assert_eq!(
1756 out.len(),
1757 1,
1758 "expected 1 real statement, got {}: {:?}",
1759 out.len(),
1760 out
1761 );
1762 assert!(out[0].starts_with("CREATE TABLE foo"));
1763 }
1764
1765 #[test]
1766 fn sanitise_schema_drops_indented_comment_lines() {
1767 let input = " -- indented comment\n\tCREATE TABLE foo (x INT);\n";
1768 let out = sanitise_schema(input);
1769 assert_eq!(out.len(), 1);
1770 assert!(out[0].contains("CREATE TABLE foo"));
1771 }
1772
1773 #[test]
1774 fn sanitise_schema_real_constant_produces_only_ddl() {
1775 for stmt in sanitise_schema(SCHEMA) {
1779 let first_word = stmt
1780 .split_whitespace()
1781 .next()
1782 .expect("non-empty statement")
1783 .to_uppercase();
1784 assert!(
1785 matches!(
1786 first_word.as_str(),
1787 "CREATE" | "INSERT" | "UPDATE" | "DROP" | "ALTER" | "WITH"
1788 ),
1789 "SCHEMA produced non-DDL statement starting with {first_word:?}: {stmt:?}"
1790 );
1791 }
1792 }
1793}