diff --git a/src/cluster.rs b/src/cluster.rs index f5bfffc..b3cea4f 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -4,10 +4,12 @@ //! transport trait (control connection), framed codec, and the TCP + //! loopback impls. c5: the handshake state machine. c6: the connection //! [`manager`] (registry) and per-peer connection actors ([`conn`]), started -//! as an explicit supervision subtree. Everything above them lands in later +//! as an explicit supervision subtree, plus the handshake on the +//! accept/connect path ([`connect`]). Everything above them lands in later //! chunks. pub mod conn; +pub mod connect; pub mod envelope; pub mod handshake; pub mod manager; @@ -21,6 +23,7 @@ use crate::scheduler::{sleep, spawn, JoinHandle}; use crate::supervisor::{ChildSpec, OneForOne, Restart}; pub use conn::{spawn_established, ConnHandle}; +pub use connect::{dial, spawn_acceptor, AcceptorHandle}; pub use manager::{Manager, MANAGER}; /// A running cluster subtree: an explicitly-started supervisor over the diff --git a/src/cluster/conn.rs b/src/cluster/conn.rs index 78122e7..43d99cb 100644 --- a/src/cluster/conn.rs +++ b/src/cluster/conn.rs @@ -10,60 +10,85 @@ //! so no read/write split is needed. //! //! The handshake completes *before* this actor exists (on the accept/connect -//! path — c6b) and produces the [`Peer`]; the actor registers that peer with -//! the [`manager`](crate::cluster::manager), which monitors it so any exit -//! deregisters the connection. Heartbeat send and fixed-timeout liveness join -//! the loop in c6c (the timeout arm of the same `select`). +//! path — c6b) and produces the [`Peer`]; the *path* then registers the +//! connection with the [`manager`](crate::cluster::manager), which takes +//! ownership of its [`ConnHandle`] and monitors the actor, so any exit +//! deregisters the connection. The actor itself holds no authority over its +//! own lifetime: it runs until the manager drops its handle (deregistration, +//! `Disconnect`, or manager shutdown) or the connection ends. Heartbeat send +//! and fixed-timeout liveness join the loop in c6c (the timeout arm of the +//! same `select`). use crate::channel::{channel, Receiver, Selectable, Sender}; use crate::cluster::handshake::Peer; use crate::cluster::manager::{Call, Registered, Reply, MANAGER}; use crate::cluster::transport::FramedConn; use crate::gen_server; -use crate::scheduler::{self, spawn}; +use crate::pid::Pid; +use crate::scheduler::spawn; /// Commands to a running connection actor. enum Cmd { Shutdown, } -/// A handle to a running connection actor. +/// The manager's authority over one connection actor: while this handle +/// lives the connection lives, and dropping it stops the actor and closes +/// the socket. Only the [`manager`](crate::cluster::manager) holds one — +/// callers of [`spawn_established`] get a [`Pid`] and no lifetime authority, +/// so a connection can never outlive, or die with, whichever actor happened +/// to establish it. pub struct ConnHandle { cmd_tx: Sender, } +impl std::fmt::Debug for ConnHandle { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.write_str("ConnHandle") + } +} + impl ConnHandle { /// Ask the connection to close and exit. Idempotent, and a no-op if the - /// actor has already gone. + /// actor has already gone. Dropping the handle does the same thing; this + /// exists for the manager's explicit `Disconnect` path. pub fn shutdown(&self) { let _ = self.cmd_tx.send(Cmd::Shutdown); } } -/// Spawn a connection actor for an **already-established** connection: the -/// handshake has completed elsewhere and produced `peer`. Returns as soon as -/// the actor is spawned; the actor's first act is to register with the manager. -pub fn spawn_established(framed: FramedConn, peer: Peer) -> ConnHandle { - let (cmd_tx, cmd_rx) = channel(); - spawn(move || run(framed, peer, cmd_rx)); - ConnHandle { cmd_tx } -} +/// The name was already claimed by a live connection, so this one was +/// refused; its actor has been stopped and its socket closed. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct RegisterRefused; -fn run(mut framed: FramedConn, peer: Peer, cmd_rx: Receiver) { - let me = scheduler::self_pid(); +/// Spawn a connection actor for an **already-established** connection (the +/// handshake completed on the path and produced `peer`) and register it with +/// the manager, synchronously, before returning. The manager takes the +/// actor's [`ConnHandle`]; the caller gets only the [`Pid`], because +/// connection lifetime belongs to the table and not to the establishing +/// actor. A refusal has already stopped the actor and closed the socket. +pub fn spawn_established(framed: FramedConn, peer: Peer) -> Result { + let (cmd_tx, cmd_rx) = channel(); + let name = peer.node_name.clone(); + let pid = spawn(move || run(framed, peer, cmd_rx)).pid(); match gen_server::call( MANAGER, Call::Register { - name: peer.node_name.clone(), - pid: me, + name, + pid, + handle: ConnHandle { cmd_tx }, }, ) { - Ok(Reply::Registered(Registered::Ok)) => {} - // Duplicate name, or the manager is unreachable: do not run. The - // connection drops (closing the socket) as `framed` falls out of scope. - _ => return, + Ok(Reply::Registered(Registered::Ok)) => Ok(pid), + // Duplicate name, or the manager is unreachable. Either way the + // handle went with the call and is dropped there (or never arrived + // and dropped with it), which stops the actor and closes the socket. + _ => Err(RegisterRefused), } +} +fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver) { loop { match framed.readable_arm() { Some(arm) => { diff --git a/src/cluster/connect.rs b/src/cluster/connect.rs new file mode 100644 index 0000000..503b37c --- /dev/null +++ b/src/cluster/connect.rs @@ -0,0 +1,345 @@ +//! RFC 010 c6b — the handshake on the accept/connect path. +//! +//! Per D8 (re-amended): the c5 machines are driven by **straight-line code +//! on the path**, not by an actor. The dial side runs [`Initiator`]; the +//! acceptor loop runs [`Responder`]. A connection actor is spawned only +//! *after* a successful handshake ([`spawn_established`]); every reject, +//! protocol failure, timeout, and tie-break loss is resolved right here, +//! on the path, by closing — no actor ever exists for a connection that +//! didn't establish. +//! +//! Buffer trap (binding): the path reader and the steady-state actor share +//! ONE [`FramedConn`]. Its decode buffer may hold read-ahead past the +//! handshake frames, so the *whole* `FramedConn` travels into +//! [`spawn_established`] — never a fresh codec over the same socket. +//! +//! Layering: [`dial_handshake`] and [`accept_handshake`] are the bare path +//! steps — IO on a `FramedConn`, no manager, no actors — testable over the +//! loopback transport on plain threads. [`dial`] and [`spawn_acceptor`] are +//! the manager-integrated layer (actor context required): they keep the +//! [`manager`](crate::cluster::manager)'s dial-intent set honest and spawn +//! the connection actor on success. + +use std::io; +use std::time::{Duration, Instant}; + +use crate::channel::{channel, Receiver, Selectable, Sender}; +use crate::cluster::conn::spawn_established; +use crate::cluster::envelope::{Frame, RejectReason}; +use crate::cluster::handshake::{ + HelloCtx, Initiator, InitiatorOutcome, Local, Peer, Responder, ResponderOutcome, +}; +use crate::cluster::manager::{Call, Reply, MANAGER}; +use crate::cluster::transport::{FramedConn, Listener, RecvError, SendError, Transport}; +use crate::gen_server; +use crate::pid::Pid; +use crate::scheduler::{self, spawn}; + +/// How long either side waits for the peer's handshake frame before giving +/// up and closing. Enforced on the path via [`FramedConn::recv_deadline`], +/// so a peer that connects and goes silent cannot wedge the acceptor. +pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5); + +/// Why a handshake did not establish. In every case the connection has +/// already been closed on the path by the time this is returned. +#[derive(Debug)] +pub enum HandshakeError { + /// A `HelloReject` travelled — sent by us (accept side) or received by + /// us (dial side). + Rejected(RejectReason), + /// Accept side only: the inbound dial lost the simultaneous-connect + /// tie-break (D7) and was closed silently, no frame sent. + TieBreakLoss, + /// The peer spoke a valid frame that is wrong here (non-`Hello` first + /// frame; non-response to our `Hello`), or an undecodable byte stream. + Protocol, + /// EOF before the handshake resolved. On the dial side this is also + /// what losing the tie-break looks like: the peer closes silently. + Closed, + /// [`HANDSHAKE_TIMEOUT`] (or the caller's deadline) passed first. + TimedOut, + /// The transport failed mid-handshake. + Transport(io::Error), +} + +fn from_send(e: SendError) -> HandshakeError { + match e { + // Handshake frames are small and self-made; an encode failure is a + // protocol-level impossibility, not a transport fault. + SendError::Encode(_) => HandshakeError::Protocol, + SendError::Io(e) => HandshakeError::Transport(e), + } +} + +fn from_recv(e: RecvError) -> HandshakeError { + match e { + RecvError::Corrupt(_) => HandshakeError::Protocol, + RecvError::TruncatedByPeer => HandshakeError::Closed, + RecvError::Io(e) => HandshakeError::Transport(e), + RecvError::TimedOut => HandshakeError::TimedOut, + } +} + +/// Dial-side path step: send our `Hello`, interpret the one response. On +/// `Ok` the connection is established and `framed` is live (with any +/// read-ahead intact in its buffer); on `Err` the connection is closed. +pub fn dial_handshake( + framed: &mut FramedConn, + local: &Local, + deadline: Instant, +) -> Result { + let (initiator, hello) = Initiator::new(local); + if let Err(e) = framed.send(&hello) { + framed.close(); + return Err(from_send(e)); + } + let outcome = match framed.recv_deadline(deadline) { + Ok(Some(frame)) => initiator.on_frame(frame), + Ok(None) => { + framed.close(); + return Err(HandshakeError::Closed); + } + Err(e) => { + framed.close(); + return Err(from_recv(e)); + } + }; + match outcome { + InitiatorOutcome::Established(peer) => Ok(peer), + InitiatorOutcome::Rejected(reason) => { + framed.close(); + Err(HandshakeError::Rejected(reason)) + } + InitiatorOutcome::Failed(_) => { + framed.close(); + Err(HandshakeError::Protocol) + } + } +} + +/// Accept-side path step: read the first frame, judge it, answer or close. +/// +/// `ctx_for` supplies the [`HelloCtx`] for the *offered* name — knowledge +/// only the frame reveals, which is why it is a callback and not a value +/// (the integrated acceptor asks the manager; loopback tests fabricate). +/// It is not called when the first frame is not a `Hello`. +/// +/// On `Ok` the ack has been sent and `framed` is live (read-ahead intact); +/// on `Err` any owed reject has been sent and the connection is closed. +pub fn accept_handshake( + framed: &mut FramedConn, + local: Local, + ctx_for: impl FnOnce(&str) -> HelloCtx, + deadline: Instant, +) -> Result { + let frame = match framed.recv_deadline(deadline) { + Ok(Some(frame)) => frame, + Ok(None) => { + framed.close(); + return Err(HandshakeError::Closed); + } + Err(e) => { + framed.close(); + return Err(from_recv(e)); + } + }; + let ctx = match &frame { + Frame::Hello { node_name, .. } => ctx_for(node_name), + _ => HelloCtx::default(), + }; + match Responder::new(local).on_frame(frame, ctx) { + ResponderOutcome::Accepted { reply, peer } => { + if let Err(e) = framed.send(&reply) { + framed.close(); + return Err(from_send(e)); + } + Ok(peer) + } + ResponderOutcome::Rejected { reply, reason } => { + // Best effort: the reject is the cross-version compatibility + // anchor, but if the write fails the peer sees a bare close, + // which it must survive anyway. + let _ = framed.send(&reply); + framed.close(); + Err(HandshakeError::Rejected(reason)) + } + ResponderOutcome::TieBreakLoss => { + // D7: close silently — the peer computes the same verdict. + framed.close(); + Err(HandshakeError::TieBreakLoss) + } + ResponderOutcome::Failed(_) => { + framed.close(); + Err(HandshakeError::Protocol) + } + } +} + +// --------------------------------------------------------------------------- +// Manager-integrated layer +// --------------------------------------------------------------------------- + +/// Why an integrated [`dial`] did not produce a connection. +#[derive(Debug)] +pub enum DialError { + /// Another dial to this peer name is already in flight. + AlreadyDialing, + /// The manager is not running (or answered nonsense). + ManagerUnavailable, + /// The transport could not connect. + Connect(io::Error), + /// Connected, but the handshake did not establish. + Handshake(HandshakeError), + /// The peer at `addr` established, but answered as a different name + /// than the one we dialed — the tie-break bookkeeping (keyed by the + /// dialed name) would be unsound, so the connection is closed. + PeerNameMismatch { expected: String, got: String }, + /// The handshake established, but the manager refused the registration: + /// a connection to this peer already exists. The loser has been closed. + Duplicate, +} + +/// Dial `peer_name` at `addr` and run the handshake, keeping the manager's +/// dial-intent set honest around it: the intent is registered *before* +/// connecting (so a crossing inbound `Hello` sees it) and cleared the +/// moment the handshake resolves, before the connection actor is spawned. +/// Must run inside an actor. Retrying is the caller's business (c7's dial +/// loop); a lost tie-break surfaces as `Handshake(Closed)` — the peer's +/// accepted connection is already on its way. +pub fn dial( + transport: &dyn Transport, + addr: &str, + peer_name: &str, + local: &Local, +) -> Result { + let me = scheduler::self_pid(); + match gen_server::call( + MANAGER, + Call::DialBegin { + name: peer_name.to_string(), + pid: me, + }, + ) { + Ok(Reply::DialBegan(true)) => {} + Ok(Reply::DialBegan(false)) => return Err(DialError::AlreadyDialing), + _ => return Err(DialError::ManagerUnavailable), + } + let result = connect_and_shake(transport, addr, local); + // Cleared immediately on outcome — a stale intent during the established + // window would corrupt later tie-breaks. Synchronous (a call): the + // intent is provably gone before anything else happens. + let _ = gen_server::call( + MANAGER, + Call::DialEnd { + name: peer_name.to_string(), + }, + ); + let (mut framed, peer) = result?; + if peer.node_name != peer_name { + framed.close(); + return Err(DialError::PeerNameMismatch { + expected: peer_name.to_string(), + got: peer.node_name, + }); + } + spawn_established(framed, peer).map_err(|_| DialError::Duplicate) +} + +fn connect_and_shake( + transport: &dyn Transport, + addr: &str, + local: &Local, +) -> Result<(FramedConn, Peer), DialError> { + let conn = transport.dial(addr).map_err(DialError::Connect)?; + let mut framed = FramedConn::new(conn); + let deadline = Instant::now() + HANDSHAKE_TIMEOUT; + let peer = dial_handshake(&mut framed, local, deadline).map_err(DialError::Handshake)?; + Ok((framed, peer)) +} + +/// A running acceptor. [`shutdown`](AcceptorHandle::shutdown) (or dropping +/// the last handle) stops the accept loop only: connections it established +/// belong to the [`manager`](crate::cluster::manager) and keep running, to +/// be torn down through the table (`Disconnect`, a peer close, or manager +/// shutdown). +pub struct AcceptorHandle { + cmd_tx: Sender<()>, + addr: String, +} + +impl AcceptorHandle { + /// Ask the acceptor to stop. Idempotent; a no-op if it already has. + pub fn shutdown(&self) { + let _ = self.cmd_tx.send(()); + } + + /// The concrete bound address, dialable as-is. + pub fn local_addr(&self) -> &str { + &self.addr + } +} + +/// Spawn the acceptor actor over a bound listener. Each inbound connection +/// is handshaken **inline in the loop** (a deliberate serialization: the +/// per-frame deadline bounds how long any one peer can hold the line, and +/// nothing concurrent exists to be starved before c7). The listener must be +/// fd-backed ([`Listener::readable_arm`]); the loopback listener is not, +/// and its acceptor exits immediately — loopback handshakes are driven +/// synchronously through the path fns instead, per D8. +pub fn spawn_acceptor(listener: Box, local: Local) -> AcceptorHandle { + let addr = listener.local_addr(); + let (cmd_tx, cmd_rx) = channel(); + spawn(move || accept_loop(listener, local, cmd_rx)); + AcceptorHandle { cmd_tx, addr } +} + +fn accept_loop(mut listener: Box, local: Local, cmd_rx: Receiver<()>) { + loop { + let Some(arm) = listener.readable_arm() else { + return; + }; + let arms: [&dyn Selectable; 2] = [&cmd_rx, &arm]; + match crate::channel::try_select(&arms) { + Ok(0) => match cmd_rx.try_recv() { + Ok(Some(())) => return, + Ok(None) => continue, // spurious wake + Err(_) => return, // all handles dropped + }, + Ok(_) => { + // The listener is readable: accept completes without parking. + let conn = match listener.accept() { + Ok(conn) => conn, + Err(_) => return, // listener itself is broken + }; + handle_inbound(FramedConn::new(conn), &local); + } + Err(_) => return, // fd arm failed to register: listener is gone + } + } +} + +/// Run the accept-side handshake for one inbound connection, asking the +/// manager for the [`HelloCtx`], and hand the established connection to the +/// manager. Every failure was already resolved on the path (reject sent / +/// closed, or the registration refused and the actor stopped), so there is +/// nothing for the acceptor to carry forward. +fn handle_inbound(mut framed: FramedConn, local: &Local) { + let deadline = Instant::now() + HANDSHAKE_TIMEOUT; + let ctx_for = |name: &str| match gen_server::call( + MANAGER, + Call::HelloCtx { + peer_name: name.to_string(), + }, + ) { + Ok(Reply::HelloCtx(ctx)) => ctx, + // Manager unreachable: nobody could register this connection anyway, + // so claim the name taken and reject rather than accept an orphan. + _ => HelloCtx { + name_claimed: true, + dialing_this_peer: false, + }, + }; + if let Ok(peer) = accept_handshake(&mut framed, local.clone(), ctx_for, deadline) { + let _ = spawn_established(framed, peer); + } +} diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 9ab3218..b5153b1 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -1,11 +1,14 @@ //! RFC 010 c6 — the cluster connection manager. //! //! One manager per runtime: the single registry of live peer connections and -//! the source of truth for whether a peer name is already claimed. Every -//! connection actor registers here as its first act and is *monitored* by the -//! manager, so the table self-heals on any exit path — a connection that -//! panics, is cancelled, or closes cleanly is removed without cooperation from -//! the dying actor. +//! the source of truth for whether a peer name is already claimed. The +//! accept/connect path registers each established connection here, handing +//! over its [`ConnHandle`] — **the manager owns connection lifetime**. A +//! connection lives as long as its table entry, so it neither outlives nor +//! dies with whichever actor happened to establish it. Registered actors are +//! also *monitored*, so the table self-heals on any exit path — a connection +//! that panics, is cancelled, or closes cleanly is removed without +//! cooperation from the dying actor. //! //! Membership as consumers will see it (the `node_up`/`node_down` interface and //! the view) and the connector dial loop are c7, built on top of this table. @@ -14,6 +17,8 @@ use std::collections::HashMap; +use crate::cluster::conn::ConnHandle; +use crate::cluster::handshake::HelloCtx; use crate::gen_server::{GenServer, GenServerCtx, GenServerName, Watcher}; use crate::monitor::{monitor, Down}; use crate::pid::Pid; @@ -23,10 +28,23 @@ use crate::pid::Pid; /// manager is always found at the same key. pub const MANAGER: GenServerName = GenServerName::new("smarm.cluster.manager"); +/// One live connection's entry: the actor running it, and the handle whose +/// lifetime *is* the connection's (see the module docs). +struct ConnEntry { + pid: Pid, + _handle: ConnHandle, +} + /// The connection registry: peer name → the connection actor that owns that /// peer's control connection. pub struct Manager { - conns: HashMap, + conns: HashMap, + /// In-flight dial intents: peer name -> the actor performing the dial. + /// Registered *before* connecting so a crossing inbound `Hello` sees it + /// (the `dialing_this_peer` half of [`HelloCtx`]); cleared the moment + /// the dial resolves, and — because the dialer is monitored — on the + /// dialer's death, so a panicking dial can never wedge the tie-break. + dials: HashMap, watcher: Option>, } @@ -34,6 +52,7 @@ impl Manager { pub fn new() -> Self { Manager { conns: HashMap::new(), + dials: HashMap::new(), watcher: None, } } @@ -57,18 +76,43 @@ pub enum Registered { /// Requests to the manager. pub enum Call { - /// A freshly-established connection actor claims its peer's name. The pid - /// is the calling connection actor, which the manager then monitors. - Register { name: String, pid: Pid }, + /// The path claims its peer's name for a freshly-established connection, + /// handing the manager the actor's [`ConnHandle`]. The manager monitors + /// `pid` and holds the handle for as long as the entry lives; a + /// [`Registered::Duplicate`] verdict drops the handle here, which stops + /// the refused actor. + Register { + name: String, + pid: Pid, + handle: ConnHandle, + }, + /// Tear down the connection to `name`: the manager drops its handle, the + /// actor stops, and the monitor removes the entry. A no-op if no such + /// connection is live. + Disconnect { name: String }, /// The current peer names, sorted. For observation and tests. Peers, + /// A dialer declares an in-flight dial to `name` before connecting. The + /// pid is the dialing actor, monitored so the intent dies with it. + DialBegin { name: String, pid: Pid }, + /// The dial to `name` resolved (either way): drop the intent. A call, + /// not a cast, so the intent is provably gone before the dialer moves on. + DialEnd { name: String }, + /// The [`HelloCtx`] for an inbound `Hello` offering `peer_name` — the + /// accept path asks this between reading the frame and judging it. + HelloCtx { peer_name: String }, } /// Replies from the manager. #[derive(Debug)] pub enum Reply { Registered(Registered), + Disconnected, Peers(Vec), + /// `false`: another dial to this name is already in flight — do not dial. + DialBegan(bool), + DialEnded, + HelloCtx(HelloCtx), } impl GenServer for Manager { @@ -84,27 +128,58 @@ impl GenServer for Manager { fn handle_call(&mut self, request: Call) -> Reply { match request { - Call::Register { name, pid } => { + Call::Register { name, pid, handle } => { if self.conns.contains_key(&name) { + // `handle` drops here: the refused actor stops itself. return Reply::Registered(Registered::Duplicate); } if let Some(w) = &self.watcher { w.watch(monitor(pid)); } - self.conns.insert(name, pid); + self.conns.insert( + name, + ConnEntry { + pid, + _handle: handle, + }, + ); Reply::Registered(Registered::Ok) } + Call::Disconnect { name } => { + // Dropping the entry drops the handle, which stops the actor. + self.conns.remove(&name); + Reply::Disconnected + } Call::Peers => { let mut names: Vec = self.conns.keys().cloned().collect(); names.sort(); Reply::Peers(names) } + Call::DialBegin { name, pid } => { + if self.dials.contains_key(&name) { + return Reply::DialBegan(false); + } + if let Some(w) = &self.watcher { + w.watch(monitor(pid)); + } + self.dials.insert(name, pid); + Reply::DialBegan(true) + } + Call::DialEnd { name } => { + self.dials.remove(&name); + Reply::DialEnded + } + Call::HelloCtx { peer_name } => Reply::HelloCtx(HelloCtx { + name_claimed: self.conns.contains_key(&peer_name), + dialing_this_peer: self.dials.contains_key(&peer_name), + }), } } fn handle_cast(&mut self, _request: ()) {} fn handle_down(&mut self, down: Down) { - self.conns.retain(|_, pid| *pid != down.pid); + self.conns.retain(|_, entry| entry.pid != down.pid); + self.dials.retain(|_, pid| *pid != down.pid); } } diff --git a/src/cluster/transport.rs b/src/cluster/transport.rs index d62d786..5564c47 100644 --- a/src/cluster/transport.rs +++ b/src/cluster/transport.rs @@ -71,6 +71,15 @@ pub trait Listener: Send { /// The concrete bound address, dialable as-is (e.g. the real port when /// bound with port 0). fn local_addr(&self) -> String; + + /// Readiness as a [`select`](crate::select) arm, mirroring + /// [`Conn::readable_arm`]: `Some` lets an acceptor wait on "an inbound + /// connection is pending" alongside a command inbox in one `select`, so + /// it can be told to stop without a poll loop. Default `None` (the + /// loopback listener has no fd and must be driven synchronously). + fn readable_arm(&self) -> Option { + None + } } /// A way of establishing control connections. Object-safe on purpose: the @@ -128,6 +137,10 @@ pub enum RecvError { TruncatedByPeer, /// The transport failed mid-read. Io(io::Error), + /// The deadline passed before a full frame arrived + /// ([`FramedConn::recv_deadline`] only; plain [`recv`](FramedConn::recv) + /// never returns this). + TimedOut, } impl std::fmt::Display for RecvError { @@ -136,6 +149,7 @@ impl std::fmt::Display for RecvError { RecvError::Corrupt(e) => write!(f, "frame stream corrupt: {e:?}"), RecvError::TruncatedByPeer => write!(f, "peer closed mid-frame"), RecvError::Io(e) => write!(f, "transport read failed: {e}"), + RecvError::TimedOut => write!(f, "deadline passed mid-receive"), } } } @@ -193,6 +207,49 @@ impl FramedConn { } } + /// Like [`recv`](FramedConn::recv), but gives up with + /// [`RecvError::TimedOut`] once `deadline` passes without a full frame. + /// The deadline is enforced between reads via the connection's fd arm + /// (so the caller must be an actor); a transport with no fd (loopback) + /// cannot be timed out and this degrades to a plain blocking `recv` — + /// the same caveat as liveness. + pub fn recv_deadline( + &mut self, + deadline: std::time::Instant, + ) -> Result, RecvError> { + loop { + match Frame::decode(&self.rbuf) { + Ok(Some((frame, consumed))) => { + self.rbuf.drain(..consumed); + return Ok(Some(frame)); + } + Ok(None) => {} + Err(e) => return Err(RecvError::Corrupt(e)), + } + if let Some(arm) = self.conn.readable_arm() { + let left = deadline.saturating_duration_since(std::time::Instant::now()); + if left.is_zero() { + return Err(RecvError::TimedOut); + } + match crate::channel::try_select_timeout(&[&arm], left) { + Ok(Some(_)) => {} + Ok(None) => return Err(RecvError::TimedOut), + Err(e) => return Err(RecvError::Io(e)), + } + } + let mut chunk = [0u8; READ_CHUNK]; + let n = self.conn.read(&mut chunk).map_err(RecvError::Io)?; + if n == 0 { + return if self.rbuf.is_empty() { + Ok(None) + } else { + Err(RecvError::TruncatedByPeer) + }; + } + self.rbuf.extend_from_slice(&chunk[..n]); + } + } + /// Close the underlying connection (idempotent, see [`Conn::close`]). pub fn close(&mut self) { self.conn.close(); diff --git a/src/cluster/transport/tcp.rs b/src/cluster/transport/tcp.rs index d3764e6..2fdc421 100644 --- a/src/cluster/transport/tcp.rs +++ b/src/cluster/transport/tcp.rs @@ -197,6 +197,10 @@ pub struct TcpListener { } impl Listener for TcpListener { + fn readable_arm(&self) -> Option { + Some(crate::scheduler::FdArm::readable(self.inner.as_raw_fd())) + } + fn accept(&mut self) -> io::Result> { loop { wait_readable(self.inner.as_raw_fd())?; diff --git a/tests/cluster_conn_lifecycle.rs b/tests/cluster_conn_lifecycle.rs index 040afbc..eff1b93 100644 --- a/tests/cluster_conn_lifecycle.rs +++ b/tests/cluster_conn_lifecycle.rs @@ -2,11 +2,12 @@ //! //! 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. The actor registers with the manager, which -//! 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 shutdown removes exactly that one, and a peer close -//! (EOF, no command) removes the other. +//! 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. @@ -80,14 +81,23 @@ fn connection_up_commanded_shutdown_and_eof_all_reflected_in_table() { // 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. - let h1 = spawn_established(FramedConn::new(a1), peer("node-b")); - let _h2 = spawn_established(FramedConn::new(a2), peer("node-c")); + 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 shutdown reaps exactly its own connection. - h1.shutdown(); + // 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. @@ -95,7 +105,7 @@ fn connection_up_commanded_shutdown_and_eof_all_reflected_in_table() { wait_peers(&[]); // node-b's far end stayed open until here, so its removal above was the - // shutdown command and not an EOF. + // disconnect command and not an EOF. drop(b1); // All connection actors have exited; stop the manager so `run` returns. diff --git a/tests/cluster_connect.rs b/tests/cluster_connect.rs new file mode 100644 index 0000000..405ddc4 --- /dev/null +++ b/tests/cluster_connect.rs @@ -0,0 +1,472 @@ +//! 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::{HelloCtx, Local}; +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::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"); + HelloCtx::default() + }, + 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, |_| HelloCtx::default(), 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"), + |_| HelloCtx { + name_claimed: false, + dialing_this_peer: true, + }, + 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"), + |_| HelloCtx::default(), + 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"), + |_| HelloCtx::default(), + 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")); + 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::HelloCtx { + peer_name: "ghost".into(), + }, + ) { + Ok(Reply::HelloCtx(ctx)) => assert!(ctx.dialing_this_peer), + other => panic!("HelloCtx 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::HelloCtx { + peer_name: "ghost".into(), + }, + ) { + Ok(Reply::HelloCtx(ctx)) if !ctx.dialing_this_peer => 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")); + 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")); + 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")); + 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(); + }); +}