From 8b396b662db5a8f60acb6a1b08736855fc390dbb Mon Sep 17 00:00:00 2001 From: Johnathan Corgan Date: Thu, 1 Oct 2026 22:40:40 +0000 Subject: [PATCH] Keep the daemon's copy of a native API descriptor until the client holds it A native API flow's descriptor reaches the client inside a message: an arrival for a flow a listener accepts, or a connect or listen reply. The daemon closed its own copy once the message was written, so until the client read it, the message was the only reference to the socket. xnu's descriptor collector flushes a socket in that state, and the client then receives a flow that reads as end of file with the datagrams the daemon held for it gone. That is the intermittent macOS failure of the two listener tests. A listener now keeps the daemon's copy of each flow it hands over until the client's first write on the flow, the listener's close, or the flow's end, on every platform the listener builds on. The flow is recorded before its reader starts, so a write already queued cannot race the record. The connection's serving loop, now a method on the connection so a test can run it over a real socket, keeps the copy sent in a connect or listen reply until the client's next command on that connection or the connection's end of file. A client sends its next command only after reading the reply, so either event means the descriptor has left the message. The shipped client closes the connection as soon as it has the reply, so it sees no change. The cost is accepted and documented: a flow a client accepts and closes without ever writing stays open, holding its port and a flow slot, until the listener closes, so a server that refuses flows by dropping them pays for each one until then; a client speaking the protocol directly that leaves its setup connection open sees a flow or listener it closes stay open until its next command or the connection's close. The reference and how-to pages, the client rustdoc on FipsStream, FipsListener and accept, and the design note say so, and the security reference records that a remote peer opening flows from many source ports to such a server can exhaust the node-wide max_flows ceiling. Tests cover an arrival surviving a provoked collection (deterministic on macOS), a held flow outliving its dropped descriptor until the listener closes, a client's write releasing it, a flow its client still holds working after the listener has closed and let its copy go, and a reply's copy kept until the next command and let go when the connection ends. The native API harness asserts the new lifetime of a refused flow. On macOS and FreeBSD the daemon notices a client's close only when a reader retries its read, up to a quarter second later, so the test helpers that wait for a close (forgotten, rebind, settle_closed) retry for up to five seconds by the clock rather than for a count of yields, and still_open waits two retry intervals there before asserting a flow is still open. --- docs/design/fips-native-api.md | 9 +- docs/how-to/use-the-native-datagram-api.md | 10 + docs/how-to/write-a-native-api-client.md | 12 +- docs/reference/native-api.md | 37 +- docs/reference/security.md | 26 +- src/native/client/mod.rs | 15 +- src/native/mod.rs | 581 +++++++++++++++++++-- src/native/seqpacket.rs | 13 + testing/native-api/client.py | 4 +- testing/native-api/test.sh | 26 +- 10 files changed, 661 insertions(+), 72 deletions(-) diff --git a/docs/design/fips-native-api.md b/docs/design/fips-native-api.md index 26ce5ea6..c16fcc55 100644 --- a/docs/design/fips-native-api.md +++ b/docs/design/fips-native-api.md @@ -133,9 +133,12 @@ and silent drops on its listeners. **Not a connection in the TCP sense.** A successful `connect` is a local registration and contacts no peer. There is no handshake, no keepalive and no notification that a peer went away. A flow ends when its descriptor closes, and -in no other way. In particular **a peer cannot end your flow: it has no close to -send.** That single fact shapes every program written against this interface, -and the consequences are drawn out in +in no other way. The one delay is an accepted flow never sent on, which ends +only once its listener has closed as well (see +[../reference/native-api.md](../reference/native-api.md#fipslistener)). In +particular **a peer cannot end your flow: it has no close to send.** That +single fact shapes every program written against this interface, and the +consequences are drawn out in [../how-to/use-the-native-datagram-api.md](../how-to/use-the-native-datagram-api.md#four-things-that-will-bite-you). ## See also diff --git a/docs/how-to/use-the-native-datagram-api.md b/docs/how-to/use-the-native-datagram-api.md index 355ba093..b7a19493 100644 --- a/docs/how-to/use-the-native-datagram-api.md +++ b/docs/how-to/use-the-native-datagram-api.md @@ -227,6 +227,16 @@ leaves the flows already accepted from it untouched. A program that parks streams in a `Vec` and never removes them holds ports and flow slots exactly as if it had leaked descriptors. +**An accepted flow you drop without ever sending on stays open until you +drop its listener.** Until your program has sent on an accepted flow, the +daemon keeps its own copy of the flow's descriptor, because on macOS the +kernel can otherwise destroy the flow while its descriptor is still on the +way to you. Your first `send` on the flow, or dropping the listener, lets +that copy go. So a server that refuses flows by dropping them unanswered +holds a port and a flow slot for each one until its listener goes, and a +long-lived listener that refuses many flows can walk the node into its flow +ceiling. + **Nothing peer-driven ever ends a flow, so your program has to.** The v1 wire carries no half-close. Nothing closes the daemon's half of a live accepted flow, so a loop written as "echo until the flow closes", or one diff --git a/docs/how-to/write-a-native-api-client.md b/docs/how-to/write-a-native-api-client.md index a128874a..a45c96ad 100644 --- a/docs/how-to/write-a-native-api-client.md +++ b/docs/how-to/write-a-native-api-client.md @@ -124,7 +124,7 @@ one producer on the surface, which is what lets a caller read it. ## Step 6: Keep descriptor hygiene -Five rules. Each one leaks a flow or loses one when broken. +Six rules. Each one leaks a flow or loses one when broken. **Request close-on-exec** with `MSG_CMSG_CLOEXEC` on the `recvmsg`, rather than setting it afterwards. Without it the descriptor survives an `exec` into a @@ -140,7 +140,15 @@ descriptors rather than dropping them on the floor. **Lift the descriptor out of an arrival you cannot parse** before discarding the message. Refusing a flow is closing its descriptor; discarding the message -without taking it leaks the flow instead. +without taking it leaks the flow instead. A refused flow ends only when the +listener closes, though, unless you wrote on it first: the daemon keeps its +own copy of an accepted flow's descriptor until your first write or the +listener's close. + +**Close the setup connection once you have the reply.** The daemon keeps its +own copy of the descriptor in its last reply until your next command on that +connection or the connection's close. A flow or listener you close while the +connection sits idle stays open until one of those happens. **Bound the partial line.** A daemon that stopped sending newlines would otherwise grow your buffer without end. The shipped client caps it at 64 KiB, diff --git a/docs/reference/native-api.md b/docs/reference/native-api.md index cea846d8..a6a21501 100644 --- a/docs/reference/native-api.md +++ b/docs/reference/native-api.md @@ -195,7 +195,10 @@ That is what lets the blanket reference implementation cover `&str` and One datagram flow, and the descriptor it rides on. A flow is an exact match of both ends and both ports. **The descriptor is the flow**: it lives while a -process holds that descriptor and ends when the last one closes it. +process holds that descriptor and ends when the last one closes it. The one +exception is an accepted flow that has never been sent on: the daemon keeps +its own copy of that descriptor until the first `send` or until the listener +is dropped. See `accept` below. `Send + Sync + 'static`, with no `Arc` and no borrow. There is **no `Clone` and no `try_clone`**. Because `send` and `recv` both take `&self`, a shared borrow @@ -340,6 +343,15 @@ no other way to refuse one. An unparseable arrival is therefore reported only after the descriptor it carried has been taken into ownership, so a parse failure refuses the flow rather than leaking it. +**A refused flow stays open until its listener is dropped**, unless the +program sent on it first. Until the first `send` on an accepted flow, the +daemon keeps its own copy of the flow's descriptor, because on macOS the +kernel can otherwise destroy a socket whose descriptor is still in an unread +arrival. That copy goes at the first `send`, when the listener is dropped, or +when the flow ends any other way, and the flow then ends with the program's +own close. Until then a dropped flow holds its port and its slot against +`max_flows`. + **`incoming()`** returns an `Incoming<'_>`, which borrows the listener for the iterator's lifetime, so the listener cannot be moved or dropped mid-iteration. @@ -360,7 +372,8 @@ none either. A bounded accept is `set_nonblocking` plus a wait of the caller's own on the descriptor. **Dropping** closes the descriptor and unbinds the port. Flows already accepted -from it are untouched; flows still pending on it go with it. +from it and still held are untouched; flows still pending on it go with it, and +so do flows accepted from it and dropped without ever being sent on. ### Incoming @@ -613,6 +626,20 @@ own half non-blocking and leaves the client's half blocking. `SOCK_SEQPACKET` is what preserves message boundaries in both directions, which is why the payload needs no framing. +**The daemon keeps a copy of the client's half after sending it.** While a +descriptor sits unread in a message, the message can be its only reference, +and the macOS kernel's descriptor collector destroys a socket in that state: +the client then receives a flow that reads as end of file with its datagrams +gone. So the daemon keeps its copy until one of these: + +- For a `connect` or `listen` reply, at the client's next command on the same + connection, or when that connection closes. +- For an arrival on a listener, at the client's first write on the flow, when + the listener closes, or when the flow ends any other way. + +A flow or listener the client closes before then ends when the daemon's copy +goes, not at the client's close. + A refused `connect` leaves the port free: the socket pair is built before the port is claimed, so a failure to build it needs no rollback. @@ -626,7 +653,11 @@ returning. **The connection owns nothing.** Closing it releases no flow and no listener, and a descriptor kept across the close keeps working. What owns the flow is the -descriptor. +descriptor. The connection does delay one thing: a descriptor from its last +reply that the client closes while the connection is still open, with no +further command sent, stays open until the next command or the connection's +close (see Passing the descriptor). The shipped client closes the connection +as soon as it has the reply, so it never meets this. ## Command reference diff --git a/docs/reference/security.md b/docs/reference/security.md index b69b5ea7..b140f893 100644 --- a/docs/reference/security.md +++ b/docs/reference/security.md @@ -245,12 +245,14 @@ machine. **The file descriptor carries the grant, not the connection.** A setup call hands the client a socket descriptor and the connection it was made on is then -closed; the flow or the held port lives until that descriptor is closed. A -descriptor is an ordinary kernel object, so it survives `fork`, survives -`exec` unless the client asked for it close-on-exec when it received it, and -can be handed to another process over `SCM_RIGHTS`. A process holding one can -send as this node on that flow, or receive on that port, without ever opening -the API socket and without being in the `fips` group. +closed; the flow or the held port lives until that descriptor is closed and +the daemon has let go of the copy it keeps while the descriptor is being +handed over. A descriptor is an ordinary kernel object, so it survives +`fork`, survives `exec` unless the client asked for it close-on-exec when it +received it, and can be handed to another process over `SCM_RIGHTS`. A +process holding one can send as this node on that flow, or receive on that +port, without ever opening the API socket and without being in the `fips` +group. Nothing revokes a descriptor already handed out. Restarting the daemon closes its own halves and ends every flow and listener at once, and that is the only revocation there is. @@ -269,6 +271,18 @@ peer had sent it, reaching any listener on this node under any peer identity the caller names. Leave it off outside a test harness; a packaged node does not enable it. +**A remote peer can fill the node's flow ceiling through a server that +refuses flows by dropping them.** Until a program first sends on a flow it +accepted, the daemon keeps its own copy of that flow's descriptor, so a flow +accepted and dropped unanswered keeps its slot against the node-wide +`node.native_api.max_flows` until its listener is dropped. A peer that opens +flows to such a listener from many source ports can therefore exhaust the +ceiling, and every other program on the node then gets `EMFILE` on `connect` +and silently loses arrivals on its listeners. This is the accepted cost of +keeping a flow alive while its descriptor is on the way to the program; see +[../how-to/use-the-native-datagram-api.md](../how-to/use-the-native-datagram-api.md) +for what releases the daemon's copy. + The socket is local only. It is not reachable over the network, and nothing about it changes the mesh's own authentication: a peer still verifies the node's signature, which is precisely why a local caller that can send through diff --git a/src/native/client/mod.rs b/src/native/client/mod.rs index ee4c8ade..f02ec7f1 100644 --- a/src/native/client/mod.rs +++ b/src/native/client/mod.rs @@ -373,7 +373,10 @@ fn expired(error: io::Error) -> io::Error { /// /// The protocol has no close command. Dropping the stream closes its /// descriptor, and that is what releases the flow and its local port at the -/// daemon. +/// daemon. An accepted flow that was never sent on is the exception: the daemon +/// keeps its own copy of its descriptor until the first [`FipsStream::send`] or +/// until the listener is dropped, so dropping the stream before either releases +/// nothing yet. #[derive(Debug)] pub struct FipsStream { fd: OwnedFd, @@ -627,8 +630,9 @@ impl AsFd for FipsStream { /// **The listener is a descriptor**, which is what makes it pollable: it joins /// an existing `poll`, `select` or `epoll` loop with no new mechanism, and /// [`FipsListener::accept`] is one `recvmsg` on it. Dropping the listener closes -/// that descriptor, which unbinds the port; flows already accepted from it are -/// untouched. +/// that descriptor, which unbinds the port; flows already accepted from it and +/// still held are untouched, and those accepted and dropped without ever being +/// sent on end with it. #[derive(Debug)] pub struct FipsListener { fd: OwnedFd, @@ -676,7 +680,10 @@ impl FipsListener { /// Refusing a flow is dropping the stream, which closes its descriptor. /// There is no other way to refuse one, which is why an unreadable arrival /// message is reported after the descriptor it carried has been taken: the - /// flow is then refused rather than leaked. + /// flow is then refused rather than leaked. A refused flow that was never + /// sent on still holds its port and its slot against the node's flow limit + /// until this listener is dropped, because the daemon keeps its own copy of + /// the descriptor until then. pub fn accept(&self) -> io::Result<(FipsStream, FipsAddr)> { let mut buf = [0u8; CHUNK]; let chunk = fdpass::recv(self.fd.as_raw_fd(), &mut buf)?; diff --git a/src/native/mod.rs b/src/native/mod.rs index 57a772f4..305c5760 100644 --- a/src/native/mod.rs +++ b/src/native/mod.rs @@ -18,7 +18,20 @@ //! replies and then has no further part in anything it opened. A flow lives //! until its own descriptor reaches end of file and a listener until its own //! does, whichever task holds them, which is what makes a descriptor this API -//! hands back behave like one a syscall would have. +//! hands back behave like one a syscall would have. The one exception: the +//! connection keeps a copy of the descriptor in its last reply until the +//! client's next command or its close, for the same reason a listener keeps a +//! copy of a flow it hands over (below). The shipped client closes the +//! connection as soon as it has the reply, so for it the copy is gone at once. +//! +//! **A listener keeps a copy of a flow it handed over.** While a flow's +//! descriptor sits unread in an arrival message, the message can be the only +//! reference to it, and xnu's descriptor collector flushes a socket in that +//! state. So `hand_over` keeps the daemon's copy until the client has written +//! on the flow or closed the listener, and a flow then reaches end of file once +//! both the client and that copy have let go. The cost is that a client which +//! accepts a flow and closes it without writing leaves it open until the +//! listener closes. See `dgram_probe.rs` for the measurement. //! //! **The wire is connected.** A datagram a client writes leaves this node over //! FSP, and one arriving on a held port reaches its flow. `max_payload` is the @@ -233,6 +246,15 @@ mod unix_impl { /// writes into. sock: Arc, counts: Arc, + /// The daemon's copy of the client's half, kept while that half may + /// still be in flight to a listener's client. + /// + /// Holding it keeps the socket reachable from outside the message that + /// carries it, which is what stops xnu's collector flushing it before + /// the client reads the arrival. `None` once released, and always for + /// a connected flow, whose descriptor went back in an RPC reply and is + /// held by the connection that sent it instead; see `Connection::run`. + pin: Option, } /// What `stats` reports about one flow. @@ -283,6 +305,31 @@ mod unix_impl { self.table().remove(&id); } + /// Let go of the daemon's copy of a flow's client half. + /// + /// The copy is dropped after the lock is released: if the client has + /// already closed its own, this is the last reference, and the close it + /// causes is the flow's end of file. + fn unpin(&self, id: u64) { + let pin = self.table().get_mut(&id).and_then(|flow| flow.pin.take()); + drop(pin); + } + + /// Let go of the daemon's copy of every flow on a local port. + /// + /// Called when a listener ends. Every pinned flow on its port came from + /// it: a connected flow carries no pin, and no later listener can hold + /// the port until this one's release has been served. + fn unpin_port(&self, port: u16) { + let pins: Vec = self + .table() + .values_mut() + .filter(|flow| flow.local == port) + .filter_map(|flow| flow.pin.take()) + .collect(); + drop(pins); + } + /// What the debug `stats` command reports, or `None` for a flow this /// node does not hold. fn stats(&self, id: u64) -> Option { @@ -503,6 +550,7 @@ mod unix_impl { key, peer, wiring, + None, &self.outbound, &self.node, &self.flows, @@ -700,6 +748,19 @@ mod unix_impl { #[cfg(test)] impl Connection { + /// Another connection to the same node, sharing its flow table, as a + /// second client of one daemon has. + pub(super) fn sibling(&self) -> Self { + Self::new( + self.node.clone(), + self.outbound.clone(), + Arc::clone(&self.flows), + self.limits, + Arc::clone(&self.npub), + self.debug, + ) + } + /// Build a connection wired to a channel a test serves. pub(super) fn for_test( node: mpsc::Sender, @@ -730,14 +791,24 @@ mod unix_impl { } } - /// Wait until a flow's reader task has observed the client's close. + /// Wait until a flow's reader task has observed the client's close, + /// failing by name if it never does. + /// + /// A flow the node no longer holds counts as closed: the reader sets + /// the flag and forgets the flow in the same turn, so the flag alone is + /// almost never there to be seen. Bounded by the clock rather than by + /// turns of the runtime, for the reason given at `CLOSE_WAIT`. pub(super) async fn settle_closed(&self, flow: u64) { - for _ in 0..1000 { - match self.flows.stats(flow) { - Some(stats) if stats.closed => return, - _ => tokio::task::yield_now().await, - } - } + let seen = super::tests::eventually(async || match self.flows.stats(flow) { + Some(stats) if !stats.closed => None, + _ => Some(()), + }) + .await; + assert!( + seen.is_some(), + "the reader never saw flow {flow} close within {:?}", + super::tests::CLOSE_WAIT + ); } } @@ -780,35 +851,48 @@ mod unix_impl { /// The reader cannot start any earlier: it stamps every datagram it /// forwards with the flow's key and the peer's address, and the local port /// is not known until the registry has answered. + /// + /// The flow is recorded before the reader is spawned. A handed-over client + /// can already have written, and on a multi-thread runtime the reader can + /// run at once, so recording afterwards would let its unpin find no flow + /// and leave the pin in place for the flow's whole life. + #[allow( + clippy::too_many_arguments, + reason = "one hand-off of plumbing to two tasks, with no state to group" + )] fn start( id: u64, key: FlowKey, peer: XOnlyPublicKey, wiring: Wiring, + pin: Option, outbound: &mpsc::Sender, node: &mpsc::Sender, flows: &Arc, ) { let counts = Arc::new(Counts::default()); + let pinned = pin.is_some(); + flows.record( + id, + Flow { + local: key.local, + sock: Arc::clone(&wiring.sock), + counts: Arc::clone(&counts), + pin, + }, + ); tokio::spawn(drain( id, key, peer, Arc::clone(&wiring.sock), - Arc::clone(&counts), + counts, + pinned, outbound.clone(), node.clone(), Arc::clone(flows), )); - tokio::spawn(feed(Arc::clone(&wiring.sock), wiring.inbound)); - flows.record( - id, - Flow { - local: key.local, - sock: wiring.sock, - counts, - }, - ); + tokio::spawn(feed(wiring.sock, wiring.inbound)); } /// Serve one listener until its client closes the descriptor. @@ -872,6 +956,15 @@ mod unix_impl { } debug!(port, "Native API listener closed by its client"); + // Before the release, so the port cannot have passed to another + // listener whose flows this would also let go of. When the client has + // closed the listener, an arrival it never read went with the + // listener's receive queue, so no descriptor from this port is still + // in flight to it. The other two ways out of the loop, the node going + // away and a failed read, can leave the client's half open with + // arrivals unread on it. Both are teardown, and the hold goes with the + // listener rather than outliving it. + flows.unpin_port(port); let _ = node .send(NativeMessage::Release { flows: Vec::new(), @@ -891,6 +984,15 @@ mod unix_impl { /// Both writes are try-sends. They go onto a socket pair whose client half /// has not been sent yet, so no process can read either one and a task that /// parked on one would stop serving this listener entirely. + /// + /// The daemon keeps its copy of the client's half after the arrival is + /// written, until the client writes on the flow or closes the listener. + /// Until the client reads the arrival, the message carrying the descriptor + /// is otherwise its only reference, and xnu's collector flushes a socket in + /// that state: the client then receives a flow that reads as end of file, + /// with its held datagrams gone. Nothing in the protocol says when the + /// client has read the arrival, so a client that closes the flow without + /// writing leaves it open until the listener closes. async fn hand_over( arrival: Arrival, listener: &Seqpacket, @@ -968,13 +1070,11 @@ mod unix_impl { accepted.key, accepted.peer, wiring, + Some(theirs), outbound, node, flows, ); - // Dropping our copy leaves the client holding the only reference to its - // half, so its close tears the flow down. - drop(theirs); } /// What a failed hand-off write says about the client, for the counter. @@ -1098,6 +1198,11 @@ mod unix_impl { /// Counting continues alongside the forwarding: `stats` is how a check /// observes that a datagram reached the daemon, independently of whether it /// then reached a peer. + /// + /// `pinned` says the daemon still holds a copy of the client's half. The + /// first datagram the client writes proves it holds the descriptor, so the + /// copy is let go then, once, and the client's close ends the flow from + /// there on. #[allow( clippy::too_many_arguments, reason = "one hand-off of plumbing to a task, with no state to group" @@ -1108,6 +1213,7 @@ mod unix_impl { peer: XOnlyPublicKey, sock: Arc, counts: Arc, + mut pinned: bool, outbound: mpsc::Sender, node: mpsc::Sender, flows: Arc, @@ -1116,6 +1222,10 @@ mod unix_impl { loop { match sock.recv(&mut buf).await { Ok(Received::Datagram(len)) => { + if pinned { + flows.unpin(id); + pinned = false; + } counts.datagrams.fetch_add(1, Ordering::Relaxed); counts.bytes.fetch_add(len as u64, Ordering::Relaxed); let sent = outbound @@ -1210,21 +1320,42 @@ mod unix_impl { npub: Arc, debug: bool, ) -> Result<(), std::io::Error> { - let mut connection = Connection::new(node, outbound, flows, limits, npub, debug); - let mut reader = BufReader::new(stream); - let mut line = Vec::new(); + Connection::new(node, outbound, flows, limits, npub, debug) + .run(stream) + .await + } - while read_command(&mut reader, &mut line).await? { - let (response, fd) = connection.answer(&line).await; - let mut json = serde_json::to_vec(&response)?; - json.push(b'\n'); - fdpass::reply(reader.get_ref(), &json, fd.as_ref().map(AsFd::as_fd)).await?; - // Dropping our copy leaves the client holding the only reference to - // its half, so its close tears the flow or the listener down. - drop(fd); + impl Connection { + /// Answer the commands on one client connection, in order, until it + /// closes or misbehaves. + /// + /// **The daemon keeps its copy of the descriptor in the last reply** + /// until the client's next command arrives or the connection ends. + /// Until the client reads the reply, the message carrying the + /// descriptor can be its only reference, and xnu's collector flushes a + /// socket in that state, as it does an unread arrival's. A client that + /// sends another command has read the reply first, unless it pipelined, + /// which the shipped client never does. Once the copy is gone the + /// client holds the only reference, so its close tears the flow or the + /// listener down; until then a close it makes waits for the copy. + pub(super) async fn run(mut self, stream: UnixStream) -> Result<(), std::io::Error> { + let mut reader = BufReader::new(stream); + let mut line = Vec::new(); + let mut kept: Option = None; + + while read_command(&mut reader, &mut line).await? { + drop(kept.take()); + let (response, fd) = self.answer(&line).await; + let mut json = serde_json::to_vec(&response)?; + json.push(b'\n'); + fdpass::reply(reader.get_ref(), &json, fd.as_ref().map(AsFd::as_fd)).await?; + kept = fd; + } + + // End of file, or an early return above: either way `kept` goes + // with this frame, and with it the last copy the daemon holds. + Ok(()) } - - Ok(()) } /// Read one newline-terminated command into `line`, refusing an oversized @@ -1365,6 +1496,47 @@ mod tests { sock } + /// How long a test waits for the daemon to notice that a client closed a + /// descriptor. + /// + /// On Linux the close wakes the task reading the daemon's half, so a wait + /// that will succeed does so within a few turns of the runtime. On macOS + /// and FreeBSD nothing wakes it: the reader sees the close only when its + /// bounded wait expires and it retries the read, as much as one + /// `CLOSE_RETRY` interval (`seqpacket.rs`) after the close, and a + /// listener's close that ends a flow's hold costs two of those in a row. + /// Counting turns of the runtime bounds nothing there: a thousand of them + /// ran out well inside one interval on a macOS runner. A wait that is + /// going to succeed still returns as soon as it does; this is how long one + /// that is not takes to say so. + pub(super) const CLOSE_WAIT: std::time::Duration = std::time::Duration::from_secs(5); + + // Four times the longest run of unnoticed closes a test waits through, so + // that lengthening `CLOSE_RETRY` cannot quietly use up the margin. + const _: () = assert!( + CLOSE_WAIT.as_millis() >= 8 * super::seqpacket::CLOSE_LATENCY.as_millis(), + "CLOSE_WAIT no longer covers two unnoticed closes four times over" + ); + + /// Retry `attempt` until it produces a value, for up to [`CLOSE_WAIT`]. + /// + /// `None` means the bound ran out. Attempts are a millisecond apart rather + /// than a yield apart, so a wait that lasts a quarter second on macOS is + /// not spent spinning, which for [`rebind`] would mean opening and closing + /// a socket pair on every turn. + pub(super) async fn eventually(mut attempt: impl AsyncFnMut() -> Option) -> Option { + tokio::time::timeout(CLOSE_WAIT, async { + loop { + if let Some(value) = attempt().await { + return value; + } + tokio::time::sleep(std::time::Duration::from_millis(1)).await; + } + }) + .await + .ok() + } + /// 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; @@ -1726,6 +1898,332 @@ mod tests { assert_eq!(&buf[..3], &[0x00, 0xff, 0x10]); } + /// Wait until a listener descriptor has an arrival to read. + /// + /// Polled without blocking and yielding in between, because the arrival is + /// written by the listener's task on this same runtime and a blocking wait + /// would stop the task it is waiting for. A readable listener means + /// `hand_over` has finished: it writes the arrival and settles what happens + /// to the daemon's copy of the descriptor in one synchronous step. + async fn readable(listener: &StdUnixStream) { + for _ in 0..1000 { + let mut poll = libc::pollfd { + fd: listener.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) }; + if rc > 0 && (poll.revents & libc::POLLIN) != 0 { + return; + } + tokio::task::yield_now().await; + } + panic!("no arrival became readable on the listener"); + } + + /// Wait until the node no longer holds a flow, failing if it never lets go. + /// + /// The reader marks a flow closed before it gives the registry entry back + /// and forgets the flow, so `stats` can still find a closed flow for a few + /// turns of the runtime, and on macOS and FreeBSD the reader notices the + /// close itself only when it next retries. Bounded by [`CLOSE_WAIT`], so a + /// flow that is never released names itself. + async fn forgotten(connection: &mut Connection, flow: u64) { + let line = format!(r#"{{"command":"stats","params":{{"flow_id":{flow}}}}}"#); + let mut last = serde_json::Value::Null; + let gone = eventually(async || { + last = ask(connection, &line).await; + (last["data"]["errno"] == "ENOENT").then_some(()) + }) + .await; + assert!( + gone.is_some(), + "the node still holds flow {flow} after {CLOSE_WAIT:?}: {last}" + ); + } + + /// Bind a listener on a port a closed one held, retrying until the closed + /// one has given it back. + /// + /// The release reaches the registry from the closed listener's own task + /// once that task has noticed the close, which takes a few turns of the + /// runtime, or on macOS and FreeBSD up to a retry interval. Bounded by + /// [`CLOSE_WAIT`], so a port that is never given back fails here by name. + async fn rebind(connection: &mut Connection, port: u16) -> OwnedFd { + let line = format!(r#"{{"command":"listen","params":{{"local_port":{port}}}}}"#); + eventually(async || { + let (response, fd) = connection.answer(line.as_bytes()).await; + (serde_json::to_value(response).unwrap()["status"] == "ok") + .then(|| fd.expect("a rebound listener still gets a descriptor")) + }) + .await + .unwrap_or_else(|| { + panic!("the closed listener never gave port {port} back within {CLOSE_WAIT:?}") + }) + } + + /// Bind a listener on 4242, deliver one datagram to it from a new peer, + /// and accept the flow that announced, returning everything a test needs. + async fn arrive_and_accept( + connection: &mut Connection, + ) -> (StdUnixStream, u64, serde_json::Value, StdUnixStream) { + let (_value, listener) = listen(connection, 4242).await; + let value = ask(connection, &arrival(5000, 4242, "00ff10")).await; + assert_eq!(value["data"]["outcome"], "announced", "{value}"); + let flow = value["data"]["flow_id"].as_u64().unwrap(); + let (message, client) = accept(&listener); + (listener, flow, message, client) + } + + #[tokio::test] + async fn an_arrival_survives_a_kernel_collection_before_its_client_reads_it() { + // The macOS failure, made deterministic. Between the daemon writing the + // arrival and the client reading it, the flow's descriptor exists only + // inside the arrival message. xnu's descriptor collector flushes a + // socket in that state, so unless the daemon still holds its own copy, + // the client receives a flow that reads as end of file with its held + // datagram gone. On Linux the kernel keeps the socket either way, so + // this test can only fail on macOS. + let (mut connection, _outbound) = connect(); + let (_value, listener) = listen(&mut connection, 4242).await; + let value = ask(&mut connection, &arrival(5000, 4242, "00ff10")).await; + assert_eq!(value["data"]["outcome"], "announced", "{value}"); + + readable(&listener).await; + super::dgram_probe::provoke_collection(); + + let (message, mut client) = accept(&listener); + assert_eq!(message["held"], 1); + let mut buf = [0u8; 64]; + assert_eq!( + client.read(&mut buf).unwrap(), + 3, + "the accepted flow lost its held datagram to the kernel's collector" + ); + assert_eq!(&buf[..3], &[0x00, 0xff, 0x10]); + } + + #[tokio::test] + async fn a_handed_over_flow_outlives_its_dropped_descriptor_until_its_listener_closes() { + // The daemon keeps its copy of a handed-over descriptor until the + // client has shown it holds one, because dropping it is what exposes + // the descriptor to the macOS collector. On Linux the hold is + // observable this way: a client that closes the flow without ever + // writing does not end it while the listener that produced it is open. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, client) = arrive_and_accept(&mut connection).await; + + drop(client); + still_open( + &mut connection, + flow, + "the flow ended while the daemon should still hold its descriptor", + ) + .await; + + // Closing the listener ends the hold: whatever the client did with the + // arrival, the descriptor is no longer in flight in a live socket. + drop(listener); + forgotten(&mut connection, flow).await; + } + + #[tokio::test] + async fn a_handed_over_flow_closes_with_its_descriptor_once_its_client_has_written() { + // A datagram from the client proves it holds the descriptor, so the + // daemon lets its copy go and the client's close ends the flow at once, + // with the listener still open. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, mut client) = arrive_and_accept(&mut connection).await; + + client.write_all(b"x").unwrap(); + connection.settle(flow, 1).await; + drop(client); + forgotten(&mut connection, flow).await; + drop(listener); + } + + #[tokio::test] + async fn a_handed_over_flow_its_client_still_holds_outlives_its_listener_and_closes_with_its_descriptor() + { + // Closing the listener lets the daemon's copy go, and the client's own + // copy then carries the flow by itself: it keeps working after the + // listener has gone, and the client's close ends it with no write ever + // made. Letting the copy go must not end a flow the client still holds. + let (mut connection, _outbound) = connect(); + let (listener, flow, _message, mut client) = arrive_and_accept(&mut connection).await; + + drop(listener); + // The port comes back only after the listener's task has let its + // flows' copies go, so a rebound port means that has happened. + let _rebound = rebind(&mut connection, 4242).await; + + let mut buf = [0u8; 64]; + assert_eq!(client.read(&mut buf).unwrap(), 3); + let value = ask( + &mut connection, + &format!(r#"{{"command":"inject","params":{{"flow_id":{flow},"data":"ab"}}}}"#), + ) + .await; + assert_eq!(value["status"], "ok", "{value}"); + assert_eq!(client.read(&mut buf).unwrap(), 1); + assert_eq!(buf[0], 0xab); + + drop(client); + forgotten(&mut connection, flow).await; + } + + /// Serve a sibling of `connection` over a real socket, the way the daemon + /// serves a client, returning the client's end and the serving task. + /// + /// Through `run` rather than `answer`, because what these tests observe is + /// what the serving loop does with the descriptor in a reply it has sent. + fn serve_socket( + connection: &Connection, + ) -> (StdUnixStream, tokio::task::JoinHandle>) { + let (ours, theirs) = StdUnixStream::pair().expect("AF_UNIX socketpair"); + ours.set_nonblocking(true).expect("the socket is open"); + let ours = tokio::net::UnixStream::from_std(ours).expect("inside a runtime"); + let task = tokio::spawn(connection.sibling().run(ours)); + (bounded(theirs), task) + } + + /// Write one command on a client socket and read its reply line, with the + /// descriptor it carried. + async fn call(client: &StdUnixStream, line: &str) -> (serde_json::Value, Option) { + let mut writer = client; + writer.write_all(line.as_bytes()).unwrap(); + writer.write_all(b"\n").unwrap(); + readable(client).await; + let mut buf = [0u8; 4096]; + let chunk = super::fdpass::recv(client.as_raw_fd(), &mut buf) + .expect("a reply should be readable on the connection"); + let reply = buf[..chunk.len] + .strip_suffix(b"\n") + .expect("one whole reply line per read"); + (serde_json::from_slice(reply).unwrap(), chunk.fd) + } + + /// Open a flow through a served socket, returning its id and descriptor. + async fn connect_over(client: &StdUnixStream) -> (u64, OwnedFd) { + let line = format!( + r#"{{"command":"connect","params":{{"peer":"{PEER}","remote_port":4242,"local_port":4243}}}}"# + ); + let (value, fd) = call(client, &line).await; + assert_eq!(value["status"], "ok", "{value}"); + let flow = value["data"]["flow_id"].as_u64().unwrap(); + ( + flow, + fd.expect("a connect reply carries the flow's descriptor"), + ) + } + + /// Assert that the node still holds a flow open, after giving its reader + /// every chance to notice a close. + /// + /// Where a close wakes the reader, a few turns of the runtime are that + /// chance. On macOS and FreeBSD the reader notices a close only when its + /// bounded wait expires and it retries, so a check made sooner passes + /// whether or not the flow has closed. There this also waits out two of + /// those intervals: a reader that parked just before a close has retried + /// by then, with one interval to spare for a loaded runner. Elsewhere the + /// interval is zero and the wait costs nothing. + async fn still_open(connection: &mut Connection, flow: u64, why: &str) { + for _ in 0..1000 { + tokio::task::yield_now().await; + } + tokio::time::sleep(super::seqpacket::CLOSE_LATENCY * 2).await; + let value = ask( + connection, + &format!(r#"{{"command":"stats","params":{{"flow_id":{flow}}}}}"#), + ) + .await; + assert_eq!(value["status"], "ok", "{why}: {value}"); + assert_eq!(value["data"]["closed"], false, "{why}: {value}"); + } + + /// Any command at all, sent only so the serving loop reads one. + const NEXT: &str = r#"{"command":"stats","params":{"flow_id":0}}"#; + + #[tokio::test] + async fn a_descriptor_sent_in_a_reply_is_kept_until_the_clients_next_command() { + // Until the client reads a reply, the message carrying its descriptor + // can be the only reference to it, which is what xnu's collector + // flushes. So the daemon keeps its copy until the client sends another + // command, which it does only after reading the reply. On Linux the + // hold is observable as a flow that outlives the client's close. + let (mut probe, _outbound) = connect(); + let (client, _task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(fd); + still_open( + &mut probe, + flow, + "the flow ended while the connection should still hold its descriptor", + ) + .await; + + call(&client, NEXT).await; + forgotten(&mut probe, flow).await; + } + + #[tokio::test] + async fn a_descriptor_sent_in_a_reply_is_let_go_when_the_connection_ends() { + // The connection's close ends its hold. The client's own copy then + // carries the flow alone, so the flow survives the connection and ends + // with the client's close. + let (mut probe, _outbound) = connect(); + let (client, task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(client); + task.await + .unwrap() + .expect("a connection closed between commands ends cleanly"); + still_open( + &mut probe, + flow, + "the flow ended with the connection while the client still holds it", + ) + .await; + + drop(fd); + forgotten(&mut probe, flow).await; + } + + #[tokio::test] + async fn a_flow_its_client_closes_after_its_next_command_ends_at_once() { + // Once the client has sent another command the hold is over, so the + // client's close is the flow's end of file with the connection still + // open. + let (mut probe, _outbound) = connect(); + let (client, _task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + call(&client, NEXT).await; + drop(fd); + forgotten(&mut probe, flow).await; + drop(client); + } + + #[tokio::test] + async fn a_flow_whose_reply_is_never_followed_by_a_command_ends_with_the_connection() { + // A client that closes the flow and then the connection, sending + // nothing more, still ends the flow: the connection's end of file is + // the last point at which the daemon lets its copy go. + let (mut probe, _outbound) = connect(); + let (client, task) = serve_socket(&probe); + let (flow, fd) = connect_over(&client).await; + + drop(fd); + drop(client); + task.await + .unwrap() + .expect("a connection closed between commands ends cleanly"); + forgotten(&mut probe, flow).await; + } + #[test] fn a_listener_that_closed_and_one_that_stopped_reading_are_counted_apart() { // Both take the same cleanup, so the counter is the only place the @@ -1946,19 +2444,8 @@ mod tests { // The unbind is what `close(listen_fd)` means in Berkeley, and the // daemon can only observe it by reading its own half. Without the read - // arm the port is held for the node's lifetime and every attempt below - // fails. - for _ in 0..1000 { - let (response, fd) = connection - .answer(br#"{"command":"listen","params":{"local_port":4242}}"#) - .await; - if serde_json::to_value(response).unwrap()["status"] == "ok" { - assert!(fd.is_some(), "a rebound listener still gets a descriptor"); - return; - } - tokio::task::yield_now().await; - } - panic!("the closed listener never gave its port back"); + // arm the port is held for the node's lifetime and every attempt fails. + rebind(&mut connection, 4242).await; } #[tokio::test] diff --git a/src/native/seqpacket.rs b/src/native/seqpacket.rs index d2537a27..a9bd60b8 100644 --- a/src/native/seqpacket.rs +++ b/src/native/seqpacket.rs @@ -180,6 +180,19 @@ pub fn set_sndbuf(fd: &OwnedFd, bytes: usize) -> io::Result<()> { #[cfg(any(target_os = "macos", target_os = "freebsd"))] const CLOSE_RETRY: Duration = Duration::from_millis(250); +/// How long a reader can leave a closed peer unnoticed, for the native API +/// tests that assert a flow is still open and so must wait long enough to have +/// seen it close. +/// +/// `CLOSE_RETRY` where the reactor cannot see a close, because the reader then +/// notices one only when that bound expires and it retries the read. +#[cfg(all(test, any(target_os = "macos", target_os = "freebsd")))] +pub(super) const CLOSE_LATENCY: Duration = CLOSE_RETRY; + +/// Zero where a close wakes the reader itself. +#[cfg(all(test, not(any(target_os = "macos", target_os = "freebsd"))))] +pub(super) const CLOSE_LATENCY: Duration = Duration::ZERO; + /// The daemon's half of a flow's or a listener's socket pair, registered with /// the reactor. pub struct Seqpacket { diff --git a/testing/native-api/client.py b/testing/native-api/client.py index f5f41af8..221a0e0c 100755 --- a/testing/native-api/client.py +++ b/testing/native-api/client.py @@ -7,7 +7,9 @@ connection owns nothing — a flow lives until its own descriptor is closed, and listener until its own is — so the single connection is a convenience for the checks rather than a lifetime the daemon respects. Descriptors are what keep things alive, and this tool holds them until the step that closes them or until -it exits. +it exits. The daemon does keep its own copy of the descriptor in its last reply +until the next command arrives, so a check that closes one must send a command +before it expects the close to have taken effect. Kinds of step: diff --git a/testing/native-api/test.sh b/testing/native-api/test.sh index 2badf346..7470324c 100755 --- a/testing/native-api/test.sh +++ b/testing/native-api/test.sh @@ -695,12 +695,20 @@ check_backlog_is_not_the_clients_bound() { } check_refusing_a_flow_frees_it() { - log "Refusing an accepted flow is closing its descriptor, and that frees it" + log "Refusing an accepted flow is closing its descriptor, and its listener's close frees it" # There is no reject command: Berkeley has exactly one way to refuse a # connection and so does this. Keeping a second would let a client refuse a # flow two indistinguishable ways. # - # Two things have to follow the close, and the second is what makes this + # The daemon keeps its own copy of an accepted flow's descriptor until the + # client writes on the flow or closes the listener, because on macOS the + # kernel can destroy a socket whose descriptor is still in an unread + # arrival. So a flow refused without a write outlives its close, and goes + # when the listener does. Both halves are asserted: a flow freed at the + # close would mean the daemon let its copy go early, and one that outlived + # the listener would be a leak. + # + # Two things have to follow the release, and the second is what makes this # more than a repeat of the connected-flow close: the node forgets the flow, # and the registry entry and its key go with it, so the very same key # announces a new flow afterwards rather than delivering into the dead one. @@ -713,16 +721,22 @@ check_refusing_a_flow_frees_it() { {"command":"stats","params":{"flow_id":"@a"},"expect":{"status":"ok"}}, {"fd":"a","close":true}, {"sleep":1}, - {"command":"stats","params":{"flow_id":"@a"},"expect":{"status":"error"}}, + {"command":"stats","params":{"flow_id":"@a"}, + "expect":{"status":"ok","data.closed":false}}, + {"fd":"L","close":true}, + {"command":"stats","params":{"flow_id":"@a"},"settle":true, + "expect":{"status":"error"}}, + {"command":"listen","params":{"local_port":4304},"keep_listener":"M", + "expect":{"status":"ok"}}, {"command":"arrive","params":{"peer":"'"$PEER"'","src_port":5000,"dst_port":4304,"data":"bb"}, "expect":{"status":"ok","data.outcome":"announced"}}, - {"accept":"L","keep_fd":"b","expect":{"local_port":4304,"remote_port":5000}}, + {"accept":"M","keep_fd":"b","expect":{"local_port":4304,"remote_port":5000}}, {"fd":"b","read":1,"expect_bytes":"bb"} ]' if run_client "$script"; then - pass "a refused flow is gone and its key is free to arrive again" + pass "a refused flow lasts until its listener closes, then is gone and its key is free to arrive again" else - fail "closing a refused flow did not release it" + fail "a refused flow did not last until its listener closed, or was not released then" fi }