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