Skip to content
Merged
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
62 changes: 21 additions & 41 deletions crates/ember-cluster/src/gossip.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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 {
Expand All @@ -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;
}
}

Expand Down Expand Up @@ -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;
}
}
}
Expand Down Expand Up @@ -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;
}
}
}
Expand All @@ -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;
}
}
}
Expand All @@ -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;
}
}
}
Expand Down Expand Up @@ -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);
}
Expand Down
20 changes: 10 additions & 10 deletions crates/ember-cluster/src/raft.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand All @@ -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 {
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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
}
Expand Down Expand Up @@ -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]
Expand Down Expand Up @@ -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]
Expand Down
15 changes: 15 additions & 0 deletions crates/ember-cluster/src/topology.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,15 @@ impl NodeId {
pub fn parse(s: &str) -> Result<Self, uuid::Error> {
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<String, ...>` storage.
pub(crate) fn as_key(&self) -> String {
self.0.to_string()
}
}

impl Default for NodeId {
Expand Down Expand Up @@ -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();

Expand Down