Skip to main content

ironflow_store/memory/
mod.rs

1//! In-memory [`Store`](crate::store::Store) implementation for development and testing.
2//!
3//! [`InMemoryStore`] uses `Arc<RwLock<..>>` internally, making it safe to share
4//! across tasks. Data is lost when the process exits.
5//!
6//! # Examples
7//!
8//! ```no_run
9//! use std::collections::HashMap;
10//! use ironflow_store::prelude::*;
11//! use serde_json::json;
12//!
13//! # async fn example() -> Result<(), ironflow_store::error::StoreError> {
14//! let store = InMemoryStore::new();
15//!
16//! let run = store.create_run(NewRun {
17//!     workflow_name: "test".to_string(),
18//!     trigger: TriggerKind::Manual,
19//!     payload: json!({}),
20//!     max_retries: 3,
21//!     handler_version: None,
22//!     labels: HashMap::new(),
23//!     scheduled_at: None,
24//!     created_by: None,
25//!     idempotency_key: None,
26//!     concurrency_key: None,
27//!     priority: 0,
28//!     concurrency_limits: Vec::new(),
29//!     max_cost_usd: None,
30//!     worker_tags: Vec::new(),
31//! }).await?.into_run();
32//!
33//! assert_eq!(run.status.state, RunStatus::Pending);
34//! # Ok(())
35//! # }
36//! ```
37
38use std::collections::{BTreeSet, HashMap};
39use std::sync::Arc;
40
41use chrono::{DateTime, Utc};
42use tokio::sync::RwLock;
43use uuid::Uuid;
44
45use crate::entities::User;
46
47mod api_key_store;
48mod approval_delegation_store;
49mod artifact_store;
50mod audit_log_store;
51mod descendants;
52mod log_store;
53mod provider_account_store;
54mod run_store;
55mod schedule_store;
56mod secret_store;
57mod signal_store;
58mod stats_history;
59mod user_store;
60
61#[derive(Debug, Default)]
62pub(super) struct State {
63    pub(super) runs: HashMap<Uuid, crate::entities::Run>,
64    /// Idempotency key -> run holding it. Guarded by the same lock as `runs`,
65    /// so check-then-insert is atomic.
66    pub(super) idempotency_keys: HashMap<String, Uuid>,
67    pub(super) steps: HashMap<Uuid, crate::entities::Step>,
68    pub(super) step_dependencies: Vec<crate::entities::StepDependency>,
69    pub(super) artifacts: HashMap<Uuid, crate::entities::Artifact>,
70    pub(super) users: HashMap<Uuid, User>,
71    /// Group membership per user, kept sorted and deduplicated.
72    pub(super) user_groups: HashMap<Uuid, BTreeSet<String>>,
73    pub(super) api_keys: HashMap<Uuid, crate::entities::ApiKey>,
74    pub(super) secrets: HashMap<String, EncryptedSecret>,
75    pub(super) schedules: HashMap<Uuid, crate::entities::Schedule>,
76    pub(super) approval_delegations: HashMap<Uuid, crate::entities::ApprovalDelegation>,
77    pub(super) audit_logs: Vec<crate::entities::AuditLogEntry>,
78    pub(super) log_entries: Vec<crate::entities::LogEntry>,
79    pub(super) provider_accounts: HashMap<Uuid, crate::entities::ProviderAccount>,
80    /// Latest window per `(account_id, window, model_scope)`, `""` for no scope.
81    pub(super) provider_account_windows:
82        HashMap<(Uuid, String, String), crate::entities::ProviderAccountWindow>,
83    pub(super) provider_account_usage: Vec<crate::entities::ProviderAccountUsagePoint>,
84    pub(super) signals: Vec<crate::entities::Signal>,
85    /// Idempotency ID -> signal holding it. Guarded by the same lock as
86    /// `signals`, so check-then-insert is atomic.
87    pub(super) signal_idempotency: HashMap<String, Uuid>,
88    /// Issued refresh tokens, keyed by their SHA-256 hash.
89    pub(super) refresh_tokens: HashMap<String, StoredRefreshToken>,
90    /// Paused workflows, keyed by workflow name.
91    pub(super) workflow_pauses: HashMap<String, crate::entities::WorkflowPause>,
92}
93
94#[derive(Debug, Clone)]
95pub(super) struct StoredRefreshToken {
96    pub(super) user_id: Uuid,
97    pub(super) expires_at: DateTime<Utc>,
98}
99
100#[derive(Debug, Clone)]
101pub(super) struct EncryptedSecret {
102    pub(super) id: Uuid,
103    pub(super) key: String,
104    #[cfg(feature = "secret-store")]
105    pub(super) encrypted_value: Vec<u8>,
106    #[cfg(feature = "secret-store")]
107    pub(super) nonce: Vec<u8>,
108    #[cfg(feature = "secret-store")]
109    pub(super) key_version: i32,
110    pub(super) created_at: chrono::DateTime<chrono::Utc>,
111    pub(super) updated_at: chrono::DateTime<chrono::Utc>,
112}
113
114/// In-memory store backed by `Arc<RwLock<..>>`.
115///
116/// Thread-safe and cheap to clone. All data is held in memory and lost on drop.
117/// Implements [`Store`](crate::store::Store) so a single `Arc<InMemoryStore>`
118/// covers runs, users, API keys, and secrets.
119///
120/// # Examples
121///
122/// ```
123/// use ironflow_store::memory::InMemoryStore;
124///
125/// let store = InMemoryStore::new();
126/// let store2 = store.clone(); // cheap Arc clone
127/// ```
128#[derive(Debug, Clone)]
129pub struct InMemoryStore {
130    pub(super) state: Arc<RwLock<State>>,
131    #[cfg(feature = "secret-store")]
132    pub(super) key_ring: Option<Arc<crate::crypto::KeyRing>>,
133}
134
135impl InMemoryStore {
136    /// Create a new empty in-memory store.
137    ///
138    /// # Examples
139    ///
140    /// ```
141    /// use ironflow_store::memory::InMemoryStore;
142    ///
143    /// let store = InMemoryStore::new();
144    /// ```
145    pub fn new() -> Self {
146        Self {
147            state: Arc::new(RwLock::new(State::default())),
148            #[cfg(feature = "secret-store")]
149            key_ring: None,
150        }
151    }
152
153    /// Set a single, unversioned master key for secret encryption.
154    ///
155    /// Shorthand for a key ring holding this key alone at
156    /// [`LEGACY_KEY_VERSION`](crate::crypto::LEGACY_KEY_VERSION).
157    ///
158    /// Required before using [`SecretStore`](crate::secret_store::SecretStore)
159    /// methods that read/write secret values. Without a key, those methods
160    /// return [`StoreError::Crypto`](crate::error::StoreError::Crypto).
161    ///
162    /// Listing and deleting secrets works without a key.
163    ///
164    /// # Examples
165    ///
166    /// ```
167    /// use ironflow_store::memory::InMemoryStore;
168    /// use ironflow_store::crypto::MasterKey;
169    ///
170    /// # fn example() -> Result<(), ironflow_store::crypto::CryptoError> {
171    /// let mut store = InMemoryStore::new();
172    /// let key = MasterKey::from_hex(
173    ///     "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"
174    /// )?;
175    /// store.set_master_key(key);
176    /// # Ok(())
177    /// # }
178    /// ```
179    #[cfg(feature = "secret-store")]
180    pub fn set_master_key(&mut self, key: crate::crypto::MasterKey) {
181        self.set_key_ring(crate::crypto::KeyRing::single(key));
182    }
183
184    /// Set the versioned key ring for secret encryption.
185    ///
186    /// New secrets are encrypted with the ring's active version; existing ones
187    /// are decrypted with whichever version they were written with.
188    ///
189    /// # Examples
190    ///
191    /// ```
192    /// use ironflow_store::memory::InMemoryStore;
193    /// use ironflow_store::crypto::KeyRing;
194    ///
195    /// # fn example() -> Result<(), ironflow_store::crypto::CryptoError> {
196    /// let mut store = InMemoryStore::new();
197    /// let spec = format!("1:{},2:{}", "aa".repeat(32), "bb".repeat(32));
198    /// store.set_key_ring(KeyRing::from_spec(&spec, Some(2))?);
199    /// # Ok(())
200    /// # }
201    /// ```
202    #[cfg(feature = "secret-store")]
203    pub fn set_key_ring(&mut self, ring: crate::crypto::KeyRing) {
204        self.key_ring = Some(Arc::new(ring));
205    }
206
207    /// Override a run's `created_at` timestamp for testing retention policies.
208    ///
209    /// # Examples
210    ///
211    /// ```no_run
212    /// use chrono::{Utc, TimeDelta};
213    /// use ironflow_store::memory::InMemoryStore;
214    /// use uuid::Uuid;
215    ///
216    /// # async fn example(store: &InMemoryStore, run_id: Uuid) {
217    /// let old = Utc::now() - TimeDelta::days(100);
218    /// store.set_run_created_at(run_id, old).await;
219    /// # }
220    /// ```
221    pub async fn set_run_created_at(
222        &self,
223        run_id: Uuid,
224        created_at: chrono::DateTime<chrono::Utc>,
225    ) {
226        let mut state = self.state.write().await;
227        if let Some(run) = state.runs.get_mut(&run_id) {
228            run.created_at = created_at;
229        }
230    }
231}
232
233impl Default for InMemoryStore {
234    fn default() -> Self {
235        Self::new()
236    }
237}
238
239#[cfg(test)]
240mod tests {
241    use std::collections::HashMap;
242
243    use serde_json::json;
244
245    use crate::entities::{NewRun, TriggerKind};
246
247    use super::InMemoryStore;
248
249    pub(crate) fn new_run_req(name: &str) -> NewRun {
250        NewRun {
251            created_by: None,
252            workflow_name: name.to_string(),
253            trigger: TriggerKind::Manual,
254            payload: json!({}),
255            max_retries: 3,
256            handler_version: None,
257            labels: HashMap::new(),
258            scheduled_at: None,
259            idempotency_key: None,
260            concurrency_key: None,
261            priority: 0,
262            concurrency_limits: Vec::new(),
263            max_cost_usd: None,
264            worker_tags: Vec::new(),
265        }
266    }
267
268    pub(crate) async fn create_terminal_run(
269        store: &InMemoryStore,
270        name: &str,
271        status: crate::entities::RunStatus,
272    ) -> crate::entities::Run {
273        use crate::store::RunStore;
274
275        let run = store
276            .create_run(new_run_req(name))
277            .await
278            .unwrap()
279            .into_run();
280        store
281            .update_run_status(run.id, crate::entities::RunStatus::Running)
282            .await
283            .unwrap();
284        store.update_run_status(run.id, status).await.unwrap();
285        store.get_run(run.id).await.unwrap().unwrap()
286    }
287}