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