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 }