#![allow(dead_code, unreachable_pub)]
use spate_coordination::store::memory::MemoryStore;
use spate_coordination::{
Clock, CoordinationConfig, CoordinationError, CoordinationEvent, MemoryCoordinator,
PlanContext, PlanFinality, PlannedSplit, SplitCoordinator, SplitId, SplitPlan, SplitPlanner,
SplitProgress, SplitSpec, StoreCoordinator,
};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
pub const LEASE: Duration = Duration::from_millis(1500);
pub const POLL_INTERVAL: Duration = Duration::from_millis(5);
pub const DEADLINE: Duration = Duration::from_secs(20);
pub fn config(instance_id: Option<&str>) -> CoordinationConfig {
CoordinationConfig {
lease_duration: LEASE,
op_timeout: Duration::from_millis(200),
instance_id: instance_id.map(str::to_string),
replan_interval: LEASE,
reconcile_interval: Duration::from_millis(300),
drain_deadline: LEASE / 2,
rebalance_delay: Duration::ZERO,
..CoordinationConfig::default()
}
}
pub fn runtime() -> tokio::runtime::Runtime {
tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.expect("test runtime")
}
pub fn store() -> MemoryStore {
MemoryStore::new(LEASE)
}
pub use spate_coordination::clock::TestClock;
pub fn store_with_clock(clock: Arc<dyn Clock>) -> MemoryStore {
MemoryStore::with_clock(LEASE, clock)
}
pub fn worker_with_clock(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
clock: Arc<dyn Clock>,
) -> MemoryCoordinator {
StoreCoordinator::with_clock(store.clone(), config(instance_id), io.clone(), None, clock)
.expect("coordinator")
}
pub fn worker(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
) -> MemoryCoordinator {
StoreCoordinator::new(store.clone(), config(instance_id), io.clone(), None)
.expect("coordinator")
}
pub fn worker_drain_deadline(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
drain_deadline: Duration,
) -> MemoryCoordinator {
let mut config = config(instance_id);
config.drain_deadline = drain_deadline;
StoreCoordinator::new(store.clone(), config, io.clone(), None).expect("coordinator")
}
pub fn worker_rebalance_delay(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
rebalance_delay: Duration,
) -> MemoryCoordinator {
let mut config = config(instance_id);
config.rebalance_delay = rebalance_delay;
StoreCoordinator::new(store.clone(), config, io.clone(), None).expect("coordinator")
}
pub fn worker_rebalance_delay_clock(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
rebalance_delay: Duration,
clock: Arc<dyn Clock>,
) -> MemoryCoordinator {
let mut config = config(instance_id);
config.rebalance_delay = rebalance_delay;
StoreCoordinator::with_clock(store.clone(), config, io.clone(), None, clock)
.expect("coordinator")
}
pub fn worker_max_in_flight(
store: &MemoryStore,
io: &tokio::runtime::Handle,
instance_id: Option<&str>,
max_in_flight: u32,
) -> MemoryCoordinator {
let mut config = config(instance_id);
config.max_in_flight = max_in_flight;
StoreCoordinator::new(store.clone(), config, io.clone(), None).expect("coordinator")
}
pub struct PhasedPlanner {
pub fingerprint: String,
pub phases: Vec<(Vec<PlannedSplit>, PlanFinality)>,
}
impl PhasedPlanner {
pub fn one_final(fingerprint: &str, ids: &[&str]) -> PhasedPlanner {
PhasedPlanner {
fingerprint: fingerprint.to_string(),
phases: vec![(splits(ids), PlanFinality::Final)],
}
}
}
impl SplitPlanner for PhasedPlanner {
fn fingerprint(&self) -> String {
self.fingerprint.clone()
}
fn plan(&mut self, ctx: PlanContext<'_>) -> Result<SplitPlan, CoordinationError> {
let index: usize = ctx
.planner_state
.map(|bytes| String::from_utf8_lossy(bytes).parse().expect("cursor"))
.unwrap_or(0);
let (splits, finality) = match self.phases.get(index) {
Some((splits, finality)) => (splits.clone(), *finality),
None => (
Vec::new(),
self.phases.last().map_or(PlanFinality::Open, |(_, f)| *f),
),
};
let next = (index + 1).min(self.phases.len());
Ok(SplitPlan::new(splits, finality).with_planner_state(next.to_string().into_bytes()))
}
}
pub fn splits(ids: &[&str]) -> Vec<PlannedSplit> {
ids.iter()
.map(|id| {
PlannedSplit::new(SplitSpec::new(
SplitId::new(*id).expect("test split id"),
format!("descriptor:{id}").into_bytes(),
))
})
.collect()
}
pub fn split_id(id: &str) -> SplitId {
SplitId::new(id).expect("test split id")
}
#[derive(Default)]
pub struct Held {
pub splits: BTreeMap<String, (u64, Option<SplitProgress>)>,
pub all_complete: bool,
pub stalled: Option<(u64, u64)>,
pub quarantined: Vec<(String, u32)>,
pub revoke_requests: Vec<String>,
}
impl Held {
pub fn fold(&mut self, events: Vec<CoordinationEvent>) {
for event in events {
match event {
CoordinationEvent::Gained {
split,
epoch,
progress,
} => {
self.splits
.insert(split.id.as_str().to_string(), (epoch.0, progress));
}
CoordinationEvent::Lost { split } => {
self.splits.remove(split.as_str());
}
CoordinationEvent::Quarantined { split, attempts } => {
self.splits.remove(split.as_str());
self.quarantined
.push((split.as_str().to_string(), attempts));
}
CoordinationEvent::RevokeRequested { split } => {
let id = split.as_str().to_string();
if !self.revoke_requests.contains(&id) {
self.revoke_requests.push(id);
}
}
CoordinationEvent::AllComplete => self.all_complete = true,
CoordinationEvent::Stalled {
completed,
quarantined,
} => self.stalled = Some((completed, quarantined)),
_ => {}
}
}
}
}
pub fn drive(
coordinator: &mut impl SplitCoordinator,
held: &mut Held,
what: &str,
mut done: impl FnMut(&Held) -> bool,
) {
let deadline = Instant::now() + DEADLINE;
while !done(held) {
assert!(Instant::now() < deadline, "timed out: {what}");
let events = coordinator
.poll()
.unwrap_or_else(|e| panic!("poll failed while {what}: {e}"));
if events.is_empty() {
std::thread::sleep(POLL_INTERVAL);
}
held.fold(events);
}
}
pub fn drive_clocked(
coordinator: &mut impl SplitCoordinator,
clock: &TestClock,
held: &mut Held,
what: &str,
mut done: impl FnMut(&Held) -> bool,
) {
let step = LEASE / 12;
let deadline = Instant::now() + DEADLINE;
while !done(held) {
assert!(Instant::now() < deadline, "timed out: {what}");
clock.advance(step);
std::thread::sleep(POLL_INTERVAL);
held.fold(
coordinator
.poll()
.unwrap_or_else(|e| panic!("poll failed while {what}: {e}")),
);
}
}
pub fn wait_until(timeout: Duration, what: &str, mut check: impl FnMut() -> bool) {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if check() {
return;
}
std::thread::sleep(POLL_INTERVAL);
}
panic!("timed out after {timeout:?} waiting for: {what}");
}
pub const BASE_WATERMARK: i64 = 1;
pub const DRAINED_WATERMARK: i64 = 42;
pub fn consent_to_revocations(coordinator: &mut impl SplitCoordinator, held: &mut Held) {
for id in std::mem::take(&mut held.revoke_requests) {
if !held.splits.contains_key(&id) {
continue; }
let watermark = held.splits[&id]
.1
.as_ref()
.map_or(DRAINED_WATERMARK, |p| p.watermark.max(DRAINED_WATERMARK));
let _ = coordinator.commit(&split_id(&id), &SplitProgress::new(watermark, vec![]));
let _ = coordinator.release_drained(&[split_id(&id)]);
held.splits.remove(&id);
}
}
pub fn commit_held_except(coordinator: &mut impl SplitCoordinator, held: &Held, skip: &[String]) {
for (id, (_, carried)) in &held.splits {
if skip.contains(id) {
continue;
}
let watermark = carried
.as_ref()
.map_or(BASE_WATERMARK, |p| p.watermark.max(BASE_WATERMARK));
match coordinator.commit(&split_id(id), &SplitProgress::new(watermark, vec![])) {
Ok(()) => {}
Err(e) if e.kind == spate_coordination::CoordinationErrorKind::Fenced => {}
Err(e) => panic!("commit failed: {e}"),
}
}
}
pub fn commit_held(coordinator: &mut impl SplitCoordinator, held: &Held) {
for (id, (_, carried)) in &held.splits {
let watermark = carried
.as_ref()
.map_or(BASE_WATERMARK, |p| p.watermark.max(BASE_WATERMARK));
match coordinator.commit(&split_id(id), &SplitProgress::new(watermark, vec![])) {
Ok(()) => {}
Err(e) if e.kind == spate_coordination::CoordinationErrorKind::Fenced => {}
Err(e) => panic!("commit failed: {e}"),
}
}
}
pub fn drive_pair<C: SplitCoordinator>(
a: (&mut C, &mut Held),
b: (&mut C, &mut Held),
what: &str,
mut done: impl FnMut(&Held, &Held) -> bool,
) {
let deadline = Instant::now() + DEADLINE;
while !done(a.1, b.1) {
assert!(Instant::now() < deadline, "timed out: {what}");
for (coordinator, held) in [(&mut *a.0, &mut *a.1), (&mut *b.0, &mut *b.1)] {
let events = coordinator
.poll()
.unwrap_or_else(|e| panic!("poll failed while {what}: {e}"));
if events.is_empty() {
std::thread::sleep(POLL_INTERVAL);
}
held.fold(events);
}
}
}
pub fn settle_pair_clocked<C: SplitCoordinator>(
a: (&mut C, &mut Held),
b: (&mut C, &mut Held),
clock: &TestClock,
what: &str,
mut done: impl FnMut(&Held, &Held) -> bool,
) {
let step = LEASE / 12;
let deadline = Instant::now() + DEADLINE;
let round = |a: (&mut C, &mut Held), b: (&mut C, &mut Held), consent: bool| {
clock.advance(step);
std::thread::sleep(POLL_INTERVAL);
for (coordinator, held) in [(a.0, a.1), (b.0, b.1)] {
held.fold(
coordinator
.poll()
.unwrap_or_else(|e| panic!("poll failed while {what}: {e}")),
);
commit_held(coordinator, held);
if consent {
consent_to_revocations(coordinator, held);
}
}
};
loop {
assert!(Instant::now() < deadline, "timed out: {what}");
if !done(a.1, b.1) {
round((&mut *a.0, &mut *a.1), (&mut *b.0, &mut *b.1), true);
continue;
}
for _ in 0..QUIET_ROUNDS {
round((&mut *a.0, &mut *a.1), (&mut *b.0, &mut *b.1), false);
}
let disturbed =
!a.1.revoke_requests.is_empty() || !b.1.revoke_requests.is_empty() || !done(a.1, b.1);
if !disturbed {
return;
}
for (coordinator, held) in [(&mut *a.0, &mut *a.1), (&mut *b.0, &mut *b.1)] {
consent_to_revocations(coordinator, held);
}
}
}
pub const QUIET_ROUNDS: usize = 4;
pub fn crash<C: SplitCoordinator + 'static>(runtime: tokio::runtime::Runtime, coordinator: C) {
runtime.shutdown_background();
std::mem::forget(coordinator);
}