Stop using SOCK_SEQPACKET on FreeBSD, which is not a record socket there

The FreeBSD package job had stalled seven consecutive runs, the last two
killed by the thirty-minute bound and the five before it burning up to
six hours each. The job builds and gets into the unit tests; libtest
names the culprit itself in the run log, in a line the earlier reading of
that log missed:

  test native::client::tests::an_empty_datagram_is_not_reported_as_a_closed_flow
  has been running for over 60 seconds

That test sends a zero-length datagram and then reads it back with an
unbounded blocking recv. FreeBSD accepts AF_UNIX SOCK_SEQPACKET and
returns a socket that is not an atomic-record socket: seqpacketproto
carries no PR_ATOMIC and shares its send and receive handlers with
SOCK_STREAM. A zero-length send with no control data therefore queues
nothing, wakes nobody, and returns success, so the recv waits for a
message that was never delivered. Nothing bounds it: plain cargo test has
no per-test deadline, so one blocked thread holds the whole binary open.
The same kernel fact accounts for the five boundary tests that failed
within a third of a second in the same run, which needed no separate
explanation.

macOS already takes SOCK_DGRAM here because it has no AF_UNIX
SOCK_SEQPACKET at all. FreeBSD needs the same substitution for a
different reason, so it joins that arm, along with the close-detection
retry the arm carries. The cfg is written by exclusion so that a new unix
target gets the Linux arm, which fails loudly by refusing the socketpair
rather than quietly losing messages.

Two things guard against a repeat rather than fixing this instance. Every
blocking read in these tests now carries a deadline, so a platform that
swallows a message fails with an assertion naming the flow instead of
wedging the suite. And the FreeBSD job runs its tests under a wall-clock
bound, so a hang is a named step failure in fifteen minutes rather than a
job cancelled at the ceiling, with the output kept up to the kill because
that output is what names the blocked test.

