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