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