From 620f88f35fa5830432412f770c6b2965b07b845b Mon Sep 17 00:00:00 2001 From: Kacy Fortner Date: Tue, 17 Feb 2026 18:36:40 -0500 Subject: [PATCH] refactor(cluster): extract emit helper, add NodeId::as_key, clarify format_flags gossip.rs - extract `async fn emit(&self, event: GossipEvent)` that sends on event_tx and warns when the channel is closed. replaces 6 copies of the inline send-and-check pattern in apply_updates and mark_alive. MemberJoined and MemberSuspected were previously fire-and-forget (let _ = ...); they now consistently warn on channel closure like the other events. - add comment to queue_update explaining why dropping the oldest pending updates is safe (convergence doesn't require every update to be delivered). topology.rs - add NodeId::as_key() returning the full UUID string for use as a map key. NodeId::Display intentionally shows only 8 chars for readability, which is too short to use as a unique key in storage. - add doc comment to format_flags explaining the two intentional divergences from NodeFlags::Display: role is included, and pfail renders as "fail?" per the Redis cluster protocol spec rather than "pfail". raft.rs - replace node_id.0.to_string() with node_id.as_key() throughout apply_command. the UUID string access is now named rather than relying on tuple field indexing. --- crates/ember-cluster/src/gossip.rs | 62 ++++++++++------------------ crates/ember-cluster/src/raft.rs | 20 ++++----- crates/ember-cluster/src/topology.rs | 15 +++++++ 3 files changed, 46 insertions(+), 51 deletions(-) diff --git a/crates/ember-cluster/src/gossip.rs b/crates/ember-cluster/src/gossip.rs index 827f3fcf..9cafbbd3 100644 --- a/crates/ember-cluster/src/gossip.rs +++ b/crates/ember-cluster/src/gossip.rs @@ -374,6 +374,17 @@ impl GossipEngine { }); } + /// Sends a gossip event to the external event channel. + /// + /// Logs a warning when the channel is closed (receiver dropped). This + /// normally only happens during shutdown, so seeing the message in steady + /// state indicates a bug in the event consumer. + async fn emit(&self, event: GossipEvent) { + if self.event_tx.send(event).await.is_err() { + warn!("gossip event channel closed, dropping event"); + } + } + async fn apply_updates(&mut self, updates: &[NodeUpdate]) { for update in updates { match update { @@ -400,14 +411,7 @@ impl GossipEngine { if member.state != MemberStatus::Alive { member.state = MemberStatus::Alive; member.state_change = Instant::now(); - if self - .event_tx - .send(GossipEvent::MemberAlive(*node)) - .await - .is_err() - { - warn!("event channel closed, cannot send MemberAlive event"); - } + self.emit(GossipEvent::MemberAlive(*node)).await; } } } else { @@ -423,10 +427,7 @@ impl GossipEngine { slots: Vec::new(), }, ); - let _ = self - .event_tx - .send(GossipEvent::MemberJoined(*node, *addr)) - .await; + self.emit(GossipEvent::MemberJoined(*node, *addr)).await; } } @@ -455,10 +456,7 @@ impl GossipEngine { { member.state = MemberStatus::Suspect; member.state_change = Instant::now(); - let _ = self - .event_tx - .send(GossipEvent::MemberSuspected(*node)) - .await; + self.emit(GossipEvent::MemberSuspected(*node)).await; } } } @@ -486,14 +484,7 @@ impl GossipEngine { { member.state = MemberStatus::Dead; member.state_change = Instant::now(); - if self - .event_tx - .send(GossipEvent::MemberFailed(*node)) - .await - .is_err() - { - warn!("event channel closed, cannot send MemberFailed event"); - } + self.emit(GossipEvent::MemberFailed(*node)).await; } } } @@ -506,14 +497,7 @@ impl GossipEngine { if member.state != MemberStatus::Left { member.state = MemberStatus::Left; member.state_change = Instant::now(); - if self - .event_tx - .send(GossipEvent::MemberLeft(*node)) - .await - .is_err() - { - warn!("event channel closed, cannot send MemberLeft event"); - } + self.emit(GossipEvent::MemberLeft(*node)).await; } } } @@ -526,14 +510,7 @@ impl GossipEngine { if member.state == MemberStatus::Suspect { member.state = MemberStatus::Alive; member.state_change = Instant::now(); - if self - .event_tx - .send(GossipEvent::MemberAlive(node)) - .await - .is_err() - { - warn!("event channel closed, cannot send MemberAlive event"); - } + self.emit(GossipEvent::MemberAlive(node)).await; } } } @@ -603,7 +580,10 @@ impl GossipEngine { fn queue_update(&mut self, update: NodeUpdate) { self.pending_updates.push(update); - // Keep bounded + // When the queue overflows, drop the oldest pending updates. + // This is safe: gossip convergence doesn't require every update to be + // delivered. Members re-gossip their state on each protocol period, so + // a dropped update will be re-sent in the next round. if self.pending_updates.len() > self.config.max_piggyback * 2 { self.pending_updates.drain(0..self.config.max_piggyback); } diff --git a/crates/ember-cluster/src/raft.rs b/crates/ember-cluster/src/raft.rs index 1dc6c215..8f650d08 100644 --- a/crates/ember-cluster/src/raft.rs +++ b/crates/ember-cluster/src/raft.rs @@ -157,7 +157,7 @@ impl Storage { addr, is_primary, } => { - let key = node_id.0.to_string(); + let key = node_id.as_key(); state.nodes.insert( key.clone(), NodeInfo { @@ -172,7 +172,7 @@ impl Storage { } ClusterCommand::RemoveNode { node_id } => { - let key = node_id.0.to_string(); + let key = node_id.as_key(); state.nodes.remove(&key); state.slots.retain(|_, owner| owner != &key); ClusterResponse::Ok @@ -190,7 +190,7 @@ impl Storage { )); } } - let key = node_id.0.to_string(); + let key = node_id.as_key(); if let Some(node) = state.nodes.get_mut(&key) { node.slots = slots.clone(); for slot_range in slots { @@ -205,7 +205,7 @@ impl Storage { } ClusterCommand::PromoteReplica { replica_id } => { - let key = replica_id.0.to_string(); + let key = replica_id.as_key(); if let Some(node) = state.nodes.get_mut(&key) { node.is_primary = true; ClusterResponse::Ok @@ -224,8 +224,8 @@ impl Storage { state.migrations.insert( *slot, MigrationState { - from: from.0.to_string(), - to: to.0.to_string(), + from: from.as_key(), + to: to.as_key(), }, ); ClusterResponse::Ok @@ -238,7 +238,7 @@ impl Storage { )); } state.migrations.remove(slot); - let key = new_owner.0.to_string(); + let key = new_owner.as_key(); state.slots.insert(*slot, key); ClusterResponse::Ok } @@ -474,7 +474,7 @@ mod tests { let state_arc = storage.state(); let state = state_arc.read().await; - assert!(state.nodes.contains_key(&node_id.0.to_string())); + assert!(state.nodes.contains_key(&node_id.as_key())); } #[tokio::test] @@ -515,8 +515,8 @@ mod tests { let state_arc = storage.state(); let state = state_arc.read().await; - assert_eq!(state.slots.get(&0), Some(&node_id.0.to_string())); - assert_eq!(state.slots.get(&5460), Some(&node_id.0.to_string())); + assert_eq!(state.slots.get(&0), Some(&node_id.as_key())); + assert_eq!(state.slots.get(&5460), Some(&node_id.as_key())); } #[tokio::test] diff --git a/crates/ember-cluster/src/topology.rs b/crates/ember-cluster/src/topology.rs index bb079850..51f4405e 100644 --- a/crates/ember-cluster/src/topology.rs +++ b/crates/ember-cluster/src/topology.rs @@ -29,6 +29,15 @@ impl NodeId { pub fn parse(s: &str) -> Result { Ok(Self(Uuid::parse_str(s)?)) } + + /// Returns the full UUID string for use as a map key. + /// + /// Note: `Display` for `NodeId` shows only the first 8 characters (for + /// readability in logs). This method returns the full UUID needed when + /// `NodeId` is used as a key in `BTreeMap` storage. + pub(crate) fn as_key(&self) -> String { + self.0.to_string() + } } impl Default for NodeId { @@ -243,6 +252,12 @@ impl ClusterNode { .to_string() } + /// Formats the node flags for the CLUSTER NODES wire response. + /// + /// Intentionally diverges from `NodeFlags::Display` in two ways: + /// 1. Role (`master`/`slave`) is included here but is not a flag field. + /// 2. `pfail` renders as `fail?` per the Redis cluster protocol spec, + /// whereas `NodeFlags::Display` uses the field name `pfail`. fn format_flags(&self) -> String { let mut flags = Vec::new();