1use alloc::collections::BTreeMap;
2use alloc::string::{String, ToString};
3use alloc::vec;
4use alloc::vec::Vec;
5
6use crate::router::encode_slice_le;
7use crate::{
8 DataEndpoint, DataType, E2eEncryptionPolicy, MessageElement, TelemetryError, TelemetryResult,
9 config::{
10 OwnedDataTypeDefinition, OwnedEndpointDefinition, OwnedRuntimeSchemaSnapshot,
11 RuntimeSchemaSnapshot, e2e_encryption_policy_code, e2e_encryption_policy_from_code,
12 export_schema, message_class_code, message_class_from_code, message_data_type_code,
13 message_data_type_from_code, reliable_code, reliable_from_code,
14 },
15 packet::Packet,
16 try_enum_from_u32,
17};
18
19pub const DISCOVERY_ROUTE_TTL_MS: u64 = 30_000;
20pub const DISCOVERY_FAST_INTERVAL_MS: u64 = 250;
21pub const DISCOVERY_SLOW_INTERVAL_MS: u64 = 5_000;
22pub const DISCOVERY_SLOW_LINK_CAPACITY_BPS: u64 = 512;
23pub const DISCOVERY_SLOW_LINK_PING_INTERVAL_MS: u64 = 15_000;
24pub const DISCOVERY_SLOW_LINK_FULL_INTERVAL_MS: u64 = 120_000;
25pub const TIMESYNC_SLOW_LINK_MIN_INTERVAL_MS: u64 = 30_000;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct DiscoveryCadenceState {
29 pub current_interval_ms: u64,
30 pub next_announce_ms: u64,
31}
32
33impl Default for DiscoveryCadenceState {
34 fn default() -> Self {
35 Self {
36 current_interval_ms: DISCOVERY_FAST_INTERVAL_MS,
37 next_announce_ms: 0,
38 }
39 }
40}
41
42impl DiscoveryCadenceState {
43 pub fn on_topology_change(&mut self, now_ms: u64) {
45 self.current_interval_ms = DISCOVERY_FAST_INTERVAL_MS;
46 self.next_announce_ms = now_ms;
47 }
48
49 pub fn on_announce_sent(&mut self, now_ms: u64) {
51 self.next_announce_ms = now_ms.saturating_add(self.current_interval_ms);
52 self.current_interval_ms = core::cmp::min(
53 self.current_interval_ms.saturating_mul(2),
54 DISCOVERY_SLOW_INTERVAL_MS,
55 );
56 }
57
58 pub fn due(&self, now_ms: u64) -> bool {
60 now_ms >= self.next_announce_ms
61 }
62}
63
64#[derive(Debug, Clone, PartialEq, Eq)]
65pub struct TopologyBoardNode {
66 pub sender_id: String,
67 pub reachable_endpoints: Vec<DataEndpoint>,
68 pub reachable_timesync_sources: Vec<String>,
69 pub connections: Vec<String>,
70}
71
72#[derive(Debug, Clone, PartialEq, Eq)]
73pub struct TopologyLink {
74 pub source: String,
75 pub target: String,
76}
77
78#[derive(Debug, Clone, PartialEq, Eq)]
79pub struct TopologyAnnouncerRoute {
80 pub sender_id: String,
81 pub reachable_endpoints: Vec<DataEndpoint>,
82 pub reachable_timesync_sources: Vec<String>,
83 pub routers: Vec<TopologyBoardNode>,
84 pub last_seen_ms: u64,
85 pub age_ms: u64,
86}
87
88#[derive(Debug, Clone, PartialEq, Eq)]
89pub struct TopologySideRoute {
90 pub side_id: usize,
91 pub side_name: String,
92 pub reachable_endpoints: Vec<DataEndpoint>,
93 pub reachable_timesync_sources: Vec<String>,
94 pub announcers: Vec<TopologyAnnouncerRoute>,
95 pub last_seen_ms: u64,
96 pub age_ms: u64,
97}
98
99#[derive(Debug, Clone, PartialEq, Eq)]
100pub struct TopologySnapshot {
101 pub advertised_endpoints: Vec<DataEndpoint>,
102 pub advertised_timesync_sources: Vec<String>,
103 pub routers: Vec<TopologyBoardNode>,
104 pub links: Vec<TopologyLink>,
105 pub routes: Vec<TopologySideRoute>,
106 pub current_announce_interval_ms: u64,
107 pub next_announce_ms: u64,
108}
109
110#[derive(Debug, Clone, PartialEq, Eq)]
111pub struct ClientStatsSnapshot {
112 pub sender_id: String,
113 pub connected: bool,
114 pub side_ids: Vec<usize>,
115 pub side_names: Vec<String>,
116 pub last_seen_ms: Option<u64>,
117 pub age_ms: Option<u64>,
118 pub reachable_endpoints: Vec<DataEndpoint>,
119 pub reachable_timesync_sources: Vec<String>,
120 pub packets_sent: u64,
121 pub packets_received: u64,
122 pub bytes_sent: u64,
123 pub bytes_received: u64,
124}
125
126pub const LINK_CAPABILITY_HEADER_TEMPLATES: u32 = 0x0000_0001;
127pub const LINK_CAPABILITY_CHUNKING: u32 = 0x0000_0002;
128pub const LINK_CAPABILITY_RELIABILITY: u32 = 0x0000_0004;
129pub const LINK_CAPABILITY_CRYPTO: u32 = 0x0000_0008;
130pub const LINK_CAPABILITY_END_TO_END_RELIABILITY: u32 = 0x0000_0010;
131pub const LINK_CAPABILITY_OMIT_UNCHANGED_TIMESTAMPS: u32 = 0x0000_0020;
132
133pub const LINK_PROFILE_CANONICAL: u8 = 0;
134pub const LINK_PROFILE_TEMPLATE: u8 = 1;
135pub const LINK_PROFILE_IPV6_LIKE: u8 = 2;
136pub const LINK_PROFILE_IPV4_LIKE: u8 = 3;
137
138pub const ADDRESS_MODE_DYNAMIC: u8 = 0;
139pub const ADDRESS_MODE_REQUESTED: u8 = 1;
140pub const ADDRESS_MODE_STATIC: u8 = 2;
141pub const ADDRESS_STATE_REQUEST: u8 = 0;
142pub const ADDRESS_STATE_APPROVED: u8 = 1;
143
144#[derive(Debug, Clone, Copy, PartialEq, Eq)]
145pub struct LinkCapabilities {
146 pub version: u8,
147 pub flags: u32,
148 pub profile: u8,
149 pub max_frame_bytes: u32,
150 pub compact_header_target_bytes: u32,
151 pub max_side_transport_templates: u32,
152}
153
154#[derive(Debug, Clone, PartialEq, Eq)]
155pub struct AddressAdvertisement {
156 pub hostname: String,
157 pub address: u32,
158 pub requested_address: u32,
159 pub mode: u8,
160 pub state: u8,
161 pub birth_ms: u64,
162 pub owner_hash: u64,
163 pub reachable_endpoints: Vec<DataEndpoint>,
164 pub reachable_network_variables: Vec<DataType>,
168 pub reachable_timesync_sources: Vec<String>,
169 pub link_capabilities: LinkCapabilities,
170}
171
172#[inline]
173pub const fn is_router_control_endpoint(ep: DataEndpoint) -> bool {
174 matches!(
175 ep,
176 DataEndpoint::TelemetryError | DataEndpoint::TimeSync | DataEndpoint::Discovery
177 )
178}
179
180#[inline]
182pub const fn is_discovery_endpoint(ep: DataEndpoint) -> bool {
183 matches!(ep, DataEndpoint::Discovery)
184}
185
186#[inline]
188pub const fn is_discovery_type(ty: DataType) -> bool {
189 matches!(
190 ty,
191 DataType::DiscoveryAnnounce
192 | DataType::DiscoveryTimeSyncSources
193 | DataType::DiscoveryTopology
194 | DataType::DiscoverySchema
195 | DataType::DiscoveryTopologyRequest
196 | DataType::DiscoverySchemaRequest
197 | DataType::ManagedVariableRequest
198 | DataType::ManagedVariableValue
199 | DataType::DiscoveryLeave
200 | DataType::DiscoveryLinkCapabilities
201 | DataType::DiscoveryAddress
202 )
203}
204
205#[inline]
206pub const fn is_discovery_request_type(ty: DataType) -> bool {
207 matches!(
208 ty,
209 DataType::DiscoveryTopologyRequest
210 | DataType::DiscoverySchemaRequest
211 | DataType::ManagedVariableRequest
212 | DataType::DiscoveryLeave
213 | DataType::DiscoveryAddress
214 )
215}
216
217fn sort_dedup_strings(items: &mut Vec<String>) {
218 items.sort_unstable();
219 items.dedup();
220}
221
222pub fn normalize_topology_boards(boards: &mut Vec<TopologyBoardNode>) {
224 for board in boards.iter_mut() {
225 board
226 .reachable_endpoints
227 .retain(|ep| !is_router_control_endpoint(*ep));
228 board.reachable_endpoints.sort_unstable();
229 board.reachable_endpoints.dedup();
230 sort_dedup_strings(&mut board.reachable_timesync_sources);
231 board.connections.retain(|peer| peer != &board.sender_id);
232 sort_dedup_strings(&mut board.connections);
233 }
234 boards.sort_unstable_by(|a, b| a.sender_id.cmp(&b.sender_id));
235 boards.dedup_by(|a, b| a.sender_id == b.sender_id);
236}
237
238pub fn topology_links_from_boards(boards: &[TopologyBoardNode]) -> Vec<TopologyLink> {
239 let mut links = Vec::new();
240 for board in boards {
241 for peer in &board.connections {
242 if peer == &board.sender_id {
243 continue;
244 }
245 let (source, target) = if board.sender_id <= *peer {
246 (board.sender_id.clone(), peer.clone())
247 } else {
248 (peer.clone(), board.sender_id.clone())
249 };
250 links.push(TopologyLink { source, target });
251 }
252 }
253 links.sort_unstable_by(|a, b| (&a.source, &a.target).cmp(&(&b.source, &b.target)));
254 links.dedup_by(|a, b| a.source == b.source && a.target == b.target);
255 links
256}
257
258pub fn merge_topology_boards(dst: &mut Vec<TopologyBoardNode>, src: &[TopologyBoardNode]) {
260 let mut merged: BTreeMap<String, TopologyBoardNode> = dst
261 .iter()
262 .cloned()
263 .map(|board| (board.sender_id.clone(), board))
264 .collect();
265 for board in src {
266 let entry = merged
267 .entry(board.sender_id.clone())
268 .or_insert_with(|| TopologyBoardNode {
269 sender_id: board.sender_id.clone(),
270 reachable_endpoints: Vec::new(),
271 reachable_timesync_sources: Vec::new(),
272 connections: Vec::new(),
273 });
274 entry
275 .reachable_endpoints
276 .extend(board.reachable_endpoints.iter().copied());
277 entry
278 .reachable_timesync_sources
279 .extend(board.reachable_timesync_sources.iter().cloned());
280 entry.connections.extend(board.connections.iter().cloned());
281 }
282 let mut out: Vec<TopologyBoardNode> = merged.into_values().collect();
283 normalize_topology_boards(&mut out);
284 *dst = out;
285}
286
287pub fn summarize_topology_boards(boards: &[TopologyBoardNode]) -> (Vec<DataEndpoint>, Vec<String>) {
289 let mut reachable_endpoints = Vec::new();
290 let mut reachable_timesync_sources = Vec::new();
291 for board in boards {
292 reachable_endpoints.extend(board.reachable_endpoints.iter().copied());
293 reachable_timesync_sources.extend(board.reachable_timesync_sources.iter().cloned());
294 }
295 reachable_endpoints.sort_unstable();
296 reachable_endpoints.dedup();
297 reachable_endpoints.retain(|ep| !is_router_control_endpoint(*ep));
298 sort_dedup_strings(&mut reachable_timesync_sources);
299 (reachable_endpoints, reachable_timesync_sources)
300}
301
302pub fn build_discovery_announce(
304 sender: &str,
305 timestamp_ms: u64,
306 endpoints: &[DataEndpoint],
307) -> TelemetryResult<Packet> {
308 let payload_words: Vec<u32> = endpoints.iter().copied().map(|ep| ep.as_u32()).collect();
309 Packet::new(
310 DataType::DiscoveryAnnounce,
311 &[DataEndpoint::Discovery],
312 sender,
313 timestamp_ms,
314 encode_slice_le(payload_words.as_slice()),
315 )
316}
317
318pub fn decode_discovery_announce(pkt: &Packet) -> TelemetryResult<Vec<DataEndpoint>> {
320 if pkt.data_type() != DataType::DiscoveryAnnounce {
321 return Err(TelemetryError::InvalidType);
322 }
323 decode_discovery_payload(pkt.payload())
324}
325
326pub fn decode_discovery_payload(payload: &[u8]) -> TelemetryResult<Vec<DataEndpoint>> {
328 if !payload.len().is_multiple_of(4) {
329 return Err(TelemetryError::Unpack("discovery payload width"));
330 }
331
332 let mut endpoints = Vec::with_capacity(payload.len() / 4);
333 for chunk in payload.as_chunks::<4>().0 {
334 let raw = u32::from_le_bytes(*chunk);
335 let ep = try_enum_from_u32(raw).ok_or(TelemetryError::Unpack("bad discovery endpoint"))?;
336 if is_discovery_endpoint(ep) {
337 continue;
338 }
339 endpoints.push(ep);
340 }
341 endpoints.sort_unstable();
342 endpoints.dedup();
343 Ok(endpoints)
344}
345
346pub fn build_discovery_timesync_sources<S: AsRef<str>>(
348 sender: &str,
349 timestamp_ms: u64,
350 sources: &[S],
351) -> TelemetryResult<Packet> {
352 let mut payload = Vec::new();
353 let mut deduped: Vec<&str> = sources.iter().map(|s| s.as_ref()).collect();
354 deduped.sort_unstable();
355 deduped.dedup();
356
357 payload.extend_from_slice(&(deduped.len() as u32).to_le_bytes());
358 for source in deduped {
359 let bytes = source.as_bytes();
360 let len = u32::try_from(bytes.len())
361 .map_err(|_| TelemetryError::Pack("discovery source id too long"))?;
362 payload.extend_from_slice(&len.to_le_bytes());
363 payload.extend_from_slice(bytes);
364 }
365
366 Packet::new(
367 DataType::DiscoveryTimeSyncSources,
368 &[DataEndpoint::Discovery],
369 sender,
370 timestamp_ms,
371 payload.into(),
372 )
373}
374
375pub fn build_discovery_topology_request(
376 sender: &str,
377 timestamp_ms: u64,
378) -> TelemetryResult<Packet> {
379 Packet::new(
380 DataType::DiscoveryTopologyRequest,
381 &[DataEndpoint::Discovery],
382 sender,
383 timestamp_ms,
384 Vec::<u8>::new().into(),
385 )
386}
387
388pub fn build_discovery_schema_request(sender: &str, timestamp_ms: u64) -> TelemetryResult<Packet> {
389 Packet::new(
390 DataType::DiscoverySchemaRequest,
391 &[DataEndpoint::Discovery],
392 sender,
393 timestamp_ms,
394 Vec::<u8>::new().into(),
395 )
396}
397
398pub fn build_discovery_leave(sender: &str, timestamp_ms: u64) -> TelemetryResult<Packet> {
399 Packet::new(
400 DataType::DiscoveryLeave,
401 &[DataEndpoint::Discovery],
402 sender,
403 timestamp_ms,
404 Vec::<u8>::new().into(),
405 )
406}
407
408pub fn build_discovery_link_capabilities(
409 sender: &str,
410 timestamp_ms: u64,
411 capabilities: LinkCapabilities,
412) -> TelemetryResult<Packet> {
413 let mut payload = Vec::with_capacity(18);
414 payload.push(capabilities.version);
415 payload.extend_from_slice(&capabilities.flags.to_le_bytes());
416 payload.push(capabilities.profile);
417 payload.extend_from_slice(&capabilities.max_frame_bytes.to_le_bytes());
418 payload.extend_from_slice(&capabilities.compact_header_target_bytes.to_le_bytes());
419 payload.extend_from_slice(&capabilities.max_side_transport_templates.to_le_bytes());
420 Packet::new(
421 DataType::DiscoveryLinkCapabilities,
422 &[DataEndpoint::Discovery],
423 sender,
424 timestamp_ms,
425 payload.into(),
426 )
427}
428
429pub fn decode_discovery_link_capabilities(pkt: &Packet) -> TelemetryResult<LinkCapabilities> {
430 if pkt.data_type() != DataType::DiscoveryLinkCapabilities {
431 return Err(TelemetryError::InvalidType);
432 }
433 let payload = pkt.payload();
434 if payload.len() != 18 {
435 return Err(TelemetryError::Unpack("discovery link capabilities width"));
436 }
437 Ok(LinkCapabilities {
438 version: payload[0],
439 flags: u32::from_le_bytes(payload[1..5].try_into().expect("4-byte flags")),
440 profile: payload[5],
441 max_frame_bytes: u32::from_le_bytes(payload[6..10].try_into().expect("4-byte max frame")),
442 compact_header_target_bytes: u32::from_le_bytes(
443 payload[10..14].try_into().expect("4-byte target"),
444 ),
445 max_side_transport_templates: u32::from_le_bytes(
446 payload[14..18].try_into().expect("4-byte templates"),
447 ),
448 })
449}
450
451pub fn build_discovery_address(
452 sender: &str,
453 timestamp_ms: u64,
454 ad: &AddressAdvertisement,
455) -> TelemetryResult<Packet> {
456 let mut payload = Vec::new();
457 payload.push(2);
458 payload.push(ad.mode);
459 payload.push(ad.state);
460 payload.extend_from_slice(&ad.address.to_le_bytes());
461 payload.extend_from_slice(&ad.requested_address.to_le_bytes());
462 payload.extend_from_slice(&ad.birth_ms.to_le_bytes());
463 payload.extend_from_slice(&ad.owner_hash.to_le_bytes());
464 encode_string(&mut payload, &ad.hostname)?;
465 let mut endpoints = ad.reachable_endpoints.clone();
466 endpoints.retain(|ep| !is_discovery_endpoint(*ep));
467 endpoints.sort_unstable();
468 endpoints.dedup();
469 let endpoint_count = u32::try_from(endpoints.len())
470 .map_err(|_| TelemetryError::Pack("discovery address endpoint count"))?;
471 payload.extend_from_slice(&endpoint_count.to_le_bytes());
472 for ep in endpoints {
473 payload.extend_from_slice(&ep.as_u32().to_le_bytes());
474 }
475 let mut network_variables = ad.reachable_network_variables.clone();
476 network_variables.sort_unstable();
477 network_variables.dedup();
478 let network_variable_count = u32::try_from(network_variables.len())
479 .map_err(|_| TelemetryError::Pack("discovery address network variable count"))?;
480 payload.extend_from_slice(&network_variable_count.to_le_bytes());
481 for ty in network_variables {
482 payload.extend_from_slice(&ty.as_u32().to_le_bytes());
483 }
484 let mut sources = ad.reachable_timesync_sources.clone();
485 sort_dedup_strings(&mut sources);
486 let source_count = u32::try_from(sources.len())
487 .map_err(|_| TelemetryError::Pack("discovery address source count"))?;
488 payload.extend_from_slice(&source_count.to_le_bytes());
489 for source in sources {
490 encode_string(&mut payload, &source)?;
491 }
492 payload.push(ad.link_capabilities.version);
493 payload.extend_from_slice(&ad.link_capabilities.flags.to_le_bytes());
494 payload.push(ad.link_capabilities.profile);
495 payload.extend_from_slice(&ad.link_capabilities.max_frame_bytes.to_le_bytes());
496 payload.extend_from_slice(
497 &ad.link_capabilities
498 .compact_header_target_bytes
499 .to_le_bytes(),
500 );
501 payload.extend_from_slice(
502 &ad.link_capabilities
503 .max_side_transport_templates
504 .to_le_bytes(),
505 );
506 Packet::new(
507 DataType::DiscoveryAddress,
508 &[DataEndpoint::Discovery],
509 sender,
510 timestamp_ms,
511 payload.into(),
512 )
513}
514
515pub fn decode_discovery_address(pkt: &Packet) -> TelemetryResult<AddressAdvertisement> {
516 if pkt.data_type() != DataType::DiscoveryAddress {
517 return Err(TelemetryError::InvalidType);
518 }
519 let payload = pkt.payload();
520 let mut cursor = 0usize;
521 let version = read_u8(payload, &mut cursor, "discovery address version")?;
522 if version != 1 && version != 2 {
523 return Err(TelemetryError::Unpack("discovery address version"));
524 }
525 let mode = read_u8(payload, &mut cursor, "discovery address mode")?;
526 let state = read_u8(payload, &mut cursor, "discovery address state")?;
527 let address = read_u32(payload, &mut cursor, "discovery address current")?;
528 let requested_address = read_u32(payload, &mut cursor, "discovery address requested")?;
529 if payload.len().saturating_sub(cursor) < 8 {
530 return Err(TelemetryError::Unpack("discovery address birth"));
531 }
532 let birth_ms = u64::from_le_bytes(
533 payload[cursor..cursor + 8]
534 .try_into()
535 .expect("8-byte birth"),
536 );
537 cursor += 8;
538 if payload.len().saturating_sub(cursor) < 8 {
539 return Err(TelemetryError::Unpack("discovery address owner"));
540 }
541 let owner_hash = u64::from_le_bytes(
542 payload[cursor..cursor + 8]
543 .try_into()
544 .expect("8-byte owner"),
545 );
546 cursor += 8;
547 let hostname = decode_string(payload, &mut cursor, "discovery address hostname")?;
548 let endpoint_count =
549 read_u32(payload, &mut cursor, "discovery address endpoint count")? as usize;
550 let mut reachable_endpoints = Vec::with_capacity(endpoint_count);
551 for _ in 0..endpoint_count {
552 let raw = read_u32(payload, &mut cursor, "discovery address endpoint")?;
553 let ep = try_enum_from_u32(raw).ok_or(TelemetryError::Unpack("bad discovery endpoint"))?;
554 if !is_discovery_endpoint(ep) {
555 reachable_endpoints.push(ep);
556 }
557 }
558 reachable_endpoints.sort_unstable();
559 reachable_endpoints.dedup();
560 let mut reachable_network_variables = Vec::new();
561 if version >= 2 {
562 let count = read_u32(
563 payload,
564 &mut cursor,
565 "discovery address network variable count",
566 )? as usize;
567 reachable_network_variables.reserve(count);
568 for _ in 0..count {
569 let raw = read_u32(payload, &mut cursor, "discovery address network variable")?;
570 let ty = DataType::try_from_u32(raw)
571 .ok_or(TelemetryError::Unpack("bad discovery network variable"))?;
572 if !is_discovery_type(ty) {
573 reachable_network_variables.push(ty);
574 }
575 }
576 reachable_network_variables.sort_unstable();
577 reachable_network_variables.dedup();
578 }
579 let source_count = read_u32(payload, &mut cursor, "discovery address source count")? as usize;
580 let mut reachable_timesync_sources = Vec::with_capacity(source_count);
581 for _ in 0..source_count {
582 let source = decode_string(payload, &mut cursor, "discovery address source")?;
583 if !source.is_empty() {
584 reachable_timesync_sources.push(source);
585 }
586 }
587 sort_dedup_strings(&mut reachable_timesync_sources);
588 let version = read_u8(payload, &mut cursor, "discovery address link version")?;
589 let flags = read_u32(payload, &mut cursor, "discovery address link flags")?;
590 let profile = read_u8(payload, &mut cursor, "discovery address link profile")?;
591 let max_frame_bytes = read_u32(payload, &mut cursor, "discovery address link max frame")?;
592 let compact_header_target_bytes = read_u32(
593 payload,
594 &mut cursor,
595 "discovery address link compact target",
596 )?;
597 let max_side_transport_templates =
598 read_u32(payload, &mut cursor, "discovery address link templates")?;
599 if cursor != payload.len() {
600 return Err(TelemetryError::Unpack("discovery address trailing bytes"));
601 }
602 if hostname.is_empty() || address == 0 {
603 return Err(TelemetryError::Unpack("bad discovery address"));
604 }
605 Ok(AddressAdvertisement {
606 hostname,
607 address,
608 requested_address,
609 mode,
610 state,
611 birth_ms,
612 owner_hash,
613 reachable_endpoints,
614 reachable_network_variables,
615 reachable_timesync_sources,
616 link_capabilities: LinkCapabilities {
617 version,
618 flags,
619 profile,
620 max_frame_bytes,
621 compact_header_target_bytes,
622 max_side_transport_templates,
623 },
624 })
625}
626
627pub fn build_managed_variable_request(
628 sender: &str,
629 timestamp_ms: u64,
630 ty: DataType,
631) -> TelemetryResult<Packet> {
632 Packet::new(
633 DataType::ManagedVariableRequest,
634 &[DataEndpoint::Discovery],
635 sender,
636 timestamp_ms,
637 encode_slice_le(&[ty.as_u32()]),
638 )
639}
640
641pub fn decode_managed_variable_request(pkt: &Packet) -> TelemetryResult<DataType> {
642 if pkt.data_type() != DataType::ManagedVariableRequest {
643 return Err(TelemetryError::InvalidType);
644 }
645 let payload = pkt.payload();
646 if payload.len() != 4 {
647 return Err(TelemetryError::Unpack("managed variable request width"));
648 }
649 let raw = u32::from_le_bytes(payload.try_into().expect("4-byte payload"));
650 try_enum_from_u32(raw).ok_or(TelemetryError::Unpack("bad managed variable data type"))
651}
652
653pub fn decode_discovery_timesync_sources(pkt: &Packet) -> TelemetryResult<Vec<String>> {
655 if pkt.data_type() != DataType::DiscoveryTimeSyncSources {
656 return Err(TelemetryError::InvalidType);
657 }
658 decode_discovery_timesync_sources_payload(pkt.payload())
659}
660
661pub fn decode_discovery_timesync_sources_payload(payload: &[u8]) -> TelemetryResult<Vec<String>> {
663 if payload.len() < 4 {
664 return Err(TelemetryError::Unpack("discovery timesync source count"));
665 }
666
667 let count = u32::from_le_bytes(payload[..4].try_into().expect("4-byte count")) as usize;
668 let mut cursor = 4usize;
669 let mut out = Vec::with_capacity(count);
670
671 for _ in 0..count {
672 if payload.len().saturating_sub(cursor) < 4 {
673 return Err(TelemetryError::Unpack("discovery timesync source len"));
674 }
675 let len = u32::from_le_bytes(payload[cursor..cursor + 4].try_into().expect("4-byte len"))
676 as usize;
677 cursor += 4;
678 if payload.len().saturating_sub(cursor) < len {
679 return Err(TelemetryError::Unpack("discovery timesync source bytes"));
680 }
681 let raw = &payload[cursor..cursor + len];
682 cursor += len;
683 let source = core::str::from_utf8(raw)
684 .map_err(|_| TelemetryError::Unpack("discovery timesync source utf8"))?;
685 if !source.is_empty() {
686 out.push(source.to_string());
687 }
688 }
689
690 if cursor != payload.len() {
691 return Err(TelemetryError::Unpack("discovery timesync trailing bytes"));
692 }
693
694 out.sort_unstable();
695 out.dedup();
696 Ok(out)
697}
698
699pub fn build_discovery_topology(
701 sender: &str,
702 timestamp_ms: u64,
703 boards: &[TopologyBoardNode],
704) -> TelemetryResult<Packet> {
705 let mut payload = Vec::new();
706 let mut normalized = boards.to_vec();
707 normalize_topology_boards(&mut normalized);
708
709 payload.extend_from_slice(&(normalized.len() as u32).to_le_bytes());
710 for board in normalized {
711 let sender_bytes = board.sender_id.as_bytes();
712 let sender_len = u32::try_from(sender_bytes.len())
713 .map_err(|_| TelemetryError::Pack("discovery topology sender id too long"))?;
714 payload.extend_from_slice(&sender_len.to_le_bytes());
715 payload.extend_from_slice(sender_bytes);
716
717 payload.extend_from_slice(&(board.reachable_endpoints.len() as u32).to_le_bytes());
718 for ep in board.reachable_endpoints {
719 payload.extend_from_slice(&ep.as_u32().to_le_bytes());
720 }
721
722 payload.extend_from_slice(&(board.reachable_timesync_sources.len() as u32).to_le_bytes());
723 for source in board.reachable_timesync_sources {
724 let bytes = source.as_bytes();
725 let len = u32::try_from(bytes.len())
726 .map_err(|_| TelemetryError::Pack("discovery topology source id too long"))?;
727 payload.extend_from_slice(&len.to_le_bytes());
728 payload.extend_from_slice(bytes);
729 }
730
731 payload.extend_from_slice(&(board.connections.len() as u32).to_le_bytes());
732 for peer in board.connections {
733 let bytes = peer.as_bytes();
734 let len = u32::try_from(bytes.len())
735 .map_err(|_| TelemetryError::Pack("discovery topology connection id too long"))?;
736 payload.extend_from_slice(&len.to_le_bytes());
737 payload.extend_from_slice(bytes);
738 }
739 }
740
741 Packet::new(
742 DataType::DiscoveryTopology,
743 &[DataEndpoint::Discovery],
744 sender,
745 timestamp_ms,
746 payload.into(),
747 )
748}
749
750fn decode_string(
751 payload: &[u8],
752 cursor: &mut usize,
753 label: &'static str,
754) -> TelemetryResult<String> {
755 if payload.len().saturating_sub(*cursor) < 4 {
756 return Err(TelemetryError::Unpack(label));
757 }
758 let len = u32::from_le_bytes(
759 payload[*cursor..*cursor + 4]
760 .try_into()
761 .expect("4-byte len"),
762 ) as usize;
763 *cursor += 4;
764 if payload.len().saturating_sub(*cursor) < len {
765 return Err(TelemetryError::Unpack(label));
766 }
767 let raw = &payload[*cursor..*cursor + len];
768 *cursor += len;
769 core::str::from_utf8(raw)
770 .map(|s| s.to_string())
771 .map_err(|_| TelemetryError::Unpack(label))
772}
773
774pub fn decode_discovery_topology(pkt: &Packet) -> TelemetryResult<Vec<TopologyBoardNode>> {
776 if pkt.data_type() != DataType::DiscoveryTopology {
777 return Err(TelemetryError::InvalidType);
778 }
779 decode_discovery_topology_payload(pkt.payload())
780}
781
782pub fn decode_discovery_topology_payload(
784 payload: &[u8],
785) -> TelemetryResult<Vec<TopologyBoardNode>> {
786 if payload.len() < 4 {
787 return Err(TelemetryError::Unpack("discovery topology board count"));
788 }
789
790 let count = u32::from_le_bytes(payload[..4].try_into().expect("4-byte count")) as usize;
791 let mut cursor = 4usize;
792 let mut boards = Vec::with_capacity(count);
793
794 for _ in 0..count {
795 let sender_id = decode_string(payload, &mut cursor, "discovery topology sender id")?;
796
797 if payload.len().saturating_sub(cursor) < 4 {
798 return Err(TelemetryError::Unpack("discovery topology endpoint count"));
799 }
800 let endpoint_count = u32::from_le_bytes(
801 payload[cursor..cursor + 4]
802 .try_into()
803 .expect("4-byte count"),
804 ) as usize;
805 cursor += 4;
806 let mut reachable_endpoints = Vec::with_capacity(endpoint_count);
807 for _ in 0..endpoint_count {
808 if payload.len().saturating_sub(cursor) < 4 {
809 return Err(TelemetryError::Unpack("discovery topology endpoint"));
810 }
811 let raw =
812 u32::from_le_bytes(payload[cursor..cursor + 4].try_into().expect("4-byte ep"));
813 cursor += 4;
814 let ep =
815 try_enum_from_u32(raw).ok_or(TelemetryError::Unpack("bad discovery endpoint"))?;
816 if !is_discovery_endpoint(ep) {
817 reachable_endpoints.push(ep);
818 }
819 }
820
821 if payload.len().saturating_sub(cursor) < 4 {
822 return Err(TelemetryError::Unpack(
823 "discovery topology timesync source count",
824 ));
825 }
826 let source_count = u32::from_le_bytes(
827 payload[cursor..cursor + 4]
828 .try_into()
829 .expect("4-byte count"),
830 ) as usize;
831 cursor += 4;
832 let mut reachable_timesync_sources = Vec::with_capacity(source_count);
833 for _ in 0..source_count {
834 let source = decode_string(payload, &mut cursor, "discovery topology timesync source")?;
835 if !source.is_empty() {
836 reachable_timesync_sources.push(source);
837 }
838 }
839
840 if payload.len().saturating_sub(cursor) < 4 {
841 return Err(TelemetryError::Unpack(
842 "discovery topology connection count",
843 ));
844 }
845 let connection_count = u32::from_le_bytes(
846 payload[cursor..cursor + 4]
847 .try_into()
848 .expect("4-byte count"),
849 ) as usize;
850 cursor += 4;
851 let mut connections = Vec::with_capacity(connection_count);
852 for _ in 0..connection_count {
853 let peer = decode_string(payload, &mut cursor, "discovery topology connection")?;
854 if !peer.is_empty() {
855 connections.push(peer);
856 }
857 }
858
859 boards.push(TopologyBoardNode {
860 sender_id,
861 reachable_endpoints,
862 reachable_timesync_sources,
863 connections,
864 });
865 }
866
867 if cursor != payload.len() {
868 return Err(TelemetryError::Unpack("discovery topology trailing bytes"));
869 }
870
871 normalize_topology_boards(&mut boards);
872 Ok(boards)
873}
874
875fn encode_string(payload: &mut Vec<u8>, value: &str) -> TelemetryResult<()> {
876 let len = u32::try_from(value.len())
877 .map_err(|_| TelemetryError::Pack("discovery schema string too long"))?;
878 payload.extend_from_slice(&len.to_le_bytes());
879 payload.extend_from_slice(value.as_bytes());
880 Ok(())
881}
882
883fn read_u8(payload: &[u8], cursor: &mut usize, label: &'static str) -> TelemetryResult<u8> {
884 if payload.len().saturating_sub(*cursor) < 1 {
885 return Err(TelemetryError::Unpack(label));
886 }
887 let out = payload[*cursor];
888 *cursor += 1;
889 Ok(out)
890}
891
892fn read_u32(payload: &[u8], cursor: &mut usize, label: &'static str) -> TelemetryResult<u32> {
893 if payload.len().saturating_sub(*cursor) < 4 {
894 return Err(TelemetryError::Unpack(label));
895 }
896 let out = u32::from_le_bytes(
897 payload[*cursor..*cursor + 4]
898 .try_into()
899 .expect("4-byte u32"),
900 );
901 *cursor += 4;
902 Ok(out)
903}
904
905pub fn build_discovery_schema(sender: &str, timestamp_ms: u64) -> TelemetryResult<Packet> {
907 #[cfg(feature = "std")]
908 return build_discovery_schema_from_owned_snapshot(sender, timestamp_ms, export_schema());
909 #[cfg(not(feature = "std"))]
910 return build_discovery_schema_from_snapshot(sender, timestamp_ms, export_schema());
911}
912
913pub fn elect_discovery_master(local_sender: &str, boards: &[TopologyBoardNode]) -> String {
921 let mut nodes: BTreeMap<String, Vec<String>> = BTreeMap::new();
922 for board in boards {
923 nodes.entry(board.sender_id.clone()).or_default();
924 for peer in board.connections.iter() {
925 nodes
926 .entry(peer.clone())
927 .or_default()
928 .push(board.sender_id.clone());
929 nodes
930 .entry(board.sender_id.clone())
931 .or_default()
932 .push(peer.clone());
933 }
934 }
935 nodes.entry(local_sender.to_string()).or_default();
936 for peers in nodes.values_mut() {
937 peers.sort_unstable();
938 peers.dedup();
939 }
940
941 let mut best_sender = local_sender.to_string();
942 let mut best_unreachable = usize::MAX;
943 let mut best_max_hops = usize::MAX;
944 let mut best_total_hops = usize::MAX;
945
946 for sender in nodes.keys() {
947 let mut frontier = vec![sender.clone()];
948 let mut seen: BTreeMap<String, usize> = BTreeMap::new();
949 seen.insert(sender.clone(), 0);
950 let mut idx = 0;
951 while idx < frontier.len() {
952 let cur = frontier[idx].clone();
953 idx += 1;
954 let cur_dist = seen[&cur];
955 if let Some(peers) = nodes.get(&cur) {
956 for peer in peers {
957 if !seen.contains_key(peer) {
958 seen.insert(peer.clone(), cur_dist + 1);
959 frontier.push(peer.clone());
960 }
961 }
962 }
963 }
964
965 let unreachable = nodes.len().saturating_sub(seen.len());
966 let max_hops = seen.values().copied().max().unwrap_or(0);
967 let total_hops = seen.values().copied().sum();
968 let better = unreachable < best_unreachable
969 || (unreachable == best_unreachable && max_hops < best_max_hops)
970 || (unreachable == best_unreachable
971 && max_hops == best_max_hops
972 && total_hops < best_total_hops)
973 || (unreachable == best_unreachable
974 && max_hops == best_max_hops
975 && total_hops == best_total_hops
976 && sender < &best_sender);
977 if better {
978 best_sender = sender.clone();
979 best_unreachable = unreachable;
980 best_max_hops = max_hops;
981 best_total_hops = total_hops;
982 }
983 }
984
985 best_sender
986}
987
988pub fn build_discovery_schema_from_snapshot(
990 sender: &str,
991 timestamp_ms: u64,
992 mut schema: RuntimeSchemaSnapshot,
993) -> TelemetryResult<Packet> {
994 let mut payload = Vec::new();
995 payload.extend_from_slice(&3u32.to_le_bytes());
996
997 schema.endpoints.sort_unstable_by_key(|def| def.id.as_u32());
998 schema.types.sort_unstable_by_key(|def| def.id.as_u32());
999
1000 payload.extend_from_slice(&(schema.endpoints.len() as u32).to_le_bytes());
1001 for ep in schema.endpoints {
1002 payload.extend_from_slice(&ep.id.as_u32().to_le_bytes());
1003 payload.push(ep.link_local_only as u8);
1004 encode_string(&mut payload, ep.name)?;
1005 #[cfg(feature = "std")]
1006 encode_string(&mut payload, ep.description)?;
1007 #[cfg(not(feature = "std"))]
1008 encode_string(&mut payload, "")?;
1009 }
1010
1011 payload.extend_from_slice(&(schema.types.len() as u32).to_le_bytes());
1012 for ty in schema.types {
1013 payload.extend_from_slice(&ty.id.as_u32().to_le_bytes());
1014 encode_string(&mut payload, ty.name)?;
1015 #[cfg(feature = "std")]
1016 encode_string(&mut payload, ty.description)?;
1017 #[cfg(not(feature = "std"))]
1018 encode_string(&mut payload, "")?;
1019 match ty.element {
1020 MessageElement::Static(count, data_type, class) => {
1021 payload.push(0);
1022 payload.extend_from_slice(&(count as u32).to_le_bytes());
1023 payload.push(message_data_type_code(data_type));
1024 payload.push(message_class_code(class));
1025 }
1026 MessageElement::Dynamic(data_type, class) => {
1027 payload.push(1);
1028 payload.extend_from_slice(&0u32.to_le_bytes());
1029 payload.push(message_data_type_code(data_type));
1030 payload.push(message_class_code(class));
1031 }
1032 }
1033 payload.push(reliable_code(ty.reliable));
1034 payload.push(ty.priority);
1035 payload.push(e2e_encryption_policy_code(ty.e2e_encryption));
1036 payload.extend_from_slice(&(ty.endpoints.len() as u32).to_le_bytes());
1037 for ep in ty.endpoints {
1038 payload.extend_from_slice(&ep.as_u32().to_le_bytes());
1039 }
1040 }
1041
1042 Packet::new(
1043 DataType::DiscoverySchema,
1044 &[DataEndpoint::Discovery],
1045 sender,
1046 timestamp_ms,
1047 payload.into(),
1048 )
1049}
1050
1051pub fn build_discovery_schema_from_owned_snapshot(
1053 sender: &str,
1054 timestamp_ms: u64,
1055 mut schema: OwnedRuntimeSchemaSnapshot,
1056) -> TelemetryResult<Packet> {
1057 let mut payload = Vec::new();
1058 payload.extend_from_slice(&3u32.to_le_bytes());
1059
1060 schema.endpoints.sort_unstable_by_key(|def| def.id.as_u32());
1061 schema.types.sort_unstable_by_key(|def| def.id.as_u32());
1062
1063 payload.extend_from_slice(&(schema.endpoints.len() as u32).to_le_bytes());
1064 for ep in schema.endpoints {
1065 payload.extend_from_slice(&ep.id.as_u32().to_le_bytes());
1066 payload.push(ep.link_local_only as u8);
1067 encode_string(&mut payload, &ep.name)?;
1068 #[cfg(feature = "std")]
1072 encode_string(&mut payload, &ep.description)?;
1073 #[cfg(not(feature = "std"))]
1074 encode_string(&mut payload, "")?;
1075 }
1076
1077 payload.extend_from_slice(&(schema.types.len() as u32).to_le_bytes());
1078 for ty in schema.types {
1079 payload.extend_from_slice(&ty.id.as_u32().to_le_bytes());
1080 encode_string(&mut payload, &ty.name)?;
1081 #[cfg(feature = "std")]
1082 encode_string(&mut payload, &ty.description)?;
1083 #[cfg(not(feature = "std"))]
1084 encode_string(&mut payload, "")?;
1085 match ty.element {
1086 MessageElement::Static(count, data_type, class) => {
1087 payload.push(0);
1088 payload.extend_from_slice(&(count as u32).to_le_bytes());
1089 payload.push(message_data_type_code(data_type));
1090 payload.push(message_class_code(class));
1091 }
1092 MessageElement::Dynamic(data_type, class) => {
1093 payload.push(1);
1094 payload.extend_from_slice(&0u32.to_le_bytes());
1095 payload.push(message_data_type_code(data_type));
1096 payload.push(message_class_code(class));
1097 }
1098 }
1099 payload.push(reliable_code(ty.reliable));
1100 payload.push(ty.priority);
1101 payload.push(e2e_encryption_policy_code(ty.e2e_encryption));
1102 payload.extend_from_slice(&(ty.endpoints.len() as u32).to_le_bytes());
1103 for ep in ty.endpoints {
1104 payload.extend_from_slice(&ep.as_u32().to_le_bytes());
1105 }
1106 }
1107
1108 Packet::new(
1109 DataType::DiscoverySchema,
1110 &[DataEndpoint::Discovery],
1111 sender,
1112 timestamp_ms,
1113 payload.into(),
1114 )
1115}
1116
1117pub fn decode_discovery_schema(pkt: &Packet) -> TelemetryResult<OwnedRuntimeSchemaSnapshot> {
1119 if pkt.data_type() != DataType::DiscoverySchema {
1120 return Err(TelemetryError::InvalidType);
1121 }
1122 decode_discovery_schema_payload(pkt.payload())
1123}
1124
1125pub fn decode_discovery_schema_payload(
1127 payload: &[u8],
1128) -> TelemetryResult<OwnedRuntimeSchemaSnapshot> {
1129 let mut cursor = 0usize;
1130 let version = read_u32(payload, &mut cursor, "discovery schema version")?;
1131 if version != 1 && version != 2 && version != 3 {
1132 return Err(TelemetryError::Unpack("discovery schema version"));
1133 }
1134
1135 let endpoint_count =
1136 read_u32(payload, &mut cursor, "discovery schema endpoint count")? as usize;
1137 let mut endpoints = Vec::with_capacity(endpoint_count);
1138 for _ in 0..endpoint_count {
1139 let id = DataEndpoint(read_u32(
1140 payload,
1141 &mut cursor,
1142 "discovery schema endpoint id",
1143 )?);
1144 let link_local_only =
1145 read_u8(payload, &mut cursor, "discovery schema endpoint flags")? != 0;
1146 let name = decode_string(payload, &mut cursor, "discovery schema endpoint name")?;
1147 let description = if version >= 2 {
1148 decode_string(
1149 payload,
1150 &mut cursor,
1151 "discovery schema endpoint description",
1152 )?
1153 } else {
1154 String::new()
1155 };
1156 endpoints.push(OwnedEndpointDefinition {
1157 id,
1158 name,
1159 description,
1160 link_local_only,
1161 });
1162 }
1163
1164 let type_count = read_u32(payload, &mut cursor, "discovery schema type count")? as usize;
1165 let mut types = Vec::with_capacity(type_count);
1166 for _ in 0..type_count {
1167 let id = DataType(read_u32(payload, &mut cursor, "discovery schema type id")?);
1168 let name = decode_string(payload, &mut cursor, "discovery schema type name")?;
1169 let description = if version >= 2 {
1170 decode_string(payload, &mut cursor, "discovery schema type description")?
1171 } else {
1172 String::new()
1173 };
1174 let element_kind = read_u8(payload, &mut cursor, "discovery schema element kind")?;
1175 let count = read_u32(payload, &mut cursor, "discovery schema element count")? as usize;
1176 let data_type = message_data_type_from_code(read_u8(
1177 payload,
1178 &mut cursor,
1179 "discovery schema data type",
1180 )?)
1181 .ok_or(TelemetryError::Unpack("discovery schema data type"))?;
1182 let class =
1183 message_class_from_code(read_u8(payload, &mut cursor, "discovery schema class")?)
1184 .ok_or(TelemetryError::Unpack("discovery schema class"))?;
1185 let element = match element_kind {
1186 0 => MessageElement::Static(count, data_type, class),
1187 1 => MessageElement::Dynamic(data_type, class),
1188 _ => return Err(TelemetryError::Unpack("discovery schema element kind")),
1189 };
1190 let reliable =
1191 reliable_from_code(read_u8(payload, &mut cursor, "discovery schema reliable")?)
1192 .ok_or(TelemetryError::Unpack("discovery schema reliable"))?;
1193 let priority = read_u8(payload, &mut cursor, "discovery schema priority")?;
1194 let e2e_encryption = if version >= 3 {
1195 e2e_encryption_policy_from_code(read_u8(
1196 payload,
1197 &mut cursor,
1198 "discovery schema e2e cryptography",
1199 )?)
1200 .ok_or(TelemetryError::Unpack("discovery schema e2e cryptography"))?
1201 } else {
1202 E2eEncryptionPolicy::PreferOff
1203 };
1204 let endpoint_count =
1205 read_u32(payload, &mut cursor, "discovery schema type endpoint count")? as usize;
1206 let mut type_endpoints = Vec::with_capacity(endpoint_count);
1207 for _ in 0..endpoint_count {
1208 type_endpoints.push(DataEndpoint(read_u32(
1209 payload,
1210 &mut cursor,
1211 "discovery schema type endpoint",
1212 )?));
1213 }
1214 types.push(OwnedDataTypeDefinition {
1215 id,
1216 name,
1217 description,
1218 element,
1219 endpoints: type_endpoints,
1220 reliable,
1221 priority,
1222 e2e_encryption,
1223 });
1224 }
1225
1226 if cursor != payload.len() {
1227 return Err(TelemetryError::Unpack("discovery schema trailing bytes"));
1228 }
1229 Ok(OwnedRuntimeSchemaSnapshot { endpoints, types })
1230}