269 lines
12 KiB
Rust
269 lines
12 KiB
Rust
//! Find out when another actor dies, without it knowing or caring that you're
|
|
//! watching.
|
|
//!
|
|
//! Say one actor manages a pool of workers and needs to know when a worker
|
|
//! exits, so it can replace it. The worker does not need to know it is being
|
|
//! watched, and nothing about the worker's own behavior should change because
|
|
//! someone is watching it. That is what [`monitor`] is for: call
|
|
//! `monitor(target)` to get a [`Monitor`], and read exactly one [`Down`]
|
|
//! message off `monitor.rx` whenever `target` terminates, however it
|
|
//! terminates.
|
|
//!
|
|
//! ```
|
|
//! use smarm::{monitor, run, spawn, DownReason};
|
|
//!
|
|
//! run(|| {
|
|
//! let worker = spawn(|| {
|
|
//! // does some work, then returns
|
|
//! });
|
|
//! let pid = worker.pid();
|
|
//!
|
|
//! let m = monitor(pid);
|
|
//! let _ = worker.join();
|
|
//!
|
|
//! let down = m.rx.recv().expect("monitor channel closed before Down");
|
|
//! assert_eq!(down.pid, pid);
|
|
//! assert_eq!(down.reason, DownReason::Exit);
|
|
//! });
|
|
//! ```
|
|
//!
|
|
//! A monitor is one-directional and one-shot:
|
|
//!
|
|
//! - **One-directional**: the watcher learns that the target died, but the
|
|
//! target is completely unaffected. It never learns it was being watched,
|
|
//! and its own behavior and lifetime do not change because of the monitor.
|
|
//! This is the opposite of a [`link`](mod@crate::link), which is bidirectional:
|
|
//! linking two actors means an abnormal death on either side can bring the
|
|
//! other down too. Reach for a monitor when you just want to *know*; reach
|
|
//! for a link when a peer's crash should actually stop you.
|
|
//! - **One-shot**: you get exactly one [`Down`] per `monitor()` call, then the
|
|
//! channel closes. Calling `monitor` again on the same target (or a
|
|
//! different one) gives you an independent registration with its own
|
|
//! [`Monitor`] and its own one-shot channel; nothing stops you from
|
|
//! monitoring the same actor many times over; each call is watched and
|
|
//! fires on its own.
|
|
//!
|
|
//! ## Why a monitor never hands you the panic value
|
|
//!
|
|
//! If the target panicked, [`Down`] tells you *that* it panicked
|
|
//! ([`DownReason::Panic`]), but not the panic's payload. The payload has a
|
|
//! single owner: it is handed to whichever caller `join()`s the actor's
|
|
//! [`JoinHandle`](crate::JoinHandle), as a `JoinError`. A monitor only needs
|
|
//! to know that something went wrong, not reproduce the exact value that
|
|
//! caused it, so it gets the reason and nothing else.
|
|
//!
|
|
//! Monitoring a target that is already gone (it finished and was cleaned up,
|
|
//! or the pid never pointed at a real actor) is not an error: you get a
|
|
//! [`Down`] with [`DownReason::NoProc`] right away, instead of waiting
|
|
//! forever for something that already happened.
|
|
//!
|
|
//! ## Stopping a monitor early
|
|
//!
|
|
//! [`demonitor`] cancels a monitor before it fires. If the registration was
|
|
//! still live, it removes it and returns `Some` of the monitor's id: no
|
|
//! `Down` will arrive on that channel from here on. If the target had already
|
|
//! died and its `Down` already sent, there is nothing left to cancel and
|
|
//! `demonitor` returns `None`; the `Down` you already have (or that is
|
|
//! already sitting in the channel) is unaffected.
|
|
//!
|
|
//! If you want to cancel *and* make sure a `Down` that already arrived is
|
|
//! discarded without reading it, just drop the [`Monitor`]: dropping it closes
|
|
//! its receiver, and any queued `Down` is dropped along with it.
|
|
//!
|
|
//! ## Correctness notes for implementers
|
|
//!
|
|
//! A target that is still alive at the moment `monitor()` registers is
|
|
//! guaranteed to eventually produce a real `Down`: registration and the
|
|
//! target's own termination bookkeeping run under the same lock, so there is
|
|
//! no window in which the target could die without the just-added
|
|
//! registration seeing it. `demonitor` is similarly race-free against a target
|
|
//! that has since died and had its slot reused by a new, unrelated actor: it
|
|
//! is checked against the exact monitored incarnation, so it can never remove
|
|
//! a different actor's registration by accident, it simply reports `None`.
|
|
|
|
use crate::channel::{channel, Receiver, Sender};
|
|
use crate::pid::Pid;
|
|
use crate::scheduler::with_runtime;
|
|
|
|
/// Why a monitored actor went down.
|
|
///
|
|
/// Carries no payload: see the module docs for why a monitor never receives
|
|
/// the panic value itself.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum DownReason {
|
|
/// The target returned normally.
|
|
Exit,
|
|
/// The target panicked. The payload is delivered to the actor's joiner,
|
|
/// not to monitors.
|
|
Panic,
|
|
/// The target was cooperatively cancelled via `request_stop`.
|
|
Stopped,
|
|
/// The target was already gone (finished and reclaimed, or never alive)
|
|
/// at the moment `monitor()` was called.
|
|
NoProc,
|
|
}
|
|
|
|
/// A monitored actor's termination notice.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub struct Down {
|
|
/// The pid that was being monitored.
|
|
pub pid: Pid,
|
|
/// How it went down.
|
|
pub reason: DownReason,
|
|
}
|
|
|
|
/// A unique identifier for one [`monitor`] registration.
|
|
///
|
|
/// Opaque and `Copy`. Never reused for the life of the runtime, so if you
|
|
/// monitor the same target more than once, each call's id is distinct. This
|
|
/// is what lets [`demonitor`] tear down exactly one of several monitors on
|
|
/// the same target without disturbing the others.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
|
pub struct MonitorId(pub(crate) u64);
|
|
|
|
/// A live monitor: the receiving end of the one-shot [`Down`] channel, plus the
|
|
/// identity needed to [`demonitor`] it.
|
|
///
|
|
/// Read the notification from [`Monitor::rx`]. Not `Clone`, since only one
|
|
/// side is meant to consume it. Dropping a `Monitor` closes the receiving
|
|
/// end; if a `Down` had already arrived but was never read, it is discarded
|
|
/// along with it.
|
|
pub struct Monitor {
|
|
/// This registration's process-unique id.
|
|
pub id: MonitorId,
|
|
/// The pid being monitored.
|
|
pub target: Pid,
|
|
/// The one-shot channel the `Down` arrives on.
|
|
pub rx: Receiver<Down>,
|
|
}
|
|
|
|
/// Monitor `target`. Returns a [`Monitor`] whose `rx` receives exactly one
|
|
/// [`Down`].
|
|
///
|
|
/// If `target` is still live, the `Down` arrives when it terminates. If
|
|
/// `target` is already gone, a [`DownReason::NoProc`] `Down` is queued
|
|
/// immediately so the caller's `rx.recv()` returns without parking.
|
|
pub fn monitor<A>(target: Pid<A>) -> Monitor {
|
|
let target = target.erase();
|
|
let (tx, rx) = channel::<Down>();
|
|
|
|
// Implementation note: registration happens under the target's cold
|
|
// lock. `tx.clone()` takes the channel's own lock, a Channel-class
|
|
// RawMutex, which is explicitly permitted under a Leaf (cold) lock by
|
|
// the lock order documented in raw_mutex.rs. We must still not *send*
|
|
// under the lock, since `Sender::send` can unpark a parked receiver,
|
|
// and there's no reason to nest that.
|
|
let (id, registered) = with_runtime(|inner| {
|
|
let id = inner.alloc_monitor_id();
|
|
let registered = match inner.slot_at(target) {
|
|
Some(slot) => {
|
|
let mut cold = slot.cold.lock();
|
|
if slot.is_live_for(target) {
|
|
cold.monitors.push((id, tx.clone()));
|
|
true
|
|
} else {
|
|
false
|
|
}
|
|
}
|
|
None => false,
|
|
};
|
|
(id, registered)
|
|
});
|
|
|
|
if !registered {
|
|
let _ = tx.send(Down {
|
|
pid: target,
|
|
reason: DownReason::NoProc,
|
|
});
|
|
}
|
|
|
|
Monitor { id, target, rx }
|
|
}
|
|
|
|
/// Flag `target`'s tenancy as watchable: its death will stamp the slot's
|
|
/// terminal record (see [`terminal_reason`]), exactly as registering a name
|
|
/// does. The bridge calls this wherever a smarm pid is *encoded across the
|
|
/// boundary* — a contract reply, an introspection listing — because BEAM can
|
|
/// only watch pids it holds, and can only hold pids that crossed. Keeping the
|
|
/// bit rare is what keeps the record alive: anonymous never-exported churn
|
|
/// (holder threads, egress tasks) stays ineligible and cannot evict a
|
|
/// watchable tenancy's record from a LIFO-recycled slot.
|
|
///
|
|
/// Generation-checked and live-screened: marking a pid whose tenancy already
|
|
/// ended is a no-op — its record either exists (it was flagged before dying)
|
|
/// or is honestly unknowable. Same `Runtime::run()` context contract as
|
|
/// [`monitor`].
|
|
pub fn mark_watchable<A>(target: Pid<A>) {
|
|
let target = target.erase();
|
|
with_runtime(|inner| {
|
|
if let Some(slot) = inner.slot_at(target) {
|
|
// Cold lock FIRST: finalize publishes Done and checks the
|
|
// watchable bit under this same lock, so the mark either lands
|
|
// before finalize reads it (the death stamps) or observes the
|
|
// tenancy already dead (no-op). No lost-stamp window between an
|
|
// unlocked liveness read and the flag set.
|
|
let mut cold = slot.cold.lock();
|
|
if slot.is_live_for(target) {
|
|
cold.watchable = true;
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
/// The terminal [`DownReason`] of the tenancy `target` names, if that tenancy
|
|
/// ever registered a name and is the *most recent named* death of its slot:
|
|
/// finalize stamps the slot with `(generation, reason)` for once-registered
|
|
/// tenancies (anonymous green-thread churn does not stamp — nor evict), and
|
|
/// the record survives reclaim and the next tenant's install, until the next
|
|
/// *named* tenant of the slot itself dies. `None` means the pid never lived,
|
|
/// is still alive, never held a name, or its record was overwritten by a
|
|
/// later named 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
|
|
/// 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
|
|
/// down and its `Down` was already sent (or is already sitting in the
|
|
/// channel, unread).
|
|
///
|
|
/// This only stops a *future* `Down`. If you also want to discard a `Down`
|
|
/// that already arrived (or is about to, in a race with this call), drop `m`
|
|
/// instead of, or in addition to, calling this: dropping the [`Monitor`]
|
|
/// closes its receiver and any queued notice is discarded with it.
|
|
pub fn demonitor(m: &Monitor) -> Option<MonitorId> {
|
|
// Implementation note: the registration is removed under the target's
|
|
// cold lock, but the `Sender` is moved *out* and dropped only after the
|
|
// lock is released. Dropping the last sender runs `Sender::drop`, which
|
|
// may unpark a parked receiver; legal under a cold lock, but pointless
|
|
// to nest.
|
|
let removed: Option<(MonitorId, Sender<Down>)> = with_runtime(|inner| {
|
|
let slot = inner.slot_at(m.target)?;
|
|
let mut cold = slot.cold.lock();
|
|
if slot.generation() != m.target.generation() {
|
|
return None; // slot reused; the Down already fired
|
|
}
|
|
let pos = cold.monitors.iter().position(|(mid, _)| *mid == m.id)?;
|
|
Some(cold.monitors.remove(pos))
|
|
});
|
|
// `removed`'s sender drops here, outside the lock.
|
|
removed.map(|(id, _sender)| id)
|
|
}
|