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