Tree snapshot of d9c62a8 (2026-08-18). The 20 source commits between
16ef583 (c9) and d9c62a8 were never pushed and the clone that held them
was lost; this commit carries their combined tree verbatim so the build
history stays auditable from the c1–c9 commits below it. Original
hashes as recorded in the session handoff:
c10 f03e94d pid targeting + auto-serialization (RemotePid, D14 name
on the wire); Phase 3 gate
c11 7ef4bad DownReason::Disconnected, wire tag 5
c12 d124162 remote monitors (Monitor/Demonitor/Down frames)
c13 9de967b connection-loss synthesis (A+B: Monitors::teardown +
unread-command Disconnected); Phase 4 gate
c14 7e822b7 eager pg eviction (reaper actor, ReaperInboxes)
dbe1a22 InboundVerdict::label(), trace::Event::ClusterInbound
31a9877 tests/channel.rs monitor-churn target gated on `go`
653559e Discovery::Withdrawn{name, addr}
c15 b41d76e distributed pg: Sync on NodeUp, Join/Leave broadcast,
NodeDown sweep, members_all; PgMsg wire type
c16 fafa881 pick_any / dispatch_any; Phase 5 complete
Phase 6 Tier A:
195c73e p4 NodeEvent::NodeDown(NodeInfo)
48fd766 p1 connector Candidate{name, addr, state}
ce8cf99 p2+p7 conn.rs select arms as Vec<Arm>; Outbound::Drained
bf24988 p6 RemotePid::from_local -> Option
9ae0380 p3 PeerStanding{Free, Claimed, Dialing}
c7d62a1 p11 cluster::Timing knobs, threaded by value
46f171d p11 cluster_disconnect un-ignored on SMARM_FAST_TIMING
Phase 6 Tier B:
391a9ae p5 cluster::RemoteDownReason{Local, Disconnected};
DownReason::Disconnected removed from core
7ddd908 p9 pg ctl channel unconditional, one cfg seam at spawn
d9c62a8 PeerNameMismatch parks the candidate; ClusterDial trace
Verified at d9c62a8: default 361/0, cluster 448/0, clippy --lib on
default / cluster / cluster+smarm-trace, fmt, 10x flake on
cluster_dial_mismatch, 5x on cluster_pg.
935 lines
35 KiB
Rust
935 lines
35 KiB
Rust
//! Process groups: one name, many actors.
|
|
//!
|
|
//! A process group is a named set of actors that you can look up, fan out to,
|
|
//! or pick a worker from. It is the natural home for a *worker pool* (several
|
|
//! interchangeable actors doing the same job), for *service discovery* (find
|
|
//! everyone currently offering some capability), and for *broadcast* (reach
|
|
//! every member of a group at once).
|
|
//!
|
|
//! The set is *live*: members [`join`] it, and a member that dies is removed
|
|
//! automatically. You never deregister a dead actor — there is no bookkeeping
|
|
//! to get wrong, and [`members`] / [`pick`] never hand you an actor that has
|
|
//! already gone.
|
|
//!
|
|
//! ## Joining and reading a group
|
|
//!
|
|
//! ```
|
|
//! use smarm::{channel, join, leave, members, pick, run, spawn};
|
|
//!
|
|
//! run(|| {
|
|
//! let (tx1, rx1) = channel::<()>();
|
|
//! let (tx2, rx2) = channel::<()>();
|
|
//! let w1 = spawn(move || { rx1.recv().unwrap(); });
|
|
//! let w2 = spawn(move || { rx2.recv().unwrap(); });
|
|
//!
|
|
//! // Workers join the group; `members` is the live view of it.
|
|
//! join("pool", w1.pid());
|
|
//! join("pool", w2.pid());
|
|
//! assert_eq!(members("pool").len(), 2);
|
|
//!
|
|
//! // One worker dies. Nothing tells the group — smarm evicts it
|
|
//! // automatically, so it is gone from `members` and never picked.
|
|
//! tx1.send(()).unwrap();
|
|
//! w1.join().unwrap();
|
|
//! assert_eq!(members("pool"), vec![w2.pid()]);
|
|
//! assert_eq!(pick("pool"), Some(w2.pid()));
|
|
//!
|
|
//! // Voluntary departure works too.
|
|
//! leave("pool", w2.pid());
|
|
//! assert!(pick("pool").is_none());
|
|
//!
|
|
//! tx2.send(()).unwrap();
|
|
//! w2.join().unwrap();
|
|
//! });
|
|
//! ```
|
|
//!
|
|
//! [`members`] returns every live member in the order they joined; [`pick`]
|
|
//! returns one of them, or `None` when the group is empty. The same actor can
|
|
//! belong to any number of groups at once, and joining a group it is already in
|
|
//! is a harmless no-op.
|
|
//!
|
|
//! ## Membership ends on its own
|
|
//!
|
|
//! You do not have to clean up after a member that dies. When an actor exits —
|
|
//! for any reason — it is removed from every group it had joined, before any
|
|
//! later read or send can observe it. [`leave`] is only for *voluntary*
|
|
//! departure, when a still-living actor wants out of a group.
|
|
//!
|
|
//! This is the main difference from keeping your own `Vec<Pid>`: a plain list
|
|
//! goes stale the instant a member dies, and you would have to notice and prune
|
|
//! it yourself. A group prunes itself.
|
|
//!
|
|
//! ## Sending to a group
|
|
//!
|
|
//! For a worker pool you usually want to hand a job to *one* available member.
|
|
//! [`dispatch`] picks a live member and sends it a message in a single step,
|
|
//! returning the member it reached:
|
|
//!
|
|
//! ```ignore
|
|
//! use smarm::{dispatch, join, Addressable};
|
|
//!
|
|
//! struct Job(String);
|
|
//! struct Worker;
|
|
//! impl Addressable for Worker { type Msg = Job; }
|
|
//!
|
|
//! // Each worker has published a `Pid<Worker>` inbox and joined the pool.
|
|
//! join("workers", worker_a);
|
|
//! join("workers", worker_b);
|
|
//!
|
|
//! // Route one job to whichever live worker `pick` lands on.
|
|
//! match dispatch::<Worker>("workers", Job("resize image".into())) {
|
|
//! Ok(who) => println!("sent to {who:?}"),
|
|
//! Err(returned) => println!("no worker took it: {returned:?}"),
|
|
//! }
|
|
//! ```
|
|
//!
|
|
//! When you want the pids themselves rather than to send right away, [`pick_as`]
|
|
//! and [`members_as`] return typed [`Pid<A>`](Pid)s for a homogeneous group, so
|
|
//! the follow-up send stays compile-checked. The untyped [`pick`] and
|
|
//! [`members`] are for mixed groups, where all you can rely on is identity.
|
|
//!
|
|
//! ## Groups vs. the registry
|
|
//!
|
|
//! A group is the many-actors counterpart to the [`registry`](crate::registry).
|
|
//! The registry binds a name to *at most one* actor and re-resolves it on every
|
|
//! send — what you want for a single well-known service. A group binds a name to
|
|
//! *many* actors, and one actor may sit in many groups. Reach for the registry
|
|
//! when there is exactly one of something; reach for a group when there is a set.
|
|
//!
|
|
//! ## Identity and clustering
|
|
//!
|
|
//! A group member is described by a [`Member`] — a [`Pid`] plus a [`NodeId`] and
|
|
//! an [`Incarnation`]. Everything on this page is **local**: you pass and
|
|
//! receive plain [`Pid`]s, and [`members`] / [`pick`] / [`dispatch`] only ever
|
|
//! name actors on this node (Erlang's `get_local_members`). With the `cluster`
|
|
//! feature a group also holds the members other nodes have announced, carried
|
|
//! under their [`NodeId`]; those never surface here — the cluster-wide reads
|
|
//! live in [`cluster::pg`](crate::cluster::pg) (`members_all` and friends) and
|
|
//! return a `Local | Remote` member type, since a [`Pid`] cannot hold a remote.
|
|
//!
|
|
//! ## Running context
|
|
//!
|
|
//! Every function here addresses the current runtime, so each must be called
|
|
//! from inside [`run`](crate::run) (that is, on an actor thread). Calling one
|
|
//! from outside a running runtime panics.
|
|
|
|
use crate::channel::{channel, Sender};
|
|
use crate::monitor::{register_monitor, unregister_monitor, Down, DownReason, MonitorId};
|
|
use crate::pid::{assert_type, Addressable, Pid};
|
|
use crate::registry::{send_to, SendError};
|
|
use crate::scheduler::{spawn_under, with_runtime};
|
|
use std::collections::HashMap;
|
|
|
|
/// A cluster node handle. A `u32` integer handle, *not* an interned atom — the
|
|
/// single deliberate divergence from the BEAM wire shape.
|
|
#[derive(Copy, Clone, PartialEq, Eq, Hash, Debug)]
|
|
pub struct NodeId(u32);
|
|
|
|
impl NodeId {
|
|
#[inline]
|
|
pub const fn new(v: u32) -> Self {
|
|
Self(v)
|
|
}
|
|
#[inline]
|
|
pub const fn get(self) -> u32 {
|
|
self.0
|
|
}
|
|
}
|
|
|
|
impl From<u32> for NodeId {
|
|
#[inline]
|
|
fn from(v: u32) -> Self {
|
|
Self(v)
|
|
}
|
|
}
|
|
|
|
/// A node's incarnation epoch — the BEAM `Creation` field adopted verbatim: it
|
|
/// separates a crashed node from its restart. Fixed for the life of a run
|
|
/// until clustering supplies a real one.
|
|
#[derive(Copy, Clone, PartialEq, Eq, Hash, Debug)]
|
|
pub struct Incarnation(u32);
|
|
|
|
impl Incarnation {
|
|
#[inline]
|
|
pub const fn new(v: u32) -> Self {
|
|
Self(v)
|
|
}
|
|
#[inline]
|
|
pub const fn get(self) -> u32 {
|
|
self.0
|
|
}
|
|
}
|
|
|
|
impl From<u32> for Incarnation {
|
|
#[inline]
|
|
fn from(v: u32) -> Self {
|
|
Self(v)
|
|
}
|
|
}
|
|
|
|
/// The fixed single-node identity used until clustering supplies real values.
|
|
/// Carried like `wake_slot` so the public API never has to change to acquire it.
|
|
pub const DEFAULT_NODE_ID: NodeId = NodeId(0);
|
|
/// The fixed incarnation for the single-node default. Non-zero so it never
|
|
/// collides with a BEAM "any creation" wildcard at interop time.
|
|
pub const DEFAULT_INCARNATION: Incarnation = Incarnation(1);
|
|
|
|
/// A group member's full identity: `(node, incarnation, pid)`.
|
|
///
|
|
/// Deliberately a field-for-field image of a modern BEAM pid (`NEW_PID_EXT`).
|
|
/// In memory it is a plain struct — no wire packing; the packed representation
|
|
/// belongs to the remote-reference boundary, not here.
|
|
#[derive(Copy, Clone, PartialEq, Eq, Hash, Debug)]
|
|
pub struct Member {
|
|
/// Which node the pid lives on. `DEFAULT_NODE_ID` while single-node.
|
|
pub node: NodeId,
|
|
/// The node's incarnation epoch at the time of joining.
|
|
pub incarnation: Incarnation,
|
|
/// Pure local slot identity — unchanged; cluster identity is layered
|
|
/// *around* it here rather than overloading `Pid::generation`.
|
|
pub pid: Pid,
|
|
}
|
|
|
|
/// One membership: a [`Member`] and the id of the monitor that watches its
|
|
/// liveness. The monitor's `Down` is delivered to the group reaper's single
|
|
/// inbox (see [`ProcessGroups::deaths`]), so the membership carries only what
|
|
/// [`leave`] needs to tear the registration down: the id.
|
|
pub(crate) struct Membership {
|
|
pub(crate) member: Member,
|
|
/// `None` for a remote member (cluster): the origin node is its liveness
|
|
/// authority; nothing here watches it.
|
|
pub(crate) monitor: Option<MonitorId>,
|
|
}
|
|
|
|
/// The store: `name → multiset<Member>`. Within a single group a `Member`
|
|
/// appears at most once (`join` is idempotent); the *multiset* framing is for
|
|
/// cluster-readiness — the same pid is freely a member of many groups, and the
|
|
/// width admits multiples in general.
|
|
///
|
|
/// Locking discipline. Held under one Leaf-class `RawMutex` on `RuntimeInner`,
|
|
/// mirroring the registry, and never held together with another Leaf lock (it
|
|
/// never touches the registry or a slot's cold lock). Monitor registration and
|
|
/// removal take the target's cold lock (also Leaf), so they run *before* /
|
|
/// *after* the group lock, never under it — see [`join`] for the ordering that
|
|
/// makes that safe. Nothing under this lock ever touches a channel.
|
|
///
|
|
/// Eviction is *eager*: every membership's monitor delivers to the one
|
|
/// `deaths` channel, drained by a per-run reaper actor that sweeps the dead
|
|
/// pid out of every group the moment its `Down` is scheduled. The read path
|
|
/// keeps a slot-liveness backstop for the window between a death and the
|
|
/// reaper's turn.
|
|
pub(crate) struct ProcessGroups {
|
|
groups: HashMap<String, Vec<Membership>>,
|
|
/// The reaper's inboxes: every membership monitor is registered against
|
|
/// a clone of `deaths`. `None` until the first `join` of a run spawns
|
|
/// the reaper; a stale one (receiver gone with the previous run's
|
|
/// teardown) is detected via `receiver_alive` and replaced.
|
|
reaper: Option<ReaperInboxes>,
|
|
/// `NodeId → node name` for every peer with members in the store, kept
|
|
/// by the pg actor under this lock, so a stored remote member can be
|
|
/// rendered back to its wire identity without asking anyone.
|
|
#[cfg(feature = "cluster")]
|
|
node_names: HashMap<NodeId, String>,
|
|
}
|
|
|
|
impl ProcessGroups {
|
|
pub(crate) fn new() -> Self {
|
|
Self {
|
|
groups: HashMap::new(),
|
|
reaper: None,
|
|
#[cfg(feature = "cluster")]
|
|
node_names: HashMap::new(),
|
|
}
|
|
}
|
|
|
|
/// Forget the reaper. Called at the start of every `run()` so a stopped
|
|
/// reaper from a previous run is never sent to; `join` respawns.
|
|
pub(crate) fn reset_reaper(&mut self) {
|
|
self.reaper = None;
|
|
}
|
|
|
|
/// Insert `ms` into `group`. Idempotent on the *member*: `false` means the
|
|
/// member was already present and nothing changed; `true` means inserted.
|
|
pub(crate) fn join(&mut self, group: &str, ms: Membership) -> bool {
|
|
let v = self.groups.entry(group.to_owned()).or_default();
|
|
if v.iter().any(|e| e.member == ms.member) {
|
|
return false;
|
|
}
|
|
v.push(ms);
|
|
true
|
|
}
|
|
|
|
/// Remove `member`'s membership from `group`, returning it (so the caller
|
|
/// can unregister its monitor outside the lock). An emptied group is pruned.
|
|
pub(crate) fn leave(&mut self, group: &str, member: Member) -> Option<Membership> {
|
|
let v = self.groups.get_mut(group)?;
|
|
let pos = v.iter().position(|e| e.member == member)?;
|
|
let removed = v.remove(pos);
|
|
if v.is_empty() {
|
|
self.groups.remove(group);
|
|
}
|
|
Some(removed)
|
|
}
|
|
|
|
/// The one dumb eviction primitive: drop every member matching `pred` from
|
|
/// every group, pruning emptied groups, and return the evicted
|
|
/// memberships with the group each was in. The primitive does not know
|
|
/// *why* a member leaves; that is the caller's concern. Its callers are
|
|
/// the reaper (a local death) and the cluster's node-down / re-sync
|
|
/// sweeps — all over this same predicate path, which is the whole reason
|
|
/// to shape eviction as a predicate. Insertion order within a group is
|
|
/// preserved (`members` / `pick` are order-stable).
|
|
pub(crate) fn remove_where(
|
|
&mut self,
|
|
mut pred: impl FnMut(&Member) -> bool,
|
|
) -> Vec<(String, Membership)> {
|
|
let mut evicted = Vec::new();
|
|
self.groups.retain(|g, v| {
|
|
let mut i = 0;
|
|
while i < v.len() {
|
|
if pred(&v[i].member) {
|
|
evicted.push((g.clone(), v.remove(i)));
|
|
} else {
|
|
i += 1;
|
|
}
|
|
}
|
|
!v.is_empty()
|
|
});
|
|
evicted
|
|
}
|
|
|
|
/// Raw enumeration of a group's members — no liveness filtering. Used by
|
|
/// tests to assert storage state independently of the read-path backstop.
|
|
#[cfg(test)]
|
|
fn members_of(&self, group: &str) -> Vec<Member> {
|
|
self.groups
|
|
.get(group)
|
|
.map(|v| v.iter().map(|e| e.member).collect())
|
|
.unwrap_or_default()
|
|
}
|
|
|
|
/// Live members of `group` **on `node`**, in insertion order. The
|
|
/// `is_live` oracle is the read-path backstop: a member whose slot is
|
|
/// already dead is dropped from the *result* even if the reaper has not
|
|
/// swept it yet. Backstop only — the entry stays in storage; eviction is
|
|
/// the reaper's job. The node filter is what keeps the local API local:
|
|
/// a remote member's `pid` is another node's slot bits, meaningless to
|
|
/// `is_live` and to any local send.
|
|
fn members_where(
|
|
&self,
|
|
group: &str,
|
|
node: NodeId,
|
|
mut is_live: impl FnMut(Pid) -> bool,
|
|
) -> Vec<Pid> {
|
|
self.groups
|
|
.get(group)
|
|
.map(|v| {
|
|
v.iter()
|
|
.filter(|e| e.member.node == node)
|
|
.map(|e| e.member.pid)
|
|
.filter(|&p| is_live(p))
|
|
.collect()
|
|
})
|
|
.unwrap_or_default()
|
|
}
|
|
|
|
/// The first live member of `group` on `node` in insertion order —
|
|
/// stateless first-live `pick`, with the same read-path backstop and node
|
|
/// filter as `members_where`.
|
|
fn first_member_where(
|
|
&self,
|
|
group: &str,
|
|
node: NodeId,
|
|
mut is_live: impl FnMut(Pid) -> bool,
|
|
) -> Option<Pid> {
|
|
self.groups
|
|
.get(group)?
|
|
.iter()
|
|
.filter(|e| e.member.node == node)
|
|
.map(|e| e.member.pid)
|
|
.find(|&p| is_live(p))
|
|
}
|
|
}
|
|
|
|
/// The store's cluster-side surface: raw reads the pg actor needs to speak
|
|
/// for this node (`Sync`, membership checks) and the peer-name memo. One
|
|
/// `cfg` block: everything here exists only when there is a mesh.
|
|
#[cfg(feature = "cluster")]
|
|
impl ProcessGroups {
|
|
/// Does `group` hold `member` right now? (Raw storage, no liveness.)
|
|
pub(crate) fn contains(&self, group: &str, member: &Member) -> bool {
|
|
self.groups
|
|
.get(group)
|
|
.is_some_and(|v| v.iter().any(|e| e.member == *member))
|
|
}
|
|
|
|
/// Every stored member of `group`, any node, insertion order. Raw storage.
|
|
pub(crate) fn all_of(&self, group: &str) -> Vec<Member> {
|
|
self.groups
|
|
.get(group)
|
|
.map(|v| v.iter().map(|e| e.member).collect())
|
|
.unwrap_or_default()
|
|
}
|
|
|
|
/// `(group, [pid])` for every group with a member on `node` — the
|
|
/// `Sync` payload. Raw storage; groups with no such member are omitted.
|
|
pub(crate) fn groups_on(&self, node: NodeId) -> Vec<(String, Vec<Pid>)> {
|
|
let mut out: Vec<(String, Vec<Pid>)> = self
|
|
.groups
|
|
.iter()
|
|
.filter_map(|(g, v)| {
|
|
let pids: Vec<Pid> = v
|
|
.iter()
|
|
.filter(|e| e.member.node == node)
|
|
.map(|e| e.member.pid)
|
|
.collect();
|
|
(!pids.is_empty()).then(|| (g.clone(), pids))
|
|
})
|
|
.collect();
|
|
out.sort_by(|a, b| a.0.cmp(&b.0));
|
|
out
|
|
}
|
|
|
|
/// Record / forget the name behind a peer's `NodeId`.
|
|
pub(crate) fn set_node_name(&mut self, node: NodeId, name: String) {
|
|
self.node_names.insert(node, name);
|
|
}
|
|
pub(crate) fn forget_node_name(&mut self, node: NodeId) {
|
|
self.node_names.remove(&node);
|
|
}
|
|
pub(crate) fn node_name(&self, node: NodeId) -> Option<&str> {
|
|
self.node_names.get(&node).map(String::as_str)
|
|
}
|
|
}
|
|
|
|
/// The group reaper: one detached actor per run, spawned by the first `join`,
|
|
/// parked on the shared `deaths` inbox. Every local membership's monitor
|
|
/// delivers here, so a death is swept out of *every* group it joined as soon
|
|
/// as the reaper is scheduled — no group operation has to happen first.
|
|
/// Sweeps by `(node, pid)`: only local members, since a remote member's pid
|
|
/// bits are meaningless here. Exits when the last sender is gone, i.e. never
|
|
/// during a run (the store holds one); the run's teardown stops it like any
|
|
/// other parked actor. Spawned under `ROOT_PID` so its exit signal is absorbed
|
|
/// rather than delivered to whichever supervisor's child happened to join
|
|
/// first.
|
|
///
|
|
/// Under `cluster` the same actor is the node's **pg actor** (RFC 010 Phase
|
|
/// 5, c15): it also drains a control inbox of local join/leave announcements,
|
|
/// the membership stream and the exposed `"pg"` inbox — see
|
|
/// [`crate::cluster::pg`]. Its store-side sweep is unchanged.
|
|
#[cfg(not(feature = "cluster"))]
|
|
fn reaper(rx: crate::channel::Receiver<Down>, ctl: crate::channel::Receiver<PgEvent>) {
|
|
// No mesh: nothing to tell about joins/leaves. Drop the control inbox
|
|
// so announcements are refused at the sender rather than queued.
|
|
drop(ctl);
|
|
while let Ok(down) = rx.recv() {
|
|
sweep_local_death(down.pid);
|
|
// Evicted memberships hold only ids; their monitors have fired.
|
|
}
|
|
}
|
|
|
|
/// What the local API tells the reaper besides deaths (which arrive as
|
|
/// [`Down`] on their own inbox — that channel's type is fixed by the monitor
|
|
/// primitive, so the two cannot be one enum). The default reaper has no use
|
|
/// for these; the cluster's pg actor broadcasts them (RFC 010 Phase 5).
|
|
// The default reaper never looks inside — that is the point, not a bug.
|
|
#[cfg_attr(not(feature = "cluster"), allow(dead_code))]
|
|
pub(crate) enum PgEvent {
|
|
/// `join` inserted `pid` into `group`. The consumer re-checks the store
|
|
/// before acting on it.
|
|
Joined { group: String, pid: Pid },
|
|
/// `leave` removed `pid` from `group`.
|
|
Left { group: String, pid: Pid },
|
|
/// `cluster::start` has the manager up and the local identity set: take
|
|
/// a membership subscription, register + expose the `"pg"` name, and
|
|
/// start speaking to peers.
|
|
#[cfg(feature = "cluster")]
|
|
Attach,
|
|
}
|
|
|
|
/// Evict the local member `pid` from every group. The reaper's one store
|
|
/// operation; returns what was evicted with its group (the cluster's
|
|
/// `Leave` broadcast wants both).
|
|
pub(crate) fn sweep_local_death(pid: Pid) -> Vec<(String, Membership)> {
|
|
with_runtime(|inner| {
|
|
let node = inner.node_id;
|
|
inner
|
|
.process_groups
|
|
.lock()
|
|
.remove_where(|m| m.node == node && m.pid == pid)
|
|
})
|
|
}
|
|
|
|
/// The reaper's inboxes. `deaths` is the liveness authority for the set
|
|
/// (`ctl` is created and dropped with it, on the same actor).
|
|
#[derive(Clone)]
|
|
pub(crate) struct ReaperInboxes {
|
|
pub(crate) deaths: Sender<Down>,
|
|
/// The control inbox: local `join`/`leave` announce here (see
|
|
/// [`PgEvent`]). The default reaper closes it on entry.
|
|
pub(crate) ctl: Sender<PgEvent>,
|
|
}
|
|
|
|
impl ReaperInboxes {
|
|
fn alive(&self) -> bool {
|
|
self.deaths.receiver_alive()
|
|
}
|
|
}
|
|
|
|
/// Live senders for the reaper's inboxes, spawning the reaper if this run has
|
|
/// none yet. Two racing first-spawns may both spawn; the loser's senders drop
|
|
/// on return, its spare reaper sees a closed inbox and exits.
|
|
pub(crate) fn reaper_inboxes() -> ReaperInboxes {
|
|
let existing = with_runtime(|inner| {
|
|
let pg = inner.process_groups.lock();
|
|
pg.reaper.clone().filter(ReaperInboxes::alive)
|
|
});
|
|
if let Some(r) = existing {
|
|
return r;
|
|
}
|
|
let (tx, rx) = channel::<Down>();
|
|
let (ctl_tx, ctl_rx) = channel::<PgEvent>();
|
|
// Detached: the handle drops here. The reaper's lifetime is the run's.
|
|
// The ONE seam between the local store and the cluster: same inboxes,
|
|
// different body.
|
|
#[cfg(not(feature = "cluster"))]
|
|
let _ = spawn_under(crate::runtime::ROOT_PID, move || reaper(rx, ctl_rx));
|
|
#[cfg(feature = "cluster")]
|
|
let _ = spawn_under(crate::runtime::ROOT_PID, move || {
|
|
crate::cluster::pg::actor(rx, ctl_rx)
|
|
});
|
|
let fresh = ReaperInboxes {
|
|
deaths: tx,
|
|
ctl: ctl_tx,
|
|
};
|
|
with_runtime(|inner| {
|
|
let mut pg = inner.process_groups.lock();
|
|
match &pg.reaper {
|
|
Some(r) if r.alive() => r.clone(),
|
|
_ => {
|
|
pg.reaper = Some(fresh.clone());
|
|
fresh
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
/// A live sender for the reaper's `deaths` inbox (spawning it if needed).
|
|
fn deaths_sender() -> Sender<Down> {
|
|
reaper_inboxes().deaths
|
|
}
|
|
|
|
/// Announce a local group change to the reaper, if this run has one. A
|
|
/// closed inbox is the default reaper (uninterested) or a run tearing down.
|
|
fn announce(msg: PgEvent) {
|
|
let ctl = with_runtime(|inner| {
|
|
inner
|
|
.process_groups
|
|
.lock()
|
|
.reaper
|
|
.as_ref()
|
|
.map(|r| r.ctl.clone())
|
|
});
|
|
if let Some(ctl) = ctl {
|
|
let _ = ctl.send(msg);
|
|
}
|
|
}
|
|
|
|
/// Build the full member identity for `pid` from runtime identity.
|
|
pub(crate) fn member_for(inner: &crate::runtime::RuntimeInner, pid: Pid) -> Member {
|
|
Member {
|
|
node: inner.node_id,
|
|
incarnation: inner.incarnation,
|
|
pid,
|
|
}
|
|
}
|
|
|
|
/// Is `pid` a live actor right now? Generation-checked atomic slot-word read,
|
|
/// no lock — identical to the registry's guard. The read-path backstop: a
|
|
/// generation is never reused, so a dead member is detectable independently of
|
|
/// whether its monitor `Down` has been drained yet.
|
|
pub(crate) fn live(inner: &crate::runtime::RuntimeInner, pid: Pid) -> bool {
|
|
inner.slot_at(pid).is_some_and(|s| s.is_live_for(pid))
|
|
}
|
|
|
|
/// Add `pid` to `group`. The same pid may join many groups; within one group a
|
|
/// pid is a member at most once (idempotent). Returns `true` if this call newly
|
|
/// added the membership, `false` if it was already a member.
|
|
///
|
|
/// Installs a monitor on `pid` so the actor's death evicts it from the group
|
|
/// automatically — you never have to remove a dead member yourself. Joining a
|
|
/// pid that is already dead is accepted and evicted the same way (via a
|
|
/// `NoProc` notice), so it never shows up in a read.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn join<A>(group: impl Into<String>, pid: Pid<A>) -> bool {
|
|
let group = group.into();
|
|
let pid = pid.erase();
|
|
let deaths = deaths_sender();
|
|
// Record the membership BEFORE arming its monitor: the reaper sweeps by
|
|
// pid on the first `Down`, so a `Down` that could precede the entry would
|
|
// leave a corpse in storage forever (visible to no read — the backstop
|
|
// hides it — but a leak, and once groups are clustered a member that
|
|
// would be announced). Arming after insertion means every `Down` finds
|
|
// its entry. The monitor id is allocated up front so `leave` can tear the
|
|
// registration down even if it lands in the tiny window before arming (an
|
|
// orphaned registration is harmless: its `Down` names a pid whose
|
|
// membership is gone, and the sweep finds nothing).
|
|
let id = with_runtime(|inner| inner.alloc_monitor_id());
|
|
let inserted = with_runtime(|inner| {
|
|
let ms = Membership {
|
|
member: member_for(inner, pid),
|
|
monitor: Some(id),
|
|
};
|
|
inner.process_groups.lock().join(&group, ms)
|
|
});
|
|
if !inserted {
|
|
return false;
|
|
}
|
|
// Tell the reaper (the cluster's pg actor re-checks the store before it
|
|
// broadcasts, so a `leave`/death that overtakes this announcement is
|
|
// never advertised as a join).
|
|
announce(PgEvent::Joined {
|
|
group: group.clone(),
|
|
pid,
|
|
});
|
|
|
|
// Outside the group lock: registration takes the target's cold lock (Leaf).
|
|
// The registration races `finalize_actor` under that cold lock exactly as
|
|
// every other monitor does, so no death can slip between the join and the
|
|
// monitor being in place.
|
|
if !register_monitor(pid, id, &deaths) {
|
|
// Already gone: queue the notice ourselves, exactly as `monitor` does.
|
|
let _ = deaths.send(Down {
|
|
pid,
|
|
reason: DownReason::NoProc,
|
|
});
|
|
}
|
|
true
|
|
}
|
|
|
|
/// Drop `pid`'s membership of `group`. Returns whether a membership was
|
|
/// removed. The membership's monitor registration is torn down.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn leave<A>(group: &str, pid: Pid<A>) -> bool {
|
|
let pid = pid.erase();
|
|
let removed = with_runtime(|inner| {
|
|
let member = member_for(inner, pid);
|
|
inner.process_groups.lock().leave(group, member)
|
|
});
|
|
match removed {
|
|
Some(ms) => {
|
|
if let Some(id) = ms.monitor {
|
|
unregister_monitor(pid, id);
|
|
}
|
|
announce(PgEvent::Left {
|
|
group: group.to_owned(),
|
|
pid,
|
|
});
|
|
true
|
|
}
|
|
None => false,
|
|
}
|
|
}
|
|
|
|
/// Every live member of `group`, in the order they joined. Returns an empty
|
|
/// vector if the group does not exist or has no live members.
|
|
///
|
|
/// Dead members are never returned: the reaper evicts a member as soon as its
|
|
/// death is processed, and as a backstop a member whose slot is already dead
|
|
/// is dropped from the result even in the brief window before the reaper's
|
|
/// turn.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn members(group: &str) -> Vec<Pid> {
|
|
with_runtime(|inner| {
|
|
inner
|
|
.process_groups
|
|
.lock()
|
|
.members_where(group, inner.node_id, |pid| live(inner, pid))
|
|
})
|
|
}
|
|
|
|
/// One live member of `group`, or `None` if the group is empty (or every
|
|
/// member has died). Selection is a stateless first-live scan in join order,
|
|
/// with the same dead-member backstop as [`members`]; smarter, load-aware
|
|
/// routing is a later, clustered concern.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn pick(group: &str) -> Option<Pid> {
|
|
with_runtime(|inner| {
|
|
inner
|
|
.process_groups
|
|
.lock()
|
|
.first_member_where(group, inner.node_id, |pid| live(inner, pid))
|
|
})
|
|
}
|
|
|
|
/// Typed [`pick`]: one live member of `group` as a [`Pid<A>`](Pid).
|
|
/// For a homogeneous pool every member is an `A`, so the picked member comes
|
|
/// back typed and dispatch is an ordinary compile-checked [`send_to`] rather
|
|
/// than the [`send_dyn`](crate::send_dyn) escape hatch. Re-types via the
|
|
/// unchecked `assert_type` primitive — a wrong `A` degrades to
|
|
/// [`SendError::NoChannel`] on the next send, never a misdelivery.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn pick_as<A: Addressable>(group: &str) -> Option<Pid<A>> {
|
|
pick(group).map(assert_type::<A>)
|
|
}
|
|
|
|
/// Typed `members`: every live member of `group` as a [`Pid<A>`], same
|
|
/// unchecked re-type as [`pick_as`]. Fan-out stays compile-checked end to end.
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn members_as<A: Addressable>(group: &str) -> Vec<Pid<A>> {
|
|
members(group).into_iter().map(assert_type::<A>).collect()
|
|
}
|
|
|
|
/// Pick a live member of `group` and send it `msg` in one step, returning the
|
|
/// member it reached on success. The pick-a-live-member-and-send combinator
|
|
/// over [`pick_as`] + [`send_to`].
|
|
///
|
|
/// Errors hand `msg` back undelivered: [`SendError::NoMember`] if the pool is
|
|
/// empty (or all-dead), otherwise whatever the underlying [`send_to`] returns
|
|
/// (e.g. the picked member died in the window between pick and send →
|
|
/// [`SendError::Dead`]).
|
|
///
|
|
/// Panics if called outside `Runtime::run()`.
|
|
pub fn dispatch<A: Addressable>(group: &str, msg: A::Msg) -> Result<Pid<A>, SendError<A::Msg>> {
|
|
match pick_as::<A>(group) {
|
|
Some(pid) => send_to::<A>(pid, msg).map(|()| pid),
|
|
None => Err(SendError::NoMember(msg)),
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::scheduler::spawn;
|
|
use std::time::{Duration, Instant};
|
|
|
|
fn member(index: u32, generation: u32) -> Member {
|
|
Member {
|
|
node: DEFAULT_NODE_ID,
|
|
incarnation: DEFAULT_INCARNATION,
|
|
pid: Pid::new(index, generation),
|
|
}
|
|
}
|
|
|
|
/// A synthetic membership: the store never looks at the id.
|
|
fn synth(index: u32, generation: u32) -> Membership {
|
|
Membership {
|
|
member: member(index, generation),
|
|
monitor: Some(MonitorId(0)),
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn join_is_idempotent_within_a_group() {
|
|
let mut pg = ProcessGroups::new();
|
|
assert!(pg.join("workers", synth(1, 0)), "first join inserts");
|
|
assert!(
|
|
!pg.join("workers", synth(1, 0)),
|
|
"second identical join is refused"
|
|
);
|
|
assert_eq!(pg.members_of("workers"), vec![member(1, 0)]);
|
|
}
|
|
|
|
#[test]
|
|
fn same_pid_in_many_groups_is_independent() {
|
|
let mut pg = ProcessGroups::new();
|
|
pg.join("a", synth(1, 0));
|
|
pg.join("b", synth(1, 0));
|
|
pg.join("b", synth(2, 0));
|
|
assert_eq!(pg.members_of("a"), vec![member(1, 0)]);
|
|
assert_eq!(pg.members_of("b"), vec![member(1, 0), member(2, 0)]);
|
|
}
|
|
|
|
#[test]
|
|
fn distinct_generations_are_distinct_members() {
|
|
// ABA guard: same slot index, different generation = different actor.
|
|
let mut pg = ProcessGroups::new();
|
|
assert!(pg.join("g", synth(1, 0)));
|
|
assert!(
|
|
pg.join("g", synth(1, 1)),
|
|
"different generation is a distinct member"
|
|
);
|
|
assert_eq!(pg.members_of("g"), vec![member(1, 0), member(1, 1)]);
|
|
}
|
|
|
|
#[test]
|
|
fn leave_removes_one_membership_and_prunes_empty_groups() {
|
|
let mut pg = ProcessGroups::new();
|
|
pg.join("g", synth(1, 0));
|
|
pg.join("g", synth(2, 0));
|
|
assert!(pg.leave("g", member(1, 0)).is_some());
|
|
assert_eq!(pg.members_of("g"), vec![member(2, 0)]);
|
|
assert!(
|
|
pg.leave("g", member(1, 0)).is_none(),
|
|
"second leave finds nothing"
|
|
);
|
|
assert!(pg.leave("g", member(2, 0)).is_some());
|
|
assert!(pg.members_of("g").is_empty(), "group is now empty");
|
|
assert!(
|
|
pg.leave("never", member(9, 0)).is_none(),
|
|
"leaving an unknown group is a no-op"
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn remove_where_sweeps_every_group() {
|
|
let mut pg = ProcessGroups::new();
|
|
for (g, m) in [
|
|
("a", synth(1, 0)),
|
|
("a", synth(2, 0)),
|
|
("b", synth(1, 0)),
|
|
("c", synth(3, 0)),
|
|
] {
|
|
pg.join(g, m);
|
|
}
|
|
// Death of pid index 1 (any generation) evicts it everywhere.
|
|
let evicted = pg.remove_where(|mem| mem.pid.index() == 1);
|
|
assert_eq!(evicted.len(), 2, "pid 1 was in a and b");
|
|
assert_eq!(pg.members_of("a"), vec![member(2, 0)]);
|
|
assert!(pg.members_of("b").is_empty(), "b held only pid 1; pruned");
|
|
assert_eq!(pg.members_of("c"), vec![member(3, 0)]);
|
|
}
|
|
|
|
#[test]
|
|
fn remove_where_can_match_an_incarnation_sweep() {
|
|
// Shape check for the node-down / incarnation sweep caller.
|
|
let mut pg = ProcessGroups::new();
|
|
let stale = Membership {
|
|
member: Member {
|
|
node: DEFAULT_NODE_ID,
|
|
incarnation: Incarnation::new(7),
|
|
pid: Pid::new(1, 0),
|
|
},
|
|
monitor: Some(MonitorId(0)),
|
|
};
|
|
pg.join("g", stale);
|
|
pg.join("g", synth(2, 0));
|
|
let evicted = pg.remove_where(|mem| mem.incarnation == Incarnation::new(7));
|
|
assert_eq!(evicted.len(), 1);
|
|
assert_eq!(pg.members_of("g"), vec![member(2, 0)]);
|
|
}
|
|
|
|
#[test]
|
|
fn read_backstop_hides_a_member_the_reaper_has_not_yet_swept() {
|
|
let mut pg = ProcessGroups::new();
|
|
pg.join("g", synth(1, 0));
|
|
pg.join("g", synth(2, 0));
|
|
|
|
// The slot-word oracle already reports pid 1 dead (finalize window),
|
|
// ahead of the reaper's turn.
|
|
let dead = Pid::new(1, 0);
|
|
let oracle = |pid: Pid| pid != dead;
|
|
|
|
assert_eq!(
|
|
pg.members_where("g", DEFAULT_NODE_ID, oracle),
|
|
vec![Pid::new(2, 0)],
|
|
"dead pid filtered from read"
|
|
);
|
|
assert_eq!(
|
|
pg.first_member_where("g", DEFAULT_NODE_ID, oracle),
|
|
Some(Pid::new(2, 0)),
|
|
"pick skips the dead first member"
|
|
);
|
|
|
|
// Backstop does not evict — that stays the reaper's job; raw storage
|
|
// still holds both until it runs.
|
|
assert_eq!(pg.members_of("g"), vec![member(1, 0), member(2, 0)]);
|
|
}
|
|
|
|
// ---- reaper: eager eviction against a live runtime ----
|
|
|
|
/// Raw storage view for a group, bypassing the read-path backstop.
|
|
fn stored(group: &str) -> Vec<Member> {
|
|
with_runtime(|inner| inner.process_groups.lock().members_of(group))
|
|
}
|
|
|
|
/// Cooperative wait (`smarm::sleep`, never an OS block) until `pred`.
|
|
fn wait_until(what: &str, mut pred: impl FnMut() -> bool) {
|
|
let deadline = Instant::now() + Duration::from_secs(2);
|
|
while !pred() {
|
|
assert!(Instant::now() < deadline, "timed out waiting for: {what}");
|
|
crate::sleep(Duration::from_millis(1));
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn a_death_is_swept_from_storage_without_any_group_operation() {
|
|
crate::run(|| {
|
|
let (tx, rx) = channel::<()>();
|
|
let w = spawn(move || {
|
|
rx.recv().unwrap();
|
|
});
|
|
let pid = w.pid();
|
|
join("a", pid);
|
|
join("b", pid);
|
|
assert_eq!(stored("a"), vec![member_for_test(pid)]);
|
|
|
|
tx.send(()).unwrap();
|
|
w.join().unwrap();
|
|
// No members()/pick()/join() on a or b from here on: the reaper
|
|
// alone must clear both.
|
|
wait_until("reaper sweeps a and b", || {
|
|
stored("a").is_empty() && stored("b").is_empty()
|
|
});
|
|
});
|
|
}
|
|
|
|
#[test]
|
|
fn a_dead_at_join_pid_is_swept_from_storage() {
|
|
crate::run(|| {
|
|
let h = spawn(|| {});
|
|
let pid = h.pid();
|
|
h.join().unwrap();
|
|
assert!(join("late", pid), "join is accepted; eviction is uniform");
|
|
wait_until("reaper sweeps the NoProc member", || {
|
|
stored("late").is_empty()
|
|
});
|
|
});
|
|
}
|
|
|
|
#[test]
|
|
fn leave_then_death_does_not_disturb_a_rejoined_group() {
|
|
// A monitor unregistered by `leave` must not fire later; the pid's
|
|
// fresh membership after re-join is swept exactly once, by its own
|
|
// monitor, on death.
|
|
crate::run(|| {
|
|
let (tx, rx) = channel::<()>();
|
|
let w = spawn(move || {
|
|
rx.recv().unwrap();
|
|
});
|
|
let pid = w.pid();
|
|
join("g", pid);
|
|
assert!(leave("g", pid));
|
|
assert!(join("g", pid));
|
|
assert_eq!(members("g"), vec![pid]);
|
|
tx.send(()).unwrap();
|
|
w.join().unwrap();
|
|
wait_until("reaper sweeps g", || stored("g").is_empty());
|
|
});
|
|
}
|
|
|
|
#[test]
|
|
fn reaper_is_respawned_for_a_second_run_of_the_same_runtime() {
|
|
let rt = crate::runtime::init(crate::runtime::Config::exact(1));
|
|
let body = || {
|
|
let h = spawn(|| {});
|
|
let pid = h.pid();
|
|
h.join().unwrap();
|
|
join("g", pid);
|
|
wait_until("reaper sweeps g", || stored("g").is_empty());
|
|
};
|
|
rt.run(body);
|
|
rt.run(body);
|
|
}
|
|
|
|
fn member_for_test(pid: Pid) -> Member {
|
|
with_runtime(|inner| member_for(inner, pid))
|
|
}
|
|
}
|