Skip to main content

sedsnet/
relay.rs

1use crate::config::{
2    RuntimeMemoryConfig, runtime_reliable_max_end_to_end_ack_cache,
3    runtime_reliable_max_end_to_end_pending, runtime_reliable_max_pending,
4    runtime_reliable_max_retries, runtime_reliable_max_return_routes,
5    runtime_reliable_retransmit_ms,
6};
7use crate::diagnostics::{
8    AdaptiveLinkStats, DiscoveryRuntimeStats, QueueRuntimeStats, ReliableRuntimeStats,
9    RouteModeStats, RouteOverrideStats, RoutePriorityStats, RouteWeightStats, RuntimeSideStats,
10    RuntimeStatsSnapshot, RuntimeTypeStats, TypedRouteOverrideStats,
11};
12#[cfg(feature = "discovery")]
13use crate::discovery::{
14    self, ClientStatsSnapshot, DISCOVERY_ROUTE_TTL_MS, DISCOVERY_SLOW_LINK_FULL_INTERVAL_MS,
15    DISCOVERY_SLOW_LINK_PING_INTERVAL_MS, DiscoveryCadenceState,
16    TIMESYNC_SLOW_LINK_MIN_INTERVAL_MS, TopologyAnnouncerRoute, TopologyBoardNode,
17    TopologySideRoute, TopologySnapshot,
18};
19#[cfg(feature = "discovery")]
20use crate::packet::sender_address_u32;
21use crate::packet::{Packet, hash_bytes_u64};
22use crate::queue::{BoundedDeque, ByteCost};
23use crate::wire_format;
24use crate::{is_reliable_type, message_meta, message_priority, reliable_mode};
25use crate::{
26    router::{Clock, CompactTimestampOmissionPolicy, SideTransportProfile},
27    {
28        RouteSelectionMode, TelemetryError, TelemetryResult,
29        lock::{ReentryGate, ReentryGuard, RouterMutex},
30    },
31};
32use alloc::borrow::ToOwned;
33use alloc::boxed::Box;
34use alloc::collections::{BTreeMap, BTreeSet, VecDeque};
35use alloc::string::{String, ToString};
36use alloc::{sync::Arc, vec, vec::Vec};
37use core::mem::size_of;
38
39const MAX_PRIORITY_BURST: u8 = 8;
40use crc32fast::Hasher as Crc32Hasher;
41
42/// Logical side index (CAN, UART, RADIO, etc.)
43pub type RelaySideId = usize;
44const SIDE_TRANSPORT_MAGIC: &[u8; 3] = b"SDT";
45const SIDE_TRANSPORT_KIND_FULL: u8 = 0x01;
46const SIDE_TRANSPORT_KIND_COMPACT: u8 = 0x02;
47const SIDE_TRANSPORT_KIND_CHUNK: u8 = 0x03;
48const SIDE_TRANSPORT_KIND_COMPACT_DELTA: u8 = 0x04;
49const SIDE_TRANSPORT_KIND_COMPACT_SAME_TIMESTAMP: u8 = 0x05;
50const SIDE_TRANSPORT_FLAG_PAYLOAD_COMPRESSED: u8 = 0x01;
51const SIDE_TRANSPORT_FLAG_WIRE_CONTRACT: u8 = 0x04;
52const SIDE_TRANSPORT_FLAG_PACKET_NONCE: u8 = 0x08;
53const SIDE_TRANSPORT_FLAG_ENDPOINT_BITMAP_PRESENT: u8 = 0x20;
54const SIDE_TRANSPORT_FLAG_COMPACT_RELIABLE_HEADER: u8 = 0x40;
55const CONTROL_SLOW_LINK_CAPACITY_BPS: u64 = 512;
56const SIDE_TRANSPORT_CHUNK_OVERHEAD: usize = 3 + 1 + 4 + 2 + 2 + wire_format::CRC32_BYTES;
57const SIDE_TRANSPORT_EP_BITMAP_BITS: usize = (crate::MAX_VALUE_DATA_ENDPOINT as usize) + 1;
58const SIDE_TRANSPORT_EP_BITMAP_BYTES: usize = SIDE_TRANSPORT_EP_BITMAP_BITS.div_ceil(8);
59pub const IPV4_LIKE_COMPACT_HEADER_TARGET_BYTES: usize =
60    crate::router::IPV4_LIKE_COMPACT_HEADER_TARGET_BYTES;
61pub const IPV6_LIKE_COMPACT_HEADER_TARGET_BYTES: usize =
62    crate::router::IPV6_LIKE_COMPACT_HEADER_TARGET_BYTES;
63pub const DEFAULT_SIDE_TRANSPORT_TEMPLATE_LIMIT: usize =
64    crate::router::DEFAULT_SIDE_TRANSPORT_TEMPLATE_LIMIT;
65/// Packet Handler function type
66type PacketHandlerFn = dyn Fn(&Packet) -> TelemetryResult<()> + Send + Sync + 'static;
67
68/// Packed Handler function type
69type PackedHandlerFn = dyn Fn(&[u8]) -> TelemetryResult<()> + Send + Sync + 'static;
70
71/// TX handler for a relay side: either packed or packet-based.
72#[derive(Clone)]
73pub enum RelayTxHandlerFn {
74    Packed(Arc<PackedHandlerFn>),
75    Packet(Arc<PacketHandlerFn>),
76}
77
78#[derive(Clone, Copy, Debug)]
79pub struct RelaySideOptions {
80    /// Enables the relay's per-link reliable transport layer on this side.
81    ///
82    /// When `true` and the side uses a packed TX handler, reliable schema traffic on this hop
83    /// gains relay-managed sequence numbers, ACKs, packet requests, and retransmits.
84    /// Packet-output sides still receive decoded packets rather than packed reliable framing.
85    pub reliable_enabled: bool,
86    /// Marks the side as eligible for link-local-only endpoints and discovery routes.
87    pub link_local_enabled: bool,
88    /// Allows packets received from this side to enter relay processing.
89    pub ingress_enabled: bool,
90    /// Allows the relay to transmit packets toward this side.
91    pub egress_enabled: bool,
92    /// Enables side-local header-template reuse for packed transport.
93    pub header_template_enabled: bool,
94    /// Maximum number of bytes to emit per packed TX callback.
95    ///
96    /// When non-zero and a packed frame would exceed this size, the relay
97    /// splits it into ordered side-transport chunks and reassembles them on RX
98    /// before normal relay processing. This is intended for fixed-size links
99    /// such as CAN or I2C while keeping the user API packet-oriented.
100    pub max_frame_bytes: usize,
101    /// Target total side-transport overhead for compact follow-up frames.
102    pub compact_header_target_bytes: usize,
103    /// Maximum side-local header templates retained for TX and RX dictionaries.
104    pub max_side_transport_templates: usize,
105    /// Omits the timestamp field from compact follow-up frames when it is unchanged.
106    pub omit_unchanged_compact_timestamps: bool,
107    /// Optional per-data-type timestamp omission policy for compact follow-up frames.
108    pub compact_timestamp_omission_types: CompactTimestampOmissionPolicy,
109    /// Declared compact-link profile for stats and future negotiation.
110    pub side_transport_profile: SideTransportProfile,
111}
112
113impl Default for RelaySideOptions {
114    fn default() -> Self {
115        Self {
116            reliable_enabled: false,
117            link_local_enabled: false,
118            ingress_enabled: true,
119            egress_enabled: true,
120            header_template_enabled: false,
121            max_frame_bytes: 0,
122            compact_header_target_bytes: 0,
123            max_side_transport_templates: DEFAULT_SIDE_TRANSPORT_TEMPLATE_LIMIT,
124            omit_unchanged_compact_timestamps: false,
125            compact_timestamp_omission_types: CompactTimestampOmissionPolicy::none(),
126            side_transport_profile: SideTransportProfile::Canonical,
127        }
128    }
129}
130
131impl RelaySideOptions {
132    /// Convenience preset for bounded packed-side transport.
133    ///
134    /// `max_frame_bytes == 0` leaves packed frames unbounded. Values greater
135    /// than zero enable relay-managed chunking/reassembly on this side.
136    #[inline]
137    pub fn with_small_packet_transport(mut self, max_frame_bytes: usize) -> Self {
138        self.header_template_enabled = true;
139        self.max_frame_bytes = max_frame_bytes;
140        self.compact_header_target_bytes = IPV6_LIKE_COMPACT_HEADER_TARGET_BYTES;
141        self.side_transport_profile = SideTransportProfile::Ipv6Like;
142        self
143    }
144
145    #[inline]
146    pub fn with_ipv4_like_compact_header_target(mut self) -> Self {
147        self.header_template_enabled = true;
148        self.compact_header_target_bytes = IPV4_LIKE_COMPACT_HEADER_TARGET_BYTES;
149        self.omit_unchanged_compact_timestamps = true;
150        self.side_transport_profile = SideTransportProfile::Ipv4Like;
151        self
152    }
153
154    #[inline]
155    pub fn with_ipv6_like_compact_header_target(mut self) -> Self {
156        self.header_template_enabled = true;
157        self.compact_header_target_bytes = IPV6_LIKE_COMPACT_HEADER_TARGET_BYTES;
158        self.side_transport_profile = SideTransportProfile::Ipv6Like;
159        self
160    }
161
162    #[inline]
163    pub fn with_template_transport(mut self) -> Self {
164        self.header_template_enabled = true;
165        self.side_transport_profile = SideTransportProfile::Template;
166        self
167    }
168
169    #[inline]
170    pub fn with_omitted_unchanged_compact_timestamps(mut self) -> Self {
171        self.header_template_enabled = true;
172        self.omit_unchanged_compact_timestamps = true;
173        self
174    }
175
176    #[inline]
177    pub fn with_omitted_unchanged_compact_timestamps_for_type(
178        mut self,
179        ty: crate::DataType,
180    ) -> Self {
181        self.header_template_enabled = true;
182        self.compact_timestamp_omission_types.insert(ty);
183        self
184    }
185
186    #[inline]
187    pub fn effective_transport_profile(self) -> SideTransportProfile {
188        if !self.header_template_enabled && self.max_frame_bytes == 0 {
189            SideTransportProfile::Canonical
190        } else if self.side_transport_profile == SideTransportProfile::Canonical {
191            SideTransportProfile::Template
192        } else {
193            self.side_transport_profile
194        }
195    }
196
197    #[cfg(feature = "discovery")]
198    #[inline]
199    pub fn link_capabilities(self) -> discovery::LinkCapabilities {
200        let mut flags = discovery::LINK_CAPABILITY_END_TO_END_RELIABILITY;
201        if self.header_template_enabled {
202            flags |= discovery::LINK_CAPABILITY_HEADER_TEMPLATES;
203        }
204        if self.max_frame_bytes != 0 {
205            flags |= discovery::LINK_CAPABILITY_CHUNKING;
206        }
207        if self.reliable_enabled {
208            flags |= discovery::LINK_CAPABILITY_RELIABILITY;
209        }
210        if self.omit_unchanged_compact_timestamps
211            || !self.compact_timestamp_omission_types.is_empty()
212        {
213            flags |= discovery::LINK_CAPABILITY_OMIT_UNCHANGED_TIMESTAMPS;
214        }
215        #[cfg(feature = "cryptography")]
216        {
217            flags |= discovery::LINK_CAPABILITY_CRYPTO;
218        }
219        discovery::LinkCapabilities {
220            version: 1,
221            flags,
222            profile: self.effective_transport_profile().discovery_code(),
223            max_frame_bytes: self.max_frame_bytes.min(u32::MAX as usize) as u32,
224            compact_header_target_bytes: self.compact_header_target_bytes.min(u32::MAX as usize)
225                as u32,
226            max_side_transport_templates: self.max_side_transport_templates.min(u32::MAX as usize)
227                as u32,
228        }
229    }
230}
231
232/// One side of the relay – a name + TX handler.
233#[derive(Clone)]
234pub struct RelaySide {
235    pub name: Arc<str>,
236    pub tx_handler: RelayTxHandlerFn,
237    pub opts: RelaySideOptions,
238}
239
240#[derive(Clone, Debug, PartialEq, Eq)]
241pub enum RelayItem {
242    Packed(Arc<[u8]>),
243    Packet(Arc<Packet>),
244}
245
246/// Item that was received by the relay from some side.
247#[derive(Clone, Debug, PartialEq, Eq)]
248struct RelayRxItem {
249    src: RelaySideId,
250    data: RelayItem,
251    priority: u8,
252}
253
254impl ByteCost for RelayRxItem {
255    fn byte_cost(&self) -> usize {
256        match &self.data {
257            RelayItem::Packed(bytes) => bytes.len(),
258            RelayItem::Packet(pkt) => pkt.byte_cost(),
259        }
260    }
261}
262
263/// Item that is ready to be transmitted out a destination side.
264#[derive(Clone, Debug, PartialEq, Eq)]
265struct RelayTxItem {
266    src: Option<RelaySideId>,
267    dst: RelaySideId,
268    data: RelayItem,
269    priority: u8,
270}
271
272#[derive(Clone, Debug, PartialEq, Eq)]
273struct RelayReplayItem {
274    dst: RelaySideId,
275    bytes: Arc<[u8]>,
276    priority: u8,
277}
278
279impl ByteCost for RelayTxItem {
280    fn byte_cost(&self) -> usize {
281        match &self.data {
282            RelayItem::Packed(bytes) => bytes.len(),
283            RelayItem::Packet(pkt) => pkt.byte_cost(),
284        }
285    }
286}
287
288impl ByteCost for RelayReplayItem {
289    fn byte_cost(&self) -> usize {
290        self.bytes.len()
291    }
292}
293
294// -------------------- Reliable delivery state (relay) --------------------
295
296#[derive(Debug, Clone)]
297struct ReliableTxState {
298    next_seq: u32,
299    sent_order: VecDeque<u32>,
300    sent: BTreeMap<u32, ReliableSent>,
301}
302
303#[derive(Debug, Clone)]
304struct ReliableSent {
305    bytes: Arc<[u8]>,
306    last_send_ms: u64,
307    retries: u32,
308    queued: bool,
309    partial_acked: bool,
310}
311
312#[derive(Debug, Clone)]
313struct ReliableRxState {
314    expected_seq: u32,
315    buffered: BTreeMap<u32, Arc<[u8]>>,
316}
317
318#[derive(Debug, Clone)]
319struct ReliableReturnRouteState {
320    side: RelaySideId,
321}
322
323#[cfg(feature = "discovery")]
324#[derive(Debug, Clone, Default, PartialEq, Eq)]
325struct DiscoverySenderState {
326    reachable: Vec<crate::DataEndpoint>,
327    reachable_timesync_sources: Vec<String>,
328    topology_boards: Vec<TopologyBoardNode>,
329    last_seen_ms: u64,
330}
331
332#[inline]
333fn is_internal_control_type(ty: crate::DataType) -> bool {
334    if matches!(
335        ty,
336        crate::DataType::ReliableAck
337            | crate::DataType::ReliablePartialAck
338            | crate::DataType::ReliablePacketRequest
339            | crate::DataType::P2pMessage
340    ) {
341        return true;
342    }
343
344    #[cfg(feature = "timesync")]
345    if matches!(
346        ty,
347        crate::DataType::TimeSyncAnnounce
348            | crate::DataType::TimeSyncRequest
349            | crate::DataType::TimeSyncResponse
350    ) {
351        return true;
352    }
353
354    #[cfg(feature = "discovery")]
355    if discovery::is_discovery_type(ty) {
356        return true;
357    }
358
359    let _ = ty;
360    false
361}
362
363#[cfg(feature = "discovery")]
364#[derive(Debug, Clone, Default, PartialEq, Eq)]
365struct DiscoverySideState {
366    reachable: Vec<crate::DataEndpoint>,
367    reachable_timesync_sources: Vec<String>,
368    last_seen_ms: u64,
369    announcers: BTreeMap<String, DiscoverySenderState>,
370}
371
372#[derive(Clone, Debug, Default, PartialEq, Eq)]
373struct SideChunkAssembly {
374    total: u16,
375    received: BTreeMap<u16, Arc<[u8]>>,
376}
377
378#[derive(Clone, Debug, Default)]
379struct SideTransportState {
380    tx_template_ids: BTreeMap<u64, u32>,
381    tx_templates: BTreeMap<u64, SideHeaderTemplate>,
382    tx_last_timestamps: BTreeMap<u32, u64>,
383    rx_templates: BTreeMap<u64, SideHeaderTemplate>,
384    rx_templates_by_id: BTreeMap<u32, SideHeaderTemplate>,
385    rx_last_timestamps: BTreeMap<u32, u64>,
386    rx_chunks: BTreeMap<u32, SideChunkAssembly>,
387    next_chunk_id: u32,
388    next_template_id: u32,
389}
390
391impl SideTransportState {
392    fn tx_template_count(&self) -> usize {
393        self.tx_template_ids.len()
394    }
395
396    fn rx_template_count(&self) -> usize {
397        self.rx_templates_by_id.len()
398    }
399
400    fn insert_tx_template(
401        &mut self,
402        template: SideHeaderTemplate,
403        template_id: u32,
404        max_templates: usize,
405    ) -> bool {
406        if max_templates == 0 {
407            return false;
408        }
409        let mut evicted = false;
410        if self.tx_template_ids.len() >= max_templates
411            && !self.tx_template_ids.contains_key(&template.hash)
412            && let Some(old_hash) = self.tx_template_ids.keys().next().copied()
413        {
414            if let Some(old_id) = self.tx_template_ids.remove(&old_hash) {
415                self.tx_last_timestamps.remove(&old_id);
416            }
417            self.tx_templates.remove(&old_hash);
418            evicted = true;
419        }
420        self.tx_template_ids.insert(template.hash, template_id);
421        self.tx_templates.insert(template.hash, template);
422        evicted
423    }
424
425    fn insert_rx_template(
426        &mut self,
427        template_id: u32,
428        template: SideHeaderTemplate,
429        max_templates: usize,
430    ) -> bool {
431        if max_templates == 0 {
432            return false;
433        }
434        let mut evicted = false;
435        if self.rx_templates_by_id.len() >= max_templates
436            && !self.rx_templates_by_id.contains_key(&template_id)
437            && let Some(old_id) = self.rx_templates_by_id.keys().next().copied()
438            && let Some(old_template) = self.rx_templates_by_id.remove(&old_id)
439        {
440            self.rx_templates.remove(&old_template.hash);
441            self.rx_last_timestamps.remove(&old_id);
442            evicted = true;
443        }
444        self.rx_templates_by_id
445            .insert(template_id, template.clone());
446        self.rx_templates.insert(template.hash, template);
447        evicted
448    }
449}
450
451#[derive(Clone, Debug, PartialEq, Eq)]
452struct SideHeaderTemplate {
453    hash: u64,
454    base_flags: u8,
455    prefix: Arc<[u8]>,
456    between: Arc<[u8]>,
457    reliable_flags: Option<u8>,
458    reliable_compact: bool,
459}
460
461type SideTemplateExtract<'a> = (
462    SideHeaderTemplate,
463    crate::DataType,
464    u8,
465    u64,
466    u64,
467    u16,
468    Option<(u32, u32)>,
469    &'a [u8],
470);
471
472#[derive(Clone, Copy, Debug, PartialEq, Eq)]
473enum SideCompactTimestampMode {
474    Absolute,
475    Delta,
476    Omitted,
477}
478
479#[derive(Debug, Clone, Default)]
480struct AdaptiveRouteStats {
481    estimated_bandwidth_bps: u64,
482    peak_bandwidth_bps: u64,
483    last_observed_ms: u64,
484    last_slow_observed_ms: u64,
485    sample_count: u64,
486    window_started_ms: u64,
487    window_bytes: u64,
488    peak_usage_bps: u64,
489}
490
491#[cfg(feature = "discovery")]
492#[derive(Debug, Clone, Default)]
493struct DiscoverySideThrottleState {
494    next_ping_ms: u64,
495    next_full_ms: u64,
496}
497
498#[cfg(all(feature = "discovery", feature = "timesync"))]
499#[derive(Debug, Clone, Default)]
500struct TimeSyncSideThrottleState {
501    next_allowed_ms: u64,
502}
503
504#[cfg(feature = "discovery")]
505#[derive(Debug, Clone, Copy, PartialEq, Eq)]
506enum DiscoveryAdvertiseLevel {
507    MinimalPing,
508    Full,
509}
510
511impl AdaptiveRouteStats {
512    #[inline]
513    fn observe(&mut self, bytes: usize, sample_bps: u64, now_ms: u64) {
514        self.estimated_bandwidth_bps = if self.estimated_bandwidth_bps == 0 {
515            sample_bps
516        } else if sample_bps >= self.estimated_bandwidth_bps {
517            self.estimated_bandwidth_bps
518                .saturating_mul(3)
519                .saturating_add(sample_bps.saturating_mul(5))
520                / 8
521        } else {
522            self.estimated_bandwidth_bps
523                .saturating_mul(7)
524                .saturating_add(sample_bps)
525                / 8
526        };
527        self.peak_bandwidth_bps = self.peak_bandwidth_bps.max(sample_bps);
528        self.last_observed_ms = now_ms;
529        if sample_bps > 0 && sample_bps <= CONTROL_SLOW_LINK_CAPACITY_BPS {
530            self.last_slow_observed_ms = now_ms;
531        }
532        self.sample_count = self.sample_count.saturating_add(1);
533        if self.window_started_ms == 0 || now_ms.saturating_sub(self.window_started_ms) > 1_000 {
534            self.window_started_ms = now_ms;
535            self.window_bytes = 0;
536        }
537        self.window_bytes = self.window_bytes.saturating_add(bytes as u64);
538        self.peak_usage_bps = self.peak_usage_bps.max(self.current_usage_bps(now_ms));
539    }
540
541    #[inline]
542    fn current_usage_bps(&self, now_ms: u64) -> u64 {
543        if self.window_started_ms == 0 {
544            return 0;
545        }
546        let elapsed_ms = now_ms.saturating_sub(self.window_started_ms).max(1);
547        (u128::from(self.window_bytes).saturating_mul(1000) / u128::from(elapsed_ms))
548            .min(u128::from(u64::MAX)) as u64
549    }
550
551    #[inline]
552    fn available_headroom_bps(&self, now_ms: u64) -> u64 {
553        let capacity = self
554            .estimated_bandwidth_bps
555            .max(self.peak_bandwidth_bps)
556            .max(1);
557        capacity.saturating_sub(self.current_usage_bps(now_ms))
558    }
559
560    #[inline]
561    fn weight(&self, now_ms: u64) -> u64 {
562        self.available_headroom_bps(now_ms).max(1)
563    }
564
565    #[inline]
566    fn snapshot(&self, now_ms: u64, auto_balancing_enabled: bool) -> AdaptiveLinkStats {
567        let current_usage_bps = self.current_usage_bps(now_ms);
568        let estimated_capacity_bps = self.estimated_bandwidth_bps.max(1);
569        let peak_capacity_bps = self.peak_bandwidth_bps.max(estimated_capacity_bps);
570        let available_headroom_bps = peak_capacity_bps.saturating_sub(current_usage_bps);
571        AdaptiveLinkStats {
572            auto_balancing_enabled,
573            estimated_capacity_bps,
574            peak_capacity_bps,
575            current_usage_bps,
576            peak_usage_bps: self.peak_usage_bps.max(current_usage_bps),
577            available_headroom_bps,
578            effective_weight: available_headroom_bps.max(1),
579            last_observed_ms: self.last_observed_ms,
580            sample_count: self.sample_count,
581        }
582    }
583}
584
585#[derive(Debug, Clone, Default)]
586struct TypeRuntimeStatsInner {
587    tx_packets: u64,
588    tx_bytes: u64,
589    rx_packets: u64,
590    rx_bytes: u64,
591    relayed_tx_packets: u64,
592    relayed_tx_bytes: u64,
593    relayed_rx_packets: u64,
594    relayed_rx_bytes: u64,
595    tx_retries: u64,
596    handler_failures: u64,
597}
598
599#[derive(Debug, Clone, Default)]
600struct SideRuntimeStatsInner {
601    tx_packets: u64,
602    tx_bytes: u64,
603    rx_packets: u64,
604    rx_bytes: u64,
605    relayed_tx_packets: u64,
606    relayed_tx_bytes: u64,
607    relayed_rx_packets: u64,
608    relayed_rx_bytes: u64,
609    tx_retries: u64,
610    tx_handler_failures: u64,
611    total_handler_retries: u64,
612    side_transport_full_frames: u64,
613    side_transport_compact_frames: u64,
614    side_transport_compact_delta_frames: u64,
615    side_transport_compact_omitted_timestamp_frames: u64,
616    side_transport_chunk_frames: u64,
617    side_transport_raw_bytes: u64,
618    side_transport_wire_bytes: u64,
619    side_transport_bytes_saved: u64,
620    side_transport_min_compact_overhead_bytes: Option<usize>,
621    side_transport_max_compact_overhead_bytes: Option<usize>,
622    side_transport_compact_target_misses: u64,
623    side_transport_template_evictions: u64,
624    data_types: BTreeMap<u32, TypeRuntimeStatsInner>,
625}
626
627impl SideRuntimeStatsInner {
628    fn type_stats_mut(&mut self, ty: crate::DataType) -> &mut TypeRuntimeStatsInner {
629        self.data_types.entry(ty.as_u32()).or_default()
630    }
631
632    fn note_tx(&mut self, ty: crate::DataType, bytes: usize, retries: usize) {
633        self.tx_packets = self.tx_packets.saturating_add(1);
634        self.tx_bytes = self.tx_bytes.saturating_add(bytes as u64);
635        self.relayed_tx_packets = self.relayed_tx_packets.saturating_add(1);
636        self.relayed_tx_bytes = self.relayed_tx_bytes.saturating_add(bytes as u64);
637        self.tx_retries = self.tx_retries.saturating_add(retries as u64);
638        self.total_handler_retries = self.total_handler_retries.saturating_add(retries as u64);
639        let stats = self.type_stats_mut(ty);
640        stats.tx_packets = stats.tx_packets.saturating_add(1);
641        stats.tx_bytes = stats.tx_bytes.saturating_add(bytes as u64);
642        stats.relayed_tx_packets = stats.relayed_tx_packets.saturating_add(1);
643        stats.relayed_tx_bytes = stats.relayed_tx_bytes.saturating_add(bytes as u64);
644        stats.tx_retries = stats.tx_retries.saturating_add(retries as u64);
645    }
646
647    fn note_rx(&mut self, ty: crate::DataType, bytes: usize) {
648        self.rx_packets = self.rx_packets.saturating_add(1);
649        self.rx_bytes = self.rx_bytes.saturating_add(bytes as u64);
650        self.relayed_rx_packets = self.relayed_rx_packets.saturating_add(1);
651        self.relayed_rx_bytes = self.relayed_rx_bytes.saturating_add(bytes as u64);
652        let stats = self.type_stats_mut(ty);
653        stats.rx_packets = stats.rx_packets.saturating_add(1);
654        stats.rx_bytes = stats.rx_bytes.saturating_add(bytes as u64);
655        stats.relayed_rx_packets = stats.relayed_rx_packets.saturating_add(1);
656        stats.relayed_rx_bytes = stats.relayed_rx_bytes.saturating_add(bytes as u64);
657    }
658
659    fn note_tx_failure(&mut self, ty: crate::DataType, retries: usize) {
660        self.tx_handler_failures = self.tx_handler_failures.saturating_add(1);
661        self.tx_retries = self.tx_retries.saturating_add(retries as u64);
662        self.total_handler_retries = self.total_handler_retries.saturating_add(retries as u64);
663        let stats = self.type_stats_mut(ty);
664        stats.handler_failures = stats.handler_failures.saturating_add(1);
665        stats.tx_retries = stats.tx_retries.saturating_add(retries as u64);
666    }
667
668    fn note_side_transport_full(&mut self, raw_bytes: usize, wire_bytes: usize) {
669        self.side_transport_full_frames = self.side_transport_full_frames.saturating_add(1);
670        self.note_side_transport_bytes(raw_bytes, wire_bytes);
671    }
672
673    fn note_side_transport_compact(
674        &mut self,
675        raw_bytes: usize,
676        wire_bytes: usize,
677        compact_overhead_bytes: usize,
678        used_timestamp_delta: bool,
679        omitted_timestamp: bool,
680    ) {
681        self.side_transport_compact_frames = self.side_transport_compact_frames.saturating_add(1);
682        if used_timestamp_delta {
683            self.side_transport_compact_delta_frames =
684                self.side_transport_compact_delta_frames.saturating_add(1);
685        }
686        if omitted_timestamp {
687            self.side_transport_compact_omitted_timestamp_frames = self
688                .side_transport_compact_omitted_timestamp_frames
689                .saturating_add(1);
690        }
691        self.note_side_transport_bytes(raw_bytes, wire_bytes);
692        self.side_transport_min_compact_overhead_bytes = Some(
693            self.side_transport_min_compact_overhead_bytes
694                .map_or(compact_overhead_bytes, |v| v.min(compact_overhead_bytes)),
695        );
696        self.side_transport_max_compact_overhead_bytes = Some(
697            self.side_transport_max_compact_overhead_bytes
698                .map_or(compact_overhead_bytes, |v| v.max(compact_overhead_bytes)),
699        );
700    }
701
702    fn note_side_transport_chunks(&mut self, chunks: usize) {
703        self.side_transport_chunk_frames = self
704            .side_transport_chunk_frames
705            .saturating_add(chunks as u64);
706    }
707
708    fn note_side_transport_bytes(&mut self, raw_bytes: usize, wire_bytes: usize) {
709        self.side_transport_raw_bytes = self
710            .side_transport_raw_bytes
711            .saturating_add(raw_bytes as u64);
712        self.side_transport_wire_bytes = self
713            .side_transport_wire_bytes
714            .saturating_add(wire_bytes as u64);
715        if raw_bytes > wire_bytes {
716            self.side_transport_bytes_saved = self
717                .side_transport_bytes_saved
718                .saturating_add((raw_bytes - wire_bytes) as u64);
719        }
720    }
721
722    fn note_side_transport_compact_target_miss(&mut self) {
723        self.side_transport_compact_target_misses =
724            self.side_transport_compact_target_misses.saturating_add(1);
725    }
726
727    fn note_side_transport_template_eviction(&mut self) {
728        self.side_transport_template_evictions =
729            self.side_transport_template_evictions.saturating_add(1);
730    }
731}
732
733#[derive(Clone, Debug)]
734pub struct RelayConfig {
735    sender: Arc<str>,
736    memory: RuntimeMemoryConfig,
737}
738
739impl RelayConfig {
740    pub fn new() -> Self {
741        Self::default()
742    }
743
744    pub fn with_sender<S: AsRef<str>>(mut self, sender: S) -> Self {
745        self.sender = Arc::from(sender.as_ref());
746        self
747    }
748
749    pub fn with_memory_config(mut self, memory: RuntimeMemoryConfig) -> TelemetryResult<Self> {
750        memory.validate()?;
751        self.memory = memory;
752        Ok(self)
753    }
754
755    fn sender(&self) -> Arc<str> {
756        self.sender.clone()
757    }
758
759    fn memory_config(&self) -> RuntimeMemoryConfig {
760        self.memory
761    }
762}
763
764impl Default for RelayConfig {
765    fn default() -> Self {
766        Self {
767            sender: Arc::from("RELAY"),
768            memory: RuntimeMemoryConfig::default(),
769        }
770    }
771}
772
773#[derive(Clone, Copy, Debug, PartialEq, Eq)]
774enum RouteSelectionOrigin {
775    Flood,
776    Discovered,
777}
778
779/// Internal state, protected by RouterMutex so all public methods can take &self.
780struct RelayInner {
781    memory: RuntimeMemoryConfig,
782    sides: Vec<Option<RelaySide>>,
783    route_overrides: BTreeMap<(Option<RelaySideId>, RelaySideId), bool>,
784    typed_route_overrides: BTreeMap<(Option<RelaySideId>, u32, RelaySideId), bool>,
785    route_weights: BTreeMap<(Option<RelaySideId>, RelaySideId), u32>,
786    route_priorities: BTreeMap<(Option<RelaySideId>, RelaySideId), u32>,
787    source_route_modes: BTreeMap<Option<RelaySideId>, RouteSelectionMode>,
788    route_selection_cursors: BTreeMap<Option<RelaySideId>, u64>,
789    adaptive_route_stats: BTreeMap<RelaySideId, AdaptiveRouteStats>,
790    side_runtime_stats: BTreeMap<RelaySideId, SideRuntimeStatsInner>,
791    side_transport: BTreeMap<RelaySideId, SideTransportState>,
792    rx_queue: BoundedDeque<RelayRxItem>,
793    tx_queue: BoundedDeque<RelayTxItem>,
794    tx_priority_burst: u8,
795    replay_queue: BoundedDeque<RelayReplayItem>,
796    recent_rx: BoundedDeque<u64>,
797    reliable_tx: BTreeMap<(RelaySideId, u32), ReliableTxState>,
798    reliable_rx: BTreeMap<(RelaySideId, u32), ReliableRxState>,
799    reliable_return_routes: BTreeMap<u64, ReliableReturnRouteState>,
800    reliable_return_route_order: VecDeque<u64>,
801    end_to_end_acked_destinations: BTreeMap<u64, BTreeSet<u64>>,
802    end_to_end_acked_destination_order: VecDeque<u64>,
803    total_handler_failures: u64,
804    total_handler_retries: u64,
805    #[cfg(feature = "discovery")]
806    discovery_routes: BTreeMap<RelaySideId, DiscoverySideState>,
807    #[cfg(feature = "discovery")]
808    discovery_cadence: DiscoveryCadenceState,
809    #[cfg(feature = "discovery")]
810    discovery_side_throttle: BTreeMap<RelaySideId, DiscoverySideThrottleState>,
811    #[cfg(all(feature = "discovery", feature = "timesync"))]
812    timesync_side_throttle: BTreeMap<RelaySideId, TimeSyncSideThrottleState>,
813}
814
815#[derive(Clone, Copy, Debug, PartialEq, Eq)]
816enum RelayQueueKind {
817    Rx,
818    Tx,
819    Replay,
820    Recent,
821    ReliableRxBuffer,
822    #[cfg(feature = "discovery")]
823    Discovery,
824}
825
826impl RelayInner {
827    #[cfg(feature = "discovery")]
828    fn topology_board_byte_cost(board: &TopologyBoardNode) -> usize {
829        board
830            .sender_id
831            .len()
832            .saturating_add(board.reachable_endpoints.len() * size_of::<crate::DataEndpoint>())
833            .saturating_add(
834                board
835                    .reachable_timesync_sources
836                    .iter()
837                    .map(|s| s.len())
838                    .sum::<usize>(),
839            )
840            .saturating_add(board.connections.iter().map(|s| s.len()).sum::<usize>())
841    }
842
843    #[cfg(feature = "discovery")]
844    fn discovery_sender_byte_cost(sender: &str, state: &DiscoverySenderState) -> usize {
845        sender
846            .len()
847            .saturating_add(state.reachable.len() * size_of::<crate::DataEndpoint>())
848            .saturating_add(
849                state
850                    .reachable_timesync_sources
851                    .iter()
852                    .map(|s| s.len())
853                    .sum::<usize>(),
854            )
855            .saturating_add(
856                state
857                    .topology_boards
858                    .iter()
859                    .map(Self::topology_board_byte_cost)
860                    .sum::<usize>(),
861            )
862            .saturating_add(size_of::<DiscoverySenderState>())
863    }
864
865    #[cfg(feature = "discovery")]
866    fn discovery_route_byte_cost(route: &DiscoverySideState) -> usize {
867        size_of::<DiscoverySideState>()
868            .saturating_add(route.reachable.len() * size_of::<crate::DataEndpoint>())
869            .saturating_add(
870                route
871                    .reachable_timesync_sources
872                    .iter()
873                    .map(|s| s.len())
874                    .sum::<usize>(),
875            )
876            .saturating_add(
877                route
878                    .announcers
879                    .iter()
880                    .map(|(sender, state)| Self::discovery_sender_byte_cost(sender, state))
881                    .sum::<usize>(),
882            )
883    }
884
885    #[cfg(feature = "discovery")]
886    fn discovery_bytes_used(&self) -> usize {
887        self.discovery_routes
888            .values()
889            .map(Self::discovery_route_byte_cost)
890            .sum()
891    }
892
893    #[inline]
894    fn reliable_rx_buffered_bytes(&self) -> usize {
895        self.reliable_rx
896            .values()
897            .flat_map(|state| state.buffered.values())
898            .map(|bytes| size_of::<Arc<[u8]>>() + bytes.len())
899            .sum()
900    }
901
902    #[inline]
903    fn shared_queue_bytes_used(&self) -> usize {
904        self.rx_queue
905            .bytes_used()
906            .saturating_add(self.tx_queue.bytes_used())
907            .saturating_add(self.replay_queue.bytes_used())
908            .saturating_add(self.recent_rx.max_bytes())
909            .saturating_add(self.reliable_rx_buffered_bytes())
910            .saturating_add(crate::config::schema_bytes_used())
911            .saturating_add({
912                #[cfg(feature = "discovery")]
913                {
914                    self.discovery_bytes_used()
915                }
916                #[cfg(not(feature = "discovery"))]
917                {
918                    0
919                }
920            })
921    }
922
923    fn reliable_rx_buffer_len(&self) -> usize {
924        self.reliable_rx
925            .values()
926            .map(|state| state.buffered.len())
927            .sum()
928    }
929
930    fn pop_reliable_rx_buffered(&mut self) -> Option<Arc<[u8]>> {
931        let key = self
932            .reliable_rx
933            .iter()
934            .find_map(|(key, state)| (!state.buffered.is_empty()).then_some(*key))?;
935        self.reliable_rx
936            .get_mut(&key)?
937            .buffered
938            .pop_first()
939            .map(|(_, v)| v)
940    }
941
942    fn pop_shared_queue_item(&mut self, preferred: RelayQueueKind) -> bool {
943        match preferred {
944            RelayQueueKind::Rx => self.rx_queue.pop_front().is_some(),
945            RelayQueueKind::Tx => self.tx_queue.pop_front().is_some(),
946            RelayQueueKind::Replay => self.replay_queue.pop_front().is_some(),
947            RelayQueueKind::Recent => self.recent_rx.pop_front().is_some(),
948            RelayQueueKind::ReliableRxBuffer => self.pop_reliable_rx_buffered().is_some(),
949            #[cfg(feature = "discovery")]
950            RelayQueueKind::Discovery => self.pop_discovery_route(),
951        }
952    }
953
954    #[cfg(feature = "discovery")]
955    fn pop_discovery_route(&mut self) -> bool {
956        let Some((&side, _)) = self
957            .discovery_routes
958            .iter()
959            .min_by_key(|(_, route)| route.last_seen_ms)
960        else {
961            return false;
962        };
963        self.discovery_routes.remove(&side);
964        Self::queue_budget_warning("topology route evicted because shared queue budget is full");
965        true
966    }
967
968    fn largest_shared_queue(&self) -> Option<RelayQueueKind> {
969        let candidates = [
970            (
971                RelayQueueKind::Rx,
972                self.rx_queue.bytes_used(),
973                self.rx_queue.len(),
974            ),
975            (
976                RelayQueueKind::Tx,
977                self.tx_queue.bytes_used(),
978                self.tx_queue.len(),
979            ),
980            (
981                RelayQueueKind::Replay,
982                self.replay_queue.bytes_used(),
983                self.replay_queue.len(),
984            ),
985            (RelayQueueKind::Recent, 0, 0),
986            (
987                RelayQueueKind::ReliableRxBuffer,
988                self.reliable_rx_buffered_bytes(),
989                self.reliable_rx_buffer_len(),
990            ),
991            #[cfg(feature = "discovery")]
992            (
993                RelayQueueKind::Discovery,
994                self.discovery_bytes_used(),
995                self.discovery_routes.len(),
996            ),
997        ];
998        candidates
999            .into_iter()
1000            .filter(|(_, bytes, len)| *bytes > 0 && *len > 0)
1001            .max_by_key(|(kind, bytes, _)| {
1002                (
1003                    *bytes,
1004                    if *kind == RelayQueueKind::ReliableRxBuffer {
1005                        0
1006                    } else {
1007                        1
1008                    },
1009                )
1010            })
1011            .map(|(kind, _, _)| kind)
1012    }
1013
1014    fn make_shared_queue_room(
1015        &mut self,
1016        incoming_cost: usize,
1017        preferred: RelayQueueKind,
1018    ) -> TelemetryResult<()> {
1019        if incoming_cost > self.memory.max_queue_budget {
1020            return Err(TelemetryError::PacketTooLarge(
1021                "Item exceeds maximum shared queue budget",
1022            ));
1023        }
1024
1025        while self.shared_queue_bytes_used().saturating_add(incoming_cost)
1026            > self.memory.max_queue_budget
1027        {
1028            let victim = self.largest_shared_queue().unwrap_or(preferred);
1029            if victim == RelayQueueKind::Discovery {
1030                Self::queue_budget_warning("topology data is using the largest queue budget share");
1031            }
1032            if !self.pop_shared_queue_item(victim) && !self.pop_shared_queue_item(preferred) {
1033                return Err(TelemetryError::PacketTooLarge(
1034                    "Item exceeds maximum shared queue budget",
1035                ));
1036            }
1037        }
1038
1039        Ok(())
1040    }
1041
1042    #[inline]
1043    fn queue_budget_warning(msg: &str) {
1044        #[cfg(feature = "std")]
1045        eprintln!("sedsnet queue budget warning: {msg}");
1046        let _ = msg;
1047    }
1048
1049    #[cfg(feature = "discovery")]
1050    fn fit_discovery_budget(&mut self) {
1051        while self.shared_queue_bytes_used() > self.memory.max_queue_budget {
1052            if !self.pop_discovery_route() {
1053                break;
1054            }
1055        }
1056    }
1057
1058    fn push_rx(&mut self, item: RelayRxItem) -> TelemetryResult<()> {
1059        self.make_shared_queue_room(item.byte_cost(), RelayQueueKind::Rx)?;
1060        self.rx_queue
1061            .push_back_prioritized(item, |queued| queued.priority)
1062    }
1063
1064    fn push_tx(&mut self, item: RelayTxItem) -> TelemetryResult<()> {
1065        self.make_shared_queue_room(item.byte_cost(), RelayQueueKind::Tx)?;
1066        self.tx_queue
1067            .push_back_prioritized(item, |queued| queued.priority)
1068    }
1069
1070    fn push_replay(&mut self, item: RelayReplayItem) -> TelemetryResult<()> {
1071        self.make_shared_queue_room(item.byte_cost(), RelayQueueKind::Replay)?;
1072        self.replay_queue
1073            .push_back_prioritized(item, |queued| queued.priority)
1074    }
1075
1076    fn push_recent_rx(&mut self, id: u64) -> TelemetryResult<()> {
1077        while self.recent_rx.len() >= self.memory.max_recent_rx_ids {
1078            let _ = self.recent_rx.pop_front();
1079        }
1080        self.make_shared_queue_room(0, RelayQueueKind::Recent)?;
1081        self.recent_rx.push_back(id)
1082    }
1083
1084    fn buffer_reliable_rx(
1085        &mut self,
1086        side: RelaySideId,
1087        ty: crate::DataType,
1088        seq: u32,
1089        bytes: Arc<[u8]>,
1090    ) -> TelemetryResult<()> {
1091        let key = Relay::reliable_key(side, ty);
1092        if self
1093            .reliable_rx
1094            .get(&key)
1095            .is_some_and(|state| state.buffered.contains_key(&seq))
1096        {
1097            return Ok(());
1098        }
1099        let cost = size_of::<Arc<[u8]>>() + bytes.len();
1100        self.make_shared_queue_room(cost, RelayQueueKind::ReliableRxBuffer)?;
1101        let rx_state = self
1102            .reliable_rx
1103            .entry(key)
1104            .or_insert_with(|| ReliableRxState {
1105                expected_seq: 1,
1106                buffered: BTreeMap::new(),
1107            });
1108        if rx_state.buffered.len() >= runtime_reliable_max_pending() {
1109            let _ = rx_state.buffered.pop_first();
1110        }
1111        rx_state.buffered.insert(seq, bytes);
1112        Ok(())
1113    }
1114}
1115
1116/// Relay that fans out packets from one side to all others.
1117/// - Supports both packed bytes and full Packet.
1118/// - Has RX & TX queues, like Router.
1119/// - Uses a Clock for the *_with_timeout APIs, same style as Router.
1120pub struct Relay {
1121    sender: RouterMutex<Arc<str>>,
1122    state: RouterMutex<RelayInner>,
1123    side_tx_gate: ReentryGate,
1124    clock: Box<dyn Clock + Send + Sync>,
1125}
1126
1127enum RemoteSidePlan {
1128    Target(Vec<RelaySideId>),
1129}
1130
1131impl Relay {
1132    const END_TO_END_ACK_SENDER: &'static str = "E2EACK";
1133    const END_TO_END_ACK_PREFIX: &'static str = "E2EACK:";
1134
1135    fn relay_item_priority(data: &RelayItem) -> TelemetryResult<u8> {
1136        let ty = match data {
1137            RelayItem::Packet(pkt) => pkt.data_type(),
1138            RelayItem::Packed(bytes) => wire_format::peek_envelope(bytes.as_ref())?.ty,
1139        };
1140        Ok(crate::scheduler_priority(ty))
1141    }
1142
1143    #[inline]
1144    fn is_side_tx_busy(err: &TelemetryError) -> bool {
1145        matches!(err, TelemetryError::Io("side tx busy"))
1146    }
1147
1148    fn process_replay_queue_item(&self) -> TelemetryResult<bool> {
1149        let Some(item) = ({
1150            let mut st = self.state.lock();
1151            st.replay_queue.pop_front()
1152        }) else {
1153            return Ok(false);
1154        };
1155        let frame = wire_format::peek_frame_info(item.bytes.as_ref())?;
1156        let ty = frame.envelope.ty;
1157        let Some(hdr) = frame.reliable else {
1158            return Ok(false);
1159        };
1160        {
1161            let mut st = self.state.lock();
1162            let tx_state = self.reliable_tx_state_mut(&mut st, item.dst, ty);
1163            if !tx_state.sent.contains_key(&hdr.seq) {
1164                return Ok(false);
1165            }
1166        }
1167        if let Err(e) = self.send_reliable_raw_to_side(item.dst, item.bytes.clone()) {
1168            if Self::is_side_tx_busy(&e) {
1169                let mut st = self.state.lock();
1170                st.push_replay(item)?;
1171                return Ok(false);
1172            }
1173            return Err(e);
1174        }
1175        let mut st = self.state.lock();
1176        let tx_state = self.reliable_tx_state_mut(&mut st, item.dst, ty);
1177        if let Some(sent) = tx_state.sent.get_mut(&hdr.seq) {
1178            sent.last_send_ms = self.clock.now_ms();
1179            sent.queued = false;
1180        }
1181        Ok(true)
1182    }
1183
1184    fn pop_ready_tx_item(
1185        &self,
1186    ) -> Option<(
1187        Option<RelaySideId>,
1188        RelaySideId,
1189        RelayTxHandlerFn,
1190        RelaySideOptions,
1191        RelayItem,
1192    )> {
1193        let mut st = self.state.lock();
1194        let mut priority_burst = st.tx_priority_burst;
1195        let item = st
1196            .tx_queue
1197            .pop_front_fair(&mut priority_burst, MAX_PRIORITY_BURST, |queued| {
1198                queued.priority
1199            });
1200        st.tx_priority_burst = priority_burst;
1201        if let Some(item) = item {
1202            let side = st.sides.get(item.dst).and_then(|side| side.clone());
1203            side.map(|s| (item.src, item.dst, s.tx_handler, s.opts, item.data))
1204        } else {
1205            None
1206        }
1207    }
1208
1209    fn send_tx_item(
1210        &self,
1211        src: Option<RelaySideId>,
1212        dst: RelaySideId,
1213        handler: RelayTxHandlerFn,
1214        opts: RelaySideOptions,
1215        data: RelayItem,
1216    ) -> TelemetryResult<bool> {
1217        let allowed = {
1218            let mut st = self.state.lock();
1219            let ty = match &data {
1220                RelayItem::Packet(pkt) => Some(pkt.data_type()),
1221                RelayItem::Packed(bytes) => Some(wire_format::peek_envelope(bytes.as_ref())?.ty),
1222            };
1223            let route_allowed = self.route_allowed_locked(&st, src, ty, dst);
1224            #[cfg(all(feature = "discovery", feature = "timesync"))]
1225            let timesync_allowed = ty
1226                .map(|ty| {
1227                    Self::timesync_allowed_for_side_locked(&mut st, dst, ty, self.clock.now_ms())
1228                })
1229                .unwrap_or(true);
1230            #[cfg(not(all(feature = "discovery", feature = "timesync")))]
1231            let timesync_allowed = true;
1232            route_allowed && timesync_allowed
1233        };
1234        if !allowed {
1235            return Ok(false);
1236        }
1237        if opts.reliable_enabled && matches!(handler, RelayTxHandlerFn::Packed(_)) {
1238            self.send_reliable_to_side(dst, data)?;
1239            Ok(true)
1240        } else if let Some(adjusted) = self.adjust_reliable_for_side(opts, data)? {
1241            self.call_tx_handler(dst, &handler, &adjusted)?;
1242            Ok(true)
1243        } else {
1244            Ok(false)
1245        }
1246    }
1247
1248    /// Create a new relay with the given clock.
1249    pub fn new(clock: Box<dyn Clock + Send + Sync>) -> Self {
1250        Self::new_with_config(RelayConfig::default(), clock)
1251    }
1252
1253    /// Create a new relay with explicit runtime configuration and clock.
1254    pub fn new_with_config(cfg: RelayConfig, clock: Box<dyn Clock + Send + Sync>) -> Self {
1255        let memory = cfg.memory_config();
1256        Self {
1257            sender: RouterMutex::new(cfg.sender()),
1258            state: RouterMutex::new(RelayInner {
1259                memory,
1260                sides: Vec::new(),
1261                route_overrides: BTreeMap::new(),
1262                typed_route_overrides: BTreeMap::new(),
1263                route_weights: BTreeMap::new(),
1264                route_priorities: BTreeMap::new(),
1265                source_route_modes: BTreeMap::new(),
1266                route_selection_cursors: BTreeMap::new(),
1267                adaptive_route_stats: BTreeMap::new(),
1268                side_runtime_stats: BTreeMap::new(),
1269                side_transport: BTreeMap::new(),
1270                rx_queue: BoundedDeque::new(
1271                    memory.max_queue_budget,
1272                    memory.starting_queue_size,
1273                    memory.queue_grow_step,
1274                ),
1275                tx_queue: BoundedDeque::new(
1276                    memory.max_queue_budget,
1277                    memory.starting_queue_size,
1278                    memory.queue_grow_step,
1279                ),
1280                tx_priority_burst: 0,
1281                replay_queue: BoundedDeque::new(
1282                    memory.max_queue_budget,
1283                    memory.starting_queue_size,
1284                    memory.queue_grow_step,
1285                ),
1286                recent_rx: BoundedDeque::new(
1287                    memory.recent_rx_queue_bytes(),
1288                    memory.recent_rx_queue_bytes(),
1289                    memory.queue_grow_step,
1290                ),
1291                reliable_tx: BTreeMap::new(),
1292                reliable_rx: BTreeMap::new(),
1293                reliable_return_routes: BTreeMap::new(),
1294                reliable_return_route_order: VecDeque::new(),
1295                end_to_end_acked_destinations: BTreeMap::new(),
1296                end_to_end_acked_destination_order: VecDeque::new(),
1297                total_handler_failures: 0,
1298                total_handler_retries: 0,
1299                #[cfg(feature = "discovery")]
1300                discovery_routes: BTreeMap::new(),
1301                #[cfg(feature = "discovery")]
1302                discovery_cadence: DiscoveryCadenceState::default(),
1303                #[cfg(feature = "discovery")]
1304                discovery_side_throttle: BTreeMap::new(),
1305                #[cfg(all(feature = "discovery", feature = "timesync"))]
1306                timesync_side_throttle: BTreeMap::new(),
1307            }),
1308            side_tx_gate: ReentryGate::new(),
1309            clock,
1310        }
1311    }
1312
1313    #[inline]
1314    fn sender_arc(&self) -> Arc<str> {
1315        self.sender.lock().clone()
1316    }
1317
1318    #[inline]
1319    pub fn sender(&self) -> Arc<str> {
1320        self.sender_arc()
1321    }
1322
1323    pub fn set_sender<S: AsRef<str>>(&self, sender: S) {
1324        *self.sender.lock() = Arc::from(sender.as_ref());
1325    }
1326
1327    #[inline]
1328    fn try_enter_side_tx(&self) -> Option<ReentryGuard<'_>> {
1329        self.side_tx_gate.try_enter()
1330    }
1331
1332    #[inline]
1333    fn side_tx_active(&self) -> bool {
1334        self.side_tx_gate.is_active()
1335    }
1336
1337    #[inline]
1338    fn side_ref(st: &RelayInner, side: RelaySideId) -> TelemetryResult<&RelaySide> {
1339        st.sides
1340            .get(side)
1341            .and_then(|side| side.as_ref())
1342            .ok_or(TelemetryError::HandlerError("relay: invalid side id"))
1343    }
1344
1345    fn note_side_tx_success(
1346        &self,
1347        side: RelaySideId,
1348        ty: crate::DataType,
1349        bytes: usize,
1350        attempts: usize,
1351    ) {
1352        let mut st = self.state.lock();
1353        let entry = st.side_runtime_stats.entry(side).or_default();
1354        entry.note_tx(ty, bytes, attempts.saturating_sub(1));
1355    }
1356
1357    fn note_side_tx_failure(&self, side: RelaySideId, ty: crate::DataType, attempts: usize) {
1358        let mut st = self.state.lock();
1359        st.total_handler_failures = st.total_handler_failures.saturating_add(1);
1360        st.total_handler_retries = st.total_handler_retries.saturating_add(attempts as u64);
1361        let entry = st.side_runtime_stats.entry(side).or_default();
1362        entry.note_tx_failure(ty, attempts);
1363    }
1364
1365    fn note_side_rx(&self, side: RelaySideId, ty: crate::DataType, bytes: usize) {
1366        let mut st = self.state.lock();
1367        let entry = st.side_runtime_stats.entry(side).or_default();
1368        entry.note_rx(ty, bytes);
1369    }
1370
1371    #[inline]
1372    fn ensure_side_ingress_enabled(&self, side: RelaySideId) -> TelemetryResult<()> {
1373        let st = self.state.lock();
1374        let side_ref = Self::side_ref(&st, side)?;
1375        if side_ref.opts.ingress_enabled {
1376            Ok(())
1377        } else {
1378            Err(TelemetryError::HandlerError(
1379                "relay: ingress disabled for side id",
1380            ))
1381        }
1382    }
1383
1384    #[inline]
1385    fn route_allowed_locked(
1386        &self,
1387        st: &RelayInner,
1388        src: Option<RelaySideId>,
1389        ty: Option<crate::DataType>,
1390        dst: RelaySideId,
1391    ) -> bool {
1392        let Ok(dst_side) = Self::side_ref(st, dst) else {
1393            return false;
1394        };
1395        if !dst_side.opts.egress_enabled {
1396            return false;
1397        }
1398        if let Some(src_id) = src {
1399            let Ok(src_side) = Self::side_ref(st, src_id) else {
1400                return false;
1401            };
1402            if !src_side.opts.ingress_enabled || src_id == dst {
1403                return false;
1404            }
1405        }
1406        let base_allowed = st.route_overrides.get(&(src, dst)).copied().unwrap_or(true);
1407        if !base_allowed {
1408            return false;
1409        }
1410
1411        let Some(ty) = ty else {
1412            return true;
1413        };
1414        if st
1415            .typed_route_overrides
1416            .keys()
1417            .any(|(typed_src, typed_ty, _)| *typed_src == src && *typed_ty == ty.as_u32())
1418        {
1419            return st
1420                .typed_route_overrides
1421                .get(&(src, ty.as_u32(), dst))
1422                .copied()
1423                .unwrap_or(false);
1424        }
1425        true
1426    }
1427
1428    fn has_typed_route_overrides_locked(
1429        st: &RelayInner,
1430        src: Option<RelaySideId>,
1431        ty: crate::DataType,
1432    ) -> bool {
1433        st.typed_route_overrides
1434            .keys()
1435            .any(|(typed_src, typed_ty, _)| *typed_src == src && *typed_ty == ty.as_u32())
1436    }
1437
1438    fn eligible_side_ids_locked(
1439        &self,
1440        st: &RelayInner,
1441        src: Option<RelaySideId>,
1442        ty: Option<crate::DataType>,
1443        restrict_link_local: bool,
1444    ) -> Vec<RelaySideId> {
1445        st.sides
1446            .iter()
1447            .enumerate()
1448            .filter_map(|(side_id, side)| {
1449                let side = side.as_ref()?;
1450                if restrict_link_local && !side.opts.link_local_enabled {
1451                    return None;
1452                }
1453                if self.route_allowed_locked(st, src, ty, side_id) {
1454                    Some(side_id)
1455                } else {
1456                    None
1457                }
1458            })
1459            .collect()
1460    }
1461
1462    fn apply_route_selection_locked(
1463        &self,
1464        st: &mut RelayInner,
1465        src: Option<RelaySideId>,
1466        mut sides: Vec<RelaySideId>,
1467        origin: RouteSelectionOrigin,
1468    ) -> Vec<RelaySideId> {
1469        if sides.len() <= 1 {
1470            return sides;
1471        }
1472
1473        let selection_mode = st.source_route_modes.get(&src).copied();
1474        if selection_mode.is_none() && origin == RouteSelectionOrigin::Discovered {
1475            return self.apply_adaptive_discovery_selection_locked(st, src, sides);
1476        }
1477
1478        match selection_mode.unwrap_or(RouteSelectionMode::Fanout) {
1479            RouteSelectionMode::Fanout => sides,
1480            RouteSelectionMode::Weighted => {
1481                sides.sort_unstable();
1482                let total_weight = sides.iter().fold(0_u64, |acc, side| {
1483                    acc + u64::from(st.route_weights.get(&(src, *side)).copied().unwrap_or(1))
1484                });
1485                if total_weight == 0 {
1486                    return Vec::new();
1487                }
1488                let cursor = st.route_selection_cursors.entry(src).or_insert(0);
1489                let pick = *cursor % total_weight;
1490                *cursor = cursor.wrapping_add(1);
1491                let mut remaining = pick;
1492                for side in sides {
1493                    let weight =
1494                        u64::from(st.route_weights.get(&(src, side)).copied().unwrap_or(1));
1495                    if remaining < weight {
1496                        return vec![side];
1497                    }
1498                    remaining -= weight;
1499                }
1500                Vec::new()
1501            }
1502            RouteSelectionMode::Failover => {
1503                sides.sort_by_key(|side| {
1504                    (
1505                        st.route_priorities.get(&(src, *side)).copied().unwrap_or(0),
1506                        *side,
1507                    )
1508                });
1509                sides.truncate(1);
1510                sides
1511            }
1512        }
1513    }
1514
1515    fn apply_adaptive_discovery_selection_locked(
1516        &self,
1517        st: &mut RelayInner,
1518        src: Option<RelaySideId>,
1519        mut sides: Vec<RelaySideId>,
1520    ) -> Vec<RelaySideId> {
1521        sides.sort_unstable();
1522        let mut unmeasured: Vec<_> = sides
1523            .iter()
1524            .copied()
1525            .filter(|side| !st.adaptive_route_stats.contains_key(side))
1526            .collect();
1527        if !unmeasured.is_empty() {
1528            let cursor = st.route_selection_cursors.entry(src).or_insert(0);
1529            let pick = (*cursor as usize) % unmeasured.len();
1530            *cursor = cursor.wrapping_add(1);
1531            return vec![unmeasured.swap_remove(pick)];
1532        }
1533
1534        let now_ms = self.clock.now_ms();
1535        let total_weight = sides.iter().fold(0_u64, |acc, side| {
1536            acc + st
1537                .adaptive_route_stats
1538                .get(side)
1539                .map(|stats| stats.weight(now_ms))
1540                .unwrap_or(1)
1541        });
1542        if total_weight == 0 {
1543            sides.truncate(1);
1544            return sides;
1545        }
1546
1547        let cursor = st.route_selection_cursors.entry(src).or_insert(0);
1548        let pick = *cursor % total_weight;
1549        *cursor = cursor.wrapping_add(1);
1550        let mut remaining = pick;
1551        for side in sides {
1552            let weight = st
1553                .adaptive_route_stats
1554                .get(&side)
1555                .map(|stats| stats.weight(now_ms))
1556                .unwrap_or(1);
1557            if remaining < weight {
1558                return vec![side];
1559            }
1560            remaining -= weight;
1561        }
1562        Vec::new()
1563    }
1564
1565    /// Prefer the discovery side with the shortest advertised topology path
1566    /// to the packet's destination. This prevents a route reflected through
1567    /// another router from competing equally with the direct route.
1568    #[cfg(feature = "discovery")]
1569    fn retain_shortest_discovery_routes_locked(
1570        st: &RelayInner,
1571        sides: &mut Vec<RelaySideId>,
1572        endpoints: &[crate::DataEndpoint],
1573        now_ms: u64,
1574    ) {
1575        if sides.len() <= 1 || endpoints.is_empty() {
1576            return;
1577        }
1578
1579        let route_distance = |side: RelaySideId| -> Option<usize> {
1580            let route = st.discovery_routes.get(&side)?;
1581            route
1582                .announcers
1583                .iter()
1584                .filter(|(_, state)| {
1585                    now_ms.saturating_sub(state.last_seen_ms) <= DISCOVERY_ROUTE_TTL_MS
1586                })
1587                .filter_map(|(announcer, state)| {
1588                    let targets: BTreeSet<&str> = state
1589                        .topology_boards
1590                        .iter()
1591                        .filter(|board| {
1592                            endpoints
1593                                .iter()
1594                                .any(|endpoint| board.reachable_endpoints.contains(endpoint))
1595                        })
1596                        .map(|board| board.sender_id.as_str())
1597                        .collect();
1598                    if targets.is_empty() {
1599                        return state
1600                            .reachable
1601                            .iter()
1602                            .any(|endpoint| endpoints.contains(endpoint))
1603                            .then_some(usize::MAX / 2);
1604                    }
1605
1606                    let mut distances: BTreeMap<&str, usize> = BTreeMap::new();
1607                    let mut pending = VecDeque::from([(announcer.as_str(), 0usize)]);
1608                    while let Some((sender, distance)) = pending.pop_front() {
1609                        if distances.contains_key(sender) {
1610                            continue;
1611                        }
1612                        distances.insert(sender, distance);
1613                        if targets.contains(sender) {
1614                            return Some(distance);
1615                        }
1616                        for board in &state.topology_boards {
1617                            if board.sender_id == sender {
1618                                for peer in &board.connections {
1619                                    if !distances.contains_key(peer.as_str()) {
1620                                        pending.push_back((peer.as_str(), distance + 1));
1621                                    }
1622                                }
1623                            } else if board.connections.iter().any(|peer| peer == sender)
1624                                && !distances.contains_key(board.sender_id.as_str())
1625                            {
1626                                pending.push_back((board.sender_id.as_str(), distance + 1));
1627                            }
1628                        }
1629                    }
1630                    None
1631                })
1632                .min()
1633        };
1634
1635        let distances: Vec<(RelaySideId, usize)> = sides
1636            .iter()
1637            .filter_map(|side| route_distance(*side).map(|distance| (*side, distance)))
1638            .collect();
1639        let Some(best) = distances.iter().map(|(_, distance)| *distance).min() else {
1640            return;
1641        };
1642        sides.retain(|side| {
1643            distances
1644                .iter()
1645                .any(|(candidate, distance)| candidate == side && *distance == best)
1646        });
1647    }
1648
1649    fn record_side_tx_sample(
1650        &self,
1651        side: RelaySideId,
1652        bytes: usize,
1653        started_ms: u64,
1654        ended_ms: u64,
1655    ) {
1656        let sample_ms = ended_ms.saturating_sub(started_ms).max(1);
1657        let sample_bps = ((bytes as u128).saturating_mul(1000) / u128::from(sample_ms))
1658            .min(u128::from(u64::MAX)) as u64;
1659        let mut st = self.state.lock();
1660        st.adaptive_route_stats
1661            .entry(side)
1662            .or_default()
1663            .observe(bytes, sample_bps, ended_ms);
1664    }
1665
1666    /// Seed adaptive route selection with a transport-measured link probe.
1667    ///
1668    /// Call this after a side-specific bring-up probe, or whenever the transport already knows the
1669    /// duration for a frame. The relay does not emit synthetic probe frames by itself.
1670    pub fn note_side_link_probe_sample(
1671        &self,
1672        side: RelaySideId,
1673        bytes: usize,
1674        duration_ms: u64,
1675    ) -> TelemetryResult<()> {
1676        {
1677            let st = self.state.lock();
1678            let _ = Self::side_ref(&st, side).map_err(|_| TelemetryError::BadArg)?;
1679        }
1680        let ended_ms = self.clock.now_ms();
1681        self.record_side_tx_sample(side, bytes, ended_ms.saturating_sub(duration_ms), ended_ms);
1682        Ok(())
1683    }
1684
1685    fn relay_item_wire_len(data: &RelayItem) -> TelemetryResult<usize> {
1686        match data {
1687            RelayItem::Packet(pkt) => Ok(wire_format::pack_packet(pkt).len()),
1688            RelayItem::Packed(bytes) => Ok(bytes.len()),
1689        }
1690    }
1691
1692    #[inline]
1693    fn decode_end_to_end_reliable_ack(payload: &[u8]) -> TelemetryResult<u64> {
1694        if payload.len() != 8 {
1695            return Err(TelemetryError::Unpack("bad reliable e2e ack payload"));
1696        }
1697        Ok(u64::from_le_bytes(payload[0..8].try_into().unwrap()))
1698    }
1699
1700    #[inline]
1701    fn is_end_to_end_ack_sender(sender: &str) -> bool {
1702        sender == Self::END_TO_END_ACK_SENDER || sender.starts_with(Self::END_TO_END_ACK_PREFIX)
1703    }
1704
1705    #[inline]
1706    fn sender_hash(sender: &str) -> u64 {
1707        if let Some(address) = sender
1708            .strip_prefix("@addr:")
1709            .and_then(|value| value.parse::<u32>().ok())
1710        {
1711            return u64::from(address);
1712        }
1713        hash_bytes_u64(0x517C_C1B7_2722_0A95, sender.as_bytes())
1714    }
1715
1716    #[cfg(feature = "discovery")]
1717    fn canonical_sender_locked(st: &RelayInner, sender: &str) -> String {
1718        let Some(address) = sender
1719            .strip_prefix("@addr:")
1720            .and_then(|value| value.parse::<u32>().ok())
1721        else {
1722            return sender.to_string();
1723        };
1724        st.discovery_routes
1725            .values()
1726            .flat_map(|route| route.announcers.iter())
1727            .flat_map(|(announcer, state)| {
1728                core::iter::once(announcer.as_str()).chain(
1729                    state
1730                        .topology_boards
1731                        .iter()
1732                        .map(|board| board.sender_id.as_str()),
1733                )
1734            })
1735            .find(|candidate| sender_address_u32(candidate) == address)
1736            .map(ToString::to_string)
1737            .unwrap_or_else(|| sender.to_string())
1738    }
1739
1740    fn decode_end_to_end_ack_sender_hash(sender: &str) -> Option<u64> {
1741        sender
1742            .strip_prefix(Self::END_TO_END_ACK_PREFIX)
1743            .filter(|sender| !sender.is_empty())
1744            .map(Self::sender_hash)
1745    }
1746
1747    #[cfg(feature = "discovery")]
1748    fn is_end_to_end_destination_sender(&self, sender: &str) -> bool {
1749        sender != self.sender_arc().as_ref() && !Self::is_end_to_end_ack_sender(sender)
1750    }
1751
1752    /// Extract the logical packet ID targeted by an end-to-end reliable ACK item.
1753    ///
1754    /// Relay queues can hold either decoded packets or packed frames. This
1755    /// helper normalizes both forms so relay ACK-routing logic can treat them
1756    /// uniformly.
1757    ///
1758    /// Only relay-visible end-to-end `ReliableAck` packets qualify here.
1759    /// Unrelated traffic returns `Ok(None)`.
1760    fn reliable_control_target_packet_id(data: &RelayItem) -> TelemetryResult<Option<u64>> {
1761        match data {
1762            RelayItem::Packet(pkt) => {
1763                if pkt.data_type() != crate::DataType::ReliableAck
1764                    || !Self::is_end_to_end_ack_sender(pkt.sender())
1765                {
1766                    return Ok(None);
1767                }
1768                Self::decode_end_to_end_reliable_ack(pkt.payload()).map(Some)
1769            }
1770            RelayItem::Packed(bytes) => {
1771                if wire_format::peek_frame_info(bytes.as_ref())
1772                    .ok()
1773                    .is_some_and(|frame| frame.ack_only())
1774                {
1775                    return Ok(None);
1776                }
1777                let pkt = wire_format::unpack_packet(bytes.as_ref())?;
1778                if pkt.data_type() != crate::DataType::ReliableAck
1779                    || !Self::is_end_to_end_ack_sender(pkt.sender())
1780                {
1781                    return Ok(None);
1782                }
1783                Self::decode_end_to_end_reliable_ack(pkt.payload()).map(Some)
1784            }
1785        }
1786    }
1787
1788    fn note_reliable_return_route(&self, side: RelaySideId, packet_id: u64) {
1789        let mut st = self.state.lock();
1790        Self::remember_reliable_return_route_locked(&mut st, packet_id);
1791        st.reliable_return_routes
1792            .insert(packet_id, ReliableReturnRouteState { side });
1793    }
1794
1795    /// Refresh or insert `packet_id` in the bounded reliable return-route cache.
1796    ///
1797    /// The relay uses this cache to route end-to-end acknowledgements back
1798    /// toward the source side that most recently forwarded the corresponding
1799    /// reliable data packet.
1800    fn remember_reliable_return_route_locked(st: &mut RelayInner, packet_id: u64) {
1801        let cap = runtime_reliable_max_return_routes().max(1);
1802        st.reliable_return_route_order
1803            .retain(|id| st.reliable_return_routes.contains_key(id) && *id != packet_id);
1804        while st.reliable_return_route_order.len() >= cap {
1805            if let Some(oldest) = st.reliable_return_route_order.pop_front() {
1806                st.reliable_return_routes.remove(&oldest);
1807            } else {
1808                break;
1809            }
1810        }
1811        st.reliable_return_route_order.push_back(packet_id);
1812    }
1813
1814    fn note_end_to_end_acked_destination_locked(
1815        st: &mut RelayInner,
1816        packet_id: u64,
1817        sender_hash: u64,
1818    ) {
1819        let entry_cap = runtime_reliable_max_end_to_end_ack_cache().max(1);
1820        st.end_to_end_acked_destination_order
1821            .retain(|id| st.end_to_end_acked_destinations.contains_key(id) && *id != packet_id);
1822        while st.end_to_end_acked_destination_order.len() >= entry_cap {
1823            if let Some(oldest) = st.end_to_end_acked_destination_order.pop_front() {
1824                st.end_to_end_acked_destinations.remove(&oldest);
1825            } else {
1826                break;
1827            }
1828        }
1829        st.end_to_end_acked_destination_order.push_back(packet_id);
1830
1831        let acked = st
1832            .end_to_end_acked_destinations
1833            .entry(packet_id)
1834            .or_default();
1835        let sender_cap = runtime_reliable_max_end_to_end_pending().max(1);
1836        if acked.len() < sender_cap || acked.contains(&sender_hash) {
1837            acked.insert(sender_hash);
1838        }
1839    }
1840
1841    #[inline]
1842    fn reliable_key(side: RelaySideId, ty: crate::DataType) -> (RelaySideId, u32) {
1843        (side, ty.as_u32())
1844    }
1845
1846    fn reliable_tx_state_mut<'a>(
1847        &'a self,
1848        st: &'a mut RelayInner,
1849        side: RelaySideId,
1850        ty: crate::DataType,
1851    ) -> &'a mut ReliableTxState {
1852        let key = Self::reliable_key(side, ty);
1853        st.reliable_tx
1854            .entry(key)
1855            .or_insert_with(|| ReliableTxState {
1856                next_seq: 1,
1857                sent_order: VecDeque::new(),
1858                sent: BTreeMap::new(),
1859            })
1860    }
1861
1862    fn reliable_rx_state_mut<'a>(
1863        &'a self,
1864        st: &'a mut RelayInner,
1865        side: RelaySideId,
1866        ty: crate::DataType,
1867    ) -> &'a mut ReliableRxState {
1868        let key = Self::reliable_key(side, ty);
1869        st.reliable_rx
1870            .entry(key)
1871            .or_insert_with(|| ReliableRxState {
1872                expected_seq: 1,
1873                buffered: BTreeMap::new(),
1874            })
1875    }
1876
1877    fn handle_reliable_ack(&self, side: RelaySideId, ty: crate::DataType, ack: u32) {
1878        let mut st = self.state.lock();
1879        let tx_state = self.reliable_tx_state_mut(&mut st, side, ty);
1880        if matches!(reliable_mode(ty), crate::ReliableMode::Unordered) {
1881            tx_state.sent.remove(&ack);
1882            tx_state.sent_order.retain(|seq| *seq != ack);
1883            return;
1884        }
1885
1886        while let Some(seq) = tx_state.sent_order.front().copied() {
1887            if seq > ack {
1888                break;
1889            }
1890            tx_state.sent_order.pop_front();
1891            tx_state.sent.remove(&seq);
1892        }
1893    }
1894
1895    fn handle_reliable_partial_ack(&self, side: RelaySideId, ty: crate::DataType, seq: u32) {
1896        let mut st = self.state.lock();
1897        let tx_state = self.reliable_tx_state_mut(&mut st, side, ty);
1898        if let Some(sent) = tx_state.sent.get_mut(&seq) {
1899            sent.partial_acked = true;
1900        }
1901    }
1902
1903    fn reliable_control_packet(
1904        &self,
1905        control_ty: crate::DataType,
1906        ty: crate::DataType,
1907        seq: u32,
1908    ) -> TelemetryResult<Packet> {
1909        let sender = self.sender_arc();
1910        Packet::new(
1911            control_ty,
1912            message_meta(control_ty).endpoints_ref(),
1913            sender.as_ref(),
1914            self.clock.now_ms(),
1915            crate::router::encode_slice_le(&[ty.as_u32(), seq]),
1916        )
1917    }
1918
1919    fn queue_reliable_ack(
1920        &self,
1921        side: RelaySideId,
1922        ty: crate::DataType,
1923        seq: u32,
1924    ) -> TelemetryResult<()> {
1925        let pkt = self.reliable_control_packet(crate::DataType::ReliableAck, ty, seq)?;
1926        let data = RelayItem::Packet(Arc::new(pkt));
1927        let priority = Self::relay_item_priority(&data)?;
1928        let mut st = self.state.lock();
1929        st.push_tx(RelayTxItem {
1930            src: None,
1931            dst: side,
1932            data,
1933            priority,
1934        })?;
1935        Ok(())
1936    }
1937
1938    fn queue_reliable_packet_request(
1939        &self,
1940        side: RelaySideId,
1941        ty: crate::DataType,
1942        seq: u32,
1943    ) -> TelemetryResult<()> {
1944        let pkt = self.reliable_control_packet(crate::DataType::ReliablePacketRequest, ty, seq)?;
1945        let data = RelayItem::Packet(Arc::new(pkt));
1946        let priority = Self::relay_item_priority(&data)?;
1947        let mut st = self.state.lock();
1948        st.push_tx(RelayTxItem {
1949            src: None,
1950            dst: side,
1951            data,
1952            priority,
1953        })?;
1954        Ok(())
1955    }
1956
1957    fn queue_reliable_partial_ack(
1958        &self,
1959        side: RelaySideId,
1960        ty: crate::DataType,
1961        seq: u32,
1962    ) -> TelemetryResult<()> {
1963        let pkt = self.reliable_control_packet(crate::DataType::ReliablePartialAck, ty, seq)?;
1964        let data = RelayItem::Packet(Arc::new(pkt));
1965        let priority = Self::relay_item_priority(&data)?;
1966        let mut st = self.state.lock();
1967        st.push_tx(RelayTxItem {
1968            src: None,
1969            dst: side,
1970            data,
1971            priority,
1972        })?;
1973        Ok(())
1974    }
1975
1976    fn queue_reliable_retransmit(
1977        &self,
1978        side: RelaySideId,
1979        ty: crate::DataType,
1980        seq: u32,
1981    ) -> TelemetryResult<()> {
1982        let mut queued = None;
1983        {
1984            let mut st = self.state.lock();
1985            let tx_state = self.reliable_tx_state_mut(&mut st, side, ty);
1986            if let Some(sent) = tx_state.sent.get_mut(&seq)
1987                && !sent.queued
1988            {
1989                sent.queued = true;
1990                sent.partial_acked = false;
1991                queued = Some(sent.bytes.clone());
1992            }
1993        }
1994
1995        if let Some(bytes) = queued {
1996            let mut st = self.state.lock();
1997            st.push_replay(RelayReplayItem {
1998                dst: side,
1999                bytes,
2000                priority: message_priority(ty).saturating_add(16),
2001            })?;
2002        }
2003        Ok(())
2004    }
2005
2006    fn send_reliable_raw_to_side(
2007        &self,
2008        side: RelaySideId,
2009        bytes: Arc<[u8]>,
2010    ) -> TelemetryResult<()> {
2011        let (handler, opts) = {
2012            let st = self.state.lock();
2013            let side_ref = Self::side_ref(&st, side)?;
2014            if !side_ref.opts.egress_enabled {
2015                return Ok(());
2016            }
2017            (side_ref.tx_handler.clone(), side_ref.opts)
2018        };
2019
2020        let Some(_side_tx_guard) = self.try_enter_side_tx() else {
2021            return Err(TelemetryError::Io("side tx busy"));
2022        };
2023        let started_ms = self.clock.now_ms();
2024        let ty = wire_format::peek_envelope(bytes.as_ref())
2025            .map(|env| env.ty)
2026            .unwrap_or(crate::DataType::ReliableAck);
2027        let result = match handler {
2028            RelayTxHandlerFn::Packed(f) => {
2029                let frames = self.encode_side_transport_frames(side, opts, bytes.clone())?;
2030                let mut sent_bytes = 0usize;
2031                for frame in frames {
2032                    f(frame.as_ref())?;
2033                    sent_bytes = sent_bytes.saturating_add(frame.len());
2034                }
2035                self.record_side_tx_sample(side, sent_bytes, started_ms, self.clock.now_ms());
2036                self.note_side_tx_success(side, ty, sent_bytes, 1);
2037                return Ok(());
2038            }
2039            RelayTxHandlerFn::Packet(f) => {
2040                if wire_format::peek_frame_info(bytes.as_ref())
2041                    .ok()
2042                    .is_some_and(|frame| frame.ack_only())
2043                {
2044                    return Ok(());
2045                }
2046                let pkt = wire_format::unpack_packet(bytes.as_ref())?;
2047                f(&pkt)
2048            }
2049        };
2050        if result.is_ok() {
2051            self.record_side_tx_sample(side, bytes.len(), started_ms, self.clock.now_ms());
2052            self.note_side_tx_success(side, ty, bytes.len(), 1);
2053        } else {
2054            self.note_side_tx_failure(side, ty, 1);
2055        }
2056        result
2057    }
2058
2059    fn send_reliable_to_side(&self, side: RelaySideId, data: RelayItem) -> TelemetryResult<()> {
2060        let (handler, opts, hop_reliable_enabled) = {
2061            let st = self.state.lock();
2062            let side_ref = Self::side_ref(&st, side)?;
2063            let opts = side_ref.opts;
2064            let hop_reliable_enabled = opts.reliable_enabled
2065                && !self.side_has_multiple_announcers_locked(&st, side, self.clock.now_ms());
2066            (side_ref.tx_handler.clone(), opts, hop_reliable_enabled)
2067        };
2068
2069        let RelayTxHandlerFn::Packed(f) = &handler else {
2070            return self.call_tx_handler(side, &handler, &data);
2071        };
2072
2073        if !hop_reliable_enabled {
2074            let mut adjusted_opts = opts;
2075            adjusted_opts.reliable_enabled = false;
2076            if let Some(adjusted) = self.adjust_reliable_for_side(adjusted_opts, data)? {
2077                return self.call_tx_handler(side, &handler, &adjusted);
2078            }
2079            return Ok(());
2080        }
2081
2082        let ty = match &data {
2083            RelayItem::Packet(pkt) => pkt.data_type(),
2084            RelayItem::Packed(bytes) => {
2085                let Ok(frame) = wire_format::peek_frame_info(bytes.as_ref()) else {
2086                    return self.call_tx_handler(side, &handler, &data);
2087                };
2088                frame.envelope.ty
2089            }
2090        };
2091
2092        if !is_reliable_type(ty) {
2093            if let Some(adjusted) = self.adjust_reliable_for_side(opts, data)? {
2094                self.call_tx_handler(side, &handler, &adjusted)?;
2095            }
2096            return Ok(());
2097        }
2098
2099        let (seq, flags) = {
2100            let mut st = self.state.lock();
2101            let tx_state = self.reliable_tx_state_mut(&mut st, side, ty);
2102            if tx_state.sent.len() >= runtime_reliable_max_pending() {
2103                return Err(TelemetryError::PacketTooLarge(
2104                    "relay reliable history full",
2105                ));
2106            }
2107            let seq = tx_state.next_seq;
2108            let next = tx_state.next_seq.wrapping_add(1);
2109            tx_state.next_seq = if next == 0 { 1 } else { next };
2110            let flags = match reliable_mode(ty) {
2111                crate::ReliableMode::Unordered => wire_format::RELIABLE_FLAG_UNORDERED,
2112                _ => 0,
2113            };
2114            (seq, flags)
2115        };
2116
2117        let bytes: Arc<[u8]> = match data {
2118            RelayItem::Packet(pkt) => wire_format::pack_packet_with_reliable(
2119                &pkt,
2120                wire_format::ReliableHeader { flags, seq, ack: 0 },
2121            ),
2122            RelayItem::Packed(bytes) => {
2123                let Some(rewritten) =
2124                    wire_format::rewrite_reliable_header_owned(bytes.as_ref(), flags, seq, 0)?
2125                else {
2126                    let Some(_side_tx_guard) = self.try_enter_side_tx() else {
2127                        return Err(TelemetryError::Io("side tx busy"));
2128                    };
2129                    let started_ms = self.clock.now_ms();
2130                    let frames = self.encode_side_transport_frames(side, opts, bytes.clone())?;
2131                    let mut sent_bytes = 0usize;
2132                    for frame in frames {
2133                        f(frame.as_ref())?;
2134                        sent_bytes = sent_bytes.saturating_add(frame.len());
2135                    }
2136                    self.record_side_tx_sample(side, sent_bytes, started_ms, self.clock.now_ms());
2137                    self.note_side_tx_success(side, ty, sent_bytes, 1);
2138                    return Ok(());
2139                };
2140                rewritten
2141            }
2142        };
2143
2144        let Some(_side_tx_guard) = self.try_enter_side_tx() else {
2145            return Err(TelemetryError::Io("side tx busy"));
2146        };
2147        let started_ms = self.clock.now_ms();
2148        let frames = self.encode_side_transport_frames(side, opts, bytes.clone())?;
2149        let mut sent_bytes = 0usize;
2150        for frame in frames {
2151            f(frame.as_ref())?;
2152            sent_bytes = sent_bytes.saturating_add(frame.len());
2153        }
2154        self.record_side_tx_sample(side, sent_bytes, started_ms, self.clock.now_ms());
2155        self.note_side_tx_success(side, ty, sent_bytes, 1);
2156
2157        {
2158            let mut st = self.state.lock();
2159            let tx_state = self.reliable_tx_state_mut(&mut st, side, ty);
2160            tx_state.sent_order.push_back(seq);
2161            tx_state.sent.insert(
2162                seq,
2163                ReliableSent {
2164                    bytes: bytes.clone(),
2165                    last_send_ms: self.clock.now_ms(),
2166                    retries: 0,
2167                    queued: false,
2168                    partial_acked: false,
2169                },
2170            );
2171        }
2172
2173        Ok(())
2174    }
2175
2176    fn item_route_info(
2177        &self,
2178        data: &RelayItem,
2179    ) -> TelemetryResult<(Vec<crate::DataEndpoint>, crate::DataType)> {
2180        match data {
2181            RelayItem::Packet(pkt) => {
2182                let mut eps = pkt.endpoints().to_vec();
2183                eps.sort_unstable();
2184                eps.dedup();
2185                Ok((eps, pkt.data_type()))
2186            }
2187            RelayItem::Packed(bytes) => {
2188                let env = wire_format::peek_envelope(bytes.as_ref())?;
2189                let mut eps: Vec<crate::DataEndpoint> = env.endpoints.iter().copied().collect();
2190                eps.sort_unstable();
2191                eps.dedup();
2192                Ok((eps, env.ty))
2193            }
2194        }
2195    }
2196
2197    fn endpoints_are_link_local_only(eps: &[crate::DataEndpoint]) -> bool {
2198        !eps.is_empty() && eps.iter().all(|ep| ep.is_link_local_only())
2199    }
2200
2201    fn item_target_senders(&self, data: &RelayItem) -> TelemetryResult<Arc<[u64]>> {
2202        match data {
2203            RelayItem::Packet(pkt) => Ok(Arc::from(pkt.wire_target_senders())),
2204            RelayItem::Packed(bytes) => {
2205                Ok(wire_format::peek_envelope(bytes.as_ref())?.target_senders)
2206            }
2207        }
2208    }
2209
2210    #[cfg(feature = "discovery")]
2211    fn has_explicit_route_policy_locked(
2212        st: &RelayInner,
2213        src: Option<RelaySideId>,
2214        ty: crate::DataType,
2215    ) -> bool {
2216        st.route_overrides
2217            .keys()
2218            .any(|(route_src, _)| *route_src == src)
2219            || Self::has_typed_route_overrides_locked(st, src, ty)
2220    }
2221
2222    #[cfg(feature = "discovery")]
2223    fn side_matches_target_senders_locked(
2224        st: &RelayInner,
2225        side: RelaySideId,
2226        target_senders: &[u64],
2227        now_ms: u64,
2228    ) -> bool {
2229        st.discovery_routes
2230            .get(&side)
2231            .map(|route| {
2232                if now_ms.saturating_sub(route.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2233                    return false;
2234                }
2235                route.announcers.values().any(|sender_state| {
2236                    if now_ms.saturating_sub(sender_state.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2237                        return false;
2238                    }
2239                    sender_state
2240                        .topology_boards
2241                        .iter()
2242                        .any(|board| target_senders.contains(&Self::sender_hash(&board.sender_id)))
2243                })
2244            })
2245            .unwrap_or(false)
2246    }
2247
2248    fn remote_side_plan(
2249        &self,
2250        data: &RelayItem,
2251        exclude: RelaySideId,
2252    ) -> TelemetryResult<RemoteSidePlan> {
2253        #[cfg(feature = "discovery")]
2254        {
2255            let (eps, ty) = self.item_route_info(data)?;
2256            let target_senders = self.item_target_senders(data)?;
2257            let preferred_packet_id = Self::reliable_control_target_packet_id(data)?;
2258            if discovery::is_discovery_type(ty) {
2259                let mut st = self.state.lock();
2260                let sides = self.eligible_side_ids_locked(&st, Some(exclude), Some(ty), false);
2261                return Ok(RemoteSidePlan::Target(self.apply_route_selection_locked(
2262                    &mut st,
2263                    Some(exclude),
2264                    sides,
2265                    RouteSelectionOrigin::Flood,
2266                )));
2267            }
2268
2269            #[cfg(feature = "timesync")]
2270            let preferred_timesync_source = self.preferred_timesync_route_source(data, ty)?;
2271            #[cfg(not(feature = "timesync"))]
2272            let preferred_timesync_source: Option<String> = None;
2273            let mut st = self.state.lock();
2274            if let Some(packet_id) = preferred_packet_id {
2275                let target_side = self.allowed_target_side_locked(
2276                    &st,
2277                    exclude,
2278                    ty,
2279                    st.reliable_return_routes
2280                        .get(&packet_id)
2281                        .map(|route| route.side),
2282                );
2283                if let Some(side) = target_side {
2284                    #[cfg(feature = "timesync")]
2285                    if !Self::timesync_allowed_for_side_locked(
2286                        &mut st,
2287                        side,
2288                        ty,
2289                        self.clock.now_ms(),
2290                    ) {
2291                        return Ok(RemoteSidePlan::Target(Vec::new()));
2292                    }
2293                    return Ok(RemoteSidePlan::Target(vec![side]));
2294                }
2295                return Ok(RemoteSidePlan::Target(Vec::new()));
2296            }
2297            let restrict_link_local = Self::endpoints_are_link_local_only(&eps);
2298            let discovered_origin = if is_reliable_type(ty) {
2299                RouteSelectionOrigin::Flood
2300            } else {
2301                RouteSelectionOrigin::Discovered
2302            };
2303            if st.discovery_routes.is_empty() {
2304                let mut fallback = self.eligible_side_ids_locked(
2305                    &st,
2306                    Some(exclude),
2307                    Some(ty),
2308                    restrict_link_local,
2309                );
2310                #[cfg(feature = "timesync")]
2311                {
2312                    fallback = Self::filter_timesync_sides_locked(
2313                        &mut st,
2314                        ty,
2315                        self.clock.now_ms(),
2316                        fallback,
2317                    );
2318                }
2319                return Ok(RemoteSidePlan::Target(if fallback.len() == 1 {
2320                    fallback
2321                } else {
2322                    Vec::new()
2323                }));
2324            }
2325            let now_ms = self.clock.now_ms();
2326            let mut had_exact = false;
2327            let mut exact_targets = Vec::new();
2328            let mut had_known = false;
2329            let mut generic_targets = Vec::new();
2330
2331            for (&side, route) in st.discovery_routes.iter() {
2332                if side == exclude
2333                    || now_ms.saturating_sub(route.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS
2334                {
2335                    continue;
2336                }
2337                if restrict_link_local
2338                    && st
2339                        .sides
2340                        .get(side)
2341                        .and_then(|side| side.as_ref())
2342                        .map(|s| !s.opts.link_local_enabled)
2343                        .unwrap_or(true)
2344                {
2345                    continue;
2346                }
2347                if !self.route_allowed_locked(&st, Some(exclude), Some(ty), side) {
2348                    continue;
2349                }
2350                if !target_senders.is_empty() {
2351                    if !Self::side_matches_target_senders_locked(&st, side, &target_senders, now_ms)
2352                    {
2353                        continue;
2354                    }
2355                    had_known = true;
2356                    generic_targets.push(side);
2357                    continue;
2358                }
2359                if preferred_timesync_source.as_deref().is_some_and(|source| {
2360                    route.reachable_timesync_sources.iter().any(|s| s == source)
2361                }) {
2362                    had_exact = true;
2363                    exact_targets.push(side);
2364                    continue;
2365                }
2366                if eps.iter().copied().any(|ep| route.reachable.contains(&ep)) {
2367                    had_known = true;
2368                    generic_targets.push(side);
2369                }
2370            }
2371
2372            if had_exact {
2373                #[cfg(feature = "timesync")]
2374                {
2375                    exact_targets = Self::filter_timesync_sides_locked(
2376                        &mut st,
2377                        ty,
2378                        self.clock.now_ms(),
2379                        exact_targets,
2380                    );
2381                }
2382                let targets = self.filter_end_to_end_satisfied_sides_locked(
2383                    &st,
2384                    data,
2385                    exact_targets,
2386                    &eps,
2387                    ty,
2388                )?;
2389                Ok(RemoteSidePlan::Target(self.apply_route_selection_locked(
2390                    &mut st,
2391                    Some(exclude),
2392                    targets,
2393                    discovered_origin,
2394                )))
2395            } else if had_known {
2396                Self::retain_shortest_discovery_routes_locked(
2397                    &st,
2398                    &mut generic_targets,
2399                    &eps,
2400                    now_ms,
2401                );
2402                #[cfg(feature = "timesync")]
2403                {
2404                    generic_targets = Self::filter_timesync_sides_locked(
2405                        &mut st,
2406                        ty,
2407                        self.clock.now_ms(),
2408                        generic_targets,
2409                    );
2410                }
2411                let targets = self.filter_end_to_end_satisfied_sides_locked(
2412                    &st,
2413                    data,
2414                    generic_targets,
2415                    &eps,
2416                    ty,
2417                )?;
2418                Ok(RemoteSidePlan::Target(self.apply_route_selection_locked(
2419                    &mut st,
2420                    Some(exclude),
2421                    targets,
2422                    discovered_origin,
2423                )))
2424            } else {
2425                if Self::has_explicit_route_policy_locked(&st, Some(exclude), ty) {
2426                    let mut sides = self.eligible_side_ids_locked(
2427                        &st,
2428                        Some(exclude),
2429                        Some(ty),
2430                        restrict_link_local,
2431                    );
2432                    #[cfg(feature = "timesync")]
2433                    {
2434                        sides = Self::filter_timesync_sides_locked(
2435                            &mut st,
2436                            ty,
2437                            self.clock.now_ms(),
2438                            sides,
2439                        );
2440                    }
2441                    Ok(RemoteSidePlan::Target(self.apply_route_selection_locked(
2442                        &mut st,
2443                        Some(exclude),
2444                        sides,
2445                        RouteSelectionOrigin::Flood,
2446                    )))
2447                } else {
2448                    Ok(RemoteSidePlan::Target(Vec::new()))
2449                }
2450            }
2451        }
2452        #[cfg(not(feature = "discovery"))]
2453        {
2454            let (_, ty) = self.item_route_info(data)?;
2455            let mut st = self.state.lock();
2456            if let Some(packet_id) = Self::reliable_control_target_packet_id(data)? {
2457                let target_side = self.allowed_target_side_locked(
2458                    &st,
2459                    exclude,
2460                    ty,
2461                    st.reliable_return_routes
2462                        .get(&packet_id)
2463                        .map(|route| route.side),
2464                );
2465                if let Some(side) = target_side {
2466                    return Ok(RemoteSidePlan::Target(vec![side]));
2467                }
2468                return Ok(RemoteSidePlan::Target(Vec::new()));
2469            }
2470            let sides = self.eligible_side_ids_locked(&st, Some(exclude), Some(ty), false);
2471            Ok(RemoteSidePlan::Target(self.apply_route_selection_locked(
2472                &mut st,
2473                Some(exclude),
2474                sides,
2475                RouteSelectionOrigin::Flood,
2476            )))
2477        }
2478    }
2479
2480    #[inline]
2481    fn allowed_target_side_locked(
2482        &self,
2483        st: &RelayInner,
2484        exclude: RelaySideId,
2485        ty: crate::DataType,
2486        target_side: Option<RelaySideId>,
2487    ) -> Option<RelaySideId> {
2488        target_side.filter(|side| self.route_allowed_locked(st, Some(exclude), Some(ty), *side))
2489    }
2490
2491    fn filter_end_to_end_satisfied_sides_locked(
2492        &self,
2493        st: &RelayInner,
2494        data: &RelayItem,
2495        sides: Vec<RelaySideId>,
2496        eps: &[crate::DataEndpoint],
2497        ty: crate::DataType,
2498    ) -> TelemetryResult<Vec<RelaySideId>> {
2499        if !is_reliable_type(ty) || Self::reliable_control_target_packet_id(data)?.is_some() {
2500            return Ok(sides);
2501        }
2502        let packet_id = match data {
2503            RelayItem::Packet(pkt) => pkt.packet_id(),
2504            RelayItem::Packed(bytes) => match wire_format::packet_id_from_wire(bytes.as_ref()) {
2505                Ok(packet_id) => packet_id,
2506                Err(TelemetryError::Unpack("reliable control frame")) => return Ok(sides),
2507                Err(err) => return Err(err),
2508            },
2509        };
2510        let Some(acked) = st.end_to_end_acked_destinations.get(&packet_id) else {
2511            return Ok(sides);
2512        };
2513        let now_ms = self.clock.now_ms();
2514        let mut filtered = Vec::new();
2515        for side in sides {
2516            let Some(route) = st.discovery_routes.get(&side) else {
2517                filtered.push(side);
2518                continue;
2519            };
2520            let mut still_pending = false;
2521            let mut had_destination_board = false;
2522            for sender_state in route.announcers.values() {
2523                if now_ms.saturating_sub(sender_state.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2524                    continue;
2525                }
2526                for board in sender_state.topology_boards.iter() {
2527                    if !self.is_end_to_end_destination_sender(&board.sender_id) {
2528                        continue;
2529                    }
2530                    had_destination_board = true;
2531                    let sender_hash = Self::sender_hash(&board.sender_id);
2532                    if acked.contains(&sender_hash) {
2533                        continue;
2534                    }
2535                    if eps
2536                        .iter()
2537                        .copied()
2538                        .any(|ep| board.reachable_endpoints.contains(&ep))
2539                    {
2540                        still_pending = true;
2541                        break;
2542                    }
2543                    // Keep forwarding while any discovered destination sender for this packet
2544                    // remains unacked, even if topology/schema metadata changed for new packets.
2545                    still_pending = true;
2546                    break;
2547                }
2548                if still_pending {
2549                    break;
2550                }
2551            }
2552            if still_pending || !had_destination_board {
2553                filtered.push(side);
2554            }
2555        }
2556        Ok(filtered)
2557    }
2558
2559    #[cfg(feature = "discovery")]
2560    fn side_has_multiple_announcers_locked(
2561        &self,
2562        st: &RelayInner,
2563        side: RelaySideId,
2564        now_ms: u64,
2565    ) -> bool {
2566        st.discovery_routes
2567            .get(&side)
2568            .map(|route| {
2569                route
2570                    .announcers
2571                    .values()
2572                    .filter(|sender| {
2573                        now_ms.saturating_sub(sender.last_seen_ms) <= DISCOVERY_ROUTE_TTL_MS
2574                    })
2575                    .take(2)
2576                    .count()
2577                    > 1
2578            })
2579            .unwrap_or(false)
2580    }
2581
2582    #[cfg(not(feature = "discovery"))]
2583    fn side_has_multiple_announcers_locked(
2584        &self,
2585        _st: &RelayInner,
2586        _side: RelaySideId,
2587        _now_ms: u64,
2588    ) -> bool {
2589        false
2590    }
2591
2592    #[cfg(feature = "discovery")]
2593    fn sender_topology_board_mut<'a>(
2594        sender_state: &'a mut DiscoverySenderState,
2595        sender_id: &str,
2596    ) -> &'a mut TopologyBoardNode {
2597        if let Some(idx) = sender_state
2598            .topology_boards
2599            .iter()
2600            .position(|board| board.sender_id == sender_id)
2601        {
2602            return &mut sender_state.topology_boards[idx];
2603        }
2604        sender_state.topology_boards.push(TopologyBoardNode {
2605            sender_id: sender_id.to_string(),
2606            reachable_endpoints: Vec::new(),
2607            reachable_timesync_sources: Vec::new(),
2608            connections: Vec::new(),
2609        });
2610        sender_state
2611            .topology_boards
2612            .last_mut()
2613            .expect("board inserted above")
2614    }
2615
2616    #[cfg(feature = "discovery")]
2617    fn refresh_sender_topology_state(sender_state: &mut DiscoverySenderState) {
2618        discovery::normalize_topology_boards(&mut sender_state.topology_boards);
2619        let (reachable, reachable_timesync_sources) =
2620            discovery::summarize_topology_boards(&sender_state.topology_boards);
2621        sender_state.reachable = reachable;
2622        sender_state.reachable_timesync_sources = reachable_timesync_sources;
2623    }
2624
2625    #[cfg(feature = "discovery")]
2626    fn recompute_discovery_side_state(route: &mut DiscoverySideState) {
2627        let mut reachable = Vec::new();
2628        let mut reachable_timesync_sources = Vec::new();
2629        let mut last_seen_ms = 0u64;
2630        for sender in route.announcers.values() {
2631            reachable.extend(sender.reachable.iter().copied());
2632            reachable_timesync_sources.extend(sender.reachable_timesync_sources.iter().cloned());
2633            last_seen_ms = last_seen_ms.max(sender.last_seen_ms);
2634        }
2635        reachable.sort_unstable();
2636        reachable.dedup();
2637        reachable_timesync_sources.sort_unstable();
2638        reachable_timesync_sources.dedup();
2639        route.reachable = reachable;
2640        route.reachable_timesync_sources = reachable_timesync_sources;
2641        route.last_seen_ms = last_seen_ms;
2642    }
2643
2644    #[cfg(feature = "discovery")]
2645    fn local_discovery_topology_board(&self, st: &RelayInner, now_ms: u64) -> TopologyBoardNode {
2646        let mut connections = Vec::new();
2647        for route in st.discovery_routes.values() {
2648            if now_ms.saturating_sub(route.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2649                continue;
2650            }
2651            for (sender, sender_state) in route.announcers.iter() {
2652                if now_ms.saturating_sub(sender_state.last_seen_ms) <= DISCOVERY_ROUTE_TTL_MS {
2653                    connections.push(sender.clone());
2654                }
2655            }
2656        }
2657        connections.sort_unstable();
2658        connections.dedup();
2659        let sender = self.sender_arc();
2660        TopologyBoardNode {
2661            sender_id: sender.to_string(),
2662            reachable_endpoints: Vec::new(),
2663            reachable_timesync_sources: Vec::new(),
2664            connections,
2665        }
2666    }
2667
2668    #[cfg(feature = "discovery")]
2669    fn advertised_discovery_topology_for_link_locked(
2670        &self,
2671        st: &RelayInner,
2672        now_ms: u64,
2673        link_local_enabled: bool,
2674    ) -> Vec<TopologyBoardNode> {
2675        let mut boards = vec![self.local_discovery_topology_board(st, now_ms)];
2676        for route in st.discovery_routes.values() {
2677            if now_ms.saturating_sub(route.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2678                continue;
2679            }
2680            for (announcer, sender_state) in route.announcers.iter() {
2681                if now_ms.saturating_sub(sender_state.last_seen_ms) > DISCOVERY_ROUTE_TTL_MS {
2682                    continue;
2683                }
2684                let mut sender_boards = sender_state.topology_boards.clone();
2685                if sender_boards.is_empty() {
2686                    let sender = self.sender_arc();
2687                    sender_boards.push(TopologyBoardNode {
2688                        sender_id: announcer.clone(),
2689                        reachable_endpoints: sender_state.reachable.clone(),
2690                        reachable_timesync_sources: sender_state.reachable_timesync_sources.clone(),
2691                        connections: vec![sender.to_string()],
2692                    });
2693                } else if let Some(board) = sender_boards
2694                    .iter_mut()
2695                    .find(|board| board.sender_id == *announcer)
2696                {
2697                    board.connections.push(self.sender_arc().to_string());
2698                }
2699                if !link_local_enabled {
2700                    for board in sender_boards.iter_mut() {
2701                        board
2702                            .reachable_endpoints
2703                            .retain(|ep| !ep.is_link_local_only());
2704                    }
2705                }
2706                discovery::merge_topology_boards(&mut boards, &sender_boards);
2707            }
2708        }
2709        discovery::normalize_topology_boards(&mut boards);
2710        boards
2711    }
2712
2713    #[cfg(feature = "discovery")]
2714    fn note_discovery_topology_change_locked(st: &mut RelayInner, now_ms: u64) {
2715        st.discovery_cadence.on_topology_change(now_ms);
2716    }
2717
2718    #[cfg(feature = "discovery")]
2719    fn prune_discovery_routes_locked(st: &mut RelayInner, now_ms: u64) -> bool {
2720        let before = st.discovery_routes.clone();
2721        st.discovery_routes.retain(|_, route| {
2722            route.announcers.retain(|_, sender| {
2723                now_ms.saturating_sub(sender.last_seen_ms) <= DISCOVERY_ROUTE_TTL_MS
2724            });
2725            Self::recompute_discovery_side_state(route);
2726            !route.announcers.is_empty()
2727        });
2728        st.discovery_routes != before
2729    }
2730
2731    #[cfg(feature = "discovery")]
2732    fn reconcile_end_to_end_acked_destinations_locked(&self, st: &mut RelayInner) {
2733        let mut active_senders = BTreeSet::new();
2734        for route in st.discovery_routes.values() {
2735            for sender_state in route.announcers.values() {
2736                for board in sender_state.topology_boards.iter() {
2737                    if self.is_end_to_end_destination_sender(&board.sender_id) {
2738                        active_senders.insert(Self::sender_hash(&board.sender_id));
2739                    }
2740                }
2741            }
2742        }
2743        st.end_to_end_acked_destinations.retain(|_, acked| {
2744            acked.retain(|sender_hash| active_senders.contains(sender_hash));
2745            !acked.is_empty()
2746        });
2747    }
2748
2749    #[cfg(feature = "discovery")]
2750    fn advertised_discovery_endpoints_for_link_locked(
2751        &self,
2752        st: &RelayInner,
2753        now_ms: u64,
2754        link_local_enabled: bool,
2755    ) -> Vec<crate::DataEndpoint> {
2756        let (reachable_endpoints, _) = discovery::summarize_topology_boards(
2757            &self.advertised_discovery_topology_for_link_locked(st, now_ms, link_local_enabled),
2758        );
2759        reachable_endpoints
2760            .into_iter()
2761            .filter(|ep| {
2762                !discovery::is_discovery_endpoint(*ep)
2763                    && (link_local_enabled || !ep.is_link_local_only())
2764            })
2765            .collect()
2766    }
2767
2768    #[cfg(feature = "discovery")]
2769    fn advertised_discovery_timesync_sources_for_link_locked(
2770        &self,
2771        st: &RelayInner,
2772        now_ms: u64,
2773    ) -> Vec<String> {
2774        let (_, sources) = discovery::summarize_topology_boards(
2775            &self.advertised_discovery_topology_for_link_locked(st, now_ms, true),
2776        );
2777        sources
2778    }
2779
2780    #[cfg(feature = "discovery")]
2781    #[cfg(feature = "timesync")]
2782    fn preferred_timesync_route_source(
2783        &self,
2784        data: &RelayItem,
2785        ty: crate::DataType,
2786    ) -> TelemetryResult<Option<String>> {
2787        if !matches!(
2788            ty,
2789            crate::DataType::TimeSyncAnnounce | crate::DataType::TimeSyncResponse
2790        ) {
2791            return Ok(None);
2792        }
2793
2794        let sender = match data {
2795            RelayItem::Packet(pkt) => pkt.sender().to_owned(),
2796            RelayItem::Packed(bytes) => {
2797                if wire_format::peek_frame_info(bytes.as_ref())
2798                    .ok()
2799                    .is_some_and(|frame| frame.ack_only())
2800                {
2801                    return Ok(None);
2802                }
2803                wire_format::unpack_packet(bytes.as_ref())?
2804                    .sender()
2805                    .to_owned()
2806            }
2807        };
2808        Ok(Some(sender))
2809    }
2810
2811    #[cfg(feature = "discovery")]
2812    #[inline]
2813    fn side_is_slow_control_link_locked(
2814        st: &RelayInner,
2815        side_id: RelaySideId,
2816        now_ms: u64,
2817    ) -> bool {
2818        st.adaptive_route_stats.get(&side_id).is_some_and(|stats| {
2819            let recent_slow = stats.last_slow_observed_ms > 0
2820                && now_ms.saturating_sub(stats.last_slow_observed_ms)
2821                    <= DISCOVERY_SLOW_LINK_FULL_INTERVAL_MS;
2822            stats.sample_count > 0
2823                && ((stats.estimated_bandwidth_bps > 0
2824                    && stats.estimated_bandwidth_bps <= CONTROL_SLOW_LINK_CAPACITY_BPS)
2825                    || recent_slow)
2826        })
2827    }
2828
2829    #[cfg(feature = "discovery")]
2830    fn discovery_level_for_side_locked(
2831        st: &mut RelayInner,
2832        side_id: RelaySideId,
2833        now_ms: u64,
2834    ) -> Option<DiscoveryAdvertiseLevel> {
2835        if !Self::side_is_slow_control_link_locked(st, side_id, now_ms) {
2836            st.discovery_side_throttle.remove(&side_id);
2837            return Some(DiscoveryAdvertiseLevel::Full);
2838        }
2839
2840        let throttle = st.discovery_side_throttle.entry(side_id).or_default();
2841        if now_ms >= throttle.next_full_ms {
2842            throttle.next_full_ms = now_ms.saturating_add(DISCOVERY_SLOW_LINK_FULL_INTERVAL_MS);
2843            throttle.next_ping_ms = now_ms.saturating_add(DISCOVERY_SLOW_LINK_PING_INTERVAL_MS);
2844            return Some(DiscoveryAdvertiseLevel::Full);
2845        }
2846        if now_ms >= throttle.next_ping_ms {
2847            throttle.next_ping_ms = now_ms.saturating_add(DISCOVERY_SLOW_LINK_PING_INTERVAL_MS);
2848            return Some(DiscoveryAdvertiseLevel::MinimalPing);
2849        }
2850        None
2851    }
2852
2853    #[cfg(all(feature = "discovery", feature = "timesync"))]
2854    #[inline]
2855    fn is_timesync_type(ty: crate::DataType) -> bool {
2856        matches!(
2857            ty,
2858            crate::DataType::TimeSyncAnnounce
2859                | crate::DataType::TimeSyncRequest
2860                | crate::DataType::TimeSyncResponse
2861        )
2862    }
2863
2864    #[cfg(all(feature = "discovery", feature = "timesync"))]
2865    fn timesync_allowed_for_side_locked(
2866        st: &mut RelayInner,
2867        side_id: RelaySideId,
2868        ty: crate::DataType,
2869        now_ms: u64,
2870    ) -> bool {
2871        if !Self::is_timesync_type(ty) {
2872            return true;
2873        }
2874        if !Self::side_is_slow_control_link_locked(st, side_id, now_ms) {
2875            st.timesync_side_throttle.remove(&side_id);
2876            return true;
2877        }
2878
2879        let throttle = st.timesync_side_throttle.entry(side_id).or_default();
2880        if now_ms >= throttle.next_allowed_ms {
2881            throttle.next_allowed_ms = now_ms.saturating_add(TIMESYNC_SLOW_LINK_MIN_INTERVAL_MS);
2882            return true;
2883        }
2884        false
2885    }
2886
2887    #[cfg(all(feature = "discovery", feature = "timesync"))]
2888    fn filter_timesync_sides_locked(
2889        st: &mut RelayInner,
2890        ty: crate::DataType,
2891        now_ms: u64,
2892        sides: Vec<RelaySideId>,
2893    ) -> Vec<RelaySideId> {
2894        sides
2895            .into_iter()
2896            .filter(|side| Self::timesync_allowed_for_side_locked(st, *side, ty, now_ms))
2897            .collect()
2898    }
2899
2900    #[cfg(feature = "discovery")]
2901    fn queue_discovery_announce(&self, include_schema: bool) -> TelemetryResult<()> {
2902        #[cfg(not(feature = "std"))]
2903        let _ = include_schema;
2904        let now_ms = self.clock.now_ms();
2905        let per_side = {
2906            let mut st = self.state.lock();
2907            if Self::prune_discovery_routes_locked(&mut st, now_ms) {
2908                self.reconcile_end_to_end_acked_destinations_locked(&mut st);
2909                Self::note_discovery_topology_change_locked(&mut st, now_ms);
2910            }
2911            st.fit_discovery_budget();
2912            if !st.sides.iter().any(|side| side.is_some()) {
2913                return Ok(());
2914            }
2915            st.discovery_cadence.on_announce_sent(now_ms);
2916            let side_entries = st
2917                .sides
2918                .iter()
2919                .enumerate()
2920                .filter_map(|(side_id, side)| {
2921                    side.as_ref()
2922                        .map(|side| (side_id, side.opts.link_local_enabled, side.opts))
2923                })
2924                .collect::<Vec<_>>();
2925            let mut per_side = Vec::new();
2926            for (side_id, link_local_enabled, opts) in side_entries {
2927                if !self.route_allowed_locked(
2928                    &st,
2929                    None,
2930                    Some(crate::DataType::DiscoveryAnnounce),
2931                    side_id,
2932                ) {
2933                    continue;
2934                }
2935                let Some(level) = Self::discovery_level_for_side_locked(&mut st, side_id, now_ms)
2936                else {
2937                    continue;
2938                };
2939                let capabilities = opts.link_capabilities();
2940                if level == DiscoveryAdvertiseLevel::MinimalPing {
2941                    per_side.push((
2942                        side_id,
2943                        level,
2944                        Vec::new(),
2945                        Vec::new(),
2946                        Vec::new(),
2947                        capabilities,
2948                    ));
2949                    continue;
2950                }
2951                let endpoints = self.advertised_discovery_endpoints_for_link_locked(
2952                    &st,
2953                    now_ms,
2954                    link_local_enabled,
2955                );
2956                let timesync_sources =
2957                    self.advertised_discovery_timesync_sources_for_link_locked(&st, now_ms);
2958                let topology = self.advertised_discovery_topology_for_link_locked(
2959                    &st,
2960                    now_ms,
2961                    link_local_enabled,
2962                );
2963                per_side.push((
2964                    side_id,
2965                    level,
2966                    endpoints,
2967                    timesync_sources,
2968                    topology,
2969                    capabilities,
2970                ));
2971            }
2972            per_side
2973        };
2974        let mut st = self.state.lock();
2975        for (dst, level, endpoints, timesync_sources, topology, capabilities) in per_side {
2976            let sender = self.sender_arc();
2977            // no_std schemas are immutable and no_std receivers discard
2978            // remote schema packets, so only hosted relays advertise them.
2979            #[cfg(feature = "std")]
2980            if include_schema && level == DiscoveryAdvertiseLevel::Full {
2981                let pkt = discovery::build_discovery_schema(sender.as_ref(), now_ms)?;
2982                let data = RelayItem::Packet(Arc::new(pkt));
2983                let priority = Self::relay_item_priority(&data)?;
2984                st.push_tx(RelayTxItem {
2985                    src: None,
2986                    dst,
2987                    data,
2988                    priority,
2989                })?;
2990            }
2991            if level == DiscoveryAdvertiseLevel::Full {
2992                let pkt = discovery::build_discovery_link_capabilities(
2993                    sender.as_ref(),
2994                    now_ms,
2995                    capabilities,
2996                )?;
2997                let data = RelayItem::Packet(Arc::new(pkt));
2998                let priority = Self::relay_item_priority(&data)?;
2999                st.push_tx(RelayTxItem {
3000                    src: None,
3001                    dst,
3002                    data,
3003                    priority,
3004                })?;
3005            }
3006            if level == DiscoveryAdvertiseLevel::MinimalPing || !endpoints.is_empty() {
3007                let pkt = discovery::build_discovery_announce(
3008                    sender.as_ref(),
3009                    now_ms,
3010                    endpoints.as_slice(),
3011                )?;
3012                let data = RelayItem::Packet(Arc::new(pkt));
3013                let priority = Self::relay_item_priority(&data)?;
3014                st.push_tx(RelayTxItem {
3015                    src: None,
3016                    dst,
3017                    data,
3018                    priority,
3019                })?;
3020            }
3021            if level == DiscoveryAdvertiseLevel::Full && !timesync_sources.is_empty() {
3022                let pkt = discovery::build_discovery_timesync_sources(
3023                    sender.as_ref(),
3024                    now_ms,
3025                    timesync_sources.as_slice(),
3026                )?;
3027                let data = RelayItem::Packet(Arc::new(pkt));
3028                let priority = Self::relay_item_priority(&data)?;
3029                st.push_tx(RelayTxItem {
3030                    src: None,
3031                    dst,
3032                    data,
3033                    priority,
3034                })?;
3035            }
3036            if level == DiscoveryAdvertiseLevel::Full && !topology.is_empty() {
3037                let pkt = discovery::build_discovery_topology(sender.as_ref(), now_ms, &topology)?;
3038                let data = RelayItem::Packet(Arc::new(pkt));
3039                let priority = Self::relay_item_priority(&data)?;
3040                st.push_tx(RelayTxItem {
3041                    src: None,
3042                    dst,
3043                    data,
3044                    priority,
3045                })?;
3046            }
3047        }
3048        Ok(())
3049    }
3050
3051    #[cfg(feature = "discovery")]
3052    fn poll_discovery_announce(&self) -> TelemetryResult<bool> {
3053        let now_ms = self.clock.now_ms();
3054        let due = {
3055            let mut st = self.state.lock();
3056            let removed = Self::prune_discovery_routes_locked(&mut st, now_ms);
3057            if removed {
3058                self.reconcile_end_to_end_acked_destinations_locked(&mut st);
3059                Self::note_discovery_topology_change_locked(&mut st, now_ms);
3060            }
3061            st.fit_discovery_budget();
3062            let has_any = st.sides.iter().enumerate().any(|(side_id, side)| {
3063                let Some(side) = side.as_ref() else {
3064                    return false;
3065                };
3066                if !self.route_allowed_locked(
3067                    &st,
3068                    None,
3069                    Some(crate::DataType::DiscoveryAnnounce),
3070                    side_id,
3071                ) {
3072                    return false;
3073                }
3074                let _ = side;
3075                true
3076            });
3077            if !st.sides.iter().any(|side| side.is_some()) || !has_any {
3078                return Ok(false);
3079            }
3080            st.discovery_cadence.due(now_ms)
3081        };
3082        if !due {
3083            return Ok(false);
3084        }
3085        // Keep periodic discovery lightweight so a hosted relay cannot fill a
3086        // constrained link with repeated schema snapshots.
3087        self.queue_discovery_announce(false)?;
3088        Ok(true)
3089    }
3090
3091    #[cfg(feature = "discovery")]
3092    fn learn_discovery_item(&self, src: RelaySideId, data: &RelayItem) -> TelemetryResult<()> {
3093        let pkt = match data {
3094            RelayItem::Packet(pkt) => {
3095                if !discovery::is_discovery_type(pkt.data_type()) {
3096                    return Ok(());
3097                }
3098                pkt.as_ref().clone()
3099            }
3100            RelayItem::Packed(bytes) => {
3101                let env = wire_format::peek_envelope(bytes.as_ref())?;
3102                if !discovery::is_discovery_type(env.ty) {
3103                    return Ok(());
3104                }
3105                if wire_format::peek_frame_info(bytes.as_ref())
3106                    .ok()
3107                    .is_some_and(|frame| frame.ack_only())
3108                {
3109                    return Ok(());
3110                }
3111                wire_format::unpack_packet(bytes.as_ref())?
3112            }
3113        };
3114
3115        let now_ms = self.clock.now_ms();
3116        if pkt.data_type() == crate::DataType::DiscoverySchema {
3117            let snapshot = discovery::decode_discovery_schema(&pkt)?;
3118            let incoming_cost = crate::config::owned_schema_byte_cost(&snapshot);
3119            let mut st = self.state.lock();
3120            st.make_shared_queue_room(incoming_cost, RelayQueueKind::Discovery)?;
3121            let budget = st.memory.max_queue_budget;
3122            drop(st);
3123            let report = crate::config::merge_owned_schema_snapshot_with_budget(snapshot, budget)?;
3124            if report.changed() {
3125                let mut st = self.state.lock();
3126                st.fit_discovery_budget();
3127                Self::note_discovery_topology_change_locked(&mut st, now_ms);
3128            }
3129            return Ok(());
3130        }
3131        if pkt.data_type() == crate::DataType::DiscoveryLinkCapabilities {
3132            let _ = discovery::decode_discovery_link_capabilities(&pkt)?;
3133            return Ok(());
3134        }
3135        let mut st = self.state.lock();
3136        if pkt.data_type() == crate::DataType::DiscoveryLeave {
3137            let leaving = pkt.sender();
3138            let before = st.discovery_routes.clone();
3139            for route in st.discovery_routes.values_mut() {
3140                route.announcers.remove(leaving);
3141                for sender_state in route.announcers.values_mut() {
3142                    sender_state
3143                        .topology_boards
3144                        .retain(|board| board.sender_id != leaving);
3145                    for board in sender_state.topology_boards.iter_mut() {
3146                        board.connections.retain(|peer| peer != leaving);
3147                    }
3148                    Self::refresh_sender_topology_state(sender_state);
3149                }
3150                Self::recompute_discovery_side_state(route);
3151            }
3152            st.discovery_routes
3153                .retain(|_, route| !route.announcers.is_empty());
3154            if st.discovery_routes != before {
3155                Self::note_discovery_topology_change_locked(&mut st, now_ms);
3156            }
3157            let _ = Self::prune_discovery_routes_locked(&mut st, now_ms);
3158            self.reconcile_end_to_end_acked_destinations_locked(&mut st);
3159            return Ok(());
3160        }
3161        let address_ad = if pkt.data_type() == crate::DataType::DiscoveryAddress {
3162            Some(discovery::decode_discovery_address(&pkt)?)
3163        } else {
3164            None
3165        };
3166        let mut topology_ad = if pkt.data_type() == crate::DataType::DiscoveryTopology {
3167            Some(discovery::decode_discovery_topology(&pkt)?)
3168        } else {
3169            None
3170        };
3171        let announcer_id = if let Some(ad) = address_ad.as_ref() {
3172            ad.hostname.clone()
3173        } else {
3174            let canonical = Self::canonical_sender_locked(&st, pkt.sender());
3175            if canonical != pkt.sender() {
3176                canonical
3177            } else if let Some(address) = pkt
3178                .sender()
3179                .strip_prefix("@addr:")
3180                .and_then(|value| value.parse::<u32>().ok())
3181            {
3182                topology_ad
3183                    .as_ref()
3184                    .and_then(|boards| {
3185                        boards
3186                            .iter()
3187                            .find(|board| sender_address_u32(&board.sender_id) == address)
3188                    })
3189                    .map(|board| board.sender_id.clone())
3190                    .unwrap_or(canonical)
3191            } else {
3192                canonical
3193            }
3194        };
3195        let mut route = st.discovery_routes.get(&src).cloned().unwrap_or_default();
3196        if pkt.sender() != announcer_id {
3197            route.announcers.remove(pkt.sender());
3198        }
3199        let side_link_local_enabled = st
3200            .sides
3201            .get(src)
3202            .and_then(|entry| entry.as_ref())
3203            .map(|side_ref| side_ref.opts.link_local_enabled)
3204            .unwrap_or(false);
3205        let mut sender_state = route
3206            .announcers
3207            .get(&announcer_id)
3208            .cloned()
3209            .unwrap_or_default();
3210        let changed = match pkt.data_type() {
3211            crate::DataType::DiscoveryAddress => {
3212                let ad = address_ad.expect("decoded above");
3213                let mut reachable = ad.reachable_endpoints;
3214                if !side_link_local_enabled {
3215                    reachable.retain(|ep| !ep.is_link_local_only());
3216                }
3217                let board = Self::sender_topology_board_mut(&mut sender_state, &ad.hostname);
3218                let changed = board.reachable_endpoints != reachable
3219                    || board.reachable_timesync_sources != ad.reachable_timesync_sources;
3220                board.reachable_endpoints = reachable;
3221                board.reachable_timesync_sources = ad.reachable_timesync_sources;
3222                Self::refresh_sender_topology_state(&mut sender_state);
3223                changed
3224            }
3225            crate::DataType::DiscoveryAnnounce => {
3226                let mut reachable = discovery::decode_discovery_announce(&pkt)?;
3227                if !side_link_local_enabled {
3228                    reachable.retain(|ep| !ep.is_link_local_only());
3229                }
3230                let board = Self::sender_topology_board_mut(&mut sender_state, &announcer_id);
3231                let changed = board.reachable_endpoints != reachable;
3232                board.reachable_endpoints = reachable;
3233                Self::refresh_sender_topology_state(&mut sender_state);
3234                changed
3235            }
3236            crate::DataType::DiscoveryTimeSyncSources => {
3237                let sources = discovery::decode_discovery_timesync_sources(&pkt)?;
3238                let board = Self::sender_topology_board_mut(&mut sender_state, &announcer_id);
3239                let changed = board.reachable_timesync_sources != sources;
3240                board.reachable_timesync_sources = sources;
3241                Self::refresh_sender_topology_state(&mut sender_state);
3242                changed
3243            }
3244            crate::DataType::DiscoveryTopology => {
3245                let mut boards = topology_ad
3246                    .take()
3247                    .expect("topology packet was decoded before route selection");
3248                for board in boards.iter_mut() {
3249                    board.sender_id = Self::canonical_sender_locked(&st, &board.sender_id);
3250                    for peer in board.connections.iter_mut() {
3251                        *peer = Self::canonical_sender_locked(&st, peer);
3252                    }
3253                }
3254                if !side_link_local_enabled {
3255                    for board in boards.iter_mut() {
3256                        board
3257                            .reachable_endpoints
3258                            .retain(|ep| !ep.is_link_local_only());
3259                    }
3260                }
3261                let changed = sender_state.topology_boards != boards;
3262                sender_state.topology_boards = boards;
3263                Self::refresh_sender_topology_state(&mut sender_state);
3264                changed
3265            }
3266            crate::DataType::DiscoverySchema => false,
3267            _ => false,
3268        };
3269        sender_state.last_seen_ms = now_ms;
3270        route.announcers.insert(announcer_id, sender_state);
3271        Self::recompute_discovery_side_state(&mut route);
3272        st.discovery_routes.insert(src, route);
3273        st.fit_discovery_budget();
3274        if changed {
3275            Self::note_discovery_topology_change_locked(&mut st, now_ms);
3276        }
3277        let _ = Self::prune_discovery_routes_locked(&mut st, now_ms);
3278        self.reconcile_end_to_end_acked_destinations_locked(&mut st);
3279        Ok(())
3280    }
3281
3282    #[cfg(not(feature = "discovery"))]
3283    fn learn_discovery_item(&self, _src: RelaySideId, _data: &RelayItem) -> TelemetryResult<()> {
3284        Ok(())
3285    }
3286
3287    #[cfg(not(feature = "discovery"))]
3288    fn queue_discovery_announce(&self) -> TelemetryResult<()> {
3289        Ok(())
3290    }
3291
3292    #[cfg(not(feature = "discovery"))]
3293    fn poll_discovery_announce(&self) -> TelemetryResult<bool> {
3294        Ok(false)
3295    }
3296
3297    fn process_reliable_timeouts(&self) -> TelemetryResult<()> {
3298        let now = self.clock.now_ms();
3299        let mut requeue: Vec<(RelaySideId, crate::DataType, u32)> = Vec::new();
3300
3301        {
3302            let mut st = self.state.lock();
3303            if st.reliable_tx.is_empty() {
3304                return Ok(());
3305            }
3306
3307            for ((side, ty_u32), tx_state) in st.reliable_tx.iter_mut() {
3308                let Some(ty) = crate::DataType::try_from_u32(*ty_u32) else {
3309                    continue;
3310                };
3311                let sent_order: Vec<u32> = tx_state.sent_order.iter().copied().collect();
3312                for seq in sent_order {
3313                    let Some(sent) = tx_state.sent.get_mut(&seq) else {
3314                        continue;
3315                    };
3316                    if sent.queued
3317                        || now.wrapping_sub(sent.last_send_ms) < runtime_reliable_retransmit_ms()
3318                    {
3319                        continue;
3320                    }
3321                    if sent.partial_acked {
3322                        continue;
3323                    }
3324                    if sent.retries >= runtime_reliable_max_retries() {
3325                        tx_state.sent.remove(&seq);
3326                        tx_state.sent_order.retain(|existing| *existing != seq);
3327                        continue;
3328                    }
3329                    sent.retries += 1;
3330                    requeue.push((*side, ty, seq));
3331                }
3332            }
3333        }
3334
3335        for (side, ty, seq) in requeue {
3336            self.queue_reliable_retransmit(side, ty, seq)?;
3337        }
3338
3339        Ok(())
3340    }
3341
3342    /// Compute a de-dupe hash for a QueueItem.
3343    /// Uses packet ID for Packet items, and attempts to extract packet ID from
3344    /// packed bytes. If extraction fails, hashes raw bytes as a fallback.
3345    fn get_hash(item: &RelayRxItem) -> u64 {
3346        match &item.data {
3347            RelayItem::Packet(pkt) => pkt.packet_id(),
3348            RelayItem::Packed(bytes) => {
3349                let reliable_seq = wire_format::peek_frame_info(bytes.as_ref())
3350                    .ok()
3351                    .and_then(|frame| frame.reliable)
3352                    .and_then(|hdr| {
3353                        if (hdr.flags & wire_format::RELIABLE_FLAG_ACK_ONLY) != 0 {
3354                            None
3355                        } else {
3356                            Some(hdr.seq)
3357                        }
3358                    });
3359
3360                match wire_format::packet_id_from_wire(bytes.as_ref()) {
3361                    Ok(id) => {
3362                        if let Some(seq) = reliable_seq {
3363                            hash_bytes_u64(id, &seq.to_le_bytes())
3364                        } else {
3365                            id
3366                        }
3367                    }
3368                    Err(_e) => {
3369                        // Fallback: if bytes are malformed (or compression feature mismatch),
3370                        // hash raw bytes so we can still dedupe identical network duplicates.
3371                        let h: u64 = 0x9E37_79B9_7F4A_7C15;
3372                        hash_bytes_u64(h, bytes.as_ref())
3373                    }
3374                }
3375            }
3376        }
3377    }
3378
3379    /// Compute a dedupe ID for an incoming RelayRxItem.
3380    /// Note: we intentionally do *not* include `src` so that the same
3381    /// packet coming from multiple sides is only processed once.
3382    fn is_duplicate_pkt(&self, item: &RelayRxItem) -> TelemetryResult<bool> {
3383        let id = Self::get_hash(item);
3384
3385        let mut st = self.state.lock();
3386        if st.recent_rx.contains(&id) {
3387            Ok(true)
3388        } else {
3389            st.push_recent_rx(id)?;
3390            Ok(false)
3391        }
3392    }
3393
3394    fn should_forward_duplicate_reliable_item(&self, item: &RelayRxItem) -> TelemetryResult<bool> {
3395        let (_, ty) = self.item_route_info(&item.data)?;
3396        if !is_reliable_type(ty)
3397            || matches!(
3398                ty,
3399                crate::DataType::ReliableAck
3400                    | crate::DataType::ReliablePartialAck
3401                    | crate::DataType::ReliablePacketRequest
3402            )
3403        {
3404            return Ok(false);
3405        }
3406
3407        let RemoteSidePlan::Target(sides) = self.remote_side_plan(&item.data, item.src)?;
3408        let st = self.state.lock();
3409        let now_ms = self.clock.now_ms();
3410        Ok(sides
3411            .into_iter()
3412            .any(|side| self.side_has_multiple_announcers_locked(&st, side, now_ms)))
3413    }
3414
3415    /// Register a side whose TX callback consumes packed packet bytes.
3416    ///
3417    /// Returns the side id later used for ingress APIs such as `rx_packed_from_side`.
3418    /// The default options disable the relay's per-link reliable framing on this side.
3419    pub fn add_side_packed<N, F>(&self, name: N, tx: F) -> RelaySideId
3420    where
3421        N: AsRef<str>,
3422        F: Fn(&[u8]) -> TelemetryResult<()> + Send + Sync + 'static,
3423    {
3424        self.add_side_packed_with_options(name, tx, RelaySideOptions::default())
3425    }
3426
3427    /// Register a packed side with bounded-frame transport enabled.
3428    ///
3429    /// `max_frame_bytes == 0` leaves frames unbounded.
3430    pub fn add_side_packed_small_packets<N, F>(
3431        &self,
3432        name: N,
3433        tx: F,
3434        max_frame_bytes: usize,
3435    ) -> RelaySideId
3436    where
3437        N: AsRef<str>,
3438        F: Fn(&[u8]) -> TelemetryResult<()> + Send + Sync + 'static,
3439    {
3440        self.add_side_packed_with_options(
3441            name,
3442            tx,
3443            RelaySideOptions::default().with_small_packet_transport(max_frame_bytes),
3444        )
3445    }
3446
3447    /// Register a packed-output side with explicit side options.
3448    ///
3449    /// `opts.reliable_enabled` enables relay-managed per-hop ACK/retransmit behavior on this side.
3450    /// `opts.link_local_enabled` gates link-local-only forwarding and discovery use of this side.
3451    /// `ingress_enabled` and `egress_enabled` set the initial directional policy.
3452    pub fn add_side_packed_with_options<N, F>(
3453        &self,
3454        name: N,
3455        tx: F,
3456        opts: RelaySideOptions,
3457    ) -> RelaySideId
3458    where
3459        N: AsRef<str>,
3460        F: Fn(&[u8]) -> TelemetryResult<()> + Send + Sync + 'static,
3461    {
3462        let mut st = self.state.lock();
3463        let side = Some(RelaySide {
3464            name: Arc::from(name.as_ref()),
3465            tx_handler: RelayTxHandlerFn::Packed(Arc::new(tx)),
3466            opts,
3467        });
3468        let id = if let Some(id) = st.sides.iter().position(Option::is_none) {
3469            st.sides[id] = side;
3470            id
3471        } else {
3472            let id = st.sides.len();
3473            st.sides.push(side);
3474            id
3475        };
3476        st.side_runtime_stats
3477            .insert(id, SideRuntimeStatsInner::default());
3478        st.side_transport.insert(id, SideTransportState::default());
3479        #[cfg(feature = "discovery")]
3480        Self::note_discovery_topology_change_locked(&mut st, self.clock.now_ms());
3481        id
3482    }
3483
3484    /// Register a side whose TX callback receives decoded [`Packet`] values.
3485    ///
3486    /// Packet-output sides do not preserve the relay's packed reliable hop framing, so use a
3487    /// packed side when this hop should participate in relay-managed per-link reliability.
3488    pub fn add_side_packet<N, F>(&self, name: N, tx: F) -> RelaySideId
3489    where
3490        N: AsRef<str>,
3491        F: Fn(&Packet) -> TelemetryResult<()> + Send + Sync + 'static,
3492    {
3493        self.add_side_packet_with_options(name, tx, RelaySideOptions::default())
3494    }
3495
3496    /// Register a packet-output side with explicit side options.
3497    pub fn add_side_packet_with_options<N, F>(
3498        &self,
3499        name: N,
3500        tx: F,
3501        opts: RelaySideOptions,
3502    ) -> RelaySideId
3503    where
3504        N: AsRef<str>,
3505        F: Fn(&Packet) -> TelemetryResult<()> + Send + Sync + 'static,
3506    {
3507        let mut st = self.state.lock();
3508        let side = Some(RelaySide {
3509            name: Arc::from(name.as_ref()),
3510            tx_handler: RelayTxHandlerFn::Packet(Arc::new(tx)),
3511            opts,
3512        });
3513        let id = if let Some(id) = st.sides.iter().position(Option::is_none) {
3514            st.sides[id] = side;
3515            id
3516        } else {
3517            let id = st.sides.len();
3518            st.sides.push(side);
3519            id
3520        };
3521        st.side_runtime_stats
3522            .insert(id, SideRuntimeStatsInner::default());
3523        st.side_transport.insert(id, SideTransportState::default());
3524        #[cfg(feature = "discovery")]
3525        Self::note_discovery_topology_change_locked(&mut st, self.clock.now_ms());
3526        id
3527    }
3528
3529    /// Remove a side while keeping existing side IDs stable.
3530    ///
3531    /// `side` must be an id returned by one of the `add_side_*` calls. Remaining side ids are not
3532    /// renumbered.
3533    pub fn remove_side(&self, side: RelaySideId) -> TelemetryResult<()> {
3534        let now_ms = self.clock.now_ms();
3535        let mut st = self.state.lock();
3536        let slot = st.sides.get_mut(side).ok_or(TelemetryError::BadArg)?;
3537        if slot.is_none() {
3538            return Err(TelemetryError::BadArg);
3539        }
3540        *slot = None;
3541        while st.sides.last().is_some_and(Option::is_none) {
3542            st.sides.pop();
3543        }
3544        // Preserve capacity for link churn. It is bounded by the maximum number
3545        // of sides seen by this relay and is reclaimed when the relay drops.
3546        st.route_overrides
3547            .retain(|(src_side, dst_side), _| *src_side != Some(side) && *dst_side != side);
3548        st.typed_route_overrides
3549            .retain(|(src_side, _, dst_side), _| *src_side != Some(side) && *dst_side != side);
3550        st.route_weights
3551            .retain(|(src_side, dst_side), _| *src_side != Some(side) && *dst_side != side);
3552        st.route_priorities
3553            .retain(|(src_side, dst_side), _| *src_side != Some(side) && *dst_side != side);
3554        st.source_route_modes.remove(&Some(side));
3555        st.route_selection_cursors.remove(&Some(side));
3556        st.adaptive_route_stats.remove(&side);
3557        #[cfg(feature = "discovery")]
3558        st.discovery_side_throttle.remove(&side);
3559        #[cfg(all(feature = "discovery", feature = "timesync"))]
3560        st.timesync_side_throttle.remove(&side);
3561        st.side_runtime_stats.remove(&side);
3562        st.reliable_return_routes
3563            .retain(|_, route| route.side != side);
3564        st.rx_queue.retain(|queued| queued.src != side);
3565        st.tx_queue
3566            .retain(|queued| queued.dst != side && queued.src != Some(side));
3567        st.replay_queue.retain(|queued| queued.dst != side);
3568        st.reliable_tx.retain(|(side_id, _), _| *side_id != side);
3569        st.reliable_rx.retain(|(side_id, _), _| *side_id != side);
3570        #[cfg(feature = "discovery")]
3571        {
3572            st.discovery_routes.remove(&side);
3573            Self::note_discovery_topology_change_locked(&mut st, now_ms);
3574        }
3575        Ok(())
3576    }
3577
3578    /// Enable or disable ingress processing for a registered side.
3579    pub fn set_side_ingress_enabled(
3580        &self,
3581        side: RelaySideId,
3582        enabled: bool,
3583    ) -> TelemetryResult<()> {
3584        let now_ms = self.clock.now_ms();
3585        let mut st = self.state.lock();
3586        let side_ref = st
3587            .sides
3588            .get_mut(side)
3589            .and_then(|side| side.as_mut())
3590            .ok_or(TelemetryError::BadArg)?;
3591        side_ref.opts.ingress_enabled = enabled;
3592        #[cfg(feature = "discovery")]
3593        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3594        Ok(())
3595    }
3596
3597    /// Enable or disable egress toward a registered side.
3598    pub fn set_side_egress_enabled(&self, side: RelaySideId, enabled: bool) -> TelemetryResult<()> {
3599        let now_ms = self.clock.now_ms();
3600        let mut st = self.state.lock();
3601        let side_ref = st
3602            .sides
3603            .get_mut(side)
3604            .and_then(|side| side.as_mut())
3605            .ok_or(TelemetryError::BadArg)?;
3606        side_ref.opts.egress_enabled = enabled;
3607        if !enabled {
3608            st.tx_queue.retain(|queued| queued.dst != side);
3609            st.replay_queue.retain(|queued| queued.dst != side);
3610        }
3611        #[cfg(feature = "discovery")]
3612        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3613        Ok(())
3614    }
3615
3616    /// Set the route-selection policy for traffic originating from `src`.
3617    ///
3618    /// `src == None` targets locally-originated relay traffic such as discovery output.
3619    pub fn set_source_route_mode(
3620        &self,
3621        src: Option<RelaySideId>,
3622        mode: RouteSelectionMode,
3623    ) -> TelemetryResult<()> {
3624        let now_ms = self.clock.now_ms();
3625        let mut st = self.state.lock();
3626        if let Some(src) = src {
3627            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3628        }
3629        // Preserve explicit Fanout so discovery does not replace it with its
3630        // adaptive single-path default. Clearing is a separate operation.
3631        st.source_route_modes.insert(src, mode);
3632        st.route_selection_cursors.remove(&src);
3633        #[cfg(feature = "discovery")]
3634        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3635        Ok(())
3636    }
3637
3638    /// Clear a source-specific route-selection override.
3639    pub fn clear_source_route_mode(&self, src: Option<RelaySideId>) -> TelemetryResult<()> {
3640        let now_ms = self.clock.now_ms();
3641        let mut st = self.state.lock();
3642        if let Some(src) = src {
3643            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3644        }
3645        st.source_route_modes.remove(&src);
3646        st.route_selection_cursors.remove(&src);
3647        #[cfg(feature = "discovery")]
3648        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3649        Ok(())
3650    }
3651
3652    /// Set the weighted-routing weight from `src` toward `dst`.
3653    pub fn set_route_weight(
3654        &self,
3655        src: Option<RelaySideId>,
3656        dst: RelaySideId,
3657        weight: u32,
3658    ) -> TelemetryResult<()> {
3659        let now_ms = self.clock.now_ms();
3660        let mut st = self.state.lock();
3661        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3662        if let Some(src) = src {
3663            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3664        }
3665        st.route_weights.insert((src, dst), weight);
3666        st.route_selection_cursors.remove(&src);
3667        #[cfg(feature = "discovery")]
3668        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3669        Ok(())
3670    }
3671
3672    /// Clear a previously configured weighted-routing weight override.
3673    pub fn clear_route_weight(
3674        &self,
3675        src: Option<RelaySideId>,
3676        dst: RelaySideId,
3677    ) -> TelemetryResult<()> {
3678        let now_ms = self.clock.now_ms();
3679        let mut st = self.state.lock();
3680        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3681        if let Some(src) = src {
3682            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3683        }
3684        st.route_weights.remove(&(src, dst));
3685        st.route_selection_cursors.remove(&src);
3686        #[cfg(feature = "discovery")]
3687        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3688        Ok(())
3689    }
3690
3691    /// Set the failover priority from `src` toward `dst`.
3692    pub fn set_route_priority(
3693        &self,
3694        src: Option<RelaySideId>,
3695        dst: RelaySideId,
3696        priority: u32,
3697    ) -> TelemetryResult<()> {
3698        let now_ms = self.clock.now_ms();
3699        let mut st = self.state.lock();
3700        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3701        if let Some(src) = src {
3702            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3703        }
3704        st.route_priorities.insert((src, dst), priority);
3705        #[cfg(feature = "discovery")]
3706        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3707        Ok(())
3708    }
3709
3710    /// Clear a previously configured failover priority override.
3711    pub fn clear_route_priority(
3712        &self,
3713        src: Option<RelaySideId>,
3714        dst: RelaySideId,
3715    ) -> TelemetryResult<()> {
3716        let now_ms = self.clock.now_ms();
3717        let mut st = self.state.lock();
3718        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3719        if let Some(src) = src {
3720            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3721        }
3722        st.route_priorities.remove(&(src, dst));
3723        #[cfg(feature = "discovery")]
3724        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3725        Ok(())
3726    }
3727
3728    /// Allow or block routing from `src` toward `dst`.
3729    pub fn set_route(
3730        &self,
3731        src: Option<RelaySideId>,
3732        dst: RelaySideId,
3733        enabled: bool,
3734    ) -> TelemetryResult<()> {
3735        let now_ms = self.clock.now_ms();
3736        let mut st = self.state.lock();
3737        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3738        if let Some(src) = src {
3739            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3740        }
3741        st.route_overrides.insert((src, dst), enabled);
3742        #[cfg(feature = "discovery")]
3743        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3744        Ok(())
3745    }
3746
3747    /// Allow or block routing for a specific `DataType` from `src` toward `dst`.
3748    pub fn set_typed_route(
3749        &self,
3750        src: Option<RelaySideId>,
3751        ty: crate::DataType,
3752        dst: RelaySideId,
3753        enabled: bool,
3754    ) -> TelemetryResult<()> {
3755        let now_ms = self.clock.now_ms();
3756        let mut st = self.state.lock();
3757        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3758        if let Some(src) = src {
3759            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3760        }
3761        st.typed_route_overrides
3762            .insert((src, ty.as_u32(), dst), enabled);
3763        #[cfg(feature = "discovery")]
3764        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3765        Ok(())
3766    }
3767
3768    /// Clear a typed route override for the `(src, ty, dst)` triple.
3769    pub fn clear_typed_route(
3770        &self,
3771        src: Option<RelaySideId>,
3772        ty: crate::DataType,
3773        dst: RelaySideId,
3774    ) -> TelemetryResult<()> {
3775        let now_ms = self.clock.now_ms();
3776        let mut st = self.state.lock();
3777        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3778        if let Some(src) = src {
3779            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3780        }
3781        st.typed_route_overrides.remove(&(src, ty.as_u32(), dst));
3782        #[cfg(feature = "discovery")]
3783        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3784        Ok(())
3785    }
3786
3787    /// Clear a non-typed route override so the relay falls back to default behavior.
3788    pub fn clear_route(&self, src: Option<RelaySideId>, dst: RelaySideId) -> TelemetryResult<()> {
3789        let now_ms = self.clock.now_ms();
3790        let mut st = self.state.lock();
3791        let _ = Self::side_ref(&st, dst).map_err(|_| TelemetryError::BadArg)?;
3792        if let Some(src) = src {
3793            let _ = Self::side_ref(&st, src).map_err(|_| TelemetryError::BadArg)?;
3794        }
3795        st.route_overrides.remove(&(src, dst));
3796        #[cfg(feature = "discovery")]
3797        Self::note_discovery_topology_change_locked(&mut st, now_ms);
3798        Ok(())
3799    }
3800
3801    #[cfg(feature = "discovery")]
3802    /// Queues an immediate discovery announcement for this relay.
3803    pub fn announce_discovery(&self) -> TelemetryResult<()> {
3804        self.queue_discovery_announce(true)
3805    }
3806
3807    /// Broadcast that this relay is leaving so peers can prune topology immediately.
3808    pub fn announce_leave(&self) -> TelemetryResult<()> {
3809        let pkt = discovery::build_discovery_leave("relay", self.clock.now_ms())?;
3810        let mut st = self.state.lock();
3811        let dsts: Vec<usize> = st
3812            .sides
3813            .iter()
3814            .enumerate()
3815            .filter_map(|(idx, side)| side.as_ref().map(|_| idx))
3816            .collect();
3817        for dst in dsts {
3818            let data = RelayItem::Packet(Arc::new(pkt.clone()));
3819            let priority = Self::relay_item_priority(&data)?;
3820            st.push_tx(RelayTxItem {
3821                src: None,
3822                dst,
3823                data,
3824                priority,
3825            })?;
3826        }
3827        Ok(())
3828    }
3829
3830    #[cfg(feature = "discovery")]
3831    /// Polls discovery state and queues an announce if the cadence says one is due.
3832    pub fn poll_discovery(&self) -> TelemetryResult<bool> {
3833        self.poll_discovery_announce()
3834    }
3835
3836    #[cfg(feature = "discovery")]
3837    /// Exports the relay's current discovered topology snapshot.
3838    pub fn export_topology(&self) -> TopologySnapshot {
3839        let now_ms = self.clock.now_ms();
3840        let mut st = self.state.lock();
3841        if Self::prune_discovery_routes_locked(&mut st, now_ms) {
3842            self.reconcile_end_to_end_acked_destinations_locked(&mut st);
3843            Self::note_discovery_topology_change_locked(&mut st, now_ms);
3844        }
3845        let routes = st
3846            .discovery_routes
3847            .iter()
3848            .filter_map(|(&side_id, route)| {
3849                let side = st.sides.get(side_id).and_then(|side| side.as_ref())?;
3850                let announcers = route
3851                    .announcers
3852                    .iter()
3853                    .map(|(sender_id, sender_state)| TopologyAnnouncerRoute {
3854                        sender_id: sender_id.clone(),
3855                        reachable_endpoints: sender_state
3856                            .reachable
3857                            .iter()
3858                            .copied()
3859                            .filter(|ep| !discovery::is_router_control_endpoint(*ep))
3860                            .collect(),
3861                        reachable_timesync_sources: sender_state.reachable_timesync_sources.clone(),
3862                        routers: sender_state.topology_boards.clone(),
3863                        last_seen_ms: sender_state.last_seen_ms,
3864                        age_ms: now_ms.saturating_sub(sender_state.last_seen_ms),
3865                    })
3866                    .collect();
3867                Some(TopologySideRoute {
3868                    side_id,
3869                    side_name: side.name.to_string(),
3870                    reachable_endpoints: route
3871                        .reachable
3872                        .iter()
3873                        .copied()
3874                        .filter(|ep| !discovery::is_router_control_endpoint(*ep))
3875                        .collect(),
3876                    reachable_timesync_sources: route.reachable_timesync_sources.clone(),
3877                    announcers,
3878                    last_seen_ms: route.last_seen_ms,
3879                    age_ms: now_ms.saturating_sub(route.last_seen_ms),
3880                })
3881            })
3882            .collect();
3883        let routers = self.advertised_discovery_topology_for_link_locked(&st, now_ms, true);
3884        let advertised_endpoints =
3885            self.advertised_discovery_endpoints_for_link_locked(&st, now_ms, true);
3886        let advertised_timesync_sources =
3887            self.advertised_discovery_timesync_sources_for_link_locked(&st, now_ms);
3888        let links = discovery::topology_links_from_boards(&routers);
3889        TopologySnapshot {
3890            advertised_endpoints,
3891            advertised_timesync_sources,
3892            routers,
3893            links,
3894            routes,
3895            current_announce_interval_ms: st.discovery_cadence.current_interval_ms,
3896            next_announce_ms: st.discovery_cadence.next_announce_ms,
3897        }
3898    }
3899
3900    #[cfg(feature = "discovery")]
3901    pub fn client_stats(&self, sender_id: &str) -> Option<ClientStatsSnapshot> {
3902        let now_ms = self.clock.now_ms();
3903        let st = self.state.lock();
3904        let mut side_ids = Vec::new();
3905        let mut side_names = Vec::new();
3906        let mut last_seen_ms = None::<u64>;
3907        let mut reachable_endpoints = Vec::new();
3908        let mut reachable_timesync_sources = Vec::new();
3909        let mut packets_sent = 0u64;
3910        let mut packets_received = 0u64;
3911        let mut bytes_sent = 0u64;
3912        let mut bytes_received = 0u64;
3913
3914        for (side_id, route) in &st.discovery_routes {
3915            let Some(sender_state) = route.announcers.get(sender_id) else {
3916                continue;
3917            };
3918            side_ids.push(*side_id);
3919            if let Some(side_name) = st
3920                .sides
3921                .get(*side_id)
3922                .and_then(|side| side.as_ref())
3923                .map(|side| side.name.clone())
3924            {
3925                side_names.push(side_name.to_string());
3926            }
3927            last_seen_ms = Some(last_seen_ms.unwrap_or(0).max(sender_state.last_seen_ms));
3928            reachable_endpoints.extend(sender_state.reachable.iter().copied());
3929            reachable_timesync_sources
3930                .extend(sender_state.reachable_timesync_sources.iter().cloned());
3931            if let Some(stats) = st.side_runtime_stats.get(side_id) {
3932                packets_sent = packets_sent.saturating_add(stats.tx_packets);
3933                packets_received = packets_received.saturating_add(stats.rx_packets);
3934                bytes_sent = bytes_sent.saturating_add(stats.tx_bytes);
3935                bytes_received = bytes_received.saturating_add(stats.rx_bytes);
3936            }
3937        }
3938
3939        if side_ids.is_empty() {
3940            return None;
3941        }
3942        reachable_endpoints.retain(|ep| !discovery::is_router_control_endpoint(*ep));
3943        reachable_endpoints.sort_unstable();
3944        reachable_endpoints.dedup();
3945        reachable_timesync_sources.sort_unstable();
3946        reachable_timesync_sources.dedup();
3947        side_ids.sort_unstable();
3948        side_ids.dedup();
3949        side_names.sort_unstable();
3950        side_names.dedup();
3951        let age_ms = last_seen_ms.map(|seen| now_ms.saturating_sub(seen));
3952        Some(ClientStatsSnapshot {
3953            sender_id: sender_id.to_string(),
3954            connected: age_ms.is_some_and(|age| age <= DISCOVERY_ROUTE_TTL_MS),
3955            side_ids,
3956            side_names,
3957            last_seen_ms,
3958            age_ms,
3959            reachable_endpoints,
3960            reachable_timesync_sources,
3961            packets_sent,
3962            packets_received,
3963            bytes_sent,
3964            bytes_received,
3965        })
3966    }
3967
3968    pub fn export_runtime_stats(&self) -> RuntimeStatsSnapshot {
3969        let now_ms = self.clock.now_ms();
3970        let st = self.state.lock();
3971
3972        let mut sides = Vec::new();
3973        for (side_id, side) in st.sides.iter().enumerate() {
3974            let Some(side) = side.as_ref() else { continue };
3975            let stats = st
3976                .side_runtime_stats
3977                .get(&side_id)
3978                .cloned()
3979                .unwrap_or_default();
3980            let adaptive = st
3981                .adaptive_route_stats
3982                .get(&side_id)
3983                .cloned()
3984                .unwrap_or_default()
3985                .snapshot(now_ms, true);
3986            let (tx_template_count, rx_template_count) = st
3987                .side_transport
3988                .get(&side_id)
3989                .map(|state| (state.tx_template_count(), state.rx_template_count()))
3990                .unwrap_or((0, 0));
3991            let mut data_types: Vec<RuntimeTypeStats> = stats
3992                .data_types
3993                .into_iter()
3994                .map(|(ty, item)| RuntimeTypeStats {
3995                    data_type: crate::DataType(ty),
3996                    tx_packets: item.tx_packets,
3997                    tx_bytes: item.tx_bytes,
3998                    rx_packets: item.rx_packets,
3999                    rx_bytes: item.rx_bytes,
4000                    relayed_tx_packets: item.relayed_tx_packets,
4001                    relayed_tx_bytes: item.relayed_tx_bytes,
4002                    relayed_rx_packets: item.relayed_rx_packets,
4003                    relayed_rx_bytes: item.relayed_rx_bytes,
4004                    tx_retries: item.tx_retries,
4005                    handler_failures: item.handler_failures,
4006                })
4007                .collect();
4008            data_types.sort_unstable_by_key(|item| item.data_type.as_u32());
4009            sides.push(RuntimeSideStats {
4010                side_id,
4011                side_name: side.name.to_string(),
4012                reliable_enabled: side.opts.reliable_enabled,
4013                link_local_enabled: side.opts.link_local_enabled,
4014                header_template_enabled: side.opts.header_template_enabled,
4015                max_frame_bytes: side.opts.max_frame_bytes,
4016                compact_header_target_bytes: side.opts.compact_header_target_bytes,
4017                side_transport_profile: side.opts.effective_transport_profile().as_str(),
4018                ingress_enabled: side.opts.ingress_enabled,
4019                egress_enabled: side.opts.egress_enabled,
4020                tx_packets: stats.tx_packets,
4021                tx_bytes: stats.tx_bytes,
4022                rx_packets: stats.rx_packets,
4023                rx_bytes: stats.rx_bytes,
4024                relayed_tx_packets: stats.relayed_tx_packets,
4025                relayed_tx_bytes: stats.relayed_tx_bytes,
4026                relayed_rx_packets: stats.relayed_rx_packets,
4027                relayed_rx_bytes: stats.relayed_rx_bytes,
4028                local_delivery_packets: 0,
4029                tx_retries: stats.tx_retries,
4030                tx_handler_failures: stats.tx_handler_failures,
4031                local_handler_failures: 0,
4032                total_handler_retries: stats.total_handler_retries,
4033                side_transport_full_frames: stats.side_transport_full_frames,
4034                side_transport_compact_frames: stats.side_transport_compact_frames,
4035                side_transport_compact_delta_frames: stats.side_transport_compact_delta_frames,
4036                side_transport_compact_omitted_timestamp_frames: stats
4037                    .side_transport_compact_omitted_timestamp_frames,
4038                side_transport_chunk_frames: stats.side_transport_chunk_frames,
4039                side_transport_raw_bytes: stats.side_transport_raw_bytes,
4040                side_transport_wire_bytes: stats.side_transport_wire_bytes,
4041                side_transport_bytes_saved: stats.side_transport_bytes_saved,
4042                side_transport_min_compact_overhead_bytes: stats
4043                    .side_transport_min_compact_overhead_bytes,
4044                side_transport_max_compact_overhead_bytes: stats
4045                    .side_transport_max_compact_overhead_bytes,
4046                side_transport_compact_target_misses: stats.side_transport_compact_target_misses,
4047                side_transport_template_evictions: stats.side_transport_template_evictions,
4048                side_transport_tx_template_count: tx_template_count,
4049                side_transport_rx_template_count: rx_template_count,
4050                max_side_transport_templates: side.opts.max_side_transport_templates,
4051                adaptive,
4052                data_types,
4053            });
4054        }
4055
4056        let mut route_modes: Vec<RouteModeStats> = st
4057            .route_selection_cursors
4058            .iter()
4059            .map(|(src, cursor)| RouteModeStats {
4060                src_side_id: *src,
4061                selection_mode: st.source_route_modes.get(src).copied(),
4062                cursor: *cursor,
4063            })
4064            .collect();
4065        for src in st.source_route_modes.keys() {
4066            if !route_modes.iter().any(|mode| mode.src_side_id == *src) {
4067                route_modes.push(RouteModeStats {
4068                    src_side_id: *src,
4069                    selection_mode: st.source_route_modes.get(src).copied(),
4070                    cursor: 0,
4071                });
4072            }
4073        }
4074        route_modes.sort_unstable_by_key(|mode| mode.src_side_id.unwrap_or(usize::MAX));
4075
4076        let mut route_overrides: Vec<RouteOverrideStats> = st
4077            .route_overrides
4078            .iter()
4079            .map(|((src, dst), enabled)| RouteOverrideStats {
4080                src_side_id: *src,
4081                dst_side_id: *dst,
4082                enabled: *enabled,
4083            })
4084            .collect();
4085        route_overrides.sort_unstable_by_key(|item| {
4086            (item.src_side_id.unwrap_or(usize::MAX), item.dst_side_id)
4087        });
4088
4089        let mut typed_route_overrides: Vec<TypedRouteOverrideStats> = st
4090            .typed_route_overrides
4091            .iter()
4092            .map(|((src, ty, dst), enabled)| TypedRouteOverrideStats {
4093                src_side_id: *src,
4094                data_type: crate::DataType(*ty),
4095                dst_side_id: *dst,
4096                enabled: *enabled,
4097            })
4098            .collect();
4099        typed_route_overrides.sort_unstable_by_key(|item| {
4100            (
4101                item.src_side_id.unwrap_or(usize::MAX),
4102                item.data_type.as_u32(),
4103                item.dst_side_id,
4104            )
4105        });
4106
4107        let mut route_weights: Vec<RouteWeightStats> = st
4108            .route_weights
4109            .iter()
4110            .map(|((src, dst), weight)| RouteWeightStats {
4111                src_side_id: *src,
4112                dst_side_id: *dst,
4113                weight: *weight,
4114            })
4115            .collect();
4116        route_weights.sort_unstable_by_key(|item| {
4117            (item.src_side_id.unwrap_or(usize::MAX), item.dst_side_id)
4118        });
4119
4120        let mut route_priorities: Vec<RoutePriorityStats> = st
4121            .route_priorities
4122            .iter()
4123            .map(|((src, dst), priority)| RoutePriorityStats {
4124                src_side_id: *src,
4125                dst_side_id: *dst,
4126                priority: *priority,
4127            })
4128            .collect();
4129        route_priorities.sort_unstable_by_key(|item| {
4130            (item.src_side_id.unwrap_or(usize::MAX), item.dst_side_id)
4131        });
4132
4133        #[cfg(feature = "discovery")]
4134        let discovery = DiscoveryRuntimeStats {
4135            route_count: st.discovery_routes.len(),
4136            announcer_count: st
4137                .discovery_routes
4138                .values()
4139                .map(|route| route.announcers.len())
4140                .sum(),
4141            current_announce_interval_ms: Some(st.discovery_cadence.current_interval_ms),
4142            next_announce_ms: Some(st.discovery_cadence.next_announce_ms),
4143        };
4144        #[cfg(not(feature = "discovery"))]
4145        let discovery = DiscoveryRuntimeStats {
4146            route_count: 0,
4147            announcer_count: 0,
4148            current_announce_interval_ms: None,
4149            next_announce_ms: None,
4150        };
4151
4152        RuntimeStatsSnapshot {
4153            sides,
4154            route_modes,
4155            route_overrides,
4156            typed_route_overrides,
4157            route_weights,
4158            route_priorities,
4159            queues: QueueRuntimeStats {
4160                rx_len: st.rx_queue.len(),
4161                rx_bytes: st.rx_queue.bytes_used(),
4162                tx_len: st.tx_queue.len(),
4163                tx_bytes: st.tx_queue.bytes_used(),
4164                replay_len: st.replay_queue.len(),
4165                replay_bytes: st.replay_queue.bytes_used(),
4166                recent_rx_len: st.recent_rx.len(),
4167                recent_rx_bytes: st.recent_rx.bytes_used(),
4168                reliable_rx_buffered_len: st.reliable_rx_buffer_len(),
4169                reliable_rx_buffered_bytes: st.reliable_rx_buffered_bytes(),
4170                shared_queue_bytes_used: st.shared_queue_bytes_used(),
4171            },
4172            reliable: ReliableRuntimeStats {
4173                reliable_return_route_count: st.reliable_return_routes.len(),
4174                end_to_end_pending_count: 0,
4175                end_to_end_pending_destination_count: 0,
4176                end_to_end_acked_cache_count: st.end_to_end_acked_destinations.len(),
4177            },
4178            discovery,
4179            total_handler_failures: st.total_handler_failures,
4180            total_handler_retries: st.total_handler_retries,
4181        }
4182    }
4183
4184    /// Export current relay memory usage/layout as JSON for profiling.
4185    pub fn export_memory_layout_json(&self) -> String {
4186        let st = self.state.lock();
4187        #[cfg(feature = "discovery")]
4188        let discovery_bytes = st.discovery_bytes_used();
4189        #[cfg(not(feature = "discovery"))]
4190        let discovery_bytes = 0usize;
4191        let schema_bytes = crate::config::schema_bytes_used();
4192        let memory = st.memory;
4193        let mut out = String::new();
4194        let _ = core::fmt::Write::write_fmt(
4195            &mut out,
4196            format_args!(
4197                "{{\"kind\":\"relay\",\
4198                 \"shared_queue_bytes_used\":{},\"shared_queue_bytes_allocated\":{},\
4199                 \"rx_queue_bytes_used\":{},\"rx_queue_bytes_allocated\":{},\"rx_queue_len\":{},\
4200                 \"tx_queue_bytes_used\":{},\"tx_queue_bytes_allocated\":{},\"tx_queue_len\":{},\
4201                 \"replay_queue_bytes_used\":{},\"replay_queue_bytes_allocated\":{},\"replay_queue_len\":{},\
4202                 \"recent_rx_bytes_used\":{},\"recent_rx_bytes_allocated\":{},\"recent_rx_len\":{},\
4203                 \"reliable_rx_buffer_bytes_used\":{},\"reliable_rx_buffer_bytes_allocated\":{},\"reliable_rx_buffer_len\":{},\
4204                 \"discovery_bytes_used\":{},\"discovery_bytes_allocated\":{},\
4205                 \"schema_bytes_used\":{},\"schema_bytes_allocated\":{}}}",
4206                st.shared_queue_bytes_used(),
4207                memory.max_queue_budget,
4208                st.rx_queue.bytes_used(),
4209                st.rx_queue.max_bytes(),
4210                st.rx_queue.len(),
4211                st.tx_queue.bytes_used(),
4212                st.tx_queue.max_bytes(),
4213                st.tx_queue.len(),
4214                st.replay_queue.bytes_used(),
4215                st.replay_queue.max_bytes(),
4216                st.replay_queue.len(),
4217                st.recent_rx.bytes_used(),
4218                st.recent_rx.max_bytes(),
4219                st.recent_rx.len(),
4220                st.reliable_rx_buffered_bytes(),
4221                memory.max_queue_budget,
4222                st.reliable_rx_buffer_len(),
4223                discovery_bytes,
4224                memory.max_queue_budget,
4225                schema_bytes,
4226                memory.max_queue_budget,
4227            ),
4228        );
4229        out
4230    }
4231
4232    #[cfg(test)]
4233    pub(crate) fn debug_end_to_end_acked_destination_count(&self, packet_id: u64) -> Option<usize> {
4234        let st = self.state.lock();
4235        st.end_to_end_acked_destinations
4236            .get(&packet_id)
4237            .map(BTreeSet::len)
4238    }
4239
4240    #[cfg(test)]
4241    pub(crate) fn debug_end_to_end_acked_packet_count(&self) -> usize {
4242        let st = self.state.lock();
4243        st.end_to_end_acked_destinations.len()
4244    }
4245
4246    #[cfg(test)]
4247    pub(crate) fn debug_reliable_return_route_count(&self) -> usize {
4248        let st = self.state.lock();
4249        st.reliable_return_routes.len()
4250    }
4251
4252    #[cfg(test)]
4253    pub(crate) fn debug_side_storage(&self) -> (usize, usize) {
4254        let st = self.state.lock();
4255        (st.sides.len(), st.sides.capacity())
4256    }
4257
4258    /// Enqueue packed bytes that originated from `src` into the relay RX queue.
4259    ///
4260    /// Note: `Arc::from(bytes)` allocates and copies `len` bytes into a new `Arc<[u8]>`.
4261    /// This is still “fast enough” for many cases, but it is not allocation-free / ISR-safe.
4262    pub fn rx_packed_from_side(&self, src: RelaySideId, bytes: &[u8]) -> TelemetryResult<()> {
4263        self.ensure_side_ingress_enabled(src)?;
4264        let Some(bytes) = self.decode_side_transport_frame(src, bytes)? else {
4265            return Ok(());
4266        };
4267        let mut st = self.state.lock();
4268
4269        let data = RelayItem::Packed(bytes);
4270        let priority = Self::relay_item_priority(&data)?;
4271        st.push_rx(RelayRxItem {
4272            src,
4273            data,
4274            priority,
4275        })
4276    }
4277
4278    /// Enqueue a full packet that originated from `src` into the relay RX queue.
4279    ///
4280    /// The packet is wrapped in `Arc<Packet>` so fanout can clone the pointer cheaply.
4281    pub fn rx_from_side(&self, src: RelaySideId, packet: Packet) -> TelemetryResult<()> {
4282        self.ensure_side_ingress_enabled(src)?;
4283        let mut st = self.state.lock();
4284
4285        let data = RelayItem::Packet(Arc::new(packet));
4286        let priority = Self::relay_item_priority(&data)?;
4287        st.push_rx(RelayRxItem {
4288            src,
4289            data,
4290            priority,
4291        })
4292    }
4293
4294    /// Clear both RX and TX queues.
4295    pub fn clear_queues(&self) {
4296        let mut st = self.state.lock();
4297        st.rx_queue.clear();
4298        st.tx_queue.clear();
4299        st.tx_priority_burst = 0;
4300    }
4301
4302    /// Clear only RX queue.
4303    pub fn clear_rx_queue(&self) {
4304        let mut st = self.state.lock();
4305        st.rx_queue.clear();
4306    }
4307
4308    /// Clear only TX queue.
4309    pub fn clear_tx_queue(&self) {
4310        let mut st = self.state.lock();
4311        st.tx_queue.clear();
4312        st.tx_priority_burst = 0;
4313        st.replay_queue.clear();
4314    }
4315
4316    /// Internal: expand one RX item into TX items for all other sides.
4317    ///
4318    /// Fanout is cheap: the `RelayItem` is cloned (Arc bump) and reused across all destinations.
4319    fn process_rx_queue_item(&self, item: RelayRxItem) -> TelemetryResult<()> {
4320        self.ensure_side_ingress_enabled(item.src)?;
4321        match &item.data {
4322            RelayItem::Packet(pkt) => {
4323                let bytes = wire_format::pack_packet(pkt).len();
4324                self.note_side_rx(item.src, pkt.data_type(), bytes);
4325            }
4326            RelayItem::Packed(bytes) => {
4327                if let Ok(env) = wire_format::peek_envelope(bytes.as_ref()) {
4328                    self.note_side_rx(item.src, env.ty, bytes.len());
4329                }
4330            }
4331        }
4332        match &item.data {
4333            RelayItem::Packet(pkt) => {
4334                if is_reliable_type(pkt.data_type()) && !is_internal_control_type(pkt.data_type()) {
4335                    self.note_reliable_return_route(item.src, pkt.packet_id());
4336                }
4337            }
4338            RelayItem::Packed(bytes) => {
4339                if let Ok(env) = wire_format::peek_envelope(bytes.as_ref())
4340                    && is_reliable_type(env.ty)
4341                    && !is_internal_control_type(env.ty)
4342                    && let Ok(packet_id) = wire_format::packet_id_from_wire(bytes.as_ref())
4343                {
4344                    self.note_reliable_return_route(item.src, packet_id);
4345                }
4346            }
4347        }
4348        let mut released_buffered: Vec<Arc<[u8]>> = Vec::new();
4349        if let RelayItem::Packed(bytes) = &item.data {
4350            let (_opts, handler_is_packed, hop_reliable_enabled) = {
4351                let st = self.state.lock();
4352                let side_ref = Self::side_ref(&st, item.src)?;
4353                let opts = side_ref.opts;
4354                (
4355                    opts,
4356                    matches!(side_ref.tx_handler, RelayTxHandlerFn::Packed(_)),
4357                    opts.reliable_enabled
4358                        && !self.side_has_multiple_announcers_locked(
4359                            &st,
4360                            item.src,
4361                            self.clock.now_ms(),
4362                        ),
4363                )
4364            };
4365
4366            let frame = match wire_format::peek_frame_info(bytes.as_ref()) {
4367                Ok(frame) => frame,
4368                Err(e) => {
4369                    if matches!(e, TelemetryError::Unpack(msg) if msg == "crc32 mismatch")
4370                        && hop_reliable_enabled
4371                        && handler_is_packed
4372                        && let Ok(frame) = wire_format::peek_frame_info_unchecked(bytes.as_ref())
4373                    {
4374                        if is_reliable_type(frame.envelope.ty)
4375                            && let Some(hdr) = frame.reliable
4376                        {
4377                            let unordered = (hdr.flags & wire_format::RELIABLE_FLAG_UNORDERED) != 0;
4378                            let unsequenced =
4379                                (hdr.flags & wire_format::RELIABLE_FLAG_UNSEQUENCED) != 0;
4380
4381                            if !unsequenced {
4382                                let requested = if unordered {
4383                                    hdr.seq
4384                                } else {
4385                                    let mut st = self.state.lock();
4386                                    let rx_state = self.reliable_rx_state_mut(
4387                                        &mut st,
4388                                        item.src,
4389                                        frame.envelope.ty,
4390                                    );
4391                                    rx_state.expected_seq.min(hdr.seq)
4392                                };
4393                                self.queue_reliable_packet_request(
4394                                    item.src,
4395                                    frame.envelope.ty,
4396                                    requested,
4397                                )?;
4398                            }
4399                        }
4400                        return Ok(());
4401                    }
4402                    return Err(e);
4403                }
4404            };
4405
4406            if hop_reliable_enabled
4407                && handler_is_packed
4408                && is_reliable_type(frame.envelope.ty)
4409                && let Some(hdr) = frame.reliable
4410            {
4411                if frame.ack_only() {
4412                    self.handle_reliable_ack(item.src, frame.envelope.ty, hdr.ack);
4413                    return Ok(());
4414                }
4415                let unordered = (hdr.flags & wire_format::RELIABLE_FLAG_UNORDERED) != 0;
4416                let unsequenced = (hdr.flags & wire_format::RELIABLE_FLAG_UNSEQUENCED) != 0;
4417
4418                if !unsequenced {
4419                    if unordered {
4420                        self.queue_reliable_ack(item.src, frame.envelope.ty, hdr.seq)?;
4421                    } else {
4422                        let mut release: Vec<Arc<[u8]>> = Vec::new();
4423                        let mut last_delivered = None;
4424                        let mut ack_old = None;
4425                        let mut request_missing = None;
4426                        let mut partial_ack = None;
4427                        {
4428                            let mut st = self.state.lock();
4429                            let rx_state =
4430                                self.reliable_rx_state_mut(&mut st, item.src, frame.envelope.ty);
4431                            let expected_seq = rx_state.expected_seq;
4432                            if hdr.seq < expected_seq {
4433                                ack_old = Some(expected_seq.saturating_sub(1));
4434                            } else if hdr.seq > expected_seq {
4435                                request_missing = Some(expected_seq);
4436                                partial_ack = Some(hdr.seq);
4437                                st.buffer_reliable_rx(
4438                                    item.src,
4439                                    frame.envelope.ty,
4440                                    hdr.seq,
4441                                    bytes.clone(),
4442                                )?;
4443                            } else {
4444                                release.push(bytes.clone());
4445                                last_delivered = Some(hdr.seq);
4446                                let mut next_expected = hdr.seq.wrapping_add(1);
4447                                while let Some(buf) = rx_state.buffered.remove(&next_expected) {
4448                                    release.push(buf);
4449                                    last_delivered = Some(next_expected);
4450                                    let next = next_expected.wrapping_add(1);
4451                                    next_expected = if next == 0 { 1 } else { next };
4452                                }
4453                                rx_state.expected_seq = next_expected;
4454                            }
4455                        }
4456
4457                        if let Some(ack_seq) = ack_old {
4458                            self.queue_reliable_ack(item.src, frame.envelope.ty, ack_seq)?;
4459                            return Ok(());
4460                        }
4461                        if let Some(request_seq) = request_missing {
4462                            if let Some(partial_seq) = partial_ack {
4463                                self.queue_reliable_partial_ack(
4464                                    item.src,
4465                                    frame.envelope.ty,
4466                                    partial_seq,
4467                                )?;
4468                            }
4469                            self.queue_reliable_packet_request(
4470                                item.src,
4471                                frame.envelope.ty,
4472                                request_seq,
4473                            )?;
4474                            return Ok(());
4475                        }
4476                        if let Some(ack_seq) = last_delivered {
4477                            self.queue_reliable_ack(item.src, frame.envelope.ty, ack_seq)?;
4478                        }
4479                        released_buffered.extend(release.into_iter().skip(1));
4480                    }
4481                }
4482            }
4483        }
4484
4485        if self.is_duplicate_pkt(&item)? && !self.should_forward_duplicate_reliable_item(&item)? {
4486            // Already fanned out this packet recently; skip.
4487            return Ok(());
4488        }
4489
4490        self.dispatch_relay_rx_item(&item)?;
4491
4492        for release_bytes in released_buffered {
4493            let release_item = RelayRxItem {
4494                src: item.src,
4495                priority: Self::relay_item_priority(&RelayItem::Packed(release_bytes.clone()))?,
4496                data: RelayItem::Packed(release_bytes),
4497            };
4498            if self.is_duplicate_pkt(&release_item)?
4499                && !self.should_forward_duplicate_reliable_item(&release_item)?
4500            {
4501                continue;
4502            }
4503            self.dispatch_relay_rx_item(&release_item)?;
4504        }
4505        Ok(())
4506    }
4507
4508    fn dispatch_relay_rx_item(&self, item: &RelayRxItem) -> TelemetryResult<()> {
4509        match &item.data {
4510            RelayItem::Packet(pkt) => {
4511                if matches!(
4512                    pkt.data_type(),
4513                    crate::DataType::ReliableAck
4514                        | crate::DataType::ReliablePartialAck
4515                        | crate::DataType::ReliablePacketRequest
4516                ) {
4517                    if pkt.data_type() == crate::DataType::ReliableAck
4518                        && Self::is_end_to_end_ack_sender(pkt.sender())
4519                        && Self::decode_end_to_end_reliable_ack(pkt.payload()).is_ok()
4520                    {
4521                        if let Ok(packet_id) = Self::decode_end_to_end_reliable_ack(pkt.payload())
4522                            && let Some(sender_hash) =
4523                                Self::decode_end_to_end_ack_sender_hash(pkt.sender())
4524                        {
4525                            let mut st = self.state.lock();
4526                            Self::note_end_to_end_acked_destination_locked(
4527                                &mut st,
4528                                packet_id,
4529                                sender_hash,
4530                            );
4531                        }
4532                    } else {
4533                        let vals = pkt.data_as_u32()?;
4534                        if vals.len() != 2 {
4535                            return Err(TelemetryError::Unpack("bad reliable control payload"));
4536                        }
4537                        let ty = crate::DataType::try_from_u32(vals[0])
4538                            .ok_or(TelemetryError::InvalidType)?;
4539                        let seq = vals[1];
4540                        match pkt.data_type() {
4541                            crate::DataType::ReliableAck => {
4542                                self.handle_reliable_ack(item.src, ty, seq)
4543                            }
4544                            crate::DataType::ReliablePartialAck => {
4545                                self.handle_reliable_partial_ack(item.src, ty, seq)
4546                            }
4547                            crate::DataType::ReliablePacketRequest => {
4548                                self.queue_reliable_retransmit(item.src, ty, seq)?
4549                            }
4550                            _ => {}
4551                        }
4552                        return Ok(());
4553                    }
4554                }
4555            }
4556            RelayItem::Packed(bytes) => {
4557                let env = wire_format::peek_envelope(bytes.as_ref())?;
4558                if matches!(
4559                    env.ty,
4560                    crate::DataType::ReliableAck
4561                        | crate::DataType::ReliablePacketRequest
4562                        | crate::DataType::ReliablePartialAck
4563                ) {
4564                    let pkt = wire_format::unpack_packet(bytes.as_ref())?;
4565                    return self.dispatch_relay_rx_item(&RelayRxItem {
4566                        src: item.src,
4567                        data: RelayItem::Packet(Arc::new(pkt)),
4568                        priority: item.priority,
4569                    });
4570                }
4571            }
4572        }
4573
4574        let src = item.src;
4575        let data = item.data.clone();
4576        self.learn_discovery_item(src, &data)?;
4577
4578        let plan = self.remote_side_plan(&data, src)?;
4579        let mut st = self.state.lock();
4580        let RemoteSidePlan::Target(sides) = plan;
4581        for dst in sides {
4582            let priority = Self::relay_item_priority(&data)?;
4583            st.push_tx(RelayTxItem {
4584                src: Some(src),
4585                dst,
4586                data: data.clone(),
4587                priority,
4588            })?;
4589        }
4590        Ok(())
4591    }
4592
4593    #[inline]
4594    fn crc32_bytes(data: &[u8]) -> u32 {
4595        let mut hasher = Crc32Hasher::new();
4596        hasher.update(data);
4597        hasher.finalize()
4598    }
4599
4600    fn wrap_side_transport_frame(kind: u8, body: &[u8]) -> Arc<[u8]> {
4601        let mut out = Vec::with_capacity(
4602            SIDE_TRANSPORT_MAGIC.len() + 1 + body.len() + wire_format::CRC32_BYTES,
4603        );
4604        out.extend_from_slice(SIDE_TRANSPORT_MAGIC);
4605        out.push(kind);
4606        out.extend_from_slice(body);
4607        let crc = Self::crc32_bytes(&out);
4608        out.extend_from_slice(&crc.to_le_bytes());
4609        Arc::from(out)
4610    }
4611
4612    fn parse_side_transport_wrapper(bytes: &[u8]) -> TelemetryResult<Option<(u8, &[u8])>> {
4613        if bytes.len() < SIDE_TRANSPORT_MAGIC.len() + 1 + wire_format::CRC32_BYTES {
4614            return Ok(None);
4615        }
4616        if &bytes[..SIDE_TRANSPORT_MAGIC.len()] != SIDE_TRANSPORT_MAGIC {
4617            return Ok(None);
4618        }
4619        let data_len = bytes.len() - wire_format::CRC32_BYTES;
4620        let expected = u32::from_le_bytes([
4621            bytes[data_len],
4622            bytes[data_len + 1],
4623            bytes[data_len + 2],
4624            bytes[data_len + 3],
4625        ]);
4626        let data = &bytes[..data_len];
4627        if Self::crc32_bytes(data) != expected {
4628            return Err(TelemetryError::Unpack("side transport crc32 mismatch"));
4629        }
4630        let kind = data[SIDE_TRANSPORT_MAGIC.len()];
4631        Ok(Some((kind, &data[SIDE_TRANSPORT_MAGIC.len() + 1..])))
4632    }
4633
4634    fn read_uleb128_local(buf: &[u8], off: &mut usize) -> TelemetryResult<u64> {
4635        let mut result = 0u64;
4636        let mut shift = 0u32;
4637        for _ in 0..10 {
4638            let byte = *buf.get(*off).ok_or(TelemetryError::Unpack("short read"))?;
4639            *off += 1;
4640            result |= u64::from(byte & 0x7F) << shift;
4641            if (byte & 0x80) == 0 {
4642                return Ok(result);
4643            }
4644            shift += 7;
4645        }
4646        Err(TelemetryError::Unpack("uleb128 too long"))
4647    }
4648
4649    fn write_uleb128_local(mut value: u64, out: &mut Vec<u8>) {
4650        loop {
4651            let mut byte = (value & 0x7F) as u8;
4652            value >>= 7;
4653            if value != 0 {
4654                byte |= 0x80;
4655            }
4656            out.push(byte);
4657            if value == 0 {
4658                break;
4659            }
4660        }
4661    }
4662
4663    fn uleb128_len_local(mut value: u64) -> usize {
4664        let mut len = 1;
4665        while value >= 0x80 {
4666            value >>= 7;
4667            len += 1;
4668        }
4669        len
4670    }
4671
4672    fn extract_side_header_template(bytes: &[u8]) -> TelemetryResult<SideTemplateExtract<'_>> {
4673        if bytes.len() < wire_format::CRC32_BYTES + 4 {
4674            return Err(TelemetryError::Unpack("short buffer"));
4675        }
4676        let data_len = bytes.len() - wire_format::CRC32_BYTES;
4677        let data = &bytes[..data_len];
4678        let mut off = 0usize;
4679        let flags = *data
4680            .get(off)
4681            .ok_or(TelemetryError::Unpack("short prelude"))?;
4682        off += 1;
4683        off += 1; // NEP
4684        let ty_u64 = Self::read_uleb128_local(data, &mut off)?;
4685        let ty_u32 = u32::try_from(ty_u64).map_err(|_| TelemetryError::Unpack("bad data type"))?;
4686        if ty_u32 > crate::MAX_VALUE_DATA_TYPE {
4687            return Err(TelemetryError::Unpack("bad data type"));
4688        }
4689        let ty = crate::DataType(ty_u32);
4690        let data_size_off = off;
4691        let data_size = Self::read_uleb128_local(data, &mut off)?;
4692        let timestamp = Self::read_uleb128_local(data, &mut off)?;
4693        let nonce = if (flags & SIDE_TRANSPORT_FLAG_PACKET_NONCE) != 0 {
4694            u16::try_from(Self::read_uleb128_local(data, &mut off)?)
4695                .map_err(|_| TelemetryError::Unpack("packet nonce too large"))?
4696        } else {
4697            0
4698        };
4699        let between_start = off;
4700        let _source_address = u32::try_from(Self::read_uleb128_local(data, &mut off)?)
4701            .map_err(|_| TelemetryError::Unpack("source address too large"))?;
4702        let endpoint_bitmap_bytes = if (flags & SIDE_TRANSPORT_FLAG_ENDPOINT_BITMAP_PRESENT) != 0 {
4703            SIDE_TRANSPORT_EP_BITMAP_BYTES
4704        } else {
4705            0
4706        };
4707        if data.len() < off + endpoint_bitmap_bytes {
4708            return Err(TelemetryError::Unpack("short buffer"));
4709        }
4710        off += endpoint_bitmap_bytes;
4711        if (flags & SIDE_TRANSPORT_FLAG_WIRE_CONTRACT) != 0 {
4712            let contract_len = usize::try_from(Self::read_uleb128_local(data, &mut off)?)
4713                .map_err(|_| TelemetryError::Unpack("wire contract length"))?;
4714            if data.len() < off + contract_len {
4715                return Err(TelemetryError::Unpack("short buffer"));
4716            }
4717            off += contract_len;
4718        }
4719        let reliable_span = wire_format::reliable_header_span(bytes)?;
4720        let (reliable_flags, reliable_seq_ack, reliable_compact, payload_off) =
4721            if let Some((rel_off, rel_len, hdr)) = reliable_span {
4722                if data.len() < rel_off + rel_len {
4723                    return Err(TelemetryError::Unpack("short buffer"));
4724                }
4725                (
4726                    Some(hdr.flags),
4727                    Some((hdr.seq, hdr.ack)),
4728                    (flags & SIDE_TRANSPORT_FLAG_COMPACT_RELIABLE_HEADER) != 0,
4729                    rel_off + rel_len,
4730                )
4731            } else {
4732                (None, None, false, off)
4733            };
4734        if payload_off > data.len() {
4735            return Err(TelemetryError::Unpack("short buffer"));
4736        }
4737        let payload = &data[payload_off..];
4738        let prefix = Arc::<[u8]>::from(&data[1..data_size_off]);
4739        let between_end = reliable_span
4740            .map(|(rel_off, _, _)| rel_off)
4741            .unwrap_or(payload_off);
4742        let between = Arc::<[u8]>::from(&data[between_start..between_end]);
4743        let base_flags =
4744            flags & !(SIDE_TRANSPORT_FLAG_PAYLOAD_COMPRESSED | SIDE_TRANSPORT_FLAG_PACKET_NONCE);
4745        let mut hash = 0xD1B5_4A32_9C7E_01F3u64;
4746        hash = hash_bytes_u64(hash, &[base_flags]);
4747        hash = hash_bytes_u64(hash, &prefix);
4748        hash = hash_bytes_u64(hash, &between);
4749        if let Some(rel_flags) = reliable_flags {
4750            hash = hash_bytes_u64(hash, &[rel_flags]);
4751        }
4752        let template = SideHeaderTemplate {
4753            hash,
4754            base_flags,
4755            prefix,
4756            between,
4757            reliable_flags,
4758            reliable_compact,
4759        };
4760        Ok((
4761            template,
4762            ty,
4763            flags,
4764            data_size,
4765            timestamp,
4766            nonce,
4767            reliable_seq_ack,
4768            payload,
4769        ))
4770    }
4771
4772    fn reconstruct_side_compact_frame(
4773        template: &SideHeaderTemplate,
4774        body: &[u8],
4775        timestamp_mode: SideCompactTimestampMode,
4776        timestamp_base: Option<u64>,
4777    ) -> TelemetryResult<(Arc<[u8]>, u64)> {
4778        if body.is_empty() {
4779            return Err(TelemetryError::Unpack("short side compact frame"));
4780        }
4781        let mut off = 0usize;
4782        let flags = body[off];
4783        off += 1;
4784        if (flags & !(SIDE_TRANSPORT_FLAG_PAYLOAD_COMPRESSED | SIDE_TRANSPORT_FLAG_PACKET_NONCE))
4785            != template.base_flags
4786        {
4787            return Err(TelemetryError::Unpack("side compact flags mismatch"));
4788        }
4789        let data_size = Self::read_uleb128_local(body, &mut off)?;
4790        let timestamp = match timestamp_mode {
4791            SideCompactTimestampMode::Absolute => Self::read_uleb128_local(body, &mut off)?,
4792            SideCompactTimestampMode::Delta => {
4793                let timestamp_field = Self::read_uleb128_local(body, &mut off)?;
4794                let base = timestamp_base.ok_or(TelemetryError::Unpack(
4795                    "missing side compact timestamp context",
4796                ))?;
4797                base.checked_add(timestamp_field)
4798                    .ok_or(TelemetryError::Unpack(
4799                        "side compact timestamp delta overflow",
4800                    ))?
4801            }
4802            SideCompactTimestampMode::Omitted => timestamp_base.ok_or(TelemetryError::Unpack(
4803                "missing side compact timestamp context",
4804            ))?,
4805        };
4806        let nonce = if (flags & SIDE_TRANSPORT_FLAG_PACKET_NONCE) != 0 {
4807            Some(Self::read_uleb128_local(body, &mut off)?)
4808        } else {
4809            None
4810        };
4811        let reliable_seq_ack = if template.reliable_flags.is_some() {
4812            let seq = u32::try_from(Self::read_uleb128_local(body, &mut off)?)
4813                .map_err(|_| TelemetryError::Unpack("side compact reliable seq too large"))?;
4814            let ack = u32::try_from(Self::read_uleb128_local(body, &mut off)?)
4815                .map_err(|_| TelemetryError::Unpack("side compact reliable ack too large"))?;
4816            Some((seq, ack))
4817        } else {
4818            None
4819        };
4820        let payload = &body[off..];
4821        let mut raw = Vec::with_capacity(
4822            1 + template.prefix.len() + template.between.len() + payload.len() + 32,
4823        );
4824        raw.push(flags);
4825        raw.extend_from_slice(&template.prefix);
4826        Self::write_uleb128_local(data_size, &mut raw);
4827        Self::write_uleb128_local(timestamp, &mut raw);
4828        if let Some(nonce) = nonce {
4829            Self::write_uleb128_local(nonce, &mut raw);
4830        }
4831        raw.extend_from_slice(&template.between);
4832        if let Some(rel_flags) = template.reliable_flags {
4833            let (seq, ack) =
4834                reliable_seq_ack.ok_or(TelemetryError::Unpack("missing side compact reliable"))?;
4835            wire_format::write_reliable_header_encoded(
4836                wire_format::ReliableHeader {
4837                    flags: rel_flags,
4838                    seq,
4839                    ack,
4840                },
4841                template.reliable_compact,
4842                &mut raw,
4843            );
4844        }
4845        raw.extend_from_slice(payload);
4846        let crc = Self::crc32_bytes(&raw);
4847        raw.extend_from_slice(&crc.to_le_bytes());
4848        Ok((Arc::from(raw), timestamp))
4849    }
4850
4851    fn split_side_transport_frame(
4852        &self,
4853        side: RelaySideId,
4854        frame: Arc<[u8]>,
4855        max_frame_bytes: usize,
4856    ) -> TelemetryResult<Vec<Arc<[u8]>>> {
4857        if max_frame_bytes <= SIDE_TRANSPORT_CHUNK_OVERHEAD {
4858            return Err(TelemetryError::BadArg);
4859        }
4860        let payload_budget = max_frame_bytes - SIDE_TRANSPORT_CHUNK_OVERHEAD;
4861        let mut st = self.state.lock();
4862        let side_state = st
4863            .side_transport
4864            .get_mut(&side)
4865            .ok_or(TelemetryError::BadArg)?;
4866        let transfer_id = side_state.next_chunk_id.wrapping_add(1).max(1);
4867        side_state.next_chunk_id = transfer_id;
4868        drop(st);
4869
4870        let total = frame.len().div_ceil(payload_budget);
4871        let total_u16 =
4872            u16::try_from(total).map_err(|_| TelemetryError::PacketTooLarge("too many chunks"))?;
4873        let mut frames = Vec::with_capacity(total);
4874        for (idx, chunk) in frame.chunks(payload_budget).enumerate() {
4875            let mut body = Vec::with_capacity(8 + chunk.len());
4876            body.extend_from_slice(&transfer_id.to_le_bytes());
4877            body.extend_from_slice(&(idx as u16).to_le_bytes());
4878            body.extend_from_slice(&total_u16.to_le_bytes());
4879            body.extend_from_slice(chunk);
4880            frames.push(Self::wrap_side_transport_frame(
4881                SIDE_TRANSPORT_KIND_CHUNK,
4882                &body,
4883            ));
4884        }
4885        Ok(frames)
4886    }
4887
4888    fn encode_side_transport_frames(
4889        &self,
4890        side: RelaySideId,
4891        opts: RelaySideOptions,
4892        raw: Arc<[u8]>,
4893    ) -> TelemetryResult<Vec<Arc<[u8]>>> {
4894        if !opts.header_template_enabled && opts.max_frame_bytes == 0 {
4895            return Ok(vec![raw]);
4896        }
4897        let raw_len = raw.len();
4898        let mut compact_payload_len = None;
4899        let mut used_compact = false;
4900        let mut used_timestamp_delta = false;
4901        let mut omitted_timestamp = false;
4902        let (template, ty, flags, data_size, timestamp, nonce, reliable_seq_ack, payload) =
4903            Self::extract_side_header_template(raw.as_ref())?;
4904        let (template_id, use_compact, previous_timestamp) = {
4905            let mut st = self.state.lock();
4906            let side_state = st
4907                .side_transport
4908                .get_mut(&side)
4909                .ok_or(TelemetryError::BadArg)?;
4910            if let Some(id) = side_state.tx_template_ids.get(&template.hash).copied() {
4911                let previous = side_state.tx_last_timestamps.get(&id).copied();
4912                (id, true, previous)
4913            } else {
4914                let next = side_state.next_template_id.wrapping_add(1).max(1);
4915                side_state.next_template_id = next;
4916                let evicted = side_state.insert_tx_template(
4917                    template,
4918                    next,
4919                    opts.max_side_transport_templates,
4920                );
4921                if evicted {
4922                    st.side_runtime_stats
4923                        .entry(side)
4924                        .or_default()
4925                        .note_side_transport_template_eviction();
4926                }
4927                if let Some(side_state) = st.side_transport.get_mut(&side) {
4928                    side_state.tx_last_timestamps.insert(next, timestamp);
4929                }
4930                (next, false, None)
4931            }
4932        };
4933        let wrapped = if use_compact {
4934            used_compact = true;
4935            compact_payload_len = Some(payload.len());
4936            let timestamp_field = if let Some(previous) = previous_timestamp {
4937                let delta = timestamp.saturating_sub(previous);
4938                let omit_timestamp = opts.omit_unchanged_compact_timestamps
4939                    || opts.compact_timestamp_omission_types.contains(ty);
4940                if omit_timestamp && timestamp == previous {
4941                    omitted_timestamp = true;
4942                    None
4943                } else if timestamp >= previous
4944                    && Self::uleb128_len_local(delta) < Self::uleb128_len_local(timestamp)
4945                {
4946                    used_timestamp_delta = true;
4947                    Some(delta)
4948                } else {
4949                    Some(timestamp)
4950                }
4951            } else {
4952                Some(timestamp)
4953            };
4954            let mut body = Vec::with_capacity(payload.len() + 32);
4955            body.push(flags);
4956            Self::write_uleb128_local(u64::from(template_id), &mut body);
4957            Self::write_uleb128_local(data_size, &mut body);
4958            if let Some(timestamp_field) = timestamp_field {
4959                Self::write_uleb128_local(timestamp_field, &mut body);
4960            }
4961            if (flags & SIDE_TRANSPORT_FLAG_PACKET_NONCE) != 0 {
4962                Self::write_uleb128_local(u64::from(nonce), &mut body);
4963            }
4964            if let Some((seq, ack)) = reliable_seq_ack {
4965                Self::write_uleb128_local(u64::from(seq), &mut body);
4966                Self::write_uleb128_local(u64::from(ack), &mut body);
4967            }
4968            body.extend_from_slice(payload);
4969            {
4970                let mut st = self.state.lock();
4971                if let Some(side_state) = st.side_transport.get_mut(&side) {
4972                    side_state.tx_last_timestamps.insert(template_id, timestamp);
4973                }
4974            }
4975            let kind = if omitted_timestamp {
4976                SIDE_TRANSPORT_KIND_COMPACT_SAME_TIMESTAMP
4977            } else if used_timestamp_delta {
4978                SIDE_TRANSPORT_KIND_COMPACT_DELTA
4979            } else {
4980                SIDE_TRANSPORT_KIND_COMPACT
4981            };
4982            Self::wrap_side_transport_frame(kind, &body)
4983        } else {
4984            let mut body = Vec::with_capacity(raw.len() + 4);
4985            Self::write_uleb128_local(u64::from(template_id), &mut body);
4986            body.extend_from_slice(raw.as_ref());
4987            Self::wrap_side_transport_frame(SIDE_TRANSPORT_KIND_FULL, &body)
4988        };
4989        let frames = if opts.max_frame_bytes != 0 && wrapped.len() > opts.max_frame_bytes {
4990            self.split_side_transport_frame(side, wrapped, opts.max_frame_bytes)
4991        } else {
4992            Ok(vec![wrapped])
4993        }?;
4994        let wire_len = frames.iter().map(|frame| frame.len()).sum::<usize>();
4995        let mut st = self.state.lock();
4996        let stats = st.side_runtime_stats.entry(side).or_default();
4997        if used_compact {
4998            let overhead = compact_payload_len
4999                .map(|payload_len| wire_len.saturating_sub(payload_len))
5000                .unwrap_or(wire_len);
5001            stats.note_side_transport_compact(
5002                raw_len,
5003                wire_len,
5004                overhead,
5005                used_timestamp_delta,
5006                omitted_timestamp,
5007            );
5008            if opts.compact_header_target_bytes != 0 && overhead > opts.compact_header_target_bytes
5009            {
5010                stats.note_side_transport_compact_target_miss();
5011            }
5012        } else {
5013            stats.note_side_transport_full(raw_len, wire_len);
5014        }
5015        if frames.len() > 1 {
5016            stats.note_side_transport_chunks(frames.len());
5017        }
5018        Ok(frames)
5019    }
5020
5021    fn decode_side_transport_frame(
5022        &self,
5023        side: RelaySideId,
5024        bytes: &[u8],
5025    ) -> TelemetryResult<Option<Arc<[u8]>>> {
5026        let Some((kind, body)) = Self::parse_side_transport_wrapper(bytes)? else {
5027            return Ok(Some(Arc::from(bytes)));
5028        };
5029        match kind {
5030            SIDE_TRANSPORT_KIND_FULL => {
5031                let mut off = 0usize;
5032                let template_id = u32::try_from(Self::read_uleb128_local(body, &mut off)?)
5033                    .map_err(|_| TelemetryError::Unpack("side template id too large"))?;
5034                let raw = Arc::<[u8]>::from(&body[off..]);
5035                if let Ok((template, _, _, _, timestamp, _, _, _)) =
5036                    Self::extract_side_header_template(raw.as_ref())
5037                {
5038                    let mut st = self.state.lock();
5039                    let max_templates = st
5040                        .sides
5041                        .get(side)
5042                        .and_then(|side| side.as_ref())
5043                        .map(|side| side.opts.max_side_transport_templates)
5044                        .unwrap_or(DEFAULT_SIDE_TRANSPORT_TEMPLATE_LIMIT);
5045                    let evicted = st.side_transport.get_mut(&side).is_some_and(|side_state| {
5046                        let evicted =
5047                            side_state.insert_rx_template(template_id, template, max_templates);
5048                        side_state.rx_last_timestamps.insert(template_id, timestamp);
5049                        evicted
5050                    });
5051                    if evicted {
5052                        st.side_runtime_stats
5053                            .entry(side)
5054                            .or_default()
5055                            .note_side_transport_template_eviction();
5056                    }
5057                }
5058                Ok(Some(raw))
5059            }
5060            SIDE_TRANSPORT_KIND_COMPACT
5061            | SIDE_TRANSPORT_KIND_COMPACT_DELTA
5062            | SIDE_TRANSPORT_KIND_COMPACT_SAME_TIMESTAMP => {
5063                if body.is_empty() {
5064                    return Err(TelemetryError::Unpack("short side compact frame"));
5065                }
5066                let mut off = 1usize;
5067                let template_id = u32::try_from(Self::read_uleb128_local(body, &mut off)?)
5068                    .map_err(|_| TelemetryError::Unpack("side template id too large"))?;
5069                let mut compact_body = Vec::with_capacity(1 + body.len().saturating_sub(off));
5070                compact_body.push(body[0]);
5071                compact_body.extend_from_slice(&body[off..]);
5072                let (template, timestamp_base) = {
5073                    let st = self.state.lock();
5074                    let state = st.side_transport.get(&side);
5075                    let template = state
5076                        .and_then(|state| state.rx_templates_by_id.get(&template_id))
5077                        .cloned();
5078                    let timestamp_base = if matches!(
5079                        kind,
5080                        SIDE_TRANSPORT_KIND_COMPACT_DELTA
5081                            | SIDE_TRANSPORT_KIND_COMPACT_SAME_TIMESTAMP
5082                    ) {
5083                        state
5084                            .and_then(|state| state.rx_last_timestamps.get(&template_id))
5085                            .copied()
5086                    } else {
5087                        None
5088                    };
5089                    (template, timestamp_base)
5090                };
5091                let template =
5092                    template.ok_or(TelemetryError::Unpack("unknown side compact template"))?;
5093                let timestamp_mode = match kind {
5094                    SIDE_TRANSPORT_KIND_COMPACT_DELTA => SideCompactTimestampMode::Delta,
5095                    SIDE_TRANSPORT_KIND_COMPACT_SAME_TIMESTAMP => SideCompactTimestampMode::Omitted,
5096                    _ => SideCompactTimestampMode::Absolute,
5097                };
5098                let (frame, timestamp) = Self::reconstruct_side_compact_frame(
5099                    &template,
5100                    &compact_body,
5101                    timestamp_mode,
5102                    timestamp_base,
5103                )?;
5104                let mut st = self.state.lock();
5105                if let Some(side_state) = st.side_transport.get_mut(&side) {
5106                    side_state.rx_last_timestamps.insert(template_id, timestamp);
5107                }
5108                Ok(Some(frame))
5109            }
5110            SIDE_TRANSPORT_KIND_CHUNK => {
5111                if body.len() < 8 {
5112                    return Err(TelemetryError::Unpack("short side chunk frame"));
5113                }
5114                let transfer_id = u32::from_le_bytes([body[0], body[1], body[2], body[3]]);
5115                let index = u16::from_le_bytes([body[4], body[5]]);
5116                let total = u16::from_le_bytes([body[6], body[7]]);
5117                let payload = Arc::<[u8]>::from(&body[8..]);
5118                let assembled = {
5119                    let mut st = self.state.lock();
5120                    let side_state = st
5121                        .side_transport
5122                        .get_mut(&side)
5123                        .ok_or(TelemetryError::BadArg)?;
5124                    let entry = side_state.rx_chunks.entry(transfer_id).or_default();
5125                    if entry.total == 0 {
5126                        entry.total = total;
5127                    } else if entry.total != total {
5128                        side_state.rx_chunks.remove(&transfer_id);
5129                        return Err(TelemetryError::Unpack("side chunk total mismatch"));
5130                    }
5131                    entry.received.entry(index).or_insert(payload);
5132                    if entry.received.len() == usize::from(total) {
5133                        let entry = side_state
5134                            .rx_chunks
5135                            .remove(&transfer_id)
5136                            .ok_or(TelemetryError::Unpack("side chunk missing"))?;
5137                        let mut out = Vec::new();
5138                        for idx in 0..entry.total {
5139                            let chunk = entry
5140                                .received
5141                                .get(&idx)
5142                                .ok_or(TelemetryError::Unpack("side chunk gap"))?;
5143                            out.extend_from_slice(chunk);
5144                        }
5145                        Some(Arc::<[u8]>::from(out))
5146                    } else {
5147                        None
5148                    }
5149                };
5150                match assembled {
5151                    Some(frame) => self.decode_side_transport_frame(side, frame.as_ref()),
5152                    None => Ok(None),
5153                }
5154            }
5155            _ => Err(TelemetryError::Unpack("unknown side transport frame")),
5156        }
5157    }
5158
5159    /// Helper: call a TX handler with the best representation we have.
5160    /// - Packet handler + Packet item: direct.
5161    /// - Packed handler + Packed item: direct.
5162    /// - Packet handler + Packed item: unpack for this call.
5163    /// - Packed handler + Packet item: pack for this call.
5164    fn call_tx_handler(
5165        &self,
5166        side: RelaySideId,
5167        handler: &RelayTxHandlerFn,
5168        data: &RelayItem,
5169    ) -> TelemetryResult<()> {
5170        let opts = {
5171            let st = self.state.lock();
5172            Self::side_ref(&st, side)?.opts
5173        };
5174        let Some(_side_tx_guard) = self.try_enter_side_tx() else {
5175            return Err(TelemetryError::Io("side tx busy"));
5176        };
5177        let started_ms = self.clock.now_ms();
5178        let ty = match data {
5179            RelayItem::Packet(pkt) => pkt.data_type(),
5180            RelayItem::Packed(bytes) => wire_format::peek_envelope(bytes.as_ref())?.ty,
5181        };
5182        let result = match (handler, data) {
5183            // Fast paths
5184            (RelayTxHandlerFn::Packed(f), RelayItem::Packed(bytes)) => {
5185                let frames = self.encode_side_transport_frames(side, opts, bytes.clone())?;
5186                let mut sent_bytes = 0usize;
5187                for frame in frames {
5188                    f(frame.as_ref())?;
5189                    sent_bytes = sent_bytes.saturating_add(frame.len());
5190                }
5191                self.record_side_tx_sample(side, sent_bytes, started_ms, self.clock.now_ms());
5192                self.note_side_tx_success(side, ty, sent_bytes, 1);
5193                return Ok(());
5194            }
5195            (RelayTxHandlerFn::Packet(f), RelayItem::Packet(pkt)) => f(pkt),
5196
5197            // Conversion paths
5198            (RelayTxHandlerFn::Packed(f), RelayItem::Packet(pkt)) => {
5199                let owned = wire_format::pack_packet(pkt);
5200                let frames = self.encode_side_transport_frames(side, opts, owned)?;
5201                let mut sent_bytes = 0usize;
5202                for frame in frames {
5203                    f(frame.as_ref())?;
5204                    sent_bytes = sent_bytes.saturating_add(frame.len());
5205                }
5206                self.record_side_tx_sample(side, sent_bytes, started_ms, self.clock.now_ms());
5207                self.note_side_tx_success(side, ty, sent_bytes, 1);
5208                return Ok(());
5209            }
5210            (RelayTxHandlerFn::Packet(f), RelayItem::Packed(bytes)) => {
5211                if wire_format::peek_frame_info(bytes.as_ref())
5212                    .ok()
5213                    .is_some_and(|frame| frame.ack_only())
5214                {
5215                    return Ok(());
5216                }
5217                let pkt = wire_format::unpack_packet(bytes.as_ref())?;
5218                f(&pkt)
5219            }
5220        };
5221        if result.is_ok()
5222            && let Ok(bytes) = Self::relay_item_wire_len(data)
5223        {
5224            self.record_side_tx_sample(side, bytes, started_ms, self.clock.now_ms());
5225            self.note_side_tx_success(side, ty, bytes, 1);
5226        } else if result.is_err() {
5227            self.note_side_tx_failure(side, ty, 1);
5228        }
5229        result
5230    }
5231
5232    fn adjust_reliable_for_side(
5233        &self,
5234        opts: RelaySideOptions,
5235        data: RelayItem,
5236    ) -> TelemetryResult<Option<RelayItem>> {
5237        if opts.reliable_enabled {
5238            return Ok(Some(data));
5239        }
5240
5241        match data {
5242            RelayItem::Packed(bytes) => {
5243                let frame = match wire_format::peek_frame_info(bytes.as_ref()) {
5244                    Ok(frame) => frame,
5245                    Err(_) => return Ok(Some(RelayItem::Packed(bytes))),
5246                };
5247                if is_reliable_type(frame.envelope.ty)
5248                    && let Some(hdr) = frame.reliable
5249                {
5250                    if (hdr.flags & wire_format::RELIABLE_FLAG_ACK_ONLY) != 0 {
5251                        return Ok(None);
5252                    }
5253                    if (hdr.flags & wire_format::RELIABLE_FLAG_UNSEQUENCED) == 0 {
5254                        let Some(rewritten) = wire_format::rewrite_reliable_header_owned(
5255                            bytes.as_ref(),
5256                            wire_format::RELIABLE_FLAG_UNSEQUENCED,
5257                            hdr.seq,
5258                            0,
5259                        )?
5260                        else {
5261                            return Ok(Some(RelayItem::Packed(bytes)));
5262                        };
5263                        return Ok(Some(RelayItem::Packed(rewritten)));
5264                    }
5265                }
5266                Ok(Some(RelayItem::Packed(bytes)))
5267            }
5268            RelayItem::Packet(pkt) => {
5269                if matches!(
5270                    pkt.data_type(),
5271                    crate::DataType::ReliableAck
5272                        | crate::DataType::ReliablePartialAck
5273                        | crate::DataType::ReliablePacketRequest
5274                ) {
5275                    return Ok(None);
5276                }
5277                Ok(Some(RelayItem::Packet(pkt)))
5278            }
5279        }
5280    }
5281
5282    /// Drain the RX queue fully, expanding to TX items.
5283    #[inline]
5284    pub fn process_rx_queue(&self) -> TelemetryResult<()> {
5285        self.process_rx_queue_with_timeout(0)
5286    }
5287
5288    /// Drain the TX queue fully, invoking per-side TX handlers.
5289    ///
5290    /// If called from inside a side TX callback, this becomes a no-op so relay TX handlers cannot
5291    /// recurse into nested queue drains on the same stack.
5292    #[inline]
5293    pub fn process_tx_queue(&self) -> TelemetryResult<()> {
5294        self.process_tx_queue_with_timeout(0)
5295    }
5296
5297    /// Drain RX then TX queues fully (one pass).
5298    #[inline]
5299    pub fn process_all_queues(&self) -> TelemetryResult<()> {
5300        self.process_all_queues_with_timeout(0)
5301    }
5302
5303    /// Process the TX queue for up to `timeout_ms` milliseconds.
5304    ///
5305    /// `timeout_ms == 0` drains fully. If called from inside a side TX callback, this becomes a
5306    /// no-op so relay TX handlers cannot recurse into nested queue drains on the same stack.
5307    pub fn process_tx_queue_with_timeout(&self, timeout_ms: u32) -> TelemetryResult<()> {
5308        if self.side_tx_active() {
5309            return Ok(());
5310        }
5311        #[cfg(feature = "discovery")]
5312        {
5313            let _ = self.poll_discovery()?;
5314        }
5315        let start = self.clock.now_ms();
5316        loop {
5317            self.process_reliable_timeouts()?;
5318            if self.process_replay_queue_item()? {
5319                if timeout_ms != 0 && self.clock.now_ms().wrapping_sub(start) >= timeout_ms as u64 {
5320                    break;
5321                }
5322                continue;
5323            }
5324            let Some((src, dst, handler, opts, data)) = self.pop_ready_tx_item() else {
5325                break;
5326            };
5327            match self.send_tx_item(src, dst, handler, opts, data.clone()) {
5328                Ok(sent) => {
5329                    if sent
5330                        && timeout_ms != 0
5331                        && self.clock.now_ms().wrapping_sub(start) >= timeout_ms as u64
5332                    {
5333                        break;
5334                    }
5335                }
5336                Err(e) if Self::is_side_tx_busy(&e) => {
5337                    let priority = Self::relay_item_priority(&data)?;
5338                    let mut st = self.state.lock();
5339                    st.push_tx(RelayTxItem {
5340                        src,
5341                        dst,
5342                        data,
5343                        priority,
5344                    })?;
5345                    break;
5346                }
5347                Err(e) => return Err(e),
5348            }
5349        }
5350        Ok(())
5351    }
5352
5353    /// Process RX queue with timeout.
5354    pub fn process_rx_queue_with_timeout(&self, timeout_ms: u32) -> TelemetryResult<()> {
5355        #[cfg(feature = "discovery")]
5356        {
5357            let _ = self.poll_discovery()?;
5358        }
5359        let start = self.clock.now_ms();
5360        loop {
5361            let item_opt = {
5362                let mut st = self.state.lock();
5363                st.rx_queue.pop_front()
5364            };
5365            let Some(item) = item_opt else { break };
5366            self.process_rx_queue_item(item)?;
5367
5368            if timeout_ms != 0 && self.clock.now_ms().wrapping_sub(start) >= timeout_ms as u64 {
5369                break;
5370            }
5371        }
5372        Ok(())
5373    }
5374
5375    /// Process RX and TX queues interleaved for up to `timeout_ms` milliseconds.
5376    ///
5377    /// `timeout_ms == 0` drains fully. If called from inside a side TX callback, this becomes a
5378    /// no-op so relay TX handlers cannot recurse into nested queue drains on the same stack.
5379    pub fn process_all_queues_with_timeout(&self, timeout_ms: u32) -> TelemetryResult<()> {
5380        if self.side_tx_active() {
5381            return Ok(());
5382        }
5383        #[cfg(feature = "discovery")]
5384        {
5385            let _ = self.poll_discovery()?;
5386        }
5387        let drain_fully = timeout_ms == 0;
5388        let start = if drain_fully { 0 } else { self.clock.now_ms() };
5389
5390        loop {
5391            let mut did_any = false;
5392            self.process_reliable_timeouts()?;
5393
5394            // First move RX → TX
5395            if let Some(item) = {
5396                let mut st = self.state.lock();
5397                st.rx_queue.pop_front()
5398            } {
5399                self.process_rx_queue_item(item)?;
5400                did_any = true;
5401            }
5402
5403            if !drain_fully && self.clock.now_ms().wrapping_sub(start) >= timeout_ms as u64 {
5404                break;
5405            }
5406
5407            if self.process_replay_queue_item()? {
5408                did_any = true;
5409            }
5410
5411            // Then send out TX
5412            let sent_one = if let Some((src, dst, handler, opts, data)) = self.pop_ready_tx_item() {
5413                self.send_tx_item(src, dst, handler, opts, data)?
5414            } else {
5415                false
5416            };
5417
5418            if sent_one {
5419                did_any = true;
5420            }
5421
5422            if !drain_fully && self.clock.now_ms().wrapping_sub(start) >= timeout_ms as u64 {
5423                break;
5424            }
5425
5426            if !did_any {
5427                break;
5428            }
5429        }
5430
5431        Ok(())
5432    }
5433
5434    /// Runs one application-loop maintenance cycle.
5435    ///
5436    /// This polls built-in discovery when that feature is compiled in, then drains queued RX/TX
5437    /// work for up to `timeout_ms` milliseconds.
5438    pub fn periodic(&self, timeout_ms: u32) -> TelemetryResult<()> {
5439        #[cfg(feature = "discovery")]
5440        {
5441            let _ = self.poll_discovery()?;
5442        }
5443
5444        self.process_all_queues_with_timeout(timeout_ms)
5445    }
5446}