diff --git a/src/app/mod.rs b/src/app/mod.rs index fc6dfe5..52e75e1 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -888,6 +888,10 @@ fn update(state: &mut AppState, message: AppMessage) -> Task { state.ever_connected.remove(&id); notify::play(Sound::PeerLeave, state.config.custom_sound_peer_leave.as_deref()); } + // Core-only recovery phase: presentation for this state lands in + // the separate UI follow-up. In particular, do not play the + // terminal ReconnectFailed chime here. + UiEvent::PeerRecoveryStarted { .. } => {} UiEvent::PeerConnectionFailed { id } => { state.peers.remove(&id); state.audio_levels.remove(&id); diff --git a/src/core/messages.rs b/src/core/messages.rs index b83fddc..11b4b7f 100644 --- a/src/core/messages.rs +++ b/src/core/messages.rs @@ -85,6 +85,9 @@ pub enum UiEvent { RoomLeft, PeerJoined { id: EndpointId, state: PeerState }, PeerLeft { id: EndpointId }, + /// The fixed reconnect grace expired and bounded background gossip recovery + /// has started. This is non-terminal and must not play the failure chime. + PeerRecoveryStarted { id: EndpointId }, PeerConnectionFailed { id: EndpointId }, PeerUpdated { id: EndpointId, state: PeerState }, /// Audio link to a peer is being (re)established — show a connecting state. diff --git a/src/core/mod.rs b/src/core/mod.rs index afd8341..1b7d9c6 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -1,5 +1,6 @@ pub mod messages; pub mod jitter; +mod recovery; use crate::audio::{AudioBackend, PlatformAudioBackend}; use crate::audio::eq::{Eq, EqSettings}; @@ -11,6 +12,7 @@ use crate::network::{ gossip::IrohGossipState, }; use crate::core::messages::{CoreCommand, UiEvent}; +use crate::core::recovery::RecoveryCoordinator; use crate::config::{NetworkMode, RecordingMode}; use crate::presence::PresenceMode; @@ -102,6 +104,39 @@ type GraceTimers = Arc>>; +type KnownPeers = + Arc>>>; + +#[derive(Clone)] +struct RecoveryContext { + coordinator: RecoveryCoordinator, + room_state: Arc, + known_peers: KnownPeers, + ticket: String, +} + +impl RecoveryContext { + fn retained_addr(&self, peer_id: &EndpointId) -> Option { + self.known_peers + .lock() + .unwrap() + .get(&self.ticket) + .and_then(|peers| peers.get(peer_id)) + .cloned() + } + + fn cancel(&self, peer_id: EndpointId) { + self.coordinator.cancel(peer_id); + } + + fn forget(&self, peer_id: EndpointId) { + if let Some(peers) = self.known_peers.lock().unwrap().get_mut(&self.ticket) { + peers.remove(&peer_id); + } + self.coordinator.cancel(peer_id); + } +} + /// Cancel and forget a peer's pending grace timer, if any. No-op if none is armed. fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) { if let Some(handle) = timers.lock().unwrap().remove(peer_id) { @@ -116,12 +151,17 @@ fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) { /// link repeatedly resetting the clock and dodging eviction forever. On firing it /// also scrubs the peer from `seen_connected` so a later rejoin isn't treated as a /// reconnect on its initial dial. +struct GraceExpiry<'a> { + transport: &'a Arc, + jitter: &'a Arc>>, + ui_tx: &'a mpsc::Sender, + recovery: Option<&'a RecoveryContext>, +} + fn arm_grace_timer( timers: &GraceTimers, seen_connected: &SeenConnected, - transport: &Arc, - jitter: &Arc>>, - ui_tx: &mpsc::Sender, + expiry: GraceExpiry<'_>, grace: Duration, peer_id: EndpointId, ) { @@ -129,15 +169,29 @@ fn arm_grace_timer( if timers_guard.contains_key(&peer_id) { return; } - let transport_evict = transport.clone(); - let jitter_evict = jitter.clone(); - let ui_evict = ui_tx.clone(); + let transport_evict = expiry.transport.clone(); + let jitter_evict = expiry.jitter.clone(); + let ui_evict = expiry.ui_tx.clone(); let timers_evict = timers.clone(); let seen_evict = seen_connected.clone(); + let recovery_evict = expiry.recovery.cloned(); let handle = tokio::spawn(async move { tokio::time::sleep(grace).await; - crate::log_msg(&format!("Reconnect grace expired; evicting peer {:?}", peer_id)); + crate::log_msg(&format!("Reconnect grace expired for peer {:?}", peer_id)); + + if let Some(recovery) = &recovery_evict + && !recovery.coordinator.begin(peer_id) + { + return; + } + transport_evict.remove_audio_sender(peer_id); + if let Some(recovery) = &recovery_evict { + // Revoke roster authority before the first await in teardown. A + // verified Announce racing after this point is then a PeerJoined and + // cancels recovery instead of being erased after it was accepted. + recovery.room_state.mark_peer_disconnected(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 @@ -146,7 +200,38 @@ fn arm_grace_timer( // reconnect. timers_evict.lock().unwrap().remove(&peer_id); seen_evict.lock().unwrap().remove(&peer_id); - let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await; + + let Some(recovery) = recovery_evict else { + let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await; + return; + }; + if !recovery.coordinator.is_active(&peer_id) { + return; + } + + let Some(addr) = recovery.retained_addr(&peer_id) else { + crate::log_msg(&format!( + "Cannot recover peer {:?}: no retained authenticated address", + peer_id + )); + recovery.cancel(peer_id); + let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await; + return; + }; + + match recovery.coordinator.activate(peer_id, addr) { + Ok(true) => { + let _ = ui_evict.send(UiEvent::PeerRecoveryStarted { id: peer_id }).await; + } + Ok(false) => {} + Err(()) => { + crate::log_msg(&format!( + "Cannot recover peer {:?}: recovery coordinator unavailable", + peer_id + )); + let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await; + } + } }); timers_guard.insert(peer_id, handle); } @@ -309,6 +394,7 @@ pub struct ConnEventHandler { seen_connected: SeenConnected, transport: Arc, jitter: Arc>>, + recovery: Option, grace: Duration, } @@ -326,6 +412,7 @@ impl ConnEventHandler { seen_connected, transport, jitter, + recovery: None, grace: RECONNECT_GRACE, } } @@ -336,6 +423,11 @@ impl ConnEventHandler { self } + fn with_recovery(mut self, recovery: RecoveryContext) -> Self { + self.recovery = Some(recovery); + self + } + pub async fn handle(&self, event: ConnEvent) { match event { ConnEvent::Connecting(id) => { @@ -349,9 +441,12 @@ impl ConnEventHandler { arm_grace_timer( &self.grace_timers, &self.seen_connected, - &self.transport, - &self.jitter, - &self.ui_tx, + GraceExpiry { + transport: &self.transport, + jitter: &self.jitter, + ui_tx: &self.ui_tx, + recovery: self.recovery.as_ref(), + }, self.grace, id, ); @@ -359,6 +454,15 @@ impl ConnEventHandler { let _ = self.ui_tx.send(UiEvent::PeerConnecting { id }).await; } ConnEvent::Connected(id) => { + // A transport event cannot readmit a grace-expired peer. Ignore a + // stale/racing link until authenticated gossip emits PeerJoined. + if self + .recovery + .as_ref() + .is_some_and(|recovery| recovery.coordinator.is_active(&id)) + { + return; + } // The audio link came back — the peer recovered within the grace // window, so cancel its eviction. cancel_grace_timer(&self.grace_timers, &id); @@ -371,6 +475,9 @@ 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); + if let Some(recovery) = &self.recovery { + recovery.forget(id); + } self.transport.remove_audio_sender(id); self.transport.disconnect_peer(id).await; self.jitter.lock().await.remove(&id); @@ -387,6 +494,7 @@ struct ActiveSession { mixer_task: tokio::task::JoinHandle<()>, event_task: tokio::task::JoinHandle<()>, conn_event_task: tokio::task::JoinHandle<()>, + recovery_task: tokio::task::JoinHandle<()>, grace_timers: GraceTimers, transport: Arc, /// Loaded PipeWire echo-cancel module (if enabled); unloads on drop. @@ -421,6 +529,7 @@ impl ActiveSession { for (_, handle) in self.grace_timers.lock().unwrap().drain() { handle.abort(); } + self.recovery_task.abort(); crate::log_msg("Aborted tasks"); let audio_backend_clone = audio_backend.clone(); @@ -730,8 +839,7 @@ async fn run_core_loop( // first room's peers — the old single-set version cleared them on any ticket // change, so an A→B→A bounce stranded the rejoiner with an empty bootstrap. // Inner map keyed by peer id so updates refresh the address. - let known_peers: Arc>>> = - Arc::new(std::sync::Mutex::new(HashMap::new())); + let known_peers: KnownPeers = Arc::new(std::sync::Mutex::new(HashMap::new())); let audio_backend = Arc::new(PlatformAudioBackend::new()); @@ -1474,6 +1582,15 @@ async fn run_core_loop( // The ticket of the room this event loop serves, so peer add/remove // updates the right per-ticket bucket in `known_peers` (A8 archive). let ticket_events = ticket_str.clone(); + let (recovery_coordinator, recovery_task) = + RecoveryCoordinator::spawn(room_state.clone()); + let recovery_context = RecoveryContext { + coordinator: recovery_coordinator, + room_state: room_state.clone(), + known_peers: known_peers.clone(), + ticket: ticket_str.clone(), + }; + let recovery_events = recovery_context.clone(); // Friends store + ui sender, so a connected peer who is a friend has // their saved address auto-healed (W7) — populates `last_addr` so the // presence scheduler can reach them later. @@ -1486,6 +1603,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); + recovery_events.cancel(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). @@ -1529,15 +1647,9 @@ async fn run_core_loop( // Graceful leave — evict immediately. cancel_grace_timer(&grace_timers_events, &peer_id); seen_connected_events.lock().unwrap().remove(&peer_id); - // Graceful leave: drop them as a rejoin dial target - // for this room (a transient PeerConnectionLost - // deliberately does NOT, so we can still re-dial a - // peer who's still up). - if let Some(peers) = - known_peers_events.lock().unwrap().get_mut(&ticket_events) - { - peers.remove(&peer_id); - } + // A signed Leave cancels background recovery and + // drops the retained target. Transient loss keeps it. + recovery_events.forget(peer_id); transport_events.remove_audio_sender(peer_id); transport_events.disconnect_peer(peer_id).await; jitter_events.lock().await.remove(&peer_id); @@ -1551,6 +1663,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); + recovery_events.cancel(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 @@ -1598,9 +1711,12 @@ async fn run_core_loop( arm_grace_timer( &grace_timers_events, &seen_connected_events, - &transport_events, - &jitter_events, - &ui_tx_events, + GraceExpiry { + transport: &transport_events, + jitter: &jitter_events, + ui_tx: &ui_tx_events, + recovery: Some(&recovery_events), + }, RECONNECT_GRACE, peer_id, ); @@ -1624,7 +1740,8 @@ async fn run_core_loop( seen_connected.clone(), transport.clone(), jitter.clone(), - ); + ) + .with_recovery(recovery_context); let conn_event_task = tokio::spawn(async move { while let Some(event) = conn_events.recv().await { conn_handler.handle(event).await; @@ -1638,6 +1755,7 @@ async fn run_core_loop( mixer_task, event_task, conn_event_task, + recovery_task, grace_timers, transport: transport.clone(), #[cfg(target_os = "linux")] diff --git a/src/core/recovery.rs b/src/core/recovery.rs new file mode 100644 index 0000000..587c28c --- /dev/null +++ b/src/core/recovery.rs @@ -0,0 +1,264 @@ +use crate::network::{RoomState, gossip::IrohGossipState}; +use iroh::{EndpointAddr, EndpointId}; +use std::collections::{HashMap, HashSet}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use tokio::sync::mpsc; +use tokio::task::JoinHandle; +use tokio::time::Instant; + +const RECOVERY_COMMAND_CAPACITY: usize = 64; +const RECOVERY_DELAYS: [Duration; 7] = [ + Duration::from_secs(1), + Duration::from_secs(2), + Duration::from_secs(4), + Duration::from_secs(8), + Duration::from_secs(15), + Duration::from_secs(30), + Duration::from_secs(60), +]; + +fn recovery_delay(attempt: usize) -> Duration { + RECOVERY_DELAYS[attempt.min(RECOVERY_DELAYS.len() - 1)] +} + +enum RecoveryCommand { + Start { + peer_id: EndpointId, + addr: EndpointAddr, + }, + Cancel(EndpointId), +} + +struct RecoveryEntry { + addr: EndpointAddr, + attempt: usize, + next_attempt: Instant, +} + +#[async_trait::async_trait] +trait RecoveryRoom: Send + Sync { + async fn rebootstrap_peers(&self, peers: Vec) -> Result<(), String>; +} + +#[async_trait::async_trait] +impl RecoveryRoom for IrohGossipState { + async fn rebootstrap_peers(&self, peers: Vec) -> Result<(), String> { + RoomState::rebootstrap_peers(self, peers) + .await + .map_err(|error| error.to_string()) + } +} + +/// Cloneable command side of the single per-session recovery coordinator. +/// `active` is shared with transport/event handlers so cancellation is visible +/// immediately even while the coordinator is awaiting an in-flight gossip call. +#[derive(Clone)] +pub(super) struct RecoveryCoordinator { + tx: mpsc::Sender, + active: Arc>>, +} + +impl RecoveryCoordinator { + pub(super) fn spawn(room_state: Arc) -> (Self, JoinHandle<()>) { + Self::spawn_inner(room_state) + } + + fn spawn_inner(room_state: Arc) -> (Self, JoinHandle<()>) { + let (tx, rx) = mpsc::channel(RECOVERY_COMMAND_CAPACITY); + let active = Arc::new(Mutex::new(HashSet::new())); + let handle = Self { + tx, + active: active.clone(), + }; + let task = tokio::spawn(run_coordinator(room_state, active, rx)); + (handle, task) + } + + /// Reserve one recovery slot before grace-expiry teardown begins. Returns + /// false when the peer is already recovering, preventing duplicate work. + pub(super) fn begin(&self, peer_id: EndpointId) -> bool { + self.active.lock().unwrap().insert(peer_id) + } + + /// Activate the reserved slot with its retained authenticated address. + /// Uses a bounded non-blocking send while holding the active-set lock so a + /// concurrent cancellation is ordered before or after this command. + pub(super) fn activate(&self, peer_id: EndpointId, addr: EndpointAddr) -> Result { + let mut active = self.active.lock().unwrap(); + if !active.contains(&peer_id) { + return Ok(false); + } + if self + .tx + .try_send(RecoveryCommand::Start { peer_id, addr }) + .is_err() + { + active.remove(&peer_id); + return Err(()); + } + Ok(true) + } + + pub(super) fn cancel(&self, peer_id: EndpointId) { + self.active.lock().unwrap().remove(&peer_id); + // Cancellation is governed by the shared active set, so it remains + // immediate even if the bounded command queue is temporarily full. + let _ = self.tx.try_send(RecoveryCommand::Cancel(peer_id)); + } + + pub(super) fn is_active(&self, peer_id: &EndpointId) -> bool { + self.active.lock().unwrap().contains(peer_id) + } +} + +async fn run_coordinator( + room_state: Arc, + active: Arc>>, + mut rx: mpsc::Receiver, +) { + let mut entries: HashMap = HashMap::new(); + + loop { + // The shared active set is the authoritative cancellation gate. Prune + // here as well as on Cancel commands so a saturated command queue cannot + // leave an inactive, past-due entry spinning the timer loop. + let active_snapshot = active.lock().unwrap().clone(); + entries.retain(|peer_id, _| active_snapshot.contains(peer_id)); + let next_deadline = entries.values().map(|entry| entry.next_attempt).min(); + let command = match next_deadline { + Some(deadline) => { + tokio::select! { + command = rx.recv() => command, + _ = tokio::time::sleep_until(deadline) => { + let now = Instant::now(); + let active_snapshot = active.lock().unwrap().clone(); + let due: Vec<(EndpointId, EndpointAddr)> = entries + .iter() + .filter(|(id, entry)| { + entry.next_attempt <= now && active_snapshot.contains(*id) + }) + .map(|(id, entry)| (*id, entry.addr.clone())) + .collect(); + + if !due.is_empty() { + let addrs = due.iter().map(|(_, addr)| addr.clone()).collect(); + if let Err(error) = room_state.rebootstrap_peers(addrs).await { + crate::log_msg(&format!( + "Background peer recovery attempt failed: {error}" + )); + } + + let scheduled_at = Instant::now(); + for (peer_id, _) in due { + if !active.lock().unwrap().contains(&peer_id) { + entries.remove(&peer_id); + continue; + } + if let Some(entry) = entries.get_mut(&peer_id) { + entry.next_attempt = scheduled_at + recovery_delay(entry.attempt); + entry.attempt = entry.attempt.saturating_add(1); + } + } + } + continue; + } + } + } + None => rx.recv().await, + }; + + match command { + Some(RecoveryCommand::Start { peer_id, addr }) => { + if active.lock().unwrap().contains(&peer_id) { + entries.entry(peer_id).or_insert(RecoveryEntry { + addr, + attempt: 0, + next_attempt: Instant::now(), + }); + } + } + Some(RecoveryCommand::Cancel(peer_id)) => { + entries.remove(&peer_id); + } + None => break, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use iroh::SecretKey; + + struct RecordingRoom { + attempts: mpsc::UnboundedSender>, + } + + #[async_trait::async_trait] + impl RecoveryRoom for RecordingRoom { + async fn rebootstrap_peers(&self, peers: Vec) -> Result<(), String> { + self.attempts.send(peers).map_err(|error| error.to_string()) + } + } + + #[test] + fn retry_backoff_reaches_and_stays_at_sixty_seconds() { + let actual: Vec = (0..10) + .map(|attempt| recovery_delay(attempt).as_secs()) + .collect(); + assert_eq!(actual, vec![1, 2, 4, 8, 15, 30, 60, 60, 60, 60]); + } + + #[test] + fn recovery_slots_are_deduplicated_and_cancel_immediately() { + let (tx, mut rx) = mpsc::channel(4); + let coordinator = RecoveryCoordinator { + tx, + active: Arc::new(Mutex::new(HashSet::new())), + }; + let peer_id = SecretKey::generate().public(); + + assert!(coordinator.begin(peer_id)); + assert!( + !coordinator.begin(peer_id), + "a peer gets only one recovery slot" + ); + assert_eq!( + coordinator.activate(peer_id, EndpointAddr::from(peer_id)), + Ok(true) + ); + assert!(matches!( + rx.try_recv(), + Ok(RecoveryCommand::Start { peer_id: id, .. }) if id == peer_id + )); + + coordinator.cancel(peer_id); + assert!(!coordinator.is_active(&peer_id)); + assert!(matches!( + rx.try_recv(), + Ok(RecoveryCommand::Cancel(id)) if id == peer_id + )); + } + + #[tokio::test] + async fn coordinator_attempts_rebootstrap_immediately() { + let (attempts_tx, mut attempts_rx) = mpsc::unbounded_channel(); + let (coordinator, task) = RecoveryCoordinator::spawn_inner(Arc::new(RecordingRoom { + attempts: attempts_tx, + })); + let peer_id = SecretKey::generate().public(); + let addr = EndpointAddr::from(peer_id); + + assert!(coordinator.begin(peer_id)); + assert_eq!(coordinator.activate(peer_id, addr.clone()), Ok(true)); + let attempted = tokio::time::timeout(Duration::from_secs(1), attempts_rx.recv()) + .await + .expect("first recovery attempt should be immediate") + .expect("recording room remains subscribed"); + assert_eq!(attempted, vec![addr]); + + coordinator.cancel(peer_id); + task.abort(); + } +} diff --git a/src/network/gossip.rs b/src/network/gossip.rs index 355b055..242c1b3 100644 --- a/src/network/gossip.rs +++ b/src/network/gossip.rs @@ -5,7 +5,7 @@ use iroh_gossip::proto::TopicId; use tokio::sync::mpsc; use tokio::sync::mpsc::Receiver; use std::sync::{Arc, Mutex}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use async_trait::async_trait; use tokio_stream::StreamExt; use serde::{Serialize, Deserialize}; @@ -198,6 +198,10 @@ pub struct IrohGossipState { secret_key: SecretKey, self_state: Arc>>, peers: Arc>>, + /// Previously verified peers whose live roster entry was removed by a + /// transient disconnect. Retained only so a later authenticated `Leave` + /// still reaches core and cancels background recovery. + disconnected_peers: Arc>>, event_tx: mpsc::Sender, event_rx: Mutex>>, active_topic: Mutex>>, @@ -223,6 +227,7 @@ impl IrohGossipState { secret_key, self_state: Arc::new(Mutex::new(None)), peers: Arc::new(Mutex::new(HashMap::new())), + disconnected_peers: Arc::new(Mutex::new(HashSet::new())), event_tx, event_rx: Mutex::new(Some(event_rx)), active_topic: Mutex::new(None), @@ -295,6 +300,7 @@ impl RoomState for IrohGossipState { let event_tx = self.event_tx.clone(); let peers = self.peers.clone(); + let disconnected_peers = self.disconnected_peers.clone(); let address_lookup = self.address_lookup.clone(); let self_state_clone = self.self_state.clone(); let gossip_sender_clone = gossip_sender.clone(); @@ -393,6 +399,7 @@ impl RoomState for IrohGossipState { // peer-supplied: cap/validate once at ingest // so invalid offers never render a Watch button. state.sharing = state.sharing.and_then(crate::screenshare::sanitize_ticket); + disconnected_peers.lock().unwrap().remove(&payload.author); let (is_new, state_changed) = { let mut peer_map = peers.lock().unwrap(); let is_new = !peer_map.contains_key(&payload.author); @@ -423,7 +430,11 @@ impl RoomState for IrohGossipState { GossipMessage::Leave => { crate::log_msg(&format!("Gossip peer leave request from author={:?}", payload.author)); let removed = peers.lock().unwrap().remove(&payload.author).is_some(); - if removed { + let was_disconnected = disconnected_peers + .lock() + .unwrap() + .remove(&payload.author); + if removed || was_disconnected { let _ = event_tx.send(RoomEvent::PeerLeft(payload.author)).await; } } @@ -472,6 +483,7 @@ impl RoomState for IrohGossipState { // cached presence entry; a rejoin re-announces as new. let removed = peers.lock().unwrap().remove(&peer_id).is_some(); if removed { + disconnected_peers.lock().unwrap().insert(peer_id); crate::log_msg(&format!("Peer connection lost (NeighborDown): {:?}", peer_id)); let _ = event_tx.send(RoomEvent::PeerConnectionLost(peer_id)).await; } @@ -547,6 +559,12 @@ impl RoomState for IrohGossipState { .map_err(|e| NetError::Gossip(e.to_string())) } + fn mark_peer_disconnected(&self, peer_id: EndpointId) { + if self.peers.lock().unwrap().remove(&peer_id).is_some() { + self.disconnected_peers.lock().unwrap().insert(peer_id); + } + } + async fn send_chat(&self, text: String) -> Result<(), NetError> { let name = { let guard = self.self_state.lock().unwrap(); @@ -602,6 +620,7 @@ impl RoomState for IrohGossipState { } self.peers.lock().unwrap().clear(); + self.disconnected_peers.lock().unwrap().clear(); Ok(()) } diff --git a/src/network/mod.rs b/src/network/mod.rs index 37d41be..e57be30 100644 --- a/src/network/mod.rs +++ b/src/network/mod.rs @@ -200,6 +200,11 @@ pub trait RoomState: Send + Sync { /// active only after its normal signed `Announce` is received and verified. async fn rebootstrap_peers(&self, peers: Vec) -> Result<(), NetError>; + /// Remove a peer from the authenticated live roster before background + /// recovery. This only revokes membership; a fresh verified `Announce` is + /// required to add the peer again. + fn mark_peer_disconnected(&self, peer_id: EndpointId); + /// Broadcasts a room text-chat message authored by us (our display name is /// taken from the current self-state). async fn send_chat(&self, text: String) -> Result<(), NetError>; diff --git a/tests/gossip_rebootstrap.rs b/tests/gossip_rebootstrap.rs index f0c5e71..6478a13 100644 --- a/tests/gossip_rebootstrap.rs +++ b/tests/gossip_rebootstrap.rs @@ -78,6 +78,31 @@ async fn await_joined(rx: &mut mpsc::Receiver, peer_id: EndpointId) - } } +async fn await_joined_all(rx: &mut mpsc::Receiver, peer_ids: &[EndpointId]) { + let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT; + let mut remaining = peer_ids.to_vec(); + while !remaining.is_empty() { + match tokio::time::timeout_at(deadline, rx.recv()).await { + Ok(Some(RoomEvent::PeerJoined(id, _))) => remaining.retain(|wanted| *wanted != id), + Ok(Some(_)) => {} + Ok(None) => panic!("room event channel closed while waiting for PeerJoined set"), + Err(_) => panic!("timed out waiting for PeerJoined set: {remaining:?}"), + } + } +} + +async fn await_left(rx: &mut mpsc::Receiver, peer_id: EndpointId) { + let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT; + loop { + match tokio::time::timeout_at(deadline, rx.recv()).await { + Ok(Some(RoomEvent::PeerLeft(id))) if id == peer_id => return, + Ok(Some(_)) => continue, + Ok(None) => panic!("room event channel closed while waiting for PeerLeft"), + Err(_) => panic!("timed out waiting for PeerLeft({peer_id:?})"), + } + } +} + async fn await_absent(room: &IrohGossipState, peer_id: EndpointId) { let deadline = tokio::time::Instant::now() + EVENT_TIMEOUT; loop { @@ -215,3 +240,115 @@ async fn rebootstrap_uses_retained_full_address_with_empty_lookup() { let recovered = await_joined(&mut events_a, b_id).await; assert_eq!(recovered.name, "Bob rebound"); } + +#[tokio::test] +async fn demoted_peer_requires_a_fresh_signed_announce_to_rejoin() { + let a = spawn_node(SecretKey::generate()).await; + let b = spawn_node(SecretKey::generate()).await; + let ticket = ticket(b.endpoint.addr()); + let mut events_a = a.room.subscribe_events().await.expect("subscribe A events"); + + establish_room(&a, &b, &ticket, &mut events_a).await; + a.room.mark_peer_disconnected(b.endpoint.id()); + assert!( + !a.room + .active_peers() + .iter() + .any(|(id, _)| *id == b.endpoint.id()), + "demotion must revoke live roster membership" + ); + + b.room + .update_self_state(state("Bob authenticated again", b.endpoint.addr())) + .await + .expect("broadcast fresh signed announce"); + + let recovered = await_joined(&mut events_a, b.endpoint.id()).await; + assert_eq!(recovered.name, "Bob authenticated again"); +} + +#[tokio::test] +async fn signed_leave_after_demotion_still_emits_peer_left() { + let a = spawn_node(SecretKey::generate()).await; + let b = spawn_node(SecretKey::generate()).await; + let ticket = ticket(b.endpoint.addr()); + let mut events_a = a.room.subscribe_events().await.expect("subscribe A events"); + + establish_room(&a, &b, &ticket, &mut events_a).await; + a.room.mark_peer_disconnected(b.endpoint.id()); + + b.room + .leave() + .await + .expect("broadcast signed Leave after demotion"); + await_left(&mut events_a, b.endpoint.id()).await; +} + +#[tokio::test] +async fn targeted_rebootstrap_preserves_healthy_peer_in_three_peer_room() { + let a = spawn_node(SecretKey::generate()).await; + let b = spawn_node(SecretKey::generate()).await; + let c_secret = SecretKey::generate(); + let c = spawn_node(c_secret.clone()).await; + let c_id = c.endpoint.id(); + let ticket = ticket(c.endpoint.addr()); + let mut events_a = a.room.subscribe_events().await.expect("subscribe A events"); + let mut events_b = b.room.subscribe_events().await.expect("subscribe B events"); + + c.room + .join(&ticket, state("Carol", c.endpoint.addr()), vec![]) + .await + .expect("C hosts topic"); + b.room + .join(&ticket, state("Bob", b.endpoint.addr()), vec![]) + .await + .expect("B joins C"); + await_joined(&mut events_b, c_id).await; + a.room + .join(&ticket, state("Alice", a.endpoint.addr()), vec![]) + .await + .expect("A joins C"); + await_joined_all(&mut events_a, &[b.endpoint.id(), c_id]).await; + await_joined(&mut events_b, a.endpoint.id()).await; + + // Remove only C. A and B keep their existing topic subscriptions and remain + // mutually present while C is rebound to a fresh address. + c.endpoint.close().await; + drop(c); + a.room.mark_peer_disconnected(c_id); + b.room.mark_peer_disconnected(c_id); + let c_rebound = spawn_node(c_secret).await; + let rebound_addr = c_rebound.endpoint.addr(); + c_rebound + .room + .join( + &ticket, + state("Carol recovered", rebound_addr.clone()), + vec![], + ) + .await + .expect("rebound C rejoins as host without bootstrap"); + + a.room + .rebootstrap_peers(vec![rebound_addr]) + .await + .expect("A targets only C for recovery"); + let recovered = await_joined(&mut events_a, c_id).await; + assert_eq!(recovered.name, "Carol recovered"); + await_joined(&mut events_b, c_id).await; + + assert!( + a.room + .active_peers() + .iter() + .any(|(id, _)| *id == b.endpoint.id()), + "healthy B must remain present at A throughout C recovery" + ); + assert!( + b.room + .active_peers() + .iter() + .any(|(id, _)| *id == a.endpoint.id()), + "healthy A must remain present at B throughout C recovery" + ); +}