Repository navigation
A pipe watch can ask whether the other end is gone, and netstack decides a client has left by it - #760
Conversation
This reverts commit 8dd55d0, which took the test out red: the guest binary, its RUST_SKIP row, its registration, its dispatch arm and body, and the `netstack` connector tests/netcase's runner holds for it. The next commit is what turns it green. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
A piped TCP stream was three things: the bridge in `piped_connections`, the smoltcp socket with its two 64 KiB buffers, and the id's entry in `sockets`. Only a client's close request removed all three. A connection that ended any other way, its client dead or simply dropping its pipe ends, lost its bridge in `bridge_piped` and kept the other two for the life of the process, outside the count `max_piped_connections` bounds. `bridge_piped` already ends a connection on what the kernel says of the client's pipe ends (a write refused `Gone`, a read at end of file) and on what the peer says on the wire. Where it lets the bridge go it now removes the socket and the entry too, so the close request is no longer what a stream's lifetime rests on. The socket stays one pass longer where netstack aborted it: `abort` leaves the socket `Closed` with a reset owed, and a socket removed before the next poll sends nothing. Before this change that reset left only because the socket was never removed. `spent` is that condition, and `TimeWait` counts as spent, as it did for the bridge. A request that names the id of a stream netstack has let go is answered as an unknown id is: a shutdown or an option `ERR_NOT_CONNECTED`, a close `done`. Measured on tests/netcase under QEMU: `netstack_socket_churn`, red at 06bbc23 on `netstack held 0 stream(s) before them and holds 4 after`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
Datagram sockets are released only by their close request, listeners only on a pass something else causes; neither is bounded in number; and logkeeper hears a reader leave only at its next write. All three are read from the code. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
The stream count reads the table alone, so a connection whose entry left and whose socket stayed in the set passed it. `net.sockets.untabled` is the count that moves then. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
|
The negative control and the two mutations, as applied at
diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs
index c816d8ae4..0e611b67c 100644
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -331,13 +331,7 @@ struct PendingUdpRecv {
}
/// A piped TCP connection: data flows through kernel pipes instead of IPC messages.
-///
-/// **The socket and its id live exactly as long as this does.** What ends a
-/// connection is what the kernel says of the client's pipe ends and what the
-/// peer says on the wire; a close request only asks for that end early, and a
-/// client that dies sends none.
struct PipedConnection {
- socket_id: u32,
handle: SocketHandle,
rx_write: Option<Pipe>,
tx_read: Option<Pipe>,
@@ -425,15 +419,8 @@ fn send_room(socket: &tcp::Socket) -> bool {
socket.can_send() && socket.send_capacity() > socket.send_queue()
}
-/// Whether `socket` has said its last to its peer: it is closed, and the reset
-/// an abort owes has left. `TimeWait` only waits.
-fn spent(socket: &tcp::Socket) -> bool {
- !socket.is_open() && !(socket.state() == tcp::State::Closed && socket.remote_endpoint().is_some())
-}
-
-fn piped_connection(socket_id: u32, handle: SocketHandle, pipes: DataPipes) -> PipedConnection {
+fn piped_connection(handle: SocketHandle, pipes: DataPipes) -> PipedConnection {
PipedConnection {
- socket_id,
handle,
rx_write: Some(pipes.to_client),
tx_read: Some(pipes.from_client),
@@ -1259,7 +1246,7 @@ impl Netstack {
let stream_id = self.alloc_id();
self.sockets.insert(stream_id, SocketKind::TcpStream(old_handle));
- self.piped_connections.push(piped_connection(stream_id, old_handle, pipes));
+ self.piped_connections.push(piped_connection(old_handle, pipes));
// Create replacement listener
let rx_buf = tcp::SocketBuffer::new(vec![0u8; TCP_SOCKET_BUFFER]);
@@ -1392,15 +1379,14 @@ impl Netstack {
}
}
- if conn.is_fully_closed() && spent(socket) {
+ // Fully clean up when both sides are done
+ if conn.is_fully_closed() && !socket.is_open() {
closed.push(i);
}
}
for &i in closed.iter().rev() {
- let conn = self.piped_connections.swap_remove(i);
- socket_set.remove(conn.handle);
- self.sockets.remove(&conn.socket_id);
+ self.piped_connections.swap_remove(i);
}
}
@@ -1488,7 +1474,7 @@ impl Netstack {
};
pc.client.result(&resp);
let pc = self.pending_piped_connects.swap_remove(i);
- self.piped_connections.push(piped_connection(pc.socket_id, pc.handle, pc.pipes));
+ self.piped_connections.push(piped_connection(pc.handle, pc.pipes));
continue;
}
if socket.state() == tcp::State::Closed {
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -1403,7 +1403,6 @@
for &i in closed.iter().rev() {
let conn = self.piped_connections.swap_remove(i);
- socket_set.remove(conn.handle);
self.sockets.remove(&conn.socket_id);
}
}
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -428,7 +428,7 @@
/// Whether `socket` has said its last to its peer: it is closed, and the reset
/// an abort owes has left. `TimeWait` only waits.
fn spent(socket: &tcp::Socket) -> bool {
- !socket.is_open() && !(socket.state() == tcp::State::Closed && socket.remote_endpoint().is_some())
+ !socket.is_open()
}
fn piped_connection(socket_id: u32, handle: SocketHandle, pipes: DataPipes) -> PipedConnection { |
|
Review round 1 of Net: 8 files, +243 −6. Production BLOCKER
NOTE
Rulings the brief asked for
SEND BACK |
A bridge whose pipe ends are both gone waited on the wire with no bound, and held its slot in piped_live meanwhile. Measured on listen/tests.rs's bench before this change, a second at a time for an hour, spent() never true: an aborted socket whose next hop answers no ARP (no segment leaves), a peer silent after the handshake (FIN-WAIT-1, 372 retransmissions), and a peer that acknowledges the FIN and never sends its own (FIN-WAIT-2). Such a connection is now `ownerless`: the wire finishing it is the event, and OWNERLESS_LIFE the ceiling, counted from the pass that found both ends gone. At the ceiling an open socket is reset and kept for the one poll that sends the reset; one already aborted whose reset never left is let go. Both say so. The number is RFC 9293 3.8.3's R2 at the 100 seconds it asks for at least. It is counted from the client's leaving and not from the peer's last word because smoltcp's own socket timeout, armed at the same moment, was measured not to end a FIN-WAIT-2 whose peer acknowledges once a second; RFC 9293 3.8.6.1 leaves a connection held open to the system's resource management. A connection with either pipe end left is not ownerless and is never cut. The loop's wait is bounded by the first ceiling, which nothing on the wire wakes a pass for. The option requests on a stream its peer reset are filed: the write is refused as macOS refuses it (measured), the read where hosts answer. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
…close Its exit is met: a connection netstack lets go leaves no socket and no table entry, whoever ended it, and netstack_socket_churn is green; the negative control is the release reverted, red on the issue's recorded line (4 streams after, 0 before). The rule it carried is PipedConnection's doc in userland/netstack/src/main.rs. No other file cites the slug or the name. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
|
Round 1's measurements, control and mutation, by
diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs
index a25aa3f6f..56b3152c2 100644
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -464,16 +464,8 @@ enum Ownerless {
/// The pass's answer for `socket`, whose client's pipe ends have both been
/// gone for `gone`.
fn ownerless(socket: &mut tcp::Socket, gone: Duration) -> Ownerless {
- if spent(socket) {
- Ownerless::Over
- } else if gone < OWNERLESS_LIFE {
- Ownerless::Waits
- } else if socket.is_open() {
- socket.abort();
- Ownerless::Cut
- } else {
- Ownerless::Unsaid
- }
+ let _ = gone;
+ if spent(socket) { Ownerless::Over } else { Ownerless::Waits }
}
fn piped_connection(socket_id: u32, handle: SocketHandle, pipes: DataPipes) -> PipedConnection {
diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs
index a25aa3f6f..de0102893 100644
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -469,7 +469,6 @@ fn ownerless(socket: &mut tcp::Socket, gone: Duration) -> Ownerless {
} else if gone < OWNERLESS_LIFE {
Ownerless::Waits
} else if socket.is_open() {
- socket.abort();
Ownerless::Cut
} else {
Ownerless::Unsaid
diff --git a/userland/netstack/src/listen/tests.rs b/userland/netstack/src/listen/tests.rs
index 79575e8f8..07508c523 100644
--- a/userland/netstack/src/listen/tests.rs
+++ b/userland/netstack/src/listen/tests.rs
@@ -84,6 +84,10 @@ struct Net {
listener: SocketHandle,
listening: Listening,
sent: Vec<Sent>,
+ /// The clock every poll reads.
+ now: Instant,
+ /// Whether the far end answers ARP.
+ arp: bool,
}
impl Net {
@@ -95,7 +99,16 @@ impl Net {
socket.listen(PORT).expect("a fresh socket listens");
let mut sockets = SocketSet::new(Vec::new());
let listener = sockets.add(socket);
- Self { iface, wire, sockets, listener, listening: Listening::new(PORT), sent: Vec::new() }
+ Self {
+ iface,
+ wire,
+ sockets,
+ listener,
+ listening: Listening::new(PORT),
+ sent: Vec::new(),
+ now: Instant::from_millis(0),
+ arp: true,
+ }
}
fn socket(&mut self) -> &mut tcp::Socket<'static> {
@@ -103,10 +116,10 @@ impl Net {
}
/// One pass of netstack's loop: everything the wire holds, in one batch, and
- /// the far end's ARP answered.
+ /// the far end's ARP answered while it answers any.
fn pass(&mut self) {
loop {
- while self.iface.poll(Instant::from_millis(0), &mut self.wire, &mut self.sockets) != PollResult::None {}
+ while self.iface.poll(self.now, &mut self.wire, &mut self.sockets) != PollResult::None {}
if self.wire.outbound.is_empty() {
return;
}
@@ -122,6 +135,9 @@ impl Net {
EthernetProtocol::Arp => {
let arp = ArpRepr::parse(&ArpPacket::new_checked(eth.payload()).unwrap()).unwrap();
let ArpRepr::EthernetIpv4 { operation: ArpOperation::Request, .. } = arp else { return };
+ if !self.arp {
+ return;
+ }
let reply = ArpRepr::EthernetIpv4 {
operation: ArpOperation::Reply,
source_hardware_addr: PEER_MAC,
@@ -347,3 +363,58 @@ fn a_connection_closed_by_both_ends_is_spent() {
assert_eq!(net.socket().state(), tcp::State::TimeWait, "the premise: the peer's FIN answered ours");
assert!(crate::spent(net.socket()), "a connection both ends closed was kept");
}
+
+/// The handshake from `from` finished, and the socket the client's.
+fn established(net: &mut Net, from: u16) -> u32 {
+ let isn = net.syn(from);
+ net.send(from, TcpControl::None, PEER_ISN + 1, Some(isn + 1));
+ net.pass();
+ assert_eq!(net.socket().state(), tcp::State::Established);
+ isn
+}
+
+/// MEASUREMENT: whether `spent` ever turns true, a second at a time for an hour.
+fn spent_within_an_hour(net: &mut Net) -> Option<u64> {
+ for second in 0..3600u64 {
+ net.pass();
+ if crate::spent(net.socket()) {
+ return Some(second);
+ }
+ net.now += smoltcp::time::Duration::from_secs(1);
+ }
+ None
+}
+
+#[test]
+fn measure_an_aborted_connection_whose_next_hop_answers_no_arp() {
+ let mut net = Net::new();
+ established(&mut net, 5001);
+ net.now += smoltcp::time::Duration::from_secs(61);
+ net.arp = false;
+ net.socket().abort();
+ let before = net.sent.len();
+ let at = spent_within_an_hour(&mut net);
+ assert!(at.is_some(), "the slot never comes back; {} segment(s) left in the hour", net.sent.len() - before);
+}
+
+#[test]
+fn measure_a_peer_that_goes_silent_after_the_handshake() {
+ let mut net = Net::new();
+ established(&mut net, 5001);
+ net.socket().close();
+ let before = net.sent.len();
+ let at = spent_within_an_hour(&mut net);
+ assert!(at.is_some(), "the slot never comes back; state {}, {} segment(s) left in the hour", net.socket().state(), net.sent.len() - before);
+}
+
+#[test]
+fn measure_a_peer_that_acknowledges_the_fin_and_goes_silent() {
+ let mut net = Net::new();
+ let isn = established(&mut net, 5001);
+ net.socket().close();
+ net.pass();
+ net.send(5001, TcpControl::None, PEER_ISN + 1, Some(isn + 2));
+ let before = net.sent.len();
+ let at = spent_within_an_hour(&mut net);
+ assert!(at.is_some(), "the slot never comes back; state {}, {} segment(s) left in the hour", net.socket().state(), net.sent.len() - before);
+}
--- /Users/jan/.claude/jobs/2280e09e/tmp/scratchpad/orch/netstackclose-r1/tests.rs.keep 2026-10-08 10:53:35
+++ userland/netstack/src/listen/tests.rs 2026-10-08 10:53:35
@@ -440,3 +440,20 @@
assert_eq!(at_the_ceiling(&mut net, answer), Ownerless::Cut, "the slot of a connection with no client never comes back");
the_cut_is_said(&mut net);
}
+
+/// MEASUREMENT, not for the tree: smoltcp's own socket timeout, armed where
+/// the client leaves, against the peer that keeps answering.
+#[test]
+fn measure_smoltcp_timeout_against_a_peer_that_answers() {
+ let mut net = Net::new();
+ let isn = established(&mut net, 5001);
+ net.socket().close();
+ net.socket().set_timeout(Some(smoltcp::time::Duration::from_secs(100)));
+ for second in 0..3600u64 {
+ net.send(5001, TcpControl::None, PEER_ISN + 1, Some(isn + 2));
+ net.pass();
+ assert!(net.socket().is_open(), "smoltcp's timeout closed it after {second}s");
+ net.now += smoltcp::time::Duration::from_secs(1);
+ }
+ panic!("after an hour the socket is still {}", net.socket().state());
+}
// What a host answers set_nodelay and nodelay on a stream its peer reset.
use std::io::Read;
use std::net::{TcpListener, TcpStream};
fn main() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let mut client = TcpStream::connect(listener.local_addr().unwrap()).unwrap();
let (peer, _) = listener.accept().unwrap();
// Linger 0 then close: the peer's reset.
socket_linger0(&peer);
drop(peer);
let mut buf = [0u8; 1];
println!("read after the peer's reset: {:?}", client.read(&mut buf));
println!("set_nodelay(true): {:?}", client.set_nodelay(true));
println!("nodelay(): {:?}", client.nodelay());
println!("shutdown(Write): {:?}", client.shutdown(std::net::Shutdown::Write));
}
fn socket_linger0(s: &TcpStream) {
use std::os::fd::AsRawFd;
#[repr(C)] struct Linger { on: i32, secs: i32 }
unsafe extern "C" { fn setsockopt(fd: i32, level: i32, name: i32, val: *const Linger, len: u32) -> i32; }
let (level, name) = if cfg!(target_os = "macos") { (0xffff, 0x0080) } else { (1, 13) };
let l = Linger { on: 1, secs: 0 };
assert_eq!(unsafe { setsockopt(s.as_raw_fd(), level, name, &l, 8) }, 0);
} |
|
Review round 2 of Net: 10 files, +447 −51. Production Round 1's BLOCKERs
BLOCKER
NOTE
Rulings the brief asked for
Round 1's other NOTEs: the closed issue is deleted and Evidence at SEND BACK |
Round 2: both BLOCKERs measured red at
|
…nerless rests on it Review round 2 of #760 measured two defects in the 100-second bound on a connection with no client. The first: netstack decided a client was gone by reading its send pipe to the end of file, and it reads that pipe only while the socket takes bytes. A client that shut its sending half down and exited, or filled the send buffer against a silent peer and exited, was never read again, never ownerless, and held its slot for the life of netstack. The kernel had no word for it: a watch on a pipe end knew READABLE (bytes, or no writer and an empty ring) and WRITABLE (room, or no reader), so "the writer is gone" could be waited for only on an empty pipe. The owner ruled the event in. toyos_abi::inbox::OTHER_END_GONE is a third OP_WATCH condition, of a pipe end only: ready when the pipe's other end has no holder left, whatever the ring holds, level-triggered like its siblings. The kernel already kept both counts and already posted the end's own watch at that transition; the condition is one more question the look asks, registered on that same watch. An object that is no pipe end is refused the question. The counts and the three answers moved to kernel/pure/pipe.rs, where the host runs their cases. toyos::poller::Poller::wait_answers hands a caller what its watch was answered beside the token. netstack watches both of a connection's pipes for it for as long as it holds them, keyed by socket id: an answer names its connection, and a connection let go since is no place in a list another has taken. A receive pipe with no reader is closed; a send pipe with no writer is still drained while the socket takes bytes. The zero-byte write that asked the receive pipe on every pass is gone. The second: at the ceiling the reset was given one poll, and a connection silent for 100 seconds has outlived its next hop's 60-second neighbour entry, so that poll sent an ARP request and the socket was removed with the reset unsent. A cut connection is now kept until the reset has left or RESET_LIFE, three seconds, is over: RFC 4861's budget for address resolution, three solicitations a second apart. The log says which. Filed: the absolute ceiling cuts a departed client's unsent tail, and a shutdown of the sending half drops what the send pipe still holds. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
# Conflicts: # kernel/pure/lib.rs
|
Round 3 patches, all against
diff --git a/kernel/src/pipe.rs b/kernel/src/pipe.rs
index de4190997..013246c7d 100644
--- a/kernel/src/pipe.rs
+++ b/kernel/src/pipe.rs
@@ -271,7 +271,7 @@ pub fn has_space(pipe_id: PipeId) -> bool {
/// Whether the end opposite `end` of this pipe has no holder left.
pub fn other_end_gone(pipe_id: PipeId, end: End) -> bool {
- with_pipes(|pipes| pipes.get(pipe_id).is_some_and(|p| p.holders.other_end_gone(end)))
+ with_pipes(|pipes| pipes.get(pipe_id).is_some_and(|p| p.holders.other_end_gone(end) && false))
}
/// Mark the pipe so the next consumer inherits RT priority.
diff --git a/kernel/pure/lib.rs b/kernel/pure/lib.rs
index f3a288bec..86f6f211f 100644
--- a/kernel/pure/lib.rs
+++ b/kernel/pure/lib.rs
@@ -1,8 +1,7 @@
//! What the kernel decides without touching the machine: the scheduler core
//! ([`sched`]), the process and thread lifecycle ([`proclife`]), which PCID an
-//! address space is handed ([`pcid`]), what type the range registers give a
-//! range ([`mtrr`]) and what a pipe's ends are told of each other ([`pipe`]).
-//! The kernel binary links it; the host runs
+//! address space is handed ([`pcid`]) and what type the range registers give a
+//! range ([`mtrr`]). The kernel binary links it; the host runs
//! its tests, because none of it reads a register, a clock or a kernel lock.
#![no_std]
@@ -14,6 +13,5 @@ extern crate std;
pub mod mtrr;
pub mod pcid;
-pub mod pipe;
pub mod proclife;
pub mod sched;
diff --git a/kernel/pure/pipe.rs b/kernel/pure/pipe.rs
deleted file mode 100644
index 76b2b39ac..000000000
--- a/kernel/pure/pipe.rs
+++ /dev/null
@@ -1,166 +0,0 @@
-//! Who holds a pipe's two ends, and what a watch on one of them is told.
-//!
-//! A pipe end is held by counted references, one per object that names it: a
-//! handle duplicated or moved to another process is the same object or another
-//! reference, so an end is gone only when the last of them is. Every answer
-//! here is a function of the two counts and of what the ring holds, read under
-//! the kernel's pipe lock.
-
-#![forbid(unsafe_code)]
-
-/// One end of a pipe.
-#[derive(Clone, Copy, PartialEq, Eq, Debug)]
-pub enum End {
- Read,
- Write,
-}
-
-impl End {
- pub fn other(self) -> Self {
- match self {
- Self::Read => Self::Write,
- Self::Write => Self::Read,
- }
- }
-}
-
-/// What a reference's release left of its pipe.
-#[derive(Clone, Copy, PartialEq, Eq, Debug)]
-#[must_use = "a release that left an end unheld owes the other end a post"]
-pub enum Released {
- /// Its end has another holder.
- Held,
- /// It was its end's last holder: the other end's watch is owed a post.
- EndGone,
- /// It was the pipe's last holder.
- PipeGone,
-}
-
-/// How many references hold each end of one pipe.
-#[derive(Default)]
-pub struct Holders {
- readers: u32,
- writers: u32,
-}
-
-impl Holders {
- pub const fn new() -> Self {
- Self { readers: 0, writers: 0 }
- }
-
- fn count(&mut self, end: End) -> &mut u32 {
- match end {
- End::Read => &mut self.readers,
- End::Write => &mut self.writers,
- }
- }
-
- pub fn hold(&mut self, end: End) {
- let count = self.count(end);
- *count = count.checked_add(1).expect("pipe holder overflow");
- }
-
- pub fn release(&mut self, end: End) -> Released {
- let count = self.count(end);
- *count = count.checked_sub(1).expect("pipe holder underflow");
- match (self.readers, self.writers) {
- (0, 0) => Released::PipeGone,
- _ if !self.held(end) => Released::EndGone,
- _ => Released::Held,
- }
- }
-
- pub fn held(&self, end: End) -> bool {
- match end {
- End::Read => self.readers > 0,
- End::Write => self.writers > 0,
- }
- }
-
- /// A read answers now: with bytes, or with the end of a stream no writer
- /// can add to.
- pub fn readable(&self, available: u32) -> bool {
- available > 0 || !self.held(End::Write)
- }
-
- /// A write answers now: it takes bytes, or is refused for want of a reader.
- pub fn writable(&self, space: u32) -> bool {
- space > 0 || !self.held(End::Read)
- }
-
- /// The end opposite `end` has no holder, whatever the ring holds.
- pub fn other_end_gone(&self, end: End) -> bool {
- !self.held(end.other())
- }
-}
-
-#[cfg(test)]
-mod tests {
- use super::*;
-
- /// A pipe as `pipe::create` leaves it: one holder of each end.
- fn made() -> Holders {
- let mut pipe = Holders::new();
- pipe.hold(End::Read);
- pipe.hold(End::Write);
- pipe
- }
-
- #[test]
- fn an_end_whose_other_end_is_held_is_not_told_it_is_gone() {
- let pipe = made();
- assert!(!pipe.other_end_gone(End::Read));
- assert!(!pipe.other_end_gone(End::Write));
- }
-
- /// **Whatever the ring holds.** `readable` answers for the bytes and says
- /// nothing of the writer; `writable` for the room and nothing of the reader.
- #[test]
- fn an_end_dropped_with_the_ring_neither_empty_nor_full_is_gone() {
- for (available, space) in [(0, 8), (3, 5), (8, 0)] {
- let mut pipe = made();
- assert_eq!(pipe.release(End::Write), Released::EndGone);
- assert!(pipe.other_end_gone(End::Read), "the writer left with {available} byte(s) in the ring");
- assert!(pipe.readable(available));
- assert!(!pipe.other_end_gone(End::Write), "a reader that is held was told gone");
-
- let mut pipe = made();
- assert_eq!(pipe.release(End::Read), Released::EndGone);
- assert!(pipe.other_end_gone(End::Write), "the reader left with room for {space} byte(s)");
- assert!(pipe.writable(space));
- assert!(!pipe.other_end_gone(End::Read), "a writer that is held was told gone");
- }
- }
-
- /// The two older conditions do not stand in for it: each is ready with the
- /// other end held.
- #[test]
- fn readable_and_writable_say_nothing_of_the_other_end() {
- let pipe = made();
- assert!(pipe.readable(1) && !pipe.other_end_gone(End::Read));
- assert!(pipe.writable(1) && !pipe.other_end_gone(End::Write));
- assert!(!pipe.readable(0) && !pipe.writable(0));
- }
-
- /// A handle duplicated, or moved to another process, is a second
- /// reference: the end is gone with the last of them and not before.
- #[test]
- fn an_end_with_a_second_holder_outlives_the_first() {
- for end in [End::Read, End::Write] {
- let mut pipe = made();
- pipe.hold(end);
- assert_eq!(pipe.release(end), Released::Held);
- assert!(!pipe.other_end_gone(end.other()), "{end:?} was told gone with a holder left");
- assert_eq!(pipe.release(end), Released::EndGone);
- assert!(pipe.other_end_gone(end.other()));
- }
- }
-
- /// The pipe's last reference owes nobody a post: no end is left to watch.
- #[test]
- fn the_last_holder_of_the_pipe_ends_it() {
- let mut pipe = made();
- assert_eq!(pipe.release(End::Read), Released::EndGone);
- assert_eq!(pipe.release(End::Write), Released::PipeGone);
- }
-}
diff --git a/kernel/src/inbox/mod.rs b/kernel/src/inbox/mod.rs
index 32310a66a..dc3851599 100644
--- a/kernel/src/inbox/mod.rs
+++ b/kernel/src/inbox/mod.rs
@@ -4,7 +4,7 @@
//! into both kernel and userspace. `OP_WATCH` fires once; userspace re-submits to re-arm.
//!
//! **A watch is a poll registered on the watched object's own
-//! [`Watch`](crate::watch::Watch)**, one entry per watch its interest names, and
+//! [`Watch`](crate::watch::Watch)**, one entry per direction it asked for, and
//! the object's post fires it. There is no table of sources here: what a
//! handle watches is `ops::read_watch`/`ops::write_watch`'s answer, and the
//! poll holds no reference to the object at all.
@@ -109,11 +109,10 @@ pub struct WatchFlags(u32);
impl WatchFlags {
pub const READABLE: Self = Self(toyos_abi::inbox::READABLE);
pub const WRITABLE: Self = Self(toyos_abi::inbox::WRITABLE);
- pub const OTHER_END_GONE: Self = Self(toyos_abi::inbox::OTHER_END_GONE);
/// Every bit `toyos_abi::inbox` defines for `Submission::op_flags`;
/// hand-copied and unchecked, for the reason `syscall/vm.rs`'s
/// `MMAP_PROT_KNOWN` gives for all four of these masks.
- const KNOWN: u32 = Self::READABLE.0 | Self::WRITABLE.0 | Self::OTHER_END_GONE.0;
+ const KNOWN: u32 = Self::READABLE.0 | Self::WRITABLE.0;
/// A bit outside [`Self::KNOWN`] is an interest this kernel would register
/// for neither direction, so the whole watch is refused rather than served.
@@ -125,7 +124,6 @@ impl WatchFlags {
}
pub fn readable(self) -> bool { self.0 & Self::READABLE.0 != 0 }
pub fn writable(self) -> bool { self.0 & Self::WRITABLE.0 != 0 }
- pub fn other_end_gone(self) -> bool { self.0 & Self::OTHER_END_GONE.0 != 0 }
pub fn raw(self) -> u32 { self.0 }
}
@@ -136,7 +134,6 @@ impl WatchFlags {
struct Readiness {
readable: bool,
writable: bool,
- other_end_gone: bool,
}
impl Readiness {
@@ -144,12 +141,11 @@ impl Readiness {
let mut flags = 0u32;
if self.readable { flags |= WatchFlags::READABLE.raw(); }
if self.writable { flags |= WatchFlags::WRITABLE.raw(); }
- if self.other_end_gone { flags |= WatchFlags::OTHER_END_GONE.raw(); }
flags
}
fn any(self) -> bool {
- self.readable || self.writable || self.other_end_gone
+ self.readable || self.writable
}
}
@@ -536,17 +532,8 @@ fn arm(
let ready = readiness_of(object, flags).any();
let read = if flags.readable() { ops::read_watch(object) } else { None };
let write = if flags.writable() { ops::write_watch(object) } else { None };
- // A pipe end's own watch is the one its other end's last holder posts;
- // an object with no other end is refused the question, whatever else the
- // watch asks. A direction asked above has taken that watch already.
- let own = match flags.other_end_gone().then(|| ops::pipe_end_watch(object)) {
- Some(None) => return Err(SyscallError::NotSupported),
- Some(Some(_)) if read.is_some() || write.is_some() => None,
- Some(own) => own,
- None => None,
- };
- // Nothing its interest names: nothing could ever answer this poll, so it is refused, not registered.
- if !ready && read.is_none() && write.is_none() && own.is_none() {
+ // No readiness in either direction: nothing could ever answer this poll, so it is refused, not registered.
+ if !ready && read.is_none() && write.is_none() {
return Err(SyscallError::NotSupported);
}
@@ -579,21 +566,17 @@ fn arm(
if let Some(watch) = &write {
watch.add_poll(PollEntry { poll: poll.clone(), direction: WatchFlags::WRITABLE });
}
- if let Some(watch) = &own {
- watch.add_poll(PollEntry { poll: poll.clone(), direction: WatchFlags::OTHER_END_GONE });
- }
if readiness_of(object, flags).any() {
poll.fire(0);
}
Ok(())
}
-/// What the object answers each condition, restricted to what was asked for.
+/// Per-direction readiness of the object, restricted to what was asked for.
fn readiness_of(object: &KObjectRef, flags: WatchFlags) -> Readiness {
Readiness {
readable: flags.readable() && ops::has_data(object),
writable: flags.writable() && ops::has_space(object),
- other_end_gone: flags.other_end_gone() && ops::other_end_gone(object),
}
}
diff --git a/kernel/src/object/ops.rs b/kernel/src/object/ops.rs
index c1b1e1e2d..3ed962a68 100644
--- a/kernel/src/object/ops.rs
+++ b/kernel/src/object/ops.rs
@@ -15,7 +15,6 @@ use crate::drivers::serial;
use crate::file_cache;
use crate::time::Deadline;
use crate::pipe::{self, PipeId};
-use kernel::pipe::End;
use crate::process::PipeMap;
use crate::user_ptr::{UserBytes, UserBytesMut};
use crate::inbox::PollEntry;
@@ -313,32 +312,6 @@ pub fn write_watch(object: &KObjectRef) -> Option<WatchRef> {
}
}
-/// The watch of a pipe end, which its other end's last holder posts as it
-/// lets go; `None` for an object that is no pipe end.
-pub fn pipe_end_watch(object: &KObjectRef) -> Option<WatchRef> {
- match object {
- KObjectRef::PipeRead(_) => read_watch(object),
- KObjectRef::PipeWrite(_) => write_watch(object),
- KObjectRef::Connection(_) | KObjectRef::Acceptor(_) | KObjectRef::Process(_)
- | KObjectRef::Console(_) | KObjectRef::Device(_) | KObjectRef::SysCap(_)
- | KObjectRef::File(_) | KObjectRef::Inbox(_) | KObjectRef::Connector(_)
- | KObjectRef::Namespace(_) | KObjectRef::SharedMem(_) => None,
- }
-}
-
-/// Whether this pipe end's other end has no holder left; `false` for an
-/// object that is no pipe end.
-pub fn other_end_gone(object: &KObjectRef) -> bool {
- match object {
- KObjectRef::PipeRead(r) => pipe::other_end_gone(r.id(), End::Read),
- KObjectRef::PipeWrite(w) => pipe::other_end_gone(w.id(), End::Write),
- KObjectRef::Connection(_) | KObjectRef::Acceptor(_) | KObjectRef::Process(_)
- | KObjectRef::Console(_) | KObjectRef::Device(_) | KObjectRef::SysCap(_)
- | KObjectRef::File(_) | KObjectRef::Inbox(_) | KObjectRef::Connector(_)
- | KObjectRef::Namespace(_) | KObjectRef::SharedMem(_) => false,
- }
-}
-
/// Whether closing one handle to this object ends what its watches watch, so
/// every poll on them — in any ring — is answered as gone. `false` for the log
/// and the keyboard, which the machine ends on its own and which other handles
diff --git a/kernel/src/pipe.rs b/kernel/src/pipe.rs
index de4190997..b0c8122a7 100644
--- a/kernel/src/pipe.rs
+++ b/kernel/src/pipe.rs
@@ -1,6 +1,5 @@
use crate::mm::pmm;
-use kernel::pipe::{End, Holders, Released};
use toyos_abi::ring::Ring;
use alloc::sync::Arc;
@@ -40,7 +39,7 @@ pub use handle::{PipeReader, PipeWriter};
/// take a counted reference, so no program point exists between the two where
/// the last other end could close and free it.
mod handle {
- use super::{release, with_pipes_mut, End, Pipe, PipeId};
+ use super::{close_read, close_write, with_pipes_mut, Pipe, PipeId};
/// One counted reference to a pipe's read end — `Arc`, for a reader slot.
pub struct PipeReader(PipeId);
@@ -51,7 +50,8 @@ mod handle {
impl PipeReader {
/// The only constructor: taking and counting the reference is one statement.
pub(super) fn acquire(pipe: &mut Pipe) -> Self {
- pipe.hold(End::Read);
+ pipe.readers = pipe.readers.checked_add(1).expect("pipe reader overflow");
+ pipe.publish_ends();
Self(pipe.id)
}
@@ -60,7 +60,8 @@ mod handle {
impl PipeWriter {
pub(super) fn acquire(pipe: &mut Pipe) -> Self {
- pipe.hold(End::Write);
+ pipe.writers = pipe.writers.checked_add(1).expect("pipe writer overflow");
+ pipe.publish_ends();
Self(pipe.id)
}
@@ -86,13 +87,13 @@ mod handle {
impl Drop for PipeReader {
fn drop(&mut self) {
- release(self.0, End::Read);
+ close_read(self.0);
}
}
impl Drop for PipeWriter {
fn drop(&mut self) {
- release(self.0, End::Write);
+ close_write(self.0);
}
}
}
@@ -112,7 +113,8 @@ struct Pipe {
id: PipeId,
/// `None` until first use — allocating eagerly would charge every pending `SYS_CONNECT` 4 MiB before either end sent a byte.
backing: Option<Backing>,
- holders: Holders,
+ readers: u32,
+ writers: u32,
/// The read end's and the write end's watches. Held by `Arc` so a blocking site or a
/// poll registration can clone one out from under the table lock and hold it across its own park.
read_watch: Arc<Watch>,
@@ -129,7 +131,8 @@ impl Pipe {
Self {
id,
backing: None,
- holders: Holders::new(),
+ readers: 0,
+ writers: 0,
read_watch: Arc::new(Watch::new()),
write_watch: Arc::new(Watch::new()),
rt_boost_pending: false,
@@ -151,13 +154,8 @@ impl Pipe {
/// Republish "is the other end gone?" into the mapped header, for netstack; the kernel itself decides from its own counts.
fn publish_ends(&mut self) {
let Some(backing) = self.backing.as_mut() else { return };
- if self.holders.held(End::Read) { backing.ring.open_reader() } else { backing.ring.close_reader() }
- if self.holders.held(End::Write) { backing.ring.open_writer() } else { backing.ring.close_writer() }
- }
-
- fn hold(&mut self, end: End) {
- self.holders.hold(end);
- self.publish_ends();
+ if self.readers == 0 { backing.ring.close_reader() } else { backing.ring.open_reader() }
+ if self.writers == 0 { backing.ring.close_writer() } else { backing.ring.open_writer() }
}
fn available(&self) -> u32 {
@@ -220,7 +218,7 @@ pub fn try_read(pipe_id: PipeId, buf: &mut UserBytesMut) -> Option<usize> {
let boost = pipe.rt_boost_pending;
pipe.rt_boost_pending = false;
(Some(n), boost)
- } else if !pipe.holders.held(End::Write) {
+ } else if pipe.writers == 0 {
(Some(0), false)
} else {
(None, false)
@@ -243,7 +241,7 @@ pub enum PipeWrite {
pub fn try_write(pipe_id: PipeId, buf: &UserBytes) -> Option<PipeWrite> {
with_pipes_mut(|pipes| {
let pipe = pipes.get_mut(pipe_id)?;
- if !pipe.holders.held(End::Read) {
+ if pipe.readers == 0 {
return Some(PipeWrite::BrokenPipe);
}
let Some(backing) = pipe.back() else {
@@ -259,21 +257,16 @@ pub fn try_write(pipe_id: PipeId, buf: &UserBytes) -> Option<PipeWrite> {
pub fn has_data(pipe_id: PipeId) -> bool {
with_pipes(|pipes| {
- pipes.get(pipe_id).is_some_and(|p| p.holders.readable(p.available()))
+ pipes.get(pipe_id).is_some_and(|p| p.available() > 0 || p.writers == 0)
})
}
pub fn has_space(pipe_id: PipeId) -> bool {
with_pipes(|pipes| {
- pipes.get(pipe_id).is_some_and(|p| p.holders.writable(p.space()))
+ pipes.get(pipe_id).is_some_and(|p| p.space() > 0 || p.readers == 0)
})
}
-/// Whether the end opposite `end` of this pipe has no holder left.
-pub fn other_end_gone(pipe_id: PipeId, end: End) -> bool {
- with_pipes(|pipes| pipes.get(pipe_id).is_some_and(|p| p.holders.other_end_gone(end)))
-}
-
/// Mark the pipe so the next consumer inherits RT priority.
pub fn set_rt_boost_pending(pipe_id: PipeId) {
with_pipes_mut(|pipes| {
@@ -283,22 +276,41 @@ pub fn set_rt_boost_pending(pipe_id: PipeId) {
});
}
-/// One reference to `end` lets go; its end's last posts the other end's watch.
-fn release(pipe_id: PipeId, end: End) {
- let released = with_pipes_mut(|pipes| {
- let pipe = pipes.get_mut(pipe_id).expect("release: pipe not found");
- let released = pipe.holders.release(end);
+fn close_read(pipe_id: PipeId) {
+ // `true` when the pipe still lives and its write end is now the one to wake.
+ let wake_writers = with_pipes_mut(|pipes| {
+ let pipe = pipes.get_mut(pipe_id).expect("close_read: pipe not found");
+ pipe.readers = pipe.readers.checked_sub(1).expect("pipe reader underflow");
+ pipe.publish_ends();
+ if pipe.readers == 0 && pipe.writers == 0 {
+ let pipe = pipes.remove(pipe_id).unwrap();
+ free_pipe(pipe);
+ false // pipe freed, no one to wake
+ } else {
+ pipe.readers == 0
+ }
+ });
+ if wake_writers {
+ crate::scheduler::wake_pipe_writers(pipe_id);
+ }
+}
+
+fn close_write(pipe_id: PipeId) {
+ // `true` when the pipe still lives and its read end is now the one to wake.
+ let wake_readers = with_pipes_mut(|pipes| {
+ let pipe = pipes.get_mut(pipe_id).expect("close_write: pipe not found");
+ pipe.writers = pipe.writers.checked_sub(1).expect("pipe writer underflow");
pipe.publish_ends();
- if released == Released::PipeGone {
+ if pipe.readers == 0 && pipe.writers == 0 {
let pipe = pipes.remove(pipe_id).unwrap();
free_pipe(pipe);
+ false // pipe freed, no one to wake
+ } else {
+ pipe.writers == 0
}
- released
});
- match (released, end) {
- (Released::Held | Released::PipeGone, _) => {}
- (Released::EndGone, End::Read) => crate::scheduler::wake_pipe_writers(pipe_id),
- (Released::EndGone, End::Write) => crate::scheduler::wake_pipe_readers(pipe_id),
+ if wake_readers {
+ crate::scheduler::wake_pipe_readers(pipe_id);
}
}
diff --git a/kernel/pure/pipe.rs b/kernel/pure/pipe.rs
index 76b2b39ac..065a4132b 100644
--- a/kernel/pure/pipe.rs
+++ b/kernel/pure/pipe.rs
@@ -90,7 +90,7 @@ impl Holders {
/// The end opposite `end` has no holder, whatever the ring holds.
pub fn other_end_gone(&self, end: End) -> bool {
- !self.held(end.other())
+ !self.held(end)
}
}
diff --git a/kernel/src/pipe.rs b/kernel/src/pipe.rs
index de4190997..a57fc19fd 100644
--- a/kernel/src/pipe.rs
+++ b/kernel/src/pipe.rs
@@ -271,7 +271,7 @@ pub fn has_space(pipe_id: PipeId) -> bool {
/// Whether the end opposite `end` of this pipe has no holder left.
pub fn other_end_gone(pipe_id: PipeId, end: End) -> bool {
- with_pipes(|pipes| pipes.get(pipe_id).is_some_and(|p| p.holders.other_end_gone(end)))
+ with_pipes(|pipes| pipes.get(pipe_id).is_some_and(|p| p.holders.other_end_gone(end) && p.available() == 0))
}
/// Mark the pipe so the next consumer inherits RT priority.
diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs
index fd74c9de5..4b4421810 100644
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -507,7 +507,7 @@ enum Ownerless {
fn ownerless(socket: &mut tcp::Socket, waited: Duration, cut: bool) -> Ownerless {
if spent(socket) {
Ownerless::Over
- } else if waited < if cut { RESET_LIFE } else { OWNERLESS_LIFE } {
+ } else if waited < if cut { Duration::ZERO } else { OWNERLESS_LIFE } {
Ownerless::Waits
} else if cut {
Ownerless::Unsaid
diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs
index fd74c9de5..cce7bd151 100644
--- a/userland/netstack/src/main.rs
+++ b/userland/netstack/src/main.rs
@@ -1876,7 +1876,7 @@ fn main() {
// room wakes the NIC. Its writer's leaving is asked until answered,
// for the same reason: it stays so.
let room = send_room(socket_set.get::<tcp::Socket>(conn.handle));
- let send = if room { READABLE } else { 0 } | if conn.writer_gone { 0 } else { OTHER_END_GONE };
+ let send = if room { READABLE } else { 0 };
if let (true, Some(pipe)) = (send != 0, &conn.tx_read) {
poller.watch(pipe, send, TOKEN_SEND_PIPE | u64::from(conn.socket_id));
}
diff --git a/tests/toyos-rust-tests/src/bin/netstack_socket_churn.rs b/tests/toyos-rust-tests/src/bin/netstack_socket_churn.rs
index ca405615d..4233af45c 100644
--- a/tests/toyos-rust-tests/src/bin/netstack_socket_churn.rs
+++ b/tests/toyos-rust-tests/src/bin/netstack_socket_churn.rs
@@ -65,6 +65,41 @@ fn main() {
.expect("usage: netstack_socket_churn <ending host port> <holding host port>")
};
let (port, holding) = (port(1), port(2));
+ // MEASUREMENT (review round 2, BLOCKER 1 and 2): the holding server keeps
+ // every connection and reads nothing. Each arm drops both pipe ends with
+ // no close request and asks `net.piped.live` until it is back.
+ let held = count("net.piped.live");
+ for arm in ["both ends dropped", "(a) shutdown(1), then dropped", "(b) written until the pipe is full, then dropped"] {
+ let conn = toyos::net::tcp_connect(HOST, holding, 0).unwrap_or_else(|e| panic!("{arm}: connect: {e:?}"));
+ if arm.starts_with("(a)") {
+ toyos::net::tcp_shutdown(conn.socket_id, 1).unwrap_or_else(|e| panic!("{arm}: shutdown: {e:?}"));
+ }
+ if arm.starts_with("(b)") {
+ let (mut written, mut stalls) = (0usize, 0);
+ while stalls < 3 {
+ match conn.tx.write_nonblock(&[0u8; 65536]) {
+ Ok(n) if n > 0 => (written, stalls) = (written + n, 0),
+ _ => {
+ stalls += 1;
+ std::thread::sleep(Duration::from_secs(1));
+ }
+ }
+ }
+ println!("MEASURE {arm}: {written} byte(s) written before three seconds without room");
+ }
+ drop(conn);
+ let asked = Instant::now();
+ let back = loop {
+ if count("net.piped.live") == held {
+ break Some(asked.elapsed());
+ }
+ if asked.elapsed() > Duration::from_secs(130) {
+ break None;
+ }
+ };
+ println!("MEASURE {arm}: piped.live back after {back:?} (None: still held 130s after the drop)");
+ }
+ assert_eq!(port, 0, "MEASUREMENT over");
let (streams, live) = (count("net.sockets.tcp"), count("net.piped.live"));
let untabled = count("net.sockets.untabled");
for round in 1..=ROUNDS {
diff --git a/tests/toyos.rs b/tests/toyos.rs
index 408feaa56..4056b4f87 100644
--- a/tests/toyos.rs
+++ b/tests/toyos.rs
@@ -2816,7 +2816,7 @@ fn netstack_socket_churn() -> Result<(), String> {
let mut console = qemu.boot_log().to_string();
await_marker(&mut qemu, &mut console, LEASED, "netstack's lease").map_err(|e| format!("{e}\n{console}"))?;
let result =
- qemu.run_test(&format!("test_rs_netstack_socket_churn {port} {holding}"), Duration::from_secs(120));
+ qemu.run_test(&format!("test_rs_netstack_socket_churn {port} {holding}"), Duration::from_secs(600));
if let Some(why) = &result.error {
return Err(format!("{why}\nthe job said:\n{}", result.stdout));
} |
|
Review round 3 of Net: 19 files, +1027 −145. Production +409 −101 ( Round 2's BLOCKERs
BLOCKER
NOTE
Rulings the brief asked for
SEND BACK |
…holders `pipe_other_end_gone` is a discovered binary on the T14's shared boot beside `poll_wake_pipe`. A thread parks on a ring holding one `OTHER_END_GONE` watch and nothing else, and the other end's last holder goes once the roster says the thread is blocked: a reader told of its writer and a writer of its reader, with the ring empty and with bytes left in it, and with a child process as the last holder. A watch made after the leaving is answered by its own submit. A connection and a file are refused the question alone and beside `READABLE`. The watch stays armed while a duplicate of the other end's handle is held, while that handle sits unreceived in a connection's queue, and while it lives in a child's table. `netstack_socket_churn` gains an arm for netstack's half: a client moves netstack an acceptor where its receive end belongs and keeps its send end, on the server that holds every connection and sends nothing. netstack never writes that end, so the kernel's refusal of its watch is all that ends the connection, and `net.piped.live` must come back. Review round 3's NOTEs: the per-pass cost of the two watches is filed (`issues/netstack-watches-both-pipes-of-every-connection-on-every-pass.md`), and three issues are made true of the tree: the two poller slots' exit now needs the condition on a connection, `clientless()` also times a client that still holds the read end of a receive pipe netstack closed, and `wait` is the form that drops a completion's result. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
|
Round 4 patches and deciding lines, at
|
|
T14 at head
Thirty rows over nine boots, none failed. No LAN row was booted: the bench's wired path to the development machine is down. |
|
Review round 4 of Net: 23 files, +1333 −168. Production +409 −101, unchanged since round 3 ( Round 3's BLOCKERs
Round 3's NOTEs
BLOCKERNone in the branch. NOTE
Rulings the brief asked for
LAND AFTER NAMED CHANGES |
The track issue conflicted in the meter paragraph, which #766 rewrote with the T14's reading and this branch had re-owned. It takes main's text whole, its two new bullets included, with the one owner line this branch changed: "the power-off stage" is this slice, whose stage paragraph the branch deletes, so the meter's owner reads "this stage". Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A
Closes
issues/netstack-keeps-the-socket-of-a-connection-whose-client-sent-no-close.md(filed by #755, deleted here). Head62ac7c817;origin/main1084ddc9a(#746, one workspace) is merged in, and it is the merge base.CI at this head, run 37770538069:
hostsuccess (job 113288745429,[ci] Host: 78 step(s), all green),toolchain / buildsuccess (113288746060),guest / suitesuccess (113291240431:PASS netstack_socket_churn (3s),test result: ok. 33 passed, 33 total,[ci] Guest: 5 step(s), all green). Neither was run locally. The T14 at this head: nine boots, thirty rows, none failed (#760 (comment)),pipe_other_end_goneamong them on the shipping kernel; the sections are listed below.What changed, and why
The ABI: a pipe watch can ask whether the other end is gone
The owner's ruling, verbatim from the option he chose: "Yes, add the event (Recommended)" — "A pipe watch gains one condition: ready when the other end has no holder left, whatever the pipe still holds. The kernel already tracks both counts and wakes watchers at that moment, so it is a small ABI addition. The network stack then knows a client is gone in every case, and the 100-second bound covers all of them. It lands in #760 with the work that needs it, through review as an ABI change."
toyos_abi::inbox::OTHER_END_GONE(16) is a thirdOP_WATCHcondition besideREADABLEandWRITABLE, interest going in and result coming back. Of a pipe end: the pipe's other end has no holder left in any process, whatever the ring holds. Level-triggered as its siblings are: a watch armed after the last holder left is answered at once (arm's first look), and one armed before is fired by the postreleasealready made on that end's own watch and answered by the submitter's look. No new object, no new syscall.NotSupported, whatever else it asks. A connection is not given the condition: no server needs it, each reads its connection to end of file.READABLEisavailable > 0 || no writerandWRITABLEisspace > 0 || no reader, so "the writer is gone" could be waited for only on an empty pipe and "the reader is gone" only on a full one.pipe_flag_forgery), and every other way to decide it is a poll, a timer, or draining a possibly-dead client's bytes into memory.kernel/pure/pipe.rsnow holds the two counts (Holders) and the three answers, pure and host-tested;kernel/src/pipe.rskeeps the ring, the lock and the posts, and its two close functions are onerelease.ops::pipe_end_watchandops::other_end_gonedispatch with no_arm.toyos::poller::Poller::wait_answershands a caller what each watch was answered beside its token: the conditions met, or the kernel's refusal.waitis that with the answer dropped.netstack: a client is gone when the kernel says so of both its ends
Round 2's BLOCKER 1: netstack decided a client was gone by reading its send pipe to end of file, and it reads that pipe only while the socket takes bytes. A client that (a) shut its sending half down and exited, or (b) wrote more than the send buffer to a silent peer and exited, was never read again, never ownerless, and held its slot for the life of netstack.
OTHER_END_GONEfor as long as netstack holds them, withREADABLEwhile the socket has send room andWRITABLEwhile the receive pipe is holding bytes back. So a bridge is never without a watch armed for its client's leaving, which also closes round 2's NOTE about a client that dies with none armed.writer_goneand is still drained while the socket takes bytes; end of file still closes the socket after the last byte, and no longer decides anything about the client.clientless(): the receive pipe closed, and the send pipe closed or its writer gone. A client holding an end of a direction still open is never timed.swap_removemoved there. An answer for a pipe netstack no longer holds is ignored.netstack_socket_churnnow stages it.net.piped.ownerlessin netstack's inspect snapshot counts connections whose client is gone.Both bounds
OWNERLESS_LIFE, 100 seconds, from the pass that first found the client gone: R2 of RFC 9293 §3.8.3 at the 100 seconds SHLD-11 asks for at least; §3.8.6.1 leaves a connection its peer holds open "subject to the implementation's resource management concerns". Absolute, by the orchestrator's ruling in round 2, not the owner's; what it costs is filed (below).RESET_LIFE, 3 seconds, from the cut (round 2's BLOCKER 2). A connection silent for 100 s has outlived its next hop's 60 s neighbour entry (smoltcp 0.12.0,iface/neighbor.rs), so the poll after the cut sends an ARP request and not the reset. The cut socket is now kept until the reset has left orRESET_LIFEis over, whichever is first. The number is address resolution's own budget: RFC 4861 §7.2.2, "If no Neighbor Advertisement is received after MAX_MULTICAST_SOLICIT solicitations, address resolution has failed", withMAX_MULTICAST_SOLICIT3 transmissions andRETRANS_TIMER1,000 milliseconds in §10. That is IPv6's; RFC 1122 §2.3.2.1 gives ARP "1 per second per destination" and no count, and smoltcp re-asks at that rate (SILENT_TIME) with no count either. Both RFC texts were fetched and read this round.netstack: the reset has left, ornetstack: letting a connection go with its reset unsent — no next hop took it in 3s.ownerless_wake_inbounds the loop's wait by whichever bound each connection is under.From rounds 0 and 1, unchanged
A piped connection's socket and its id leave with its bridge; the socket stays until it has said its last (
spent); a request naming a stream netstack has let go is answered as an unknown id is;netstack_socket_churnis back and readsnet.sockets.untabled. Their measurements and controls are atc9a2df7efandb27bb7f0a, in the earlier comments on this pull request.Measured at
e7a40645f: the ceiling through netstack's own loopNot taken again at this head: the merge and this round's commit leave
userland/netstack/and every kernel file of the change byte for byte (git diff e7a40645f 62ac7c817 -- kernel/pure/pipe.rs kernel/src/pipe.rs kernel/src/inbox kernel/src/object/ops.rs userland/netstack toyos/src toyos-abi/src/inbox.rsprints nothing).measure-guest.patchonnetstack_socket_churn(posted below; exit 1 by construction,measure-guest.log): a host server holds every connection and reads nothing; each arm drops both pipe ends with no close request and asksnet.piped.liveuntil it is back. Run alone on the machine.resetting a connection — its client left 100s ago and its peer has not finished it, thenthe reset has lefttcp_shutdown(id, 1), then droppedAt
b27bb7f0aarm (a) was still held 130 s after the drop, and the plain drop loggedletting a connection go with its reset unsentagainst a gateway that answers ARP. These are a measurement on one run, not a test: no QEMU test asserts a duration.Gates
All at
62ac7c817, run alone on the machine by the orchestrator frombuild-request.sh; logs under the scratchpad,orch/netstackclose-r4/, exits inresults.txt.cargo metadata --locked --format-version 1metadata.logcargo run -- --build-onlybuild.logcargo test -p kernel --lib --features sched-check pipe::(assrc/ci.rsruns the kernel's library, narrowed)kernel-pipe.log, 5 passedcargo test --manifest-path userland/netstack/Cargo.toml(assrc/ci.rsruns a userland crate)netstack-host.log, 27 passedcargo test --test toyos-build -- netstack_socket_churnguest-churn.logcargo test --test toyos-build -- iommu_virtio_platformguest-iommu.logcargo run -- --ci hostThe T14: read, thirty rows green
pipe_other_end_goneis a discovered binary, so its green is a T14 reading: every section below was booted by the orchestrator at this head and judged green (comment 6058892938). Staged from62ac7c817withcargo test --test toyos-build -- --metal --metal-readback <dir> <filter>, each exit 2 (staged, machine untouched);t14/request.txtlists every image with its sha256. A metal filter is one substring and selects the shared boots' members too, so the rows review round 3 named take nine images:pipesharedpipe_other_end_gone,poll_wake_pipe,abuse_pipe_map,abuse_pipe_owner,abuse_pipe_ring,pipe_flag_forgeryinboxsharedinbox_cancel_wakes,inbox_empty_write,abuse_inbox,inbox_log_poststd_sharedstd_io,std_processand the otherstd_membersevery_wait,wait_storm,blocking_read,poller,process_lifecycleshared, one eachkill_ends_every_wait,exit_wait_storm,blocking_read_stress,poller_capacity,process_lifecyclehandle_lifetimeshared-debug, the actuator kernelhandle_lifetimeChecks of high-risk code
Each a checked patch, applied, run and reversed, tree clean after each; the patches are in the comments below.
This round, at
62ac7c817(results.txt):qemu-armalone, the green armpipe_other_end_goneon a QEMU boot oftests/testcasesOTHER_END_GONEpoll on the end's own watch deletednetstack_socket_churncount()the test loops on wakes a pass, and each pass submits the watches again. The churn test holds the answer and not the wake.pipe_other_end_gone(withqemu-arm)STALLED: 123s of guard expired, and the guest had said nothing: the binary never printed its first line, which follows the four parked arms. The log does not say which thread stood still; the green arm on the same tree is what makes it the mutation.Some(None) => return Err(NotSupported)toSome(None) => Nonepipe_other_end_gone(withqemu-arm)a watch of 0x11 on a connection, which has no other end,left: None,right: Some(Err(NotSupported)): the connection was servedREADABLEand kept armed. The parked arms had passed.Err(e) => self.refuse(..)to an arm that does nothingnetstack_socket_churnnetstack still counts connection 942118025 live 20s after the kernel refused the watch of its receive end, which is no pipe endFrom round 3, at
e7a40645f, not rerun: every file those six patches touch is unchanged since (thegit diffabove), and all six stillgit apply --checkat this head. The churn test they ran against gained one arm, placed before the two that turned red. c1 (the kernel never answers; churn, 1), c2 (the kernel reverted tomain; churn, 1, red onInvalidArgumentand not on the missing wake), m1 (the pure answer asks its own end;pipe::, 101), m2 (gone only when the ring is empty; churn, 1 on the 10-byte arm), m3 (noRESET_LIFE; netstack's host tests, 101), m5 (netstack never asks its send pipe; churn, 1). Their patches and deciding lines are in the round 3 comment.What holds what. The answer's truth table is the host's (
kernel/pure/pipe.rs). The wake ispipe_other_end_gone's: a thread parked inwait_answers(1, u64::MAX, ..)on a ring holding oneOTHER_END_GONEwatch and nothing else, the other end's last holder let go once the roster says the thread is blocked, a reader told of its writer and a writer of its reader, with the ring empty and with bytes left in it, and with a child process's exit as the leaving. The same binary holds the refusal on a connection and on a file, alone and besideREADABLE; a watch made after the leaving, answered by its own submit; and the three places a false "gone" would reset a live connection: a duplicate held, the handle unreceived in a connection's queue, the handle in a living child's table. netstack's decision on a refused watch is the churn test's new arm. The reset's second bound is netstack's host bench.Independent oracles: smoltcp's own state machine and neighbour cache on a wire; RFC 9293, RFC 4861 and RFC 1122 for the numbers; a host TCP server behind QEMU's user network; and the T14, which has answered (comment 6058892938).
Why the new binary is a T14 member, and where its mutations ran
A discovered binary under
tests/toyos-rust-tests/src/bin/runs on the T14's shared boots and on no QEMU or CI run. I did not also register it inMACHINE_TESTS: the behaviour needs no machine shape, a metal row reaches it, and rootCLAUDE.mdputs a metal row before a guest test. The price is real and stays: CI does not hold the wake, and a regression of it is seen at the next T14 run and not at the pull request that makes it.w1 parks a thread for good and r1 panics, so neither was run on the T14.
qemu-arm.patchis a throwawayMACHINE_TESTSentry that bootstests/testcasesunder QEMU and runs the binary; it was applied for the green arm and under each mutation, and ships nothing.Why the changed guest test needs QEMU
netstack_socket_churnhas three arms on a host server that holds every connection. Two from round 3: a client shuts its sending half down, leaves nothing or ten bytes in its send pipe, drops both ends, andnet.piped.ownerlessmust rise by one. One new: a client moves netstack an acceptor where its receive end belongs and keeps its send end; the server sends nothing, so netstack never writes that end, the kernel's refusal of its watch is all that ends the connection, andnet.piped.livemust come back. Counts under a hang ceiling, never a duration. No type holds them. A host test cannot:pipe_answered's refusal drops kernel pipe handles and runs in netstack's loop, and neither has a host build. A metal row cannot: netstack owns its NIC and the T14's peer is the bench's network, with no server the harness controls. It is one boot the suite already makes.Filed and corrected
issues/netstack-watches-both-pipes-of-every-connection-on-every-pass.md. Two watches are submitted per live connection on every pass, and every client read or write through the kernel posts the poll netstack keeps on that pipe. Read, not measured; owner the network-stack track; exit a watch submitted only on a change of interest, or a T14 throughput reading. Not fixed here: it would changeuserland/netstack/src/main.rs, and the six controls above stand on that file as it is.issues/netstack-spends-two-poller-slots-per-piped-connection.md(both pipes are watched on every pass; its exit now needs the condition on a connection),issues/netstack-cuts-a-departed-clients-unsent-tail-at-the-ceiling.md(clientless()also times a living client that holds only the read end of a receive pipe netstack closed),issues/a-close-of-one-handle-ends-every-rings-poll-on-its-object.md(waitis the form that drops a result;wait_answershands it over).Size
Against
origin/main: 23 files, +1333 −168. Production is as review round 3 counted it (+409 −101) and unchanged this round. Tests this round:pipe_other_end_gone.rs+198,netstack_socket_churn.rs+41,tests/toyos.rs+2. Issues: one filed, three corrected.Unsure of
host, the guest suite and the T14 are read at this head, as said at the top; the merge queue runs the first two again on the merged tree.MACHINE_TESTSregistration would, at one boot per run; it is the orchestrator's or the owner's to ask for.e7a40645f.🤖 Generated with Claude Code
https://claude.ai/code/session_01RvnWQFcMuGqTHYhvSnTe8A