Skip to main content

sedsnet/
discovery.rs

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    /// Switches discovery back to fast cadence and schedules an immediate announce.
44    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    /// Advances the cadence after sending an announce, backing off toward the slow interval.
50    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    /// Returns `true` when discovery should emit another announce at `now_ms`.
59    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    /// Network-variable data types available through this sender and the
165    /// routers behind it. This lets routing distinguish variables that share
166    /// a broad schema endpoint without falling back to link fanout.
167    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/// Returns `true` when the endpoint is reserved for discovery control traffic.
181#[inline]
182pub const fn is_discovery_endpoint(ep: DataEndpoint) -> bool {
183    matches!(ep, DataEndpoint::Discovery)
184}
185
186/// Returns `true` when the data type is a discovery control packet type.
187#[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
222/// Normalizes a topology-board list in place so it can be compared, exported, or encoded.
223pub 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
258/// Merges board topology views keyed by sender ID.
259pub 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
287/// Summarizes a board topology list into aggregated endpoint and time-source reachability.
288pub 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
302/// Builds a discovery announce packet advertising reachable non-discovery endpoints.
303pub 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
318/// Decodes a discovery announce packet into its advertised endpoints.
319pub 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
326/// Decodes a discovery announce payload into a sorted, deduplicated endpoint list.
327pub 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
346/// Builds a discovery packet advertising reachable time sync source identifiers.
347pub 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
653/// Decodes a discovery time sync source packet into source identifiers.
654pub 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
661/// Decodes a discovery time sync source payload into a sorted, deduplicated source list.
662pub 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
699/// Builds a discovery packet advertising the sender's current board/edge topology graph.
700pub 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
774/// Decodes a discovery topology packet into board-node records.
775pub 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
782/// Decodes a discovery topology payload into normalized board-node records.
783pub 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
905/// Builds a discovery packet containing the complete runtime schema snapshot.
906pub 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
913/// Elects the authoritative discovery/schema master for the current topology view.
914///
915/// Tie-breaks are deterministic:
916/// 1. Fewest unreachable boards.
917/// 2. Lowest maximum hop distance to any reachable board.
918/// 3. Lowest total hop distance across all reachable boards.
919/// 4. Lexicographically smallest sender ID.
920pub 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
988/// Builds a discovery schema packet from an explicit snapshot.
989pub 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
1051/// Builds a discovery schema packet from an owned snapshot.
1052pub 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        // Descriptions are documentation, not routing metadata. Embedded nodes
1069        // retain them in flash for local introspection, but transmitting them
1070        // duplicates several KiB into the constrained discovery queue.
1071        #[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
1117/// Decodes a discovery schema packet.
1118pub 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
1125/// Decodes a discovery schema payload into runtime definitions.
1126pub 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}