use std::collections::VecDeque;
use std::io;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::mpsc::{Receiver, Sender, SyncSender, channel, sync_channel};
use std::sync::{Arc, LazyLock, Mutex};
use std::thread::{self, JoinHandle};
use crate::runtime::log::{ErrlogSevEnum, errlog_sev_printf};
use crate::runtime::task::{InheritedRuntime, StackSizeClass, ThreadPriority, enter_ioc_thread};
#[derive(Clone, Copy, Debug)]
pub struct WorkerRole {
pub suffix: &'static str,
pub stack: StackSizeClass,
pub priority: ThreadPriority,
}
enum Assignment {
Joinable {
body: Box<dyn FnOnce() + Send + 'static>,
ambient: InheritedRuntime,
done: SyncSender<thread::Result<()>>,
},
Detached {
body: Box<dyn FnOnce() + Send + 'static>,
ambient: InheritedRuntime,
label: String,
},
Stop,
}
struct SetState {
leased: bool,
running: usize,
parked: bool,
live_workers: usize,
}
impl SetState {
fn became_free(&mut self) -> bool {
if !self.leased && self.running == 0 && !self.parked {
self.parked = true;
true
} else {
false
}
}
}
struct SetHandle {
index: usize,
senders: Vec<Sender<Assignment>>,
dead: AtomicBool,
state: Mutex<SetState>,
joins: Mutex<Vec<JoinHandle<()>>>,
}
struct Registry {
idle: VecDeque<Arc<SetHandle>>,
all: Vec<Arc<SetHandle>>,
created: usize,
stopping: bool,
}
struct PoolInner {
roster: Box<[WorkerRole]>,
name_prefix: &'static str,
capacity: usize,
set_reservation: usize,
reservation: &'static Reservation,
materialise: fn(&Mutex<SetState>) -> bool,
reg: Mutex<Registry>,
}
impl PoolInner {
fn lock(&self) -> std::sync::MutexGuard<'_, Registry> {
self.reg.lock().unwrap_or_else(|e| e.into_inner())
}
}
fn lock_set(set: &SetHandle) -> std::sync::MutexGuard<'_, SetState> {
set.state.lock().unwrap_or_else(|e| e.into_inner())
}
fn lock_joins(set: &SetHandle) -> std::sync::MutexGuard<'_, Vec<JoinHandle<()>>> {
set.joins.lock().unwrap_or_else(|e| e.into_inner())
}
fn free_if_idle(inner: &Arc<PoolInner>, set: &Arc<SetHandle>, freed: bool) {
if !freed {
return;
}
let mut reg = inner.lock();
if reg.stopping {
return;
}
if set.dead.load(Ordering::SeqCst) {
return;
}
reg.idle.push_back(set.clone());
}
pub struct SetLease {
inner: Arc<PoolInner>,
set: Arc<SetHandle>,
}
impl Drop for SetLease {
fn drop(&mut self) {
let freed = {
let mut st = lock_set(&self.set);
st.leased = false;
st.became_free()
};
free_if_idle(&self.inner, &self.set, freed);
}
}
pub struct Worker {
inner: Arc<PoolInner>,
set: Arc<SetHandle>,
tx: Sender<Assignment>,
}
const NEVER_DISPATCHED: &str =
"worker pool: the job was never dispatched — its worker thread had already exited";
impl Worker {
fn charge(&self) {
lock_set(&self.set).running += 1;
}
pub fn run<F>(self, body: F) -> Job
where
F: FnOnce() + Send + 'static,
{
let (done, done_rx) = sync_channel(1);
self.charge();
if self
.tx
.send(Assignment::Joinable {
body: Box::new(body),
ambient: InheritedRuntime::capture(),
done,
})
.is_err()
{
finish_job(&self.inner, &self.set);
errlog_sev_printf(
ErrlogSevEnum::Major,
&format!(
"{} worker pool: set {} took no job — {NEVER_DISPATCHED}.",
self.inner.name_prefix, self.set.index
),
);
return Job { done: None };
}
Job {
done: Some(done_rx),
}
}
pub fn run_detached<F>(self, label: String, body: F)
where
F: FnOnce() + Send + 'static,
{
self.charge();
if self
.tx
.send(Assignment::Detached {
body: Box::new(body),
ambient: InheritedRuntime::capture(),
label: label.clone(),
})
.is_err()
{
finish_job(&self.inner, &self.set);
errlog_sev_printf(
ErrlogSevEnum::Major,
&format!("{label}: {NEVER_DISPATCHED}. This connection is being torn down."),
);
}
}
}
pub struct Job {
done: Option<Receiver<thread::Result<()>>>,
}
impl Job {
pub fn join(self) -> thread::Result<()> {
match self.done {
Some(done) => done.recv().unwrap_or(Ok(())),
None => Err(Box::new(NEVER_DISPATCHED)),
}
}
}
fn announce_panic(label: &str) {
errlog_sev_printf(
ErrlogSevEnum::Major,
&format!(
"{label}: the connection thread panicked; this connection is being \
torn down. Other connections are unaffected."
),
);
}
fn announce_worker_death(prefix: &str, index: usize, roles: usize) {
errlog_sev_printf(
ErrlogSevEnum::Major,
&format!(
"{prefix} worker pool: a thread of set {index} exited unexpectedly. \
The set's {roles} threads are being retired and its slot returned; \
other connections are unaffected."
),
);
}
struct WorkerExit {
inner: Arc<PoolInner>,
set: Arc<SetHandle>,
reserved: usize,
clean: bool,
}
impl Drop for WorkerExit {
fn drop(&mut self) {
self.inner.reservation.release(self.reserved);
let (first_death, last_gone) = {
let mut st = lock_set(&self.set);
let already_dead = self.set.dead.load(Ordering::SeqCst);
if self.clean && !already_dead {
return;
}
st.live_workers -= 1;
self.set.dead.store(true, Ordering::SeqCst);
(!already_dead, st.live_workers == 0)
};
if first_death {
for tx in &self.set.senders {
let _ = tx.send(Assignment::Stop);
}
}
{
let mut reg = self.inner.lock();
if !reg.stopping {
reg.idle.retain(|s| !Arc::ptr_eq(s, &self.set));
if last_gone {
reg.all.retain(|s| !Arc::ptr_eq(s, &self.set));
reg.created -= 1;
lock_joins(&self.set).clear();
}
}
}
if first_death {
announce_worker_death(
self.inner.name_prefix,
self.set.index,
self.inner.roster.len(),
);
}
}
}
fn worker_loop(inner: Arc<PoolInner>, set: Arc<SetHandle>, rx: Receiver<Assignment>) {
while let Ok(assignment) = rx.recv() {
match assignment {
Assignment::Stop => break,
Assignment::Joinable {
body,
ambient,
done,
} => {
let outcome = ambient.run(|| catch_unwind(AssertUnwindSafe(body)));
let _ = done.send(outcome);
finish_job(&inner, &set);
}
Assignment::Detached {
body,
ambient,
label,
} => {
let outcome = ambient.run(|| catch_unwind(AssertUnwindSafe(body)));
if outcome.is_err() {
announce_panic(&label);
}
finish_job(&inner, &set);
}
}
}
}
fn finish_job(inner: &Arc<PoolInner>, set: &Arc<SetHandle>) {
let freed = {
let mut st = lock_set(set);
st.running -= 1;
st.became_free()
};
free_if_idle(inner, set, freed);
}
pub const POOL_RESERVATION_ENV: &str = "EPICS_RS_POOL_RESERVATION_MB";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ThreadMemoryTarget {
Host,
VxWorks,
Rtems,
}
impl ThreadMemoryTarget {
const fn current() -> Self {
if cfg!(target_os = "vxworks") {
ThreadMemoryTarget::VxWorks
} else if cfg!(target_os = "rtems") {
ThreadMemoryTarget::Rtems
} else {
ThreadMemoryTarget::Host
}
}
}
const fn per_thread_overhead(target: ThreadMemoryTarget) -> usize {
match target {
ThreadMemoryTarget::VxWorks => 1 << 20,
ThreadMemoryTarget::Rtems | ThreadMemoryTarget::Host => 0,
}
}
const PER_THREAD_OBJECT_ARENA: usize = 0;
const fn default_reservation_budget(target: ThreadMemoryTarget) -> usize {
match target {
ThreadMemoryTarget::VxWorks | ThreadMemoryTarget::Rtems => 160 << 20,
ThreadMemoryTarget::Host => usize::MAX,
}
}
fn resolve_reservation_budget(raw: Option<&str>, default: usize) -> usize {
let Some(raw) = raw else {
return default;
};
match raw.trim().parse::<usize>() {
Ok(mb) if mb > 0 => mb.saturating_mul(1 << 20),
_ => {
errlog_sev_printf(
ErrlogSevEnum::Minor,
&format!(
"{POOL_RESERVATION_ENV}={raw:?} is not a positive whole number of MiB; \
keeping the built-in worker-pool reservation budget"
),
);
default
}
}
}
const RESERVATION_PROBE_FLOOR: usize = 8 << 20;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TargetAnswer {
granted: bool,
basis: &'static str,
}
#[cfg(target_os = "vxworks")]
fn target_admits(bytes: usize) -> Option<TargetAnswer> {
const MAP_ANON_FD: libc::c_int = -1;
const BASIS: &str = "would not reserve that much address space in one mapping";
let addr = unsafe {
libc::mmap(
std::ptr::null_mut(),
bytes,
libc::PROT_NONE,
libc::MAP_PRIVATE | libc::MAP_ANON,
MAP_ANON_FD,
0,
)
};
if addr == libc::MAP_FAILED {
return Some(TargetAnswer {
granted: false,
basis: BASIS,
});
}
unsafe { libc::munmap(addr, bytes) };
Some(TargetAnswer {
granted: true,
basis: BASIS,
})
}
#[cfg(target_os = "rtems")]
fn target_admits(bytes: usize) -> Option<TargetAnswer> {
unsafe extern "C" {
fn malloc_free_space() -> libc::size_t;
}
Some(TargetAnswer {
granted: unsafe { malloc_free_space() } >= bytes,
basis: "has less than that free in the heap its thread stacks come from",
})
}
#[cfg(not(any(target_os = "vxworks", target_os = "rtems")))]
fn target_admits(_bytes: usize) -> Option<TargetAnswer> {
None
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BudgetVerdict {
Confirmed(usize),
Clamped {
asked: usize,
adopted: usize,
basis: &'static str,
},
Unverifiable { adopted: usize, from_env: bool },
FloorHeld { asked: usize, basis: &'static str },
}
impl BudgetVerdict {
const fn budget(self) -> usize {
match self {
BudgetVerdict::Confirmed(bytes)
| BudgetVerdict::Unverifiable { adopted: bytes, .. } => bytes,
BudgetVerdict::Clamped { adopted, .. } => adopted,
BudgetVerdict::FloorHeld { .. } => RESERVATION_PROBE_FLOOR,
}
}
fn notice(self) -> Option<(ErrlogSevEnum, String)> {
match self {
BudgetVerdict::Confirmed(_) => None,
BudgetVerdict::Unverifiable {
from_env: false, ..
} => None,
BudgetVerdict::Clamped {
asked,
adopted,
basis,
} => Some((
ErrlogSevEnum::Major,
format!(
"worker-pool reservation budget clamped from {} MiB to {} MiB: at {} MiB this \
target {}. {POOL_RESERVATION_ENV} names a ceiling the target still has to \
confirm; it does not add memory",
asked >> 20,
adopted >> 20,
asked >> 20,
basis
),
)),
BudgetVerdict::Unverifiable { adopted, .. } => Some((
ErrlogSevEnum::Minor,
format!(
"{POOL_RESERVATION_ENV} sets the worker-pool reservation budget to {} MiB, \
and this target has no measurement that tracks its thread-memory ceiling: \
this value cannot be verified and is taken as given",
adopted >> 20
),
)),
BudgetVerdict::FloorHeld { asked, basis } => Some((
ErrlogSevEnum::Major,
format!(
"this target confirmed no worker-pool reservation budget down to {} MiB \
(asked for {} MiB): even at {} MiB it {}. Holding that floor, so the pool \
refuses nearly every client instead of exhausting the target",
RESERVATION_PROBE_FLOOR >> 20,
asked >> 20,
RESERVATION_PROBE_FLOOR >> 20,
basis
),
)),
}
}
}
fn decide_reservation_budget(
requested: usize,
from_env: bool,
mut admits: impl FnMut(usize) -> Option<TargetAnswer>,
) -> BudgetVerdict {
if requested == usize::MAX {
return BudgetVerdict::Confirmed(requested);
}
let mut candidate = requested;
loop {
let Some(TargetAnswer { granted, basis }) = admits(candidate) else {
return BudgetVerdict::Unverifiable {
adopted: requested,
from_env,
};
};
match granted {
true if candidate == requested => return BudgetVerdict::Confirmed(candidate),
true => {
return BudgetVerdict::Clamped {
asked: requested,
adopted: candidate,
basis,
};
}
false if candidate <= RESERVATION_PROBE_FLOOR => {
return BudgetVerdict::FloorHeld {
asked: requested,
basis,
};
}
false => candidate = (candidate / 2).max(RESERVATION_PROBE_FLOOR),
}
}
}
fn announce_reservation_budget(verdict: BudgetVerdict) -> usize {
if let Some((severity, message)) = verdict.notice() {
errlog_sev_printf(severity, &message);
}
verdict.budget()
}
struct Reservation {
budget: usize,
held: AtomicUsize,
}
impl Reservation {
const fn new(budget: usize) -> Self {
Self {
budget,
held: AtomicUsize::new(0),
}
}
fn try_reserve(&self, bytes: usize) -> Result<(), (usize, usize)> {
let mut held = self.held.load(Ordering::SeqCst);
loop {
let Some(next) = held.checked_add(bytes).filter(|n| *n <= self.budget) else {
return Err((held, self.budget));
};
match self
.held
.compare_exchange_weak(held, next, Ordering::SeqCst, Ordering::SeqCst)
{
Ok(_) => return Ok(()),
Err(actual) => held = actual,
}
}
}
fn charge(&self, bytes: usize) {
self.held.fetch_add(bytes, Ordering::SeqCst);
}
fn release(&self, bytes: usize) {
self.held.fetch_sub(bytes, Ordering::SeqCst);
}
#[cfg(test)]
fn held(&self) -> usize {
self.held.load(Ordering::SeqCst)
}
}
static PROCESS_RESERVATION: LazyLock<Reservation> = LazyLock::new(|| {
let default = default_reservation_budget(ThreadMemoryTarget::current());
let requested =
resolve_reservation_budget(std::env::var(POOL_RESERVATION_ENV).ok().as_deref(), default);
Reservation::new(announce_reservation_budget(decide_reservation_budget(
requested,
requested != default,
target_admits,
)))
});
fn thread_reservation_bytes(stack: StackSizeClass) -> usize {
stack.bytes() + per_thread_overhead(ThreadMemoryTarget::current()) + PER_THREAD_OBJECT_ARENA
}
fn thread_reservation(role: &WorkerRole) -> usize {
thread_reservation_bytes(role.stack)
}
#[derive(Debug, Clone, Copy)]
pub struct ObjectArenaExhausted {
pub objects: usize,
}
impl std::fmt::Display for ObjectArenaExhausted {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"the target could not create the kernel mutex objects for a set of \
{} workers; this is transient, and a client that retries will be \
admitted once the target has objects again",
self.objects
)
}
}
impl std::error::Error for ObjectArenaExhausted {}
fn materialise_set_mutex(state: &Mutex<SetState>) -> bool {
state.try_lock().is_ok()
}
pub struct ThreadCharge {
bytes: usize,
}
impl ThreadCharge {
pub fn fixed(stack: StackSizeClass) -> Self {
let bytes = thread_reservation_bytes(stack);
PROCESS_RESERVATION.charge(bytes);
Self { bytes }
}
}
impl Drop for ThreadCharge {
fn drop(&mut self) {
PROCESS_RESERVATION.release(self.bytes);
}
}
#[derive(Debug)]
pub enum AcquireError {
AtCapacity {
capacity: usize,
},
OutOfReservation {
requested: usize,
reserved: usize,
budget: usize,
},
SpawnFailed(io::Error),
ShuttingDown,
}
impl std::fmt::Display for AcquireError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
AcquireError::AtCapacity { capacity } => {
write!(f, "worker pool at capacity ({capacity} sets)")
}
AcquireError::OutOfReservation {
requested,
reserved,
budget,
} => write!(
f,
"worker pool at its thread-memory budget: this set needs {} KiB, \
{} of {} MiB already reserved — raise {POOL_RESERVATION_ENV} \
if the target has the memory",
requested / 1024,
reserved >> 20,
budget >> 20,
),
AcquireError::SpawnFailed(e) => write!(f, "cannot create a worker set: {e}"),
AcquireError::ShuttingDown => write!(f, "worker pool is shutting down"),
}
}
}
impl std::error::Error for AcquireError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
AcquireError::SpawnFailed(e) => Some(e),
AcquireError::AtCapacity { .. }
| AcquireError::OutOfReservation { .. }
| AcquireError::ShuttingDown => None,
}
}
}
impl From<AcquireError> for io::Error {
fn from(cause: AcquireError) -> io::Error {
let kind = match &cause {
AcquireError::AtCapacity { .. } | AcquireError::OutOfReservation { .. } => {
io::ErrorKind::WouldBlock
}
AcquireError::SpawnFailed(e) => e.kind(),
AcquireError::ShuttingDown => io::ErrorKind::BrokenPipe,
};
io::Error::new(kind, cause)
}
}
pub struct WorkerPool<const N: usize> {
inner: Arc<PoolInner>,
}
impl<const N: usize> WorkerPool<N> {
pub fn new(name_prefix: &'static str, roster: [WorkerRole; N], capacity: usize) -> Self {
Self::with_reservation(name_prefix, roster, capacity, &PROCESS_RESERVATION)
}
fn with_reservation(
name_prefix: &'static str,
roster: [WorkerRole; N],
capacity: usize,
reservation: &'static Reservation,
) -> Self {
Self::with_reservation_and_gate(
name_prefix,
roster,
capacity,
reservation,
materialise_set_mutex,
)
}
fn with_reservation_and_gate(
name_prefix: &'static str,
roster: [WorkerRole; N],
capacity: usize,
reservation: &'static Reservation,
materialise: fn(&Mutex<SetState>) -> bool,
) -> Self {
let set_reservation = roster.iter().map(thread_reservation).sum();
Self {
inner: Arc::new(PoolInner {
roster: Box::new(roster),
name_prefix,
capacity,
set_reservation,
reservation,
materialise,
reg: Mutex::new(Registry {
idle: VecDeque::new(),
all: Vec::new(),
created: 0,
stopping: false,
}),
}),
}
}
pub fn acquire(&self) -> Result<(SetLease, [Worker; N]), AcquireError> {
enum Decision {
Reuse(Arc<SetHandle>),
Grow(usize),
Full,
}
let decision = {
let mut reg = self.inner.lock();
if reg.stopping {
return Err(AcquireError::ShuttingDown);
}
if let Some(set) = reg.idle.pop_front() {
Decision::Reuse(set)
} else if reg.created < self.inner.capacity {
let index = reg.created;
reg.created += 1;
Decision::Grow(index)
} else {
Decision::Full
}
};
let set = match decision {
Decision::Full => {
return Err(AcquireError::AtCapacity {
capacity: self.inner.capacity,
});
}
Decision::Reuse(set) => set,
Decision::Grow(index) => {
if let Err((reserved, budget)) = self
.inner
.reservation
.try_reserve(self.inner.set_reservation)
{
self.inner.lock().created -= 1;
return Err(AcquireError::OutOfReservation {
requested: self.inner.set_reservation,
reserved,
budget,
});
}
match self.spawn_set(index) {
Ok(set) => {
let mut reg = self.inner.lock();
reg.all.push(set.clone());
set
}
Err(e) => {
self.inner.lock().created -= 1;
return Err(AcquireError::SpawnFailed(e));
}
}
}
};
{
let mut st = lock_set(&set);
st.leased = true;
st.parked = false;
}
let workers: Vec<Worker> = (0..N)
.map(|slot| Worker {
inner: self.inner.clone(),
set: set.clone(),
tx: set.senders[slot].clone(),
})
.collect();
let workers: [Worker; N] = workers
.try_into()
.unwrap_or_else(|_| unreachable!("N workers for an N-role set"));
let lease = SetLease {
inner: self.inner.clone(),
set,
};
Ok((lease, workers))
}
fn spawn_set(&self, index: usize) -> io::Result<Arc<SetHandle>> {
let mut senders = Vec::with_capacity(N);
let mut receivers = Vec::with_capacity(N);
for _ in 0..N {
let (tx, rx) = channel::<Assignment>();
senders.push(tx);
receivers.push(rx);
}
let set = Arc::new(SetHandle {
index,
senders,
dead: AtomicBool::new(false),
state: Mutex::new(SetState {
leased: false,
running: 0,
parked: false,
live_workers: N,
}),
joins: Mutex::new(Vec::with_capacity(N)),
});
if !(self.inner.materialise)(&set.state) {
let unspawned: usize = self.inner.roster.iter().map(thread_reservation).sum();
self.inner.reservation.release(unspawned);
return Err(io::Error::new(
io::ErrorKind::WouldBlock,
ObjectArenaExhausted { objects: N },
));
}
let mut joins: Vec<JoinHandle<()>> = Vec::with_capacity(N);
for (slot, rx) in receivers.into_iter().enumerate() {
let role = self.inner.roster[slot];
let name = format!("{}-{} {index}", self.inner.name_prefix, role.suffix);
let inner = self.inner.clone();
let set_for_worker = set.clone();
let spawned = thread::Builder::new()
.name(name)
.stack_size(role.stack.bytes())
.spawn(move || {
let mut exit = WorkerExit {
inner: inner.clone(),
set: set_for_worker.clone(),
reserved: thread_reservation(&role),
clean: false,
};
let _ = enter_ioc_thread(role.priority);
worker_loop(inner, set_for_worker, rx);
exit.clean = true;
});
match spawned {
Ok(handle) => joins.push(handle),
Err(e) => {
let unspawned: usize = self.inner.roster[joins.len()..]
.iter()
.map(thread_reservation)
.sum();
self.inner.reservation.release(unspawned);
for tx in &set.senders {
let _ = tx.send(Assignment::Stop);
}
for handle in joins {
let _ = handle.join();
}
return Err(e);
}
}
}
lock_joins(&set).extend(joins);
Ok(set)
}
pub fn worker_count(&self) -> usize {
self.inner.lock().created * N
}
pub fn set_usage(&self) -> (usize, usize, usize) {
let reg = self.inner.lock();
let busy = reg.created - reg.idle.len();
(busy, reg.created, self.inner.capacity)
}
}
impl<const N: usize> Drop for WorkerPool<N> {
fn drop(&mut self) {
let (senders, joins) = {
let mut reg = self.inner.lock();
reg.stopping = true;
reg.idle.clear();
let mut senders: Vec<Sender<Assignment>> = Vec::new();
let mut joins: Vec<JoinHandle<()>> = Vec::new();
for set in ®.all {
senders.extend(set.senders.iter().cloned());
joins.append(&mut lock_joins(set));
}
(senders, joins)
};
for tx in senders {
let _ = tx.send(Assignment::Stop);
}
for handle in joins {
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::{Duration, Instant};
fn roster2() -> [WorkerRole; 2] {
[
WorkerRole {
suffix: "reader",
stack: StackSizeClass::Small,
priority: ThreadPriority::Low,
},
WorkerRole {
suffix: "writer",
stack: StackSizeClass::Small,
priority: ThreadPriority::Low,
},
]
}
#[test]
fn sequential_borrows_reuse_one_set() {
let pool: WorkerPool<2> = WorkerPool::new("test-pool", roster2(), 4);
const BORROWS: usize = 8;
for i in 0..BORROWS {
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
let ran = Arc::new(AtomicUsize::new(0));
let r = ran.clone();
let jr = reader.run(move || {
r.fetch_add(1, Ordering::SeqCst);
});
let w = ran.clone();
let jw = writer.run(move || {
w.fetch_add(1, Ordering::SeqCst);
});
assert!(jr.join().is_ok());
assert!(jw.join().is_ok());
drop(lease);
let deadline = Instant::now() + Duration::from_secs(5);
while pool.set_usage().0 != 0 {
assert!(Instant::now() < deadline, "set never returned to idle");
thread::yield_now();
}
assert_eq!(ran.load(Ordering::SeqCst), 2);
assert_eq!(
pool.worker_count(),
2,
"borrow {i} created new threads instead of reusing the idle set"
);
}
assert_eq!(
pool.worker_count(),
2,
"{BORROWS} sequential borrows must have created exactly one set"
);
}
#[test]
fn a_set_is_not_reidled_while_a_job_runs() {
let pool: WorkerPool<2> = WorkerPool::new("test-hold", roster2(), 4);
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
let gate = Arc::new((Mutex::new(false), std::sync::Condvar::new()));
let g = gate.clone();
let blocking = reader.run(move || {
let (m, cv) = &*g;
let mut open = m.lock().unwrap();
while !*open {
open = cv.wait(open).unwrap();
}
});
let quick = writer.run(|| {});
assert!(quick.join().is_ok());
drop(lease);
assert_eq!(
pool.set_usage().0,
1,
"a running job must keep its set busy"
);
{
let (m, cv) = &*gate;
*m.lock().unwrap() = true;
cv.notify_all();
}
assert!(blocking.join().is_ok());
let deadline = Instant::now() + Duration::from_secs(5);
while pool.set_usage().0 != 0 {
assert!(
Instant::now() < deadline,
"set never returned after its last job"
);
thread::yield_now();
}
assert_eq!(pool.worker_count(), 2);
}
#[test]
fn acquire_refuses_at_capacity_without_creating_a_thread() {
let pool: WorkerPool<2> = WorkerPool::new("test-cap", roster2(), 1);
let (lease, _workers) = pool.acquire().expect("first borrow");
let before = pool.worker_count();
let refused = pool.acquire().err();
assert!(
matches!(refused, Some(AcquireError::AtCapacity { capacity: 1 })),
"a full pool must refuse by naming the bound it reached, not queue \
or grow: {refused:?}"
);
assert_eq!(
pool.worker_count(),
before,
"a refusal must create no thread"
);
drop(lease);
}
#[test]
fn a_full_pool_and_a_refused_spawn_are_not_the_same_refusal() {
let pool: WorkerPool<2> = WorkerPool::new("test-gate", roster2(), 1);
let (lease, _workers) = pool.acquire().expect("first borrow");
let full = pool.acquire().err().expect("the pool is full");
let spawn_failed = AcquireError::SpawnFailed(io::ErrorKind::WouldBlock.into());
assert!(
matches!(full, AcquireError::AtCapacity { .. }),
"a full pool is a capacity refusal: {full:?}"
);
assert!(
!matches!(spawn_failed, AcquireError::AtCapacity { .. }),
"a refused spawn must never present as the capacity gate: it is \
the difference between 'raise the bound' and 'add memory'"
);
let as_io: io::Error = full.into();
assert_eq!(
as_io.kind(),
io::ErrorKind::WouldBlock,
"the historical kind is preserved for callers that only propagate"
);
assert!(
matches!(
as_io
.get_ref()
.and_then(|e| e.downcast_ref::<AcquireError>()),
Some(AcquireError::AtCapacity { capacity: 1 })
),
"the gate must survive the io::Error conversion: {as_io:?}"
);
drop(lease);
}
#[cfg(unix)]
#[test]
fn eagain_still_decodes_as_would_block() {
assert_eq!(
io::Error::from_raw_os_error(libc::EAGAIN).kind(),
io::ErrorKind::WouldBlock,
"EAGAIN decodes as WouldBlock — the collapse `AcquireError` exists \
to undo. If this ever stops holding, say so here rather than in a \
comment."
);
}
#[test]
fn a_panicked_job_returns_its_set_and_the_worker_survives() {
let pool: WorkerPool<2> = WorkerPool::new("test-panic", roster2(), 2);
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
let boom = reader.run(|| panic!("job blew up"));
let ok = writer.run(|| {});
assert!(boom.join().is_err(), "the panic must reach the joiner");
assert!(ok.join().is_ok());
drop(lease);
let deadline = Instant::now() + Duration::from_secs(5);
while pool.set_usage().0 != 0 {
assert!(Instant::now() < deadline, "panicked set never returned");
thread::yield_now();
}
let created_before = pool.worker_count();
let (lease2, [r2, w2]) = pool.acquire().expect("borrow after panic");
assert!(r2.run(|| {}).join().is_ok());
assert!(w2.run(|| {}).join().is_ok());
drop(lease2);
assert_eq!(
pool.worker_count(),
created_before,
"a lost worker is never recreated, and a survivor needs no replacement"
);
}
#[test]
fn a_detached_job_returns_its_set() {
let pool: WorkerPool<2> = WorkerPool::new("test-detach", roster2(), 2);
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
let ran = Arc::new(AtomicUsize::new(0));
let r = ran.clone();
reader.run_detached("conn".into(), move || {
r.fetch_add(1, Ordering::SeqCst);
});
let done = writer.run(|| {});
assert!(done.join().is_ok());
drop(lease);
let deadline = Instant::now() + Duration::from_secs(5);
while pool.set_usage().0 != 0 {
assert!(Instant::now() < deadline, "detached set never returned");
thread::yield_now();
}
assert_eq!(ran.load(Ordering::SeqCst), 1);
}
struct PanicOnDrop;
impl Drop for PanicOnDrop {
fn drop(&mut self) {
panic!("payload drop: the worker thread dies here, outside catch_unwind");
}
}
#[test]
fn a_worker_that_dies_retires_its_set_instead_of_leaking_it() {
let pool: WorkerPool<2> = WorkerPool::new("test-dead", roster2(), 2);
let (lease, [reader, _writer]) = pool.acquire().expect("borrow");
let (go, wait) = channel::<()>();
let job = reader.run(move || {
let _ = wait.recv();
std::panic::panic_any(PanicOnDrop);
});
drop(job);
drop(lease);
go.send(()).expect("the worker is waiting on this");
drop(go);
let deadline = Instant::now() + Duration::from_secs(10);
loop {
let (busy, created, _cap) = pool.set_usage();
if busy == 0 {
assert_eq!(
created, 0,
"a set with a dead worker must not stay countable: its \
threads are gone, so its slot must return to the bound"
);
break;
}
assert!(
Instant::now() < deadline,
"the set is still busy with no lease and no live job: a worker \
that died took its set out of circulation permanently, which \
is the one-set-per-death leak measured on target"
);
thread::yield_now();
}
let (lease2, [r2, w2]) = pool.acquire().expect("borrow after a death");
assert!(
r2.run(|| {}).join().is_ok(),
"a fresh set must actually run"
);
assert!(w2.run(|| {}).join().is_ok());
drop(lease2);
}
#[test]
fn a_death_under_a_live_lease_returns_the_slot_and_never_repools_the_set() {
let pool: WorkerPool<2> = WorkerPool::new("test-dead-leased", roster2(), 2);
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
let (go, wait) = channel::<()>();
let job = reader.run(move || {
let _ = wait.recv();
std::panic::panic_any(PanicOnDrop);
});
drop(job);
go.send(()).expect("the worker is waiting on this");
drop(go);
let deadline = Instant::now() + Duration::from_secs(10);
while pool.set_usage().1 != 0 {
assert!(
Instant::now() < deadline,
"a set that lost a thread stayed countable while its lease was \
held; the slot must return as soon as the threads are gone"
);
thread::yield_now();
}
assert!(
writer.run(|| {}).join().is_err(),
"a job dispatched into a retired set must be reported as not run"
);
drop(lease);
assert_eq!(
pool.set_usage(),
(0, 0, 2),
"the lease drop must not re-pool a retired set"
);
assert!(
pool.acquire().is_ok(),
"the pool must still admit after a death under lease"
);
}
#[test]
fn dropping_the_pool_joins_its_workers() {
let pool: WorkerPool<2> = WorkerPool::new("test-drop", roster2(), 2);
let (lease, [reader, writer]) = pool.acquire().expect("borrow");
assert!(reader.run(|| {}).join().is_ok());
assert!(writer.run(|| {}).join().is_ok());
drop(lease);
let deadline = Instant::now() + Duration::from_secs(5);
while pool.set_usage().0 != 0 {
assert!(Instant::now() < deadline, "set never returned before drop");
thread::yield_now();
}
drop(pool);
}
const HOST_SET: usize = 2 * 512 * 1024;
#[test]
fn admission_refuses_at_the_memory_budget_before_the_count_bound() {
static TWO_SETS: Reservation = Reservation::new(2 * HOST_SET);
let pool: WorkerPool<2> =
WorkerPool::with_reservation("test-budget", roster2(), 8, &TWO_SETS);
let (l1, _w1) = pool.acquire().expect("first set fits");
let (l2, _w2) = pool.acquire().expect("second set fits exactly");
assert_eq!(pool.worker_count(), 4, "two sets, two threads each");
let refused = pool.acquire().err().expect("the third set does not fit");
assert!(
matches!(
refused,
AcquireError::OutOfReservation {
requested,
reserved,
budget,
} if requested == HOST_SET
&& reserved == 2 * HOST_SET
&& budget == 2 * HOST_SET
),
"the refusal must name what was asked for, what is held and the \
budget — the three numbers the remedy needs: {refused:?}"
);
assert_eq!(
pool.worker_count(),
4,
"a refusal must not have created the threads it refused"
);
assert_eq!(
pool.set_usage(),
(2, 2, 8),
"the refused grow must leave the slot reservation exactly as it \
found it"
);
drop(l1);
drop(l2);
let (l3, _w3) = pool.acquire().expect("an idle set is reused, not grown");
assert_eq!(pool.worker_count(), 4);
drop(l3);
drop(pool);
assert_eq!(
TWO_SETS.held.load(Ordering::SeqCst),
0,
"every thread's reservation must come back when the pool is dropped"
);
}
#[test]
fn a_dead_set_gives_its_memory_back_to_the_budget() {
static ONE_SET: Reservation = Reservation::new(HOST_SET);
let pool: WorkerPool<2> = WorkerPool::with_reservation("test-rel", roster2(), 4, &ONE_SET);
let (lease, [reader, _writer]) = pool.acquire().expect("the one set fits");
let (go, wait) = channel::<()>();
let job = reader.run(move || {
let _ = wait.recv();
std::panic::panic_any(PanicOnDrop);
});
drop(job);
drop(lease);
go.send(()).expect("the worker is waiting on this");
drop(go);
let deadline = Instant::now() + Duration::from_secs(10);
while ONE_SET.held.load(Ordering::SeqCst) != 0 {
assert!(
Instant::now() < deadline,
"a set that lost a worker kept its reservation: held {} of {}",
ONE_SET.held.load(Ordering::SeqCst),
HOST_SET
);
thread::yield_now();
}
pool.acquire()
.expect("the budget freed by the dead set must admit a new one");
}
#[test]
fn a_dead_set_retires_its_thread_handles_with_its_slot() {
let pool: WorkerPool<2> = WorkerPool::new("test-joins", roster2(), 4);
let (lease, [reader, _writer]) = pool.acquire().expect("borrow");
let set = lease.set.clone();
assert_eq!(
lock_joins(&set).len(),
2,
"a live set owns one handle per thread"
);
let (go, wait) = channel::<()>();
let job = reader.run(move || {
let _ = wait.recv();
std::panic::panic_any(PanicOnDrop);
});
drop(job);
drop(lease);
go.send(()).expect("the worker is waiting on this");
drop(go);
let deadline = Instant::now() + Duration::from_secs(10);
while pool.set_usage().1 != 0 {
assert!(
Instant::now() < deadline,
"the dead set never gave its slot back: {:?}",
pool.set_usage()
);
thread::yield_now();
}
assert!(
lock_joins(&set).is_empty(),
"the slot came back and the threads did not — {} handle(s) still \
neither joined nor detached",
lock_joins(&set).len()
);
}
#[test]
fn each_target_is_charged_the_figure_measured_on_it() {
for target in [
ThreadMemoryTarget::Host,
ThreadMemoryTarget::VxWorks,
ThreadMemoryTarget::Rtems,
] {
let (overhead, budget) = (
per_thread_overhead(target),
default_reservation_budget(target),
);
match target {
ThreadMemoryTarget::Host => {
assert_eq!(overhead, 0);
assert_eq!(budget, usize::MAX, "a host meets no thread-memory wall");
}
ThreadMemoryTarget::VxWorks => {
assert_eq!(
overhead,
1024 * 1024,
"the flat per-thread reservation measured on VxWorks 7: \
three stack classes, walls within 2.1 %"
);
assert_eq!(budget, 160 * 1024 * 1024);
}
ThreadMemoryTarget::Rtems => {
assert_eq!(
overhead, 0,
"RTEMS spends less than the stacks it is asked for — \
1,167,383 B per client against 1,572,864 B declared — \
so there is no flat term to charge"
);
assert_eq!(budget, 160 * 1024 * 1024);
}
}
}
let default = default_reservation_budget(ThreadMemoryTarget::VxWorks);
assert_eq!(resolve_reservation_budget(None, default), default);
assert_eq!(resolve_reservation_budget(Some("8"), default), 8 << 20);
assert_eq!(resolve_reservation_budget(Some(" 12 "), default), 12 << 20);
assert_eq!(resolve_reservation_budget(Some("0"), default), default);
assert_eq!(resolve_reservation_budget(Some("lots"), default), default);
assert_eq!(resolve_reservation_budget(Some(""), default), default);
}
const TEST_BASIS: &str = "answered from the table this test wrote";
fn probe<'a>(
answers: &'static [(usize, Option<bool>)],
asked: &'a mut Vec<usize>,
) -> impl FnMut(usize) -> Option<TargetAnswer> + 'a {
move |bytes| {
asked.push(bytes);
answers
.iter()
.find(|(size, _)| *size == bytes)
.map(|(_, answer)| *answer)
.unwrap_or(Some(false))
.map(|granted| TargetAnswer {
granted,
basis: TEST_BASIS,
})
}
}
#[test]
fn a_configured_budget_is_confirmed_clamped_or_declared_unverifiable() {
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(usize::MAX, false, probe(&[], &mut asked)),
BudgetVerdict::Confirmed(usize::MAX)
);
assert!(asked.is_empty(), "no mapping is asked for on a host");
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(
160 << 20,
false,
probe(&[(160 << 20, Some(true))], &mut asked)
),
BudgetVerdict::Confirmed(160 << 20)
);
assert_eq!(asked, vec![160 << 20], "an honest value costs one mapping");
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(
320 << 20,
true,
probe(
&[(320 << 20, Some(false)), (160 << 20, Some(true))],
&mut asked
)
),
BudgetVerdict::Clamped {
asked: 320 << 20,
adopted: 160 << 20,
basis: TEST_BASIS
}
);
assert_eq!(asked, vec![320 << 20, 160 << 20]);
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(320 << 20, true, probe(&[(320 << 20, None)], &mut asked)),
BudgetVerdict::Unverifiable {
adopted: 320 << 20,
from_env: true
}
);
assert_eq!(asked, vec![320 << 20], "one question, then no more");
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(160 << 20, false, probe(&[(160 << 20, None)], &mut asked)),
BudgetVerdict::Unverifiable {
adopted: 160 << 20,
from_env: false
}
);
}
#[test]
fn the_budget_descent_terminates_on_the_floor() {
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(64 << 20, true, probe(&[], &mut asked)),
BudgetVerdict::FloorHeld {
asked: 64 << 20,
basis: TEST_BASIS
}
);
assert_eq!(
asked,
vec![64 << 20, 32 << 20, 16 << 20, 8 << 20],
"halving, and the last question is the floor itself"
);
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(
10 << 20,
true,
probe(&[(RESERVATION_PROBE_FLOOR, Some(true))], &mut asked)
),
BudgetVerdict::Clamped {
asked: 10 << 20,
adopted: RESERVATION_PROBE_FLOOR,
basis: TEST_BASIS
}
);
assert_eq!(asked, vec![10 << 20, RESERVATION_PROBE_FLOOR]);
let mut asked = Vec::new();
assert_eq!(
decide_reservation_budget(4 << 20, true, probe(&[], &mut asked)),
BudgetVerdict::FloorHeld {
asked: 4 << 20,
basis: TEST_BASIS
}
);
assert_eq!(asked, vec![4 << 20]);
}
#[test]
fn every_verdict_but_confirmation_is_announced() {
for verdict in [
BudgetVerdict::Confirmed(160 << 20),
BudgetVerdict::Clamped {
asked: 320 << 20,
adopted: 160 << 20,
basis: TEST_BASIS,
},
BudgetVerdict::Unverifiable {
adopted: 320 << 20,
from_env: true,
},
BudgetVerdict::Unverifiable {
adopted: 160 << 20,
from_env: false,
},
BudgetVerdict::FloorHeld {
asked: 320 << 20,
basis: TEST_BASIS,
},
] {
let notice = verdict.notice();
match verdict {
BudgetVerdict::Confirmed(bytes) => {
assert_eq!(notice, None, "an honoured budget is not news");
assert_eq!(verdict.budget(), bytes);
}
BudgetVerdict::Unverifiable {
adopted,
from_env: false,
} => {
assert_eq!(
notice, None,
"the built-in default carries its own measurement"
);
assert_eq!(verdict.budget(), adopted);
}
BudgetVerdict::Unverifiable { adopted, .. } => {
let (severity, message) = notice.expect("an unverified value must say so");
assert_eq!(severity, ErrlogSevEnum::Minor);
assert!(
message.contains("cannot be verified")
&& message.contains(POOL_RESERVATION_ENV),
"the notice must name the switch and its own uncertainty: {message}"
);
assert_eq!(verdict.budget(), adopted);
}
BudgetVerdict::Clamped {
asked,
adopted,
basis,
} => {
let (severity, message) = notice.expect("a clamp must say so");
assert_eq!(
severity,
ErrlogSevEnum::Major,
"the IOC is not doing what the switch said"
);
assert!(
message.contains(&format!("{} MiB", asked >> 20))
&& message.contains(&format!("{} MiB", adopted >> 20)),
"both numbers, or the operator cannot tell what happened: {message}"
);
assert!(
message.contains(basis),
"the notice must give the target's own account of the refusal, not a \
mechanism it did not use: {message}"
);
assert_eq!(verdict.budget(), adopted);
}
BudgetVerdict::FloorHeld { asked, basis } => {
let (severity, message) = notice.expect("a held floor must say so");
assert_eq!(severity, ErrlogSevEnum::Major);
assert!(
message.contains(&format!("{} MiB", asked >> 20)),
"the notice must name what was asked for: {message}"
);
assert!(
message.contains(basis),
"the notice must give the target's own account of the refusal: {message}"
);
assert_eq!(verdict.budget(), RESERVATION_PROBE_FLOOR);
}
}
}
}
#[test]
fn a_fixed_thread_charges_the_process_account_and_gives_it_back() {
let before = PROCESS_RESERVATION.held();
let expect = thread_reservation_bytes(StackSizeClass::Small);
{
let _charge = ThreadCharge::fixed(StackSizeClass::Small);
assert_eq!(
PROCESS_RESERVATION.held(),
before + expect,
"a fixed thread must appear in the account the pool divides"
);
}
assert_eq!(
PROCESS_RESERVATION.held(),
before,
"the charge is released by the guard's `Drop`, not by a caller"
);
}
#[test]
fn the_spawn_helper_holds_its_charge_for_the_thread_and_not_the_call() {
use crate::runtime::task::spawn_dedicated_thread;
let before = PROCESS_RESERVATION.held();
let expect = thread_reservation_bytes(StackSizeClass::Small);
let (release, wait) = channel::<()>();
let (started, running) = channel::<()>();
let handle = spawn_dedicated_thread(
"charged-fixed-thread".to_string(),
ThreadPriority::Low,
StackSizeClass::Small,
move || {
let _ = started.send(());
let _ = wait.recv();
},
)
.expect("the host can create one thread");
running.recv().expect("the thread starts");
assert_eq!(
PROCESS_RESERVATION.held(),
before + expect,
"the account must hold the stack while the thread runs"
);
drop(release);
handle.join().expect("the thread ends cleanly");
assert_eq!(
PROCESS_RESERVATION.held(),
before,
"and must be back where it started once the thread is gone"
);
}
#[test]
fn a_target_that_refuses_a_mutex_object_refuses_the_connection() {
static ARENA: Reservation = Reservation::new(8 * HOST_SET);
fn arena_empty(_: &Mutex<SetState>) -> bool {
false
}
let pool: WorkerPool<2> =
WorkerPool::with_reservation_and_gate("test-arena", roster2(), 8, &ARENA, arena_empty);
let before = ARENA.held();
let refused = pool.acquire().err().expect("the arena refuses the set");
let AcquireError::SpawnFailed(ref e) = refused else {
panic!("an arena refusal is the target saying no: {refused:?}");
};
let arena = e
.get_ref()
.and_then(|src| src.downcast_ref::<ObjectArenaExhausted>())
.expect("the cause must be recognisable by type, not by prose");
assert_eq!(arena.objects, 2, "one object per worker in the set");
assert_eq!(
e.kind(),
io::ErrorKind::WouldBlock,
"a transient refusal is retryable, and a client's retry is the pacing"
);
assert_eq!(pool.worker_count(), 0, "a refusal must create no thread");
assert_eq!(
ARENA.held(),
before,
"the set's memory must go back: it was reserved for threads that do \
not exist"
);
let ok: WorkerPool<2> = WorkerPool::with_reservation_and_gate(
"test-arena-recovers",
roster2(),
1,
&ARENA,
materialise_set_mutex,
);
let _lease = ok.acquire().expect("a target with objects admits");
assert!(pool.acquire().is_err(), "still refusing");
}
}