feat(cluster): RFC 010 c6a — connection actor + manager subtree
Per-peer connection actor as a single select-loop plain actor owning the whole FramedConn: one select folds its command inbox and the transport's readable arm, so reads and control share one execution context — no reader thread, no read/write split. The handshake is bypassed here (c6b wires it); the actor is spawned already-established and self-registers with the manager. Manager gen_server: the peer-name -> conn-pid registry and the uniqueness source the handshake's NameTaken depends on. It monitors each connection, so the table self-heals on any exit path. Explicit supervision subtree keeps the manager up; connections are dynamic and monitored, never restarted (c7 re-dials). Transport gains an additive Conn::readable_arm -> Option<FdArm> (default None; TCP returns its fd's arm, loopback stays None). Existing c3 transport tests unchanged. Lifecycle test over localhost TCP: up reflected in the table, commanded shutdown reaps exactly one, peer EOF reaps the other.
This commit is contained in:
+55
-2
@@ -2,9 +2,62 @@
|
|||||||
//!
|
//!
|
||||||
//! c1: feature flag + optional deps. c2: the owned envelope. c3: the
|
//! c1: feature flag + optional deps. c2: the owned envelope. c3: the
|
||||||
//! transport trait (control connection), framed codec, and the TCP +
|
//! transport trait (control connection), framed codec, and the TCP +
|
||||||
//! loopback impls. c5: the handshake state machine. Everything above them
|
//! loopback impls. c5: the handshake state machine. c6: the connection
|
||||||
//! lands in later chunks.
|
//! [`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 envelope;
|
||||||
pub mod handshake;
|
pub mod handshake;
|
||||||
|
pub mod manager;
|
||||||
pub mod transport;
|
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();
|
||||||
|
}
|
||||||
|
|||||||
@@ -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<Cmd>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<Cmd>) {
|
||||||
|
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<Cmd>) -> 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
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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<Manager> = 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<String, Pid>,
|
||||||
|
watcher: Option<Watcher<Manager>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
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<String>),
|
||||||
|
}
|
||||||
|
|
||||||
|
impl GenServer for Manager {
|
||||||
|
type Call = Call;
|
||||||
|
type Reply = Reply;
|
||||||
|
type Cast = ();
|
||||||
|
type Info = ();
|
||||||
|
type Timer = ();
|
||||||
|
|
||||||
|
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
||||||
|
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<String> = 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);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -50,6 +50,17 @@ pub trait Conn: Send {
|
|||||||
/// Diagnostic label for logs only. Mesh identity comes from the
|
/// Diagnostic label for logs only. Mesh identity comes from the
|
||||||
/// handshake (`Hello`/`HelloAck`), never from the transport.
|
/// handshake (`Hello`/`HelloAck`), never from the transport.
|
||||||
fn peer_addr(&self) -> String;
|
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<crate::scheduler::FdArm> {
|
||||||
|
None
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A bound listen point producing inbound [`Conn`]s.
|
/// A bound listen point producing inbound [`Conn`]s.
|
||||||
@@ -191,4 +202,10 @@ impl FramedConn {
|
|||||||
pub fn peer_addr(&self) -> String {
|
pub fn peer_addr(&self) -> String {
|
||||||
self.conn.peer_addr()
|
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<crate::scheduler::FdArm> {
|
||||||
|
self.conn.readable_arm()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -180,6 +180,10 @@ impl Conn for TcpConn {
|
|||||||
Err(_) => "<disconnected>".to_string(),
|
Err(_) => "<disconnected>".to_string(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn readable_arm(&self) -> Option<crate::scheduler::FdArm> {
|
||||||
|
Some(crate::scheduler::FdArm::readable(self.fd()))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
@@ -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<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)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 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<String> = 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();
|
||||||
|
});
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user