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:
@@ -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,
|
||||
|
||||
@@ -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
@@ -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
@@ -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);
|
||||
|
||||
@@ -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),
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user