//! RFC 010 c6b — the handshake on the accept/connect path. //! //! Path-level tests drive [`dial_handshake`]/[`accept_handshake`] over the //! loopback transport on plain threads (its intended use — synchronous, no //! runtime). Integration tests run the manager-backed [`dial`] and //! [`spawn_acceptor`] over real localhost TCP inside `smarm::run`, and the //! two-node case as subprocesses via the c4 harness. Flake budget: see //! tests/common/mod.rs. #![cfg(feature = "cluster")] mod common; use std::sync::mpsc; use std::time::{Duration, Instant}; use common::{maybe_child, spawn_node, WAIT}; use smarm::cluster::connect::{ accept_handshake, dial, dial_handshake, spawn_acceptor, DialError, HandshakeError, HANDSHAKE_TIMEOUT, }; use smarm::cluster::envelope::{Frame, NodeMeta, RejectReason}; use smarm::cluster::handshake::{Local, PeerStanding}; use smarm::cluster::manager::{Call, Manager, Reply, MANAGER}; use smarm::cluster::transport::loopback::LoopbackTransport; use smarm::cluster::transport::tcp::TcpTransport; use smarm::cluster::transport::{FramedConn, Transport}; use smarm::cluster::Timing; use smarm::gen_server::{self, GenServerBuilder}; use smarm::pg::Incarnation; use smarm::{run, sleep}; const ROLES: &[(&str, fn())] = &[ ("hs_listener", role_hs_listener), ("hs_dialer", role_hs_dialer), ]; const HASH: u64 = 0xC6B0_C6B0_C6B0_C6B0; fn local(name: &str) -> Local { Local { node_name: name.into(), incarnation: Incarnation::new(3), build_hash: HASH, meta: NodeMeta { role: "test".into(), region: "test".into(), }, } } /// A loopback conn pair as `FramedConn`s, ready for a threaded handshake. fn loopback_pair() -> (FramedConn, FramedConn) { let t = LoopbackTransport::default(); let mut l = t.listen("hs").unwrap(); let dialer = FramedConn::new(t.dial("hs").unwrap()); let accepted = FramedConn::new(l.accept().unwrap()); (dialer, accepted) } /// Far-future deadline for loopback paths, where it cannot fire anyway. fn no_deadline() -> Instant { Instant::now() + Duration::from_secs(3600) } /// Park the node forever: it has announced everything the parent asserts on, /// and must now hold its connection open until SIGKILLed. fn park() -> ! { loop { sleep(Duration::from_secs(1)); } } /// Cooperative bounded receive across the closure/actor boundary. A blocking /// `std::mpsc` wait would park the OS thread and starve the single-threaded /// scheduler, so every wait inside `run` polls with [`sleep`] instead. fn poll_recv(rx: &mpsc::Receiver, what: &str) -> T { let deadline = Instant::now() + WAIT; loop { match rx.try_recv() { Ok(v) => return v, Err(mpsc::TryRecvError::Empty) => { assert!(Instant::now() < deadline, "timed out waiting for {what}"); sleep(Duration::from_millis(1)); } Err(mpsc::TryRecvError::Disconnected) => panic!("channel closed waiting for {what}"), } } } // --------------------------------------------------------------------------- // Path level, over loopback on plain threads // --------------------------------------------------------------------------- #[test] fn loopback_happy_path_establishes_both_ends() { maybe_child(ROLES); let (mut dialer, mut accepted) = loopback_pair(); let responder = std::thread::spawn(move || { accept_handshake( &mut accepted, local("node-b"), |name| { assert_eq!(name, "node-a"); PeerStanding::Free }, no_deadline(), ) }); let peer_of_dialer = dial_handshake(&mut dialer, &local("node-a"), no_deadline()).unwrap(); let peer_of_acceptor = responder.join().unwrap().unwrap(); assert_eq!(peer_of_dialer.node_name, "node-b"); assert_eq!(peer_of_acceptor.node_name, "node-a"); } #[test] fn loopback_hash_mismatch_rejected_with_frame_then_eof() { maybe_child(ROLES); let (mut dialer, mut accepted) = loopback_pair(); let mut wrong = local("node-b"); wrong.build_hash ^= 1; let responder = std::thread::spawn(move || { accept_handshake(&mut accepted, wrong, |_| PeerStanding::Free, no_deadline()) }); // The dial side receives the reject frame — the compatibility anchor. match dial_handshake(&mut dialer, &local("node-a"), no_deadline()) { Err(HandshakeError::Rejected(RejectReason::HashMismatch)) => {} other => panic!("expected HashMismatch reject, got {other:?}"), } match responder.join().unwrap() { Err(HandshakeError::Rejected(RejectReason::HashMismatch)) => {} other => panic!("expected accept side to report the reject, got {other:?}"), } } #[test] fn loopback_tie_break_loser_closed_silently() { maybe_child(ROLES); // The inbound dial is from "node-z"; we are "node-a" with our own dial to // node-z in flight. dial_wins("node-z", "node-a") is false, so the // inbound loses: closed with no frame at all. let (mut dialer, mut accepted) = loopback_pair(); let responder = std::thread::spawn(move || { accept_handshake( &mut accepted, local("node-a"), |_| PeerStanding::Dialing, no_deadline(), ) }); // Silent close: the dial side sees EOF, never a frame. match dial_handshake(&mut dialer, &local("node-z"), no_deadline()) { Err(HandshakeError::Closed) => {} other => panic!("expected silent close (Closed), got {other:?}"), } match responder.join().unwrap() { Err(HandshakeError::TieBreakLoss) => {} other => panic!("expected TieBreakLoss on the accept side, got {other:?}"), } } #[test] fn loopback_read_ahead_past_hello_survives_into_established_conn() { maybe_child(ROLES); // The buffer trap, proven: the dialer coalesces Hello + Heartbeat before // the responder's first read, so the Heartbeat lands in the shared // FramedConn's decode buffer during the handshake. The dialer sends // nothing afterwards — the post-handshake recv can only succeed if the // read-ahead travelled with the FramedConn. let (mut dialer, mut accepted) = loopback_pair(); let (_init, hello) = smarm::cluster::handshake::Initiator::new(&local("node-a")); dialer.send(&hello).unwrap(); dialer.send(&Frame::Heartbeat).unwrap(); // Both frames are buffered before the responder reads at all. let (tx, rx) = mpsc::channel(); std::thread::spawn(move || { let peer = accept_handshake( &mut accepted, local("node-b"), |_| PeerStanding::Free, no_deadline(), ) .unwrap(); let next = accepted.recv(); let _ = tx.send((peer, next)); }); // A bounded wait: if the Heartbeat were NOT carried in the buffer, the // recv above would block forever (the dialer stays open and silent). let (peer, next) = rx .recv_timeout(Duration::from_secs(5)) .expect("read-ahead lost: post-handshake recv blocked"); assert_eq!(peer.node_name, "node-a"); match next { Ok(Some(Frame::Heartbeat)) => {} other => panic!("expected the read-ahead Heartbeat, got {other:?}"), } drop(dialer); } // --------------------------------------------------------------------------- // Deadline + manager integration, over TCP inside the runtime // --------------------------------------------------------------------------- #[test] fn tcp_silent_peer_times_out_on_the_accept_path() { maybe_child(ROLES); run(|| { let t = TcpTransport; let mut l = t.listen("127.0.0.1:0").unwrap(); // Connect and then say nothing at all. let silent = t.dial(&l.local_addr()).unwrap(); let mut accepted = FramedConn::new(l.accept().unwrap()); let (tx, rx) = mpsc::channel(); smarm::spawn(move || { let r = accept_handshake( &mut accepted, local("node-b"), |_| PeerStanding::Free, Instant::now() + Duration::from_millis(200), ); let _ = tx.send(r); }); match poll_recv(&rx, "accept-path outcome") { Err(HandshakeError::TimedOut) => {} other => panic!("expected TimedOut, got {other:?}"), } drop(silent); }); } /// Poll the manager until its peer set matches `expected` (sorted), or fail. fn wait_peers(expected: &[&str]) { let want: Vec = expected.iter().map(|s| s.to_string()).collect(); for _ in 0..5000 { if let Ok(Reply::Peers(got)) = gen_server::call(MANAGER, Call::Peers) { if got == want { return; } } sleep(Duration::from_millis(1)); } let got = gen_server::call(MANAGER, Call::Peers); panic!("timed out waiting for peers == {want:?}; last = {got:?}"); } #[test] fn tcp_duplicate_name_rejected_by_acceptor() { maybe_child(ROLES); run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let listener = TcpTransport.listen("127.0.0.1:0").unwrap(); let acceptor = spawn_acceptor(listener, local("node-b"), Timing::default()); let addr = acceptor.local_addr().to_string(); // First dial offering "dup-node": establishes and registers. let mut first = FramedConn::new(TcpTransport.dial(&addr).unwrap()); let peer = dial_handshake( &mut first, &local("dup-node"), Instant::now() + HANDSHAKE_TIMEOUT, ) .unwrap(); assert_eq!(peer.node_name, "node-b"); wait_peers(&["dup-node"]); // Second dial offering the same name: deterministic NameTaken. let mut second = FramedConn::new(TcpTransport.dial(&addr).unwrap()); match dial_handshake( &mut second, &local("dup-node"), Instant::now() + HANDSHAKE_TIMEOUT, ) { Err(HandshakeError::Rejected(RejectReason::NameTaken)) => {} other => panic!("expected NameTaken, got {other:?}"), } // The established connection was untouched by the rejected one. wait_peers(&["dup-node"]); // Teardown: the acceptor owns no connections, so the established one // is torn down through the table. acceptor.shutdown(); assert!(matches!( gen_server::call( MANAGER, Call::Disconnect { name: "dup-node".to_string() } ), Ok(Reply::Disconnected) )); wait_peers(&[]); first.close(); mgr.shutdown(); }); } #[test] fn dial_intent_cleared_when_dialer_dies() { maybe_child(ROLES); run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let (begun_tx, begun_rx) = mpsc::channel(); let (go_tx, go_rx) = mpsc::channel::<()>(); smarm::spawn(move || { let me = smarm::self_pid(); match gen_server::call( MANAGER, Call::DialBegin { name: "ghost".into(), pid: me, }, ) { Ok(Reply::DialBegan(true)) => {} other => panic!("DialBegin failed: {other:?}"), } let _ = begun_tx.send(()); let () = poll_recv(&go_rx, "go signal"); panic!("dialer dies mid-dial"); }); poll_recv(&begun_rx, "DialBegin done"); // While the dialer lives, the intent is visible. match gen_server::call( MANAGER, Call::Standing { peer_name: "ghost".into(), }, ) { Ok(Reply::Standing(s)) => assert_eq!(s, PeerStanding::Dialing), other => panic!("PeerStanding failed: {other:?}"), } // Kill it; the monitor must clear the intent without cooperation. go_tx.send(()).unwrap(); let deadline = Instant::now() + WAIT; loop { match gen_server::call( MANAGER, Call::Standing { peer_name: "ghost".into(), }, ) { Ok(Reply::Standing(s)) if s != PeerStanding::Dialing => break, _ if Instant::now() > deadline => { panic!("dial intent not cleared after dialer death") } _ => sleep(Duration::from_millis(1)), } } mgr.shutdown(); }); } // --------------------------------------------------------------------------- // Two nodes, two processes: the integrated dial against a real acceptor // --------------------------------------------------------------------------- /// Announce, then park forever. Neither role ever tears its connection /// down: a table entry only exists while the *peer* holds its side open, so /// any teardown here would retract the other node's observation before it /// had made it. The parent reaps both with SIGKILL once it has both /// announcements (see [`common::Node`]'s `Drop`). fn role_hs_listener() { run(|| { let _mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let listener = TcpTransport.listen("127.0.0.1:0").unwrap(); let acceptor = spawn_acceptor(listener, local("node-b"), Timing::default()); println!("LISTENING {}", acceptor.local_addr()); wait_peers(&["node-a"]); println!("PEERS node-a"); park(); }); } fn role_hs_dialer() { let addr = std::env::var("SMARM_PEER_ADDR").expect("SMARM_PEER_ADDR not set"); run(move || { let _mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let (tx, rx) = mpsc::channel(); smarm::spawn(move || { let r = dial( &TcpTransport, &addr, "node-b", &local("node-a"), Timing::default(), ); let _ = tx.send(r); }); if let Err(e) = poll_recv(&rx, "dial outcome") { println!("DIAL failed: {e:?}"); std::process::exit(3); } wait_peers(&["node-b"]); println!("PEERS node-b"); park(); }); } #[test] fn two_node_integrated_handshake_over_tcp() { maybe_child(ROLES); let mut listener = spawn_node("hs_listener", &[]); let addr = listener.wait_listening(); let mut dialer = spawn_node("hs_dialer", &[("SMARM_PEER_ADDR", &addr)]); // Each node reports its own table naming the other: a real dial against a // real acceptor established in both directions. Both nodes then park — // clean-exit behaviour is the c4 harness's own smoke test, and demanding // it here would mean a teardown, which is exactly what cannot be ordered // safely across two processes. Dropping the nodes SIGKILLs them. dialer.wait_line("PEERS node-b", |l| l == "PEERS node-b"); listener.wait_line("PEERS node-a", |l| l == "PEERS node-a"); } // --------------------------------------------------------------------------- // Integrated-dial guardrails (no acceptor involved) // --------------------------------------------------------------------------- #[test] fn concurrent_dial_to_same_name_refused() { maybe_child(ROLES); run(|| { let mgr = GenServerBuilder::new(Manager::new()) .named(MANAGER) .start() .expect("manager name is free"); let (begun_tx, begun_rx) = mpsc::channel(); let (go_tx, go_rx) = mpsc::channel::<()>(); // First dialer parks with the intent held (it never connects — // 'holding the intent' is all this test needs from it). smarm::spawn(move || { let me = smarm::self_pid(); assert!(matches!( gen_server::call( MANAGER, Call::DialBegin { name: "node-x".into(), pid: me, } ), Ok(Reply::DialBegan(true)) )); let _ = begun_tx.send(()); let () = poll_recv(&go_rx, "go signal"); let _ = gen_server::call( MANAGER, Call::DialEnd { name: "node-x".into(), }, ); }); poll_recv(&begun_rx, "DialBegin done"); // Second integrated dial to the same name: refused before connecting // (the addr is unroutable on purpose — it must never be dialed). let (tx, rx) = mpsc::channel(); smarm::spawn(move || { let r = dial( &TcpTransport, "127.0.0.1:1", "node-x", &local("node-a"), Timing::default(), ); let _ = tx.send(r); }); match poll_recv(&rx, "second dial outcome") { Err(DialError::AlreadyDialing) => {} other => panic!("expected AlreadyDialing, got {other:?}"), } go_tx.send(()).unwrap(); mgr.shutdown(); }); }