tidaldb/tidal-server/tests/cluster_election.rs
jx12n 95461d3cf8 feat(m11): Raft leader election over WAL stream (m11p4)
Kind-3 term markers in the WAL stream, STREAM-relative vote frontiers,
heartbeat-only divergence detection + quarantine, and fenced promote.
Elections converge in 0.6–1.0s; zero acked-write loss across all kill points.
Closes G5 (leaderless recovery) from the v0.9 wave.
2026-06-11 23:30:24 -06:00

531 lines
21 KiB
Rust
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//! m11p4 exit gates — automatic failover, fencing, bounded churn — over a
//! REAL 3-process cluster (tier 3, `cluster-e2e` feature).
//!
//! The three gates from docs/roadmap-to-cluster.md §4/m11p4:
//!
//! 1. **Auto-failover, zero acked loss** (`mp_auto_failover_*`): SIGKILL the
//! leader under `ack=quorum` load with ZERO operator verbs → a survivor is
//! elected and writes resume in <10s, and the m11p3 ledger invariants
//! (frontier + content) hold on the new leader, across repeated rounds
//! with pseudo-random kill points.
//! 2. **Fencing under partition + restart** (`mp_fenced_ex_leader_*`): the
//! leader is partitioned away (real TCP severs), the survivors elect, the
//! old leader RESTARTS while still partitioned — and cannot accept a
//! single write (its durable boot state forbids self-leadership, the
//! §1.4-1 fix), then rejoins as a follower on heal.
//! 3. **Bounded churn** (`mp_flapping_links_bounded_churn`): repeated
//! sever/heal cycles on the leader's links produce bounded elections (the
//! pre-vote absorbs flaps; terms never explode) and the cluster converges
//! to exactly one leader that serves quorum writes.
//!
//! Election timings are tuned fast (5001000ms timeouts) so the suite stays
//! within the tier-3 budget; the production defaults scale the same
//! machinery up, not a different protocol.
#![cfg(feature = "cluster-e2e")]
#![allow(clippy::unwrap_used, clippy::significant_drop_tightening)]
mod support;
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use std::time::{Duration, Instant};
use support::{
multiproc::{ClusterOptions, MultiProcCluster, convergence_budget},
partition::{ProxyController, proxied_rewrite},
};
/// All three roster regions, so every directed edge gets its own relay and a
/// node can be isolated in BOTH directions (inbound edges via
/// `region(r).sever_all()` cut what reaches it; outbound edges via
/// `edge(r, peer)` cut its own heartbeats/ships — without the outbound cut a
/// "partitioned" LEADER keeps resetting every follower's election timer).
const ALL_REGIONS: [&str; 3] = ["us-east", "eu-west", "ap-south"];
/// Fully isolate `region` from `peers` (both directions, gRPC + HTTP).
fn isolate(proxies: &ProxyController, region: &str, peers: &[&str]) {
proxies.region(region).sever_all();
for peer in peers {
proxies.edge(region, peer).sever_all();
}
}
/// Undo [`isolate`].
fn rejoin(proxies: &ProxyController, region: &str, peers: &[&str]) {
proxies.region(region).heal_all();
for peer in peers {
proxies.edge(region, peer).heal_all();
}
}
/// The fast election block every test in this suite runs with. Constraint:
/// `lease (350) + heartbeat (100) < timeout_min (500)` — the C2 bound.
const FAST_ELECTION_YAML: &str = "election:\n heartbeat_interval_ms: 100\n election_timeout_min_ms: 500\n election_timeout_max_ms: 1000\n leader_lease_ms: 350";
/// The exit gate's failover budget: detect → elect → writes resume.
const FAILOVER_BUDGET: Duration = Duration::from_secs(10);
/// A unique all-alpha search token for an entity id (mirrors the m11p3
/// ledger checker's probe).
fn item_token(entity_id: u64) -> String {
let mut token = String::from("elx");
for d in entity_id.to_string().bytes() {
token.push(char::from(b'a' + (d - b'0')));
}
token
}
/// POST with `x-tidal-ack` through a dedicated client. `Some(seq)` only for
/// a 2xx carrying `x-tidal-seq` — the ledger's definition of "acknowledged".
fn post_acked(
client: &reqwest::blocking::Client,
base: &str,
path: &str,
ack: &str,
body: &serde_json::Value,
) -> Option<u64> {
let resp = client
.post(format!("{base}{path}"))
.header("x-tidal-ack", ack)
.json(body)
.send()
.ok()?;
if !resp.status().is_success() {
return None;
}
resp.headers()
.get("x-tidal-seq")?
.to_str()
.ok()?
.parse()
.ok()
}
/// Poll the LIVE nodes for an elected leader: a node whose local status
/// reports `role == "leader"` at a term above `after_term`. Returns
/// `(node_idx, term, elapsed)`.
fn await_elected_leader(
cluster: &MultiProcCluster,
candidates: &[usize],
after_term: u64,
budget: Duration,
) -> (usize, u64, Duration) {
let started = Instant::now();
let deadline = started + budget;
loop {
for &idx in candidates {
if let Some(status) = cluster.local_status(idx) {
let term = status["term"].as_u64().unwrap_or(0);
if status["role"].as_str() == Some("leader") && term > after_term {
return (idx, term, started.elapsed());
}
}
}
assert!(
Instant::now() < deadline,
"no leader elected among {candidates:?} within {budget:?} \
(terms must move past {after_term})"
);
std::thread::sleep(Duration::from_millis(50));
}
}
/// Write through `base` until one `ack=quorum` write succeeds (or the
/// deadline passes). Returns the first acked seq.
fn await_write_resumes(
client: &reqwest::blocking::Client,
base: &str,
entity_id: u64,
deadline: Instant,
) -> u64 {
loop {
if let Some(seq) = post_acked(
client,
base,
"/signals",
"quorum",
&serde_json::json!({ "entity_id": entity_id, "signal": "view", "weight": 1.0 }),
) {
return seq;
}
assert!(
Instant::now() < deadline,
"quorum writes did not resume before the failover budget expired"
);
std::thread::sleep(Duration::from_millis(50));
}
}
/// Gate 1: kill the leader under `ack=quorum` load — zero operator verbs —
/// across repeated rounds. Every round asserts: a survivor is elected and
/// writes resume within the 10s budget; the new leader's durable frontier
/// covers every acknowledged seq (INVARIANT A); every acknowledged item is
/// searchable on the new leader (INVARIANT B). The killed node restarts and
/// rejoins (reseeded if it quarantined with a divergent suffix — writes that
/// were leader-staged but never quorum-acked).
// One linear multi-round drill (load -> kill -> elect -> invariants ->
// rejoin); splitting it would scatter the round's ordering rules.
#[allow(clippy::too_many_lines)]
#[test]
fn mp_auto_failover_writes_resume_zero_acked_loss() {
let rounds: usize = std::env::var("TIDAL_ELECTION_KILLPOINTS")
.ok()
.and_then(|v| v.parse().ok())
.filter(|&n| n > 0)
.unwrap_or(5);
let mut cluster = MultiProcCluster::start_with(
ClusterOptions::new(3).with_topology_extra(FAST_ELECTION_YAML),
);
let client = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(3))
.build()
.unwrap();
let mut entity_cursor: u64 = 1;
let mut current_leader = 0usize;
let mut last_term = 0u64;
for round in 0..rounds {
// ── Load: two writer threads against the current leader. ──────────
let stop = Arc::new(AtomicBool::new(false));
let leader_base = cluster.node(current_leader);
let mut writers = Vec::new();
for w in 0..2u64 {
let stop = Arc::clone(&stop);
let base = leader_base.clone();
let first_entity = entity_cursor + w * 10_000;
writers.push(std::thread::spawn(move || {
let client = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(3))
.build()
.unwrap();
let mut acked: Vec<(u64, u64)> = Vec::new();
let mut entity = first_entity;
while !stop.load(Ordering::Acquire) {
let item_seq = post_acked(
&client,
&base,
"/items",
"quorum",
&serde_json::json!({
"entity_id": entity,
"metadata": { "title": item_token(entity) },
}),
);
let view_seq = post_acked(
&client,
&base,
"/signals",
"quorum",
&serde_json::json!({
"entity_id": entity, "signal": "view", "weight": 1.0
}),
);
if let Some(seq) = item_seq {
acked.push((entity, seq.max(view_seq.unwrap_or(0))));
}
entity += 1;
}
acked
}));
}
// Pseudo-random kill point per round (reproducible).
std::thread::sleep(Duration::from_millis(150 + (round as u64 * 97) % 400));
cluster.kill_hard(current_leader);
stop.store(true, Ordering::Release);
let mut ledger: Vec<(u64, u64)> = Vec::new();
for w in writers {
ledger.extend(w.join().expect("writer thread"));
}
let max_acked_seq = ledger.iter().map(|&(_, s)| s).max().unwrap_or(0);
// ── ZERO operator verbs: the survivors elect on their own. ─────────
let survivors: Vec<usize> = (0..3).filter(|&i| i != current_leader).collect();
let (new_leader, new_term, elapsed) =
await_elected_leader(&cluster, &survivors, last_term, FAILOVER_BUDGET);
// Writes must RESUME (not just leadership exist) inside the budget.
let resume_deadline = Instant::now() + FAILOVER_BUDGET;
let _ = await_write_resumes(
&client,
&cluster.node(new_leader),
900_000 + round as u64,
resume_deadline,
);
println!(
"round {round}: leader {current_leader} killed -> {new_leader} elected at \
term {new_term} in {elapsed:?}; {} acked writes (max seq {max_acked_seq})",
ledger.len()
);
// ── INVARIANT A (frontier): the elected leader's ELECTION-TIME
// position — in the killed leader's stream numbering, the only
// numbering the acked seqs live in — covers every acknowledged
// write (the vote restriction guarantees it). The leader's own
// `last_seq` is its NEW stream's numbering and is NOT comparable.
let status = cluster.local_status(new_leader).expect("leader status");
let prev_term = status["prev_log_term"].as_u64().unwrap();
let prev_seq = status["prev_log_seq"].as_u64().unwrap();
assert_eq!(
prev_term, last_term,
"round {round}: the elected leader's election-time tail term must be \
the killed leader's term (same stream numbering as the acked seqs)"
);
assert!(
prev_seq >= max_acked_seq,
"round {round}: elected leader's election-time frontier {prev_seq} is \
below an acknowledged seq {max_acked_seq} — acked-write loss"
);
// ── INVARIANT B (content): every acked item is searchable on the
// new leader (the text index auto-commits within ~2s; allow 10).
let search_deadline = Instant::now() + Duration::from_secs(10);
for &(entity, _) in &ledger {
let token = item_token(entity);
loop {
let found: serde_json::Value = client
.get(format!(
"{}/search?query={token}&limit=5",
cluster.node(new_leader)
))
.send()
.unwrap()
.json()
.unwrap();
let hit = found["items"]
.as_array()
.is_some_and(|r| r.iter().any(|x| x["entity_id"].as_u64() == Some(entity)));
if hit {
break;
}
assert!(
Instant::now() < search_deadline,
"round {round}: acked item {entity} (token {token}) not found on \
the elected leader — acked-write loss"
);
std::thread::sleep(Duration::from_millis(200));
}
}
// ── Bring the killed node back for the next round. A divergent
// suffix (leader-staged, never quorum-acked writes) legitimately
// quarantines — the documented recovery is a reseed.
cluster.restart(current_leader, &[]);
let rejoin_deadline = Instant::now() + convergence_budget() + Duration::from_secs(10);
loop {
let status = cluster.local_status(current_leader);
let quarantined = status
.as_ref()
.and_then(|s| s["quarantined"].as_bool())
.unwrap_or(false);
if quarantined {
println!(
"round {round}: restarted node {current_leader} quarantined \
(divergent suffix) — reseeding, the documented recovery"
);
cluster.kill_hard(current_leader);
cluster.wipe_data_dir(current_leader);
cluster.restart(current_leader, &[]);
}
let caught_up = cluster.local_status(current_leader).is_some_and(|s| {
s["term"].as_u64().unwrap_or(0) >= new_term
&& s["role"].as_str() == Some("follower")
&& s["lag_events"].as_u64() == Some(0)
});
if caught_up {
break;
}
assert!(
Instant::now() < rejoin_deadline,
"round {round}: killed node {current_leader} did not rejoin/converge"
);
std::thread::sleep(Duration::from_millis(200));
}
current_leader = new_leader;
last_term = new_term;
entity_cursor += 100_000;
}
}
/// Gate 2: partition the leader away with real TCP severs, let the survivors
/// elect, RESTART the old leader while still partitioned — it must boot as a
/// follower (durable election state, never the topology file) and cannot
/// accept a single write; on heal it rejoins the new term as a follower.
#[test]
fn mp_fenced_ex_leader_restart_cannot_write() {
let (rewrite, proxies) = proxied_rewrite(&ALL_REGIONS);
let mut opts = ClusterOptions::new(3)
.with_topology_extra(FAST_ELECTION_YAML)
.with_rewrite(rewrite);
opts.log = "info".into();
let mut cluster = MultiProcCluster::start_with(opts);
let client = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(2))
.build()
.unwrap();
// Baseline data, fully converged BEFORE the partition so the old leader
// carries no divergent suffix (this gate is about fencing, not reseed).
for entity in 1..=5u64 {
post_acked(
&client,
&cluster.node(0),
"/signals",
"quorum",
&serde_json::json!({ "entity_id": entity, "signal": "view", "weight": 1.0 }),
)
.expect("baseline quorum write");
}
cluster.wait_converged_all(convergence_budget());
// ── Partition the leader away (BOTH directions); the survivors elect. ──
isolate(&proxies, "us-east", &["eu-west", "ap-south"]);
let (new_leader, new_term, elapsed) =
await_elected_leader(&cluster, &[1, 2], 0, FAILOVER_BUDGET);
println!("survivors elected node {new_leader} at term {new_term} in {elapsed:?}");
// The new leadership serves quorum writes (2 of 3 replicas).
let _ = await_write_resumes(
&client,
&cluster.node(new_leader),
700_001,
Instant::now() + FAILOVER_BUDGET,
);
// ── Restart the old leader while STILL partitioned. ────────────────────
cluster.kill_hard(0);
cluster.restart(0, &[]);
// The §1.4-1 assertion: across a multi-second window, the restarted
// ex-leader REFUSES every write — its durable state boots it as a
// follower and the topology file's `leader: us-east` is dead weight.
let fence_window = Instant::now() + Duration::from_secs(3);
let mut attempts = 0u32;
while Instant::now() < fence_window {
let accepted = post_acked(
&client,
&cluster.node(0),
"/signals",
"leader",
&serde_json::json!({ "entity_id": 700_100, "signal": "view", "weight": 1.0 }),
);
assert!(
accepted.is_none(),
"the restarted, partitioned ex-leader ACCEPTED a write (seq {accepted:?}) — \
the §1.4-1 split-brain hole is open"
);
attempts += 1;
std::thread::sleep(Duration::from_millis(100));
}
let status = cluster.local_status(0).expect("ex-leader serves status");
assert_eq!(
status["is_leader"].as_bool(),
Some(false),
"restarted ex-leader must not claim leadership: {status}"
);
println!("fencing held across {attempts} write attempts: {status}");
// ── Heal: the ex-leader joins the new term as a follower and converges.
rejoin(&proxies, "us-east", &["eu-west", "ap-south"]);
let rejoin_deadline = Instant::now() + convergence_budget() + Duration::from_secs(30);
loop {
let status = cluster.local_status(0).expect("status");
assert_ne!(
status["quarantined"].as_bool(),
Some(true),
"a fully-converged-then-partitioned ex-leader must rejoin CLEAN \
(no divergent suffix existed): {status}"
);
if status["term"].as_u64().unwrap_or(0) >= new_term
&& status["role"].as_str() == Some("follower")
&& status["lag_events"].as_u64() == Some(0)
{
break;
}
assert!(
Instant::now() < rejoin_deadline,
"healed ex-leader did not rejoin term {new_term} and converge: {status}"
);
std::thread::sleep(Duration::from_millis(200));
}
}
/// Gate 3: flapping links produce BOUNDED churn. Short flaps (below the
/// election timeout) are absorbed by the pre-vote/lease machinery; long
/// flaps elect; terms never explode; the cluster converges to exactly one
/// leader that serves quorum writes.
#[test]
fn mp_flapping_links_bounded_churn() {
let (rewrite, proxies) = proxied_rewrite(&ALL_REGIONS);
let mut cluster = MultiProcCluster::start_with(
ClusterOptions::new(3)
.with_topology_extra(FAST_ELECTION_YAML)
.with_rewrite(rewrite),
);
let _ = &mut cluster;
let client = reqwest::blocking::Client::builder()
.timeout(Duration::from_secs(2))
.build()
.unwrap();
// Short flaps: sever 250ms (below the 500ms timeout floor), heal 400ms.
// The lease refusals + timer resets must absorb these without elections.
for _ in 0..4 {
isolate(&proxies, "us-east", &["eu-west", "ap-south"]);
std::thread::sleep(Duration::from_millis(250));
rejoin(&proxies, "us-east", &["eu-west", "ap-south"]);
std::thread::sleep(Duration::from_millis(400));
}
// Long flaps: sever past the timeout so real elections happen, then heal.
for _ in 0..3 {
isolate(&proxies, "us-east", &["eu-west", "ap-south"]);
std::thread::sleep(Duration::from_millis(1_500));
rejoin(&proxies, "us-east", &["eu-west", "ap-south"]);
std::thread::sleep(Duration::from_millis(800));
}
// Convergence: exactly one leader, agreed term, bounded churn.
let deadline = Instant::now() + Duration::from_secs(15);
let (leaders, term) = loop {
let statuses: Vec<serde_json::Value> =
(0..3).filter_map(|i| cluster.local_status(i)).collect();
let leaders: Vec<usize> = statuses
.iter()
.enumerate()
.filter(|(_, s)| s["role"].as_str() == Some("leader"))
.map(|(i, _)| i)
.collect();
let terms: Vec<u64> = statuses
.iter()
.map(|s| s["term"].as_u64().unwrap_or(0))
.collect();
let agreed = terms.iter().max() == terms.iter().min();
if statuses.len() == 3 && leaders.len() == 1 && agreed {
break (leaders, terms[0]);
}
assert!(
Instant::now() < deadline,
"cluster did not converge to one leader on one term: \
leaders={leaders:?} terms={terms:?}"
);
std::thread::sleep(Duration::from_millis(100));
};
// Bounded churn: 3 long flaps + slop can justify a handful of terms,
// never dozens (a term explosion = pre-vote regression / livelock).
assert!(
term <= 15,
"term {term} after 7 flaps — election churn is unbounded"
);
println!("converged: leader node {} at term {term}", leaders[0]);
// The survivor of all that chaos still serves quorum writes.
let _ = await_write_resumes(
&client,
&cluster.node(leaders[0]),
800_001,
Instant::now() + FAILOVER_BUDGET,
);
}