//! RFC 010 c13 — connection-loss synthesis. //! //! Local suite (`run()`, no network): the read-side backstop. A //! `RemoteMonitor` whose channel closes without a notice reads as //! `Disconnected` exactly once (a `Monitor` command that reached the conn //! actor's inbox but was never processed — the drain gap); after //! `demonitor_remote` a closed channel stays a plain `Err`, never a notice. //! //! Cross-process: the headline contrast — an actor's own death gives its //! TRUE reason, loss of the LINK gives `Disconnected` (both a commanded //! `Disconnect` and a SIGKILLed peer process are `Disconnected` from the //! monitor's view: nobody is left to say otherwise). Reconnect does not //! resurrect: the old monitor yields nothing more, proven by stream ORDER //! (a fresh monitor over the new link delivers first). The ignored test //! trips liveness by SIGSTOP and then drops the link too, asserting exactly //! one notice for one monitor. #![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::manager::{Call, Reply, MANAGER}; 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, gen_server, 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 ========================================= /// A `Monitor` command handed to the connection but never processed (its /// receiver dropped unread) reads as `Disconnected` — once. A second read /// is the ordinary closed-channel `Err`, so "exactly one notice" holds. #[test] fn unread_command_reads_as_disconnected_once() { maybe_child(ROLES); run(|| { remote::set_local_identity("me", Incarnation::new(7)); let (probe_tx, _probe_rx) = channel(); let inbox = remote::bind_outbound_probe_with_monitors("peer", Incarnation::new(5), probe_tx); let target = RemotePid::::from_parts("peer", Incarnation::new(5), 9, 1); let m = monitor_remote(target.clone()); assert!( matches!(m.try_recv(), Ok(None)), "command is in flight, no notice yet" ); drop(inbox); // the conn actor died with the command unread let d = m.recv().unwrap(); assert_eq!(d.pid, target); assert_eq!(d.reason, RemoteDownReason::Disconnected); assert!( m.recv().is_err(), "second read is closed, not a second notice" ); assert!(m.try_recv().is_err()); }); } /// After `demonitor_remote`, a closed channel is a closed channel: no /// notice is synthesized for a monitor the caller cancelled. #[test] fn cancelled_monitor_never_synthesizes() { maybe_child(ROLES); run(|| { remote::set_local_identity("me", Incarnation::new(7)); let (probe_tx, _probe_rx) = channel(); let inbox = remote::bind_outbound_probe_with_monitors("peer", Incarnation::new(5), probe_tx); let target = RemotePid::::from_parts("peer", Incarnation::new(5), 9, 1); let m = monitor_remote(target); demonitor_remote(&m); drop(inbox); assert!(m.recv().is_err()); assert!(m.try_recv().is_err()); }); } // ================= cross-process ====================================== const ROLES: &[(&str, fn())] = &[ ("server", role_server), ("client", role_client), ("client_stop", role_client_stop), ]; const CTL: Name = Name::new("c13.ctl"); fn cfg(name: &str, seeds: Vec<(String, String)>) -> Config { Config { node_name: name.to_string(), meta: NodeMeta { role: "c13".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(), } } /// The p11 knobs make the liveness test fast: both roles of that test are /// spawned with `SMARM_FAST_TIMING=1` and agree on a 100ms heartbeat / /// 500ms liveness window. Everything else runs the shipping defaults. fn timing() -> Timing { if std::env::var_os("SMARM_FAST_TIMING").is_some() { Timing { heartbeat_interval: Duration::from_millis(100), liveness_timeout: Duration::from_millis(500), initial_backoff: Duration::from_millis(50), max_backoff: Duration::from_millis(500), ..Timing::default() } } else { 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"), } } } fn disconnect(name: &str) { assert!(matches!( gen_server::call( MANAGER, Call::Disconnect { name: name.to_string() } ), Ok(Reply::Disconnected) )); } /// Server: `spawn` ⇒ a parked worker (answer carries its pid); /// `kill:` releases it, whereupon it returns (Exit). 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" => { let (go_tx, go_rx) = channel::<()>(); let p: Pid = spawn(move || { let _ = go_rx.recv(); }) .pid(); workers.insert(p.index(), go_tx); ( "ok".into(), Some(RemotePid::from_local(p).expect("identity set")), ) } 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(); } }); } /// Client-side setup shared by both client roles: join, expose the reply /// path, hand back an `ask` closure and the membership stream. fn client_setup() -> ( smarm::cluster::Cluster, smarm::cluster::membership::MembershipEvents, impl Fn(&str) -> Answer, ) { let server_addr = std::env::var("SMARM_SERVER_ADDR").expect("SMARM_SERVER_ADDR"); 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 = move |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() }; (cluster, ev, ask) } fn role_client() { smarm::run(move || { let (_cluster, ev, ask) = client_setup(); // 1. Headline: actor death ⇒ TRUE reason; link cut ⇒ Disconnected. let a = ask("spawn").pid.unwrap(); let b = ask("spawn").pid.unwrap(); let ma = monitor_remote(a.clone()); let mb = monitor_remote(b.clone()); ask(&format!("kill:{}", a.index())); let d = ma.recv().unwrap(); assert_eq!(d.pid, a); println!("DOWN actor {:?}", d.reason); disconnect("server"); let d = mb.recv().unwrap(); assert_eq!(d.pid, b); println!("DOWN link {:?}", d.reason); // 2. Reconnect does not resurrect. The connector redials on // node_down; over the NEW link a fresh monitor delivers, while // the old one (already answered) yields nothing further — order // proves it, and `b` is even still alive on the server. wait_up(&ev, "server"); println!("RECONNECTED"); let c = ask("spawn").pid.unwrap(); let mc = monitor_remote(c.clone()); ask(&format!("kill:{}", b.index())); ask(&format!("kill:{}", c.index())); assert_eq!(mc.recv().unwrap().reason, DownReason::Exit.into()); let stray = matches!(mb.try_recv(), Ok(Some(_))); println!("RESURRECT stray={stray}"); // 3. Peer PROCESS killed ⇒ Disconnected too (nobody is left to send // Down): the parent SIGKILLs the server once it sees the marker. let e = ask("spawn").pid.unwrap(); let me_ = monitor_remote(e.clone()); println!("KILL SERVER NOW"); let d = me_.recv().unwrap(); assert_eq!(d.pid, e); println!("DOWN procdeath {:?}", d.reason); println!("CLIENT DONE"); loop { smarm::sleep(Duration::from_secs(3600)); } }); } /// The slow role: liveness expiry (peer SIGSTOPped) followed by the link /// dropping for real (peer SIGKILLed) — one monitor, exactly one notice. fn role_client_stop() { smarm::run(move || { let (_cluster, _ev, ask) = client_setup(); let a = ask("spawn").pid.unwrap(); let ma = monitor_remote(a.clone()); println!("STOP SERVER NOW"); let d = ma.recv().unwrap(); // liveness expiry, ~liveness_timeout assert_eq!(d.pid, a); println!("DOWN stopped {:?}", d.reason); println!("KILL SERVER NOW"); // Give the drop every chance to produce a second notice, then look. smarm::sleep(Duration::from_secs(1)); let dup = matches!(ma.try_recv(), Ok(Some(_))); println!("DUPLICATE dup={dup}"); println!("CLIENT DONE"); loop { smarm::sleep(Duration::from_secs(3600)); } }); } /// The Phase 4 c13 gate: partition vs. death distinguishable; nothing /// survives reconnect; a dead peer process is a Disconnected too. #[test] fn link_loss_is_disconnected_and_does_not_survive_reconnect() { 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 actor Local(Exit)", |l| l == "DOWN actor Local(Exit)"); client.wait_line("DOWN link Disconnected", |l| l == "DOWN link Disconnected"); client.wait_line("RECONNECTED", |l| l == "RECONNECTED"); client.wait_line("RESURRECT stray=false", |l| l == "RESURRECT stray=false"); client.wait_line("KILL SERVER NOW", |l| l == "KILL SERVER NOW"); server.kill(); client.wait_line("DOWN procdeath Disconnected", |l| { l == "DOWN procdeath Disconnected" }); client.wait_line("CLIENT DONE", |l| l == "CLIENT DONE"); } /// Covers the invariant the headline test cannot: liveness expiry and the /// transport drop both firing for the same connection yield ONE notice. /// Runs on the fast [`timing`] (both roles) — was `#[ignore]`d at the 4s /// default until the p11 knobs landed. #[test] fn timeout_then_drop_yields_one_notice() { maybe_child(ROLES); let fast = ("SMARM_FAST_TIMING", "1"); let mut server = spawn_node("server", &[fast]); let saddr = server.wait_listening(); server.wait_line("READY", |l| l == "READY"); let mut client = spawn_node("client_stop", &[("SMARM_SERVER_ADDR", &saddr), fast]); client.wait_line("STOP SERVER NOW", |l| l == "STOP SERVER NOW"); let spid = server.pid().expect("server alive") as libc::pid_t; assert_eq!(unsafe { libc::kill(spid, libc::SIGSTOP) }, 0); client.wait_line("DOWN stopped Disconnected", |l| { l == "DOWN stopped Disconnected" }); client.wait_line("KILL SERVER NOW", |l| l == "KILL SERVER NOW"); server.kill(); // SIGKILL works on a stopped process; Drop would too client.wait_line("DUPLICATE dup=false", |l| l == "DUPLICATE dup=false"); client.wait_line("CLIENT DONE", |l| l == "CLIENT DONE"); }