diff --git a/src/cluster.rs b/src/cluster.rs index 286ed0c..851a011 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -10,21 +10,32 @@ pub mod conn; pub mod connect; +pub mod connector; +pub mod discovery; pub mod envelope; pub mod handshake; pub mod manager; pub mod membership; pub mod transport; -use std::time::Duration; +use std::io; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; use crate::gen_server::{self, GenServerBuilder}; use crate::monitor::monitor; +use crate::pg::Incarnation; use crate::scheduler::{sleep, spawn, JoinHandle}; use crate::supervisor::{ChildSpec, OneForOne, Restart}; +use envelope::NodeMeta; +use handshake::Local; +use transport::tcp::TcpTransport; +use transport::Transport; + pub use conn::{spawn_established, ConnHandle}; pub use connect::{dial, spawn_acceptor, AcceptorHandle}; +pub use connector::{spawn_connector, ConnectorHandle}; +pub use discovery::{Discovery, StaticSeeds, Strategy}; pub use manager::{Manager, MANAGER}; pub use membership::{subscribe, view, MembershipEvents, NodeEvent, NodeInfo}; @@ -65,20 +76,67 @@ const fn fold_u32(mut h: u64, v: u32) -> u64 { h } -/// 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, +/// How to run this node: its identity and how it finds peers. +pub struct Config { + /// This node's claimed name — the mesh-wide identity peers dial by and + /// the tie-break input. Must be unique across the mesh. + pub node_name: String, + /// Metadata offered in this node's `Hello`. + pub meta: NodeMeta, + /// The control-connection listen address (e.g. `"127.0.0.1:0"`; the + /// concrete bound address is [`Cluster::local_addr`]). + pub listen_addr: String, + /// The peer-discovery strategy — [`StaticSeeds`] until richer ones land. + pub strategy: Box, } -/// 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 { +/// A running cluster node: the supervised [`Manager`], the acceptor over the +/// bound listener, and the connector driving its [`Strategy`]. Roles will +/// eventually mount this; until the role mechanism lands it is started by +/// hand (RFC 010 §7). +/// +/// Dropping the handle stops the acceptor and connector loops (no new +/// connections in either direction) but detaches the manager subtree, which +/// — with every established connection — keeps running for the life of the +/// runtime, the same split as [`AcceptorHandle`] alone. +pub struct Cluster { + _sup: JoinHandle, + acceptor: AcceptorHandle, + connector: ConnectorHandle, + local: Local, +} + +impl Cluster { + /// The concrete bound listen address, dialable as-is. + pub fn local_addr(&self) -> &str { + self.acceptor.local_addr() + } + + /// This node's handshake identity (name, incarnation, build hash, meta). + pub fn local(&self) -> &Local { + &self.local + } + + /// Stop accepting and dialing. Established connections stay up (they + /// belong to the manager); tear those down via the manager. + pub fn shutdown(&self) { + self.acceptor.shutdown(); + self.connector.shutdown(); + } +} + +/// Start a cluster node: the supervised manager (blocking until it is +/// registered and ready to answer), the acceptor bound per +/// [`Config::listen_addr`], and the connector running [`Config::strategy`]. +/// The node's identity is completed here: `incarnation` is +/// [`self_incarnation`] and `build_hash` is [`BUILD_HASH`] — c7 is its first +/// consumer. Errs only if the listener cannot bind. +/// +/// 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 the +/// connector's dial loop, never resurrected onto a stale socket. +pub fn start(config: Config) -> io::Result { let sup = spawn(|| { OneForOne::new() .child(ChildSpec::new(Restart::Permanent, manager_child)) @@ -87,7 +145,35 @@ pub fn start() -> Cluster { while gen_server::whereis_server(MANAGER).is_none() { sleep(Duration::from_millis(1)); } - Cluster { _sup: sup } + let local = Local { + node_name: config.node_name, + incarnation: self_incarnation(), + build_hash: BUILD_HASH, + meta: config.meta, + }; + let listener = TcpTransport.listen(&config.listen_addr)?; + let acceptor = spawn_acceptor(listener, local.clone()); + let connector = spawn_connector(Box::new(TcpTransport), local.clone(), config.strategy); + Ok(Cluster { + _sup: sup, + acceptor, + connector, + local, + }) +} + +/// This process's incarnation epoch: milliseconds since the Unix epoch, +/// truncated to `u32`. Not a clock — its one job is separating a node from +/// its own restart (two starts of the same name land on the same value only +/// if they happen within the same millisecond modulo ~49.7 days). Seconds +/// would be too coarse: a crash-and-restart inside one second is routine +/// under supervision. +pub fn self_incarnation() -> Incarnation { + let ms = SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_millis()) + .unwrap_or(0); + Incarnation::new(ms as u32) } /// The supervised manager child body. It *is* the child actor: it starts the diff --git a/src/cluster/connector.rs b/src/cluster/connector.rs new file mode 100644 index 0000000..b4633ea --- /dev/null +++ b/src/cluster/connector.rs @@ -0,0 +1,229 @@ +//! RFC 010 c7b — the connector: the dial loop that turns discovered +//! candidates into a full mesh. +//! +//! A plain select-loop actor (the c6 shape). It spawns its [`Strategy`] as a +//! child actor and receives [`Discovery`] events from it; it tracks which +//! peers are up by **subscribing to membership like any other consumer** — +//! no privileged channel into the manager, the same snapshot-then-stream +//! surface c8 will use. One `select` folds the command inbox, the discovery +//! stream, the membership stream, and the earliest retry deadline into a +//! single wait. +//! +//! Per-candidate state: dial on arrival; on failure retry with capped +//! exponential backoff ([`INITIAL_BACKOFF`] doubling to [`MAX_BACKOFF`]); +//! on the peer's `node_up` stop dialing and reset the backoff; on its +//! `node_down` resume immediately (a fresh sequence — the reconnect case is +//! the one backoff exists to pace, but the *first* retry after a death +//! should be prompt). A candidate bearing our own name is parked permanently +//! — that seed is us. Every other failure retries: in particular a +//! `NameTaken` reject can be our own ghost at the peer, not yet reaped by +//! its liveness timer, so it must not park. +//! +//! Dials run **inline in the loop** — the same deliberate serialization as +//! the acceptor (c6b): each attempt is bounded by the connect + handshake +//! deadlines, and nothing concurrent exists to be starved. A wall of slow +//! unreachable seeds would stretch the loop's latency; revisit if a real +//! deployment ever hits that shape. + +use std::collections::{HashMap, HashSet}; +use std::time::{Duration, Instant}; + +use crate::channel::{channel, select, select_timeout, Receiver, Selectable, Sender}; +use crate::cluster::connect::dial; +use crate::cluster::discovery::{Discovery, Strategy}; +use crate::cluster::handshake::Local; +use crate::cluster::membership::{subscribe, NodeEvent}; +use crate::cluster::transport::Transport; +use crate::pg::NodeId; +use crate::scheduler::spawn; + +/// First retry delay after a failed dial attempt. +pub const INITIAL_BACKOFF: Duration = Duration::from_millis(250); +/// Backoff ceiling: an unreachable seed is retried this often, forever. +pub const MAX_BACKOFF: Duration = Duration::from_secs(5); + +enum Cmd { + Shutdown, +} + +/// A running connector. `shutdown` (or dropping the last handle) stops the +/// dial loop and its strategy only — established connections belong to the +/// manager, exactly as with the acceptor. +pub struct ConnectorHandle { + cmd_tx: Sender, +} + +impl ConnectorHandle { + /// Ask the connector to stop. Idempotent; a no-op if it already has. + pub fn shutdown(&self) { + let _ = self.cmd_tx.send(Cmd::Shutdown); + } +} + +/// One discovered `(name, addr)` and its dial state. +struct Candidate { + name: String, + addr: String, + /// This seed is the local node itself: never dialed. + parked: bool, + /// Delay to apply after the *next* failure. + backoff: Duration, + next_attempt: Instant, +} + +/// Spawn the connector actor. The strategy is spawned as its child; the +/// membership subscription is taken inside the actor. Must be called from +/// inside an actor (the same requirement as `dial`). +pub fn spawn_connector( + transport: Box, + local: Local, + strategy: Box, +) -> ConnectorHandle { + let (cmd_tx, cmd_rx) = channel(); + spawn(move || run(transport, local, strategy, cmd_rx)); + ConnectorHandle { cmd_tx } +} + +fn run( + transport: Box, + local: Local, + strategy: Box, + cmd_rx: Receiver, +) { + // Membership is the connector's source of truth for "who is up" — the + // snapshot seeds `connected` before any candidate arrives. + let Some(events) = subscribe() else { + return; // no manager, no cluster to connect + }; + let (disc_tx, disc_rx) = channel(); + spawn(move || strategy.run(disc_tx)); + + let mut cands: Vec = Vec::new(); + let mut connected: HashSet = HashSet::new(); + let mut names: HashMap = HashMap::new(); // NodeDown carries only the id + let mut strategy_done = false; + + loop { + // Drain every input, then act. Order does not matter: acting is + // idempotent against the resulting state. + match drain_cmd(&cmd_rx) { + Drained::Stop => return, + Drained::Open => {} + } + if !strategy_done { + strategy_done = drain_discoveries(&disc_rx, &local, &mut cands); + } + match drain_events(&events.rx, &mut connected, &mut names, &mut cands) { + Drained::Stop => return, // manager gone: the cluster is tearing down + Drained::Open => {} + } + + // Dial everything due, inline (see the module docs on serialization). + let now = Instant::now(); + for c in cands + .iter_mut() + .filter(|c| !c.parked && !connected.contains(&c.name) && c.next_attempt <= now) + { + // The outcome does not branch the bookkeeping: on success the + // manager's node_up is on its way and flips `connected` (backing + // off meanwhile keeps a racing re-attempt from spinning); every + // failure retries — see the module docs. + let _ = dial(&*transport, &c.addr, &c.name, &local); + c.next_attempt = Instant::now() + c.backoff; + c.backoff = (c.backoff * 2).min(MAX_BACKOFF); + } + + // Wait: until the earliest retry deadline among actionable + // candidates, or indefinitely if none is pending. + let deadline = cands + .iter() + .filter(|c| !c.parked && !connected.contains(&c.name)) + .map(|c| c.next_attempt) + .min(); + let mut arms: Vec<&dyn Selectable> = vec![&cmd_rx, &events.rx]; + if !strategy_done { + arms.push(&disc_rx); + } + match deadline { + Some(d) => { + let wait = d.saturating_duration_since(Instant::now()); + let _ = select_timeout(&arms, wait); + } + None => { + let _ = select(&arms); + } + } + } +} + +enum Drained { + Open, + Stop, +} + +fn drain_cmd(rx: &Receiver) -> Drained { + match rx.try_recv() { + Ok(Some(Cmd::Shutdown)) => Drained::Stop, + Ok(None) => Drained::Open, + Err(_) => Drained::Stop, // all handles dropped + } +} + +/// Pull every pending discovery into the candidate set (deduplicated by +/// `(name, addr)`; a candidate bearing the local name is parked). Returns +/// `true` once the strategy's channel closes — it has said all it will. +fn drain_discoveries(rx: &Receiver, local: &Local, cands: &mut Vec) -> bool { + loop { + match rx.try_recv() { + Ok(Some(Discovery::Candidate { name, addr })) => { + if cands.iter().any(|c| c.name == name && c.addr == addr) { + continue; + } + let parked = name == local.node_name; + cands.push(Candidate { + name, + addr, + parked, + backoff: INITIAL_BACKOFF, + next_attempt: Instant::now(), + }); + } + Ok(None) => return false, + Err(_) => return true, // strategy done; its candidates live on here + } + } +} + +/// Fold pending membership events into `connected` (and the id→name map). +/// `node_up` resets its candidates' backoff; `node_down` schedules a prompt +/// redial with a fresh sequence. +fn drain_events( + rx: &Receiver, + connected: &mut HashSet, + names: &mut HashMap, + cands: &mut [Candidate], +) -> Drained { + loop { + match rx.try_recv() { + Ok(Some(NodeEvent::NodeUp(info))) => { + names.insert(info.node, info.name.clone()); + for c in cands.iter_mut().filter(|c| c.name == info.name) { + c.backoff = INITIAL_BACKOFF; + } + connected.insert(info.name); + } + Ok(Some(NodeEvent::NodeDown { node })) => { + if let Some(name) = names.remove(&node) { + connected.remove(&name); + let now = Instant::now(); + for c in cands.iter_mut().filter(|c| c.name == name) { + c.backoff = INITIAL_BACKOFF; + c.next_attempt = now; + } + } + } + Ok(None) => return Drained::Open, + Err(_) => return Drained::Stop, + } + } +} diff --git a/src/cluster/discovery.rs b/src/cluster/discovery.rs new file mode 100644 index 0000000..b7c0e29 --- /dev/null +++ b/src/cluster/discovery.rs @@ -0,0 +1,67 @@ +//! RFC 010 c7b — peer discovery: the [`Strategy`] seam and the static-seeds +//! implementation. +//! +//! A strategy is **push-based and runs as its own actor**: the +//! [`connector`](crate::cluster::connector) spawns it with the sending end of +//! a channel, and the strategy emits [`Discovery`] events whenever it learns +//! something — once at startup for a static list, continuously for a future +//! mDNS/DNS strategy — for as long as it cares to run. Returning ends the +//! strategy actor; the candidates it pushed live on in the connector (the +//! connector owns all retry/backoff state, so a strategy never re-announces). +//! +//! A candidate is a **`(node_name, addr)` pair**, not a bare address: the +//! dial path and the D7 tie-break are keyed by peer *name* (the dial intent +//! must be registered before connecting so a crossing inbound `Hello` sees +//! it), so an anonymous dial would reintroduce exactly the +//! simultaneous-connect flap D7 exists to prevent. Discovery mechanisms know +//! names — that is what they discover. + +use crate::channel::Sender; + +/// A discovery event, as pushed by a [`Strategy`]. +/// +/// Additive-only for now (candidates are announced, never withdrawn); +/// `#[non_exhaustive]` so expiry can land later without breaking strategies. +#[derive(Debug, Clone, PartialEq, Eq)] +#[non_exhaustive] +pub enum Discovery { + /// A peer worth dialing: its claimed node name and a dialable address. + Candidate { name: String, addr: String }, +} + +/// A source of peers to dial. Implementations are spawned as actors by the +/// connector — see the module docs for the contract. +pub trait Strategy: Send + 'static { + /// Run the strategy: push [`Discovery`] events into `out` as they are + /// learned; return when done discovering (or when `out` reports closed — + /// the connector is gone). Runs inside an actor, so blocking + /// cooperatively is fine. + fn run(self: Box, out: Sender); +} + +/// The static-seeds strategy: a fixed `(name, addr)` list, announced once. +#[derive(Debug, Clone, Default)] +pub struct StaticSeeds { + seeds: Vec<(String, String)>, +} + +impl StaticSeeds { + pub fn new(seeds: impl IntoIterator, impl Into)>) -> Self { + StaticSeeds { + seeds: seeds + .into_iter() + .map(|(n, a)| (n.into(), a.into())) + .collect(), + } + } +} + +impl Strategy for StaticSeeds { + fn run(self: Box, out: Sender) { + for (name, addr) in self.seeds { + if out.send(Discovery::Candidate { name, addr }).is_err() { + return; // connector gone; nobody to discover for + } + } + } +} diff --git a/tests/cluster_mesh.rs b/tests/cluster_mesh.rs new file mode 100644 index 0000000..6b6fb18 --- /dev/null +++ b/tests/cluster_mesh.rs @@ -0,0 +1,179 @@ +//! RFC 010 c7 — the Phase 2 gate: a 3-node mesh under the subprocess +//! harness, repeatable. +//! +//! Each node process runs the integrated `cluster::start` (manager + +//! acceptor + connector + static seeds), subscribes to membership like any +//! consumer, and announces protocol-visible facts as lines: +//! `LISTENING `, `MEMBER-UP inc=`, `MEMBER-DOWN `. +//! Then it **parks forever** — cross-process teardown is retractable state +//! (binding trap), so the parent SIGKILLs via `Node`'s `Drop` and clean exit +//! stays the c4 harness's own smoke test. +//! +//! Ports: nodes bind `:0` and report, so the mesh is built by seeding each +//! node with the previously-reported addresses (n1: no seeds; n2: n1; +//! n3: n1+n2 — inbound covers the reverse edges). The late-seed test is the +//! one exception: the parent pre-reserves a port by binding-and-closing it, +//! seeds one node with it, then starts the second node on that exact +//! address. In principle another process could steal the port in the gap; +//! in practice the window is microseconds on a local runner — accepted, and +//! confined to that one test. +#![cfg(feature = "cluster")] + +mod common; + +use common::{maybe_child, spawn_node, Node}; +use smarm::cluster::envelope::NodeMeta; +use smarm::cluster::membership::{subscribe, NodeEvent}; +use smarm::cluster::{start, Config, StaticSeeds}; +use smarm::pg::NodeId; +use std::collections::HashMap; +use std::time::Duration; + +const ROLES: &[(&str, fn())] = &[("node", role_node)]; + +/// A mesh node: identity and seeds from env, membership events to stdout, +/// park forever (the parent reaps). +fn role_node() { + let name = std::env::var("SMARM_NODE_NAME").expect("SMARM_NODE_NAME not set"); + let listen = std::env::var("SMARM_LISTEN_ADDR").unwrap_or_else(|_| "127.0.0.1:0".to_string()); + // Seeds: comma-separated `name=addr` pairs; empty or unset means none. + let seeds: Vec<(String, String)> = std::env::var("SMARM_SEEDS") + .unwrap_or_default() + .split(',') + .filter(|s| !s.is_empty()) + .map(|s| { + let (n, a) = s.split_once('=').expect("seed must be name=addr"); + (n.to_string(), a.to_string()) + }) + .collect(); + + smarm::run(move || { + let cluster = start(Config { + node_name: name, + meta: NodeMeta { + role: "mesh-test".to_string(), + region: "local".to_string(), + }, + listen_addr: listen, + strategy: Box::new(StaticSeeds::new(seeds)), + }) + .expect("listener binds"); + println!("LISTENING {}", cluster.local_addr()); + + let events = subscribe().expect("manager is up"); + let mut names: HashMap = HashMap::new(); + loop { + match events.rx.recv() { + Ok(NodeEvent::NodeUp(info)) => { + names.insert(info.node, info.name.clone()); + println!("MEMBER-UP {} inc={}", info.name, info.incarnation.get()); + } + Ok(NodeEvent::NodeDown { node }) => { + let name = names.remove(&node).unwrap_or_else(|| "?".to_string()); + println!("MEMBER-DOWN {name}"); + } + Err(_) => break, // manager gone; park below regardless + } + } + loop { + smarm::sleep(Duration::from_secs(3600)); + } + }); +} + +fn spawn_mesh_node(name: &str, seeds: &str, listen: Option<&str>) -> Node { + let mut env: Vec<(&str, &str)> = vec![("SMARM_NODE_NAME", name), ("SMARM_SEEDS", seeds)]; + if let Some(addr) = listen { + env.push(("SMARM_LISTEN_ADDR", addr)); + } + spawn_node("node", &env) +} + +/// Wait for `MEMBER-UP inc=` and return the incarnation. +fn wait_member_up(node: &mut Node, peer: &str) -> u32 { + let prefix = format!("MEMBER-UP {peer} inc="); + let line = node.wait_line(&format!("MEMBER-UP {peer}"), |l| l.starts_with(&prefix)); + line[prefix.len()..].parse().expect("incarnation parses") +} + +fn wait_member_down(node: &mut Node, peer: &str) { + let want = format!("MEMBER-DOWN {peer}"); + node.wait_line(&want, |l| l == want); +} + +/// The gate, plus the kill and restart facts, as one mesh's life: three +/// nodes form a full mesh (every node sees both others up); killing one +/// yields `node_down` at both survivors; its restart under the same name +/// arrives as a NEW incarnation — the ghost and its successor are +/// distinguishable at every observer. +#[test] +fn three_node_mesh_forms_then_kill_then_restart_distinguishable() { + maybe_child(ROLES); + + let mut n1 = spawn_mesh_node("node-1", "", None); + let a1 = n1.wait_listening(); + let mut n2 = spawn_mesh_node("node-2", &format!("node-1={a1}"), None); + let a2 = n2.wait_listening(); + let mut n3 = spawn_mesh_node("node-3", &format!("node-1={a1},node-2={a2}"), None); + let _a3 = n3.wait_listening(); + + // Full mesh: each node reports both peers up (dialed or inbound alike). + wait_member_up(&mut n1, "node-2"); + let inc3_at_n1 = wait_member_up(&mut n1, "node-3"); + wait_member_up(&mut n2, "node-1"); + let inc3_at_n2 = wait_member_up(&mut n2, "node-3"); + wait_member_up(&mut n3, "node-1"); + wait_member_up(&mut n3, "node-2"); + assert_eq!( + inc3_at_n1, inc3_at_n2, + "one node, one incarnation, all observers" + ); + + // Kill node-3 (SIGKILL via Drop): node_down at both survivors. + drop(n3); + wait_member_down(&mut n1, "node-3"); + wait_member_down(&mut n2, "node-3"); + + // Restart node-3 under the same name: it re-dials its seeds and comes + // up everywhere as a new incarnation — never the ghost's. + let mut n3b = spawn_mesh_node("node-3", &format!("node-1={a1},node-2={a2}"), None); + let _ = n3b.wait_listening(); + let inc3b_at_n1 = wait_member_up(&mut n1, "node-3"); + let inc3b_at_n2 = wait_member_up(&mut n2, "node-3"); + assert_eq!(inc3b_at_n1, inc3b_at_n2); + assert_ne!( + inc3_at_n1, inc3b_at_n1, + "a restarted node must be distinguishable from its ghost" + ); + wait_member_up(&mut n3b, "node-1"); + wait_member_up(&mut n3b, "node-2"); +} + +/// A seed that is unreachable at start is not fatal: the connector retries +/// on backoff, and when a node finally appears at that address, the mesh +/// edge forms. +#[test] +fn seed_unreachable_at_start_then_arriving_later() { + maybe_child(ROLES); + + // Pre-reserve an address by binding and immediately closing it (see the + // module docs for the accepted steal window). Dials to it are refused + // until node-b starts there. + let reserved = { + let l = std::net::TcpListener::bind("127.0.0.1:0").expect("bind"); + l.local_addr().expect("addr").to_string() + }; + + let mut a = spawn_mesh_node("node-a", &format!("node-b={reserved}"), None); + let _ = a.wait_listening(); + + // Let a few refused attempts happen before the seed comes up, so the + // retry path is what forms the edge (backoff cap 5s < harness WAIT 10s). + std::thread::sleep(Duration::from_millis(600)); + + let mut b = spawn_mesh_node("node-b", "", Some(&reserved)); + let _ = b.wait_listening(); + + wait_member_up(&mut a, "node-b"); + wait_member_up(&mut b, "node-a"); +}