Tree snapshot of d9c62a8 (2026-08-18). The 20 source commits between
16ef583 (c9) and d9c62a8 were never pushed and the clone that held them
was lost; this commit carries their combined tree verbatim so the build
history stays auditable from the c1–c9 commits below it. Original
hashes as recorded in the session handoff:
c10 f03e94d pid targeting + auto-serialization (RemotePid, D14 name
on the wire); Phase 3 gate
c11 7ef4bad DownReason::Disconnected, wire tag 5
c12 d124162 remote monitors (Monitor/Demonitor/Down frames)
c13 9de967b connection-loss synthesis (A+B: Monitors::teardown +
unread-command Disconnected); Phase 4 gate
c14 7e822b7 eager pg eviction (reaper actor, ReaperInboxes)
dbe1a22 InboundVerdict::label(), trace::Event::ClusterInbound
31a9877 tests/channel.rs monitor-churn target gated on `go`
653559e Discovery::Withdrawn{name, addr}
c15 b41d76e distributed pg: Sync on NodeUp, Join/Leave broadcast,
NodeDown sweep, members_all; PgMsg wire type
c16 fafa881 pick_any / dispatch_any; Phase 5 complete
Phase 6 Tier A:
195c73e p4 NodeEvent::NodeDown(NodeInfo)
48fd766 p1 connector Candidate{name, addr, state}
ce8cf99 p2+p7 conn.rs select arms as Vec<Arm>; Outbound::Drained
bf24988 p6 RemotePid::from_local -> Option
9ae0380 p3 PeerStanding{Free, Claimed, Dialing}
c7d62a1 p11 cluster::Timing knobs, threaded by value
46f171d p11 cluster_disconnect un-ignored on SMARM_FAST_TIMING
Phase 6 Tier B:
391a9ae p5 cluster::RemoteDownReason{Local, Disconnected};
DownReason::Disconnected removed from core
7ddd908 p9 pg ctl channel unconditional, one cfg seam at spawn
d9c62a8 PeerNameMismatch parks the candidate; ClusterDial trace
Verified at d9c62a8: default 361/0, cluster 448/0, clippy --lib on
default / cluster / cluster+smarm-trace, fmt, 10x flake on
cluster_dial_mismatch, 5x on cluster_pg.
176 lines
6.9 KiB
Rust
176 lines
6.9 KiB
Rust
//! 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, Timing};
|
|
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)),
|
|
timing: Timing::default(),
|
|
})
|
|
.expect("listener binds");
|
|
println!("LISTENING {}", cluster.local_addr());
|
|
|
|
let events = subscribe().expect("manager is up");
|
|
loop {
|
|
match events.rx.recv() {
|
|
Ok(NodeEvent::NodeUp(info)) => {
|
|
println!("MEMBER-UP {} inc={}", info.name, info.incarnation.get());
|
|
}
|
|
Ok(NodeEvent::NodeDown(info)) => {
|
|
println!("MEMBER-DOWN {}", info.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");
|
|
}
|