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