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