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