Skip to main content

qualia_core_db/net/acoustic_ble_mesh/
manager.rs

1//! Top-level mesh network manager: orchestrates the acoustic and BLE networks,
2//! the router, the data store, and the performance monitor; owns node discovery,
3//! message send/receive, routing, and interface selection.
4
5use super::*;
6
7/// Acoustic & BLE Mesh Network Manager
8pub struct MeshNetworkManager {
9    acoustic_network: AcousticNetwork,
10    ble_network: BleNetwork,
11    mesh_router: MeshRouter,
12    data_store: MeshDataStore,
13    performance_monitor: MeshPerformanceMonitor,
14}
15
16impl MeshNetworkManager {
17    /// Create new mesh network manager
18    pub fn new() -> Self {
19        Self {
20            acoustic_network: AcousticNetwork::new(),
21            ble_network: BleNetwork::new(),
22            mesh_router: MeshRouter::new(),
23            data_store: MeshDataStore::new(),
24            performance_monitor: MeshPerformanceMonitor::new(),
25        }
26    }
27
28    /// Initialize mesh networks
29    pub fn initialize(&mut self) -> Result<(), MeshError> {
30        // Initialize acoustic network
31        self.acoustic_network.initialize()?;
32
33        // Initialize BLE network
34        self.ble_network.initialize()?;
35
36        // Initialize mesh router
37        self.mesh_router.initialize()?;
38
39        // Initialize data store
40        self.data_store.initialize()?;
41
42        Ok(())
43    }
44
45    /// Discover nearby nodes
46    pub fn discover_nodes(&mut self) -> Result<Vec<DiscoveredNode>, MeshError> {
47        let mut discovered_nodes = Vec::new();
48
49        // Discover acoustic nodes
50        let acoustic_nodes = self.acoustic_network.discover_nodes()?;
51        for node in acoustic_nodes {
52            discovered_nodes.push(DiscoveredNode {
53                node_id: node.node_id.clone(),
54                node_type: node.node_type.clone(),
55                interface: NetworkInterface::Acoustic,
56                capabilities: NodeCapabilities::Acoustic(node.capabilities.clone()),
57                signal_strength: node.signal_strength,
58                location: node.location,
59            });
60        }
61
62        // Discover BLE nodes
63        let ble_nodes = self.ble_network.discover_nodes()?;
64        for node in ble_nodes {
65            discovered_nodes.push(DiscoveredNode {
66                node_id: node.node_id.clone(),
67                node_type: NodeType::Mobile, // BLE nodes are typically mobile
68                interface: NetworkInterface::Ble,
69                capabilities: NodeCapabilities::Ble(node.capabilities.clone()),
70                signal_strength: node.rssi as f64,
71                location: None, // BLE nodes typically don't have location info
72            });
73        }
74
75        Ok(discovered_nodes)
76    }
77
78    /// Discover nearby nodes into a caller-owned zero-heap buffer.
79    pub fn discover_nodes_into(
80        &mut self,
81        out: &mut [DiscoveredNodeHandle],
82    ) -> Result<usize, MeshError> {
83        let mut written = 0usize;
84
85        let acoustic_nodes = self.acoustic_network.discover_nodes()?;
86        for node in acoustic_nodes {
87            if written >= out.len() {
88                return Err(MeshError::BufferTooSmall(
89                    "discovered node output buffer exhausted".to_string(),
90                ));
91            }
92            out[written] = Self::discovered_node_handle(
93                &node.node_id,
94                node.node_type,
95                NetworkInterface::Acoustic,
96                Self::acoustic_capability_tag(&node.capabilities),
97                node.signal_strength,
98                node.location,
99            );
100            written += 1;
101        }
102
103        let ble_nodes = self.ble_network.discover_nodes()?;
104        for node in ble_nodes {
105            if written >= out.len() {
106                return Err(MeshError::BufferTooSmall(
107                    "discovered node output buffer exhausted".to_string(),
108                ));
109            }
110            out[written] = Self::discovered_node_handle(
111                &node.node_id,
112                NodeType::Mobile,
113                NetworkInterface::Ble,
114                Self::ble_capability_tag(&node.capabilities),
115                node.rssi as f64,
116                None,
117            );
118            written += 1;
119        }
120
121        Ok(written)
122    }
123
124    /// Send message through mesh network
125    pub fn send_message(
126        &mut self,
127        destination: String,
128        payload: Vec<u8>,
129        priority: MessagePriority,
130    ) -> Result<String, MeshError> {
131        // Create message
132        let message_id = self.generate_message_id();
133        let message = StoredMessage {
134            message_id: message_id.clone(),
135            source: "local_node".to_string(),
136            destination: destination.clone(),
137            payload,
138            priority,
139            timestamp: Instant::now(),
140            ttl: Duration::from_secs(3600), // 1 hour TTL
141            delivery_attempts: 0,
142            status: MessageStatus::Pending,
143        };
144
145        // Store message
146        self.data_store.store_message(message.clone())?;
147
148        // Route message
149        self.route_message(&message)?;
150
151        Ok(message_id)
152    }
153
154    /// Route a transient payload without cloning it into the heap-backed persistence layer.
155    pub fn send_message_ephemeral(
156        &mut self,
157        destination: &str,
158        payload: &[u8],
159        priority: MessagePriority,
160    ) -> Result<u64, MeshError> {
161        let message_hash = self.generate_message_hash(destination, payload, priority);
162        self.route_payload(destination, payload, priority)?;
163        Ok(message_hash)
164    }
165
166    /// Receive message from mesh network
167    pub fn receive_message(&mut self, message: StoredMessage) -> Result<(), MeshError> {
168        // Store received message
169        self.data_store.store_message(message.clone())?;
170
171        // Update performance metrics
172        self.performance_monitor.update_receive_metrics(&message);
173
174        // Forward if not destined for this node
175        if message.destination != "local_node" && message.ttl.as_secs() > 0 {
176            self.route_message(&message)?;
177        }
178
179        Ok(())
180    }
181
182    /// Get network status
183    pub fn get_network_status(&self) -> NetworkStatus {
184        NetworkStatus {
185            acoustic_nodes: self.acoustic_network.get_node_count(),
186            ble_nodes: self.ble_network.get_node_count(),
187            total_nodes: self.acoustic_network.get_node_count() + self.ble_network.get_node_count(),
188            active_routes: self.mesh_router.get_route_count(),
189            pending_messages: self.data_store.get_pending_message_count(),
190            network_uptime: self.performance_monitor.get_uptime(),
191        }
192    }
193
194    /// Get performance statistics
195    pub fn get_performance_stats(&self) -> MeshGlobalMetrics {
196        self.performance_monitor.get_global_stats()
197    }
198
199    /// Optimize network performance
200    pub fn optimize_network(&mut self) -> Result<(), MeshError> {
201        // Optimize routing
202        self.mesh_router.optimize_routes()?;
203
204        // Optimize buffer management
205        self.data_store.optimize_buffers()?;
206
207        // Optimize discovery
208        self.acoustic_network.optimize_discovery()?;
209        self.ble_network.optimize_discovery()?;
210
211        Ok(())
212    }
213
214    // Internal methods
215
216    /// Route message through network
217    fn route_message(&mut self, message: &StoredMessage) -> Result<(), MeshError> {
218        // Determine best interface for routing
219        let interface =
220            self.select_best_interface_for_payload(message.payload.len(), message.priority)?;
221
222        match interface {
223            NetworkInterface::Acoustic => {
224                self.acoustic_network.send_message(message)?;
225            }
226            NetworkInterface::Ble => {
227                self.ble_network.send_message(message)?;
228            }
229            NetworkInterface::Hybrid => {
230                // Use both interfaces for redundancy
231                self.acoustic_network.send_message(message)?;
232                self.ble_network.send_message(message)?;
233            }
234        }
235
236        Ok(())
237    }
238
239    /// Route a transient payload without constructing a heap-backed message envelope.
240    fn route_payload(
241        &mut self,
242        destination: &str,
243        payload: &[u8],
244        priority: MessagePriority,
245    ) -> Result<(), MeshError> {
246        let interface = self.select_best_interface_for_payload(payload.len(), priority)?;
247
248        match interface {
249            NetworkInterface::Acoustic => {
250                self.acoustic_network
251                    .send_payload(destination, payload, priority)?
252            }
253            NetworkInterface::Ble => {
254                self.ble_network
255                    .send_payload(destination, payload, priority)?
256            }
257            NetworkInterface::Hybrid => {
258                self.acoustic_network
259                    .send_payload(destination, payload, priority)?;
260                self.ble_network
261                    .send_payload(destination, payload, priority)?;
262            }
263        }
264
265        Ok(())
266    }
267
268    /// Select best interface for message routing
269    fn select_best_interface_for_payload(
270        &self,
271        payload_len: usize,
272        priority: MessagePriority,
273    ) -> Result<NetworkInterface, MeshError> {
274        // Simple selection logic - in real implementation would be more sophisticated
275        if payload_len > 1000 {
276            // Large payload - use acoustic
277            Ok(NetworkInterface::Acoustic)
278        } else if priority == MessagePriority::Critical {
279            // Critical message - use both for redundancy
280            Ok(NetworkInterface::Hybrid)
281        } else {
282            // Default to BLE for small messages
283            Ok(NetworkInterface::Ble)
284        }
285    }
286
287    /// Generate unique message ID
288    fn generate_message_id(&self) -> String {
289        use std::sync::atomic::{AtomicU64, Ordering};
290        static COUNTER: AtomicU64 = AtomicU64::new(1);
291        format!("msg_{}", COUNTER.fetch_add(1, Ordering::SeqCst))
292    }
293
294    fn generate_message_hash(
295        &self,
296        destination: &str,
297        payload: &[u8],
298        priority: MessagePriority,
299    ) -> u64 {
300        let mut payload_hash = 0xcbf2_9ce4_8422_2325u64;
301        for byte in payload {
302            payload_hash ^= *byte as u64;
303            payload_hash = payload_hash.wrapping_mul(0x0000_0100_0000_01b3);
304        }
305        q_hash(destination) ^ payload_hash ^ (priority as u64)
306    }
307
308    fn discovered_node_handle(
309        node_id: &str,
310        node_type: NodeType,
311        interface: NetworkInterface,
312        capability_tag: u8,
313        signal_strength: f64,
314        location: Option<Location>,
315    ) -> DiscoveredNodeHandle {
316        DiscoveredNodeHandle {
317            node_id_hash: q_hash(node_id),
318            node_type,
319            interface,
320            capability_tag,
321            signal_strength,
322            location,
323        }
324    }
325
326    fn acoustic_capability_tag(capabilities: &AcousticCapabilities) -> u8 {
327        match capabilities.modulation {
328            ModulationType::FSK => 0x01,
329            ModulationType::PSK => 0x02,
330            ModulationType::OFDM => 0x03,
331            ModulationType::DSSS => 0x04,
332            ModulationType::Chirp => 0x05,
333        }
334    }
335
336    fn ble_capability_tag(capabilities: &BleCapabilities) -> u8 {
337        let mut tag = 0u8;
338        if capabilities
339            .features
340            .iter()
341            .any(|feature| *feature == BleFeature::ExtendedAdvertising)
342        {
343            tag |= 0x01;
344        }
345        if capabilities
346            .features
347            .iter()
348            .any(|feature| *feature == BleFeature::LE2MPHY)
349        {
350            tag |= 0x02;
351        }
352        if capabilities
353            .features
354            .iter()
355            .any(|feature| *feature == BleFeature::LEDataPacketLengthExtension)
356        {
357            tag |= 0x04;
358        }
359        if tag == 0 {
360            0x10
361        } else {
362            tag
363        }
364    }
365}