What this does not establish, stated plainly because the code comments
would otherwise imply otherwise: nobody ran any of this on FreeBSD. The
mechanism is read out of the FreeBSD kernel source and inferred from a CI
log whose shape matches it. The freebsd-gated probes added alongside are
what would turn that into a measurement, and the next run of the job is
the first real test.
This commit is contained in:
Johnathan Corgan
2026-08-22 12:55:52 +01:00
parent d69325a2a3
commit cd1957c894
5 changed files with 312 additions and 48 deletions
+22 -1
View File
@@ -102,7 +102,28 @@ jobs:
# CI — the main CI matrix is Linux-only, and a release build
# compiles no #[cfg(test)] code (AF-prefix strip round-trips,
# platform module, config path gates).
cargo test
#
# Bounded, because libtest has no per-test deadline and this job
# runs plain `cargo test` rather than nextest, so one test that
# blocks on a syscall holds the whole binary open with nothing to
# end it. Six runs in August 2026 did exactly that. The bound is on
# the suite rather than on the job so the failure is a named step
# failure at fifteen minutes instead of the job ceiling at thirty,
# and `timeout` leaves the test binary's stdout untouched up to the
# kill: that stdout is what carries libtest's "has been running for
# over N seconds" lines, which are what name the blocked test.
#
# This bounds the damage; it does not fix anything. A test that can
# block for ever is a defect at the test, and the ones this suite
# has are bounded where they are written.
if ! timeout -s KILL 900 cargo test; then
echo "FAIL: cargo test failed, or did not finish within its 900s bound." >&2
echo " If the output above stops mid-run, look for libtest's" >&2
echo " 'has been running for over' lines: they name the test" >&2
echo " that blocked, and the tests that never reported at all" >&2
echo " are the rest of the answer." >&2
exit 1
fi
packaging/freebsd/build-pkg.sh \
--version "$FREEBSD_PACKAGE_VERSION" \
+76 -29
View File
@@ -811,8 +811,70 @@ mod tests {
}
}
/// How long a test waits for something that should already be there.
///
/// Every assertion in this module about a datagram or an arrival has the
/// same failure mode: the thing never comes. A blocking `recv` against that
/// defect parks for ever, so the suite wedges with no named failure and no
/// diagnostic instead of reporting one, and **a hang is not a red**. That is
/// not hypothetical: it is what a kernel whose `AF_UNIX SOCK_SEQPACKET`
/// drops a zero-length message did to this suite on a FreeBSD runner, six
/// runs in a row, until the job's own ceiling killed it.
///
/// The bound is not a workaround for a slow machine. It is what makes the
/// assertion decidable: the descriptors here are socket pairs this thread
/// already wrote to, so anything that is coming has arrived, and a wait past
/// this is a wait that will not end. Loose enough that a loaded runner does
/// not trip it, short enough that a real block reds within one test.
const DEADLINE: Duration = Duration::from_secs(5);
/// One receive on a flow, bounded by [`DEADLINE`].
///
/// The deadline is set for the call and put back afterwards, so a test can
/// still assert what the flow's own deadline is. A blocked receive panics
/// naming the flow rather than returning `WouldBlock`, because no caller
/// here asked for a deadline and a `WouldBlock` they did not ask for would
/// be read as the defect it is hiding.
fn recv_bounded(flow: &FipsStream, buf: &mut [u8]) -> io::Result<usize> {
let previous = flow.read_timeout().expect("the descriptor is open");
flow.set_read_timeout(Some(DEADLINE))
.expect("the descriptor is open");
let outcome = flow.recv(buf);
flow.set_read_timeout(previous)
.expect("the descriptor is open");
if let Err(error) = &outcome {
assert!(
error.kind() != io::ErrorKind::WouldBlock,
"nothing arrived on the flow within {DEADLINE:?} and the recv \
was still waiting: the datagram this asserts on was never \
delivered, or the peer closed without this platform saying so"
);
}
outcome
}
/// A listener over `fd`, with [`DEADLINE`] on its accepts.
///
/// Every test here writes the arrival onto the far half before accepting,
/// so a blocked accept means the message was not delivered. Bounded for the
/// same reason [`recv_bounded`] is: without it that is a hang rather than a
/// failure.
fn listener(fd: OwnedFd) -> FipsListener {
seqpacket::set_timeout(fd.as_raw_fd(), libc::SO_RCVTIMEO, Some(DEADLINE))
.expect("the descriptor is open");
FipsListener {
fd,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
}
}
/// A stream whose far end is a socket the test drives, standing in for the
/// daemon's half of a real flow.
///
/// The far half carries [`DEADLINE`] from the start, because the tests read
/// it directly rather than through [`recv_bounded`]. The flow half does not:
/// two tests assert that a fresh flow has no deadline, which is the API's
/// documented default, so the flow's bound is applied per call instead.
fn stream(max: usize) -> (FipsStream, UnixStream) {
let (ours, theirs) = seqpacket::pair().expect("a socket pair should be available");
let stream = FipsStream::wired(
@@ -825,7 +887,10 @@ mod tests {
max,
},
);
(stream, UnixStream::from(theirs))
let far = UnixStream::from(theirs);
far.set_read_timeout(Some(DEADLINE))
.expect("the descriptor is open");
(stream, far)
}
#[test]
@@ -972,10 +1037,7 @@ mod tests {
#[test]
fn a_nonblocking_listener_reports_would_block_rather_than_waiting_for_a_flow() {
let (ours, _theirs) = seqpacket::pair().expect("a socket pair should be available");
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
listener
.set_nonblocking(true)
@@ -991,10 +1053,7 @@ mod tests {
assert_eq!(flow.as_fd().as_raw_fd(), flow.as_raw_fd());
let (ours, _theirs) = seqpacket::pair().unwrap();
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
assert_eq!(listener.as_fd().as_raw_fd(), listener.as_raw_fd());
}
@@ -1216,10 +1275,7 @@ mod tests {
let (passed, held) = seqpacket::pair().unwrap();
let held = UnixStream::from(held);
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
// The daemon writes the peer's opening datagram onto the flow's own
// half first, and the arrival that carries that flow's descriptor
@@ -1238,7 +1294,7 @@ mod tests {
assert_eq!(flow.max_payload(), 1362);
let mut buf = [0u8; 64];
assert_eq!(flow.recv(&mut buf).unwrap(), 7);
assert_eq!(recv_bounded(&flow, &mut buf).unwrap(), 7);
assert_eq!(&buf[..7], b"opening");
// The descriptor is the flow's own half and nothing else's.
@@ -1253,10 +1309,7 @@ mod tests {
// could only be discovered by blocking in it.
let (ours, theirs) = seqpacket::pair().unwrap();
let (passed, _held) = seqpacket::pair().unwrap();
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
let mut poll = libc::pollfd {
fd: listener.as_raw_fd(),
@@ -1278,10 +1331,7 @@ mod tests {
#[test]
fn an_arrival_with_no_descriptor_is_reported_rather_than_taken_as_a_flow() {
let (ours, theirs) = seqpacket::pair().unwrap();
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
fdpass::send_once(theirs.as_raw_fd(), ARRIVAL, None).unwrap();
assert_eq!(
@@ -1294,10 +1344,7 @@ mod tests {
fn incoming_yields_the_same_flows_accept_would() {
let (ours, theirs) = seqpacket::pair().unwrap();
let (passed, _held) = seqpacket::pair().unwrap();
let listener = FipsListener {
fd: ours,
local: FipsAddr::new(codec::pton(NODE).unwrap(), 4242),
};
let listener = listener(ours);
fdpass::send_once(theirs.as_raw_fd(), ARRIVAL, Some(passed.as_fd())).unwrap();
let flow = listener.incoming().next().unwrap().unwrap();
@@ -1318,7 +1365,7 @@ mod tests {
(&far).write_all(b"back").unwrap();
let mut buf = [0u8; 64];
assert_eq!(flow.recv(&mut buf).unwrap(), 4);
assert_eq!(recv_bounded(&flow, &mut buf).unwrap(), 4);
assert_eq!(&buf[..4], b"back");
}
@@ -1338,10 +1385,10 @@ mod tests {
assert_eq!(sent, 0, "{}", io::Error::last_os_error());
let mut buf = [0u8; 64];
assert_eq!(flow.recv(&mut buf).unwrap(), 0);
assert_eq!(recv_bounded(&flow, &mut buf).unwrap(), 0);
drop(far);
let error = flow.recv(&mut buf).unwrap_err();
let error = recv_bounded(&flow, &mut buf).unwrap_err();
assert_eq!(error.raw_os_error(), Some(libc::EPIPE));
}
+143
View File
@@ -319,3 +319,146 @@ fn an_empty_datagram_queued_before_the_close_is_not_read_as_the_close() {
datagram would be swallowed and reported as a close"
);
}
/// Whether FreeBSD signals a closed `SOCK_DGRAM` peer, and how.
///
/// The open question this platform's arm of [`super::seqpacket`] rests on, and
/// the counterpart of [`darwin_reports_a_closed_dgram_peer_somehow`]. The three
/// kernels do three different things: Linux reports nothing at all, Darwin
/// reports `ECONNRESET`, and FreeBSD has not been measured. `recv_once` accepts
/// either `ECONNRESET` or a zero-byte read plus `POLLHUP`, so either signal is
/// enough and the assertion does not care which.
///
/// **A failure here is the measurement coming back negative, not a
/// regression.** It says FreeBSD behaves as Linux does, that a `SOCK_DGRAM`
/// descriptor carries no close signal, and that the flow's end of file has to
/// come from the client's line-protocol connection instead of from the flow
/// descriptor. That would be a real gap in the port, and it should red here,
/// with the errno in the message, rather than present as a flow that never
/// reports `EPIPE`.
///
/// The message carries what was actually read either way, because a negative
/// result is the useful one and it has to name what the kernel did.
#[cfg(target_os = "freebsd")]
#[test]
fn freebsd_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,
"FreeBSD reports nothing when a connected SOCK_DGRAM peer closes: \
POLLHUP unset and recv gave {read:?}, which is what an idle socket \
with a live peer gives. The flow descriptor cannot carry the close on \
this platform, and seqpacket's end-of-file rule needs another source."
);
}
/// Create a connected `AF_UNIX` `SOCK_SEQPACKET` pair, for the one platform
/// that accepts the constant without implementing its semantics.
///
/// Deliberately separate from [`dgram_pair`] and used only by the two tests
/// below: nothing in the daemon opens one of these on FreeBSD any more, and
/// these tests exist to record why.
#[cfg(target_os = "freebsd")]
fn seqpacket_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_SEQPACKET, 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])) })
}
/// FreeBSD's `AF_UNIX` `SOCK_SEQPACKET` does not keep message boundaries.
///
/// This is why the socket type is chosen per platform rather than taken from
/// the constant's name. `seqpacketproto` in `uipc_usrreq.c` carries
/// `PR_CONNREQUIRED | PR_CAPATTACH | PR_SOCKBUF` and no `PR_ATOMIC`, and its
/// `pr_sosend` and `pr_soreceive` are the `SOCK_STREAM` handlers, which split a
/// read only at a control mbuf, at `M_EOR`, or when the buffer fills. Nothing
/// in this daemon ever sends `MSG_EOR`, so consecutive sends coalesce.
///
/// **A failure here is the measurement coming back positive, not a
/// regression.** It says this kernel now gives `AF_UNIX` `SOCK_SEQPACKET` the
/// record semantics Linux gives it, and that FreeBSD could move back onto the
/// `SOCK_SEQPACKET` arm and regain the close signal that comes with it.
#[cfg(target_os = "freebsd")]
#[test]
fn freebsd_seqpacket_does_not_keep_message_boundaries() {
let (a, b) = seqpacket_pair().expect("AF_UNIX SOCK_SEQPACKET socketpair");
assert_eq!(send(a.as_raw_fd(), &[1, 2, 3]).unwrap(), 3);
assert_eq!(send(a.as_raw_fd(), &[4, 5]).unwrap(), 2);
let mut buf = [0u8; 64];
let first = recv(b.as_raw_fd(), &mut buf);
assert_eq!(
first.as_ref().ok().copied(),
Some(5),
"expected the two sends to coalesce into one five-byte read, which is \
what a SOCK_STREAM-backed seqpacket does; got {first:?}. If this is \
now 3, FreeBSD has gained record semantics and seqpacket::SOCK_TYPE \
should be reconsidered."
);
}
/// FreeBSD's `AF_UNIX` `SOCK_SEQPACKET` accepts a zero-length send and queues
/// nothing, which is the mechanism behind six stalled FreeBSD CI runs.
///
/// The send handler's whole body is inside `while (mc.mc_len + cmc.mc_len > 0)`,
/// so a zero-length send with no control data enters no iteration: nothing is
/// appended to the receive buffer, `sorwakeup_locked` is never reached, and
/// `send` returns 0 as success. A client that then does a blocking `recv`
/// waiting for that message waits for ever, and libtest has no per-test
/// deadline to end it.
///
/// The measurement is what makes that a fact rather than a reading of someone
/// else's source tree. The recv is non-blocking, so this test cannot itself
/// become the hang it describes.
///
/// **A failure here is the measurement coming back positive, not a
/// regression**, and it would mean the diagnosis behind the socket-type change
/// is wrong and should be reopened.
#[cfg(target_os = "freebsd")]
#[test]
fn freebsd_seqpacket_drops_a_zero_length_message_instead_of_delivering_it() {
let (a, b) = seqpacket_pair().expect("AF_UNIX SOCK_SEQPACKET socketpair");
// SAFETY: the descriptor is open and owned by `a`; a null pointer with a
// zero length is what a zero-length send is.
let sent = unsafe { libc::send(a.as_raw_fd(), std::ptr::null(), 0, 0) };
assert_eq!(
sent,
0,
"a zero-length send was refused outright: {}",
io::Error::last_os_error()
);
// A second, non-empty send follows so the test can tell "the empty message
// was dropped" from "the empty message has not arrived yet". If the empty
// one were queued it would be read first, before these five bytes.
assert_eq!(send(a.as_raw_fd(), b"after").unwrap(), 5);
let mut buf = [0u8; 64];
let first = recv(b.as_raw_fd(), &mut buf);
assert_eq!(
first.as_ref().ok().copied(),
Some(5),
"expected the zero-length message to have been dropped and the next \
message to be first in the queue; got {first:?}. A 0 here means \
FreeBSD does deliver an empty seqpacket message after all, and the \
diagnosis of the stalled FreeBSD runs is wrong."
);
}
+37 -3
View File
@@ -550,9 +550,21 @@ mod unix_impl {
// The buffer is the only bound on arrivals a client has stopped
// reading, so a failure to size it is worth a warning rather than a
// refusal: the listener still works, with the system default.
if let Err(error) = set_sndbuf(&ours, queue * ARRIVAL_ALLOWANCE) {
//
// **Both halves, because the two kernels charge different ones.**
// Linux accounts an `AF_UNIX` message against the sender's
// `SO_SNDBUF`; BSD queues it in the receiver's `so_rcv`, whose
// `SOCK_DGRAM` default is small. Sizing only the sender leaves the
// declared backlog inert on FreeBSD and Darwin, where arrivals are
// then dropped well below it. This is the same correction
// `wire_flow` already makes for a flow's pair.
let budget = queue * ARRIVAL_ALLOWANCE;
if let Err(error) = set_sndbuf(&ours, budget) {
warn!(error = %error, "Could not size a native API listener's send buffer");
}
if let Err(error) = set_rcvbuf(&theirs, budget) {
warn!(error = %error, "Could not size a native API listener's receive buffer");
}
let sock = match Seqpacket::new(ours) {
Ok(sock) => sock,
@@ -1333,12 +1345,31 @@ mod tests {
serde_json::to_value(response).unwrap()
}
/// How long a test waits for something that should already be on a
/// descriptor.
///
/// Every read below is of a message this test already caused to be written,
/// so anything that is coming has arrived and a longer wait is one that will
/// not end. These tests run inside a current-thread runtime and read with
/// blocking syscalls, so an unbounded read here does not merely fail to
/// return: it stops the runtime that would have produced what it is waiting
/// for. Without a bound that is a wedged suite with no named failure, which
/// is how a FreeBSD portability bug turned into six cancelled CI runs.
const DEADLINE: std::time::Duration = std::time::Duration::from_secs(5);
/// Put [`DEADLINE`] on a descriptor a test reads with blocking calls.
fn bounded(sock: StdUnixStream) -> StdUnixStream {
sock.set_read_timeout(Some(DEADLINE))
.expect("the descriptor is open");
sock
}
/// Send one command that opens a flow, returning the reply and descriptor.
async fn open(connection: &mut Connection, line: &str) -> (serde_json::Value, StdUnixStream) {
let (response, fd) = connection.answer(line.as_bytes()).await;
let value = serde_json::to_value(response).unwrap();
let fd = fd.expect("this command should carry a descriptor");
(value, StdUnixStream::from(fd))
(value, bounded(StdUnixStream::from(fd)))
}
/// Bind a listener and keep its descriptor, the way a client does.
@@ -1355,13 +1386,16 @@ mod tests {
/// Goes through the same `recvmsg` a client process runs, so what these
/// tests assert about the ancillary framing is what a client would see.
fn accept(listener: &StdUnixStream) -> (serde_json::Value, StdUnixStream) {
// The listener descriptor comes back from `listen` already bounded, so
// this recvmsg cannot park: an arrival the test wrote is there, or the
// read fails and names the listener.
let mut buf = [0u8; 4096];
let chunk = super::fdpass::recv(listener.as_raw_fd(), &mut buf)
.expect("an arrival should be readable on the listener");
let fd: OwnedFd = chunk.fd.expect("an arrival carries the flow's descriptor");
let value: serde_json::Value =
serde_json::from_slice(&buf[..chunk.len]).expect("the arrival is one JSON object");
(value, StdUnixStream::from(fd))
(value, bounded(StdUnixStream::from(fd)))
}
/// Deliver a datagram as though a peer had sent it.
+34 -15
View File
@@ -10,11 +10,22 @@
//! [`AsyncFd`], which is tokio's supported way to put an arbitrary file
//! descriptor under the reactor.
//!
//! **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.
//! **macOS and FreeBSD use `SOCK_DGRAM` instead**, for two different reasons
//! that come to the same thing. macOS does not implement `SOCK_SEQPACKET` for
//! `AF_UNIX` at all. FreeBSD accepts the constant and gives back a socket that
//! is not an atomic-record socket: `seqpacketproto` carries no `PR_ATOMIC` and
//! shares its send and receive handlers with `SOCK_STREAM`, so consecutive
//! messages coalesce and a zero-length message is queued nowhere and delivered
//! never. **This was read out of the FreeBSD kernel source and inferred from a
//! CI log; nobody has run it on FreeBSD.** What was observed is a CI run in
//! which the zero-length-message test hung until the job was killed and five
//! boundary tests failed within 360 ms, which is the shape the source predicts.
//! The probes in [`super::dgram_probe`] are what would turn that into a
//! measurement, and they have not been run on FreeBSD either.
//! Everything else this module needs works on both: `SCM_RIGHTS` is supported,
//! `AsyncFd` is backed by kqueue, and a connected `SOCK_DGRAM` pair keeps
//! message boundaries and delivers an empty datagram, which `SOCK_SEQPACKET`
//! does on Linux and does not here.
//!
//! 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
@@ -42,13 +53,21 @@ pub enum Received {
/// 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"))]
/// `SOCK_SEQPACKET` where `AF_UNIX` genuinely implements it as an atomic-record
/// socket, `SOCK_DGRAM` where it does not. macOS has no `AF_UNIX`
/// `SOCK_SEQPACKET`; FreeBSD has the name without the record semantics. On both
/// the replacement is `SOCK_DGRAM`, which keeps message boundaries and delivers
/// a zero-length message, and those two are the properties the API's contract
/// with its clients rests on. See the module header for what changes with it and
/// what does not, and [`super::dgram_probe`] for the measurements.
///
/// The list is by exclusion rather than by inclusion so that a new unix target
/// gets Linux's arm, which is the one that fails loudly: a kernel without
/// `AF_UNIX` `SOCK_SEQPACKET` refuses the `socketpair` outright rather than
/// returning something that silently loses messages.
#[cfg(not(any(target_os = "macos", target_os = "freebsd")))]
const SOCK_TYPE: libc::c_int = libc::SOCK_SEQPACKET;
#[cfg(target_os = "macos")]
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
const SOCK_TYPE: libc::c_int = libc::SOCK_DGRAM;
/// Create a connected socket pair of this platform's [`SOCK_TYPE`].
@@ -142,7 +161,7 @@ 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
/// How long a receive waits on the reactor before retrying the syscall, on a
/// platform whose reactor cannot see a closed `SOCK_DGRAM` peer.
///
/// This is a detection latency, not a poll interval for data: an arriving
@@ -153,7 +172,7 @@ pub fn set_sndbuf(fd: &OwnedFd, bytes: usize) -> io::Result<()> {
/// 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")]
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
const CLOSE_RETRY: Duration = Duration::from_millis(250);
/// The daemon's half of a flow's or a listener's socket pair, registered with
@@ -215,14 +234,14 @@ impl Seqpacket {
// 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")]
#[cfg(any(target_os = "macos", target_os = "freebsd"))]
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"))]
#[cfg(not(any(target_os = "macos", target_os = "freebsd")))]
let mut guard = self.inner.readable().await?;
let attempt = guard.try_io(|inner| recv_once(inner.get_ref().as_raw_fd(), buf));
match attempt {