feat: distribution simulation tests #34
2 changed files with 19 additions and 15 deletions
|
|
@ -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<NodeAction> {
|
||||
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<NodeAction> {
|
||||
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<NodeAction> {
|
||||
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<u8> {
|
||||
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<RegistryEntry>) {
|
||||
if !entries.is_empty() {
|
||||
self.registry.merge_batch(entries, self.cluster_size());
|
||||
}
|
||||
membership_bytes
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
);
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue