1use anyhow::Result;
2use sqlx::SqlitePool;
3use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
4
5use crate::store::{NamespaceRecord, NamespaceStats, QueueStats, WorkflowStore};
6use crate::types::*;
7
8const 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
176const LOCK_STALE_SECS: f64 = 60.0;
179const LOCK_HEARTBEAT_SECS: u64 = 15;
181
182#[derive(Clone)]
189pub struct SqliteStore {
190 pool: SqlitePool,
191 instance_id: String,
192}
193
194async 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 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 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 pub async fn from_pool(pool: SqlitePool) -> Result<Self> {
275 Self::from_attached_pool(pool).await
276 }
277
278 pub fn pool(&self) -> &SqlitePool {
281 &self.pool
282 }
283
284 pub async fn acquire_engine_lock(&self) -> Result<()> {
287 let now = timestamp_now();
288
289 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 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 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 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 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 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 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 Ok(())
390 }
391
392 #[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 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 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 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 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 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 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 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 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 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 async fn create_timer(&self, timer: &WorkflowTimer) -> Result<i64> {
940 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 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 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 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 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 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 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 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 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 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 async fn get_queue_stats(&self, namespace: &str) -> Result<Vec<QueueStats>> {
1391 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 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 async fn try_acquire_scheduler_lock(&self) -> Result<bool> {
1443 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
1457pub(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#[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}