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, PurgePolicy, PurgeableRun, ReapedRun,
19    Run, RunCreation, RunFilter, RunStats, RunStatus, RunUpdate, StatsHistoryBucket,
20    StatsHistoryFilter, Step, StepDependency, StepUpdate,
21};
22use crate::error::StoreError;
23use crate::log_store::LogStore;
24use crate::schedule_store::ScheduleStore;
25use crate::secret_store::SecretStore;
26use crate::user_store::UserStore;
27
28/// Boxed future for [`RunStore`] methods — ensures object safety for `dyn RunStore`.
29pub type StoreFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, StoreError>> + Send + 'a>>;
30
31/// Error recorded on a run that exhausted its retries through lease expiries.
32///
33/// Set by [`RunStore::reap_expired_leases`] when a run has been recovered more
34/// than `max_retries` times.
35pub const LEASE_EXPIRED_ERROR: &str = "worker lease expired";
36
37/// Async storage abstraction for workflow runs and steps.
38///
39/// All methods return a [`StoreFuture`] (boxed future) to maintain object safety,
40/// allowing the store to be used as `Arc<dyn RunStore>`.
41///
42/// # Examples
43///
44/// ```no_run
45/// use std::collections::HashMap;
46/// use ironflow_store::prelude::*;
47/// use serde_json::json;
48/// use uuid::Uuid;
49///
50/// # async fn example() -> Result<(), ironflow_store::error::StoreError> {
51/// let store = InMemoryStore::new();
52///
53/// let run = store.create_run(NewRun {
54///     workflow_name: "deploy".to_string(),
55///     trigger: TriggerKind::Manual,
56///     payload: json!({}),
57///     max_retries: 3,
58///     handler_version: None,
59///     labels: HashMap::new(),
60///     scheduled_at: None,
61///     created_by: None,
62///     idempotency_key: None,
63///     max_cost_usd: None,
64/// }).await?.into_run();
65///
66/// let fetched = store.get_run(run.id).await?;
67/// assert!(fetched.is_some());
68/// # Ok(())
69/// # }
70/// ```
71pub trait RunStore: Send + Sync {
72    /// Create a new run in `Pending` status.
73    ///
74    /// When [`NewRun::idempotency_key`] is set and already bound to a run created
75    /// within [`IDEMPOTENCY_WINDOW`](crate::entities::IDEMPOTENCY_WINDOW), nothing is
76    /// inserted and that run is returned as [`RunCreation::Existing`]. A key bound to
77    /// an older run is released and reused for the new one.
78    ///
79    /// Concurrent calls sharing the same key resolve to a single run: exactly one
80    /// receives [`RunCreation::Created`], the others [`RunCreation::Existing`].
81    fn create_run(&self, req: NewRun) -> StoreFuture<'_, RunCreation>;
82
83    /// Look up the run bound to an idempotency key.
84    ///
85    /// Returns `None` when the key is unknown, or when the run holding it is older
86    /// than [`IDEMPOTENCY_WINDOW`](crate::entities::IDEMPOTENCY_WINDOW).
87    fn find_run_by_idempotency_key(&self, key: &str) -> StoreFuture<'_, Option<Run>>;
88
89    /// Get a run by ID. Returns `None` if not found.
90    fn get_run(&self, id: Uuid) -> StoreFuture<'_, Option<Run>>;
91
92    /// List runs matching the given filter, with pagination.
93    ///
94    /// Results are ordered by `created_at` descending (newest first).
95    fn list_runs(&self, filter: RunFilter, page: u32, per_page: u32) -> StoreFuture<'_, Page<Run>>;
96
97    /// Update a run's status with FSM validation.
98    ///
99    /// # Errors
100    ///
101    /// Returns [`StoreError::InvalidTransition`] if the transition is not allowed.
102    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
103    fn update_run_status(&self, id: Uuid, new_status: RunStatus) -> StoreFuture<'_, ()>;
104
105    /// Apply a partial update to a run.
106    ///
107    /// # Errors
108    ///
109    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
110    fn update_run(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, ()>;
111
112    /// Atomically pick the oldest pending run and transition it to `Running`.
113    ///
114    /// In PostgreSQL, this uses `SELECT FOR UPDATE SKIP LOCKED` for safe
115    /// multi-worker concurrency. The in-memory implementation uses a write lock.
116    ///
117    /// When `lease` is `Some`, the worker lease is attached in the same
118    /// transaction as the status change, so a run is never `Running` without an
119    /// owner. Pass `None` for callers that execute runs in-process and cannot
120    /// refresh a lease (inline execution, API-side resume): those runs are never
121    /// recovered by [`reap_expired_leases`](Self::reap_expired_leases).
122    ///
123    /// Returns `None` if no pending runs are available.
124    fn pick_next_pending(&self, lease: Option<LeaseRequest>) -> StoreFuture<'_, Option<Run>>;
125
126    /// Extend the worker lease on a run and return the new expiry.
127    ///
128    /// # Errors
129    ///
130    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
131    /// Returns [`StoreError::LeaseLost`] if the run is no longer `Running` or if
132    /// the lease belongs to another worker — the caller must stop executing it.
133    fn renew_lease(&self, id: Uuid, lease: LeaseRequest) -> StoreFuture<'_, DateTime<Utc>>;
134
135    /// Recover runs whose worker lease expired, at most `limit` per call.
136    ///
137    /// Each recovered run has its retry count incremented and its lease cleared,
138    /// then goes back to `Pending` — or to `Failed` with `worker lease expired`
139    /// once `max_retries` is exhausted. Runs without a lease are never touched.
140    ///
141    /// The whole batch is atomic per run (`FOR UPDATE SKIP LOCKED` in
142    /// PostgreSQL), so concurrent reapers never recover the same run twice.
143    ///
144    /// Callers are responsible for the side effects that follow a recovery:
145    /// failing orphaned steps and publishing status-change events.
146    fn reap_expired_leases(&self, limit: u32) -> StoreFuture<'_, Vec<ReapedRun>>;
147
148    /// Atomically claim approval steps whose SLA deadline has passed.
149    ///
150    /// Returns the claimed steps with their *pre-claim* `approval_deadline_at`
151    /// still populated, so the caller can report which deadline fired. The
152    /// timer is cleared in the same transaction, so a deadline fires at most
153    /// once even with several API instances running the escalator (the
154    /// PostgreSQL implementation uses `FOR UPDATE SKIP LOCKED`).
155    ///
156    /// Only steps still in [`StepStatus::AwaitingApproval`](crate::entities::StepStatus::AwaitingApproval)
157    /// are returned.
158    ///
159    /// Delivery is at most once: a caller that crashes between the claim and
160    /// the escalation leaves the gate open with no timer, the same trade-off
161    /// [`reap_expired_leases`](Self::reap_expired_leases) accepts.
162    fn claim_due_approval_deadlines(&self, limit: u32) -> StoreFuture<'_, Vec<Step>>;
163
164    /// Create a new step for a run.
165    ///
166    /// # Errors
167    ///
168    /// Returns [`StoreError::RunNotFound`] if the parent run does not exist.
169    fn create_step(&self, step: NewStep) -> StoreFuture<'_, Step>;
170
171    /// Apply a partial update to a step after execution.
172    ///
173    /// # Errors
174    ///
175    /// Returns [`StoreError::StepNotFound`] if the step does not exist.
176    fn update_step(&self, id: Uuid, update: StepUpdate) -> StoreFuture<'_, ()>;
177
178    /// Get a single step by ID. Returns `None` if not found.
179    fn get_step(&self, id: Uuid) -> StoreFuture<'_, Option<Step>>;
180
181    /// List all steps for a run, ordered by position ascending.
182    fn list_steps(&self, run_id: Uuid) -> StoreFuture<'_, Vec<Step>>;
183
184    /// Get aggregated statistics across runs matching the filter.
185    ///
186    /// Returns counts of runs by terminal state, counts of active runs,
187    /// and totals for cost and duration. Computed efficiently by the store
188    /// implementation (single SQL query in PostgreSQL).
189    ///
190    /// Pass [`RunFilter::default()`] to get stats across all runs.
191    fn get_stats(&self, filter: RunFilter) -> StoreFuture<'_, RunStats>;
192
193    /// Get time-bucketed historical statistics for trend charts.
194    ///
195    /// Aggregates runs into time buckets based on the filter's granularity,
196    /// counting terminal statuses and computing duration percentiles.
197    /// Returns buckets ordered by time ascending; empty periods are omitted.
198    ///
199    /// # Errors
200    ///
201    /// Returns [`StoreError::Database`] on underlying store failures.
202    ///
203    /// # Examples
204    ///
205    /// ```no_run
206    /// use ironflow_store::entities::{StatsHistoryFilter, HistoryPeriod, HistoryGranularity};
207    /// use ironflow_store::store::RunStore;
208    ///
209    /// # async fn example(store: &dyn RunStore) -> Result<(), ironflow_store::error::StoreError> {
210    /// let filter = StatsHistoryFilter {
211    ///     workflow_name: None,
212    ///     period: HistoryPeriod::SevenDays,
213    ///     granularity: HistoryGranularity::OneDay,
214    /// };
215    /// let buckets = store.get_stats_history(filter).await?;
216    /// for b in &buckets {
217    ///     println!("{}: {} completed, {} failed", b.time, b.completed, b.failed);
218    /// }
219    /// # Ok(())
220    /// # }
221    /// ```
222    fn get_stats_history(
223        &self,
224        filter: StatsHistoryFilter,
225    ) -> StoreFuture<'_, Vec<StatsHistoryBucket>>;
226
227    /// Create step dependency edges in batch.
228    ///
229    /// Each entry records that `step_id` depends on `depends_on`.
230    /// Duplicate edges are silently ignored.
231    ///
232    /// # Errors
233    ///
234    /// Returns [`StoreError`] if a referenced step does not exist.
235    fn create_step_dependencies(&self, deps: Vec<NewStepDependency>) -> StoreFuture<'_, ()>;
236
237    /// List all step dependencies for a given run.
238    ///
239    /// Returns every edge where either `step_id` or `depends_on` belongs
240    /// to the run. Ordered by `created_at` ascending.
241    fn list_step_dependencies(&self, run_id: Uuid) -> StoreFuture<'_, Vec<StepDependency>>;
242
243    /// List runs eligible for purging according to the given policy.
244    ///
245    /// A run is eligible when it is in a terminal state ([`RunStatus::is_terminal`])
246    /// **and** either older than `policy.max_age_days` or exceeding
247    /// `policy.max_runs_per_workflow` for its workflow (oldest first).
248    ///
249    /// Runs in non-terminal states (`Pending`, `Running`, `Retrying`,
250    /// `AwaitingApproval`) are never returned.
251    ///
252    /// # Examples
253    ///
254    /// ```no_run
255    /// use ironflow_store::entities::PurgePolicy;
256    /// use ironflow_store::store::RunStore;
257    ///
258    /// # async fn example(store: &dyn RunStore) -> Result<(), ironflow_store::error::StoreError> {
259    /// let policy = PurgePolicy { max_age_days: 90, max_runs_per_workflow: 1000, dry_run: false };
260    /// let purgeable = store.list_purgeable_runs(&policy, 100).await?;
261    /// for p in &purgeable {
262    ///     println!("purge {} ({}): {}", p.run_id, p.workflow_name, p.reason);
263    /// }
264    /// # Ok(())
265    /// # }
266    /// ```
267    fn list_purgeable_runs(
268        &self,
269        policy: &PurgePolicy,
270        batch_size: u32,
271    ) -> StoreFuture<'_, Vec<PurgeableRun>>;
272
273    /// Delete a run and all its associated data (steps, step dependencies).
274    ///
275    /// Returns the `storage_key` of every artifact that belonged to the run,
276    /// so the caller can delete the corresponding blobs from the blob store.
277    ///
278    /// # Errors
279    ///
280    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
281    ///
282    /// # Examples
283    ///
284    /// ```no_run
285    /// use ironflow_store::store::RunStore;
286    /// use uuid::Uuid;
287    ///
288    /// # async fn example(store: &dyn RunStore, run_id: Uuid) -> Result<(), ironflow_store::error::StoreError> {
289    /// let storage_keys = store.delete_run(run_id).await?;
290    /// // Caller deletes blobs from the blob store using these keys.
291    /// # Ok(())
292    /// # }
293    /// ```
294    fn delete_run(&self, id: Uuid) -> StoreFuture<'_, Vec<String>>;
295
296    /// Apply a partial update to a run and return the updated run.
297    ///
298    /// Combines [`update_run`](Self::update_run) and [`get_run`](Self::get_run) in
299    /// a single operation to avoid an extra round-trip. Store implementations
300    /// may override this for efficiency (e.g. reading within the same transaction).
301    ///
302    /// The default implementation calls `update_run` followed by `get_run`.
303    ///
304    /// # Errors
305    ///
306    /// Returns [`StoreError::RunNotFound`] if the run does not exist.
307    /// Returns [`StoreError::InvalidTransition`] if the status transition is not allowed.
308    fn update_run_returning(&self, id: Uuid, update: RunUpdate) -> StoreFuture<'_, Run> {
309        Box::pin(async move {
310            self.update_run(id, update).await?;
311            self.get_run(id).await?.ok_or(StoreError::RunNotFound(id))
312        })
313    }
314}
315
316/// Unified storage abstraction combining all store capabilities.
317///
318/// Implementors provide runs, steps, users, API keys, and secrets
319/// through a single type. Pick one backend (in-memory or PostgreSQL)
320/// and it handles everything.
321///
322/// Both [`InMemoryStore`](crate::memory::InMemoryStore) and
323/// [`PostgresStore`](crate::postgres::PostgresStore) implement this trait.
324///
325/// # Examples
326///
327/// ```no_run
328/// use std::collections::HashMap;
329/// use std::sync::Arc;
330/// use ironflow_store::prelude::*;
331///
332/// # async fn example() -> Result<(), ironflow_store::error::StoreError> {
333/// let store: Arc<dyn Store> = Arc::new(InMemoryStore::new());
334///
335/// // All capabilities through one reference
336/// let _run = store.create_run(NewRun {
337///     workflow_name: "deploy".to_string(),
338///     trigger: TriggerKind::Manual,
339///     payload: serde_json::json!({}),
340///     max_retries: 3,
341///     handler_version: None,
342///     labels: HashMap::new(),
343///     scheduled_at: None,
344///     created_by: None,
345///     idempotency_key: None,
346///     max_cost_usd: None,
347/// }).await?.into_run();
348/// let _users = store.count_users().await?;
349/// # Ok(())
350/// # }
351/// ```
352pub trait Store:
353    RunStore
354    + UserStore
355    + ApiKeyStore
356    + SecretStore
357    + AuditLogStore
358    + ArtifactStore
359    + LogStore
360    + ScheduleStore
361{
362}
363
364impl<
365    T: RunStore
366        + UserStore
367        + ApiKeyStore
368        + SecretStore
369        + AuditLogStore
370        + ArtifactStore
371        + LogStore
372        + ScheduleStore,
373> Store for T
374{
375}