Skip to main content

ironflow_store/
store.rs

1//! The [`RunStore`] trait — async storage abstraction for runs and steps.
2//!
3//! Implement this trait to plug in any backing store. Built-in implementations:
4//!
5//! - [`InMemoryStore`](crate::memory::InMemoryStore) — development and testing.
6//! - `PostgresStore` — production (behind the `store-postgres` feature).
7
8use std::future::Future;
9use std::pin::Pin;
10
11use chrono::{DateTime, Utc};
12use uuid::Uuid;
13
14use crate::api_key_store::ApiKeyStore;
15use crate::artifact_store::ArtifactStore;
16use crate::audit_log_store::AuditLogStore;
17use crate::entities::{
18    LeaseRequest, NewRun, NewStep, NewStepDependency, Page, ReapedRun, Run, RunCreation, RunFilter,
19    RunStats, RunStatus, RunUpdate, Step, StepDependency, StepUpdate,
20};
21use crate::error::StoreError;
22use crate::secret_store::SecretStore;
23use crate::user_store::UserStore;
24
25/// Boxed future for [`RunStore`] methods — ensures object safety for `dyn RunStore`.
26pub type StoreFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, StoreError>> + Send + 'a>>;
27
28/// Error recorded on a run that exhausted its retries through lease expiries.
29///
30/// Set by [`RunStore::reap_expired_leases`] when a run has been recovered more
31/// than `max_retries` times.
32pub const LEASE_EXPIRED_ERROR: &str = "worker lease expired";
33
34/// Async storage abstraction for workflow runs and steps.
35///
36/// All methods return a [`StoreFuture`] (boxed future) to maintain object safety,
37/// allowing the store to be used as `Arc<dyn RunStore>`.
38///
39/// # Examples
40///
41/// ```no_run
42/// use std::collections::HashMap;
43/// use ironflow_store::prelude::*;
44/// use serde_json::json;
45/// use uuid::Uuid;
46///
47/// # async fn example() -> Result<(), ironflow_store::error::StoreError> {
48/// let store = InMemoryStore::new();
49///
50/// let run = store.create_run(NewRun {
51///     workflow_name: "deploy".to_string(),
52///     trigger: TriggerKind::Manual,
53///     payload: json!({}),
54///     max_retries: 3,
55///     handler_version: None,
56///     labels: HashMap::new(),
57///     scheduled_at: None,
58///     created_by: None,
59///     idempotency_key: None,
60///     max_cost_usd: None,
61/// }).await?.into_run();
62///
63/// let fetched = store.get_run(run.id).await?;
64/// assert!(fetched.is_some());
65/// # Ok(())
66/// # }
67/// ```
68pub trait RunStore: Send + Sync {
69    /// Create a new run in `Pending` status.
70    ///
71    /// When [`NewRun::idempotency_key`] is set and already bound to a run created
72    /// within [`IDEMPOTENCY_WINDOW`](crate::entities::IDEMPOTENCY_WINDOW), nothing is
73    /// inserted and that run is returned as [`RunCreation::Existing`]. A key bound to
74    /// an older run is released and reused for the new one.
75    ///
76    /// Concurrent calls sharing the same key resolve to a single run: exactly one
77    /// receives [`RunCreation::Created`], the others [`RunCreation::Existing`].
78    fn create_run(&self, req: NewRun) -> StoreFuture<'_, RunCreation>;
79
80    /// Look up the run bound to an idempotency key.
81    ///
82    /// Returns `None` when the key is unknown, or when the run holding it is older
83    /// than [`IDEMPOTENCY_WINDOW`](crate::entities::IDEMPOTENCY_WINDOW).
84    fn find_run_by_idempotency_key(&self, key: &str) -> StoreFuture<'_, Option<Run>>;
85
86    /// Get a run by ID. Returns `None` if not found.
87    fn get_run(&self, id: Uuid) -> StoreFuture<'_, Option<Run>>;
88
89    /// List runs matching the given filter, with pagination.
90    ///
91    /// Results are ordered by `created_at` descending (newest first).
92    fn list_runs(&self, filter: RunFilter, page: u32, per_page: u32) -> StoreFuture<'_, Page<Run>>;
93
94    /// Update a run's status with FSM validation.
95    ///
96    /// # Errors
97    ///
98    /// Returns [`StoreError::InvalidTransition`] if the transition is not allowed.
99    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
100    fn update_run_status(&self, id: Uuid, new_status: RunStatus) -> StoreFuture<'_, ()>;
101
102    /// Apply a partial update to a run.
103    ///
104    /// # Errors
105    ///
106    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
107    fn update_run(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, ()>;
108
109    /// Atomically pick the oldest pending run and transition it to `Running`.
110    ///
111    /// In PostgreSQL, this uses `SELECT FOR UPDATE SKIP LOCKED` for safe
112    /// multi-worker concurrency. The in-memory implementation uses a write lock.
113    ///
114    /// When `lease` is `Some`, the worker lease is attached in the same
115    /// transaction as the status change, so a run is never `Running` without an
116    /// owner. Pass `None` for callers that execute runs in-process and cannot
117    /// refresh a lease (inline execution, API-side resume): those runs are never
118    /// recovered by [`reap_expired_leases`](Self::reap_expired_leases).
119    ///
120    /// Returns `None` if no pending runs are available.
121    fn pick_next_pending(&self, lease: Option<LeaseRequest>) -> StoreFuture<'_, Option<Run>>;
122
123    /// Extend the worker lease on a run and return the new expiry.
124    ///
125    /// # Errors
126    ///
127    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
128    /// Returns [`StoreError::LeaseLost`] if the run is no longer `Running` or if
129    /// the lease belongs to another worker — the caller must stop executing it.
130    fn renew_lease(&self, id: Uuid, lease: LeaseRequest) -> StoreFuture<'_, DateTime<Utc>>;
131
132    /// Recover runs whose worker lease expired, at most `limit` per call.
133    ///
134    /// Each recovered run has its retry count incremented and its lease cleared,
135    /// then goes back to `Pending` — or to `Failed` with `worker lease expired`
136    /// once `max_retries` is exhausted. Runs without a lease are never touched.
137    ///
138    /// The whole batch is atomic per run (`FOR UPDATE SKIP LOCKED` in
139    /// PostgreSQL), so concurrent reapers never recover the same run twice.
140    ///
141    /// Callers are responsible for the side effects that follow a recovery:
142    /// failing orphaned steps and publishing status-change events.
143    fn reap_expired_leases(&self, limit: u32) -> StoreFuture<'_, Vec<ReapedRun>>;
144
145    /// Create a new step for a run.
146    ///
147    /// # Errors
148    ///
149    /// Returns [`StoreError::RunNotFound`] if the parent run does not exist.
150    fn create_step(&self, step: NewStep) -> StoreFuture<'_, Step>;
151
152    /// Apply a partial update to a step after execution.
153    ///
154    /// # Errors
155    ///
156    /// Returns [`StoreError::StepNotFound`] if the step does not exist.
157    fn update_step(&self, id: Uuid, update: StepUpdate) -> StoreFuture<'_, ()>;
158
159    /// Get a single step by ID. Returns `None` if not found.
160    fn get_step(&self, id: Uuid) -> StoreFuture<'_, Option<Step>>;
161
162    /// List all steps for a run, ordered by position ascending.
163    fn list_steps(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Step>>;
164
165    /// Get aggregated statistics across runs matching the filter.
166    ///
167    /// Returns counts of runs by terminal state, counts of active runs,
168    /// and totals for cost and duration. Computed efficiently by the store
169    /// implementation (single SQL query in PostgreSQL).
170    ///
171    /// Pass [`RunFilter::default()`] to get stats across all runs.
172    fn get_stats(&self, filter: RunFilter) -> StoreFuture<'_, RunStats>;
173
174    /// Create step dependency edges in batch.
175    ///
176    /// Each entry records that `step_id` depends on `depends_on`.
177    /// Duplicate edges are silently ignored.
178    ///
179    /// # Errors
180    ///
181    /// Returns [`StoreError`] if a referenced step does not exist.
182    fn create_step_dependencies(&self, deps: Vec<NewStepDependency>) -> StoreFuture<'_, ()>;
183
184    /// List all step dependencies for a given run.
185    ///
186    /// Returns every edge where either `step_id` or `depends_on` belongs
187    /// to the run. Ordered by `created_at` ascending.
188    fn list_step_dependencies(&self, run_id: Uuid) -> StoreFuture<'_, Vec<StepDependency>>;
189
190    /// Apply a partial update to a run and return the updated run.
191    ///
192    /// Combines [`update_run`](Self::update_run) and [`get_run`](Self::get_run) in
193    /// a single operation to avoid an extra round-trip. Store implementations
194    /// may override this for efficiency (e.g. reading within the same transaction).
195    ///
196    /// The default implementation calls `update_run` followed by `get_run`.
197    ///
198    /// # Errors
199    ///
200    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
201    /// Returns [`StoreError::InvalidTransition`] if the status transition is not allowed.
202    fn update_run_returning(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, Run> {
203        Box::pin(async move {
204            self.update_run(id, update).await?;
205            self.get_run(id).await?.ok_or(StoreError::RunNotFound(id))
206        })
207    }
208}
209
210/// Unified storage abstraction combining all store capabilities.
211///
212/// Implementors provide runs, steps, users, API keys, and secrets
213/// through a single type. Pick one backend (in-memory or PostgreSQL)
214/// and it handles everything.
215///
216/// Both [`InMemoryStore`](crate::memory::InMemoryStore) and
217/// [`PostgresStore`](crate::postgres::PostgresStore) implement this trait.
218///
219/// # Examples
220///
221/// ```no_run
222/// use std::collections::HashMap;
223/// use std::sync::Arc;
224/// use ironflow_store::prelude::*;
225///
226/// # async fn example() -> Result<(), ironflow_store::error::StoreError> {
227/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
228///
229/// // All capabilities through one reference
230/// let _run = store.create_run(NewRun {
231///     workflow_name: "deploy".to_string(),
232///     trigger: TriggerKind::Manual,
233///     payload: serde_json::json!({}),
234///     max_retries: 3,
235///     handler_version: None,
236///     labels: HashMap::new(),
237///     scheduled_at: None,
238///     created_by: None,
239///     idempotency_key: None,
240///     max_cost_usd: None,
241/// }).await?.into_run();
242/// let _users = store.count_users().await?;
243/// # Ok(())
244/// # }
245/// ```
246pub trait Store:
247    RunStore + UserStore + ApiKeyStore + SecretStore + AuditLogStore + ArtifactStore
248{
249}
250
251impl<T: RunStore + UserStore + ApiKeyStore + SecretStore + AuditLogStore + ArtifactStore> Store
252    for T
253{
254}