diff --git a/CHANGELOG.md b/CHANGELOG.md index cb148ae3..521f13a4 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -245,8 +245,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 counterpart in Berkeley sockets and a client author must know it: the v1 wire carries no half-close, so nothing peer-driven ever closes a flow, and a server written to read until the flow ends waits for a signal that cannot - arrive. The listener is built on `SOCK_SEQPACKET`, so it is available on - Linux and FreeBSD. + arrive. The listener uses `SOCK_SEQPACKET` on Linux and FreeBSD and + `SOCK_DGRAM` on macOS, which does not implement `SOCK_SEQPACKET` for + `AF_UNIX`; both keep the message boundaries the API's contract with its + clients rests on. The two kernels signal a closed peer differently and were + measured rather than reasoned about, so the receive path treats a Darwin + `ECONNRESET` as end of file alongside the `POLLHUP` and zero-byte read that + Linux gives. `EAGAIN` is deliberately not in that company: it means the + socket is empty and the peer alive, so it stays an error and the caller + waits again. ### Changed diff --git a/docs/how-to/use-the-native-datagram-api.md b/docs/how-to/use-the-native-datagram-api.md index cd65e605..f572671f 100644 --- a/docs/how-to/use-the-native-datagram-api.md +++ b/docs/how-to/use-the-native-datagram-api.md @@ -9,7 +9,23 @@ descriptor it reads and writes datagrams on. It is not a stable API surface, not a reliability layer, and not the v2 external process API. No compatibility promise is made: the socket protocol, the Rust client, and the configuration keys may change or be -withdrawn in any release. It is built on Linux and FreeBSD only. +withdrawn in any release. It is built on Linux, FreeBSD and macOS only. + +**The macOS build is exercised less than the others, and you should know +by how much.** The socket type differs there: macOS has no +`SOCK_SEQPACKET` for `AF_UNIX`, so a flow descriptor is a `SOCK_DGRAM` +socket, which keeps message boundaries just the same but reports a +closed peer differently. The unit tests run on macOS in the project's +own automation and pass in full. The end-to-end suite does not: it +drives a client container against a node container over a shared +volume, and that arrangement is Linux only. So on macOS the socket +lifecycle, the descriptor hand-off across a process boundary, and the +reclaiming of a port when a client exits are covered by unit tests +rather than by anything that runs the daemon and a client as two real +processes. + +That is a gap in testing, not a known defect. It is written here +because you are the one who would meet it first. Before enabling it, read the security posture below. It is short and it is the whole of the access control. diff --git a/docs/reference/configuration.md b/docs/reference/configuration.md index a31b61e1..73eec368 100644 --- a/docs/reference/configuration.md +++ b/docs/reference/configuration.md @@ -404,7 +404,7 @@ tuning under high load or on memory-constrained devices. ### Native Datagram API (`node.native_api.*`) -**Experimental, off by default, and built on Linux and FreeBSD only.** A +**Experimental, off by default, and built on Linux, FreeBSD and macOS only.** A client process connects to a Unix socket and asks either to open a flow to a remote pubkey or to hold a local port. Both answers carry a file descriptor: a flow's, which the client sends and receives datagrams on, or a listener's, @@ -1072,7 +1072,7 @@ node: enabled: true socket_path: null # null = auto (platform runtime dir → XDG → /tmp) # native_api: # uncomment to enable the experimental native datagram API - # enabled: true # opt-in, default false; Linux and FreeBSD only + # enabled: true # opt-in, default false; not on Windows # socket_path: /run/fips/api.sock # omit the key for the resolution above # pending_per_flow: 16 # datagrams held for one flow; 1..=64 # backlog: 16 # flows announced on one listener, awaiting its task; at least 1 diff --git a/docs/reference/native-api.md b/docs/reference/native-api.md index 7f0614c4..fcad6899 100644 --- a/docs/reference/native-api.md +++ b/docs/reference/native-api.md @@ -41,10 +41,12 @@ fixed string. The resolver takes the first of these that applies and appends `api.sock`: 1. `/run/fips/api.sock`, when `/run/fips` exists as a directory. This is the - packaged Linux convention and the constant the shipped Rust client compiles - in. -2. `/var/run/fips/api.sock` on FreeBSD, whose service scripts create that - directory. Linux skips this step. + packaged Linux convention, and the constant the shipped Rust client compiles + in on Linux. +2. `/var/run/fips/api.sock` on FreeBSD and macOS, whose service scripts and + LaunchDaemons create that directory. Linux skips this step, and macOS has no + `/run` for step 1 to find, so this is where a packaged macOS daemon lands and + is the constant the client compiles in there. 3. `$XDG_RUNTIME_DIR/fips/api.sock`, when that variable names an existing directory. A development run usually gets this one. 4. `/tmp/fips-api.sock`, the last resort. @@ -374,7 +376,7 @@ oversize datagram from anything else. Compare against the `libc` constant, never against a literal. The client maps each name onto `libc::` for the platform it was built for, so the number -differs between Linux and FreeBSD. +differs between Linux, FreeBSD and macOS. | errno | `kind()` | What causes it | | ----- | -------- | -------------- | diff --git a/docs/reference/security.md b/docs/reference/security.md index 8f391cf3..7043a258 100644 --- a/docs/reference/security.md +++ b/docs/reference/security.md @@ -200,7 +200,7 @@ that is enabled. ## Native Datagram API **Experimental. Disabled by default** (`node.native_api.enabled`, default -`false`), and built on Linux and FreeBSD only. It is not a stable API +`false`), and built on Linux, FreeBSD and macOS only. It is not a stable API surface, not a reliability layer, and not the v2 external process API. No compatibility promise is made about it. diff --git a/docs/tutorials/native-api-walkthrough.md b/docs/tutorials/native-api-walkthrough.md index 8baa5364..132cd2bf 100644 --- a/docs/tutorials/native-api-walkthrough.md +++ b/docs/tutorials/native-api-walkthrough.md @@ -14,7 +14,7 @@ you can build the daemon from source and can read Rust; it does not assume you have worked through the tutorials. **The API is experimental.** Names, fields and the command set may change -without a deprecation cycle. It is Linux and FreeBSD only. +without a deprecation cycle. It is Linux, FreeBSD and macOS only. ## What you will end up with diff --git a/src/config/node.rs b/src/config/node.rs index fbc6b67a..578491b9 100644 --- a/src/config/node.rs +++ b/src/config/node.rs @@ -936,7 +936,7 @@ impl ControlConfig { /// Native datagram API socket (`node.native_api.*`). /// -/// **Experimental, and built on Linux and FreeBSD only.** The API hands a +/// **Experimental, and built on Linux, FreeBSD and macOS only.** The API hands a /// client a file descriptor over `SCM_RIGHTS`, which Windows has no equivalent /// of, and does it over an `AF_UNIX` `SOCK_SEQPACKET` socket, which macOS does /// not implement. No listener is built on either, and this section is ignored @@ -1294,7 +1294,7 @@ pub struct NodeConfig { pub control: ControlConfig, /// Native datagram API (`node.native_api.*`). Experimental; the listener - /// is built on Linux and FreeBSD only. + /// is built on Linux, FreeBSD and macOS only. #[serde(default, skip_serializing_if = "NativeApiConfig::is_default")] pub native_api: NativeApiConfig, diff --git a/src/native/client/mod.rs b/src/native/client/mod.rs index d46946a3..da2ea09a 100644 --- a/src/native/client/mod.rs +++ b/src/native/client/mod.rs @@ -49,7 +49,23 @@ pub use secp256k1::XOnlyPublicKey; /// The path-taking constructors exist for a program that is told where its /// daemon is; everything else uses this, the way a Berkeley call needs no /// argument to find the kernel. +/// +/// **Platform-conditional, because the daemon's own default is.** The daemon +/// resolves its path at startup by looking for a directory rather than by +/// compiling a string, and macOS has no `/run` at all, so a client that +/// compiled the Linux path there would look somewhere that cannot exist. The +/// value here is the first branch of that resolver which applies to the +/// platform: `/run/fips` on Linux, `/var/run/fips` on macOS and FreeBSD, whose +/// packaged services create it. +/// +/// A daemon that fell through to `$XDG_RUNTIME_DIR` or `/tmp`, which a +/// development run usually does, is not at this path on any platform. That is +/// what `connect_at` and `bind_at` are for, and it is why the resolver is +/// documented rather than hidden. +#[cfg(not(any(target_os = "macos", target_os = "freebsd")))] pub const SOCKET: &str = "/run/fips/api.sock"; +#[cfg(any(target_os = "macos", target_os = "freebsd"))] +pub const SOCKET: &str = "/var/run/fips/api.sock"; /// Bytes taken from the RPC connection per `recvmsg`. /// @@ -442,7 +458,7 @@ impl FipsStream { ) }; if sent < 0 { - let error = io::Error::last_os_error(); + let error = seqpacket::peer_gone_as_epipe(io::Error::last_os_error()); if error.kind() == io::ErrorKind::Interrupted { continue; } @@ -458,8 +474,11 @@ impl FipsStream { /// is not end of file, and this is the one place where mimicking Berkeley /// exactly would be wrong: reading a zero-byte datagram as a close would let /// a peer tear down a live flow by sending nothing. A closed daemon half is - /// `EPIPE`, discriminated by `POLLHUP`, measured for this socket pair in - /// [`seqpacket`](super::seqpacket). + /// `EPIPE` on every platform, but the platforms disagree on how the kernel + /// says so: Linux discriminates a zero-byte read with `POLLHUP`, and Darwin + /// returns `ECONNRESET` outright. Both are measured for this socket pair in + /// [`seqpacket`](super::seqpacket), and both are translated to `EPIPE` here + /// so a caller never sees the difference. /// /// A datagram longer than `buf` is truncated and the remainder discarded, /// which is `SOCK_SEQPACKET` behaviour. Size `buf` at @@ -471,7 +490,7 @@ impl FipsStream { let received = unsafe { libc::recv(self.fd.as_raw_fd(), buf.as_mut_ptr().cast(), buf.len(), 0) }; if received < 0 { - let error = io::Error::last_os_error(); + let error = seqpacket::peer_gone_as_epipe(io::Error::last_os_error()); if error.kind() == io::ErrorKind::Interrupted { continue; } @@ -985,22 +1004,36 @@ mod tests { let (passed, held) = seqpacket::pair().unwrap(); let held = UnixStream::from(held); - // Both written before the client reads, so one recvmsg carries the - // plain line, the descriptor-bearing line and the descriptor. This is - // the measured rule: the ancillary data belongs to the sendmsg the read - // ended on, not to the first line the reader completes. + // Both are written before the client reads, so everything below is + // already queued and no read here waits on anything. (&daemon).write_all(REFUSAL).unwrap(); fdpass::send_once(daemon.as_raw_fd(), CONNECT_REPLY, Some(passed.as_fd())).unwrap(); drop(passed); + // **How many reads this takes is the platform's business, and asserting + // it was wrong.** Linux coalesces the plain write with the sendmsg that + // follows, so one recvmsg returns both lines and the descriptor. Darwin + // stops a stream read at the ancillary boundary, so the plain line + // arrives by itself and the descriptor-bearing line comes on the next + // read. Measured on macos-latest, where the earlier form of this test + // failed on exactly that difference. + // + // The rule under test is the same on both and is what the assertions + // below check: a descriptor belongs to the last complete line of the + // read that carried it. Darwin satisfies it more easily than Linux + // does, since the read it arrives on holds nothing later. let mut wire = Wire::new(client); - wire.fill().unwrap(); + for _ in 0..4 { + if wire.lines.len() >= 2 { + break; + } + wire.fill().unwrap(); + } assert_eq!( wire.lines.len(), 2, - "the plain line and the descriptor-bearing one should arrive as one \ - recvmsg; if they did not, the measured association rule this \ - module rests on does not hold on this platform" + "both lines should have arrived within four reads of a socket that \ + already held them" ); let (first, first_fd) = wire.line().unwrap(); diff --git a/src/native/dgram_probe.rs b/src/native/dgram_probe.rs new file mode 100644 index 00000000..1ee7b522 --- /dev/null +++ b/src/native/dgram_probe.rs @@ -0,0 +1,321 @@ +//! What a connected `AF_UNIX` `SOCK_DGRAM` pair does at end of file. +//! +//! This module is tests only. It exists to answer, by measurement on each +//! kernel rather than from the manual pages, the one question a macOS port of +//! this API turns on. +//! +//! [`super::seqpacket`] uses `SOCK_SEQPACKET`, which macOS does not implement +//! for `AF_UNIX`. The candidate replacement there is `SOCK_DGRAM`, which macOS +//! does implement and which also keeps message boundaries. What is not +//! transferable is the rule that tells a close apart from an empty datagram: +//! both produce a zero-byte read, and `seqpacket` resolves them with a latched +//! `POLLHUP` measured on Linux 6.8. Datagram poll semantics differ between +//! kernels, so that rule has to be re-established on Darwin before anything is +//! built on it. +//! +//! **The tests below assert the properties an implementation would need.** A +//! failure here is the measurement coming back negative, not a regression: it +//! says this kernel cannot support the `seqpacket` close rule on `SOCK_DGRAM` +//! and that a macOS port needs a different close signal. The same code runs on +//! every unix so the platforms can be compared without the test itself being a +//! variable. +//! +//! `SOCK_CLOEXEC` is deliberately not passed in the type argument, though +//! `super::seqpacket::pair` does pass it. Linux and FreeBSD accept it there and +//! macOS does not, and that difference belongs to the port rather than to this +//! measurement. + +use std::io; +use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd}; + +/// Create a connected `AF_UNIX` `SOCK_DGRAM` pair. +fn dgram_pair() -> io::Result<(OwnedFd, OwnedFd)> { + let mut fds = [0 as libc::c_int; 2]; + // SAFETY: `fds` is a two-element array of the type socketpair writes, and + // the call either fills both entries or reports failure. + let rc = unsafe { libc::socketpair(libc::AF_UNIX, libc::SOCK_DGRAM, 0, fds.as_mut_ptr()) }; + if rc != 0 { + return Err(io::Error::last_os_error()); + } + // SAFETY: socketpair reported success, so both entries are open descriptors + // this frame now owns. + Ok(unsafe { (OwnedFd::from_raw_fd(fds[0]), OwnedFd::from_raw_fd(fds[1])) }) +} + +/// One non-blocking receive, returning the byte count or the errno. +fn recv(fd: RawFd, buf: &mut [u8]) -> io::Result { + // SAFETY: the descriptor is open and the pointer and length describe `buf`. + let n = unsafe { libc::recv(fd, buf.as_mut_ptr().cast(), buf.len(), libc::MSG_DONTWAIT) }; + if n < 0 { + return Err(io::Error::last_os_error()); + } + Ok(n as usize) +} + +/// Send one datagram, returning the byte count or the errno. +fn send(fd: RawFd, buf: &[u8]) -> io::Result { + // SAFETY: the descriptor is open and the pointer and length describe `buf`. + let n = unsafe { libc::send(fd, buf.as_ptr().cast(), buf.len(), 0) }; + if n < 0 { + return Err(io::Error::last_os_error()); + } + Ok(n as usize) +} + +/// Whether `POLLHUP` is set, by the same rule `seqpacket::peer_hung_up` uses. +/// +/// `POLLIN` is requested rather than nothing. An empty `events` registers no +/// filter on Darwin, so a poll asking for nothing reports nothing there and +/// every assertion built on this helper would pass whatever the kernel did. +/// That is the shape of a guard that executes and cannot fail, so it is worth +/// more than a comment: with `POLLIN` requested, a Darwin that began reporting +/// `POLLHUP` would red the tests below instead of slipping past them. +fn hung_up(fd: RawFd) -> bool { + let mut poll = libc::pollfd { + fd, + events: libc::POLLIN, + revents: 0, + }; + // SAFETY: `poll` points at one live pollfd and the call cannot block. + let rc = unsafe { libc::poll(&mut poll, 1, 0) }; + rc > 0 && (poll.revents & libc::POLLHUP) != 0 +} + +#[test] +fn a_connected_dgram_pair_keeps_message_boundaries() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + assert_eq!(send(a.as_raw_fd(), &[1, 2, 3]).unwrap(), 3); + assert_eq!(send(a.as_raw_fd(), &[4, 5]).unwrap(), 2); + + // Two sends must read back as two messages of their own lengths. A stream + // socket would hand back all five bytes in one read, which is the failure + // this discriminates. + let mut buf = [0u8; 64]; + assert_eq!(recv(b.as_raw_fd(), &mut buf).unwrap(), 3); + assert_eq!(&buf[..3], &[1, 2, 3]); + assert_eq!(recv(b.as_raw_fd(), &mut buf).unwrap(), 2); + assert_eq!(&buf[..2], &[4, 5]); +} + +#[test] +fn an_empty_datagram_reads_as_zero_bytes_and_is_not_a_hangup() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + assert_eq!(send(a.as_raw_fd(), &[]).unwrap(), 0); + + let mut buf = [0u8; 64]; + assert_eq!( + recv(b.as_raw_fd(), &mut buf).unwrap(), + 0, + "an empty datagram must be delivered as a zero-byte message" + ); + assert!( + !hung_up(b.as_raw_fd()), + "an empty datagram must not look like a closed peer: if this fails, a \ + client could tear down its own flow by sending nothing" + ); +} + +/// Linux 6.8, measured 2026-08-20 with this test and cross-checked with a C +/// probe over both socket types: a connected `AF_UNIX` `SOCK_DGRAM` pair gives +/// **no close signal at all**. The peer closing leaves `revents` empty and +/// leaves `recv` returning `EAGAIN`, which is what an idle socket with a live +/// peer also does. `SOCK_SEQPACKET` on the same kernel sets `POLLHUP` and +/// returns a zero-byte read. +/// +/// This is asserted rather than merely written down so that a kernel which +/// starts reporting the close reds this test and reopens the question. +#[cfg(target_os = "linux")] +#[test] +fn linux_gives_a_dgram_pair_no_close_signal_at_all() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + drop(a); + + let mut buf = [0u8; 64]; + let read = recv(b.as_raw_fd(), &mut buf); + let err = read.as_ref().err().map(|e| e.raw_os_error()); + assert_eq!( + err, + Some(Some(libc::EAGAIN)), + "expected the closed peer to be indistinguishable from an idle socket, got {read:?}" + ); + assert!( + !hung_up(b.as_raw_fd()), + "POLLHUP is now set on a closed SOCK_DGRAM peer: this kernel has gained \ + the close signal Linux 6.8 did not have, and the macOS port's design \ + question should be reopened" + ); +} + +/// The open question, and the only thing a Mac is needed for. +/// +/// BSD kernels differ from Linux on datagram close reporting, so Darwin may +/// return `ECONNRESET`, or set `POLLHUP`, where Linux reports nothing. Either +/// would give the receive path something to key on. +/// +/// **A failure here is the measurement coming back negative, not a +/// regression.** It says Darwin behaves as Linux does, that a `SOCK_DGRAM` +/// descriptor carries no close signal, and that a macOS port must take the +/// close from the client's line-protocol connection instead of from the flow +/// descriptor. +#[cfg(target_os = "macos")] +#[test] +fn darwin_reports_a_closed_dgram_peer_somehow() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + drop(a); + + let mut buf = [0u8; 64]; + let read = recv(b.as_raw_fd(), &mut buf); + let hup = hung_up(b.as_raw_fd()); + let signalled = hup + || matches!(&read, Ok(0)) + || read.as_ref().err().is_some_and(|e| { + e.raw_os_error() != Some(libc::EAGAIN) && e.raw_os_error() != Some(libc::EWOULDBLOCK) + }); + assert!( + signalled, + "Darwin reports nothing when a connected SOCK_DGRAM peer closes: \ + POLLHUP unset and recv gave {read:?}, which is what an idle socket \ + gives. The flow descriptor cannot carry the close on this platform." + ); +} + +/// Which signal Darwin gives: `ECONNRESET`, and not `POLLHUP`. +/// +/// Measured 2026-08-20 on `macos-latest`, run 32353220389, and identical across +/// all three of nextest's attempts, so it is the kernel's behaviour and not a +/// race. The exact reading was `poll` returning 0 with an empty `revents`, and +/// `recv` returning errno 54, `ECONNRESET`. +/// +/// This is the opposite of `SOCK_SEQPACKET` on Linux, which sets `POLLHUP` and +/// returns a zero-byte read, and it is why +/// [`super::seqpacket::recv_once`](super::seqpacket) treats `ECONNRESET` as end +/// of file alongside the `POLLHUP` rule rather than choosing between them by +/// platform: one rule that accepts either signal is correct on both kernels. +/// +/// Asserted rather than only written down, so that a Darwin release which moves +/// to `POLLHUP`, or stops reporting the close at all, reds this test instead of +/// silently changing what the receive path depends on. +#[cfg(target_os = "macos")] +#[test] +fn darwin_signals_a_closed_dgram_peer_with_econnreset() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + drop(a); + + let mut poll = libc::pollfd { + fd: b.as_raw_fd(), + events: libc::POLLIN, + revents: 0, + }; + // SAFETY: `poll` points at one live pollfd and the call cannot block. + let rc = unsafe { libc::poll(&mut poll, 1, 0) }; + + let mut buf = [0u8; 64]; + let read = recv(b.as_raw_fd(), &mut buf); + let errno = read.as_ref().err().and_then(|e| e.raw_os_error()); + + assert_eq!( + errno, + Some(libc::ECONNRESET), + "Darwin no longer reports a closed SOCK_DGRAM peer as ECONNRESET. \ + Measured: poll rc={rc}, revents=0x{:04x}, recv={read:?}. The receive \ + path treats ECONNRESET as end of file and would now hang instead.", + poll.revents, + ); + assert!( + (poll.revents & libc::POLLHUP) == 0, + "Darwin has gained POLLHUP on a closed SOCK_DGRAM peer, revents=0x{:04x}. \ + Nothing breaks, since the receive path accepts either signal, but the \ + record here is now wrong and the SOCK_SEQPACKET comparison it rests on \ + should be re-read.", + poll.revents, + ); +} + +/// Whether a `SOCK_DGRAM` pair can carry the largest payload the API offers. +/// +/// Darwin bounds a unix-domain datagram with the `net.local.dgram.maxdgram` +/// sysctl, whose default is small, and it is a system tunable rather than +/// something this process can rely on. Linux has no equivalent ceiling on an +/// `AF_UNIX` datagram beyond the socket buffer. The API advertises a payload +/// limit of 1362 bytes to its clients, so a kernel that refuses a datagram that +/// size would break the contract the client was told. +/// +/// Written to answer on failure as well as on success: the assertion message +/// carries the largest size that did cross, so a negative result names the +/// actual ceiling rather than only saying the hoped-for one was not reached. +#[test] +fn a_dgram_pair_carries_the_largest_payload_the_api_advertises() { + /// The `max_payload` the API reports to a client, from `super::mod`'s + /// wire-derived limit. Duplicated rather than imported because the constant + /// lives inside a module gated to the platforms that have the listener, and + /// this test runs where that module does not exist. + const ADVERTISED_PAYLOAD: usize = 1362; + + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + + // Walk up rather than testing one size, so a failure reports the ceiling. + let mut largest = 0usize; + let mut buf = vec![0u8; ADVERTISED_PAYLOAD * 4]; + for size in [64, 256, 1024, ADVERTISED_PAYLOAD, ADVERTISED_PAYLOAD * 2] { + let payload = vec![0xA5u8; size]; + if send(a.as_raw_fd(), &payload).is_err() { + break; + } + match recv(b.as_raw_fd(), &mut buf) { + Ok(n) if n == size => largest = size, + _ => break, + } + } + + assert!( + largest >= ADVERTISED_PAYLOAD, + "a SOCK_DGRAM pair carried at most {largest} bytes, below the \ + {ADVERTISED_PAYLOAD} the API advertises to clients. On Darwin this is \ + the net.local.dgram.maxdgram ceiling and the port has to raise it, or \ + lower what it advertises, rather than let a client send what it was \ + told it could." + ); +} + +#[test] +fn a_datagram_queued_before_the_close_is_still_readable() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + assert_eq!(send(a.as_raw_fd(), &[7, 7, 7]).unwrap(), 3); + drop(a); + + // Data sent before the close must survive it. A kernel that discards the + // queue on close would lose a client's last datagram. + let mut buf = [0u8; 64]; + assert_eq!( + recv(b.as_raw_fd(), &mut buf).unwrap(), + 3, + "a datagram queued before the peer closed must still be delivered" + ); + assert_eq!(&buf[..3], &[7, 7, 7]); +} + +#[test] +fn an_empty_datagram_queued_before_the_close_is_not_read_as_the_close() { + let (a, b) = dgram_pair().expect("AF_UNIX SOCK_DGRAM socketpair"); + assert_eq!(send(a.as_raw_fd(), &[]).unwrap(), 0); + drop(a); + + // The ordering case the close rule is weakest against. The peer has closed, + // so POLLHUP is set, and the queued empty datagram also reads as zero + // bytes, so `zero && hung_up` cannot tell them apart. Whatever the kernel + // does here, the implementation has to handle it; this test records which + // kernel loses the datagram. + let mut buf = [0u8; 64]; + let read = recv(b.as_raw_fd(), &mut buf); + let hup = hung_up(b.as_raw_fd()); + assert!( + matches!(read, Ok(0)), + "expected the queued empty datagram to read as zero bytes, got {read:?}" + ); + assert!( + !hup, + "POLLHUP is set while an empty datagram is still queued, so a zero-byte \ + read plus POLLHUP cannot mean end of file on this kernel: the queued \ + datagram would be swallowed and reported as a close" + ); +} diff --git a/src/native/fdpass.rs b/src/native/fdpass.rs index 97400ed0..0bda4ccd 100644 --- a/src/native/fdpass.rs +++ b/src/native/fdpass.rs @@ -23,12 +23,26 @@ use tokio::net::UnixStream; /// Space for the control message, sized at runtime and checked against this. /// -/// `CMSG_SPACE(4)` is 24 bytes on the platforms this builds for. The array is -/// `u64` so it carries the alignment `cmsghdr` requires, and is larger than -/// needed so a platform with a wider header is caught by the assertion rather -/// than by memory corruption. +/// `CMSG_SPACE(4)` is 24 bytes on Linux and 16 on Darwin, whose `cmsghdr` is +/// 12 bytes and whose alignment is 4 rather than 8. The array is `u64` so it +/// carries the alignment `cmsghdr` requires, and at 64 bytes is larger than +/// either, so a platform with a wider header is caught by the assertion rather +/// than by memory corruption. Nothing computes from the number: the send path +/// checks the runtime `CMSG_SPACE` against this buffer's size and the receive +/// path offers the whole buffer. type CmsgBuf = [u64; 8]; +/// `recvmsg` flags that make a received descriptor close-on-exec. +/// +/// Linux and FreeBSD do it atomically with `MSG_CMSG_CLOEXEC`, which is the +/// only way to be certain no `fork` in another thread wins the race. macOS has +/// no equivalent flag, so there is nothing to pass and [`recv`] sets +/// `FD_CLOEXEC` on each descriptor afterwards instead. +#[cfg(not(target_os = "macos"))] +const RECV_FLAGS: libc::c_int = libc::MSG_CMSG_CLOEXEC; +#[cfg(target_os = "macos")] +const RECV_FLAGS: libc::c_int = 0; + /// Send `line` on `stream`, with `fd` in the ancillary data when given. /// /// The whole reply goes in one datagram-shaped `sendmsg`. A short write is @@ -154,9 +168,10 @@ pub(super) struct Chunk { /// the kernel closes the descriptor rather than queueing it, so the reply looks /// right and the flow is silently gone. /// -/// `MSG_CMSG_CLOEXEC` keeps a received descriptor out of a child the client -/// forks later. `EINTR` is retried, because a signal delivered during the wait -/// says nothing about the connection. +/// A received descriptor is kept out of a child the client forks later, by +/// [`RECV_FLAGS`] where the platform has a flag for it and by an `fcntl` on +/// each descriptor where it does not. `EINTR` is retried, because a signal +/// delivered during the wait says nothing about the connection. pub(super) fn recv(sock: RawFd, buf: &mut [u8]) -> io::Result { loop { let mut iov = libc::iovec { @@ -173,12 +188,21 @@ pub(super) fn recv(sock: RawFd, buf: &mut [u8]) -> io::Result { // SAFETY: `sock` is the caller's open socket, and `msg` describes // buffers that outlive the call. - let received = unsafe { libc::recvmsg(sock, &mut msg, libc::MSG_CMSG_CLOEXEC) }; + let received = unsafe { libc::recvmsg(sock, &mut msg, RECV_FLAGS) }; if received < 0 { let error = io::Error::last_os_error(); if error.kind() == io::ErrorKind::Interrupted { continue; } + // A listener's descriptor is a socket pair half, and Darwin reports + // its closed peer as ECONNRESET where Linux returns a zero-byte + // message. They are the same event, so it is reported as the empty + // chunk every caller here already reads as the far end going away. + // Doing it here rather than in each caller keeps `accept`'s + // documented EPIPE true on both platforms. + if error.raw_os_error() == Some(libc::ECONNRESET) { + return Ok(Chunk { len: 0, fd: None }); + } return Err(error); } @@ -186,6 +210,17 @@ pub(super) fn recv(sock: RawFd, buf: &mut [u8]) -> io::Result { // `msg_controllen`, and `control` is still alive. let mut fds = unsafe { take_fds(&msg) }; + // Darwin has no `MSG_CMSG_CLOEXEC`, so the flag is set here instead. + // Later than the atomic form and with the same window `seqpacket::pair` + // documents: a concurrent `fork` and `exec` in these few instructions + // would inherit the descriptor. A failure to set it is reported rather + // than ignored, because the descriptor is live either way and the + // caller must not be told the receive was clean. + #[cfg(target_os = "macos")] + for fd in &fds { + super::seqpacket::set_cloexec(fd.as_raw_fd())?; + } + // Whatever did arrive is taken before the truncation check, so nothing // leaks on that path: dropping an `OwnedFd` closes it. The connection // cannot continue either way, because a descriptor the kernel dropped diff --git a/src/native/mod.rs b/src/native/mod.rs index 87ec2f16..346fa7af 100644 --- a/src/native/mod.rs +++ b/src/native/mod.rs @@ -25,13 +25,15 @@ //! real limit: the transport MTU less the FIPS encapsulation and the four-byte //! port header. //! -//! **Platform support.** The listener is built on Linux and FreeBSD only, and -//! two separate things bound that. Windows has no `SCM_RIGHTS` and so no way to -//! hand a descriptor to another process at all. macOS has `SCM_RIGHTS` but does -//! not implement `SOCK_SEQPACKET` for `AF_UNIX`, so only the descriptor's -//! socket type is missing there; see [`seqpacket`] for what a macOS port would -//! have to settle. The gate is explicit rather than `cfg(unix)` so macOS fails -//! to build here instead of failing at `socketpair` on a running node. +//! **Platform support.** The listener is built on Linux, FreeBSD and macOS. +//! Windows is excluded and cannot be included: it has no `SCM_RIGHTS`, so there +//! is no way to hand a descriptor to another process at all, which is the whole +//! mechanism. macOS needed only the descriptor's socket type, since it has +//! `SCM_RIGHTS` but does not implement `SOCK_SEQPACKET` for `AF_UNIX`; see +//! [`seqpacket`] for the type it uses instead and for the measured difference +//! in how the two report a close. The gate is an explicit platform list rather +//! than `cfg(unix)` so a platform nobody has measured fails to build here +//! instead of failing at `socketpair` on a running node. //! //! - `protocol.rs` — the command types and the pure decisions over them. No //! I/O, no node state. @@ -53,23 +55,29 @@ pub mod link; pub mod protocol; pub mod registry; -#[cfg(any(target_os = "linux", target_os = "freebsd"))] +// Tests only, and compiled on every unix rather than only where the listener +// is, because its whole purpose is to compare one kernel's answer with +// another's. See the module header. +#[cfg(all(test, unix))] +mod dgram_probe; + +#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] pub mod client; -#[cfg(any(target_os = "linux", target_os = "freebsd"))] +#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] pub mod fdpass; -#[cfg(any(target_os = "linux", target_os = "freebsd"))] +#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] pub mod seqpacket; -#[cfg(any(target_os = "linux", target_os = "freebsd"))] +#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] pub use unix_impl::NativeApi; -#[cfg(any(target_os = "linux", target_os = "freebsd"))] +#[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] mod unix_impl { use super::fdpass; use super::link::{Accepted, DropReason, NativeMessage, Outbound, Outcome}; use super::protocol::{self, Command, Connect, Inject, Listen}; use super::registry::{Arrival, Datagram, FlowKey, Limits, RegistryError}; - use super::seqpacket::{Received, Seqpacket, pair, set_sndbuf}; + use super::seqpacket::{Received, Seqpacket, pair, set_rcvbuf, set_sndbuf}; use crate::config::NativeApiConfig; use crate::control::protocol::{Request, Response}; use crate::identity::{NodeAddr, decode_npub, encode_npub}; @@ -726,6 +734,29 @@ mod unix_impl { /// it too, and a listener has no connection to reach it through. fn wire_flow(per_flow: usize) -> std::io::Result<(Wiring, OwnedFd, mpsc::Sender)> { let (ours, theirs) = pair()?; + + // `hand_over` writes a flow's whole held batch onto this pair before the + // client has the descriptor, so nothing is reading while it is written + // and the batch has to fit in the kernel's buffer. + // + // **Both halves are sized because the two kernels charge different + // ones.** Linux accounts an `AF_UNIX` message against the sender's + // `SO_SNDBUF`, which is generous by default; BSD queues it in the + // receiver's `so_rcv`, whose `SOCK_DGRAM` default on Darwin is small + // enough that a two or three datagram batch fills it. Sizing only the + // sender, as the listener pair does, leaves the flow pair unbounded by + // anything this code sets on Darwin, and the first batch past the + // ceiling destroys the whole arriving flow before its client ever sees + // it. Neither call is fatal: a flow with a default-sized buffer works, + // it just holds less. + let budget = per_flow.max(1) * ARRIVAL_ALLOWANCE; + if let Err(error) = set_sndbuf(&ours, budget) { + warn!(error = %error, "Could not size a native API flow's send buffer"); + } + if let Err(error) = set_rcvbuf(&theirs, budget) { + warn!(error = %error, "Could not size a native API flow's receive buffer"); + } + let sock = Arc::new(Seqpacket::new(ours)?); let (sink, inbound) = mpsc::channel::(per_flow.max(1)); Ok((Wiring { sock, inbound }, theirs, sink)) @@ -936,13 +967,25 @@ mod unix_impl { /// What a failed hand-off write says about the client, for the counter. /// /// Both take the same three-part cleanup, so this decides only what an - /// operator is told: `EPIPE` is a client that closed its listener between - /// the arrival being taken off the queue and this write, which a healthy - /// client can lose, and anything else is a full send buffer, which is a - /// client that stopped reading. Reporting the two alike would leave a - /// normal close looking like a fault. + /// operator is told: a gone listener is a client that closed between the + /// arrival being taken off the queue and this write, which a healthy client + /// can lose, and anything else is a full send buffer, which is a client that + /// stopped reading. Reporting the two alike would leave a normal close + /// looking like a fault. + /// + /// **Three errnos mean "gone", because the platforms do not agree.** Linux + /// `SOCK_SEQPACKET` reports a closed peer as `EPIPE`. Darwin disconnects the + /// survivor of a `SOCK_DGRAM` pair instead, so its first send gives + /// `ECONNRESET` and later ones `EDESTADDRREQ`, and neither is `BrokenPipe`. + /// Matching on the kind alone would file every ordinary macOS listener close + /// under the counter an operator reads to find a wedged client. pub(super) fn why(error: &std::io::Error) -> DropReason { - if error.kind() == std::io::ErrorKind::BrokenPipe { + let gone = error.kind() == std::io::ErrorKind::BrokenPipe + || matches!( + error.raw_os_error(), + Some(libc::ECONNRESET) | Some(libc::EDESTADDRREQ) + ); + if gone { DropReason::ListenerGone } else { DropReason::ListenerNotReading @@ -1114,9 +1157,24 @@ mod unix_impl { } /// Write datagrams the node delivered onto the client's descriptor. + /// + /// **A full client buffer drops the datagram and keeps the flow.** On Linux + /// this never arrives here: the send reports `EAGAIN` and `Seqpacket::send` + /// waits for the client to drain. Darwin's `SOCK_DGRAM` has no sender-side + /// queue to wait on, so an unread client surfaces as `ENOBUFS` on the send + /// itself. Returning on it would end this flow's only writer while the + /// registration, the port and the reader all stayed alive, so every later + /// inbound datagram would be counted as a full queue for the rest of the + /// flow's life, and a client that resumed reading would never recover. + /// Dropping the datagram is what a datagram API does when the far end + /// cannot take it. async fn feed(sock: Arc, mut inbound: mpsc::Receiver) { while let Some(datagram) = inbound.recv().await { if let Err(error) = sock.send(&datagram).await { + if error.raw_os_error() == Some(libc::ENOBUFS) { + debug!(error = %error, "Native API flow write dropped a datagram"); + continue; + } debug!(error = %error, "Native API flow write failed"); return; } @@ -1204,7 +1262,10 @@ mod unix_impl { } } -#[cfg(all(test, any(target_os = "linux", target_os = "freebsd")))] +#[cfg(all( + test, + any(target_os = "linux", target_os = "freebsd", target_os = "macos") +))] mod tests { use super::link::{self, NativeMessage, Outbound}; use super::registry::{Limits, Registry}; @@ -1735,6 +1796,41 @@ mod tests { assert_eq!(reason, link::DropReason::ListenerGone); } + #[test] + fn a_gone_listener_is_recognised_by_every_errno_a_platform_uses_for_it() { + // The same event, spelled three ways. Linux SOCK_SEQPACKET reports a + // closed peer as EPIPE; Darwin disconnects the survivor of a SOCK_DGRAM + // pair, so its first send gives ECONNRESET and later ones EDESTADDRREQ. + // Matching on ErrorKind::BrokenPipe alone recognises only the first, and + // would file every ordinary macOS listener close under the counter an + // operator reads to find a client that has stopped reading. + // + // Built from raw errnos rather than ErrorKind, because that is the only + // form that distinguishes them: ECONNRESET maps to ConnectionReset and + // EDESTADDRREQ to Uncategorized, and neither is BrokenPipe. + use super::unix_impl::why; + use std::io::Error; + + for errno in [libc::EPIPE, libc::ECONNRESET, libc::EDESTADDRREQ] { + assert_eq!( + why(&Error::from_raw_os_error(errno)), + link::DropReason::ListenerGone, + "errno {errno} should count as a listener that went away" + ); + } + + // The discrimination has to survive: a full buffer is still a client + // that stopped reading, and folding everything into ListenerGone would + // pass the loop above while destroying what the counter is for. + for errno in [libc::ENOBUFS, libc::EAGAIN] { + assert_eq!( + why(&Error::from_raw_os_error(errno)), + link::DropReason::ListenerNotReading, + "errno {errno} should still count as a client not reading" + ); + } + } + #[test] fn the_hand_off_writes_what_a_flow_held_before_the_arrival_that_carries_it() { // The plan, which is the half of the ordering guarantee this test owns: diff --git a/src/native/seqpacket.rs b/src/native/seqpacket.rs index fe54cb60..321cf24f 100644 --- a/src/native/seqpacket.rs +++ b/src/native/seqpacket.rs @@ -10,14 +10,20 @@ //! [`AsyncFd`], which is tokio's supported way to put an arbitrary file //! descriptor under the reactor. //! -//! **Not available on macOS**, which does not implement `SOCK_SEQPACKET` for -//! `AF_UNIX`. Everything else this module needs does work there: `SCM_RIGHTS` -//! is supported, and `AsyncFd` is backed by kqueue. So a macOS port is a change -//! of socket type to `SOCK_DGRAM`, which macOS does support and which also -//! keeps message boundaries, plus one thing that has to be **measured on a Mac -//! rather than assumed**: whether a closed peer on a connected `SOCK_DGRAM` -//! pair is distinguishable from an empty datagram. The `POLLHUP` result below -//! was measured on Linux and BSD poll semantics for datagram sockets differ. +//! **macOS uses `SOCK_DGRAM` instead**, because it does not implement +//! `SOCK_SEQPACKET` for `AF_UNIX`. Everything else this module needs works +//! there: `SCM_RIGHTS` is supported, `AsyncFd` is backed by kqueue, and a +//! connected `SOCK_DGRAM` pair keeps message boundaries just as `SOCK_SEQPACKET` +//! does. +//! +//! The two types are **not** interchangeable at end of file, and that is the +//! whole difficulty of the port. Measured 2026-08-20 and recorded in +//! [`super::dgram_probe`]: on Linux 6.8 a closed peer on a connected +//! `SOCK_DGRAM` pair produces no signal whatsoever, leaving `revents` empty and +//! `recv` returning `EAGAIN`, which is exactly what an idle socket with a live +//! peer does. Darwin does report the close. So the socket type is chosen per +//! platform and the close rule below accepts either signal rather than assuming +//! the one this kernel happens to use. use std::io; use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd}; @@ -34,31 +40,78 @@ pub enum Received { Eof, } -/// Create a connected `SOCK_SEQPACKET` pair. +/// The socket type a flow's descriptor uses on this platform. +/// +/// `SOCK_SEQPACKET` everywhere it exists for `AF_UNIX`. macOS does not +/// implement it there, and `SOCK_DGRAM` is the replacement: it keeps message +/// boundaries, which is the property the API's contract with its clients rests +/// on. See the module header for what changes with it and what does not. +#[cfg(not(target_os = "macos"))] +const SOCK_TYPE: libc::c_int = libc::SOCK_SEQPACKET; +#[cfg(target_os = "macos")] +const SOCK_TYPE: libc::c_int = libc::SOCK_DGRAM; + +/// Create a connected socket pair of this platform's [`SOCK_TYPE`]. /// /// Both descriptors are close-on-exec so neither leaks into a child process. -/// Neither is set non-blocking here: the two halves are independent sockets, so -/// [`Seqpacket::new`] can make the daemon's half non-blocking for the reactor -/// while the half handed to the client stays blocking, which is what a client -/// calling `recv` in a loop expects. +/// Linux and FreeBSD take `SOCK_CLOEXEC` in `socketpair`'s type argument and +/// set it atomically; **macOS rejects it there**, so Darwin sets `FD_CLOEXEC` +/// with `fcntl` afterwards instead. The Darwin path has a window between the +/// two calls in which a concurrent `fork` and `exec` would inherit the +/// descriptors. It is accepted rather than closed because the alternative needs +/// a lock this module has no business holding, and the daemon spawns no child +/// on this path. +/// +/// Neither descriptor is set non-blocking here: the two halves are independent +/// sockets, so [`Seqpacket::new`] can make the daemon's half non-blocking for +/// the reactor while the half handed to the client stays blocking, which is +/// what a client calling `recv` in a loop expects. pub fn pair() -> io::Result<(OwnedFd, OwnedFd)> { + #[cfg(not(target_os = "macos"))] + let sock_type = SOCK_TYPE | libc::SOCK_CLOEXEC; + #[cfg(target_os = "macos")] + let sock_type = SOCK_TYPE; + let mut fds = [0 as libc::c_int; 2]; // SAFETY: `fds` is a two-element array of the type socketpair writes, and // the return value is checked before either descriptor is read. - let rc = unsafe { - libc::socketpair( - libc::AF_UNIX, - libc::SOCK_SEQPACKET | libc::SOCK_CLOEXEC, - 0, - fds.as_mut_ptr(), - ) - }; + let rc = unsafe { libc::socketpair(libc::AF_UNIX, sock_type, 0, fds.as_mut_ptr()) }; if rc != 0 { return Err(io::Error::last_os_error()); } // SAFETY: socketpair reported success, so both entries are open - // descriptors this process now owns. - Ok(unsafe { (OwnedFd::from_raw_fd(fds[0]), OwnedFd::from_raw_fd(fds[1])) }) + // descriptors this process now owns. Wrapping them before any further + // syscall means an error below closes them rather than leaking them. + let pair = unsafe { (OwnedFd::from_raw_fd(fds[0]), OwnedFd::from_raw_fd(fds[1])) }; + + #[cfg(target_os = "macos")] + { + set_cloexec(pair.0.as_raw_fd())?; + set_cloexec(pair.1.as_raw_fd())?; + } + + Ok(pair) +} + +/// Mark `fd` close-on-exec, for the platform that will not do it atomically. +/// +/// Read-modify-write rather than a bare set, so that any other flag the +/// descriptor carries survives. Used both here, where `socketpair` refuses +/// `SOCK_CLOEXEC`, and by [`super::fdpass`], where `recvmsg` has no +/// `MSG_CMSG_CLOEXEC` to pass. +#[cfg(target_os = "macos")] +pub(super) fn set_cloexec(fd: RawFd) -> io::Result<()> { + // SAFETY: the descriptor is open for the call. + let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) }; + if flags < 0 { + return Err(io::Error::last_os_error()); + } + // SAFETY: the descriptor is open and the flag word is the one just read. + let rc = unsafe { libc::fcntl(fd, libc::F_SETFD, flags | libc::FD_CLOEXEC) }; + if rc < 0 { + return Err(io::Error::last_os_error()); + } + Ok(()) } /// Size the send buffer of `fd`, in bytes. @@ -89,6 +142,20 @@ pub fn set_sndbuf(fd: &OwnedFd, bytes: usize) -> io::Result<()> { Ok(()) } +/// How long a receive waits on the reactor before retrying the syscall, on the +/// platform whose reactor cannot see a closed `SOCK_DGRAM` peer. +/// +/// This is a detection latency, not a poll interval for data: an arriving +/// datagram does wake the reactor normally, so ordinary traffic is unaffected +/// and this bound is only reached on an idle flow. What it bounds is how long a +/// closed flow keeps its port and its registry entry before the daemon notices. +/// +/// A quarter second is chosen against the cost, which is one timer per idle +/// flow and listener. Shortening it buys a faster reclaim of something nothing +/// is waiting on; lengthening it holds a dead flow's port longer. +#[cfg(target_os = "macos")] +const CLOSE_RETRY: Duration = Duration::from_millis(250); + /// The daemon's half of a flow's or a listener's socket pair, registered with /// the reactor. pub struct Seqpacket { @@ -125,6 +192,37 @@ impl Seqpacket { /// payload the API accepts, so a truncation means the client exceeded it. pub async fn recv(&self, buf: &mut [u8]) -> io::Result { loop { + // **A read before the wait, because on one platform the reactor + // cannot see a close.** Darwin's `unp_disconnect` takes a different + // branch for `SOCK_DGRAM` than for `SOCK_STREAM`: it clears + // `SS_ISCONNECTED` and latches `so_error`, and calls none of + // `sorwakeup`, `socantrcvmore` or `soisdisconnected`. So a closed + // peer wakes no knote. The registration is edge-triggered and was + // made while the socket was healthy, so waiting first would park + // for ever on exactly the event this call exists to report. The + // latched error is visible to a syscall, and only to a syscall. + // + // On Linux this attempt costs one `recv` returning `EAGAIN` before + // the wait, and changes nothing else. + match recv_once(self.raw(), buf) { + Err(error) if error.kind() == io::ErrorKind::WouldBlock => {} + other => return other, + } + + // **The wait is bounded where the close raises no event**, so a peer + // that goes away while this task is parked is noticed on the next + // attempt rather than never. `so_error` is latched until a read + // consumes it, so the bound sets detection latency and cannot lose + // the event. Everywhere else the readiness is authoritative and the + // wait stays unbounded. + #[cfg(target_os = "macos")] + let mut guard = match tokio::time::timeout(CLOSE_RETRY, self.inner.readable()).await { + Ok(ready) => ready?, + // The bound expired. The loop retries the syscall, which is the + // only thing that can see this platform's close. + Err(_elapsed) => continue, + }; + #[cfg(not(target_os = "macos"))] let mut guard = self.inner.readable().await?; let attempt = guard.try_io(|inner| recv_once(inner.get_ref().as_raw_fd(), buf)); match attempt { @@ -184,11 +282,23 @@ impl Seqpacket { /// This matters because reading an empty datagram as a close would let a client /// tear down its own flow by sending nothing, and the defect would present as a /// spurious disconnect. +/// +/// **`ECONNRESET` is treated as end of file too**, for the platform whose +/// datagram sockets report a close that way rather than through `POLLHUP`. Both +/// arms are compiled everywhere rather than split by `cfg`, because a rule that +/// accepts either signal is correct on both kernels and a rule that assumes one +/// would be silently wrong on the other. `EAGAIN` is deliberately not in that +/// company: it means the socket is empty and the peer is alive, so it stays an +/// error and the caller waits for readiness again. fn recv_once(fd: RawFd, buf: &mut [u8]) -> io::Result { // SAFETY: the descriptor is open and the pointer and length describe `buf`. let n = unsafe { libc::recv(fd, buf.as_mut_ptr().cast(), buf.len(), 0) }; if n < 0 { - return Err(io::Error::last_os_error()); + let err = io::Error::last_os_error(); + if err.raw_os_error() == Some(libc::ECONNRESET) { + return Ok(Received::Eof); + } + return Err(err); } if n == 0 && peer_hung_up(fd) { return Ok(Received::Eof); @@ -196,11 +306,37 @@ fn recv_once(fd: RawFd, buf: &mut [u8]) -> io::Result { Ok(Received::Datagram(n as usize)) } +/// Translate a platform's spelling of "the peer is gone" into this API's. +/// +/// Darwin reports a closed `SOCK_DGRAM` peer as `ECONNRESET`, on the read path +/// and the write path alike, where Linux `SOCK_SEQPACKET` reports `EPIPE` on a +/// write and a zero-byte read plus `POLLHUP` on a read. The client's contract +/// names `EPIPE` for that condition and says so in its documentation, so the +/// difference is translated once here rather than at each call site. Every +/// other errno passes through untouched, because only this one condition has +/// two spellings. +pub(super) fn peer_gone_as_epipe(error: io::Error) -> io::Error { + if error.raw_os_error() == Some(libc::ECONNRESET) { + return io::Error::from_raw_os_error(libc::EPIPE); + } + error +} + /// Whether the peer has closed its half of the pair. /// -/// `events` is left empty on purpose: `POLLHUP` is reported in `revents` -/// whether or not it was requested, which was confirmed by measurement rather -/// than assumed. The poll does not block, and runs only on the zero-byte path. +/// `POLLIN` is requested even though only `POLLHUP` is read. `POLLHUP` is +/// reported in `revents` whether or not it was asked for, which holds on Linux +/// and was measured there; **an empty `events` is not enough on Darwin**, where +/// a poll that requests nothing registers no filter and returns 0 with an empty +/// `revents` for a peer that has in fact closed. That was measured too, in +/// [`super::dgram_probe`], and asking for `POLLIN` costs nothing on either +/// platform because the result is masked to `POLLHUP` regardless. +/// +/// This function is not what detects a close on Darwin: `ECONNRESET` arrives +/// first and both callers act on it before reaching the zero-byte path. It is +/// corrected so that it answers truthfully wherever it is called from, rather +/// than being left as a function whose contract holds on one platform. +/// The poll does not block, and runs only on the zero-byte path. /// /// Visible within [`super`] because the client half needs the same /// discrimination on the same socket pair: the rule belongs to the pair, not to @@ -208,7 +344,7 @@ fn recv_once(fd: RawFd, buf: &mut [u8]) -> io::Result { pub(super) fn peer_hung_up(fd: RawFd) -> bool { let mut poll = libc::pollfd { fd, - events: 0, + events: libc::POLLIN, revents: 0, }; // SAFETY: `poll` points at one live pollfd and the call cannot block. @@ -216,6 +352,35 @@ pub(super) fn peer_hung_up(fd: RawFd) -> bool { rc > 0 && (poll.revents & libc::POLLHUP) != 0 } +/// Size the receive buffer of `fd`, in bytes. +/// +/// The counterpart to [`set_sndbuf`], and needed because the two kernels charge +/// a queued `AF_UNIX` message to different ends: Linux to the sender's +/// `SO_SNDBUF`, BSD to the receiver's `so_rcv`. Sizing only one leaves the +/// queue bounded by a system default on the other platform. As with the send +/// buffer, the value is not read back and asserted, because a kernel is free to +/// double it or clamp it to its own minimum. +pub(super) fn set_rcvbuf(fd: &OwnedFd, bytes: usize) -> io::Result<()> { + let size = libc::c_int::try_from(bytes).map_err(|_| { + io::Error::new(io::ErrorKind::InvalidInput, "receive buffer size too large") + })?; + // SAFETY: the descriptor is open for the call, and the pointer and length + // describe one `c_int` that outlives it. + let rc = unsafe { + libc::setsockopt( + fd.as_raw_fd(), + libc::SOL_SOCKET, + libc::SO_RCVBUF, + std::ptr::addr_of!(size).cast(), + std::mem::size_of::() as libc::socklen_t, + ) + }; + if rc != 0 { + return Err(io::Error::last_os_error()); + } + Ok(()) +} + /// Set or clear a receive or send timeout on a socket. /// /// `option` is `SO_RCVTIMEO` or `SO_SNDTIMEO`. `None` clears the timeout, which @@ -232,9 +397,15 @@ pub(super) fn set_timeout(fd: RawFd, option: libc::c_int, dur: Option) } // Saturating rather than wrapping: a duration past `time_t` becomes the // longest wait the kernel can express, which is the caller's intent. + // + // `tv_usec` is cast rather than converted because `suseconds_t` is + // `i64` on Linux and `i32` on Darwin. `From` does not exist for the + // Darwin width and `try_from` is a clippy error on the Linux one, so a + // cast is the only form that compiles on both. It cannot truncate: + // `subsec_micros` is below 1_000_000 by construction and fits either. Some(d) => libc::timeval { tv_sec: d.as_secs().min(libc::time_t::MAX as u64) as libc::time_t, - tv_usec: libc::suseconds_t::from(d.subsec_micros()), + tv_usec: d.subsec_micros() as libc::suseconds_t, }, None => libc::timeval { tv_sec: 0, @@ -322,6 +493,22 @@ mod tests { StdUnixStream::from(fd) } + /// One receive, bounded in time. + /// + /// Every assertion in this module is about something arriving on the + /// descriptor, and the failure mode of each is that it never does. An + /// unbounded await against that defect parks for ever: the suite wedges + /// with no named failure and no diagnostic, and a hang is not a red. This + /// matters more since the socket type became platform-dependent, because a + /// kernel that does not report a close makes the end-of-file test the one + /// that hangs, and it would hang on a runner rather than here. + async fn recv_bounded(sock: &Seqpacket, buf: &mut [u8]) -> Received { + tokio::time::timeout(Duration::from_secs(5), sock.recv(buf)) + .await + .expect("nothing arrived on the descriptor within 5s: it is still waiting") + .expect("recv failed") + } + #[tokio::test] async fn message_boundaries_survive_in_both_directions() { let (daemon, theirs) = pair().unwrap(); @@ -335,11 +522,11 @@ mod tests { theirs.write_all(b"three").unwrap(); let mut buf = [0u8; 64]; - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Datagram(3)); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Datagram(3)); assert_eq!(&buf[..3], b"one"); - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Datagram(3)); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Datagram(3)); assert_eq!(&buf[..3], b"two"); - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Datagram(5)); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Datagram(5)); assert_eq!(&buf[..5], b"three"); daemon.send(b"alpha").await.unwrap(); @@ -358,7 +545,7 @@ mod tests { drop(client(theirs)); let mut buf = [0u8; 64]; - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Eof); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Eof); } #[tokio::test] @@ -380,8 +567,31 @@ mod tests { theirs.write_all(b"after").unwrap(); let mut buf = [0u8; 64]; - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Datagram(0)); - assert_eq!(daemon.recv(&mut buf).await.unwrap(), Received::Datagram(5)); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Datagram(0)); + assert_eq!(recv_bounded(&daemon, &mut buf).await, Received::Datagram(5)); + } + + #[test] + fn both_halves_of_a_pair_are_close_on_exec() { + // Asserted rather than assumed because the platforms disagree on how it + // is set: Linux and FreeBSD get it atomically from socketpair's type + // argument, macOS rejects it there and has to use a second fcntl. A + // missing flag leaks both descriptors into any child the daemon spawns, + // which is silent, so nothing else would report it. + let (daemon, theirs) = pair().unwrap(); + for (label, fd) in [ + ("daemon", daemon.as_raw_fd()), + ("client", theirs.as_raw_fd()), + ] { + // SAFETY: the descriptor is open and owned by this frame. + let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) }; + assert!(flags >= 0, "{label} half: F_GETFD failed"); + assert_eq!( + flags & libc::FD_CLOEXEC, + libc::FD_CLOEXEC, + "{label} half of the pair is not close-on-exec" + ); + } } #[tokio::test] diff --git a/src/node/dataplane/rx_loop.rs b/src/node/dataplane/rx_loop.rs index 43723ccf..8dee06ad 100644 --- a/src/node/dataplane/rx_loop.rs +++ b/src/node/dataplane/rx_loop.rs @@ -135,8 +135,8 @@ impl Node { drop(control_tx); // Native datagram API socket. Experimental, off by default, and built - // on Linux and FreeBSD only: Windows has no way to pass a descriptor - // between processes, and macOS has no AF_UNIX SOCK_SEQPACKET. Bound + // on Linux, FreeBSD and macOS: Windows has no way to pass a descriptor + // between processes at all, which is the mechanism itself. Bound // synchronously so a bad path or a socket already in use is reported // here, before the node starts serving, rather than at whatever later // moment a spawned bind happened to run. @@ -154,7 +154,7 @@ impl Node { // there is no listener the channel exists and nothing ever sends, // and the guard keeps it open so the arm never sees a closed // receiver. - #[cfg(any(target_os = "linux", target_os = "freebsd"))] + #[cfg(any(target_os = "linux", target_os = "freebsd", target_os = "macos"))] let guard = { let mut guard = Some(tx.clone()); if self.config().node.native_api.enabled { @@ -174,7 +174,7 @@ impl Node { } guard }; - #[cfg(not(any(target_os = "linux", target_os = "freebsd")))] + #[cfg(not(any(target_os = "linux", target_os = "freebsd", target_os = "macos")))] let guard = Some(tx.clone()); (rx, guard) };