From 5d8e413db251e0fe954c279f6cec214f6bd4e106 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 13 Feb 2026 08:39:14 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20three=20SWIM=20notification=20bugs=20?= =?UTF-8?q?=E2=80=94=20death=20dissemination,=20piggyback=20notifications,?= =?UTF-8?q?=20partition=20recovery?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bug 1 (swim/node.rs): SwimProbe::check_suspicion_timeouts() calls members.declare_dead() before translate_probe_actions() processes the DeclareDead action. The second declare_dead() returned false (already dead), so MembershipChanged{Dead} was never emitted and the death was never enqueued for dissemination. Fix: remove the redundant declare_dead() call in translate_probe_actions since the probe already performed the mutation. Bug 2 (swim/node.rs): apply_membership_update() — which processes piggyback on every ping/ack/ping_req — updated the internal member list but never emitted NodeAction::MembershipChanged. This meant DistributedNode was blind to all state transitions learned via gossip piggyback (e.g., a dead node refuting via incarnation bump). Fix: return MembershipChanged actions from apply_piggyback and propagate through handle_ping/handle_ack/handle_ping_req. Bug 3 (node.rs): DistributedNode::handle_ping/handle_ack/handle_ping_req never processed MembershipChanged actions from SwimNode — only tick() did. Fix: extract process_membership_changes() helper and call it from all four message paths (tick, handle_ping, handle_ack, handle_ping_req). Additional fixes: - registry.rs: add re_disseminate_all() for anti-entropy on partition heal - node.rs: call re_disseminate_all on MemberState::Alive transitions so registry state accumulated during partition reaches recovering nodes - cluster_scenarios: enable dead_reprobe in 10% message loss test, since correct death dissemination (now working) causes cascading false deaths without a recovery mechanism Authored by Claude, lovingly guided by Zachery Aaron Shores-Chmielewski --- crates/distribution/src/node.rs | 26 ++++---- crates/distribution/src/registry.rs | 11 ++++ crates/distribution/src/swim/node.rs | 59 +++++++++++-------- crates/simulation/tests/cluster_scenarios.rs | 11 ++-- .../simulation/tests/distribution_registry.rs | 2 - 5 files changed, 68 insertions(+), 41 deletions(-) diff --git a/crates/distribution/src/node.rs b/crates/distribution/src/node.rs index 6dc9a5f..a11d18c 100644 --- a/crates/distribution/src/node.rs +++ b/crates/distribution/src/node.rs @@ -137,17 +137,7 @@ impl DistributedNode { let actions = self.swim.tick(); // Process membership changes from SWIM - let membership_changes: Vec<_> = actions - .iter() - .filter_map(|a| match a { - NodeAction::MembershipChanged { node_id, state, .. } => Some((*node_id, *state)), - _ => None, - }) - .collect(); - - for (node_id, state) in membership_changes { - self.handle_membership_change(node_id, state); - } + self.process_membership_changes(&actions); // Periodic republish let to_republish = self.republish.tick(self.tick_count); @@ -169,6 +159,7 @@ impl DistributedNode { pub fn handle_ping(&mut self, from: NodeId, from_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec { let membership_bytes = self.extract_registry_piggyback(piggyback); let actions = self.swim.handle_ping(from, from_addr, sequence, &membership_bytes); + self.process_membership_changes(&actions); self.maybe_update_routing_table(from, from_addr); self.inject_registry_piggyback(actions) } @@ -176,12 +167,14 @@ impl DistributedNode { pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec { let membership_bytes = self.extract_registry_piggyback(piggyback); let actions = self.swim.handle_ack(from, sequence, &membership_bytes); + self.process_membership_changes(&actions); self.inject_registry_piggyback(actions) } pub fn handle_ping_req(&mut self, from: NodeId, target: NodeId, target_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec { let membership_bytes = self.extract_registry_piggyback(piggyback); let actions = self.swim.handle_ping_req(from, target, target_addr, sequence, &membership_bytes); + self.process_membership_changes(&actions); self.inject_registry_piggyback(actions) } @@ -310,12 +303,23 @@ impl DistributedNode { self.routing_table.insert(node_id, addr); } + fn process_membership_changes(&mut self, actions: &[NodeAction]) { + for action in actions { + if let NodeAction::MembershipChanged { node_id, state, .. } = action { + self.handle_membership_change(*node_id, *state); + } + } + } + fn handle_membership_change(&mut self, node_id: NodeId, state: MemberState) { match state { MemberState::Alive => { if let Some(entry) = self.swim.members().get(&node_id) { self.routing_table.insert(node_id, entry.addr); } + // Re-disseminate registry entries so the recovering node + // catches up on state accumulated during the partition. + self.registry.re_disseminate_all(self.cluster_size()); } MemberState::Dead => { self.routing_table.remove(&node_id); diff --git a/crates/distribution/src/registry.rs b/crates/distribution/src/registry.rs index 377c257..72f3690 100644 --- a/crates/distribution/src/registry.rs +++ b/crates/distribution/src/registry.rs @@ -235,6 +235,17 @@ impl ClusterRegistry { } } + /// Re-enqueue all entries for dissemination (anti-entropy on membership change). + /// + /// Called when a previously-dead node comes back alive, ensuring that + /// registry state accumulated during a partition is gossiped to the + /// recovering node. + pub fn re_disseminate_all(&mut self, cluster_size: usize) { + for entry in self.entries.values().cloned().collect::>() { + self.enqueue(entry, cluster_size); + } + } + /// Periodic GC: remove tombstones past TTL with exhausted dissemination budgets. pub fn gc_tick(&mut self) { self.tick_count += 1; diff --git a/crates/distribution/src/swim/node.rs b/crates/distribution/src/swim/node.rs index 94c3dc8..27696bc 100644 --- a/crates/distribution/src/swim/node.rs +++ b/crates/distribution/src/swim/node.rs @@ -85,29 +85,31 @@ impl SwimNode { /// Handle a received ping. pub fn handle_ping(&mut self, from: NodeId, from_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec { - self.apply_piggyback(piggyback); + let mut actions = self.apply_piggyback(piggyback); // Ensure the sender is in our member list self.members.apply(from, from_addr, MemberState::Alive, 0); // Reply with ack let pb = self.dissemination.pack_piggyback(self.max_piggyback); - vec![NodeAction::SendAck { + actions.push(NodeAction::SendAck { to: from, to_addr: from_addr, sequence, piggyback: pb, - }] + }); + actions } /// Handle a received ack. pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec { - self.apply_piggyback(piggyback); + let mut actions = self.apply_piggyback(piggyback); let probe_actions = self.probe.step( SwimEvent::AckReceived { from, sequence }, &mut self.members, ); - self.translate_probe_actions(probe_actions) + actions.extend(self.translate_probe_actions(probe_actions)); + actions } /// Handle a received indirect ping request. @@ -119,16 +121,17 @@ impl SwimNode { sequence: u64, piggyback: &[u8], ) -> Vec { - self.apply_piggyback(piggyback); + let mut actions = self.apply_piggyback(piggyback); // Forward a ping to the target on behalf of the requester let pb = self.dissemination.pack_piggyback(self.max_piggyback); - vec![NodeAction::SendPing { + actions.push(NodeAction::SendPing { to: target, to_addr: target_addr, sequence, piggyback: pb, - }] + }); + actions } /// Handle a join request from a new node. @@ -214,14 +217,16 @@ impl SwimNode { self.members.alive_count() + 1 // +1 for self } - fn apply_piggyback(&mut self, bytes: &[u8]) { + fn apply_piggyback(&mut self, bytes: &[u8]) -> Vec { let updates = DisseminationQueue::unpack_piggyback(bytes); + let mut actions = Vec::new(); for update in updates { - self.apply_membership_update(update); + actions.extend(self.apply_membership_update(update)); } + actions } - fn apply_membership_update(&mut self, update: MembershipUpdate) { + fn apply_membership_update(&mut self, update: MembershipUpdate) -> Vec { // Check if this is about us if update.node_id == self.members.self_id() { if update.state == MemberState::Suspect || update.state == MemberState::Dead { @@ -237,7 +242,7 @@ impl SwimNode { self.cluster_size(), ); } - return; + return Vec::new(); } let changed = self.members.apply( @@ -252,6 +257,13 @@ impl SwimNode { membership_update(update.node_id, update.addr, update.state, update.incarnation), self.cluster_size(), ); + vec![NodeAction::MembershipChanged { + node_id: update.node_id, + state: update.state, + incarnation: update.incarnation, + }] + } else { + Vec::new() } } @@ -307,20 +319,21 @@ impl SwimNode { } } SwimAction::DeclareDead(node_id) => { + // Note: declare_dead() was already called by SwimProbe::check_suspicion_timeouts(), + // so we must NOT call it again (it would return false since state is already Dead). + // We just need to disseminate the update and emit the MembershipChanged action. if let Some(entry) = self.members.get(&node_id) { let inc = entry.incarnation; let addr = entry.addr; - if self.members.declare_dead(node_id) { - self.dissemination.enqueue( - membership_update(node_id, addr, MemberState::Dead, inc), - self.cluster_size(), - ); - actions.push(NodeAction::MembershipChanged { - node_id, - state: MemberState::Dead, - incarnation: inc, - }); - } + self.dissemination.enqueue( + membership_update(node_id, addr, MemberState::Dead, inc), + self.cluster_size(), + ); + actions.push(NodeAction::MembershipChanged { + node_id, + state: MemberState::Dead, + incarnation: inc, + }); } } SwimAction::Refute { new_incarnation } => { diff --git a/crates/simulation/tests/cluster_scenarios.rs b/crates/simulation/tests/cluster_scenarios.rs index 6c256c9..e858963 100644 --- a/crates/simulation/tests/cluster_scenarios.rs +++ b/crates/simulation/tests/cluster_scenarios.rs @@ -126,8 +126,9 @@ fn cluster_converges_under_10_percent_message_loss() { // Given: 5 nodes with 10% message loss from the start. // 10% loss is significant for SWIM because it can hit both direct probe // AND indirect probes in the same cycle, causing false suspicions. - // We verify the cluster degrades but doesn't crash, and at least some - // membership information survives. + // Dead reprobe is enabled so false deaths can self-correct — without it, + // correct death dissemination (via piggyback) causes cascading false deaths + // that collapse the entire cluster under even modest message loss. let config = DistributionSimConfig { name: "message-loss-10pct".into(), num_nodes: 5, @@ -139,7 +140,7 @@ fn cluster_converges_under_10_percent_message_loss() { probe_timeout: 5, indirect_probes: 2, suspicion_timeout: 20, - dead_reprobe_interval: 0, + dead_reprobe_interval: 30, }, network_faults: vec![NetworkFault::SetDropRate { round: 1, @@ -151,8 +152,8 @@ fn cluster_converges_under_10_percent_message_loss() { let trace = run_simulation(config); let metrics = analyze(&trace); - // With 10% loss and the deterministic LCG, SWIM's probe cycle is disrupted - // enough to cause false deaths. The test verifies: + // With 10% loss, SWIM's probe cycle is disrupted enough to cause + // false suspicions. Dead reprobe allows recovery. We verify: // 1. The simulation completes without panic (implicit — we got here) // 2. At least partial membership is maintained (some nodes still know about others) let result = check_membership_accuracy(&metrics, 0.15); diff --git a/crates/simulation/tests/distribution_registry.rs b/crates/simulation/tests/distribution_registry.rs index ec1ae96..d5babda 100644 --- a/crates/simulation/tests/distribution_registry.rs +++ b/crates/simulation/tests/distribution_registry.rs @@ -72,7 +72,6 @@ fn registry_name_converges_across_cluster() { // ──────────────────────────────────────────────────────────────────────────── #[test] -#[ignore = "BUG: SwimProbe::check_suspicion_timeouts calls declare_dead before translate_probe_actions, so MembershipChanged{Dead} is never emitted"] fn split_brain_naming_converges_after_partition_heals() { // Given: 6-node cluster, partition {0,1,2} vs {3,4,5} at round 10 // Node 0 registers "leader" at round 12, node 3 registers "leader" at round 12 @@ -142,7 +141,6 @@ fn split_brain_naming_converges_after_partition_heals() { // ──────────────────────────────────────────────────────────────────────────── #[test] -#[ignore = "BUG: SwimProbe::check_suspicion_timeouts calls declare_dead before translate_probe_actions, so MembershipChanged{Dead} is never emitted"] fn tombstone_propagates_when_name_owner_dies() { // Given: 5-node cluster, node 0 registers "svc" at round 5, killed at round 15 let config = DistributionSimConfig {