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}