use std::collections::BTreeSet;
use std::net::SocketAddr;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::Duration;
use tokio_util::sync::CancellationToken;
use super::membership::{
ClusterState, MemberRecord, PRUNE_MEMORY_WINDOW_MULTIPLE, TOMBSTONE_TIMEOUT_MULTIPLE,
};
use super::node::{ClusterNode, ClusterRuntimeConfig};
use super::transport::{IncomingFrames, LoopbackRouter, PeerTransport};
use super::wire::{self, ClusterMessage, RejectReason};
use super::{ClusterHandle, ClusterMemberStatus, ClusterMetrics, REJECT_REASONS};
use crate::entropy::{Entropy, SeededEntropy};
use crate::time::TickingClock;
const CLUSTER: &str = "autumn";
const SECRET: &[u8] = b"a-shared-cluster-secret-value-32";
const OTHER_SECRET: &[u8] = b"a-different-cluster-secret-value";
const COUNTER: &str = "boids_sighted";
const PUSH: Duration = Duration::from_millis(500);
const SUSPICION: Duration = Duration::from_millis(2_500);
struct TestNode {
id: String,
addr: String,
handle: ClusterHandle,
token: CancellationToken,
}
fn test_clock() -> TickingClock {
TickingClock::starting_at(
chrono::DateTime::<chrono::Utc>::from_timestamp(1_765_430_000, 0).unwrap_or_default(),
)
}
async fn advance_time(clock: &TickingClock, dur: Duration) {
clock.advance(dur);
tokio::time::sleep(dur).await;
}
async fn settle(clock: &TickingClock, rounds: u32) {
for _ in 0..rounds {
advance_time(clock, PUSH).await;
}
}
fn window_rounds() -> u32 {
u32::try_from(
SUSPICION
.saturating_mul(TOMBSTONE_TIMEOUT_MULTIPLE)
.as_millis()
.saturating_div(PUSH.as_millis().max(1))
.saturating_add(2),
)
.unwrap_or(u32::MAX)
}
fn start_node(
router: &LoopbackRouter,
clock: &TickingClock,
secret: &[u8],
node_id: &str,
seed: u64,
seed_peers: Vec<String>,
) -> TestNode {
start_node_on(
router.endpoint() as Arc<dyn PeerTransport>,
clock,
secret,
node_id,
seed,
seed_peers,
)
}
fn start_node_on(
transport: Arc<dyn PeerTransport>,
clock: &TickingClock,
secret: &[u8],
node_id: &str,
seed: u64,
seed_peers: Vec<String>,
) -> TestNode {
let addr = transport.local_addr().to_string();
let token = CancellationToken::new();
let handle = ClusterNode::start(
ClusterRuntimeConfig {
cluster_name: CLUSTER.to_owned(),
secret: secret.to_vec(),
node_id: Some(node_id.to_owned()),
advertise_addr: None,
seed_peers,
push_interval: PUSH,
suspicion_timeout: SUSPICION,
},
Arc::new(SeededEntropy::new(seed)),
Arc::new(clock.clone()),
token.clone(),
transport,
)
.expect("a cluster node must start on a loopback transport");
TestNode {
id: node_id.to_owned(),
addr,
handle,
token,
}
}
fn record_of(node: &TestNode, id: &str) -> Option<MemberRecord> {
let state = node.handle.inner.lock_state();
let record = state.members.get(id).cloned();
drop(state);
record
}
fn signed_self_push(id: &str, addr: &str, incarnation: u64, seq: u64) -> Vec<u8> {
let mut document = ClusterState::default();
document.members.insert(
id.to_owned(),
MemberRecord::alive(addr.to_owned(), incarnation),
);
wire::sign_envelope(
SECRET,
CLUSTER,
id,
incarnation,
seq,
&ClusterMessage::StatePush { state: document },
)
.as_ref()
.and_then(wire::encode_frame)
.unwrap_or_default()
}
fn member_ids(handle: &ClusterHandle) -> Vec<String> {
let mut ids: Vec<String> = handle.members().into_iter().map(|m| m.id).collect();
ids.sort();
ids
}
fn assert_converged(a: &TestNode, b: &TestNode, expected: &[&str], why: &str) {
let expected: Vec<String> = expected.iter().map(|id| (*id).to_owned()).collect();
assert_eq!(
member_ids(&a.handle),
expected,
"{why}; {} | {}",
view_of(a),
view_of(b)
);
assert_eq!(
member_ids(&b.handle),
expected,
"…and both nodes must converge on the SAME view ({why}); {} | {}",
view_of(a),
view_of(b)
);
}
fn pushes_sent(node: &TestNode) -> u64 {
node.handle
.inner
.metrics
.pushes_sent
.load(std::sync::atomic::Ordering::Relaxed)
}
fn pushes_unsendable(node: &TestNode) -> u64 {
node.handle
.inner
.metrics
.pushes_unsendable
.load(std::sync::atomic::Ordering::Relaxed)
}
fn grow_document_past_the_frame_cap(node: &TestNode, cells: u64) {
let mut filler = super::counter::CounterShards::default();
for boot in 0..cells {
filler.increment_cell(
&super::counter::cell_key("node-a-long-enough-id-to-model-a-real-one", boot),
1,
);
}
let mut state = node.handle.inner.lock_state();
state
.counters
.entry(COUNTER.to_owned())
.or_default()
.merge(&filler);
drop(state);
}
fn collect_families(node: &TestNode) -> Vec<crate::actuator::MetricFamily> {
use crate::actuator::MetricsSource as _;
super::ClusterMetricsSource {
handle: node.handle.clone(),
}
.collect()
}
fn family_value(families: &[crate::actuator::MetricFamily], name: &str) -> Option<f64> {
families
.iter()
.find(|family| family.name == name)
.and_then(|family| family.samples.first())
.map(|sample| sample.value)
}
fn rejected_series(families: &[crate::actuator::MetricFamily]) -> Vec<(String, f64)> {
families
.iter()
.find(|family| family.name == "autumn_cluster_frames_rejected_total")
.map(|family| {
family
.samples
.iter()
.map(|sample| {
let reason = sample
.labels
.first()
.map_or_else(String::new, |(_, reason)| reason.clone());
(reason, sample.value)
})
.collect()
})
.unwrap_or_default()
}
fn zeroed_series() -> Vec<(String, f64)> {
REJECT_REASONS
.iter()
.map(|reason| (reason.label().to_owned(), 0.0))
.collect()
}
struct RetainSpy {
inner: Arc<dyn PeerTransport>,
retained: std::sync::Mutex<Vec<BTreeSet<String>>>,
}
impl RetainSpy {
fn new(inner: Arc<dyn PeerTransport>) -> Self {
Self {
inner,
retained: std::sync::Mutex::new(Vec::new()),
}
}
fn last_retained(&self) -> BTreeSet<String> {
self.retained
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.last()
.cloned()
.unwrap_or_default()
}
}
impl PeerTransport for RetainSpy {
fn send(&self, to: &str, frame: Vec<u8>) {
self.inner.send(to, frame);
}
fn send_farewell(&self, to: &str, frame: Vec<u8>) -> bool {
self.inner.send_farewell(to, frame)
}
fn take_incoming(&self) -> Option<IncomingFrames> {
self.inner.take_incoming()
}
fn local_addr(&self) -> SocketAddr {
self.inner.local_addr()
}
fn start(&self, shutdown: &CancellationToken, entropy: &Arc<dyn Entropy>) {
self.inner.start(shutdown, entropy);
}
fn pending_frames(&self) -> usize {
self.inner.pending_frames()
}
fn dropped_frames(&self) -> u64 {
self.inner.dropped_frames()
}
fn framing_rejections(&self) -> u64 {
self.inner.framing_rejections()
}
fn note_unauthenticated_frame(&self, from: &str) {
self.inner.note_unauthenticated_frame(from);
}
fn retain_peers(&self, live: &BTreeSet<String>) {
self.retained
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(live.clone());
self.inner.retain_peers(live);
}
}
struct RejectionSpy {
inner: Arc<dyn PeerTransport>,
reported: std::sync::Mutex<Vec<String>>,
}
impl RejectionSpy {
fn new(inner: Arc<dyn PeerTransport>) -> Self {
Self {
inner,
reported: std::sync::Mutex::new(Vec::new()),
}
}
fn reported(&self) -> Vec<String> {
self.reported
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
impl PeerTransport for RejectionSpy {
fn send(&self, to: &str, frame: Vec<u8>) {
self.inner.send(to, frame);
}
fn send_farewell(&self, to: &str, frame: Vec<u8>) -> bool {
self.inner.send_farewell(to, frame)
}
fn take_incoming(&self) -> Option<IncomingFrames> {
self.inner.take_incoming()
}
fn local_addr(&self) -> SocketAddr {
self.inner.local_addr()
}
fn start(&self, shutdown: &CancellationToken, entropy: &Arc<dyn Entropy>) {
self.inner.start(shutdown, entropy);
}
fn note_unauthenticated_frame(&self, from: &str) {
self.reported
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(from.to_owned());
self.inner.note_unauthenticated_frame(from);
}
}
struct StalledPeerTransport {
inner: Arc<dyn PeerTransport>,
stalled: AtomicBool,
deaf: AtomicBool,
refused_farewells: AtomicU64,
farewells: std::sync::Mutex<Vec<String>>,
}
impl StalledPeerTransport {
fn new(inner: Arc<dyn PeerTransport>) -> Self {
Self {
inner,
stalled: AtomicBool::new(false),
deaf: AtomicBool::new(false),
refused_farewells: AtomicU64::new(0),
farewells: std::sync::Mutex::new(Vec::new()),
}
}
fn stall(&self) {
self.stalled.store(true, Ordering::Relaxed);
}
fn go_deaf(&self) {
self.deaf.store(true, Ordering::Relaxed);
}
fn farewell_targets(&self) -> Vec<String> {
self.farewells
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
impl PeerTransport for StalledPeerTransport {
fn send(&self, to: &str, frame: Vec<u8>) {
if self.stalled.load(Ordering::Relaxed) {
return;
}
self.inner.send(to, frame);
}
fn send_farewell(&self, to: &str, frame: Vec<u8>) -> bool {
self.farewells
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(to.to_owned());
if self.deaf.load(Ordering::Relaxed) {
self.refused_farewells.fetch_add(1, Ordering::Relaxed);
return false;
}
self.inner.send_farewell(to, frame)
}
fn take_incoming(&self) -> Option<IncomingFrames> {
self.inner.take_incoming()
}
fn local_addr(&self) -> SocketAddr {
self.inner.local_addr()
}
fn start(&self, shutdown: &CancellationToken, entropy: &Arc<dyn Entropy>) {
self.inner.start(shutdown, entropy);
}
fn dropped_frames(&self) -> u64 {
self.refused_farewells.load(Ordering::Relaxed)
}
}
fn view_of(node: &TestNode) -> String {
format!(
"{} sees {:?} (counter {} = {})",
node.id,
node.handle.members(),
COUNTER,
node.handle.counter(COUNTER).get()
)
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_two_nodes_converge_to_two_member_view() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"seeding B at A's address must give A a two-member view",
);
assert!(
a.handle
.members()
.iter()
.all(|m| m.status == ClusterMemberStatus::Alive),
"a freshly converged view must be entirely Alive; {}",
view_of(&a)
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_increment_on_a_reads_on_b() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
a.handle.counter(COUNTER).increment();
assert_eq!(
a.handle.counter(COUNTER).get(),
1,
"a local increment must be visible immediately on the writer; {}",
view_of(&a)
);
settle(&clock, 4).await;
assert_eq!(
b.handle.counter(COUNTER).get(),
1,
"an increment on A must be readable on B after a few push intervals; {} | {}",
view_of(&a),
view_of(&b)
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_concurrent_increments_converge() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
for _ in 0..3 {
a.handle.counter(COUNTER).increment();
b.handle.counter(COUNTER).increment();
advance_time(&clock, PUSH).await;
}
a.handle.counter(COUNTER).increment_by(4);
settle(&clock, 6).await;
assert_eq!(
a.handle.counter(COUNTER).get(),
10,
"3 + 3 interleaved increments plus 4 on A must total 10 on A; {} | {}",
view_of(&a),
view_of(&b)
);
assert_eq!(
b.handle.counter(COUNTER).get(),
a.handle.counter(COUNTER).get(),
"both nodes must converge on the identical total; {} | {}",
view_of(&a),
view_of(&b)
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_wrong_secret_peer_never_joins() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
let intruder = start_node(
&router,
&clock,
OTHER_SECRET,
"node-intruder",
3,
vec![a.addr.clone(), b.addr.clone()],
);
settle(&clock, 10).await;
assert!(
a.handle.frames_rejected_total() > 0,
"A must actually have seen and REFUSED the intruder's frames — a view \
that stays at two because nothing arrived proves nothing; {}",
view_of(&a)
);
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"a peer signing with a different secret must never enter the view",
);
a.token.cancel();
b.token.cancel();
intruder.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_clean_leave_converges_to_one_member_view() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
a.handle.counter(COUNTER).increment();
settle(&clock, 2).await;
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, 1).await;
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"a clean leave must converge A to a one-member view well before the \
suspicion timeout ({}ms); {}",
SUSPICION.as_millis(),
view_of(&a)
);
a.handle.counter(COUNTER).increment();
assert_eq!(
a.handle.counter(COUNTER).get(),
2,
"the surviving node must keep incrementing and reading its counter; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_kill_without_leave_converges_after_suspicion_timeout() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
a.handle.counter(COUNTER).increment();
settle(&clock, 2).await;
router.disconnect(&b.addr);
b.token.cancel();
advance_time(&clock, PUSH.saturating_mul(2)).await;
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned(), "node-b".to_owned()],
"a silent peer must stay in the view until the suspicion timeout — \
suspicion is the correctness path, not an instant eviction; {}",
view_of(&a)
);
advance_time(&clock, SUSPICION).await;
settle(&clock, 1).await;
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"past the suspicion timeout the killed peer must leave the view; {}",
view_of(&a)
);
a.handle.counter(COUNTER).increment();
assert_eq!(
a.handle.counter(COUNTER).get(),
2,
"the survivor must keep serving the counter after a peer is killed; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_departure_carries_the_final_document() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
b.handle.counter(COUNTER).increment_by(7);
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, 1).await;
assert_eq!(
a.handle.counter(COUNTER).get(),
7,
"the survivor must end up with the increments the departing node held \
at cancellation — a counter that reads 0 here is a write accepted and \
then thrown away; {} | {}",
view_of(&a),
view_of(&b)
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"…and the departure must still converge the view to one member; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_departure_survives_a_full_peer_queue() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let stalled = Arc::new(StalledPeerTransport::new(
router.endpoint() as Arc<dyn PeerTransport>
));
let b = start_node_on(
Arc::clone(&stalled) as Arc<dyn PeerTransport>,
&clock,
SECRET,
"node-b",
2,
vec![a.addr.clone()],
);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"the pair must converge while B's queue still drains, or the departure \
below proves nothing",
);
stalled.stall();
b.handle.counter(COUNTER).increment_by(7);
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, 1).await;
assert!(
stalled.farewell_targets().contains(&a.addr),
"the departure must be offered to the transport's departure lane, not \
to the queue it would be dropped from; observed {:?}",
stalled.farewell_targets()
);
assert_eq!(
a.handle.counter(COUNTER).get(),
7,
"the survivor must still receive the increments the departing node held \
at cancellation — a stalled peer must not turn a clean leave into lost \
writes; {} | {}",
view_of(&a),
view_of(&b)
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"…and the departure must still converge the view to one member well \
inside the suspicion timeout; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_a_refused_departure_is_counted_and_leaves_the_suspicion_path() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let stalled = Arc::new(StalledPeerTransport::new(
router.endpoint() as Arc<dyn PeerTransport>
));
let b = start_node_on(
Arc::clone(&stalled) as Arc<dyn PeerTransport>,
&clock,
SECRET,
"node-b",
2,
vec![a.addr.clone()],
);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"the pair must converge before B's transport goes dead",
);
assert_eq!(
b.handle
.inner
.metrics
.frames_dropped
.load(Ordering::Relaxed),
0,
"sanity: nothing has been dropped yet"
);
stalled.stall();
stalled.go_deaf();
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, 1).await;
assert!(
b.handle
.inner
.metrics
.frames_dropped
.load(Ordering::Relaxed)
> 0,
"a departure the transport refused must be counted on the way out — it \
is the last thing this process can say about itself, and a silent \
degradation here reads exactly like a clean leave"
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned(), "node-b".to_owned()],
"sanity: nothing reached A, so its view is still the pair; {}",
view_of(&a)
);
advance_time(&clock, SUSPICION).await;
settle(&clock, 1).await;
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"past the suspicion timeout the survivor must converge without ever \
having heard the departure; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_write_storm_is_bounded_by_the_push_cadence_floor() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
advance_time(&clock, PUSH.checked_div(2).unwrap_or(PUSH)).await;
let before = pushes_sent(&a);
for _ in 0..50 {
a.handle.counter(COUNTER).increment();
tokio::task::yield_now().await;
}
let storm_pushes = pushes_sent(&a).saturating_sub(before);
assert!(
storm_pushes >= 1,
"the first nudge after a quiet gap must still push promptly — a floor \
that swallows it would make every write wait out the interval; {}",
view_of(&a)
);
assert!(
storm_pushes <= 4,
"50 increments inside one floor window must collapse into a handful of \
pushes, not 50: observed {storm_pushes} pushes; {}",
view_of(&a)
);
settle(&clock, 4).await;
assert_eq!(
b.handle.counter(COUNTER).get(),
50,
"every increment in the storm must reach B once the pushes resume; {} | {}",
view_of(&a),
view_of(&b)
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn metrics_source_emits_every_cluster_family() {
use crate::actuator::MetricKind;
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
let families = collect_families(&a);
let names: Vec<&str> = families.iter().map(|family| family.name.as_str()).collect();
assert_eq!(
names,
vec![
"autumn_cluster_members",
"autumn_cluster_pushes_sent_total",
"autumn_cluster_pushes_unsendable_total",
"autumn_cluster_pushes_received_total",
"autumn_cluster_merges_applied_total",
"autumn_cluster_frames_dropped_total",
"autumn_cluster_frames_rejected_total",
],
"the seven documented families must all be emitted, under exactly the \
names docs/guide/clustering.md tells operators to alert on"
);
assert_eq!(
family_value(&families, "autumn_cluster_pushes_unsendable_total"),
Some(0.0),
"a healthy node publishes the unsendable series at zero — an alert on a \
label that only appears once the document outgrows the frame cap is an \
alert nobody wrote"
);
let members = families
.iter()
.find(|family| family.name == "autumn_cluster_members")
.expect("the members family is in the list asserted above");
let expected = f64::from(u32::try_from(a.handle.members().len()).unwrap_or(u32::MAX));
assert_eq!(
members.kind,
MetricKind::Gauge,
"the view's size is a gauge"
);
assert_eq!(
members.samples.iter().map(|s| s.value).collect::<Vec<_>>(),
vec![expected],
"the members gauge must track the local view exactly; {} | {}",
view_of(&a),
view_of(&b)
);
assert!(
expected > 1.0,
"sanity: the gauge must be read against a converged two-member view, \
or a broken gauge reading 1 would pass; {}",
view_of(&a)
);
assert_eq!(
rejected_series(&families),
zeroed_series(),
"every reason label must be published from boot, and on a healthy \
cluster every one of them must read ZERO — asserting only the labels \
leaves the values, which are the thing an alert fires on, unpinned"
);
a.token.cancel();
b.token.cancel();
}
#[test]
fn rejection_counters_are_per_reason_and_fold_framing_into_oversize() {
let metrics = ClusterMetrics::default();
assert_eq!(
metrics.rejections_by_reason(),
REJECT_REASONS
.iter()
.map(|reason| (reason.label(), 0))
.collect::<Vec<_>>(),
"a fresh node must publish every documented series at zero"
);
for reason in REJECT_REASONS {
metrics.record_rejection(reason);
}
assert_eq!(
metrics.rejections_by_reason(),
REJECT_REASONS
.iter()
.map(|reason| (reason.label(), 1))
.collect::<Vec<_>>(),
"each reason must increment its OWN series, never a shared or shifted one"
);
metrics.record_rejection(RejectReason::Mac);
metrics.framing_rejected.store(3, Ordering::Relaxed);
let expected: Vec<(&str, u64)> = REJECT_REASONS
.iter()
.map(|reason| match *reason {
RejectReason::Oversize => (reason.label(), 4),
RejectReason::Mac => (reason.label(), 2),
other => (other.label(), 1),
})
.collect();
assert_eq!(
metrics.rejections_by_reason(),
expected,
"the transport's framing rejections must land in `oversize` (not \
`malformed`, and not in a series of their own), while every other \
reason keeps its own count"
);
assert_eq!(
metrics.rejected_total(),
13,
"the unlabelled aggregate must be the sum over every reason INCLUDING \
the folded framing rejections"
);
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn forged_frames_are_attributed_to_their_own_reason_series() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let hostile = router.endpoint();
let hostile_addr = hostile.local_addr().to_string();
settle(&clock, 2).await;
assert_eq!(
rejected_series(&collect_families(&a)),
zeroed_series(),
"sanity: nothing has been rejected yet, or the deltas below prove nothing"
);
let forge = |secret: &[u8], cluster: &str, sender: &str, incarnation: u64, seq: u64| {
wire::sign_envelope(
secret,
cluster,
sender,
incarnation,
seq,
&ClusterMessage::Leave,
)
.as_ref()
.and_then(wire::encode_frame)
.unwrap_or_default()
};
let wrong_secret = forge(OTHER_SECRET, CLUSTER, "node-intruder", 9, 1);
let wrong_cluster = forge(SECRET, "some-other-cluster", "node-stranger", 9, 1);
let replayed = forge(SECRET, CLUSTER, "node-ghost", 9, 1);
for frame in [
wrong_secret,
wrong_cluster,
replayed.clone(),
replayed.clone(),
] {
assert!(
router.deliver(&hostile_addr, &a.addr, frame),
"every forged frame must reach node A, or this test proves nothing"
);
settle(&clock, 1).await;
}
let expected: Vec<(String, f64)> = REJECT_REASONS
.iter()
.map(|reason| {
let value = match *reason {
RejectReason::Mac | RejectReason::Cluster | RejectReason::Replay => 1.0,
_ => 0.0,
};
(reason.label().to_owned(), value)
})
.collect();
assert_eq!(
rejected_series(&collect_families(&a)),
expected,
"a wrong-secret frame must move ONLY `mac`, a foreign cluster name only \
`cluster`, and a replayed frame only `replay` — an operator alerting \
on `mac` must not be looking at a series somebody else's traffic \
increments; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn an_oversized_document_is_counted_as_unsendable() {
const FILLER_CELLS: u64 = 2_000;
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
assert_eq!(
pushes_unsendable(&a),
0,
"sanity: a healthy node must be able to send; {}",
view_of(&a)
);
grow_document_past_the_frame_cap(&a, FILLER_CELLS);
let sent_before = pushes_sent(&a);
settle(&clock, 3).await;
assert!(
pushes_unsendable(&a) > 0,
"a document past the frame cap must increment the unsendable counter — \
silently dropping every heartbeat is the one failure an operator has no \
other way to see; {}",
view_of(&a)
);
assert_eq!(
pushes_sent(&a),
sent_before,
"…and it must NOT be counted as a push that was sent; {}",
view_of(&a)
);
assert_eq!(
family_value(
&collect_families(&a),
"autumn_cluster_pushes_unsendable_total"
),
Some(f64::from(
u32::try_from(pushes_unsendable(&a)).unwrap_or(u32::MAX)
)),
"the counter must be published as autumn_cluster_pushes_unsendable_total"
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_pruned_member_rejoins_at_a_lower_incarnation() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"the pair must converge before B departs, or the prune below proves nothing",
);
let departed_incarnation = b.handle.incarnation();
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, window_rounds()).await;
assert_eq!(
record_of(&a, &b.id),
None,
"sanity: A must have pruned B's tombstone after ten suspicion timeouts, \
or the rejoin below is testing the pre-prune path instead; {}",
view_of(&a)
);
let rejoin_incarnation = 1;
assert!(
rejoin_incarnation < departed_incarnation,
"the rejoin must really be lower than the watermark A recorded \
({departed_incarnation}), or this test asserts nothing"
);
let push_from_b = |seq: u64| signed_self_push(&b.id, &b.addr, rejoin_incarnation, seq);
assert!(
router.deliver(&b.addr, &a.addr, push_from_b(0)),
"the returning node's push must reach A to prove anything"
);
settle(&clock, 1).await;
let replay_drops = a
.handle
.inner
.metrics
.rejections_by_reason()
.into_iter()
.find(|(reason, _)| *reason == "replay")
.map_or(u64::MAX, |(_, count)| count);
assert_eq!(
replay_drops,
0,
"a pruned member is a forgotten member: its next push must be judged as \
a FRESH sender, whatever incarnation it carries, or the returning node \
is replay-dropped for good; {}",
view_of(&a)
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"a record at an incarnation this node has already collected must not be \
re-adopted while it still remembers collecting it; {}",
view_of(&a)
);
settle(
&clock,
window_rounds().saturating_mul(PRUNE_MEMORY_WINDOW_MULTIPLE),
)
.await;
assert!(
router.deliver(&b.addr, &a.addr, push_from_b(1)),
"the returning node's later push must reach A too"
);
settle(&clock, 1).await;
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned(), "node-b".to_owned()],
"the memory is a bounded local note, not a permanent refusal: once it \
lapses the returning node rejoins on its own, with no operator and no \
restart; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_replayed_record_of_a_departed_member_leaves_the_document_again() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"the pair must converge before B departs, or the replay below is not a \
replay of anything",
);
let departed_incarnation = b.handle.incarnation();
b.token.cancel();
advance_time(&clock, Duration::from_millis(250)).await;
settle(&clock, window_rounds()).await;
assert_eq!(
record_of(&a, &b.id),
None,
"sanity: A must have pruned B's tombstone, or the replay below is not \
hitting the forgotten-watermark path at all; {}",
view_of(&a)
);
let captured_push = |seq: u64| signed_self_push(&b.id, &b.addr, departed_incarnation, seq);
assert!(
router.deliver(&b.addr, &a.addr, captured_push(0)),
"the replayed frame must reach A to prove anything"
);
settle(&clock, 1).await;
assert_eq!(
record_of(&a, &b.id),
None,
"a replayed record at an incarnation this node has already collected \
must not put the member back in the document; {}",
view_of(&a)
);
settle(
&clock,
window_rounds().saturating_mul(PRUNE_MEMORY_WINDOW_MULTIPLE),
)
.await;
assert!(
router.deliver(&b.addr, &a.addr, captured_push(1)),
"the second replayed frame must reach A too"
);
settle(&clock, 1).await;
assert_eq!(
record_of(&a, &b.id),
Some(MemberRecord::alive(b.addr.clone(), departed_incarnation)),
"sanity: past its memory window A learns the replayed record — that is \
the join doing its job, and it is what the rest of this test has to \
undo; {}",
view_of(&a)
);
settle(&clock, window_rounds()).await;
assert_eq!(
record_of(&a, &b.id),
Some(MemberRecord::left(b.addr.clone(), departed_incarnation)),
"a member Down for a whole tombstone window must be recorded as Left, \
at its own incarnation, so the tombstone lifecycle can reach it; {}",
view_of(&a)
);
settle(&clock, window_rounds()).await;
assert_eq!(
record_of(&a, &b.id),
None,
"the converted record must then prune: a replayed frame may cost one \
window of churn, but it must never leave a record in the document \
that nothing can ever take out; {}",
view_of(&a)
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned()],
"…and the survivor's view is its own again; {}",
view_of(&a)
);
a.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn push_round_retires_writers_for_addresses_that_left_the_target_set() {
const OLD_ADDR: &str = "127.0.0.1:47901";
const NEW_ADDR: &str = "127.0.0.1:47902";
let router = LoopbackRouter::new();
let clock = test_clock();
let spy = Arc::new(RetainSpy::new(router.endpoint()));
let token = CancellationToken::new();
let handle = ClusterNode::start(
ClusterRuntimeConfig {
cluster_name: CLUSTER.to_owned(),
secret: SECRET.to_vec(),
node_id: Some("node-a".to_owned()),
advertise_addr: None,
seed_peers: Vec::new(),
push_interval: PUSH,
suspicion_timeout: SUSPICION,
},
Arc::new(SeededEntropy::new(11)),
Arc::new(clock.clone()),
token.clone(),
Arc::clone(&spy) as Arc<dyn PeerTransport>,
)
.expect("a cluster node must start on a spying transport");
handle
.inner
.lock_state()
.members
.insert("node-x".to_owned(), MemberRecord::alive(OLD_ADDR, 1));
settle(&clock, 2).await;
assert!(
spy.last_retained().contains(OLD_ADDR),
"a known member's address must be in the retained set; observed {:?}",
spy.last_retained()
);
handle
.inner
.lock_state()
.members
.insert("node-x".to_owned(), MemberRecord::alive(NEW_ADDR, 2));
settle(&clock, 2).await;
let retained = spy.last_retained();
assert!(
retained.contains(NEW_ADDR),
"the member's new address must become a target; observed {retained:?}"
);
assert!(
!retained.contains(OLD_ADDR),
"the abandoned address must be retired through the push round, or its \
writer task and queue outlive the address forever; observed {retained:?}"
);
token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn receive_loop_reports_only_unauthenticated_frames_to_the_transport() {
let router = LoopbackRouter::new();
let clock = test_clock();
let spy = Arc::new(RejectionSpy::new(router.endpoint()));
let addr = spy.local_addr().to_string();
let token = CancellationToken::new();
let handle = ClusterNode::start(
ClusterRuntimeConfig {
cluster_name: CLUSTER.to_owned(),
secret: SECRET.to_vec(),
node_id: Some("node-a".to_owned()),
advertise_addr: None,
seed_peers: Vec::new(),
push_interval: PUSH,
suspicion_timeout: SUSPICION,
},
Arc::new(SeededEntropy::new(11)),
Arc::new(clock.clone()),
token.clone(),
Arc::clone(&spy) as Arc<dyn PeerTransport>,
)
.expect("a cluster node must start on a spying transport");
let peer = start_node(&router, &clock, SECRET, "node-b", 2, vec![addr.clone()]);
let intruder = start_node(
&router,
&clock,
OTHER_SECRET,
"node-intruder",
3,
vec![addr.clone()],
);
settle(&clock, 8).await;
assert!(
handle.frames_rejected_total() > 0,
"sanity: the intruder's frames must actually have arrived and been \
refused, or this test proves nothing"
);
let reported = spy.reported();
assert!(
reported.contains(&intruder.addr),
"a frame that failed the MAC must be reported to the transport, which \
is the only way the connection holding a slot can be closed; the node \
reported {reported:?}"
);
assert!(
!reported.contains(&peer.addr),
"a peer whose frames verify must NEVER be reported: its connection has \
been earned, and charging it for a `replay` or an unknown `payload` \
would close a healthy peer's socket; the node reported {reported:?}"
);
token.cancel();
peer.token.cancel();
intruder.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_silent_peer_reads_as_suspect_before_it_drops_out() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"the pair must be Alive before B goes silent",
);
router.disconnect(&b.addr);
advance_time(&clock, PUSH.saturating_mul(3)).await;
let view = a.handle.members();
let peer = view.iter().find(|member| member.id == b.id);
assert_eq!(
peer.map(|member| member.status),
Some(ClusterMemberStatus::Suspect),
"a peer silent past two push intervals must be reported Suspect — not \
quietly Alive, and not evicted; observed {view:?}"
);
assert_eq!(
view.iter()
.find(|member| member.id == a.id)
.map(|member| member.status),
Some(ClusterMemberStatus::Alive),
"…while this node itself stays Alive in its own view; observed {view:?}"
);
assert_eq!(
member_ids(&a.handle),
vec!["node-a".to_owned(), "node-b".to_owned()],
"a Suspect member is a warning, not an eviction: it must still be in \
the view; {}",
view_of(&a)
);
a.token.cancel();
b.token.cancel();
}
#[tokio::test(flavor = "current_thread", start_paused = true)]
async fn loopback_replayed_leave_is_refuted_by_live_node() {
let router = LoopbackRouter::new();
let clock = test_clock();
let a = start_node(&router, &clock, SECRET, "node-a", 1, Vec::new());
let b = start_node(&router, &clock, SECRET, "node-b", 2, vec![a.addr.clone()]);
settle(&clock, 6).await;
let captured_incarnation = a.handle.incarnation();
let frame = wire::sign_envelope(
SECRET,
CLUSTER,
&a.id,
captured_incarnation,
u64::MAX / 2,
&ClusterMessage::Leave,
)
.as_ref()
.and_then(wire::encode_frame)
.unwrap_or_default();
assert!(
!frame.is_empty(),
"the test must be able to forge a signed leave frame to replay"
);
let delivered = router.deliver(&a.addr, &b.addr, frame);
assert!(
delivered,
"the forged frame must reach node B to prove anything"
);
settle(&clock, 6).await;
assert!(
a.handle.incarnation() > captured_incarnation,
"A must refute the replayed leave by bumping its incarnation \
(was {captured_incarnation}, now {}); {}",
a.handle.incarnation(),
view_of(&a)
);
assert_converged(
&a,
&b,
&["node-a", "node-b"],
"a replayed leave must not evict a live node: B's view must return to \
two members",
);
a.token.cancel();
b.token.cancel();
}