fix: three SWIM notification bugs — death dissemination, piggyback notifications, partition recovery
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
This commit is contained in:
parent
c84e708846
commit
5d8e413db2
5 changed files with 68 additions and 41 deletions
|
|
@ -137,17 +137,7 @@ impl DistributedNode {
|
||||||
let actions = self.swim.tick();
|
let actions = self.swim.tick();
|
||||||
|
|
||||||
// Process membership changes from SWIM
|
// Process membership changes from SWIM
|
||||||
let membership_changes: Vec<_> = actions
|
self.process_membership_changes(&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);
|
|
||||||
}
|
|
||||||
|
|
||||||
// Periodic republish
|
// Periodic republish
|
||||||
let to_republish = self.republish.tick(self.tick_count);
|
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<NodeAction> {
|
pub fn handle_ping(&mut self, from: NodeId, from_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
||||||
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
||||||
let actions = self.swim.handle_ping(from, from_addr, sequence, &membership_bytes);
|
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.maybe_update_routing_table(from, from_addr);
|
||||||
self.inject_registry_piggyback(actions)
|
self.inject_registry_piggyback(actions)
|
||||||
}
|
}
|
||||||
|
|
@ -176,12 +167,14 @@ impl DistributedNode {
|
||||||
pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
||||||
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
||||||
let actions = self.swim.handle_ack(from, sequence, &membership_bytes);
|
let actions = self.swim.handle_ack(from, sequence, &membership_bytes);
|
||||||
|
self.process_membership_changes(&actions);
|
||||||
self.inject_registry_piggyback(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<NodeAction> {
|
pub fn handle_ping_req(&mut self, from: NodeId, target: NodeId, target_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
||||||
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
let membership_bytes = self.extract_registry_piggyback(piggyback);
|
||||||
let actions = self.swim.handle_ping_req(from, target, target_addr, sequence, &membership_bytes);
|
let actions = self.swim.handle_ping_req(from, target, target_addr, sequence, &membership_bytes);
|
||||||
|
self.process_membership_changes(&actions);
|
||||||
self.inject_registry_piggyback(actions)
|
self.inject_registry_piggyback(actions)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -310,12 +303,23 @@ impl DistributedNode {
|
||||||
self.routing_table.insert(node_id, addr);
|
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) {
|
fn handle_membership_change(&mut self, node_id: NodeId, state: MemberState) {
|
||||||
match state {
|
match state {
|
||||||
MemberState::Alive => {
|
MemberState::Alive => {
|
||||||
if let Some(entry) = self.swim.members().get(&node_id) {
|
if let Some(entry) = self.swim.members().get(&node_id) {
|
||||||
self.routing_table.insert(node_id, entry.addr);
|
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 => {
|
MemberState::Dead => {
|
||||||
self.routing_table.remove(&node_id);
|
self.routing_table.remove(&node_id);
|
||||||
|
|
|
||||||
|
|
@ -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::<Vec<_>>() {
|
||||||
|
self.enqueue(entry, cluster_size);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Periodic GC: remove tombstones past TTL with exhausted dissemination budgets.
|
/// Periodic GC: remove tombstones past TTL with exhausted dissemination budgets.
|
||||||
pub fn gc_tick(&mut self) {
|
pub fn gc_tick(&mut self) {
|
||||||
self.tick_count += 1;
|
self.tick_count += 1;
|
||||||
|
|
|
||||||
|
|
@ -85,29 +85,31 @@ impl SwimNode {
|
||||||
|
|
||||||
/// Handle a received ping.
|
/// Handle a received ping.
|
||||||
pub fn handle_ping(&mut self, from: NodeId, from_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
pub fn handle_ping(&mut self, from: NodeId, from_addr: SocketAddr, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
||||||
self.apply_piggyback(piggyback);
|
let mut actions = self.apply_piggyback(piggyback);
|
||||||
|
|
||||||
// Ensure the sender is in our member list
|
// Ensure the sender is in our member list
|
||||||
self.members.apply(from, from_addr, MemberState::Alive, 0);
|
self.members.apply(from, from_addr, MemberState::Alive, 0);
|
||||||
|
|
||||||
// Reply with ack
|
// Reply with ack
|
||||||
let pb = self.dissemination.pack_piggyback(self.max_piggyback);
|
let pb = self.dissemination.pack_piggyback(self.max_piggyback);
|
||||||
vec![NodeAction::SendAck {
|
actions.push(NodeAction::SendAck {
|
||||||
to: from,
|
to: from,
|
||||||
to_addr: from_addr,
|
to_addr: from_addr,
|
||||||
sequence,
|
sequence,
|
||||||
piggyback: pb,
|
piggyback: pb,
|
||||||
}]
|
});
|
||||||
|
actions
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Handle a received ack.
|
/// Handle a received ack.
|
||||||
pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
pub fn handle_ack(&mut self, from: NodeId, sequence: u64, piggyback: &[u8]) -> Vec<NodeAction> {
|
||||||
self.apply_piggyback(piggyback);
|
let mut actions = self.apply_piggyback(piggyback);
|
||||||
let probe_actions = self.probe.step(
|
let probe_actions = self.probe.step(
|
||||||
SwimEvent::AckReceived { from, sequence },
|
SwimEvent::AckReceived { from, sequence },
|
||||||
&mut self.members,
|
&mut self.members,
|
||||||
);
|
);
|
||||||
self.translate_probe_actions(probe_actions)
|
actions.extend(self.translate_probe_actions(probe_actions));
|
||||||
|
actions
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Handle a received indirect ping request.
|
/// Handle a received indirect ping request.
|
||||||
|
|
@ -119,16 +121,17 @@ impl SwimNode {
|
||||||
sequence: u64,
|
sequence: u64,
|
||||||
piggyback: &[u8],
|
piggyback: &[u8],
|
||||||
) -> Vec<NodeAction> {
|
) -> Vec<NodeAction> {
|
||||||
self.apply_piggyback(piggyback);
|
let mut actions = self.apply_piggyback(piggyback);
|
||||||
|
|
||||||
// Forward a ping to the target on behalf of the requester
|
// Forward a ping to the target on behalf of the requester
|
||||||
let pb = self.dissemination.pack_piggyback(self.max_piggyback);
|
let pb = self.dissemination.pack_piggyback(self.max_piggyback);
|
||||||
vec![NodeAction::SendPing {
|
actions.push(NodeAction::SendPing {
|
||||||
to: target,
|
to: target,
|
||||||
to_addr: target_addr,
|
to_addr: target_addr,
|
||||||
sequence,
|
sequence,
|
||||||
piggyback: pb,
|
piggyback: pb,
|
||||||
}]
|
});
|
||||||
|
actions
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Handle a join request from a new node.
|
/// Handle a join request from a new node.
|
||||||
|
|
@ -214,14 +217,16 @@ impl SwimNode {
|
||||||
self.members.alive_count() + 1 // +1 for self
|
self.members.alive_count() + 1 // +1 for self
|
||||||
}
|
}
|
||||||
|
|
||||||
fn apply_piggyback(&mut self, bytes: &[u8]) {
|
fn apply_piggyback(&mut self, bytes: &[u8]) -> Vec<NodeAction> {
|
||||||
let updates = DisseminationQueue::unpack_piggyback(bytes);
|
let updates = DisseminationQueue::unpack_piggyback(bytes);
|
||||||
|
let mut actions = Vec::new();
|
||||||
for update in updates {
|
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<NodeAction> {
|
||||||
// Check if this is about us
|
// Check if this is about us
|
||||||
if update.node_id == self.members.self_id() {
|
if update.node_id == self.members.self_id() {
|
||||||
if update.state == MemberState::Suspect || update.state == MemberState::Dead {
|
if update.state == MemberState::Suspect || update.state == MemberState::Dead {
|
||||||
|
|
@ -237,7 +242,7 @@ impl SwimNode {
|
||||||
self.cluster_size(),
|
self.cluster_size(),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
return;
|
return Vec::new();
|
||||||
}
|
}
|
||||||
|
|
||||||
let changed = self.members.apply(
|
let changed = self.members.apply(
|
||||||
|
|
@ -252,6 +257,13 @@ impl SwimNode {
|
||||||
membership_update(update.node_id, update.addr, update.state, update.incarnation),
|
membership_update(update.node_id, update.addr, update.state, update.incarnation),
|
||||||
self.cluster_size(),
|
self.cluster_size(),
|
||||||
);
|
);
|
||||||
|
vec![NodeAction::MembershipChanged {
|
||||||
|
node_id: update.node_id,
|
||||||
|
state: update.state,
|
||||||
|
incarnation: update.incarnation,
|
||||||
|
}]
|
||||||
|
} else {
|
||||||
|
Vec::new()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -307,10 +319,12 @@ impl SwimNode {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
SwimAction::DeclareDead(node_id) => {
|
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) {
|
if let Some(entry) = self.members.get(&node_id) {
|
||||||
let inc = entry.incarnation;
|
let inc = entry.incarnation;
|
||||||
let addr = entry.addr;
|
let addr = entry.addr;
|
||||||
if self.members.declare_dead(node_id) {
|
|
||||||
self.dissemination.enqueue(
|
self.dissemination.enqueue(
|
||||||
membership_update(node_id, addr, MemberState::Dead, inc),
|
membership_update(node_id, addr, MemberState::Dead, inc),
|
||||||
self.cluster_size(),
|
self.cluster_size(),
|
||||||
|
|
@ -322,7 +336,6 @@ impl SwimNode {
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
SwimAction::Refute { new_incarnation } => {
|
SwimAction::Refute { new_incarnation } => {
|
||||||
self.dissemination.enqueue(
|
self.dissemination.enqueue(
|
||||||
membership_update(
|
membership_update(
|
||||||
|
|
|
||||||
|
|
@ -126,8 +126,9 @@ fn cluster_converges_under_10_percent_message_loss() {
|
||||||
// Given: 5 nodes with 10% message loss from the start.
|
// Given: 5 nodes with 10% message loss from the start.
|
||||||
// 10% loss is significant for SWIM because it can hit both direct probe
|
// 10% loss is significant for SWIM because it can hit both direct probe
|
||||||
// AND indirect probes in the same cycle, causing false suspicions.
|
// AND indirect probes in the same cycle, causing false suspicions.
|
||||||
// We verify the cluster degrades but doesn't crash, and at least some
|
// Dead reprobe is enabled so false deaths can self-correct — without it,
|
||||||
// membership information survives.
|
// correct death dissemination (via piggyback) causes cascading false deaths
|
||||||
|
// that collapse the entire cluster under even modest message loss.
|
||||||
let config = DistributionSimConfig {
|
let config = DistributionSimConfig {
|
||||||
name: "message-loss-10pct".into(),
|
name: "message-loss-10pct".into(),
|
||||||
num_nodes: 5,
|
num_nodes: 5,
|
||||||
|
|
@ -139,7 +140,7 @@ fn cluster_converges_under_10_percent_message_loss() {
|
||||||
probe_timeout: 5,
|
probe_timeout: 5,
|
||||||
indirect_probes: 2,
|
indirect_probes: 2,
|
||||||
suspicion_timeout: 20,
|
suspicion_timeout: 20,
|
||||||
dead_reprobe_interval: 0,
|
dead_reprobe_interval: 30,
|
||||||
},
|
},
|
||||||
network_faults: vec![NetworkFault::SetDropRate {
|
network_faults: vec![NetworkFault::SetDropRate {
|
||||||
round: 1,
|
round: 1,
|
||||||
|
|
@ -151,8 +152,8 @@ fn cluster_converges_under_10_percent_message_loss() {
|
||||||
let trace = run_simulation(config);
|
let trace = run_simulation(config);
|
||||||
let metrics = analyze(&trace);
|
let metrics = analyze(&trace);
|
||||||
|
|
||||||
// With 10% loss and the deterministic LCG, SWIM's probe cycle is disrupted
|
// With 10% loss, SWIM's probe cycle is disrupted enough to cause
|
||||||
// enough to cause false deaths. The test verifies:
|
// false suspicions. Dead reprobe allows recovery. We verify:
|
||||||
// 1. The simulation completes without panic (implicit — we got here)
|
// 1. The simulation completes without panic (implicit — we got here)
|
||||||
// 2. At least partial membership is maintained (some nodes still know about others)
|
// 2. At least partial membership is maintained (some nodes still know about others)
|
||||||
let result = check_membership_accuracy(&metrics, 0.15);
|
let result = check_membership_accuracy(&metrics, 0.15);
|
||||||
|
|
|
||||||
|
|
@ -72,7 +72,6 @@ fn registry_name_converges_across_cluster() {
|
||||||
// ────────────────────────────────────────────────────────────────────────────
|
// ────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
#[test]
|
#[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() {
|
fn split_brain_naming_converges_after_partition_heals() {
|
||||||
// Given: 6-node cluster, partition {0,1,2} vs {3,4,5} at round 10
|
// 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
|
// 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]
|
#[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() {
|
fn tombstone_propagates_when_name_owner_dies() {
|
||||||
// Given: 5-node cluster, node 0 registers "svc" at round 5, killed at round 15
|
// Given: 5-node cluster, node 0 registers "svc" at round 5, killed at round 15
|
||||||
let config = DistributionSimConfig {
|
let config = DistributionSimConfig {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue