feat(cluster): RFC 010 c5 — handshake as a pure state machine
Frames in, actions out — no IO, no clocks, no actors; the c6 connection actor will drive it. Initiator (dial: emit Hello, interpret the single response) and Responder (accept: judge the first frame) as consuming-self machines; check order proto -> hash -> name -> tie-break. Driver-supplied HelloCtx carries the two facts the pure machine cannot know (name claimed, own dial in flight). Tie-break ratified as a wire fact: the smaller name's dial survives; the losing inbound closes silently (both ends compute the same verdict, no reject frame needed). Peer's own name offered => NameTaken. build_hash is config-supplied; derivation lands with c6.
This commit is contained in:
+3
-1
@@ -2,7 +2,9 @@
|
|||||||
//!
|
//!
|
||||||
//! 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. Everything above them lands in later chunks.
|
//! loopback impls. c5: the handshake state machine. Everything above them
|
||||||
|
//! lands in later chunks.
|
||||||
|
|
||||||
pub mod envelope;
|
pub mod envelope;
|
||||||
|
pub mod handshake;
|
||||||
pub mod transport;
|
pub mod transport;
|
||||||
|
|||||||
@@ -0,0 +1,173 @@
|
|||||||
|
//! RFC 010 c5 — the handshake as a pure state machine.
|
||||||
|
//!
|
||||||
|
//! Frames in, actions out — no IO, no clocks, no actors. The c6 connection
|
||||||
|
//! actor drives these machines and executes their actions; everything
|
||||||
|
//! time-shaped (handshake deadline, heartbeats) lives there.
|
||||||
|
|
||||||
|
use crate::cluster::envelope::{Frame, NodeMeta, RejectReason, PROTO_VERSION};
|
||||||
|
use crate::pg::Incarnation;
|
||||||
|
|
||||||
|
/// This node's identity and metadata, as offered in (or checked against) a
|
||||||
|
/// `Hello`.
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
|
pub struct Local {
|
||||||
|
pub node_name: String,
|
||||||
|
pub incarnation: Incarnation,
|
||||||
|
pub build_hash: u64,
|
||||||
|
pub meta: NodeMeta,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The peer identity a successful handshake yields (what c7 feeds `node_up`).
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct Peer {
|
||||||
|
pub node_name: String,
|
||||||
|
pub incarnation: Incarnation,
|
||||||
|
pub meta: NodeMeta,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Driver-supplied context for an inbound `Hello` — knowledge the pure
|
||||||
|
/// machine cannot have (c6 owns the connection table and dial set).
|
||||||
|
#[derive(Debug, Clone, Copy, Default)]
|
||||||
|
pub struct HelloCtx {
|
||||||
|
/// The offered name is already claimed by an established peer.
|
||||||
|
pub name_claimed: bool,
|
||||||
|
/// We have our own dial in flight to this peer name.
|
||||||
|
pub dialing_this_peer: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Simultaneous-connect tie-break: does the connection dialed by
|
||||||
|
/// `dialer_name` survive against the reverse dial?
|
||||||
|
/// The rule (ratified 2026-08-14, a wire-protocol fact): the connection
|
||||||
|
/// dialed by the lexicographically **smaller** name survives. Both ends know
|
||||||
|
/// both names, so both compute the same verdict — which is why the losing
|
||||||
|
/// side may close silently instead of sending a reject.
|
||||||
|
pub fn dial_wins(dialer_name: &str, acceptor_name: &str) -> bool {
|
||||||
|
dialer_name < acceptor_name
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Dial side: emits `Hello` at construction, interprets the single response.
|
||||||
|
#[must_use]
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct Initiator(());
|
||||||
|
|
||||||
|
/// What the dial side's response frame meant.
|
||||||
|
#[must_use]
|
||||||
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
|
pub enum InitiatorOutcome {
|
||||||
|
Established(Peer),
|
||||||
|
Rejected(RejectReason),
|
||||||
|
/// Protocol violation before the ack — close. Carries the offending frame.
|
||||||
|
Failed(Frame),
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Initiator {
|
||||||
|
/// Start a dial-side handshake: the returned frame is the `Hello` to
|
||||||
|
/// send; the returned machine is the right to interpret the response.
|
||||||
|
pub fn new(local: &Local) -> (Self, Frame) {
|
||||||
|
let hello = Frame::Hello {
|
||||||
|
proto_version: PROTO_VERSION,
|
||||||
|
build_hash: local.build_hash,
|
||||||
|
node_name: local.node_name.clone(),
|
||||||
|
incarnation: local.incarnation,
|
||||||
|
meta: local.meta.clone(),
|
||||||
|
};
|
||||||
|
(Initiator(()), hello)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Interpret the response. The `HelloAck` carries no hash or version —
|
||||||
|
/// the responder already checked ours against its own, and equality is
|
||||||
|
/// symmetric, so a one-sided check is sound.
|
||||||
|
pub fn on_frame(self, frame: Frame) -> InitiatorOutcome {
|
||||||
|
match frame {
|
||||||
|
Frame::HelloAck {
|
||||||
|
node_name,
|
||||||
|
incarnation,
|
||||||
|
meta,
|
||||||
|
} => InitiatorOutcome::Established(Peer {
|
||||||
|
node_name,
|
||||||
|
incarnation,
|
||||||
|
meta,
|
||||||
|
}),
|
||||||
|
Frame::HelloReject { reason } => InitiatorOutcome::Rejected(reason),
|
||||||
|
other => InitiatorOutcome::Failed(other),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Accept side: awaits exactly one `Hello`, answers or closes.
|
||||||
|
#[must_use]
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct Responder {
|
||||||
|
local: Local,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// What to do with an inbound connection's first frame.
|
||||||
|
#[must_use]
|
||||||
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
|
pub enum ResponderOutcome {
|
||||||
|
/// Send the ack; the connection is established.
|
||||||
|
Accepted { reply: Frame, peer: Peer },
|
||||||
|
/// Send the reject, then close.
|
||||||
|
Rejected { reply: Frame, reason: RejectReason },
|
||||||
|
/// Lost the simultaneous-connect tie-break: close silently, no frame.
|
||||||
|
TieBreakLoss,
|
||||||
|
/// Protocol violation before Hello — close, no reply. Carries the frame.
|
||||||
|
Failed(Frame),
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Responder {
|
||||||
|
pub fn new(local: Local) -> Self {
|
||||||
|
Responder { local }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Judge the connection's first frame. Check order is proto → hash →
|
||||||
|
/// name → tie-break: validity before identity. `HelloReject` is the
|
||||||
|
/// cross-version compatibility anchor, so a version-mismatched peer
|
||||||
|
/// still gets one.
|
||||||
|
pub fn on_frame(self, frame: Frame, ctx: HelloCtx) -> ResponderOutcome {
|
||||||
|
let Frame::Hello {
|
||||||
|
proto_version,
|
||||||
|
build_hash,
|
||||||
|
node_name,
|
||||||
|
incarnation,
|
||||||
|
meta,
|
||||||
|
} = frame
|
||||||
|
else {
|
||||||
|
return ResponderOutcome::Failed(frame);
|
||||||
|
};
|
||||||
|
|
||||||
|
let reject = |reason| ResponderOutcome::Rejected {
|
||||||
|
reply: Frame::HelloReject { reason },
|
||||||
|
reason,
|
||||||
|
};
|
||||||
|
|
||||||
|
if proto_version != PROTO_VERSION {
|
||||||
|
return reject(RejectReason::ProtoVersion);
|
||||||
|
}
|
||||||
|
if build_hash != self.local.build_hash {
|
||||||
|
return reject(RejectReason::HashMismatch);
|
||||||
|
}
|
||||||
|
if node_name == self.local.node_name || ctx.name_claimed {
|
||||||
|
return reject(RejectReason::NameTaken);
|
||||||
|
}
|
||||||
|
// Simultaneous connect: the inbound frame is the peer's dial. If our
|
||||||
|
// own in-flight dial wins instead, drop this one silently — the peer
|
||||||
|
// computes the same verdict (see `dial_wins`).
|
||||||
|
if ctx.dialing_this_peer && !dial_wins(&node_name, &self.local.node_name) {
|
||||||
|
return ResponderOutcome::TieBreakLoss;
|
||||||
|
}
|
||||||
|
|
||||||
|
ResponderOutcome::Accepted {
|
||||||
|
reply: Frame::HelloAck {
|
||||||
|
node_name: self.local.node_name,
|
||||||
|
incarnation: self.local.incarnation,
|
||||||
|
meta: self.local.meta,
|
||||||
|
},
|
||||||
|
peer: Peer {
|
||||||
|
node_name,
|
||||||
|
incarnation,
|
||||||
|
meta,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,249 @@
|
|||||||
|
//! RFC 010 c5 — handshake state-machine tests (roadmap: happy path; hash
|
||||||
|
//! mismatch; proto-version mismatch; name already claimed; simultaneous-connect
|
||||||
|
//! tie-break; garbage before Hello). Pure — no IO, no actors, no runtime.
|
||||||
|
#![cfg(feature = "cluster")]
|
||||||
|
|
||||||
|
use smarm::cluster::envelope::{Frame, NodeMeta, RejectReason, PROTO_VERSION};
|
||||||
|
use smarm::cluster::handshake::{
|
||||||
|
dial_wins, HelloCtx, Initiator, InitiatorOutcome, Local, Responder, ResponderOutcome,
|
||||||
|
};
|
||||||
|
use smarm::pg::Incarnation;
|
||||||
|
|
||||||
|
const HASH: u64 = 0xDEAD_BEEF_CAFE_F00D;
|
||||||
|
|
||||||
|
fn local(name: &str) -> Local {
|
||||||
|
Local {
|
||||||
|
node_name: name.into(),
|
||||||
|
incarnation: Incarnation::new(7),
|
||||||
|
build_hash: HASH,
|
||||||
|
meta: NodeMeta {
|
||||||
|
role: "worker".into(),
|
||||||
|
region: "eu-west".into(),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The Hello that `Initiator::new(&local(name))` emits, built by hand.
|
||||||
|
fn hello_from(name: &str) -> Frame {
|
||||||
|
let l = local(name);
|
||||||
|
Frame::Hello {
|
||||||
|
proto_version: PROTO_VERSION,
|
||||||
|
build_hash: l.build_hash,
|
||||||
|
node_name: l.node_name,
|
||||||
|
incarnation: l.incarnation,
|
||||||
|
meta: l.meta,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn happy_path_establishes_both_ends() {
|
||||||
|
// alpha dials beta.
|
||||||
|
let (initiator, hello) = Initiator::new(&local("alpha"));
|
||||||
|
assert_eq!(hello, hello_from("alpha"), "initiator emits its identity");
|
||||||
|
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
let (reply, peer) = match responder.on_frame(hello, HelloCtx::default()) {
|
||||||
|
ResponderOutcome::Accepted { reply, peer } => (reply, peer),
|
||||||
|
other => panic!("expected Accepted, got {other:?}"),
|
||||||
|
};
|
||||||
|
assert_eq!(peer.node_name, "alpha");
|
||||||
|
assert_eq!(peer.incarnation, Incarnation::new(7));
|
||||||
|
assert_eq!(peer.meta.role, "worker");
|
||||||
|
|
||||||
|
// The ack carries the responder's identity, no hash/version (one-sided
|
||||||
|
// check — sound because equality is symmetric).
|
||||||
|
let l = local("beta");
|
||||||
|
assert_eq!(
|
||||||
|
reply,
|
||||||
|
Frame::HelloAck {
|
||||||
|
node_name: l.node_name,
|
||||||
|
incarnation: l.incarnation,
|
||||||
|
meta: l.meta,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
|
||||||
|
match initiator.on_frame(reply) {
|
||||||
|
InitiatorOutcome::Established(peer) => {
|
||||||
|
assert_eq!(peer.node_name, "beta");
|
||||||
|
assert_eq!(peer.incarnation, Incarnation::new(7));
|
||||||
|
assert_eq!(peer.meta.region, "eu-west");
|
||||||
|
}
|
||||||
|
other => panic!("expected Established, got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn hash_mismatch_rejected() {
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
let hello = Frame::Hello {
|
||||||
|
proto_version: PROTO_VERSION,
|
||||||
|
build_hash: HASH ^ 1,
|
||||||
|
node_name: "alpha".into(),
|
||||||
|
incarnation: Incarnation::new(7),
|
||||||
|
meta: local("alpha").meta,
|
||||||
|
};
|
||||||
|
match responder.on_frame(hello, HelloCtx::default()) {
|
||||||
|
ResponderOutcome::Rejected { reply, reason } => {
|
||||||
|
assert_eq!(reason, RejectReason::HashMismatch);
|
||||||
|
assert_eq!(reply, Frame::HelloReject { reason });
|
||||||
|
}
|
||||||
|
other => panic!("expected Rejected, got {other:?}"),
|
||||||
|
}
|
||||||
|
|
||||||
|
// The dialer side of the same story: a reject frame comes back.
|
||||||
|
let (initiator, _hello) = Initiator::new(&local("alpha"));
|
||||||
|
match initiator.on_frame(Frame::HelloReject {
|
||||||
|
reason: RejectReason::HashMismatch,
|
||||||
|
}) {
|
||||||
|
InitiatorOutcome::Rejected(RejectReason::HashMismatch) => {}
|
||||||
|
other => panic!("expected Rejected(HashMismatch), got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn proto_version_mismatch_rejected_and_checked_first() {
|
||||||
|
// Both proto and hash wrong: proto wins — nothing after the version can
|
||||||
|
// be trusted, and HelloReject is the cross-version compatibility anchor.
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
let hello = Frame::Hello {
|
||||||
|
proto_version: PROTO_VERSION + 1,
|
||||||
|
build_hash: HASH ^ 1,
|
||||||
|
node_name: "alpha".into(),
|
||||||
|
incarnation: Incarnation::new(7),
|
||||||
|
meta: local("alpha").meta,
|
||||||
|
};
|
||||||
|
match responder.on_frame(hello, HelloCtx::default()) {
|
||||||
|
ResponderOutcome::Rejected { reason, .. } => {
|
||||||
|
assert_eq!(reason, RejectReason::ProtoVersion);
|
||||||
|
}
|
||||||
|
other => panic!("expected Rejected, got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn claimed_name_rejected() {
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
let ctx = HelloCtx {
|
||||||
|
name_claimed: true,
|
||||||
|
dialing_this_peer: false,
|
||||||
|
};
|
||||||
|
match responder.on_frame(hello_from("alpha"), ctx) {
|
||||||
|
ResponderOutcome::Rejected { reason, .. } => {
|
||||||
|
assert_eq!(reason, RejectReason::NameTaken);
|
||||||
|
}
|
||||||
|
other => panic!("expected Rejected, got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn own_name_offered_rejected_as_name_taken() {
|
||||||
|
// Self-connect or genuine collision: the responder's own name arrives.
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
match responder.on_frame(hello_from("beta"), HelloCtx::default()) {
|
||||||
|
ResponderOutcome::Rejected { reason, .. } => {
|
||||||
|
assert_eq!(reason, RejectReason::NameTaken);
|
||||||
|
}
|
||||||
|
other => panic!("expected Rejected, got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn hash_checked_before_name() {
|
||||||
|
// Wrong hash AND claimed name: hash wins (validity before identity).
|
||||||
|
let responder = Responder::new(local("beta"));
|
||||||
|
let hello = Frame::Hello {
|
||||||
|
proto_version: PROTO_VERSION,
|
||||||
|
build_hash: HASH ^ 1,
|
||||||
|
node_name: "alpha".into(),
|
||||||
|
incarnation: Incarnation::new(7),
|
||||||
|
meta: local("alpha").meta,
|
||||||
|
};
|
||||||
|
let ctx = HelloCtx {
|
||||||
|
name_claimed: true,
|
||||||
|
dialing_this_peer: false,
|
||||||
|
};
|
||||||
|
match responder.on_frame(hello, ctx) {
|
||||||
|
ResponderOutcome::Rejected { reason, .. } => {
|
||||||
|
assert_eq!(reason, RejectReason::HashMismatch);
|
||||||
|
}
|
||||||
|
other => panic!("expected Rejected, got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn dial_wins_is_deterministic_and_antisymmetric() {
|
||||||
|
// The smaller name's dial survives; both ends compute the same verdict.
|
||||||
|
assert!(dial_wins("alpha", "beta"));
|
||||||
|
assert!(!dial_wins("beta", "alpha"));
|
||||||
|
for (a, b) in [("a", "b"), ("node-1", "node-2"), ("x", "xx")] {
|
||||||
|
assert_ne!(dial_wins(a, b), dial_wins(b, a), "({a}, {b})");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn simultaneous_connect_exactly_one_side_accepts() {
|
||||||
|
// alpha and beta dial each other at once. Each responder sees the peer's
|
||||||
|
// Hello while its own dial is in flight.
|
||||||
|
let ctx = HelloCtx {
|
||||||
|
name_claimed: false,
|
||||||
|
dialing_this_peer: true,
|
||||||
|
};
|
||||||
|
|
||||||
|
// On beta: inbound is alpha's dial; alpha < beta, so the inbound wins.
|
||||||
|
let on_beta = Responder::new(local("beta")).on_frame(hello_from("alpha"), ctx);
|
||||||
|
assert!(
|
||||||
|
matches!(on_beta, ResponderOutcome::Accepted { .. }),
|
||||||
|
"beta must accept alpha's dial, got {on_beta:?}"
|
||||||
|
);
|
||||||
|
|
||||||
|
// On alpha: inbound is beta's dial; it loses — close silently, no frame
|
||||||
|
// (ratified: both ends can compute the outcome, a reject adds nothing).
|
||||||
|
let on_alpha = Responder::new(local("alpha")).on_frame(hello_from("beta"), ctx);
|
||||||
|
assert!(
|
||||||
|
matches!(on_alpha, ResponderOutcome::TieBreakLoss),
|
||||||
|
"alpha must silently drop beta's dial, got {on_alpha:?}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn tiebreak_loss_only_applies_when_dialing() {
|
||||||
|
// Same inbound Hello, no dial in flight: plain accept.
|
||||||
|
let on_alpha = Responder::new(local("alpha")).on_frame(hello_from("beta"), HelloCtx::default());
|
||||||
|
assert!(matches!(on_alpha, ResponderOutcome::Accepted { .. }));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn garbage_before_hello_fails_without_reply() {
|
||||||
|
// Any valid-but-wrong frame before Hello is a protocol violation: close,
|
||||||
|
// no reject frame. (Undecodable bytes are the codec's Err, not ours.)
|
||||||
|
for frame in [
|
||||||
|
Frame::Heartbeat,
|
||||||
|
Frame::HelloAck {
|
||||||
|
node_name: "alpha".into(),
|
||||||
|
incarnation: Incarnation::new(7),
|
||||||
|
meta: local("alpha").meta,
|
||||||
|
},
|
||||||
|
Frame::Demonitor { monitor_id: 3 },
|
||||||
|
] {
|
||||||
|
let out = Responder::new(local("beta")).on_frame(frame.clone(), HelloCtx::default());
|
||||||
|
match out {
|
||||||
|
ResponderOutcome::Failed(f) => assert_eq!(f, frame),
|
||||||
|
other => panic!("expected Failed({frame:?}), got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn garbage_before_ack_fails_the_initiator() {
|
||||||
|
for frame in [
|
||||||
|
Frame::Heartbeat,
|
||||||
|
hello_from("beta"),
|
||||||
|
Frame::Demonitor { monitor_id: 3 },
|
||||||
|
] {
|
||||||
|
let (initiator, _hello) = Initiator::new(&local("alpha"));
|
||||||
|
match initiator.on_frame(frame.clone()) {
|
||||||
|
InitiatorOutcome::Failed(f) => assert_eq!(f, frame),
|
||||||
|
other => panic!("expected Failed({frame:?}), got {other:?}"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user