feat(cluster): RFC 010 c9 — remote Name sends: outbound table + the one inbound seam

Outbound (D13, ratified 2026-08-15): a module-private table node → dedicated
Sender<Frame>, 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<M> (node, Name<M>), 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.
This commit is contained in:
Claude
2026-08-15 21:22:56 +00:00
parent a7f98f8d48
commit 16ef583455
10 changed files with 688 additions and 16 deletions
+10
View File
@@ -241,6 +241,16 @@ impl<T> Sender<T> {
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<T>) -> 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,
+2
View File
@@ -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
+94 -12
View File
@@ -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<Frame>` 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<Cmd>,
/// 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<Sender<Frame>>,
}
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<Sender<Frame>> {
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<Pid, RegisterRefused> {
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<Cmd>) {
fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver<Cmd>, out_rx: Receiver<Frame>) {
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<Cmd>) {
/// `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<Cmd>,
out_rx: &Receiver<Frame>,
) {
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<Frame>) -> 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
}
+24 -4
View File
@@ -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);
+240
View File
@@ -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<Frame>` — 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<M>)`. Sendable via
/// [`send`]; typed, so the payload is `M` and the wire hash is
/// [`type_hash::<M>()`](type_hash).
pub struct RemoteName<M> {
node: String,
name: Name<M>,
_marker: PhantomData<fn() -> M>,
}
impl<M> RemoteName<M> {
pub fn new(node: impl Into<String>, name: Name<M>) -> Self {
RemoteName {
node: node.into(),
name,
_marker: PhantomData,
}
}
pub fn node(&self) -> &str {
&self.node
}
pub fn name(&self) -> Name<M> {
self.name
}
}
impl<M> Clone for RemoteName<M> {
fn clone(&self) -> Self {
RemoteName {
node: self.node.clone(),
name: self.name,
_marker: PhantomData,
}
}
}
impl<M> std::fmt::Debug for RemoteName<M> {
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<M> {
/// 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<M> RemoteSendError<M> {
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<String, Sender<Frame>>,
}
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<Frame>) {
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<M>(target: RemoteName<M>, msg: M) -> Result<(), RemoteSendError<M>>
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::<M>(),
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),
}
}
+37
View File
@@ -215,6 +215,7 @@ impl<M> std::error::Error for SendError<M> {}
trait ErasedSender: Send {
fn as_any(&self) -> &dyn Any;
fn queued_len(&self) -> usize;
fn receiver_alive(&self) -> bool;
}
impl<M: Send + 'static> ErasedSender for Sender<M> {
@@ -224,6 +225,9 @@ impl<M: Send + 'static> ErasedSender for Sender<M> {
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<M>` filed
@@ -442,6 +446,18 @@ pub(crate) fn register_with<M: Send + 'static>(
/// 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<String>`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<M: Send + 'static>(reg: &mut Registry, me: Pid, tx: Sender<M>) {
let mb = reg
.by_index
@@ -450,6 +466,16 @@ fn publish_channel<M: Send + 'static>(reg: &mut Registry, me: Pid, tx: Sender<M>
if mb.pid != me {
*mb = Mailbox::new(me);
}
if let Some(existing) = mb.channels.get(&TypeId::of::<M>()) {
assert!(
!existing.sender.receiver_alive() || same_channel::<M>(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::<M>()
);
}
mb.channels.insert(
TypeId::of::<M>(),
Channel {
@@ -459,6 +485,17 @@ fn publish_channel<M: Send + 'static>(reg: &mut Registry, me: Pid, tx: Sender<M>
);
}
/// 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<M: Send + 'static>(existing: &Channel, tx: &Sender<M>) -> bool {
existing
.sender
.as_any()
.downcast_ref::<Sender<M>>()
.is_some_and(|old| old.same_channel(tx))
}
/// Publish the current actor's `Sender<A::Msg>` into its mailbox **without**
/// binding a name, and hand back the typed [`Pid<A>`] that addresses this
/// actor directly.
+7
View File
@@ -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<crate::cluster::expose::ExposureState>,
/// RFC 010 c9: the outbound table (node name → the connection's
/// dedicated `Sender<Frame>`), manager-maintained. Leaf; the send happens
/// outside the lock. cfg-gated like `exposure`.
#[cfg(feature = "cluster")]
pub(crate) outbound: RawMutex<crate::cluster::remote::Outbound>,
/// Recycled stacks waiting to be reused by the next spawn.
pub(crate) stack_pool: RawMutex<Vec<crate::stack::Stack>>,
/// 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),
+220
View File
@@ -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 <addr>`,
//! `MEMBER-UP <name>`, `GOT <payload>`, `SEND-RESULT <case> <verdict>`.
//! 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<String> = Name::new("c9.inbox");
const HIDDEN: Name<String> = 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::<String>();
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::<String>();
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<String> = 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<String> = 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::<u64>(),
&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<String> = 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"
);
}
+11
View File
@@ -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,
+43
View File
@@ -258,3 +258,46 @@ fn send_dyn_to_dead_pid_is_dead() {
assert!(matches!(send_dyn::<u64>(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::<u64>();
let (tx2, _rx2) = channel::<u64>();
register(Name::<u64>::new("dup-a"), tx1).unwrap();
register(Name::<u64>::new("dup-b"), tx2).unwrap(); // panics
});
}
#[test]
fn two_names_on_one_cloned_sender_is_fine() {
run(|| {
let (tx, rx) = channel::<u64>();
register(Name::<u64>::new("twin-a"), tx.clone()).unwrap();
register(Name::<u64>::new("twin-b"), tx).unwrap();
send(Name::<u64>::new("twin-a"), 1).unwrap();
send(Name::<u64>::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::<u64>();
register(Name::<u64>::new("reborn"), tx1).unwrap();
drop(rx1); // old inbox gone: replacement is the honest thing to do
let (tx2, rx2) = channel::<u64>();
register(Name::<u64>::new("reborn-2"), tx2).unwrap();
send(Name::<u64>::new("reborn-2"), 9).unwrap();
assert_eq!(rx2.recv().unwrap(), 9);
});
}