Skip to main content

assay_workflow/store/
sqlite.rs

1use anyhow::Result;
2use sqlx::SqlitePool;
3use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
4
5use crate::store::{NamespaceRecord, NamespaceStats, QueueStats, WorkflowStore};
6use crate::types::*;
7
8/// Workflow-module DDL. v0.1.2 schema-qualifies every table to the
9/// `workflow` schema, which on SQLite is an attached database (one
10/// `workflow.db` file per data dir, attached on connect). On PG the
11/// same DDL targets the `workflow` schema.
12///
13/// `engine.events` and `engine.lock` (engine-core infrastructure) live
14/// in the `engine` attachment; engine-core DDL is owned by
15/// `assay_domain::engine::SqliteEngineSchema`. We still bootstrap them
16/// here for v0.1.2 because the workflow store is the embedder for the
17/// `engine.events` notification outbox and for the SQLite single-instance
18/// lock — both pre-date the engine-core schema and stay co-located on
19/// SQLite to keep `SqliteStore::new(url)` self-sufficient for tests.
20const SCHEMA: &str = r#"
21CREATE TABLE IF NOT EXISTS workflow.namespaces (
22    name            TEXT PRIMARY KEY,
23    created_at      REAL NOT NULL
24);
25
26INSERT OR IGNORE INTO workflow.namespaces (name, created_at)
27    VALUES ('main', strftime('%s', 'now'));
28
29CREATE TABLE IF NOT EXISTS workflow.workflows (
30    id              TEXT PRIMARY KEY,
31    namespace       TEXT NOT NULL DEFAULT 'main',
32    run_id          TEXT NOT NULL,
33    workflow_type   TEXT NOT NULL,
34    task_queue      TEXT NOT NULL DEFAULT 'main',
35    status          TEXT NOT NULL DEFAULT 'PENDING',
36    input           TEXT,
37    result          TEXT,
38    error           TEXT,
39    parent_id       TEXT,
40    claimed_by      TEXT,
41    search_attributes TEXT,
42    archived_at     REAL,
43    archive_uri     TEXT,
44    -- Workflow-task dispatch (Phase 9): a workflow is "dispatchable" when
45    -- it has new events a worker needs to replay against. Set true on
46    -- start, on activity completion, on timer fire, on signal arrival.
47    -- Cleared when a worker claims the dispatch lease.
48    needs_dispatch  INTEGER NOT NULL DEFAULT 0,
49    dispatch_claimed_by    TEXT,
50    dispatch_last_heartbeat REAL,
51    created_at      REAL NOT NULL,
52    updated_at      REAL NOT NULL,
53    completed_at    REAL
54);
55CREATE INDEX IF NOT EXISTS workflow.idx_wf_status_queue ON workflows(status, task_queue);
56CREATE INDEX IF NOT EXISTS workflow.idx_wf_namespace ON workflows(namespace);
57CREATE INDEX IF NOT EXISTS workflow.idx_wf_dispatch ON workflows(task_queue, needs_dispatch, dispatch_claimed_by);
58
59CREATE TABLE IF NOT EXISTS workflow.events (
60    id              INTEGER PRIMARY KEY AUTOINCREMENT,
61    workflow_id     TEXT NOT NULL REFERENCES workflows(id),
62    seq             INTEGER NOT NULL,
63    event_type      TEXT NOT NULL,
64    payload         TEXT,
65    timestamp       REAL NOT NULL
66);
67CREATE INDEX IF NOT EXISTS workflow.idx_wf_events_lookup ON events(workflow_id, seq);
68
69CREATE TABLE IF NOT EXISTS workflow.activities (
70    id              INTEGER PRIMARY KEY AUTOINCREMENT,
71    workflow_id     TEXT NOT NULL REFERENCES workflows(id),
72    seq             INTEGER NOT NULL,
73    name            TEXT NOT NULL,
74    task_queue      TEXT NOT NULL DEFAULT 'main',
75    input           TEXT,
76    status          TEXT NOT NULL DEFAULT 'PENDING',
77    result          TEXT,
78    error           TEXT,
79    attempt         INTEGER NOT NULL DEFAULT 1,
80    max_attempts    INTEGER NOT NULL DEFAULT 3,
81    initial_interval_secs   REAL NOT NULL DEFAULT 1,
82    backoff_coefficient     REAL NOT NULL DEFAULT 2,
83    start_to_close_secs     REAL NOT NULL DEFAULT 300,
84    heartbeat_timeout_secs  REAL,
85    claimed_by      TEXT,
86    scheduled_at    REAL NOT NULL,
87    started_at      REAL,
88    completed_at    REAL,
89    last_heartbeat  REAL,
90    UNIQUE (workflow_id, seq)
91);
92CREATE INDEX IF NOT EXISTS workflow.idx_wf_act_pending ON activities(task_queue, status, scheduled_at);
93
94CREATE TABLE IF NOT EXISTS workflow.timers (
95    id              INTEGER PRIMARY KEY AUTOINCREMENT,
96    workflow_id     TEXT NOT NULL REFERENCES workflows(id),
97    seq             INTEGER NOT NULL,
98    fire_at         REAL NOT NULL,
99    fired           INTEGER NOT NULL DEFAULT 0,
100    UNIQUE (workflow_id, seq)
101);
102CREATE INDEX IF NOT EXISTS workflow.idx_wf_timers_due ON timers(fire_at);
103
104CREATE TABLE IF NOT EXISTS workflow.signals (
105    id              INTEGER PRIMARY KEY AUTOINCREMENT,
106    workflow_id     TEXT NOT NULL REFERENCES workflows(id),
107    name            TEXT NOT NULL,
108    payload         TEXT,
109    consumed        INTEGER NOT NULL DEFAULT 0,
110    received_at     REAL NOT NULL
111);
112CREATE INDEX IF NOT EXISTS workflow.idx_wf_signals_lookup ON signals(workflow_id, name, consumed);
113
114CREATE TABLE IF NOT EXISTS workflow.schedules (
115    name            TEXT NOT NULL,
116    namespace       TEXT NOT NULL DEFAULT 'main',
117    workflow_type   TEXT NOT NULL,
118    cron_expr       TEXT NOT NULL,
119    timezone        TEXT NOT NULL DEFAULT 'UTC',
120    input           TEXT,
121    task_queue      TEXT NOT NULL DEFAULT 'main',
122    overlap_policy  TEXT NOT NULL DEFAULT 'skip',
123    paused          INTEGER NOT NULL DEFAULT 0,
124    last_run_at     REAL,
125    next_run_at     REAL,
126    last_workflow_id TEXT,
127    created_at      REAL NOT NULL,
128    PRIMARY KEY (namespace, name)
129);
130
131CREATE TABLE IF NOT EXISTS workflow.workers (
132    id              TEXT PRIMARY KEY,
133    namespace       TEXT NOT NULL DEFAULT 'main',
134    identity        TEXT NOT NULL,
135    task_queue      TEXT NOT NULL,
136    workflows       TEXT,
137    activities      TEXT,
138    max_concurrent_workflows  INTEGER NOT NULL DEFAULT 10,
139    max_concurrent_activities INTEGER NOT NULL DEFAULT 10,
140    active_tasks    INTEGER NOT NULL DEFAULT 0,
141    last_heartbeat  REAL NOT NULL,
142    registered_at   REAL NOT NULL
143);
144
145CREATE TABLE IF NOT EXISTS workflow.snapshots (
146    workflow_id     TEXT NOT NULL REFERENCES workflows(id),
147    event_seq       INTEGER NOT NULL,
148    state_json      TEXT NOT NULL,
149    created_at      REAL NOT NULL,
150    PRIMARY KEY (workflow_id, event_seq)
151);
152
153-- workflow.api_keys retired in plan-15 slice 3 (auth tokens come from
154-- the auth module).
155DROP TABLE IF EXISTS workflow.api_keys;
156
157CREATE TABLE IF NOT EXISTS engine.lock (
158    id              INTEGER PRIMARY KEY CHECK (id = 1),
159    instance_id     TEXT NOT NULL,
160    started_at      REAL NOT NULL,
161    last_heartbeat  REAL NOT NULL
162);
163
164CREATE TABLE IF NOT EXISTS engine.events (
165    id              INTEGER PRIMARY KEY AUTOINCREMENT,
166    ts              REAL NOT NULL DEFAULT (CAST(strftime('%s','now') AS REAL)),
167    namespace       TEXT NOT NULL,
168    subsystem       TEXT NOT NULL,
169    kind            TEXT NOT NULL,
170    payload         TEXT NOT NULL DEFAULT '{}'
171);
172CREATE INDEX IF NOT EXISTS engine.idx_engine_events_ns_id ON events(namespace, id);
173CREATE INDEX IF NOT EXISTS engine.idx_engine_events_ts_prune ON events(ts);
174"#;
175
176/// Stale lock timeout — if the lock holder hasn't heartbeated in this
177/// many seconds, assume it's dead and allow takeover.
178const LOCK_STALE_SECS: f64 = 60.0;
179/// How often to refresh the lock heartbeat.
180const LOCK_HEARTBEAT_SECS: u64 = 15;
181
182/// `Clone` is derived because the underlying `SqlitePool` is itself
183/// `Clone` (it's `Arc<PoolInner>` internally) — cloning the store hands
184/// back a new wrapper around the same connection pool. The
185/// `instance_id` is per-store identity (heartbeat row tag), shared
186/// across clones so all clones look like the same instance to
187/// `engine.lock`.
188#[derive(Clone)]
189pub struct SqliteStore {
190    pool: SqlitePool,
191    instance_id: String,
192}
193
194/// Build a fresh [`SqlitePool`] with `engine` + `workflow` ATTACHed to
195/// in-memory shared-cache databases (one alias per pool). Each connection
196/// in the pool inherits the same ATTACHed databases via `after_connect`.
197///
198/// This is the test-friendly path for `SqliteStore::new(url)` callers that
199/// pass `sqlite::memory:` or any path-based URL — every connection sees
200/// the same `engine.*` / `workflow.*` data because the shared-cache URI
201/// pins the in-memory DB to a process-global name.
202///
203/// Production embedders (the engine binary) build their own pool with
204/// file-backed ATTACHes (`<data_dir>/engine.db`, `<data_dir>/workflow.db`)
205/// and call [`SqliteStore::from_attached_pool`] instead.
206async fn build_default_pool(url: &str) -> Result<SqlitePool> {
207    use std::str::FromStr;
208    use std::sync::atomic::{AtomicU64, Ordering};
209
210    static SEQ: AtomicU64 = AtomicU64::new(0);
211    let suffix = format!(
212        "{}_{}",
213        std::process::id(),
214        SEQ.fetch_add(1, Ordering::Relaxed)
215    );
216    let engine_alias = format!("file:assay_engine_{suffix}?mode=memory&cache=shared");
217    let workflow_alias = format!("file:assay_workflow_{suffix}?mode=memory&cache=shared");
218
219    let opts = SqliteConnectOptions::from_str(url)?.create_if_missing(true);
220
221    let pool = SqlitePoolOptions::new()
222        .max_connections(1)
223        .after_connect(move |conn, _meta| {
224            let engine_alias = engine_alias.clone();
225            let workflow_alias = workflow_alias.clone();
226            Box::pin(async move {
227                use sqlx::Executor;
228                conn.execute(format!("ATTACH DATABASE '{engine_alias}' AS engine").as_str())
229                    .await?;
230                conn.execute(format!("ATTACH DATABASE '{workflow_alias}' AS workflow").as_str())
231                    .await?;
232                Ok(())
233            })
234        })
235        .connect_with(opts)
236        .await?;
237    Ok(pool)
238}
239
240impl SqliteStore {
241    /// Open a SqliteStore at `url`. Provisions an in-memory `engine` +
242    /// `workflow` ATTACH automatically — convenient for tests and
243    /// embedders that don't need persistent module isolation. Production
244    /// deployments use [`SqliteStore::from_attached_pool`] with the
245    /// engine-controlled pool that ATTACHes to `<data_dir>/*.db` files.
246    pub async fn new(url: &str) -> Result<Self> {
247        let pool = build_default_pool(url).await?;
248        Self::from_attached_pool(pool).await
249    }
250
251    /// Construct from an externally-managed pool that already has the
252    /// `engine` and `workflow` databases ATTACHed. The engine binary
253    /// uses this — its pool's `after_connect` hook ATTACHes the
254    /// per-module file paths from `[backend].data_dir`.
255    pub async fn from_attached_pool(pool: SqlitePool) -> Result<Self> {
256        let instance_id = format!("assay-{:016x}", {
257            use std::collections::hash_map::DefaultHasher;
258            use std::hash::{Hash, Hasher};
259            let mut h = DefaultHasher::new();
260            std::time::SystemTime::now().hash(&mut h);
261            std::process::id().hash(&mut h);
262            h.finish()
263        });
264        let store = Self { pool, instance_id };
265        store.migrate().await?;
266        Ok(store)
267    }
268
269    /// Backward-compat alias for [`SqliteStore::from_attached_pool`].
270    /// Older call sites passed a bare pool from `SqlitePool::connect()`;
271    /// after v0.1.2 the pool must already have the engine + workflow
272    /// databases attached. The implementation is identical, kept under
273    /// the legacy name so external embedders don't break on upgrade.
274    pub async fn from_pool(pool: SqlitePool) -> Result<Self> {
275        Self::from_attached_pool(pool).await
276    }
277
278    /// Expose the underlying pool (used by the engine to build an
279    /// `SqliteEngineEventBus` that shares the same connection).
280    pub fn pool(&self) -> &SqlitePool {
281        &self.pool
282    }
283
284    /// Acquire the single-instance engine lock.
285    /// Returns an error if another instance is already running.
286    pub async fn acquire_engine_lock(&self) -> Result<()> {
287        let now = timestamp_now();
288
289        // Try to insert the lock
290        let result = sqlx::query(
291            "INSERT INTO engine.lock (id, instance_id, started_at, last_heartbeat) VALUES (1, ?, ?, ?)",
292        )
293        .bind(&self.instance_id)
294        .bind(now)
295        .bind(now)
296        .execute(&self.pool)
297        .await;
298
299        match result {
300            Ok(_) => Ok(()),
301            Err(_) => {
302                // Lock exists — check if it's stale
303                let row: Option<(String, f64)> = sqlx::query_as(
304                    "SELECT instance_id, last_heartbeat FROM engine.lock WHERE id = 1",
305                )
306                .fetch_optional(&self.pool)
307                .await?;
308
309                if let Some((existing_id, last_hb)) = row {
310                    if now - last_hb > LOCK_STALE_SECS {
311                        // Stale lock — take over
312                        sqlx::query(
313                            "UPDATE engine.lock SET instance_id = ?, started_at = ?, last_heartbeat = ? WHERE id = 1",
314                        )
315                        .bind(&self.instance_id)
316                        .bind(now)
317                        .bind(now)
318                        .execute(&self.pool)
319                        .await?;
320                        tracing::warn!(
321                            "Took over stale engine lock from {existing_id} (last heartbeat {:.0}s ago)",
322                            now - last_hb
323                        );
324                        Ok(())
325                    } else {
326                        let age = now - last_hb;
327                        anyhow::bail!(
328                            "Another assay engine instance is already running (id: {existing_id}, \
329                             last heartbeat {age:.0}s ago).\n\n\
330                             SQLite only supports a single engine instance. For multi-instance \
331                             deployment (Kubernetes, Docker Swarm), use PostgreSQL:\n\n\
332                             \x20 assay serve --backend postgres://user:pass@host:5432/dbname"
333                        );
334                    }
335                } else {
336                    anyhow::bail!("Unexpected engine lock state");
337                }
338            }
339        }
340    }
341
342    /// Refresh the engine lock heartbeat. Called periodically by the engine.
343    pub async fn refresh_engine_lock(&self) -> Result<()> {
344        sqlx::query("UPDATE engine.lock SET last_heartbeat = ? WHERE id = 1 AND instance_id = ?")
345            .bind(timestamp_now())
346            .bind(&self.instance_id)
347            .execute(&self.pool)
348            .await?;
349        Ok(())
350    }
351
352    /// Release the engine lock on shutdown.
353    pub async fn release_engine_lock(&self) -> Result<()> {
354        sqlx::query("DELETE FROM engine.lock WHERE id = 1 AND instance_id = ?")
355            .bind(&self.instance_id)
356            .execute(&self.pool)
357            .await?;
358        Ok(())
359    }
360
361    /// Start background task to keep the lock alive.
362    pub fn spawn_lock_heartbeat(self: &std::sync::Arc<Self>) {
363        let store = std::sync::Arc::clone(self);
364        tokio::spawn(async move {
365            let mut tick =
366                tokio::time::interval(std::time::Duration::from_secs(LOCK_HEARTBEAT_SECS));
367            loop {
368                tick.tick().await;
369                if let Err(e) = store.refresh_engine_lock().await {
370                    tracing::error!("Engine lock heartbeat failed: {e}");
371                }
372            }
373        });
374    }
375
376    /// Apply the baseline schema. SCHEMA's `CREATE TABLE IF NOT EXISTS`
377    /// statements are the source of truth — pre-1.0 we don't carry
378    /// `ALTER TABLE ADD COLUMN` history. For additive migrations later,
379    /// chain a `Self::add_column_if_missing(&self.pool, "<table>",
380    /// "<column>", "<type_def>")` call here before returning.
381    async fn migrate(&self) -> Result<()> {
382        for statement in SCHEMA.split(';') {
383            let trimmed = statement.trim();
384            if !trimmed.is_empty() {
385                sqlx::query(trimmed).execute(&self.pool).await?;
386            }
387        }
388        // Future additive migrations go here; see doc-comment above.
389        Ok(())
390    }
391
392    /// Add a column to an existing table if it's not already there.
393    ///
394    /// SQLite (unlike Postgres) doesn't support `ADD COLUMN IF NOT EXISTS`,
395    /// so we check via `pragma_table_info` before issuing the ALTER. Each
396    /// call is idempotent across startups.
397    ///
398    /// Currently unused — kept as the documented pattern for the first
399    /// additive migration after v0.11.3. Remove `#[allow(dead_code)]` when
400    /// a caller is added.
401    #[allow(dead_code)]
402    async fn add_column_if_missing(
403        pool: &SqlitePool,
404        table: &str,
405        column: &str,
406        type_def: &str,
407    ) -> Result<()> {
408        let exists: Option<(String,)> =
409            sqlx::query_as("SELECT name FROM pragma_table_info(?) WHERE name = ?")
410                .bind(table)
411                .bind(column)
412                .fetch_optional(pool)
413                .await?;
414        if exists.is_none() {
415            let sql = format!("ALTER TABLE {table} ADD COLUMN {column} {type_def}");
416            sqlx::query(&sql).execute(pool).await?;
417        }
418        Ok(())
419    }
420}
421
422impl WorkflowStore for SqliteStore {
423    // ── Namespaces ─────────────────────────────────────────
424
425    async fn create_namespace(&self, name: &str) -> Result<()> {
426        sqlx::query("INSERT INTO workflow.namespaces (name, created_at) VALUES (?, ?)")
427            .bind(name)
428            .bind(timestamp_now())
429            .execute(&self.pool)
430            .await?;
431        Ok(())
432    }
433
434    async fn list_namespaces(&self) -> Result<Vec<NamespaceRecord>> {
435        let rows = sqlx::query_as::<_, (String, f64)>(
436            "SELECT name, created_at FROM workflow.namespaces ORDER BY name",
437        )
438        .fetch_all(&self.pool)
439        .await?;
440        Ok(rows
441            .into_iter()
442            .map(|(name, created_at)| NamespaceRecord { name, created_at })
443            .collect())
444    }
445
446    async fn delete_namespace(&self, name: &str) -> Result<bool> {
447        // Mirror PG: 'main' is always available, can't be deleted.
448        let res = sqlx::query("DELETE FROM workflow.namespaces WHERE name = ? AND name != 'main'")
449            .bind(name)
450            .execute(&self.pool)
451            .await?;
452        Ok(res.rows_affected() > 0)
453    }
454
455    async fn get_namespace_stats(&self, namespace: &str) -> Result<NamespaceStats> {
456        let total: (i64,) =
457            sqlx::query_as("SELECT COUNT(*) FROM workflow.workflows WHERE namespace = ?")
458                .bind(namespace)
459                .fetch_one(&self.pool)
460                .await?;
461        let running: (i64,) = sqlx::query_as(
462            "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = ? AND status = 'RUNNING'",
463        )
464        .bind(namespace)
465        .fetch_one(&self.pool)
466        .await?;
467        let pending: (i64,) = sqlx::query_as(
468            "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = ? AND status = 'PENDING'",
469        )
470        .bind(namespace)
471        .fetch_one(&self.pool)
472        .await?;
473        let completed: (i64,) = sqlx::query_as(
474            "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = ? AND status = 'COMPLETED'",
475        )
476        .bind(namespace)
477        .fetch_one(&self.pool)
478        .await?;
479        let failed: (i64,) = sqlx::query_as(
480            "SELECT COUNT(*) FROM workflow.workflows WHERE namespace = ? AND status = 'FAILED'",
481        )
482        .bind(namespace)
483        .fetch_one(&self.pool)
484        .await?;
485        let schedules: (i64,) =
486            sqlx::query_as("SELECT COUNT(*) FROM workflow.schedules WHERE namespace = ?")
487                .bind(namespace)
488                .fetch_one(&self.pool)
489                .await?;
490        let workers: (i64,) =
491            sqlx::query_as("SELECT COUNT(*) FROM workflow.workers WHERE namespace = ?")
492                .bind(namespace)
493                .fetch_one(&self.pool)
494                .await?;
495
496        Ok(NamespaceStats {
497            namespace: namespace.to_string(),
498            total_workflows: total.0,
499            running: running.0,
500            pending: pending.0,
501            completed: completed.0,
502            failed: failed.0,
503            schedules: schedules.0,
504            workers: workers.0,
505        })
506    }
507
508    // ── Workflows ──────────────────────────────────────────
509
510    async fn create_workflow(&self, wf: &WorkflowRecord) -> Result<()> {
511        sqlx::query(
512            "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)
513             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
514        )
515        .bind(&wf.id)
516        .bind(&wf.namespace)
517        .bind(&wf.run_id)
518        .bind(&wf.workflow_type)
519        .bind(&wf.task_queue)
520        .bind(&wf.status)
521        .bind(&wf.input)
522        .bind(&wf.result)
523        .bind(&wf.error)
524        .bind(&wf.parent_id)
525        .bind(&wf.claimed_by)
526        .bind(&wf.search_attributes)
527        .bind(wf.archived_at)
528        .bind(&wf.archive_uri)
529        .bind(wf.created_at)
530        .bind(wf.updated_at)
531        .bind(wf.completed_at)
532        .execute(&self.pool)
533        .await?;
534        Ok(())
535    }
536
537    async fn get_workflow(&self, id: &str) -> Result<Option<WorkflowRecord>> {
538        let row = sqlx::query_as::<_, SqliteWorkflowRow>(
539            "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 = ?",
540        )
541        .bind(id)
542        .fetch_optional(&self.pool)
543        .await?;
544        Ok(row.map(Into::into))
545    }
546
547    async fn list_workflows(
548        &self,
549        namespace: &str,
550        status: Option<WorkflowStatus>,
551        workflow_type: Option<&str>,
552        search_attrs_filter: Option<&str>,
553        limit: i64,
554        offset: i64,
555    ) -> Result<Vec<WorkflowRecord>> {
556        let status_str = status.map(|s| s.to_string());
557
558        // Parse search filter into (key, value) pairs. Each pair adds a
559        // `json_extract(search_attributes, '$.key') = value` predicate so
560        // matches require every filter key to be present in the stored
561        // attributes. Invalid/empty JSON → no filter (all pass).
562        let filter_pairs: Vec<(String, serde_json::Value)> = search_attrs_filter
563            .and_then(|s| serde_json::from_str::<serde_json::Value>(s).ok())
564            .and_then(|v| v.as_object().cloned())
565            .map(|m| m.into_iter().collect())
566            .unwrap_or_default();
567
568        let mut sql = String::from(
569            "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
570             FROM workflow.workflows
571             WHERE namespace = ?
572               AND (? IS NULL OR status = ?)
573               AND (? IS NULL OR workflow_type = ?)",
574        );
575        for _ in &filter_pairs {
576            sql.push_str(" AND json_extract(search_attributes, '$.' || ?) = ?");
577        }
578        sql.push_str(" ORDER BY created_at DESC LIMIT ? OFFSET ?");
579
580        let mut q = sqlx::query_as::<_, SqliteWorkflowRow>(&sql)
581            .bind(namespace)
582            .bind(&status_str)
583            .bind(&status_str)
584            .bind(workflow_type)
585            .bind(workflow_type);
586        for (key, value) in &filter_pairs {
587            q = q.bind(key.clone());
588            // Bind the JSON value as its string/number representation.
589            // json_extract on a stored JSON string returns its "natural"
590            // SQLite type (text for strings, numeric for numbers), so we
591            // match by the same type.
592            match value {
593                serde_json::Value::String(s) => q = q.bind(s.clone()),
594                serde_json::Value::Number(n) => {
595                    if let Some(i) = n.as_i64() {
596                        q = q.bind(i);
597                    } else if let Some(f) = n.as_f64() {
598                        q = q.bind(f);
599                    } else {
600                        q = q.bind(n.to_string());
601                    }
602                }
603                serde_json::Value::Bool(b) => q = q.bind(*b as i64),
604                _ => q = q.bind(value.to_string()),
605            }
606        }
607        let rows = q.bind(limit).bind(offset).fetch_all(&self.pool).await?;
608        Ok(rows.into_iter().map(Into::into).collect())
609    }
610
611    async fn update_workflow_status(
612        &self,
613        id: &str,
614        status: WorkflowStatus,
615        result: Option<&str>,
616        error: Option<&str>,
617    ) -> Result<()> {
618        let now = timestamp_now();
619        let completed_at = if status.is_terminal() {
620            Some(now)
621        } else {
622            None
623        };
624        sqlx::query(
625            "UPDATE workflow.workflows SET status = ?, result = COALESCE(?, result), error = COALESCE(?, error), updated_at = ?, completed_at = COALESCE(?, completed_at) WHERE id = ?",
626        )
627        .bind(status.to_string())
628        .bind(result)
629        .bind(error)
630        .bind(now)
631        .bind(completed_at)
632        .bind(id)
633        .execute(&self.pool)
634        .await?;
635        Ok(())
636    }
637
638    async fn claim_workflow(&self, id: &str, worker_id: &str) -> Result<bool> {
639        let res = sqlx::query(
640            "UPDATE workflow.workflows SET claimed_by = ?, status = 'RUNNING', updated_at = ? WHERE id = ? AND claimed_by IS NULL",
641        )
642        .bind(worker_id)
643        .bind(timestamp_now())
644        .bind(id)
645        .execute(&self.pool)
646        .await?;
647        Ok(res.rows_affected() > 0)
648    }
649
650    async fn mark_workflow_dispatchable(&self, workflow_id: &str) -> Result<()> {
651        sqlx::query("UPDATE workflow.workflows SET needs_dispatch = 1 WHERE id = ?")
652            .bind(workflow_id)
653            .execute(&self.pool)
654            .await?;
655        Ok(())
656    }
657
658    async fn claim_workflow_task(
659        &self,
660        task_queue: &str,
661        worker_id: &str,
662    ) -> Result<Option<WorkflowRecord>> {
663        let now = timestamp_now();
664        // Atomic: pick the oldest dispatchable + unclaimed workflow on the queue
665        let row = sqlx::query_as::<_, SqliteWorkflowRow>(
666            "UPDATE workflow.workflows
667             SET dispatch_claimed_by = ?, dispatch_last_heartbeat = ?, needs_dispatch = 0
668             WHERE id = (
669                SELECT id FROM workflow.workflows
670                WHERE task_queue = ?
671                  AND needs_dispatch = 1
672                  AND dispatch_claimed_by IS NULL
673                  AND status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')
674                ORDER BY updated_at ASC
675                LIMIT 1
676             )
677             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",
678        )
679        .bind(worker_id)
680        .bind(now)
681        .bind(task_queue)
682        .fetch_optional(&self.pool)
683        .await?;
684        Ok(row.map(Into::into))
685    }
686
687    async fn release_workflow_task(&self, workflow_id: &str, worker_id: &str) -> Result<()> {
688        sqlx::query(
689            "UPDATE workflow.workflows
690             SET dispatch_claimed_by = NULL, dispatch_last_heartbeat = NULL
691             WHERE id = ? AND dispatch_claimed_by = ?",
692        )
693        .bind(workflow_id)
694        .bind(worker_id)
695        .execute(&self.pool)
696        .await?;
697        Ok(())
698    }
699
700    async fn release_stale_dispatch_leases(&self, now: f64, timeout_secs: f64) -> Result<u64> {
701        // Re-arm needs_dispatch so the work goes back into the pool. Don't
702        // touch workflows that have reached a terminal state — those should
703        // never be re-dispatched.
704        let res = sqlx::query(
705            "UPDATE workflow.workflows
706             SET dispatch_claimed_by = NULL,
707                 dispatch_last_heartbeat = NULL,
708                 needs_dispatch = 1
709             WHERE dispatch_claimed_by IS NOT NULL
710               AND (? - dispatch_last_heartbeat) > ?
711               AND status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')",
712        )
713        .bind(now)
714        .bind(timeout_secs)
715        .execute(&self.pool)
716        .await?;
717        Ok(res.rows_affected())
718    }
719
720    // ── Events ─────────────────────────────────────────────
721
722    async fn append_event(&self, ev: &WorkflowEvent) -> Result<i64> {
723        let res = sqlx::query(
724            "INSERT INTO workflow.events (workflow_id, seq, event_type, payload, timestamp) VALUES (?, ?, ?, ?, ?)",
725        )
726        .bind(&ev.workflow_id)
727        .bind(ev.seq)
728        .bind(&ev.event_type)
729        .bind(&ev.payload)
730        .bind(ev.timestamp)
731        .execute(&self.pool)
732        .await?;
733        Ok(res.last_insert_rowid())
734    }
735
736    async fn list_events(&self, workflow_id: &str) -> Result<Vec<WorkflowEvent>> {
737        let rows = sqlx::query_as::<_, SqliteEventRow>(
738            "SELECT id, workflow_id, seq, event_type, payload, timestamp FROM workflow.events WHERE workflow_id = ? ORDER BY seq ASC",
739        )
740        .bind(workflow_id)
741        .fetch_all(&self.pool)
742        .await?;
743        Ok(rows.into_iter().map(Into::into).collect())
744    }
745
746    async fn list_events_page(
747        &self,
748        workflow_id: &str,
749        cursor: Option<i32>,
750        limit: i64,
751        descending: bool,
752    ) -> Result<Vec<WorkflowEvent>> {
753        let limit = limit.clamp(0, 1_000);
754        if limit == 0 {
755            return Ok(Vec::new());
756        }
757        let rows = if descending {
758            sqlx::query_as::<_, SqliteEventRow>(
759                "SELECT id, workflow_id, seq, event_type, payload, timestamp
760                 FROM workflow.events
761                 WHERE workflow_id = ? AND (? IS NULL OR seq < ?)
762                 ORDER BY seq DESC LIMIT ?",
763            )
764            .bind(workflow_id)
765            .bind(cursor)
766            .bind(cursor)
767            .bind(limit)
768            .fetch_all(&self.pool)
769            .await?
770        } else {
771            sqlx::query_as::<_, SqliteEventRow>(
772                "SELECT id, workflow_id, seq, event_type, payload, timestamp
773                 FROM workflow.events
774                 WHERE workflow_id = ? AND (? IS NULL OR seq > ?)
775                 ORDER BY seq ASC LIMIT ?",
776            )
777            .bind(workflow_id)
778            .bind(cursor)
779            .bind(cursor)
780            .bind(limit)
781            .fetch_all(&self.pool)
782            .await?
783        };
784        Ok(rows.into_iter().map(Into::into).collect())
785    }
786
787    async fn get_event_count(&self, workflow_id: &str) -> Result<i64> {
788        let row: (i64,) =
789            sqlx::query_as("SELECT COUNT(*) FROM workflow.events WHERE workflow_id = ?")
790                .bind(workflow_id)
791                .fetch_one(&self.pool)
792                .await?;
793        Ok(row.0)
794    }
795
796    // ── Activities ──────────────────────────────────────────
797
798    async fn create_activity(&self, act: &WorkflowActivity) -> Result<i64> {
799        let res = sqlx::query(
800            "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)
801             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
802        )
803        .bind(&act.workflow_id)
804        .bind(act.seq)
805        .bind(&act.name)
806        .bind(&act.task_queue)
807        .bind(&act.input)
808        .bind(&act.status)
809        .bind(act.attempt)
810        .bind(act.max_attempts)
811        .bind(act.initial_interval_secs)
812        .bind(act.backoff_coefficient)
813        .bind(act.start_to_close_secs)
814        .bind(act.heartbeat_timeout_secs)
815        .bind(act.scheduled_at)
816        .execute(&self.pool)
817        .await?;
818        Ok(res.last_insert_rowid())
819    }
820
821    async fn get_activity(&self, id: i64) -> Result<Option<WorkflowActivity>> {
822        let row = sqlx::query_as::<_, SqliteActivityRow>(
823            "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
824             FROM workflow.activities WHERE id = ?",
825        )
826        .bind(id)
827        .fetch_optional(&self.pool)
828        .await?;
829        Ok(row.map(Into::into))
830    }
831
832    async fn get_activity_by_workflow_seq(
833        &self,
834        workflow_id: &str,
835        seq: i32,
836    ) -> Result<Option<WorkflowActivity>> {
837        let row = sqlx::query_as::<_, SqliteActivityRow>(
838            "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
839             FROM workflow.activities WHERE workflow_id = ? AND seq = ?",
840        )
841        .bind(workflow_id)
842        .bind(seq)
843        .fetch_optional(&self.pool)
844        .await?;
845        Ok(row.map(Into::into))
846    }
847
848    async fn claim_activity(
849        &self,
850        task_queue: &str,
851        worker_id: &str,
852    ) -> Result<Option<WorkflowActivity>> {
853        let now = timestamp_now();
854        let row = sqlx::query_as::<_, SqliteActivityRow>(
855            "UPDATE workflow.activities SET status = 'RUNNING', claimed_by = ?, started_at = ?
856             WHERE id = (
857                SELECT id FROM workflow.activities
858                WHERE task_queue = ? AND status = 'PENDING'
859                ORDER BY scheduled_at ASC
860                LIMIT 1
861             )
862             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",
863        )
864        .bind(worker_id)
865        .bind(now)
866        .bind(task_queue)
867        .fetch_optional(&self.pool)
868        .await?;
869        Ok(row.map(Into::into))
870    }
871
872    async fn requeue_activity_for_retry(
873        &self,
874        id: i64,
875        next_attempt: i32,
876        next_scheduled_at: f64,
877    ) -> Result<()> {
878        sqlx::query(
879            "UPDATE workflow.activities
880             SET status = 'PENDING', attempt = ?, scheduled_at = ?,
881                 claimed_by = NULL, started_at = NULL, last_heartbeat = NULL,
882                 error = NULL
883             WHERE id = ?",
884        )
885        .bind(next_attempt)
886        .bind(next_scheduled_at)
887        .bind(id)
888        .execute(&self.pool)
889        .await?;
890        Ok(())
891    }
892
893    async fn complete_activity(
894        &self,
895        id: i64,
896        result: Option<&str>,
897        error: Option<&str>,
898        failed: bool,
899    ) -> Result<()> {
900        let status = if failed { "FAILED" } else { "COMPLETED" };
901        sqlx::query(
902            "UPDATE workflow.activities SET status = ?, result = ?, error = ?, completed_at = ? WHERE id = ?",
903        )
904        .bind(status)
905        .bind(result)
906        .bind(error)
907        .bind(timestamp_now())
908        .bind(id)
909        .execute(&self.pool)
910        .await?;
911        Ok(())
912    }
913
914    async fn heartbeat_activity(&self, id: i64, _details: Option<&str>) -> Result<()> {
915        sqlx::query("UPDATE workflow.activities SET last_heartbeat = ? WHERE id = ?")
916            .bind(timestamp_now())
917            .bind(id)
918            .execute(&self.pool)
919            .await?;
920        Ok(())
921    }
922
923    async fn get_timed_out_activities(&self, now: f64) -> Result<Vec<WorkflowActivity>> {
924        let rows = sqlx::query_as::<_, SqliteActivityRow>(
925            "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
926             FROM workflow.activities
927             WHERE status = 'RUNNING'
928               AND heartbeat_timeout_secs IS NOT NULL
929               AND (? - COALESCE(last_heartbeat, started_at)) > heartbeat_timeout_secs",
930        )
931        .bind(now)
932        .fetch_all(&self.pool)
933        .await?;
934        Ok(rows.into_iter().map(Into::into).collect())
935    }
936
937    // ── Timers ──────────────────────────────────────────────
938
939    async fn create_timer(&self, timer: &WorkflowTimer) -> Result<i64> {
940        // Idempotent: INSERT OR IGNORE on UNIQUE (workflow_id, seq).
941        // If the row already existed, last_insert_rowid() is 0 — fall back to SELECT.
942        let res = sqlx::query(
943            "INSERT OR IGNORE INTO workflow.timers (workflow_id, seq, fire_at, fired) VALUES (?, ?, ?, 0)",
944        )
945        .bind(&timer.workflow_id)
946        .bind(timer.seq)
947        .bind(timer.fire_at)
948        .execute(&self.pool)
949        .await?;
950
951        let id = res.last_insert_rowid();
952        if id != 0 {
953            return Ok(id);
954        }
955
956        // Row already existed — return its id.
957        let (existing_id,): (i64,) =
958            sqlx::query_as("SELECT id FROM workflow.timers WHERE workflow_id = ? AND seq = ?")
959                .bind(&timer.workflow_id)
960                .bind(timer.seq)
961                .fetch_one(&self.pool)
962                .await?;
963        Ok(existing_id)
964    }
965
966    async fn cancel_pending_activities(&self, workflow_id: &str) -> Result<u64> {
967        let res = sqlx::query(
968            "UPDATE workflow.activities SET status = 'CANCELLED', completed_at = ?
969             WHERE workflow_id = ? AND status = 'PENDING'",
970        )
971        .bind(timestamp_now())
972        .bind(workflow_id)
973        .execute(&self.pool)
974        .await?;
975        Ok(res.rows_affected())
976    }
977
978    async fn cancel_pending_timers(&self, workflow_id: &str) -> Result<u64> {
979        let res = sqlx::query(
980            "UPDATE workflow.timers SET fired = 1
981             WHERE workflow_id = ? AND fired = 0",
982        )
983        .bind(workflow_id)
984        .execute(&self.pool)
985        .await?;
986        Ok(res.rows_affected())
987    }
988
989    async fn get_timer_by_workflow_seq(
990        &self,
991        workflow_id: &str,
992        seq: i32,
993    ) -> Result<Option<WorkflowTimer>> {
994        let row = sqlx::query_as::<_, SqliteTimerRow>(
995            "SELECT id, workflow_id, seq, fire_at, fired
996             FROM workflow.timers WHERE workflow_id = ? AND seq = ?",
997        )
998        .bind(workflow_id)
999        .bind(seq)
1000        .fetch_optional(&self.pool)
1001        .await?;
1002        Ok(row.map(Into::into))
1003    }
1004
1005    async fn fire_due_timers(&self, now: f64) -> Result<Vec<WorkflowTimer>> {
1006        let rows = sqlx::query_as::<_, SqliteTimerRow>(
1007            "UPDATE workflow.timers SET fired = 1
1008             WHERE fired = 0 AND fire_at <= ?
1009             RETURNING id, workflow_id, seq, fire_at, fired",
1010        )
1011        .bind(now)
1012        .fetch_all(&self.pool)
1013        .await?;
1014        Ok(rows.into_iter().map(Into::into).collect())
1015    }
1016
1017    // ── Signals ─────────────────────────────────────────────
1018
1019    async fn send_signal(&self, sig: &WorkflowSignal) -> Result<i64> {
1020        let res = sqlx::query(
1021            "INSERT INTO workflow.signals (workflow_id, name, payload, consumed, received_at) VALUES (?, ?, ?, 0, ?)",
1022        )
1023        .bind(&sig.workflow_id)
1024        .bind(&sig.name)
1025        .bind(&sig.payload)
1026        .bind(sig.received_at)
1027        .execute(&self.pool)
1028        .await?;
1029        Ok(res.last_insert_rowid())
1030    }
1031
1032    async fn consume_signals(&self, workflow_id: &str, name: &str) -> Result<Vec<WorkflowSignal>> {
1033        let rows = sqlx::query_as::<_, SqliteSignalRow>(
1034            "UPDATE workflow.signals SET consumed = 1
1035             WHERE workflow_id = ? AND name = ? AND consumed = 0
1036             RETURNING id, workflow_id, name, payload, consumed, received_at",
1037        )
1038        .bind(workflow_id)
1039        .bind(name)
1040        .fetch_all(&self.pool)
1041        .await?;
1042        Ok(rows.into_iter().map(Into::into).collect())
1043    }
1044
1045    // ── Schedules ───────────────────────────────────────────
1046
1047    async fn create_schedule(&self, sched: &WorkflowSchedule) -> Result<()> {
1048        sqlx::query(
1049            "INSERT INTO workflow.schedules (name, namespace, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at)
1050             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1051        )
1052        .bind(&sched.name)
1053        .bind(&sched.namespace)
1054        .bind(&sched.workflow_type)
1055        .bind(&sched.cron_expr)
1056        .bind(&sched.timezone)
1057        .bind(&sched.input)
1058        .bind(&sched.task_queue)
1059        .bind(&sched.overlap_policy)
1060        .bind(sched.paused)
1061        .bind(sched.last_run_at)
1062        .bind(sched.next_run_at)
1063        .bind(&sched.last_workflow_id)
1064        .bind(sched.created_at)
1065        .execute(&self.pool)
1066        .await?;
1067        Ok(())
1068    }
1069
1070    async fn get_schedule(&self, namespace: &str, name: &str) -> Result<Option<WorkflowSchedule>> {
1071        let row = sqlx::query_as::<_, SqliteScheduleRow>(
1072            "SELECT name, namespace, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at
1073             FROM workflow.schedules WHERE namespace = ? AND name = ?",
1074        )
1075        .bind(namespace)
1076        .bind(name)
1077        .fetch_optional(&self.pool)
1078        .await?;
1079        Ok(row.map(Into::into))
1080    }
1081
1082    async fn list_schedules(&self, namespace: &str) -> Result<Vec<WorkflowSchedule>> {
1083        let rows = sqlx::query_as::<_, SqliteScheduleRow>(
1084            "SELECT name, namespace, workflow_type, cron_expr, timezone, input, task_queue, overlap_policy, paused, last_run_at, next_run_at, last_workflow_id, created_at
1085             FROM workflow.schedules WHERE namespace = ? ORDER BY name",
1086        )
1087        .bind(namespace)
1088        .fetch_all(&self.pool)
1089        .await?;
1090        Ok(rows.into_iter().map(Into::into).collect())
1091    }
1092
1093    async fn update_schedule_last_run(
1094        &self,
1095        namespace: &str,
1096        name: &str,
1097        last_run_at: f64,
1098        next_run_at: f64,
1099        workflow_id: &str,
1100    ) -> Result<()> {
1101        sqlx::query(
1102            "UPDATE workflow.schedules SET last_run_at = ?, next_run_at = ?, last_workflow_id = ? WHERE namespace = ? AND name = ?",
1103        )
1104        .bind(last_run_at)
1105        .bind(next_run_at)
1106        .bind(workflow_id)
1107        .bind(namespace)
1108        .bind(name)
1109        .execute(&self.pool)
1110        .await?;
1111        Ok(())
1112    }
1113
1114    async fn delete_schedule(&self, namespace: &str, name: &str) -> Result<bool> {
1115        let res = sqlx::query("DELETE FROM workflow.schedules WHERE namespace = ? AND name = ?")
1116            .bind(namespace)
1117            .bind(name)
1118            .execute(&self.pool)
1119            .await?;
1120        Ok(res.rows_affected() > 0)
1121    }
1122
1123    async fn list_archivable_workflows(
1124        &self,
1125        cutoff: f64,
1126        limit: i64,
1127    ) -> Result<Vec<WorkflowRecord>> {
1128        let rows = sqlx::query_as::<_, SqliteWorkflowRow>(
1129            "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
1130             FROM workflow.workflows
1131             WHERE status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'TIMED_OUT')
1132               AND completed_at IS NOT NULL
1133               AND completed_at < ?
1134               AND archived_at IS NULL
1135             ORDER BY completed_at ASC
1136             LIMIT ?",
1137        )
1138        .bind(cutoff)
1139        .bind(limit)
1140        .fetch_all(&self.pool)
1141        .await?;
1142        Ok(rows.into_iter().map(Into::into).collect())
1143    }
1144
1145    async fn mark_archived_and_purge(
1146        &self,
1147        workflow_id: &str,
1148        archive_uri: &str,
1149        archived_at: f64,
1150    ) -> Result<()> {
1151        let mut tx = self.pool.begin().await?;
1152        sqlx::query("DELETE FROM workflow.events WHERE workflow_id = ?")
1153            .bind(workflow_id)
1154            .execute(&mut *tx)
1155            .await?;
1156        sqlx::query("DELETE FROM workflow.activities WHERE workflow_id = ?")
1157            .bind(workflow_id)
1158            .execute(&mut *tx)
1159            .await?;
1160        sqlx::query("DELETE FROM workflow.timers WHERE workflow_id = ?")
1161            .bind(workflow_id)
1162            .execute(&mut *tx)
1163            .await?;
1164        sqlx::query("DELETE FROM workflow.signals WHERE workflow_id = ?")
1165            .bind(workflow_id)
1166            .execute(&mut *tx)
1167            .await?;
1168        sqlx::query("DELETE FROM workflow.snapshots WHERE workflow_id = ?")
1169            .bind(workflow_id)
1170            .execute(&mut *tx)
1171            .await?;
1172        sqlx::query("UPDATE workflow.workflows SET archived_at = ?, archive_uri = ? WHERE id = ?")
1173            .bind(archived_at)
1174            .bind(archive_uri)
1175            .bind(workflow_id)
1176            .execute(&mut *tx)
1177            .await?;
1178        tx.commit().await?;
1179        Ok(())
1180    }
1181
1182    async fn upsert_search_attributes(&self, workflow_id: &str, patch_json: &str) -> Result<()> {
1183        // Merge at the application layer so we don't depend on SQLite's
1184        // `json_patch`, which is only available with the json1 extension.
1185        let current: Option<(Option<String>,)> =
1186            sqlx::query_as("SELECT search_attributes FROM workflow.workflows WHERE id = ?")
1187                .bind(workflow_id)
1188                .fetch_optional(&self.pool)
1189                .await?;
1190        let merged = merge_search_attrs(current.and_then(|(s,)| s).as_deref(), patch_json)?;
1191        sqlx::query("UPDATE workflow.workflows SET search_attributes = ? WHERE id = ?")
1192            .bind(merged)
1193            .bind(workflow_id)
1194            .execute(&self.pool)
1195            .await?;
1196        Ok(())
1197    }
1198
1199    async fn update_schedule(
1200        &self,
1201        namespace: &str,
1202        name: &str,
1203        patch: &SchedulePatch,
1204    ) -> Result<Option<WorkflowSchedule>> {
1205        // Build the UPDATE dynamically so unchanged fields aren't touched
1206        // and NULL from `serde_json::Value::Null` round-trips cleanly.
1207        let mut sets: Vec<&'static str> = Vec::new();
1208        if patch.cron_expr.is_some() {
1209            sets.push("cron_expr = ?");
1210        }
1211        if patch.timezone.is_some() {
1212            sets.push("timezone = ?");
1213        }
1214        if patch.input.is_some() {
1215            sets.push("input = ?");
1216        }
1217        if patch.task_queue.is_some() {
1218            sets.push("task_queue = ?");
1219        }
1220        if patch.overlap_policy.is_some() {
1221            sets.push("overlap_policy = ?");
1222        }
1223        // Updating last_run_at/next_run_at is internal only (update_schedule_last_run).
1224        if sets.is_empty() {
1225            return self.get_schedule(namespace, name).await;
1226        }
1227
1228        let sql = format!(
1229            "UPDATE workflow.schedules SET {} WHERE namespace = ? AND name = ?",
1230            sets.join(", ")
1231        );
1232        let mut q = sqlx::query(&sql);
1233        if let Some(ref v) = patch.cron_expr {
1234            q = q.bind(v);
1235        }
1236        if let Some(ref v) = patch.timezone {
1237            q = q.bind(v);
1238        }
1239        if let Some(ref v) = patch.input {
1240            q = q.bind(v.to_string());
1241        }
1242        if let Some(ref v) = patch.task_queue {
1243            q = q.bind(v);
1244        }
1245        if let Some(ref v) = patch.overlap_policy {
1246            q = q.bind(v);
1247        }
1248        let res = q.bind(namespace).bind(name).execute(&self.pool).await?;
1249        if res.rows_affected() == 0 {
1250            return Ok(None);
1251        }
1252        self.get_schedule(namespace, name).await
1253    }
1254
1255    async fn set_schedule_paused(
1256        &self,
1257        namespace: &str,
1258        name: &str,
1259        paused: bool,
1260    ) -> Result<Option<WorkflowSchedule>> {
1261        let res = sqlx::query(
1262            "UPDATE workflow.schedules SET paused = ? WHERE namespace = ? AND name = ?",
1263        )
1264        .bind(paused)
1265        .bind(namespace)
1266        .bind(name)
1267        .execute(&self.pool)
1268        .await?;
1269        if res.rows_affected() == 0 {
1270            return Ok(None);
1271        }
1272        self.get_schedule(namespace, name).await
1273    }
1274
1275    // ── Workers ─────────────────────────────────────────────
1276
1277    async fn register_worker(&self, w: &WorkflowWorker) -> Result<()> {
1278        sqlx::query(
1279            "INSERT OR REPLACE INTO workflow.workers (id, namespace, identity, task_queue, workflows, activities, max_concurrent_workflows, max_concurrent_activities, active_tasks, last_heartbeat, registered_at)
1280             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1281        )
1282        .bind(&w.id)
1283        .bind(&w.namespace)
1284        .bind(&w.identity)
1285        .bind(&w.task_queue)
1286        .bind(&w.workflows)
1287        .bind(&w.activities)
1288        .bind(w.max_concurrent_workflows)
1289        .bind(w.max_concurrent_activities)
1290        .bind(w.active_tasks)
1291        .bind(w.last_heartbeat)
1292        .bind(w.registered_at)
1293        .execute(&self.pool)
1294        .await?;
1295        Ok(())
1296    }
1297
1298    async fn heartbeat_worker(&self, id: &str, now: f64) -> Result<()> {
1299        sqlx::query("UPDATE workflow.workers SET last_heartbeat = ? WHERE id = ?")
1300            .bind(now)
1301            .bind(id)
1302            .execute(&self.pool)
1303            .await?;
1304        Ok(())
1305    }
1306
1307    async fn list_workers(&self, namespace: &str) -> Result<Vec<WorkflowWorker>> {
1308        let rows = sqlx::query_as::<_, SqliteWorkerRow>(
1309            "SELECT id, namespace, identity, task_queue, workflows, activities, max_concurrent_workflows, max_concurrent_activities, active_tasks, last_heartbeat, registered_at
1310             FROM workflow.workers WHERE namespace = ? ORDER BY registered_at",
1311        )
1312        .bind(namespace)
1313        .fetch_all(&self.pool)
1314        .await?;
1315        Ok(rows.into_iter().map(Into::into).collect())
1316    }
1317
1318    async fn remove_dead_workers(&self, cutoff: f64) -> Result<Vec<String>> {
1319        let rows: Vec<(String,)> =
1320            sqlx::query_as("SELECT id FROM workflow.workers WHERE last_heartbeat < ?")
1321                .bind(cutoff)
1322                .fetch_all(&self.pool)
1323                .await?;
1324        let ids: Vec<String> = rows.into_iter().map(|r| r.0).collect();
1325        if !ids.is_empty() {
1326            sqlx::query("DELETE FROM workflow.workers WHERE last_heartbeat < ?")
1327                .bind(cutoff)
1328                .execute(&self.pool)
1329                .await?;
1330        }
1331        Ok(ids)
1332    }
1333
1334    // ── Child Workflows ─────────────────────────────────────
1335
1336    async fn list_child_workflows(&self, parent_id: &str) -> Result<Vec<WorkflowRecord>> {
1337        let rows = sqlx::query_as::<_, SqliteWorkflowRow>(
1338            "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
1339             FROM workflow.workflows WHERE parent_id = ? ORDER BY created_at ASC",
1340        )
1341        .bind(parent_id)
1342        .fetch_all(&self.pool)
1343        .await?;
1344        Ok(rows.into_iter().map(Into::into).collect())
1345    }
1346
1347    // ── Snapshots ───────────────────────────────────────────
1348
1349    async fn create_snapshot(
1350        &self,
1351        workflow_id: &str,
1352        event_seq: i32,
1353        state_json: &str,
1354    ) -> Result<()> {
1355        sqlx::query(
1356            "INSERT OR REPLACE INTO workflow.snapshots (workflow_id, event_seq, state_json, created_at)
1357             VALUES (?, ?, ?, ?)",
1358        )
1359        .bind(workflow_id)
1360        .bind(event_seq)
1361        .bind(state_json)
1362        .bind(timestamp_now())
1363        .execute(&self.pool)
1364        .await?;
1365        Ok(())
1366    }
1367
1368    async fn get_latest_snapshot(&self, workflow_id: &str) -> Result<Option<WorkflowSnapshot>> {
1369        let row = sqlx::query_as::<_, (String, i32, String, f64)>(
1370            "SELECT workflow_id, event_seq, state_json, created_at
1371             FROM workflow.snapshots WHERE workflow_id = ?
1372             ORDER BY event_seq DESC LIMIT 1",
1373        )
1374        .bind(workflow_id)
1375        .fetch_optional(&self.pool)
1376        .await?;
1377
1378        Ok(row.map(
1379            |(workflow_id, event_seq, state_json, created_at)| WorkflowSnapshot {
1380                workflow_id,
1381                event_seq,
1382                state_json,
1383                created_at,
1384            },
1385        ))
1386    }
1387
1388    // ── Queue Stats ─────────────────────────────────────────
1389
1390    async fn get_queue_stats(&self, namespace: &str) -> Result<Vec<QueueStats>> {
1391        // Gather activity stats per queue for workflows in this namespace
1392        let rows = sqlx::query_as::<_, (String, i64, i64)>(
1393            "SELECT a.task_queue,
1394                    SUM(CASE WHEN a.status = 'PENDING' THEN 1 ELSE 0 END),
1395                    SUM(CASE WHEN a.status = 'RUNNING' THEN 1 ELSE 0 END)
1396             FROM workflow.activities a
1397             INNER JOIN workflow.workflows w ON w.id = a.workflow_id
1398             WHERE w.namespace = ?
1399             GROUP BY a.task_queue",
1400        )
1401        .bind(namespace)
1402        .fetch_all(&self.pool)
1403        .await?;
1404
1405        let mut stats: Vec<QueueStats> = rows
1406            .into_iter()
1407            .map(|(queue, pending, running)| QueueStats {
1408                queue,
1409                pending_activities: pending,
1410                running_activities: running,
1411                workers: 0,
1412            })
1413            .collect();
1414
1415        // Gather worker counts per queue in this namespace
1416        let worker_rows = sqlx::query_as::<_, (String, i64)>(
1417            "SELECT task_queue, COUNT(*) FROM workflow.workers WHERE namespace = ? GROUP BY task_queue",
1418        )
1419        .bind(namespace)
1420        .fetch_all(&self.pool)
1421        .await?;
1422
1423        for (queue, count) in worker_rows {
1424            if let Some(s) = stats.iter_mut().find(|s| s.queue == queue) {
1425                s.workers = count;
1426            } else {
1427                stats.push(QueueStats {
1428                    queue,
1429                    pending_activities: 0,
1430                    running_activities: 0,
1431                    workers: count,
1432                });
1433            }
1434        }
1435
1436        stats.sort_by(|a, b| a.queue.cmp(&b.queue));
1437        Ok(stats)
1438    }
1439
1440    // ── Leader Election ─────────────────────────────────────
1441
1442    async fn try_acquire_scheduler_lock(&self) -> Result<bool> {
1443        // SQLite is single-instance — always the leader.
1444        // Also refresh the engine lock heartbeat on each scheduler tick.
1445        self.refresh_engine_lock().await.ok();
1446        Ok(true)
1447    }
1448}
1449
1450fn timestamp_now() -> f64 {
1451    std::time::SystemTime::now()
1452        .duration_since(std::time::UNIX_EPOCH)
1453        .unwrap()
1454        .as_secs_f64()
1455}
1456
1457/// Merge a JSON-object patch into a (possibly-null) current JSON object,
1458/// returning the serialised result. Shared by SQLite and Postgres stores.
1459pub(crate) fn merge_search_attrs(current: Option<&str>, patch_json: &str) -> Result<String> {
1460    let mut current_map: serde_json::Map<String, serde_json::Value> = current
1461        .and_then(|s| serde_json::from_str::<serde_json::Value>(s).ok())
1462        .and_then(|v| v.as_object().cloned())
1463        .unwrap_or_default();
1464    let patch: serde_json::Value = serde_json::from_str(patch_json)
1465        .map_err(|e| anyhow::anyhow!("invalid search_attributes patch: {e}"))?;
1466    let patch_obj = patch
1467        .as_object()
1468        .ok_or_else(|| anyhow::anyhow!("search_attributes patch must be a JSON object"))?;
1469    for (k, v) in patch_obj {
1470        current_map.insert(k.clone(), v.clone());
1471    }
1472    Ok(serde_json::Value::Object(current_map).to_string())
1473}
1474
1475// ── SQLite row types (sqlx::FromRow) ────────────────────────
1476
1477#[derive(sqlx::FromRow)]
1478struct SqliteWorkflowRow {
1479    id: String,
1480    namespace: String,
1481    run_id: String,
1482    workflow_type: String,
1483    task_queue: String,
1484    status: String,
1485    input: Option<String>,
1486    result: Option<String>,
1487    error: Option<String>,
1488    parent_id: Option<String>,
1489    claimed_by: Option<String>,
1490    search_attributes: Option<String>,
1491    archived_at: Option<f64>,
1492    archive_uri: Option<String>,
1493    created_at: f64,
1494    updated_at: f64,
1495    completed_at: Option<f64>,
1496}
1497
1498impl From<SqliteWorkflowRow> for WorkflowRecord {
1499    fn from(r: SqliteWorkflowRow) -> Self {
1500        Self {
1501            id: r.id,
1502            namespace: r.namespace,
1503            run_id: r.run_id,
1504            workflow_type: r.workflow_type,
1505            task_queue: r.task_queue,
1506            status: r.status,
1507            input: r.input,
1508            result: r.result,
1509            error: r.error,
1510            parent_id: r.parent_id,
1511            claimed_by: r.claimed_by,
1512            search_attributes: r.search_attributes,
1513            archived_at: r.archived_at,
1514            archive_uri: r.archive_uri,
1515            created_at: r.created_at,
1516            updated_at: r.updated_at,
1517            completed_at: r.completed_at,
1518        }
1519    }
1520}
1521
1522#[derive(sqlx::FromRow)]
1523struct SqliteEventRow {
1524    id: i64,
1525    workflow_id: String,
1526    seq: i32,
1527    event_type: String,
1528    payload: Option<String>,
1529    timestamp: f64,
1530}
1531
1532impl From<SqliteEventRow> for WorkflowEvent {
1533    fn from(r: SqliteEventRow) -> Self {
1534        Self {
1535            id: Some(r.id),
1536            workflow_id: r.workflow_id,
1537            seq: r.seq,
1538            event_type: r.event_type,
1539            payload: r.payload,
1540            timestamp: r.timestamp,
1541        }
1542    }
1543}
1544
1545#[derive(sqlx::FromRow)]
1546struct SqliteActivityRow {
1547    id: i64,
1548    workflow_id: String,
1549    seq: i32,
1550    name: String,
1551    task_queue: String,
1552    input: Option<String>,
1553    status: String,
1554    result: Option<String>,
1555    error: Option<String>,
1556    attempt: i32,
1557    max_attempts: i32,
1558    initial_interval_secs: f64,
1559    backoff_coefficient: f64,
1560    start_to_close_secs: f64,
1561    heartbeat_timeout_secs: Option<f64>,
1562    claimed_by: Option<String>,
1563    scheduled_at: f64,
1564    started_at: Option<f64>,
1565    completed_at: Option<f64>,
1566    last_heartbeat: Option<f64>,
1567}
1568
1569impl From<SqliteActivityRow> for WorkflowActivity {
1570    fn from(r: SqliteActivityRow) -> Self {
1571        Self {
1572            id: Some(r.id),
1573            workflow_id: r.workflow_id,
1574            seq: r.seq,
1575            name: r.name,
1576            task_queue: r.task_queue,
1577            input: r.input,
1578            status: r.status,
1579            result: r.result,
1580            error: r.error,
1581            attempt: r.attempt,
1582            max_attempts: r.max_attempts,
1583            initial_interval_secs: r.initial_interval_secs,
1584            backoff_coefficient: r.backoff_coefficient,
1585            start_to_close_secs: r.start_to_close_secs,
1586            heartbeat_timeout_secs: r.heartbeat_timeout_secs,
1587            claimed_by: r.claimed_by,
1588            scheduled_at: r.scheduled_at,
1589            started_at: r.started_at,
1590            completed_at: r.completed_at,
1591            last_heartbeat: r.last_heartbeat,
1592        }
1593    }
1594}
1595
1596#[derive(sqlx::FromRow)]
1597struct SqliteTimerRow {
1598    id: i64,
1599    workflow_id: String,
1600    seq: i32,
1601    fire_at: f64,
1602    fired: bool,
1603}
1604
1605impl From<SqliteTimerRow> for WorkflowTimer {
1606    fn from(r: SqliteTimerRow) -> Self {
1607        Self {
1608            id: Some(r.id),
1609            workflow_id: r.workflow_id,
1610            seq: r.seq,
1611            fire_at: r.fire_at,
1612            fired: r.fired,
1613        }
1614    }
1615}
1616
1617#[derive(sqlx::FromRow)]
1618struct SqliteSignalRow {
1619    id: i64,
1620    workflow_id: String,
1621    name: String,
1622    payload: Option<String>,
1623    consumed: bool,
1624    received_at: f64,
1625}
1626
1627impl From<SqliteSignalRow> for WorkflowSignal {
1628    fn from(r: SqliteSignalRow) -> Self {
1629        Self {
1630            id: Some(r.id),
1631            workflow_id: r.workflow_id,
1632            name: r.name,
1633            payload: r.payload,
1634            consumed: r.consumed,
1635            received_at: r.received_at,
1636        }
1637    }
1638}
1639
1640#[derive(sqlx::FromRow)]
1641struct SqliteScheduleRow {
1642    name: String,
1643    namespace: String,
1644    workflow_type: String,
1645    cron_expr: String,
1646    timezone: String,
1647    input: Option<String>,
1648    task_queue: String,
1649    overlap_policy: String,
1650    paused: bool,
1651    last_run_at: Option<f64>,
1652    next_run_at: Option<f64>,
1653    last_workflow_id: Option<String>,
1654    created_at: f64,
1655}
1656
1657impl From<SqliteScheduleRow> for WorkflowSchedule {
1658    fn from(r: SqliteScheduleRow) -> Self {
1659        Self {
1660            name: r.name,
1661            namespace: r.namespace,
1662            workflow_type: r.workflow_type,
1663            cron_expr: r.cron_expr,
1664            timezone: r.timezone,
1665            input: r.input,
1666            task_queue: r.task_queue,
1667            overlap_policy: r.overlap_policy,
1668            paused: r.paused,
1669            last_run_at: r.last_run_at,
1670            next_run_at: r.next_run_at,
1671            last_workflow_id: r.last_workflow_id,
1672            created_at: r.created_at,
1673        }
1674    }
1675}
1676
1677#[derive(sqlx::FromRow)]
1678struct SqliteWorkerRow {
1679    id: String,
1680    namespace: String,
1681    identity: String,
1682    task_queue: String,
1683    workflows: Option<String>,
1684    activities: Option<String>,
1685    max_concurrent_workflows: i32,
1686    max_concurrent_activities: i32,
1687    active_tasks: i32,
1688    last_heartbeat: f64,
1689    registered_at: f64,
1690}
1691
1692impl From<SqliteWorkerRow> for WorkflowWorker {
1693    fn from(r: SqliteWorkerRow) -> Self {
1694        Self {
1695            id: r.id,
1696            namespace: r.namespace,
1697            identity: r.identity,
1698            task_queue: r.task_queue,
1699            workflows: r.workflows,
1700            activities: r.activities,
1701            max_concurrent_workflows: r.max_concurrent_workflows,
1702            max_concurrent_activities: r.max_concurrent_activities,
1703            active_tasks: r.active_tasks,
1704            last_heartbeat: r.last_heartbeat,
1705            registered_at: r.registered_at,
1706        }
1707    }
1708}