use std::collections::HashMap as StdHashMap;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
use super::*;
use crate::atom::Atom;
use crate::loader::Instruction;
use crate::loader::decode::compact::Operand;
use crate::module::{Module, ModuleOrigin, ModuleRegistry, ResolvedImport, ResolvedImportTarget};
use crate::native::native_process::{
NativeContext, NativeHandler, NativeHandlerFactory, NativeOutcome,
};
use crate::native::{Capability, NativeEntry, ProcessContext};
use crate::process::ExitReason;
use crate::term::Term;
const MARKER: Atom = Atom::OK;
const DONE: Atom = Atom::TRUE;
const BACKSTOP: Atom = Atom::INFO;
const NEVER: Atom = Atom::BADKEY;
#[derive(Clone)]
struct Latch(Arc<(Mutex<bool>, Condvar)>);
impl Latch {
fn new() -> Self {
Latch(Arc::new((Mutex::new(false), Condvar::new())))
}
fn raise(&self) {
let (lock, cvar) = &*self.0;
*lock_or_recover(lock) = true;
cvar.notify_all();
}
fn wait(&self) {
let (lock, cvar) = &*self.0;
let deadline = Instant::now() + Duration::from_secs(30);
let mut raised = lock_or_recover(lock);
while !*raised {
let now = Instant::now();
assert!(now < deadline, "latch handshake timed out");
let (guard, _timeout) = cvar
.wait_timeout(raised, deadline - now)
.unwrap_or_else(|poisoned| poisoned.into_inner());
raised = guard;
}
}
}
fn contract_scheduler() -> Arc<Scheduler> {
let config = SchedulerConfig {
thread_count: Some(1),
dirty_cpu_threads: Some(1),
dirty_io_threads: Some(1),
dirty_queue_depth: Some(8),
..SchedulerConfig::default()
};
Arc::new(
Scheduler::new(config, Arc::new(ModuleRegistry::new()))
.unwrap_or_else(|error| panic!("scheduler starts: {error}")),
)
}
fn wait_until(deadline_ms: u64, mut predicate: impl FnMut() -> bool) {
let deadline = Instant::now() + Duration::from_millis(deadline_ms);
while !predicate() {
assert!(Instant::now() <= deadline, "condition timed out");
thread::sleep(Duration::from_millis(2));
}
}
fn assert_absent(window_ms: u64, mut predicate: impl FnMut() -> bool) {
let deadline = Instant::now() + Duration::from_millis(window_ms);
while Instant::now() < deadline {
assert!(!predicate(), "observed a state the contract forbids");
thread::sleep(Duration::from_millis(2));
}
}
fn wait_parked(scheduler: &Scheduler, pid: u64) {
wait_until(10_000, || {
lock_or_recover(&scheduler.shared.wait_set)
.waiting
.contains_key(&pid)
});
}
fn wait_exit(scheduler: &Scheduler, pid: u64) {
wait_until(10_000, || {
scheduler.shared.exit_tombstones.contains_key(&pid)
});
}
fn exit_value(scheduler: &Scheduler, pid: u64) -> Option<Term> {
scheduler.shared.exit_results.get(&pid).map(|r| r.root())
}
fn observed(sink: &Arc<Mutex<Vec<Atom>>>) -> Vec<Atom> {
lock_or_recover(sink).clone()
}
fn hold_gap_and<F>(scheduler: &Arc<Scheduler>, gap: ParkGap, action: F) -> JoinHandle<()>
where
F: Fn(&Scheduler, u64) + Send + 'static,
{
let go = Latch::new();
let done = Latch::new();
let pid_cell = Arc::new(AtomicU64::new(0));
let fired = Arc::new(AtomicBool::new(false));
{
let go_h = go.clone();
let done_h = done.clone();
let cell_h = Arc::clone(&pid_cell);
let fired_h = Arc::clone(&fired);
*lock_or_recover(&scheduler.shared.park_gap_hook) =
Some(Box::new(move |shared, g, pid| {
if g != gap || pid == shared.standard_io_pid || fired_h.swap(true, Ordering::AcqRel)
{
return;
}
cell_h.store(pid, Ordering::Release);
go_h.raise();
done_h.wait();
}));
}
let sched = Arc::clone(scheduler);
thread::spawn(move || {
go.wait();
let pid = pid_cell.load(Ordering::Acquire);
action(&sched, pid);
done.raise();
})
}
struct RecordingHandler {
stop_atom: Atom,
observed: Arc<Mutex<Vec<Atom>>>,
slices: Arc<AtomicUsize>,
}
impl NativeHandler for RecordingHandler {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
self.slices.fetch_add(1, Ordering::AcqRel);
let mut stop = false;
while let Some(term) = ctx.recv() {
if let Some(atom) = term.as_atom() {
lock_or_recover(&self.observed).push(atom);
if atom == self.stop_atom {
stop = true;
}
}
}
if stop {
NativeOutcome::Stop(ExitReason::Normal)
} else {
NativeOutcome::Wait
}
}
}
fn recording_factory(
stop_atom: Atom,
observed: &Arc<Mutex<Vec<Atom>>>,
slices: &Arc<AtomicUsize>,
) -> NativeHandlerFactory {
let observed = Arc::clone(observed);
let slices = Arc::clone(slices);
Box::new(move || {
Box::new(RecordingHandler {
stop_atom,
observed: Arc::clone(&observed),
slices: Arc::clone(&slices),
})
})
}
struct ExecutingHandler {
marker: Atom,
observed: Arc<Mutex<Vec<Atom>>>,
in_slice: Latch,
release: Latch,
first_done: bool,
}
impl NativeHandler for ExecutingHandler {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
if !self.first_done {
self.first_done = true;
self.in_slice.raise();
self.release.wait();
return NativeOutcome::Wait;
}
let mut stop = false;
while let Some(term) = ctx.recv() {
if let Some(atom) = term.as_atom() {
lock_or_recover(&self.observed).push(atom);
if atom == self.marker {
stop = true;
}
}
}
if stop {
NativeOutcome::Stop(ExitReason::Normal)
} else {
NativeOutcome::Wait
}
}
}
struct ConsumerHandler {
socket: Option<Arc<Mutex<Vec<Atom>>>>,
stop_atom: Atom,
observed: Arc<Mutex<Vec<Atom>>>,
slices: Arc<AtomicUsize>,
armed: Latch,
ack: Latch,
found_by_probe: Arc<AtomicBool>,
}
impl ConsumerHandler {
fn drain(&self, ctx: &mut NativeContext<'_>) -> Vec<Atom> {
let mut got = Vec::new();
while let Some(term) = ctx.recv() {
if let Some(atom) = term.as_atom() {
got.push(atom);
}
}
if let Some(socket) = &self.socket {
got.append(&mut lock_or_recover(socket));
}
got
}
}
impl NativeHandler for ConsumerHandler {
fn handle(&mut self, ctx: &mut NativeContext<'_>) -> NativeOutcome {
self.slices.fetch_add(1, Ordering::AcqRel);
let opening = self.drain(ctx);
lock_or_recover(&self.observed).extend_from_slice(&opening);
if opening.contains(&self.stop_atom) {
return NativeOutcome::Stop(ExitReason::Normal);
}
self.armed.raise();
self.ack.wait();
let probe = self.drain(ctx);
lock_or_recover(&self.observed).extend_from_slice(&probe);
if probe.contains(&MARKER) {
self.found_by_probe.store(true, Ordering::Release);
}
if probe.contains(&self.stop_atom) {
return NativeOutcome::Stop(ExitReason::Normal);
}
NativeOutcome::Wait
}
}
#[test]
fn c1_marker_to_a_parked_process_wakes_and_is_observed() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let pid = scheduler
.spawn_native(recording_factory(MARKER, &seen, &slices))
.expect("spawn native");
wait_parked(&scheduler, pid);
assert!(
scheduler.enqueue_atom_message(pid, MARKER),
"delivery to a live parked process returns true"
);
wait_exit(&scheduler, pid);
assert!(observed(&seen).contains(&MARKER), "marker observed");
scheduler.shutdown();
}
#[test]
fn c1_delivery_in_the_store_to_register_gap_is_observed_before_sleep() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let delivered = Arc::new(AtomicBool::new(false));
let delivered_hook = Arc::clone(&delivered);
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, pid| {
delivered_hook.store(s.enqueue_atom_message(pid, MARKER), Ordering::Release);
});
let pid = scheduler
.spawn_native(recording_factory(DONE, &seen, &slices))
.expect("spawn native");
wait_until(10_000, || observed(&seen).contains(&MARKER));
assert!(
delivered.load(Ordering::Acquire),
"in-gap delivery returned true"
);
wait_parked(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 2, "scheduled exactly once");
assert!(scheduler.enqueue_atom_message(pid, DONE));
wait_exit(&scheduler, pid);
assert_eq!(
slices.load(Ordering::Acquire),
3,
"no lost or double wakeup"
);
handle.join().expect("helper joins");
scheduler.shutdown();
}
#[test]
fn c1_delivery_in_the_register_to_recheck_gap_schedules_exactly_once() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let delivered = Arc::new(AtomicBool::new(false));
let delivered_hook = Arc::clone(&delivered);
let handle = hold_gap_and(&scheduler, ParkGap::WaitRegistered, move |s, pid| {
delivered_hook.store(s.enqueue_atom_message(pid, MARKER), Ordering::Release);
});
let pid = scheduler
.spawn_native(recording_factory(DONE, &seen, &slices))
.expect("spawn native");
wait_until(10_000, || observed(&seen).contains(&MARKER));
assert!(
delivered.load(Ordering::Acquire),
"in-gap delivery returned true"
);
wait_parked(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 2, "scheduled exactly once");
assert!(scheduler.enqueue_atom_message(pid, DONE));
wait_exit(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 3, "single additional slice");
handle.join().expect("helper joins");
scheduler.shutdown();
}
#[test]
fn c1_delivery_while_executing_merges_at_store_back_and_wakes() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let in_slice = Latch::new();
let release = Latch::new();
let (seen_f, in_slice_f, release_f) = (Arc::clone(&seen), in_slice.clone(), release.clone());
let pid = scheduler
.spawn_native(Box::new(move || {
Box::new(ExecutingHandler {
marker: MARKER,
observed: Arc::clone(&seen_f),
in_slice: in_slice_f.clone(),
release: release_f.clone(),
first_done: false,
})
}))
.expect("spawn native");
in_slice.wait();
let delivered = scheduler.enqueue_atom_message(pid, MARKER);
release.raise();
assert!(delivered, "delivery to an executing process returns true");
wait_exit(&scheduler, pid);
assert!(
observed(&seen).contains(&MARKER),
"executing-position marker observed after store-back merge + recheck"
);
scheduler.shutdown();
}
static C2_AWAIT_RUNS: AtomicUsize = AtomicUsize::new(0);
static C2_PARKED: AtomicBool = AtomicBool::new(false);
fn c2_gated_await_native(_args: &[Term], context: &mut ProcessContext) -> Result<Term, Term> {
C2_AWAIT_RUNS.fetch_add(1, Ordering::AcqRel);
let _call_id = context.request_await_suspend(None);
C2_PARKED.store(true, Ordering::Release);
Ok(Term::NIL)
}
fn build_module(name: Atom, code: Vec<Instruction>) -> Module {
let label_index = code
.iter()
.enumerate()
.filter_map(|(ip, instruction)| match instruction {
Instruction::Label { label } => Some((*label, ip)),
_ => None,
})
.collect();
Module {
name,
generation: 0,
origin: ModuleOrigin::Preloaded,
exports: StdHashMap::new(),
label_index,
code,
literals: Vec::new(),
constant_pool: Default::default(),
resolved_imports: Vec::new(),
lambdas: Vec::new(),
string_table: Vec::new(),
function_table: Vec::new(),
line_table: Vec::new(),
line_info: Vec::new(),
}
}
#[test]
fn c2_gated_suspension_retains_marker_and_observes_at_completion() {
C2_AWAIT_RUNS.store(0, Ordering::Release);
C2_PARKED.store(false, Ordering::Release);
let registry = Arc::new(ModuleRegistry::new());
let scheduler = Arc::new(
Scheduler::new(
SchedulerConfig {
thread_count: Some(1),
dirty_cpu_threads: Some(1),
dirty_io_threads: Some(1),
dirty_queue_depth: Some(8),
..SchedulerConfig::default()
},
Arc::clone(®istry),
)
.unwrap_or_else(|error| panic!("scheduler starts: {error}")),
);
let name = scheduler.shared.atom_table.intern("c2_gated");
let mut module = build_module(
name,
vec![
Instruction::CallExt {
arity: Operand::Unsigned(0),
import: Operand::Unsigned(0),
},
Instruction::Label { label: 1 },
Instruction::LoopRec {
fail: Operand::Label(2),
destination: Operand::X(1),
},
Instruction::RemoveMessage,
Instruction::Move {
source: Operand::X(1),
destination: Operand::X(0),
},
Instruction::Return,
Instruction::Label { label: 2 },
Instruction::Wait {
fail: Operand::Label(1),
},
],
);
module.resolved_imports.push(ResolvedImport {
module: name,
function: name,
arity: 0,
target: ResolvedImportTarget::Native(NativeEntry {
function: c2_gated_await_native,
dirty_kind: None,
capability: Capability::Pure,
}),
});
let module = registry.insert(module);
let pid = scheduler.spawn_process(&module);
wait_until(10_000, || C2_PARKED.load(Ordering::Acquire));
wait_parked(&scheduler, pid);
assert!(
scheduler.enqueue_atom_message(pid, MARKER),
"marker delivered to the gated-parked process returns true"
);
assert_absent(60, || {
scheduler.shared.exit_tombstones.contains_key(&pid)
|| C2_AWAIT_RUNS.load(Ordering::Acquire) != 1
});
assert!(
scheduler.wake_with_result(pid, Term::small_int(42)),
"completion resumes the gated await"
);
wait_exit(&scheduler, pid);
assert_eq!(
exit_value(&scheduler, pid),
Some(Term::atom(MARKER)),
"retained marker observed at completion"
);
assert_eq!(
C2_AWAIT_RUNS.load(Ordering::Acquire),
1,
"the await native was not re-executed by the marker"
);
scheduler.shutdown();
}
#[test]
fn c3_enqueue_to_a_dead_or_absent_pid_returns_false() {
let scheduler = contract_scheduler();
assert!(
!scheduler.enqueue_atom_message(999_999, MARKER),
"delivery to an absent pid returns false"
);
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let pid = scheduler
.spawn_native(recording_factory(MARKER, &seen, &slices))
.expect("spawn native");
wait_parked(&scheduler, pid);
assert!(scheduler.enqueue_atom_message(pid, MARKER));
wait_exit(&scheduler, pid);
wait_until(10_000, || {
!scheduler.shared.process_bodies.contains_key(&pid)
});
assert!(
!scheduler.enqueue_atom_message(pid, MARKER),
"delivery to a reaped pid returns false"
);
scheduler.shutdown();
}
#[test]
fn c3_true_then_death_before_next_slice_drops_the_marker_harmlessly() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let delivered = Arc::new(AtomicBool::new(false));
let delivered_hook = Arc::clone(&delivered);
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, pid| {
delivered_hook.store(s.enqueue_atom_message(pid, MARKER), Ordering::Release);
s.exit_signal(0, pid, ExitReason::Kill)
.expect("kill delivered");
});
let pid = scheduler
.spawn_native(recording_factory(NEVER, &seen, &slices))
.expect("spawn native");
wait_exit(&scheduler, pid);
handle.join().expect("helper joins");
assert!(
delivered.load(Ordering::Acquire),
"enqueue to the live (stored, not yet parked) process returned true"
);
assert!(
observed(&seen).is_empty(),
"marker dropped with the dead mailbox, never observed"
);
assert_absent(60, || slices.load(Ordering::Acquire) != 1);
assert!(
!scheduler.shared.process_bodies.contains_key(&pid),
"process body reaped"
);
let seen2 = Arc::new(Mutex::new(Vec::new()));
let slices2 = Arc::new(AtomicUsize::new(0));
let pid2 = scheduler
.spawn_native(recording_factory(MARKER, &seen2, &slices2))
.expect("spawn native");
wait_parked(&scheduler, pid2);
assert!(scheduler.enqueue_atom_message(pid2, MARKER));
wait_exit(&scheduler, pid2);
assert!(
observed(&seen2).contains(&MARKER),
"scheduler stayed healthy"
);
scheduler.shutdown();
}
#[test]
fn c4_delivery_between_arm_and_final_probe_is_seen_by_the_probe() {
let scheduler = contract_scheduler();
let socket = Arc::new(Mutex::new(Vec::new()));
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let armed = Latch::new();
let ack = Latch::new();
let found_by_probe = Arc::new(AtomicBool::new(false));
let (socket_f, seen_f, slices_f) =
(Arc::clone(&socket), Arc::clone(&seen), Arc::clone(&slices));
let (armed_f, ack_f, probe_f) = (armed.clone(), ack.clone(), Arc::clone(&found_by_probe));
let pid = scheduler
.spawn_native(Box::new(move || {
Box::new(ConsumerHandler {
socket: Some(Arc::clone(&socket_f)),
stop_atom: DONE,
observed: Arc::clone(&seen_f),
slices: Arc::clone(&slices_f),
armed: armed_f.clone(),
ack: ack_f.clone(),
found_by_probe: Arc::clone(&probe_f),
})
}))
.expect("spawn native");
armed.wait();
lock_or_recover(&socket).push(MARKER);
let backstop_delivered = scheduler.enqueue_atom_message(pid, BACKSTOP);
ack.raise();
assert!(
backstop_delivered,
"the durable backstop marker delivers to the executing process"
);
wait_until(10_000, || observed(&seen).contains(&BACKSTOP));
assert!(
found_by_probe.load(Ordering::Acquire),
"the final probe observed the in-window delivery"
);
assert_eq!(
observed(&seen),
vec![MARKER, BACKSTOP],
"probe caught the socket event in slice 1, before the mid-slice \
backstop became visible in slice 2"
);
wait_parked(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 2, "scheduled exactly once");
assert!(scheduler.enqueue_atom_message(pid, DONE));
wait_exit(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 3, "single additional slice");
scheduler.shutdown();
}
#[test]
fn c4_delivery_after_the_final_probe_is_caught_by_the_park_recheck() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let armed = Latch::new();
let ack = Latch::new();
let found_by_probe = Arc::new(AtomicBool::new(false));
let delivered = Arc::new(AtomicBool::new(false));
let delivered_hook = Arc::clone(&delivered);
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, pid| {
delivered_hook.store(s.enqueue_atom_message(pid, MARKER), Ordering::Release);
});
let (seen_f, slices_f) = (Arc::clone(&seen), Arc::clone(&slices));
let (armed_f, ack_f, probe_f) = (armed.clone(), ack.clone(), Arc::clone(&found_by_probe));
let pid = scheduler
.spawn_native(Box::new(move || {
Box::new(ConsumerHandler {
socket: None,
stop_atom: DONE,
observed: Arc::clone(&seen_f),
slices: Arc::clone(&slices_f),
armed: armed_f.clone(),
ack: ack_f.clone(),
found_by_probe: Arc::clone(&probe_f),
})
}))
.expect("spawn native");
armed.wait();
ack.raise();
wait_until(10_000, || observed(&seen).contains(&MARKER));
assert!(
delivered.load(Ordering::Acquire),
"in-gap delivery returned true"
);
assert!(
!found_by_probe.load(Ordering::Acquire),
"the marker was caught by the recheck, not the probe"
);
wait_parked(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 2, "scheduled exactly once");
assert!(scheduler.enqueue_atom_message(pid, DONE));
wait_exit(&scheduler, pid);
assert_eq!(slices.load(Ordering::Acquire), 3, "single additional slice");
handle.join().expect("helper joins");
scheduler.shutdown();
}
#[test]
fn bare_wake_before_registration_is_lost_which_is_why_markers_are_durable() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, pid| {
s.wake_process(pid);
});
let pid = scheduler
.spawn_native(recording_factory(MARKER, &seen, &slices))
.expect("spawn native");
wait_parked(&scheduler, pid);
assert_absent(60, || slices.load(Ordering::Acquire) != 1);
handle.join().expect("helper joins");
assert!(scheduler.enqueue_atom_message(pid, MARKER));
wait_exit(&scheduler, pid);
assert!(observed(&seen).contains(&MARKER), "durable marker observed");
assert_eq!(slices.load(Ordering::Acquire), 2, "exactly one wake slice");
scheduler.shutdown();
}
#[test]
fn kill_in_the_store_to_register_gap_leaves_no_stale_wait_set_entry() {
let scheduler = contract_scheduler();
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, pid| {
s.exit_signal(0, pid, ExitReason::Kill)
.expect("kill delivered");
});
let pid = scheduler
.spawn_native(recording_factory(NEVER, &seen, &slices))
.expect("spawn native");
wait_exit(&scheduler, pid);
handle.join().expect("helper joins");
scheduler.shutdown();
let ws = lock_or_recover(&scheduler.shared.wait_set);
assert!(
!ws.waiting.contains_key(&pid),
"stale registration withdrawn by the post-registration death recheck"
);
assert!(
!ws.woken.iter().any(|(woken, _)| *woken == pid),
"no dead pid stranded in woken"
);
drop(ws);
assert!(
!scheduler.shared.process_bodies.contains_key(&pid),
"process body reaped"
);
}
#[test]
fn cascade_kill_in_the_store_to_register_gap_finalizes_and_leaves_no_residue() {
let scheduler = contract_scheduler();
let seen_a = Arc::new(Mutex::new(Vec::new()));
let slices_a = Arc::new(AtomicUsize::new(0));
let source = scheduler
.spawn_native(recording_factory(NEVER, &seen_a, &slices_a))
.expect("spawn source");
wait_parked(&scheduler, source);
let seen_b = Arc::new(Mutex::new(Vec::new()));
let slices_b = Arc::new(AtomicUsize::new(0));
let scheduler_for_gap = Arc::clone(&scheduler);
let handle = hold_gap_and(&scheduler, ParkGap::WaitStored, move |s, target| {
let facility = supervision_integration::SchedulerLinkFacility {
shared: Arc::clone(&scheduler_for_gap.shared),
};
{
use crate::native::links::LinkFacility as _;
facility.link(source, target).expect("link established");
}
s.exit_signal(0, source, ExitReason::Kill)
.expect("kill delivered");
});
let target = scheduler
.spawn_native(recording_factory(NEVER, &seen_b, &slices_b))
.expect("spawn target");
wait_exit(&scheduler, source);
wait_exit(&scheduler, target);
handle.join().expect("helper joins");
scheduler.shutdown();
let ws = lock_or_recover(&scheduler.shared.wait_set);
for pid in [source, target] {
assert!(
!ws.waiting.contains_key(&pid),
"no stale registration for {pid}"
);
assert!(
!ws.woken.iter().any(|(woken, _)| *woken == pid),
"no dead pid {pid} stranded in woken"
);
}
drop(ws);
assert!(
!scheduler.shared.process_bodies.contains_key(&target),
"cascade-killed target's body reaped (previously stranded)"
);
assert!(
!scheduler.shared.process_bodies.contains_key(&source),
"source's body reaped"
);
}
#[test]
fn kill_storm_returns_the_wait_set_to_baseline() {
let scheduler = contract_scheduler();
let (baseline_waiting, baseline_woken) = {
let ws = lock_or_recover(&scheduler.shared.wait_set);
(
ws.waiting
.keys()
.copied()
.collect::<std::collections::BTreeSet<u64>>(),
ws.woken
.iter()
.map(|(pid, _)| *pid)
.collect::<std::collections::BTreeSet<u64>>(),
)
};
let mut pids = Vec::new();
for _ in 0..16 {
let seen = Arc::new(Mutex::new(Vec::new()));
let slices = Arc::new(AtomicUsize::new(0));
let pid = scheduler
.spawn_native(recording_factory(NEVER, &seen, &slices))
.expect("spawn native");
pids.push(pid);
}
for pid in &pids {
wait_parked(&scheduler, *pid);
}
for pid in &pids {
scheduler
.exit_signal(0, *pid, ExitReason::Kill)
.expect("kill delivered");
}
for pid in &pids {
wait_exit(&scheduler, *pid);
}
wait_until(10_000, || {
let ws = lock_or_recover(&scheduler.shared.wait_set);
let waiting: std::collections::BTreeSet<u64> = ws.waiting.keys().copied().collect();
let woken: std::collections::BTreeSet<u64> = ws.woken.iter().map(|(pid, _)| *pid).collect();
waiting == baseline_waiting && woken == baseline_woken
});
scheduler.shutdown();
}