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
42pub 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;
65type PacketHandlerFn = dyn Fn(&Packet) -> TelemetryResult<()> + Send + Sync + 'static;
67
68type PackedHandlerFn = dyn Fn(&[u8]) -> TelemetryResult<()> + Send + Sync + 'static;
70
71#[derive(Clone)]
73pub enum RelayTxHandlerFn {
74 Packed(Arc<PackedHandlerFn>),
75 Packet(Arc<PacketHandlerFn>),
76}
77
78#[derive(Clone, Copy, Debug)]
79pub struct RelaySideOptions {
80 pub reliable_enabled: bool,
86 pub link_local_enabled: bool,
88 pub ingress_enabled: bool,
90 pub egress_enabled: bool,
92 pub header_template_enabled: bool,
94 pub max_frame_bytes: usize,
101 pub compact_header_target_bytes: usize,
103 pub max_side_transport_templates: usize,
105 pub omit_unchanged_compact_timestamps: bool,
107 pub compact_timestamp_omission_types: CompactTimestampOmissionPolicy,
109 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 #[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#[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#[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#[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#[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
779struct 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
1116pub 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 pub fn new(clock: Box<dyn Clock + Send + Sync>) -> Self {
1250 Self::new_with_config(RelayConfig::default(), clock)
1251 }
1252
1253 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 #[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 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 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 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 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 #[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 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 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 let h: u64 = 0x9E37_79B9_7F4A_7C15;
3372 hash_bytes_u64(h, bytes.as_ref())
3373 }
3374 }
3375 }
3376 }
3377 }
3378
3379 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn announce_discovery(&self) -> TelemetryResult<()> {
3804 self.queue_discovery_announce(true)
3805 }
3806
3807 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 pub fn poll_discovery(&self) -> TelemetryResult<bool> {
3833 self.poll_discovery_announce()
3834 }
3835
3836 #[cfg(feature = "discovery")]
3837 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 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 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 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 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 pub fn clear_rx_queue(&self) {
4304 let mut st = self.state.lock();
4305 st.rx_queue.clear();
4306 }
4307
4308 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 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 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; 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 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 (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 (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 #[inline]
5284 pub fn process_rx_queue(&self) -> TelemetryResult<()> {
5285 self.process_rx_queue_with_timeout(0)
5286 }
5287
5288 #[inline]
5293 pub fn process_tx_queue(&self) -> TelemetryResult<()> {
5294 self.process_tx_queue_with_timeout(0)
5295 }
5296
5297 #[inline]
5299 pub fn process_all_queues(&self) -> TelemetryResult<()> {
5300 self.process_all_queues_with_timeout(0)
5301 }
5302
5303 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 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 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 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 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 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}