From 9e490384740c4a051ff986a6859ea3b8c3fe2280 Mon Sep 17 00:00:00 2001 From: claude Date: Fri, 14 Aug 2026 17:17:14 +0000 Subject: [PATCH] =?UTF-8?q?feat(cluster):=20RFC=20010=20c5=20=E2=80=94=20h?= =?UTF-8?q?andshake=20as=20a=20pure=20state=20machine?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- src/cluster.rs | 4 +- src/cluster/handshake.rs | 173 ++++++++++++++++++++++++++ tests/cluster_handshake.rs | 249 +++++++++++++++++++++++++++++++++++++ 3 files changed, 425 insertions(+), 1 deletion(-) create mode 100644 src/cluster/handshake.rs create mode 100644 tests/cluster_handshake.rs diff --git a/src/cluster.rs b/src/cluster.rs index 950900c..16b78b8 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -2,7 +2,9 @@ //! //! c1: feature flag + optional deps. c2: the owned envelope. c3: the //! 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 handshake; pub mod transport; diff --git a/src/cluster/handshake.rs b/src/cluster/handshake.rs new file mode 100644 index 0000000..9e05f32 --- /dev/null +++ b/src/cluster/handshake.rs @@ -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, + }, + } + } +} diff --git a/tests/cluster_handshake.rs b/tests/cluster_handshake.rs new file mode 100644 index 0000000..0eede87 --- /dev/null +++ b/tests/cluster_handshake.rs @@ -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:?}"), + } + } +}