feat(cluster): RFC 010 c10–c16, follow-ups and Phase 6 (squash of 16ef583..d9c62a8)
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.
This commit is contained in:
@@ -99,10 +99,13 @@
|
||||
//! ## Identity and clustering
|
||||
//!
|
||||
//! A group member is described by a [`Member`] — a [`Pid`] plus a [`NodeId`] and
|
||||
//! an [`Incarnation`]. Today everything is single-node, those two fields are
|
||||
//! fixed defaults, and you only ever pass and receive a plain [`Pid`]: the extra
|
||||
//! identity is carried so this API will not have to change when groups learn to
|
||||
//! span a cluster.
|
||||
//! 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
|
||||
//!
|
||||
@@ -110,10 +113,11 @@
|
||||
//! from inside [`run`](crate::run) (that is, on an actor thread). Calling one
|
||||
//! from outside a running runtime panics.
|
||||
|
||||
use crate::monitor::{demonitor, monitor, Monitor};
|
||||
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::with_runtime;
|
||||
use crate::scheduler::{spawn_under, with_runtime};
|
||||
use std::collections::HashMap;
|
||||
|
||||
/// A cluster node handle. A `u32` integer handle, *not* an interned atom — the
|
||||
@@ -186,13 +190,15 @@ pub struct Member {
|
||||
pub pid: Pid,
|
||||
}
|
||||
|
||||
/// One membership: a [`Member`] and the [`Monitor`] that watches its liveness.
|
||||
/// The monitor lives *alongside* the group entry so a group is
|
||||
/// self-contained: draining the membership tells us whether the member is
|
||||
/// still alive, and dropping the membership drops its monitor.
|
||||
struct Membership {
|
||||
member: Member,
|
||||
monitor: Monitor,
|
||||
/// 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`
|
||||
@@ -202,46 +208,60 @@ struct Membership {
|
||||
///
|
||||
/// 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). The two operations that
|
||||
/// do need another lock are kept off the group-lock path:
|
||||
/// 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.
|
||||
///
|
||||
/// - `monitor()` / `demonitor()` take the target's cold lock (also Leaf), so
|
||||
/// they run *before* / *after* the group lock, never under it.
|
||||
/// - draining a monitor with `try_recv` takes the channel's Channel-class
|
||||
/// lock, which the lock order permits *under* a Leaf; a channel critical
|
||||
/// section only does the lock-free unpark protocol, so no Leaf ever nests
|
||||
/// under it.
|
||||
///
|
||||
/// Evicted and rejected [`Monitor`]s are therefore dropped only *after* the
|
||||
/// group lock is released, so a receiver-drop never runs a wakeup under the
|
||||
/// lock — the same discipline as `demonitor`.
|
||||
/// 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(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Insert `ms` into `group`. Idempotent on the *member*: if the member is
|
||||
/// already present the new membership is handed back (`Some`) so the caller
|
||||
/// can tear its now-redundant monitor down outside the lock; `None` means
|
||||
/// it was inserted.
|
||||
fn join(&mut self, group: &str, ms: Membership) -> Option<Membership> {
|
||||
/// 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 Some(ms);
|
||||
return false;
|
||||
}
|
||||
v.push(ms);
|
||||
None
|
||||
true
|
||||
}
|
||||
|
||||
/// Remove `member`'s membership from `group`, returning it (so the caller
|
||||
/// can `demonitor` it outside the lock). An emptied group is pruned.
|
||||
fn leave(&mut self, group: &str, member: Member) -> Option<Membership> {
|
||||
/// 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);
|
||||
@@ -252,20 +272,23 @@ impl ProcessGroups {
|
||||
}
|
||||
|
||||
/// The one dumb eviction primitive: drop every member matching `pred` from
|
||||
/// every group, pruning emptied groups, and return the evicted memberships'
|
||||
/// monitors for the caller to drop outside the lock. The primitive does not
|
||||
/// know *why* a member leaves; that is the caller's concern. Its callers are
|
||||
/// the death hook (`reap_group`) and, once clustering lands, an
|
||||
/// incarnation-eviction sweep — both 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).
|
||||
fn remove_where(&mut self, mut pred: impl FnMut(&Member) -> bool) -> Vec<Monitor> {
|
||||
/// 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(|_, v| {
|
||||
self.groups.retain(|g, v| {
|
||||
let mut i = 0;
|
||||
while i < v.len() {
|
||||
if pred(&v[i].member) {
|
||||
evicted.push(v.remove(i).monitor);
|
||||
evicted.push((g.clone(), v.remove(i)));
|
||||
} else {
|
||||
i += 1;
|
||||
}
|
||||
@@ -275,38 +298,6 @@ impl ProcessGroups {
|
||||
evicted
|
||||
}
|
||||
|
||||
/// Drain-on-contact death hook. The registry can prune a stale binding
|
||||
/// lazily, on contact, because it only ever resolves one binding at a time;
|
||||
/// a group is *iterated* — `members` fans out to everyone — so it must not
|
||||
/// carry a dead member across a broadcast. Every group operation reaps the
|
||||
/// group it touches first.
|
||||
///
|
||||
/// Drains every membership monitor in `group` with a non-blocking
|
||||
/// `try_recv`: a delivered `Down` (any reason) or a closed channel means
|
||||
/// that member is dead. On the first death detected, sweep *all* of the
|
||||
/// dead pids out of *every* group via [`remove_where`] — a death is removed
|
||||
/// from each group it joined, not just the one being touched. Returns the
|
||||
/// evicted monitors to drop outside the lock.
|
||||
fn reap_group(&mut self, group: &str) -> Vec<Monitor> {
|
||||
let dead: Vec<Pid> = {
|
||||
let Some(v) = self.groups.get(group) else {
|
||||
return Vec::new();
|
||||
};
|
||||
v.iter()
|
||||
.filter_map(|e| match e.monitor.rx.try_recv() {
|
||||
// A Down arrived, or the channel closed and drained: dead.
|
||||
Ok(Some(_)) | Err(_) => Some(e.member.pid),
|
||||
// Empty but open — the sender still lives in the slot: alive.
|
||||
Ok(None) => None,
|
||||
})
|
||||
.collect()
|
||||
};
|
||||
if dead.is_empty() {
|
||||
return Vec::new();
|
||||
}
|
||||
self.remove_where(|m| dead.contains(&m.pid))
|
||||
}
|
||||
|
||||
/// 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)]
|
||||
@@ -317,16 +308,24 @@ impl ProcessGroups {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// Live members of `group`, 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 its `Down` has not been drained yet.
|
||||
/// Backstop only — the entry stays in storage; eviction is the monitor's
|
||||
/// job (`reap_group`).
|
||||
fn members_where(&self, group: &str, mut is_live: impl FnMut(Pid) -> bool) -> Vec<Pid> {
|
||||
/// 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()
|
||||
@@ -334,19 +333,210 @@ impl ProcessGroups {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// The first live member of `group` in insertion order — stateless
|
||||
/// first-live `pick`, with the same read-path backstop as `members_where`.
|
||||
fn first_member_where(&self, group: &str, mut is_live: impl FnMut(Pid) -> bool) -> Option<Pid> {
|
||||
/// 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.
|
||||
fn member_for(inner: &crate::runtime::RuntimeInner, pid: Pid) -> Member {
|
||||
pub(crate) fn member_for(inner: &crate::runtime::RuntimeInner, pid: Pid) -> Member {
|
||||
Member {
|
||||
node: inner.node_id,
|
||||
incarnation: inner.incarnation,
|
||||
@@ -358,7 +548,7 @@ fn member_for(inner: &crate::runtime::RuntimeInner, pid: Pid) -> Member {
|
||||
/// 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.
|
||||
fn live(inner: &crate::runtime::RuntimeInner, pid: Pid) -> bool {
|
||||
pub(crate) fn live(inner: &crate::runtime::RuntimeInner, pid: Pid) -> bool {
|
||||
inner.slot_at(pid).is_some_and(|s| s.is_live_for(pid))
|
||||
}
|
||||
|
||||
@@ -367,61 +557,76 @@ fn live(inner: &crate::runtime::RuntimeInner, pid: Pid) -> bool {
|
||||
/// 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. A redundant
|
||||
/// (idempotent) join tears its extra monitor back down.
|
||||
/// 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();
|
||||
// Install the monitor BEFORE taking the group lock: monitor() acquires the
|
||||
// target's cold lock (Leaf), and two Leaf locks are never held at once. 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.
|
||||
let mon = monitor(pid);
|
||||
|
||||
let (rejected, reaped) = with_runtime(|inner| {
|
||||
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: mon,
|
||||
monitor: Some(id),
|
||||
};
|
||||
let mut pg = inner.process_groups.lock();
|
||||
let reaped = pg.reap_group(&group);
|
||||
let rejected = pg.join(&group, ms);
|
||||
(rejected, reaped)
|
||||
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: drop the reaped (dead) monitors, and if this join
|
||||
// was redundant, demonitor + drop the extra monitor we just installed.
|
||||
drop(reaped);
|
||||
match rejected {
|
||||
Some(dup) => {
|
||||
demonitor(&dup.monitor);
|
||||
false
|
||||
}
|
||||
None => true,
|
||||
// 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 is demonitored and dropped.
|
||||
/// 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, reaped) = with_runtime(|inner| {
|
||||
let removed = with_runtime(|inner| {
|
||||
let member = member_for(inner, pid);
|
||||
let mut pg = inner.process_groups.lock();
|
||||
let reaped = pg.reap_group(group);
|
||||
let removed = pg.leave(group, member);
|
||||
(removed, reaped)
|
||||
inner.process_groups.lock().leave(group, member)
|
||||
});
|
||||
|
||||
drop(reaped);
|
||||
match removed {
|
||||
Some(ms) => {
|
||||
demonitor(&ms.monitor);
|
||||
if let Some(id) = ms.monitor {
|
||||
unregister_monitor(pid, id);
|
||||
}
|
||||
announce(PgEvent::Left {
|
||||
group: group.to_owned(),
|
||||
pid,
|
||||
});
|
||||
true
|
||||
}
|
||||
None => false,
|
||||
@@ -431,21 +636,19 @@ pub fn leave<A>(group: &str, pid: Pid<A>) -> bool {
|
||||
/// 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 group is pruned of anything that has
|
||||
/// died before the read, and as a backstop a member whose slot is already dead
|
||||
/// is dropped from the result even in the brief window before its death has
|
||||
/// been fully processed.
|
||||
/// 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> {
|
||||
let (pids, reaped) = with_runtime(|inner| {
|
||||
let mut pg = inner.process_groups.lock();
|
||||
let reaped = pg.reap_group(group);
|
||||
let pids = pg.members_where(group, |pid| live(inner, pid));
|
||||
(pids, reaped)
|
||||
});
|
||||
drop(reaped);
|
||||
pids
|
||||
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
|
||||
@@ -455,14 +658,12 @@ pub fn members(group: &str) -> Vec<Pid> {
|
||||
///
|
||||
/// Panics if called outside `Runtime::run()`.
|
||||
pub fn pick(group: &str) -> Option<Pid> {
|
||||
let (picked, reaped) = with_runtime(|inner| {
|
||||
let mut pg = inner.process_groups.lock();
|
||||
let reaped = pg.reap_group(group);
|
||||
let picked = pg.first_member_where(group, |pid| live(inner, pid));
|
||||
(picked, reaped)
|
||||
});
|
||||
drop(reaped);
|
||||
picked
|
||||
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).
|
||||
@@ -505,8 +706,8 @@ pub fn dispatch<A: Addressable>(group: &str, msg: A::Msg) -> Result<Pid<A>, Send
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::channel::{channel, Sender};
|
||||
use crate::monitor::{Down, DownReason, MonitorId};
|
||||
use crate::scheduler::spawn;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
fn member(index: u32, generation: u32) -> Member {
|
||||
Member {
|
||||
@@ -516,33 +717,21 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
/// A synthetic membership with a real (but slot-less) monitor channel. The
|
||||
/// returned `Sender` stands in for the slot's `Down` sender: hold it to
|
||||
/// keep the member "alive" (`try_recv` → `Ok(None)`), `send` a `Down` to
|
||||
/// simulate death, or `drop` it to simulate a drained/closed channel.
|
||||
fn synth(index: u32, generation: u32) -> (Membership, Sender<Down>) {
|
||||
let pid = Pid::new(index, generation);
|
||||
let (tx, rx) = channel::<Down>();
|
||||
let ms = Membership {
|
||||
/// A synthetic membership: the store never looks at the id.
|
||||
fn synth(index: u32, generation: u32) -> Membership {
|
||||
Membership {
|
||||
member: member(index, generation),
|
||||
monitor: Monitor {
|
||||
id: MonitorId(0),
|
||||
target: pid,
|
||||
rx,
|
||||
},
|
||||
};
|
||||
(ms, tx)
|
||||
monitor: Some(MonitorId(0)),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn join_is_idempotent_within_a_group() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, _ta) = synth(1, 0);
|
||||
let (b, _tb) = synth(1, 0);
|
||||
assert!(pg.join("workers", a).is_none(), "first join inserts");
|
||||
assert!(pg.join("workers", synth(1, 0)), "first join inserts");
|
||||
assert!(
|
||||
pg.join("workers", b).is_some(),
|
||||
"second identical join is handed back"
|
||||
!pg.join("workers", synth(1, 0)),
|
||||
"second identical join is refused"
|
||||
);
|
||||
assert_eq!(pg.members_of("workers"), vec![member(1, 0)]);
|
||||
}
|
||||
@@ -550,12 +739,9 @@ mod tests {
|
||||
#[test]
|
||||
fn same_pid_in_many_groups_is_independent() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, _ta) = synth(1, 0);
|
||||
let (b, _tb) = synth(1, 0);
|
||||
let (c, _tc) = synth(2, 0);
|
||||
pg.join("a", a);
|
||||
pg.join("b", b);
|
||||
pg.join("b", c);
|
||||
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)]);
|
||||
}
|
||||
@@ -564,11 +750,9 @@ mod tests {
|
||||
fn distinct_generations_are_distinct_members() {
|
||||
// ABA guard: same slot index, different generation = different actor.
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, _ta) = synth(1, 0);
|
||||
let (b, _tb) = synth(1, 1);
|
||||
assert!(pg.join("g", a).is_none());
|
||||
assert!(pg.join("g", synth(1, 0)));
|
||||
assert!(
|
||||
pg.join("g", b).is_none(),
|
||||
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)]);
|
||||
@@ -577,10 +761,8 @@ mod tests {
|
||||
#[test]
|
||||
fn leave_removes_one_membership_and_prunes_empty_groups() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, _ta) = synth(1, 0);
|
||||
let (b, _tb) = synth(2, 0);
|
||||
pg.join("g", a);
|
||||
pg.join("g", b);
|
||||
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!(
|
||||
@@ -598,7 +780,7 @@ mod tests {
|
||||
#[test]
|
||||
fn remove_where_sweeps_every_group() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
for (g, (m, _t)) in [
|
||||
for (g, m) in [
|
||||
("a", synth(1, 0)),
|
||||
("a", synth(2, 0)),
|
||||
("b", synth(1, 0)),
|
||||
@@ -616,104 +798,137 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn remove_where_can_match_an_incarnation_sweep() {
|
||||
// Shape check for the later evict_incarnation(node, inc) caller.
|
||||
// Shape check for the node-down / incarnation sweep caller.
|
||||
let mut pg = ProcessGroups::new();
|
||||
let pid = Pid::new(1, 0);
|
||||
let (tx, rx) = channel::<Down>();
|
||||
let dead = Membership {
|
||||
let stale = Membership {
|
||||
member: Member {
|
||||
node: DEFAULT_NODE_ID,
|
||||
incarnation: Incarnation::new(7),
|
||||
pid,
|
||||
},
|
||||
monitor: Monitor {
|
||||
id: MonitorId(0),
|
||||
target: pid,
|
||||
rx,
|
||||
pid: Pid::new(1, 0),
|
||||
},
|
||||
monitor: Some(MonitorId(0)),
|
||||
};
|
||||
let _keep = tx;
|
||||
let (live, _tl) = synth(2, 0);
|
||||
pg.join("g", dead);
|
||||
pg.join("g", live);
|
||||
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 reap_keeps_live_members() {
|
||||
fn read_backstop_hides_a_member_the_reaper_has_not_yet_swept() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, _ta) = synth(1, 0); // sender held: member stays alive
|
||||
pg.join("a", a);
|
||||
assert!(pg.reap_group("a").is_empty(), "no deaths");
|
||||
assert_eq!(pg.members_of("a"), vec![member(1, 0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reap_evicts_a_dead_member_and_sweeps_all_its_groups() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a1, ta1) = synth(1, 0); // pid 1 in group a
|
||||
let (a2, _ta2) = synth(2, 0); // pid 2 in group a (stays alive)
|
||||
let (b1, _tb1) = synth(1, 0); // pid 1 in group b
|
||||
pg.join("a", a1);
|
||||
pg.join("a", a2);
|
||||
pg.join("b", b1);
|
||||
// pid 1 dies: its group-a monitor receives a Down. Its group-b monitor
|
||||
// has not — reap must still sweep pid 1 out of b by the pid predicate.
|
||||
ta1.send(Down {
|
||||
pid: Pid::new(1, 0),
|
||||
reason: DownReason::Exit,
|
||||
})
|
||||
.unwrap();
|
||||
let evicted = pg.reap_group("a");
|
||||
assert_eq!(
|
||||
evicted.len(),
|
||||
2,
|
||||
"pid 1's memberships in both a and b are evicted"
|
||||
);
|
||||
assert_eq!(pg.members_of("a"), vec![member(2, 0)]);
|
||||
assert!(pg.members_of("b").is_empty(), "swept from b too; pruned");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reap_treats_a_closed_channel_as_dead() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
let (a, ta) = synth(1, 0);
|
||||
pg.join("a", a);
|
||||
drop(ta); // sender gone, queue empty → try_recv = Err(RecvError) = dead
|
||||
let evicted = pg.reap_group("a");
|
||||
assert_eq!(evicted.len(), 1);
|
||||
assert!(pg.members_of("a").is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_backstop_hides_a_member_the_monitor_has_not_yet_reaped() {
|
||||
let mut pg = ProcessGroups::new();
|
||||
// Both senders held: reap_group would see Ok(None) and evict neither.
|
||||
let (a, _ta) = synth(1, 0);
|
||||
let (b, _tb) = synth(2, 0);
|
||||
pg.join("g", a);
|
||||
pg.join("g", b);
|
||||
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 any Down delivery.
|
||||
// ahead of the reaper's turn.
|
||||
let dead = Pid::new(1, 0);
|
||||
let oracle = |pid: Pid| pid != dead;
|
||||
|
||||
assert_eq!(
|
||||
pg.members_where("g", oracle),
|
||||
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", oracle),
|
||||
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 monitor's job; raw storage
|
||||
// still holds both until reap runs.
|
||||
// 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))
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user