//! 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, } #[derive(Debug)] struct Answer { text: String, pid: Option>, } struct Client; impl Addressable for Client { type Msg = Answer; } impl serde::Serialize for Ctl { fn serialize(&self, s: S) -> Result { 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: D) -> Result { let (cmd, reply_to) = <(String, RemotePid)>::deserialize(d)?; Ok(Ctl { cmd, reply_to }) } } impl serde::Serialize for Answer { fn serialize(&self, s: S) -> Result { 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: D) -> Result { let (text, pid) = <(String, Option>)>::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::::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::::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 = 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:` 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::"`. fn role_server() { smarm::run(move || { let cluster = start(cfg("server", vec![])).expect("binds"); println!("LISTENING {}", cluster.local_addr()); let (tx, rx) = channel::(); register(CTL, tx).unwrap(); expose(CTL); println!("READY"); let mut workers: HashMap> = HashMap::new(); loop { let ctl = rx.recv().unwrap(); println!("CTL {}", ctl.cmd); let (text, pid): (String, Option>) = 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::(); let me: Pid = install::(tx); expose_type::(); 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::::from_parts("server", server_inc, idx, gen); println!( "DOWN hidden {:?}", monitor_remote(hidden).recv().unwrap().reason ); let bogus = RemotePid::::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"); }