From 16ef583455d2d12998037fba1d6e30a33272477b Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 15 Aug 2026 21:22:56 +0000 Subject: [PATCH] =?UTF-8?q?feat(cluster):=20RFC=20010=20c9=20=E2=80=94=20r?= =?UTF-8?q?emote=20Name=20sends:=20outbound=20table=20+=20the=20one=20inbo?= =?UTF-8?q?und=20seam?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Outbound (D13, ratified 2026-08-15): a module-private table node → dedicated Sender, populated/torn down by the manager inside the same serialized handlers that own connection lifetime (Register / Disconnect / reap / terminate), on RuntimeInner beside the exposure state. remote::send is one leaf-lock lookup + one channel send: no gen_server on the data plane (rejected send-through-manager — worse than the c8 argument, it is the data plane), no published channel a leaked conn pid could inject raw frames into (rejected send_dyn-at-conn-pid — module privacy cannot fence a published channel). Ok(()) = handed to the connection's inbox, local knowledge only (RFC §3); missing entry / closed channel = NotConnected. The outbound sender is a SEPARATE channel from cmd_tx on purpose: a clone would let the table hold lifetime authority (D9 violation); its closure is not a stop signal. Buffering toward a slow peer is unbounded (BEAM busy_dist_port shape); backpressure out of scope, documented not silently absent. Inbound: remote::deliver_named is THE resolution seam (RFC v2) — exposed-set check (unexposed = unreachable, the safety) → hash check against what the name was exposed with → registry whereis → c8's decode_deliver. Verdicts are local-only (InboundVerdict, discarded by the conn actor for now; a trace hook is the place). The conn actor gains a third select arm (outbound inbox → wire) and interprets SendNamed; Send/Monitor/Down are consumed for liveness and ignored until c10/c11. Public surface: RemoteName (node, Name), RemoteSendError, NotConnected, send, send_remote_raw (untyped escape hatch so tests can put deliberately wrong frames on the wire — the typed API cannot express a hash mismatch). CORE (found the hard way this chunk): publishing a second channel of the same message type on one live actor silently replaced and CLOSED the first, so a select/recv on it returned 'closed' immediately forever — a hot loop starving the single-threaded scheduler. publish_channel now asserts when an existing same-type channel's receiver is still alive AND the new sender is not a clone of it (Sender::same_channel via Arc::ptr_eq); replacing a dead-receiver channel stays silent (an actor re-registering after dropping its inbox is legit). Message names the sanctioned shapes. Three registry tests pin panic / cloned-sender-ok / dead-receiver-ok. Full default suite clean. Harness: wait_line drains stderr before dumping on timeout/EOF (eprintln! diagnostics no longer vanish); Node::transcript() accessor for ordering-proof assertions. tests/cluster_remote_send.rs 1/0, 10/10 flake runs, subprocess harness: cross-node name-send delivers; unexposed name unreachable (registered locally, never delivered); wrong hash never misroutes (unknown-hash and known-hash-wrong-channel flavours); send to an unconnected node = NotConnected locally; every frame-bearing send is Ok — 'handed to transport', asserted and documented at the test. Negatives are proven by STREAM ORDERING (they precede the positive on one in-order connection), not by sleeping. All cluster suites regression-clean (envelope 15, handshake 11, transport 11, lifecycle 1, liveness 3, connect 9, two_node 3, membership 4, mesh 2, expose 5); clippy --lib green both configs; fmt clean; default build compiles. --- src/channel.rs | 10 ++ src/cluster.rs | 2 + src/cluster/conn.rs | 106 ++++++++++++++-- src/cluster/manager.rs | 28 +++- src/cluster/remote.rs | 240 +++++++++++++++++++++++++++++++++++ src/registry.rs | 37 ++++++ src/runtime.rs | 7 + tests/cluster_remote_send.rs | 220 ++++++++++++++++++++++++++++++++ tests/common/mod.rs | 11 ++ tests/registry.rs | 43 +++++++ 10 files changed, 688 insertions(+), 16 deletions(-) create mode 100644 src/cluster/remote.rs create mode 100644 tests/cluster_remote_send.rs diff --git a/src/channel.rs b/src/channel.rs index a2af72a..9c50081 100644 --- a/src/channel.rs +++ b/src/channel.rs @@ -241,6 +241,16 @@ impl Sender { self.inner.lock().queue.len() } + /// Whether the [`Receiver`] is still alive (a send would be accepted). + pub(crate) fn receiver_alive(&self) -> bool { + self.inner.lock().receiver_alive + } + + /// Whether `other` is a sender of this very channel (a clone). + pub(crate) fn same_channel(&self, other: &Sender) -> bool { + Arc::ptr_eq(&self.inner, &other.inner) + } + /// Push `value` onto the channel. Succeeds unconditionally as long as /// the [`Receiver`] is still alive: the queue has no capacity limit, so /// this never blocks and never fails except when the channel is closed, diff --git a/src/cluster.rs b/src/cluster.rs index e2b2225..a56cc6c 100644 --- a/src/cluster.rs +++ b/src/cluster.rs @@ -17,6 +17,7 @@ pub mod expose; pub mod handshake; pub mod manager; pub mod membership; +pub mod remote; pub mod transport; use std::io; @@ -40,6 +41,7 @@ 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}; +pub use remote::{NotConnected, RemoteName, RemoteSendError}; /// 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 diff --git a/src/cluster/conn.rs b/src/cluster/conn.rs index b2953b5..a9abad4 100644 --- a/src/cluster/conn.rs +++ b/src/cluster/conn.rs @@ -21,6 +21,16 @@ //! [`Frame::Heartbeat`](crate::cluster::envelope::Frame::Heartbeat)s, and a //! [`LIVENESS_TIMEOUT`] window — reset by any inbound frame — tears the //! connection down when it empties. +//! +//! c9 adds the third arm — the connection's dedicated **outbound inbox** +//! (`Sender` bound in the manager-maintained outbound table, D13), +//! drained onto the wire in the same loop — and inbound *interpretation*: +//! `SendNamed` goes to the one resolution seam, +//! [`remote::deliver_named`](crate::cluster::remote::deliver_named). +//! Frames the connection actor has no business with yet (`Send`, `Monitor`, +//! …: c10/c11) are consumed for liveness and otherwise ignored. The outbound +//! sender is a separate channel from `cmd_tx` on purpose: closing it is not +//! a stop signal — lifetime authority stays with the [`ConnHandle`] (D9). use std::time::{Duration, Instant}; @@ -28,6 +38,7 @@ use crate::channel::{channel, try_select_timeout, Receiver, Selectable, Sender}; use crate::cluster::envelope::Frame; use crate::cluster::handshake::Peer; use crate::cluster::manager::{Call, Registered, Reply, MANAGER}; +use crate::cluster::remote::deliver_named; use crate::cluster::transport::FramedConn; use crate::gen_server; use crate::pid::Pid; @@ -46,6 +57,11 @@ enum Cmd { /// to establish it. pub struct ConnHandle { cmd_tx: Sender, + /// The connection's dedicated outbound inbox. The manager moves this + /// into the outbound table on `Register` (see + /// [`take_outbound`](ConnHandle::take_outbound)); a `Duplicate` verdict + /// drops it with the handle. + out_tx: Option>, } impl std::fmt::Debug for ConnHandle { @@ -61,6 +77,12 @@ impl ConnHandle { pub fn shutdown(&self) { let _ = self.cmd_tx.send(Cmd::Shutdown); } + + /// Manager-only: take the outbound sender to bind into the outbound + /// table. Once, at registration. + pub(crate) fn take_outbound(&mut self) -> Option> { + self.out_tx.take() + } } /// The name was already claimed by a live connection, so this one was @@ -76,14 +98,18 @@ pub struct RegisterRefused; /// actor. A refusal has already stopped the actor and closed the socket. pub fn spawn_established(framed: FramedConn, peer: Peer) -> Result { let (cmd_tx, cmd_rx) = channel(); + let (out_tx, out_rx) = channel(); let reg_peer = peer.clone(); - let pid = spawn(move || run(framed, peer, cmd_rx)).pid(); + let pid = spawn(move || run(framed, peer, cmd_rx, out_rx)).pid(); match gen_server::call( MANAGER, Call::Register { peer: reg_peer, pid, - handle: ConnHandle { cmd_tx }, + handle: ConnHandle { + cmd_tx, + out_tx: Some(out_tx), + }, }, ) { Ok(Reply::Registered(Registered::Ok)) => Ok(pid), @@ -107,20 +133,30 @@ pub const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(1); /// honest detector. pub const LIVENESS_TIMEOUT: Duration = Duration::from_secs(4); -fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver) { +fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver, out_rx: Receiver) { match framed.readable_arm() { - Some(arm) => run_live(&mut framed, arm, &cmd_rx), + Some(arm) => run_live(&mut framed, arm, &cmd_rx, &out_rx), None => run_inert(&cmd_rx), } framed.close(); } /// The steady-state loop over an fd-backed connection: one -/// `select_timeout` folds the command inbox, socket readability, and the -/// nearer of the two deadlines (`hb_send`, `liveness`) into a single wait. -fn run_live(framed: &mut FramedConn, arm: crate::scheduler::FdArm, cmd_rx: &Receiver) { +/// `select_timeout` folds the command inbox, the outbound inbox, socket +/// readability, and the nearer of the two deadlines (`hb_send`, +/// `liveness`) into a single wait. +fn run_live( + framed: &mut FramedConn, + arm: crate::scheduler::FdArm, + cmd_rx: &Receiver, + out_rx: &Receiver, +) { let mut next_hb = Instant::now(); let mut live_until = Instant::now() + LIVENESS_TIMEOUT; + // The outbound sender lives in the manager's table and is dropped on + // unbind; after that this arm would wake forever, so it drops out of + // the select (not a stop signal — see the module docs). + let mut out_open = true; loop { let now = Instant::now(); if now >= live_until { @@ -133,14 +169,18 @@ fn run_live(framed: &mut FramedConn, arm: crate::scheduler::FdArm, cmd_rx: &Rece next_hb = now + HEARTBEAT_INTERVAL; } let wait = next_hb.min(live_until).saturating_duration_since(now); - let arms: [&dyn Selectable; 2] = [cmd_rx, &arm]; + // Arm indices: 0 cmd, 1 fd, 2 outbound (when open). + let mut arms: Vec<&dyn Selectable> = vec![cmd_rx, &arm]; + if out_open { + arms.push(out_rx); + } match try_select_timeout(&arms, wait) { Ok(Some(0)) => { if should_stop(cmd_rx) { break; } } - Ok(Some(_)) => match pump_readable(framed) { + Ok(Some(1)) => match pump_readable(framed) { Pump::Ended => break, Pump::Frames(n) => { if n > 0 { @@ -148,6 +188,11 @@ fn run_live(framed: &mut FramedConn, arm: crate::scheduler::FdArm, cmd_rx: &Rece } } }, + Ok(Some(_)) => match pump_outbound(framed, out_rx) { + Outbound::Sent => {} + Outbound::Closed => out_open = false, + Outbound::WireFailed => break, + }, // A deadline passed; the top of the loop acts on whichever. Ok(None) => {} // The fd arm failed to register — the connection is gone. @@ -156,6 +201,30 @@ fn run_live(framed: &mut FramedConn, arm: crate::scheduler::FdArm, cmd_rx: &Rece } } +/// What one outbound wake yielded. +enum Outbound { + Sent, + /// The manager unbound this connection's sender; nothing more will come. + Closed, + /// The socket refused a write: the connection is gone. + WireFailed, +} + +/// Drain every queued outbound frame onto the wire. +fn pump_outbound(framed: &mut FramedConn, out_rx: &Receiver) -> Outbound { + loop { + match out_rx.try_recv() { + Ok(Some(frame)) => { + if framed.send(&frame).is_err() { + return Outbound::WireFailed; + } + } + Ok(None) => return Outbound::Sent, + Err(_) => return Outbound::Closed, + } + } +} + /// No fd to select on (loopback): only a command can end the wait, and /// neither heartbeats nor liveness run — a transport that can't report /// readiness can't be timed either (same caveat as @@ -196,8 +265,11 @@ enum Pump { /// after a level-triggered readable indication), then drain every complete /// frame the buffer now holds. A blocking `recv` here would park the actor /// past its heartbeat and liveness deadlines whenever a frame arrives split. -/// Frames are not interpreted yet — a heartbeat's entire job is the liveness -/// reset, and everything else waits for c8. +/// Every consumed frame counts for liveness; `SendNamed` additionally goes +/// to the one inbound resolution seam. Its verdict is local knowledge only +/// — nothing goes back on the wire (RFC §3) — and is currently discarded +/// (a future trace hook is the place to surface it). Frames for later +/// chunks (`Send` c10, `Monitor`/`Down` c11) are consumed and ignored. fn pump_readable(framed: &mut FramedConn) -> Pump { let eof = match framed.read_once() { Ok(n) => n == 0, @@ -206,7 +278,17 @@ fn pump_readable(framed: &mut FramedConn) -> Pump { let mut got = 0; loop { match framed.next_buffered() { - Ok(Some(_frame)) => got += 1, + Ok(Some(frame)) => { + got += 1; + if let Frame::SendNamed { + name, + type_hash, + payload, + } = frame + { + let _verdict = deliver_named(&name, type_hash, &payload); + } + } Ok(None) => break, Err(_) => return Pump::Ended, // corrupt stream } diff --git a/src/cluster/manager.rs b/src/cluster/manager.rs index 4a40152..2588023 100644 --- a/src/cluster/manager.rs +++ b/src/cluster/manager.rs @@ -26,6 +26,7 @@ use crate::channel::Sender; use crate::cluster::conn::ConnHandle; use crate::cluster::handshake::{HelloCtx, Peer}; use crate::cluster::membership::{NodeEvent, NodeInfo}; +use crate::cluster::remote::{bind_outbound, unbind_outbound}; use crate::gen_server::{GenServer, GenServerCtx, GenServerName, Watcher}; use crate::monitor::{monitor, Down}; use crate::pg::NodeId; @@ -181,9 +182,21 @@ impl GenServer for Manager { self.watcher = Some(ctx.watcher()); } + /// Manager shutdown drops every entry (and with it every ConnHandle); + /// the outbound table must not outlive the connections it names. + fn terminate(&mut self) { + for name in self.conns.keys() { + unbind_outbound(name); + } + } + fn handle_call(&mut self, request: Call) -> Reply { match request { - Call::Register { peer, pid, handle } => { + Call::Register { + peer, + pid, + mut handle, + } => { if self.conns.contains_key(&peer.node_name) { // `handle` drops here: the refused actor stops itself. return Reply::Registered(Registered::Duplicate); @@ -191,6 +204,11 @@ impl GenServer for Manager { if let Some(w) = &self.watcher { w.watch(monitor(pid)); } + // The outbound table (c9) is maintained here, inside the same + // serialized handlers that own the connection's lifetime. + if let Some(out) = handle.take_outbound() { + bind_outbound(&peer.node_name, out); + } let info = NodeInfo { node: self.node_id(&peer.node_name, peer.incarnation.get()), name: peer.node_name.clone(), @@ -211,6 +229,7 @@ impl GenServer for Manager { Call::Disconnect { name } => { // Dropping the entry drops the handle, which stops the actor. if let Some(entry) = self.conns.remove(&name) { + unbind_outbound(&name); self.emit(&NodeEvent::NodeDown { node: entry.info.node, }); @@ -257,14 +276,15 @@ impl GenServer for Manager { fn handle_down(&mut self, down: Down) { let mut downs = Vec::new(); - self.conns.retain(|_, entry| { + self.conns.retain(|name, entry| { let dead = entry.pid == down.pid; if dead { - downs.push(entry.info.node); + downs.push((name.clone(), entry.info.node)); } !dead }); - for node in downs { + for (name, node) in downs { + unbind_outbound(&name); self.emit(&NodeEvent::NodeDown { node }); } self.dials.retain(|_, pid| *pid != down.pid); diff --git a/src/cluster/remote.rs b/src/cluster/remote.rs new file mode 100644 index 0000000..33e4191 --- /dev/null +++ b/src/cluster/remote.rs @@ -0,0 +1,240 @@ +//! RFC 010 c9 — remote `Name` sends: the outbound path and the single +//! inbound name-resolution seam. +//! +//! ## Outbound (D13, ratified 2026-08-15) +//! +//! A module-private table `node name → Sender` — one dedicated +//! outbound channel per live connection, populated and torn down by the +//! manager inside the same serialized handlers that own the connection's +//! lifetime (Register / Disconnect / reap), living on `RuntimeInner` beside +//! the exposure state. [`send`] is one leaf-lock lookup + one channel send: +//! no gen_server on the data plane, no published channel anyone holding a +//! pid could inject raw frames into. `Ok(())` means **handed to the +//! connection's inbox** — local knowledge only, exactly the BEAM contract +//! (RFC §3): a missing entry or a closed channel is +//! [`RemoteSendError::NotConnected`]; delivery confirmation is the monitor's +//! job (c11). The entry-present/actor-dying-mid-send window is *honest* +//! under that contract, not a bug. +//! +//! The outbound sender is deliberately **separate from the conn actor's +//! command channel**: if it were a clone of `cmd_tx`, the manager dropping +//! its `ConnHandle` would no longer close that channel and connection +//! lifetime would leak to whoever holds a sender — a D9 violation. +//! +//! Buffering is unbounded toward a slow peer (the BEAM `busy_dist_port` +//! shape); backpressure is out of c9's scope and noted here rather than +//! silently absent. +//! +//! ## Inbound — the ONE resolution seam (RFC v2) +//! +//! Every wire-name → local-pid resolution goes through [`deliver_named`], +//! and nothing else: the conn actor hands it the three fields of a +//! `SendNamed` and gets back a verdict. It checks the exposed set first (an +//! unexposed name is unreachable — the gun's safety), then the type hash +//! against what the name was exposed with, then resolves the name through +//! the registry and delivers via c8's [`decode_deliver`]. When an owned-name +//! table lands beside the `&'static str` registry, it slots in here without +//! touching call sites. Module privacy enforces the funnel: the exposed and +//! outbound tables are `pub(crate)`, and no other module resolves names for +//! the wire. +//! +//! Refusals are silent to the sender by design (§3: send failure reflects +//! local knowledge only); they are observable locally as the returned +//! [`InboundVerdict`], which the conn actor may log or count. + +use std::collections::HashMap; +use std::marker::PhantomData; + +use crate::channel::Sender; +use crate::cluster::envelope::{encode_payload, Frame, PayloadError}; +use crate::cluster::expose::{decode_deliver, exposed_hash, type_hash, DeliverError}; +use crate::pid::Name; +use crate::registry::whereis; +use crate::scheduler::with_runtime; + +/// A name on a specific remote node: `(node_name, Name)`. Sendable via +/// [`send`]; typed, so the payload is `M` and the wire hash is +/// [`type_hash::()`](type_hash). +pub struct RemoteName { + node: String, + name: Name, + _marker: PhantomData M>, +} + +impl RemoteName { + pub fn new(node: impl Into, name: Name) -> Self { + RemoteName { + node: node.into(), + name, + _marker: PhantomData, + } + } + pub fn node(&self) -> &str { + &self.node + } + pub fn name(&self) -> Name { + self.name + } +} + +impl Clone for RemoteName { + fn clone(&self) -> Self { + RemoteName { + node: self.node.clone(), + name: self.name, + _marker: PhantomData, + } + } +} + +impl std::fmt::Debug for RemoteName { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}@{}", self.name.as_str(), self.node) + } +} + +/// Why a remote send did not leave this node. Local knowledge only. +#[derive(Debug)] +pub enum RemoteSendError { + /// No live connection to that node right now (never connected, or gone + /// and not yet re-dialed). The message is handed back. + NotConnected(M), + /// The payload did not serialize. + Encode(M, PayloadError), +} + +impl RemoteSendError { + pub fn into_inner(self) -> M { + match self { + RemoteSendError::NotConnected(m) | RemoteSendError::Encode(m, _) => m, + } + } +} + +/// The outbound table, one per runtime (a `RuntimeInner` field). +pub(crate) struct Outbound { + by_node: HashMap>, +} + +impl Outbound { + pub(crate) fn new() -> Self { + Outbound { + by_node: HashMap::new(), + } + } +} + +/// Manager-only: bind `node`'s outbound channel. Called inside `Register`. +pub(crate) fn bind_outbound(node: &str, tx: Sender) { + with_runtime(|inner| { + inner.outbound.lock().by_node.insert(node.to_string(), tx); + }); +} + +/// Manager-only: unbind `node`'s outbound channel. Called on `Disconnect`, +/// reap, and manager shutdown. Dropping the sender is what closes the conn +/// actor's outbound arm — but that arm's closure is NOT a stop signal (the +/// cmd channel is, per D9); the actor simply stops selecting on it. +pub(crate) fn unbind_outbound(node: &str) { + with_runtime(|inner| { + inner.outbound.lock().by_node.remove(node); + }); +} + +/// Send `msg` to `target`. `Ok(())` = handed to the connection's inbox, and +/// nothing more — see the module docs. Must run inside +/// [`run`](crate::run). +pub fn send(target: RemoteName, msg: M) -> Result<(), RemoteSendError> +where + M: serde::Serialize + Send + 'static, +{ + let payload = match encode_payload(&msg) { + Ok(p) => p, + Err(e) => return Err(RemoteSendError::Encode(msg, e)), + }; + let frame = Frame::SendNamed { + name: target.name.as_str().to_string(), + type_hash: type_hash::(), + payload, + }; + match hand_to_connection(&target.node, frame) { + Ok(()) => Ok(()), + Err(NotConnected) => Err(RemoteSendError::NotConnected(msg)), + } +} + +/// No live connection to the named node — the payload-free form of +/// [`RemoteSendError::NotConnected`], for the raw path. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct NotConnected; + +/// The untyped escape hatch: send pre-encoded `payload` under an explicit +/// `type_hash`. Exists so tests (and future codecs) can put deliberately +/// wrong frames on the wire; the typed [`send`] cannot express a hash/type +/// mismatch, by design. Same `Ok` semantics as [`send`]. +pub fn send_remote_raw( + node: &str, + name: &str, + type_hash: u64, + payload: &[u8], +) -> Result<(), NotConnected> { + hand_to_connection( + node, + Frame::SendNamed { + name: name.to_string(), + type_hash, + payload: payload.to_vec(), + }, + ) +} + +/// One lookup, one send. Clone the sender out under the lock and send +/// outside it (a channel send can unpark the conn actor). +fn hand_to_connection(node: &str, frame: Frame) -> Result<(), NotConnected> { + let tx = with_runtime(|inner| inner.outbound.lock().by_node.get(node).cloned()); + match tx { + Some(tx) => tx.send(frame).map_err(|_| NotConnected), + None => Err(NotConnected), + } +} + +/// What the inbound seam did with a `SendNamed`. Local observability only; +/// nothing goes back on the wire (RFC §3). +#[derive(Debug)] +pub enum InboundVerdict { + /// Decoded and handed to the name's holder. + Delivered, + /// The name is not in this node's exposed set. + NotExposed, + /// The frame's hash is not the hash the name was exposed with. + HashMismatch { expected: u64, got: u64 }, + /// Exposed, but no live holder right now (unbound, or its holder died + /// and the binding is being pruned). + Unresolved, + /// Resolved, but the delivery half refused it (decode failure, or the + /// holder's channel does not accept the exposed type — a local + /// re-registration under a different type; never a misroute). + Refused(DeliverError), +} + +/// THE inbound resolution seam: exposed-set check → hash check → registry +/// resolution → c8 delivery. See the module docs. Must run inside +/// [`run`](crate::run) — the conn actor's context. +pub fn deliver_named(name: &str, type_hash: u64, payload: &[u8]) -> InboundVerdict { + let Some(expected) = exposed_hash(name) else { + return InboundVerdict::NotExposed; + }; + if expected != type_hash { + return InboundVerdict::HashMismatch { + expected, + got: type_hash, + }; + } + let Some(pid) = whereis(name) else { + return InboundVerdict::Unresolved; + }; + match decode_deliver(type_hash, pid, payload) { + Ok(()) => InboundVerdict::Delivered, + Err(e) => InboundVerdict::Refused(e), + } +} diff --git a/src/registry.rs b/src/registry.rs index 7c20b1d..76be1db 100644 --- a/src/registry.rs +++ b/src/registry.rs @@ -215,6 +215,7 @@ impl std::error::Error for SendError {} trait ErasedSender: Send { fn as_any(&self) -> &dyn Any; fn queued_len(&self) -> usize; + fn receiver_alive(&self) -> bool; } impl ErasedSender for Sender { @@ -224,6 +225,9 @@ impl ErasedSender for Sender { fn queued_len(&self) -> usize { Sender::queued_len(self) } + fn receiver_alive(&self) -> bool { + Sender::receiver_alive(self) + } } /// One typed channel of an actor, type-erased. Concretely a `Sender` filed @@ -442,6 +446,18 @@ pub(crate) fn register_with( /// name) and [`install`] (which does not). A leftover mailbox at this slot /// index from a dead prior incarnation (pid mismatch) is replaced wholesale. /// Caller holds the registry lock and has established that `me` is live. +/// +/// **One channel per message type per actor.** Publishing a second `M` +/// channel on the same live actor replaces the first — and if the first's +/// receiver is still alive, that replacement drops its last sender, closing +/// it, and any `recv`/`select` on it then returns "closed" immediately and +/// forever: a silent hot loop that starves the scheduler. That is never +/// intended, so it panics here (found the hard way in RFC 010 c9, where two +/// `Name`s registered on one actor did exactly this). Replacing a +/// channel whose receiver is already gone is fine (an actor re-registering +/// after dropping its old inbox) and stays silent. To hold two names of the +/// same type, register them from two actors, or bind both names to one +/// cloned sender. fn publish_channel(reg: &mut Registry, me: Pid, tx: Sender) { let mb = reg .by_index @@ -450,6 +466,16 @@ fn publish_channel(reg: &mut Registry, me: Pid, tx: Sender if mb.pid != me { *mb = Mailbox::new(me); } + if let Some(existing) = mb.channels.get(&TypeId::of::()) { + assert!( + !existing.sender.receiver_alive() || same_channel::(existing, &tx), + "smarm: actor {me:?} already publishes a live channel for message type `{}`; \ + a second one would replace and CLOSE the first (its receiver would then \ + read as closed forever). Register the second name from another actor, or \ + bind both names to a clone of the same sender.", + type_name::() + ); + } mb.channels.insert( TypeId::of::(), Channel { @@ -459,6 +485,17 @@ fn publish_channel(reg: &mut Registry, me: Pid, tx: Sender ); } +/// True if `existing` and `tx` are senders of the very same channel (a +/// cloned sender bound under a second name is the sanctioned way to hold two +/// names of one type on one actor). +fn same_channel(existing: &Channel, tx: &Sender) -> bool { + existing + .sender + .as_any() + .downcast_ref::>() + .is_some_and(|old| old.same_channel(tx)) +} + /// Publish the current actor's `Sender` into its mailbox **without** /// binding a name, and hand back the typed [`Pid`] that addresses this /// actor directly. diff --git a/src/runtime.rs b/src/runtime.rs index d25e167..911114c 100644 --- a/src/runtime.rs +++ b/src/runtime.rs @@ -968,6 +968,11 @@ pub(crate) struct RuntimeInner { /// takes `registry`, never this). cfg-gated: zero-cost-when-off (c1). #[cfg(feature = "cluster")] pub(crate) exposure: RawMutex, + /// RFC 010 c9: the outbound table (node name → the connection's + /// dedicated `Sender`), manager-maintained. Leaf; the send happens + /// outside the lock. cfg-gated like `exposure`. + #[cfg(feature = "cluster")] + pub(crate) outbound: 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. @@ -1030,6 +1035,8 @@ impl RuntimeInner { process_groups: RawMutex::new(crate::pg::ProcessGroups::new()), #[cfg(feature = "cluster")] exposure: RawMutex::new(crate::cluster::expose::ExposureState::new()), + #[cfg(feature = "cluster")] + outbound: RawMutex::new(crate::cluster::remote::Outbound::new()), stack_pool: RawMutex::new(Vec::new()), stack_pool_cap, stack_reserve: crate::stack::round_to_pages(stack_reserve), diff --git a/tests/cluster_remote_send.rs b/tests/cluster_remote_send.rs new file mode 100644 index 0000000..2f6e1e1 --- /dev/null +++ b/tests/cluster_remote_send.rs @@ -0,0 +1,220 @@ +//! RFC 010 c9 — remote `Name` sends: the outbound seam and the single +//! inbound name-resolution seam, cross-process. +//! +//! Two node processes each run the integrated `cluster::start`. The +//! *receiver* registers a `String` inbox under a name and exposes it (and +//! registers a second name it does NOT expose); the *sender* waits for +//! `node_up`, then sends. Facts cross as stdout lines: `LISTENING `, +//! `MEMBER-UP `, `GOT `, `SEND-RESULT `. +//! Roles park forever afterwards (retractable-state trap); the parent +//! SIGKILLs via `Drop`. +//! +//! What is asserted at each end (roadmap-binding): +//! - cross-node name-send delivers the payload; +//! - an unexposed name is unreachable — the receiver's inbox stays empty +//! even though the name IS registered locally; +//! - a wrong type hash is a decode failure at the receiver, never a +//! misroute — the `String` inbox does not see a `u64` delivered under a +//! made-up hash, nor a `u64` under `u64`'s hash; +//! - a send to a disconnected (never-connected) node fails locally with +//! `NotConnected`, and `Ok(())` means only "handed to the transport". +//! +//! Timing note for the "stays empty" assertions: they are proven by +//! ORDERING, not by waiting — the sender emits the negative-case frames +//! BEFORE the positive one on the same connection (in-order stream), so when +//! the receiver has seen the positive payload, the negatives have already +//! been processed and refused. No sleep-and-hope. +#![cfg(feature = "cluster")] + +mod common; + +use common::{maybe_child, spawn_node, Node}; +use smarm::cluster::envelope::NodeMeta; +use smarm::cluster::expose::expose; +use smarm::cluster::membership::{subscribe, NodeEvent}; +use smarm::cluster::remote::{send_remote_raw, RemoteName, RemoteSendError}; +use smarm::cluster::{start, Config, StaticSeeds}; +use smarm::{channel, register, Name}; +use std::time::Duration; + +const ROLES: &[(&str, fn())] = &[("receiver", role_receiver), ("sender", role_sender)]; + +const INBOX: Name = Name::new("c9.inbox"); +const HIDDEN: Name = Name::new("c9.hidden"); + +fn base_config(name: &str, seeds: Vec<(String, String)>) -> Config { + Config { + node_name: name.to_string(), + meta: NodeMeta { + role: "c9".to_string(), + region: "local".to_string(), + }, + listen_addr: "127.0.0.1:0".to_string(), + strategy: Box::new(StaticSeeds::new(seeds)), + } +} + +/// Receiver: register + expose INBOX; register HIDDEN unexposed **in a +/// separate actor** (one actor holds one channel per message type — a +/// second `register` of the same `M` on one actor silently replaces the +/// first, closing it); print every payload that lands in either. +fn role_receiver() { + smarm::run(|| { + let cluster = start(base_config("recv", vec![])).expect("listener binds"); + println!("LISTENING {}", cluster.local_addr()); + + // HIDDEN's holder: its own actor, so its String channel does not + // displace INBOX's on the root actor. + let (hidden_ready_tx, hidden_ready_rx) = channel::<()>(); + smarm::spawn(move || { + let (hid_tx, hid_rx) = channel::(); + register(HIDDEN, hid_tx).unwrap(); + hidden_ready_tx.send(()).unwrap(); + loop { + match hid_rx.recv() { + Ok(s) => println!("GOT-HIDDEN {s}"), + Err(_) => break, + } + } + }); + hidden_ready_rx.recv().unwrap(); + + let (in_tx, in_rx) = channel::(); + register(INBOX, in_tx).unwrap(); + expose(INBOX); + println!("READY"); + loop { + match in_rx.recv() { + Ok(s) => println!("GOT {s}"), + Err(_) => break, + } + } + loop { + smarm::sleep(Duration::from_secs(3600)); + } + }); +} + +/// Sender: connect to recv, wait for node_up, then in this ORDER on the one +/// connection: hidden-name send, wrong-hash sends (two flavours), then the +/// positive send. Plus a send to a node that is not connected at all. +fn role_sender() { + let recv_addr = std::env::var("SMARM_RECV_ADDR").expect("SMARM_RECV_ADDR"); + smarm::run(move || { + let _cluster = start(base_config("send", vec![("recv".to_string(), recv_addr)])) + .expect("listener binds"); + let events = subscribe().expect("manager is up"); + loop { + match events.rx.recv() { + Ok(NodeEvent::NodeUp(info)) if info.name == "recv" => break, + Ok(_) => continue, + Err(_) => panic!("manager gone"), + } + } + println!("MEMBER-UP recv"); + + // Not connected: purely local knowledge, no frame leaves. + let ghost: RemoteName = RemoteName::new("nowhere", INBOX); + let r = smarm::cluster::remote::send(ghost, "lost".to_string()); + println!( + "SEND-RESULT not-connected {}", + match r { + Err(RemoteSendError::NotConnected(_)) => "NotConnected", + Ok(()) => "Ok", + Err(_) => "OtherErr", + } + ); + + // Unexposed name at the peer: the frame goes (local knowledge can't + // know the peer's exposed set) and the peer refuses it. + let hidden: RemoteName = RemoteName::new("recv", HIDDEN); + let r = smarm::cluster::remote::send(hidden, "should not land".to_string()); + println!( + "SEND-RESULT hidden {}", + if r.is_ok() { "Ok" } else { "Err" } + ); + + // Wrong hash, two flavours: (a) a u64 payload under a made-up hash + // (unknown type at the peer); (b) a u64 payload under u64's real + // hash against a String-typed name (decoder known, wrong channel). + // Both are raw sends — the typed API cannot express them, by design. + let bogus = 0xdead_beef_u64; + let r = send_remote_raw( + "recv", + "c9.inbox", + bogus, + &smarm::cluster::envelope::encode_payload(&7u64).unwrap(), + ); + println!( + "SEND-RESULT wrong-hash-unknown {}", + if r.is_ok() { "Ok" } else { "Err" } + ); + let r = send_remote_raw( + "recv", + "c9.inbox", + smarm::cluster::expose::type_hash::(), + &smarm::cluster::envelope::encode_payload(&7u64).unwrap(), + ); + println!( + "SEND-RESULT wrong-hash-known {}", + if r.is_ok() { "Ok" } else { "Err" } + ); + + // Positive: last on the stream, so its arrival proves the negatives + // were already processed. + let inbox: RemoteName = RemoteName::new("recv", INBOX); + let r = smarm::cluster::remote::send(inbox, "hello from send".to_string()); + println!( + "SEND-RESULT positive {}", + if r.is_ok() { "Ok" } else { "Err" } + ); + + loop { + smarm::sleep(Duration::from_secs(3600)); + } + }); +} + +fn wait_send_result(node: &mut Node, case: &str) -> String { + let prefix = format!("SEND-RESULT {case} "); + let line = node.wait_line(&prefix, |l| l.starts_with(&prefix)); + line[prefix.len()..].to_string() +} + +#[test] +fn remote_name_send_delivers_and_refusals_never_misroute() { + maybe_child(ROLES); + + let mut recv = spawn_node("receiver", &[]); + let addr = recv.wait_listening(); + recv.wait_line("READY", |l| l == "READY"); + + let mut send = spawn_node("sender", &[("SMARM_RECV_ADDR", &addr)]); + send.wait_line("MEMBER-UP recv", |l| l == "MEMBER-UP recv"); + + // Local-knowledge-only failure for an unknown node. + assert_eq!(wait_send_result(&mut send, "not-connected"), "NotConnected"); + // Every frame-bearing send is Ok — Ok means "handed to the transport", + // nothing about what the peer does with it (RFC §3, documented here). + assert_eq!(wait_send_result(&mut send, "hidden"), "Ok"); + assert_eq!(wait_send_result(&mut send, "wrong-hash-unknown"), "Ok"); + assert_eq!(wait_send_result(&mut send, "wrong-hash-known"), "Ok"); + assert_eq!(wait_send_result(&mut send, "positive"), "Ok"); + + // The positive payload lands... + recv.wait_line("GOT hello from send", |l| l == "GOT hello from send"); + // ...and, by stream ordering, every negative before it was refused: no + // GOT for the wrong-hash frames, no GOT-HIDDEN at all. The transcript + // up to this point is the proof. + let transcript = recv.transcript(); + let gots: Vec<&str> = transcript + .iter() + .map(|s| s.as_str()) + .filter(|l| l.starts_with("GOT")) + .collect(); + assert_eq!( + gots, + ["GOT hello from send"], + "exactly one delivery, the exposed one" + ); +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index ba522e7..36ae39f 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -133,6 +133,13 @@ impl Node { } } + /// Every stdout/stderr line seen so far, in arrival order. For + /// ordering-proof assertions ("by the time X arrived, Y had not"). + #[allow(dead_code)] + pub fn transcript(&self) -> &[String] { + &self.transcript + } + /// Wait until a stdout line satisfies `pred`; return it. Panics with the /// full transcript after [`WAIT`]. `what` names the expectation in the /// panic message. @@ -149,6 +156,9 @@ impl Node { } } Err(RecvTimeoutError::Timeout) => { + // Pull in whatever stderr arrived since the last drain, + // so a role's eprintln! diagnostics survive into the dump. + self.drain_stderr(); panic!( "node {:?}: timed out waiting for {what} after {WAIT:?}; transcript:\n{}", self.role, @@ -156,6 +166,7 @@ impl Node { ); } Err(RecvTimeoutError::Disconnected) => { + self.drain_stderr(); panic!( "node {:?}: output closed while waiting for {what}; transcript:\n{}", self.role, diff --git a/tests/registry.rs b/tests/registry.rs index d1e23f3..2ca251b 100644 --- a/tests/registry.rs +++ b/tests/registry.rs @@ -258,3 +258,46 @@ fn send_dyn_to_dead_pid_is_dead() { assert!(matches!(send_dyn::(p, 1u64), Err(SendError::Dead(_)))); }); } + +// --- one channel per message type per actor ----------------------------------- + +/// Registering a second name of the same message type on one actor, with a +/// *fresh* channel, would silently replace and close the first — so it +/// panics (found in RFC 010 c9). The sanctioned shapes stay quiet: bind both +/// names to a clone of one sender, or use two actors. +#[test] +#[should_panic(expected = "already publishes a live channel")] +fn second_live_channel_of_same_type_on_one_actor_panics() { + run(|| { + let (tx1, _rx1) = channel::(); + let (tx2, _rx2) = channel::(); + register(Name::::new("dup-a"), tx1).unwrap(); + register(Name::::new("dup-b"), tx2).unwrap(); // panics + }); +} + +#[test] +fn two_names_on_one_cloned_sender_is_fine() { + run(|| { + let (tx, rx) = channel::(); + register(Name::::new("twin-a"), tx.clone()).unwrap(); + register(Name::::new("twin-b"), tx).unwrap(); + send(Name::::new("twin-a"), 1).unwrap(); + send(Name::::new("twin-b"), 2).unwrap(); + assert_eq!(rx.recv().unwrap(), 1); + assert_eq!(rx.recv().unwrap(), 2); + }); +} + +#[test] +fn replacing_a_channel_whose_receiver_is_gone_is_fine() { + run(|| { + let (tx1, rx1) = channel::(); + register(Name::::new("reborn"), tx1).unwrap(); + drop(rx1); // old inbox gone: replacement is the honest thing to do + let (tx2, rx2) = channel::(); + register(Name::::new("reborn-2"), tx2).unwrap(); + send(Name::::new("reborn-2"), 9).unwrap(); + assert_eq!(rx2.recv().unwrap(), 9); + }); +}