#![cfg_attr(not(test), deny(clippy::disallowed_methods))]
#![cfg_attr(
not(test),
deny(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::unreachable,
clippy::todo,
clippy::unimplemented,
clippy::indexing_slicing,
clippy::string_slice,
clippy::arithmetic_side_effects,
)
)]
use std::collections::BTreeMap;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use super::counter::CounterShards;
use super::{Incarnation, NodeId};
use crate::time::MonotonicInstant;
pub const TOMBSTONE_TIMEOUT_MULTIPLE: u32 = 10;
pub const PRUNE_MEMORY_WINDOW_MULTIPLE: u32 = 2;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum MemberStatus {
Alive,
Left,
}
impl MemberStatus {
const fn rank(self) -> u8 {
match self {
Self::Alive => 0,
Self::Left => 1,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct MemberRecord {
#[serde(default)]
pub addr: String,
#[serde(default)]
pub incarnation: Incarnation,
pub status: MemberStatus,
}
impl MemberRecord {
pub fn alive(addr: impl Into<String>, incarnation: Incarnation) -> Self {
Self {
addr: addr.into(),
incarnation,
status: MemberStatus::Alive,
}
}
pub fn left(addr: impl Into<String>, incarnation: Incarnation) -> Self {
Self {
addr: addr.into(),
incarnation,
status: MemberStatus::Left,
}
}
const fn merge_key(&self) -> (Incarnation, u8, &str) {
(self.incarnation, self.status.rank(), self.addr.as_str())
}
pub fn merge(&mut self, other: &Self) {
if other.merge_key() > self.merge_key() {
self.clone_from(other);
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ClusterState {
#[serde(default)]
pub members: BTreeMap<NodeId, MemberRecord>,
#[serde(default)]
pub counters: BTreeMap<String, CounterShards>,
}
impl ClusterState {
pub fn merge(&mut self, other: &Self, overlay: &LivenessOverlay, now: MonotonicInstant) {
for (id, theirs) in &other.members {
if let Some(ours) = self.members.get_mut(id) {
ours.merge(theirs);
} else if !overlay.refuses_readoption(id, theirs, now) {
self.members.insert(id.clone(), theirs.clone());
}
}
for (name, theirs) in &other.counters {
self.counters.entry(name.clone()).or_default().merge(theirs);
}
}
pub fn refute(&mut self, me: &str, own: &MemberRecord) -> Option<Incarnation> {
let observed = self.members.get_mut(me)?;
if observed == own {
return None;
}
if observed.incarnation < own.incarnation {
return None;
}
let bumped = observed.incarnation.saturating_add(1);
observed.incarnation = bumped;
observed.status = MemberStatus::Alive;
observed.addr.clone_from(&own.addr);
Some(bumped)
}
pub fn convert_down_members(
&mut self,
me: &str,
overlay: &mut LivenessOverlay,
now: MonotonicInstant,
) -> Vec<NodeId> {
let window = overlay.tombstone_timeout();
let mut converted = Vec::new();
for (id, record) in &mut self.members {
if id == me || record.status == MemberStatus::Left {
continue;
}
let quiet_since = overlay.quiet_since(id, now);
if now.saturating_duration_since(quiet_since) > window {
record.status = MemberStatus::Left;
converted.push(id.clone());
}
}
converted
}
pub fn observe_tombstones(&self, overlay: &mut LivenessOverlay, now: MonotonicInstant) {
for (id, record) in &self.members {
match record.status {
MemberStatus::Left => overlay.mark_tombstoned(id, now),
MemberStatus::Alive => overlay.clear_tombstone(id),
}
}
}
pub fn prune_tombstones(
&mut self,
overlay: &mut LivenessOverlay,
now: MonotonicInstant,
) -> Vec<NodeId> {
let window = overlay.tombstone_timeout();
let expired: Vec<(NodeId, Incarnation)> = self
.members
.iter()
.filter_map(|(id, record)| {
let lapsed = record.status == MemberStatus::Left
&& overlay
.tombstoned_at(id)
.is_some_and(|seen| now.saturating_duration_since(seen) > window);
lapsed.then(|| (id.clone(), record.incarnation))
})
.collect();
for (id, incarnation) in &expired {
self.members.remove(id);
overlay.forget(id);
overlay.remember_pruned(id, *incarnation, now);
}
overlay.expire_pruned(now);
expired.into_iter().map(|(id, _)| id).collect()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Liveness {
Alive,
Suspect,
Down,
}
impl Liveness {
pub const fn in_view(self) -> bool {
matches!(self, Self::Alive | Self::Suspect)
}
}
#[derive(Debug)]
pub struct LivenessOverlay {
push_interval: Duration,
suspicion_timeout: Duration,
last_seen: BTreeMap<NodeId, MonotonicInstant>,
tombstoned_at: BTreeMap<NodeId, MonotonicInstant>,
unheard_since: BTreeMap<NodeId, MonotonicInstant>,
recently_pruned: BTreeMap<NodeId, (Incarnation, MonotonicInstant)>,
}
impl LivenessOverlay {
pub const fn new(push_interval: Duration, suspicion_timeout: Duration) -> Self {
Self {
push_interval,
suspicion_timeout,
last_seen: BTreeMap::new(),
tombstoned_at: BTreeMap::new(),
unheard_since: BTreeMap::new(),
recently_pruned: BTreeMap::new(),
}
}
pub const fn suspicion_timeout(&self) -> Duration {
self.suspicion_timeout
}
pub fn record_receipt(&mut self, node: &str, at: MonotonicInstant) {
self.last_seen.insert(node.to_owned(), at);
}
pub fn forget(&mut self, node: &str) {
self.last_seen.remove(node);
self.tombstoned_at.remove(node);
self.unheard_since.remove(node);
}
fn quiet_since(&mut self, node: &str, now: MonotonicInstant) -> MonotonicInstant {
if let Some(seen) = self.last_receipt(node) {
self.unheard_since.remove(node);
return seen;
}
*self.unheard_since.entry(node.to_owned()).or_insert(now)
}
fn remember_pruned(&mut self, node: &str, incarnation: Incarnation, at: MonotonicInstant) {
let entry = self
.recently_pruned
.entry(node.to_owned())
.or_insert((incarnation, at));
if incarnation >= entry.0 {
*entry = (incarnation, at);
}
}
fn expire_pruned(&mut self, now: MonotonicInstant) {
let window = self.prune_memory_window();
self.recently_pruned
.retain(|_, (_, at)| now.saturating_duration_since(*at) <= window);
}
fn refuses_readoption(&self, node: &str, record: &MemberRecord, now: MonotonicInstant) -> bool {
let window = self.prune_memory_window();
self.recently_pruned
.get(node)
.is_some_and(|(collected, at)| {
record.incarnation <= *collected && now.saturating_duration_since(*at) <= window
})
}
fn mark_tombstoned(&mut self, node: &str, at: MonotonicInstant) {
self.tombstoned_at.entry(node.to_owned()).or_insert(at);
}
fn clear_tombstone(&mut self, node: &str) {
self.tombstoned_at.remove(node);
}
fn tombstoned_at(&self, node: &str) -> Option<MonotonicInstant> {
self.tombstoned_at.get(node).copied()
}
fn last_receipt(&self, node: &str) -> Option<MonotonicInstant> {
self.last_seen.get(node).copied()
}
const fn tombstone_timeout(&self) -> Duration {
self.suspicion_timeout()
.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE)
}
const fn prune_memory_window(&self) -> Duration {
self.tombstone_timeout()
.saturating_mul(PRUNE_MEMORY_WINDOW_MULTIPLE)
}
pub fn liveness(&self, node: &str, now: MonotonicInstant) -> Liveness {
let Some(last) = self.last_receipt(node) else {
return Liveness::Down;
};
let silence = now.saturating_duration_since(last);
if silence > self.suspicion_timeout {
Liveness::Down
} else if silence > self.push_interval.saturating_mul(2) {
Liveness::Suspect
} else {
Liveness::Alive
}
}
}
#[cfg(test)]
mod tests {
use super::{
ClusterState, Liveness, LivenessOverlay, MemberRecord, MemberStatus,
PRUNE_MEMORY_WINDOW_MULTIPLE, TOMBSTONE_TIMEOUT_MULTIPLE,
};
use crate::time::MonotonicInstant;
use std::time::Duration;
const PUSH: Duration = Duration::from_millis(500);
const SUSPICION: Duration = Duration::from_millis(2_500);
fn at(millis: u64) -> MonotonicInstant {
MonotonicInstant::from_origin_elapsed(Duration::from_millis(millis))
}
fn after_collecting_node_b(pruned_at: u64) -> (ClusterState, LivenessOverlay) {
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
let mut state = ClusterState::default();
state
.members
.insert("node-b".to_owned(), MemberRecord::left("127.0.0.1:7002", 9));
state.observe_tombstones(&mut overlay, at(0));
assert_eq!(
state.prune_tombstones(&mut overlay, at(pruned_at)),
vec!["node-b".to_owned()],
"sanity: the tombstone must prune, or nothing below is tested"
);
(state, overlay)
}
#[test]
fn member_merge_prefers_higher_incarnation() {
let mut older = MemberRecord::alive("127.0.0.1:7001", 3);
let newer = MemberRecord::alive("127.0.0.1:7002", 4);
older.merge(&newer);
assert_eq!(
older.incarnation, 4,
"a higher incarnation must win outright; observed {older:?}"
);
assert_eq!(
older.addr, "127.0.0.1:7002",
"the winning record's address must win with it; observed {older:?}"
);
let stale = MemberRecord::left("127.0.0.1:7000", 2);
older.merge(&stale);
assert_eq!(
older.incarnation, 4,
"a lower incarnation must never win; observed {older:?}"
);
assert_eq!(
older.status,
MemberStatus::Alive,
"a stale Left must not evict a live member; observed {older:?}"
);
}
#[test]
fn leave_supersedes_alive_at_same_incarnation() {
let mut alive = MemberRecord::alive("127.0.0.1:7001", 7);
let left = MemberRecord::left("127.0.0.1:7001", 7);
alive.merge(&left);
assert_eq!(
alive.status,
MemberStatus::Left,
"at an equal incarnation Left must beat Alive; observed {alive:?}"
);
let mut left_first = MemberRecord::left("127.0.0.1:7001", 7);
left_first.merge(&MemberRecord::alive("127.0.0.1:7001", 7));
assert_eq!(
left_first.status,
MemberStatus::Left,
"…in either merge direction; observed {left_first:?}"
);
let mut lower_addr = MemberRecord::alive("127.0.0.1:7001", 7);
lower_addr.merge(&MemberRecord::alive("127.0.0.1:7009", 7));
assert_eq!(
lower_addr.addr, "127.0.0.1:7009",
"the lexicographically greater addr must win the tie; observed {lower_addr:?}"
);
}
#[test]
fn refutation_bumps_incarnation_over_stale_leave() {
let own = MemberRecord::alive("127.0.0.1:7001", 4);
let mut buried = ClusterState::default();
buried
.members
.insert("node-a".to_owned(), MemberRecord::left("127.0.0.1:7001", 4));
assert_eq!(
buried.refute("node-a", &own),
Some(5),
"a Left about ourselves must be refuted one incarnation higher"
);
let record = buried
.members
.get("node-a")
.expect("the refuting node must stay in its own document");
assert_eq!(
(record.status, record.incarnation),
(MemberStatus::Alive, 5),
"refutation must restore Alive at the bumped incarnation; observed {record:?}"
);
let mut stale_alive = ClusterState::default();
stale_alive.members.insert(
"node-a".to_owned(),
MemberRecord::alive("127.0.0.1:9999", 7),
);
assert_eq!(
stale_alive.refute("node-a", &own),
Some(8),
"a stale Alive about ourselves at a higher incarnation must also be refuted"
);
let mut echoed = ClusterState::default();
echoed.members.insert("node-a".to_owned(), own.clone());
assert_eq!(
echoed.refute("node-a", &own),
None,
"an exact echo of our own record must not be refuted"
);
let mut older = ClusterState::default();
older
.members
.insert("node-a".to_owned(), MemberRecord::left("127.0.0.1:7001", 2));
assert_eq!(
older.refute("node-a", &own),
None,
"a record older than ours loses the merge already; no bump needed"
);
}
#[test]
fn missed_pushes_transition_member_to_down() {
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
overlay.record_receipt("node-b", at(0));
assert_eq!(
overlay.liveness("node-b", at(900)),
Liveness::Alive,
"within two push intervals a member is Alive"
);
assert_eq!(
overlay.liveness("node-b", at(1_100)),
Liveness::Suspect,
"past two push intervals a member is Suspect — a warning, not an eviction"
);
assert!(
overlay.liveness("node-b", at(1_100)).in_view(),
"a Suspect member must still appear in the view"
);
assert_eq!(
overlay.liveness("node-b", at(2_600)),
Liveness::Down,
"past the suspicion timeout a member is Down"
);
assert!(
!overlay.liveness("node-b", at(2_600)).in_view(),
"a Down member must leave the view"
);
assert_eq!(
overlay.liveness("never-seen", at(0)),
Liveness::Down,
"a peer never heard from is Down, not Alive"
);
}
#[test]
fn push_receipt_resets_suspicion() {
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
overlay.record_receipt("node-b", at(0));
assert_eq!(
overlay.liveness("node-b", at(1_500)),
Liveness::Suspect,
"sanity: the member must be Suspect before the reset, or this test \
proves nothing"
);
overlay.record_receipt("node-b", at(1_500));
assert_eq!(
overlay.liveness("node-b", at(1_900)),
Liveness::Alive,
"a fresh receipt must reset the suspicion clock"
);
assert_eq!(
overlay.liveness("node-b", at(2_700)),
Liveness::Suspect,
"…and silence must then be measured from the NEW receipt"
);
assert_eq!(
overlay.liveness("node-b", at(4_100)),
Liveness::Down,
"…including the Down threshold"
);
}
#[test]
fn left_tombstones_prune_after_ten_suspicion_timeouts() {
let window = SUSPICION.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE);
let window_ms = u64::try_from(window.as_millis()).expect("the window fits in a u64 of ms");
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
overlay.record_receipt("node-a", at(0));
overlay.record_receipt("node-b", at(0));
let mut state = ClusterState::default();
state.members.insert(
"node-a".to_owned(),
MemberRecord::alive("127.0.0.1:7001", 4),
);
state
.members
.insert("node-b".to_owned(), MemberRecord::left("127.0.0.1:7002", 9));
state.observe_tombstones(&mut overlay, at(0));
assert_eq!(
state.prune_tombstones(&mut overlay, at(window_ms)),
Vec::<String>::new(),
"a tombstone must survive its whole window — pruning it early would \
let a straggling push resurrect the departed member; observed {state:?}"
);
overlay.record_receipt("node-b", at(window_ms));
assert_eq!(
state.prune_tombstones(&mut overlay, at(window_ms.saturating_add(1))),
vec!["node-b".to_owned()],
"past ten suspicion timeouts the tombstone must be pruned, and the \
pruned id must be REPORTED so the caller can forget the sender's \
replay watermark too; observed {state:?}"
);
assert!(
!state.members.contains_key("node-b"),
"the pruned tombstone must leave the document; observed {state:?}"
);
assert_eq!(
overlay.liveness("node-b", at(window_ms.saturating_add(1))),
Liveness::Down,
"pruning must forget the overlay row too, so a pruned member reads \
as never-seen rather than as a one-millisecond-old receipt"
);
assert!(
state.members.contains_key("node-a"),
"an Alive member must never be pruned, however long it has been \
silent — silence is the overlay's business, not the document's; \
observed {state:?}"
);
}
#[test]
fn reintroduced_tombstone_expires_again() {
let window = SUSPICION.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE);
let window_ms = u64::try_from(window.as_millis()).expect("the window fits in a u64 of ms");
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
let mut state = ClusterState::default();
state
.members
.insert("node-b".to_owned(), MemberRecord::left("127.0.0.1:7002", 9));
state.observe_tombstones(&mut overlay, at(0));
assert_eq!(
state.prune_tombstones(&mut overlay, at(window_ms.saturating_add(1))),
vec!["node-b".to_owned()],
"sanity: the first tombstone must prune, or this test proves nothing"
);
let reintroduced =
at(window_ms.saturating_mul(u64::from(PRUNE_MEMORY_WINDOW_MULTIPLE).saturating_add(2)));
let mut lagging = ClusterState::default();
lagging
.members
.insert("node-b".to_owned(), MemberRecord::left("127.0.0.1:7002", 9));
state.merge(&lagging, &overlay, reintroduced);
assert!(
state.members.contains_key("node-b"),
"sanity: past the memory window the record must be re-learned, or \
the expiry assertions below prove nothing; observed {state:?}"
);
state.observe_tombstones(&mut overlay, reintroduced);
overlay.record_receipt("node-c", reintroduced);
assert_eq!(
state.prune_tombstones(&mut overlay, reintroduced.saturating_add(window)),
Vec::<String>::new(),
"a freshly re-learned tombstone must survive a full window; observed {state:?}"
);
assert_eq!(
state.prune_tombstones(
&mut overlay,
reintroduced
.saturating_add(window)
.saturating_add(Duration::from_millis(1)),
),
vec!["node-b".to_owned()],
"a REINTRODUCED tombstone must expire one window after it was \
re-observed — otherwise repeated departures grow the document \
without bound; observed {state:?}"
);
assert!(
!state.members.contains_key("node-b"),
"the re-expired tombstone must leave the document; observed {state:?}"
);
}
#[test]
fn pruned_tombstones_are_not_re_adopted_during_the_memory_window() {
let window = SUSPICION.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE);
let window_ms = u64::try_from(window.as_millis()).expect("the window fits in a u64 of ms");
let pruned_at = window_ms.saturating_add(1);
let mut lagging = ClusterState::default();
lagging
.members
.insert("node-b".to_owned(), MemberRecord::left("127.0.0.1:7002", 9));
let (mut state, overlay) = after_collecting_node_b(pruned_at);
state.merge(&lagging, &overlay, at(pruned_at.saturating_add(1)));
assert!(
!state.members.contains_key("node-b"),
"a Left record this node already collected must not come back from \
a peer that prunes later — that exchange is the recycling loop \
that keeps departed ids in the document forever; observed {state:?}"
);
let inside = pruned_at
.saturating_add(window_ms.saturating_mul(u64::from(PRUNE_MEMORY_WINDOW_MULTIPLE)));
state.merge(&lagging, &overlay, at(inside));
assert!(
!state.members.contains_key("node-b"),
"the refusal must hold for the whole memory window — the peer keeps \
pushing every interval; observed {state:?}"
);
state.merge(&lagging, &overlay, at(inside.saturating_add(1)));
assert!(
state.members.contains_key("node-b"),
"the memory is bounded: past its window a tombstone a peer still \
holds must merge normally, or this is a permanent local override \
of the join rather than garbage collection; observed {state:?}"
);
let (mut rejoining, overlay) = after_collecting_node_b(pruned_at);
let mut returned = ClusterState::default();
returned.members.insert(
"node-b".to_owned(),
MemberRecord::alive("127.0.0.1:7002", 11),
);
rejoining.merge(&returned, &overlay, at(pruned_at.saturating_add(1)));
assert_eq!(
rejoining.members.get("node-b").map(|record| record.status),
Some(MemberStatus::Alive),
"a node that came back must never be held out by what this node \
collected about its previous life; observed {rejoining:?}"
);
let (mut departing_again, overlay) = after_collecting_node_b(pruned_at);
let mut left_higher = ClusterState::default();
left_higher.members.insert(
"node-b".to_owned(),
MemberRecord::left("127.0.0.1:7002", 12),
);
departing_again.merge(&left_higher, &overlay, at(pruned_at.saturating_add(1)));
assert_eq!(
departing_again
.members
.get("node-b")
.map(|record| record.incarnation),
Some(12),
"a LATER departure outranks the collected one and must propagate \
normally; refusing it would strand a real leave; observed \
{departing_again:?}"
);
let replayed_alive = {
let mut document = ClusterState::default();
document.members.insert(
"node-b".to_owned(),
MemberRecord::alive("127.0.0.1:7002", 9),
);
document
};
let (mut replayed, overlay) = after_collecting_node_b(pruned_at);
replayed.merge(&replayed_alive, &overlay, at(pruned_at.saturating_add(1)));
assert!(
!replayed.members.contains_key("node-b"),
"a replayed Alive record at the collected incarnation must be \
refused too: re-admitting it puts a member back in the document \
that no push refreshes and no prune collects, so replaying one \
captured frame per departed id grows the document toward the \
frame cap without knowing the secret; observed {replayed:?}"
);
replayed.merge(&replayed_alive, &overlay, at(inside.saturating_add(1)));
assert!(
replayed.members.contains_key("node-b"),
"…and the guard is still only a memory: past its window the join \
is a join again, and what it re-learns leaves on its own (see \
`members_down_for_a_tombstone_window_are_recorded_as_left`); \
observed {replayed:?}"
);
}
#[test]
fn members_down_for_a_tombstone_window_are_recorded_as_left() {
let window = SUSPICION.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE);
let window_ms = u64::try_from(window.as_millis()).expect("the window fits in a u64 of ms");
let out_of_view = 3_000;
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
overlay.record_receipt("node-b", at(0));
let mut state = ClusterState::default();
state.members.insert(
"node-a".to_owned(),
MemberRecord::alive("127.0.0.1:7001", 4),
);
state.members.insert(
"node-b".to_owned(),
MemberRecord::alive("127.0.0.1:7002", 9),
);
state.members.insert(
"node-c".to_owned(),
MemberRecord::alive("127.0.0.1:7003", 2),
);
assert_eq!(
overlay.liveness("node-b", at(out_of_view)),
Liveness::Down,
"sanity: the member must already be Down here, or the assertions \
below cannot tell the view's timeout from the document's window"
);
assert_eq!(
state.convert_down_members("node-a", &mut overlay, at(out_of_view)),
Vec::<String>::new(),
"a member that has only just gone Down is out of the VIEW, not out \
of the document — converting at the suspicion timeout would \
tombstone a peer that missed three pushes; observed {state:?}"
);
assert_eq!(
state.convert_down_members("node-a", &mut overlay, at(window_ms)),
Vec::<String>::new(),
"…and it keeps its record for the whole window; observed {state:?}"
);
let converted_at = window_ms.saturating_add(1);
assert_eq!(
state.convert_down_members("node-a", &mut overlay, at(converted_at)),
vec!["node-b".to_owned()],
"past a whole tombstone window of silence the record must become a \
tombstone, or it never leaves the document at all; observed {state:?}"
);
assert_eq!(
state.members.get("node-b"),
Some(&MemberRecord::left("127.0.0.1:7002", 9)),
"the conversion must keep the record's incarnation and address: a \
live node wrongly converted refutes at incarnation + 1, which is \
the argument that undoes this everywhere; observed {state:?}"
);
assert_eq!(
state.members.get("node-a").map(|record| record.status),
Some(MemberStatus::Alive),
"this node's OWN record must never be converted — a node records no \
receipts from itself, so its own record is permanently silent and \
would bury itself on the first round; observed {state:?}"
);
assert_eq!(
state.members.get("node-c").map(|record| record.status),
Some(MemberStatus::Alive),
"a member with no receipt at all must still get its whole window; \
observed {state:?}"
);
assert_eq!(
state.convert_down_members(
"node-a",
&mut overlay,
at(out_of_view.saturating_add(window_ms).saturating_add(1)),
),
vec!["node-c".to_owned()],
"…and must then convert too: a record nothing has ever refreshed is \
precisely the record no prune can reach; observed {state:?}"
);
state.observe_tombstones(&mut overlay, at(converted_at));
assert_eq!(
state.prune_tombstones(&mut overlay, at(converted_at.saturating_add(window_ms))),
Vec::<String>::new(),
"sanity: a converted record ages like any tombstone, from the \
observation; observed {state:?}"
);
assert_eq!(
state.prune_tombstones(
&mut overlay,
at(converted_at
.saturating_add(window_ms)
.saturating_add(window_ms)),
),
vec!["node-b".to_owned(), "node-c".to_owned()],
"a converted record must then PRUNE like any other tombstone — \
conversion without pruning only renames the leak; observed {state:?}"
);
assert_eq!(
state.members.keys().collect::<Vec<_>>(),
vec!["node-a"],
"only this node's own record may outlive a full window of silence; \
observed {state:?}"
);
}
#[test]
fn a_member_that_refreshes_inside_the_window_is_never_converted() {
let window = SUSPICION.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE);
let window_ms = u64::try_from(window.as_millis()).expect("the window fits in a u64 of ms");
let refreshed_at = 3_500;
let mut overlay = LivenessOverlay::new(PUSH, SUSPICION);
overlay.record_receipt("node-b", at(0));
let mut state = ClusterState::default();
state.members.insert(
"node-b".to_owned(),
MemberRecord::alive("127.0.0.1:7002", 9),
);
assert_eq!(
state.convert_down_members("node-a", &mut overlay, at(3_000)),
Vec::<String>::new(),
"sanity: the member must be silent here, or the refresh below \
proves nothing; observed {state:?}"
);
overlay.record_receipt("node-b", at(refreshed_at));
assert_eq!(
state.convert_down_members("node-a", &mut overlay, at(window_ms.saturating_add(1))),
Vec::<String>::new(),
"a member that answered inside the window must not be converted on \
the deadline its earlier silence set: the clock is the silence \
since this node last heard from it, not the uptime of the cluster; \
observed {state:?}"
);
assert_eq!(
state.members.get("node-b").map(|record| record.status),
Some(MemberStatus::Alive),
"…and its record must still be Alive; observed {state:?}"
);
assert_eq!(
state.convert_down_members(
"node-a",
&mut overlay,
at(refreshed_at.saturating_add(window_ms).saturating_add(1)),
),
vec!["node-b".to_owned()],
"the receipt must MOVE the deadline, not remove it — a member that \
answers once and then vanishes still has to leave the document; \
observed {state:?}"
);
}
}