1use std::collections::HashMap;
52use std::sync::Arc;
53use tokio::sync::{RwLock, broadcast};
54use serde::{Deserialize, Serialize};
55
56#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
62#[serde(untagged)]
63pub enum ConfigValue {
64 String(String),
65 Integer(i64),
66 Float(f64),
67 Boolean(bool),
68 Array(Vec<ConfigValue>),
69 Object(HashMap<String, ConfigValue>),
70 Null,
71}
72
73impl ConfigValue {
74 pub fn as_string(&self) -> Option<&str> {
76 match self {
77 ConfigValue::String(s) => Some(s),
78 _ => None,
79 }
80 }
81
82 pub fn as_integer(&self) -> Option<i64> {
84 match self {
85 ConfigValue::Integer(i) => Some(*i),
86 _ => None,
87 }
88 }
89
90 pub fn as_float(&self) -> Option<f64> {
92 match self {
93 ConfigValue::Float(f) => Some(*f),
94 ConfigValue::Integer(i) => Some(*i as f64),
95 _ => None,
96 }
97 }
98
99 pub fn as_boolean(&self) -> Option<bool> {
101 match self {
102 ConfigValue::Boolean(b) => Some(*b),
103 _ => None,
104 }
105 }
106
107 pub fn is_null(&self) -> bool {
109 matches!(self, ConfigValue::Null)
110 }
111}
112
113impl From<String> for ConfigValue {
114 fn from(s: String) -> Self {
115 ConfigValue::String(s)
116 }
117}
118
119impl From<&str> for ConfigValue {
120 fn from(s: &str) -> Self {
121 ConfigValue::String(s.to_string())
122 }
123}
124
125impl From<i64> for ConfigValue {
126 fn from(i: i64) -> Self {
127 ConfigValue::Integer(i)
128 }
129}
130
131impl From<i32> for ConfigValue {
132 fn from(i: i32) -> Self {
133 ConfigValue::Integer(i as i64)
134 }
135}
136
137impl From<f64> for ConfigValue {
138 fn from(f: f64) -> Self {
139 ConfigValue::Float(f)
140 }
141}
142
143impl From<bool> for ConfigValue {
144 fn from(b: bool) -> Self {
145 ConfigValue::Boolean(b)
146 }
147}
148
149#[derive(Debug, Clone)]
155pub struct ConfigChangeEvent {
156 pub key: String,
158 pub old_value: Option<ConfigValue>,
160 pub new_value: Option<ConfigValue>,
162 pub timestamp: chrono::DateTime<chrono::Utc>,
164}
165
166#[derive(Clone)]
175pub struct ConfigurationBus {
176 store: Arc<RwLock<HashMap<String, ConfigValue>>>,
178 change_tx: broadcast::Sender<ConfigChangeEvent>,
180}
181
182impl ConfigurationBus {
183 pub fn new() -> Self {
185 let (change_tx, _) = broadcast::channel(1000);
186 Self {
187 store: Arc::new(RwLock::new(HashMap::new())),
188 change_tx,
189 }
190 }
191
192 pub async fn set(&self, key: impl Into<String>, value: impl Into<ConfigValue>) {
196 let key = key.into();
197 let value = value.into();
198
199 let old_value = {
200 let mut store = self.store.write().await;
201 store.insert(key.clone(), value.clone())
202 };
203
204 let event = ConfigChangeEvent {
206 key,
207 old_value,
208 new_value: Some(value),
209 timestamp: chrono::Utc::now(),
210 };
211 let _ = self.change_tx.send(event);
212 }
213
214 pub async fn set_many(&self, values: HashMap<String, ConfigValue>) {
216 let mut store = self.store.write().await;
217 for (key, value) in values {
218 let old_value = store.insert(key.clone(), value.clone());
219 let event = ConfigChangeEvent {
220 key,
221 old_value,
222 new_value: Some(value),
223 timestamp: chrono::Utc::now(),
224 };
225 let _ = self.change_tx.send(event);
226 }
227 }
228
229 pub async fn get(&self, key: &str) -> Option<ConfigValue> {
231 self.store.read().await.get(key).cloned()
232 }
233
234 pub async fn get_string(&self, key: &str) -> Option<String> {
236 self.get(key).await.and_then(|v| v.as_string().map(|s| s.to_string()))
237 }
238
239 pub async fn get_integer(&self, key: &str) -> Option<i64> {
241 self.get(key).await.and_then(|v| v.as_integer())
242 }
243
244 pub async fn get_float(&self, key: &str) -> Option<f64> {
246 self.get(key).await.and_then(|v| v.as_float())
247 }
248
249 pub async fn get_boolean(&self, key: &str) -> Option<bool> {
251 self.get(key).await.and_then(|v| v.as_boolean())
252 }
253
254 pub async fn get_or_default(&self, key: &str, default: ConfigValue) -> ConfigValue {
256 self.get(key).await.unwrap_or(default)
257 }
258
259 pub async fn delete(&self, key: &str) -> Option<ConfigValue> {
263 let old_value = {
264 let mut store = self.store.write().await;
265 store.remove(key)
266 };
267
268 if let Some(ref old) = old_value {
269 let event = ConfigChangeEvent {
270 key: key.to_string(),
271 old_value: Some(old.clone()),
272 new_value: None,
273 timestamp: chrono::Utc::now(),
274 };
275 let _ = self.change_tx.send(event);
276 }
277
278 old_value
279 }
280
281 pub async fn contains(&self, key: &str) -> bool {
283 self.store.read().await.contains_key(key)
284 }
285
286 pub async fn keys(&self) -> Vec<String> {
288 self.store.read().await.keys().cloned().collect()
289 }
290
291 pub async fn keys_with_prefix(&self, prefix: &str) -> Vec<String> {
295 self.store
296 .read()
297 .await
298 .keys()
299 .filter(|k| k.starts_with(prefix))
300 .cloned()
301 .collect()
302 }
303
304 pub async fn get_with_prefix(&self, prefix: &str) -> HashMap<String, ConfigValue> {
306 self.store
307 .read()
308 .await
309 .iter()
310 .filter(|(k, _)| k.starts_with(prefix))
311 .map(|(k, v)| (k.clone(), v.clone()))
312 .collect()
313 }
314
315 pub fn subscribe(&self) -> broadcast::Receiver<ConfigChangeEvent> {
319 self.change_tx.subscribe()
320 }
321
322 pub async fn len(&self) -> usize {
324 self.store.read().await.len()
325 }
326
327 pub async fn is_empty(&self) -> bool {
329 self.store.read().await.is_empty()
330 }
331
332 pub async fn clear(&self) {
334 let mut store = self.store.write().await;
335 store.clear();
336 }
337
338 pub async fn dump(&self) -> HashMap<String, ConfigValue> {
340 self.store.read().await.clone()
341 }
342
343 pub async fn load(&self, values: HashMap<String, ConfigValue>) {
347 let mut store = self.store.write().await;
348 for (key, value) in values {
349 store.insert(key, value);
350 }
351 }
352}
353
354impl Default for ConfigurationBus {
355 fn default() -> Self {
356 Self::new()
357 }
358}
359
360#[cfg(test)]
365mod tests {
366 use super::*;
367
368 #[tokio::test]
369 async fn test_set_and_get() {
370 let bus = ConfigurationBus::new();
371
372 bus.set("test.key", "value").await;
373 let value = bus.get_string("test.key").await;
374
375 assert_eq!(value, Some("value".to_string()));
376 }
377
378 #[tokio::test]
379 async fn test_typed_getters() {
380 let bus = ConfigurationBus::new();
381
382 bus.set("string", ConfigValue::String("hello".into())).await;
383 bus.set("integer", ConfigValue::Integer(42)).await;
384 bus.set("float", ConfigValue::Float(3.14)).await;
385 bus.set("boolean", ConfigValue::Boolean(true)).await;
386
387 assert_eq!(bus.get_string("string").await, Some("hello".to_string()));
388 assert_eq!(bus.get_integer("integer").await, Some(42));
389 assert_eq!(bus.get_float("float").await, Some(3.14));
390 assert_eq!(bus.get_boolean("boolean").await, Some(true));
391 }
392
393 #[tokio::test]
394 async fn test_keys_with_prefix() {
395 let bus = ConfigurationBus::new();
396
397 bus.set("sapiens.auth.jwt_ttl", 3600).await;
398 bus.set("sapiens.auth.jwt_secret", "secret").await;
399 bus.set("postman.smtp.host", "smtp.example.com").await;
400
401 let sapiens_keys = bus.keys_with_prefix("sapiens.").await;
402 assert_eq!(sapiens_keys.len(), 2);
403
404 let postman_keys = bus.keys_with_prefix("postman.").await;
405 assert_eq!(postman_keys.len(), 1);
406 }
407
408 #[tokio::test]
409 async fn test_delete() {
410 let bus = ConfigurationBus::new();
411
412 bus.set("to.delete", "value").await;
413 assert!(bus.contains("to.delete").await);
414
415 let deleted = bus.delete("to.delete").await;
416 assert!(deleted.is_some());
417 assert!(!bus.contains("to.delete").await);
418 }
419
420 #[tokio::test]
421 async fn test_subscription() {
422 let bus = ConfigurationBus::new();
423 let mut rx = bus.subscribe();
424
425 bus.set("new.key", "new_value").await;
426
427 let event = rx.try_recv().unwrap();
428 assert_eq!(event.key, "new.key");
429 assert!(event.old_value.is_none());
430 assert_eq!(event.new_value, Some(ConfigValue::String("new_value".into())));
431 }
432
433 #[tokio::test]
434 async fn test_get_with_prefix() {
435 let bus = ConfigurationBus::new();
436
437 bus.set("app.name", "MyApp").await;
438 bus.set("app.version", "1.0.0").await;
439 bus.set("db.url", "postgres://localhost").await;
440
441 let app_config = bus.get_with_prefix("app.").await;
442 assert_eq!(app_config.len(), 2);
443 assert_eq!(app_config.get("app.name"), Some(&ConfigValue::String("MyApp".into())));
444 }
445
446 #[test]
447 fn test_config_value_conversions() {
448 let s: ConfigValue = "hello".into();
449 assert_eq!(s.as_string(), Some("hello"));
450
451 let i: ConfigValue = 42i64.into();
452 assert_eq!(i.as_integer(), Some(42));
453
454 let f: ConfigValue = 3.14f64.into();
455 assert_eq!(f.as_float(), Some(3.14));
456
457 let b: ConfigValue = true.into();
458 assert_eq!(b.as_boolean(), Some(true));
459 }
460}