diff --git a/src/cluster.rs b/src/cluster.rs index 851a011..e2b2225 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -13,6 +13,7 @@ pub mod connect; pub mod connector; pub mod discovery; pub mod envelope; +pub mod expose; pub mod handshake; pub mod manager; pub mod membership; @@ -36,6 +37,7 @@ 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}; diff --git a/src/cluster/expose.rs b/src/cluster/expose.rs new file mode 100644 index 0000000..e35de5a --- /dev/null +++ b/src/cluster/expose.rs @@ -0,0 +1,238 @@ +//! RFC 010 c8 — explicit exposure: the node's remote surface, and the +//! fixed-seed type hash. +//! +//! Nothing local is remotely reachable by default (RFC §4 — "a gun needs a +//! safety"). [`expose`] marks a registered name remotely addressable and +//! registers `M`'s decoder under [`type_hash::()`](type_hash); +//! [`expose_type`] registers only the decoder (the reply-to path: a +//! `RemotePid` received in a message is sendable only if `A::Msg`'s +//! decoder was explicitly registered). The exposed set is the node's +//! visible, auditable remote surface ([`exposed_names`]). +//! +//! ## Where the state lives +//! +//! On `RuntimeInner`, the [`pg`](crate::pg) pattern: a leaf-locked table, +//! cfg-gated behind the `cluster` feature (zero-cost-when-off, per c1). +//! Chosen over manager-held state because c9's inbound decode consults it +//! per frame — a hot path that must not serialize every remote delivery +//! through one gen_server. The state resets with the runtime, like every +//! registry. +//! +//! ## The watchable fold (D3), against the code as it stands +//! +//! RFC §4: the exposed set is not a new registry — it folds into the +//! existing `watchable` machinery, one set, two set-sites (a pid crossing +//! the membrane, and expose). Reading the code: `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 here — the guarantee holds by registration, and re-registration +//! after a holder's death re-stamps the new holder for free (a per-tenancy +//! mark taken at expose time could not do that). The cluster's own +//! `mark_watchable` set-site is therefore the **pid crossing the wire** — +//! serialization of a pid into a frame, c10 — the exact analog of the +//! membrane crossing. What lives here is only the name/type-level state +//! neither the registry nor the slot bits can carry: which names are +//! exposed, and how to decode each type hash. +//! +//! ## The hash +//! +//! [`type_hash`] is FNV-1a 64 (fixed seed: the FNV offset basis) over +//! `TypeId`, so it is a constant of the binary: stable across runs of the +//! same build — exactly the scope the build-hash handshake reduces the mesh +//! to — and deliberately *not* stable across builds (scope guard: no +//! cross-version wire compatibility). A collision between two exposed types +//! degrades to a decode error or a refused channel, never a misroute — the +//! local `SendError::NoChannel` guarantee survives the network (RFC §3). +//! +//! ## The decoder contract +//! +//! A decoder is **decode-and-deliver-to-pid**: it captures `M` (the one +//! typed site), decodes the payload, and hands the value to the target's +//! published channel via the registry's own dynamic send. Wire-name → +//! local-pid resolution deliberately stays *outside* — that is c9's single +//! resolution seam, and it calls [`decode_deliver`]. + +use std::any::TypeId; +use std::collections::HashMap; +use std::hash::{Hash, Hasher}; + +use crate::cluster::envelope::{decode_payload, PayloadError}; +use crate::pid::{Name, Pid}; +use crate::registry::{send_dyn, SendError}; +use crate::scheduler::with_runtime; + +/// The fixed-seed `TypeId` → `u64` hash: FNV-1a 64 over the `TypeId`'s hash +/// bytes, seeded with the FNV offset basis. A constant of the binary — see +/// the module docs for scope. +pub fn type_hash() -> u64 { + let mut h = Fnv1a64::new(); + TypeId::of::().hash(&mut h); + h.finish() +} + +/// FNV-1a 64 as a `Hasher`, so `TypeId` (opaque, `Hash`-only) can feed it. +/// Same constants as the const fns in [`crate::cluster`] (BUILD_HASH). +struct Fnv1a64(u64); + +impl Fnv1a64 { + fn new() -> Self { + Fnv1a64(0xcbf2_9ce4_8422_2325) + } +} + +impl Hasher for Fnv1a64 { + fn write(&mut self, bytes: &[u8]) { + for &b in bytes { + self.0 ^= b as u64; + self.0 = self.0.wrapping_mul(0x0000_0100_0000_01b3); + } + } + fn finish(&self) -> u64 { + self.0 + } +} + +/// Why a [`decode_deliver`] did not deliver. Payload-free mirror of the +/// registry's `SendError` where relevant — the caller (c9's inbound path) +/// has only bytes to give back, not a typed message. +#[derive(Debug)] +pub enum DeliverError { + /// No decoder is registered under this hash — the type was never + /// exposed here. + UnknownType, + /// The bytes did not decode as the registered type. + Decode(PayloadError), + /// The target actor is dead (or was never alive). + Dead, + /// The target is live but has no channel for this message type, or that + /// channel is closed — the `NoChannel` guarantee: a decoded value is + /// refused, never misrouted. + WrongChannel, +} + +/// A registered decoder: decode `bytes` as the captured type and deliver to +/// `pid`'s published channel. `Arc`, so [`decode_deliver`] can clone it out +/// from under the exposure lock and call it lock-free — the decoder's +/// `send_dyn` takes the registry lock, and the two are mutual Leaves that +/// must never nest. +type Decoder = std::sync::Arc Result<(), DeliverError> + Send + Sync>; + +/// The exposure state, one per runtime (a `RuntimeInner` field, pg-style). +pub(crate) struct ExposureState { + /// The exposed names: registry key → the type hash it expects. + exposed: HashMap<&'static str, u64>, + /// The decoders: type hash → decode-and-deliver. + decoders: HashMap, +} + +impl ExposureState { + pub(crate) fn new() -> Self { + ExposureState { + exposed: HashMap::new(), + decoders: HashMap::new(), + } + } +} + +/// Mark `name` remotely addressable and register `M`'s decoder under its +/// type hash (so both name-sends and pid-sends of `M` work — RFC §4). +/// Returns the hash. +/// +/// Exposure is a **name-level fact**, independent of who currently holds the +/// name (names late-bind: the registry re-resolves on every send, and c9's +/// seam resolves per delivery). Exposing an unregistered name is therefore +/// valid — deliveries fail with "unresolved" until someone registers it. +/// Idempotent. Must run inside [`run`](crate::run). +pub fn expose(name: Name) -> u64 +where + M: serde::de::DeserializeOwned + Send + 'static, +{ + let h = ensure_decoder::(); + with_runtime(|inner| { + inner.exposure.lock().exposed.insert(name.as_str(), h); + }); + h +} + +/// Register only `M`'s decoder (no name): the reply-to path. Returns the +/// hash. Idempotent. Must run inside [`run`](crate::run). +pub fn expose_type() -> u64 +where + M: serde::de::DeserializeOwned + Send + 'static, +{ + ensure_decoder::() +} + +fn ensure_decoder() -> u64 +where + M: serde::de::DeserializeOwned + Send + 'static, +{ + let h = type_hash::(); + with_runtime(|inner| { + inner + .exposure + .lock() + .decoders + .entry(h) + .or_insert_with(decoder::); + }); + h +} + +/// The one typed site: decode as `M`, deliver via the registry's dynamic +/// send. See the module docs for the error mapping. +fn decoder() -> Decoder +where + M: serde::de::DeserializeOwned + Send + 'static, +{ + std::sync::Arc::new(|pid, bytes| { + let m: M = decode_payload(bytes).map_err(DeliverError::Decode)?; + send_dyn(pid, m).map_err(|e| match e { + SendError::Dead(_) | SendError::Unresolved(_) | SendError::NoMember(_) => { + DeliverError::Dead + } + SendError::NoChannel(_) | SendError::Closed(_) => DeliverError::WrongChannel, + }) + }) +} + +/// The type hash `name` was exposed with, or `None` if it is not exposed. +/// Must run inside [`run`](crate::run). +pub fn exposed_hash(name: &str) -> Option { + with_runtime(|inner| inner.exposure.lock().exposed.get(name).copied()) +} + +/// Whether a decoder is registered under `hash`. Must run inside +/// [`run`](crate::run). +pub fn decoder_registered(hash: u64) -> bool { + with_runtime(|inner| inner.exposure.lock().decoders.contains_key(&hash)) +} + +/// Decode `bytes` under `hash`'s registered decoder and deliver to `pid`. +/// This is the delivery half c9's single resolution seam calls after it has +/// resolved a wire name to a local pid. Must run inside [`run`](crate::run). +pub fn decode_deliver(hash: u64, to: Pid, bytes: &[u8]) -> Result<(), DeliverError> { + // Clone the Arc under the lock, call outside it: the decoder's + // `send_dyn` takes the registry lock — a mutual Leaf with the exposure + // lock (the runtime asserts if Leaves nest). This also keeps unrelated + // deliveries uncoupled from a slow decode. + let d = with_runtime(|inner| inner.exposure.lock().decoders.get(&hash).cloned()); + match d { + Some(d) => d(to, bytes), + None => Err(DeliverError::UnknownType), + } +} + +/// The auditable remote surface: every exposed name and its type hash, +/// unordered. Must run inside [`run`](crate::run). +pub fn exposed_names() -> Vec<(&'static str, u64)> { + with_runtime(|inner| { + inner + .exposure + .lock() + .exposed + .iter() + .map(|(&n, &h)| (n, h)) + .collect() + }) +} diff --git a/src/runtime.rs b/src/runtime.rs index be56e5e..d25e167 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -962,6 +962,12 @@ pub(crate) struct RuntimeInner { /// checks under it read only the atomic slot word, and the eviction path /// keeps it off the send path. pub(crate) process_groups: RawMutex, + /// RFC 010 c8: the exposure registry (exposed names + type-hash decoders). + /// RawMutex Leaf, same discipline as `process_groups`; decoders run under + /// it and are leaf-only by contract (they decode and send — `send_dyn` + /// takes `registry`, never this). cfg-gated: zero-cost-when-off (c1). + #[cfg(feature = "cluster")] + pub(crate) exposure: RawMutex, /// Recycled stacks waiting to be reused by the next spawn. pub(crate) stack_pool: RawMutex>, /// Maximum number of stacks to retain in the pool. @@ -1022,6 +1028,8 @@ impl RuntimeInner { node_id, incarnation, process_groups: RawMutex::new(crate::pg::ProcessGroups::new()), + #[cfg(feature = "cluster")] + exposure: RawMutex::new(crate::cluster::expose::ExposureState::new()), stack_pool: RawMutex::new(Vec::new()), stack_pool_cap, stack_reserve: crate::stack::round_to_pages(stack_reserve), diff --git a/tests/cluster_expose.rs b/tests/cluster_expose.rs new file mode 100644 index 0000000..9aad3af --- /dev/null +++ b/tests/cluster_expose.rs @@ -0,0 +1,161 @@ +//! RFC 010 c8 — exposure registry + type hashing. Purely local, no network. +//! +//! Payload types are std types (`String`, `u64`) because the crate's serde is +//! deliberately derive-less (`default-features = false`) — user crates bring +//! their own derive; the contract here is `DeserializeOwned`. +//! +//! The hash-stability test re-execs the current binary (the c4 harness): the +//! guarantee under test is "stable across runs in the SAME binary" — exactly +//! what the build-hash handshake reduces the mesh to — not stability across +//! builds, which the scope guard explicitly rejects. +#![cfg(feature = "cluster")] + +mod common; + +use common::{maybe_child, spawn_node}; +use smarm::cluster::envelope::encode_payload; +use smarm::cluster::expose::{ + decode_deliver, decoder_registered, expose, expose_type, exposed_hash, exposed_names, + type_hash, DeliverError, +}; +use smarm::monitor::{monitor, terminal_reason, DownReason}; +use smarm::{channel, register, run, spawn, Name}; + +const ROLES: &[(&str, fn())] = &[("hasher", role_hasher)]; + +/// Print the hashes this process computes; the parent (a different run of +/// the same binary) compares against its own. +fn role_hasher() { + println!("HASH-STRING {}", type_hash::()); + println!("HASH-U64 {}", type_hash::()); +} + +const GREETER: Name = Name::new("expose-test.greeter"); + +/// Exposed and unexposed lookup, the returned hash, and the audit listing. +#[test] +fn exposed_and_unexposed_lookup() { + maybe_child(ROLES); + run(|| { + let h = expose(GREETER); + assert_eq!(h, type_hash::()); + assert_eq!(exposed_hash("expose-test.greeter"), Some(h)); + assert_eq!(exposed_hash("never-exposed"), None); + assert!(exposed_names().contains(&("expose-test.greeter", h))); + }); +} + +/// Distinct types land on distinct hashes (FNV over distinct TypeIds — a +/// smoke assertion; a collision would degrade to a decode error, never a +/// misroute, per RFC §3). +#[test] +fn distinct_types_distinct_hashes() { + maybe_child(ROLES); + run(|| { + assert_ne!(type_hash::(), type_hash::()); + assert_ne!(type_hash::(), type_hash::>()); + }); +} + +/// The decode-and-deliver contract: a registered hash decodes into the +/// target's typed channel; an unknown hash, corrupt bytes, and a missing +/// channel each fail without delivering — `WrongChannel`, never a misroute. +#[test] +fn decoder_registration_and_delivery() { + maybe_child(ROLES); + run(|| { + let h_string = expose_type::(); + let h_u64 = expose_type::(); + assert!(decoder_registered(h_string)); + assert!(!decoder_registered(h_string.wrapping_add(1))); + + // A live actor with a String channel (registered from its own body, + // announced via a ready signal — the tests/registry.rs idiom). + let (ready_tx, ready_rx) = channel::<()>(); + let (stop_tx, stop_rx) = channel::<()>(); + let (msg_tx, msg_rx) = channel::(); + let pid = spawn(move || { + register(Name::::new("expose-test.sink"), msg_tx).unwrap(); + ready_tx.send(()).unwrap(); + let _ = stop_rx.recv(); + }) + .pid(); + ready_rx.recv().unwrap(); + + // Happy path: decode + deliver through the published channel. + let bytes = encode_payload("hello across the seam").unwrap(); + decode_deliver(h_string, pid, &bytes).unwrap(); + assert_eq!(msg_rx.recv().unwrap(), "hello across the seam"); + + // Unknown hash: nothing was registered under it. + assert!(matches!( + decode_deliver(h_string.wrapping_add(1), pid, &bytes), + Err(DeliverError::UnknownType) + )); + + // Corrupt bytes: the decoder fails before any send. + assert!(matches!( + decode_deliver(h_string, pid, &[0xff; 3]), + Err(DeliverError::Decode(_)) + )); + + // Right decoder, wrong channel: the actor has no u64 channel, so the + // decoded value is refused — the NoChannel guarantee. + let u64_bytes = encode_payload(&7u64).unwrap(); + assert!(matches!( + decode_deliver(h_u64, pid, &u64_bytes), + Err(DeliverError::WrongChannel) + )); + + stop_tx.send(()).unwrap(); + }); +} + +/// `expose` and the bridge crossing agree on the resulting set: both funnel +/// the pid-boundary mark through the watchable machinery, so an exposed +/// name's holder dies with a terminal record — the exact observable +/// `mark_watchable` guarantees the membrane. (For named holders the mark is +/// already stamped by `register` itself; this pins the shared contract.) +#[test] +fn expose_and_bridge_crossing_agree_on_the_set() { + maybe_child(ROLES); + run(|| { + let (ready_tx, ready_rx) = channel::<()>(); + let (stop_tx, stop_rx) = channel::<()>(); + let (msg_tx, _msg_rx) = channel::(); + let pid = spawn(move || { + register(GREETER, msg_tx).unwrap(); + ready_tx.send(()).unwrap(); + let _ = stop_rx.recv(); + }) + .pid(); + ready_rx.recv().unwrap(); + + expose(GREETER); + let m = monitor(pid); + stop_tx.send(()).unwrap(); + assert_eq!(m.rx.recv().unwrap().reason, DownReason::Exit); + assert_eq!(terminal_reason(pid), Some(DownReason::Exit)); + }); +} + +/// Hash stability across runs in the same binary: a re-exec of this binary +/// computes the same hashes this process does. +#[test] +fn hash_stable_across_runs_in_same_binary() { + maybe_child(ROLES); + let (mine_string, mine_u64) = { + // Computing a TypeId hash needs no runtime, but keep the contract + // uniform with real call sites. + (type_hash::(), type_hash::()) + }; + let mut child = spawn_node("hasher", &[]); + let line = child.wait_line("HASH-STRING", |l| l.starts_with("HASH-STRING ")); + assert_eq!( + line["HASH-STRING ".len()..].parse::().unwrap(), + mine_string + ); + let line = child.wait_line("HASH-U64", |l| l.starts_with("HASH-U64 ")); + assert_eq!(line["HASH-U64 ".len()..].parse::().unwrap(), mine_u64); + child.wait_exit(); +}