Skip to main content

moirai_core/communication/
router.rs

1#![expect(
2    clippy::unwrap_used,
3    reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
4)]
5
6use crate::channel::{ChannelError, MpmcSender};
7use std::collections::HashMap;
8use std::hash::Hash;
9use std::sync::{Arc, RwLock};
10
11/// Router for message-based communication patterns
12pub struct MessageRouter<K: Hash + Eq + Clone, V: Send + 'static> {
13    /// Routes mapped by key
14    routes: Arc<RwLock<HashMap<K, MpmcSender<V>>>>,
15}
16
17impl<K: Hash + Eq + Clone, V: Send + 'static> MessageRouter<K, V> {
18    /// Create a new message router
19    pub fn new() -> Self {
20        Self {
21            routes: Arc::new(RwLock::new(HashMap::new())),
22        }
23    }
24
25    /// Register a route
26    pub fn register(&self, key: K, sender: MpmcSender<V>) {
27        let mut routes = self.routes.write().unwrap();
28        routes.insert(key, sender);
29    }
30
31    /// Route a message to the appropriate channel
32    pub fn route(&self, key: &K, message: V) -> Result<(), ChannelError> {
33        let routes = self.routes.read().unwrap();
34
35        if let Some(sender) = routes.get(key) {
36            sender.try_send(message)
37        } else {
38            Err(ChannelError::Closed)
39        }
40    }
41
42    /// Remove a route
43    pub fn unregister(&self, key: &K) -> bool {
44        let mut routes = self.routes.write().unwrap();
45        routes.remove(key).is_some()
46    }
47}
48
49impl<K: Hash + Eq + Clone, V: Send + 'static> Default for MessageRouter<K, V> {
50    fn default() -> Self {
51        Self::new()
52    }
53}