Tree snapshot of d9c62a8 (2026-08-18). The 20 source commits between
16ef583 (c9) and d9c62a8 were never pushed and the clone that held them
was lost; this commit carries their combined tree verbatim so the build
history stays auditable from the c1–c9 commits below it. Original
hashes as recorded in the session handoff:
c10 f03e94d pid targeting + auto-serialization (RemotePid, D14 name
on the wire); Phase 3 gate
c11 7ef4bad DownReason::Disconnected, wire tag 5
c12 d124162 remote monitors (Monitor/Demonitor/Down frames)
c13 9de967b connection-loss synthesis (A+B: Monitors::teardown +
unread-command Disconnected); Phase 4 gate
c14 7e822b7 eager pg eviction (reaper actor, ReaperInboxes)
dbe1a22 InboundVerdict::label(), trace::Event::ClusterInbound
31a9877 tests/channel.rs monitor-churn target gated on `go`
653559e Discovery::Withdrawn{name, addr}
c15 b41d76e distributed pg: Sync on NodeUp, Join/Leave broadcast,
NodeDown sweep, members_all; PgMsg wire type
c16 fafa881 pick_any / dispatch_any; Phase 5 complete
Phase 6 Tier A:
195c73e p4 NodeEvent::NodeDown(NodeInfo)
48fd766 p1 connector Candidate{name, addr, state}
ce8cf99 p2+p7 conn.rs select arms as Vec<Arm>; Outbound::Drained
bf24988 p6 RemotePid::from_local -> Option
9ae0380 p3 PeerStanding{Free, Claimed, Dialing}
c7d62a1 p11 cluster::Timing knobs, threaded by value
46f171d p11 cluster_disconnect un-ignored on SMARM_FAST_TIMING
Phase 6 Tier B:
391a9ae p5 cluster::RemoteDownReason{Local, Disconnected};
DownReason::Disconnected removed from core
7ddd908 p9 pg ctl channel unconditional, one cfg seam at spawn
d9c62a8 PeerNameMismatch parks the candidate; ClusterDial trace
Verified at d9c62a8: default 361/0, cluster 448/0, clippy --lib on
default / cluster / cluster+smarm-trace, fmt, 10x flake on
cluster_dial_mismatch, 5x on cluster_pg.
483 lines
17 KiB
Rust
483 lines
17 KiB
Rust
//! 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<T>(rx: &mpsc::Receiver<T>, 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<String> = 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();
|
|
});
|
|
}
|