Skip to main content

assay_workflow/store/
postgres.rs

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
10/// v0.1.2 schema layout: workflow tables live in the `workflow` schema;
11/// the engine-events outbox lives in the `engine` schema (created
12/// alongside the engine-core tables by `assay_domain::engine`). The
13/// store creates the workflow schema first, then runs DDL against it
14/// schema-qualified.
15const 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
166/// One-shot relocation: moves v0.13.1's prefixed `public.workflow_*`
167/// tables into the `workflow` schema. Idempotent — each step gates on
168/// `to_regclass(public.<old>) IS NOT NULL` so fresh installs and
169/// already-migrated DBs are no-ops.
170///
171/// SCHEMA above already created empty `workflow.*` tables; we DROP
172/// them with CASCADE here before ALTER TABLE … SET SCHEMA so the move
173/// has somewhere to land. RESTRICT would fail on the FKs between
174/// workflow.events / .activities / .timers / .signals / .snapshots
175/// and workflow.workflows.
176const 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
277/// Split a Postgres DDL script into individual statements ready for `sqlx::query`.
278///
279/// Drops pure-comment lines (those starting with `--` after optional whitespace)
280/// *before* splitting on `;`. Without this step, a semicolon inside a line comment
281/// (e.g. `-- Idempotent across startups; fresh installs pick the column up`) would
282/// split the surrounding comment into fragments — one of which is naked prose that
283/// Postgres tries to parse as SQL and rejects with `syntax error at or near "<word>"`.
284///
285/// The filter only drops *pure-comment* lines (leading whitespace then `--`), leaving
286/// `--`-after-code untouched. That keeps string literals safe (could legally contain
287/// `--`) and is conservative enough to remain correct if the SCHEMA grows more prose.
288fn 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/// `Clone` is derived because the underlying `PgPool` is itself `Clone`
303/// (it's `Arc<PoolInner>` internally) — cloning the store hands back a
304/// new wrapper around the same connection pool. Required so engine
305/// composition (`EngineState<S>`) can derive `Clone` and pass through
306/// axum `with_state`.
307#[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    /// Build a store from an existing pool. Runs migrations on the target
319    /// database. Useful when the engine owns the pool (shared with other
320    /// modules) and hands a clone to the workflow module, or for tests that
321    /// point many stores at different databases in the same Postgres server.
322    pub async fn from_pool(pool: PgPool) -> Result<Self> {
323        let store = Self { pool };
324        store.migrate().await?;
325        Ok(store)
326    }
327
328    /// Expose the underlying pool (used by the engine to build a
329    /// `PgEngineEventBus` that shares the same connection pool).
330    pub fn pool(&self) -> &PgPool {
331        &self.pool
332    }
333
334    async fn migrate(&self) -> Result<()> {
335        // Apply the base schema (tables + indexes) statement-by-statement.
336        // This creates the workflow + engine schemas and the v0.13.2
337        // schema-qualified tables. On a fresh install this is the only
338        // step that runs; on an upgrade from v0.13.1 the empty new
339        // tables are dropped + replaced by the legacy public.* tables
340        // in the relocation block below.
341        for statement in sanitise_schema(SCHEMA) {
342            sqlx::query(&statement).execute(&self.pool).await?;
343        }
344        // v0.13.1 → v0.13.2 relocation. Idempotent: on fresh installs
345        // the public.* tables don't exist and every branch is a no-op.
346        sqlx::raw_sql(V0_13_2_RELOCATION_SQL)
347            .execute(&self.pool)
348            .await?;
349        // Drop the v0.13.0 LISTEN/NOTIFY triggers if they still exist on
350        // the target database. The Rust-managed CDC outbox in
351        // assay_domain::events is the replacement; leaving stale
352        // triggers in place would double-publish NOTIFYs with channels
353        // no one listens to. Triggers reference the post-relocation
354        // table names (workflow.workflows, workflow.activities) so they
355        // execute regardless of which side of the migration we're on.
356        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    /// Try to acquire pg_advisory_lock for leader election.
370    /// Returns true if this instance is the leader (scheduler should run).
371    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    // ── Namespaces ─────────────────────────────────────────
381
382    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    // ── Workflows ──────────────────────────────────────────
464
465    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        // Bind placeholders for the filter follow $3; next index is 4.
527        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            // JSONB ->> always returns TEXT; compare by stringified value.
549            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        // Atomic claim with FOR UPDATE SKIP LOCKED so multiple engine
613        // replicas don't fight over the same workflow task.
614        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    // ── Events ─────────────────────────────────────────────
668
669    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    // ── Activities ──────────────────────────────────────────
742
743    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        // Atomic claim using FOR UPDATE SKIP LOCKED — prevents contention
800        // between multiple assay serve instances claiming the same activity
801        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    // ── Timers ──────────────────────────────────────────────
980
981    async fn create_timer(&self, timer: &WorkflowTimer) -> Result<i64> {
982        // Idempotent: ON CONFLICT (workflow_id, seq) DO NOTHING.
983        // If a row already exists, RETURNING produces no rows — fall back to SELECT.
984        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        // Row already existed — return its id.
1001        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    // ── Signals ─────────────────────────────────────────────
1062
1063    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    // ── Schedules ───────────────────────────────────────────
1090
1091    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    // ── Workers ─────────────────────────────────────────────
1325
1326    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    // ── Child Workflows ─────────────────────────────────────
1384
1385    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    // ── Snapshots ───────────────────────────────────────────
1397
1398    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    // ── Queue Stats ─────────────────────────────────────────
1439
1440    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    // ── Leader Election ─────────────────────────────────────
1469
1470    async fn try_acquire_scheduler_lock(&self) -> Result<bool> {
1471        // pg_try_advisory_lock is session-scoped — only one connection
1472        // in the pool will hold the lock. In a multi-replica Kubernetes
1473        // deployment, only one pod's connection wins.
1474        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// ── Postgres row types (sqlx::FromRow) ──────────────────────
1489
1490#[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        // Regression: the exact shape that broke v0.11.3–v0.11.5 in production.
1747        // `-- foo; bar` used to split into "foo" and " bar" fragments, the second
1748        // of which was executed as SQL and rejected with `syntax error at or near "bar"`.
1749        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        // The real SCHEMA constant must not produce any statement whose first
1776        // token isn't a recognised SQL keyword. A prose fragment leaking in
1777        // (e.g. "fresh installs...") means the filter regressed.
1778        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}