use super::{Admitted, PoolState, Refused, Woken};
use super::{rng, seeded_sweep};
use rng::Rng;
use seeded_sweep::{
SWEEP_SEEDS, assert_distinct_seeds_diverge, assert_same_seed_replays, fold, fold_usize,
initial_trace, next_seed,
};
const SEEDS: usize = 2_000;
const STEPS: usize = 2_000;
const DRAIN_STEPS: usize = 1_000_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Thread {
Free,
Running(usize),
Waiting,
Waking {
timed_out: bool,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Fate {
Admitted,
AtCapacity,
NoThread,
Draining,
}
const CONFIGURE_INVALID: u64 = 20;
const CONFIGURE_IN_FORCE: u64 = 21;
const CONFIGURE_ALREADY: u64 = 22;
const CONFIGURE_IN_USE: u64 = 23;
struct Sim {
state: PoolState<usize>,
threads: Vec<Thread>,
ceiling: usize,
limit: usize,
configured_at_birth: bool,
drain_at: usize,
fates: Vec<Fate>,
runs: Vec<u32>,
rng: Rng,
trace: u64,
seed: u64,
step: usize,
}
impl Sim {
fn new(seed: u64) -> Self {
let mut rng = Rng::new(seed);
let ceiling = rng.below(4).saturating_add(1);
let limit = rng.below(8).saturating_add(1);
let configured_at_birth = rng.below(4) == 0;
let drain_at = rng.below(STEPS);
let mut trace = initial_trace();
fold_usize(&mut trace, ceiling);
fold_usize(&mut trace, limit);
fold(&mut trace, u64::from(configured_at_birth));
fold_usize(&mut trace, drain_at);
Self {
state: PoolState::new(),
threads: Vec::new(),
ceiling,
limit,
configured_at_birth,
drain_at,
fates: Vec::new(),
runs: Vec::new(),
rng,
trace,
seed,
step: 0,
}
}
fn at(&self) -> String {
format!(
"seed {:#018x} step {} (ceiling {}, limit {}, draining {}, threads {:?}, queued {}, live {}, idle {}, wakeups {})",
self.seed,
self.step,
self.ceiling,
self.limit,
self.state.draining,
self.threads,
self.state.queue.len(),
self.state.live,
self.state.idle,
self.state.wakeups,
)
}
fn submit(&mut self) {
let job = self.fates.len();
self.fates.push(Fate::Admitted);
self.runs.push(0);
let bounded = self.rng.below(4) != 0;
let limit = bounded.then_some(self.limit);
let queued = self.state.queue.len();
let admitted = self.state.admit(job, limit, self.ceiling);
fold_usize(&mut self.trace, 1);
match admitted {
Err(Refused::AtCapacity(back)) => {
assert!(
back == job && bounded && queued >= self.limit,
"refused job {back} with {queued} queued: {}",
self.at()
);
if let Some(fate) = self.fates.get_mut(job) {
*fate = Fate::AtCapacity;
}
fold_usize(&mut self.trace, 2);
}
Err(Refused::Draining(back)) => {
assert!(
back == job && self.state.draining,
"job {back} was refused as draining while the pool still admits: {}",
self.at()
);
if let Some(fate) = self.fates.get_mut(job) {
*fate = Fate::Draining;
}
fold_usize(&mut self.trace, 8);
}
Ok(Admitted::Woke) => {
let waiting: Vec<usize> = self
.threads
.iter()
.enumerate()
.filter(|&(_, thread)| *thread == Thread::Waiting)
.map(|(index, _)| index)
.collect();
let woken = waiting
.get(self.rng.below(waiting.len()))
.and_then(|&index| self.threads.get_mut(index));
if let Some(thread) = woken {
*thread = Thread::Waking { timed_out: false };
}
fold_usize(&mut self.trace, 3);
}
Ok(Admitted::Queued) => fold_usize(&mut self.trace, 4),
Ok(Admitted::Start(first)) => {
assert!(
first == job,
"admission started job {first}, not {job}: {}",
self.at()
);
if self.rng.below(8) == 0 {
let live = self.state.live;
match self.state.start_failed(first) {
None => assert!(live > 0, "a job was left to no thread: {}", self.at()),
Some(back) => {
assert!(
back == job && live == 0,
"a failed start returned job {back} with {live} live: {}",
self.at()
);
if let Some(fate) = self.fates.get_mut(job) {
*fate = Fate::NoThread;
}
}
}
fold_usize(&mut self.trace, 5);
} else {
self.state.started();
if let Some(runs) = self.runs.get_mut(first) {
*runs = runs.saturating_add(1);
}
self.threads.push(Thread::Running(first));
fold_usize(&mut self.trace, 6);
}
}
}
}
fn configure(&mut self) {
let proposed = self.rng.below(9);
let before = (
self.state.queue.len(),
self.state.live,
self.state.idle,
self.state.wakeups,
);
let arm = if proposed < 1 {
CONFIGURE_INVALID
} else if self.configured_at_birth {
if proposed == self.ceiling {
CONFIGURE_IN_FORCE
} else {
CONFIGURE_ALREADY
}
} else {
CONFIGURE_IN_USE
};
assert_eq!(
before,
(
self.state.queue.len(),
self.state.live,
self.state.idle,
self.state.wakeups
),
"a configure attempt moved work it must not touch: {}",
self.at()
);
fold(&mut self.trace, arm);
}
fn advance(&mut self, index: usize, settling: bool) {
let Some(&thread) = self.threads.get(index) else {
return;
};
let next = match thread {
Thread::Free => match self.state.take() {
Some(job) => {
if let Some(runs) = self.runs.get_mut(job) {
*runs = runs.saturating_add(1);
}
fold_usize(&mut self.trace, job);
Some(Thread::Running(job))
}
None => {
if self.state.park_or_retire() {
None
} else {
Some(Thread::Waiting)
}
}
},
Thread::Running(_) => Some(Thread::Free),
Thread::Waiting => Some(Thread::Waking {
timed_out: settling || self.rng.below(3) != 0,
}),
Thread::Waking { timed_out } => match self.state.woken(timed_out) {
Woken::Resume => Some(Thread::Free),
Woken::Wait => match self.state.drain_wake() {
Woken::Resume => Some(Thread::Free),
Woken::Exit => None,
Woken::Wait => Some(Thread::Waiting),
},
Woken::Exit => None,
},
};
fold_usize(&mut self.trace, index.saturating_add(16));
match next {
Some(next) => {
if let Some(slot) = self.threads.get_mut(index) {
*slot = next;
}
}
None => {
self.threads.swap_remove(index);
}
}
}
fn check(&self) {
assert_eq!(
self.state.live,
self.threads.len(),
"live threads counted: {}",
self.at()
);
assert!(
self.state.live <= self.ceiling,
"over the ceiling: {}",
self.at()
);
let parked = self
.threads
.iter()
.filter(|thread| matches!(thread, Thread::Waiting | Thread::Waking { .. }))
.count();
assert_eq!(
parked,
self.state.idle.saturating_add(self.state.wakeups),
"parked threads are idle or claimed: {}",
self.at()
);
let working = self
.threads
.iter()
.any(|thread| matches!(thread, Thread::Free | Thread::Running(_)));
assert!(
self.state.queue.is_empty() || working || self.state.wakeups > 0,
"a queued job has no thread that will reach it: {}",
self.at()
);
assert!(
self.runs.iter().all(|&runs| runs <= 1),
"a job ran twice: {}",
self.at()
);
}
fn run(mut self) -> u64 {
for step in 0..STEPS {
self.step = step;
if step == self.drain_at {
self.state.begin_drain();
for thread in &mut self.threads {
if *thread == Thread::Waiting {
*thread = Thread::Waking { timed_out: false };
}
}
fold_usize(&mut self.trace, 9);
}
let pick = self.rng.below(self.threads.len().saturating_add(3));
match pick {
0 => self.submit(),
1 => self.configure(),
index => self.advance(index.saturating_sub(2), false),
}
self.check();
}
let mut settled = 0usize;
while !self.threads.is_empty() {
assert!(
settled < DRAIN_STEPS,
"the pool never settles: {}",
self.at()
);
settled = settled.saturating_add(1);
self.step = STEPS.saturating_add(settled);
let index = self.rng.below(self.threads.len());
self.advance(index, true);
self.check();
}
assert!(
self.state.draining,
"shutdown began at step {} but the pool never closed admission: {}",
self.drain_at,
self.at()
);
assert!(
self.state.queue.is_empty()
&& self.state.idle == 0
&& self.state.wakeups == 0
&& self.state.live == 0,
"a settled pool holds nothing: {}",
self.at()
);
for (job, (&fate, &runs)) in self.fates.iter().zip(&self.runs).enumerate() {
let expected = u32::from(fate == Fate::Admitted);
assert_eq!(
runs,
expected,
"job {job} ({fate:?}) ran {runs} times: {}",
self.at()
);
fold(&mut self.trace, u64::from(runs));
fold(
&mut self.trace,
match fate {
Fate::Admitted => 30,
Fate::AtCapacity => 31,
Fate::NoThread => 32,
Fate::Draining => 33,
},
);
}
self.trace
}
}
fn sweep(seed: u64) -> u64 {
Sim::new(seed).run()
}
#[test]
fn sim_every_admitted_job_runs_once_under_every_interleaving() {
let [first_seed, ..] = SWEEP_SEEDS;
let mut state = first_seed;
for _ in 0..SEEDS {
sweep(next_seed(&mut state));
}
}
#[test]
fn sim_a_seed_replays_its_trace_and_distinct_seeds_diverge() {
for seed in SWEEP_SEEDS {
assert_same_seed_replays(sweep, seed);
}
let [first, second, ..] = SWEEP_SEEDS;
assert_distinct_seeds_diverge(sweep, first, second);
}