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}