use std::any::Any;
use std::collections::{BTreeMap, VecDeque};
use std::fmt;
use std::future::{Future, poll_fn};
use std::num::{NonZeroU64, NonZeroUsize};
use std::pin::pin;
use std::process::ExitStatus;
use std::sync::Arc;
#[cfg(all(unix, feature = "process"))]
use std::sync::Mutex;
use std::task::{Context, Poll};
use std::time::Duration;
#[cfg(feature = "process")]
use std::io;
#[cfg(all(unix, feature = "process"))]
use crate::rt::sync::Notify;
use crate::rt::sync::{OwnedSemaphorePermit, Semaphore};
use lgwks_deps::tokio::task::Id;
#[cfg(all(unix, feature = "process"))]
use lgwks_deps::tokio::process::{ChildStderr, ChildStdout, Command};
#[cfg(feature = "process")]
pub use lgwks_std::process::ContainmentMechanism;
#[cfg(feature = "process")]
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct IdentifiedSpawn {
task: TaskId,
leader: lgwks_std::process::ProcessIdentity,
}
#[cfg(feature = "process")]
impl IdentifiedSpawn {
#[must_use]
pub const fn task(&self) -> TaskId {
self.task
}
#[must_use]
pub const fn leader(&self) -> &lgwks_std::process::ProcessIdentity {
&self.leader
}
}
#[cfg(all(unix, feature = "process"))]
fn identify_leader(pid: i32) -> io::Result<lgwks_std::process::ProcessIdentity> {
if let Some(leader) = lgwks_std::process::identify_process(pid)? {
return Ok(leader);
}
let refusal = Err(io::Error::new(
io::ErrorKind::NotFound,
format!("lgwks_bot: the unreaped leader {pid} was absent from the process table"),
));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "identify_leader: returning an error to the caller");
refusal
}
use super::cancel::CancellationToken;
use super::clock::{Clock, TimeSource};
#[cfg(all(unix, feature = "process"))]
use super::io::{AsyncRead, AsyncReadExt};
#[cfg(feature = "process")]
use super::process::ProcessSpec;
#[cfg(all(unix, feature = "process"))]
use super::process::{CapturedStream, ProcessRun, ProcessRunError};
use super::task::{JoinSet, yield_now};
#[cfg(feature = "script")]
pub use super::tenancy::{SpawnRefused, TenancyPolicy};
#[cfg(feature = "script")]
use crate::script::Tenant;
const YIELD_INTERVAL: u64 = 32;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Budget {
Iterations(NonZeroU64),
For(Duration),
Ongoing,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Outcome {
Exhausted {
iterations: u64,
},
Cancelled {
iterations: u64,
},
}
impl Outcome {
#[must_use]
pub const fn iterations(self) -> u64 {
match self {
Self::Exhausted { iterations } | Self::Cancelled { iterations } => iterations,
}
}
#[must_use]
pub const fn was_cancelled(self) -> bool {
matches!(self, Self::Cancelled { .. })
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum TrySpawnRefusal {
AtCapacity,
Cancelled,
}
impl fmt::Display for TrySpawnRefusal {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::AtCapacity => formatter.write_str("the supervisor is at its in-flight bound"),
Self::Cancelled => {
formatter.write_str("the supervisor is cancelled and admits no new work")
}
}
}
}
impl std::error::Error for TrySpawnRefusal {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct SupervisorCancelled;
impl fmt::Display for SupervisorCancelled {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("the supervisor is cancelled and admits no new process")
}
}
impl std::error::Error for SupervisorCancelled {}
#[cfg(all(feature = "process", feature = "script"))]
#[derive(Debug)]
#[non_exhaustive]
pub enum RunForError {
Refused(SpawnRefused),
Run(ProcessRunError),
}
#[cfg(all(feature = "process", feature = "script"))]
impl fmt::Display for RunForError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Refused(ref refusal) => write!(formatter, "{refusal}"),
Self::Run(ref error) => write!(formatter, "{error}"),
}
}
}
#[cfg(all(feature = "process", feature = "script"))]
impl std::error::Error for RunForError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Refused(ref refusal) => Some(refusal),
Self::Run(ref error) => Some(error),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[non_exhaustive]
pub struct TaskId(u64);
impl TaskId {
#[must_use]
pub const fn get(self) -> u64 {
self.0
}
}
impl fmt::Display for TaskId {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "task {}", self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum TaskOutcome {
Completed {
task: TaskId,
cleanup: Option<CleanupReceipt>,
#[cfg(feature = "process")]
containment: Option<Containment>,
},
Cancelled {
task: TaskId,
cleanup: Option<CleanupReceipt>,
#[cfg(feature = "process")]
containment: Option<Containment>,
},
Aborted {
task: TaskId,
},
Panicked {
task: TaskId,
message: String,
},
Failed {
task: TaskId,
status: Option<ExitStatus>,
cleanup: Option<CleanupReceipt>,
#[cfg(feature = "process")]
containment: Option<Containment>,
},
#[cfg(all(unix, feature = "process"))]
CleanupSettled {
task: TaskId,
cleanup: CleanupReceipt,
containment: Containment,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum CleanupReceipt {
CleanupConfirmed,
CleanupPending,
CleanupFailed,
CleanupSurvivors {
survivors: Vec<i32>,
},
}
impl CleanupReceipt {
#[must_use]
pub const fn is_complete(&self) -> bool {
matches!(self, Self::CleanupConfirmed)
}
#[must_use]
pub fn survivors(&self) -> &[i32] {
match *self {
Self::CleanupSurvivors { ref survivors } => survivors.as_slice(),
Self::CleanupConfirmed | Self::CleanupPending | Self::CleanupFailed => &[],
}
}
}
#[cfg(feature = "process")]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum ResidualRisk {
TableUnreadable,
CaptureTruncated,
LeaderExited,
}
#[cfg(feature = "process")]
#[derive(Debug, Clone, Default, PartialEq, Eq)]
#[non_exhaustive]
pub struct Containment {
mechanism: ContainmentMechanism,
captured: usize,
signalled: usize,
survivors: Vec<i32>,
residual: Option<ResidualRisk>,
}
#[cfg(feature = "process")]
impl Containment {
#[must_use]
pub const fn mechanism(&self) -> ContainmentMechanism {
self.mechanism
}
#[must_use]
pub const fn captured(&self) -> usize {
self.captured
}
#[must_use]
pub const fn signalled(&self) -> usize {
self.signalled
}
#[must_use]
pub fn survivors(&self) -> &[i32] {
&self.survivors
}
#[must_use]
pub const fn residual_risk(&self) -> Option<ResidualRisk> {
self.residual
}
#[must_use]
pub fn is_complete(&self) -> bool {
self.survivors.is_empty() && self.residual.is_none() && self.mechanism.read_a_table()
}
#[cfg(all(unix, feature = "process"))]
fn record_capture(&mut self, set: &Capture) {
self.mechanism = set.mechanism.strongest(self.mechanism);
if set.truncated {
self.note(ResidualRisk::CaptureTruncated);
}
self.captured = self.captured.max(set.pids.len());
}
fn record_signal(&mut self) {
self.signalled = self.signalled.saturating_add(1);
}
fn record_survivors(&mut self, running: &[i32]) {
self.survivors = running.to_vec();
}
fn note(&mut self, risk: ResidualRisk) {
if self.residual.is_none() {
self.residual = Some(risk);
}
}
}
impl TaskOutcome {
#[must_use]
pub const fn task(&self) -> TaskId {
match *self {
Self::Completed { task, .. }
| Self::Cancelled { task, .. }
| Self::Aborted { task }
| Self::Failed { task, .. }
| Self::Panicked { task, .. } => task,
#[cfg(all(unix, feature = "process"))]
Self::CleanupSettled { task, .. } => task,
}
}
#[must_use]
pub const fn is_success(&self) -> bool {
#[cfg(all(unix, feature = "process"))]
{
matches!(self, Self::Completed { .. } | Self::CleanupSettled { .. })
}
#[cfg(not(all(unix, feature = "process")))]
{
matches!(self, Self::Completed { .. })
}
}
#[must_use]
pub const fn is_panic(&self) -> bool {
matches!(self, Self::Panicked { .. })
}
#[must_use]
pub const fn is_process_failure(&self) -> bool {
matches!(self, Self::Failed { .. })
}
#[must_use]
pub fn exit_status(&self) -> Option<ExitStatus> {
match *self {
Self::Failed { status, .. } => status,
Self::Completed { .. }
| Self::Cancelled { .. }
| Self::Aborted { .. }
| Self::Panicked { .. } => None,
#[cfg(all(unix, feature = "process"))]
Self::CleanupSettled { .. } => None,
}
}
#[must_use]
pub const fn cleanup(&self) -> Option<&CleanupReceipt> {
match *self {
Self::Completed { ref cleanup, .. }
| Self::Cancelled { ref cleanup, .. }
| Self::Failed { ref cleanup, .. } => cleanup.as_ref(),
Self::Aborted { .. } | Self::Panicked { .. } => None,
#[cfg(all(unix, feature = "process"))]
Self::CleanupSettled { ref cleanup, .. } => Some(cleanup),
}
}
#[cfg(feature = "process")]
#[must_use]
pub const fn containment(&self) -> Option<&Containment> {
match *self {
Self::Completed {
ref containment, ..
}
| Self::Cancelled {
ref containment, ..
}
| Self::Failed {
ref containment, ..
} => containment.as_ref(),
Self::Aborted { .. } | Self::Panicked { .. } => None,
#[cfg(all(unix, feature = "process"))]
Self::CleanupSettled {
ref containment, ..
} => Some(containment),
}
}
#[must_use]
pub fn panic_message(&self) -> Option<&str> {
match *self {
Self::Panicked { ref message, .. } => Some(message.as_str()),
Self::Completed { .. }
| Self::Cancelled { .. }
| Self::Aborted { .. }
| Self::Failed { .. } => None,
#[cfg(all(unix, feature = "process"))]
Self::CleanupSettled { .. } => None,
}
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct ShutdownReport {
outcomes: Vec<TaskOutcome>,
stats: Stats,
#[cfg(all(unix, feature = "process"))]
cleanup_owners: Arc<CleanupOwners>,
}
impl ShutdownReport {
#[must_use]
pub fn outcomes(&self) -> &[TaskOutcome] {
&self.outcomes
}
pub fn into_outcomes(self) -> Result<Vec<TaskOutcome>, Self> {
#[cfg(all(unix, feature = "process"))]
if self.pending_cleanup_count() > 0 {
let refusal = Err(self);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "into_outcomes: returning an error to the caller");
return refusal;
}
Ok(self.outcomes)
}
#[must_use]
pub const fn stats(&self) -> Stats {
self.stats
}
pub fn panicked(&self) -> impl Iterator<Item = &TaskOutcome> {
self.outcomes.iter().filter(|outcome| outcome.is_panic())
}
#[must_use]
pub fn is_clean(&self) -> bool {
self.outcomes.iter().all(|outcome| {
outcome.is_success() && outcome.cleanup().is_none_or(CleanupReceipt::is_complete)
}) && {
#[cfg(feature = "process")]
{
self.pending_cleanup_count() == 0
&& self
.outcomes
.iter()
.all(|outcome| outcome.containment().is_none_or(Containment::is_complete))
}
#[cfg(not(feature = "process"))]
{
true
}
}
}
pub fn failed(&self) -> impl Iterator<Item = &TaskOutcome> {
self.outcomes
.iter()
.filter(|outcome| outcome.is_process_failure())
}
#[cfg(all(unix, feature = "process"))]
#[must_use]
pub fn pending_cleanup_count(&self) -> usize {
self.cleanup_owners.pending_count()
}
#[cfg(all(unix, feature = "process"))]
pub fn reap_pending_cleanups(&mut self) -> usize {
self.reap_pending_cleanups_with(&NATIVE_GROUP_OBSERVER)
}
#[cfg(all(unix, feature = "process"))]
fn reap_pending_cleanups_with(&mut self, observer: &dyn GroupObserver) -> usize {
let settled = self.cleanup_owners.drive(observer, |_| false);
let count = settled.len();
self.outcomes
.extend(
settled
.into_iter()
.map(|(task, containment)| TaskOutcome::CleanupSettled {
task,
cleanup: settled_receipt(&containment),
containment,
}),
);
count
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct Stats {
pub spawned: u64,
pub completed: u64,
pub succeeded: u64,
pub failed: u64,
pub cancelled: u64,
pub aborted: u64,
pub panicked: u64,
pub refused: u64,
pub reports_dropped: u64,
}
impl Stats {
#[must_use]
pub const fn in_flight(self) -> u64 {
self.spawned.saturating_sub(self.completed)
}
}
#[cfg(feature = "script")]
struct Lease {
permit: Option<OwnedSemaphorePermit>,
charge: Option<TenantCharge>,
}
#[cfg(not(feature = "script"))]
struct Lease {
_permit: OwnedSemaphorePermit,
}
#[cfg(feature = "script")]
struct TenantCharge {
shell: Arc<tenancy_support::TenancyShell>,
tenant: Tenant,
}
#[cfg(feature = "script")]
impl Lease {
fn plain(permit: OwnedSemaphorePermit) -> Self {
Self {
permit: Some(permit),
charge: None,
}
}
fn tenanted(
permit: OwnedSemaphorePermit,
shell: Arc<tenancy_support::TenancyShell>,
tenant: Tenant,
) -> Self {
Self {
permit: Some(permit),
charge: Some(TenantCharge { shell, tenant }),
}
}
}
impl Lease {
#[cfg(not(feature = "script"))]
fn plain(permit: OwnedSemaphorePermit) -> Self {
Self { _permit: permit }
}
}
#[cfg(feature = "script")]
impl Drop for Lease {
fn drop(&mut self) {
if let Some(charge) = self.charge.take()
&& let Some(permit) = self.permit.take()
{
charge.shell.release(&charge.tenant, permit);
}
}
}
#[cfg(feature = "script")]
const IMPLICIT_TENANT: &str = "_implicit";
#[cfg(feature = "script")]
fn implicit_tenant() -> Tenant {
Tenant::assumed(IMPLICIT_TENANT)
}
#[cfg(feature = "script")]
mod tenancy_support {
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll, Waker};
use crate::rt::sync::{OwnedSemaphorePermit, Semaphore};
use super::Lease;
use crate::rt::tenancy::{Arrival, DeficitRoundRobin, Grant, GrantOutcome, TryArrival};
use crate::script::Tenant;
pub(super) struct TenancyShell {
permits: Arc<Semaphore>,
round: std::sync::Mutex<DeficitRoundRobin<Arc<WaitSlot>>>,
}
struct WaitSlot {
state: std::sync::Mutex<SlotState>,
}
struct SlotState {
permit: Option<OwnedSemaphorePermit>,
abandoned: bool,
handed: bool,
withdrawn: bool,
waker: Option<Waker>,
}
impl WaitSlot {
fn fresh() -> Arc<Self> {
Arc::new(Self {
state: std::sync::Mutex::new(SlotState {
permit: None,
abandoned: false,
handed: false,
withdrawn: false,
waker: None,
}),
})
}
fn lock(&self) -> std::sync::MutexGuard<'_, SlotState> {
crate::journal::owner::lock(&self.state)
}
fn is_live(&self) -> bool {
let state = self.lock();
!state.abandoned && !state.handed
}
fn deliver(&self, permit: OwnedSemaphorePermit) -> Option<OwnedSemaphorePermit> {
let mut state = self.lock();
if state.abandoned {
state.withdrawn = true;
return Some(permit);
}
if state.handed {
return Some(permit);
}
state.permit = Some(permit);
state.handed = true;
if let Some(waker) = state.waker.take() {
waker.wake();
}
None
}
fn take_permit(&self) -> Option<OwnedSemaphorePermit> {
self.lock().permit.take()
}
fn leave(&self) -> Departure {
let mut state = self.lock();
if state.handed {
return Departure::Handed(state.permit.take());
}
state.abandoned = true;
Departure::Abandoned
}
fn register(&self, waker: &Waker) {
let mut state = self.lock();
let stale = state
.waker
.as_ref()
.is_none_or(|current| !current.will_wake(waker));
if stale {
state.waker = Some(waker.clone());
}
}
}
enum Departure {
Handed(Option<OwnedSemaphorePermit>),
Abandoned,
}
pub(super) struct WaitPermit {
slot: Arc<WaitSlot>,
shell: Arc<TenancyShell>,
tenant: Tenant,
}
impl WaitPermit {
fn new(slot: Arc<WaitSlot>, shell: Arc<TenancyShell>, tenant: Tenant) -> Self {
Self {
slot,
shell,
tenant,
}
}
}
impl Future for WaitPermit {
type Output = OwnedSemaphorePermit;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
let this = self.get_mut();
this.slot.register(context.waker());
match this.slot.take_permit() {
Some(permit) => Poll::Ready(permit),
None => Poll::Pending,
}
}
}
impl Drop for WaitPermit {
fn drop(&mut self) {
match self.slot.leave() {
Departure::Handed(Some(permit)) => self.shell.release(&self.tenant, permit),
Departure::Handed(None) => {}
Departure::Abandoned => self.shell.report_abandoned(&self.slot, &self.tenant),
}
}
}
pub(super) enum Admission {
Admitted(Lease),
Queued(WaitPermit),
Refused {
limit: usize,
},
SupervisorQueueFull {
limit: usize,
},
}
impl TenancyShell {
pub(super) fn new(permits: Arc<Semaphore>, policy: super::TenancyPolicy) -> Self {
Self {
permits,
round: std::sync::Mutex::new(DeficitRoundRobin::new(policy)),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, DeficitRoundRobin<Arc<WaitSlot>>> {
crate::journal::owner::lock(&self.round)
}
fn take_from_pool(&self) -> Option<OwnedSemaphorePermit> {
Arc::clone(&self.permits).try_acquire_owned().ok()
}
fn lease(self: &Arc<Self>, permit: OwnedSemaphorePermit, tenant: &Tenant) -> Lease {
Lease::tenanted(permit, Arc::clone(self), tenant.clone())
}
pub(super) fn admit(self: &Arc<Self>, tenant: &Tenant) -> Admission {
let slot = WaitSlot::fresh();
let mut round = self.lock();
let outcome = round.arrive(tenant, Arc::clone(&slot), || self.take_from_pool());
match outcome {
Arrival::Immediate(permit) => Admission::Admitted(self.lease(permit, tenant)),
Arrival::Queued => {
self.pump(&mut round, None);
Admission::Queued(WaitPermit::new(
Arc::clone(&slot),
Arc::clone(self),
tenant.clone(),
))
}
Arrival::Refused { limit } => Admission::Refused { limit },
Arrival::SupervisorFull { limit } => Admission::SupervisorQueueFull { limit },
}
}
pub(super) fn try_admit(self: &Arc<Self>, tenant: &Tenant) -> Admission {
let mut round = self.lock();
let outcome = round.try_arrive(tenant, || self.take_from_pool());
match outcome {
TryArrival::Immediate(permit) => Admission::Admitted(self.lease(permit, tenant)),
TryArrival::Contended => Admission::Refused {
limit: round.policy().queue_per_tenant(),
},
}
}
pub(super) fn release(&self, tenant: &Tenant, permit: OwnedSemaphorePermit) {
let mut round = self.lock();
round.note_release(tenant);
self.pump(&mut round, Some(permit));
}
fn pump(
&self,
round: &mut std::sync::MutexGuard<'_, DeficitRoundRobin<Arc<WaitSlot>>>,
spare: Option<OwnedSemaphorePermit>,
) {
let mut spare = spare;
loop {
let permit = match spare.take() {
Some(permit) => permit,
None if round.has_eligible() => match self.take_from_pool() {
Some(permit) => permit,
None => return,
},
None => return,
};
let mut is_live = |slot: &Arc<WaitSlot>| slot.is_live();
match round.grant(permit, &mut is_live) {
GrantOutcome::Idle(permit) => {
drop(permit);
return;
}
GrantOutcome::Granted(grant) => spare = Self::settle(round, grant),
}
}
}
fn settle(
round: &mut std::sync::MutexGuard<'_, DeficitRoundRobin<Arc<WaitSlot>>>,
grant: Grant<Arc<WaitSlot>, OwnedSemaphorePermit>,
) -> Option<OwnedSemaphorePermit> {
let returned = grant.waiter.deliver(grant.permit)?;
round.note_release(&grant.tenant);
Some(returned)
}
fn report_abandoned(&self, slot: &WaitSlot, tenant: &Tenant) {
let mut round = self.lock();
if slot.lock().withdrawn {
return;
}
let mut is_live = |waiter: &Arc<WaitSlot>| waiter.is_live();
round.note_abandoned(tenant, &mut is_live);
}
pub(super) fn counts_of(&self, tenant: &Tenant) -> (usize, usize) {
let round = self.lock();
(round.in_flight_of(tenant), round.queued_of(tenant))
}
#[cfg(test)]
pub(super) fn retained_waiters(&self, tenant: &Tenant) -> usize {
self.lock().retained_of(tenant)
}
pub(super) fn tenants_waiting(&self) -> usize {
self.lock().tenants_with_queues()
}
}
#[cfg(test)]
mod tests {
use std::error::Error;
use std::sync::Arc;
use crate::rt::sync::Semaphore;
use super::{Departure, TenancyShell, WaitSlot};
use crate::rt::tenancy::{Arrival, GrantOutcome, TenancyPolicy};
use crate::script::Tenant;
#[test]
fn a_late_abandonment_leaves_the_live_count_exact() -> Result<(), Box<dyn Error>> {
let shell = TenancyShell::new(Arc::new(Semaphore::new(0)), TenancyPolicy::new(4, 4));
let tenant = Tenant::new("late")?;
let leaving = WaitSlot::fresh();
let staying = [WaitSlot::fresh(), WaitSlot::fresh(), WaitSlot::fresh()];
{
let mut round = shell.lock();
for slot in std::iter::once(&leaving).chain(&staying) {
let arrival = round.arrive(&tenant, Arc::clone(slot), || None::<()>);
assert!(
matches!(arrival, Arrival::Queued),
"a spent pool queues every arrival"
);
}
}
let spare = Arc::new(Semaphore::new(1)).try_acquire_owned()?;
let mut round = shell.lock();
let mut is_live = |slot: &Arc<WaitSlot>| slot.is_live();
let grant = match round.grant(spare, &mut is_live) {
GrantOutcome::Granted(grant) => Some(grant),
GrantOutcome::Idle(_) => None,
}
.ok_or("a queued waiter takes the permit")?;
assert!(
Arc::ptr_eq(&grant.waiter, &leaving),
"the round chose the oldest waiter"
);
assert!(
matches!(leaving.leave(), Departure::Abandoned),
"the owner leaves before any permit reached it"
);
let recycled = TenancyShell::settle(&mut round, grant);
drop(round);
shell.report_abandoned(&leaving, &tenant);
assert!(
recycled.is_some(),
"the failed delivery recycled its permit"
);
let (in_flight, queued) = shell.counts_of(&tenant);
assert_eq!(
in_flight, 0,
"the failed delivery charged its admission back"
);
assert_eq!(
queued,
staying.len(),
"three owners are still waiting, and the live count must say so"
);
assert_eq!(
shell.retained_waiters(&tenant),
staying.len(),
"the round retains exactly the waiter still owned"
);
Ok(())
}
}
}
pub struct Supervisor {
token: CancellationToken,
set: JoinSet<TaskEnd>,
permits: Arc<Semaphore>,
#[cfg(all(unix, feature = "process"))]
cleanup_owners: Arc<CleanupOwners>,
identities: BTreeMap<Id, TaskId>,
reports: VecDeque<TaskOutcome>,
report_cap: usize,
max_in_flight: usize,
next_task: u64,
spawned: u64,
completed: u64,
succeeded: u64,
failed: u64,
cancelled: u64,
aborted: u64,
panicked: u64,
refused: u64,
reports_dropped: u64,
clock: Clock,
#[cfg(feature = "script")]
tenancy: Option<Arc<tenancy_support::TenancyShell>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum TaskEnd {
Completed,
#[cfg(feature = "process")]
CompletedWithCleanup {
cleanup: ProcessCleanup,
},
Cancelled,
#[cfg(feature = "process")]
CancelledWithCleanup {
cleanup: ProcessCleanup,
},
#[cfg(feature = "process")]
Failed {
status: Option<ExitStatus>,
cleanup: ProcessCleanup,
},
}
#[derive(Debug)]
enum Wait {
Joined(Result<(Id, TaskEnd), super::task::JoinError>),
Permit(OwnedSemaphorePermit),
Closed,
Cancelled,
Recheck,
}
#[cfg(feature = "process")]
#[derive(Clone, Debug, PartialEq, Eq)]
struct ProcessCleanup {
receipt: CleanupReceipt,
containment: Containment,
}
#[cfg(feature = "process")]
impl ProcessCleanup {
fn into_parts(self) -> (CleanupReceipt, Containment) {
(self.receipt, self.containment)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Retention {
Capped,
Draining,
}
const MAX_PANIC_MESSAGE_CHARS: usize = 512;
const COOPERATIVE_DRAIN_GRACE: Duration = Duration::from_millis(50);
#[cfg(feature = "script")]
const REAP_PER_ADMISSION: usize = 2;
#[cfg(all(unix, feature = "process"))]
const PROCESS_CLEANUP_GRACE: Duration = Duration::from_secs(2);
#[cfg(all(unix, feature = "process"))]
const CONTAINMENT_ROUNDS: usize = 4;
#[cfg(all(unix, feature = "process"))]
const GROUP_SIGNAL_ATTEMPTS: usize = 4;
#[cfg(all(unix, feature = "process"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum GroupSignal {
Delivered,
Absent,
Present,
Refused,
}
#[cfg(all(unix, feature = "process"))]
fn is_absence(error: &io::Error) -> bool {
error.kind() == io::ErrorKind::NotFound || error.raw_os_error() == Some(ESRCH)
}
#[cfg(all(unix, feature = "process"))]
const ESRCH: i32 = 3;
#[cfg(all(unix, feature = "process"))]
const EPERM: i32 = 1;
#[cfg(all(unix, feature = "process"))]
fn is_present_but_unsignalable(error: &io::Error) -> bool {
error.kind() == io::ErrorKind::PermissionDenied || error.raw_os_error() == Some(EPERM)
}
#[cfg(all(unix, feature = "process"))]
const CLEANUP_RECHECK: Duration = Duration::from_millis(100);
const FALLBACK_MAX_IN_FLIGHT: usize = 4;
impl Default for Supervisor {
fn default() -> Self {
let bound =
std::thread::available_parallelism().map_or(FALLBACK_MAX_IN_FLIGHT, NonZeroUsize::get);
Self::new(bound)
}
}
impl Supervisor {
#[must_use]
pub fn new(max_in_flight: usize) -> Self {
Self::with_clock(max_in_flight, Clock::wall())
}
#[must_use]
pub fn with_clock(max_in_flight: usize, clock: Clock) -> Self {
let bound = max_in_flight.clamp(1, Semaphore::MAX_PERMITS);
Self::assembled(bound, Arc::new(Semaphore::new(bound)), clock)
}
#[cfg(feature = "script")]
#[must_use]
pub fn with_tenancy(max_in_flight: usize, policy: TenancyPolicy) -> Self {
let bound = max_in_flight.clamp(1, Semaphore::MAX_PERMITS);
let permits = Arc::new(Semaphore::new(bound));
let mut supervisor = Self::assembled(bound, Arc::clone(&permits), Clock::wall());
supervisor.tenancy = Some(Arc::new(tenancy_support::TenancyShell::new(
permits, policy,
)));
supervisor
}
fn assembled(bound: usize, permits: Arc<Semaphore>, clock: Clock) -> Self {
Self {
token: CancellationToken::new(),
set: JoinSet::new(),
permits,
#[cfg(all(unix, feature = "process"))]
cleanup_owners: Arc::new(CleanupOwners::default()),
identities: BTreeMap::new(),
reports: VecDeque::new(),
report_cap: bound,
max_in_flight: bound,
next_task: 0,
spawned: 0,
completed: 0,
succeeded: 0,
failed: 0,
cancelled: 0,
aborted: 0,
panicked: 0,
refused: 0,
reports_dropped: 0,
clock,
#[cfg(feature = "script")]
tenancy: None,
}
}
#[must_use]
pub fn clock(&self) -> &Clock {
&self.clock
}
#[must_use]
pub fn child_token(&self) -> CancellationToken {
self.token.child_token()
}
#[cfg(all(unix, feature = "process"))]
fn prepare_process(&self, spec: &ProcessSpec) -> (CancellationToken, RunBounds) {
(
self.child_token(),
RunBounds {
deadline: spec.deadline_duration(),
out_limit: spec.stdout_capture(),
err_limit: spec.stderr_capture(),
},
)
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.token.is_cancelled()
}
pub fn cancel(&self) {
self.token.cancel();
}
#[must_use]
pub fn stats(&self) -> Stats {
Stats {
spawned: self.spawned,
completed: self.completed,
succeeded: self.succeeded,
failed: self.failed,
cancelled: self.cancelled,
aborted: self.aborted,
panicked: self.panicked,
refused: self.refused,
reports_dropped: self.reports_dropped,
}
}
#[cfg(all(unix, feature = "process"))]
#[must_use]
pub fn pending_cleanup_count(&self) -> usize {
self.cleanup_owners.pending_count()
}
#[must_use]
pub fn snapshot(&self) -> SupervisorSnapshot {
let live_limit = self.max_in_flight;
let mut live: Vec<LiveTask> = Vec::with_capacity(self.identities.len().min(live_limit));
let mut live_truncated: usize = 0;
let mut ordered: Vec<TaskId> = self.identities.values().copied().collect();
ordered.sort_unstable_by_key(|task| task.get());
for (index, task) in ordered.into_iter().enumerate() {
if index < live_limit {
live.push(LiveTask {
task,
state: TaskState::Running { position: index },
});
} else {
live_truncated = live_truncated.saturating_add(1);
}
}
#[cfg(all(unix, feature = "process"))]
let pending_cleanups = self.cleanup_owners.pending_count();
SupervisorSnapshot {
max_in_flight: self.max_in_flight,
stats: self.stats(),
live,
live_truncated,
cancelled: self.token.is_cancelled(),
reports_pending: self.reports.len(),
#[cfg(all(unix, feature = "process"))]
pending_cleanups,
}
}
pub fn next_report(&mut self) -> Option<TaskOutcome> {
#[cfg(all(unix, feature = "process"))]
self.refresh_cleanup_owners();
self.reports.pop_front()
}
pub fn reap(&mut self) -> usize {
self.reap_at_most(usize::MAX)
}
fn reap_at_most(&mut self, limit: usize) -> usize {
let mut reaped: usize = 0;
while reaped < limit {
let Some(joined) = self.set.try_join_next_with_id() else {
break;
};
self.absorb(joined, Retention::Capped);
reaped = reaped.saturating_add(1);
}
#[cfg(all(unix, feature = "process"))]
self.refresh_cleanup_owners();
reaped
}
pub async fn wait_idle(&mut self) -> usize {
let mut joined: usize = 0;
while let Some(next) = self.set.join_next_with_id().await {
self.absorb(next, Retention::Capped);
joined = joined.saturating_add(1);
}
#[cfg(all(unix, feature = "process"))]
self.refresh_cleanup_owners();
joined
}
pub async fn spawn<F, Fut>(&mut self, body: F)
where
F: FnOnce(CancellationToken) -> Fut,
Fut: Future<Output = ()> + Send + 'static,
{
let Some(lease) = self.claim().await else {
return;
};
self.place_body(lease, body);
}
pub fn try_spawn<F, Fut>(&mut self, body: F) -> Result<(), TrySpawnRefusal>
where
F: FnOnce(CancellationToken) -> Fut,
Fut: Future<Output = ()> + Send + 'static,
{
let Some(lease) = self.claim_now() else {
if self.token.is_cancelled() {
let refusal = Err(TrySpawnRefusal::Cancelled);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "try_spawn: returning an error to the caller");
return refusal;
}
let refusal = Err(TrySpawnRefusal::AtCapacity);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "try_spawn: returning an error to the caller");
return refusal;
};
self.place_body(lease, body);
Ok(())
}
#[cfg(feature = "script")]
pub async fn spawn_for<F, Fut>(&mut self, tenant: &Tenant, body: F) -> Result<(), SpawnRefused>
where
F: FnOnce(CancellationToken) -> Fut,
Fut: Future<Output = ()> + Send + 'static,
{
let Some(shell) = self.tenancy.clone() else {
self.spawn(body).await;
return Ok(());
};
let lease = match self.claim_tenanted(&shell, tenant).await {
Ok(lease) => lease,
Err(refusal) => {
self.refused = self.refused.saturating_add(1);
let refused = Err(refusal);
lgwks_std::trace::debug!(error = ?refused.as_ref().err(), "spawn_for: returning an error to the caller");
return refused;
}
};
self.place_body(lease, body);
Ok(())
}
#[cfg(feature = "script")]
#[must_use]
pub fn tenant_capacity(&self, tenant: &Tenant) -> (usize, usize) {
self.tenancy
.as_ref()
.map_or((0, 0), |shell| shell.counts_of(tenant))
}
#[cfg(feature = "script")]
#[must_use]
pub fn tenants_waiting(&self) -> usize {
self.tenancy
.as_ref()
.map_or(0, |shell| shell.tenants_waiting())
}
pub async fn spawn_repeating<F, Fut>(&mut self, budget: Budget, body: F)
where
F: FnMut(u64) -> Fut + Send + 'static,
Fut: Future<Output = ()> + Send + 'static,
{
let Some(permit) = self.claim().await else {
return;
};
let token = self.child_token();
let clock = self.clock.clone();
self.place(permit, async move {
match repeat_on(&clock, &token, budget, body).await {
Outcome::Exhausted { .. } => TaskEnd::Completed,
Outcome::Cancelled { .. } => TaskEnd::Cancelled,
}
});
}
#[cfg(all(unix, feature = "process"))]
pub async fn spawn_process(&mut self, spec: &ProcessSpec) -> io::Result<TaskId> {
let (task, ()) = self.spawn_process_reading(spec, |_| Ok(())).await?;
Ok(task)
}
#[cfg(all(unix, feature = "process"))]
pub async fn spawn_process_identified(
&mut self,
spec: &ProcessSpec,
) -> io::Result<IdentifiedSpawn> {
let (task, leader) = self.spawn_process_reading(spec, identify_leader).await?;
Ok(IdentifiedSpawn { task, leader })
}
#[cfg(all(unix, feature = "process"))]
async fn spawn_process_reading<T>(
&mut self,
spec: &ProcessSpec,
read: impl FnOnce(i32) -> io::Result<T>,
) -> io::Result<(TaskId, T)> {
let Some(permit) = self.claim().await else {
if self.token.is_cancelled() {
let refusal = Err(io::Error::other(SupervisorCancelled));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "spawn_process: returning an error to the caller");
return refusal;
}
let refusal = Err(io::Error::other(
"lgwks_bot: the supervisor's in-flight semaphore was closed",
));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "spawn_process: returning an error to the caller");
return refusal;
};
if self.token.is_cancelled() {
self.refused = self.refused.saturating_add(1);
let refusal = Err(io::Error::other(SupervisorCancelled));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "spawn_process: returning an error to the caller");
return refusal;
}
let (token, bounds) = self.prepare_process(spec);
let clock = self.clock.clone();
let (child, group_id) = start(spec)?;
let task = self.allocate_task_id();
let group = ProcessGroup::of(group_id, task, permit, Arc::clone(&self.cleanup_owners));
let live = LiveProcess::enter(&self.cleanup_owners);
let read = read(group_id)?;
let placed = self.place_owned(task, async move {
let end = drive_process(&clock, child, &token, group, bounds).await;
drop(live);
task_end(end)
});
Ok((placed, read))
}
#[cfg(feature = "process")]
pub async fn run_process(&mut self, spec: &ProcessSpec) -> Result<ProcessRun, ProcessRunError> {
self.run_process_observed(spec, None).await
}
#[cfg(all(unix, feature = "process"))]
pub async fn run_process_observed(
&mut self,
spec: &ProcessSpec,
on_line: Option<LineObserver<'_>>,
) -> Result<ProcessRun, ProcessRunError> {
let Some(lease) = self.claim().await else {
let refusal = Err(ProcessRunError::Refused);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "run_process_observed: returning an error to the caller");
return refusal;
};
self.drive_run(lease, spec, on_line).await
}
#[cfg(all(unix, feature = "process", feature = "script"))]
pub async fn run_process_for(
&mut self,
tenant: &Tenant,
spec: &ProcessSpec,
) -> Result<ProcessRun, RunForError> {
let Some(shell) = self.tenancy.clone() else {
return self.run_process(spec).await.map_err(RunForError::Run);
};
self.reap();
let lease = match self.claim_tenanted(&shell, tenant).await {
Ok(lease) => lease,
Err(refusal) => {
self.refused = self.refused.saturating_add(1);
let refused = Err(RunForError::Refused(refusal));
lgwks_std::trace::debug!(error = ?refused.as_ref().err(), "run_process_for: returning an error to the caller");
return refused;
}
};
self.drive_run(lease, spec, None)
.await
.map_err(RunForError::Run)
}
#[cfg(all(unix, feature = "process"))]
async fn drive_run(
&mut self,
lease: Lease,
spec: &ProcessSpec,
on_line: Option<LineObserver<'_>>,
) -> Result<ProcessRun, ProcessRunError> {
if self.token.is_cancelled() {
self.refused = self.refused.saturating_add(1);
let refusal = Err(ProcessRunError::Refused);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "drive_run: returning an error to the caller");
return refusal;
}
let (token, bounds) = self.prepare_process(spec);
let clock = self.clock.clone();
let (child, group_id) = match start(spec) {
Ok(started) => started,
Err(source) => {
let refusal = Err(ProcessRunError::NotStarted { source });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "drive_run: returning an error to the caller");
return refusal;
}
};
let task = self.allocate_task_id();
let group = ProcessGroup::of(group_id, task, lease, Arc::clone(&self.cleanup_owners));
self.spawned = self.spawned.saturating_add(1);
let end = drive_process_observed(&clock, child, &token, group, bounds, on_line).await;
self.completed = self.completed.saturating_add(1);
let deadline_fired = matches!(end.observation, ProcessObservation::Deadline);
let settled = match end.observation {
ProcessObservation::Exited => {
if end.status.is_some_and(|status| status.success()) {
self.succeeded = self.succeeded.saturating_add(1);
} else {
self.failed = self.failed.saturating_add(1);
}
true
}
ProcessObservation::Deadline | ProcessObservation::Unobservable => {
self.failed = self.failed.saturating_add(1);
true
}
ProcessObservation::Cancelled => {
self.cancelled = self.cancelled.saturating_add(1);
false
}
};
if !settled {
let refusal = Err(ProcessRunError::AfterStart {
source: io::Error::new(
io::ErrorKind::Interrupted,
"the supervisor was cancelled while the process ran",
),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "drive_run: returning an error to the caller");
return refusal;
}
let status = if matches!(end.observation, ProcessObservation::Exited) {
end.status
} else {
None
};
Ok(ProcessRun {
status,
deadline_fired,
stdout: end.stdout,
stderr: end.stderr,
cleanup: end.cleanup.receipt,
containment: end.cleanup.containment,
})
}
#[cfg(all(not(unix), feature = "process"))]
pub async fn run_process_observed(
&mut self,
spec: &ProcessSpec,
on_line: Option<LineObserver<'_>>,
) -> Result<ProcessRun, ProcessRunError> {
let _ = (spec, on_line);
Err(ProcessRunError::NotStarted {
source: Self::no_process_group(),
})
}
#[cfg(all(not(unix), feature = "process", feature = "script"))]
pub async fn run_process_for(
&mut self,
tenant: &Tenant,
spec: &ProcessSpec,
) -> Result<ProcessRun, RunForError> {
let _ = (tenant, spec);
let refusal = Err(RunForError::Run(ProcessRunError::NotStarted {
source: Self::no_process_group(),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "run_process_for: returning an error to the caller");
refusal
}
#[cfg(all(not(unix), feature = "process"))]
pub async fn spawn_process(&mut self, spec: &ProcessSpec) -> io::Result<TaskId> {
let _ = spec;
Err(Self::no_process_group())
}
#[cfg(all(not(unix), feature = "process"))]
pub async fn spawn_process_identified(
&mut self,
spec: &ProcessSpec,
) -> io::Result<IdentifiedSpawn> {
let _ = spec;
Err(Self::no_process_group())
}
#[cfg(all(not(unix), feature = "process"))]
fn no_process_group() -> io::Error {
io::Error::new(
io::ErrorKind::Unsupported,
"lgwks_bot: supervised process-group cleanup is Unix-only",
)
}
pub async fn shutdown(mut self) -> ShutdownReport {
self.token.cancel();
let watchdog = self.clock.wall_watchdog();
let deadline = watchdog.elapsed().saturating_add(COOPERATIVE_DRAIN_GRACE);
loop {
while let Some(joined) = self.set.try_join_next_with_id() {
self.absorb(joined, Retention::Draining);
}
let elapsed = watchdog.elapsed();
if self.set.is_empty()
|| (elapsed >= deadline
&& !self.cleanup_holds_grace(elapsed.saturating_sub(deadline)))
{
break;
}
yield_now().await;
}
self.set.abort_all();
while let Some(joined) = self.set.join_next_with_id().await {
self.absorb(joined, Retention::Draining);
}
#[cfg(all(unix, feature = "process"))]
self.refresh_cleanup_owners();
let stats = self.stats();
let outcomes: Vec<TaskOutcome> = self.reports.drain(..).collect();
ShutdownReport {
outcomes,
stats,
#[cfg(all(unix, feature = "process"))]
cleanup_owners: Arc::clone(&self.cleanup_owners),
}
}
#[cfg(all(unix, feature = "process"))]
fn cleanup_holds_grace(&self, overrun: Duration) -> bool {
overrun < PROCESS_CLEANUP_GRACE
&& self
.cleanup_owners
.live
.load(std::sync::atomic::Ordering::Acquire)
> 0
}
#[cfg(not(all(unix, feature = "process")))]
const fn cleanup_holds_grace(&self, _overrun: Duration) -> bool {
false
}
async fn claim(&mut self) -> Option<Lease> {
if self.token.is_cancelled() {
return self.refuse();
}
#[cfg(feature = "script")]
if let Some(shell) = self.tenancy.clone() {
return match self.claim_tenanted(&shell, &implicit_tenant()).await {
Ok(lease) => Some(lease),
Err(_) => {
self.refused = self.refused.saturating_add(1);
None
}
};
}
self.reap();
loop {
if let Some(permit) = self.try_take() {
return self.admit(permit);
}
match self.wait_at_bound().await {
Wait::Joined(joined) => self.absorb(joined, Retention::Capped),
Wait::Recheck => {
self.reap();
}
Wait::Permit(permit) => return self.admit(permit),
Wait::Closed | Wait::Cancelled => return self.refuse(),
}
}
}
#[cfg(feature = "script")]
async fn claim_tenanted(
&mut self,
shell: &Arc<tenancy_support::TenancyShell>,
tenant: &Tenant,
) -> Result<Lease, SpawnRefused> {
if self.token.is_cancelled() {
let refusal = Err(SpawnRefused::Cancelled);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "claim_tenanted: returning an error to the caller");
return refusal;
}
match shell.admit(tenant) {
tenancy_support::Admission::Admitted(lease) => {
self.reap_at_most(REAP_PER_ADMISSION);
if self.token.is_cancelled() {
Err(SpawnRefused::Cancelled)
} else {
Ok(lease)
}
}
tenancy_support::Admission::Refused { limit } => Err(SpawnRefused::TenantAtCapacity {
tenant: tenant.clone(),
limit,
}),
tenancy_support::Admission::SupervisorQueueFull { limit } => {
Err(SpawnRefused::SupervisorQueueFull { limit })
}
tenancy_support::Admission::Queued(waiter) => {
let mut waiter = std::pin::pin!(waiter);
self.reap_at_most(REAP_PER_ADMISSION);
let token = self.token.clone();
match token.run_until_cancelled(waiter.as_mut()).await {
None => Err(SpawnRefused::Cancelled),
Some(permit) if !self.token.is_cancelled() => {
Ok(Lease::tenanted(permit, Arc::clone(shell), tenant.clone()))
}
Some(permit) => {
drop(Lease::tenanted(permit, Arc::clone(shell), tenant.clone()));
Err(SpawnRefused::Cancelled)
}
}
}
}
}
fn admit(&mut self, permit: OwnedSemaphorePermit) -> Option<Lease> {
if self.token.is_cancelled() {
drop(permit);
return self.refuse();
}
Some(Lease::plain(permit))
}
const fn refuse(&mut self) -> Option<Lease> {
self.refused = self.refused.saturating_add(1);
None
}
async fn wait_at_bound(&mut self) -> Wait {
let permits = &self.permits;
let set = &mut self.set;
let mut cancelled = pin!(self.token.cancelled());
#[cfg(all(unix, feature = "process"))]
let mut recheck = pin!(recheck(&self.cleanup_owners));
#[cfg(not(all(unix, feature = "process")))]
let mut recheck = pin!(std::future::pending::<()>());
if set.is_empty() {
let mut acquire = pin!(Arc::clone(permits).acquire_owned());
return poll_fn(|context: &mut Context<'_>| {
if let Poll::Ready(acquired) = acquire.as_mut().poll(context) {
return Poll::Ready(acquired.map_or(Wait::Closed, Wait::Permit));
}
if cancelled.as_mut().poll(context).is_ready() {
return Poll::Ready(Wait::Cancelled);
}
recheck.as_mut().poll(context).map(|()| Wait::Recheck)
})
.await;
}
let mut joined = pin!(set.join_next_with_id());
poll_fn(|context: &mut Context<'_>| {
if let Poll::Ready(next) = joined.as_mut().poll(context) {
return Poll::Ready(next.map_or(Wait::Recheck, Wait::Joined));
}
if cancelled.as_mut().poll(context).is_ready() {
return Poll::Ready(Wait::Cancelled);
}
recheck.as_mut().poll(context).map(|()| Wait::Recheck)
})
.await
}
fn claim_now(&mut self) -> Option<Lease> {
if self.token.is_cancelled() {
self.refused = self.refused.saturating_add(1);
return None;
}
self.reap();
#[cfg(feature = "script")]
if let Some(shell) = self.tenancy.clone() {
return match shell.try_admit(&implicit_tenant()) {
tenancy_support::Admission::Admitted(lease) if !self.token.is_cancelled() => {
Some(lease)
}
tenancy_support::Admission::Admitted(_) => {
self.refused = self.refused.saturating_add(1);
None
}
tenancy_support::Admission::Refused { .. }
| tenancy_support::Admission::SupervisorQueueFull { .. } => {
self.refused = self.refused.saturating_add(1);
None
}
tenancy_support::Admission::Queued(_) => {
self.refused = self.refused.saturating_add(1);
None
}
};
}
if self.token.is_cancelled() {
self.refused = self.refused.saturating_add(1);
return None;
}
match self.try_take() {
Some(permit) => Some(Lease::plain(permit)),
None => {
self.refused = self.refused.saturating_add(1);
None
}
}
}
fn place_body<F, Fut>(&mut self, lease: Lease, body: F) -> TaskId
where
F: FnOnce(CancellationToken) -> Fut,
Fut: Future<Output = ()> + Send + 'static,
{
let token = self.child_token();
let future = body(token.clone());
self.place(lease, async move {
future.await;
if token.is_cancelled() {
TaskEnd::Cancelled
} else {
TaskEnd::Completed
}
})
}
fn try_take(&mut self) -> Option<OwnedSemaphorePermit> {
Arc::clone(&self.permits).try_acquire_owned().ok()
}
fn place<Fut>(&mut self, lease: Lease, future: Fut) -> TaskId
where
Fut: Future<Output = TaskEnd> + Send + 'static,
{
let task = self.allocate_task_id();
self.place_owned(task, async move {
let _lease = lease;
future.await
});
task
}
fn place_owned<Fut>(&mut self, task: TaskId, future: Fut) -> TaskId
where
Fut: Future<Output = TaskEnd> + Send + 'static,
{
self.spawned = self.spawned.saturating_add(1);
let handle = self.set.spawn(future);
self.identities.insert(handle.id(), task);
task
}
fn allocate_task_id(&mut self) -> TaskId {
let task = TaskId(self.next_task);
self.next_task = self.next_task.saturating_add(1);
task
}
#[cfg(all(unix, feature = "process"))]
fn refresh_cleanup_owners(&mut self) {
for (task, containment) in self.cleanup_owners.drive(&NATIVE_GROUP_OBSERVER, |task| {
self.identities.values().any(|live| *live == task)
}) {
if self.reports.len() < self.report_cap {
self.reports.push_back(TaskOutcome::CleanupSettled {
task,
cleanup: settled_receipt(&containment),
containment,
});
} else {
self.reports_dropped = self.reports_dropped.saturating_add(1);
}
}
}
fn absorb(
&mut self,
joined: Result<(Id, TaskEnd), super::task::JoinError>,
retention: Retention,
) {
let raw = joined
.as_ref()
.map_or_else(|error| error.id(), |pair| pair.0);
let outcome = match joined {
Ok((_, TaskEnd::Completed)) => {
self.succeeded = self.succeeded.saturating_add(1);
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Completed {
task,
cleanup: None,
#[cfg(feature = "process")]
containment: None,
})
}
#[cfg(feature = "process")]
Ok((_, TaskEnd::CompletedWithCleanup { cleanup })) => {
self.succeeded = self.succeeded.saturating_add(1);
let (receipt, containment) = cleanup.into_parts();
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Completed {
task,
cleanup: Some(receipt),
#[cfg(feature = "process")]
containment: Some(containment),
})
}
Ok((_, TaskEnd::Cancelled)) => {
self.cancelled = self.cancelled.saturating_add(1);
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Cancelled {
task,
cleanup: None,
#[cfg(feature = "process")]
containment: None,
})
}
#[cfg(feature = "process")]
Ok((_, TaskEnd::CancelledWithCleanup { cleanup })) => {
self.cancelled = self.cancelled.saturating_add(1);
let (receipt, containment) = cleanup.into_parts();
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Cancelled {
task,
cleanup: Some(receipt),
#[cfg(feature = "process")]
containment: Some(containment),
})
}
#[cfg(feature = "process")]
Ok((_, TaskEnd::Failed { status, cleanup })) => {
let (receipt, containment) = cleanup.into_parts();
self.failed = self.failed.saturating_add(1);
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Failed {
task,
status,
cleanup: Some(receipt),
#[cfg(feature = "process")]
containment: Some(containment),
})
}
Err(error) if error.is_panic() => {
self.panicked = self.panicked.saturating_add(1);
let message = panic_message(error.into_panic());
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Panicked { task, message })
}
Err(_) => {
self.aborted = self.aborted.saturating_add(1);
self.identities
.remove(&raw)
.map(|task| TaskOutcome::Aborted { task })
}
};
self.completed = self.completed.saturating_add(1);
let Some(outcome) = outcome else {
self.reports_dropped = self.reports_dropped.saturating_add(1);
return;
};
match retention {
Retention::Capped if self.reports.len() >= self.report_cap => {
self.reports_dropped = self.reports_dropped.saturating_add(1);
}
Retention::Capped | Retention::Draining => self.reports.push_back(outcome),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum TaskState {
Running {
position: usize,
},
}
impl TaskState {
#[must_use]
pub const fn is_running(self) -> bool {
matches!(self, Self::Running { .. })
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct SupervisorSnapshot {
max_in_flight: usize,
stats: Stats,
live: Vec<LiveTask>,
live_truncated: usize,
cancelled: bool,
reports_pending: usize,
#[cfg(all(unix, feature = "process"))]
pending_cleanups: usize,
}
impl SupervisorSnapshot {
#[must_use]
pub const fn max_in_flight(&self) -> usize {
self.max_in_flight
}
#[must_use]
pub const fn stats(&self) -> Stats {
self.stats
}
#[must_use]
pub fn live(&self) -> &[LiveTask] {
&self.live
}
#[must_use]
pub const fn live_truncated(&self) -> usize {
self.live_truncated
}
#[must_use]
pub const fn is_cancelled(&self) -> bool {
self.cancelled
}
#[must_use]
pub const fn reports_pending(&self) -> usize {
self.reports_pending
}
#[cfg(all(unix, feature = "process"))]
#[must_use]
pub const fn pending_cleanups(&self) -> usize {
self.pending_cleanups
}
#[must_use]
pub const fn live_limit(&self) -> usize {
self.max_in_flight
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.live.is_empty() && self.live_truncated == 0
}
#[must_use]
pub fn occupied(&self) -> usize {
self.live
.iter()
.filter(|task| task.state.is_running())
.count()
}
#[must_use]
pub fn free(&self) -> usize {
self.max_in_flight.saturating_sub(self.occupied())
}
#[must_use]
pub fn next_action(&self) -> NextAction {
if self.cancelled {
NextAction::Stopped
} else if self.free() > 0 {
NextAction::Admit
} else {
NextAction::WaitForCapacity
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum NextAction {
Admit,
WaitForCapacity,
Stopped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct LiveTask {
pub task: TaskId,
pub state: TaskState,
}
impl Drop for Supervisor {
fn drop(&mut self) {
self.token.cancel();
self.set.abort_all();
}
}
#[cfg(all(unix, feature = "process"))]
fn start(spec: &ProcessSpec) -> io::Result<(OwnedChild, i32)> {
let mut command = Command::new(spec.program());
spec.configure(&mut command)?;
command.process_group(0);
let child = OwnedChild::own(command.into_std().spawn()?)?;
let group = child.group()?;
Ok((child, group))
}
#[cfg(all(unix, feature = "process"))]
const REAP_POLL: Duration = Duration::from_millis(5);
#[cfg(all(unix, feature = "process"))]
static ORPHANS: std::sync::Mutex<Vec<std::process::Child>> = std::sync::Mutex::new(Vec::new());
#[cfg(all(unix, feature = "process"))]
fn reap_orphans() {
crate::journal::owner::lock(&ORPHANS)
.retain_mut(|orphan| matches!(orphan.try_wait(), Ok(None)));
}
#[cfg(all(unix, feature = "process"))]
struct OwnedChild {
inner: Option<std::process::Child>,
pid: u32,
stdout: Option<ChildStdout>,
stderr: Option<ChildStderr>,
}
#[cfg(all(unix, feature = "process"))]
impl OwnedChild {
fn own(mut inner: std::process::Child) -> io::Result<Self> {
reap_orphans();
let stdout = inner.stdout.take();
let stderr = inner.stderr.take();
let mut child = Self {
pid: inner.id(),
inner: Some(inner),
stdout: None,
stderr: None,
};
child.stdout = stdout.map(ChildStdout::from_std).transpose()?;
child.stderr = stderr.map(ChildStderr::from_std).transpose()?;
Ok(child)
}
fn group(&self) -> io::Result<i32> {
i32::try_from(self.pid).map_err(io::Error::other)
}
async fn reap(&mut self) -> io::Result<std::process::ExitStatus> {
loop {
let Some(inner) = self.inner.as_mut() else {
let refusal = Err(io::Error::other("the child was already reaped"));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), pid = self.pid, "owned child: reap after reap");
return refusal;
};
if let Some(status) = inner.try_wait()? {
self.inner = None;
reap_orphans();
return Ok(status);
}
crate::rt::time::sleep(REAP_POLL).await;
}
}
}
#[cfg(all(unix, feature = "process"))]
impl Drop for OwnedChild {
fn drop(&mut self) {
let Some(mut inner) = self.inner.take() else {
return;
};
if let Err(error) = inner.kill() {
lgwks_std::trace::debug!(%error, pid = self.pid, "owned child: direct kill on drop refused");
}
if !matches!(inner.try_wait(), Ok(Some(_))) {
crate::journal::owner::lock(&ORPHANS).push(inner);
}
reap_orphans();
}
}
#[cfg(all(unix, feature = "process"))]
struct ProcessGroup<'ops> {
group: i32,
task: TaskId,
permit: Option<Lease>,
owners: Arc<CleanupOwners>,
armed: bool,
leader_reaped: bool,
leader_exited: bool,
group_absent: bool,
containment: Containment,
signalled_pids: std::collections::BTreeSet<i32>,
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
capture: &'ops dyn DescendantCapture,
}
#[cfg(all(unix, feature = "process"))]
trait GroupSignaller: Sync {
fn signal(&self, group: i32) -> io::Result<()>;
}
#[cfg(all(unix, feature = "process"))]
trait GroupObserver: Sync {
fn exists(&self, group: i32) -> io::Result<bool>;
}
#[cfg(all(unix, feature = "process"))]
#[derive(Clone, Debug, Default, Eq, PartialEq)]
struct Capture {
mechanism: ContainmentMechanism,
pids: Vec<i32>,
truncated: bool,
}
#[cfg(all(unix, feature = "process"))]
impl Capture {
fn from_set(set: &lgwks_std::process::DescendantSet) -> Self {
Self {
mechanism: set.mechanism(),
pids: set.pids().to_vec(),
truncated: set.is_truncated(),
}
}
fn is_empty(&self) -> bool {
self.pids.is_empty()
}
}
#[cfg(all(unix, feature = "process"))]
trait DescendantCapture: Sync {
fn capture(&self, root: i32) -> io::Result<Capture>;
fn signal(&self, pid: i32) -> io::Result<()>;
fn running(&self, pids: &[i32]) -> io::Result<Vec<i32>>;
}
#[cfg(all(unix, feature = "process"))]
struct PendingCleanup {
task: TaskId,
group: i32,
_lease: Lease,
absence_observed: bool,
containment: Containment,
}
#[cfg(all(unix, feature = "process"))]
#[derive(Default)]
struct CleanupOwners {
pending: Mutex<Vec<PendingCleanup>>,
registered: Notify,
live: std::sync::atomic::AtomicUsize,
}
#[cfg(all(unix, feature = "process"))]
struct LiveProcess(Arc<CleanupOwners>);
#[cfg(all(unix, feature = "process"))]
impl LiveProcess {
fn enter(owners: &Arc<CleanupOwners>) -> Self {
owners
.live
.fetch_add(1, std::sync::atomic::Ordering::AcqRel);
Self(Arc::clone(owners))
}
}
#[cfg(all(unix, feature = "process"))]
impl Drop for LiveProcess {
fn drop(&mut self) {
self.0
.live
.fetch_sub(1, std::sync::atomic::Ordering::AcqRel);
}
}
#[cfg(all(unix, feature = "process"))]
impl CleanupOwners {
fn register(&self, task: TaskId, group: i32, lease: Lease, containment: Containment) {
crate::journal::owner::lock(&self.pending).push(PendingCleanup {
task,
group,
_lease: lease,
absence_observed: false,
containment,
});
self.registered.notify_one();
}
fn pending_count(&self) -> usize {
crate::journal::owner::lock(&self.pending).len()
}
fn drive(
&self,
observer: &dyn GroupObserver,
task_is_live: impl Fn(TaskId) -> bool,
) -> Vec<(TaskId, Containment)> {
let mut pending = crate::journal::owner::lock(&self.pending);
let mut retained = Vec::with_capacity(pending.len());
let mut settled = Vec::new();
for mut owner in pending.drain(..) {
let absent =
owner.absence_observed || matches!(observer.exists(owner.group), Ok(false));
if absent {
if task_is_live(owner.task) {
owner.absence_observed = true;
retained.push(owner);
} else {
settled.push((owner.task, owner.containment));
}
} else {
retained.push(owner);
}
}
*pending = retained;
settled
}
}
#[cfg(all(unix, feature = "process"))]
async fn recheck(owners: &CleanupOwners) {
if owners.pending_count() > 0 {
crate::rt::time::sleep(CLEANUP_RECHECK).await;
} else {
owners.registered.notified().await;
}
}
#[cfg(all(unix, feature = "process"))]
impl fmt::Debug for CleanupOwners {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("CleanupOwners")
.field("pending", &self.pending_count())
.finish()
}
}
#[cfg(feature = "process")]
#[cfg(all(unix, feature = "process"))]
struct NativeGroupSignaller;
#[cfg(all(unix, feature = "process"))]
impl GroupSignaller for NativeGroupSignaller {
fn signal(&self, group: i32) -> io::Result<()> {
lgwks_std::process::kill_process_group(group)
}
}
#[cfg(all(unix, feature = "process"))]
struct NativeGroupObserver;
#[cfg(all(unix, feature = "process"))]
impl GroupObserver for NativeGroupObserver {
fn exists(&self, group: i32) -> io::Result<bool> {
lgwks_std::process::process_group_exists(group)
}
}
#[cfg(feature = "process")]
#[cfg(all(unix, feature = "process"))]
static NATIVE_GROUP_SIGNALLER: NativeGroupSignaller = NativeGroupSignaller;
#[cfg(all(unix, feature = "process"))]
static NATIVE_GROUP_OBSERVER: NativeGroupObserver = NativeGroupObserver;
#[cfg(all(unix, feature = "process"))]
struct NativeDescendantCapture;
#[cfg(all(unix, feature = "process"))]
impl DescendantCapture for NativeDescendantCapture {
fn capture(&self, root: i32) -> io::Result<Capture> {
lgwks_std::process::capture_descendants(root).map(|set| Capture::from_set(&set))
}
fn signal(&self, pid: i32) -> io::Result<()> {
lgwks_std::process::kill_process(pid)
}
fn running(&self, pids: &[i32]) -> io::Result<Vec<i32>> {
Ok(lgwks_std::process::running_processes(pids)?
.into_iter()
.collect())
}
}
#[cfg(all(unix, feature = "process"))]
static NATIVE_DESCENDANT_CAPTURE: NativeDescendantCapture = NativeDescendantCapture;
#[cfg(all(unix, feature = "process"))]
fn settled_receipt(containment: &Containment) -> CleanupReceipt {
if containment.survivors().is_empty() {
CleanupReceipt::CleanupConfirmed
} else {
CleanupReceipt::CleanupSurvivors {
survivors: containment.survivors().to_vec(),
}
}
}
#[derive(Clone, Copy)]
#[cfg(all(unix, feature = "process"))]
enum ProcessObservation {
Exited,
Cancelled,
Deadline,
Unobservable,
}
#[cfg(all(unix, feature = "process"))]
async fn observe_pid_without_reaping(
clock: &Clock,
pid: i32,
deadline: Option<Duration>,
token: &CancellationToken,
) -> ProcessObservation {
if pid <= 0 {
return ProcessObservation::Unobservable;
}
let started_at = clock.now();
let deadline = deadline.map(|duration| started_at.saturating_add(duration));
loop {
if token.is_cancelled() {
return ProcessObservation::Cancelled;
}
if deadline.is_some_and(|at| clock.now() >= at) {
return ProcessObservation::Deadline;
}
match lgwks_std::process::child_has_exited_without_reaping(pid) {
Ok(true) => return ProcessObservation::Exited,
Ok(false) => {}
Err(_) => return ProcessObservation::Unobservable,
}
if token
.run_until_cancelled(crate::rt::time::sleep(Duration::from_millis(5)))
.await
.is_none()
{
return ProcessObservation::Cancelled;
}
}
}
#[cfg(feature = "process")]
pub type LineObserver<'a> = &'a (dyn Fn(&[u8]) + Send + Sync + 'a);
#[cfg(all(unix, feature = "process"))]
struct ProcessEnd {
observation: ProcessObservation,
status: Option<ExitStatus>,
cleanup: ProcessCleanup,
stdout: CapturedStream,
stderr: CapturedStream,
}
#[cfg(all(unix, feature = "process"))]
fn task_end(end: ProcessEnd) -> TaskEnd {
let ProcessEnd {
observation,
status,
cleanup,
..
} = end;
match observation {
ProcessObservation::Exited => match status {
Some(ref observed) if observed.success() => TaskEnd::CompletedWithCleanup { cleanup },
status => TaskEnd::Failed { status, cleanup },
},
ProcessObservation::Cancelled => TaskEnd::CancelledWithCleanup { cleanup },
ProcessObservation::Deadline | ProcessObservation::Unobservable => TaskEnd::Failed {
status: None,
cleanup,
},
}
}
#[cfg(all(unix, feature = "process"))]
async fn drive_process(
clock: &Clock,
child: OwnedChild,
token: &CancellationToken,
group: ProcessGroup<'static>,
bounds: RunBounds,
) -> ProcessEnd {
drive_process_observed(clock, child, token, group, bounds, None).await
}
#[cfg(all(unix, feature = "process"))]
#[derive(Debug, Clone, Copy)]
struct RunBounds {
deadline: Option<Duration>,
out_limit: Option<NonZeroUsize>,
err_limit: Option<NonZeroUsize>,
}
#[cfg(all(unix, feature = "process"))]
async fn drive_process_observed(
clock: &Clock,
mut child: OwnedChild,
token: &CancellationToken,
mut group: ProcessGroup<'static>,
bounds: RunBounds,
on_line: Option<LineObserver<'_>>,
) -> ProcessEnd {
let RunBounds {
deadline,
out_limit,
err_limit,
} = bounds;
let pid = child.group().ok();
let capture_out = capture(child.stdout.take(), out_limit, on_line);
let capture_err = capture(child.stderr.take(), err_limit, None);
let wait = async {
let observation = match pid {
Some(pid) => observe_pid_without_reaping(clock, pid, deadline, token).await,
None => ProcessObservation::Unobservable,
};
if matches!(observation, ProcessObservation::Exited) {
group.mark_exited();
}
let cleanup = group.cleanup().await;
(observation, cleanup)
};
let (stdout, stderr, (observation, mut receipt)) = join3(capture_out, capture_err, wait).await;
if matches!(receipt, CleanupReceipt::CleanupFailed)
&& !matches!(observation, ProcessObservation::Exited)
{
let containment = group.containment().clone();
drop(group);
return ProcessEnd {
observation,
status: None,
cleanup: ProcessCleanup {
receipt,
containment,
},
stdout,
stderr,
};
}
let (status, cleanup) = if matches!(receipt, CleanupReceipt::CleanupFailed) {
let containment = group.containment().clone();
drop(group);
let status = child.reap().await.ok();
(
status,
ProcessCleanup {
receipt,
containment,
},
)
} else {
let status = child.reap().await.ok();
group.mark_reaped();
if !matches!(receipt, CleanupReceipt::CleanupConfirmed) {
receipt = group.confirm_absence().await;
}
let containment = group.containment().clone();
(
status,
ProcessCleanup {
receipt,
containment,
},
)
};
ProcessEnd {
observation,
status,
cleanup,
stdout,
stderr,
}
}
#[cfg(all(unix, feature = "process"))]
async fn capture<R>(
reader: Option<R>,
limit: Option<NonZeroUsize>,
on_line: Option<LineObserver<'_>>,
) -> CapturedStream
where
R: AsyncRead + Unpin,
{
let (Some(mut reader), Some(limit)) = (reader, limit) else {
return CapturedStream::default();
};
let cap = limit.get();
let mut bytes: Vec<u8> = Vec::with_capacity(cap);
let mut total: u64 = 0;
let mut truncated = false;
let mut buffer = [0_u8; 8192];
let mut fragment: Vec<u8> = Vec::new();
let mut dropped_from_fragment: u64 = 0;
loop {
match reader.read(&mut buffer).await {
Ok(0) => break,
Ok(read) => {
total = match u64::try_from(read) {
Ok(count) => total.saturating_add(count),
Err(_) => u64::MAX,
};
let Some(chunk) = buffer.get(..read) else {
continue;
};
if let Some(observe) = on_line {
for line in lines_of(chunk, &mut fragment, &mut dropped_from_fragment) {
observe(line);
}
}
let room = cap.saturating_sub(bytes.len());
truncated |= read > room;
if room > 0 {
let take = room.min(read);
if let Some(head) = buffer.get(..take) {
bytes.extend_from_slice(head);
}
}
}
Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
Err(_) => break,
}
}
CapturedStream::from_parts(bytes, total, truncated)
}
#[cfg(all(unix, feature = "process"))]
const MAX_OBSERVED_LINE_BYTES: usize = 8 * 1024;
#[cfg(all(unix, feature = "process"))]
fn lines_of<'a>(
chunk: &'a [u8],
fragment: &'a mut Vec<u8>,
dropped: &mut u64,
) -> impl Iterator<Item = &'a [u8]> {
let mut complete = Vec::new();
for (index, byte) in chunk.iter().enumerate() {
if *byte != b'\n' {
if fragment.len() < MAX_OBSERVED_LINE_BYTES {
fragment.push(*byte);
} else {
*dropped = dropped.saturating_add(1);
}
continue;
}
let end = index.saturating_sub(fragment.len());
if let Some(line) = chunk.get(end..index) {
complete.push(strip_carriage_return(line));
}
fragment.clear();
}
complete.into_iter()
}
#[cfg(all(unix, feature = "process"))]
fn strip_carriage_return(line: &[u8]) -> &[u8] {
match line.split_last() {
Some((&b'\r', head)) => head,
_ => line,
}
}
#[cfg(all(unix, feature = "process"))]
async fn join3<A, B, C>(first: A, second: B, third: C) -> (A::Output, B::Output, C::Output)
where
A: Future,
B: Future,
C: Future,
{
let mut first = std::pin::pin!(first);
let mut second = std::pin::pin!(second);
let mut third = std::pin::pin!(third);
let mut first_out: Option<A::Output> = None;
let mut second_out: Option<B::Output> = None;
let mut third_out: Option<C::Output> = None;
std::future::poll_fn(|context| {
if first_out.is_none()
&& let std::task::Poll::Ready(value) = first.as_mut().poll(context)
{
first_out = Some(value);
}
if second_out.is_none()
&& let std::task::Poll::Ready(value) = second.as_mut().poll(context)
{
second_out = Some(value);
}
if third_out.is_none()
&& let std::task::Poll::Ready(value) = third.as_mut().poll(context)
{
third_out = Some(value);
}
if first_out.is_some() && second_out.is_some() && third_out.is_some() {
let first_value = first_out.take();
let second_value = second_out.take();
let third_value = third_out.take();
return match (first_value, second_value, third_value) {
(Some(first_value), Some(second_value), Some(third_value)) => {
std::task::Poll::Ready((first_value, second_value, third_value))
}
_ => std::task::Poll::Pending,
};
}
std::task::Poll::Pending
})
.await
}
#[cfg(all(unix, feature = "process"))]
impl<'ops> ProcessGroup<'ops> {
fn of(
group: i32,
task: TaskId,
lease: Lease,
owners: Arc<CleanupOwners>,
) -> ProcessGroup<'static> {
ProcessGroup {
group,
task,
permit: Some(lease),
owners,
armed: true,
leader_reaped: false,
leader_exited: false,
group_absent: false,
containment: Containment::default(),
signalled_pids: std::collections::BTreeSet::new(),
signaller: &NATIVE_GROUP_SIGNALLER,
observer: &NATIVE_GROUP_OBSERVER,
capture: &NATIVE_DESCENDANT_CAPTURE,
}
}
async fn cleanup(&mut self) -> CleanupReceipt {
if self.leader_reaped {
return CleanupReceipt::CleanupFailed;
}
if self.group <= 0 {
return CleanupReceipt::CleanupFailed;
}
for round in 0..CONTAINMENT_ROUNDS {
let tree = self.read_tree();
if matches!(self.signal_group(), GroupSignal::Refused) {
return CleanupReceipt::CleanupFailed;
}
if let Some(ref captured) = tree {
self.signal_captured(captured);
self.observe_captured(captured);
}
if self.no_survivors() {
break;
}
if round.saturating_add(1) < CONTAINMENT_ROUNDS {
yield_now().await;
}
}
self.receipt()
}
fn read_tree(&mut self) -> Option<Capture> {
if self.leader_exited {
return None;
}
match self.capture.capture(self.group) {
Ok(set) => {
self.containment.record_capture(&set);
Some(set)
}
Err(_) => {
self.containment.note(ResidualRisk::TableUnreadable);
None
}
}
}
fn signal_group(&mut self) -> GroupSignal {
match self.signaller.signal(self.group) {
Ok(()) => GroupSignal::Delivered,
Err(error) if is_absence(&error) => {
self.group_absent = true;
GroupSignal::Absent
}
Err(error) if is_present_but_unsignalable(&error) => GroupSignal::Present,
Err(_) => GroupSignal::Refused,
}
}
fn signal_captured(&mut self, captured: &Capture) {
let observed = self.containment.survivors.clone();
for pid in &captured.pids {
let fresh = !self.signalled_pids.contains(pid);
if !fresh && !observed.contains(pid) {
continue;
}
match self.capture.signal(*pid) {
Ok(()) => {
self.signalled_pids.insert(*pid);
self.containment.record_signal();
}
Err(error) if is_absence(&error) => {
self.signalled_pids.remove(pid);
}
Err(_) => {}
}
}
}
fn observe_captured(&mut self, captured: &Capture) {
if captured.is_empty() {
self.containment.record_survivors(&[]);
return;
}
self.observe_running(&captured.pids);
}
fn observe_running(&mut self, pids: &[i32]) {
match self.capture.running(pids) {
Ok(running) => self.containment.record_survivors(&running),
Err(_) => self.containment.record_survivors(pids),
}
}
fn no_survivors(&self) -> bool {
self.containment.survivors.is_empty()
}
fn receipt(&mut self) -> CleanupReceipt {
if !self.containment.survivors.is_empty() {
return CleanupReceipt::CleanupSurvivors {
survivors: self.containment.survivors.clone(),
};
}
if self.group_absent {
self.disarm();
return CleanupReceipt::CleanupConfirmed;
}
CleanupReceipt::CleanupPending
}
async fn confirm_absence(&mut self) -> CleanupReceipt {
if !self.leader_reaped || self.group <= 0 {
return CleanupReceipt::CleanupFailed;
}
for round in 0..CONTAINMENT_ROUNDS {
if matches!(self.observer.exists(self.group), Ok(false)) {
self.group_absent = true;
}
if !self.no_survivors() {
let survivors = self.containment.survivors.clone();
self.observe_running(&survivors);
}
if self.group_absent && self.no_survivors() {
return self.receipt();
}
if round.saturating_add(1) < CONTAINMENT_ROUNDS {
yield_now().await;
}
}
self.receipt()
}
fn kill(&mut self) {
if self.leader_reaped || self.group <= 0 {
return;
}
if let Some(ref captured) = self.read_tree() {
for pid in &captured.pids {
if self.capture.signal(*pid).is_ok() {
self.signalled_pids.insert(*pid);
self.containment.record_signal();
}
}
}
for _ in 0..GROUP_SIGNAL_ATTEMPTS {
match self.signal_group() {
GroupSignal::Absent => break,
GroupSignal::Delivered | GroupSignal::Present => {
std::thread::yield_now();
}
GroupSignal::Refused => break,
}
}
}
fn containment(&self) -> &Containment {
&self.containment
}
fn disarm(&mut self) {
self.armed = false;
}
fn mark_reaped(&mut self) {
self.leader_reaped = true;
}
fn mark_exited(&mut self) {
self.leader_exited = true;
self.containment.note(ResidualRisk::LeaderExited);
}
}
#[cfg(all(unix, feature = "process"))]
impl Drop for ProcessGroup<'_> {
fn drop(&mut self) {
if self.armed && self.group > 0 {
if !self.leader_reaped {
self.kill();
}
if let Some(permit) = self.permit.take() {
self.owners
.register(self.task, self.group, permit, self.containment.clone());
}
}
}
}
fn panic_message(payload: Box<dyn Any + Send + 'static>) -> String {
if let Some(text) = payload.downcast_ref::<&'static str>() {
return truncate_panic_message(text);
}
if let Some(text) = payload.downcast_ref::<String>() {
return truncate_panic_message(text);
}
String::from("<panic payload was neither &str nor String>")
}
fn truncate_panic_message(text: &str) -> String {
let mut kept: String = text.chars().take(MAX_PANIC_MESSAGE_CHARS).collect();
if kept.len() < text.len() {
kept.push('…');
}
kept
}
pub async fn repeat<F, Fut>(token: &CancellationToken, budget: Budget, body: F) -> Outcome
where
F: FnMut(u64) -> Fut,
Fut: Future<Output = ()>,
{
repeat_on(&Clock::wall(), token, budget, body).await
}
pub async fn repeat_on<F, Fut>(
clock: &Clock,
token: &CancellationToken,
budget: Budget,
mut body: F,
) -> Outcome
where
F: FnMut(u64) -> Fut,
Fut: Future<Output = ()>,
{
let iteration_limit = match budget {
Budget::Iterations(limit) => Some(limit.get()),
Budget::For(_) | Budget::Ongoing => None,
};
let started_at = clock.now();
let deadline = match budget {
Budget::For(limit) => Some(started_at.saturating_add(limit)),
Budget::Iterations(_) | Budget::Ongoing => None,
};
let mut iterations: u64 = 0;
let mut until_yield: u64 = YIELD_INTERVAL;
loop {
if token.is_cancelled() {
return Outcome::Cancelled { iterations };
}
let spent = iteration_limit.is_some_and(|limit| iterations >= limit)
|| deadline.is_some_and(|deadline| clock.now() >= deadline);
if spent {
return Outcome::Exhausted { iterations };
}
if let Some(deadline) = deadline {
if clock.source() != TimeSource::Wall {
match wait_for_bound(clock, token, deadline, body(iterations)).await {
Bounded::Spent => return Outcome::Exhausted { iterations },
Bounded::Stopped => return Outcome::Cancelled { iterations },
Bounded::Finished => iterations = iterations.saturating_add(1),
}
} else {
let remaining = deadline.saturating_sub(clock.now());
match crate::rt::time::timeout(
remaining,
token.run_until_cancelled(body(iterations)),
)
.await
{
Err(_elapsed) => return Outcome::Exhausted { iterations },
Ok(None) => return Outcome::Cancelled { iterations },
Ok(Some(())) => iterations = iterations.saturating_add(1),
}
continue;
}
} else {
match token.run_until_cancelled(body(iterations)).await {
Some(()) => iterations = iterations.saturating_add(1),
None => return Outcome::Cancelled { iterations },
}
}
until_yield = until_yield.saturating_sub(1);
if until_yield == 0 {
until_yield = YIELD_INTERVAL;
yield_now().await;
}
}
}
enum Bounded {
Finished,
Spent,
Stopped,
}
async fn wait_for_bound<Fut>(
clock: &Clock,
token: &CancellationToken,
deadline: Duration,
body: Fut,
) -> Bounded
where
Fut: Future<Output = ()>,
{
const LOGICAL_POLL: Duration = Duration::from_millis(1);
let mut body = std::pin::pin!(body);
loop {
if clock.now() >= deadline {
return Bounded::Spent;
}
if crate::rt::time::timeout(LOGICAL_POLL, &mut body)
.await
.is_ok()
{
return Bounded::Finished;
}
if token.is_cancelled() {
return Bounded::Stopped;
}
}
}
impl core::fmt::Debug for Supervisor {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("Supervisor")
.field("stats", &self.stats())
.field("available_permits", &self.permits.available_permits())
.field("tracked", &self.set.len())
.field("reports_queued", &self.reports.len())
.field("report_cap", &self.report_cap)
.field("cancelled", &self.token.is_cancelled())
.finish()
}
}
#[cfg(test)]
mod tests {
#[cfg(all(unix, feature = "process"))]
use super::CleanupOwners;
#[cfg(feature = "process")]
use super::CleanupReceipt;
#[cfg(all(not(unix), feature = "process"))]
use super::ProcessSpec;
#[cfg(feature = "script")]
use super::SpawnRefused;
#[cfg(all(unix, feature = "process"))]
use super::SupervisorCancelled;
use super::{
Budget, Clock, Lease, MAX_PANIC_MESSAGE_CHARS, Outcome, Supervisor, TaskId, TaskOutcome,
TrySpawnRefusal, repeat, truncate_panic_message,
};
#[cfg(all(unix, feature = "process"))]
use super::{
CONTAINMENT_ROUNDS, Containment, ContainmentMechanism, DescendantCapture, EPERM, ESRCH,
GroupObserver, GroupSignaller, ProcessGroup, ResidualRisk,
};
use crate::rt::cancel::CancellationToken;
use crate::rt::runtime::block_on;
use crate::rt::task::yield_now;
use std::future::Future;
use std::future::pending;
use std::num::NonZeroU64;
use std::sync::Arc;
#[cfg(any(feature = "process", feature = "script"))]
use std::sync::atomic::AtomicUsize;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::task::Poll;
use std::time::Duration;
#[cfg(all(not(unix), feature = "process"))]
#[test]
fn non_unix_process_group_spawn_is_refused_as_unsupported() {
let mut supervisor = Supervisor::new(1);
let result = block_on(supervisor.spawn_process(&ProcessSpec::new("unused")));
assert_eq!(
result.as_ref().err().map(std::io::Error::kind),
Some(std::io::ErrorKind::Unsupported),
"non-Unix callers receive an explicit safe refusal before starting a process"
);
assert_eq!(supervisor.stats().spawned, 0);
}
#[cfg(all(not(unix), feature = "process"))]
#[test]
fn non_unix_process_group_run_is_refused_as_unsupported() {
let mut supervisor = Supervisor::new(1);
let result = block_on(supervisor.run_process(&ProcessSpec::new("unused")));
assert!(
result.as_ref().err().is_some_and(|error| match *error {
crate::rt::process::ProcessRunError::NotStarted { ref source } => {
source.kind() == std::io::ErrorKind::Unsupported
}
_ => false,
}),
"the result-bearing runner must refuse on non-Unix before the fork, got {result:?}"
);
assert_eq!(
supervisor.stats().spawned,
0,
"a refused run must start nothing"
);
}
fn explode(message: &'static str) {
std::panic::resume_unwind(Box::new(message));
}
fn budget_of(iterations: u64) -> Budget {
match NonZeroU64::new(iterations) {
Some(limit) => Budget::Iterations(limit),
None => Budget::Ongoing,
}
}
#[cfg(all(unix, feature = "process"))]
struct RecordingGroupSignaller {
calls: AtomicUsize,
}
#[cfg(all(unix, feature = "process"))]
impl GroupSignaller for RecordingGroupSignaller {
fn signal(&self, _group: i32) -> std::io::Result<()> {
self.calls.fetch_add(1, Ordering::Relaxed);
Ok(())
}
}
#[cfg(all(unix, feature = "process"))]
struct EpermGroupSignaller;
#[cfg(all(unix, feature = "process"))]
impl GroupSignaller for EpermGroupSignaller {
fn signal(&self, _group: i32) -> std::io::Result<()> {
Err(std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"the group leader is a zombie",
))
}
}
#[cfg(all(unix, feature = "process"))]
struct UnexpectedGroupSignaller;
#[cfg(all(unix, feature = "process"))]
impl GroupSignaller for UnexpectedGroupSignaller {
fn signal(&self, _group: i32) -> std::io::Result<()> {
Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"injected unexpected failure",
))
}
}
#[cfg(all(unix, feature = "process"))]
struct SequenceGroupObserver {
calls: AtomicUsize,
present_before_absent: usize,
}
#[cfg(all(unix, feature = "process"))]
impl GroupObserver for SequenceGroupObserver {
fn exists(&self, _group: i32) -> std::io::Result<bool> {
let call = self.calls.fetch_add(1, Ordering::Relaxed);
Ok(call < self.present_before_absent)
}
}
#[cfg(all(unix, feature = "process"))]
#[derive(Default)]
struct ScriptedCapture {
tree: Vec<i32>,
gone: Vec<i32>,
immortal: Vec<i32>,
unreadable: bool,
truncated: bool,
reads: AtomicUsize,
signals: AtomicUsize,
}
#[cfg(all(unix, feature = "process"))]
impl ScriptedCapture {
fn of(tree: &[i32]) -> Self {
Self {
tree: tree.to_vec(),
..Self::default()
}
}
fn reads(&self) -> usize {
self.reads.load(Ordering::Relaxed)
}
fn signals(&self) -> usize {
self.signals.load(Ordering::Relaxed)
}
}
#[cfg(all(unix, feature = "process"))]
impl DescendantCapture for ScriptedCapture {
fn capture(&self, _root: i32) -> std::io::Result<super::Capture> {
self.reads.fetch_add(1, Ordering::Relaxed);
if self.unreadable {
let refusal = Err(std::io::Error::other("injected unreadable process table"));
lgwks_std::trace::debug!(
error = ?refusal.as_ref().err(),
"scripted capture: returning an error to the caller"
);
return refusal;
}
Ok(super::Capture {
mechanism: ContainmentMechanism::ProcessTableSnapshot,
pids: self.tree.clone(),
truncated: self.truncated,
})
}
fn signal(&self, pid: i32) -> std::io::Result<()> {
if self.gone.contains(&pid) {
let refusal = Err(std::io::Error::from_raw_os_error(3));
lgwks_std::trace::debug!(
pid,
error = ?refusal.as_ref().err(),
"scripted signal: returning an error to the caller"
);
return refusal;
}
self.signals.fetch_add(1, Ordering::Relaxed);
Ok(())
}
fn running(&self, pids: &[i32]) -> std::io::Result<Vec<i32>> {
Ok(pids
.iter()
.copied()
.filter(|pid| self.immortal.contains(pid))
.collect())
}
}
#[cfg(all(unix, feature = "process"))]
fn whole_capture() -> Containment {
Containment {
mechanism: ContainmentMechanism::ProcessTableSnapshot,
captured: 0,
signalled: 0,
survivors: Vec::new(),
residual: None,
}
}
#[cfg(all(unix, feature = "process"))]
fn settled_tasks(settled: &[(TaskId, Containment)]) -> Vec<TaskId> {
settled.iter().map(|pair| pair.0).collect()
}
#[cfg(all(unix, feature = "process"))]
fn armed_group<'ops>(
group: i32,
task: TaskId,
leader_reaped: bool,
permit: Option<Lease>,
owners: Arc<CleanupOwners>,
seams: Seams<'ops>,
) -> ProcessGroup<'ops> {
ProcessGroup {
group,
task,
permit,
owners,
armed: true,
leader_reaped,
leader_exited: false,
group_absent: false,
containment: Containment::default(),
signalled_pids: std::collections::BTreeSet::new(),
signaller: seams.signaller,
observer: seams.observer,
capture: seams.capture,
}
}
#[cfg(all(unix, feature = "process"))]
#[derive(Copy, Clone)]
struct Seams<'ops> {
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
capture: &'ops dyn DescendantCapture,
}
#[cfg(all(unix, feature = "process"))]
fn group_seams<'ops>(
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
) -> Seams<'ops> {
Seams {
signaller,
observer,
capture: &NO_CAPTURE,
}
}
#[cfg(all(unix, feature = "process"))]
fn capture_seams<'ops>(
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
capture: &'ops dyn DescendantCapture,
) -> Seams<'ops> {
Seams {
signaller,
observer,
capture,
}
}
#[cfg(all(unix, feature = "process"))]
static NO_CAPTURE: ScriptedCapture = ScriptedCapture {
tree: Vec::new(),
gone: Vec::new(),
immortal: Vec::new(),
unreadable: false,
truncated: false,
reads: AtomicUsize::new(0),
signals: AtomicUsize::new(0),
};
#[cfg(all(unix, feature = "process"))]
struct FailingGroupObserver;
#[cfg(all(unix, feature = "process"))]
impl GroupObserver for FailingGroupObserver {
fn exists(&self, _group: i32) -> std::io::Result<bool> {
Err(std::io::Error::other("injected observation failure"))
}
}
#[cfg(all(unix, feature = "process"))]
fn test_group<'ops>(
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
leader_reaped: bool,
) -> ProcessGroup<'ops> {
armed_group(
42,
TaskId(0),
leader_reaped,
test_lease(),
Arc::new(CleanupOwners::default()),
group_seams(signaller, observer),
)
}
#[cfg(all(unix, feature = "process"))]
fn capture_group<'ops>(
signaller: &'ops dyn GroupSignaller,
observer: &'ops dyn GroupObserver,
capture: &'ops dyn DescendantCapture,
) -> ProcessGroup<'ops> {
armed_group(
42,
TaskId(0),
false,
test_lease(),
Arc::new(CleanupOwners::default()),
capture_seams(signaller, observer, capture),
)
}
#[cfg(all(unix, feature = "process"))]
const fn sequence(present_before_absent: usize) -> SequenceGroupObserver {
SequenceGroupObserver {
calls: AtomicUsize::new(0),
present_before_absent,
}
}
#[cfg(all(unix, feature = "process"))]
const fn recording() -> RecordingGroupSignaller {
RecordingGroupSignaller {
calls: AtomicUsize::new(0),
}
}
#[cfg(all(unix, feature = "process"))]
fn test_lease() -> Option<Lease> {
let permit = Arc::new(crate::rt::sync::Semaphore::new(1))
.try_acquire_owned()
.ok()?;
Some(Lease::plain(permit))
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_reaped_or_recycled_group_id_is_never_signalled() {
let signaller = recording();
let observer = sequence(0);
let mut group = test_group(&signaller, &observer, true);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupFailed,
"the stale identity is refused rather than signalled"
);
drop(group);
assert_eq!(
signaller.calls.load(Ordering::Relaxed),
0,
"cleanup and Drop must not signal after the owner has reaped the leader"
);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn an_unsignalable_present_group_stays_pending_rather_than_failed() {
let signaller = EpermGroupSignaller;
let observer = sequence(usize::MAX);
let mut group = test_group(&signaller, &observer, false);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"an unsignalable but still-present group is pending, not a failed kill"
);
assert!(
group.armed,
"an unresolved group must stay armed for the drop-time fallback"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn an_unexpected_signal_error_is_a_failed_cleanup() {
let signaller = UnexpectedGroupSignaller;
let observer = sequence(0);
let mut group = test_group(&signaller, &observer, false);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupFailed,
"an unexpected errno remains a refused termination"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_descendant_that_keeps_running_is_named_in_the_receipt() {
let signaller = recording();
let observer = sequence(usize::MAX);
let capture = ScriptedCapture {
tree: vec![70, 71],
immortal: vec![71],
..ScriptedCapture::default()
};
let mut group = capture_group(&signaller, &observer, &capture);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupSurvivors {
survivors: vec![71],
},
"a captured descendant that ignored every signal must be named, not absorbed \
into a pending group"
);
let report = group.containment();
assert_eq!(
report.captured(),
2,
"the report must account for both captured pids"
);
assert_eq!(
report.signalled(),
1 + CONTAINMENT_ROUNDS,
"one signal for the descendant that stopped and one per round for the one that \
did not: a repeat is only ever sent while an observation proves the pid alive"
);
assert_eq!(
report.survivors(),
[71],
"the survivor list is the running pid and nothing else"
);
assert_eq!(
report.mechanism(),
ContainmentMechanism::ProcessTableSnapshot,
"the mechanism that read the table is the one reported"
);
assert_eq!(
report.residual_risk(),
None,
"a whole reading with no survivor carries no residual risk"
);
assert!(
!report.is_complete(),
"a named survivor is by definition an incomplete containment"
);
let pinned_reads = capture.reads();
group.mark_reaped();
assert_eq!(
block_on(group.confirm_absence()),
CleanupReceipt::CleanupSurvivors {
survivors: vec![71],
},
"a survivor is never promoted to a clean cleanup by a later observation"
);
assert_eq!(
capture.reads(),
pinned_reads,
"the post-reap pass observes the captured survivor and reads no tree"
);
assert_eq!(
capture.signals(),
1 + CONTAINMENT_ROUNDS,
"the survivor is signalled once in the round that captured it and once per later \
round, and only because an observation proves it alive; the descendant that \
stopped is signalled exactly once"
);
assert!(
group.armed,
"an unresolved cleanup stays armed for its registry transfer"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_captured_pid_that_had_ended_is_not_reported_as_a_survivor() {
let signaller = recording();
let observer = sequence(0);
let capture = ScriptedCapture {
tree: vec![70],
gone: vec![70],
..ScriptedCapture::default()
};
let mut group = capture_group(&signaller, &observer, &capture);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"the group is still pinned by its unreaped leader, so the pinned phase is pending"
);
let pinned_reads = capture.reads();
group.mark_reaped();
assert_eq!(
block_on(group.confirm_absence()),
CleanupReceipt::CleanupConfirmed,
"an absent group and no running descendant is a complete cleanup"
);
assert_eq!(
capture.reads(),
pinned_reads,
"a clean post-reap pass reads no tree: on Linux a walk from a reaped leader \
was a whole-table `ps` spawn per process"
);
assert!(
group.containment().is_complete(),
"the report must name the mechanism it ran and carry no survivor"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_leader_that_exited_on_its_own_reads_no_table_and_claims_no_tree() {
let signaller = recording();
let observer = sequence(0);
let capture = ScriptedCapture {
tree: vec![70],
..ScriptedCapture::default()
};
let mut group = capture_group(&signaller, &observer, &capture);
group.mark_exited();
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"the exited leader is unreaped, so its group is still pinned and pending"
);
assert_eq!(
capture.reads(),
0,
"a walk from an exited leader can name nothing, so no table is read"
);
assert_eq!(
capture.signals(),
0,
"nothing was named, so no pid is signalled"
);
assert!(
signaller.calls.load(Ordering::Relaxed) >= 1,
"the group is still signalled: a member that stayed in it is still reachable"
);
assert_eq!(
group.containment().residual_risk(),
Some(ResidualRisk::LeaderExited),
"the report names the limit rather than presenting an empty capture as an empty tree"
);
assert_eq!(
group.containment().mechanism(),
ContainmentMechanism::ProcessGroupOnly,
"no table was read"
);
group.mark_reaped();
assert_eq!(
block_on(group.confirm_absence()),
CleanupReceipt::CleanupConfirmed,
"the group itself is observed gone after the reap"
);
assert!(
!group.containment().is_complete(),
"a confirmed group is not a confirmed tree when the tree was never seen"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_captured_pid_is_signalled_once_and_not_once_per_observation() {
let signaller = recording();
let observer = sequence(usize::MAX);
let capture = ScriptedCapture::of(&[70, 71, 72]);
let mut group = capture_group(&signaller, &observer, &capture);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"nothing is running, so the pinned phase is pending on the group alone"
);
assert_eq!(
capture.signals(),
3,
"three captured pids, three signals: a repeat would be a signal to an id whose \
process has since ended"
);
assert_eq!(
capture.reads(),
1,
"the drain stops on the first round once nothing captured is running"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn an_unreadable_process_table_leaves_the_group_as_the_only_mechanism() {
let signaller = recording();
let observer = sequence(usize::MAX);
let capture = ScriptedCapture {
unreadable: true,
..ScriptedCapture::default()
};
let mut group = capture_group(&signaller, &observer, &capture);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"the group was signalled and nothing else could be read"
);
let report = group.containment();
assert_eq!(
report.mechanism(),
ContainmentMechanism::ProcessGroupOnly,
"no table was read, so the mechanism must say the group was all there was"
);
assert_eq!(
report.residual_risk(),
Some(ResidualRisk::TableUnreadable),
"the limit is named rather than left to be inferred from an empty capture"
);
assert_eq!(
report.captured(),
0,
"nothing was captured, so nothing may be reported as captured"
);
assert!(
!report.is_complete(),
"a group-only cleanup never claims the whole tree"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_truncated_capture_is_reported_as_a_prefix_and_not_as_the_tree() {
let signaller = recording();
let observer = sequence(usize::MAX);
let capture = ScriptedCapture {
tree: vec![70, 71],
truncated: true,
..ScriptedCapture::default()
};
let mut group = capture_group(&signaller, &observer, &capture);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"the group is pinned by its unreaped leader"
);
assert_eq!(
group.containment().residual_risk(),
Some(ResidualRisk::CaptureTruncated),
"a reading stopped at its bound must be reported as a prefix"
);
assert!(
!group.containment().is_complete(),
"a prefix cannot claim the whole tree"
);
drop(group);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn bounded_group_termination_keeps_cleanup_owned_until_absence_is_observed() {
let signaller = recording();
let observer = sequence(1);
let mut group = test_group(&signaller, &observer, false);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"the bounded termination receipt must not claim observed absence"
);
let before_reap = signaller.calls.load(Ordering::Relaxed);
assert_eq!(
before_reap, 1,
"the drain stops as soon as nothing captured is still running, so a group with no \
descendants is signalled once and not once per configured round"
);
assert!(
group.armed,
"exhausting the signal budget is not proof that the group disappeared"
);
group.mark_reaped();
assert_eq!(
block_on(group.confirm_absence()),
CleanupReceipt::CleanupConfirmed,
"a signal-zero probe after reaping observes eventual group absence"
);
assert_eq!(observer.calls.load(Ordering::Relaxed), 2);
drop(group);
assert_eq!(
signaller.calls.load(Ordering::Relaxed),
before_reap,
"no signal may follow reaping and release of the numeric id"
);
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn pending_cleanup_owner_keeps_its_permit_until_later_absence_receipt()
-> Result<(), Box<dyn std::error::Error>> {
use super::CleanupOwners;
use crate::rt::sync::Semaphore;
let semaphore = Arc::new(Semaphore::new(1));
let Some(permit) = Arc::clone(&semaphore).try_acquire_owned().ok() else {
return Err("the test semaphore started without its one permit".into());
};
let owners = Arc::new(CleanupOwners::default());
let task = TaskId(9);
let cleanup_observer = sequence(usize::MAX);
let signaller = recording();
let mut group = armed_group(
42,
task,
false,
Some(Lease::plain(permit)),
Arc::clone(&owners),
group_seams(&signaller, &cleanup_observer),
);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupPending,
"exhausting bounded signal attempts remains pending"
);
group.mark_reaped();
assert_eq!(
block_on(group.confirm_absence()),
CleanupReceipt::CleanupPending,
"a bounded post-reap pass cannot claim absence"
);
let later_observer = sequence(1);
drop(group);
assert_eq!(
semaphore.available_permits(),
0,
"dropping a pending task-local guard must transfer, not release, its lease"
);
assert_eq!(owners.pending_count(), 1, "cleanup remains owned");
assert!(
owners.drive(&FailingGroupObserver, |_| false).is_empty(),
"an observation error cannot release cleanup ownership"
);
assert_eq!(semaphore.available_permits(), 0);
assert!(
owners.drive(&later_observer, |_| false).is_empty(),
"a still-present group cannot produce a terminal receipt"
);
assert_eq!(semaphore.available_permits(), 0);
assert!(
owners
.drive(&later_observer, |candidate| candidate == task)
.is_empty(),
"absence before the task report must retain attribution and capacity"
);
assert_eq!(semaphore.available_permits(), 0);
let settled = owners.drive(&later_observer, |_| false);
assert_eq!(
settled_tasks(&settled),
vec![task],
"later absence is attributed to its task"
);
assert_eq!(
semaphore.available_permits(),
1,
"terminal proof releases capacity"
);
assert_eq!(owners.pending_count(), 0);
assert_eq!(
signaller.calls.load(Ordering::Relaxed),
1,
"one signal from the drain, and none after the reap: the id is released then, and \
the drop-time fallback must not signal a number the OS may reissue"
);
Ok(())
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn cleanup_failed_transfers_its_lease_until_absence_is_observed()
-> Result<(), Box<dyn std::error::Error>> {
use super::CleanupOwners;
use crate::rt::sync::Semaphore;
let semaphore = Arc::new(Semaphore::new(1));
let Some(permit) = Arc::clone(&semaphore).try_acquire_owned().ok() else {
return Err("the test semaphore started without its one permit".into());
};
let owners = Arc::new(CleanupOwners::default());
let task = TaskId(10);
let observer = sequence(0);
let signaller = recording();
let mut group = armed_group(
43,
task,
true,
Some(Lease::plain(permit)),
Arc::clone(&owners),
group_seams(&signaller, &observer),
);
assert_eq!(
block_on(group.cleanup()),
CleanupReceipt::CleanupFailed,
"a guard whose leader was reaped cannot signal a possibly reused id"
);
drop(group);
assert_eq!(owners.pending_count(), 1, "failed cleanup remains owned");
assert_eq!(semaphore.available_permits(), 0);
assert_eq!(
settled_tasks(&owners.drive(&observer, |_| false)),
vec![task],
"an absent group settles the obligation it was transferred for"
);
assert_eq!(semaphore.available_permits(), 1);
assert_eq!(signaller.calls.load(Ordering::Relaxed), 0);
Ok(())
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn shutdown_report_keeps_pending_cleanup_and_emits_later_terminal_receipt()
-> Result<(), Box<dyn std::error::Error>> {
use crate::rt::sync::Semaphore;
let semaphore = Arc::new(Semaphore::new(1));
let Some(permit) = Arc::clone(&semaphore).try_acquire_owned().ok() else {
return Err("the test semaphore started without its one permit".into());
};
let task = TaskId(11);
let report = block_on(Supervisor::new(1).shutdown());
report
.cleanup_owners
.register(task, 44, Lease::plain(permit), whole_capture());
assert_eq!(report.pending_cleanup_count(), 1);
assert!(
!report.is_clean(),
"an unresolved native resource is not clean"
);
let Err(mut report) = report.into_outcomes() else {
return Err("into_outcomes consumed a report with pending cleanup".into());
};
assert_eq!(
report.pending_cleanup_count(),
1,
"Err must retain ownership"
);
assert_eq!(semaphore.available_permits(), 0);
assert_eq!(
report.reap_pending_cleanups_with(&FailingGroupObserver),
0,
"an observer error is not a terminal receipt"
);
let observer = sequence(1);
assert_eq!(report.reap_pending_cleanups_with(&observer), 0);
assert_eq!(report.pending_cleanup_count(), 1);
assert_eq!(report.reap_pending_cleanups_with(&observer), 1);
assert_eq!(report.pending_cleanup_count(), 0);
assert_eq!(semaphore.available_permits(), 1);
assert!(report.is_clean());
let outcomes = match report.into_outcomes() {
Ok(outcomes) => outcomes,
Err(_) => return Err("settled cleanup report remained unavailable".into()),
};
assert_eq!(
outcomes,
vec![TaskOutcome::CleanupSettled {
task,
cleanup: CleanupReceipt::CleanupConfirmed,
containment: whole_capture(),
}]
);
Ok(())
}
async fn settle(supervisor: &mut Supervisor) -> bool {
let mut spins: u32 = 0;
loop {
supervisor.reap();
if supervisor.stats().in_flight() == 0 {
return true;
}
if spins >= 100_000 {
return false;
}
spins = spins.saturating_add(1);
yield_now().await;
}
}
fn drain(supervisor: &mut Supervisor) -> Vec<TaskOutcome> {
let mut outcomes = Vec::new();
while let Some(outcome) = supervisor.next_report() {
outcomes.push(outcome);
}
outcomes
}
fn tasks_where(outcomes: &[TaskOutcome], keep: impl Fn(&TaskOutcome) -> bool) -> Vec<TaskId> {
outcomes
.iter()
.filter(|outcome| keep(outcome))
.map(TaskOutcome::task)
.collect()
}
async fn yield_until(mut predicate: impl FnMut() -> bool, limit: u32) -> bool {
let mut yields: u32 = 0;
while !predicate() {
if yields >= limit {
return false;
}
yields = yields.saturating_add(1);
yield_now().await;
}
true
}
async fn spawn_heartbeat(supervisor: &mut Supervisor, beats: &Arc<AtomicU64>) {
let beats = Arc::clone(beats);
supervisor
.spawn(move |token| async move {
while !token.is_cancelled() {
beats.fetch_add(1, Ordering::SeqCst);
yield_now().await;
}
})
.await;
}
#[test]
fn a_large_finite_budget_still_yields_to_a_sibling_task() {
const BUDGET: u64 = 1_000_000;
block_on(async {
let beats = Arc::new(AtomicU64::new(0));
let ran = Arc::new(AtomicU64::new(0));
let first_seen = Arc::new(AtomicU64::new(u64::MAX));
let later_seen = Arc::new(AtomicU64::new(0));
let mut supervisor = Supervisor::new(4);
spawn_heartbeat(&mut supervisor, &beats).await;
let body_ran = Arc::clone(&ran);
let body_first = Arc::clone(&first_seen);
let body_later = Arc::clone(&later_seen);
let body_beats = Arc::clone(&beats);
supervisor
.spawn_repeating(budget_of(BUDGET), move |_tick| {
let body_ran = Arc::clone(&body_ran);
let body_first = Arc::clone(&body_first);
let body_later = Arc::clone(&body_later);
let body_beats = Arc::clone(&body_beats);
async move {
body_ran.fetch_add(1, Ordering::SeqCst);
let now = body_beats.load(Ordering::SeqCst);
if now < body_first.load(Ordering::SeqCst) {
body_first.store(now, Ordering::SeqCst);
}
if now > body_later.load(Ordering::SeqCst) {
body_later.store(now, Ordering::SeqCst);
}
}
})
.await;
assert!(
yield_until(|| ran.load(Ordering::SeqCst) > 0, 1_000).await,
"the repeating worker never took a single iteration"
);
assert!(
yield_until(
|| later_seen.load(Ordering::SeqCst) > first_seen.load(Ordering::SeqCst),
1_000
)
.await,
"the repeating body ran {} of its {BUDGET} iterations seeing the heartbeat count \
frozen at {}, so one poll of it held the executor for the whole run",
ran.load(Ordering::SeqCst),
first_seen.load(Ordering::SeqCst)
);
supervisor.shutdown().await;
});
}
#[test]
fn a_cancel_from_a_sibling_task_stops_an_ongoing_ready_body() {
const BRAKE: u64 = 1_000_000;
block_on(async {
let root = CancellationToken::new();
let ran = Arc::new(AtomicU64::new(0));
let mut supervisor = Supervisor::new(2);
let canceller_root = root.clone();
let canceller_ran = Arc::clone(&ran);
supervisor
.spawn(move |_task_token| async move {
while canceller_ran.load(Ordering::SeqCst) == 0 {
yield_now().await;
}
canceller_root.cancel();
})
.await;
let loop_ran = Arc::clone(&ran);
let loop_root = root.clone();
let brake = root.clone();
let outcome = repeat(&loop_root, Budget::Ongoing, move |_tick| {
let brake = brake.clone();
let loop_ran = Arc::clone(&loop_ran);
async move {
let count = loop_ran.fetch_add(1, Ordering::SeqCst).saturating_add(1);
if count >= BRAKE {
brake.cancel();
}
}
})
.await;
assert!(
outcome.was_cancelled(),
"the loop ended as {outcome:?} rather than as cancelled"
);
let iterations = outcome.iterations();
assert!(
iterations < BRAKE,
"the loop ran to its own {BRAKE}-iteration brake ({iterations} iterations) \
instead of observing the cancel from the sibling task"
);
assert!(
settle(&mut supervisor).await,
"the canceller task never finished"
);
let settled = ran.load(Ordering::SeqCst);
for _ in 0..64 {
yield_now().await;
}
assert_eq!(
ran.load(Ordering::SeqCst),
settled,
"the body was invoked again after the loop observed the cancel"
);
});
}
#[test]
fn shutdown_lands_on_an_ongoing_ready_body() {
const BRAKE: u64 = 1_000_000;
block_on(async {
let ran = Arc::new(AtomicU64::new(0));
let mut supervisor = Supervisor::new(1);
let body_ran = Arc::clone(&ran);
supervisor
.spawn(move |token| async move {
let brake = token.clone();
let _outcome = repeat(&token, Budget::Ongoing, move |_tick| {
let brake = brake.clone();
let body_ran = Arc::clone(&body_ran);
async move {
let count = body_ran.fetch_add(1, Ordering::SeqCst).saturating_add(1);
if count >= BRAKE {
brake.cancel();
}
}
})
.await;
})
.await;
assert!(
yield_until(|| ran.load(Ordering::SeqCst) > 0, 1_000).await,
"the repeating worker never took a single iteration"
);
supervisor.shutdown().await;
let iterations = ran.load(Ordering::SeqCst);
assert!(
iterations < BRAKE,
"the loop ran to its own {BRAKE}-iteration brake ({iterations} iterations) \
instead of being stopped by the shutdown"
);
});
}
#[test]
fn an_already_cancelled_token_never_invokes_the_body() {
let token = CancellationToken::new();
token.cancel();
let ran = AtomicU64::new(0);
let body_ran = &ran;
let outcome = block_on(repeat(&token, Budget::Ongoing, move |_tick| async move {
body_ran.fetch_add(1, Ordering::SeqCst);
}));
assert_eq!(
outcome,
Outcome::Cancelled { iterations: 0 },
"a token cancelled before the loop starts must stop it with no iterations"
);
assert_eq!(
ran.load(Ordering::SeqCst),
0,
"the body must not be invoked at all past a cancellation boundary"
);
}
#[test]
fn a_budget_of_iterations_stops_at_the_limit() {
let token = CancellationToken::new();
let ticks = Arc::new(AtomicU64::new(0));
let counter = Arc::clone(&ticks);
let outcome = block_on(repeat(&token, budget_of(3), move |_tick| {
let counter = Arc::clone(&counter);
async move {
counter.fetch_add(1, Ordering::SeqCst);
}
}));
assert_eq!(
outcome,
Outcome::Exhausted { iterations: 3 },
"a three-iteration budget must run the body exactly three times"
);
assert_eq!(
ticks.load(Ordering::SeqCst),
3,
"the body must have been entered once per iteration"
);
}
fn repeat_cancelling_itself<Rest, Fut>(rest: Rest) -> Outcome
where
Rest: Fn() -> Fut,
Fut: Future<Output = ()>,
{
let token = CancellationToken::new();
let trigger = token.clone();
block_on(repeat(&token, Budget::Ongoing, move |_tick| {
trigger.cancel();
rest()
}))
}
#[test]
fn an_ongoing_budget_stops_when_the_token_is_cancelled() {
let outcome = repeat_cancelling_itself(|| async {});
assert_eq!(
outcome,
Outcome::Cancelled { iterations: 1 },
"a cancel during the first iteration must end the loop after it"
);
}
#[test]
fn cancellation_interrupts_a_body_that_is_still_awaiting() {
let outcome = repeat_cancelling_itself(pending::<()>);
assert!(
outcome.was_cancelled(),
"an in-flight body must be dropped at the cancel, got {outcome:?}"
);
}
#[test]
fn a_budget_of_duration_expires() {
let token = CancellationToken::new();
let outcome = block_on(repeat(
&token,
Budget::For(Duration::ZERO),
move |_tick| async {},
));
assert_eq!(
outcome,
Outcome::Exhausted { iterations: 0 },
"a zero duration must expire before the first iteration"
);
}
#[test]
fn a_supervisor_reports_what_it_has_run() {
block_on(async {
let mut supervisor = Supervisor::new(2);
supervisor.spawn(|_token| async {}).await;
supervisor.spawn(|_token| async {}).await;
assert!(settle(&mut supervisor).await, "both tasks must finish");
let stats = supervisor.stats();
assert_eq!(stats.spawned, 2, "both tasks must be counted as spawned");
assert_eq!(
stats.completed, 2,
"both tasks must be counted as completed"
);
assert_eq!(
stats.succeeded, 2,
"both tasks returned, so both must be counted as succeeded"
);
assert_eq!(stats.in_flight(), 0, "nothing must remain in flight");
assert_eq!(stats.refused, 0, "nothing must have been refused");
});
}
#[test]
fn a_cancelled_supervisor_refuses_admission_without_constructing_the_body() {
block_on(async {
let mut supervisor = Supervisor::new(2);
supervisor.cancel();
let built = Arc::new(AtomicU64::new(0));
supervisor
.spawn({
let built = Arc::clone(&built);
move |_token| {
built.fetch_add(1, Ordering::Relaxed);
async {}
}
})
.await;
let stats = supervisor.stats();
assert_eq!(
built.load(Ordering::Relaxed),
0,
"a cancelled supervisor must never construct the body it refuses"
);
assert_eq!(stats.spawned, 0, "a refused spawn must not be placed");
assert_eq!(
stats.refused, 1,
"the waiting-spawn refusal must be counted"
);
let built_now = Arc::new(AtomicU64::new(0));
let refused = supervisor.try_spawn({
let built_now = Arc::clone(&built_now);
move |_token| {
built_now.fetch_add(1, Ordering::Relaxed);
async {}
}
});
assert_eq!(
refused,
Err(TrySpawnRefusal::Cancelled),
"the immediate path must name cancellation, not capacity"
);
assert_eq!(
built_now.load(Ordering::Relaxed),
0,
"the immediate path must also refuse before the constructor"
);
assert_eq!(supervisor.stats().refused, 2, "both refusals counted");
});
}
async fn occupy(supervisor: &mut Supervisor) -> (Arc<AtomicBool>, Arc<AtomicBool>) {
let gate = Arc::new(AtomicBool::new(false));
let done = Arc::new(AtomicBool::new(false));
supervisor
.spawn({
let gate = Arc::clone(&gate);
let done = Arc::clone(&done);
move |_token| async move {
while !gate.load(Ordering::Relaxed) {
yield_now().await;
}
done.store(true, Ordering::Relaxed);
}
})
.await;
(gate, done)
}
#[test]
fn cancellation_closes_admission_even_after_a_slot_frees() {
block_on(async {
let mut supervisor = Supervisor::new(1);
let (gate, done) = occupy(&mut supervisor).await;
supervisor.cancel();
gate.store(true, Ordering::Relaxed);
for _ in 0..256 {
if done.load(Ordering::Relaxed) {
break;
}
yield_now().await;
}
assert!(
done.load(Ordering::Relaxed),
"the occupier must have ended before the admission attempt"
);
let built = Arc::new(AtomicU64::new(0));
let refused = supervisor.try_spawn({
let built = Arc::clone(&built);
move |_token| {
built.fetch_add(1, Ordering::Relaxed);
async {}
}
});
assert_eq!(
refused,
Err(TrySpawnRefusal::Cancelled),
"a freed slot is not an admission offer after cancellation"
);
assert_eq!(
built.load(Ordering::Relaxed),
0,
"no body may be constructed after cancellation, slot or not"
);
assert_eq!(supervisor.stats().spawned, 1, "only the occupier ran");
});
}
#[test]
fn a_cancelled_supervisor_with_a_full_pool_refuses_without_constructing() {
block_on(async {
let mut supervisor = Supervisor::new(1);
let (gate, done) = occupy(&mut supervisor).await;
supervisor.cancel();
let built = Arc::new(AtomicU64::new(0));
let refused_before = supervisor.stats().refused;
supervisor
.spawn({
let built = Arc::clone(&built);
move |_token| {
built.fetch_add(1, Ordering::Relaxed);
async {}
}
})
.await;
assert_eq!(
built.load(Ordering::Relaxed),
0,
"no body may be constructed for a cancelled supervisor, full pool or not"
);
assert_eq!(
supervisor.stats().refused,
refused_before + 1,
"the waiting-spawn refusal must be counted"
);
assert_eq!(supervisor.stats().spawned, 1, "only the occupier ran");
gate.store(true, Ordering::Relaxed);
for _ in 0..256 {
if done.load(Ordering::Relaxed) {
break;
}
yield_now().await;
}
});
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_cancelled_supervisor_starts_no_process() {
use super::ProcessSpec;
block_on(async {
static MARK_SEQ: AtomicU64 = AtomicU64::new(0);
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or_else(
|before| before.duration().as_nanos(),
|since| since.as_nanos(),
);
let seq = MARK_SEQ.fetch_add(1, Ordering::Relaxed);
let marker = std::env::temp_dir().join(format!("lgwks-bot-cancel-{nanos}-{seq}.mark"));
let mut supervisor = Supervisor::new(2);
supervisor.cancel();
let spawned_before = supervisor.stats().spawned;
let refused_before = supervisor.stats().refused;
let mut spec = ProcessSpec::new("sh");
spec.arg("-c").arg(format!("touch {}", marker.display()));
let refusal = supervisor.spawn_process(&spec).await;
assert!(
refusal.is_err(),
"a cancelled supervisor must not start the process it was asked for"
);
assert!(
refusal
.as_ref()
.err()
.and_then(std::io::Error::get_ref)
.and_then(|cause| cause.downcast_ref::<SupervisorCancelled>())
.is_some(),
"the refusal must name cancellation, got {refusal:?}"
);
assert_eq!(
supervisor.stats().spawned,
spawned_before,
"no process may be placed after cancellation"
);
assert_eq!(
supervisor.stats().refused,
refused_before + 1,
"the process refusal must be counted"
);
assert!(
!marker.exists(),
"no child ran: the marker {} must be absent",
marker.display()
);
let removed = std::fs::remove_file(&marker);
assert!(
removed.is_ok() || !marker.exists(),
"the marker path must be gone: {}",
marker.display()
);
});
}
#[test]
fn try_spawn_refuses_at_the_bound_instead_of_growing() {
block_on(async {
let mut supervisor = Supervisor::new(1);
let held = supervisor.try_spawn(|_token| async {
pending::<()>().await;
});
assert!(held.is_ok(), "the first spawn must take the only slot");
let refused = supervisor.try_spawn(|_token| async {});
assert_eq!(
refused,
Err(TrySpawnRefusal::AtCapacity),
"a second spawn past the bound must be refused, not queued"
);
let stats = supervisor.stats();
assert_eq!(stats.refused, 1, "the refusal must be counted");
assert_eq!(stats.spawned, 1, "a refused body must never be spawned");
supervisor.cancel();
});
}
#[test]
fn the_retained_set_stays_at_the_bound_across_many_spawns() {
block_on(async {
let mut supervisor = Supervisor::new(1);
for _ in 0..64 {
supervisor.spawn(|_token| async {}).await;
}
let stats = supervisor.stats();
assert_eq!(stats.spawned, 64, "every spawn must be counted");
assert!(
stats.in_flight() <= 2,
"the retained set must stay at the bound, not grow with total \
spawns; got {} in flight of {} spawned",
stats.in_flight(),
stats.spawned
);
assert!(
settle(&mut supervisor).await,
"the remaining tasks must finish"
);
assert_eq!(
supervisor.stats().in_flight(),
0,
"settling must leave nothing in flight"
);
});
}
#[test]
fn dropping_a_supervisor_cancels_the_tasks_it_owns() {
block_on(async {
let probe = {
let mut supervisor = Supervisor::new(1);
supervisor
.spawn(|_token| async {
pending::<()>().await;
})
.await;
supervisor.child_token()
};
assert!(
probe.is_cancelled(),
"a task's token must be cancelled when the supervisor is dropped"
);
});
}
#[test]
fn shutdown_drains_every_task() {
block_on(async {
let mut supervisor = Supervisor::new(4);
supervisor
.spawn(|token| async move {
let _outcome = repeat(&token, Budget::Ongoing, |_tick| async {}).await;
})
.await;
supervisor
.spawn(|token| async move {
let _outcome = repeat(&token, Budget::Ongoing, |_tick| async {}).await;
})
.await;
let report = supervisor.shutdown().await;
assert_eq!(
report.stats().spawned,
2,
"shutdown must report the counters it drained against"
);
});
}
#[test]
fn a_repeating_task_stops_at_its_budget_under_supervision() {
block_on(async {
let ticks = Arc::new(AtomicU64::new(0));
let counter = Arc::clone(&ticks);
let mut supervisor = Supervisor::new(1);
supervisor
.spawn_repeating(budget_of(5), move |_tick| {
let counter = Arc::clone(&counter);
async move {
counter.fetch_add(1, Ordering::SeqCst);
}
})
.await;
assert!(
settle(&mut supervisor).await,
"a repeating task must finish once its budget is spent"
);
assert_eq!(
ticks.load(Ordering::SeqCst),
5,
"a supervised repeating task must stop at its iteration budget"
);
supervisor.shutdown().await;
});
}
#[test]
fn a_panicking_task_is_not_reported_as_a_success() {
block_on(async {
let mut supervisor = Supervisor::new(2);
supervisor.spawn(|_token| async {}).await;
supervisor
.spawn(|_token| async {
explode("worker failed before producing its receipt");
})
.await;
assert!(
settle(&mut supervisor).await,
"a panicking task still ends, so nothing is left in flight"
);
let stats = supervisor.stats();
assert_eq!(stats.spawned, 2, "both tasks were started");
assert_eq!(stats.completed, 2, "both tasks ended");
assert_eq!(
stats.succeeded, 1,
"only the task that returned is a success: {stats:?}"
);
assert_eq!(stats.panicked, 1, "one task panicked: {stats:?}");
let outcomes = drain(&mut supervisor);
assert_eq!(outcomes.len(), 2, "every task reports exactly once");
assert_eq!(
tasks_where(&outcomes, TaskOutcome::is_panic),
vec![TaskId(1)],
"the report names the task that panicked, in spawn order: {outcomes:?}"
);
assert_eq!(
tasks_where(&outcomes, TaskOutcome::is_success),
vec![TaskId(0)],
"the task that returned is reported as the one that returned: {outcomes:?}"
);
let messages: Vec<&str> = outcomes
.iter()
.filter_map(TaskOutcome::panic_message)
.collect();
assert_eq!(
messages,
vec!["worker failed before producing its receipt"],
"the panic payload is preserved rather than replaced by a marker"
);
});
}
#[test]
fn a_panic_after_an_observable_effect_keeps_both_facts() {
block_on(async {
let effects = Arc::new(AtomicU64::new(0));
let counter = Arc::clone(&effects);
let mut supervisor = Supervisor::new(1);
supervisor
.spawn(move |_token| {
let counter = Arc::clone(&counter);
async move {
counter.fetch_add(1, Ordering::SeqCst);
explode("died after the effect");
}
})
.await;
assert!(settle(&mut supervisor).await, "the task must end");
assert_eq!(
effects.load(Ordering::SeqCst),
1,
"the effect happened, so the test's premise holds"
);
assert_eq!(
supervisor.stats().succeeded,
0,
"an effect that happened is not a task that finished"
);
let outcomes = drain(&mut supervisor);
assert_eq!(
outcomes,
vec![TaskOutcome::Panicked {
task: TaskId(0),
message: String::from("died after the effect"),
}],
"the report carries the panic, and does not claim success"
);
});
}
#[test]
fn cooperative_cancellation_and_abort_are_told_apart() {
block_on(async {
let mut supervisor = Supervisor::new(2);
supervisor
.spawn(|token| async move {
token.cancelled().await;
})
.await;
supervisor
.spawn(|_token| async {
loop {
yield_now().await;
}
})
.await;
let report = supervisor.shutdown().await;
let stats = report.stats();
assert_eq!(stats.spawned, 2, "both tasks were started");
assert_eq!(stats.completed, 2, "both tasks ended");
assert_eq!(
stats.cancelled, 1,
"the cooperative body returned under cancellation: {stats:?}"
);
assert_eq!(
stats.aborted, 1,
"the body that ignored its token had to be dropped: {stats:?}"
);
assert_eq!(
stats.succeeded, 0,
"neither task ran to the end of its own work: {stats:?}"
);
assert!(
!report.is_clean(),
"a shutdown that cancelled or aborted anything is not a clean finish"
);
assert_eq!(
tasks_where(report.outcomes(), TaskOutcome::is_success),
Vec::<TaskId>::new(),
"no failed task may appear as successful: {:?}",
report.outcomes()
);
assert_eq!(
report
.outcomes()
.iter()
.map(TaskOutcome::task)
.collect::<Vec<TaskId>>(),
vec![TaskId(0), TaskId(1)],
"both terminal outcomes are reported, each attributed to its task"
);
});
}
#[test]
fn shutdown_reports_a_clean_finish_when_every_task_returned() {
block_on(async {
let mut supervisor = Supervisor::new(2);
supervisor.spawn(|_token| async {}).await;
supervisor.spawn(|_token| async {}).await;
assert!(settle(&mut supervisor).await, "both tasks must finish");
let report = supervisor.shutdown().await;
assert_eq!(
report.stats().succeeded,
2,
"both tasks returned before the cancel"
);
assert!(
report.panicked().next().is_none(),
"a clean finish has no panics to report"
);
});
}
#[test]
fn an_undrained_report_buffer_is_capped_and_counted() {
block_on(async {
let mut supervisor = Supervisor::new(1);
for _ in 0..32 {
supervisor.spawn(|_token| async {}).await;
}
assert!(settle(&mut supervisor).await, "every task must finish");
let stats = supervisor.stats();
assert_eq!(stats.spawned, 32, "every spawn is still counted");
assert_eq!(stats.succeeded, 32, "every success is still counted");
assert_eq!(
stats
.succeeded
.saturating_add(stats.cancelled)
.saturating_add(stats.aborted)
.saturating_add(stats.panicked),
stats.completed,
"the four outcomes must account for every task that ended"
);
let retained = drain(&mut supervisor).len();
assert_eq!(
retained, 1,
"the buffer retains at most the in-flight bound"
);
assert_eq!(
stats.reports_dropped, 31,
"and every report it did not retain is counted: {stats:?}"
);
});
}
#[test]
fn a_capped_report_buffer_still_attributes_the_report_it_keeps() {
block_on(async {
let mut supervisor = Supervisor::new(1);
assert!(
supervisor
.try_spawn(|_token| async {
pending::<()>().await;
})
.is_ok(),
"the first spawn must take the only slot"
);
assert_eq!(
supervisor.try_spawn(|_token| async {}),
Err(TrySpawnRefusal::AtCapacity),
"the premise: the bound is genuinely reached"
);
let report = supervisor.shutdown().await;
assert_eq!(
report.outcomes().len(),
1,
"one task ended, so one report is retained"
);
assert_eq!(
report.outcomes().first().map(TaskOutcome::task),
Some(TaskId(0)),
"the retained report still names its task"
);
assert_eq!(
report.stats().aborted,
1,
"the body never observed its token, so it was dropped"
);
});
}
#[test]
fn shutdown_retains_every_report_even_past_the_cap() {
block_on(async {
let mut supervisor = Supervisor::new(1);
supervisor.spawn(|_token| async {}).await;
assert!(settle(&mut supervisor).await, "the first task must finish");
assert_eq!(
supervisor.stats().reports_dropped,
0,
"the premise: the one report so far fit in the buffer"
);
supervisor
.spawn(|_token| async {
pending::<()>().await;
})
.await;
let report = supervisor.shutdown().await;
assert_eq!(
report.outcomes().len(),
2,
"shutdown must hand over both outcomes, not just the ones that \
fit: {:?}",
report.outcomes()
);
assert_eq!(
report.stats().reports_dropped,
0,
"the drain path is not capped"
);
assert_eq!(
tasks_where(report.outcomes(), TaskOutcome::is_success),
vec![TaskId(0)],
"the task that returned is reported as the one that returned"
);
assert_eq!(
tasks_where(report.outcomes(), |outcome| !outcome.is_success()),
vec![TaskId(1)],
"and the one that did not is attributed to the other task"
);
});
}
#[cfg(all(unix, feature = "process"))]
fn reduce(value: u64, bound: NonZeroU64) -> Result<u64, std::num::TryFromIntError> {
let wide = u128::from(value).wrapping_mul(u128::from(bound.get()));
u64::try_from(wide >> 64)
}
#[cfg(all(unix, feature = "process"))]
struct Seed(u64);
#[cfg(all(unix, feature = "process"))]
impl Seed {
fn draw(&self, field: u64, bound: NonZeroU64) -> Result<u64, std::num::TryFromIntError> {
let mut mixed = self
.0
.wrapping_add(field.wrapping_mul(0x9E37_79B9_7F4A_7C15));
mixed = (mixed ^ (mixed >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
mixed = (mixed ^ (mixed >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
reduce(mixed ^ (mixed >> 31), bound)
}
fn below(&self, field: u64, max: u64) -> Result<u64, std::num::TryFromIntError> {
self.draw(field, NonZeroU64::MIN.saturating_add(max))
}
}
#[cfg(all(unix, feature = "process"))]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum WorldGroup {
Delivered,
Absent,
Present,
Refused,
}
#[cfg(all(unix, feature = "process"))]
struct World {
tree: Vec<i32>,
gone: Vec<i32>,
immortal: Vec<i32>,
unreadable: bool,
truncated: bool,
exited: bool,
group: WorldGroup,
present_probes: usize,
expected: CleanupReceipt,
containment: Containment,
}
#[cfg(all(unix, feature = "process"))]
impl World {
fn draw(seed: u64) -> Result<Self, std::num::TryFromIntError> {
let draws = Seed(seed);
let length = draws.below(0, 3)?;
let mut tree = Vec::with_capacity(usize::try_from(length)?);
let mut gone = Vec::new();
let mut immortal = Vec::new();
for index in 0..length {
let pid_field = index.saturating_add(1);
let pid = i32::try_from(draws.below(pid_field, 60)?)?;
if tree.contains(&pid) {
continue;
}
tree.push(pid);
let ends = draws.below(index.saturating_add(20), 1)? == 0;
let ignores = draws.below(index.saturating_add(40), 1)? == 0;
if ends {
gone.push(pid);
} else if ignores {
immortal.push(pid);
}
}
tree.sort_unstable();
gone.sort_unstable();
immortal.sort_unstable();
let unreadable = draws.below(3, 3)? == 0;
let truncated = !unreadable && draws.below(4, 3)? == 0;
let group = match draws.below(5, 3)? {
0 => WorldGroup::Absent,
1 => WorldGroup::Present,
2 => WorldGroup::Refused,
_ => WorldGroup::Delivered,
};
let present_probes = usize::try_from(draws.below(6, 2)?)?;
let exited = draws.below(7, 3)? == 0;
let read = !unreadable && !exited;
let named: Vec<i32> = if read { tree.clone() } else { Vec::new() };
let refused = group == WorldGroup::Refused;
let per_pid = read && !refused;
let survivors: Vec<i32> = if per_pid {
named
.iter()
.copied()
.filter(|pid| immortal.contains(pid))
.collect()
} else {
Vec::new()
};
let rounds = if survivors.is_empty() {
1
} else {
CONTAINMENT_ROUNDS
};
let delivered = named.iter().filter(|pid| !gone.contains(pid)).count();
let signalled = if per_pid {
delivered.saturating_add(survivors.len().saturating_mul(rounds.saturating_sub(1)))
} else {
0
};
let residual = if exited {
Some(ResidualRisk::LeaderExited)
} else if unreadable {
Some(ResidualRisk::TableUnreadable)
} else if truncated {
Some(ResidualRisk::CaptureTruncated)
} else {
None
};
let expected = match group {
WorldGroup::Refused => CleanupReceipt::CleanupFailed,
_ if !survivors.is_empty() => CleanupReceipt::CleanupSurvivors {
survivors: survivors.clone(),
},
WorldGroup::Absent => CleanupReceipt::CleanupConfirmed,
WorldGroup::Delivered | WorldGroup::Present => CleanupReceipt::CleanupPending,
};
Ok(Self {
tree,
gone,
immortal,
unreadable,
truncated,
exited,
group,
present_probes,
expected,
containment: Containment {
mechanism: if read {
ContainmentMechanism::ProcessTableSnapshot
} else {
ContainmentMechanism::ProcessGroupOnly
},
captured: named.len(),
signalled,
survivors,
residual,
},
})
}
fn settled(&self) -> CleanupReceipt {
if !self.containment.survivors.is_empty() {
return CleanupReceipt::CleanupSurvivors {
survivors: self.containment.survivors.clone(),
};
}
if self.group == WorldGroup::Absent || self.present_probes < CONTAINMENT_ROUNDS {
CleanupReceipt::CleanupConfirmed
} else {
CleanupReceipt::CleanupPending
}
}
fn capture(&self) -> ScriptedCapture {
ScriptedCapture {
tree: self.tree.clone(),
gone: self.gone.clone(),
immortal: self.immortal.clone(),
unreadable: self.unreadable,
truncated: self.truncated,
reads: AtomicUsize::new(0),
signals: AtomicUsize::new(0),
}
}
fn signaller(&self) -> WorldSignaller {
WorldSignaller(self.group)
}
}
#[cfg(all(unix, feature = "process"))]
struct WorldSignaller(WorldGroup);
#[cfg(all(unix, feature = "process"))]
impl GroupSignaller for WorldSignaller {
fn signal(&self, _group: i32) -> std::io::Result<()> {
match self.0 {
WorldGroup::Absent => Err(std::io::Error::from_raw_os_error(ESRCH)),
WorldGroup::Present => Err(std::io::Error::from_raw_os_error(EPERM)),
WorldGroup::Refused => Err(std::io::Error::new(
std::io::ErrorKind::InvalidInput,
"injected refusal",
)),
WorldGroup::Delivered => Ok(()),
}
}
}
#[cfg(all(unix, feature = "process"))]
fn receipt_arm(receipt: &CleanupReceipt) -> u64 {
match *receipt {
CleanupReceipt::CleanupConfirmed => 1,
CleanupReceipt::CleanupPending => 2,
CleanupReceipt::CleanupFailed => 3,
CleanupReceipt::CleanupSurvivors { .. } => 4,
}
}
#[cfg(all(unix, feature = "process"))]
fn receipt_code(receipt: &CleanupReceipt) -> Result<u64, std::num::TryFromIntError> {
Ok(match *receipt {
CleanupReceipt::CleanupSurvivors { ref survivors } => {
u64::try_from(survivors.len())?.saturating_add(4)
}
_ => receipt_arm(receipt),
})
}
#[cfg(all(unix, feature = "process"))]
fn sim_drain(seed: u64) -> Result<u64, Box<dyn std::error::Error>> {
let world = World::draw(seed)?;
let at = format!("seed {seed:#018x}");
let capture = world.capture();
let signaller = world.signaller();
let observer = sequence(world.present_probes);
let mut group = capture_group(&signaller, &observer, &capture);
if world.exited {
group.mark_exited();
}
let pinned = block_on(group.cleanup());
assert_eq!(
pinned, world.expected,
"{at}: the pinned phase must report the model's receipt"
);
assert_eq!(
*group.containment(),
world.containment,
"{at}: the containment report must carry the model's counts and limits"
);
assert_eq!(
capture.signals(),
world.containment.signalled(),
"{at}: the drain must deliver exactly the signals the model counts"
);
if world.exited {
assert_eq!(
capture.reads(),
0,
"{at}: a leader that had exited has no tree to walk, so no table is read"
);
assert!(
!group.containment().is_complete(),
"{at}: a cleanup that never saw the tree must not claim it"
);
}
group.mark_reaped();
let settled = block_on(group.confirm_absence());
let expected = world.settled();
assert_eq!(
settled, expected,
"{at}: only an observed absence settles the group, and a survivor is never \
promoted to a clean cleanup"
);
let captured = u64::try_from(group.containment().captured())?;
let signalled = u64::try_from(group.containment().signalled())?;
let pinned_code = receipt_code(&pinned)?;
let settled_code = receipt_code(&settled)?;
Ok(pinned_code
.wrapping_add(captured)
.wrapping_add(signalled)
.wrapping_add(settled_code))
}
#[cfg(all(unix, feature = "process"))]
const DRAIN_SEEDS: [u64; 12] = [
0x5EED_2630_0000_0001,
0x5EED_2630_0000_0002,
0x5EED_2630_0000_0003,
0x5EED_2630_0000_0004,
0x5EED_2630_0000_0005,
0x5EED_2630_0000_0006,
0x5EED_2630_0000_000E,
0x5EED_2630_0000_0019,
0x5EED_2630_0000_0025,
0x5EED_2630_FFFF_FFFF,
0xDEAD_BEEF_0263_0001,
0xC0FF_EE00_2630_0001,
];
#[cfg(all(unix, feature = "process"))]
#[test]
fn sim_a_seeded_drain_reaches_the_models_receipt() -> Result<(), Box<dyn std::error::Error>> {
for seed in DRAIN_SEEDS {
sim_drain(seed)?;
}
Ok(())
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn sim_the_same_seed_replays_the_same_drain_trace() -> Result<(), Box<dyn std::error::Error>> {
for seed in DRAIN_SEEDS {
assert_eq!(
sim_drain(seed)?,
sim_drain(seed)?,
"seed {seed:#018x}: the same seed must drive the same world and the same report"
);
}
Ok(())
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn sim_every_receipt_arm_is_reachable_in_the_family() -> Result<(), Box<dyn std::error::Error>>
{
let arms = DRAIN_SEEDS
.iter()
.map(|seed| World::draw(*seed).map(|world| receipt_arm(&world.expected)))
.collect::<Result<Vec<u64>, _>>()?;
for arm in [1_u64, 2, 3, 4] {
assert!(
arms.contains(&arm),
"the family must reach receipt arm {arm}; it drew {arms:?} over {} seeds",
DRAIN_SEEDS.len()
);
}
Ok(())
}
#[cfg(all(unix, feature = "process"))]
const EXITED_SWEEP_SEEDS: u64 = 4_096;
#[cfg(all(unix, feature = "process"))]
#[test]
fn sim_an_exited_leader_reads_no_table_and_claims_no_tree()
-> Result<(), Box<dyn std::error::Error>> {
let mut exited_arms = std::collections::BTreeSet::new();
let mut alive = 0_u64;
for offset in 0..EXITED_SWEEP_SEEDS {
let seed = 0x5EED_0347_0000_0000_u64.wrapping_add(offset);
sim_drain(seed)?;
let world = World::draw(seed)?;
if world.exited {
exited_arms.insert(receipt_arm(&world.expected));
} else {
alive = alive.saturating_add(1);
}
}
assert_eq!(
exited_arms,
std::collections::BTreeSet::from([1_u64, 2, 3]),
"an exited leader must be driven to the confirmed, pending and failed receipts; \
it can never have a named survivor, because nothing was named"
);
assert!(
alive > 0,
"the sweep must still hold leaders that were alive at the cleanup"
);
Ok(())
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn sim_distinct_seeds_drive_distinct_drains() -> Result<(), Box<dyn std::error::Error>> {
let [first, second, ..] = DRAIN_SEEDS;
assert_ne!(
sim_drain(first)?,
sim_drain(second)?,
"seeds {first:#018x} and {second:#018x} must draw different worlds"
);
Ok(())
}
#[test]
fn an_over_long_panic_message_is_truncated_not_dropped() {
let long = "x".repeat(MAX_PANIC_MESSAGE_CHARS.saturating_add(1));
let kept = truncate_panic_message(&long);
assert_eq!(
kept.chars().count(),
MAX_PANIC_MESSAGE_CHARS.saturating_add(1),
"a truncated message keeps the cap plus the single-character marker"
);
assert!(
kept.ends_with('…'),
"a truncated message says so, so it is never read as complete"
);
assert_eq!(
truncate_panic_message("short"),
"short",
"a message within the cap is carried unchanged"
);
}
#[cfg(feature = "script")]
#[test]
fn a_cancelled_admission_hands_its_granted_permit_back_to_the_round()
-> Result<(), Box<dyn std::error::Error>> {
use crate::rt::sync::Notify;
use crate::rt::tenancy::TenancyPolicy;
use crate::script::Tenant;
block_on(async {
let mut supervisor = Supervisor::with_tenancy(2, TenancyPolicy::new(1, 8));
let shell = Arc::clone(
supervisor
.tenancy
.as_ref()
.ok_or("a policy must install a tenancy shell")?,
);
let tenant = Tenant::new("held")?;
let helper = Tenant::new("helper")?;
let gate = Arc::new(Notify::new());
let parked = Arc::new(AtomicUsize::new(0));
let root = supervisor.token.clone();
supervisor
.spawn_for(&tenant, {
let gate = Arc::clone(&gate);
let parked = Arc::clone(&parked);
move |_token| async move {
parked.fetch_add(1, Ordering::SeqCst);
gate.notified().await;
}
})
.await
.map_err(|refusal| format!("the occupier was refused: {refusal}"))?;
assert!(
yield_until(|| parked.load(Ordering::SeqCst) >= 1, 1_000).await,
"the occupier never parked, so the premise of this test did not hold"
);
let helper_gate = Arc::clone(&gate);
let helper_shell = Arc::clone(&shell);
let helper_tenant = tenant.clone();
let helper_root = root.clone();
supervisor
.spawn_for(&helper, move |_token| async move {
while helper_shell.counts_of(&helper_tenant).1 == 0 {
yield_now().await;
}
helper_gate.notify_waiters();
helper_root.cancel();
})
.await
.map_err(|refusal| format!("the helper was refused: {refusal}"))?;
let outcome = supervisor.claim_tenanted(&shell, &tenant).await;
assert_eq!(
outcome.err(),
Some(SpawnRefused::Cancelled),
"a supervisor cancelled while a grant was in flight refuses the admission"
);
assert!(
settle(&mut supervisor).await,
"the occupier and the helper must both finish"
);
let (in_flight, queued) = shell.counts_of(&tenant);
assert_eq!(
(in_flight, queued),
(0, 0),
"the permit the round had already charged to {tenant} must come back \
through the round; a bare permit returns to the pool and leaves the \
tenant charged for an admission it does not hold ({in_flight} in flight)"
);
assert_eq!(
supervisor.permits.available_permits(),
2,
"both permits are back in the pool, so nothing leaked"
);
Ok::<(), Box<dyn std::error::Error>>(())
})
}
#[cfg(feature = "script")]
#[test]
fn a_tenanted_admission_joins_a_bounded_backlog_and_keeps_the_set_at_the_ceiling()
-> Result<(), Box<dyn std::error::Error>> {
use super::REAP_PER_ADMISSION;
use crate::rt::tenancy::TenancyPolicy;
use crate::script::Tenant;
const LIMIT: usize = 64;
const FLOOD: usize = 32;
const ADMISSIONS: usize = 10_000;
const _: () = assert!(
REAP_PER_ADMISSION < FLOOD,
"a per-admission reap must be smaller than the flood it bounds"
);
block_on(async {
let mut supervisor = Supervisor::with_tenancy(LIMIT, TenancyPolicy::new(FLOOD, LIMIT));
let attacker = Tenant::new("attacker")?;
let neighbour = Tenant::new("neighbour")?;
let ended = Arc::new(AtomicUsize::new(0));
for _ in 0..FLOOD {
let ended = Arc::clone(&ended);
supervisor
.spawn_for(&attacker, move |_token| async move {
ended.fetch_add(1, Ordering::SeqCst);
})
.await
.map_err(|refusal| format!("a flood body was refused: {refusal}"))?;
}
assert!(
yield_until(|| ended.load(Ordering::SeqCst) >= FLOOD, 10_000).await,
"the flood bodies never ended, so the premise of this test did not hold"
);
for _ in 0..FLOOD {
yield_now().await;
}
let before = supervisor.stats().succeeded;
supervisor
.spawn_for(&neighbour, |_token| async {})
.await
.map_err(|refusal| format!("the neighbour was refused: {refusal}"))?;
let joined = supervisor.stats().succeeded.saturating_sub(before);
assert!(
joined <= u64::try_from(REAP_PER_ADMISSION)?,
"the neighbour's admission joined {joined} of the flood's finished tasks; \
at most {REAP_PER_ADMISSION} may land on its critical path"
);
let mut largest = 0_usize;
for admission in 0..ADMISSIONS {
let tenant = if admission % 2 == 0 {
&attacker
} else {
&neighbour
};
supervisor
.spawn_for(tenant, |_token| async {})
.await
.map_err(|refusal| format!("admission {admission} was refused: {refusal}"))?;
largest = largest.max(supervisor.set.len());
yield_now().await;
}
assert!(
largest <= LIMIT,
"the retained set reached {largest} tasks against an in-flight ceiling of {LIMIT}"
);
let report = supervisor.shutdown().await;
assert_eq!(
report.stats.succeeded,
u64::try_from(FLOOD + 1 + ADMISSIONS)?,
"every body is joined and counted by shutdown"
);
Ok(())
})
}
#[cfg(feature = "script")]
#[test]
fn a_parked_tenant_admission_keeps_its_place_in_the_queue()
-> Result<(), Box<dyn std::error::Error>> {
use crate::rt::sync::OwnedSemaphorePermit;
use crate::rt::tenancy::TenancyPolicy;
use crate::script::Tenant;
use std::future::Future;
use std::task::{Context, Poll, Waker};
const SLICES: u32 = 24;
const SLICE: Duration = Duration::from_millis(50);
block_on(async {
let mut supervisor = Supervisor::with_tenancy(1, TenancyPolicy::new(3, 8));
let shell = Arc::clone(
supervisor
.tenancy
.as_ref()
.ok_or("a policy must install a tenancy shell")?,
);
let pool = Arc::clone(&supervisor.permits);
let tenant = Tenant::new("queued")?;
let mut held = Some(match Arc::clone(&pool).try_acquire_owned() {
Ok(permit) => permit,
Err(_) => return Err("the pool did not start with its one permit".into()),
});
let mut claim = std::pin::pin!(supervisor.claim_tenanted(&shell, &tenant));
let mut probe = Context::from_waker(Waker::noop());
assert!(
claim.as_mut().poll(&mut probe).is_pending(),
"the claim must park on a spent pool rather than be admitted"
);
let mut second = queued_waiter(&shell, &tenant)?;
let mut third = queued_waiter(&shell, &tenant)?;
let mut slices = 0_u32;
while slices < SLICES {
assert!(
crate::rt::time::timeout(SLICE, claim.as_mut())
.await
.is_err(),
"the claim cannot resolve while the pool is spent"
);
assert_eq!(
shell.counts_of(&tenant).1,
3,
"all three waiters stay live while they are parked"
);
assert_eq!(
shell.retained_waiters(&tenant),
3,
"a waiter that keeps its place leaves no abandoned entry behind; \
the deque grew to {} entries over {slices} polls",
shell.retained_waiters(&tenant)
);
slices = slices.saturating_add(1);
}
let mut order: Vec<u8> = Vec::new();
let mut claim_out: Option<Result<Lease, SpawnRefused>> = None;
let mut claim_done = false;
for _ in 0..3 {
if let Some(permit) = held.take() {
shell.release(&tenant, permit);
}
let mut second_out: Option<OwnedSemaphorePermit> = None;
let mut third_out: Option<OwnedSemaphorePermit> = None;
crate::rt::time::timeout(
Duration::from_secs(5),
std::future::poll_fn(|context| {
if !claim_done
&& claim_out.is_none()
&& let Poll::Ready(value) = claim.as_mut().poll(context)
{
claim_done = true;
claim_out = Some(value);
}
if second_out.is_none()
&& let Poll::Ready(permit) = second.as_mut().poll(context)
{
second_out = Some(permit);
}
if third_out.is_none()
&& let Poll::Ready(permit) = third.as_mut().poll(context)
{
third_out = Some(permit);
}
if claim_done || second_out.is_some() || third_out.is_some() {
return Poll::Ready(());
}
Poll::Pending
}),
)
.await
.map_err(|_| "no waiter resolved after a permit was released into the round")?;
let index = if claim_out.is_some() {
match claim_out.take() {
Some(Ok(lease)) => {
drop(lease);
0
}
Some(Err(refusal)) => {
return Err(format!("the parked claim was refused: {refusal}").into());
}
None => return Err("the claim reported no outcome at all".into()),
}
} else if let Some(permit) = second_out {
held = Some(permit);
1
} else if let Some(permit) = third_out {
held = Some(permit);
2
} else {
return Err("a resolution was reported with no waiter behind it".into());
};
order.push(index);
}
assert_eq!(
order,
vec![0, 1, 2],
"one tenant's three waiters are served in arrival order: the claim \
that parked first, then the two behind it"
);
assert_eq!(
shell.retained_waiters(&tenant),
0,
"every waiter left the queue once it was served"
);
Ok::<(), Box<dyn std::error::Error>>(())
})
}
#[test]
fn a_waiting_spawn_keeps_its_place_in_the_pool_queue() -> Result<(), Box<dyn std::error::Error>>
{
const REGISTER_AT: Duration = Duration::from_millis(90);
const RELEASE_AT: Duration = Duration::from_millis(150);
const LIMIT: Duration = Duration::from_secs(2);
block_on(async {
let pool = Arc::new(crate::rt::sync::Semaphore::new(1));
let mut first = Supervisor::assembled(1, Arc::clone(&pool), Clock::wall());
let mut second = Supervisor::assembled(1, Arc::clone(&pool), Clock::wall());
let mut held = match Arc::clone(&pool).try_acquire_owned() {
Ok(permit) => Some(permit),
Err(_) => return Err("the shared pool did not start with its one permit".into()),
};
let mut first_claim = std::pin::pin!(first.claim());
let mut second_claim = std::pin::pin!(second.claim());
let mut first_out: Option<Option<Lease>> = None;
let mut second_out: Option<Option<Lease>> = None;
let started = std::time::Instant::now();
let mut second_started = false;
let mut released = false;
std::future::poll_fn(|context| {
let now = started.elapsed();
if !second_started && now >= REGISTER_AT {
second_started = true;
}
if first_out.is_none()
&& let Poll::Ready(value) = first_claim.as_mut().poll(context)
{
first_out = Some(value);
}
if second_started
&& second_out.is_none()
&& let Poll::Ready(value) = second_claim.as_mut().poll(context)
{
second_out = Some(value);
}
if !released && now >= RELEASE_AT {
if let Some(permit) = held.take() {
drop(permit);
}
released = true;
}
if first_out.is_some() || second_out.is_some() || now >= LIMIT {
return Poll::Ready(());
}
context.waker().wake_by_ref();
Poll::Pending
})
.await;
assert!(
released,
"the driver never reached the release moment within {LIMIT:?}"
);
assert!(
first_out.is_some(),
"the caller that arrived first never received the released permit, so \
its place in the pool's queue was forfeited; the caller that arrived \
second was served instead (second resolved: {})",
second_out.is_some()
);
assert!(
second_out.is_none(),
"the permit went to the caller that arrived second: a waiter that keeps \
its place is served in arrival order"
);
Ok::<(), Box<dyn std::error::Error>>(())
})
}
#[cfg(feature = "script")]
fn queued_waiter(
shell: &Arc<super::tenancy_support::TenancyShell>,
tenant: &crate::script::Tenant,
) -> Result<std::pin::Pin<Box<super::tenancy_support::WaitPermit>>, Box<dyn std::error::Error>>
{
match shell.admit(tenant) {
super::tenancy_support::Admission::Queued(waiter) => Ok(Box::pin(waiter)),
super::tenancy_support::Admission::Admitted(_) => {
Err("the pool was spent, so this arrival must have parked".into())
}
super::tenancy_support::Admission::Refused { limit } => {
Err(format!("the queue bound of {limit} was already reached").into())
}
super::tenancy_support::Admission::SupervisorQueueFull { limit } => {
Err(format!("the supervisor's waiting bound of {limit} was already reached").into())
}
}
}
#[test]
fn wait_idle_joins_every_task_without_cancelling_and_leaves_the_supervisor_usable() {
block_on(async {
let ran = Arc::new(AtomicU64::new(0));
let mut supervisor = Supervisor::new(3);
for _ in 0..40 {
let ran = Arc::clone(&ran);
supervisor
.spawn(move |token| async move {
yield_now().await;
if !token.is_cancelled() {
ran.fetch_add(1, Ordering::SeqCst);
}
})
.await;
}
supervisor.wait_idle().await;
let stats = supervisor.stats();
assert_eq!(stats.in_flight(), 0, "nothing is left in flight");
assert_eq!(stats.succeeded, 40, "every task completed, none cancelled");
assert_eq!(stats.cancelled, 0, "waiting is not stopping");
assert_eq!(ran.load(Ordering::SeqCst), 40, "every body ran to its end");
assert!(
!supervisor.is_cancelled(),
"the supervisor was not cancelled"
);
supervisor.spawn(|_token| async {}).await;
assert_eq!(
supervisor.wait_idle().await,
1,
"the next wait joins the next task"
);
assert_eq!(
supervisor.wait_idle().await,
0,
"an idle supervisor joins nothing"
);
});
}
#[test]
fn wait_idle_counts_reports_past_the_retention_cap_rather_than_losing_them() {
block_on(async {
let mut supervisor = Supervisor::new(2);
for _ in 0..2 {
supervisor.spawn(|_token| async {}).await;
}
supervisor.wait_idle().await;
for _ in 0..2 {
supervisor.spawn(|_token| async {}).await;
}
supervisor.wait_idle().await;
let stats = supervisor.stats();
let retained = drain(&mut supervisor).len();
assert_eq!(
u64::try_from(retained)
.ok()
.map(|kept| kept.saturating_add(stats.reports_dropped)),
Some(4),
"every outcome is either retained or counted as dropped"
);
});
}
#[cfg(all(unix, feature = "process"))]
#[test]
fn a_spawn_waiting_on_a_cleanup_owned_permit_is_admitted_once_the_owner_settles()
-> Result<(), Box<dyn std::error::Error>> {
block_on(async {
let mut supervisor = Supervisor::new(1);
let Some(permit) = supervisor.try_take() else {
return Err("a fresh supervisor of bound one has its permit".into());
};
let task = TaskId(u64::MAX);
supervisor.cleanup_owners.register(
task,
0x7fff_fff0,
Lease::plain(permit),
whole_capture(),
);
let admitted = crate::rt::time::timeout(
Duration::from_secs(5),
supervisor.spawn(|_token| async {}),
)
.await;
if admitted.is_err() {
return Err("the spawn never left the bound: the recheck did not release the owner's permit".into());
}
supervisor.wait_idle().await;
let settled = drain(&mut supervisor)
.into_iter()
.any(|outcome| matches!(outcome, TaskOutcome::CleanupSettled { task: settled, .. } if settled == task));
if !settled {
return Err("the owner's settlement was not reported".into());
}
assert_eq!(supervisor.stats().succeeded, 1, "the waiting body ran");
Ok(())
})
}
}