feat(cluster): RFC 010 c7b — discovery Strategy, static seeds, connector dial loop
Phase 2 gate: 3-node mesh under the subprocess harness, repeatable (10/10).
Strategy (ratified): push-based, spawned as its own actor by the connector —
it emits Discovery events into a channel whenever it learns something and
may run forever; the connector owns all retry/backoff state. StaticSeeds
announces its list once and exits. Discovery is #[non_exhaustive] and
additive-only (candidates announced, never withdrawn) so expiry can land
later without breaking strategies.
One-viable correction to the ratified Discovery shape, flagged: a candidate
is a (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.
Connector: plain select-loop actor (the c6 shape) folding cmd inbox,
discovery stream, membership stream, and the earliest retry deadline into
one wait. It tracks who is up by SUBSCRIBING TO MEMBERSHIP like any
consumer — first consumer of c7a's snapshot-then-stream surface, no
privileged channel into the manager. Backoff: 250ms doubling to a 5s cap
(the c6c class of one-viable constants), reset on node_up; node_down
schedules a prompt redial with a fresh sequence. A candidate bearing the
local name is parked (that seed is us); every other failure retries — in
particular NameTaken 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
acceptor's deliberate serialization (each attempt bounded by the connect +
handshake deadlines).
cluster::start(Config {node_name, meta, listen_addr, strategy}) is now the
integrated node start: supervised manager + acceptor + connector. It
completes the node identity: build_hash = cluster::BUILD_HASH (first
consumer, closing the c6d loose end) and incarnation = self_incarnation()
— unix-epoch MILLIS truncated to u32, not seconds: a supervised
crash-and-restart inside one second is routine, and seconds would collide
the ghost with its successor. Cluster handle: local_addr()/local()/
shutdown(); drop stops acceptor+connector loops, manager subtree detaches
(same split as AcceptorHandle alone).
Roadmap-binding, asserted in review: no consumer touches the connection
table — Manager.conns and ConnEntry stay private; the only exposures are
Call::Peers (sorted names, pre-existing) and the membership surface.
tests/cluster_mesh.rs 2/0, 10/10 flake runs: (1) 3-node mesh forms; kill
one (SIGKILL via Drop, per the retractable-state trap: roles park forever)
=> node_down at both survivors; restart same name => new incarnation at
every observer, distinguishable from the ghost; (2) seed unreachable at
start (pre-reserved closed port; accepted micro steal-window, documented)
then arriving later => edge forms via the retry path. All cluster suites
regression-clean (envelope 15, handshake 11, transport 11, lifecycle 1,
liveness 3, connect 9, two_node 3, membership 4); clippy --lib green both
configs; fmt clean; default build compiles.
This commit is contained in:
+100
-14
@@ -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<dyn Strategy>,
|
||||
}
|
||||
|
||||
/// 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<Cluster> {
|
||||
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
|
||||
|
||||
@@ -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<Cmd>,
|
||||
}
|
||||
|
||||
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<dyn Transport>,
|
||||
local: Local,
|
||||
strategy: Box<dyn Strategy>,
|
||||
) -> ConnectorHandle {
|
||||
let (cmd_tx, cmd_rx) = channel();
|
||||
spawn(move || run(transport, local, strategy, cmd_rx));
|
||||
ConnectorHandle { cmd_tx }
|
||||
}
|
||||
|
||||
fn run(
|
||||
transport: Box<dyn Transport>,
|
||||
local: Local,
|
||||
strategy: Box<dyn Strategy>,
|
||||
cmd_rx: Receiver<Cmd>,
|
||||
) {
|
||||
// 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<Candidate> = Vec::new();
|
||||
let mut connected: HashSet<String> = HashSet::new();
|
||||
let mut names: HashMap<NodeId, String> = 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<Cmd>) -> 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<Discovery>, local: &Local, cands: &mut Vec<Candidate>) -> 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<NodeEvent>,
|
||||
connected: &mut HashSet<String>,
|
||||
names: &mut HashMap<NodeId, String>,
|
||||
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,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<Self>, out: Sender<Discovery>);
|
||||
}
|
||||
|
||||
/// 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<Item = (impl Into<String>, impl Into<String>)>) -> Self {
|
||||
StaticSeeds {
|
||||
seeds: seeds
|
||||
.into_iter()
|
||||
.map(|(n, a)| (n.into(), a.into()))
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Strategy for StaticSeeds {
|
||||
fn run(self: Box<Self>, out: Sender<Discovery>) {
|
||||
for (name, addr) in self.seeds {
|
||||
if out.send(Discovery::Candidate { name, addr }).is_err() {
|
||||
return; // connector gone; nobody to discover for
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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 <addr>`, `MEMBER-UP <name> inc=<n>`, `MEMBER-DOWN <name>`.
|
||||
//! 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<NodeId, String> = 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 <peer> inc=<n>` 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");
|
||||
}
|
||||
Reference in New Issue
Block a user