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