Files
smarm/src/cluster/conn.rs
T
Claude 160967939b feat(cluster): RFC 010 c7a — membership events + view at the manager
node_up/node_down are derived facts of the manager's own register/remove
events, so the membership state lives in the manager (no cross-actor race
between 'connection exists' and 'node is up'); src/cluster/membership.rs is
the consumer surface: NodeEvent/NodeInfo, subscribe(), view(). The conn
table stays private — no consumer touches it (roadmap-binding).

Ratified semantics: subscribe() is snapshot-then-stream — one NodeUp per
live peer is queued before the subscription joins the list, exact because
gen_server handlers are serialized. Dropped subscribers are pruned on the
next emit (closed channel), no monitor needed.

One-viable call, flagged: NodeId is memoized per (name, incarnation) — a
compact local alias for the wire identity, per pg.rs's framing. A reconnect
blip at the same incarnation keeps its id; a restart (new incarnation) gets
a fresh one, so a ghost and its successor are always distinguishable.
Allocation starts at 1; NodeId(0) stays pg::DEFAULT_NODE_ID (self).

Call::Register now carries the whole handshake Peer (the path already has
it; node_up needs incarnation + meta).

tests/cluster_membership.rs 4/0 stable over 5 runs (live up/down over
localhost TCP with commanded and EOF teardown; late-subscriber snapshot +
view agreement; restart-vs-blip id identity; dead-subscriber pruning). All
cluster suites regression-clean; clippy --lib green both configs; fmt
clean; default build compiles.
2026-08-15 07:07:14 +00:00

220 lines
8.7 KiB
Rust

//! RFC 010 c6 — the per-peer connection actor.
//!
//! One actor per established control connection. It owns the whole
//! [`FramedConn`] and, in a single [`select`](crate::select), waits on two
//! things at once: its command inbox and the connection becoming readable (the
//! [`FdArm`](crate::scheduler::FdArm) the transport hands back). That is why it
//! is a plain select-loop actor rather than a `gen_server` or `gen_statem` —
//! neither of those can fold fd-readiness into its wait, and folding it in is
//! the whole job. The single owner sends and receives on the one `FramedConn`,
//! 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 *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), 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 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;
use crate::gen_server;
use crate::pid::Pid;
use crate::scheduler::spawn;
/// Commands to a running connection actor.
enum Cmd {
Shutdown,
}
/// 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<Cmd>,
}
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. 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);
}
}
/// 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;
/// 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<Pid, RegisterRefused> {
let (cmd_tx, cmd_rx) = channel();
let reg_peer = peer.clone();
let pid = spawn(move || run(framed, peer, cmd_rx)).pid();
match gen_server::call(
MANAGER,
Call::Register {
peer: reg_peer,
pid,
handle: ConnHandle { cmd_tx },
},
) {
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),
}
}
/// 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 {
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;
}
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,
}
}
}
/// 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
/// shutdown was requested, or the last handle was dropped.
fn should_stop(cmd_rx: &Receiver<Cmd>) -> bool {
match cmd_rx.try_recv() {
Ok(Some(Cmd::Shutdown)) => true,
Ok(None) => false, // spurious wake
Err(_) => true, // all senders dropped
}
}
/// 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)
}
}