//! RFC 010 c15 — distributed pg: sync on `NodeUp`, incremental //! `Join`/`Leave`, eager eviction announced, `NodeDown` sweep. //! //! Two nodes. The *origin* joins two local workers to `"pool"` before the //! *observer* connects (so the observer's view comes from `Sync`), exposes a //! `"go"` command inbox and then does exactly what the observer tells it: //! kill one worker, join a third, leave with the second. The observer drives //! that script through the cluster itself and asserts every step from //! `members_all` — never touching the group on its own side, except once to //! prove a mixed local+remote group reads correctly and that `members` stays //! local. `dispatch_any` is exercised both ways: into the origin's worker //! (remote pick, `send_to_remote`) and, once the origin is gone, into the //! observer's own (local pick, `send_to`). Finally the parent SIGKILLs the //! origin: the observer must sweep every remote member on `NodeDown`. #![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, RemoteName}; use smarm::cluster::{ dispatch_any, members_all, pick_any, start, Config, DispatchAnyError, GroupMember, StaticSeeds, Timing, }; use smarm::{channel, join, leave, members, register, send_to, spawn_addr, Addressable, Name, Pid}; use std::time::{Duration, Instant}; const GO: Name = Name::new("go"); const POOL: &str = "pool"; /// A pool worker's message: `"die"` stops it, anything else is printed. #[derive(Debug, PartialEq)] struct Job(String); struct Worker; impl Addressable for Worker { type Msg = Job; } impl serde::Serialize for Job { fn serialize(&self, s: S) -> Result { self.0.serialize(s) } } impl<'de> serde::Deserialize<'de> for Job { fn deserialize>(d: D) -> Result { String::deserialize(d).map(Job) } } const ROLES: &[(&str, fn())] = &[("origin", role_origin), ("observer", role_observer)]; fn cfg(name: &str, seeds: Vec<(String, String)>) -> Config { Config { node_name: name.into(), meta: NodeMeta { role: "c15".into(), region: "local".into(), }, listen_addr: "127.0.0.1:0".into(), strategy: Box::new(StaticSeeds::new(seeds)), timing: Timing::default(), } } /// A pool worker: prints every job it is handed, exits on `"die"`. fn worker() -> Pid { spawn_addr::(|rx| { while let Ok(Job(s)) = rx.recv() { if s == "die" { return; } println!("JOB {s}"); } }) } fn role_origin() { smarm::run(|| { let cluster = start(cfg("origin", vec![])).expect("binds"); // Remote dispatch lands here only for a type this node accepts. expose_type::(); let w1 = worker(); let w2 = worker(); assert!(join(POOL, w1)); assert!(join(POOL, w2)); let (go_tx, go_rx) = channel::(); register(GO, go_tx).unwrap(); expose(GO); println!("LISTENING {}", cluster.local_addr()); println!("JOINED 2"); loop { match go_rx.recv().unwrap() { 1 => { send_to(w1, Job("die".into())).unwrap(); println!("KILLED w1"); } 2 => { assert!(leave(POOL, w2)); println!("LEFT w2"); } 3 => { let w3 = worker(); assert!(join(POOL, w3)); println!("JOINED w3"); } n => panic!("unknown command {n}"), } } }); } fn remote_count(group: &str) -> usize { members_all(group) .iter() .filter(|m| matches!(m, GroupMember::Remote(_))) .count() } /// Cooperative poll until `pred`; panics (with the last view) on timeout. fn wait_view(what: &str, group: &str, pred: impl Fn(&[GroupMember]) -> bool) { let deadline = Instant::now() + Duration::from_secs(5); loop { let v = members_all(group); if pred(&v) { return; } assert!( Instant::now() < deadline, "timed out waiting for {what}; view = {v:?}" ); smarm::sleep(Duration::from_millis(5)); } } fn role_observer() { let origin_addr = std::env::var("SMARM_ORIGIN_ADDR").expect("SMARM_ORIGIN_ADDR"); smarm::run(move || { let _cluster = start(cfg("observer", vec![("origin".into(), origin_addr)])).expect("binds"); let ev = subscribe().unwrap(); loop { match ev.rx.recv().unwrap() { NodeEvent::NodeUp(i) if i.name == "origin" => break, _ => {} } } let go = |n: u8| remote::send(RemoteName::new("origin", GO), n).unwrap(); // Sync: both pre-existing members arrive with no join on this side. wait_view("sync of 2 remote members", POOL, |v| { v.len() == 2 && v.iter().all(|m| matches!(m, GroupMember::Remote(_))) }); let synced = members_all(POOL); assert!(synced.iter().all(|m| match m { GroupMember::Remote(p) => p.node() == "origin", GroupMember::Local(_) => false, })); println!("SEES 2"); // Origin-side death: the origin's reaper announces the leave. go(1); wait_view("death evicted on observer", POOL, |v| v.len() == 1); println!("SEES 1 after death"); // Incremental Join. go(3); wait_view("incremental join", POOL, |v| v.len() == 2); println!("SEES 2 after join"); // Voluntary Leave. go(2); wait_view("incremental leave", POOL, |v| v.len() == 1); println!("SEES 1 after leave"); // Mixed group: our own member sits beside the remote one in // `members_all`; `members` stays local-only. let me = worker(); assert!(join(POOL, me)); wait_view("mixed local+remote", POOL, |v| { v.len() == 2 && v.contains(&GroupMember::Local(me.erase())) }); assert_eq!( members(POOL), vec![me.erase()], "local API never shows remotes" ); assert_eq!(remote_count(POOL), 1); println!("MIXED ok"); // dispatch_any: the store's first entry is the origin's w3 (it was // announced before we joined), so the pick is remote and the job // crosses the wire — the origin's worker prints it. let picked = pick_any(POOL).expect("pool has members"); assert!( matches!(picked, GroupMember::Remote(_)), "first entry is remote: {picked:?}" ); let reached = dispatch_any::(POOL, Job("from-observer".into())).unwrap(); assert_eq!(reached, picked); println!("DISPATCHED remote"); println!("PARK"); // Parent SIGKILLs the origin now: NodeDown must sweep its member, // ours must survive. wait_view("node_down sweep", POOL, |v| { v == [GroupMember::Local(me.erase())] }); assert_eq!(members(POOL), vec![me.erase()]); println!("SWEPT"); // Now the only member is ours: a local pick, a local send. let reached = dispatch_any::(POOL, Job("local".into())).unwrap(); assert_eq!(reached, GroupMember::Local(me.erase())); // And an empty group hands the message back. match dispatch_any::("nobody", Job("lost".into())) { Err(DispatchAnyError::NoMember(Job(s))) => assert_eq!(s, "lost"), other => panic!("expected NoMember, got {other:?}"), } println!("DISPATCHED local"); loop { smarm::sleep(Duration::from_secs(3600)); } }); } /// The Phase 5 gate: sync, join, leave, death, node_down — all observed from /// the peer, none of them a group operation on the peer — plus dispatch_any /// reaching a remote member and a local one. #[test] fn groups_span_two_nodes() { maybe_child(ROLES); let mut origin = spawn_node("origin", &[]); let addr = origin.wait_listening(); origin.wait_line("JOINED 2", |l| l == "JOINED 2"); let mut observer = spawn_node("observer", &[("SMARM_ORIGIN_ADDR", &addr)]); observer.wait_line("SEES 2", |l| l == "SEES 2"); origin.wait_line("KILLED w1", |l| l == "KILLED w1"); observer.wait_line("SEES 1 after death", |l| l == "SEES 1 after death"); origin.wait_line("JOINED w3", |l| l == "JOINED w3"); observer.wait_line("SEES 2 after join", |l| l == "SEES 2 after join"); origin.wait_line("LEFT w2", |l| l == "LEFT w2"); observer.wait_line("SEES 1 after leave", |l| l == "SEES 1 after leave"); observer.wait_line("MIXED ok", |l| l == "MIXED ok"); observer.wait_line("DISPATCHED remote", |l| l == "DISPATCHED remote"); origin.wait_line("JOB from-observer", |l| l == "JOB from-observer"); observer.wait_line("PARK", |l| l == "PARK"); origin.kill(); observer.wait_line("SWEPT", |l| l == "SWEPT"); // Order between the root's line and the worker's is scheduling; wait // for the later one to be certain both happened. observer.wait_line("DISPATCHED local", |l| l == "DISPATCHED local"); observer.wait_line("JOB local", |l| l == "JOB local"); }