//! RFC 010 c7a — membership events and the view, at the manager. //! //! Same construction as the c6a lifecycle suite: the handshake is bypassed, //! connections are built already-established over localhost TCP pairs with //! fabricated `Peer`s, and the manager is started plainly so the test can //! terminate. What is under test is the membership layer that c7 adds to the //! manager: `node_up`/`node_down` events to subscribers (snapshot-then-stream), //! the view, and NodeId identity — memoized per `(name, incarnation)`, so a //! reconnect blip keeps its id and a restart (new incarnation) gets a fresh one. #![cfg(feature = "cluster")] use std::time::Duration; use smarm::cluster::envelope::NodeMeta; use smarm::cluster::handshake::Peer; use smarm::cluster::manager::{Call, Manager, Reply, MANAGER}; use smarm::cluster::membership::{subscribe, view, MembershipEvents, NodeEvent}; use smarm::cluster::spawn_established; use smarm::cluster::transport::tcp::TcpTransport; use smarm::cluster::transport::{Conn, FramedConn, Transport}; use smarm::cluster::Timing; use smarm::gen_server::{self, GenServerBuilder}; use smarm::pg::{Incarnation, NodeId}; use smarm::run; /// A fabricated post-handshake peer identity, with the incarnation under the /// test's control (it is identity-bearing here, unlike in the c6a suite). fn peer(name: &str, inc: u32) -> Peer { Peer { node_name: name.to_string(), incarnation: Incarnation::new(inc), meta: NodeMeta { role: "test".to_string(), region: "test".to_string(), }, } } /// One established transport pair over localhost (TCP backlog covers the /// sequential dial-then-accept, as in the c3 conformance suite). fn pair(t: &dyn Transport) -> (Box, Box) { let mut l = t.listen("127.0.0.1:0").unwrap(); let a = t.dial(&l.local_addr()).unwrap(); let b = l.accept().unwrap(); (a, b) } /// The next event, or a panic naming the wait. The bound is generous against /// a sub-millisecond real cost. fn next_event(ev: &MembershipEvents, waiting_for: &str) -> NodeEvent { ev.rx .recv_timeout(Duration::from_secs(5)) .unwrap_or_else(|e| panic!("timed out waiting for {waiting_for}: {e:?}")) } /// Assert the subscription is drained: no event is pending. fn assert_quiet(ev: &MembershipEvents) { assert!(matches!(ev.rx.try_recv(), Ok(None))); } fn disconnect(name: &str) { assert!(matches!( gen_server::call( MANAGER, Call::Disconnect { name: name.to_string() } ), Ok(Reply::Disconnected) )); } /// Live subscription: an empty snapshot, then `NodeUp` on registration and /// `NodeDown` (same id) on commanded disconnect and on peer EOF alike. #[test] fn subscriber_sees_up_and_down() { run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let ev = subscribe().expect("manager is up"); assert_quiet(&ev); // nothing live: the snapshot is empty let t = TcpTransport; let (a1, b1) = pair(&t); let (a2, b2) = pair(&t); spawn_established(FramedConn::new(a1), peer("node-b", 1), Timing::default()) .expect("node-b registers"); spawn_established(FramedConn::new(a2), peer("node-c", 1), Timing::default()) .expect("node-c registers"); let up_b = match next_event(&ev, "node_up(node-b)") { NodeEvent::NodeUp(info) => { assert_eq!(info.name, "node-b"); assert_eq!(info.incarnation, Incarnation::new(1)); assert_eq!(info.meta.role, "test"); info } other => panic!("expected node_up(node-b), got {other:?}"), }; let up_c = match next_event(&ev, "node_up(node-c)") { NodeEvent::NodeUp(info) => { assert_eq!(info.name, "node-c"); info } other => panic!("expected node_up(node-c), got {other:?}"), }; assert_ne!(up_b.node, up_c.node, "distinct peers get distinct ids"); // Commanded disconnect: down with node-b's id. disconnect("node-b"); assert_eq!( next_event(&ev, "node_down(node-b)"), NodeEvent::NodeDown(up_b.clone()) ); // Peer EOF, no command: down with node-c's id. drop(b2); assert_eq!( next_event(&ev, "node_down(node-c)"), NodeEvent::NodeDown(up_c.clone()) ); assert_quiet(&ev); drop(b1); mgr.shutdown(); }); } /// Snapshot-then-stream: a subscriber arriving after connections established /// receives one `NodeUp` per live peer before anything else, and the view /// call agrees with it. #[test] fn late_subscriber_gets_snapshot_and_view_agrees() { run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let t = TcpTransport; let (a1, b1) = pair(&t); let (a2, b2) = pair(&t); spawn_established(FramedConn::new(a1), peer("node-b", 1), Timing::default()) .expect("node-b registers"); spawn_established(FramedConn::new(a2), peer("node-c", 1), Timing::default()) .expect("node-c registers"); let ev = subscribe().expect("manager is up"); let mut names = Vec::new(); for _ in 0..2 { match next_event(&ev, "a snapshot node_up") { NodeEvent::NodeUp(info) => names.push(info.name), other => panic!("expected a snapshot node_up, got {other:?}"), } } names.sort(); assert_eq!(names, ["node-b", "node-c"]); assert_quiet(&ev); // the snapshot is exactly the live set let mut v = view().expect("manager is up"); v.sort_by(|a, b| a.name.cmp(&b.name)); assert_eq!(v.len(), 2); assert_eq!(v[0].name, "node-b"); assert_eq!(v[1].name, "node-c"); disconnect("node-b"); disconnect("node-c"); drop((b1, b2)); // Drain the two downs so the subscription ends quiet. let _ = next_event(&ev, "node_down"); let _ = next_event(&ev, "node_down"); mgr.shutdown(); }); } /// NodeId identity: a restart (same name, new incarnation) is a NEW id — the /// ghost and its successor are distinguishable — while a reconnect blip (same /// name, same incarnation) keeps its id. #[test] fn restart_gets_new_id_blip_keeps_id() { run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let ev = subscribe().expect("manager is up"); let t = TcpTransport; let id = |e: NodeEvent, what: &str| -> NodeId { match e { NodeEvent::NodeUp(info) => info.node, other => panic!("expected node_up ({what}), got {other:?}"), } }; let down_id = |e: NodeEvent, what: &str| -> NodeId { match e { NodeEvent::NodeDown(info) => info.node, other => panic!("expected node_down ({what}), got {other:?}"), } }; // Up at incarnation 1, then the peer dies (EOF). let (a1, b1) = pair(&t); spawn_established(FramedConn::new(a1), peer("node-b", 1), Timing::default()) .expect("registers"); let id1 = id(next_event(&ev, "node_up inc 1"), "inc 1"); drop(b1); assert_eq!(down_id(next_event(&ev, "node_down inc 1"), "inc 1"), id1); // Restart: new incarnation, new id — the ghost's id is not reused. let (a2, b2) = pair(&t); spawn_established(FramedConn::new(a2), peer("node-b", 2), Timing::default()) .expect("registers"); let id2 = id(next_event(&ev, "node_up inc 2"), "inc 2"); assert_ne!( id1, id2, "a restarted node must be distinguishable from its ghost" ); // Blip: the same incarnation reconnects and keeps its id. disconnect("node-b"); assert_eq!(down_id(next_event(&ev, "node_down inc 2"), "inc 2"), id2); let (a3, b3) = pair(&t); spawn_established(FramedConn::new(a3), peer("node-b", 2), Timing::default()) .expect("registers"); let id3 = id(next_event(&ev, "node_up after blip"), "blip"); assert_eq!( id2, id3, "a reconnect at the same incarnation is the same node" ); disconnect("node-b"); let _ = next_event(&ev, "final node_down"); drop((b2, b3)); mgr.shutdown(); }); } /// A dropped subscriber is pruned on the next emit and never disturbs the /// manager or a live subscriber. #[test] fn dead_subscriber_is_pruned() { run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let dead = subscribe().expect("manager is up"); drop(dead); let live = subscribe().expect("manager is up"); let t = TcpTransport; let (a1, b1) = pair(&t); spawn_established(FramedConn::new(a1), peer("node-b", 1), Timing::default()) .expect("registers"); match next_event(&live, "node_up despite a dead co-subscriber") { NodeEvent::NodeUp(info) => assert_eq!(info.name, "node-b"), other => panic!("expected node_up, got {other:?}"), } disconnect("node-b"); let _ = next_event(&live, "node_down"); drop(b1); mgr.shutdown(); }); }