mod support;
use spate_coordination::store::memory::MemoryStore;
use spate_coordination::store::{CasOutcome, CoordinationStore as _, Keyspace};
use spate_coordination::{CoordinationErrorKind, SplitCoordinator, SplitProgress};
use std::time::Instant;
use support::{
Held, LEASE, PhasedPlanner, crash, drive, drive_pair, runtime, split_id, store, worker,
worker_drain_deadline, worker_max_in_flight,
};
#[test]
fn racing_workers_partition_the_splits_and_complete_collectively() {
let rt = runtime();
let store = store();
let ids = ["s0", "s1", "s2", "s3", "s4", "s5", "s6", "s7"];
let planner = || Box::new(PhasedPlanner::one_final("partition:v1", &ids));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
let mut b = worker(&store, rt.handle(), Some("worker-b"));
a.start(planner()).unwrap();
b.start(planner()).unwrap();
let (mut held_a, mut held_b) = (Held::default(), Held::default());
let deadline = Instant::now() + support::DEADLINE;
while !(held_a.splits.len() + held_b.splits.len() == ids.len()
&& held_a.splits.keys().all(|k| !held_b.splits.contains_key(k))
&& held_a.splits.len() == 4)
{
assert!(
Instant::now() < deadline,
"timed out: waiting for a full disjoint partition"
);
for (coordinator, held) in [(&mut a, &mut held_a), (&mut b, &mut held_b)] {
held.fold(coordinator.poll().unwrap());
support::commit_held(coordinator, held);
}
std::thread::sleep(support::POLL_INTERVAL);
}
let settle_until = Instant::now() + LEASE * 2;
while Instant::now() < settle_until {
held_a.fold(a.poll().unwrap());
held_b.fold(b.poll().unwrap());
}
assert_eq!(held_a.splits.len(), 4, "balance must hold once reached");
assert_eq!(held_b.splits.len(), 4);
assert!(held_a.splits.keys().all(|k| !held_b.splits.contains_key(k)));
for held in [&held_a, &held_b] {
for (_, progress) in held.splits.values() {
if let Some(progress) = progress {
assert_eq!(progress.watermark, 1, "carried progress is the checkpoint");
}
}
}
let deadline = Instant::now() + support::DEADLINE;
while !(held_a.all_complete && held_b.all_complete) {
assert!(Instant::now() < deadline, "collective completion timed out");
for (held, coordinator) in [(&mut held_a, &mut a), (&mut held_b, &mut b)] {
let ids: Vec<String> = held.splits.keys().cloned().collect();
for id in ids {
match coordinator.commit(&split_id(&id), &SplitProgress::completed(100, vec![])) {
Ok(()) => {
held.splits.remove(&id);
}
Err(e) if e.kind == CoordinationErrorKind::Fenced => {
held.splits.remove(&id); }
Err(e) => panic!("commit failed: {e}"),
}
}
held.fold(coordinator.poll().unwrap());
}
}
}
#[test]
fn death_takeover_waits_out_the_lease_and_carries_progress() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("takeover:v1", &["t0", "t1"]));
let rt_a = runtime();
let mut a = worker(&store, rt_a.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "A claiming both splits", |h| {
h.splits.len() == 2
});
a.commit(
&split_id("t0"),
&SplitProgress::new(42, b"resume-here".to_vec()),
)
.unwrap();
let first_epochs: Vec<u64> = held_a.splits.values().map(|(e, _)| *e).collect();
let died_at = Instant::now();
crash(rt_a, a);
let mut b = worker(&store, rt.handle(), Some("worker-b"));
b.start(planner()).unwrap();
let mut held_b = Held::default();
drive(
&mut b,
&mut held_b,
"B taking over the dead worker's splits",
|h| h.splits.len() == 2,
);
assert!(
died_at.elapsed() >= LEASE / 4,
"takeover before the dead worker's lease could have expired: {:?}",
died_at.elapsed()
);
let (epoch, progress) = &held_b.splits["t0"];
let progress = progress
.as_ref()
.expect("progress carried to the new owner");
assert_eq!(progress.watermark, 42);
assert_eq!(progress.state, b"resume-here");
assert!(
first_epochs.iter().all(|first| epoch > first),
"takeover must bump the epoch: {first_epochs:?} -> {epoch}"
);
assert!(held_b.splits["t1"].1.is_none(), "t1 never committed");
}
#[test]
fn graceful_release_hands_off_without_waiting_out_the_lease() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("release:v1", &["r0", "r1"]));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "A claiming both splits", |h| {
h.splits.len() == 2
});
let mut b = worker(&store, rt.handle(), Some("worker-b"));
b.start(planner()).unwrap();
let mut held_b = Held::default();
let released_at = Instant::now();
a.release(&[split_id("r0"), split_id("r1")]).unwrap();
drop(a);
drive(&mut b, &mut held_b, "B claiming the released splits", |h| {
h.splits.len() == 2
});
assert!(
released_at.elapsed() < LEASE,
"released splits must hand off without a lease wait: {:?}",
released_at.elapsed()
);
}
#[test]
fn stable_instance_id_reclaims_fast_after_a_restart() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("reclaim:v1", &["m0"]));
let rt_old = runtime();
let mut old = worker(&store, rt_old.handle(), Some("pod-1"));
old.start(planner()).unwrap();
let mut held_old = Held::default();
drive(&mut old, &mut held_old, "predecessor claiming", |h| {
h.splits.len() == 1
});
old.commit(&split_id("m0"), &SplitProgress::new(7, vec![]))
.unwrap();
crash(rt_old, old);
let restarted_at = Instant::now();
let mut new = worker(&store, rt.handle(), Some("pod-1"));
new.start(planner()).unwrap();
let mut held_new = Held::default();
drive(
&mut new,
&mut held_new,
"restart reclaiming its split",
|h| h.splits.len() == 1,
);
assert!(
restarted_at.elapsed() < LEASE,
"reclaim must not wait out the lease: {:?}",
restarted_at.elapsed()
);
assert_eq!(
held_new.splits["m0"].1.as_ref().map(|p| p.watermark),
Some(7),
"progress survives the restart"
);
}
#[test]
fn live_twins_sharing_an_instance_id_are_fatal() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("twin:v1", &["w0"]));
let mut first = worker(&store, rt.handle(), Some("pod-1"));
first.start(planner()).unwrap();
let mut held_first = Held::default();
drive(&mut first, &mut held_first, "first twin claiming", |h| {
h.splits.len() == 1
});
let mut second = worker(&store, rt.handle(), Some("pod-1"));
second.start(planner()).unwrap();
let mut held_second = Held::default();
drive(
&mut second,
&mut held_second,
"second twin reclaiming",
|h| h.splits.len() == 1,
);
let deadline = Instant::now() + support::DEADLINE;
let error = 'outer: loop {
assert!(
Instant::now() < deadline,
"neither twin reported the shared instance_id"
);
for (twin, held) in [
(&mut first, &mut held_first),
(&mut second, &mut held_second),
] {
match twin.poll() {
Ok(events) => held.fold(events),
Err(e) => break 'outer e,
}
}
};
assert_eq!(error.kind, CoordinationErrorKind::Fatal);
assert!(error.to_string().contains("instance_id"), "{error}");
assert!(error.to_string().contains("pod-1"), "{error}");
}
#[test]
fn a_fenced_zombie_commit_writes_nothing() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("fence:v1", &["z0"]));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "A claiming the split", |h| {
h.splits.len() == 1
});
a.commit(&split_id("z0"), &SplitProgress::new(10, b"a-10".to_vec()))
.unwrap();
let thief_watermark = steal_as_thief(&rt, &store, "z0", "thief");
let error = a
.commit(&split_id("z0"), &SplitProgress::new(11, b"a-11".to_vec()))
.unwrap_err();
assert_eq!(error.kind, CoordinationErrorKind::Fenced, "{error}");
drive(&mut a, &mut held_a, "A observing the loss", |h| {
!h.splits.contains_key("z0")
});
let record = rt
.block_on(store.get(Keyspace::Durable, "split.z0"))
.unwrap()
.unwrap();
let json: serde_json::Value = serde_json::from_slice(&record.value).unwrap();
assert_eq!(json["owner"], "thief");
assert_eq!(json["watermark"], thief_watermark);
}
fn steal_as_thief(rt: &tokio::runtime::Runtime, store: &MemoryStore, id: &str, thief: &str) -> i64 {
rt.block_on(async {
let lease_key = format!("split.{id}");
let lease = store
.get(Keyspace::Ephemeral, &lease_key)
.await
.unwrap()
.expect("victim lease");
let mut lease_val: serde_json::Value = serde_json::from_slice(&lease.value).unwrap();
let record = store
.get(Keyspace::Durable, &lease_key)
.await
.unwrap()
.expect("record");
let mut record_val: serde_json::Value = serde_json::from_slice(&record.value).unwrap();
let epoch = record_val["epoch"].as_u64().unwrap() + 1;
lease_val["owner"] = thief.into();
lease_val["nonce"] = "thief-nonce".into();
lease_val["epoch"] = epoch.into();
assert!(matches!(
store
.update(
Keyspace::Ephemeral,
&lease_key,
serde_json::to_vec(&lease_val).unwrap(),
lease.revision,
)
.await
.unwrap(),
CasOutcome::Won(_)
));
record_val["epoch"] = epoch.into();
record_val["owner"] = thief.into();
let watermark = record_val["watermark"].as_i64().unwrap();
assert!(matches!(
store
.update(
Keyspace::Durable,
&lease_key,
serde_json::to_vec(&record_val).unwrap(),
record.revision,
)
.await
.unwrap(),
CasOutcome::Won(_)
));
watermark
})
}
#[test]
fn poison_splits_quarantine_and_stall_instead_of_false_success() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("poison:v1", &["good", "bad"]));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held = Held::default();
drive(&mut a, &mut held, "claiming both splits", |h| {
h.splits.len() == 2
});
a.commit(&split_id("good"), &SplitProgress::completed(5, vec![]))
.unwrap();
let mut failures = 0;
let deadline = Instant::now() + support::DEADLINE;
while held.quarantined.is_empty() {
assert!(Instant::now() < deadline, "quarantine never happened");
if held.splits.contains_key("bad") {
a.fail(&split_id("bad"), "undecodable descriptor").unwrap();
held.splits.remove("bad");
failures += 1;
}
held.fold(a.poll().unwrap());
}
assert_eq!(held.quarantined[0].0, "bad");
assert_eq!(
(failures, held.quarantined[0].1),
(4, 4),
"exactly max_attempts failing tenancies ran"
);
drive(&mut a, &mut held, "waiting for the stall verdict", |h| {
h.stalled.is_some()
});
assert_eq!(held.stalled, Some((1, 1)));
assert!(!held.all_complete, "quarantine must block AllComplete");
}
#[test]
fn all_complete_reaches_a_late_standby_that_never_owned_anything() {
let rt = runtime();
let store = store();
let planner = || Box::new(PhasedPlanner::one_final("standby:v1", &["only"]));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "A claiming the split", |h| {
h.splits.len() == 1
});
let mut standby = worker(&store, rt.handle(), Some("worker-s"));
standby.start(planner()).unwrap();
let mut held_s = Held::default();
a.commit(&split_id("only"), &SplitProgress::completed(9, vec![]))
.unwrap();
drive_pair(
(&mut a, &mut held_a),
(&mut standby, &mut held_s),
"completion reaching both the owner and the standby",
|ha, hs| ha.all_complete && hs.all_complete,
);
assert!(held_s.splits.is_empty(), "the standby never owned work");
}
#[test]
fn a_newcomer_is_assigned_a_share_of_a_loaded_fleet() {
let rt = runtime();
let store = store();
let ids = ["s0", "s1", "s2", "s3"];
let planner = || Box::new(PhasedPlanner::one_final("scaleout:v1", &ids));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "worker-a takes the whole plan", |h| {
h.splits.len() == ids.len()
});
support::commit_held(&mut a, &held_a);
let mut b = worker(&store, rt.handle(), Some("worker-b"));
b.start(planner()).unwrap();
let mut held_b = Held::default();
let deadline = Instant::now() + support::DEADLINE;
while !(held_a.splits.len() == 2 && held_b.splits.len() == 2) {
assert!(
Instant::now() < deadline,
"timed out balancing: a={:?} b={:?}",
held_a.splits.keys().collect::<Vec<_>>(),
held_b.splits.keys().collect::<Vec<_>>()
);
held_a.fold(a.poll().unwrap());
held_b.fold(b.poll().unwrap());
support::commit_held(&mut a, &held_a);
support::consent_to_revocations(&mut a, &mut held_a);
std::thread::sleep(support::POLL_INTERVAL);
}
assert!(
held_a.splits.keys().all(|k| !held_b.splits.contains_key(k)),
"the two halves must be disjoint"
);
for (id, (_, progress)) in &held_b.splits {
let watermark = progress
.as_ref()
.unwrap_or_else(|| panic!("{id} arrived without a resume point, so it will replay"))
.watermark;
assert_eq!(
watermark,
support::DRAINED_WATERMARK,
"{id} resumed from the pre-rebalance watermark ({}), not the drained tail — \
the move was forced, not cooperative",
support::BASE_WATERMARK
);
}
}
#[test]
fn a_worker_joining_during_a_revocation_still_converges() {
let rt = runtime();
let store = store();
let ids = ["s0", "s1", "s2", "s3", "s4", "s5"];
let planner = || Box::new(PhasedPlanner::one_final("midjoin:v1", &ids));
let mut a = worker(&store, rt.handle(), Some("worker-a"));
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "worker-a takes the whole plan", |h| {
h.splits.len() == ids.len()
});
support::commit_held(&mut a, &held_a);
let mut b = worker(&store, rt.handle(), Some("worker-b"));
b.start(planner()).unwrap();
let mut held_b = Held::default();
drive(&mut a, &mut held_a, "a revocation is requested", |h| {
!h.revoke_requests.is_empty()
});
let mut c = worker(&store, rt.handle(), Some("worker-c"));
c.start(planner()).unwrap();
let mut held_c = Held::default();
let deadline = Instant::now() + support::DEADLINE;
loop {
let counts = (
held_a.splits.len(),
held_b.splits.len(),
held_c.splits.len(),
);
if counts == (2, 2, 2) {
break;
}
assert!(
Instant::now() < deadline,
"timed out converging: a/b/c = {counts:?}"
);
for (co, held) in [
(&mut a, &mut held_a),
(&mut b, &mut held_b),
(&mut c, &mut held_c),
] {
held.fold(co.poll().unwrap());
support::commit_held(co, held);
support::consent_to_revocations(co, held);
}
std::thread::sleep(support::POLL_INTERVAL);
}
let total: usize = held_a.splits.len() + held_b.splits.len() + held_c.splits.len();
assert_eq!(
total,
ids.len(),
"every split must still be held by someone"
);
}
#[test]
fn a_zero_rebalance_delay_reassigns_immediately() {
assert!(
reassignment_delay(std::time::Duration::ZERO) < LEASE * 3,
"zero delay must not withhold the split"
);
}
#[test]
fn a_departed_workers_splits_are_withheld_for_the_grace_window() {
let withheld = reassignment_delay(LEASE * 5);
assert!(
withheld >= LEASE * 3,
"a grace window must actually delay reassignment, took {withheld:?}"
);
}
fn reassignment_delay(delay: std::time::Duration) -> std::time::Duration {
let rt = runtime();
let clock = support::TestClock::frozen();
let store = support::store_with_clock(clock.clone());
let ids = ["s0", "s1"];
let planner = || Box::new(PhasedPlanner::one_final("delay:v1", &ids));
let mut a = support::worker_rebalance_delay_clock(
&store,
rt.handle(),
Some("worker-a"),
delay,
clock.clone(),
);
let rt_b = runtime();
let mut b = support::worker_rebalance_delay_clock(
&store,
rt_b.handle(),
Some("worker-b"),
delay,
clock.clone(),
);
a.start(planner()).unwrap();
b.start(planner()).unwrap();
let (mut held_a, mut held_b) = (Held::default(), Held::default());
let step = LEASE / 6;
let claim_deadline = Instant::now() + support::DEADLINE;
while !(held_a.splits.len() == 1 && held_b.splits.len() == 1) {
assert!(
Instant::now() < claim_deadline,
"both workers never held a split"
);
clock.advance(step);
std::thread::sleep(support::POLL_INTERVAL);
held_a.fold(a.poll().unwrap());
held_b.fold(b.poll().unwrap());
}
support::commit_held(&mut a, &held_a);
crash(rt_b, b);
let cap = LEASE * 10;
let mut advanced = std::time::Duration::ZERO;
while held_a.splits.len() < 2 {
assert!(advanced < cap, "the survivor never picked up the split");
clock.advance(step);
advanced += step;
std::thread::sleep(support::POLL_INTERVAL);
held_a.fold(a.poll().unwrap());
support::commit_held(&mut a, &held_a);
}
advanced
}
#[test]
fn a_drain_that_never_completes_is_forced_out() {
let rt = runtime();
let store = store();
let ids = ["s0", "s1", "s2", "s3"];
let planner = || Box::new(PhasedPlanner::one_final("forced:v1", &ids));
let mut a = worker_drain_deadline(&store, rt.handle(), Some("worker-a"), LEASE / 4);
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "worker-a takes the whole plan", |h| {
h.splits.len() == ids.len()
});
support::commit_held(&mut a, &held_a);
let mut b = worker_drain_deadline(&store, rt.handle(), Some("worker-b"), LEASE / 4);
b.start(planner()).unwrap();
let mut held_b = Held::default();
let started = Instant::now();
let deadline = Instant::now() + support::DEADLINE;
while held_b.splits.is_empty() {
assert!(
Instant::now() < deadline,
"a declining source blocked the rebalance forever"
);
held_a.fold(a.poll().unwrap());
held_b.fold(b.poll().unwrap());
support::commit_held(&mut a, &held_a);
std::thread::sleep(support::POLL_INTERVAL);
}
assert!(
!held_a.revoke_requests.is_empty(),
"the split should have been asked for before it was taken"
);
assert!(
held_a.splits.len() < ids.len(),
"the forced split must have left worker-a"
);
assert!(
started.elapsed() < LEASE,
"a split moved only after a full lease ({:?}) — that is an expiry takeover, \
not a forced revocation",
started.elapsed()
);
let moved: Vec<&String> = held_b.splits.keys().collect();
assert!(
moved.iter().all(|id| held_a.revoke_requests.contains(id)),
"worker-b holds {moved:?} but the revocations asked for {:?}",
held_a.revoke_requests
);
}
#[test]
fn a_withdrawn_assignment_record_does_not_release_anything() {
let rt = runtime();
let clock = support::TestClock::frozen();
let store = support::store_with_clock(clock.clone());
let ids = ["w0", "w1", "w2", "w3"];
let planner = || Box::new(PhasedPlanner::one_final("withdrawn:v1", &ids));
let mut a = support::worker_with_clock(&store, rt.handle(), Some("worker-a"), clock.clone());
a.start(planner()).unwrap();
let mut held_a = Held::default();
let mut b = support::worker_with_clock(&store, rt.handle(), Some("worker-b"), clock.clone());
b.start(planner()).unwrap();
let mut held_b = Held::default();
support::settle_pair_clocked(
(&mut a, &mut held_a),
(&mut b, &mut held_b),
&clock,
"both workers settle on half the plan",
|x, y| x.splits.len() == 2 && y.splits.len() == 2,
);
let before: Vec<String> = held_b.splits.keys().cloned().collect();
rt.block_on(async {
let outcome = store
.delete(Keyspace::Durable, "assign.worker-b", None)
.await
.expect("delete");
assert!(
matches!(outcome, CasOutcome::Won(_)),
"the record has to have existed, or this tests nothing"
);
});
clock.advance_stepped(LEASE, LEASE / 12, || {
std::thread::sleep(support::POLL_INTERVAL);
for event in a.poll().expect("poll a") {
held_a.fold(vec![event]);
}
for event in b.poll().expect("poll b") {
assert!(
!matches!(event, spate_coordination::CoordinationEvent::Lost { .. }),
"worker-b released a split because its assignment record vanished: {event:?}"
);
held_b.fold(vec![event]);
}
assert!(
held_b.revoke_requests.is_empty(),
"an absent assignment record was read as an instruction to hold nothing: {:?}",
held_b.revoke_requests
);
support::commit_held(&mut a, &held_a);
support::commit_held(&mut b, &held_b);
});
let after: Vec<String> = held_b.splits.keys().cloned().collect();
assert_eq!(
before, after,
"worker-b's working set changed while it had no assignment record"
);
}
#[test]
fn a_reassigned_split_cancels_its_own_revocation() {
let rt = runtime();
let clock = support::TestClock::frozen();
let store = support::store_with_clock(clock.clone());
let ids = ["c0", "c1", "c2", "c3"];
let planner = || Box::new(PhasedPlanner::one_final("cancel:v1", &ids));
let mut a = support::worker_with_clock(&store, rt.handle(), Some("worker-a"), clock.clone());
a.start(planner()).unwrap();
let mut held_a = Held::default();
support::drive_clocked(
&mut a,
&clock,
&mut held_a,
"worker-a takes the plan",
|h| h.splits.len() == ids.len(),
);
support::commit_held(&mut a, &held_a);
let rt_b = runtime();
let mut b = support::worker_with_clock(&store, rt_b.handle(), Some("worker-b"), clock.clone());
b.start(planner()).unwrap();
let mut held_b = Held::default();
let deadline = Instant::now() + support::DEADLINE;
while held_a.revoke_requests.is_empty() {
assert!(
Instant::now() < deadline,
"the leader never revoked anything from worker-a"
);
clock.advance(LEASE / 12);
std::thread::sleep(support::POLL_INTERVAL);
held_a.fold(a.poll().expect("poll a"));
held_b.fold(b.poll().expect("poll b"));
support::commit_held(&mut a, &held_a);
}
let asked = held_a.revoke_requests.clone();
let before: Vec<String> = held_a.splits.keys().cloned().collect();
assert_eq!(
before.len(),
ids.len(),
"worker-a must still hold everything when the window opens: the drain \
has been asked for, not answered"
);
crash(rt_b, b);
rt.block_on(async {
let outcome = store
.delete(Keyspace::Ephemeral, "worker.worker-b", None)
.await
.expect("delete");
assert!(
matches!(outcome, CasOutcome::Won(_)),
"the presence key has to have existed, or this tests nothing"
);
});
clock.advance_stepped(LEASE, LEASE / 12, || {
std::thread::sleep(support::POLL_INTERVAL);
for event in a.poll().expect("poll a") {
if let spate_coordination::CoordinationEvent::Lost { split } = &event {
assert!(
!asked.contains(&split.as_str().to_string()),
"worker-a gave up {split}, which the leader had already assigned \
back to it: the revocation was forced instead of cancelled"
);
}
held_a.fold(vec![event]);
}
support::commit_held(&mut a, &held_a);
});
let after: Vec<String> = held_a.splits.keys().cloned().collect();
assert_eq!(
before, after,
"worker-a's working set moved while every split was assigned to it"
);
}
#[test]
fn a_stalled_cancelled_drain_is_still_released() {
let rt = runtime();
let clock = support::TestClock::frozen();
let store = support::store_with_clock(clock.clone());
let ids = ["w0", "w1", "w2", "w3"];
let planner = || Box::new(PhasedPlanner::one_final("stalled-cancel:v1", &ids));
let mut a = support::worker_with_clock(&store, rt.handle(), Some("worker-a"), clock.clone());
a.start(planner()).unwrap();
let mut held_a = Held::default();
support::drive_clocked(
&mut a,
&clock,
&mut held_a,
"worker-a takes the plan",
|h| h.splits.len() == ids.len(),
);
support::commit_held(&mut a, &held_a);
let rt_b = runtime();
let mut b = support::worker_with_clock(&store, rt_b.handle(), Some("worker-b"), clock.clone());
b.start(planner()).unwrap();
let mut held_b = Held::default();
let deadline = Instant::now() + support::DEADLINE;
while held_a.revoke_requests.is_empty() {
assert!(
Instant::now() < deadline,
"the leader never revoked anything from worker-a"
);
clock.advance(LEASE / 12);
std::thread::sleep(support::POLL_INTERVAL);
held_a.fold(a.poll().expect("poll a"));
held_b.fold(b.poll().expect("poll b"));
support::commit_held(&mut a, &held_a);
}
let asked = held_a.revoke_requests.clone();
crash(rt_b, b);
rt.block_on(async {
let outcome = store
.delete(Keyspace::Ephemeral, "worker.worker-b", None)
.await
.expect("delete");
assert!(
matches!(outcome, CasOutcome::Won(_)),
"the presence key has to have existed, or this tests nothing"
);
});
clock.advance_stepped(LEASE, LEASE / 12, || {
std::thread::sleep(support::POLL_INTERVAL);
for event in a.poll().expect("poll a") {
if let spate_coordination::CoordinationEvent::Lost { split } = &event {
assert!(
!asked.contains(&split.as_str().to_string()),
"worker-a gave up {split} while still committing it: the revocation \
was forced instead of cancelled, so phase 2 would prove nothing"
);
}
held_a.fold(vec![event]);
}
support::commit_held(&mut a, &held_a);
});
assert_eq!(
held_a.splits.len(),
ids.len(),
"worker-a must still hold everything before the stall begins"
);
let mut lost: Option<String> = None;
let deadline = Instant::now() + support::DEADLINE;
while lost
.as_ref()
.is_none_or(|id| !held_a.splits.contains_key(id))
{
assert!(
Instant::now() < deadline,
"the stalled drain was never released: {asked:?} stayed owned with nothing \
reading them"
);
clock.advance(LEASE / 12);
std::thread::sleep(support::POLL_INTERVAL);
for event in a.poll().expect("poll a") {
if let spate_coordination::CoordinationEvent::Lost { split } = &event
&& asked.contains(&split.as_str().to_string())
&& lost.is_none()
{
lost = Some(split.as_str().to_string());
}
held_a.fold(vec![event]);
}
support::commit_held_except(&mut a, &held_a, &asked);
}
let lost = lost.expect("a stalled split was released");
assert!(
asked.contains(&lost),
"the released split must be one the leader had asked for"
);
assert_eq!(
held_a.splits.len(),
ids.len(),
"worker-a must be whole again: the stalled split was released and re-claimed"
);
}
#[test]
fn a_revocation_of_the_last_split_is_not_a_departure() {
let rt = runtime();
let store = store();
let ids = ["r0", "r1"];
let planner = || Box::new(PhasedPlanner::one_final("lastsplit:v1", &ids));
let mut a = worker_max_in_flight(&store, rt.handle(), Some("worker-a"), 1);
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "worker-a takes its one lane", |h| {
h.splits.len() == 1
});
support::commit_held(&mut a, &held_a);
let mut b = worker_max_in_flight(&store, rt.handle(), Some("worker-b"), 1);
b.start(planner()).unwrap();
let mut held_b = Held::default();
let deadline = Instant::now() + support::DEADLINE;
while !(held_a.splits.len() == 1 && held_b.splits.len() == 1) {
assert!(
Instant::now() < deadline,
"the fleet never reached one split each: a={:?} b={:?}",
held_a.splits.keys().collect::<Vec<_>>(),
held_b.splits.keys().collect::<Vec<_>>()
);
held_a.fold(a.poll().expect("poll a"));
held_b.fold(b.poll().expect("poll b"));
support::commit_held(&mut a, &held_a);
support::commit_held(&mut b, &held_b);
support::consent_to_revocations(&mut a, &mut held_a);
std::thread::sleep(support::POLL_INTERVAL);
}
let presence = rt.block_on(async {
store
.get(Keyspace::Ephemeral, "worker.worker-a")
.await
.expect("get")
});
assert!(
presence.is_some(),
"worker-a left the fleet after a revocation took its last split"
);
}
#[test]
fn splits_beyond_the_fleets_lane_budget_wait_in_the_queue() {
let rt = runtime();
let store = store();
let ids = ["q0", "q1", "q2", "q3", "q4"];
let planner = || Box::new(PhasedPlanner::one_final("queue:v1", &ids));
let mut a = worker_max_in_flight(&store, rt.handle(), Some("worker-a"), 1);
a.start(planner()).unwrap();
let mut held_a = Held::default();
drive(&mut a, &mut held_a, "worker-a takes its single lane", |h| {
h.splits.len() == 1
});
let until = Instant::now() + LEASE;
while Instant::now() < until {
held_a.fold(a.poll().expect("poll a"));
assert_eq!(
held_a.splits.len(),
1,
"a one-lane worker took {} splits",
held_a.splits.len()
);
support::commit_held(&mut a, &held_a);
std::thread::sleep(support::POLL_INTERVAL);
}
assert!(
held_a.quarantined.is_empty(),
"queued splits were parked instead of waiting: {:?}",
held_a.quarantined
);
let mut b = worker_max_in_flight(&store, rt.handle(), Some("worker-b"), 1);
b.start(planner()).unwrap();
let mut held_b = Held::default();
drive_pair(
(&mut a, &mut held_a),
(&mut b, &mut held_b),
"the second lane picks up queued work",
|_, y| y.splits.len() == 1,
);
}