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