Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
110 changes: 102 additions & 8 deletions crates/buzz-relay-mesh/src/membership.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,14 +7,20 @@
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;
use crate::status::{ConnectionState, MeshCounters, MeshPeerCounters, MeshPeerStatus, MeshStatus};
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 {
Expand All @@ -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 {
Expand All @@ -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,
}
}

Expand All @@ -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()
Expand Down Expand Up @@ -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
}
Expand All @@ -144,14 +170,35 @@ impl MeshMembership {
},
record,
phi,
connection_state: ConnectionState::Connected,
connection_state: ConnectionState::Disconnected,
},
);
true
}
}
}

/// 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
Expand Down Expand Up @@ -389,7 +436,7 @@ impl RelayMeshMembership for MeshMembership {

#[cfg(test)]
mod tests {
use std::time::{Duration, UNIX_EPOCH};
use std::time::Duration;

use super::*;

Expand All @@ -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),
}
}

Expand Down Expand Up @@ -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 changes: 5 additions & 0 deletions crates/buzz-relay-mesh/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
42 changes: 28 additions & 14 deletions crates/buzz-relay/src/mesh_boot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading