moirai_core/communication/
router.rs1#![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
11pub struct MessageRouter<K: Hash + Eq + Clone, V: Send + 'static> {
13 routes: Arc<RwLock<HashMap<K, MpmcSender<V>>>>,
15}
16
17impl<K: Hash + Eq + Clone, V: Send + 'static> MessageRouter<K, V> {
18 pub fn new() -> Self {
20 Self {
21 routes: Arc::new(RwLock::new(HashMap::new())),
22 }
23 }
24
25 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 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 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}