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