mailbox: Phase 2 — registry resolves name/pid to a messageable actor (RFC 014)

Rip out the old name<->pid bimap; rebuild registry.rs as the live mailbox
directory. A name resolves to a SINGLE actor (many-per-label is pg's job); that
actor owns a SET of typed channels (channels are typed, so no single mailbox),
keyed by message TypeId. Resolution: name -> pid (one actor) -> Mailbox ->
channel for type M. Same name + different type parameter selects different
channels of the one actor, so capability separation (§4.7) needs no new type.

- register<M>(Name<M>, Sender<M>): captures the current actor's channel under a
  name (one step, per decision). Adds channels cumulatively; NameTaken only vs a
  different LIVE holder; stale/dangling bindings pruned on contact.
- send<M>(Name<M>, msg): the payoff — resolve, clone the Sender UNDER the Leaf
  lock, release, then send (Leaf -> Channel; a send can unpark). Returns the msg
  on Unresolved / NoChannel / Closed.
- whereis -> single live pid; unregister frees only the name (mailbox/other
  names survive). name_of (reverse lookup) dropped — deferred to introspection.
- Contained erasure: each channel is Box<dyn Any+Send> filed under TypeId::of M
  and downcast to its own keying type, so the downcast can't fail on good data
  (debug-asserted via the stored type_name). Phantom M on Name re-imposes type.
- One Leaf RawMutex (the fold), so a name send stays on a single Leaf — two
  Leaves at once is a hard panic. runtime.rs field/init unchanged (same Registry
  name + new()).

tests/registry.rs rewritten for the new semantics: send-by-name delivery,
many typed channels on one actor, same-name/different-type routing, NameTaken,
takeover across slot reuse, Unresolved/NoChannel. Also fold in a Phase 1 fixup:
the phantom test key was an enum whose variants tripped dead-code under the test
build (only caught now that I run -D warnings on --tests). Full suite +
order-checker green, warning-free across all targets.
This commit is contained in:
smarm-agent
2026-06-16 19:17:45 +00:00
parent 48bc636552
commit 48bdbada2b
4 changed files with 340 additions and 278 deletions
+100 -177
View File
@@ -1,221 +1,144 @@
//! Named registry tests. Run under the scheduler: registration requires a
//! live runtime and live actors.
//! Mailbox-registry tests (RFC 014). Run under the scheduler: registration
//! captures a live actor's channel, resolution checks liveness.
//!
//! Workers register their *own* inbox (`register` claims the current actor);
//! the root closure resolves and sends by name. A `ready` handshake closes the
//! register-then-send race without busy-waiting on `whereis`.
use smarm::{channel, name_of, register, run, spawn, unregister, whereis, RegisterError};
use smarm::{channel, register, run, send, spawn, unregister, whereis, Name, RegisterError, SendError};
const SVC: Name<u64> = Name::new("svc");
#[test]
fn register_whereis_roundtrip() {
fn register_then_send_by_name_delivers() {
run(|| {
let (tx, rx) = channel::<()>();
let (ready_tx, ready_rx) = channel::<()>();
let (tx, rx) = channel::<u64>();
let h = spawn(move || {
rx.recv().unwrap();
register(SVC, tx).unwrap();
ready_tx.send(()).unwrap();
assert_eq!(rx.recv().unwrap(), 42);
});
register("worker", h.pid()).unwrap();
assert_eq!(whereis("worker"), Some(h.pid()));
assert_eq!(name_of(h.pid()).as_deref(), Some("worker"));
tx.send(()).unwrap();
ready_rx.recv().unwrap(); // worker has registered
assert_eq!(whereis("svc"), Some(h.pid()));
send(SVC, 42).unwrap();
h.join().unwrap();
});
}
#[test]
fn register_is_idempotent_for_same_binding() {
fn one_actor_many_typed_channels_route_by_type() {
// Same name, two message types: capability separation falls out of the
// type parameter — `Name<u64>` and `Name<&str>` hit different channels of
// the one actor (RFC 014 §4.7).
run(|| {
let (tx, rx) = channel::<()>();
let (ready_tx, ready_rx) = channel::<()>();
let (cmd_tx, cmd_rx) = channel::<u64>();
let (adm_tx, adm_rx) = channel::<&'static str>();
let h = spawn(move || {
rx.recv().unwrap();
register(Name::<u64>::new("port"), cmd_tx).unwrap();
register(Name::<&'static str>::new("port"), adm_tx).unwrap();
ready_tx.send(()).unwrap();
assert_eq!(cmd_rx.recv().unwrap(), 7);
assert_eq!(adm_rx.recv().unwrap(), "halt");
});
register("svc", h.pid()).unwrap();
// Same name, same pid: a no-op Ok, not NameTaken.
register("svc", h.pid()).unwrap();
tx.send(()).unwrap();
ready_rx.recv().unwrap();
send(Name::<u64>::new("port"), 7u64).unwrap();
send(Name::<&'static str>::new("port"), "halt").unwrap();
h.join().unwrap();
});
}
#[test]
fn duplicate_name_on_live_holder_is_rejected() {
fn name_held_by_live_actor_is_taken() {
run(|| {
let (tx, rx) = channel::<()>();
let (tx2, rx2) = channel::<()>();
let (ready_tx, ready_rx) = channel::<()>();
let (tx_a, rx_a) = channel::<u64>();
let a = spawn(move || {
rx.recv().unwrap();
register(SVC, tx_a).unwrap();
ready_tx.send(()).unwrap();
assert_eq!(rx_a.recv().unwrap(), 0); // wait to be released
});
let b = spawn(move || {
rx2.recv().unwrap();
});
register("svc", a.pid()).unwrap();
assert_eq!(
register("svc", b.pid()),
Err(RegisterError::NameTaken { holder: a.pid() })
);
tx.send(()).unwrap();
tx2.send(()).unwrap();
ready_rx.recv().unwrap();
// Root tries to claim a live actor's name for itself -> NameTaken.
let (tx_b, _rx_b) = channel::<u64>();
assert_eq!(register(SVC, tx_b), Err(RegisterError::NameTaken { holder: a.pid() }));
send(SVC, 0).unwrap(); // release a (delivers to the holder, a)
a.join().unwrap();
});
}
#[test]
fn dead_holder_is_pruned_and_name_taken_over() {
// a registers, signals, then dies. Its slot index is typically reused by b;
// b's registration must prune the stale binding and take the name over.
run(|| {
let (rt1, rr1) = channel::<()>();
let (tx1, _rx1) = channel::<u64>();
let a = spawn(move || {
register(SVC, tx1).unwrap();
rt1.send(()).unwrap(); // then return -> die, _rx1 dropped
});
rr1.recv().unwrap();
let a_pid = a.pid();
a.join().unwrap(); // a is dead; "svc" now points at a stale pid
let (rt2, rr2) = channel::<()>();
let (tx2, rx2) = channel::<u64>();
let b = spawn(move || {
register(SVC, tx2).unwrap(); // takes over the freed name
rt2.send(()).unwrap();
assert_eq!(rx2.recv().unwrap(), 9);
});
rr2.recv().unwrap();
assert_ne!(b.pid(), a_pid); // distinct incarnation even if slot reused
assert_eq!(whereis("svc"), Some(b.pid()));
send(SVC, 9).unwrap();
b.join().unwrap();
});
}
#[test]
fn one_name_per_pid() {
fn send_errors_unresolved_and_no_channel() {
run(|| {
let (tx, rx) = channel::<()>();
// No actor at all.
assert!(matches!(send(Name::<u64>::new("ghost"), 1u64), Err(SendError::Unresolved(_))));
let (ready_tx, ready_rx) = channel::<()>();
let (tx, rx) = channel::<u64>();
let h = spawn(move || {
rx.recv().unwrap();
register(Name::<u64>::new("svc2"), tx).unwrap();
ready_tx.send(()).unwrap();
assert_eq!(rx.recv().unwrap(), 0);
});
register("first", h.pid()).unwrap();
assert_eq!(
register("second", h.pid()),
Err(RegisterError::PidAlreadyRegistered { name: "first".into() })
);
tx.send(()).unwrap();
ready_rx.recv().unwrap();
// Right actor, wrong message type: it has a u64 channel, not a String.
let e = send(Name::<String>::new("svc2"), "x".to_string());
assert!(matches!(e, Err(SendError::NoChannel(_))));
assert_eq!(e.unwrap_err().into_inner(), "x"); // message handed back
send(Name::<u64>::new("svc2"), 0u64).unwrap();
h.join().unwrap();
});
}
#[test]
fn registering_a_dead_pid_is_noproc() {
fn unregister_frees_the_name_only() {
run(|| {
let h = spawn(|| {});
let pid = h.pid();
h.join().unwrap();
assert_eq!(register("ghost", pid), Err(RegisterError::NoProc));
});
}
#[test]
fn binding_evaporates_on_death_and_name_is_reusable() {
run(|| {
let a = spawn(|| {});
let pid_a = a.pid();
register("svc", pid_a).unwrap();
a.join().unwrap();
// Dead holder: both lookup directions report unbound.
assert_eq!(whereis("svc"), None);
assert_eq!(name_of(pid_a), None);
// And the name is free for a successor.
let (tx, rx) = channel::<()>();
let b = spawn(move || {
rx.recv().unwrap();
});
register("svc", b.pid()).unwrap();
assert_eq!(whereis("svc"), Some(b.pid()));
tx.send(()).unwrap();
b.join().unwrap();
});
}
#[test]
fn register_over_a_dead_holder_succeeds_without_lookup_in_between() {
// The eviction path inside register() itself (not via whereis pruning).
run(|| {
let a = spawn(|| {});
let pid_a = a.pid();
register("svc", pid_a).unwrap();
a.join().unwrap();
let (tx, rx) = channel::<()>();
let b = spawn(move || {
rx.recv().unwrap();
});
register("svc", b.pid()).unwrap();
assert_eq!(whereis("svc"), Some(b.pid()));
assert_eq!(name_of(pid_a), None);
tx.send(()).unwrap();
b.join().unwrap();
});
}
#[test]
fn unregister_frees_both_directions() {
run(|| {
let (tx, rx) = channel::<()>();
let (ready_tx, ready_rx) = channel::<()>();
let (done_tx, done_rx) = channel::<()>();
let (tx, _rx) = channel::<u64>(); // _rx moves into the actor, kept open
let h = spawn(move || {
rx.recv().unwrap();
register(SVC, tx).unwrap();
let _keep_open = _rx;
ready_tx.send(()).unwrap();
done_rx.recv().unwrap(); // released over a separate channel
});
register("svc", h.pid()).unwrap();
ready_rx.recv().unwrap();
assert_eq!(whereis("svc"), Some(h.pid()));
assert_eq!(unregister("svc"), Some(h.pid()));
assert_eq!(whereis("svc"), None);
assert_eq!(name_of(h.pid()), None);
assert_eq!(unregister("svc"), None);
// The pid may take a new name afterwards.
register("svc2", h.pid()).unwrap();
assert_eq!(name_of(h.pid()).as_deref(), Some("svc2"));
tx.send(()).unwrap();
assert!(matches!(send(SVC, 1u64), Err(SendError::Unresolved(_))));
done_tx.send(()).unwrap();
h.join().unwrap();
});
}
#[test]
fn whereis_from_another_actor_and_usable_with_runtime_apis() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
let stopped = Arc::new(AtomicBool::new(false));
let stopped2 = stopped.clone();
run(move || {
let (tx, rx) = channel::<()>();
let svc = spawn(move || {
// Parks forever; only a request_stop ends it.
let _ = rx.recv();
});
register("stoppable", svc.pid()).unwrap();
let stopped3 = stopped2.clone();
let client = spawn(move || {
let pid = whereis("stoppable").expect("name must resolve cross-actor");
// The registry's pids plug into the rest of the runtime API.
smarm::request_stop(pid);
stopped3.store(true, Ordering::Relaxed);
});
client.join().unwrap();
svc.join().unwrap(); // a stopped actor joins Ok
drop(tx);
});
assert!(stopped.load(Ordering::Relaxed));
}
#[test]
fn registry_under_concurrent_churn_multi_thread() {
// Many actors racing to claim the same names while holders die: the
// bimap invariant (debug-asserted internally) and Erlang semantics must
// hold under real parallelism.
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
let wins = Arc::new(AtomicU32::new(0));
let wins2 = wins.clone();
smarm::init(smarm::Config::exact(4)).run(move || {
let mut handles = Vec::new();
for round in 0..8 {
let name = format!("contested-{}", round % 2);
for _ in 0..8 {
let name = name.clone();
let wins = wins2.clone();
handles.push(spawn(move || {
let me = smarm::self_pid();
match register(&name, me) {
Ok(()) => {
wins.fetch_add(1, Ordering::Relaxed);
smarm::yield_now();
// May have been pruned-by-contact never; we are
// alive, so our binding must still resolve to us.
assert_eq!(whereis(&name), Some(me));
assert_eq!(unregister(&name), Some(me));
}
Err(RegisterError::NameTaken { .. }) => {}
Err(e) => panic!("unexpected register error: {e}"),
}
}));
}
}
for h in handles {
h.join().unwrap();
}
});
// At least one registration per name must have succeeded.
assert!(wins.load(Ordering::Relaxed) >= 2);
}