From 1c07f02be2754191464927e5be0deea11d197151 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 13 Feb 2026 09:21:59 +0000 Subject: [PATCH] fix: SWIM Suspect recovery + piggyback ordering race condition Two bugs found by simulation testing: 1. Suspect nodes could never recover after dissemination budget expired. When probing a target, the Suspect state was not re-enqueued into the dissemination queue (only Dead was). After budget exhaustion, the Suspect node never received a piggyback telling it it was suspected, so it could never refute via incarnation bump. Fix: re-enqueue both Suspect and Dead state on outgoing probes. 2. Registry entries merged before membership changes in same piggyback. When a message carried both a death notification and registry entries, the entries were merged first, then immediately tombstoned. Fix: split extract_registry_piggyback into unpack + deferred merge, processing membership changes before merging registry entries. Authored by Claude, lovingly guided by Zachery Aaron Shores-Chmielewski --- crates/distribution/src/node.rs | 24 ++++++++++++++---------- crates/distribution/src/swim/node.rs | 10 +++++----- 2 files changed, 19 insertions(+), 15 deletions(-) diff --git a/crates/distribution/src/node.rs b/crates/distribution/src/node.rs index a11d18c..3d7d87b 100644 --- a/crates/distribution/src/node.rs +++ b/crates/distribution/src/node.rs @@ -14,7 +14,7 @@ use crate::kademlia::repair::{RepairQueue, RepublishTracker}; use crate::kademlia::routing_table::RoutingTable; use crate::registry::{ pack_combined_piggyback, unpack_combined_piggyback, ClusterRegistry, RegistryConfig, - RegistryEvent, + RegistryEntry, RegistryEvent, }; use crate::swim::node::{NodeAction, SwimNode}; use crate::swim::probe::SwimConfig; @@ -157,24 +157,30 @@ impl DistributedNode { // ─── SWIM message handling (delegate to SwimNode) ─────────────────── 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 (membership_bytes, registry_entries) = unpack_combined_piggyback(piggyback); let actions = self.swim.handle_ping(from, from_addr, sequence, &membership_bytes); + // Process membership BEFORE merging registry — otherwise a death + // notification in this same piggyback would immediately tombstone + // freshly received registry entries instead of pre-existing ones. self.process_membership_changes(&actions); + self.merge_registry_entries(registry_entries); self.maybe_update_routing_table(from, from_addr); self.inject_registry_piggyback(actions) } pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec { - let membership_bytes = self.extract_registry_piggyback(piggyback); + let (membership_bytes, registry_entries) = unpack_combined_piggyback(piggyback); let actions = self.swim.handle_ack(from, sequence, &membership_bytes); self.process_membership_changes(&actions); + self.merge_registry_entries(registry_entries); 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 (membership_bytes, registry_entries) = unpack_combined_piggyback(piggyback); let actions = self.swim.handle_ping_req(from, target, target_addr, sequence, &membership_bytes); self.process_membership_changes(&actions); + self.merge_registry_entries(registry_entries); self.inject_registry_piggyback(actions) } @@ -362,13 +368,11 @@ impl DistributedNode { .collect() } - /// Extract registry entries from incoming piggyback, merge them, return membership-only bytes. - fn extract_registry_piggyback(&mut self, bytes: &[u8]) -> Vec { - let (membership_bytes, registry_entries) = unpack_combined_piggyback(bytes); - if !registry_entries.is_empty() { - self.registry.merge_batch(registry_entries, self.cluster_size()); + /// Merge registry entries received from a piggyback payload. + fn merge_registry_entries(&mut self, entries: Vec) { + if !entries.is_empty() { + self.registry.merge_batch(entries, self.cluster_size()); } - membership_bytes } } diff --git a/crates/distribution/src/swim/node.rs b/crates/distribution/src/swim/node.rs index 27696bc..6abffc6 100644 --- a/crates/distribution/src/swim/node.rs +++ b/crates/distribution/src/swim/node.rs @@ -272,14 +272,14 @@ impl SwimNode { for pa in probe_actions { match pa { SwimAction::SendPing { to, to_addr, sequence } => { - // If the target is dead, re-enqueue the death declaration + // If the target is suspect or dead, re-enqueue its state // so it piggybacks on this message. This is the key mechanism - // for partition-heal recovery: the dead node learns it was - // declared dead and refutes by bumping its incarnation. + // for partition-heal recovery: the target learns it was + // suspected/declared dead and refutes by bumping its incarnation. if let Some(entry) = self.members.get(&to) { - if entry.state == MemberState::Dead { + if entry.state == MemberState::Dead || entry.state == MemberState::Suspect { self.dissemination.enqueue( - membership_update(to, to_addr, MemberState::Dead, entry.incarnation), + membership_update(to, to_addr, entry.state, entry.incarnation), self.cluster_size(), ); }