Nothing local is remotely reachable by default (RFC §4). expose(Name<M>)
marks a name remotely addressable and registers M's decoder under
type_hash::<M>(); expose_type::<M>() registers only the decoder (the
reply-to path). exposed_names() is the auditable remote surface.
D3's watchable fold, resolved against the code as it stands (stated in the
module docs): register() ALREADY stamps every named holder watchable ('no
successfully-registered actor can die unflagged', registry.rs), so an
exposed name's holder needs no extra mark — and re-registration after a
holder's death re-stamps the new holder for free, which a per-tenancy mark
taken at expose time could not do. The cluster's own mark_watchable
set-site is therefore the pid crossing the wire (frame serialization, c10)
— the exact analog of the membrane crossing. c8 adds only the name/type
state neither the registry nor slot bits can carry. No new pid registry;
RFC §4 honored.
One-viable calls, flagged:
- State lives on RuntimeInner (the pg pattern: leaf RawMutex field,
cfg-gated behind cluster, zero-cost-when-off per c1) — c9's inbound
decode consults it per frame; manager-held state would serialize every
remote delivery through one gen_server.
- type_hash = FNV-1a 64 (fixed seed: the offset basis) over TypeId: a
constant of the binary — stable across runs of the same build (the scope
the build-hash handshake reduces the mesh to), deliberately not across
builds. Collisions degrade to decode error / refused channel, never a
misroute (the NoChannel guarantee, RFC §3).
- Decoder = decode-and-deliver-to-pid Arc closure capturing M (the one
typed site): decode_payload then send_dyn. Wire-name → pid resolution
stays OUTSIDE — that is c9's single seam, which calls decode_deliver.
Arc so the call happens with the exposure lock RELEASED: send_dyn takes
the registry lock, a mutual Leaf (the runtime asserts on nesting — caught
live by the first test run).
- expose is a name-level fact, valid for an unregistered name (names
late-bind; c9 resolves per delivery).
tests/cluster_expose.rs 5/0 stable x5, purely local per roadmap:
exposed/unexposed lookup + audit listing; decoder registration and the
delivery contract (happy path into a registered String channel; unknown
hash; corrupt bytes; wrong channel refused — never misrouted); distinct
types distinct hashes; expose/bridge-crossing agreement via the shared
watchable observable (terminal_reason after holder death); hash stability
across runs in the same binary via a c4-harness re-exec. Payload types are
std types — the crate's serde is derive-less by design, user crates bring
their own derive.
All cluster suites regression-clean (envelope 15, handshake 11, transport
11, lifecycle 1, liveness 3, connect 9, two_node 3, membership 4, mesh 2);
clippy --lib green both configs; fmt clean; default build compiles.
228 lines
8.3 KiB
Rust
228 lines
8.3 KiB
Rust
//! RFC 010 — clustering (smarm⇄smarm, explicit remote boundary).
|
|
//!
|
|
//! c1: feature flag + optional deps. c2: the owned envelope. c3: the
|
|
//! transport trait (control connection), framed codec, and the TCP +
|
|
//! loopback impls. c5: the handshake state machine. c6: the connection
|
|
//! [`manager`] (registry) and per-peer connection actors ([`conn`]), started
|
|
//! as an explicit supervision subtree, plus the handshake on the
|
|
//! accept/connect path ([`connect`]). Everything above them lands in later
|
|
//! chunks.
|
|
|
|
pub mod conn;
|
|
pub mod connect;
|
|
pub mod connector;
|
|
pub mod discovery;
|
|
pub mod envelope;
|
|
pub mod expose;
|
|
pub mod handshake;
|
|
pub mod manager;
|
|
pub mod membership;
|
|
pub mod transport;
|
|
|
|
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 expose::{expose, expose_type, type_hash, DeliverError};
|
|
pub use manager::{Manager, MANAGER};
|
|
pub use membership::{subscribe, view, MembershipEvents, NodeEvent, NodeInfo};
|
|
|
|
/// c6d — the derived build hash for [`handshake::LocalNode::build_hash`]:
|
|
/// two builds may mesh only when this matches, and it is a pure function of
|
|
/// the compile-time inputs that define wire compatibility today — the exact
|
|
/// toolchain (`rustc -V`), the declared feature set, and
|
|
/// [`envelope::PROTO_VERSION`]. FNV-1a 64 over the build-script string, then
|
|
/// the proto version folded byte-wise, so a proto bump moves the hash even
|
|
/// on an identical toolchain. The domain is deliberately lean and
|
|
/// tightenable later without a wire change — it is just a `u64`.
|
|
pub const BUILD_HASH: u64 = fold_u32(
|
|
fnv1a64(env!("SMARM_BUILD_HASH_INPUTS").as_bytes()),
|
|
envelope::PROTO_VERSION,
|
|
);
|
|
|
|
/// FNV-1a 64 (const so [`BUILD_HASH`] is a compile-time fact).
|
|
const fn fnv1a64(bytes: &[u8]) -> u64 {
|
|
let mut h: u64 = 0xcbf2_9ce4_8422_2325;
|
|
let mut i = 0;
|
|
while i < bytes.len() {
|
|
h ^= bytes[i] as u64;
|
|
h = h.wrapping_mul(0x0000_0100_0000_01b3);
|
|
i += 1;
|
|
}
|
|
h
|
|
}
|
|
|
|
/// Continue an FNV-1a state over a `u32`'s little-endian bytes.
|
|
const fn fold_u32(mut h: u64, v: u32) -> u64 {
|
|
let b = v.to_le_bytes();
|
|
let mut i = 0;
|
|
while i < b.len() {
|
|
h ^= b[i] as u64;
|
|
h = h.wrapping_mul(0x0000_0100_0000_01b3);
|
|
i += 1;
|
|
}
|
|
h
|
|
}
|
|
|
|
/// 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>,
|
|
}
|
|
|
|
/// 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))
|
|
.run()
|
|
});
|
|
while gen_server::whereis_server(MANAGER).is_none() {
|
|
sleep(Duration::from_millis(1));
|
|
}
|
|
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
|
|
/// named manager, then parks on the manager's own termination so this actor's
|
|
/// lifetime tracks the manager's — the supervisor's restart accounting keys off
|
|
/// this actor exiting.
|
|
fn manager_child() {
|
|
let m = match GenServerBuilder::new(Manager::new()).named(MANAGER).start() {
|
|
Ok(m) => m,
|
|
// Name still held by a not-yet-reaped prior instance: return and let
|
|
// the supervisor retry under its restart policy.
|
|
Err(_) => return,
|
|
};
|
|
let _ = monitor(m.pid()).rx.recv();
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
/// The hash core against the published FNV-1a 64 test vectors — the
|
|
/// contract is "this is FNV-1a", not "whatever the fn does".
|
|
#[test]
|
|
fn fnv1a64_known_vectors() {
|
|
assert_eq!(fnv1a64(b""), 0xcbf2_9ce4_8422_2325);
|
|
assert_eq!(fnv1a64(b"a"), 0xaf63_dc4c_8601_ec8c);
|
|
assert_eq!(fnv1a64(b"foobar"), 0x85944171f73967e8);
|
|
}
|
|
|
|
/// Folding the proto version continues the same FNV state: identical
|
|
/// inputs with a different version must land on a different hash.
|
|
#[test]
|
|
fn proto_version_moves_the_hash() {
|
|
let base = fnv1a64(b"same-toolchain;features=CLUSTER");
|
|
assert_ne!(fold_u32(base, 1), fold_u32(base, 2));
|
|
// And it equals hashing the bytes in one pass — the fold is a
|
|
// continuation, not a second construction.
|
|
let mut all = b"same-toolchain;features=CLUSTER".to_vec();
|
|
all.extend_from_slice(&1u32.to_le_bytes());
|
|
assert_eq!(fold_u32(base, 1), fnv1a64(&all));
|
|
}
|
|
|
|
/// The derived constant exists, is compile-time, and is not degenerate.
|
|
#[test]
|
|
fn build_hash_is_nonzero() {
|
|
const H: u64 = BUILD_HASH;
|
|
assert_ne!(H, 0);
|
|
}
|
|
}
|