use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::thread;
use std::time::{Duration, Instant};
#[path = "../tests/support/gate.rs"]
mod gate;
use std::io;
use std::sync::OnceLock;
use super::{
BLOCKING_KEEP_ALIVE, Handoff, Pool, PoolConfigError, PoolShutdown, SpawnError, block_on,
configure_decision, join_all, lock, prepare, rng, seeded_sweep, serve_pool, spawn_blocking_on,
start_os_thread,
};
use gate::{Gate, Release};
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 = 200;
const EXPIRED: Duration = Duration::from_millis(20);
const GENEROUS: Duration = Duration::from_secs(30);
const WAKE_BOUND: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Scenario {
BurstDrain,
DeadlineThenDrain,
ParkedThenShutdown,
ConfigureAfterUse,
}
const SCENARIOS: [Scenario; 4] = [
Scenario::BurstDrain,
Scenario::DeadlineThenDrain,
Scenario::ParkedThenShutdown,
Scenario::ConfigureAfterUse,
];
impl Scenario {
const fn tag(self) -> u64 {
match self {
Self::BurstDrain => 40,
Self::DeadlineThenDrain => 41,
Self::ParkedThenShutdown => 42,
Self::ConfigureAfterUse => 43,
}
}
}
const ARM_DRAINED: u64 = 50;
const ARM_EXPIRED: u64 = 51;
const ARM_WOKEN: u64 = 52;
const ARM_IN_USE: u64 = 53;
const ARM_CLOSED: u64 = 54;
fn journey(seed: u64) -> u64 {
let mut rng = Rng::new(seed);
let ceiling = rng.below(4).saturating_add(1);
let jobs = rng.below(8);
let [
burst,
deadline_then_drain,
parked_then_shutdown,
configure_after_use,
] = SCENARIOS;
let scenario = match rng.below(SCENARIOS.len()) {
0 => burst,
1 => deadline_then_drain,
2 => parked_then_shutdown,
_ => configure_after_use,
};
let live = jobs.min(ceiling);
let mut trace = initial_trace();
fold(&mut trace, seed);
fold_usize(&mut trace, ceiling);
fold_usize(&mut trace, jobs);
fold(&mut trace, scenario.tag());
let pool = Arc::new(Pool::new(ceiling, start_os_thread));
match scenario {
Scenario::BurstDrain | Scenario::DeadlineThenDrain => {
let gate = Arc::new(Gate::new());
let mut handles = Vec::with_capacity(jobs);
for index in 0..jobs {
let gate = Arc::clone(&gate);
handles.push(spawn_blocking_on(&pool, move || {
gate.pass();
index
}));
}
{
let state = lock(&pool.state);
assert_eq!(
state.live, live,
"seed {seed:#x}: {} threads alive, not the {live} a burst of {jobs} starts at a ceiling of {ceiling}",
state.live
);
assert_eq!(
state.queue.len(),
jobs.saturating_sub(live),
"seed {seed:#x}: the rest of the burst waits its turn"
);
assert_eq!(
state.handles.len(),
live,
"seed {seed:#x}: every started thread is registered to be joined"
);
}
fold_usize(&mut trace, live);
fold_usize(&mut trace, jobs.saturating_sub(live));
if scenario == Scenario::DeadlineThenDrain && jobs > 0 {
assert_eq!(
pool.shutdown(EXPIRED),
PoolShutdown::DeadlineExceeded {
joined: 0,
running: live,
queued: jobs.saturating_sub(live)
},
"seed {seed:#x}: the deadline names the threads still working and the jobs still waiting"
);
fold(&mut trace, ARM_EXPIRED);
fold_usize(&mut trace, live);
fold_usize(&mut trace, jobs.saturating_sub(live));
}
gate.release();
assert_eq!(
pool.shutdown(GENEROUS),
PoolShutdown::Drained { threads: live },
"seed {seed:#x}: shutdown let every queued and running job finish and joined every thread"
);
fold(&mut trace, ARM_DRAINED);
for (index, value) in block_on(join_all(handles)).into_iter().enumerate() {
assert_eq!(
value, index,
"seed {seed:#x}: job {index} returned another job's value"
);
fold_usize(&mut trace, value);
}
}
Scenario::ParkedThenShutdown => {
assert_eq!(block_on(spawn_blocking_on(&pool, || 5u32)), 5);
thread::park_timeout(Duration::from_millis(30));
let began = Instant::now();
let report = pool.shutdown(WAKE_BOUND.saturating_add(BLOCKING_KEEP_ALIVE));
let waited = began.elapsed();
assert_eq!(
report,
PoolShutdown::Drained { threads: 1 },
"seed {seed:#x}: a parked thread was woken and joined"
);
assert!(
waited < WAKE_BOUND,
"seed {seed:#x}: a parked thread was waited out ({waited:?}) rather than woken, \
against a {BLOCKING_KEEP_ALIVE:?} keep-alive"
);
fold(&mut trace, ARM_WOKEN);
}
Scenario::ConfigureAfterUse => {
assert_eq!(block_on(spawn_blocking_on(&pool, || 6u32)), 6);
assert_eq!(
configure_decision(&pool, ceiling.saturating_add(1)),
Err(PoolConfigError::InUse {
requested: ceiling.saturating_add(1),
running: ceiling
}),
"seed {seed:#x}: a running pool refuses another ceiling by name"
);
assert_eq!(
configure_decision(&pool, ceiling),
Err(PoolConfigError::InUse {
requested: ceiling,
running: ceiling
}),
"seed {seed:#x}: the same ceiling is still refused as unconfigured"
);
fold(&mut trace, ARM_IN_USE);
assert_eq!(
pool.shutdown(GENEROUS),
PoolShutdown::Drained { threads: 1 },
"seed {seed:#x}: a pool that refused a configure still drains"
);
fold(&mut trace, ARM_DRAINED);
}
}
let ran = Arc::new(AtomicBool::new(false));
let witness = Arc::clone(&ran);
let (refused, work) = prepare(move || witness.store(true, Ordering::SeqCst));
assert!(
matches!(pool.submit(work, Some(8)), Err(SpawnError::Shutdown)),
"seed {seed:#x}: a drained pool refuses admission with the typed reason"
);
drop(refused);
let handle = spawn_blocking_on(&pool, || 9u32);
let payload = catch_unwind(AssertUnwindSafe(|| block_on(handle)))
.err()
.map(|payload| payload.downcast::<SpawnError>());
assert!(
matches!(payload, Some(Ok(ref error)) if matches!(**error, SpawnError::Shutdown)),
"seed {seed:#x}: the never-refusing entry point fails its awaiter with the refusal"
);
assert!(
!ran.load(Ordering::SeqCst),
"seed {seed:#x}: a refused job never ran"
);
fold(&mut trace, ARM_CLOSED);
trace
}
const CYCLES: usize = 6;
const CYCLE_KEEP_ALIVE: Duration = Duration::from_millis(1);
const IDLE_POLLS: usize = 500;
fn wait_until_empty(pool: &Arc<Pool>, seed: u64) {
for _ in 0..IDLE_POLLS {
if lock(&pool.state).live == 0 {
return;
}
thread::park_timeout(Duration::from_millis(1));
}
let state = lock(&pool.state);
assert!(
state.live == 0,
"seed {seed:#x}: {} threads were still alive after an idle period of {} polls at a \
{CYCLE_KEEP_ALIVE:?} keep-alive",
state.live,
IDLE_POLLS
);
}
fn assert_handles_accounted(pool: &Arc<Pool>, where_: &str, seed: u64) {
let state = lock(&pool.state);
let live = state.live;
let held = state.handles.len();
let unreturned = state
.handles
.iter()
.filter(|handle| !handle.is_finished())
.count();
assert!(
unreturned >= live,
"seed {seed:#x} {where_}: {live} live threads but only {unreturned} unreturned handles \
of {held}; a live thread cannot have returned, since it decrements `live` first"
);
}
fn cycles(seed: u64) -> u64 {
let mut rng = Rng::new(seed);
let ceiling = rng.below(4).saturating_add(1);
let cycles = rng.below(CYCLES).saturating_add(2);
let mut trace = initial_trace();
fold(&mut trace, seed);
fold_usize(&mut trace, ceiling);
fold_usize(&mut trace, cycles);
let pool = Arc::new(Pool::tuned(ceiling, start_os_thread, CYCLE_KEEP_ALIVE));
for cycle in 0..cycles {
let burst = ceiling.saturating_add(1);
let mut handles = Vec::with_capacity(burst);
for index in 0..burst {
handles.push(spawn_blocking_on(&pool, move || index));
assert_handles_accounted(&pool, &format!("cycle {cycle} submit {index}"), seed);
}
for (index, value) in block_on(join_all(handles)).into_iter().enumerate() {
assert_eq!(
value, index,
"seed {seed:#x}: cycle {cycle} job {index} ran elsewhere"
);
fold_usize(&mut trace, value);
}
assert_handles_accounted(&pool, &format!("cycle {cycle} drained"), seed);
wait_until_empty(&pool, seed);
let held = lock(&pool.state).handles.len();
assert!(
held <= ceiling,
"seed {seed:#x}: cycle {cycle} left {held} join handles for a ceiling of {ceiling}; \
a thread that exits on its keep-alive must be reaped, not held"
);
fold(&mut trace, u64::from(held <= ceiling));
assert_handles_accounted(&pool, &format!("cycle {cycle} idle"), seed);
wait_until_all_returned(&pool, seed);
let head = spawn_blocking_on(&pool, || 0u32);
assert_eq!(block_on(head), 0);
let after_reap = lock(&pool.state).handles.len();
assert_eq!(
after_reap, 1,
"seed {seed:#x}: cycle {cycle} started a thread holding {after_reap} handles; \
the ones that had already left were not joined"
);
fold_usize(&mut trace, after_reap);
assert_handles_accounted(&pool, &format!("cycle {cycle} after the reap"), seed);
}
let report = pool.shutdown(GENEROUS);
let state = lock(&pool.state);
assert!(
matches!(report, PoolShutdown::Drained { .. }) && state.handles.is_empty(),
"seed {seed:#x}: the cycles ended holding handles the pool never joined: {report:?}"
);
trace
}
static LINGER: OnceLock<Arc<Release>> = OnceLock::new();
fn start_lingering(pool: Arc<Pool>, first: Handoff) -> io::Result<super::ThreadHandle> {
let hold = Arc::clone(LINGER.get_or_init(|| Arc::new(Release::new())));
let held = hold.generation();
thread::Builder::new()
.name("lgwks-blocking".into())
.spawn(move || {
serve_pool(&pool, first);
HELD_IN_WINDOW.fetch_add(1, Ordering::SeqCst);
hold.wait_past(held);
HELD_IN_WINDOW.fetch_sub(1, Ordering::SeqCst);
})
}
static HELD_IN_WINDOW: AtomicUsize = AtomicUsize::new(0);
fn wait_until_held(count: usize, seed: u64) {
for _ in 0..IDLE_POLLS {
if HELD_IN_WINDOW.load(Ordering::SeqCst) == count {
return;
}
thread::park_timeout(Duration::from_millis(1));
}
let held = HELD_IN_WINDOW.load(Ordering::SeqCst);
assert!(
held == count,
"seed {seed:#x}: {held} threads were held in the mid-exit window, not {count}, after \
{} polls",
IDLE_POLLS
);
}
fn assert_handles_equal_live_plus_held(pool: &Arc<Pool>, where_: &str, seed: u64) {
let state = lock(&pool.state);
let live = state.live;
let held = state.handles.len();
let in_window = HELD_IN_WINDOW.load(Ordering::SeqCst);
assert_eq!(
held,
live.saturating_add(in_window),
"seed {seed:#x} {where_}: {held} handles are not {live} live plus {in_window} threads \
held between leaving the accounting and returning"
);
}
fn wait_until_all_returned(pool: &Arc<Pool>, seed: u64) {
for _ in 0..IDLE_POLLS {
let returned = lock(&pool.state)
.handles
.iter()
.all(super::ThreadHandle::is_finished);
if returned {
return;
}
thread::park_timeout(Duration::from_millis(1));
}
let held = lock(&pool.state).handles.len();
assert!(
held == 0,
"seed {seed:#x}: {held} registered threads had not all returned after {} polls",
IDLE_POLLS
);
}
fn mid_exit(seed: u64) -> u64 {
let mut trace = initial_trace();
fold(&mut trace, seed);
let hold = Arc::clone(LINGER.get_or_init(|| Arc::new(Release::new())));
let pool = Arc::new(Pool::tuned(1, start_lingering, CYCLE_KEEP_ALIVE));
let first = spawn_blocking_on(&pool, || 1u32);
assert_eq!(block_on(first), 1, "seed {seed:#x}: the first job must run");
wait_until_empty(&pool, seed);
wait_until_held(1, seed);
{
let state = lock(&pool.state);
assert_eq!(
(state.live, state.handles.len()),
(0, 1),
"seed {seed:#x}: the fixture must hold one thread mid-exit — gone from the \
accounting, still unjoined"
);
}
assert_handles_equal_live_plus_held(&pool, "in the mid-exit window", seed);
fold(&mut trace, 60);
let second = spawn_blocking_on(&pool, || 2u32);
{
let state = lock(&pool.state);
assert_eq!(
(state.live, state.handles.len()),
(1, 2),
"seed {seed:#x}: a start inside the mid-exit window keeps the departing handle \
and adds its own, so the list outgrows the ceiling"
);
}
assert_eq!(
block_on(second),
2,
"seed {seed:#x}: the second job must run"
);
wait_until_held(2, seed);
assert_handles_equal_live_plus_held(&pool, "after the start in the window", seed);
fold(&mut trace, 61);
hold.release();
wait_until_empty(&pool, seed);
wait_until_all_returned(&pool, seed);
let third = spawn_blocking_on(&pool, || 3u32);
assert_eq!(block_on(third), 3, "seed {seed:#x}: the third job must run");
{
let state = lock(&pool.state);
assert_eq!(
state.handles.len(),
1,
"seed {seed:#x}: a start after the threads returned reaps every departed handle"
);
}
wait_until_empty(&pool, seed);
wait_until_held(1, seed);
assert_handles_equal_live_plus_held(&pool, "after the reap", seed);
fold(&mut trace, 62);
hold.release();
wait_until_all_returned(&pool, seed);
let report = pool.shutdown(GENEROUS);
let state = lock(&pool.state);
assert_eq!(
report,
PoolShutdown::Drained { threads: 1 },
"seed {seed:#x}: the lingering thread returned, so the shutdown joined it"
);
assert!(
state.handles.is_empty(),
"seed {seed:#x}: the window case ended holding handles the pool never joined"
);
trace
}
#[test]
fn sim_a_start_inside_the_mid_exit_window_keeps_both_handles_and_the_next_reap_takes_one() {
for seed in SWEEP_SEEDS {
assert_same_seed_replays(mid_exit, seed);
}
let [first, second, ..] = SWEEP_SEEDS;
assert_distinct_seeds_diverge(mid_exit, first, second);
}
#[test]
fn sim_every_scenario_drains_joins_and_never_loses_a_job() {
let [first_seed, ..] = SWEEP_SEEDS;
let mut state = first_seed;
for _ in 0..SEEDS {
journey(next_seed(&mut state));
}
}
#[test]
fn sim_many_burst_and_idle_cycles_never_outgrow_the_ceiling_in_join_handles() {
let [first_seed, ..] = SWEEP_SEEDS;
let mut state = first_seed;
for _ in 0..SEEDS {
cycles(next_seed(&mut state));
}
}
#[test]
fn sim_a_seed_replays_its_pool_lifetime_trace() {
for seed in SWEEP_SEEDS {
assert_same_seed_replays(journey, seed);
assert_same_seed_replays(cycles, seed);
assert_same_seed_replays(mid_exit, seed);
}
}
#[test]
fn sim_distinct_seeds_diverge_in_their_pool_lifetime_trace() {
let [first, second, ..] = SWEEP_SEEDS;
assert_distinct_seeds_diverge(journey, first, second);
assert_distinct_seeds_diverge(cycles, first, second);
}