fix(relay): restore mesh echo and evict stale peers

Co-authored-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co>
Signed-off-by: npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr <5a968df9a7494b4e019b9ecf739e088ba61097b4312124e9a88ae5b42e3f5f3e@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
npub1t2tgm7d8f995uqvmnm8h88sg3wnpp9a5xysjf6dg3tjmgt3ltulqdp8ehr
2026-07-17 15:50:45 -04:00
parent 0e03836018
commit 87a159b735
3 changed files with 135 additions and 22 deletions
+102 -8
View File
@@ -7,7 +7,7 @@
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use std::time::SystemTime;
use std::time::{Duration, SystemTime};
use crate::gossip::{system_time_from_millis, GossipRecord, PhiAccrual};
use crate::registry::ReadyRecord;
@@ -15,6 +15,12 @@ use crate::status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerSta
use crate::{PeerInfo, RelayMeshMembership, RuntimeId};
pub const DEFAULT_PHI_SUSPECT_THRESHOLD: f64 = 8.0;
/// Gossip records older than this are no longer useful as routing hints.
///
/// This is deliberately much longer than the ready-record TTL (45 seconds at
/// the default refresh) and the gossip interval (2 seconds), so brief registry
/// or network stalls do not churn healthy peers.
pub const DEFAULT_STALE_PEER_TIMEOUT: Duration = Duration::from_secs(5 * 60);
#[derive(Clone, Debug)]
struct PeerState {
@@ -40,6 +46,7 @@ pub struct MeshMembership {
/// unanchored state is fail-closed, not accept-any.
expected_relay_pubkey: Option<String>,
phi_suspect_threshold: f64,
stale_peer_timeout: Duration,
}
impl MeshMembership {
@@ -53,6 +60,7 @@ impl MeshMembership {
foreign_relay_rejections: Arc::new(AtomicU64::new(0)),
expected_relay_pubkey: None,
phi_suspect_threshold: DEFAULT_PHI_SUSPECT_THRESHOLD,
stale_peer_timeout: DEFAULT_STALE_PEER_TIMEOUT,
}
}
@@ -68,6 +76,14 @@ impl MeshMembership {
self
}
/// Override the maximum heartbeat age retained in membership.
/// Primarily useful for deterministic tests; production uses
/// [`DEFAULT_STALE_PEER_TIMEOUT`].
pub fn with_stale_peer_timeout(mut self, timeout: Duration) -> Self {
self.stale_peer_timeout = timeout;
self
}
pub fn local_record(&self) -> GossipRecord {
self.local_record
.read()
@@ -123,12 +139,22 @@ impl MeshMembership {
}
let heartbeat = system_time_from_millis(record.heartbeat_millis);
if heartbeat
.elapsed()
.is_ok_and(|age| age >= self.stale_peer_timeout)
{
tracing::debug!(
runtime_id = %record.runtime_id,
heartbeat_millis = record.heartbeat_millis,
"mesh membership ignored stale gossip record"
);
return false;
}
let mut peers = self.peers.write().expect("membership lock poisoned");
match peers.get_mut(&record.runtime_id) {
Some(peer) if record.version <= peer.record.version => false,
Some(peer) => {
peer.record = record;
peer.connection_state = ConnectionState::Connected;
peer.phi.observe(heartbeat);
true
}
@@ -144,7 +170,7 @@ impl MeshMembership {
},
record,
phi,
connection_state: ConnectionState::Connected,
connection_state: ConnectionState::Disconnected,
},
);
true
@@ -152,6 +178,27 @@ impl MeshMembership {
}
}
/// Remove peers whose last gossiped heartbeat is too old to be useful.
/// Returns their runtime ids so the runtime can tear down lingering
/// connection state at the same reconcile boundary.
pub fn evict_stale_peers(&self) -> Vec<RuntimeId> {
let now = SystemTime::now();
let mut peers = self.peers.write().expect("membership lock poisoned");
let stale: Vec<_> = peers
.iter()
.filter_map(|(runtime_id, peer)| {
let heartbeat = system_time_from_millis(peer.record.heartbeat_millis);
now.duration_since(heartbeat)
.is_ok_and(|age| age >= self.stale_peer_timeout)
.then_some(*runtime_id)
})
.collect();
for runtime_id in &stale {
peers.remove(runtime_id);
}
stale
}
pub fn mark_connection_state(&self, runtime_id: RuntimeId, state: ConnectionState) {
if let Some(peer) = self
.peers
@@ -389,7 +436,7 @@ impl RelayMeshMembership for MeshMembership {
#[cfg(test)]
mod tests {
use std::time::{Duration, UNIX_EPOCH};
use std::time::Duration;
use super::*;
@@ -406,10 +453,8 @@ mod tests {
draining: false,
capabilities: vec!["reliable-stream".to_string()],
version,
heartbeat_millis: (UNIX_EPOCH + Duration::from_secs(heartbeat_secs))
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis() as u64,
heartbeat_millis: crate::gossip::now_millis()
.saturating_add(Duration::from_secs(heartbeat_secs).as_millis() as u64),
}
}
@@ -498,4 +543,53 @@ mod tests {
assert!(membership.local_record().draining);
assert_eq!(membership.local_record().version, 2);
}
#[test]
fn stale_peers_are_evicted_and_cannot_be_resurrected() {
let membership =
MeshMembership::new(record(1, 1, 0)).with_stale_peer_timeout(Duration::from_secs(5));
let mut stale = record(2, 1, 0);
stale.heartbeat_millis = crate::gossip::now_millis() - 10_000;
assert!(!membership.apply_gossip_record(stale.clone()));
assert!(membership.peers().is_empty());
let mut live = record(2, 2, 0);
assert!(membership.apply_gossip_record(live.clone()));
assert_eq!(membership.status().peer_count, 1);
live.heartbeat_millis = crate::gossip::now_millis() - 10_000;
// Simulate a previously admitted record aging out in place.
membership
.peers
.write()
.expect("membership lock poisoned")
.get_mut(&rid(2))
.expect("peer present")
.record
.heartbeat_millis = live.heartbeat_millis;
assert_eq!(membership.evict_stale_peers(), vec![rid(2)]);
assert!(membership.peers().is_empty());
assert!(!membership.apply_gossip_record(live));
}
#[test]
fn gossip_does_not_claim_transport_connection() {
let membership = MeshMembership::new(record(1, 1, 0));
assert!(membership.apply_gossip_record(record(2, 1, 0)));
assert_eq!(
membership.status().peers[0].connection_state,
ConnectionState::Disconnected
);
membership.mark_connection_state(rid(2), ConnectionState::Connected);
assert_eq!(
membership.status().peers[0].connection_state,
ConnectionState::Connected
);
assert!(membership.apply_gossip_record(record(2, 2, 1)));
assert_eq!(
membership.status().peers[0].connection_state,
ConnectionState::Connected
);
}
}
+5
View File
@@ -327,6 +327,11 @@ async fn reconcile_loop(inner: Arc<Inner>) {
}
async fn reconcile_once(inner: &Arc<Inner>) {
for runtime_id in inner.membership.evict_stale_peers() {
remove_peer(inner, runtime_id);
tracing::info!(peer = %runtime_id, "mesh: evicted stale peer");
}
if let Some(registry) = &inner.registry {
match registry.scan_ready().await {
Ok(records) => inner.membership.apply_ready_records(records),
+28 -14
View File
@@ -315,25 +315,39 @@ pub(crate) async fn run_demo_echo(
tracing::info!(%session_id, %peer, "mesh demo echo: session open");
let mut drain_tick = tokio::time::interval(std::time::Duration::from_millis(100));
loop {
let frame = tokio::select! {
_ = drain_tick.tick() => {
if shutting_down.load(Ordering::Relaxed) {
if let Some(community_id) = stream.community_id() {
if let Err(e) = stream.send_goodbye(community_id, GoodbyeReason::Draining).await {
tracing::warn!(%session_id, "mesh demo echo: draining goodbye failed: {e}");
} else {
tracing::info!(%session_id, "mesh demo echo: sent draining goodbye");
// recv_validated reads the frame before asynchronously checking its
// Redis fence. Keep that future alive across drain polls: cancelling
// and recreating it after a tick would discard an already-read frame.
let frame = {
let recv = stream.recv_validated(&directory);
tokio::pin!(recv);
loop {
tokio::select! {
frame = &mut recv => break frame,
_ = drain_tick.tick() => {
if shutting_down.load(Ordering::Relaxed) {
break Ok(None);
}
} else {
let _ = stream.finish();
tracing::info!(%session_id, "mesh demo echo: drain before community latch — closing");
}
return;
}
continue;
}
frame = stream.recv_validated(&directory) => frame,
};
if shutting_down.load(Ordering::Relaxed) {
if let Some(community_id) = stream.community_id() {
if let Err(e) = stream
.send_goodbye(community_id, GoodbyeReason::Draining)
.await
{
tracing::warn!(%session_id, "mesh demo echo: draining goodbye failed: {e}");
} else {
tracing::info!(%session_id, "mesh demo echo: sent draining goodbye");
}
} else {
let _ = stream.finish();
tracing::info!(%session_id, "mesh demo echo: drain before community latch — closing");
}
return;
}
match frame {
Ok(Some(ReliableFrame::Data(payload))) => {
// recv_validated latched the community from the frame it just