From 10707152a3b5017147de7bb304c27bb748998e28 Mon Sep 17 00:00:00 2001 From: Mollusk Date: Thu, 18 Jun 2026 04:24:01 -0400 Subject: [PATCH 1/2] S8 (design-first): pure audio_sender_admitted membership seam (unwired) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Design-first checkpoint for S8 (authorize inbound audio against live room membership). Codex's design note (in the handoff task-report.md) establishes the authoritative roster = gossip IrohGossipState.peers, NOT the audio transport connection list, and recommends mirroring it into an audio-admission snapshot consulted at AudioRouter::accept + datagram ingest. This commit lands ONLY the pure decision seam + tests; wiring is deliberately paused for a senior decision on the reconnect-grace policy (gossip drops a peer from the roster on transient NeighborDown, but core keeps the audio supervisor alive for RECONNECT_GRACE — a strict roster-only gate would cut audio on blips). - audio_sender_admitted(remote, roster) -> bool (pub(crate), #[allow(dead_code)]). - 4 tests: member admitted, stranger rejected, former member rejected after roster removal, mid-join peer rejected until authenticated Announce inserts it. - No behavior change: accept/datagram/mixer paths untouched. S8 remains OPEN. 310 lib tests / clippy --all-targets / release all green (re-run by senior). Co-Authored-By: Claude Opus 4.8 --- src/network/iroh_impl.rs | 59 +++++++++++++++++++++++++++++++++++++++- 1 file changed, 58 insertions(+), 1 deletion(-) diff --git a/src/network/iroh_impl.rs b/src/network/iroh_impl.rs index 20ddf04..6c733cb 100644 --- a/src/network/iroh_impl.rs +++ b/src/network/iroh_impl.rs @@ -5,7 +5,7 @@ use bytes::Bytes; use tokio::sync::mpsc; use tokio::sync::mpsc::Receiver; use std::sync::{Arc, Mutex as StdMutex}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::time::Duration; use async_trait::async_trait; @@ -135,6 +135,13 @@ 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. +#[allow(dead_code)] // Design-first S8 seam; wiring waits for senior review. +pub(crate) fn audio_sender_admitted(remote: EndpointId, roster: &HashSet) -> bool { + roster.contains(&remote) +} + /// 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 @@ -443,3 +450,53 @@ impl NetworkTransport for IrohTransport { .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)); + } +} From 1adf8a97bb52657203b2eddcaa569ececbf09784 Mon Sep 17 00:00:00 2001 From: Mollusk Date: Thu, 18 Jun 2026 04:41:37 -0400 Subject: [PATCH 2/2] S8 (Pass 2): wire grace-aware audio-membership admission MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes S8 (High): inbound audio was authenticated by identity (remote_id) but NOT by room membership, so a former member who knew a current member's endpoint could reconnect on the audio ALPN and inject into / eavesdrop on the mix while invisible in the roster. Now audio is admitted only for live gossip-roster members (the senior+user-resolved GRACE-AWARE policy). Transport (src/network/iroh_impl.rs): - New per-session admitted_audio: HashSet on Shared (internal state, no wire/serialization change). Cleared on disconnect_all. - AudioRouter::accept consults audio_sender_admitted BEFORE ensure_supervisor — a non-member never gets a supervisor, sender handle, datagram reader, or outbound mix. Brief StdMutex check, released before the await (no RT lock). - Pure apply_audio_admission_event(roster, peer, event) with AudioAdmissionEvent {RosterPresent insert, TransientDropGrace no-op, Remove}. Grace deliberately cannot ADD membership — it only preserves an already-admitted peer — so an unknown peer can't sneak in via a grace event. +3 lifecycle tests (on top of Pass-1's 4 predicate tests). - admit/keep_for_reconnect_grace/remove/query methods for core to drive. Core (src/core/mod.rs) — authority is core's VERIFIED gossip-roster events, not transport connect/disconnect: - PeerJoined / PeerUpdated: admit_audio_sender before connect_peer. - PeerConnectionLost: keep_audio_sender_for_reconnect_grace (preserve through the existing RECONNECT_GRACE window — no audio cut on transient blips). - gossip PeerLeft, transport ConnEvent::Left, grace-timer expiry: remove_audio_sender before disconnect_peer + jitter removal (removal-before-teardown bounds the in-flight-datagram race). - datagram receiver: audio_sender_admitted gate before any jitter buffer (defense in depth against a datagram racing a removal). Mixer stays off the hot path. Mid-join: a peer who dials audio before we've verified their signed Announce is dropped (no "pending" admission, which would reintroduce the eavesdrop); their reconnect loop recovers once the Announce admits them. tests/transport_loopback.rs: admit both ends before connecting, mirroring the production room-event order. 313 lib / clippy --all-targets / transport_loopback 4 / reconnect_eviction 6 / release — all re-run green by the senior. Former-member-rejection + mid-join recovery are verifiable only in a 2-machine call (senior's to run). Co-Authored-By: Claude Opus 4.8 --- src/core/mod.rs | 9 +++ src/network/iroh_impl.rs | 114 ++++++++++++++++++++++++++++++++++-- tests/transport_loopback.rs | 3 + 3 files changed, 122 insertions(+), 4 deletions(-) diff --git a/src/core/mod.rs b/src/core/mod.rs index 2ad052d..c2029ba 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -137,6 +137,7 @@ fn arm_grace_timer( let handle = tokio::spawn(async move { tokio::time::sleep(grace).await; crate::log_msg(&format!("Reconnect grace expired; evicting peer {:?}", peer_id)); + transport_evict.remove_audio_sender(peer_id); transport_evict.disconnect_peer(peer_id).await; jitter_evict.lock().await.remove(&peer_id); // Scrub our internal state *before* announcing the eviction, so anything @@ -370,6 +371,7 @@ impl ConnEventHandler { // until the grace timer or the slow gossip Leave. cancel_grace_timer(&self.grace_timers, &id); self.seen_connected.lock().unwrap().remove(&id); + self.transport.remove_audio_sender(id); self.transport.disconnect_peer(id).await; self.jitter.lock().await.remove(&id); let _ = self.ui_tx.send(UiEvent::PeerLeft { id }).await; @@ -1239,6 +1241,9 @@ async fn run_core_loop( }; while let Some((from_peer, bytes)) = datagram_rx.recv().await { + if !transport_recv.audio_sender_admitted(from_peer) { + continue; + } if !audio_datagram_len_ok(bytes.len()) { // Malformed (< sequence header) or oversized Opus payload. continue; @@ -1473,6 +1478,7 @@ async fn run_core_loop( // A (re)join means the peer is back — cancel any // pending reconnect grace timer before re-adding it. cancel_grace_timer(&grace_timers_events, &peer_id); + transport_events.admit_audio_sender(peer_id); // Establish the audio connection as soon as the peer // is known (the transport dedupes the full-mesh race). // Hand over the full address so reconnects can dial @@ -1524,6 +1530,7 @@ async fn run_core_loop( { peers.remove(&peer_id); } + transport_events.remove_audio_sender(peer_id); transport_events.disconnect_peer(peer_id).await; jitter_events.lock().await.remove(&peer_id); let _ = ui_tx_events.send(UiEvent::PeerLeft { id: peer_id }).await; @@ -1536,6 +1543,7 @@ async fn run_core_loop( // it. Idempotent: an ordinary mute/unmute update just // re-records the same address. cancel_grace_timer(&grace_timers_events, &peer_id); + transport_events.admit_audio_sender(peer_id); transport_events.connect_peer(state.addr.clone()).await; // Auto-heal a friend's saved address (W7) on the // re-announce too — this is the path that catches a @@ -1577,6 +1585,7 @@ async fn run_core_loop( // hasn't recovered within RECONNECT_GRACE. A gossip // rejoin (PeerJoined/PeerUpdated) or a transport // reconnect (ConnEvent::Connected) cancels it first. + transport_events.keep_audio_sender_for_reconnect_grace(peer_id); let _ = ui_tx_events.send(UiEvent::PeerConnecting { id: peer_id }).await; arm_grace_timer( &grace_timers_events, diff --git a/src/network/iroh_impl.rs b/src/network/iroh_impl.rs index 6c733cb..2267207 100644 --- a/src/network/iroh_impl.rs +++ b/src/network/iroh_impl.rs @@ -56,6 +56,10 @@ struct Shared { /// 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>, incoming_tx: mpsc::Sender<(EndpointId, Bytes)>, /// Best-effort link-state notifications for the UI (connecting / connected). conn_events_tx: mpsc::Sender, @@ -112,6 +116,16 @@ impl Shared { 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. @@ -137,11 +151,39 @@ fn is_graceful_leave(err: &ConnectionError) -> bool { /// Pure S8 membership decision: iroh already authenticated `remote` as the /// connection's endpoint id, so audio admission is exactly live roster membership. -#[allow(dead_code)] // Design-first S8 seam; wiring waits for senior review. 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 @@ -340,10 +382,17 @@ impl iroh::protocol::ProtocolHandler for AudioRouter { 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 beat the gossip join event). try_send keeps the - // protocol handler from ever blocking; a full queue only happens if - // links are churning, and the supervisor will get the next one. + // 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)); @@ -376,6 +425,7 @@ impl IrohTransport { addrs: StdMutex::new(HashMap::new()), peers: tokio::sync::Mutex::new(HashMap::new()), live_conns: StdMutex::new(HashMap::new()), + admitted_audio: StdMutex::new(HashSet::new()), incoming_tx, conn_events_tx, }); @@ -404,11 +454,35 @@ impl IrohTransport { } self.shared.senders.lock().unwrap().clear(); self.shared.addrs.lock().unwrap().clear(); + self.shared.admitted_audio.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) + } } #[async_trait] @@ -499,4 +573,36 @@ mod tests { 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)); + } } diff --git a/tests/transport_loopback.rs b/tests/transport_loopback.rs index 29beb78..7638375 100644 --- a/tests/transport_loopback.rs +++ b/tests/transport_loopback.rs @@ -169,6 +169,9 @@ async fn loopback_sequenced_audio_reaches_peer_and_decodes() { 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); // Subscribe to incoming datagrams on B before any are sent. let mut b_rx = b.transport.receive_datagrams().await.expect("subscribe B");