feat(cluster): RFC 010 c6b — handshake on the accept/connect path

Drive the c5 machines as straight-line code on the path (D8): dial_handshake
and accept_handshake do the IO on a shared FramedConn, and a connection actor
is spawned only after a successful handshake. Rejects, tie-break losses (D7),
protocol faults and timeouts are all resolved on the path by closing, so no
actor ever exists for a connection that did not establish. The whole
FramedConn travels into spawn_established, carrying any read-ahead past the
handshake frames.

Handshake deadlines land here rather than in c6c: FramedConn::recv_deadline
enforces them between reads via the connection's fd arm, so a peer that
connects and goes silent cannot wedge the acceptor.

Connection lifetime moves to the manager (pulled forward from c7). The path
registers each established connection and hands over its ConnHandle; the
manager owns it, monitors the actor, and tears the connection down on
Disconnect, on peer close, or at manager shutdown. spawn_established returns
a Pid, so a connection neither outlives nor dies with whichever actor
established it — the ownership that made two-node teardown unorderable.

The manager also tracks in-flight dial intents, monitored so a panicking
dial cannot wedge the tie-break, and answers HelloCtx for the accept path.
This commit is contained in:
Claude
2026-08-14 21:09:08 +00:00
parent bbaaa062e3
commit c8ed858e4c
8 changed files with 1037 additions and 46 deletions
+4 -1
View File
@@ -4,10 +4,12 @@
//! transport trait (control connection), framed codec, and the TCP +
//! loopback impls. c5: the handshake state machine. c6: the connection
//! [`manager`] (registry) and per-peer connection actors ([`conn`]), started
//! as an explicit supervision subtree. Everything above them lands in later
//! as an explicit supervision subtree, plus the handshake on the
//! accept/connect path ([`connect`]). Everything above them lands in later
//! chunks.
pub mod conn;
pub mod connect;
pub mod envelope;
pub mod handshake;
pub mod manager;
@@ -21,6 +23,7 @@ use crate::scheduler::{sleep, spawn, JoinHandle};
use crate::supervisor::{ChildSpec, OneForOne, Restart};
pub use conn::{spawn_established, ConnHandle};
pub use connect::{dial, spawn_acceptor, AcceptorHandle};
pub use manager::{Manager, MANAGER};
/// A running cluster subtree: an explicitly-started supervisor over the
+48 -23
View File
@@ -10,60 +10,85 @@
//! so no read/write split is needed.
//!
//! The handshake completes *before* this actor exists (on the accept/connect
//! path — c6b) and produces the [`Peer`]; the actor registers that peer with
//! the [`manager`](crate::cluster::manager), which monitors it so any exit
//! deregisters the connection. Heartbeat send and fixed-timeout liveness join
//! the loop in c6c (the timeout arm of the same `select`).
//! path — c6b) and produces the [`Peer`]; the *path* then registers the
//! connection with the [`manager`](crate::cluster::manager), which takes
//! ownership of its [`ConnHandle`] and monitors the actor, so any exit
//! deregisters the connection. The actor itself holds no authority over its
//! own lifetime: it runs until the manager drops its handle (deregistration,
//! `Disconnect`, or manager shutdown) or the connection ends. Heartbeat send
//! and fixed-timeout liveness join the loop in c6c (the timeout arm of the
//! same `select`).
use crate::channel::{channel, Receiver, Selectable, Sender};
use crate::cluster::handshake::Peer;
use crate::cluster::manager::{Call, Registered, Reply, MANAGER};
use crate::cluster::transport::FramedConn;
use crate::gen_server;
use crate::scheduler::{self, spawn};
use crate::pid::Pid;
use crate::scheduler::spawn;
/// Commands to a running connection actor.
enum Cmd {
Shutdown,
}
/// A handle to a running connection actor.
/// The manager's authority over one connection actor: while this handle
/// lives the connection lives, and dropping it stops the actor and closes
/// the socket. Only the [`manager`](crate::cluster::manager) holds one —
/// callers of [`spawn_established`] get a [`Pid`] and no lifetime authority,
/// so a connection can never outlive, or die with, whichever actor happened
/// to establish it.
pub struct ConnHandle {
cmd_tx: Sender<Cmd>,
}
impl std::fmt::Debug for ConnHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("ConnHandle")
}
}
impl ConnHandle {
/// Ask the connection to close and exit. Idempotent, and a no-op if the
/// actor has already gone.
/// actor has already gone. Dropping the handle does the same thing; this
/// exists for the manager's explicit `Disconnect` path.
pub fn shutdown(&self) {
let _ = self.cmd_tx.send(Cmd::Shutdown);
}
}
/// Spawn a connection actor for an **already-established** connection: the
/// handshake has completed elsewhere and produced `peer`. Returns as soon as
/// the actor is spawned; the actor's first act is to register with the manager.
pub fn spawn_established(framed: FramedConn, peer: Peer) -> ConnHandle {
let (cmd_tx, cmd_rx) = channel();
spawn(move || run(framed, peer, cmd_rx));
ConnHandle { cmd_tx }
}
/// The name was already claimed by a live connection, so this one was
/// refused; its actor has been stopped and its socket closed.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RegisterRefused;
fn run(mut framed: FramedConn, peer: Peer, cmd_rx: Receiver<Cmd>) {
let me = scheduler::self_pid();
/// Spawn a connection actor for an **already-established** connection (the
/// handshake completed on the path and produced `peer`) and register it with
/// the manager, synchronously, before returning. The manager takes the
/// actor's [`ConnHandle`]; the caller gets only the [`Pid`], because
/// connection lifetime belongs to the table and not to the establishing
/// 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 name = peer.node_name.clone();
let pid = spawn(move || run(framed, peer, cmd_rx)).pid();
match gen_server::call(
MANAGER,
Call::Register {
name: peer.node_name.clone(),
pid: me,
name,
pid,
handle: ConnHandle { cmd_tx },
},
) {
Ok(Reply::Registered(Registered::Ok)) => {}
// Duplicate name, or the manager is unreachable: do not run. The
// connection drops (closing the socket) as `framed` falls out of scope.
_ => return,
Ok(Reply::Registered(Registered::Ok)) => Ok(pid),
// Duplicate name, or the manager is unreachable. Either way the
// handle went with the call and is dropped there (or never arrived
// and dropped with it), which stops the actor and closes the socket.
_ => Err(RegisterRefused),
}
}
fn run(mut framed: FramedConn, _peer: Peer, cmd_rx: Receiver<Cmd>) {
loop {
match framed.readable_arm() {
Some(arm) => {
+345
View File
@@ -0,0 +1,345 @@
//! RFC 010 c6b — the handshake on the accept/connect path.
//!
//! Per D8 (re-amended): the c5 machines are driven by **straight-line code
//! on the path**, not by an actor. The dial side runs [`Initiator`]; the
//! acceptor loop runs [`Responder`]. A connection actor is spawned only
//! *after* a successful handshake ([`spawn_established`]); every reject,
//! protocol failure, timeout, and tie-break loss is resolved right here,
//! on the path, by closing — no actor ever exists for a connection that
//! didn't establish.
//!
//! Buffer trap (binding): the path reader and the steady-state actor share
//! ONE [`FramedConn`]. Its decode buffer may hold read-ahead past the
//! handshake frames, so the *whole* `FramedConn` travels into
//! [`spawn_established`] — never a fresh codec over the same socket.
//!
//! Layering: [`dial_handshake`] and [`accept_handshake`] are the bare path
//! steps — IO on a `FramedConn`, no manager, no actors — testable over the
//! loopback transport on plain threads. [`dial`] and [`spawn_acceptor`] are
//! the manager-integrated layer (actor context required): they keep the
//! [`manager`](crate::cluster::manager)'s dial-intent set honest and spawn
//! the connection actor on success.
use std::io;
use std::time::{Duration, Instant};
use crate::channel::{channel, Receiver, Selectable, Sender};
use crate::cluster::conn::spawn_established;
use crate::cluster::envelope::{Frame, RejectReason};
use crate::cluster::handshake::{
HelloCtx, Initiator, InitiatorOutcome, Local, Peer, Responder, ResponderOutcome,
};
use crate::cluster::manager::{Call, Reply, MANAGER};
use crate::cluster::transport::{FramedConn, Listener, RecvError, SendError, Transport};
use crate::gen_server;
use crate::pid::Pid;
use crate::scheduler::{self, spawn};
/// How long either side waits for the peer's handshake frame before giving
/// up and closing. Enforced on the path via [`FramedConn::recv_deadline`],
/// so a peer that connects and goes silent cannot wedge the acceptor.
pub const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5);
/// Why a handshake did not establish. In every case the connection has
/// already been closed on the path by the time this is returned.
#[derive(Debug)]
pub enum HandshakeError {
/// A `HelloReject` travelled — sent by us (accept side) or received by
/// us (dial side).
Rejected(RejectReason),
/// Accept side only: the inbound dial lost the simultaneous-connect
/// tie-break (D7) and was closed silently, no frame sent.
TieBreakLoss,
/// The peer spoke a valid frame that is wrong here (non-`Hello` first
/// frame; non-response to our `Hello`), or an undecodable byte stream.
Protocol,
/// EOF before the handshake resolved. On the dial side this is also
/// what losing the tie-break looks like: the peer closes silently.
Closed,
/// [`HANDSHAKE_TIMEOUT`] (or the caller's deadline) passed first.
TimedOut,
/// The transport failed mid-handshake.
Transport(io::Error),
}
fn from_send(e: SendError) -> HandshakeError {
match e {
// Handshake frames are small and self-made; an encode failure is a
// protocol-level impossibility, not a transport fault.
SendError::Encode(_) => HandshakeError::Protocol,
SendError::Io(e) => HandshakeError::Transport(e),
}
}
fn from_recv(e: RecvError) -> HandshakeError {
match e {
RecvError::Corrupt(_) => HandshakeError::Protocol,
RecvError::TruncatedByPeer => HandshakeError::Closed,
RecvError::Io(e) => HandshakeError::Transport(e),
RecvError::TimedOut => HandshakeError::TimedOut,
}
}
/// Dial-side path step: send our `Hello`, interpret the one response. On
/// `Ok` the connection is established and `framed` is live (with any
/// read-ahead intact in its buffer); on `Err` the connection is closed.
pub fn dial_handshake(
framed: &mut FramedConn,
local: &Local,
deadline: Instant,
) -> Result<Peer, HandshakeError> {
let (initiator, hello) = Initiator::new(local);
if let Err(e) = framed.send(&hello) {
framed.close();
return Err(from_send(e));
}
let outcome = match framed.recv_deadline(deadline) {
Ok(Some(frame)) => initiator.on_frame(frame),
Ok(None) => {
framed.close();
return Err(HandshakeError::Closed);
}
Err(e) => {
framed.close();
return Err(from_recv(e));
}
};
match outcome {
InitiatorOutcome::Established(peer) => Ok(peer),
InitiatorOutcome::Rejected(reason) => {
framed.close();
Err(HandshakeError::Rejected(reason))
}
InitiatorOutcome::Failed(_) => {
framed.close();
Err(HandshakeError::Protocol)
}
}
}
/// Accept-side path step: read the first frame, judge it, answer or close.
///
/// `ctx_for` supplies the [`HelloCtx`] for the *offered* name — knowledge
/// only the frame reveals, which is why it is a callback and not a value
/// (the integrated acceptor asks the manager; loopback tests fabricate).
/// It is not called when the first frame is not a `Hello`.
///
/// On `Ok` the ack has been sent and `framed` is live (read-ahead intact);
/// on `Err` any owed reject has been sent and the connection is closed.
pub fn accept_handshake(
framed: &mut FramedConn,
local: Local,
ctx_for: impl FnOnce(&str) -> HelloCtx,
deadline: Instant,
) -> Result<Peer, HandshakeError> {
let frame = match framed.recv_deadline(deadline) {
Ok(Some(frame)) => frame,
Ok(None) => {
framed.close();
return Err(HandshakeError::Closed);
}
Err(e) => {
framed.close();
return Err(from_recv(e));
}
};
let ctx = match &frame {
Frame::Hello { node_name, .. } => ctx_for(node_name),
_ => HelloCtx::default(),
};
match Responder::new(local).on_frame(frame, ctx) {
ResponderOutcome::Accepted { reply, peer } => {
if let Err(e) = framed.send(&reply) {
framed.close();
return Err(from_send(e));
}
Ok(peer)
}
ResponderOutcome::Rejected { reply, reason } => {
// Best effort: the reject is the cross-version compatibility
// anchor, but if the write fails the peer sees a bare close,
// which it must survive anyway.
let _ = framed.send(&reply);
framed.close();
Err(HandshakeError::Rejected(reason))
}
ResponderOutcome::TieBreakLoss => {
// D7: close silently — the peer computes the same verdict.
framed.close();
Err(HandshakeError::TieBreakLoss)
}
ResponderOutcome::Failed(_) => {
framed.close();
Err(HandshakeError::Protocol)
}
}
}
// ---------------------------------------------------------------------------
// Manager-integrated layer
// ---------------------------------------------------------------------------
/// Why an integrated [`dial`] did not produce a connection.
#[derive(Debug)]
pub enum DialError {
/// Another dial to this peer name is already in flight.
AlreadyDialing,
/// The manager is not running (or answered nonsense).
ManagerUnavailable,
/// The transport could not connect.
Connect(io::Error),
/// Connected, but the handshake did not establish.
Handshake(HandshakeError),
/// The peer at `addr` established, but answered as a different name
/// than the one we dialed — the tie-break bookkeeping (keyed by the
/// dialed name) would be unsound, so the connection is closed.
PeerNameMismatch { expected: String, got: String },
/// The handshake established, but the manager refused the registration:
/// a connection to this peer already exists. The loser has been closed.
Duplicate,
}
/// Dial `peer_name` at `addr` and run the handshake, keeping the manager's
/// dial-intent set honest around it: the intent is registered *before*
/// connecting (so a crossing inbound `Hello` sees it) and cleared the
/// moment the handshake resolves, before the connection actor is spawned.
/// Must run inside an actor. Retrying is the caller's business (c7's dial
/// loop); a lost tie-break surfaces as `Handshake(Closed)` — the peer's
/// accepted connection is already on its way.
pub fn dial(
transport: &dyn Transport,
addr: &str,
peer_name: &str,
local: &Local,
) -> Result<Pid, DialError> {
let me = scheduler::self_pid();
match gen_server::call(
MANAGER,
Call::DialBegin {
name: peer_name.to_string(),
pid: me,
},
) {
Ok(Reply::DialBegan(true)) => {}
Ok(Reply::DialBegan(false)) => return Err(DialError::AlreadyDialing),
_ => return Err(DialError::ManagerUnavailable),
}
let result = connect_and_shake(transport, addr, local);
// Cleared immediately on outcome — a stale intent during the established
// window would corrupt later tie-breaks. Synchronous (a call): the
// intent is provably gone before anything else happens.
let _ = gen_server::call(
MANAGER,
Call::DialEnd {
name: peer_name.to_string(),
},
);
let (mut framed, peer) = result?;
if peer.node_name != peer_name {
framed.close();
return Err(DialError::PeerNameMismatch {
expected: peer_name.to_string(),
got: peer.node_name,
});
}
spawn_established(framed, peer).map_err(|_| DialError::Duplicate)
}
fn connect_and_shake(
transport: &dyn Transport,
addr: &str,
local: &Local,
) -> Result<(FramedConn, Peer), DialError> {
let conn = transport.dial(addr).map_err(DialError::Connect)?;
let mut framed = FramedConn::new(conn);
let deadline = Instant::now() + HANDSHAKE_TIMEOUT;
let peer = dial_handshake(&mut framed, local, deadline).map_err(DialError::Handshake)?;
Ok((framed, peer))
}
/// A running acceptor. [`shutdown`](AcceptorHandle::shutdown) (or dropping
/// the last handle) stops the accept loop only: connections it established
/// belong to the [`manager`](crate::cluster::manager) and keep running, to
/// be torn down through the table (`Disconnect`, a peer close, or manager
/// shutdown).
pub struct AcceptorHandle {
cmd_tx: Sender<()>,
addr: String,
}
impl AcceptorHandle {
/// Ask the acceptor to stop. Idempotent; a no-op if it already has.
pub fn shutdown(&self) {
let _ = self.cmd_tx.send(());
}
/// The concrete bound address, dialable as-is.
pub fn local_addr(&self) -> &str {
&self.addr
}
}
/// Spawn the acceptor actor over a bound listener. Each inbound connection
/// is handshaken **inline in the loop** (a deliberate serialization: the
/// per-frame deadline bounds how long any one peer can hold the line, and
/// nothing concurrent exists to be starved before c7). The listener must be
/// fd-backed ([`Listener::readable_arm`]); the loopback listener is not,
/// and its acceptor exits immediately — loopback handshakes are driven
/// synchronously through the path fns instead, per D8.
pub fn spawn_acceptor(listener: Box<dyn Listener>, local: Local) -> AcceptorHandle {
let addr = listener.local_addr();
let (cmd_tx, cmd_rx) = channel();
spawn(move || accept_loop(listener, local, cmd_rx));
AcceptorHandle { cmd_tx, addr }
}
fn accept_loop(mut listener: Box<dyn Listener>, local: Local, cmd_rx: Receiver<()>) {
loop {
let Some(arm) = listener.readable_arm() else {
return;
};
let arms: [&dyn Selectable; 2] = [&cmd_rx, &arm];
match crate::channel::try_select(&arms) {
Ok(0) => match cmd_rx.try_recv() {
Ok(Some(())) => return,
Ok(None) => continue, // spurious wake
Err(_) => return, // all handles dropped
},
Ok(_) => {
// The listener is readable: accept completes without parking.
let conn = match listener.accept() {
Ok(conn) => conn,
Err(_) => return, // listener itself is broken
};
handle_inbound(FramedConn::new(conn), &local);
}
Err(_) => return, // fd arm failed to register: listener is gone
}
}
}
/// Run the accept-side handshake for one inbound connection, asking the
/// manager for the [`HelloCtx`], and hand the established connection to the
/// manager. Every failure was already resolved on the path (reject sent /
/// closed, or the registration refused and the actor stopped), so there is
/// nothing for the acceptor to carry forward.
fn handle_inbound(mut framed: FramedConn, local: &Local) {
let deadline = Instant::now() + HANDSHAKE_TIMEOUT;
let ctx_for = |name: &str| match gen_server::call(
MANAGER,
Call::HelloCtx {
peer_name: name.to_string(),
},
) {
Ok(Reply::HelloCtx(ctx)) => ctx,
// Manager unreachable: nobody could register this connection anyway,
// so claim the name taken and reject rather than accept an orphan.
_ => HelloCtx {
name_claimed: true,
dialing_this_peer: false,
},
};
if let Ok(peer) = accept_handshake(&mut framed, local.clone(), ctx_for, deadline) {
let _ = spawn_established(framed, peer);
}
}
+87 -12
View File
@@ -1,11 +1,14 @@
//! RFC 010 c6 — the cluster connection manager.
//!
//! One manager per runtime: the single registry of live peer connections and
//! the source of truth for whether a peer name is already claimed. Every
//! connection actor registers here as its first act and is *monitored* by the
//! manager, so the table self-heals on any exit path — a connection that
//! panics, is cancelled, or closes cleanly is removed without cooperation from
//! the dying actor.
//! the source of truth for whether a peer name is already claimed. The
//! accept/connect path registers each established connection here, handing
//! over its [`ConnHandle`] — **the manager owns connection lifetime**. A
//! connection lives as long as its table entry, so it neither outlives nor
//! dies with whichever actor happened to establish it. Registered actors are
//! also *monitored*, so the table self-heals on any exit path — a connection
//! that panics, is cancelled, or closes cleanly is removed without
//! cooperation from the dying actor.
//!
//! Membership as consumers will see it (the `node_up`/`node_down` interface and
//! the view) and the connector dial loop are c7, built on top of this table.
@@ -14,6 +17,8 @@
use std::collections::HashMap;
use crate::cluster::conn::ConnHandle;
use crate::cluster::handshake::HelloCtx;
use crate::gen_server::{GenServer, GenServerCtx, GenServerName, Watcher};
use crate::monitor::{monitor, Down};
use crate::pid::Pid;
@@ -23,10 +28,23 @@ use crate::pid::Pid;
/// manager is always found at the same key.
pub const MANAGER: GenServerName<Manager> = GenServerName::new("smarm.cluster.manager");
/// One live connection's entry: the actor running it, and the handle whose
/// lifetime *is* the connection's (see the module docs).
struct ConnEntry {
pid: Pid,
_handle: ConnHandle,
}
/// The connection registry: peer name → the connection actor that owns that
/// peer's control connection.
pub struct Manager {
conns: HashMap<String, Pid>,
conns: HashMap<String, ConnEntry>,
/// In-flight dial intents: peer name -> the actor performing the dial.
/// Registered *before* connecting so a crossing inbound `Hello` sees it
/// (the `dialing_this_peer` half of [`HelloCtx`]); cleared the moment
/// the dial resolves, and — because the dialer is monitored — on the
/// dialer's death, so a panicking dial can never wedge the tie-break.
dials: HashMap<String, Pid>,
watcher: Option<Watcher<Manager>>,
}
@@ -34,6 +52,7 @@ impl Manager {
pub fn new() -> Self {
Manager {
conns: HashMap::new(),
dials: HashMap::new(),
watcher: None,
}
}
@@ -57,18 +76,43 @@ pub enum Registered {
/// Requests to the manager.
pub enum Call {
/// A freshly-established connection actor claims its peer's name. The pid
/// is the calling connection actor, which the manager then monitors.
Register { name: String, pid: Pid },
/// The path claims its peer's name for a freshly-established connection,
/// handing the manager the actor's [`ConnHandle`]. The manager monitors
/// `pid` and holds the handle for as long as the entry lives; a
/// [`Registered::Duplicate`] verdict drops the handle here, which stops
/// the refused actor.
Register {
name: String,
pid: Pid,
handle: ConnHandle,
},
/// Tear down the connection to `name`: the manager drops its handle, the
/// actor stops, and the monitor removes the entry. A no-op if no such
/// connection is live.
Disconnect { name: String },
/// The current peer names, sorted. For observation and tests.
Peers,
/// A dialer declares an in-flight dial to `name` before connecting. The
/// pid is the dialing actor, monitored so the intent dies with it.
DialBegin { name: String, pid: Pid },
/// The dial to `name` resolved (either way): drop the intent. A call,
/// not a cast, so the intent is provably gone before the dialer moves on.
DialEnd { name: String },
/// The [`HelloCtx`] for an inbound `Hello` offering `peer_name` — the
/// accept path asks this between reading the frame and judging it.
HelloCtx { peer_name: String },
}
/// Replies from the manager.
#[derive(Debug)]
pub enum Reply {
Registered(Registered),
Disconnected,
Peers(Vec<String>),
/// `false`: another dial to this name is already in flight — do not dial.
DialBegan(bool),
DialEnded,
HelloCtx(HelloCtx),
}
impl GenServer for Manager {
@@ -84,27 +128,58 @@ impl GenServer for Manager {
fn handle_call(&mut self, request: Call) -> Reply {
match request {
Call::Register { name, pid } => {
Call::Register { name, pid, handle } => {
if self.conns.contains_key(&name) {
// `handle` drops here: the refused actor stops itself.
return Reply::Registered(Registered::Duplicate);
}
if let Some(w) = &self.watcher {
w.watch(monitor(pid));
}
self.conns.insert(name, pid);
self.conns.insert(
name,
ConnEntry {
pid,
_handle: handle,
},
);
Reply::Registered(Registered::Ok)
}
Call::Disconnect { name } => {
// Dropping the entry drops the handle, which stops the actor.
self.conns.remove(&name);
Reply::Disconnected
}
Call::Peers => {
let mut names: Vec<String> = self.conns.keys().cloned().collect();
names.sort();
Reply::Peers(names)
}
Call::DialBegin { name, pid } => {
if self.dials.contains_key(&name) {
return Reply::DialBegan(false);
}
if let Some(w) = &self.watcher {
w.watch(monitor(pid));
}
self.dials.insert(name, pid);
Reply::DialBegan(true)
}
Call::DialEnd { name } => {
self.dials.remove(&name);
Reply::DialEnded
}
Call::HelloCtx { peer_name } => Reply::HelloCtx(HelloCtx {
name_claimed: self.conns.contains_key(&peer_name),
dialing_this_peer: self.dials.contains_key(&peer_name),
}),
}
}
fn handle_cast(&mut self, _request: ()) {}
fn handle_down(&mut self, down: Down) {
self.conns.retain(|_, pid| *pid != down.pid);
self.conns.retain(|_, entry| entry.pid != down.pid);
self.dials.retain(|_, pid| *pid != down.pid);
}
}
+57
View File
@@ -71,6 +71,15 @@ pub trait Listener: Send {
/// The concrete bound address, dialable as-is (e.g. the real port when
/// bound with port 0).
fn local_addr(&self) -> String;
/// Readiness as a [`select`](crate::select) arm, mirroring
/// [`Conn::readable_arm`]: `Some` lets an acceptor wait on "an inbound
/// connection is pending" alongside a command inbox in one `select`, so
/// it can be told to stop without a poll loop. Default `None` (the
/// loopback listener has no fd and must be driven synchronously).
fn readable_arm(&self) -> Option<crate::scheduler::FdArm> {
None
}
}
/// A way of establishing control connections. Object-safe on purpose: the
@@ -128,6 +137,10 @@ pub enum RecvError {
TruncatedByPeer,
/// The transport failed mid-read.
Io(io::Error),
/// The deadline passed before a full frame arrived
/// ([`FramedConn::recv_deadline`] only; plain [`recv`](FramedConn::recv)
/// never returns this).
TimedOut,
}
impl std::fmt::Display for RecvError {
@@ -136,6 +149,7 @@ impl std::fmt::Display for RecvError {
RecvError::Corrupt(e) => write!(f, "frame stream corrupt: {e:?}"),
RecvError::TruncatedByPeer => write!(f, "peer closed mid-frame"),
RecvError::Io(e) => write!(f, "transport read failed: {e}"),
RecvError::TimedOut => write!(f, "deadline passed mid-receive"),
}
}
}
@@ -193,6 +207,49 @@ impl FramedConn {
}
}
/// Like [`recv`](FramedConn::recv), but gives up with
/// [`RecvError::TimedOut`] once `deadline` passes without a full frame.
/// The deadline is enforced between reads via the connection's fd arm
/// (so the caller must be an actor); a transport with no fd (loopback)
/// cannot be timed out and this degrades to a plain blocking `recv` —
/// the same caveat as liveness.
pub fn recv_deadline(
&mut self,
deadline: std::time::Instant,
) -> Result<Option<Frame>, RecvError> {
loop {
match Frame::decode(&self.rbuf) {
Ok(Some((frame, consumed))) => {
self.rbuf.drain(..consumed);
return Ok(Some(frame));
}
Ok(None) => {}
Err(e) => return Err(RecvError::Corrupt(e)),
}
if let Some(arm) = self.conn.readable_arm() {
let left = deadline.saturating_duration_since(std::time::Instant::now());
if left.is_zero() {
return Err(RecvError::TimedOut);
}
match crate::channel::try_select_timeout(&[&arm], left) {
Ok(Some(_)) => {}
Ok(None) => return Err(RecvError::TimedOut),
Err(e) => return Err(RecvError::Io(e)),
}
}
let mut chunk = [0u8; READ_CHUNK];
let n = self.conn.read(&mut chunk).map_err(RecvError::Io)?;
if n == 0 {
return if self.rbuf.is_empty() {
Ok(None)
} else {
Err(RecvError::TruncatedByPeer)
};
}
self.rbuf.extend_from_slice(&chunk[..n]);
}
}
/// Close the underlying connection (idempotent, see [`Conn::close`]).
pub fn close(&mut self) {
self.conn.close();
+4
View File
@@ -197,6 +197,10 @@ pub struct TcpListener {
}
impl Listener for TcpListener {
fn readable_arm(&self) -> Option<crate::scheduler::FdArm> {
Some(crate::scheduler::FdArm::readable(self.inner.as_raw_fd()))
}
fn accept(&mut self) -> io::Result<Box<dyn Conn>> {
loop {
wait_readable(self.inner.as_raw_fd())?;