Drive the c5 machines as straight-line code on the path (D8): dial_handshake and accept_handshake do the IO on a shared FramedConn, and a connection actor is spawned only after a successful handshake. Rejects, tie-break losses (D7), protocol faults and timeouts are all resolved on the path by closing, so no actor ever exists for a connection that did not establish. The whole FramedConn travels into spawn_established, carrying any read-ahead past the handshake frames. Handshake deadlines land here rather than in c6c: FramedConn::recv_deadline enforces them between reads via the connection's fd arm, so a peer that connects and goes silent cannot wedge the acceptor. Connection lifetime moves to the manager (pulled forward from c7). The path registers each established connection and hands over its ConnHandle; the manager owns it, monitors the actor, and tears the connection down on Disconnect, on peer close, or at manager shutdown. spawn_established returns a Pid, so a connection neither outlives nor dies with whichever actor established it — the ownership that made two-node teardown unorderable. The manager also tracks in-flight dial intents, monitored so a panicking dial cannot wedge the tie-break, and answers HelloCtx for the accept path.
115 lines
4.3 KiB
Rust
115 lines
4.3 KiB
Rust
//! RFC 010 c6a — connection-actor lifecycle against the manager table.
|
|
//!
|
|
//! The handshake is bypassed here (c6b wires it): each connection is
|
|
//! constructed already-established over a real localhost TCP pair, handed a
|
|
//! fabricated `Peer`, and spawned. `spawn_established` registers it with the
|
|
//! manager, which takes its handle and monitors it, so the table reflects the
|
|
//! connection while it lives and reaps it on any exit path. This proves three
|
|
//! things at once: a live connection shows up, a commanded `Disconnect`
|
|
//! removes exactly that one, and a peer close (EOF, no command) removes the
|
|
//! other.
|
|
//!
|
|
//! TCP parks the calling actor, so everything runs inside `smarm::run`; the
|
|
//! single-threaded runtime is fine because every wait is a cooperative fd park.
|
|
#![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::spawn_established;
|
|
use smarm::cluster::transport::tcp::TcpTransport;
|
|
use smarm::cluster::transport::{Conn, FramedConn, Transport};
|
|
use smarm::gen_server::{self, GenServerBuilder};
|
|
use smarm::pg::Incarnation;
|
|
use smarm::{run, sleep};
|
|
|
|
/// A fabricated post-handshake peer identity. Only `node_name` matters to the
|
|
/// manager table; the rest is filler until c7 consumes it.
|
|
fn peer(name: &str) -> Peer {
|
|
Peer {
|
|
node_name: name.to_string(),
|
|
incarnation: Incarnation::new(1),
|
|
meta: NodeMeta {
|
|
role: "test".to_string(),
|
|
region: "test".to_string(),
|
|
},
|
|
}
|
|
}
|
|
|
|
/// One established transport pair over localhost. Relies on TCP backlog so the
|
|
/// sequential dial-then-accept needs no concurrent acceptor (same assumption as
|
|
/// the c3 conformance suite).
|
|
fn pair(t: &dyn Transport) -> (Box<dyn Conn>, Box<dyn Conn>) {
|
|
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)
|
|
}
|
|
|
|
/// Poll the manager until its peer set matches `expected` (sorted), or fail.
|
|
/// The bound is generous against a sub-millisecond real cost.
|
|
fn wait_peers(expected: &[&str]) {
|
|
let want: Vec<String> = expected.iter().map(|s| s.to_string()).collect();
|
|
for _ in 0..2000 {
|
|
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 connection_up_commanded_shutdown_and_eof_all_reflected_in_table() {
|
|
run(|| {
|
|
// The manager, started plainly and reachable at its well-known name.
|
|
// (The supervised subtree in `cluster::start` is permanent by design;
|
|
// a plainly-started manager lets this test terminate cleanly.)
|
|
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);
|
|
|
|
// Manage the `a` ends as peers node-b and node-c; keep the `b` far ends
|
|
// open so neither socket is closed from the far side yet.
|
|
spawn_established(FramedConn::new(a1), peer("node-b")).expect("node-b registers");
|
|
spawn_established(FramedConn::new(a2), peer("node-c")).expect("node-c registers");
|
|
|
|
// Up: both connections register and the table shows them.
|
|
wait_peers(&["node-b", "node-c"]);
|
|
|
|
// A commanded disconnect reaps exactly its own connection: the
|
|
// manager drops that entry's handle and the actor stops.
|
|
assert!(matches!(
|
|
gen_server::call(
|
|
MANAGER,
|
|
Call::Disconnect {
|
|
name: "node-b".to_string()
|
|
}
|
|
),
|
|
Ok(Reply::Disconnected)
|
|
));
|
|
wait_peers(&["node-c"]);
|
|
|
|
// A peer close (EOF) reaps the other with no command at all.
|
|
drop(b2);
|
|
wait_peers(&[]);
|
|
|
|
// node-b's far end stayed open until here, so its removal above was the
|
|
// disconnect command and not an EOF.
|
|
drop(b1);
|
|
|
|
// All connection actors have exited; stop the manager so `run` returns.
|
|
mgr.shutdown();
|
|
});
|
|
}
|