use crate::network::{ConnEvent, NetError, NetworkTransport}; use async_trait::async_trait; use bytes::Bytes; use iroh::endpoint::{Connection, ConnectionError, VarInt}; use iroh::{Endpoint, EndpointId}; use std::collections::{HashMap, HashSet}; use std::sync::{Arc, Mutex as StdMutex}; use std::time::Duration; use tokio::sync::mpsc; use tokio::sync::mpsc::Receiver; use crate::files::{AttachmentId, ChatAttachment}; use crate::protocol::{AUDIO_ALPN, FILES_ALPN}; /// Per-peer datagram send queue depth. Audio is real-time, so a backlog is /// useless latency — keep it shallow and drop the oldest frame when full. const SEND_QUEUE_DEPTH: usize = 8; /// How many inbound connections we'll buffer for a single peer's supervisor. /// During a healthy link the supervisor drains this promptly; a small buffer /// just covers the brief window between detecting a drop and re-awaiting. const INBOUND_QUEUE_DEPTH: usize = 4; /// Reconnect backoff bounds for the dialing side. Starts gentle so a momentary /// blip recovers almost instantly, and caps so a long outage doesn't hammer. const INITIAL_BACKOFF: Duration = Duration::from_millis(250); const MAX_BACKOFF: Duration = Duration::from_secs(5); /// QUIC application error code sent when we leave/quit intentionally, so the /// remote supervisor can tell a graceful departure from a transient drop. The /// specific value is informational; any *application* close (vs a timeout/reset) /// already means "the peer closed this on purpose." const GOODBYE_CODE: u32 = 1; /// Bound on each phase (connect, read) of a chat-attachment fetch, so a slow or /// stalled sender can't hang the fetch indefinitely. const FILE_FETCH_TIMEOUT: Duration = Duration::from_secs(30); /// State shared between the transport, its protocol handler, and every per-peer /// supervisor task. One supervisor owns a peer's whole connection lifecycle. struct Shared { /// Used by dialer-side supervisors to (re)establish connections. endpoint: Endpoint, self_id: EndpointId, /// Sync-lockable send handles, so `broadcast` can fan out from the (non-async) /// capture/encode thread without touching the Tokio runtime. A supervisor /// inserts its handle when a link comes up and removes it when the link dies, /// so `broadcast` only ever targets peers that are currently reachable. senders: StdMutex>>, /// Last-known full address (relay + direct addrs) per peer, retained for the /// supervisor's lifetime and refreshed whenever the peer (re)announces. The /// dialer dials this directly so a reconnect never depends on the address /// lookup, which a transient gossip `Leave`/`NeighborDown` can purge — the /// bug that left the dialer redial-looping with "no address" forever. addrs: StdMutex>, /// One supervisor task per known peer. peers: tokio::sync::Mutex>, /// The currently-live connection per peer, tracked so an intentional leave /// can explicitly close each one with [`GOODBYE_CODE`] (a prompt, reliable /// "I'm leaving" signal the remote can distinguish from a network drop). A /// supervisor inserts its connection when the link comes up and removes it /// when the link dies. live_conns: StdMutex>, /// Core-owned audio admission snapshot for this room session. It mirrors the /// verified gossip roster plus peers still inside reconnect grace; transport /// connections alone never mutate this set. admitted_audio: StdMutex>, /// Chat file attachments we're serving to room members this session, keyed by /// the random attachment id. Populated when we send a chat file; read by the /// file protocol handler to answer a member's fetch. Cleared on leave. Each /// blob is already byte-capped at send time. /// Blobs we serve to room members, bounded by count and byte budgets /// (Phase 3C) — an evicted id reads as "sender no longer has the file". served_files: StdMutex, incoming_tx: mpsc::Sender<(EndpointId, Bytes)>, /// Best-effort link-state notifications for the UI (connecting / connected). conn_events_tx: mpsc::Sender, } struct PeerHandle { /// The supervisor owning this peer's connect → run → reconnect loop. supervisor: tokio::task::JoinHandle<()>, /// Inbound connections accepted for this peer are routed here so the /// acceptor-side supervisor can pick them up (initial and reconnect). inbound_tx: mpsc::Sender, } /// Aborts a spawned task when dropped. The supervisor holds its send/read loops /// in these so that cancelling the supervisor (via `remove`'s `abort`) tears the /// loops down too — otherwise they'd keep their `Connection` clones alive and /// the link would never actually close. struct AbortOnDrop(tokio::task::JoinHandle<()>); impl Drop for AbortOnDrop { fn drop(&mut self) { self.0.abort(); } } impl Shared { /// Ensure a supervisor task exists for `peer_id`, returning its inbound /// connection sender. Idempotent: a second call for a known peer just hands /// back the existing sender. Both the reactive `connect_peer` path and the /// inbound `accept()` path funnel through here, so whichever fires first /// spins the peer up and the roles stay deterministic by id. async fn ensure_supervisor(self: &Arc, peer_id: EndpointId) -> mpsc::Sender { let mut peers = self.peers.lock().await; if let Some(handle) = peers.get(&peer_id) { return handle.inbound_tx.clone(); } let (inbound_tx, inbound_rx) = mpsc::channel::(INBOUND_QUEUE_DEPTH); let shared = self.clone(); let supervisor = tokio::spawn(supervise(shared, peer_id, inbound_rx)); let inbound_tx_ret = inbound_tx.clone(); peers.insert( peer_id, PeerHandle { supervisor, inbound_tx, }, ); crate::log_msg(&format!("Transport: supervising peer {:?}", peer_id)); inbound_tx_ret } /// Permanently retire a peer (it left the room). Stops its supervisor — which /// cancels any in-flight dial or live link — and drops its send handle. async fn remove(&self, peer_id: EndpointId) { self.senders.lock().unwrap().remove(&peer_id); self.addrs.lock().unwrap().remove(&peer_id); if let Some(handle) = self.peers.lock().await.remove(&peer_id) { handle.supervisor.abort(); crate::log_msg(&format!( "Transport: stopped supervising peer {:?}", peer_id )); } } fn audio_sender_admitted(&self, peer_id: EndpointId) -> bool { let roster = self.admitted_audio.lock().unwrap(); audio_sender_admitted(peer_id, &roster) } fn apply_audio_admission(&self, peer_id: EndpointId, event: AudioAdmissionEvent) { let mut roster = self.admitted_audio.lock().unwrap(); apply_audio_admission_event(&mut roster, peer_id, event); } } /// Why a peer's live-link wait woke up. enum Wake { /// The QUIC connection closed; the error says how. An *application* close /// means the peer left on purpose (reconnect would be wrong); a timeout/reset /// is a transient drop we should reconnect through. Closed(ConnectionError), /// A fresh inbound connection arrived (acceptor side) — switch to it. Replacement(Connection), /// The peer was retired; tear down. Shutdown, } /// True when the peer closed the connection with our [`GOODBYE_CODE`] — i.e. an /// intentional leave/quit via [`IrohTransport::leave`]. Any other close (idle /// timeout, reset, or an application close with a different code) is treated as a /// transient drop to reconnect through, so this never mistakes a blip for a /// departure. fn is_graceful_leave(err: &ConnectionError) -> bool { matches!(err, ConnectionError::ApplicationClosed(frame) if frame.error_code == VarInt::from_u32(GOODBYE_CODE)) } /// Pure S8 membership decision: iroh already authenticated `remote` as the /// connection's endpoint id, so audio admission is exactly live roster membership. pub(crate) fn audio_sender_admitted(remote: EndpointId, roster: &HashSet) -> bool { roster.contains(&remote) } #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub(crate) enum AudioAdmissionEvent { /// A signed gossip Announce/Update says the peer is in the live room roster. RosterPresent, /// Gossip reported a transient drop; keep admission during reconnect grace. TransientDropGrace, /// Graceful leave, transport Left eviction, or reconnect-grace expiry. Remove, } pub(crate) fn apply_audio_admission_event( roster: &mut HashSet, peer_id: EndpointId, event: AudioAdmissionEvent, ) { match event { AudioAdmissionEvent::RosterPresent => { roster.insert(peer_id); } AudioAdmissionEvent::TransientDropGrace => { // Grace is not an authority to add membership; it only preserves an // already-admitted peer until either rejoin or grace expiry. } AudioAdmissionEvent::Remove => { roster.remove(&peer_id); } } } /// Owns a single peer's connection lifecycle for as long as the peer is in the /// room: obtain a link, run the send/read loops, and on loss obtain a new one — /// with capped backoff on the dialing side. The deterministic-initiator rule /// (lower id dials, higher id accepts) is the same on first connect and every /// reconnect, so exactly one shared connection forms each time. async fn supervise( shared: Arc, peer_id: EndpointId, mut inbound_rx: mpsc::Receiver, ) { let is_dialer = shared.self_id.to_string() < peer_id.to_string(); let mut backoff = INITIAL_BACKOFF; // Show "connecting" until the first link is actually up. let _ = shared .conn_events_tx .try_send(ConnEvent::Connecting(peer_id)); let mut conn = match obtain_conn(&shared, peer_id, is_dialer, &mut inbound_rx, &mut backoff).await { Some(conn) => conn, None => return, // retired before we ever connected }; loop { // A healthy link resets the dialer's backoff for the next outage. backoff = INITIAL_BACKOFF; let (send_tx, mut send_rx) = mpsc::channel::(SEND_QUEUE_DEPTH); shared.senders.lock().unwrap().insert(peer_id, send_tx); // Publish the live connection so an intentional leave can close it with // the goodbye code. shared .live_conns .lock() .unwrap() .insert(peer_id, conn.clone()); let _ = shared .conn_events_tx .try_send(ConnEvent::Connected(peer_id)); crate::log_msg(&format!("Transport: peer {:?} link up", peer_id)); // Run until the link dies, a replacement arrives, or we're retired. The // send/read loops are held in abort-on-drop guards scoped to this block, // so however we leave it — normal exit or the supervisor being aborted // mid-`select` — both loops stop and release their `Connection` clones. let wake = { let conn_send = conn.clone(); let _send_guard = AbortOnDrop(tokio::spawn(async move { while let Some(data) = send_rx.recv().await { if conn_send.send_datagram(data).is_err() { break; } } })); let conn_read = conn.clone(); let incoming_tx = shared.incoming_tx.clone(); let _read_guard = AbortOnDrop(tokio::spawn(async move { while let Ok(bytes) = conn_read.read_datagram().await { if incoming_tx.send((peer_id, bytes)).await.is_err() { break; } } })); // Only the acceptor watches for replacements: if the dialer notices // the drop and re-dials before our own `closed()` fires, we switch to // the new link immediately instead of waiting out the idle timeout. tokio::select! { reason = conn.closed() => Wake::Closed(reason), inbound = inbound_rx.recv(), if !is_dialer => match inbound { Some(new_conn) => Wake::Replacement(new_conn), None => Wake::Shutdown, }, } }; // Dropping the send handle makes `broadcast` skip this peer until the // link is back, instead of piling frames into a dead connection. The // live connection is gone too until a new one is established. shared.senders.lock().unwrap().remove(&peer_id); shared.live_conns.lock().unwrap().remove(&peer_id); match wake { Wake::Shutdown => return, Wake::Replacement(new_conn) => { crate::log_msg(&format!( "Transport: peer {:?} replaced with new inbound link", peer_id )); let _ = shared .conn_events_tx .try_send(ConnEvent::Connecting(peer_id)); conn = new_conn; } Wake::Closed(reason) => { // A graceful application close means the peer left on purpose — // don't reconnect; tell the core to evict it now. if is_graceful_leave(&reason) { crate::log_msg(&format!( "Transport: peer {:?} left gracefully ({:?})", peer_id, reason )); let _ = shared.conn_events_tx.try_send(ConnEvent::Left(peer_id)); return; } crate::log_msg(&format!( "Transport: peer {:?} link dropped; reconnecting", peer_id )); let _ = shared .conn_events_tx .try_send(ConnEvent::Connecting(peer_id)); conn = match obtain_conn(&shared, peer_id, is_dialer, &mut inbound_rx, &mut backoff) .await { Some(conn) => conn, None => return, // retired while reconnecting }; } } } } /// Block until we hold a live connection to the peer, or `None` if the peer was /// retired first. The dialer redials with capped exponential backoff; the /// acceptor waits for the dialer to come to it (a closed `inbound_rx` means the /// supervisor was retired). async fn obtain_conn( shared: &Arc, peer_id: EndpointId, is_dialer: bool, inbound_rx: &mut mpsc::Receiver, backoff: &mut Duration, ) -> Option { if is_dialer { loop { // Dial the retained full address (relay + direct addrs) so reconnects // don't depend on the address lookup; it's re-read each attempt so a // refreshed address (peer back on a new one) is picked up next retry. // Fall back to the bare id (lookup/relay resolution) only if we somehow // never recorded an address — the dialer path always carries one. let target = shared.addrs.lock().unwrap().get(&peer_id).cloned(); let result = match target { Some(addr) => shared.endpoint.connect(addr, AUDIO_ALPN).await, None => shared.endpoint.connect(peer_id, AUDIO_ALPN).await, }; match result { Ok(conn) => return Some(conn), Err(e) => { crate::log_msg(&format!( "Transport: dial to {:?} failed: {:?}; retrying in {:?}", peer_id, e, backoff )); tokio::time::sleep(*backoff).await; *backoff = (*backoff * 2).min(MAX_BACKOFF); } } } } else { inbound_rx.recv().await } } /// The persistent audio protocol handler. Registered ONCE on the app's single /// long-lived [`Router`](iroh::protocol::Router), it delegates each inbound audio /// connection to whatever session's [`Shared`] is currently bound, or drops it /// when idle (no active call). This is what lets the endpoint + router outlive any /// individual room session — the W7 friends listener needs the endpoint reachable /// *between* calls — while the per-peer connection state stays per-session. #[derive(Clone, Default)] pub struct AudioRouter { /// The active session's transport state, swapped in on join and cleared on /// leave. `None` means "no call" → inbound audio links are dropped. current: Arc>>>, } impl std::fmt::Debug for AudioRouter { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("AudioRouter").finish_non_exhaustive() } } impl AudioRouter { pub fn new() -> Self { Self::default() } /// Route inbound audio connections to `transport`'s session (called on join). pub fn bind(&self, transport: &IrohTransport) { *self.current.lock().unwrap() = Some(transport.shared.clone()); } /// Stop routing — drop inbound audio links until the next [`bind`](Self::bind) /// (called on leave). The endpoint/router stay alive; we just have nowhere to /// hand a fresh audio connection. pub fn clear(&self) { *self.current.lock().unwrap() = None; } } impl iroh::protocol::ProtocolHandler for AudioRouter { fn accept( &self, connection: Connection, ) -> impl std::future::Future> + Send { let peer_id = connection.remote_id(); // Snapshot the current session (clone the Arc) and release the lock before // any await — the guard isn't held across suspension points. let shared = self.current.lock().unwrap().clone(); async move { // No active call → nothing to route this inbound link to; drop it. let Some(shared) = shared else { return Ok(()); }; // Deterministic roles: the lower id dials, the higher id accepts. // If we're the dialer for this peer we never consume inbound links // (we make our own), so just let this one drop. if shared.self_id.to_string() < peer_id.to_string() { return Ok(()); } if !shared.audio_sender_admitted(peer_id) { crate::log_msg(&format!( "Transport: rejected inbound audio from non-member {}", crate::short_id(&peer_id.to_string()) )); return Ok(()); } // Route the connection to this peer's supervisor (creating it if the // inbound link arrives after the signed gossip Announce admitted it). // try_send keeps the protocol handler from ever blocking; a full queue // only happens if links are churning, and the supervisor gets the next one. let inbound_tx = shared.ensure_supervisor(peer_id).await; if inbound_tx.try_send(connection).is_err() { crate::log_msg(&format!( "Transport: dropped inbound link from {:?} (queue full)", peer_id )); } Ok(()) } } } /// Protocol handler for the file-transfer plane (`FILES_ALPN`). Mirrors /// [`AudioRouter`]: it's persistent on the router and bound to the active /// session's [`Shared`] on join. On an inbound stream it authenticates the peer /// (iroh ALPN handshake gives us `remote_id`), gates on **live room membership** /// (same invariant as audio admission, so a former member can't pull files), /// reads a single 32-byte attachment id, and streams back the matching blob from /// the session serve store — or nothing if the id is unknown. #[derive(Clone, Default)] pub struct FileRouter { current: Arc>>>, } impl std::fmt::Debug for FileRouter { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("FileRouter").finish_non_exhaustive() } } impl FileRouter { pub fn new() -> Self { Self::default() } /// Route inbound file fetches to `transport`'s session (called on join). pub fn bind(&self, transport: &IrohTransport) { *self.current.lock().unwrap() = Some(transport.shared.clone()); } /// Stop serving files until the next [`bind`](Self::bind) (called on leave). pub fn clear(&self) { *self.current.lock().unwrap() = None; } } /// Max bytes we'll read for a fetch *request* frame. A request is exactly one /// 32-byte id; this small ceiling rejects a peer trying to stream us a huge /// "request" as a cheap DoS. const FILE_REQUEST_MAX: usize = 64; impl iroh::protocol::ProtocolHandler for FileRouter { fn accept( &self, connection: Connection, ) -> impl std::future::Future> + Send { let peer_id = connection.remote_id(); let shared = self.current.lock().unwrap().clone(); async move { // No active call → nothing to serve. let Some(shared) = shared else { return Ok(()); }; // Member gating: only current room members may fetch our files. Reuses // the audio admission roster (the authoritative room membership set). if !shared.audio_sender_admitted(peer_id) { crate::log_msg(&format!( "Transport: rejected file fetch from non-member {}", crate::short_id(&peer_id.to_string()) )); return Ok(()); } // Accept one bidirectional stream: read the id, write the bytes. let Ok((mut send, mut recv)) = connection.accept_bi().await else { return Ok(()); }; let Ok(req) = recv.read_to_end(FILE_REQUEST_MAX).await else { return Ok(()); }; let Some(id) = crate::files::parse_request(&req) else { return Ok(()); }; let blob = shared.served_files.lock().unwrap().get(&id); if let Some(blob) = blob { let _ = send.write_all(&blob).await; } // Finish either way: an unknown id closes with an empty body, which // the fetcher reads as a zero-length result and treats as "gone". let _ = send.finish(); // CRITICAL: `finish()` only marks the stream's EOF — it does NOT wait // for the written bytes to be delivered and acknowledged. If we return // here the `connection` drops, and its CONNECTION_CLOSE can race ahead // of the still-in-flight stream data, so the fetcher's read aborts with // "connection lost". Wait for the fetcher to receive everything and // close the connection itself (it drops `conn` right after read_to_end); // that close is our signal the transfer landed. Bounded so a fetcher // that vanishes can't pin this task forever. let _ = tokio::time::timeout(FILE_FETCH_TIMEOUT, connection.closed()).await; Ok(()) } } } pub struct IrohTransport { shared: Arc, incoming_rx: tokio::sync::Mutex>>, conn_events_rx: tokio::sync::Mutex>>, } impl IrohTransport { /// Build a per-session transport over the app's persistent `endpoint`. The /// inbound audio handler is the persistent [`AudioRouter`] (registered once on /// the router) — call [`AudioRouter::bind`] with this transport to route this /// session's inbound links to it. pub fn new(endpoint: Endpoint) -> Self { let (incoming_tx, incoming_rx) = mpsc::channel(1000); let (conn_events_tx, conn_events_rx) = mpsc::channel(100); let self_id = endpoint.id(); let shared = Arc::new(Shared { endpoint, self_id, senders: StdMutex::new(HashMap::new()), addrs: StdMutex::new(HashMap::new()), peers: tokio::sync::Mutex::new(HashMap::new()), live_conns: StdMutex::new(HashMap::new()), admitted_audio: StdMutex::new(HashSet::new()), served_files: StdMutex::new(crate::files::ServeStore::default()), incoming_tx, conn_events_tx, }); Self { shared, incoming_rx: tokio::sync::Mutex::new(Some(incoming_rx)), conn_events_rx: tokio::sync::Mutex::new(Some(conn_events_rx)), } } /// Tear down for an intentional leave/quit: close every live connection with /// the graceful goodbye code so each remote peer evicts us promptly (its /// supervisor sees an application close, not a reconnectable drop), then stop /// all supervisors so none linger redialing the about-to-close endpoint. /// Call this before shutting the router down. pub async fn leave(&self) { let conns: Vec = self .shared .live_conns .lock() .unwrap() .drain() .map(|(_, c)| c) .collect(); for conn in &conns { conn.close(VarInt::from_u32(GOODBYE_CODE), b"leave"); } // Stop every supervisor so it can't re-dial during/after teardown. let mut peers = self.shared.peers.lock().await; for (_, handle) in peers.drain() { handle.supervisor.abort(); } self.shared.senders.lock().unwrap().clear(); self.shared.addrs.lock().unwrap().clear(); self.shared.admitted_audio.lock().unwrap().clear(); self.shared.served_files.lock().unwrap().clear(); // Give the CONNECTION_CLOSE frames a moment to flush before the caller // shuts the endpoint/router down (the `conns` clones are still alive // here, so the endpoint can still transmit them). tokio::time::sleep(Duration::from_millis(150)).await; } /// Admit a peer to this session's audio plane. Core calls this from verified /// gossip roster events; the transport never derives membership on its own. pub fn admit_audio_sender(&self, peer_id: EndpointId) { self.shared .apply_audio_admission(peer_id, AudioAdmissionEvent::RosterPresent); } /// Preserve an already-admitted peer through the reconnect grace window. pub fn keep_audio_sender_for_reconnect_grace(&self, peer_id: EndpointId) { self.shared .apply_audio_admission(peer_id, AudioAdmissionEvent::TransientDropGrace); } /// Remove a peer from audio admission before tearing down transport/jitter state. pub fn remove_audio_sender(&self, peer_id: EndpointId) { self.shared .apply_audio_admission(peer_id, AudioAdmissionEvent::Remove); } pub fn audio_sender_admitted(&self, peer_id: EndpointId) -> bool { self.shared.audio_sender_admitted(peer_id) } /// Make `bytes` available to room members under `id` for the rest of this /// session (served by the [`FileRouter`] handler). Called by core when we /// send a chat file. The blob is cleared on leave. pub fn serve_attachment(&self, id: AttachmentId, bytes: Arc>) { if !self.shared.served_files.lock().unwrap().insert(id, bytes) { crate::log_msg("Transport: refused to serve an over-budget blob"); } } /// Drop a previously-served blob (e.g. a music track no longer current-or-next). pub fn forget_attachment(&self, id: AttachmentId) { self.shared.served_files.lock().unwrap().remove(&id); } /// Fetch `size` bytes stored under `id` from peer `from` over the files plane. /// Shared core of `fetch_attachment` and music-track fetching: dials /// `FILES_ALPN`, writes the 32-byte id, and reads bounded by `size`. pub async fn fetch_blob( &self, from: EndpointId, id: AttachmentId, size: u64, ) -> Result, NetError> { if !crate::files::size_within_cap(size) { return Err(NetError::Other("attachment size out of range".to_string())); } let addr = self.shared.addrs.lock().unwrap().get(&from).cloned(); let connect = async { match addr { Some(addr) => self.shared.endpoint.connect(addr, FILES_ALPN).await, None => self.shared.endpoint.connect(from, FILES_ALPN).await, } }; let conn = tokio::time::timeout(FILE_FETCH_TIMEOUT, connect) .await .map_err(|_| NetError::Other("file fetch: connect timed out".to_string()))? .map_err(|e| NetError::Other(format!("file fetch: connect failed: {e}")))?; let (mut send, mut recv) = conn .open_bi() .await .map_err(|e| NetError::Other(format!("file fetch: open stream failed: {e}")))?; send.write_all(&id) .await .map_err(|e| NetError::Other(format!("file fetch: request write failed: {e}")))?; send.finish() .map_err(|e| NetError::Other(format!("file fetch: request finish failed: {e}")))?; // `read_to_end(size)` errors if the stream exceeds `size`, rejecting an // overlong transfer; the exact-length check below rejects a short one. let read = recv.read_to_end(size as usize); let bytes = tokio::time::timeout(FILE_FETCH_TIMEOUT, read) .await .map_err(|_| NetError::Other("file fetch: read timed out".to_string()))? .map_err(|e| NetError::Other(format!("file fetch: read failed: {e}")))?; if bytes.is_empty() { return Err(NetError::Other( "file fetch: sender no longer has the file".to_string(), )); } // Exact transfer required (Phase 3C): a truncated body must not be // cached/saved/decoded as if it were the declared attachment. if bytes.len() as u64 != size { return Err(NetError::Other(format!( "file fetch: incomplete transfer ({} of {size} bytes)", bytes.len() ))); } 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, from: EndpointId, att: &ChatAttachment, ) -> Result, NetError> { self.fetch_blob(from, att.id, att.size).await } } #[async_trait] impl NetworkTransport for IrohTransport { async fn connect_peer(&self, addr: iroh::EndpointAddr) { let peer_id = addr.id; // Retain (or refresh) the peer's full address so the dialer can re-reach // it across reconnects without leaning on the address lookup. self.shared.addrs.lock().unwrap().insert(peer_id, addr); // Start (or confirm) the peer's supervisor. It owns dialing/accepting and // reconnection from here on; this call is just the reactive trigger. self.shared.ensure_supervisor(peer_id).await; } async fn disconnect_peer(&self, peer_id: EndpointId) { self.shared.remove(peer_id).await; } fn broadcast(&self, data: Bytes) { let senders = self.shared.senders.lock().unwrap(); for tx in senders.values() { // Drop on a full queue: stale audio is worthless, and we must never // block the encode thread on a slow peer. let _ = tx.try_send(data.clone()); } } async fn receive_datagrams(&self) -> Result, NetError> { let mut rx_guard = self.incoming_rx.lock().await; rx_guard .take() .ok_or_else(|| NetError::Other("Datagram receiver already subscribed".to_string())) } async fn subscribe_conn_events(&self) -> Result, NetError> { let mut rx_guard = self.conn_events_rx.lock().await; rx_guard .take() .ok_or_else(|| NetError::Other("Connection events already subscribed".to_string())) } } #[cfg(test)] mod tests { use super::*; use iroh::SecretKey; fn endpoint_id() -> EndpointId { SecretKey::generate().public() } #[test] fn audio_sender_admission_accepts_roster_member() { let member = endpoint_id(); let roster = HashSet::from([member]); assert!(audio_sender_admitted(member, &roster)); } #[test] fn audio_sender_admission_rejects_unknown_sender() { let member = endpoint_id(); let stranger = endpoint_id(); let roster = HashSet::from([member]); assert!(!audio_sender_admitted(stranger, &roster)); } #[test] fn audio_sender_admission_rejects_former_member_after_roster_removal() { let former = endpoint_id(); let mut roster = HashSet::from([former]); assert!(audio_sender_admitted(former, &roster)); roster.remove(&former); assert!(!audio_sender_admitted(former, &roster)); } #[test] fn audio_sender_admission_waits_for_mid_join_announce() { let joining_peer = endpoint_id(); let mut roster = HashSet::new(); assert!(!audio_sender_admitted(joining_peer, &roster)); roster.insert(joining_peer); assert!(audio_sender_admitted(joining_peer, &roster)); } #[test] fn audio_admission_lifecycle_keeps_peer_through_transient_grace() { let peer = endpoint_id(); let mut roster = HashSet::new(); apply_audio_admission_event(&mut roster, peer, AudioAdmissionEvent::RosterPresent); assert!(audio_sender_admitted(peer, &roster)); apply_audio_admission_event(&mut roster, peer, AudioAdmissionEvent::TransientDropGrace); assert!(audio_sender_admitted(peer, &roster)); } #[test] fn audio_admission_lifecycle_does_not_add_unknown_peer_on_grace_event() { let peer = endpoint_id(); let mut roster = HashSet::new(); apply_audio_admission_event(&mut roster, peer, AudioAdmissionEvent::TransientDropGrace); assert!(!audio_sender_admitted(peer, &roster)); } #[test] fn audio_admission_lifecycle_removes_peer_on_leave_or_grace_expiry() { let peer = endpoint_id(); let mut roster = HashSet::from([peer]); apply_audio_admission_event(&mut roster, peer, AudioAdmissionEvent::Remove); assert!(!audio_sender_admitted(peer, &roster)); } }