fix: arm reconnect grace timer from the transport, not just gossip
A second sustained outage after a reconnect never evicted the peer: the 45s RECONNECT_GRACE timer was armed only by the gossip PeerConnectionLost path, but a transport-only retained-addr reconnect leaves gossip's neighbor state stale, so the second drop produced no new PeerConnectionLost and no timer. Arm the grace timer from the transport ConnEvent::Connecting too (the supervisor reliably re-emits it on every outage), gated on a new seen_connected set so a first-ever dial isn't given an eviction clock. Route both arming sites through a shared arm_grace_timer helper that is a no-op if a timer is already pending (earliest drop notice sets one hard deadline; a flapping link can't reset it), and scrub seen_connected on eviction/leave so a later rejoin starts clean. Compiles, clippy-clean, all tests green incl. transport reconnect suite. NOT yet field-verified — pending the two-outage laptop test. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+73
-20
@@ -15,7 +15,7 @@ use crate::config::NetworkMode;
|
|||||||
use iroh::{Endpoint, EndpointId, RelayMode, endpoint::presets, protocol::Router};
|
use iroh::{Endpoint, EndpointId, RelayMode, endpoint::presets, protocol::Router};
|
||||||
use iroh_gossip::net::Gossip;
|
use iroh_gossip::net::Gossip;
|
||||||
use tokio::sync::{mpsc, Mutex};
|
use tokio::sync::{mpsc, Mutex};
|
||||||
use std::collections::HashMap;
|
use std::collections::{HashMap, HashSet};
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
@@ -58,6 +58,11 @@ const RECONNECT_GRACE: Duration = Duration::from_secs(45);
|
|||||||
/// actually comes back).
|
/// actually comes back).
|
||||||
type GraceTimers = Arc<std::sync::Mutex<HashMap<EndpointId, tokio::task::JoinHandle<()>>>>;
|
type GraceTimers = Arc<std::sync::Mutex<HashMap<EndpointId, tokio::task::JoinHandle<()>>>>;
|
||||||
|
|
||||||
|
/// Peers we've completed at least one audio link with. Lets the conn-event task
|
||||||
|
/// tell a genuine reconnect (arm an eviction timer) from a first-ever dial (don't).
|
||||||
|
/// Scrubbed whenever a peer is evicted or leaves so a later rejoin starts clean.
|
||||||
|
type SeenConnected = Arc<std::sync::Mutex<HashSet<EndpointId>>>;
|
||||||
|
|
||||||
/// Cancel and forget a peer's pending grace timer, if any. No-op if none is armed.
|
/// 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) {
|
fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) {
|
||||||
if let Some(handle) = timers.lock().unwrap().remove(peer_id) {
|
if let Some(handle) = timers.lock().unwrap().remove(peer_id) {
|
||||||
@@ -65,6 +70,42 @@ fn cancel_grace_timer(timers: &GraceTimers, peer_id: &EndpointId) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Arm a per-peer reconnect grace timer that evicts the peer if its link hasn't
|
||||||
|
/// recovered within [`RECONNECT_GRACE`]. No-op if a timer is already pending for
|
||||||
|
/// the peer, so the earliest drop notice — whether the gossip `PeerConnectionLost`
|
||||||
|
/// or the transport `Connecting` — sets one hard deadline, rather than a flapping
|
||||||
|
/// 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.
|
||||||
|
fn arm_grace_timer(
|
||||||
|
timers: &GraceTimers,
|
||||||
|
seen_connected: &SeenConnected,
|
||||||
|
transport: &Arc<IrohTransport>,
|
||||||
|
jitter: &Arc<Mutex<HashMap<EndpointId, JitterBuffer>>>,
|
||||||
|
ui_tx: &mpsc::Sender<UiEvent>,
|
||||||
|
peer_id: EndpointId,
|
||||||
|
) {
|
||||||
|
let mut timers_guard = timers.lock().unwrap();
|
||||||
|
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 timers_evict = timers.clone();
|
||||||
|
let seen_evict = seen_connected.clone();
|
||||||
|
let handle = tokio::spawn(async move {
|
||||||
|
tokio::time::sleep(RECONNECT_GRACE).await;
|
||||||
|
crate::log_msg(&format!("Reconnect grace expired; evicting peer {:?}", peer_id));
|
||||||
|
transport_evict.disconnect_peer(peer_id).await;
|
||||||
|
jitter_evict.lock().await.remove(&peer_id);
|
||||||
|
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
|
||||||
|
timers_evict.lock().unwrap().remove(&peer_id);
|
||||||
|
seen_evict.lock().unwrap().remove(&peer_id);
|
||||||
|
});
|
||||||
|
timers_guard.insert(peer_id, handle);
|
||||||
|
}
|
||||||
|
|
||||||
struct ActiveSession {
|
struct ActiveSession {
|
||||||
endpoint: Endpoint,
|
endpoint: Endpoint,
|
||||||
router: Router,
|
router: Router,
|
||||||
@@ -458,6 +499,8 @@ async fn run_core_loop(
|
|||||||
let transport_events = transport.clone();
|
let transport_events = transport.clone();
|
||||||
let grace_timers: GraceTimers = Arc::new(std::sync::Mutex::new(HashMap::new()));
|
let grace_timers: GraceTimers = Arc::new(std::sync::Mutex::new(HashMap::new()));
|
||||||
let grace_timers_events = grace_timers.clone();
|
let grace_timers_events = grace_timers.clone();
|
||||||
|
let seen_connected: SeenConnected = Arc::new(std::sync::Mutex::new(HashSet::new()));
|
||||||
|
let seen_connected_events = seen_connected.clone();
|
||||||
let event_task = tokio::spawn(async move {
|
let event_task = tokio::spawn(async move {
|
||||||
while let Some(event) = room_events.recv().await {
|
while let Some(event) = room_events.recv().await {
|
||||||
match event {
|
match event {
|
||||||
@@ -475,6 +518,7 @@ async fn run_core_loop(
|
|||||||
RoomEvent::PeerLeft(peer_id) => {
|
RoomEvent::PeerLeft(peer_id) => {
|
||||||
// Graceful leave — evict immediately.
|
// Graceful leave — evict immediately.
|
||||||
cancel_grace_timer(&grace_timers_events, &peer_id);
|
cancel_grace_timer(&grace_timers_events, &peer_id);
|
||||||
|
seen_connected_events.lock().unwrap().remove(&peer_id);
|
||||||
transport_events.disconnect_peer(peer_id).await;
|
transport_events.disconnect_peer(peer_id).await;
|
||||||
jitter_events.lock().await.remove(&peer_id);
|
jitter_events.lock().await.remove(&peer_id);
|
||||||
let _ = ui_tx_events.send(UiEvent::PeerLeft { id: peer_id }).await;
|
let _ = ui_tx_events.send(UiEvent::PeerLeft { id: peer_id }).await;
|
||||||
@@ -499,25 +543,14 @@ async fn run_core_loop(
|
|||||||
// rejoin (PeerJoined/PeerUpdated) or a transport
|
// rejoin (PeerJoined/PeerUpdated) or a transport
|
||||||
// reconnect (ConnEvent::Connected) cancels it first.
|
// reconnect (ConnEvent::Connected) cancels it first.
|
||||||
let _ = ui_tx_events.send(UiEvent::PeerConnecting { id: peer_id }).await;
|
let _ = ui_tx_events.send(UiEvent::PeerConnecting { id: peer_id }).await;
|
||||||
let transport_evict = transport_events.clone();
|
arm_grace_timer(
|
||||||
let jitter_evict = jitter_events.clone();
|
&grace_timers_events,
|
||||||
let ui_evict = ui_tx_events.clone();
|
&seen_connected_events,
|
||||||
let timers_evict = grace_timers_events.clone();
|
&transport_events,
|
||||||
let handle = tokio::spawn(async move {
|
&jitter_events,
|
||||||
tokio::time::sleep(RECONNECT_GRACE).await;
|
&ui_tx_events,
|
||||||
crate::log_msg(&format!(
|
peer_id,
|
||||||
"Reconnect grace expired; evicting peer {:?}", peer_id
|
);
|
||||||
));
|
|
||||||
transport_evict.disconnect_peer(peer_id).await;
|
|
||||||
jitter_evict.lock().await.remove(&peer_id);
|
|
||||||
let _ = ui_evict.send(UiEvent::PeerConnectionFailed { id: peer_id }).await;
|
|
||||||
timers_evict.lock().unwrap().remove(&peer_id);
|
|
||||||
});
|
|
||||||
// Replace (and abort) any timer already pending for
|
|
||||||
// this peer so repeated drops don't stack up.
|
|
||||||
if let Some(old) = grace_timers_events.lock().unwrap().insert(peer_id, handle) {
|
|
||||||
old.abort();
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -534,18 +567,37 @@ async fn run_core_loop(
|
|||||||
};
|
};
|
||||||
let ui_tx_conn = ui_tx.clone();
|
let ui_tx_conn = ui_tx.clone();
|
||||||
let grace_timers_conn = grace_timers.clone();
|
let grace_timers_conn = grace_timers.clone();
|
||||||
|
let seen_connected_conn = seen_connected.clone();
|
||||||
let transport_conn = transport.clone();
|
let transport_conn = transport.clone();
|
||||||
let jitter_conn = jitter.clone();
|
let jitter_conn = jitter.clone();
|
||||||
let conn_event_task = tokio::spawn(async move {
|
let conn_event_task = tokio::spawn(async move {
|
||||||
while let Some(event) = conn_events.recv().await {
|
while let Some(event) = conn_events.recv().await {
|
||||||
match event {
|
match event {
|
||||||
ConnEvent::Connecting(id) => {
|
ConnEvent::Connecting(id) => {
|
||||||
|
// A reconnect (we've linked with this peer before):
|
||||||
|
// arm an eviction timer so a peer that never comes
|
||||||
|
// back is cleared even when gossip doesn't re-report
|
||||||
|
// the drop — the transport reliably re-emits this on
|
||||||
|
// every outage, gossip's NeighborDown does not. A
|
||||||
|
// first-ever dial (not yet in seen_connected) gets no
|
||||||
|
// timer; ConnEvent::Connected cancels it on recovery.
|
||||||
|
if seen_connected_conn.lock().unwrap().contains(&id) {
|
||||||
|
arm_grace_timer(
|
||||||
|
&grace_timers_conn,
|
||||||
|
&seen_connected_conn,
|
||||||
|
&transport_conn,
|
||||||
|
&jitter_conn,
|
||||||
|
&ui_tx_conn,
|
||||||
|
id,
|
||||||
|
);
|
||||||
|
}
|
||||||
let _ = ui_tx_conn.send(UiEvent::PeerConnecting { id }).await;
|
let _ = ui_tx_conn.send(UiEvent::PeerConnecting { id }).await;
|
||||||
}
|
}
|
||||||
ConnEvent::Connected(id) => {
|
ConnEvent::Connected(id) => {
|
||||||
// The audio link came back — the peer recovered
|
// The audio link came back — the peer recovered
|
||||||
// within the grace window, so cancel its eviction.
|
// within the grace window, so cancel its eviction.
|
||||||
cancel_grace_timer(&grace_timers_conn, &id);
|
cancel_grace_timer(&grace_timers_conn, &id);
|
||||||
|
seen_connected_conn.lock().unwrap().insert(id);
|
||||||
let _ = ui_tx_conn.send(UiEvent::PeerConnected { id }).await;
|
let _ = ui_tx_conn.send(UiEvent::PeerConnected { id }).await;
|
||||||
}
|
}
|
||||||
ConnEvent::Left(id) => {
|
ConnEvent::Left(id) => {
|
||||||
@@ -554,6 +606,7 @@ async fn run_core_loop(
|
|||||||
// instead of leaving it "reconnecting" until the
|
// instead of leaving it "reconnecting" until the
|
||||||
// grace timer or the slow gossip Leave.
|
// grace timer or the slow gossip Leave.
|
||||||
cancel_grace_timer(&grace_timers_conn, &id);
|
cancel_grace_timer(&grace_timers_conn, &id);
|
||||||
|
seen_connected_conn.lock().unwrap().remove(&id);
|
||||||
transport_conn.disconnect_peer(id).await;
|
transport_conn.disconnect_peer(id).await;
|
||||||
jitter_conn.lock().await.remove(&id);
|
jitter_conn.lock().await.remove(&id);
|
||||||
let _ = ui_tx_conn.send(UiEvent::PeerLeft { id }).await;
|
let _ = ui_tx_conn.send(UiEvent::PeerLeft { id }).await;
|
||||||
|
|||||||
Reference in New Issue
Block a user