From d2432740c1718edf3bc6bcd11be4f8122dd98e76 Mon Sep 17 00:00:00 2001 From: Mollusk Date: Wed, 8 Jul 2026 15:15:22 -0400 Subject: [PATCH] network: per-peer connection badge (direct/relay, RTT, loss, bitrate) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Answer "am I actually P2P right now?" per peer. A 1 Hz session task snapshots the selected QUIC path of every live audio connection (IrohTransport::connection_stats), core::connstats::derive turns consecutive snapshots into RTT/loss/bitrate (path switches and counter resets invalidate the rate window), and the peer card shows a Direct/Relay badge with a hover tooltip for address, loss, and up/down bitrate. No new dependencies, no wire change. Loopback-integration-tested against real iroh endpoints; not yet field-verified on a 2-machine call (FEATURES.md row marked πŸ§ͺ). Co-Authored-By: Claude Fable 5 --- docs/FEATURES.md | 1 + src/app/mod.rs | 121 +++++++++++++++++++++++ src/core/connstats.rs | 187 ++++++++++++++++++++++++++++++++++++ src/core/messages.rs | 5 + src/core/mod.rs | 38 ++++++++ src/network/iroh_impl.rs | 51 ++++++++++ src/network/mod.rs | 26 +++++ tests/transport_loopback.rs | 65 +++++++++++++ 8 files changed, 494 insertions(+) create mode 100644 src/core/connstats.rs diff --git a/docs/FEATURES.md b/docs/FEATURES.md index e799e55..92d34b1 100644 --- a/docs/FEATURES.md +++ b/docs/FEATURES.md @@ -102,6 +102,7 @@ covers internals). When you ship a feature, add it here. | iroh QUIC transport | βœ… | | | Network mode picker | βœ… | `RelayNoDiscovery` (default), `N0Full`, `DirectOnly`. Takes effect next join. | | Retained-address reconnect | βœ… | Dials last-known full addr before falling back to bare id. | +| Per-peer connection badge (direct/relay + RTT, hover for addr/loss/bitrate) | πŸ§ͺ | Peer-card badge fed by a 1 Hz poll of the live audio link's selected QUIC path (`connection_stats` β†’ `core::connstats::derive`). Loopback-integration-tested; not field-tested on a real 2-machine call. | | Reconnect + eviction model | βœ… | Incl. two-outage reconnect-eviction fix + regression test. | | Self-hosted relay | ❌ | Decided against β€” rely on n0 relays, `RelayNoDiscovery` default. | diff --git a/src/app/mod.rs b/src/app/mod.rs index f95bc5e..e0c7179 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -869,6 +869,10 @@ pub struct AppState { game_override: GameOverrideChoice, peers: HashMap, audio_levels: HashMap, + /// Latest per-peer connection transparency info (direct/relay, RTT, window + /// loss/bitrate), replaced wholesale by each `UiEvent::ConnectionStats` + /// (~1/sec). A peer with no entry has no live audio link right now. + conn_stats: HashMap, /// Peers we've locally muted (their audio isn't mixed into our output). locally_muted: HashSet, /// When we joined the current room, for the in-room call-duration timer. @@ -1046,6 +1050,7 @@ impl AppState { self.music_prefetch_inflight = None; self.peers.clear(); self.audio_levels.clear(); + self.conn_stats.clear(); self.locally_muted.clear(); self.chat_messages.clear(); self.chat_input.clear(); @@ -1209,6 +1214,7 @@ impl Default for AppState { game_override: GameOverrideChoice::Auto, peers: HashMap::new(), audio_levels: HashMap::new(), + conn_stats: HashMap::new(), locally_muted: HashSet::new(), call_started: None, recording: false, @@ -1686,6 +1692,37 @@ fn pan_label(pan: f32) -> String { } } +/// Connection badge text on the peer card: path type + RTT ("Direct Β· 12 ms"). +fn conn_badge_label(info: &crate::core::connstats::PeerConnInfo) -> String { + let kind = if info.relay { "Relay" } else { "Direct" }; + format!("{kind} Β· {} ms", info.rtt_ms) +} + +/// First tooltip line: path type + remote address ("Direct (1.2.3.4:5)" / +/// "Relay (https://relay.example./)"). +fn conn_tooltip_path(info: &crate::core::connstats::PeerConnInfo) -> String { + let kind = if info.relay { "Relay" } else { "Direct" }; + format!("{kind} ({})", info.remote_addr) +} + +/// Loss line for the tooltip. `None` (first poll / idle window) reads as clean. +fn conn_loss_label(loss_pct: Option) -> String { + match loss_pct { + Some(pct) => format!("Loss {:.1}% (last second)", pct.clamp(0.0, 100.0)), + None => "Loss β€” (last second)".to_string(), + } +} + +/// One direction of the bitrate line ("32 kbps", "1.5 Mbps", or "β€”" until a +/// full poll window has elapsed on the current path). +fn conn_rate_label(kbps: Option) -> String { + match kbps { + None => "β€”".to_string(), + Some(k) if k >= 1000.0 => format!("{:.1} Mbps", k / 1000.0), + Some(k) => format!("{k:.0} kbps"), + } +} + fn update(state: &mut AppState, message: AppMessage) -> Task { match message { AppMessage::NicknameChanged(val) => { @@ -1934,6 +1971,11 @@ fn update(state: &mut AppState, message: AppMessage) -> Task { state.audio_levels.insert(id, val); } } + UiEvent::ConnectionStats(infos) => { + // Full replacement: a peer missing from this round has no + // live link, so its (stale) badge must go away too. + state.conn_stats = infos.into_iter().collect(); + } UiEvent::MicLevel(level) => { state.mic_level = level; } @@ -6192,6 +6234,41 @@ fn view(state: &AppState) -> Element<'_, AppMessage> { name_col = name_col .push(text(format!("Playing {game}")).size(11).color(color_blue)); } + // Connection-transparency badge: path type + RTT, with + // the full story (address, loss, bitrate) on hover. + // Rendered only while the audio link is live β€” the + // Connecting/Reconnecting indicator covers the rest. + if let Some(info) = state.conn_stats.get(peer_id) { + let dot_color = if info.relay { color_yellow } else { color_green }; + let badge = row![ + text("●").size(9).color(dot_color), + text(conn_badge_label(info)).size(11).color(color_subtext), + ] + .spacing(4) + .align_y(iced::alignment::Vertical::Center); + let detail = column![ + text(conn_tooltip_path(info)).size(11).color(color_text), + text(conn_loss_label(info.loss_pct)).size(11).color(color_subtext), + text(format!( + "↑ {} ↓ {}", + conn_rate_label(info.up_kbps), + conn_rate_label(info.down_kbps) + )) + .size(11) + .color(color_subtext), + ] + .spacing(2); + name_col = name_col.push( + tooltip( + badge, + container(detail) + .padding(8) + .style(c_style(color_crust, color_surface, 6.0)), + iced::widget::tooltip::Position::Bottom, + ) + .gap(6), + ); + } name_col }, add_friend_el, @@ -9100,6 +9177,17 @@ mod tests { state.invalid_audio.insert(attachment_id); state.connecting.insert(peer); state.ever_connected.insert(peer); + state.conn_stats.insert( + peer, + crate::core::connstats::PeerConnInfo { + relay: false, + remote_addr: "1.2.3.4:5".to_string(), + rtt_ms: 12, + loss_pct: None, + up_kbps: None, + down_kbps: None, + }, + ); state.recording = true; state.recording_started = Some(now); state.call_started = Some(now); @@ -9139,6 +9227,7 @@ mod tests { assert!(state.invalid_audio.is_empty()); assert!(state.connecting.is_empty()); assert!(state.ever_connected.is_empty()); + assert!(state.conn_stats.is_empty()); assert!(!state.recording); assert!(state.recording_started.is_none()); assert!(state.call_started.is_none()); @@ -9786,6 +9875,38 @@ mod tests { assert!(ch.is_finite() && ch >= CHAT_MIN_H); } + #[test] + fn conn_badge_and_tooltip_labels() { + use super::{conn_badge_label, conn_loss_label, conn_rate_label, conn_tooltip_path}; + let direct = crate::core::connstats::PeerConnInfo { + relay: false, + remote_addr: "192.168.1.7:53340".to_string(), + rtt_ms: 12, + loss_pct: Some(0.44), + up_kbps: Some(32.4), + down_kbps: None, + }; + assert_eq!(conn_badge_label(&direct), "Direct Β· 12 ms"); + assert_eq!(conn_tooltip_path(&direct), "Direct (192.168.1.7:53340)"); + + let relay = crate::core::connstats::PeerConnInfo { + relay: true, + remote_addr: "https://relay.example./".to_string(), + ..direct.clone() + }; + assert_eq!(conn_badge_label(&relay), "Relay Β· 12 ms"); + assert_eq!(conn_tooltip_path(&relay), "Relay (https://relay.example./)"); + + assert_eq!(conn_loss_label(Some(0.44)), "Loss 0.4% (last second)"); + // Out-of-range inputs clamp instead of reading nonsense. + assert_eq!(conn_loss_label(Some(250.0)), "Loss 100.0% (last second)"); + assert_eq!(conn_loss_label(None), "Loss β€” (last second)"); + + assert_eq!(conn_rate_label(Some(32.4)), "32 kbps"); + assert_eq!(conn_rate_label(Some(1500.0)), "1.5 Mbps"); + assert_eq!(conn_rate_label(None), "β€”"); + } + #[test] fn controls_and_drawer_width_clamps() { use super::{CHAT_MIN_W, CONTROLS_MIN_W, clamp_chat_drawer_width, clamp_controls_width}; diff --git a/src/core/connstats.rs b/src/core/connstats.rs new file mode 100644 index 0000000..4810b66 --- /dev/null +++ b/src/core/connstats.rs @@ -0,0 +1,187 @@ +//! Per-peer connection-transparency derivation. +//! +//! The transport hands us cumulative counters for each peer's selected QUIC +//! path ([`PathSnapshot`]); this module turns two consecutive snapshots into +//! the human-facing [`PeerConnInfo`] the UI renders (badge + tooltip): path +//! type, RTT, and loss/bitrate over the poll window. Pure functions only β€” +//! the polling task in `core::mod` owns the clock and the previous-snapshot +//! map. + +use crate::network::PathSnapshot; +use std::time::Duration; + +/// How often the core polls the transport for path snapshots. +pub const POLL_INTERVAL: Duration = Duration::from_secs(1); + +/// Derived, display-ready connection info for one peer, sent to the UI via +/// `UiEvent::ConnectionStats`. Window-relative fields are `None` when they +/// can't be derived yet (first poll, path switch, or an idle window). +#[derive(Debug, Clone, PartialEq)] +pub struct PeerConnInfo { + /// True = relayed path, false = direct IP path. + pub relay: bool, + /// `ip:port` for a direct path, the relay URL for a relayed one. + pub remote_addr: String, + /// Path round-trip time, rounded to whole milliseconds. + pub rtt_ms: u32, + /// Percentage of packets sent in the window that were detected lost. + pub loss_pct: Option, + /// Outbound bitrate over the window, kilobits per second. + pub up_kbps: Option, + /// Inbound bitrate over the window, kilobits per second. + pub down_kbps: Option, +} + +/// Derive display info from the current snapshot and (when comparable) the +/// previous one. `prev` is comparable only if it's the same path β€” a relayβ†’ +/// direct migration or a reconnect resets the counters, so those windows +/// yield `None` rates rather than garbage (negative deltas show up as +/// `cur < prev` and are treated the same way). +pub fn derive(prev: Option<&PathSnapshot>, cur: &PathSnapshot, elapsed: Duration) -> PeerConnInfo { + let rates = prev + .filter(|p| comparable(p, cur)) + .and_then(|p| window_rates(p, cur, elapsed)); + PeerConnInfo { + relay: cur.is_relay, + remote_addr: cur.remote_addr.clone(), + rtt_ms: cur.rtt.as_millis().min(u128::from(u32::MAX)) as u32, + loss_pct: rates.and_then(|r| r.loss_pct), + up_kbps: rates.map(|r| r.up_kbps), + down_kbps: rates.map(|r| r.down_kbps), + } +} + +/// True when `cur`'s counters continue `prev`'s: same path (address) and +/// monotonically non-decreasing counters (a reconnect on the same address +/// restarts them from zero). +fn comparable(prev: &PathSnapshot, cur: &PathSnapshot) -> bool { + prev.remote_addr == cur.remote_addr + && cur.tx_bytes >= prev.tx_bytes + && cur.rx_bytes >= prev.rx_bytes + && cur.tx_datagrams >= prev.tx_datagrams + && cur.lost_packets >= prev.lost_packets +} + +#[derive(Debug, Clone, Copy)] +struct WindowRates { + loss_pct: Option, + up_kbps: f32, + down_kbps: f32, +} + +fn window_rates(prev: &PathSnapshot, cur: &PathSnapshot, elapsed: Duration) -> Option { + let secs = elapsed.as_secs_f64(); + if secs <= 0.0 { + return None; + } + let sent = cur.tx_datagrams - prev.tx_datagrams; + let lost = cur.lost_packets - prev.lost_packets; + // Loss detection lags sending (it needs ACK timeouts), so a window can see + // more losses than sends; clamp to 100% rather than exceeding it. An idle + // window (nothing sent or lost) has no loss story to tell. + let loss_pct = if sent == 0 && lost == 0 { + None + } else { + Some(((lost as f64 / (sent.max(lost)) as f64) * 100.0) as f32) + }; + let kbps = |bytes: u64| ((bytes as f64 * 8.0 / 1000.0) / secs) as f32; + Some(WindowRates { + loss_pct, + up_kbps: kbps(cur.tx_bytes - prev.tx_bytes), + down_kbps: kbps(cur.rx_bytes - prev.rx_bytes), + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn snap(addr: &str, tx_b: u64, rx_b: u64, tx_d: u64, lost: u64) -> PathSnapshot { + PathSnapshot { + is_relay: false, + remote_addr: addr.to_string(), + rtt: Duration::from_millis(12), + tx_bytes: tx_b, + rx_bytes: rx_b, + tx_datagrams: tx_d, + lost_packets: lost, + } + } + + #[test] + fn first_poll_has_type_and_rtt_but_no_rates() { + let cur = snap("1.2.3.4:5", 1000, 2000, 50, 0); + let info = derive(None, &cur, POLL_INTERVAL); + assert_eq!(info.rtt_ms, 12); + assert!(!info.relay); + assert_eq!(info.remote_addr, "1.2.3.4:5"); + assert_eq!(info.loss_pct, None); + assert_eq!(info.up_kbps, None); + assert_eq!(info.down_kbps, None); + } + + #[test] + fn steady_window_yields_rates_and_loss() { + let prev = snap("1.2.3.4:5", 0, 0, 0, 0); + // 1s window: 4000 bytes up (32 kbps), 2000 down (16 kbps), 2 of 100 lost. + let cur = snap("1.2.3.4:5", 4000, 2000, 100, 2); + let info = derive(Some(&prev), &cur, Duration::from_secs(1)); + assert_eq!(info.up_kbps, Some(32.0)); + assert_eq!(info.down_kbps, Some(16.0)); + assert_eq!(info.loss_pct, Some(2.0)); + } + + #[test] + fn idle_window_has_no_loss_story() { + let prev = snap("1.2.3.4:5", 4000, 2000, 100, 2); + let cur = prev.clone(); + let info = derive(Some(&prev), &cur, Duration::from_secs(1)); + assert_eq!(info.loss_pct, None); + assert_eq!(info.up_kbps, Some(0.0)); + } + + #[test] + fn loss_detected_in_an_idle_window_clamps_to_full() { + // Losses can be *detected* after sending stops (ACK timeouts fire late). + let prev = snap("1.2.3.4:5", 4000, 2000, 100, 0); + let cur = snap("1.2.3.4:5", 4000, 2000, 100, 3); + let info = derive(Some(&prev), &cur, Duration::from_secs(1)); + assert_eq!(info.loss_pct, Some(100.0)); + } + + #[test] + fn path_switch_resets_the_window() { + let prev = snap("relay.example:443", 9000, 9000, 900, 5); + let cur = snap("1.2.3.4:5", 100, 100, 10, 0); + let info = derive(Some(&prev), &cur, Duration::from_secs(1)); + assert_eq!(info.up_kbps, None); + assert_eq!(info.loss_pct, None); + } + + #[test] + fn counter_reset_on_same_address_resets_the_window() { + // Same address but the connection was rebuilt β†’ counters restarted. + let prev = snap("1.2.3.4:5", 9000, 9000, 900, 5); + let cur = snap("1.2.3.4:5", 100, 100, 10, 0); + let info = derive(Some(&prev), &cur, Duration::from_secs(1)); + assert_eq!(info.up_kbps, None); + assert_eq!(info.loss_pct, None); + } + + #[test] + fn zero_elapsed_yields_no_rates() { + let prev = snap("1.2.3.4:5", 0, 0, 0, 0); + let cur = snap("1.2.3.4:5", 4000, 2000, 100, 2); + let info = derive(Some(&prev), &cur, Duration::ZERO); + assert_eq!(info.up_kbps, None); + assert_eq!(info.loss_pct, None); + } + + #[test] + fn oversized_rtt_saturates_instead_of_wrapping() { + let mut cur = snap("1.2.3.4:5", 0, 0, 0, 0); + cur.rtt = Duration::from_secs(u64::MAX); + let info = derive(None, &cur, POLL_INTERVAL); + assert_eq!(info.rtt_ms, u32::MAX); + } +} diff --git a/src/core/messages.rs b/src/core/messages.rs index 417a5e4..fd97e62 100644 --- a/src/core/messages.rs +++ b/src/core/messages.rs @@ -401,6 +401,11 @@ pub enum UiEvent { id: EndpointId, }, AudioLevels(Vec<(EndpointId, f32)>), + /// Periodic per-peer connection transparency snapshot (~1/sec): path type + /// (direct/relay), RTT, and window loss/bitrate for every peer with a live + /// audio link. A FULL replacement each time β€” a peer absent from the list + /// has no live link right now, so its badge should disappear. + ConnectionStats(Vec<(EndpointId, crate::core::connstats::PeerConnInfo)>), /// Raw (pre-gate, pre-mute) normalized RMS of the local mic, `0.0..=1.0`, /// for the settings level meter. Throttled to ~10/sec. MicLevel(f32), diff --git a/src/core/mod.rs b/src/core/mod.rs index c63285b..2fde72d 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -1,3 +1,4 @@ +pub mod connstats; pub mod jitter; pub mod messages; mod recovery; @@ -652,6 +653,7 @@ struct ActiveSession { mixer_task: tokio::task::JoinHandle<()>, event_task: tokio::task::JoinHandle<()>, conn_event_task: tokio::task::JoinHandle<()>, + conn_stats_task: tokio::task::JoinHandle<()>, recovery_task: tokio::task::JoinHandle<()>, recovery_terminal_task: tokio::task::JoinHandle<()>, grace_timers: GraceTimers, @@ -685,6 +687,7 @@ impl ActiveSession { self.mixer_task.abort(); self.event_task.abort(); self.conn_event_task.abort(); + self.conn_stats_task.abort(); // Abort any pending reconnect grace timers so they can't fire a stray // eviction (or touch a torn-down transport) after the session is gone. for (_, handle) in self.grace_timers.lock().unwrap().drain() { @@ -2485,6 +2488,40 @@ async fn run_core_loop( } }); + // Connection-transparency poll: ~1/sec, snapshot every live audio + // link's selected path and hand the UI derived badge info (path + // type, RTT, window loss/bitrate). Read-only against the + // transport; owns the previous-snapshot map the derivation diffs + // against. + let transport_stats = transport.clone(); + let ui_tx_stats = ui_tx.clone(); + let conn_stats_task = tokio::spawn(async move { + let mut prev: HashMap = + HashMap::new(); + let mut last = tokio::time::Instant::now(); + let mut ticker = tokio::time::interval(connstats::POLL_INTERVAL); + ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + loop { + ticker.tick().await; + let now = tokio::time::Instant::now(); + let elapsed = now - last; + last = now; + let snaps = transport_stats.connection_stats(); + let infos = snaps + .iter() + .map(|(id, cur)| (*id, connstats::derive(prev.get(id), cur, elapsed))) + .collect(); + prev = snaps.into_iter().collect(); + if ui_tx_stats + .send(UiEvent::ConnectionStats(infos)) + .await + .is_err() + { + break; + } + } + }); + let session = ActiveSession { room_state: room_state.clone(), capture_thread, @@ -2492,6 +2529,7 @@ async fn run_core_loop( mixer_task, event_task, conn_event_task, + conn_stats_task, recovery_task, recovery_terminal_task, grace_timers, diff --git a/src/network/iroh_impl.rs b/src/network/iroh_impl.rs index af9903c..a98e42e 100644 --- a/src/network/iroh_impl.rs +++ b/src/network/iroh_impl.rs @@ -689,6 +689,57 @@ impl IrohTransport { Ok(bytes) } + /// Snapshot the selected QUIC path of every live audio connection, for the + /// UI's per-peer connection badge (direct/relay, RTT, loss, bitrate). + /// Cheap and lock-light: the `live_conns` guard is released before touching + /// any connection, and `Connection::paths()` reads shared state without I/O. + pub fn connection_stats(&self) -> Vec<(EndpointId, crate::network::PathSnapshot)> { + // Clone the connections out so the map lock isn't held while we inspect + // paths (a supervisor inserts/removes entries as links come and go). + let conns: Vec<(EndpointId, Connection)> = self + .shared + .live_conns + .lock() + .unwrap() + .iter() + .map(|(id, conn)| (*id, conn.clone())) + .collect(); + conns + .into_iter() + .filter_map(|(id, conn)| { + let paths = conn.paths(); + // The selected path is the one carrying application data. In the + // brief window where none is flagged (e.g. mid-migration), fall + // back to the first open path rather than dropping the badge. + let path = paths + .iter() + .find(|p| p.is_selected()) + .or_else(|| paths.iter().next())?; + let stats = path.stats(); + // Per-variant display: `TransportAddr`'s own `Display` prefixes + // a scheme ("ip:1.2.3.4:5") that's noise next to the badge's + // Direct/Relay label. + let remote_addr = match path.remote_addr() { + iroh::TransportAddr::Ip(sock) => sock.to_string(), + iroh::TransportAddr::Relay(url) => url.to_string(), + other => other.to_string(), + }; + Some(( + id, + crate::network::PathSnapshot { + is_relay: path.remote_addr().is_relay(), + remote_addr, + rtt: stats.rtt, + tx_bytes: stats.udp_tx.bytes, + rx_bytes: stats.udp_rx.bytes, + tx_datagrams: stats.udp_tx.datagrams, + lost_packets: stats.lost_packets, + }, + )) + }) + .collect() + } + /// Fetch a chat attachment's bytes from its sender over the file plane. pub async fn fetch_attachment( &self, diff --git a/src/network/mod.rs b/src/network/mod.rs index 0fba1f2..aab3f78 100644 --- a/src/network/mod.rs +++ b/src/network/mod.rs @@ -184,6 +184,32 @@ pub enum ConnEvent { Left(EndpointId), } +/// Owned snapshot of a peer's *selected* QUIC path (the one currently carrying +/// application data), taken from the live audio connection for the UI's +/// connection-transparency badge. Counters are cumulative for the path's +/// lifetime; rate/loss derivation over a poll window happens in +/// `core::connstats` (which also detects path switches via `remote_addr`). +#[derive(Debug, Clone, PartialEq)] +pub struct PathSnapshot { + /// True when the path runs through a relay server, false for a direct + /// (holepunched or local) IP path. + pub is_relay: bool, + /// The path's remote transport address: `ip:port` for a direct path, the + /// relay URL for a relayed one. + pub remote_addr: String, + /// Current QUIC round-trip-time estimate for the path. + pub rtt: std::time::Duration, + /// Cumulative bytes sent in UDP datagrams on the path. + pub tx_bytes: u64, + /// Cumulative bytes received in UDP datagrams on the path. + pub rx_bytes: u64, + /// Cumulative UDP datagrams sent on the path (the loss denominator: for our + /// small voice frames these map ~1:1 to QUIC packets). + pub tx_datagrams: u64, + /// Cumulative packets detected lost on the path. + pub lost_packets: u64, +} + #[derive(Serialize, Deserialize, Clone, Debug)] pub struct PeerSpeakTicket { pub host_addr: iroh::EndpointAddr, diff --git a/tests/transport_loopback.rs b/tests/transport_loopback.rs index d2cea47..731c912 100644 --- a/tests/transport_loopback.rs +++ b/tests/transport_loopback.rs @@ -244,6 +244,71 @@ async fn loopback_sequenced_audio_reaches_peer_and_decodes() { ); } +/// Connection transparency: over a real loopback link, `connection_stats()` +/// must report the peer's selected path as direct (relay disabled here), with +/// an IP remote address and counters that advance while audio flows β€” and the +/// `connstats::derive` seam must turn two such snapshots into badge info with +/// live rates. +#[tokio::test] +async fn connection_stats_report_a_direct_path_with_live_counters() { + let a = spawn_node().await; + let b = spawn_node().await; + + a.lookup.add_endpoint_info(b.endpoint.addr()); + b.lookup.add_endpoint_info(a.endpoint.addr()); + + let a_id = a.endpoint.id(); + let b_id = b.endpoint.id(); + a.transport.admit_audio_sender(b_id); + b.transport.admit_audio_sender(a_id); + + // Keep B's receive path subscribed like production (drained implicitly). + let _b_rx = b.transport.receive_datagrams().await.expect("subscribe B"); + + a.transport.connect_peer(b.endpoint.addr()).await; + b.transport.connect_peer(a.endpoint.addr()).await; + tokio::time::sleep(Duration::from_millis(500)).await; + + let snap = |stats: Vec<(iroh::EndpointId, peerspeak::network::PathSnapshot)>| { + stats + .into_iter() + .find(|(id, _)| *id == b_id) + .map(|(_, s)| s) + .expect("peer B should appear in A's connection stats") + }; + let s1 = snap(a.transport.connection_stats()); + assert!(!s1.is_relay, "loopback with relay disabled must be direct"); + assert!( + s1.remote_addr.parse::().is_ok(), + "direct path address should be ip:port, got {}", + s1.remote_addr + ); + + // Stream real audio so the path counters move. + let mut enc = OpusEncoder::new(48000, Channels::Mono, Application::Voip).unwrap(); + for seq in 0..25u32 { + a.transport.broadcast(packet(&mut enc, seq)); + tokio::time::sleep(Duration::from_millis(5)).await; + } + + let s2 = snap(a.transport.connection_stats()); + assert!(s2.tx_bytes > s1.tx_bytes, "sent bytes should advance"); + assert!( + s2.tx_datagrams > s1.tx_datagrams, + "sent datagrams should advance" + ); + + // The derivation seam turns the two snapshots into live badge info. + let info = peerspeak::core::connstats::derive(Some(&s1), &s2, Duration::from_millis(200)); + assert!(!info.relay); + assert_eq!(info.remote_addr, s2.remote_addr); + assert!(info.rtt_ms < 1000, "localhost RTT should be sane"); + assert!( + info.up_kbps.expect("same path + positive window has a rate") > 0.0, + "audio was flowing, so the upstream rate must be non-zero" + ); +} + /// Read datagrams off a raw connection until `target` arrive or the deadline /// passes, asserting each carries the 4-byte sequence header. async fn count_audio(conn: &Connection, target: u32, deadline: tokio::time::Instant) -> u32 {