Add post-grace peer recovery coordinator

This commit is contained in:
2026-06-20 15:37:33 -04:00
parent 5564af02f9
commit 2c93c1c24f
7 changed files with 578 additions and 28 deletions
+4
View File
@@ -888,6 +888,10 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
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);
+3
View File
@@ -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.
+144 -26
View File
@@ -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<std::sync::Mutex<HashMap<EndpointId, tokio::task::JoinHan
/// Scrubbed whenever a peer is evicted or leaves so a later rejoin starts clean.
type SeenConnected = Arc<std::sync::Mutex<HashSet<EndpointId>>>;
type KnownPeers =
Arc<std::sync::Mutex<HashMap<String, HashMap<EndpointId, EndpointAddr>>>>;
#[derive(Clone)]
struct RecoveryContext {
coordinator: RecoveryCoordinator,
room_state: Arc<IrohGossipState>,
known_peers: KnownPeers,
ticket: String,
}
impl RecoveryContext {
fn retained_addr(&self, peer_id: &EndpointId) -> Option<EndpointAddr> {
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<IrohTransport>,
jitter: &'a Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
ui_tx: &'a mpsc::Sender<UiEvent>,
recovery: Option<&'a RecoveryContext>,
}
fn arm_grace_timer(
timers: &GraceTimers,
seen_connected: &SeenConnected,
transport: &Arc<IrohTransport>,
jitter: &Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
ui_tx: &mpsc::Sender<UiEvent>,
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<IrohTransport>,
jitter: Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
recovery: Option<RecoveryContext>,
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<IrohTransport>,
/// 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<std::sync::Mutex<HashMap<String, HashMap<EndpointId, EndpointAddr>>>> =
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")]
+264
View File
@@ -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<EndpointAddr>) -> Result<(), String>;
}
#[async_trait::async_trait]
impl RecoveryRoom for IrohGossipState {
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> 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<RecoveryCommand>,
active: Arc<Mutex<HashSet<EndpointId>>>,
}
impl RecoveryCoordinator {
pub(super) fn spawn(room_state: Arc<IrohGossipState>) -> (Self, JoinHandle<()>) {
Self::spawn_inner(room_state)
}
fn spawn_inner(room_state: Arc<dyn RecoveryRoom>) -> (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<bool, ()> {
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<dyn RecoveryRoom>,
active: Arc<Mutex<HashSet<EndpointId>>>,
mut rx: mpsc::Receiver<RecoveryCommand>,
) {
let mut entries: HashMap<EndpointId, RecoveryEntry> = 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<Vec<EndpointAddr>>,
}
#[async_trait::async_trait]
impl RecoveryRoom for RecordingRoom {
async fn rebootstrap_peers(&self, peers: Vec<EndpointAddr>) -> 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<u64> = (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();
}
}
+21 -2
View File
@@ -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<Mutex<Option<PeerState>>>,
peers: Arc<Mutex<HashMap<EndpointId, PeerState>>>,
/// 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<Mutex<HashSet<EndpointId>>>,
event_tx: mpsc::Sender<RoomEvent>,
event_rx: Mutex<Option<mpsc::Receiver<RoomEvent>>>,
active_topic: Mutex<Option<tokio::task::JoinHandle<()>>>,
@@ -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(())
}
+5
View File
@@ -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<EndpointAddr>) -> 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>;
+137
View File
@@ -78,6 +78,31 @@ async fn await_joined(rx: &mut mpsc::Receiver<RoomEvent>, peer_id: EndpointId) -
}
}
async fn await_joined_all(rx: &mut mpsc::Receiver<RoomEvent>, 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<RoomEvent>, 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"
);
}