Files
smarm/tests/introspect.rs
T

512 lines
16 KiB
Rust
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! RFC 016 Chunk 1 — the read primitive. These exercise exactly what the RFC
//! promised: tests that assert an actor's state, parentage, names, mailbox
//! depth, and lifecycle counts directly off `snapshot()` / `actor_info()`
//! instead of sleeping-and-hoping.
use smarm::{
actor_info, channel, monitor, register, run, self_pid, send, snapshot, spawn, tree, tree_from,
ActorInfo, ActorState, Name, Pid, RuntimeSnapshot, SNAPSHOT_FORMAT_VERSION,
};
const SVC: Name<u64> = Name::new("svc");
/// Bounded poll on the introspection result itself (not a wall-clock sleep):
/// yield until `pred` holds for the given pid, panicking if it never does.
fn spin_until(pid: Pid, mut pred: impl FnMut(&smarm::ActorInfo) -> bool) -> smarm::ActorInfo {
for _ in 0..100_000 {
if let Some(info) = actor_info(pid) {
if pred(&info) {
return info;
}
}
smarm::yield_now();
}
panic!("actor {pid:?} never reached the expected state");
}
#[test]
fn snapshot_lists_actors_with_parent_edge() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let (_cmd_tx, cmd_rx) = channel::<u64>();
let h = spawn(move || {
register(Name::<u64>::new("w"), _cmd_tx).unwrap();
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap(); // park here until released
drop(cmd_rx);
});
ready_rx.recv().unwrap();
let me = self_pid();
let snap = snapshot();
assert_eq!(snap.format_version, SNAPSHOT_FORMAT_VERSION);
let worker = snap
.actors
.iter()
.find(|a| a.pid == h.pid())
.expect("worker present in snapshot");
assert_eq!(worker.names, vec!["w"]);
// `spawn` records the spawning actor as the parent (D9).
assert_eq!(worker.supervisor, me);
assert!(!worker.trap_exit);
assert_eq!((worker.monitors, worker.links, worker.joiners), (0, 0, 0));
// The root itself is on-CPU (it's running this code) and rooted under
// the forest sentinel.
let root = snap
.actors
.iter()
.find(|a| a.pid == me)
.expect("root present");
assert_eq!(root.state, ActorState::Running);
assert_eq!(root.supervisor, smarm::Pid::new(u32::MAX, u32::MAX));
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn parked_state_is_observable() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
// The worker has nothing to do but block on the empty gate channel, so
// it must reach Parked.
let info = spin_until(h.pid(), |a| a.state == ActorState::Parked);
assert_eq!(info.state, ActorState::Parked);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn mailbox_depth_counts_queued_messages() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let (cmd_tx, _cmd_rx) = channel::<u64>();
let h = spawn(move || {
// Publish the command inbox, then block on an unrelated gate so the
// queued commands are never drained while we observe them.
register(SVC, cmd_tx).unwrap();
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
drop(_cmd_rx);
});
ready_rx.recv().unwrap();
for i in 0..3 {
send(SVC, i).unwrap();
}
// Depth is a property of the queue, set synchronously by `send`, so it
// reads 3 regardless of the worker's scheduling state.
let via_snapshot = snapshot()
.actors
.into_iter()
.find(|a| a.pid == h.pid())
.expect("worker present");
assert_eq!(via_snapshot.mailbox_depth, 3);
assert_eq!(actor_info(h.pid()).unwrap().mailbox_depth, 3);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn monitor_count_is_visible() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
let _m = monitor(h.pid());
let info = actor_info(h.pid()).expect("worker present");
assert_eq!(info.monitors, 1);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn stale_and_forged_pids_return_none() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
let live = h.pid();
// Same slot index, wrong generation → stale, no such incarnation.
let stale = Pid::new(live.index(), live.generation().wrapping_add(7));
assert!(actor_info(stale).is_none());
// Out-of-range index → not in the slab at all.
let forged = Pid::new(u32::MAX - 1, 0);
assert!(actor_info(forged).is_none());
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn done_actor_is_a_tombstone() {
run(|| {
// Hold the join handle so the slot is NOT reclaimed when the actor
// exits: outstanding_handles stays > 0, leaving a Done tombstone to
// observe.
let h = spawn(|| {});
let pid = h.pid();
let info = spin_until(pid, |a| a.state == ActorState::Done);
assert_eq!(info.state, ActorState::Done);
// The Actor record is gone at finalize, so a tombstone reports root-less
// with empty lifecycle counts.
assert_eq!(info.supervisor, Pid::new(u32::MAX, u32::MAX));
assert_eq!((info.monitors, info.links, info.joiners), (0, 0, 0));
h.join().unwrap();
});
}
#[test]
fn tree_places_child_under_its_spawner() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
let me = self_pid();
let t = tree();
assert_eq!(t.format_version, SNAPSHOT_FORMAT_VERSION);
// The root is parented at the forest sentinel, so it's a genuine root,
// and the worker it spawned hangs beneath it.
let root = t
.roots
.iter()
.find(|n| n.info.pid == me)
.expect("root in forest");
assert!(!root.orphaned);
assert!(
root.children.iter().any(|c| c.info.pid == h.pid()),
"spawned worker should be a child of its spawner"
);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
/// D8 re-rooting and nesting, exercised on a synthetic snapshot via the public
/// `tree_from` — the live lifecycle race (parent reclaimed while child lives)
/// is exactly what's awkward to stage deterministically, which is why the fold
/// is testable in isolation.
#[test]
fn tree_from_nests_children_and_reroots_orphans() {
let root_sentinel = Pid::new(u32::MAX, u32::MAX);
let root_pid = Pid::new(0, 1);
let child = Pid::new(1, 1);
let orphan = Pid::new(2, 1);
let absent_parent = Pid::new(99, 1);
let mk = |pid: Pid, supervisor: Pid| ActorInfo {
pid,
names: Vec::new(),
state: ActorState::Running,
supervisor,
trap_exit: false,
monitors: 0,
links: 0,
joiners: 0,
mailbox_depth: 0,
overruns: 0,
messages_received: 0,
budget_cycles: 0,
stack: smarm::StackInfo {
reserve: 0,
guard: 0,
depth_high_water: 0,
parks_since_shrink: 0,
shrinks: 0,
},
};
let snap = RuntimeSnapshot {
format_version: SNAPSHOT_FORMAT_VERSION,
actors: vec![
mk(root_pid, root_sentinel),
mk(child, root_pid),
mk(orphan, absent_parent),
],
};
let t = tree_from(snap);
assert_eq!(t.roots.len(), 2);
let root = t
.roots
.iter()
.find(|n| n.info.pid == root_pid)
.expect("root present");
assert!(!root.orphaned);
assert_eq!(root.children.len(), 1);
assert_eq!(root.children[0].info.pid, child);
assert!(!root.children[0].orphaned);
let o = t
.roots
.iter()
.find(|n| n.info.pid == orphan)
.expect("orphan re-rooted");
assert!(
o.orphaned,
"an actor whose parent is absent must be flagged orphaned"
);
assert!(o.children.is_empty());
}
#[test]
fn overrun_count_increments_on_forced_preemption() {
run(|| {
// A worker that forces its slice to expire, then hits an observation
// point so the slice-expiry site fires and tallies one overrun.
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
smarm::preempt::expire_timeslice_for_test();
smarm::check!(); // preempt-yield here → one overrun tallied
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap(); // worker is past the forced preemption
let info = actor_info(h.pid()).expect("worker present");
assert!(
info.overruns >= 1,
"forced timeslice expiry should tally at least one overrun, got {}",
info.overruns
);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn messages_received_counts_dequeues() {
const MQ: Name<u64> = Name::new("mq");
const N: u64 = 5;
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (done_tx, done_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let (cmd_tx, cmd_rx) = channel::<u64>();
let h = spawn(move || {
register(MQ, cmd_tx).unwrap();
ready_tx.send(()).unwrap(); // sends don't count toward received
for _ in 0..N {
cmd_rx.recv().unwrap(); // each dequeue tallies one
}
done_tx.send(()).unwrap();
gate_rx.recv().unwrap(); // happens only after we've checked
});
ready_rx.recv().unwrap();
for i in 0..N {
send(MQ, i).unwrap();
}
done_rx.recv().unwrap(); // worker has drained all N
let info = actor_info(h.pid()).expect("worker present");
assert_eq!(info.messages_received, N, "one tally per dequeued message");
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[cfg(feature = "budget-accounting")]
#[test]
fn budget_cycles_accumulate_when_enabled() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
// A little work so the consumed slice is non-trivial, then park.
let mut acc = 0u64;
for i in 0..10_000u64 {
acc = acc.wrapping_add(i);
smarm::check!();
}
std::hint::black_box(acc);
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
// Once the worker has run and yielded (here, parked), its slice is
// charged.
let info = spin_until(h.pid(), |a| a.state == ActorState::Parked);
assert!(
info.budget_cycles > 0,
"budget should accrue after the actor runs and yields"
);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
// ---------------------------------------------------------------------------
// RFC 019 §8 — the stack introspection surface.
// ---------------------------------------------------------------------------
/// Burn ~`frames` × 4 KiB of stack with a yield at max depth, so the context
/// save samples the high-water there (RFC 019 §2: hwm is SAMPLED at
/// deschedule, not tracked continuously).
#[inline(never)]
fn burn_stack_yielding(frames: usize) -> u64 {
let mut local = [0u8; 4096];
local[0] = frames as u8;
let below = if frames == 0 {
smarm::yield_now();
0
} else {
burn_stack_yielding(frames - 1)
};
std::hint::black_box(&mut local);
below.wrapping_add(local[0] as u64)
}
#[test]
fn stack_info_reports_defaults_and_sampled_depth() {
run(|| {
let (ready_tx, ready_rx) = channel::<()>();
let (gate_tx, gate_rx) = channel::<()>();
let h = spawn(move || {
// ~32 KiB deep with a yield at the bottom: the sample point.
std::hint::black_box(burn_stack_yielding(8));
ready_tx.send(()).unwrap();
gate_rx.recv().unwrap();
});
ready_rx.recv().unwrap();
let info = spin_until(h.pid(), |a| a.state == ActorState::Parked);
let s = info.stack;
assert_eq!(s.reserve, 64 * 1024, "default reserve");
assert_eq!(
s.guard,
1024 * 1024,
"default guard (kernel stack_guard_gap convention)"
);
assert!(
s.depth_high_water >= 8 * 4096,
"hwm sampled at the deep yield: expected ≥ 32 KiB, got {}",
s.depth_high_water
);
assert!(
s.depth_high_water < s.reserve,
"depth {} cannot exceed the reserve {}",
s.depth_high_water,
s.reserve
);
// Parked at the gate right now, never shrunk (64 KiB reserve cannot
// cross the shrink threshold).
assert!(s.parks_since_shrink >= 1, "the gate park must be counted");
assert_eq!(s.shrinks, 0);
gate_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn stack_info_shrink_counters_are_live() {
use smarm::runtime::{Config, SHRINK_COOLDOWN, SHRINK_THRESHOLD};
use smarm::{spawn_with, SpawnOpts};
let rt = smarm::runtime::init(Config::exact(1));
rt.run(|| {
let (park_tx, park_rx) = channel::<()>();
let spike = 768 * 4096;
assert!(spike > SHRINK_THRESHOLD);
let worker = spawn_with(
SpawnOpts {
stack_reserve: Some(8 * 1024 * 1024),
..SpawnOpts::default()
},
move || {
std::hint::black_box(burn_stack_yielding(768));
for _ in 0..(SHRINK_COOLDOWN + 8) {
park_rx.recv().unwrap();
}
},
);
let wpid = worker.pid();
// Before any parks complete: the spike depth is visible.
let info = spin_until(wpid, |a| a.state == ActorState::Parked);
assert!(
info.stack.depth_high_water >= spike,
"spike should be sampled: {} < {spike}",
info.stack.depth_high_water
);
// Cross the cooldown, then read the counters live while the worker
// is parked waiting for the remaining rounds (post-join the slot is
// reclaimed and the generation check correctly hides it).
for _ in 0..(SHRINK_COOLDOWN + 2) {
spin_until(wpid, |a| a.state == ActorState::Parked);
park_tx.send(()).unwrap();
}
let info = spin_until(wpid, |a| {
a.state == ActorState::Parked && a.stack.shrinks >= 1
});
let s = info.stack;
assert!(
s.shrinks >= 1,
"cooldown was crossed with a spike above threshold"
);
assert!(
s.parks_since_shrink < SHRINK_COOLDOWN,
"counter must reset at shrink: {}",
s.parks_since_shrink
);
assert!(
s.depth_high_water < spike,
"hwm resets to the shallow park sp at shrink; got {}",
s.depth_high_water
);
for _ in 0..6 {
spin_until(wpid, |a| a.state == ActorState::Parked);
park_tx.send(()).unwrap();
}
worker.join().unwrap();
});
}