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.
360 lines
14 KiB
Rust
360 lines
14 KiB
Rust
//! RFC 010 c12 — remote monitors.
|
|
//!
|
|
//! Local suite (`run()`, no network): the immediate answers — no connection
|
|
//! ⇒ `Disconnected`, dead incarnation ⇒ `NoProc` — and the self-node
|
|
//! collapse (a plain local monitor underneath, incl. `demonitor_remote`).
|
|
//!
|
|
//! Cross-process: a *server* exposes a control name and spawns workers on
|
|
//! request, replying with each worker's pid (via `RemotePid::from_local`,
|
|
//! the D12 set-site) or, for the deliberately unshipped one, only its raw
|
|
//! slot numbers. The *client* monitors them and asserts: kill ⇒ the true
|
|
//! reason (Exit / Panic); a corpse ⇒ its recorded terminal reason, not
|
|
//! NoProc; a live pid that never crossed the wire ⇒ NoProc (no liveness
|
|
//! leak); a demonitor racing the kill ⇒ no notice, proven by stream ORDER
|
|
//! (a later notice on the same connection arrives while the earlier slot
|
|
//! is still empty), not by sleeping.
|
|
#![cfg(feature = "cluster")]
|
|
|
|
mod common;
|
|
|
|
use common::{maybe_child, spawn_node};
|
|
use smarm::cluster::envelope::NodeMeta;
|
|
use smarm::cluster::expose::{expose, expose_type};
|
|
use smarm::cluster::membership::{subscribe, NodeEvent};
|
|
use smarm::cluster::remote::{
|
|
self, demonitor_remote, monitor_remote, send_to_remote, RemoteName, RemotePid,
|
|
};
|
|
use smarm::cluster::{start, Config, RemoteDownReason, StaticSeeds, Timing};
|
|
use smarm::pg::Incarnation;
|
|
use smarm::{channel, install, register, run, spawn, Addressable, DownReason, Erased, Name, Pid};
|
|
use std::collections::HashMap;
|
|
use std::time::Duration;
|
|
|
|
// ---- message types (hand-rolled serde; the crate is derive-less) ---------
|
|
|
|
#[derive(Debug)]
|
|
struct Ctl {
|
|
cmd: String,
|
|
reply_to: RemotePid<Client>,
|
|
}
|
|
#[derive(Debug)]
|
|
struct Answer {
|
|
text: String,
|
|
pid: Option<RemotePid<Erased>>,
|
|
}
|
|
struct Client;
|
|
impl Addressable for Client {
|
|
type Msg = Answer;
|
|
}
|
|
|
|
impl serde::Serialize for Ctl {
|
|
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
|
|
use serde::ser::SerializeTuple;
|
|
let mut t = s.serialize_tuple(2)?;
|
|
t.serialize_element(&self.cmd)?;
|
|
t.serialize_element(&self.reply_to)?;
|
|
t.end()
|
|
}
|
|
}
|
|
impl<'de> serde::Deserialize<'de> for Ctl {
|
|
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
|
|
let (cmd, reply_to) = <(String, RemotePid<Client>)>::deserialize(d)?;
|
|
Ok(Ctl { cmd, reply_to })
|
|
}
|
|
}
|
|
impl serde::Serialize for Answer {
|
|
fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
|
|
use serde::ser::SerializeTuple;
|
|
let mut t = s.serialize_tuple(2)?;
|
|
t.serialize_element(&self.text)?;
|
|
t.serialize_element(&self.pid)?;
|
|
t.end()
|
|
}
|
|
}
|
|
impl<'de> serde::Deserialize<'de> for Answer {
|
|
fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
|
|
let (text, pid) = <(String, Option<RemotePid<Erased>>)>::deserialize(d)?;
|
|
Ok(Answer { text, pid })
|
|
}
|
|
}
|
|
|
|
// ================= local suite =========================================
|
|
|
|
/// No connection to the pid's node: `Disconnected` at once — the remote
|
|
/// analog of NoProc, and the first thing c11's variant is for.
|
|
#[test]
|
|
fn unconnected_node_is_disconnected_immediately() {
|
|
maybe_child(ROLES);
|
|
run(|| {
|
|
remote::set_local_identity("me", Incarnation::new(7));
|
|
let ghost = RemotePid::<Erased>::from_parts("nowhere", Incarnation::new(1), 3, 1);
|
|
let m = monitor_remote(ghost.clone());
|
|
let d = m.recv().unwrap();
|
|
assert_eq!(d.pid, ghost);
|
|
assert_eq!(d.reason, RemoteDownReason::Disconnected);
|
|
});
|
|
}
|
|
|
|
/// The node is connected but the pid names an earlier incarnation: the
|
|
/// actor is a known corpse (RFC v2 §3), so `NoProc` at once — never
|
|
/// `Disconnected`, nothing on the wire.
|
|
#[test]
|
|
fn dead_incarnation_is_noproc_immediately() {
|
|
maybe_child(ROLES);
|
|
run(|| {
|
|
remote::set_local_identity("me", Incarnation::new(7));
|
|
let (probe_tx, probe_rx) = channel();
|
|
remote::bind_outbound_probe("peer", Incarnation::new(5), probe_tx);
|
|
let stale = RemotePid::<Erased>::from_parts("peer", Incarnation::new(4), 9, 1);
|
|
let m = monitor_remote(stale);
|
|
assert_eq!(m.recv().unwrap().reason, DownReason::NoProc.into());
|
|
assert!(probe_rx.try_recv().unwrap().is_none(), "no frame emitted");
|
|
});
|
|
}
|
|
|
|
/// A self-node pid collapses to an ordinary local monitor: the true reason
|
|
/// on exit, and `demonitor_remote` cancels it.
|
|
#[test]
|
|
fn self_node_pid_collapses_to_local_monitor() {
|
|
maybe_child(ROLES);
|
|
run(|| {
|
|
remote::set_local_identity("me", Incarnation::new(7));
|
|
let (go_tx, go_rx) = channel::<()>();
|
|
let (go2_tx, go2_rx) = channel::<()>();
|
|
let a = spawn(move || {
|
|
let _ = go_rx.recv();
|
|
})
|
|
.pid();
|
|
let b = spawn(move || {
|
|
let _ = go2_rx.recv();
|
|
})
|
|
.pid();
|
|
let ma = monitor_remote(RemotePid::from_local(a).expect("identity set"));
|
|
let mb = monitor_remote(RemotePid::from_local(b).expect("identity set"));
|
|
assert_ne!(ma.id, mb.id);
|
|
assert!(ma.target.local() == Some(a));
|
|
|
|
demonitor_remote(&mb);
|
|
go2_tx.send(()).unwrap();
|
|
go_tx.send(()).unwrap();
|
|
let d = ma.recv().unwrap();
|
|
assert_eq!(d.reason, DownReason::Exit.into());
|
|
assert_eq!(d.pid.local(), Some(a));
|
|
// `a` is down (its notice arrived), and `b` was killed first on the
|
|
// same scheduler — a notice for `b` would be here by now. After a
|
|
// demonitor the channel is closed-empty (`Err`), like the local one.
|
|
assert!(matches!(mb.try_recv(), Ok(None) | Err(_)));
|
|
});
|
|
}
|
|
|
|
// ================= cross-process ======================================
|
|
|
|
const ROLES: &[(&str, fn())] = &[("server", role_server), ("client", role_client)];
|
|
|
|
const CTL: Name<Ctl> = Name::new("c12.ctl");
|
|
|
|
fn cfg(name: &str, seeds: Vec<(String, String)>) -> Config {
|
|
Config {
|
|
node_name: name.to_string(),
|
|
meta: NodeMeta {
|
|
role: "c12".into(),
|
|
region: "local".into(),
|
|
},
|
|
listen_addr: std::env::var("SMARM_LISTEN_ADDR").unwrap_or_else(|_| "127.0.0.1:0".into()),
|
|
strategy: Box::new(StaticSeeds::new(seeds)),
|
|
timing: Timing::default(),
|
|
}
|
|
}
|
|
|
|
fn wait_up(events: &smarm::cluster::membership::MembershipEvents, who: &str) {
|
|
loop {
|
|
match events.rx.recv() {
|
|
Ok(NodeEvent::NodeUp(i)) if i.name == who => return,
|
|
Ok(_) => continue,
|
|
Err(_) => panic!("manager gone"),
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Server commands (all answered to `reply_to`):
|
|
/// - `spawn:exit` / `spawn:panic` — a parked worker; `kill:<index>` releases
|
|
/// it, whereupon it returns / panics. Answer carries its pid.
|
|
/// - `spawn:corpse` — a worker that has already exited when the answer is
|
|
/// sent; the pid was shipped (watchable) before it died.
|
|
/// - `spawn:unwatched` — a parked worker whose pid is NEVER shipped; the
|
|
/// answer carries only `text = "slot:<index>:<generation>"`.
|
|
fn role_server() {
|
|
smarm::run(move || {
|
|
let cluster = start(cfg("server", vec![])).expect("binds");
|
|
println!("LISTENING {}", cluster.local_addr());
|
|
let (tx, rx) = channel::<Ctl>();
|
|
register(CTL, tx).unwrap();
|
|
expose(CTL);
|
|
println!("READY");
|
|
let mut workers: HashMap<u32, smarm::channel::Sender<()>> = HashMap::new();
|
|
loop {
|
|
let ctl = rx.recv().unwrap();
|
|
println!("CTL {}", ctl.cmd);
|
|
let (text, pid): (String, Option<RemotePid<Erased>>) = match ctl.cmd.as_str() {
|
|
"spawn:exit" | "spawn:panic" => {
|
|
let panic = ctl.cmd == "spawn:panic";
|
|
let (go_tx, go_rx) = channel::<()>();
|
|
let p: Pid = spawn(move || {
|
|
let _ = go_rx.recv();
|
|
if panic {
|
|
panic!("worker asked to panic");
|
|
}
|
|
})
|
|
.pid();
|
|
workers.insert(p.index(), go_tx);
|
|
(
|
|
"ok".into(),
|
|
Some(RemotePid::from_local(p).expect("identity set")),
|
|
)
|
|
}
|
|
"spawn:corpse" => {
|
|
let p: Pid = spawn(|| {}).pid();
|
|
let rp = RemotePid::from_local(p).expect("identity set"); // shipped ⇒ watchable
|
|
let m = smarm::monitor(p);
|
|
let _ = m.rx.recv(); // dead before the answer goes out
|
|
("ok".into(), Some(rp))
|
|
}
|
|
"spawn:unwatched" => {
|
|
let (go_tx, go_rx) = channel::<()>();
|
|
let p: Pid = spawn(move || {
|
|
let _ = go_rx.recv();
|
|
})
|
|
.pid();
|
|
workers.insert(p.index(), go_tx);
|
|
(format!("slot:{}:{}", p.index(), p.generation()), None)
|
|
}
|
|
other => {
|
|
let idx: u32 = other.strip_prefix("kill:").unwrap().parse().unwrap();
|
|
if let Some(go) = workers.remove(&idx) {
|
|
let _ = go.send(());
|
|
}
|
|
("killed".into(), None)
|
|
}
|
|
};
|
|
send_to_remote(ctl.reply_to, Answer { text, pid }).unwrap();
|
|
}
|
|
});
|
|
}
|
|
|
|
fn role_client() {
|
|
let server_addr = std::env::var("SMARM_SERVER_ADDR").expect("SMARM_SERVER_ADDR");
|
|
smarm::run(move || {
|
|
let _cluster = start(cfg("client", vec![("server".into(), server_addr)])).expect("binds");
|
|
let ev = subscribe().unwrap();
|
|
wait_up(&ev, "server");
|
|
let (tx, rx) = channel::<Answer>();
|
|
let me: Pid<Client> = install::<Client>(tx);
|
|
expose_type::<Answer>();
|
|
let ask = |cmd: &str| -> Answer {
|
|
remote::send(
|
|
RemoteName::new("server", CTL),
|
|
Ctl {
|
|
cmd: cmd.into(),
|
|
reply_to: RemotePid::from_local(me).expect("identity set"),
|
|
},
|
|
)
|
|
.unwrap();
|
|
rx.recv().unwrap()
|
|
};
|
|
let server_inc = ev_incarnation();
|
|
|
|
// 1. kill ⇒ true reason (Exit).
|
|
let a = ask("spawn:exit").pid.unwrap();
|
|
let ma = monitor_remote(a.clone());
|
|
ask(&format!("kill:{}", a.index()));
|
|
let d = ma.recv().unwrap();
|
|
assert_eq!(d.pid, a);
|
|
println!("DOWN exit {:?}", d.reason);
|
|
|
|
// 2. kill ⇒ true reason (Panic).
|
|
let b = ask("spawn:panic").pid.unwrap();
|
|
let mb = monitor_remote(b.clone());
|
|
ask(&format!("kill:{}", b.index()));
|
|
println!("DOWN panic {:?}", mb.recv().unwrap().reason);
|
|
|
|
// 3. corpse ⇒ recorded terminal reason, not NoProc.
|
|
let c = ask("spawn:corpse").pid.unwrap();
|
|
println!("DOWN corpse {:?}", monitor_remote(c).recv().unwrap().reason);
|
|
|
|
// 4. live but never shipped/exposed ⇒ NoProc (no leak); a made-up
|
|
// slot on the same node ⇒ NoProc too, indistinguishably.
|
|
let ans = ask("spawn:unwatched");
|
|
let mut it = ans.text.strip_prefix("slot:").unwrap().split(':');
|
|
let (idx, gen): (u32, u32) = (
|
|
it.next().unwrap().parse().unwrap(),
|
|
it.next().unwrap().parse().unwrap(),
|
|
);
|
|
let hidden = RemotePid::<Erased>::from_parts("server", server_inc, idx, gen);
|
|
println!(
|
|
"DOWN hidden {:?}",
|
|
monitor_remote(hidden).recv().unwrap().reason
|
|
);
|
|
let bogus = RemotePid::<Erased>::from_parts("server", server_inc, 100_000, 1);
|
|
println!(
|
|
"DOWN bogus {:?}",
|
|
monitor_remote(bogus).recv().unwrap().reason
|
|
);
|
|
|
|
// 5. demonitor races the kill: no notice for `d1`, proven by order —
|
|
// `d2`'s notice (same connection, later) arrives while `d1`'s
|
|
// slot is still empty.
|
|
let d1 = ask("spawn:exit").pid.unwrap();
|
|
let m1 = monitor_remote(d1.clone());
|
|
demonitor_remote(&m1);
|
|
ask(&format!("kill:{}", d1.index()));
|
|
let d2 = ask("spawn:exit").pid.unwrap();
|
|
let m2 = monitor_remote(d2.clone());
|
|
ask(&format!("kill:{}", d2.index()));
|
|
assert_eq!(m2.recv().unwrap().reason, DownReason::Exit.into());
|
|
// Closed-empty (`Err`) or open-empty (`Ok(None)`) both mean no notice.
|
|
let stray = matches!(m1.try_recv(), Ok(Some(_)));
|
|
println!("DEMONITOR stray={stray}");
|
|
|
|
println!("CLIENT DONE");
|
|
loop {
|
|
smarm::sleep(Duration::from_secs(3600));
|
|
}
|
|
});
|
|
}
|
|
|
|
/// The server's incarnation as this node sees it — for building pids by hand.
|
|
fn ev_incarnation() -> Incarnation {
|
|
smarm::cluster::membership::view()
|
|
.expect("manager up")
|
|
.into_iter()
|
|
.find(|i| i.name == "server")
|
|
.map(|i| i.incarnation)
|
|
.expect("server in view")
|
|
}
|
|
|
|
/// The Phase 4 c12 gate: remote monitors report the true reason, honour
|
|
/// corpses, leak nothing for unshipped pids, and cancel cleanly.
|
|
#[test]
|
|
fn remote_monitors_report_true_reasons() {
|
|
maybe_child(ROLES);
|
|
let mut server = spawn_node("server", &[]);
|
|
let saddr = server.wait_listening();
|
|
server.wait_line("READY", |l| l == "READY");
|
|
let mut client = spawn_node("client", &[("SMARM_SERVER_ADDR", &saddr)]);
|
|
client.wait_line("DOWN exit Local(Exit)", |l| l == "DOWN exit Local(Exit)");
|
|
client.wait_line("DOWN panic Local(Panic)", |l| {
|
|
l == "DOWN panic Local(Panic)"
|
|
});
|
|
client.wait_line("DOWN corpse Local(Exit)", |l| {
|
|
l == "DOWN corpse Local(Exit)"
|
|
});
|
|
client.wait_line("DOWN hidden Local(NoProc)", |l| {
|
|
l == "DOWN hidden Local(NoProc)"
|
|
});
|
|
client.wait_line("DOWN bogus Local(NoProc)", |l| {
|
|
l == "DOWN bogus Local(NoProc)"
|
|
});
|
|
client.wait_line("DEMONITOR stray=false", |l| l == "DEMONITOR stray=false");
|
|
client.wait_line("CLIENT DONE", |l| l == "CLIENT DONE");
|
|
}
|