Skip to main content

backbone_core/config/
bus.rs

1//! Configuration Bus - Cross-Module Configuration Sharing
2//!
3//! The Configuration Bus provides a mechanism for modules to:
4//! - Share configuration values across bounded contexts
5//! - Subscribe to configuration changes
6//! - Access module-specific configuration without tight coupling
7//!
8//! # Architecture
9//!
10//! ```text
11//! ┌────────────────────────────────────────────────────────────────┐
12//! │                    Configuration Bus                           │
13//! │  ┌─────────────────────────────────────────────────────────┐  │
14//! │  │              Shared Configuration Store                  │  │
15//! │  │  sapiens.auth.jwt_ttl = 3600                            │  │
16//! │  │  postman.smtp.host = "smtp.example.com"                 │  │
17//! │  │  bucket.storage.max_size = 1073741824                 │  │
18//! │  └─────────────────────────────────────────────────────────┘  │
19//! │                              │                                 │
20//! │        ┌────────────────────┼────────────────────┐            │
21//! │        │                    │                    │            │
22//! │        ▼                    ▼                    ▼            │
23//! │   ┌─────────┐         ┌─────────┐         ┌─────────┐        │
24//! │   │ Sapiens │         │ Postman │         │Bucket │        │
25//! │   │ Module  │         │ Module  │         │ Module  │        │
26//! │   └─────────┘         └─────────┘         └─────────┘        │
27//! └────────────────────────────────────────────────────────────────┘
28//! ```
29//!
30//! # Usage
31//!
32//! ```rust,ignore
33//! use backbone_core::config::ConfigurationBus;
34//!
35//! // Create shared configuration bus
36//! let config_bus = ConfigurationBus::new();
37//!
38//! // Set module configuration
39//! config_bus.set("sapiens.auth.jwt_ttl", ConfigValue::Integer(3600)).await;
40//! config_bus.set("postman.smtp.host", ConfigValue::String("smtp.example.com".into())).await;
41//!
42//! // Get configuration from any module
43//! if let Some(ttl) = config_bus.get_integer("sapiens.auth.jwt_ttl").await {
44//!     println!("JWT TTL: {}", ttl);
45//! }
46//!
47//! // Subscribe to configuration changes
48//! let rx = config_bus.subscribe("sapiens.*").await;
49//! ```
50
51use std::collections::HashMap;
52use std::sync::Arc;
53use tokio::sync::{RwLock, broadcast};
54use serde::{Deserialize, Serialize};
55
56// ============================================================
57// Configuration Value Types
58// ============================================================
59
60/// Configuration value that can hold different types
61#[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    /// Get as string
75    pub fn as_string(&self) -> Option<&str> {
76        match self {
77            ConfigValue::String(s) => Some(s),
78            _ => None,
79        }
80    }
81
82    /// Get as integer
83    pub fn as_integer(&self) -> Option<i64> {
84        match self {
85            ConfigValue::Integer(i) => Some(*i),
86            _ => None,
87        }
88    }
89
90    /// Get as float
91    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    /// Get as boolean
100    pub fn as_boolean(&self) -> Option<bool> {
101        match self {
102            ConfigValue::Boolean(b) => Some(*b),
103            _ => None,
104        }
105    }
106
107    /// Check if null
108    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// ============================================================
150// Configuration Change Event
151// ============================================================
152
153/// Event emitted when configuration changes
154#[derive(Debug, Clone)]
155pub struct ConfigChangeEvent {
156    /// The configuration key that changed
157    pub key: String,
158    /// The old value (None if newly created)
159    pub old_value: Option<ConfigValue>,
160    /// The new value (None if deleted)
161    pub new_value: Option<ConfigValue>,
162    /// Timestamp of the change
163    pub timestamp: chrono::DateTime<chrono::Utc>,
164}
165
166// ============================================================
167// Configuration Bus
168// ============================================================
169
170/// Shared configuration bus for cross-module configuration
171///
172/// Provides a thread-safe, async-friendly way for modules to share
173/// configuration without direct dependencies.
174#[derive(Clone)]
175pub struct ConfigurationBus {
176    /// The configuration store
177    store: Arc<RwLock<HashMap<String, ConfigValue>>>,
178    /// Broadcast channel for configuration changes
179    change_tx: broadcast::Sender<ConfigChangeEvent>,
180}
181
182impl ConfigurationBus {
183    /// Create a new configuration bus
184    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    /// Set a configuration value
193    ///
194    /// Notifies all subscribers of the change.
195    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        // Emit change event
205        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    /// Set multiple configuration values at once
215    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    /// Get a configuration value
230    pub async fn get(&self, key: &str) -> Option<ConfigValue> {
231        self.store.read().await.get(key).cloned()
232    }
233
234    /// Get a string configuration value
235    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    /// Get an integer configuration value
240    pub async fn get_integer(&self, key: &str) -> Option<i64> {
241        self.get(key).await.and_then(|v| v.as_integer())
242    }
243
244    /// Get a float configuration value
245    pub async fn get_float(&self, key: &str) -> Option<f64> {
246        self.get(key).await.and_then(|v| v.as_float())
247    }
248
249    /// Get a boolean configuration value
250    pub async fn get_boolean(&self, key: &str) -> Option<bool> {
251        self.get(key).await.and_then(|v| v.as_boolean())
252    }
253
254    /// Get a configuration value with default
255    pub async fn get_or_default(&self, key: &str, default: ConfigValue) -> ConfigValue {
256        self.get(key).await.unwrap_or(default)
257    }
258
259    /// Delete a configuration value
260    ///
261    /// Returns the deleted value if it existed.
262    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    /// Check if a configuration key exists
282    pub async fn contains(&self, key: &str) -> bool {
283        self.store.read().await.contains_key(key)
284    }
285
286    /// List all configuration keys
287    pub async fn keys(&self) -> Vec<String> {
288        self.store.read().await.keys().cloned().collect()
289    }
290
291    /// List configuration keys matching a prefix
292    ///
293    /// For example, `keys_with_prefix("sapiens.")` returns all sapiens configuration.
294    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    /// Get all configuration values matching a prefix
305    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    /// Subscribe to configuration changes
316    ///
317    /// Returns a broadcast receiver that receives all configuration changes.
318    pub fn subscribe(&self) -> broadcast::Receiver<ConfigChangeEvent> {
319        self.change_tx.subscribe()
320    }
321
322    /// Get the total number of configuration entries
323    pub async fn len(&self) -> usize {
324        self.store.read().await.len()
325    }
326
327    /// Check if the configuration store is empty
328    pub async fn is_empty(&self) -> bool {
329        self.store.read().await.is_empty()
330    }
331
332    /// Clear all configuration values
333    pub async fn clear(&self) {
334        let mut store = self.store.write().await;
335        store.clear();
336    }
337
338    /// Dump all configuration as a HashMap
339    pub async fn dump(&self) -> HashMap<String, ConfigValue> {
340        self.store.read().await.clone()
341    }
342
343    /// Load configuration from a HashMap
344    ///
345    /// Merges with existing configuration (overwrites duplicates).
346    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// ============================================================
361// Tests
362// ============================================================
363
364#[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}