diff --git a/src/cluster/conn.rs b/src/cluster/conn.rs index 43d99cb..3eee6c5 100644 --- a/src/cluster/conn.rs +++ b/src/cluster/conn.rs @@ -15,11 +15,17 @@ //! 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`). +//! `Disconnect`, or manager shutdown), the connection ends, or liveness +//! expires. Heartbeat send and fixed-timeout liveness are the timeout arm of +//! the same `select` (c6c): [`HEARTBEAT_INTERVAL`] paces outbound +//! [`Frame::Heartbeat`](crate::cluster::envelope::Frame::Heartbeat)s, and a +//! [`LIVENESS_TIMEOUT`] window — reset by any inbound frame — tears the +//! connection down when it empties. -use crate::channel::{channel, Receiver, Selectable, Sender}; +use std::time::{Duration, Instant}; + +use crate::channel::{channel, try_select_timeout, Receiver, Selectable, Sender}; +use crate::cluster::envelope::Frame; use crate::cluster::handshake::Peer; use crate::cluster::manager::{Call, Registered, Reply, MANAGER}; use crate::cluster::transport::FramedConn; @@ -88,39 +94,81 @@ pub fn spawn_established(framed: FramedConn, peer: Peer) -> Result) { + match framed.readable_arm() { + Some(arm) => run_live(&mut framed, arm, &cmd_rx), + None => run_inert(&cmd_rx), + } + framed.close(); +} + +/// The steady-state loop over an fd-backed connection: one +/// `select_timeout` folds the command inbox, socket readability, and the +/// nearer of the two deadlines (`hb_send`, `liveness`) into a single wait. +fn run_live(framed: &mut FramedConn, arm: crate::scheduler::FdArm, cmd_rx: &Receiver) { + let mut next_hb = Instant::now(); + let mut live_until = Instant::now() + LIVENESS_TIMEOUT; loop { - match framed.readable_arm() { - Some(arm) => { - let arms: [&dyn Selectable; 2] = [&cmd_rx, &arm]; - match crate::channel::try_select(&arms) { - Ok(0) => { - if should_stop(&cmd_rx) { - break; - } - } - Ok(_) => { - if pump_readable(&mut framed) { - break; - } - } - // The fd arm failed to register — the connection is gone. - Err(_) => break, - } + let now = Instant::now(); + if now >= live_until { + break; // liveness expired: the peer is dead to us + } + if now >= next_hb { + if framed.send(&Frame::Heartbeat).is_err() { + break; } - None => { - // No fd to select on (loopback): only a command can end the - // wait. Liveness over such a transport is out of scope. - let arms: [&dyn Selectable; 1] = [&cmd_rx]; - let _ = crate::channel::select(&arms); - if should_stop(&cmd_rx) { + next_hb = now + HEARTBEAT_INTERVAL; + } + let wait = next_hb.min(live_until).saturating_duration_since(now); + let arms: [&dyn Selectable; 2] = [cmd_rx, &arm]; + match try_select_timeout(&arms, wait) { + Ok(Some(0)) => { + if should_stop(cmd_rx) { break; } } + Ok(Some(_)) => match pump_readable(framed) { + Pump::Ended => break, + Pump::Frames(n) => { + if n > 0 { + live_until = Instant::now() + LIVENESS_TIMEOUT; + } + } + }, + // A deadline passed; the top of the loop acts on whichever. + Ok(None) => {} + // The fd arm failed to register — the connection is gone. + Err(_) => break, } } +} - framed.close(); +/// No fd to select on (loopback): only a command can end the wait, and +/// neither heartbeats nor liveness run — a transport that can't report +/// readiness can't be timed either (same caveat as +/// [`FramedConn::recv_deadline`]). Loopback is a test transport; every real +/// connection is fd-backed. +fn run_inert(cmd_rx: &Receiver) { + loop { + let arms: [&dyn Selectable; 1] = [cmd_rx]; + let _ = crate::channel::select(&arms); + if should_stop(cmd_rx) { + break; + } + } } /// Drain the command arm. Returns `true` when the actor should exit — a @@ -133,13 +181,39 @@ fn should_stop(cmd_rx: &Receiver) -> bool { } } -/// Consume whatever is readable now. Returns `true` when the connection has -/// ended (clean EOF or an unrecoverable stream error). This chunk does not -/// interpret frames; c6c handles heartbeats and resets the liveness timer here. -fn pump_readable(framed: &mut FramedConn) -> bool { - match framed.recv() { - Ok(Some(_frame)) => false, - Ok(None) => true, // clean EOF at a frame boundary - Err(_) => true, // corrupt / truncated / io +/// What one readable wake yielded. +enum Pump { + /// The connection has ended: EOF (clean or mid-frame) or an + /// unrecoverable stream error. + Ended, + /// Still up; this many complete frames were consumed (possibly zero, if + /// the wake delivered only part of a frame). Any nonzero count resets + /// the liveness window. + Frames(usize), +} + +/// Consume one readable wake: exactly one socket read (which cannot block +/// after a level-triggered readable indication), then drain every complete +/// frame the buffer now holds. A blocking `recv` here would park the actor +/// past its heartbeat and liveness deadlines whenever a frame arrives split. +/// Frames are not interpreted yet — a heartbeat's entire job is the liveness +/// reset, and everything else waits for c8. +fn pump_readable(framed: &mut FramedConn) -> Pump { + let eof = match framed.read_once() { + Ok(n) => n == 0, + Err(_) => return Pump::Ended, + }; + let mut got = 0; + loop { + match framed.next_buffered() { + Ok(Some(_frame)) => got += 1, + Ok(None) => break, + Err(_) => return Pump::Ended, // corrupt stream + } + } + if eof { + Pump::Ended + } else { + Pump::Frames(got) } } diff --git a/src/cluster/transport.rs b/src/cluster/transport.rs index 5564c47..8f0e9fc 100644 --- a/src/cluster/transport.rs +++ b/src/cluster/transport.rs @@ -250,6 +250,36 @@ impl FramedConn { } } + /// One socket read, appended to the reassembly buffer. Returns the byte + /// count (`0` = EOF). For select-loop callers that were just told the fd + /// is readable: under the level-triggered IO thread exactly one read per + /// readable wake never blocks and never loses data — leftover socket + /// bytes re-signal on the next select, and complete frames already + /// reassembled are drained with [`next_buffered`](FramedConn::next_buffered). + /// (A plain [`recv`](FramedConn::recv) can block into the socket while + /// the buffer holds a partial frame, which a loop with deadlines to keep + /// cannot afford.) + pub fn read_once(&mut self) -> std::io::Result { + let mut chunk = [0u8; READ_CHUNK]; + let n = self.conn.read(&mut chunk)?; + self.rbuf.extend_from_slice(&chunk[..n]); + Ok(n) + } + + /// Decode the next complete frame already sitting in the reassembly + /// buffer, without touching the socket. `Ok(None)` means the buffer + /// holds no complete frame (empty, or a partial awaiting more bytes). + pub fn next_buffered(&mut self) -> Result, DecodeError> { + match Frame::decode(&self.rbuf) { + Ok(Some((frame, consumed))) => { + self.rbuf.drain(..consumed); + Ok(Some(frame)) + } + Ok(None) => Ok(None), + Err(e) => Err(e), + } + } + /// Close the underlying connection (idempotent, see [`Conn::close`]). pub fn close(&mut self) { self.conn.close(); diff --git a/tests/cluster_conn_liveness.rs b/tests/cluster_conn_liveness.rs new file mode 100644 index 0000000..ba2aea0 --- /dev/null +++ b/tests/cluster_conn_liveness.rs @@ -0,0 +1,178 @@ +//! RFC 010 c6c — heartbeat send + fixed-timeout liveness + teardown. +//! +//! Each case runs one real connection actor over an in-process localhost TCP +//! pair, with the far end held as a raw `FramedConn` (no actor) so the test +//! controls exactly what — if anything — the peer says. That gives the three +//! protocol-visible facts direct handles: heartbeats appear on the wire +//! unprompted; a mute peer is torn down (and reaped from the manager table) +//! once `LIVENESS_TIMEOUT` empties; and a peer that does nothing but send +//! heartbeats keeps the connection alive past that same window. +//! +//! Loopback has no fd and cannot drive liveness (documented on the actor), +//! so everything here is TCP. TCP parks the calling actor, so everything +//! runs inside `smarm::run`. +#![cfg(feature = "cluster")] + +use std::time::{Duration, Instant}; + +use smarm::cluster::conn::{HEARTBEAT_INTERVAL, LIVENESS_TIMEOUT}; +use smarm::cluster::envelope::{Frame, 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, spawn}; + +/// A fabricated post-handshake peer identity (same shape as the c6a suite). +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 (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) +} + +fn peers() -> Vec { + match gen_server::call(MANAGER, Call::Peers) { + Ok(Reply::Peers(p)) => p, + other => panic!("manager unreachable: {other:?}"), + } +} + +/// Poll until the manager's peer set matches `expected` (sorted) or `budget` +/// runs out. +fn wait_peers(expected: &[&str], budget: Duration) { + let want: Vec = expected.iter().map(|s| s.to_string()).collect(); + let deadline = Instant::now() + budget; + while Instant::now() < deadline { + if peers() == want { + return; + } + sleep(Duration::from_millis(10)); + } + panic!( + "timed out waiting for peers == {want:?}; last = {:?}", + peers() + ); +} + +/// The actor emits heartbeats unprompted: the raw far end, saying nothing, +/// sees a `Frame::Heartbeat` well within one interval (the first goes out at +/// spawn). +#[test] +fn heartbeats_are_sent_unprompted() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let (a, b) = pair(&TcpTransport); + spawn_established(FramedConn::new(a), peer("hb-send")).expect("register"); + let mut far = FramedConn::new(b); + + let frame = far + .recv_deadline(Instant::now() + HEARTBEAT_INTERVAL) + .expect("a heartbeat before one interval elapses"); + assert_eq!(frame, Some(Frame::Heartbeat)); + + // Teardown: closing the far end is an EOF at the actor. + far.close(); + wait_peers(&[], Duration::from_secs(2)); + mgr.shutdown(); + }); +} + +/// A mute peer is dead: no inbound frame for `LIVENESS_TIMEOUT` tears the +/// connection down and the manager's monitor reaps the table entry. The +/// entry is still present well inside the window — the teardown is the +/// timer, not an accident of setup. +#[test] +fn mute_peer_is_torn_down_after_liveness_timeout() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let (a, b) = pair(&TcpTransport); + spawn_established(FramedConn::new(a), peer("mute")).expect("register"); + // Held open and silent: no frames, no EOF. (Unread inbound + // heartbeats sit in kernel buffers; they are 5 bytes each.) + let _far = FramedConn::new(b); + + // Well inside the window the connection is still up. + sleep(LIVENESS_TIMEOUT / 2); + assert_eq!(peers(), vec!["mute".to_string()], "torn down too early"); + + // ...and once the window empties it is gone. Generous budget over + // the remaining half-window. + wait_peers(&[], LIVENESS_TIMEOUT); + mgr.shutdown(); + }); +} + +/// Heartbeats alone keep a connection alive past `LIVENESS_TIMEOUT`: a far +/// end that sends `Frame::Heartbeat` at the interval (and nothing else) +/// holds the entry; when it goes quiet, liveness finally fires. +#[test] +fn heartbeats_keep_the_connection_alive() { + run(|| { + let mgr = GenServerBuilder::new(Manager::new()) + .named(MANAGER) + .start() + .expect("manager name is free"); + + let (a, b) = pair(&TcpTransport); + spawn_established(FramedConn::new(a), peer("kept")).expect("register"); + + // The far heartbeat pump: interval-paced sends until told to stop, + // then holds the socket open, silent, so the eventual teardown is + // liveness — not EOF. + let (ctl_tx, ctl_rx) = smarm::channel::channel::<()>(); + spawn(move || { + let mut far = FramedConn::new(b); + // Phase 1: heartbeat at the interval until the first signal. + while matches!(ctl_rx.try_recv(), Ok(None)) { + far.send(&Frame::Heartbeat).expect("far send"); + sleep(HEARTBEAT_INTERVAL); + } + // Phase 2: silent but with the socket held open — dropping + // `far` here would EOF the actor and mask the liveness path. + // Exits when the test's closure ends and drops `ctl_tx` (an + // eternal park would stop `run` from ever returning). + while matches!(ctl_rx.try_recv(), Ok(None)) { + sleep(Duration::from_millis(20)); + } + }); + + // Past the liveness window with margin: still up. + sleep(LIVENESS_TIMEOUT + LIVENESS_TIMEOUT / 2); + assert_eq!( + peers(), + vec!["kept".to_string()], + "liveness fired despite heartbeats" + ); + + // Silence the pump; liveness now empties and the entry goes. + ctl_tx.send(()).expect("pump alive"); + wait_peers(&[], LIVENESS_TIMEOUT * 2); + mgr.shutdown(); + // `ctl_tx` drops here, releasing the pump's phase-2 wait. + }); +}