qualia_core_db/net/acoustic_ble_mesh/
manager.rs1use super::*;
6
7pub 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 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 pub fn initialize(&mut self) -> Result<(), MeshError> {
30 self.acoustic_network.initialize()?;
32
33 self.ble_network.initialize()?;
35
36 self.mesh_router.initialize()?;
38
39 self.data_store.initialize()?;
41
42 Ok(())
43 }
44
45 pub fn discover_nodes(&mut self) -> Result<Vec<DiscoveredNode>, MeshError> {
47 let mut discovered_nodes = Vec::new();
48
49 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 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, interface: NetworkInterface::Ble,
69 capabilities: NodeCapabilities::Ble(node.capabilities.clone()),
70 signal_strength: node.rssi as f64,
71 location: None, });
73 }
74
75 Ok(discovered_nodes)
76 }
77
78 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 pub fn send_message(
126 &mut self,
127 destination: String,
128 payload: Vec<u8>,
129 priority: MessagePriority,
130 ) -> Result<String, MeshError> {
131 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), delivery_attempts: 0,
142 status: MessageStatus::Pending,
143 };
144
145 self.data_store.store_message(message.clone())?;
147
148 self.route_message(&message)?;
150
151 Ok(message_id)
152 }
153
154 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 pub fn receive_message(&mut self, message: StoredMessage) -> Result<(), MeshError> {
168 self.data_store.store_message(message.clone())?;
170
171 self.performance_monitor.update_receive_metrics(&message);
173
174 if message.destination != "local_node" && message.ttl.as_secs() > 0 {
176 self.route_message(&message)?;
177 }
178
179 Ok(())
180 }
181
182 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 pub fn get_performance_stats(&self) -> MeshGlobalMetrics {
196 self.performance_monitor.get_global_stats()
197 }
198
199 pub fn optimize_network(&mut self) -> Result<(), MeshError> {
201 self.mesh_router.optimize_routes()?;
203
204 self.data_store.optimize_buffers()?;
206
207 self.acoustic_network.optimize_discovery()?;
209 self.ble_network.optimize_discovery()?;
210
211 Ok(())
212 }
213
214 fn route_message(&mut self, message: &StoredMessage) -> Result<(), MeshError> {
218 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 self.acoustic_network.send_message(message)?;
232 self.ble_network.send_message(message)?;
233 }
234 }
235
236 Ok(())
237 }
238
239 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 fn select_best_interface_for_payload(
270 &self,
271 payload_len: usize,
272 priority: MessagePriority,
273 ) -> Result<NetworkInterface, MeshError> {
274 if payload_len > 1000 {
276 Ok(NetworkInterface::Acoustic)
278 } else if priority == MessagePriority::Critical {
279 Ok(NetworkInterface::Hybrid)
281 } else {
282 Ok(NetworkInterface::Ble)
284 }
285 }
286
287 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}