feat(cluster): RFC 010 c6c — heartbeat send + fixed-timeout liveness + teardown
The timeout arm of the connection actor's select: HEARTBEAT_INTERVAL (1s) paces outbound Frame::Heartbeat (first at spawn, so the peer's window starts fed) and LIVENESS_TIMEOUT (4s = 4 intervals) declares the peer dead when no inbound frame arrives inside it — any frame resets the window, so heartbeats keep an idle connection alive and real traffic (c8+) counts for free. Fire => close + exit; the manager's monitor reaps the table entry as on every other exit path. Fixed timeout per RFC v2 §5 (control connection, heartbeats can't queue behind bulk). Intervals are the one-viable-answer call flagged for veto at diff review. The pump was made non-blocking to keep the deadlines honest: a plain recv() blocks into the socket while the buffer holds a partial frame, parking the actor past its timers. Two additive FramedConn methods (read_once, next_buffered): exactly one socket read per level-triggered readable wake (cannot block, cannot strand — leftovers re-signal), then drain every complete buffered frame. Liveness resets only on complete frames. No-fd transports (loopback) still get the command-only loop: no readiness means no timers, same caveat as recv_deadline. tests/cluster_conn_liveness.rs 3/0, stable over 5 runs (raw far end over localhost TCP: heartbeats appear unprompted; mute peer still up at half the window, gone after it; heartbeat-only peer survives 1.5x the window, then reaped once silenced). Cluster suites regression-clean; clippy --lib green both configs; fmt clean.
This commit is contained in:
+110
-36
@@ -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<Pid, Register
|
||||
}
|
||||
}
|
||||
|
||||
/// How often this end emits [`Frame::Heartbeat`] on an idle connection. The
|
||||
/// first one goes out immediately at spawn, so the peer's liveness window
|
||||
/// starts fed.
|
||||
pub const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(1);
|
||||
|
||||
/// How long the connection may go without a single inbound frame before it
|
||||
/// is declared dead and torn down. Any inbound frame resets the window —
|
||||
/// heartbeats keep an idle connection alive, and real traffic (c8+) counts
|
||||
/// for free. Fixed by design (RFC v2 §5): this is the control connection, a
|
||||
/// heartbeat can never queue behind bulk traffic, so a fixed timeout is an
|
||||
/// honest detector.
|
||||
pub const LIVENESS_TIMEOUT: Duration = Duration::from_secs(4);
|
||||
|
||||
fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver<Cmd>) {
|
||||
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<Cmd>) {
|
||||
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<Cmd>) {
|
||||
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<Cmd>) -> 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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<usize> {
|
||||
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<Option<Frame>, 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();
|
||||
|
||||
@@ -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<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)
|
||||
}
|
||||
|
||||
fn peers() -> Vec<String> {
|
||||
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<String> = 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.
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user