monitor/registry: terminal-outcome record — a raced watch can recover the real down reason (soak sig 4)
A watch installed after its target's death has, until now, only NoProc to report — but the bridge's proxies install their native watch asynchronously after acquire returns, so a link established before a crash (from the BEAM's view) could still lose the panic's translated reason to that blanket NoProc (width-20 soak signature 4: link_test.exs:26, 1/600 full-suite, 3/2000 link-only, all whereis-miss; deterministic repro in the bridge suite). Two primitives, no change to monitor()'s own Erlang-faithful stale-pid semantics — the upgrade is the caller's deliberate act: - finalize_actor stamps the slot with (generation, DownReason) under the same cold-lock block that publishes the outcome. The record survives reclaim, registry pruning, and the next tenant's install; only the slot's next death overwrites it. terminal_reason(pid) reads it generation-matched. - resolve_name(name) is whereis with the corpse kept: the dead-holder arm returns the stored pid it prunes (NameResolution::Corpse) instead of discarding the only evidence of who died — whereis itself prunes on the way out, so a whereis-then-lookup consumer would find the evidence already destroyed. Live/Unbound match whereis's Some/None; the name heals exactly as before. Contract pinned in tests/terminal_outcome_after_death.rs: one record per way of dying (Exit/Panic/Stopped), no record while live, corpse capture + heal on resolve_name, record independence from registry pruning, survival across slot re-tenancy, overwrite at the next tenancy's death.
This commit is contained in:
+3
-3
@@ -72,13 +72,13 @@ pub use introspect::{StackInfo,
|
|||||||
#[cfg(feature = "observer")]
|
#[cfg(feature = "observer")]
|
||||||
pub use observer::{ObserverReply, ObserverRequest};
|
pub use observer::{ObserverReply, ObserverRequest};
|
||||||
pub use link::{link, trap_exit, unlink, ExitSignal};
|
pub use link::{link, trap_exit, unlink, ExitSignal};
|
||||||
pub use monitor::{demonitor, monitor, Down, DownReason, Monitor, MonitorId};
|
pub use monitor::{demonitor, monitor, terminal_reason, Down, DownReason, Monitor, MonitorId};
|
||||||
pub use mutex::{LockTimeout, Mutex, MutexGuard};
|
pub use mutex::{LockTimeout, Mutex, MutexGuard};
|
||||||
pub use pid::{Addressable, Erased, Name, Pid, RawPid};
|
pub use pid::{Addressable, Erased, Name, Pid, RawPid};
|
||||||
pub use pg::{dispatch, join, leave, members, members_as, pick, pick_as, Incarnation, Member, NodeId};
|
pub use pg::{dispatch, join, leave, members, members_as, pick, pick_as, Incarnation, Member, NodeId};
|
||||||
pub use registry::{
|
pub use registry::{
|
||||||
install, lookup_as, register, send, send_dyn, send_to, unregister, whereis, RegisterError,
|
install, lookup_as, register, resolve_name, send, send_dyn, send_to, unregister, whereis,
|
||||||
SendError,
|
NameResolution, RegisterError, SendError,
|
||||||
};
|
};
|
||||||
pub use runtime::{init, Config, Runtime};
|
pub use runtime::{init, Config, Runtime};
|
||||||
pub use scheduler::{
|
pub use scheduler::{
|
||||||
|
|||||||
@@ -177,6 +177,32 @@ pub fn monitor<A>(target: Pid<A>) -> Monitor {
|
|||||||
Monitor { id, target, rx }
|
Monitor { id, target, rx }
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The terminal [`DownReason`] of the tenancy `target` names, if that tenancy
|
||||||
|
/// is the *most recent* death of its slot: finalize stamps the slot with
|
||||||
|
/// `(generation, reason)`, and the record survives reclaim and the next
|
||||||
|
/// tenant's install, until that next tenant itself dies. `None` means the pid
|
||||||
|
/// never lived, is still alive, or its record was overwritten by a later
|
||||||
|
/// tenancy's death — callers fall back to `NoProc` semantics.
|
||||||
|
///
|
||||||
|
/// This exists for watch-installers that raced their target's death (bridge
|
||||||
|
/// soak signature 4): a `NoProc` observed at install time can be upgraded to
|
||||||
|
/// the real reason while the record still matches, which is exactly what an
|
||||||
|
/// install that had won the race would have delivered. It does NOT change
|
||||||
|
/// [`monitor`]'s own semantics — monitoring a stale pid still queues `NoProc`,
|
||||||
|
/// the same shape Erlang gives — the upgrade is the caller's deliberate act.
|
||||||
|
/// Same context contract as [`monitor`]: must run inside `Runtime::run()`.
|
||||||
|
pub fn terminal_reason<A>(target: Pid<A>) -> Option<DownReason> {
|
||||||
|
let target = target.erase();
|
||||||
|
with_runtime(|inner| {
|
||||||
|
let slot = inner.slot_at(target)?;
|
||||||
|
let cold = slot.cold.lock();
|
||||||
|
match cold.terminal {
|
||||||
|
Some((generation, reason)) if generation == target.generation() => Some(reason),
|
||||||
|
_ => None,
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// Cancel the monitor `m`. Returns `Some(id)` if a live registration was found
|
/// Cancel the monitor `m`. Returns `Some(id)` if a live registration was found
|
||||||
/// and removed, so no `Down` will arrive on `m.rx` from here on. Returns
|
/// and removed, so no `Down` will arrive on `m.rx` from here on. Returns
|
||||||
/// `None` if there was nothing left to remove: the target had already gone
|
/// `None` if there was nothing left to remove: the target had already gone
|
||||||
|
|||||||
@@ -486,6 +486,47 @@ pub fn whereis(name: &str) -> Option<Pid> {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// What a name is bound to, three-valued (bridge soak signature 4).
|
||||||
|
///
|
||||||
|
/// [`Live`](NameResolution::Live) is [`whereis`]'s `Some`.
|
||||||
|
/// [`Corpse`](NameResolution::Corpse) carries the *stored* holder pid of a
|
||||||
|
/// dead-but-unpruned binding — a state Erlang cannot represent (its name
|
||||||
|
/// death unregisters atomically; smarm's prune is lazy), captured here before
|
||||||
|
/// the prune that `whereis` performs discards it, so the caller can consult
|
||||||
|
/// [`terminal_reason`](crate::monitor::terminal_reason) for the tenancy's
|
||||||
|
/// real down reason. [`Unbound`](NameResolution::Unbound) matches Erlang's
|
||||||
|
/// unregistered name.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum NameResolution {
|
||||||
|
/// The stored holder is live (generation-checked); the binding stands.
|
||||||
|
Live(Pid),
|
||||||
|
/// The stored holder is dead. The binding was pruned on the way out —
|
||||||
|
/// the name heals exactly as `whereis` heals it; only the evidence is
|
||||||
|
/// returned instead of discarded. A second resolve is `Unbound`.
|
||||||
|
Corpse(Pid),
|
||||||
|
/// No binding stored (never registered, or already pruned by any reader).
|
||||||
|
Unbound,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Resolve `name` like [`whereis`], but keep the corpse: the dead-holder arm
|
||||||
|
/// returns the stored pid it pruned instead of a bare `None`. Same lock
|
||||||
|
/// discipline and pruning behavior as `whereis`; same `Runtime::run()`
|
||||||
|
/// context contract.
|
||||||
|
pub fn resolve_name(name: &str) -> NameResolution {
|
||||||
|
with_runtime(|inner| {
|
||||||
|
let mut reg = inner.registry.lock();
|
||||||
|
let Some(&pid) = reg.by_name.get(name) else {
|
||||||
|
return NameResolution::Unbound;
|
||||||
|
};
|
||||||
|
if live(inner, pid) {
|
||||||
|
NameResolution::Live(pid)
|
||||||
|
} else {
|
||||||
|
reg.prune_holder(pid);
|
||||||
|
NameResolution::Corpse(pid)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// Like [`whereis`], but returns a *typed* [`Pid<A>`] instead of a bare
|
/// Like [`whereis`], but returns a *typed* [`Pid<A>`] instead of a bare
|
||||||
/// [`Pid`], so a follow-up [`send_to`] is compile-checked instead of needing
|
/// [`Pid`], so a follow-up [`send_to`] is compile-checked instead of needing
|
||||||
/// the untyped [`send_dyn`] escape hatch. `None` if the name is unbound or its
|
/// the untyped [`send_dyn`] escape hatch. `None` if the name is unbound or its
|
||||||
|
|||||||
@@ -472,6 +472,14 @@ pub(crate) struct SlotCold {
|
|||||||
/// epoch-matched unpark.
|
/// epoch-matched unpark.
|
||||||
pub(crate) waiters: Vec<(Pid, u32)>,
|
pub(crate) waiters: Vec<(Pid, u32)>,
|
||||||
pub(crate) outcome: Option<Outcome>,
|
pub(crate) outcome: Option<Outcome>,
|
||||||
|
/// The slot's most recent *death*: `(generation, reason)`, stamped by
|
||||||
|
/// `finalize_actor` and deliberately never cleared — a new tenant's
|
||||||
|
/// install leaves it standing (it describes the previous tenancy), and
|
||||||
|
/// only the next death overwrites it. Read generation-matched via
|
||||||
|
/// [`terminal_reason`](crate::monitor::terminal_reason), so a watch that
|
||||||
|
/// raced its target's death can recover the real down reason instead of
|
||||||
|
/// a blanket `NoProc` (bridge soak signature 4).
|
||||||
|
pub(crate) terminal: Option<(u32, DownReason)>,
|
||||||
pub(crate) supervisor_channel: Option<Sender<Signal>>,
|
pub(crate) supervisor_channel: Option<Sender<Signal>>,
|
||||||
/// Watchers registered via `monitor()`, each tagged with its
|
/// Watchers registered via `monitor()`, each tagged with its
|
||||||
/// `MonitorId` so `demonitor` can remove exactly one. Each receives one
|
/// `MonitorId` so `demonitor` can remove exactly one. Each receives one
|
||||||
@@ -626,6 +634,7 @@ impl Slot {
|
|||||||
actor: None,
|
actor: None,
|
||||||
waiters: Vec::new(),
|
waiters: Vec::new(),
|
||||||
outcome: None,
|
outcome: None,
|
||||||
|
terminal: None,
|
||||||
supervisor_channel: None,
|
supervisor_channel: None,
|
||||||
monitors: Vec::new(),
|
monitors: Vec::new(),
|
||||||
links: Vec::new(),
|
links: Vec::new(),
|
||||||
@@ -1667,6 +1676,11 @@ fn finalize_actor(inner: &Arc<RuntimeInner>, pid: Pid, outcome: Outcome) {
|
|||||||
None => panic!("finalize_actor: actor vanished"),
|
None => panic!("finalize_actor: actor vanished"),
|
||||||
};
|
};
|
||||||
cold.outcome = Some(joiner_outcome);
|
cold.outcome = Some(joiner_outcome);
|
||||||
|
// Terminal record (soak sig 4): stamped before the generation ever
|
||||||
|
// bumps, under the cold lock, so a reader that resolved this pid can
|
||||||
|
// recover the reason after the slot moves on. Overwritten only by the
|
||||||
|
// slot's next death.
|
||||||
|
cold.terminal = Some((pid.generation(), down_reason));
|
||||||
slot.stop_ptr.store(std::ptr::null_mut(), Ordering::Release);
|
slot.stop_ptr.store(std::ptr::null_mut(), Ordering::Release);
|
||||||
// Done is published under the cold lock, so join's
|
// Done is published under the cold lock, so join's
|
||||||
// check-Done-or-register-waiter (also under it) can never miss: it
|
// check-Done-or-register-waiter (also under it) can never miss: it
|
||||||
|
|||||||
@@ -0,0 +1,216 @@
|
|||||||
|
//! The terminal-record contract (bridge soak signature 4): a watch installed
|
||||||
|
//! *after* its target's death — the async-install race the bridge's proxies
|
||||||
|
//! live with — must be able to recover the real down reason instead of a
|
||||||
|
//! blanket `NoProc`. Two primitives carry it:
|
||||||
|
//!
|
||||||
|
//! - `finalize_actor` stamps the slot with `(generation, DownReason)`; the
|
||||||
|
//! record survives reclaim, registry pruning, and the next tenant's
|
||||||
|
//! install, and is overwritten only by the slot's next death.
|
||||||
|
//! [`terminal_reason`] reads it generation-matched.
|
||||||
|
//! - [`resolve_name`] is `whereis` with the corpse kept: the dead-holder arm
|
||||||
|
//! returns the stored pid it prunes ([`NameResolution::Corpse`]) instead
|
||||||
|
//! of discarding the only evidence of *who* died. `Unbound` stays the
|
||||||
|
//! Erlang-shaped `noproc` for names that were never (or are no longer)
|
||||||
|
//! bound.
|
||||||
|
//!
|
||||||
|
//! `monitor()` of a stale pid still queues plain `NoProc` — the upgrade is a
|
||||||
|
//! caller's deliberate act, not a semantics change.
|
||||||
|
|
||||||
|
use smarm::{
|
||||||
|
init, request_stop, resolve_name, terminal_reason, CallError, Config, DownReason, GenServer,
|
||||||
|
GenServerBuilder, GenServerName, NameResolution,
|
||||||
|
};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
use std::time::Duration;
|
||||||
|
|
||||||
|
const TARGET: GenServerName<Target> = GenServerName::new("terminal_target");
|
||||||
|
|
||||||
|
/// Named server that panics on cast — the sig-4 death.
|
||||||
|
struct Target;
|
||||||
|
|
||||||
|
impl GenServer for Target {
|
||||||
|
type Call = ();
|
||||||
|
type Reply = ();
|
||||||
|
type Cast = ();
|
||||||
|
type Info = ();
|
||||||
|
type Timer = ();
|
||||||
|
|
||||||
|
fn handle_call(&mut self, _req: ()) {}
|
||||||
|
fn handle_cast(&mut self, _op: ()) {
|
||||||
|
panic!("terminal_target: induced panic");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Slot filler for the re-tenancy phase (distinct type, held alive).
|
||||||
|
struct Filler;
|
||||||
|
|
||||||
|
impl GenServer for Filler {
|
||||||
|
type Call = ();
|
||||||
|
type Reply = ();
|
||||||
|
type Cast = ();
|
||||||
|
type Info = ();
|
||||||
|
type Timer = ();
|
||||||
|
|
||||||
|
fn handle_call(&mut self, _req: ()) {}
|
||||||
|
fn handle_cast(&mut self, _op: ()) {}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct Observed {
|
||||||
|
exit_reason: Option<DownReason>,
|
||||||
|
panic_reason: Option<DownReason>,
|
||||||
|
stopped_reason: Option<DownReason>,
|
||||||
|
live_reason: Option<DownReason>,
|
||||||
|
live_resolution_is_live: bool,
|
||||||
|
unknown_resolution: NameResolution,
|
||||||
|
/// First resolve after the named target's panic — must be Corpse(old pid).
|
||||||
|
corpse_resolution_matches: bool,
|
||||||
|
/// Second resolve — the Corpse arm pruned, so the name has healed.
|
||||||
|
resolution_after_prune: NameResolution,
|
||||||
|
/// Read AFTER the prune above: the record is slot-side, not registry-side.
|
||||||
|
corpse_reason_after_prune: Option<DownReason>,
|
||||||
|
/// Record survives the slot being re-tenanted (new tenant still alive).
|
||||||
|
corpse_reason_after_reuse: Option<DownReason>,
|
||||||
|
/// ... and dies with the next tenancy's death (overwritten).
|
||||||
|
corpse_reason_after_tenant_death: Option<DownReason>,
|
||||||
|
tenant_reason: Option<DownReason>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn terminal_record_recovers_the_reason_a_raced_watch_lost() {
|
||||||
|
let out: Arc<Mutex<Option<Observed>>> = Arc::new(Mutex::new(None));
|
||||||
|
let out_w = out.clone();
|
||||||
|
|
||||||
|
// Tiny slab: prompt slot recycling for the re-tenancy phase.
|
||||||
|
init(Config::exact(2).max_actors(32)).run(move || {
|
||||||
|
// --- Plain actors: one record per way of dying. -------------------
|
||||||
|
let h = smarm::spawn(|| {});
|
||||||
|
let pid_exit = h.pid();
|
||||||
|
let _ = h.join();
|
||||||
|
let exit_reason = terminal_reason(pid_exit);
|
||||||
|
|
||||||
|
let h = smarm::spawn(|| panic!("induced"));
|
||||||
|
let pid_panic = h.pid();
|
||||||
|
let _ = h.join();
|
||||||
|
let panic_reason = terminal_reason(pid_panic);
|
||||||
|
|
||||||
|
let h = smarm::spawn(|| loop {
|
||||||
|
smarm::sleep(Duration::from_millis(2));
|
||||||
|
});
|
||||||
|
let pid_stop = h.pid();
|
||||||
|
request_stop(pid_stop);
|
||||||
|
let _ = h.join();
|
||||||
|
let stopped_reason = terminal_reason(pid_stop);
|
||||||
|
|
||||||
|
// --- The named target: live readings first. -----------------------
|
||||||
|
let target = GenServerBuilder::new(Target)
|
||||||
|
.named(TARGET)
|
||||||
|
.start()
|
||||||
|
.expect("name free at test start");
|
||||||
|
let old_pid = target.pid();
|
||||||
|
let live_reason = terminal_reason(old_pid);
|
||||||
|
let live_resolution_is_live =
|
||||||
|
resolve_name(TARGET.as_str()) == NameResolution::Live(old_pid.erase());
|
||||||
|
let unknown_resolution = resolve_name("terminal_never_bound");
|
||||||
|
|
||||||
|
// --- Kill it by panic; confirm death via the ref, NEVER the name
|
||||||
|
// (any name reader would take the prune arm and destroy the corpse
|
||||||
|
// precondition — the same trap stale_name_slot_reuse.rs documents).
|
||||||
|
let _ = target.cast(());
|
||||||
|
loop {
|
||||||
|
match target.call(()) {
|
||||||
|
Err(CallError::ServerDown) => break,
|
||||||
|
Ok(()) => smarm::sleep(Duration::from_millis(2)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let corpse_resolution_matches =
|
||||||
|
resolve_name(TARGET.as_str()) == NameResolution::Corpse(old_pid.erase());
|
||||||
|
let resolution_after_prune = resolve_name(TARGET.as_str());
|
||||||
|
let corpse_reason_after_prune = terminal_reason(old_pid);
|
||||||
|
|
||||||
|
// --- Re-tenant the freed slot; the record must outlive the install
|
||||||
|
// and die only with the next tenancy's death.
|
||||||
|
let mut fillers = Vec::new();
|
||||||
|
let mut tenant = None;
|
||||||
|
for i in 0..24 {
|
||||||
|
let name: &'static str = Box::leak(format!("terminal_filler_{i}").into_boxed_str());
|
||||||
|
let f = GenServerBuilder::new(Filler)
|
||||||
|
.named(GenServerName::<Filler>::new(name))
|
||||||
|
.start()
|
||||||
|
.expect("filler names are fresh");
|
||||||
|
let fp = f.pid();
|
||||||
|
let landed = fp.index() == old_pid.index();
|
||||||
|
fillers.push(f);
|
||||||
|
if landed {
|
||||||
|
tenant = Some((fillers.len() - 1, fp));
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let (tenant_at, tenant_pid) = tenant.expect(
|
||||||
|
"precondition: the freed slot must be re-tenanted within the tiny slab \
|
||||||
|
(slots are recycled; every filler is held alive)",
|
||||||
|
);
|
||||||
|
let corpse_reason_after_reuse = terminal_reason(old_pid);
|
||||||
|
|
||||||
|
request_stop(tenant_pid);
|
||||||
|
loop {
|
||||||
|
match fillers[tenant_at].call(()) {
|
||||||
|
Err(CallError::ServerDown) => break,
|
||||||
|
Ok(()) => smarm::sleep(Duration::from_millis(2)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let corpse_reason_after_tenant_death = terminal_reason(old_pid);
|
||||||
|
let tenant_reason = terminal_reason(tenant_pid);
|
||||||
|
|
||||||
|
*out_w.lock().unwrap() = Some(Observed {
|
||||||
|
exit_reason,
|
||||||
|
panic_reason,
|
||||||
|
stopped_reason,
|
||||||
|
live_reason,
|
||||||
|
live_resolution_is_live,
|
||||||
|
unknown_resolution,
|
||||||
|
corpse_resolution_matches,
|
||||||
|
resolution_after_prune,
|
||||||
|
corpse_reason_after_prune,
|
||||||
|
corpse_reason_after_reuse,
|
||||||
|
corpse_reason_after_tenant_death,
|
||||||
|
tenant_reason,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
let o = out.lock().unwrap().take().expect("runtime body completed");
|
||||||
|
assert_eq!(o.exit_reason, Some(DownReason::Exit), "{o:?}");
|
||||||
|
assert_eq!(o.panic_reason, Some(DownReason::Panic), "{o:?}");
|
||||||
|
assert_eq!(o.stopped_reason, Some(DownReason::Stopped), "{o:?}");
|
||||||
|
assert_eq!(
|
||||||
|
o.live_reason, None,
|
||||||
|
"live tenancy must have no record: {o:?}"
|
||||||
|
);
|
||||||
|
assert!(o.live_resolution_is_live, "{o:?}");
|
||||||
|
assert_eq!(o.unknown_resolution, NameResolution::Unbound, "{o:?}");
|
||||||
|
assert!(
|
||||||
|
o.corpse_resolution_matches,
|
||||||
|
"first post-death resolve must carry the corpse: {o:?}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
o.resolution_after_prune,
|
||||||
|
NameResolution::Unbound,
|
||||||
|
"the Corpse arm prunes — the name heals: {o:?}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
o.corpse_reason_after_prune,
|
||||||
|
Some(DownReason::Panic),
|
||||||
|
"the record is slot-side; registry pruning must not touch it: {o:?}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
o.corpse_reason_after_reuse,
|
||||||
|
Some(DownReason::Panic),
|
||||||
|
"a new tenant's install must leave the previous tenancy's record: {o:?}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
o.corpse_reason_after_tenant_death, None,
|
||||||
|
"the next death overwrites — the old generation no longer matches: {o:?}"
|
||||||
|
);
|
||||||
|
assert_eq!(o.tenant_reason, Some(DownReason::Stopped), "{o:?}");
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user