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