diff --git a/src/cluster.rs b/src/cluster.rs index 16b78b8..f5bfffc 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -2,9 +2,62 @@ //! //! c1: feature flag + optional deps. c2: the owned envelope. c3: the //! transport trait (control connection), framed codec, and the TCP + -//! loopback impls. c5: the handshake state machine. Everything above them -//! lands in later chunks. +//! 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 +//! chunks. +pub mod conn; pub mod envelope; pub mod handshake; +pub mod manager; pub mod transport; + +use std::time::Duration; + +use crate::gen_server::{self, GenServerBuilder}; +use crate::monitor::monitor; +use crate::scheduler::{sleep, spawn, JoinHandle}; +use crate::supervisor::{ChildSpec, OneForOne, Restart}; + +pub use conn::{spawn_established, ConnHandle}; +pub use manager::{Manager, MANAGER}; + +/// A running cluster subtree: an explicitly-started supervisor over the +/// connection [`Manager`]. Roles will eventually mount this subtree; until the +/// role mechanism lands it is started by hand (RFC 010 §7). Dropping the handle +/// detaches the subtree, which keeps running for the life of the runtime. +pub struct Cluster { + _sup: JoinHandle, +} + +/// Start the cluster subtree and block until the manager is registered and +/// ready to answer. The manager is a supervised child (restarted on crash); +/// per-peer connection actors are dynamic and monitored by the manager rather +/// than statically supervised — a lost connection is re-established by dialing +/// (c7), never resurrected onto a stale socket. +pub fn start() -> Cluster { + let sup = spawn(|| { + OneForOne::new() + .child(ChildSpec::new(Restart::Permanent, manager_child)) + .run() + }); + while gen_server::whereis_server(MANAGER).is_none() { + sleep(Duration::from_millis(1)); + } + Cluster { _sup: sup } +} + +/// The supervised manager child body. It *is* the child actor: it starts the +/// named manager, then parks on the manager's own termination so this actor's +/// lifetime tracks the manager's — the supervisor's restart accounting keys off +/// this actor exiting. +fn manager_child() { + let m = match GenServerBuilder::new(Manager::new()).named(MANAGER).start() { + Ok(m) => m, + // Name still held by a not-yet-reaped prior instance: return and let + // the supervisor retry under its restart policy. + Err(_) => return, + }; + let _ = monitor(m.pid()).rx.recv(); +} diff --git a/src/cluster/conn.rs b/src/cluster/conn.rs new file mode 100644 index 0000000..78122e7 --- /dev/null +++ b/src/cluster/conn.rs @@ -0,0 +1,120 @@ +//! 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 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`). + +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}; + +/// Commands to a running connection actor. +enum Cmd { + Shutdown, +} + +/// A handle to a running connection actor. +pub struct ConnHandle { + cmd_tx: Sender, +} + +impl ConnHandle { + /// Ask the connection to close and exit. Idempotent, and a no-op if the + /// actor has already gone. + 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 } +} + +fn run(mut framed: FramedConn, peer: Peer, cmd_rx: Receiver) { + let me = scheduler::self_pid(); + match gen_server::call( + MANAGER, + Call::Register { + name: peer.node_name.clone(), + pid: me, + }, + ) { + 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, + } + + 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, + } + } + 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) { + break; + } + } + } + } + + framed.close(); +} + +/// 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) -> bool { + match cmd_rx.try_recv() { + Ok(Some(Cmd::Shutdown)) => true, + Ok(None) => false, // spurious wake + Err(_) => true, // all senders dropped + } +} + +/// 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 + } +} diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs new file mode 100644 index 0000000..9ab3218 --- /dev/null +++ b/src/cluster/manager.rs @@ -0,0 +1,110 @@ +//! 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. +//! +//! 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. +//! What lives here is only the table itself and the uniqueness rule the +//! handshake's `NameTaken` verdict depends on. + +use std::collections::HashMap; + +use crate::gen_server::{GenServer, GenServerCtx, GenServerName, Watcher}; +use crate::monitor::{monitor, Down}; +use crate::pid::Pid; + +/// Well-known name of the singleton manager within a runtime. Connection +/// actors reach it by name rather than by a passed-around ref, so a restarted +/// manager is always found at the same key. +pub const MANAGER: GenServerName = GenServerName::new("smarm.cluster.manager"); + +/// The connection registry: peer name → the connection actor that owns that +/// peer's control connection. +pub struct Manager { + conns: HashMap, + watcher: Option>, +} + +impl Manager { + pub fn new() -> Self { + Manager { + conns: HashMap::new(), + watcher: None, + } + } +} + +impl Default for Manager { + fn default() -> Self { + Manager::new() + } +} + +/// Outcome of a [`Call::Register`]. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Registered { + /// The name was free; this connection is now the peer of record. + Ok, + /// Another live connection already holds this name — the caller lost the + /// race (or is a duplicate) and must not run. + Duplicate, +} + +/// 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 current peer names, sorted. For observation and tests. + Peers, +} + +/// Replies from the manager. +#[derive(Debug)] +pub enum Reply { + Registered(Registered), + Peers(Vec), +} + +impl GenServer for Manager { + type Call = Call; + type Reply = Reply; + type Cast = (); + type Info = (); + type Timer = (); + + fn init(&mut self, ctx: &GenServerCtx) { + self.watcher = Some(ctx.watcher()); + } + + fn handle_call(&mut self, request: Call) -> Reply { + match request { + Call::Register { name, pid } => { + if self.conns.contains_key(&name) { + return Reply::Registered(Registered::Duplicate); + } + if let Some(w) = &self.watcher { + w.watch(monitor(pid)); + } + self.conns.insert(name, pid); + Reply::Registered(Registered::Ok) + } + Call::Peers => { + let mut names: Vec = self.conns.keys().cloned().collect(); + names.sort(); + Reply::Peers(names) + } + } + } + + fn handle_cast(&mut self, _request: ()) {} + + fn handle_down(&mut self, down: Down) { + self.conns.retain(|_, pid| *pid != down.pid); + } +} diff --git a/src/cluster/transport.rs b/src/cluster/transport.rs index f16057d..d62d786 100644 --- a/src/cluster/transport.rs +++ b/src/cluster/transport.rs @@ -50,6 +50,17 @@ pub trait Conn: Send { /// Diagnostic label for logs only. Mesh identity comes from the /// handshake (`Hello`/`HelloAck`), never from the transport. fn peer_addr(&self) -> String; + + /// Readiness as a [`select`](crate::select) arm, for transports backed by + /// a file descriptor. `Some` lets a driver wait on "this connection is + /// readable" alongside an ordinary command inbox in a single `select`, so + /// one actor can interleave reading with control messages without a + /// second thread. The default is `None`: a transport with no fd (the + /// in-memory loopback) cannot be selected on and must be driven another + /// way. + fn readable_arm(&self) -> Option { + None + } } /// A bound listen point producing inbound [`Conn`]s. @@ -191,4 +202,10 @@ impl FramedConn { pub fn peer_addr(&self) -> String { self.conn.peer_addr() } + + /// The underlying connection's readiness arm, if it is fd-backed (see + /// [`Conn::readable_arm`]). + pub fn readable_arm(&self) -> Option { + self.conn.readable_arm() + } } diff --git a/src/cluster/transport/tcp.rs b/src/cluster/transport/tcp.rs index c725088..d3764e6 100644 --- a/src/cluster/transport/tcp.rs +++ b/src/cluster/transport/tcp.rs @@ -180,6 +180,10 @@ impl Conn for TcpConn { Err(_) => "".to_string(), } } + + fn readable_arm(&self) -> Option { + Some(crate::scheduler::FdArm::readable(self.fd())) + } } // --------------------------------------------------------------------------- diff --git a/tests/cluster_conn_lifecycle.rs b/tests/cluster_conn_lifecycle.rs new file mode 100644 index 0000000..040afbc --- /dev/null +++ b/tests/cluster_conn_lifecycle.rs @@ -0,0 +1,104 @@ +//! 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. 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. +//! +//! 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, 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) +} + +/// 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 = 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. + let h1 = spawn_established(FramedConn::new(a1), peer("node-b")); + let _h2 = spawn_established(FramedConn::new(a2), peer("node-c")); + + // 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(); + 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 + // shutdown command and not an EOF. + drop(b1); + + // All connection actors have exited; stop the manager so `run` returns. + mgr.shutdown(); + }); +}