diff --git a/Cargo.lock b/Cargo.lock index f2824e05b5a..b014afb198d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1376,6 +1376,7 @@ name = "toyos-sched-loom" version = "0.1.0" dependencies = [ "loom", + "toyos-abi", ] [[package]] diff --git a/issues/design-debt/toyos-poller-is-sync-and-a-watch-moves-the-tail-in-two-steps.md b/issues/design-debt/toyos-poller-is-sync-and-a-watch-moves-the-tail-in-two-steps.md new file mode 100644 index 00000000000..9295124fe32 --- /dev/null +++ b/issues/design-debt/toyos-poller-is-sync-and-a-watch-moves-the-tail-in-two-steps.md @@ -0,0 +1,16 @@ +--- +status: open +kind: defect +opened: 2026-10-02 +--- + +# `toyos::Poller` is `Sync`, and a watch moves the submission tail in two steps + +`Poller::watch_raw` loads the submission tail, writes the entry at it and +stores the tail plus one, through `&self`, and `Poller` is `unsafe impl Sync` +(`toyos/src/poller.rs`). Two threads watching through one `&Poller` write one +slot at once and advance the tail once, so one watch is lost. `rg 'Arc'` +and `rg 'static.*Poller'` over `toyos`, `userland` and `tests` find no poller +shared between threads. + +**Exit**: `Poller` is not `Sync`, or a watch claims its slot in one step. diff --git a/issues/isolation/a-ring-handed-to-another-process-resolves-its-watches-in-the-receivers-table.md b/issues/isolation/a-ring-handed-to-another-process-resolves-its-watches-in-the-receivers-table.md new file mode 100644 index 00000000000..a03992186a7 --- /dev/null +++ b/issues/isolation/a-ring-handed-to-another-process-resolves-its-watches-in-the-receivers-table.md @@ -0,0 +1,20 @@ +--- +status: open +kind: defect +opened: 2026-10-02 +--- + +# A ring handed to another process resolves its watches in the receiver's table + +An `Inbox` handle carries `DUP` and `TRANSFER` (`ops::initial_rights`), and +`inbox_submit` resolves every handle a watch names in the calling process's +table (`inbox::resolve`): an `OP_WATCH` when it is submitted, and a fired +poll's when its submitter looks at the object again. The ring's page is mapped +into its creator alone, so a process handed the ring runs the submissions the +creator wrote, and looks at the creator's fired polls, against whatever its +own table holds under the creator's handle numbers. It reaches no object it +does not hold. + +**Exit**: a ring cannot leave the process that made it, or a watch names its +object by something its registrant's table decided; a test hands a ring to a +child. diff --git a/issues/kernel/a-close-of-one-handle-ends-every-rings-poll-on-its-object.md b/issues/kernel/a-close-of-one-handle-ends-every-rings-poll-on-its-object.md new file mode 100644 index 00000000000..4d0aec608cb --- /dev/null +++ b/issues/kernel/a-close-of-one-handle-ends-every-rings-poll-on-its-object.md @@ -0,0 +1,20 @@ +--- +status: open +kind: defect +opened: 2026-10-02 +--- + +# A close of one handle ends every ring's poll on its object + +`ops::close` answers every poll on a watch its object ends with `-NotFound` +(`Watch::cancel_polls`), in every ring, when any one handle to a pipe's read +end or an acceptor closes. A sibling handle from `dup` keeps the object open, +and its polls end all the same. `toyos::poller`'s `drain` hands the token of a +negative completion to its caller as it does a ready one's, so a reader that +takes that token for bytes and reads blocking parks on an object that is open +and empty. No caller in the tree reaches it: fsd's acceptors are endowed and +never duplicated. `inbox_cancel_wakes` stages the close. + +**Exit**: a poll ends only when the last handle to its source closes, or a +completion's result reaches `Poller`'s caller; a test closes a duplicate under +a watch and reads what the wait hands back. diff --git a/issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md b/issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md index b38d2a9496f..a5d16c497e7 100644 --- a/issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md +++ b/issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md @@ -10,9 +10,9 @@ Held by the small-kernel track's stage 6 step 2 (`issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md`), whose instrument is the only thing that can read it. -Every poll ring's own watch and its completions sit behind an `IrqLock` +Every poll ring's own watch sits behind an `IrqLock` (`kernel/src/inbox/mod.rs`), because a device handler's post -reaches them through the polls it fires. So any process, not only a device's +reaches it through the polls it fires. So any process, not only a device's holder, decides how long a CPU runs with interrupts masked: - **N threads parked in `submit` on one ring** @@ -25,9 +25,9 @@ holder, decides how long a CPU runs with interrupts masked: that finds the list full copies it, all with interrupts masked. - **A claim's holder polling its claim from R rings, P polls each** (up to `MAX_PENDING_WATCHES`, 1024) makes its device's - handler fire R × P entries under the claim's list lock, each taking that - ring's completions lock and posting that ring's watch, whose own N threads - it notifies. Entries a post in place fired stay in the list until + handler fire R × P entries under the claim's list lock, each posting its + ring's watch, whose own N threads it notifies. Entries a post in place + fired stay in the list until registrations sweep them four at a time. Nothing caps N: a thread costs its process a 128 KiB kernel stack diff --git a/issues/kernel/a-zero-byte-pipe-write-wakes-the-readers-watch.md b/issues/kernel/a-zero-byte-pipe-write-wakes-the-readers-watch.md deleted file mode 100644 index 5bae5e2a1de..00000000000 --- a/issues/kernel/a-zero-byte-pipe-write-wakes-the-readers-watch.md +++ /dev/null @@ -1,27 +0,0 @@ ---- -status: open -kind: defect -opened: 2026-09-26 ---- - -# A zero-byte pipe write wakes the reader's watch - -`sys_write_nonblock` (`kernel/src/syscall/io.rs`) wakes the pipe's -readers after every write that returns `Ok(n)`, `n == 0` included, and -`complete_pending_for_event` (`kernel/src/inbox/mod.rs`) completes a pending -`READABLE` watch on that wake as ready without looking at the ring. So a -zero-byte write, which moves nothing, completes the reader's watch as -readable while the reader's `read` still answers `WouldBlock`. - -netd's liveness probes are exactly such writes, once a pass, into every -piped connection's receive pipe and every piped listener's notify pipe, and -netd passes every millisecond while a piped connection lives. Every client -watching one of those pipes is woken about a thousand times a second for -nothing. Seen, not guessed: a guest waiting on its receive pipe with a -`READABLE` watch was completed as ready within `netd_refused_pipes`' first -two seconds and then read `WouldBlock` (the run with netd's send-side -refusal mutated away, before that test learned to re-check after a wake). - -Exit condition: a write that moves no bytes wakes nobody, or a wake -completes a watch only if its direction is ready, with a test whose reader -watch stays pending across a zero-byte write. diff --git a/issues/kernel/two-completions-can-name-one-arrival-and-accept-parks.md b/issues/kernel/two-completions-can-name-one-arrival-and-accept-parks.md deleted file mode 100644 index 0313fffbad5..00000000000 --- a/issues/kernel/two-completions-can-name-one-arrival-and-accept-parks.md +++ /dev/null @@ -1,33 +0,0 @@ ---- -status: open -kind: defect -opened: 2026-09-27 ---- - -# Two completions can name one arrival, and an `accept` after the second parks - -`process_watch` (`kernel/src/inbox/mod.rs`) answers a watch on a handle that -is already ready at once, and only a watch it has to register replaces the -armed one on the same handle. So a watch submitted while a connection arrives -can complete immediately beside the older armed watch that the same arrival -fires: two completions for one connection, in one drain or in two. A read -after a spurious completion is harmless where the read does not block, but -`SYS_ACCEPT` (`kernel/src/syscall/ipc.rs`, `sys_accept`) parks until a -connection is queued, so a server that accepts once per completion takes the -connection on the first and parks for good on the second — the port it serves -then answers nobody. - -Measured on `/system/bin/fsd`'s log server under `tests/quiescecase`, whose six -writers connect beside logd at boot: a blocked-task dump showed the server -`ipc parked` (the class only `sys_accept` parks in) while logd and every writer -waited on it, in one boot of four and in one of three in a second series. fsd -now asks the acceptor with a zero-timeout watch on a probe ring before it -accepts (`userland/fsd/src/main.rs`, `Server::accept`), and after the change -the same boot ran six times and never parked. Every other server that accepts after -a completion — blockd, logd's inspect and serve threads, init, -soundd, the compositor, netd, filepicker — does not. - -**Exit**: a spurious completion cannot park a server — an accept that does not -park when nothing is queued, or a watch that withdraws the armed one on its -handle whether or not it answers at once, with a test that submits a watch -while a connection arrives and counts the completions. diff --git a/issues/kernel/two-submitters-of-one-ring-race-its-last-completion-slot.md b/issues/kernel/two-submitters-of-one-ring-race-its-last-completion-slot.md new file mode 100644 index 00000000000..9ffa89392b6 --- /dev/null +++ b/issues/kernel/two-submitters-of-one-ring-race-its-last-completion-slot.md @@ -0,0 +1,17 @@ +--- +status: open +kind: defect +opened: 2026-10-02 +--- + +# Two submitters of one ring race its last completion slot + +`polls::deliver` asks the ring for room and then answers, in two holds of the +completions' lock. Two threads in `inbox_submit` on one ring can both find the +last slot free; the second answer finds the ring full, and `post_completion` +drops it and counts it in `dropped`, as it does every completion written to a +full ring. A ring one thread submits to drops no watch's answer, and +`toyos::Poller` sizes its rings past what its watches can answer. + +**Exit**: the room an answer was looked up for is the room it is written +into; a test runs two submitters against a ring with one slot. diff --git a/kernel-loom/Cargo.toml b/kernel-loom/Cargo.toml index 53a3aef831e..9111bcfc29f 100644 --- a/kernel-loom/Cargo.toml +++ b/kernel-loom/Cargo.toml @@ -49,6 +49,15 @@ wake-fence-off = [] # # Never on by default, and no kernel build can reach it. poll-fire-load-store = [] +# The negative control for a poll ring's answer. It makes `inbox/polls.rs` +# answer a fired poll with its interest and look at nothing, so a post with +# nothing behind it is an answer and `inbox_answer.rs` must red: +# +# cargo test --manifest-path kernel-loom/Cargo.toml --features post-is-an-answer \ +# --test inbox_answer +# +# Never on by default, and no kernel build can reach it. +post-is-an-answer = [] # The negative control for the ticket lock's acquire edge. It makes `sync.rs`'s # two loads of `now` — the ones that decide ownership — `Relaxed`, so the # previous owner's writes are unordered against the next owner's reads, and diff --git a/kernel-loom/src/lib.rs b/kernel-loom/src/lib.rs index afcade99636..4f2c2cdb1bf 100644 --- a/kernel-loom/src/lib.rs +++ b/kernel-loom/src/lib.rs @@ -205,6 +205,18 @@ pub mod device_irq; #[path = "../../kernel/src/inbox/once.rs"] pub mod poll_once; +/// `polls.rs` names its one-shot as `super::once`, which in the kernel is +/// `crate::inbox::once`; this is what makes that path resolve here. +pub use poll_once as once; + +extern crate alloc; + +/// A ring's polls and when one is answered, driven against a fake object by +/// `tests/inbox_answer.rs`. It names the one-shot above, `toyos-abi` and +/// `alloc`, and nothing of the kernel's. +#[path = "../../kernel/src/inbox/polls.rs"] +pub mod inbox_polls; + /// What `sleeplock.rs` names of the kernel's watch, and nothing more. /// /// **The park is shimmed, and that is the scope statement for diff --git a/kernel-loom/tests/inbox_answer.rs b/kernel-loom/tests/inbox_answer.rs new file mode 100644 index 00000000000..d51108fb5d4 --- /dev/null +++ b/kernel-loom/tests/inbox_answer.rs @@ -0,0 +1,528 @@ +//! **A post is not an answer**: a poll ring's answer is what the object held +//! when the ring's submitter looked, never what a post once announced. +//! +//! `kernel/src/inbox/polls.rs` against fake objects, standing in for +//! `inbox::submit`'s: a watch is the handle's one poll, a post fires every poll +//! armed on its object, and a submit looks at every fired poll's object again. +//! A reader that takes an answer for bytes and then reads blocking parks for +//! good on one written for nothing, which is why each case here is a reader's +//! sequence and its assertion the answers it is handed. +//! +//! In `loom::model`, because the one-shot under each poll is built from +//! loom's atomics in this crate; every model here is one thread, and a watch +//! another thread submits during a look is staged inside the look. A fire +//! racing the submitter's park is `toyos-sched-loom`'s +//! `a_fire_racing_a_submitters_park_is_never_lost`, and two submitters of one +//! ring its `an_answer_wakes_the_submitter_its_look_hid_the_poll_from`. +//! +//! The negative case is a cargo feature: +//! +//! ```text +//! cargo test --manifest-path kernel-loom/Cargo.toml --features post-is-an-answer \ +//! --test inbox_answer +//! ``` +//! +//! answers a fired poll without looking, and this file must red. + +#![cfg(feature = "loom")] + +use std::cell::{Cell, RefCell}; +use std::sync::Arc; + +use kernel_loom::inbox_polls::{awake, deliver, Look, Poll, Polls, Submitter, Wake}; +use toyos_abi::handle::RawHandle; +use toyos_abi::inbox::READABLE; +use toyos_abi::syscall::SyscallError; + +const H: RawHandle = RawHandle(5); +const G: RawHandle = RawHandle(6); +/// The log: what it holds for a reader is the reader's cursor's to say. +const L: RawHandle = RawHandle(7); + +/// The answer for bytes. +const BYTES: i32 = READABLE as i32; + +const NOTHING: [(u64, i32); 0] = []; + +/// The ring a fire tells, which has no waiter: every model is one thread. +#[derive(Clone)] +struct Ring; + +impl Wake for Ring { + fn wake(&self) {} +} + +/// One object: whether it holds bytes, and the polls armed on its watch. +#[derive(Default)] +struct Object { + bytes: Cell, + /// Whether a post on it is its readiness, the kernel holding none to look at. + posts_are_readiness: bool, + armed: RefCell>>>, +} + +/// One ring's kernel over the objects `H`, `G` and `L`. +struct Kernel { + ring: Ring, + polls: RefCell>, + /// How many polls the ring keeps. + cap: usize, + h: Object, + g: Object, + l: Object, + /// The handle the process has closed since it watched it. + closed: Cell>, + /// A watch another thread submits while the next look runs. + lands_during_look: Cell>, + /// How many more renewals a peer's post follows at once. + posts_after_renewal: Cell, + /// How many looks the submitter has made. + looks: Cell, + /// The completion ring, and how many answers it holds. + answers: RefCell>, + size: usize, +} + +impl Kernel { + fn new(size: usize) -> Self { + Self { + ring: Ring, + polls: RefCell::new(Polls::new()), + cap: 16, + h: Object::default(), + g: Object::default(), + l: Object { posts_are_readiness: true, ..Object::default() }, + closed: Cell::new(None), + lands_during_look: Cell::new(None), + posts_after_renewal: Cell::new(0), + looks: Cell::new(0), + answers: RefCell::new(Vec::new()), + size, + } + } + + fn object(&self, handle: RawHandle) -> &Object { + match handle { + H => &self.h, + G => &self.g, + L => &self.l, + other => panic!("no object behind {other:?}"), + } + } + + /// `OP_WATCH` for bytes on `handle`, answered under `token`: the handle's + /// one poll. + fn watch(&self, handle: RawHandle, token: u64) { + assert!(self.admit(handle, token), "the ring keeps no more polls"); + } + + /// The same, answering whether the ring kept the poll. + fn admit(&self, handle: RawHandle, token: u64) -> bool { + let poll = Arc::new(Poll::new(self.ring.clone(), token, handle, READABLE)); + let kept = self.polls.borrow_mut().admit(poll.clone(), self.cap); + if kept { + self.arm(poll); + } + kept + } + + /// Fired now if its object holds bytes, armed on its watch if not. + fn arm(&self, poll: Arc>) { + let object = self.object(poll.handle); + if object.bytes.get() { + poll.fire(0); + } else { + object.armed.borrow_mut().push(poll); + } + } + + /// Bytes reach the object; its post has not run. + fn fill(&self, handle: RawHandle) { + self.object(handle).bytes.set(true); + } + + /// The object's post: every poll armed on it fires, whatever it holds. + fn post(&self, handle: RawHandle) { + for poll in self.object(handle).armed.take() { + poll.fire(READABLE); + } + } + + /// The reader takes every byte the object holds. + fn take(&self, handle: RawHandle) { + self.object(handle).bytes.set(false); + } + + /// The object's source ends. + fn end(&self, handle: RawHandle) { + for poll in self.object(handle).armed.take() { + poll.end(); + } + } + + /// The process closes `handle`, whose watch other handles share: no close + /// ends its polls. + fn close(&self, handle: RawHandle) { + self.closed.set(Some(handle)); + } + + /// How many polls on `handle`'s watch a post would still fire. + fn live(&self, handle: RawHandle) -> usize { + self.object(handle).armed.borrow().iter().filter(|poll| poll.armed()).count() + } + + /// `inbox_submit`'s look at what the ring is owed. + fn submit(&self) { + deliver(self); + } + + /// Whether a submitter waiting for one more answer than the ring holds + /// would go round again instead of parking. + fn awake(&self) -> bool { + awake(self, || false) + } + + /// Every answer the reader's drain hands it. + fn drain(&self) -> Vec<(u64, i32)> { + self.answers.take() + } +} + +/// An answer's wake, which has no waiter either. +impl Wake for Kernel { + fn wake(&self) {} +} + +impl Submitter for Kernel { + fn room(&self) -> bool { + self.answers.borrow().len() < self.size + } + + fn answer(&self, user_data: u64, result: i32) { + self.answers.borrow_mut().push((user_data, result)); + } + + fn polls(&self, f: impl FnOnce(&mut Polls) -> R) -> Option { + Some(f(&mut self.polls.borrow_mut())) + } + + fn look(&self, poll: &Arc>) -> Look { + self.looks.set(self.looks.get() + 1); + if let Some((handle, token)) = self.lands_during_look.take() { + self.watch(handle, token); + } + if self.closed.get() == Some(poll.handle) { + return Look::Refused(SyscallError::NotFound); + } + let object = self.object(poll.handle); + if object.bytes.get() || (object.posts_are_readiness && poll.posted() & READABLE != 0) { + return Look::Ready(READABLE); + } + let again = Arc::new(Poll::new(self.ring.clone(), poll.user_data, poll.handle, poll.flags)); + if self.polls.borrow_mut().renew(poll, again.clone()) { + self.arm(again); + } + if let Some(left) = self.posts_after_renewal.get().checked_sub(1) { + self.posts_after_renewal.set(left); + self.post(poll.handle); + } + Look::Waits + } +} + +/// A peer's write of no bytes posts the pipe all the same +/// (`pipe::try_write` answers `Wrote(0)` and `sys_write` wakes the readers). +/// Nothing is there to read, so nothing is answered, and the poll armed again +/// answers the bytes that do come. +#[test] +fn a_post_with_nothing_to_read_answers_nothing() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 7); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING); + + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING, "a post with nothing behind it was answered"); + + kernel.fill(H); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(7, BYTES)]); + }); +} + +/// A frame is two writes, and a write posts after it has published its +/// bytes. The header's post wakes the reader, which reads header and payload +/// and watches again before the payload's post has run; that post lands on +/// the new watch with nothing left to read. +#[test] +fn a_post_that_lands_after_its_bytes_were_read_answers_nothing() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 7); + kernel.submit(); + + kernel.fill(H); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(7, BYTES)]); + + kernel.take(H); + kernel.watch(H, 7); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING, "the payload's late post was answered"); + }); +} + +/// A post fires a poll its reader then replaces before its next wait: one +/// arrival, answered once, under the newest watch's token. +#[test] +fn a_watch_replaces_a_poll_a_post_already_fired() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 1); + kernel.submit(); + + kernel.fill(H); + kernel.post(H); + kernel.watch(H, 2); + kernel.submit(); + assert_eq!(kernel.drain(), [(2, BYTES)]); + }); +} + +/// A look that finds nothing arms the poll again and goes on to the next one. +#[test] +fn a_poll_armed_again_does_not_end_the_look() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 3); + kernel.watch(G, 4); + kernel.submit(); + + kernel.post(H); + kernel.fill(G); + kernel.post(G); + kernel.submit(); + assert_eq!(kernel.drain(), [(4, BYTES)]); + }); +} + +/// A watch replaces only its own handle's poll: one the reader did not renew +/// still answers when its object fills. +#[test] +fn a_poll_left_standing_still_answers() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 3); + kernel.watch(G, 4); + kernel.submit(); + kernel.watch(G, 4); + kernel.submit(); + + kernel.fill(H); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(3, BYTES)]); + }); +} + +/// A source that ended is answered as gone: there is nothing to look at. +#[test] +fn an_ended_source_is_answered_gone_without_a_look() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 7); + kernel.submit(); + + kernel.end(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(7, -(SyscallError::NotFound as i32))]); + }); +} + +/// A full ring stops the look, and the poll it did not reach is answered by +/// the next wait's rather than dropped. +#[test] +fn a_full_ring_leaves_a_fired_poll_for_the_next_look() { + loom::model(|| { + let kernel = Kernel::new(1); + kernel.watch(H, 3); + kernel.watch(G, 4); + kernel.submit(); + + for handle in [H, G] { + kernel.fill(handle); + kernel.post(handle); + } + kernel.submit(); + assert_eq!(kernel.drain(), [(3, BYTES)]); + kernel.submit(); + assert_eq!(kernel.drain(), [(4, BYTES)]); + }); +} + +/// The log's unread records are a property of the reader's cursor, which the +/// kernel does not hold: there is nothing to look at, and its post is the +/// answer. +#[test] +fn a_post_answers_for_an_object_with_nothing_to_look_at() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(L, 7); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING); + + kernel.post(L); + kernel.submit(); + assert_eq!(kernel.drain(), [(7, BYTES)]); + }); +} + +/// A watch that replaces an armed poll takes it off its object's watch: a +/// poll left live there is one more entry for every re-watch of an idle handle. +#[test] +fn a_replaced_poll_is_withdrawn_from_its_watch() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 1); + kernel.submit(); + kernel.watch(H, 2); + kernel.submit(); + assert_eq!(kernel.live(H), 1, "the replaced poll is still live on its object's watch"); + }); +} + +/// A ring torn down leaves no poll live on any watch. +#[test] +fn a_ring_torn_down_withdraws_every_poll() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 1); + kernel.watch(G, 2); + kernel.submit(); + + kernel.polls.borrow_mut().withdraw_all(); + assert_eq!((kernel.live(H), kernel.live(G)), (0, 0)); + }); +} + +/// A handle closed since it was watched names nothing to look at: the look's +/// refusal is the answer. +#[test] +fn a_handle_closed_since_its_watch_is_answered_with_the_refusal() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 7); + kernel.submit(); + + kernel.close(H); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(7, -(SyscallError::NotFound as i32))]); + }); +} + +/// A ring keeps its cap of polls and no more, and a watch that replaces one +/// of them is inside it. +#[test] +fn a_ring_keeps_its_cap_of_polls_and_no_more() { + loom::model(|| { + let kernel = Kernel { cap: 2, ..Kernel::new(4) }; + assert!(kernel.admit(H, 1)); + assert!(kernel.admit(G, 2), "a ring under its cap refused a poll"); + assert!(!kernel.admit(L, 3), "a ring at its cap kept one more poll"); + assert!(kernel.admit(H, 4), "a ring at its cap refused a watch that replaces a poll"); + }); +} + +/// A second thread's watch lands while the look at the handle's fired poll +/// finds nothing: the newer watch holds the handle's place, and the look arms +/// nothing beside it. +#[test] +fn a_watch_during_a_look_that_finds_nothing_is_the_handles_one_poll() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 1); + kernel.submit(); + + kernel.post(H); + kernel.lands_during_look.set(Some((H, 2))); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING); + assert_eq!(kernel.live(H), 1, "the look armed the older watch again beside the newer"); + + kernel.fill(H); + kernel.post(H); + kernel.submit(); + assert_eq!(kernel.drain(), [(2, BYTES)]); + }); +} + +/// The same while the look finds bytes: one arrival, answered once, under the +/// newer watch's token, by the pass after the one it landed in. +#[test] +fn a_watch_during_a_look_that_finds_bytes_answers_alone() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 1); + kernel.submit(); + + kernel.fill(H); + kernel.post(H); + kernel.lands_during_look.set(Some((H, 2))); + kernel.submit(); + assert_eq!(kernel.drain(), NOTHING, "the look answered under the watch it was replaced by"); + assert!(kernel.awake(), "the newer watch's fire is owed the next look"); + kernel.submit(); + assert_eq!(kernel.drain(), [(2, BYTES)]); + }); +} + +/// A submitter does not park over a poll it could answer, and does over a +/// poll a full ring leaves it no room to answer: it would go round for good. +#[test] +fn a_submitter_parks_only_with_nothing_it_could_answer() { + loom::model(|| { + let kernel = Kernel::new(1); + kernel.watch(H, 3); + kernel.watch(G, 4); + kernel.submit(); + assert!(!kernel.awake(), "nothing is owed a look"); + + for handle in [H, G] { + kernel.fill(handle); + kernel.post(handle); + } + assert!(kernel.awake(), "a fired poll is owed a look and the ring has room"); + kernel.submit(); + assert!(!kernel.awake(), "the ring is full: the poll left over waits for room"); + assert_eq!(kernel.drain(), [(3, BYTES)]); + assert!(kernel.awake(), "room came back with a poll still owed its look"); + }); +} + +/// A peer that posts an empty object as fast as the submitter looks at it does +/// not hold the look: the poll renewed in this pass waits for the next, and +/// the handle behind it is answered. +#[test] +fn a_poll_fired_as_it_is_renewed_waits_for_the_next_look() { + loom::model(|| { + let kernel = Kernel::new(4); + kernel.watch(H, 3); + kernel.watch(G, 4); + kernel.submit(); + + kernel.post(H); + kernel.fill(G); + kernel.post(G); + kernel.posts_after_renewal.set(3); + kernel.submit(); + assert_eq!(kernel.looks.get(), 2, "the look went round on a poll it had renewed"); + assert_eq!(kernel.drain(), [(4, BYTES)]); + assert!(kernel.awake(), "the renewed poll's fire is owed the next look"); + }); +} diff --git a/kernel/Cargo.toml b/kernel/Cargo.toml index 38bf47fef07..6035b2d4891 100644 --- a/kernel/Cargo.toml +++ b/kernel/Cargo.toml @@ -34,6 +34,11 @@ serial-try-lock-then-some = [] # `kernel-loom/tests/poll_once.rs` must red when it does — a poll two posts # both answer. poll-fire-load-store = [] +# The same for a poll ring's answer: `kernel-loom` turns it on so +# `inbox/polls.rs` answers a fired poll without looking at its object again, +# and `kernel-loom/tests/inbox_answer.rs` must red — a post with nothing to +# read, answered. +post-is-an-answer = [] # The same arrangement, five more times, one per `kernel-loom` model that had # no control of its own until 2026-08-17. Each is declared here and never # enabled here; `kernel-loom` turns exactly one on at a time and the named diff --git a/kernel/src/inbox/mod.rs b/kernel/src/inbox/mod.rs index d2453243fa9..1d343a30029 100644 --- a/kernel/src/inbox/mod.rs +++ b/kernel/src/inbox/mod.rs @@ -5,42 +5,45 @@ //! //! **A watch is a poll registered on the watched object's own //! [`Watch`](crate::watch::Watch)**, one entry per direction it asked for, and -//! the object's post completes it. There is no table of sources here: what a +//! 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. //! +//! **Only the ring's submitter writes a watch's answer, after a look** +//! ([`polls`]): a post owes the poll a look, and `submit` looks at the object +//! again before it writes anything, so an answer is never older than the wait +//! that returned it. +//! //! **A completion is a trust boundary, and the kernel is the only writer of //! its position.** `completion_tail` lives here, never in the page; the head -//! is the process's and is read once per post, so a head the process lies +//! is the process's and is read once per write, so a head the process lies //! about makes the kernel drop the completion and count it, never write //! outside the ring. What a completion says comes from the kernel — the -//! caller's own `token`, and a result that is either the direction the object -//! posted or the refusal — so a process cannot make one appear in another -//! process's ring or say something no object said. A post writes at most one -//! entry, and a ring holds at most [`MAX_PENDING_WATCHES`] polls. +//! caller's own `token`, and a result that is either the directions the object +//! was ready in or the refusal — so a process cannot make one appear in another +//! process's ring or say something no object said. A ring holds at most +//! [`MAX_PENDING_WATCHES`] polls. //! //! **Locks.** What a completion writes, and the page it is written into, sit -//! behind an [`IrqLock`] of their own, because a device's interrupt handler -//! posts in place, firing its polls under its list lock; nothing is taken -//! under it. The rest of a ring, its submissions and its polls, is its -//! `Lock`'s, which no post reaches. A ring's own watch is an [`IrqWatch`] and -//! holds only threads, because no handle names a ring as a thing to watch. +//! behind a `Lock` of their own; nothing is taken under it. The rest of a +//! ring, its submissions and its polls, is another's. No post reaches either: +//! a fire takes its poll and posts the ring's watch. That watch is an +//! [`IrqWatch`], because a device's interrupt handler fires polls, and holds +//! only threads, because no handle names a ring as a thing to watch. use alloc::sync::Arc; -use alloc::vec::Vec; use core::sync::atomic::Ordering; -use toyos_sched::sync::CellLock; use toyos_sched::task::WaitClass; use toyos_sched::watch::{Fire, Ring}; use crate::object::shm::SharedMemObject; -use crate::object::{ops, KObjectRef}; +use crate::object::{ops, HandleError, KObjectRef}; use crate::process::{self, Pid}; use crate::scheduler; use crate::sync::Lock; use crate::time::{Deadline, Duration}; -use crate::watch::{IrqLock, IrqWatch}; +use crate::watch::IrqWatch; use crate::DirectMap; use toyos_abi::inbox::{ @@ -51,8 +54,9 @@ use toyos_abi::handle::{RawHandle, Rights}; use toyos_abi::syscall::SyscallError; mod once; +mod polls; -use once::Once; +use polls::{Look, Polls, Submitter}; /// The one owned reference to a ring, held by its handle's object; dropping it /// tears the ring down. @@ -66,18 +70,16 @@ impl InboxRef { impl Drop for InboxRef { fn drop(&mut self) { - // The page is the completions', so it goes only once no post can reach + // The page is the completions', so it goes only once nothing can reach // them. Both halves are taken out under their locks and let go of // outside them: the unmap flushes. - let Some(completions) = self.0.completions.with(Option::take) else { + let Some(completions) = self.0.completions.lock().take() else { unreachable!("an inbox is torn down by its one reference, once"); }; let Some(mut state) = self.0.state.lock().take() else { unreachable!("an inbox is torn down by its one reference, once"); }; - for poll in state.pending.drain(..) { - poll.withdraw(); - } + state.polls.withdraw_all(); // `Unmapped`'s drop flushes; the `Arc` drop after it frees the pages. drop(completions.shm.unmap_from(state.owner_pid)); } @@ -142,35 +144,17 @@ impl Readiness { if self.writable { flags |= WatchFlags::WRITABLE.raw(); } flags } -} -/// One `OP_WATCH` a ring is waiting on: one-shot across every watch it is -/// registered on and against its own registrant's recheck. -pub struct Poll { - inbox: Arc, - user_data: u64, - /// The handle the poll was submitted against; the dedup key. - handle: RawHandle, - /// Taken by exactly one of a fire and a withdrawal. - state: Once, -} - -impl Poll { - /// Post this poll's completion if nothing has answered it yet. - fn complete(&self, result: i32) { - if self.state.fire() { - self.inbox.complete(self.user_data, result); - } + fn any(self) -> bool { + self.readable || self.writable } +} - /// Answer nothing: a newer poll on the same handle replaced it, or its - /// ring went away. - fn withdraw(&self) { - let _ = self.state.withdraw(); - } +type Poll = polls::Poll>; - fn armed(&self) -> bool { - self.state.armed() +impl polls::Wake for Arc { + fn wake(&self) { + self.watch.post_in_place(); } } @@ -178,16 +162,15 @@ impl Poll { /// watch is. pub struct PollEntry { poll: Arc, - direction: Readiness, + direction: WatchFlags, } impl Ring for PollEntry { fn fire(&self, how: Fire) { - self.poll.complete(match how { - // The direction this watch is: its object posted it. - Fire::Ready => self.direction.result_flags() as i32, - Fire::Gone => -(SyscallError::NotFound as i32), - }); + match how { + Fire::Ready => self.poll.fire(self.direction.raw()), + Fire::Gone => self.poll.end(), + } } fn live(&self) -> bool { @@ -198,12 +181,11 @@ impl Ring for PollEntry { /// Hard cap on pending polls per ring. const MAX_PENDING_WATCHES: usize = 1024; -/// A ring: what a poll posts into, and what `submit` parks on. pub struct Inbox { /// `None` once the ring's one reference let go of it. state: Lock>, /// `None` from the moment that reference starts letting go of it. - completions: IrqLock>, + completions: Lock>, /// Threads parked in `submit`; never a poll — see the module header. watch: IrqWatch, } @@ -213,8 +195,7 @@ struct RingState { /// lets go of only after it has taken this. shm_phys: DirectMap, submission_size: u32, - /// Polls still armed as of the last registration, which sweeps the rest. - pending: Vec>, + polls: Polls>, owner_pid: Pid, } @@ -249,8 +230,8 @@ impl RingState { } } -/// What a poll's completion writes, which a post from an interrupt handler -/// reaches, and the ring's page, which goes only with these. +/// What a poll's completion writes, and the ring's page, which goes only +/// with these. struct Completions { /// A ring's page has no lifetime of its own; it goes with the last handle to the ring. shm: Arc, @@ -281,8 +262,6 @@ impl Completions { } /// Posts a completion, or records a drop if the ring reports itself full. - /// A full ring is not fatal here: a poll completes on the poster's thread, - /// which belongs to a different process. fn post_completion(&mut self, user_data: u64, result: i32, flags: u32) { let tail = self.completion_tail; if tail.wrapping_sub(self.completion_head().load(Ordering::Acquire)) >= self.completion_size { @@ -303,6 +282,10 @@ impl Completions { self.completion_tail.wrapping_sub(head) } + fn room(&self) -> bool { + self.completion_count() < self.completion_size + } + /// Cumulative, never cleared. fn dropped(&self) -> u32 { self.completion_dropped().load(Ordering::Relaxed) @@ -310,25 +293,12 @@ impl Completions { } impl Inbox { - /// Post one completion and wake whoever waits in `submit`. A ring already - /// torn down takes nothing and wakes nobody. - fn complete(&self, user_data: u64, result: i32) { - let posted = self.completions.with(|c| { - c.as_mut().map(|c| c.post_completion(user_data, result, 0)) - }); - if posted.is_some() { - // In place: an interrupt handler's post reaches here through the - // poll it fires. - self.watch.post_in_place(); - } - } - fn with_state(&self, f: impl FnOnce(&mut RingState) -> R) -> Result { self.state.lock().as_mut().map(f).ok_or(SyscallError::NotFound) } fn with_completions(&self, f: impl FnOnce(&Completions) -> R) -> Result { - self.completions.with(|c| c.as_ref().map(f)).ok_or(SyscallError::NotFound) + self.completions.lock().as_ref().map(f).ok_or(SyscallError::NotFound) } } @@ -387,10 +357,10 @@ pub fn create(depth: u32) -> Result<(InboxRef, u64), SyscallError> { state: Lock::new(Some(RingState { shm_phys, submission_size, - pending: Vec::new(), + polls: Polls::new(), owner_pid: pid, })), - completions: IrqLock::new(Some(Completions { + completions: Lock::new(Some(Completions { shm, page: shm_phys, completion_size, @@ -423,6 +393,7 @@ pub fn submit( } loop { + polls::deliver(inbox); let (count, dropped) = inbox.with_completions(|c| (c.completion_count(), c.dropped()))?; if count >= min_complete || min_complete == 0 { @@ -442,6 +413,12 @@ pub fn submit( return Ok(count); } + // Read here and not only at the park: a peer whose posts keep this + // thread looking keeps it from the park, and may not keep its kill. + if crate::sched::driver::current_kill_pending() { + return Err(SyscallError::Gone); + } + // The recheck closure is this ring's own condition, not mere readiness — else a waiter for `min_complete` spins. let parkable = scheduler::Parkable::at_entry(); if crate::watch::wait_until( @@ -450,7 +427,11 @@ pub fn submit( 0, WaitClass::Io, deadline, - || inbox.with_completions(|c| c.completion_count()).map_or(true, |n| n >= min_complete), + || { + polls::awake(inbox, || { + inbox.with_completions(|c| c.completion_count()).map_or(true, |n| n >= min_complete) + }) + }, ) .is_err() { @@ -495,109 +476,101 @@ fn process_submission(inbox: &Arc, submission: &Submission) { // `Submission::flags` is declared and read by nothing, so a caller setting // it is asking for a behaviour that does not exist. if submission.flags != 0 { - inbox.complete(submission.token, -(SyscallError::InvalidArgument as i32)); + polls::complete(inbox, submission.token, -(SyscallError::InvalidArgument as i32)); return; } let op = match Op::from_raw(submission.op) { Ok(op) => op, Err(_) => { - inbox.complete(submission.token, -(SyscallError::InvalidArgument as i32)); + polls::complete(inbox, submission.token, -(SyscallError::InvalidArgument as i32)); return; } }; match op { - Op::Nop => inbox.complete(submission.token, 0), + Op::Nop => polls::complete(inbox, submission.token, 0), Op::Watch => process_watch(inbox, submission), Op::Accept => process_accept(inbox, submission), } } -/// Registers an `OP_WATCH`, or answers it immediately; every refusal posts a completion rather than going silent. +/// Takes an `OP_WATCH` as its handle's one poll; every refusal writes a completion rather than going silent. fn process_watch(inbox: &Arc, submission: &Submission) { - let handle = submission.handle; let user_data = submission.token; let flags = match WatchFlags::from_raw(submission.op_flags) { Ok(flags) => flags, Err(e) => { - inbox.complete(user_data, -(e as i32)); + polls::complete(inbox, user_data, -(e as i32)); return; } }; - - // Readiness is checked on the process's table, not the thread's: a ring is process-wide. - // The object is cloned out so the registration below holds no process lock. - let resolved = process::with_process_data(|data| { - data.handles.get_ref(handle, Rights::WAIT).cloned() + // Nothing is held here: `resolve` has given the guard up. + let refused = resolve(submission.handle).map_err(HandleError::refuse_as_error).and_then(|object| { + arm(inbox, Poll::new(inbox.clone(), user_data, submission.handle, flags.raw()), None, &object) }); - let object = match resolved { - Ok(object) => object, - // Nothing is held here: `with_process_data` has given the guard up. - Err(e) => { - let refusal = e.refuse_as_error(); - inbox.complete(user_data, -(refusal as i32)); - return; - } - }; - - let readiness = readiness_of(&object, flags); - if readiness.readable || readiness.writable { - // Ready already: complete now, one-shot, with the directions that fired. - inbox.complete(user_data, readiness.result_flags() as i32); - return; + if let Err(refusal) = refused { + polls::complete(inbox, user_data, -(refusal as i32)); } +} - let read = if flags.readable() { ops::read_watch(&object) } else { None }; - let write = if flags.writable() { ops::write_watch(&object) } else { None }; - // No readiness in either direction: nothing could ever complete this poll, so it is refused, not registered. - if read.is_none() && write.is_none() { - inbox.complete(user_data, -(SyscallError::NotSupported as i32)); - return; - } +/// The object `handle` names in this process's table, not the thread's: a ring is process-wide. +/// Cloned out, so what follows holds no process lock. +fn resolve(handle: RawHandle) -> Result { + process::with_process_data(|data| data.handles.get_ref(handle, Rights::WAIT).cloned()) +} - let poll = Arc::new(Poll { - inbox: inbox.clone(), - user_data, - handle, - state: Once::new(), - }); +/// Keep `poll` as its handle's one poll, or in the place of `looked_at`, the +/// poll a look found nothing for; then fire it if `object` is ready or arm it +/// on the object's watches if not. The refusal, when nothing could ever answer +/// it. +fn arm( + inbox: &Arc, + poll: Poll, + looked_at: Option<&Arc>, + object: &KObjectRef, +) -> Result<(), SyscallError> { + let flags = WatchFlags(poll.flags); + 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 }; + // 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); + } + + let poll = Arc::new(poll); // The cap is checked before registering: registering first would leave a // watch holding a poll the ring never counted. - let admitted = inbox.with_state(|state| { - // The old poll on this handle answers nothing once this one replaces it. - if let Some(at) = state.pending.iter().position(|p| p.handle == handle && p.armed()) { - state.pending.swap_remove(at).withdraw(); - } - state.pending.retain(|p| p.armed()); - if state.pending.len() >= MAX_PENDING_WATCHES { - return false; - } - state.pending.push(poll.clone()); - true + let kept = inbox.with_state(|state| match looked_at { + None if state.polls.admit(poll.clone(), MAX_PENDING_WATCHES) => Ok(true), + None => Err(SyscallError::ResourceExhausted), + Some(old) => Ok(state.polls.renew(old, poll.clone())), }); - match admitted { - Ok(true) => {} - Ok(false) => { - inbox.complete(user_data, -(SyscallError::ResourceExhausted as i32)); - return; - } - // The ring's last handle closed under this submit. - Err(_) => return, + match kept { + Ok(Ok(true)) => {} + Ok(Err(full)) => return Err(full), + // A newer watch took the handle's place during the look, or the + // ring's last handle closed under this submit. + Ok(Ok(false)) | Err(_) => return Ok(()), + } + if ready { + poll.fire(0); + return Ok(()); } // Registered with no ring lock held, then rechecked: a post either ran // before the registration — and the recheck sees what it changed — or - // finds the entry. Whichever answers first answers alone. + // finds the entry. Whichever fires first fires alone. if let Some(watch) = &read { - watch.add_poll(PollEntry { poll: poll.clone(), direction: Readiness { readable: true, writable: false } }); + watch.add_poll(PollEntry { poll: poll.clone(), direction: WatchFlags::READABLE }); } if let Some(watch) = &write { - watch.add_poll(PollEntry { poll: poll.clone(), direction: Readiness { readable: false, writable: true } }); + watch.add_poll(PollEntry { poll: poll.clone(), direction: WatchFlags::WRITABLE }); } - let now = readiness_of(&object, flags); - if now.readable || now.writable { - poll.complete(now.result_flags() as i32); + if readiness_of(object, flags).any() { + poll.fire(0); } + Ok(()) } /// Per-direction readiness of the object, restricted to what was asked for. @@ -608,8 +581,53 @@ fn readiness_of(object: &KObjectRef, flags: WatchFlags) -> Readiness { } } +impl Submitter> for Arc { + fn room(&self) -> bool { + self.with_completions(Completions::room).unwrap_or(false) + } + + fn answer(&self, user_data: u64, result: i32) { + // A ring already torn down takes nothing. + if let Some(completions) = self.completions.lock().as_mut() { + completions.post_completion(user_data, result, 0); + } + } + + fn polls(&self, f: impl FnOnce(&mut Polls>) -> R) -> Option { + self.with_state(|state| f(&mut state.polls)).ok() + } + + fn look(&self, poll: &Arc) -> Look { + let object = match resolve(poll.handle) { + Ok(object) => object, + // Closed since it was watched, which is no bug of the process's: + // the poll is over, as a close that ends it says. + Err(HandleError::BadHandle | HandleError::Stale | HandleError::WrongType { .. }) => { + return Look::Refused(SyscallError::NotFound); + } + Err(HandleError::Rights { .. }) => return Look::Refused(SyscallError::PermissionDenied), + Err(HandleError::TableFull) => return Look::Refused(SyscallError::ResourceExhausted), + }; + let flags = WatchFlags(poll.flags); + let mut now = readiness_of(&object, flags); + // A post on its read watch is the object's readability, and nothing + // in the kernel could be looked at instead. + now.readable |= flags.readable() + && poll.posted() & WatchFlags::READABLE.raw() != 0 + && ops::read_posts_are_readiness(&object); + if now.any() { + return Look::Ready(now.result_flags()); + } + let again = Poll::new(self.clone(), poll.user_data, poll.handle, poll.flags); + match arm(self, again, Some(poll), &object) { + Ok(()) => Look::Waits, + Err(refusal) => Look::Refused(refusal), + } + } +} + /// The submission form of `SYS_ACCEPT`; refusals fold into one `-InvalidArgument` completion instead of ending the process. -fn process_accept(inbox: &Inbox, submission: &Submission) { +fn process_accept(inbox: &Arc, submission: &Submission) { let user_data = submission.token; let acceptor = process::with_process_data(|data| { @@ -621,7 +639,7 @@ fn process_accept(inbox: &Inbox, submission: &Submission) { // Nothing held: `with_process_data` has given the guard up. Err(e) => { let refusal = e.refuse_as_error(); - inbox.complete(user_data, -(refusal as i32)); + polls::complete(inbox, user_data, -(refusal as i32)); return; } }; @@ -640,10 +658,10 @@ fn process_accept(inbox: &Inbox, submission: &Submission) { ) }); match installed { - Ok(h) => inbox.complete(user_data, h.0 as i32), - Err(e) => inbox.complete(user_data, -(e as i32)), + Ok(h) => polls::complete(inbox, user_data, h.0 as i32), + Err(e) => polls::complete(inbox, user_data, -(e as i32)), } } - None => inbox.complete(user_data, -(SyscallError::WouldBlock as i32)), + None => polls::complete(inbox, user_data, -(SyscallError::WouldBlock as i32)), } } diff --git a/kernel/src/inbox/once.rs b/kernel/src/inbox/once.rs index 46097fb1ec6..7376f2bcbfa 100644 --- a/kernel/src/inbox/once.rs +++ b/kernel/src/inbox/once.rs @@ -1,6 +1,7 @@ //! The one decision a poll makes that two CPUs race for: which of a post, the -//! registrant's own recheck and a withdrawal answers it. Exactly one wins, so -//! a poll completes at most once and a withdrawn poll never does. +//! registrant's own recheck, the end of its source and a withdrawal takes it. +//! Exactly one wins, so a poll is owed at most one look and a withdrawn poll +//! none. //! //! Compiled a second time by `kernel-loom`, so it names only atomics. @@ -12,6 +13,7 @@ use loom::sync::atomic::{AtomicU8, Ordering}; const ARMED: u8 = 0; const FIRED: u8 = 1; const WITHDRAWN: u8 = 2; +const ENDED: u8 = 3; /// A poll's answer, taken once. pub struct Once(AtomicU8); @@ -34,6 +36,12 @@ impl Once { self.take(FIRED) } + /// `true` for the one caller that takes the poll because its source + /// ended, so the answer is that and not a look at the object. + pub fn end(&self) -> bool { + self.take(ENDED) + } + /// Take the poll back unanswered; `true` if it had not been answered. pub fn withdraw(&self) -> bool { self.take(WITHDRAWN) @@ -44,6 +52,11 @@ impl Once { self.0.load(Ordering::Acquire) == ARMED } + /// Whether the end of its source took the poll. Permanent once `true`. + pub fn ended(&self) -> bool { + self.0.load(Ordering::Acquire) == ENDED + } + // `poll-fire-load-store` is the negative control: the exchange split into // a load and a store lets two callers both answer, and `poll_once` reds. #[cfg(not(feature = "poll-fire-load-store"))] diff --git a/kernel/src/inbox/polls.rs b/kernel/src/inbox/polls.rs new file mode 100644 index 00000000000..974a66829bd --- /dev/null +++ b/kernel/src/inbox/polls.rs @@ -0,0 +1,279 @@ +//! A ring's polls, and when one of them is answered. +//! +//! **A post is not an answer.** An object's post fires the polls on its watch, +//! and a fire only owes the poll a look. The answer is written by the ring's +//! own submitter, in its `inbox_submit`, after it has looked at the object +//! again ([`deliver`]): an object ready at that look is answered with what it +//! holds then, and one that is not is armed again. So an answer says what the +//! object held when the wait that returned it looked — a peer's empty write, +//! and a post that lands after its bytes were read, answer nothing. +//! +//! **A handle has one poll that may answer.** A watch replaces the handle's +//! earlier poll whether a post has fired it or not, and whether or not a +//! submitter is looking at it ([`Polls::admit`]), so one look answers a handle +//! once, under the token of its newest watch. +//! +//! **A submitter parks on the polls themselves** ([`awake`]): a fire takes its +//! poll and then posts the watch the submitter parks on, and the submitter, +//! registered there, reads its polls again before it parks. The lost-wake +//! argument is that watch's, and nothing here records a fire beside the poll. +//! A poll one submitter is looking at is hidden from every other, which may +//! park over it: what wakes that one is the answer ([`complete`]), which +//! posts the same watch once the completion is written. +//! +//! Compiled a second time by `kernel-loom` and by `toyos-sched-loom`, so it +//! names nothing of the kernel's. + +use alloc::sync::Arc; +use alloc::vec::Vec; + +#[cfg(not(feature = "loom"))] +use core::sync::atomic::{AtomicU32, Ordering}; +#[cfg(feature = "loom")] +use loom::sync::atomic::{AtomicU32, Ordering}; + +use toyos_abi::handle::RawHandle; +use toyos_abi::syscall::SyscallError; + +use super::once::Once; + +/// A ring, as a fire sees it. +pub trait Wake { + /// Wake whoever waits in this ring's `inbox_submit`. Called from an + /// interrupt handler too. + fn wake(&self); +} + +/// One `OP_WATCH` a ring is waiting on: one-shot across every watch it is +/// registered on and against its own registrant's recheck. +pub struct Poll { + ring: W, + pub user_data: u64, + /// The handle the poll was submitted against; the key a watch replaces by. + pub handle: RawHandle, + /// The submission's interest, which the look asks about again. + pub flags: u32, + /// The directions whose watch posted it, for an object whose post is its + /// readiness. + posted: AtomicU32, + state: Once, +} + +impl Poll { + pub fn new(ring: W, user_data: u64, handle: RawHandle, flags: u32) -> Self { + Self { ring, user_data, handle, flags, posted: AtomicU32::new(0), state: Once::new() } + } + + /// The object may be ready: the watch of the directions in `posted` was + /// posted, or with `0` its registrant saw it ready. Owes a look, once. + pub fn fire(&self, posted: u32) { + // Before the exchange, so the look that the winning fire owes sees it. + self.posted.fetch_or(posted, Ordering::Release); + if self.state.fire() { + self.ring.wake(); + } + } + + /// The object's source ended: the poll is answered as gone, with no look. + pub fn end(&self) { + if self.state.end() { + self.ring.wake(); + } + } + + /// Answer nothing: a newer poll on the same handle replaced it, or its + /// ring went away. + pub fn withdraw(&self) { + let _ = self.state.withdraw(); + } + + pub fn armed(&self) -> bool { + self.state.armed() + } + + /// The directions whose watch posted this poll. + pub fn posted(&self) -> u32 { + self.posted.load(Ordering::Acquire) + } +} + +/// A ring's polls that may still answer: armed, or taken and not yet answered. +pub struct Polls { + polls: Vec>, + /// How many polls this ring ever kept: the next one's `order`. + kept: u64, +} + +struct Kept { + poll: Arc>, + /// A submitter took it for its look, and no other takes it. + looking: bool, + /// Its place among every poll its ring ever kept. + order: u64, +} + +impl Kept { + /// Something has taken the poll, and no submitter is looking at it. + fn owed(&self) -> bool { + !self.looking && !self.poll.armed() + } +} + +impl Polls { + pub const fn new() -> Self { + Self { polls: Vec::new(), kept: 0 } + } + + fn keep(&mut self, poll: Arc>) -> Kept { + self.kept += 1; + Kept { poll, looking: false, order: self.kept - 1 } + } + + /// Keep `poll` as its handle's one poll, unless `cap` are kept already. + /// Every earlier poll on the handle answers nothing from here on, whether + /// a post has fired it or not. + pub fn admit(&mut self, poll: Arc>, cap: usize) -> bool { + self.polls.retain(|kept| { + let other = kept.poll.handle != poll.handle; + if !other { + kept.poll.withdraw(); + } + other + }); + if self.polls.len() >= cap { + return false; + } + let kept = self.keep(poll); + self.polls.push(kept); + true + } + + /// The bound of one pass of looks: a poll kept from here on waits for the + /// next. + pub fn pass(&self) -> u64 { + self.kept + } + + /// The oldest poll owed a look among those kept before `pass`, kept while + /// it is looked at so a watch on its handle still replaces it. + pub fn take_owed(&mut self, pass: u64) -> Option>> { + let kept = self.polls.iter_mut().find(|kept| kept.order < pass && kept.owed())?; + kept.looking = true; + Some(kept.poll.clone()) + } + + /// Let go of `poll` once its look has an answer. `false` if a newer watch + /// replaced it during the look or the ring went: it answers nothing. + pub fn settle(&mut self, poll: &Arc>) -> bool { + let Some(at) = self.place(poll) else { return false }; + self.polls.remove(at); + true + } + + /// Keep `again` in `poll`'s place once its look found nothing; `false` as + /// [`Self::settle`] says it, and `again` is not kept. + pub fn renew(&mut self, poll: &Arc>, again: Arc>) -> bool { + let Some(at) = self.place(poll) else { return false }; + self.polls[at] = self.keep(again); + true + } + + /// Whether a poll is owed a look. + pub fn owed(&self) -> bool { + self.polls.iter().any(Kept::owed) + } + + fn place(&self, poll: &Arc>) -> Option { + self.polls.iter().position(|kept| Arc::ptr_eq(&kept.poll, poll)) + } + + /// The ring is going: no poll answers. + pub fn withdraw_all(&mut self) { + for kept in self.polls.drain(..) { + kept.poll.withdraw(); + } + } +} + +impl Default for Polls { + fn default() -> Self { + Self::new() + } +} + +/// What a look at a poll's object found. +pub enum Look { + /// Ready in the directions this word names, which is the answer. + Ready(u32), + /// The handle names nothing that could answer any more; the refusal is + /// the answer. + Refused(SyscallError), + /// Not ready: the poll that holds the handle's place now answers later, + /// and this one answers nothing. + Waits, +} + +/// The submitter's side of [`deliver`]; its [`Wake`] is the ring's own. +pub trait Submitter: Wake { + /// Whether the completion ring takes one more answer. + fn room(&self) -> bool; + /// Write one completion and wake nobody: [`complete`] is the one caller, + /// and owes the wake. + fn answer(&self, user_data: u64, result: i32); + /// Look at the object `poll` watches, and [`Polls::renew`] the poll if it + /// is not ready. + fn look(&self, poll: &Arc>) -> Look; + /// Run `f` on the ring's polls under their lock; `None` once the ring is + /// gone. + fn polls(&self, f: impl FnOnce(&mut Polls) -> R) -> Option; +} + +/// Write one completion, then wake every submitter parked on the ring: one +/// may wait for this count, and one may have parked over the poll it answers +/// while a look hid it. +pub fn complete(ring: &impl Submitter, user_data: u64, result: i32) { + ring.answer(user_data, result); + ring.wake(); +} + +/// Answer every poll owed a look, oldest first, while the ring has room; a +/// poll left over is answered by a later call, never dropped. One pass: a +/// poll a look renews waits for the next call, so a peer that posts an empty +/// object as fast as it is looked at holds nobody. +pub fn deliver(ring: &impl Submitter) { + let Some(pass) = ring.polls(|polls| polls.pass()) else { return }; + while ring.room() { + let Some(poll) = ring.polls(|polls| polls.take_owed(pass)).flatten() else { return }; + let result = if poll.state.ended() { + -(SyscallError::NotFound as i32) + } else { + match look(ring, &poll) { + Look::Ready(flags) => flags as i32, + Look::Refused(e) => -(e as i32), + Look::Waits => continue, + } + }; + if ring.polls(|polls| polls.settle(&poll)) == Some(true) { + complete(ring, poll.user_data, result); + } + } +} + +/// The park predicate of a submitter waiting until `enough()`, read once it is +/// registered on the watch a fire posts: it does not park over a poll +/// [`deliver`] would answer. +pub fn awake(ring: &impl Submitter, enough: impl FnOnce() -> bool) -> bool { + ring.room() && ring.polls(|polls| polls.owed()) == Some(true) || enough() +} + +/// `post-is-an-answer` is the negative control: a fired poll is answered with +/// its interest and nobody looks, and `kernel-loom`'s `inbox_answer` reds. +#[cfg(not(feature = "post-is-an-answer"))] +fn look(ring: &impl Submitter, poll: &Arc>) -> Look { + ring.look(poll) +} + +#[cfg(feature = "post-is-an-answer")] +fn look(_ring: &impl Submitter, poll: &Arc>) -> Look { + Look::Ready(poll.flags) +} diff --git a/kernel/src/object/ops.rs b/kernel/src/object/ops.rs index ac7edc38b85..8eb9d79f709 100644 --- a/kernel/src/object/ops.rs +++ b/kernel/src/object/ops.rs @@ -777,6 +777,21 @@ pub fn has_data(object: &KObjectRef) -> bool { } } +/// Whether a post on this object's read watch is its readability, which +/// [`has_data`] cannot be asked for: the log's, whose unread records are a +/// property of the reader's cursor, which the kernel does not hold; and a +/// console's, whose watch is the keyboard's while its data is the serial +/// line's (`issues/kernel/a-console-watch-waits-on-the-keyboard-not-the-serial-line.md`). +pub fn read_posts_are_readiness(object: &KObjectRef) -> bool { + match object { + KObjectRef::SysCap(_) | KObjectRef::Console(_) => true, + KObjectRef::PipeRead(_) | KObjectRef::PipeWrite(_) | KObjectRef::Connection(_) + | KObjectRef::Acceptor(_) | KObjectRef::File(_) | KObjectRef::Device(_) + | KObjectRef::Inbox(_) | KObjectRef::Connector(_) | KObjectRef::Namespace(_) + | KObjectRef::SharedMem(_) | KObjectRef::Process(_) => false, + } +} + pub fn has_space(object: &KObjectRef) -> bool { match object { KObjectRef::PipeWrite(w) => pipe::has_space(w.id()), diff --git a/kernel/src/watch.rs b/kernel/src/watch.rs index 8697bb067f2..5dc4a17cad4 100644 --- a/kernel/src/watch.rs +++ b/kernel/src/watch.rs @@ -1,6 +1,6 @@ //! The one way to wait: every waitable object holds exactly one [`Watch`], and //! a waiter on it is either a thread, which a post *wakes*, or a user poll -//! ring's entry, which a post *completes*. +//! ring's entry, which a post *fires*. //! //! The protocol is `toyos_sched::watch`'s and `toyos_sched::park`'s, and its //! lost-wake argument is theirs: a thread registers before it reads its @@ -42,14 +42,14 @@ pub struct Waitable>(toyos_sched::watch::Watch>; -/// A watch an interrupt handler posts, and the watch of a ring such a post -/// completes into. It has no `post`, which frees: [`IrqWatch::post_in_place`] +/// A watch an interrupt handler posts, and the watch of a ring whose poll such +/// a post fires. It has no `post`, which frees: [`IrqWatch::post_in_place`] /// frees nothing. pub type IrqWatch = Waitable>; -/// What an interrupt handler's post takes: an [`IrqWatch`]'s list and a poll -/// ring's completions. Held with interrupts off, so a handler never finds one -/// held by the context it interrupted; nothing allocates or frees under it. +/// What an interrupt handler's post takes: an [`IrqWatch`]'s list. Held with +/// interrupts off, so a handler never finds one held by the context it +/// interrupted; nothing allocates or frees under it. pub struct IrqLock(masked::Masked); impl IrqLock { @@ -95,7 +95,7 @@ impl Watch { } /// Something about the object changed: wake every thread waiting on it and - /// complete every poll. + /// fire every poll. pub fn post(&self) { self.post_as(WakeCause::new(WakeReason::Woken)); } @@ -142,7 +142,7 @@ impl IrqWatch { } /// Something about the object changed: wake every thread waiting on it and - /// complete every poll where it stands, freeing nothing. A handler makes it + /// fire every poll where it stands, freeing nothing. A handler makes it /// with its CPU's preempt count raised, as `device_irq_entry` holds it, so /// this post's own never reaches zero, and a pass, inside the interrupt. pub fn post_in_place(&self) { diff --git a/src/build.rs b/src/build.rs index 34429b4adf7..c97643b8640 100644 --- a/src/build.rs +++ b/src/build.rs @@ -2651,6 +2651,10 @@ mod tests { // turned on only by `kernel-loom`, to split `inbox/once.rs`'s // exchange and prove `poll_once` reds without it. "poll-fire-load-store", + // Costs no kernel build: turned on only by `kernel-loom`, so + // `inbox/polls.rs` answers a fired poll without a look and + // `inbox_answer` reds. + "post-is-an-answer", "reap-raise-relaxed", // `smp_roster.rs`'s count relaxed; `smp_bringup.rs` reds. "roster-commit-relaxed", diff --git a/src/ci.rs b/src/ci.rs index 4ff74000bfc..2538f7a4bb5 100644 --- a/src/ci.rs +++ b/src/ci.rs @@ -285,6 +285,11 @@ pub(crate) const CONTROLS: &[Control] = &[ red(KERNEL_LOOM, "log-ring-loads-swapped", Some("log_ring"), &[ Fails("a_published_record_is_whole_and_read_once"), ]), + red(KERNEL_LOOM, "post-is-an-answer", Some("inbox_answer"), &[ + Fails("a_post_with_nothing_to_read_answers_nothing"), + Fails("a_post_that_lands_after_its_bytes_were_read_answers_nothing"), + Fails("a_poll_armed_again_does_not_end_the_look"), + ]), red(KERNEL_LOOM, "poll-fire-load-store", Some("poll_once"), &[ Fails("a_post_and_a_recheck_answer_a_poll_once"), Fails("a_withdrawal_and_a_post_never_both_take_a_poll"), @@ -322,9 +327,8 @@ pub(crate) const CONTROLS: &[Control] = &[ message: "parked with the condition true and no wake owed: the post was lost", }, Says { - test: "two_posts_through_one_rings_lock_lose_no_wake", - message: "parked with both completions written and no wake owed: a ring's post was \ - lost", + test: "a_fire_racing_a_submitters_park_is_never_lost", + message: "parked over a fired poll and no wake owed: a fire was lost", }, ]), // The notify's flagged arm answering off a load: a second post reads the diff --git a/tests/toyos-rust-tests/src/bin/inbox_empty_write.rs b/tests/toyos-rust-tests/src/bin/inbox_empty_write.rs new file mode 100644 index 00000000000..29f433662a2 --- /dev/null +++ b/tests/toyos-rust-tests/src/bin/inbox_empty_write.rs @@ -0,0 +1,39 @@ +//! A watch is answered for what its object holds, never for a post. +//! +//! One poller watches two pipes, both empty and both armed. A write of no +//! bytes into the first posts its readers and leaves nothing to read; one byte +//! then goes into the second. The wait that follows answers the second pipe's +//! token alone: a reader handed the first pipe's would read it blocking and +//! park for good. + +use toyos::poller::{Poller, READABLE}; +use toyos_abi::syscall; + +const EMPTY: u64 = 1; +const FILLED: u64 = 2; + +fn main() { + let empty = syscall::pipe().expect("the pipe that stays empty"); + let filled = syscall::pipe().expect("the pipe that takes a byte"); + + let poller = Poller::new(2); + poller.watch_raw(empty.read, READABLE, EMPTY); + poller.watch_raw(filled.read, READABLE, FILLED); + // A non-blocking enter, so both polls are armed before either write: a + // write that found none would post nothing. + poller.wait(0, 0, |token| panic!("nothing is ready yet, got token {token}")); + + let byte = [0x5A]; + // A slice of a real buffer: the kernel is handed an address it can read. + assert_eq!(syscall::write(empty.write, &byte[..0]), Ok(0), "a write of no bytes"); + assert_eq!(syscall::write(filled.write, &byte), Ok(1), "a write of one byte"); + + let mut tokens = Vec::new(); + poller.wait(1, u64::MAX, |token| tokens.push(token)); + assert_eq!(tokens, [FILLED], "the wait answered a pipe with nothing to read"); + println!("inbox_empty_write: a write of no bytes answered no watch"); + + for end in [empty.read, empty.write, filled.read, filled.write] { + syscall::close(end); + } +} diff --git a/tests/toyos-rust-tests/src/bin/inbox_log_post.rs b/tests/toyos-rust-tests/src/bin/inbox_log_post.rs new file mode 100644 index 00000000000..529d02aab80 --- /dev/null +++ b/tests/toyos-rust-tests/src/bin/inbox_log_post.rs @@ -0,0 +1,77 @@ +//! A watch on the kernel's log is answered by the post a record makes. +//! +//! What the log holds for a reader is the reader's cursor's to say, which the +//! kernel does not hold, so no look can ask it: the post is the answer. A +//! reader arms its watch, a process ends and the kernel records it, and the +//! wait hands the reader its token; the record is then there to read. A watch +//! the post does not answer leaves the wait parked for good, and the harness's +//! deadline is what says so. + +use std::process::Command; + +use toyos::endow::{Endowments, SYSCAP_LABEL}; +use toyos::log::{LogTail, Record}; +use toyos::poller::{Poller, READABLE}; +use toyos::syscap::SysCap; + +const SELF_PATH: &str = "/system/bin/test_rs_inbox_log_post"; + +const LOG: u64 = 7; + +fn main() { + if std::env::args().nth(1).is_some() { + return; + } + let cap: SysCap = Endowments::get() + .take(SYSCAP_LABEL) + .expect("test-runner endows every binary it spawns a system capability"); + let poller = Poller::new(1); + let mut tail = LogTail::new(); + let mut ended: Option = None; + loop { + // Armed before the log is read, as every reader of an edge arms. An + // answer here is some other record's post: the watch is spent, and + // the round begins again. + poller.watch(&cap, READABLE, LOG); + let mut spent = false; + poller.wait(0, 0, |_| spent = true); + if recorded(&mut tail, &cap, ended) { + break; + } + if spent { + continue; + } + // Only with a watch armed and unanswered: the record is the post it + // waits for. + let pid = *ended.get_or_insert_with(end_a_process); + let mut tokens = Vec::new(); + poller.wait(1, u64::MAX, |token| tokens.push(token)); + assert_eq!(tokens, [LOG], "the wait after pid {pid}'s end was recorded"); + } + println!("inbox_log_post: a record's post answered the log's watch"); +} + +/// Run a process to its end, which the kernel records, and answer its pid. +fn end_a_process() -> u32 { + let mut child = Command::new(SELF_PATH).arg("end").spawn().expect("spawn a process to end"); + let pid = child.id(); + assert!(child.wait().expect("wait for it").success(), "the process that only ends"); + pid +} + +/// Read every record the tail has not seen; whether one of them is the end of +/// `pid`. +fn recorded(tail: &mut LogTail, cap: &SysCap, pid: Option) -> bool { + let end = pid.map(|pid| format!(" pid={pid} code=0 ")); + let mut records = vec![Record::EMPTY; 64]; + let mut found = false; + loop { + let read = tail.read(cap, &mut records).expect("read the kernel's records"); + if read.is_empty() { + return found; + } + if let Some(end) = &end { + found |= read.iter().any(|r| r.message().starts_with("exit: ") && r.message().contains(end)); + } + } +} diff --git a/tests/toyos-rust-tests/src/bin/kill_ends_every_wait.rs b/tests/toyos-rust-tests/src/bin/kill_ends_every_wait.rs index 3417cb5ac7b..13cd5d9c2d9 100644 --- a/tests/toyos-rust-tests/src/bin/kill_ends_every_wait.rs +++ b/tests/toyos-rust-tests/src/bin/kill_ends_every_wait.rs @@ -7,20 +7,29 @@ //! cannot end keeps the child's last thread in its process for ever, and //! `wait` below never returns; the harness's deadline is what says so. //! `mutual_kill` holds a kill inside a kill. +//! +//! `posted-poll` is the poll ring's wait a peer keeps from parking: this +//! process's threads write no bytes into the pipe the child watches, so every +//! post sends the child's wait round to look again. A kill those posts hold is +//! held for as long as they win a race and no longer, which the harness's +//! deadline cannot see: the posts stop at a ceiling of their own, [`HELD`], +//! and the child has to have ended before it. use std::io::{Read, Write}; use std::os::toyos::process::{ChildExt, CommandExt}; use std::process::{Child, Command, Stdio}; -use std::sync::atomic::AtomicU32; -use std::sync::OnceLock; +use std::sync::atomic::{AtomicU32, AtomicU64, Ordering}; +use std::sync::{Arc, Barrier, OnceLock}; +use std::thread::JoinHandle; use std::time::Duration; -use toyos::endow::{Endowments, SYSCAP_LABEL}; -use toyos::poller::Poller; +use toyos::endow::{Endowments, FromHandle, SYSCAP_LABEL}; +use toyos::poller::{Poller, READABLE}; use toyos::process::Process; use toyos::syscap::SysCap; use toyos::AsHandle; -use toyos_abi::syscall; +use toyos_abi::clock; +use toyos_abi::syscall::{self, SyscallError}; use toyos_abi::RawHandle; const SELF_PATH: &str = "/system/bin/test_rs_kill_ends_every_wait"; @@ -28,6 +37,25 @@ const SELF_PATH: &str = "/system/bin/test_rs_kill_ends_every_wait"; /// The label the process-wait child finds the process it waits on under. const WAITED: &str = "waited"; +/// The label the posted-poll child finds the pipe it watches under. +const POSTED: &str = "posted"; + +/// That pipe's read end, which the child holds until it ends. +struct Posted(RawHandle); + +impl FromHandle for Posted { + unsafe fn from_handle(raw: RawHandle) -> Self { + Self(raw) + } +} + +/// The threads that post the pipe the posted-poll child watches. +const POSTERS: usize = 4; + +/// How long those threads go on posting after the kill before they say the +/// child outlived it. +const HELD: Duration = Duration::from_secs(1); + /// `process::KILLED_EXIT_CODE`. const KILLED: i32 = 137; @@ -41,7 +69,7 @@ const ZOMBIE: u8 = 3; // sleep underneath (the waited process's `nanosleep`, the joined thread's // `std::thread::sleep`), so a mutation that breaks sleep would otherwise surface // under one of their names instead of its own. -const WAITS: [&str; 5] = ["futex", "poll", "sleep", "process-wait", "thread-join"]; +const WAITS: [&str; 6] = ["futex", "poll", "posted-poll", "sleep", "process-wait", "thread-join"]; fn main() { match std::env::args().nth(1).as_deref() { @@ -54,13 +82,22 @@ fn test() { for role in WAITS { // The process the process-wait child waits on, which nothing ends but this. let mut waited = (role == "process-wait").then(|| spawn("sleep", None)); - let endow = waited.as_ref().map(|w| { - let dup = syscall::dup(RawHandle(w.as_raw_handle())).expect("a handle to endow"); - (WAITED.to_string(), dup.0) - }); + // The pipe the posted-poll child watches, which no byte ever enters. + let posted = (role == "posted-poll").then(|| syscall::pipe().expect("a pipe to post")); + let endow = waited + .as_ref() + .map(|w| { + let dup = syscall::dup(RawHandle(w.as_raw_handle())).expect("a handle to endow"); + (WAITED.to_string(), dup.0) + }) + .or(posted.as_ref().map(|pipe| (POSTED.to_string(), pipe.read.0))); let mut child = spawn(role, endow); + let posts = posted.map(|pipe| Posts::hold(child.id(), pipe.write)); println!(" {role}: killing"); child.kill().expect("kill the parked child"); + if let Some(posts) = posts { + posts.until_the_child_ends(); + } let code = child.wait().expect("wait for the killed child").code(); assert_eq!(code, Some(KILLED), "a child killed in its {role} wait ended with {code:?}"); if let Some(mut waited) = waited.take() { @@ -105,6 +142,69 @@ fn spawn(role: &str, endow: Option<(String, u32)>) -> Child { child } +/// The threads posting the pipe the posted-poll child watches. +struct Posts { + /// Each answers whether the pipe lost its reader before the ceiling. + threads: Vec>, + /// When the threads stop, in nanoseconds since boot: never, until the kill. + stop_at: Arc, + write: RawHandle, +} + +impl Posts { + /// Start [`POSTERS`] threads, each writing no bytes into `write` until the + /// pipe has no reader, and return once the roster shows `pid`'s main + /// thread out of the park the posts woke it from. + fn hold(pid: u32, write: RawHandle) -> Self { + let stop_at = Arc::new(AtomicU64::new(u64::MAX)); + let posting = Arc::new(Barrier::new(POSTERS + 1)); + let threads = (0..POSTERS) + .map(|_| { + let (stop_at, posting) = (stop_at.clone(), posting.clone()); + std::thread::spawn(move || { + let byte = [0u8]; + posting.wait(); + while clock::nanos_since_boot() < stop_at.load(Ordering::Relaxed) { + // A slice of a real buffer: the kernel is handed an address it can read. + match syscall::write(write, &byte[..0]) { + Ok(0) => {} + // The child ended, and the pipe's one read end with it. + Err(SyscallError::Gone) => return true, + other => panic!("a write of no bytes answered {other:?}"), + } + } + false + }) + }) + .collect(); + // Every thread is posting before the kill: one alone loses the race + // that holds the child. + posting.wait(); + println!(" posted-poll: waiting for the roster to show it looking"); + loop { + match main_thread_status(pid) { + RosterStatus::NotParked => break, + RosterStatus::Parked => std::thread::sleep(Duration::from_millis(10)), + RosterStatus::Gone => panic!("posted-poll: a write of no bytes ended the child's wait"), + } + } + Self { threads, stop_at, write } + } + + /// After the kill: the posts go on until the child ends, and it has to + /// end within [`HELD`] of them. + fn until_the_child_ends(self) { + self.stop_at.store(clock::nanos_since_boot() + HELD.as_nanos() as u64, Ordering::Relaxed); + for thread in self.threads { + assert!( + thread.join().expect("a posting thread"), + "posted-poll: the child still watched its pipe {HELD:?} after its kill: the posts held it in its wait" + ); + } + syscall::close(self.write); + } +} + /// What the roster says about `pid`'s main thread. enum RosterStatus { /// In the roster, blocked: parked in the wait under test. @@ -157,6 +257,20 @@ fn child(role: &str) -> ! { say(role); poller.wait(1, u64::MAX, |_| {}); } + "posted-poll" => { + let Posted(read) = Endowments::get().take(POSTED).expect("the parent endowed a pipe"); + let poller = Poller::new(Poller::MAX_HANDLES); + // As many polls as a poller holds, one handle each: a post fires + // them all, and the wait parks only if no post lands while it + // looks at every one of them. + poller.watch_raw(read, READABLE, 0); + for token in 1..u64::from(Poller::MAX_HANDLES) { + let dup = syscall::dup(read).expect("another handle to the pipe"); + poller.watch_raw(dup, READABLE, token); + } + say(role); + poller.wait(1, u64::MAX, |_| {}); + } "process-wait" => { let waited: Process = Endowments::get().take(WAITED).expect("the parent endowed a process"); say(role); diff --git a/tests/toyos-rust-tests/src/netd_stream.rs b/tests/toyos-rust-tests/src/netd_stream.rs index 996fdd6c319..748c290d656 100644 --- a/tests/toyos-rust-tests/src/netd_stream.rs +++ b/tests/toyos-rust-tests/src/netd_stream.rs @@ -50,10 +50,6 @@ pub fn ring_capacity() -> u64 { /// Wait, with no deadline, until `check` answers, re-asking it each time /// `handle` reports ready for `flags`: an answer that never comes is a hang the /// harness ceiling reds. -/// -/// **A readiness completion is a reason to look again, not an answer**: a -/// zero-byte write still wakes the other end's watch, and netd's liveness -/// probes are zero-byte writes. pub fn await_until( handle: &impl AsHandle, flags: u32, diff --git a/toyos-abi/src/inbox.rs b/toyos-abi/src/inbox.rs index c9f85adf054..c454cef8a50 100644 --- a/toyos-abi/src/inbox.rs +++ b/toyos-abi/src/inbox.rs @@ -10,6 +10,15 @@ use crate::RawHandle; pub const OP_NOP: u8 = 0; +/// Answered once, when the handle is ready. +/// +/// **An answer is what the object held when the waiting `inbox_submit` looked +/// at it.** A post only owes the watch a look, so neither a post with nothing +/// behind it nor one for bytes already read is answered, and a handle no other +/// reader drains holds what its answer says. The exceptions are the objects +/// whose readiness the kernel does not hold — the log, read on the reader's own +/// cursor, and a console — whose read post is the answer. A watch replaces its +/// handle's earlier one, fired or not, so one look answers a handle once. pub const OP_WATCH: u8 = 1; // Op code 2 unused (formerly IORING_OP_POLL_REMOVE): a watch this kernel takes // is one-shot, consumed by the completion it posts, so the interest a remove @@ -78,12 +87,8 @@ pub struct RingHeader { /// Completions the kernel could not post because the completion ring /// reported itself full. Cumulative, and never cleared. /// - /// The 2x sizing makes this unreachable only for a process that keeps its - /// registrations within the depth it asked for: over-registering flushes a - /// full submission ring mid-registration, and the kernel then posts - /// completions for the handles already ready while the caller is still - /// registering the rest. `toyos`'s `Poller` sizes its rings so that cannot - /// happen and reads this on every wait. + /// `toyos`'s `Poller` sizes its rings so that cannot happen and reads this + /// on every wait. pub dropped: core::sync::atomic::AtomicU32, } diff --git a/toyos-libc-copies/src/lib.rs b/toyos-libc-copies/src/lib.rs index fd0393b3999..0de1d30f573 100644 --- a/toyos-libc-copies/src/lib.rs +++ b/toyos-libc-copies/src/lib.rs @@ -34,6 +34,9 @@ mod listing; #[path = "../../userland/libc/src/memreq.rs"] mod memreq; #[cfg(test)] +#[path = "../../userland/libc/src/pollreq.rs"] +mod pollreq; +#[cfg(test)] #[path = "../../userland/libc/src/sigmask.rs"] mod sigmask; #[cfg(test)] @@ -67,6 +70,8 @@ mod long_double; #[cfg(test)] mod memory_refusals; #[cfg(test)] +mod poll_requests; +#[cfg(test)] mod prototypes; #[cfg(test)] mod signal_masks; diff --git a/toyos-libc-copies/src/poll_requests.rs b/toyos-libc-copies/src/poll_requests.rs new file mode 100644 index 00000000000..ec1881ae069 --- /dev/null +++ b/toyos-libc-copies/src/poll_requests.rs @@ -0,0 +1,27 @@ +//! What `poll` asks a ring to watch: one watch per descriptor, which answers +//! every entry that names it. + +use toyos_abi::inbox::{READABLE, WRITABLE}; + +use crate::header; +use crate::pollreq::{watch_of, watches}; + +const POLL_H: &str = include_str!("../../userland/libc/include/poll.h"); + +fn events(names: &[&str]) -> i16 { + names.iter().fold(0, |events, name| events | header::int(POLL_H, name) as i16) +} + +#[test] +fn a_descriptor_named_twice_is_watched_once_for_both_interests() { + let entries = [(5, events(&["POLLIN"])), (7, events(&["POLLOUT"])), (5, events(&["POLLOUT"]))]; + assert_eq!(watches(&entries), [(0, 5, READABLE | WRITABLE), (1, 7, WRITABLE)]); + assert_eq!([0, 1, 2].map(|entry| watch_of(&entries, entry)), [0, 1, 0]); +} + +#[test] +fn each_descriptor_named_once_has_its_own_watch() { + let entries = [(3, events(&["POLLIN", "POLLOUT"])), (4, events(&["POLLIN"])), (9, events(&["POLLHUP"]))]; + assert_eq!(watches(&entries), [(0, 3, READABLE | WRITABLE), (1, 4, READABLE), (2, 9, 0)]); + assert_eq!([0, 1, 2].map(|entry| watch_of(&entries, entry)), [0, 1, 2]); +} diff --git a/toyos-sched/loom/Cargo.toml b/toyos-sched/loom/Cargo.toml index 1c9dcfe234c..5bec98dce47 100644 --- a/toyos-sched/loom/Cargo.toml +++ b/toyos-sched/loom/Cargo.toml @@ -114,6 +114,11 @@ fault-posted-before-it-is-set = [] [dependencies] loom = "0.7" +# `loom_watch.rs` compiles `kernel/src/inbox/polls.rs`, which names a handle and +# a refusal. +[dev-dependencies] +toyos-abi = { path = "../../toyos-abi" } + [lib] # The primitives are compiled into the lib; the models are the test targets. path = "src/lib.rs" @@ -121,3 +126,9 @@ path = "src/lib.rs" # and run there, against real atomics. Loom atomics outside a `loom::model` # panic, so this crate exposes no unit tests of its own. test = false + +# `kernel-loom`'s control for a ring's answer, which `polls.rs` carries and no +# model here turns on. A `cfg` name and not a feature, since a feature here is +# a model control. +[lints.rust] +unexpected_cfgs = { level = "warn", check-cfg = ['cfg(feature, values("post-is-an-answer"))'] } diff --git a/toyos-sched/loom/tests/loom_watch.rs b/toyos-sched/loom/tests/loom_watch.rs index 6c57e61b7f0..042924459b2 100644 --- a/toyos-sched/loom/tests/loom_watch.rs +++ b/toyos-sched/loom/tests/loom_watch.rs @@ -30,10 +30,12 @@ //! models must red with a poll completed by neither. //! //! **The ring entry is the kernel's [`Once`], compiled from -//! `kernel/src/inbox/once.rs`**, the decision a `PollEntry` makes; what else a -//! `PollEntry` is — the ring's page, its lock, the completion it writes — names -//! half the kernel and cannot be compiled here, so the model's entry counts its -//! answers instead of writing them. +//! `kernel/src/inbox/once.rs`**, the decision a `PollEntry` makes, and the last +//! two models' ring is the kernel's `kernel/src/inbox/polls.rs` whole: its +//! polls, the submitter's `deliver`, its park predicate and the wake an answer +//! owes. What else a ring is — its +//! page, the objects its looks read — names half the kernel and cannot be +//! compiled here, so the models count answers instead of writing them. //! //! [`TaskShared::notify`]: toyos_sched_loom::task::TaskShared::notify //! [`Gate`]: toyos_sched_loom::watch::Gate @@ -41,9 +43,12 @@ //! [`prepare`]: toyos_sched_loom::park::prepare //! [`Ring::fire`]: toyos_sched_loom::watch::Ring::fire +extern crate alloc; + use loom::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use loom::sync::Arc; use toyos_sched_loom::cpu::{CpuHandle, CpuHandles}; +use toyos_sched_loom::hw::CpuId; use toyos_sched_loom::mailbox::{mailbox, MailboxConsumer}; use toyos_sched_loom::model::{ model, watch_list, Kicks, LoomLock, Msg, PreemptModel, RemoteGuard, CPU0, CPU1, @@ -58,12 +63,19 @@ use toyos_sched_loom::watch::{Fire, Gate, Poster, Ring, Waiters, Watch}; #[path = "../../../kernel/src/inbox/once.rs"] mod once; +/// Names the one-shot above as `super::once`. +#[expect(dead_code, reason = "every look here finds its object ready: a renewal and a refusal are `kernel-loom`'s `inbox_answer`'s")] +#[path = "../../../kernel/src/inbox/polls.rs"] +mod polls; + /// One poll, as the kernel's is: its answer taken once, by the kernel's own /// [`once::Once`], across everything that may fire it; it counts what it /// posted. struct Poll { state: once::Once, posts: AtomicU32, + /// Whether the end of its source is what took it. + ended: AtomicBool, } #[derive(Clone)] @@ -74,6 +86,7 @@ impl Entry { Self(Arc::new(Poll { state: once::Once::new(), posts: AtomicU32::new(0), + ended: AtomicBool::new(false), })) } @@ -83,8 +96,14 @@ impl Entry { } impl Ring for Entry { - fn fire(&self, _how: Fire) { - if self.0.state.fire() { + // As the kernel's `PollEntry`: readiness fires the poll, an end ends it. + fn fire(&self, how: Fire) { + let took = match how { + Fire::Ready => self.0.state.fire(), + Fire::Gone => self.0.state.end(), + }; + if took { + self.0.ended.store(how == Fire::Gone, Ordering::Release); self.0.posts.fetch_add(1, Ordering::AcqRel); } } @@ -369,6 +388,11 @@ fn end_racing(end: fn(&World), post: fn(&World)) { assert_eq!(poll.posts(), 1); assert!(!poll.live()); + assert_eq!( + poll.0.state.ended(), + poll.0.ended.load(Ordering::Acquire), + "the one-shot names the wrong taker" + ); drop(world); } @@ -644,100 +668,180 @@ fn a_poll_on_two_watches_racing_both_posts_completes_exactly_once() { }); } -/// A poll ring as the kernel's is: its completions behind a lock of its own, -/// and a watch its submitter parks on, which holds threads and no ring. +/// A poll ring as its fires and its submitter see it: the kernel's polls, the +/// answers its submitter wrote, and the watch that submitter parks on, which +/// holds threads and no ring. struct PollRing { cpus: CpuHandles, kicks: Kicks, - written: LoomLock, + polls: LoomLock>, + answers: LoomLock, parked: RingWatch, } -/// One of that ring's polls, registered on one device's watch. -struct RingPoll { - ring: Arc, - state: once::Once, +#[derive(Clone)] +struct RingRef(Arc); + +type RingPoll = alloc::sync::Arc>; + +impl polls::Wake for RingRef { + fn wake(&self) { + let ring = &self.0; + let env = Poster { cpus: &ring.cpus, kicker: &ring.kicks, preempt: &RemoteGuard }; + ring.parked.post_in_place(WakeCause::new(WakeReason::Woken), &env); + } +} + +/// Every look finds its object ready: each device posts once, for good. +impl polls::Submitter for RingRef { + fn room(&self) -> bool { + true + } + + fn answer(&self, _user_data: u64, _result: i32) { + self.0.answers.with(|n| *n += 1); + } + + fn look(&self, poll: &RingPoll) -> polls::Look { + polls::Look::Ready(poll.flags) + } + + fn polls(&self, f: impl FnOnce(&mut polls::Polls) -> R) -> Option { + Some(self.0.polls.with(f)) + } } -struct RingEntry(Arc); +/// One of that ring's polls, as one device's watch holds it. +struct RingEntry(RingPoll); type RingWatch = Watch>>; impl Ring for RingEntry { - fn fire(&self, _how: Fire) { - if self.0.state.fire() { - let ring = &self.0.ring; - ring.written.with(|n| *n += 1); - let env = Poster { cpus: &ring.cpus, kicker: &ring.kicks, preempt: &RemoteGuard }; - ring.parked.post_in_place(WakeCause::new(WakeReason::Woken), &env); + // As the kernel's `PollEntry`. + fn fire(&self, how: Fire) { + match how { + Fire::Ready => self.0.fire(self.0.flags), + Fire::Gone => self.0.end(), } } fn live(&self) -> bool { - self.0.state.armed() + self.0.armed() } } -/// **A post in place fires its rings under its own list lock**, so beneath it -/// are the ring's lock and the ring's watch. Two devices' watches each hold a -/// poll of one ring and are posted at once, while the ring's submitter waits -/// for both completions: a submitter parked with both completions written was -/// owed the wake the second one posted. +/// `inbox::submit`'s wait for `want` answers on `cpu`, as far as its first +/// park: whether it parked. A park leaves the submitter registered, as a +/// thread blocked in `watch::wait_until` is. `siblings` is how many other +/// submitters answer into the ring. +fn submit_parks( + ring: &RingRef, + submitter: &Arc>, + cpu: CpuId, + want: u32, + siblings: u32, +) -> bool { + let enough = || ring.0.answers.with(|n| *n) >= want; + // A fire ends at most one iteration short of the last, and so does a + // sibling's look at a poll this one saw fired. + let rounds = want + siblings; + for _ in 0..=rounds { + polls::deliver(ring); + if enough() { + return false; + } + // `watch::wait_until` over `submit`'s predicate. + if polls::awake(ring, enough) { + continue; + } + ring.0.parked.register(submitter, 0); + while !polls::awake(ring, enough) { + let Ok(ticket) = + prepare(&CurrentTask::new(submitter, cpu), Cancel::Answers, WaitClass::Io) + else { + continue; + }; + match ticket.commit() { + Commit::Parked(_) => return true, + Commit::AlreadyWoken => continue, + Commit::Killed => unreachable!("nothing retires in this model"), + } + } + ring.0.parked.unregister(submitter); + } + unreachable!("{rounds} fires and looks ended more than {rounds} iterations") +} + +/// A ring nobody waits on yet, and the mailbox of each CPU a submitter of it +/// runs on. +fn poll_ring(cpus: [CpuId; CPUS]) -> (RingRef, [MailboxConsumer; CPUS]) { + let (handles, mailboxes): (Vec<_>, Vec<_>) = cpus + .into_iter() + .map(|cpu| { + let (tx, rx) = mailbox::(); + (CpuHandle::new(cpu, tx), rx) + }) + .unzip(); + let ring = RingRef(Arc::new(PollRing { + cpus: CpuHandles::new(handles), + kicks: Kicks::new(), + polls: LoomLock::new(polls::Polls::new()), + answers: LoomLock::new(0), + parked: Watch::new(watch_list()), + })); + let Ok(mailboxes) = mailboxes.try_into() else { + unreachable!("one mailbox per CPU"); + }; + (ring, mailboxes) +} + +/// `devices` watches, each holding one armed poll of `ring`. +fn devices_polled(ring: &RingRef, devices: usize) -> Vec> { + (0..devices) + .map(|handle| { + let poll = alloc::sync::Arc::new(polls::Poll::new( + ring.clone(), + handle as u64, + toyos_abi::handle::RawHandle(handle as u32), + toyos_abi::inbox::READABLE, + )); + assert!(ring.0.polls.with(|polls| polls.admit(poll.clone(), devices))); + let device = Arc::new(Watch::new(watch_list())); + device.add_ring(RingEntry(poll)); + device + }) + .collect() +} + +/// An interrupt handler's post of `device`, on its own thread. +fn posted_in_place(ring: &RingRef, device: &Arc) -> loom::thread::JoinHandle<()> { + let (device, ring) = (device.clone(), ring.clone()); + loom::thread::spawn(move || { + let env = Poster { cpus: &ring.0.cpus, kicker: &ring.0.kicks, preempt: &RemoteGuard }; + device.post_in_place(WakeCause::new(WakeReason::Woken), &env); + }) +} + +/// **A submitter parks on its polls, and a fire posts the watch it parks on.** +/// Two devices' watches each hold a poll of one ring and are posted at once, +/// from an interrupt handler's post in place, while the ring's submitter waits +/// for both answers. A fire takes its poll and posts the ring's watch; the +/// submitter answers what was fired and, registered, reads its polls again +/// before it parks. A submitter parked over a fired poll was owed the wake +/// that fire posted. #[test] -fn two_posts_through_one_rings_lock_lose_no_wake() { +fn a_fire_racing_a_submitters_park_is_never_lost() { model(|| { - let (tx, mut rx) = mailbox::(); - let ring = Arc::new(PollRing { - cpus: CpuHandles::new(vec![CpuHandle::new(CPU0, tx)]), - kicks: Kicks::new(), - written: LoomLock::new(0), - parked: Watch::new(watch_list()), - }); - let devices: Vec> = - (0..2).map(|_| Arc::new(Watch::new(watch_list()))).collect(); - for device in &devices { - device.add_ring(RingEntry(Arc::new(RingPoll { - ring: ring.clone(), - state: once::Once::new(), - }))); - } + let (ring, [mut rx]) = poll_ring([CPU0]); + let devices = devices_polled(&ring, 2); let submitter = task(1); let waiting = { let ring = ring.clone(); let submitter = submitter.clone(); - loom::thread::spawn(move || { - ring.parked.register(&submitter, 0); - // Each post ends at most one iteration. - for _ in 0..4 { - if ring.written.with(|n| *n) == 2 { - return false; - } - let Ok(ticket) = - prepare(&CurrentTask::new(&submitter, CPU0), Cancel::Answers, WaitClass::Io) - else { - continue; - }; - match ticket.commit() { - Commit::Parked(_) => return true, - Commit::AlreadyWoken => continue, - Commit::Killed => unreachable!("nothing retires in this model"), - } - } - unreachable!("two posts ended more than three iterations") - }) + loom::thread::spawn(move || submit_parks(&ring, &submitter, CPU0, 2, 0)) }; - let posters: Vec<_> = devices - .iter() - .map(|device| { - let (device, ring) = (device.clone(), ring.clone()); - loom::thread::spawn(move || { - let env = - Poster { cpus: &ring.cpus, kicker: &ring.kicks, preempt: &RemoteGuard }; - device.post_in_place(WakeCause::new(WakeReason::Woken), &env); - }) - }) - .collect(); + let posters: Vec<_> = devices.iter().map(|device| posted_in_place(&ring, device)).collect(); let parked = waiting.join().unwrap(); for poster in posters { @@ -750,11 +854,64 @@ fn two_posts_through_one_rings_lock_lose_no_wake() { assert_eq!( msgs, [Msg::Wake(TaskKey(1), WakeReason::Woken)], - "parked with both completions written and no wake owed: a ring's post was lost", + "parked over a fired poll and no wake owed: a fire was lost", ); } else { + assert_eq!(ring.0.answers.with(|n| *n), 2, "the wait ended short of its answers"); assert!(msgs.is_empty(), "a submitter that never parked is owed nothing: {msgs:?}"); } - ring.parked.unregister(&submitter); + ring.0.parked.unregister(&submitter); + // The ring's teardown: its polls hold it. + ring.0.polls.with(polls::Polls::withdraw_all); + }); +} + +/// **A submitter parked over a poll its sibling is looking at is owed the +/// answer's wake.** One device's post fires the one poll of a ring two +/// submitters wait on, each for one answer. A look hides its poll from the +/// other submitter, which reads nothing owed and parks, the fire's wake spent +/// before it registered; the looker's answer then posts the watch both park +/// on. A submitter parked with the answer written and no wake owed was lost +/// to that look. +#[test] +fn an_answer_wakes_the_submitter_its_look_hid_the_poll_from() { + model(|| { + let (ring, mut mailboxes) = poll_ring([CPU0, CPU1]); + let devices = devices_polled(&ring, 1); + let submitters = [ + (task(1), CPU0), + (Arc::new(TaskShared::new(TaskKey(2), TaskState::Running(CPU1))), CPU1), + ]; + + let waiting: Vec<_> = submitters + .iter() + .map(|(submitter, cpu)| { + let (ring, submitter, cpu) = (ring.clone(), submitter.clone(), *cpu); + loom::thread::spawn(move || submit_parks(&ring, &submitter, cpu, 1, 1)) + }) + .collect(); + let poster = posted_in_place(&ring, &devices[0]); + + let parked: Vec = waiting.into_iter().map(|wait| wait.join().unwrap()).collect(); + poster.join().unwrap(); + let guard = PreemptModel::new(); + + let answers = ring.0.answers.with(|n| *n); + for (((submitter, _), rx), parked) in submitters.iter().zip(&mut mailboxes).zip(parked) { + let msgs = drain(rx, &guard); + if parked { + assert_eq!( + msgs, + [Msg::Wake(submitter.key(), WakeReason::Woken)], + "parked with {answers} answer(s) written and no wake owed: the answer woke nobody", + ); + } else { + assert_eq!(answers, 1, "the wait ended short of its answer"); + assert!(msgs.is_empty(), "a submitter that never parked is owed nothing: {msgs:?}"); + } + ring.0.parked.unregister(submitter); + } + // The ring's teardown: its polls hold it. + ring.0.polls.with(polls::Polls::withdraw_all); }); } diff --git a/toyos-sched/src/watch.rs b/toyos-sched/src/watch.rs index 1b9fcb795e0..8a584beaf36 100644 --- a/toyos-sched/src/watch.rs +++ b/toyos-sched/src/watch.rs @@ -6,7 +6,7 @@ //! claims its word if it is parked or committing and flags it otherwise, so its //! own next commit rechecks instead of parking. A **ring** entry is one poll a //! process submitted and is *posted*: a post hands it to [`Ring::fire`], once, -//! and the environment writes that poll's completion into the ring it names. +//! and the environment owes the ring it names a look at that poll's object. //! //! **A lost wake has no expression here.** A thread registers before it reads //! its condition, under the list lock a post also takes, so a post either @@ -54,14 +54,14 @@ pub enum Fire { /// One poll a ring is waiting on, as the watch holds it. pub trait Ring { - /// Post this poll's completion. One-shot across every watch the poll is + /// Fire this poll for its ring. One-shot across every watch the poll is /// registered on: an entry that already fired, or whose poll was withdrawn, - /// posts nothing. Called with at most the posting watch's list lock held, + /// does nothing. Called with at most the posting watch's list lock held, /// and from an interrupt handler by a post in place: may take only its /// ring's own lock and post only the watch its ring's submitters park on, /// and allocates and frees nothing. fn fire(&self, how: Fire); - /// Whether a fire would still post anything. `false` is permanent. + /// Whether a fire would still do anything. `false` is permanent. fn live(&self) -> bool; } diff --git a/toyos/src/poller.rs b/toyos/src/poller.rs index bbf48209d4f..f7df1ec4e35 100644 --- a/toyos/src/poller.rs +++ b/toyos/src/poller.rs @@ -179,9 +179,7 @@ impl Rings { /// and sizes both rings from it: the submission ring holds them all, so no /// batch is ever flushed mid-registration, and the kernel's completion ring — /// always twice the submission ring — holds the most completions that can exist -/// between two [`wait`](Self::wait) calls, which is two per watched handle (a -/// registration left over from the previous round firing, and this round's -/// registration finding the handle ready). +/// between two [`wait`](Self::wait) calls. /// /// Going past the capacity is a contract violation and panics, because it is /// the caller's own bug and the alternative is the failure this replaced: the diff --git a/userland/fsd/src/main.rs b/userland/fsd/src/main.rs index 022e2e2da8f..90dc0ec6b68 100644 --- a/userland/fsd/src/main.rs +++ b/userland/fsd/src/main.rs @@ -154,7 +154,6 @@ fn main() { caps.iter().map(|c| c.dir.as_str()).collect::>().join(", "), volume.describe() ); - let caps_len = caps.len() as u32; Server { volume, caps, @@ -163,7 +162,6 @@ fn main() { streams: BTreeMap::new(), next_stream: 0, writeback: WriteBack::default(), - probe: Poller::new(caps_len), scratch: Vec::new(), } .serve() @@ -318,9 +316,6 @@ struct Server { next_stream: u64, /// When the sync nobody asked for is due. writeback: WriteBack, - /// Asks an acceptor whether a connection waits, before [`Server::accept`] - /// takes it. - probe: Poller, /// A write's bytes, copied out of the client's window or a stream's pipe /// into this process's own memory before the volume sees them. Kept, so a /// write allocates nothing. @@ -423,23 +418,7 @@ impl Server { } } - /// Take a connection that waits on `cap`'s port, if one does. - /// - /// **Asked first, because `accept` parks and a completion is a hint**: a - /// watch replaced while a connection arrives can answer beside the watch - /// that replaced it, so two completions name one connection, and the - /// second `accept` would park this server for good. This process is the - /// port's one acceptor, so a connection the probe sees is still there to - /// take. The probe's ring is drained whole each time and a completion is - /// read by its port's token, so what counts is an arrival on this port - /// since its last probe, none of them taken. fn accept(&mut self, cap: usize) { - self.probe.watch(&self.caps[cap].acceptor, READABLE, cap as u64); - let mut waiting = false; - self.probe.wait(0, 0, |token| waiting |= token == cap as u64); - if !waiting { - return; - } let conn = match self.caps[cap].acceptor.accept() { Ok(conn) => conn, Err(why) => panic!("fsd: its own acceptor refused an accept: {why:?}"), diff --git a/userland/libc/src/lib.rs b/userland/libc/src/lib.rs index 8b184115473..77da705e7b9 100644 --- a/userland/libc/src/lib.rs +++ b/userland/libc/src/lib.rs @@ -18,6 +18,7 @@ mod math; mod memory; mod memreq; mod misc; +mod pollreq; mod posix_io; mod printf; mod pthread; diff --git a/userland/libc/src/pollreq.rs b/userland/libc/src/pollreq.rs new file mode 100644 index 00000000000..dc2c7dc5b56 --- /dev/null +++ b/userland/libc/src/pollreq.rs @@ -0,0 +1,38 @@ +//! What `poll` asks a ring to watch and which watch answers each entry. It +//! reads nothing but what it is handed, so the host tests it +//! (`toyos-libc-copies`). + +use alloc::vec::Vec; + +use toyos_abi::inbox::{READABLE, WRITABLE}; + +const POLLIN: i16 = 1; +const POLLOUT: i16 = 4; + +/// The entry whose watch answers for `entries[entry]`, each a descriptor and +/// its `events`: the first that names its descriptor. +pub(crate) fn watch_of(entries: &[(i32, i16)], entry: usize) -> usize { + let (fd, _) = entries[entry]; + entries[..entry].iter().position(|&(other, _)| other == fd).unwrap_or(entry) +} + +/// The watches a `poll` of `entries` submits, each the entry it is answered +/// under, its descriptor and its interest. One per descriptor, for every +/// entry's interest: a watch replaces its handle's earlier one, so a second on +/// the descriptor would leave the first entry unanswered. +pub(crate) fn watches(entries: &[(i32, i16)]) -> Vec<(usize, i32, u32)> { + let mut interest = alloc::vec![0u32; entries.len()]; + for (entry, &(_, events)) in entries.iter().enumerate() { + let watch = watch_of(entries, entry); + if events & POLLIN != 0 { + interest[watch] |= READABLE; + } + if events & POLLOUT != 0 { + interest[watch] |= WRITABLE; + } + } + (0..entries.len()) + .filter(|&entry| watch_of(entries, entry) == entry) + .map(|entry| (entry, entries[entry].0, interest[entry])) + .collect() +} diff --git a/userland/libc/src/posix_io.rs b/userland/libc/src/posix_io.rs index 5b9da34b9ec..f88fdaeb7cb 100644 --- a/userland/libc/src/posix_io.rs +++ b/userland/libc/src/posix_io.rs @@ -573,9 +573,6 @@ pub struct pollfd { pub revents: i16, } -const POLLIN: i16 = 1; -const POLLOUT: i16 = 4; - #[no_mangle] pub unsafe extern "C" fn poll(fds: *mut pollfd, nfds: u32, timeout: i32) -> i32 { if nfds == 0 { @@ -597,13 +594,10 @@ pub unsafe extern "C" fn poll(fds: *mut pollfd, nfds: u32, timeout: i32) -> i32 let timeout_ns = if timeout < 0 { None } else { Some(timeout as u64 * 1_000_000) }; let n = nfds as usize; + let entries: alloc::vec::Vec<(i32, i16)> = (0..n).map(|i| ((*fds.add(i)).fd, (*fds.add(i)).events)).collect(); let poller = toyos::poller::Poller::new(n as u32); - for i in 0..n { - let pfd = &*fds.add(i); - let mut flags = 0u32; - if pfd.events & POLLIN != 0 { flags |= toyos::poller::READABLE; } - if pfd.events & POLLOUT != 0 { flags |= toyos::poller::WRITABLE; } - poller.watch_raw(toyos_abi::RawHandle(pfd.fd as u32), flags, i as u64); + for (entry, fd, interest) in crate::pollreq::watches(&entries) { + poller.watch_raw(toyos_abi::RawHandle(fd as u32), interest, entry as u64); } let mut ready_set = alloc::vec![false; n]; @@ -612,9 +606,10 @@ pub unsafe extern "C" fn poll(fds: *mut pollfd, nfds: u32, timeout: i32) -> i32 }); let mut ready = 0i32; for i in 0..n { + let answered = ready_set[crate::pollreq::watch_of(&entries, i)]; let pfd = &mut *fds.add(i); pfd.revents = 0; - if ready_set[i] { + if answered { pfd.revents = pfd.events; ready += 1; }