diff --git a/issues/build/the-two-watch-poll-model-counts-a-dropped-watchs-answer-as-its-completion.md b/issues/build/the-two-watch-poll-model-counts-a-dropped-watchs-answer-as-its-completion.md new file mode 100644 index 00000000000..899c9c18517 --- /dev/null +++ b/issues/build/the-two-watch-poll-model-counts-a-dropped-watchs-answer-as-its-completion.md @@ -0,0 +1,23 @@ +--- +status: open +kind: tooling +opened: 2026-09-30 +--- + +# The two-watch poll model counts a dropped watch's answer as its completion + +`toyos-sched/loom/tests/loom_watch.rs`'s +`a_poll_on_two_watches_racing_both_posts_completes_exactly_once` moves both +worlds into their producer threads, so both watches drop before its assertion, +and a watch's drop fires every live entry as `Fire::Gone`, which the model's +`Entry` counts as a post. Its "never by neither" half cannot fail: a poll no +post and no recheck completed is completed by the drop. + +**Evidence:** with each producer posting before it makes its condition true (a +lost completion by construction), `cargo test -p toyos-sched-loom --test +loom_watch a_poll_on_two_watches` is EXIT=0, 1 passed. `poll_racing`, the +single-watch model beside it, had the same shape and holds its world past the +assertion since #634. + +**Exit:** the model keeps both worlds alive past its assertions, and the +mutation above reds it. 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 new file mode 100644 index 00000000000..b38d2a9496f --- /dev/null +++ b/issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md @@ -0,0 +1,41 @@ +--- +status: assigned +kind: defect +opened: 2026-09-30 +--- + +# A process lengthens an interrupts-off walk by the threads it parks on one ring + +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` +(`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 +holder, decides how long a CPU runs with interrupts masked: + +- **N threads parked in `submit` on one ring** + are N registrations on its watch. Every completion into that ring posts the + watch in place, which notifies all N under the list lock + (`toyos-sched/src/watch.rs`), each a word exchange and, for a parked + thread, a mailbox push and perhaps an IPI (`toyos-sched/src/park.rs`). + Each woken thread's unregister is a `position` and a + `remove` over the N, and a registration + 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 + registrations sweep them four at a time. + +Nothing caps N: a thread costs its process a 128 KiB kernel stack +(`kernel/src/process.rs`) and no count. Before #634 every one of these +walks ran with interrupts open, under preemption off. By reading, not +measured: no instrument in the tree reads an interrupts-off window. + +**Exit**: step 2's interrupts-off window, read on the T14 while one process +parks 256 threads in `submit` on one ring and a sibling thread completes into +it, is no longer at step 1's head than at stage 6's first commit under the +same load. diff --git a/issues/kernel/an-irq-watchs-freeing-cancel-compiles-in-a-handler.md b/issues/kernel/an-irq-watchs-freeing-cancel-compiles-in-a-handler.md new file mode 100644 index 00000000000..9b29532b441 --- /dev/null +++ b/issues/kernel/an-irq-watchs-freeing-cancel-compiles-in-a-handler.md @@ -0,0 +1,29 @@ +--- +status: assigned +kind: defect +opened: 2026-09-30 +--- + +# An IrqWatch's freeing cancel compiles in a handler + +Held by the small-kernel track's stage 6 step 5 +(`issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md`), +whose exit reads it. + +An interrupt handler may not free: it can interrupt the allocator's holder. +`IrqWatch` has no `post`, so a thread's post written in a handler is refused +at build. But `cancel_polls` is every watch's, and it frees every entry it +takes out. Its callers are a claim's release and the close of a claim or an +audio device, both threads', and no type keeps it out of a handler. + +**Evidence**: `WATCHES[slot].cancel_polls();` written after +`pcidev::note_fault`'s post builds: the x86-64 kernel clippy shape exits 0. +The kernel holds no proof of thread context a handler cannot mint. `Parkable` +proves a context may park, and with `IrqWatch`'s cancel taking one, +`Parkable::at_entry()` written at the same line in `note_fault` builds too +(clippy 0). Its only check is at run time, an assertion on the preempt +depth: by reading, a handler that interrupted Ring 3 passes it. + +**Exit**: a proof of thread context no interrupt handler can construct, taken +by `IrqWatch`'s `cancel_polls`, so the `note_fault` line above is refused at +build. diff --git a/issues/kernel/every-wait-in-this-kernel-is-a-spin.md b/issues/kernel/every-wait-in-this-kernel-is-a-spin.md index e86fa18ee1b..b8c7cd78b4f 100644 --- a/issues/kernel/every-wait-in-this-kernel-is-a-spin.md +++ b/issues/kernel/every-wait-in-this-kernel-is-a-spin.md @@ -119,7 +119,7 @@ both before any lock conversion; the order is forced, not preferred. - A watch is a node the waiter lends to the object, and the subject is a borrowed reference, never an id. **Rejected:** a global registry, a slot - arena, two park channels, posting from interrupt context, multishot polls, + arena, two park channels, multishot polls, userspace-only blocking wrappers, a sleep lock that spins where it cannot park, poisoning, and shootdown-as-completion. A freed object cannot be named. - The park token proves the *context* may park and never encodes which locks are diff --git a/issues/kernel/nothing-fails-when-a-devices-release-or-close-stops-answering-its-polls.md b/issues/kernel/nothing-fails-when-a-devices-release-or-close-stops-answering-its-polls.md new file mode 100644 index 00000000000..47cccb8b434 --- /dev/null +++ b/issues/kernel/nothing-fails-when-a-devices-release-or-close-stops-answering-its-polls.md @@ -0,0 +1,31 @@ +--- +status: assigned +kind: defect +opened: 2026-09-30 +--- + +# Nothing fails when a device's release or close stops answering its polls + +Held by the small-kernel track's stage 6 step 5 +(`issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md`), +whose exit reads it. + +A device's watch is an `IrqWatch`, and three thread sites answer the polls on +it as gone: `pcidev::tear_down` when a claimed function is released +(`kernel/src/pcidev/mod.rs:1489`), `Claim::drop` when an audio claim goes +(`kernel/src/device.rs:86`), and the close of a claim or an audio device +through `WatchRef::cancel_polls`'s `Irq` arm (`kernel/src/object/ops.rs:263`). +A poll none of them answers waits on interrupts that are the next holder's or +nobody's. No host test compiles `pcidev`, `device` or `ops`, and no guest test +ends a polled device, so each call can go and nothing reds. + +**Evidence**: each deletion builds. The x86-64 kernel clippy shape exits 0 +with the release's `cancel_polls()` deleted, 0 with the `Irq` arm made `{}`, +and 0 with the audio claim's cancel made `{}`. + +**Exit**: `WatchRef`'s cancel is one dispatch over every variant, as main's +`Deref` was, so the `Irq` arm above cannot be written apart from the others; +and a guest test, in which a claim's holder exits with a poll of its claim on a +ring another process holds, reads that poll answered `-NotFound` for a claimed +function and for an audio device, and reds with the release's cancel deleted +and with the audio claim's. diff --git a/issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md b/issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md index 696e5048872..e7065240b7c 100644 --- a/issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md +++ b/issues/kernel/the-kernel-is-small-interrupts-post-and-threads-wait.md @@ -144,10 +144,53 @@ times: written once as straight-line code. **Exit**: no interrupts-off window longer than a register access, and keyboard input keeps flowing while a stick misbehaves. -6. **The scheduler knows nothing about devices.** Interrupt handlers only post - to their device's `Watch`, and the device's thread does the work. The - per-CPU IRQ relay, the driver list in the scheduler pass and the idle - special cases are deleted. +6. **The scheduler knows nothing about devices.** A handler posts its + device's `Watch` and ends its interrupt; the thread waiting on that watch + does the work, and no step creates a kernel thread. `irq_ring`, the driver + list in `drain_irqs` and the idle loop's device checks are gone by step 5. + Each step measures the kernel's lines, and from step 2 the longest + interrupts-off and preemption-off windows, against stage 6's first commit. + Steps 3 and 4 do not land alone: they land with #592's i8042 stage and + with usbd. Stage 6 stays open past step 5 until the owner rules on the + panel the dump paints. + 1. **Interrupts post.** A post is legal in a handler: the watches a handler + posts, and the completions of a ring they complete into, sit behind + interrupts-off locks nothing allocates or frees under, and every other + watch's lock leaves interrupts open. A claimed function's vector, the + IOMMU's refusal and both audio backends post from the handler, and + `irq_ring`'s `UserDev` and `Audio` and their arms in `drain_irqs` go. + The thread is the holder's: netd's, blockd's, soundd's mix thread, and + the `isa` claim's when #592 lands. **Exit**: `handler_post_without_a_pass`, + a vector taken on a CPU holding preemption off, inside a post of its own + watch, inside a completion into a ring polling it, or inside that ring's + own watch, posting once that section lets go and before any pass, red on + the base; the watch's loom models over the new post. + 2. **The windows, measured**: the longest interrupts-off and preemption-off + windows per CPU, reported beside the IRQ census and fed by each + architecture's masking primitives and entries, the number the ARM + track's stage 4 owes as well. Applied to stage 6's first commit for + the baseline. **Exit**: both windows read on the T14, which is x86 + metal, at stage 6's start and at step 1's head, and neither is longer + at step 1's head than at the start, under the load + `issues/kernel/a-process-lengthens-an-interrupts-off-walk-by-the-threads-it-parks-on-one-ring.md` + names as well. + 3. **The i8042's thread is ps2server's** (#592's i8042 stage): `irq_ring`'s + `I8042`, `keyboard_controller::service` and the idle loop's + `verdict_due` go with the kernel's driver. **Exit**: that stage's. + 4. **xHCI's thread is usbd's** (step 10 above): `Xhci`, `poll_if_pending` + and `port_work_pending` go with the kernel's driver, and `irq_ring` with + them. **Exit**: step 10's. + 5. **The pass is the scheduler's.** `drain_irqs` goes: the blocked-task + dump and the heartbeat become `pass`'s own, and the TCO feed stays, + since what it proves is that passes run. The dump still paints its + report on the panel and holds it there, a device the pass reaches; + whether that stays is the owner's ruling. **Exit**: `drain_irqs` and the + idle loop's device checks are gone, both windows are measured against + stage 6's start, and the exits of + `issues/kernel/an-irq-watchs-freeing-cancel-compiles-in-a-handler.md` + and + `issues/kernel/nothing-fails-when-a-devices-release-or-close-stops-answering-its-polls.md` + are met. ## Standing diff --git a/issues/kernel/toyos-runs-on-arm64.md b/issues/kernel/toyos-runs-on-arm64.md index b93f81a9881..123db6c3868 100644 --- a/issues/kernel/toyos-runs-on-arm64.md +++ b/issues/kernel/toyos-runs-on-arm64.md @@ -317,6 +317,16 @@ Each stage names its exit; "measured" means a number from a run. before its body) get a test here that reds with `put` written back as one `self.buf.write(off, trb)`; x86's TSO hides all three from every guest test until then. + **The claim's handler is the first arm of `irq()` + (`kernel/src/arch/aarch64/trap.rs:135`) that posts a watch or lets go of a + `Lock`**, and either runs `preempt::enable`, whose pass at depth zero + (`kernel/src/preempt.rs:67`) reads nothing of `DAIF`; `do_preempt`'s + `assert_baseline(BASELINE_IRQ_EXIT)` (`kernel/src/scheduler.rs:358`) passes + at depth zero, so that pass would run inside the handler, before + `irqchip::end`. Owed before that arm lands: `irq()` holds the preempt count + across every device arm, as x86-64's `device_irq_entry` does, or + `preempt::enable` refuses a pass with interrupts masked, which `IrqOff`'s + SAFETY (`kernel/src/sched/driver.rs:49`) already assumes. 7. **Userland boots.** `init`, `logd`, the compositor, netd, soundd and sshd, built for `aarch64-unknown-toyos`. C programs stay x86-only until this diff --git a/kernel-loom/Cargo.toml b/kernel-loom/Cargo.toml index 86ccc7df2c7..53a3aef831e 100644 --- a/kernel-loom/Cargo.toml +++ b/kernel-loom/Cargo.toml @@ -138,9 +138,8 @@ shard-publish-relaxed = [] # Never on by default and never reachable from a kernel build, which declares # the same name only so `cfg` checking knows it. sleeplock-acquire-off = [] -# The negative control for a claimed PCI function's interrupt record. Its two -# read-modify-writes — the reader's `swap` of the count and the scheduler pass's -# `swap` of the wake flag — become a load and a store, which is the whole of +# The negative control for a claimed PCI function's interrupt record. The +# reader's `swap` of the count becomes a load and a store, which is the whole of # what this record's design is, and `device_irq.rs` must red: # # cargo test --manifest-path kernel-loom/Cargo.toml --features device-irq-lossy \ diff --git a/kernel-loom/tests/device_irq.rs b/kernel-loom/tests/device_irq.rs index 450046d45eb..d5370bf87fb 100644 --- a/kernel-loom/tests/device_irq.rs +++ b/kernel-loom/tests/device_irq.rs @@ -2,30 +2,25 @@ //! //! The kernel programs one MSI-X vector per claimed function and accumulates //! what arrives into a record its holder reads through a syscall. So there are -//! two parties on two CPUs and two words between them: an ISR that bumps a -//! count and arms a wake, a reader that takes the count, and a scheduler pass -//! that takes the wake. +//! two parties on two CPUs and one word between them: an ISR that bumps a +//! count, and a reader that takes it. //! -//! **The invariant is that every message is counted exactly once and owes -//! exactly one wake.** Nothing here orders anything else — each word is the -//! whole of what it says, so the orderings are `Relaxed` and every property is -//! an interleaving. What makes them hold is that the taking side of both words -//! is a read-modify-write: a reader that loaded a count and then cleared it -//! drops every message the ISR recorded in between, and two scheduler passes -//! that both loaded a wake flag both wake one message's watchers. +//! **The invariant is that every message is counted exactly once.** Nothing +//! here orders anything else — the word is the whole of what it says, so the +//! orderings are `Relaxed` and the property is an interleaving. What makes it +//! hold is that the taking side is a read-modify-write: a reader that loaded a +//! count and then cleared it drops every message the ISR recorded in between. //! -//! That pair is the record's whole design, so the negative control is the pair -//! turned off — a cargo feature rather than a comment: +//! That is the record's whole design, so the negative control is it turned off +//! — a cargo feature rather than a comment: //! //! ```text //! cargo test --manifest-path kernel-loom/Cargo.toml --features device-irq-lossy \ //! --test device_irq //! ``` //! -//! makes both `swap`s a load and a store and the ISR's `fetch_add` a load, an -//! add and a store, and this file must red — at -//! [`every_message_is_counted_once`] and at [`one_message_is_one_wake`], which -//! are the two defects stated exactly. +//! makes the `swap` a load and a store and the ISR's `fetch_add` a load, an +//! add and a store, and [`every_message_is_counted_once`] must red. use kernel_loom::device_irq::Interrupt; use loom::sync::Arc; @@ -62,32 +57,6 @@ fn every_message_is_counted_once() { }); } -/// A wake is owed exactly once per message, and the pass that owes it is the -/// one that takes it. -/// -/// The wake and the ISR are the *same* CPU in the kernel — the scheduler pass -/// that drains runs after the handler that armed it, with interrupts off in -/// between — so this models the weaker thing that must also hold: two passes -/// racing each other never both wake, and never both decline. -#[test] -fn one_message_is_one_wake() { - loom::model(|| { - let irq = Arc::new(Interrupt::new()); - irq.took(); - - let other = { - let irq = irq.clone(); - loom::thread::spawn(move || irq.take_pending()) - }; - - let mine = irq.take_pending(); - let theirs = other.join().unwrap(); - - assert!(!(mine && theirs), "two passes both woke one message's watchers"); - assert!(mine || theirs, "neither pass woke a message that had already arrived"); - }); -} - /// A reader that finds nothing answers nothing, and leaves nothing behind. /// /// Single-threaded and deliberately so: this is the answer on the path a @@ -106,31 +75,3 @@ fn an_idle_record_answers_nothing() { assert!(!irq.armed(), "a drained record still reads ready"); }); } - -/// A fault owes its holder a wake, and the pass that takes it reads the fault. -/// -/// The fault handler and the scheduler pass that turns the wake into a wake-up -/// may be different CPUs, and what the woken holder reads next is the refusal: -/// a pass that took the wake and still read the claim unfaulted would wake a -/// holder into reading "no interrupt" and parking again, for a function that -/// can no longer send one. A `Relaxed` wake fails the first assertion, and a -/// fault that posts no wake the last. -#[test] -fn a_faults_wake_carries_the_fault() { - loom::model(|| { - let irq = Arc::new(Interrupt::new()); - - let handler = { - let irq = irq.clone(); - loom::thread::spawn(move || irq.fault()) - }; - - let taken_before = irq.take_pending(); - if taken_before { - assert!(irq.faulted(), "a pass took a fault's wake and read the claim unfaulted"); - } - handler.join().unwrap(); - assert!(irq.faulted()); - assert!(taken_before || irq.take_pending(), "a fault owed its holder a wake and posted none"); - }); -} diff --git a/kernel/Cargo.toml b/kernel/Cargo.toml index 1d42053dd63..5beca3f086b 100644 --- a/kernel/Cargo.toml +++ b/kernel/Cargo.toml @@ -53,10 +53,10 @@ log-commit-release-off = [] # `src/log/registry.rs`'s pointer store and load go `Relaxed`, and # `log_publish` reds. shard-publish-relaxed = [] -# `src/pcidev/record.rs`'s two read-modify-writes become a load and a store, so -# a message the ISR records between a reader's two halves is lost and a wake can -# be owed twice, and `device_irq` reds. That pair *is* the record's design, so -# turning it off reverts the whole of what the model claims. +# `src/pcidev/record.rs`'s read-modify-writes become a load and a store, so a +# message the ISR records between a reader's two halves is lost, and +# `device_irq` reds. They *are* the record's design, so turning them off reverts +# the whole of what the model claims. device-irq-lossy = [] # `src/sched/dump_request.rs`'s `take` and `end_report` go `Relaxed`, so the next report # is unordered against the last one's writes, and `dump_request` reds. diff --git a/kernel/src/actuator.rs b/kernel/src/actuator.rs index 0119f80c880..89cc9f2dc42 100644 --- a/kernel/src/actuator.rs +++ b/kernel/src/actuator.rs @@ -254,6 +254,12 @@ actuators! { /// parking, so a post lands in the window its commit must refuse the park over. watch_window = "watch-window"; + /// Raise an unheld claim slot's vector inside a post of its own watch, + /// inside a completion into a ring polling it, and inside that ring's own + /// watch, while the CPU holds preemption off, and count whether the + /// handler posted it there. + handler_post = "handler-post"; + /// Starve the four xHCI bring-up register waits in `init_one`. xhci_deaf_controller = "xhci-deaf-controller"; diff --git a/kernel/src/arch/x86_64/idt/hda.rs b/kernel/src/arch/x86_64/idt/hda.rs index 579bebedfc6..d75e99ef709 100644 --- a/kernel/src/arch/x86_64/idt/hda.rs +++ b/kernel/src/arch/x86_64/idt/hda.rs @@ -1,6 +1,6 @@ use super::device_irq::device_irq_entry; -// Lock-free, heap-free: may interrupt a CPU holding the controller lock (preemption disabled, not interrupts). +// Heap-free, and takes only interrupts-off locks: may interrupt a CPU holding the controller lock (preemption disabled, not interrupts). extern "sysv64" fn hda_handler() { crate::arch::percpu::irq_took!(Hda); crate::drivers::hda::isr_complete(); diff --git a/kernel/src/arch/x86_64/idt/user_dev.rs b/kernel/src/arch/x86_64/idt/user_dev.rs index 18142f94df5..9ad2dccbc32 100644 --- a/kernel/src/arch/x86_64/idt/user_dev.rs +++ b/kernel/src/arch/x86_64/idt/user_dev.rs @@ -6,18 +6,14 @@ //! wake every user driver in the machine on any of their interrupts, which is //! one process learning when another's device is busy. //! -//! Lock-free and heap-free like every other device entry here: the record is -//! atomics and the wake happens on the next scheduler pass. +//! Heap-free like every other device entry here: the record is atomics, and +//! the claim's watch is posted from the handler. use super::device_irq::device_irq_entry; -use crate::irq_ring::IrqSource; fn took(slot: usize) { crate::arch::percpu::irq_took!(UserDev); crate::pcidev::isr(slot); - crate::irq_ring::isr_publish(IrqSource::UserDev, crate::clock::nanos_since_boot()); - // Force resched now, so `drain_irqs` turns the record into a wake before - // the next quantum tick rather than after it. crate::preempt::set_need_resched(); crate::arch::apic::eoi(); } diff --git a/kernel/src/arch/x86_64/idt/virtio_sound.rs b/kernel/src/arch/x86_64/idt/virtio_sound.rs index 98060824d62..6e3744d624a 100644 --- a/kernel/src/arch/x86_64/idt/virtio_sound.rs +++ b/kernel/src/arch/x86_64/idt/virtio_sound.rs @@ -1,6 +1,6 @@ use super::device_irq::device_irq_entry; -// Lock-free and heap-free: may interrupt a CPU holding the controller lock, which disables preemption but not interrupts. +// Heap-free, and takes only interrupts-off locks: may interrupt a CPU holding the controller lock, which disables preemption but not interrupts. extern "sysv64" fn virtio_sound_handler() { crate::arch::percpu::irq_took!(Sound); crate::drivers::virtio_sound::isr_complete(); diff --git a/kernel/src/arch/x86_64/mod.rs b/kernel/src/arch/x86_64/mod.rs index 41ec8a0f805..1e565f53f79 100644 --- a/kernel/src/arch/x86_64/mod.rs +++ b/kernel/src/arch/x86_64/mod.rs @@ -52,10 +52,8 @@ pub const ELF_MACHINE: toyos_elf::Machine = toyos_elf::Machine::X86_64; /// Interrupts masked on this CPU for as long as the guard lives, and then put /// back as they were — restored, not enabled — so a guard nests inside a region -/// that is already masked. The one way this kernel masks and restores: the -/// scheduler's pass, a log record's reservation and publication, and the -/// console backend each hold one. `TF` is always clear in Ring 0, so the guard -/// leaves it alone. +/// that is already masked. The one way this kernel masks and restores. `TF` is +/// always clear in Ring 0, so the guard leaves it alone. /// /// Both edges are compiler barriers (no `nomem`): a memory access written /// inside the region is emitted inside it. diff --git a/kernel/src/arch/x86_64/vtd/fault.rs b/kernel/src/arch/x86_64/vtd/fault.rs index c0188b87c1e..3de98c0f382 100644 --- a/kernel/src/arch/x86_64/vtd/fault.rs +++ b/kernel/src/arch/x86_64/vtd/fault.rs @@ -1,6 +1,6 @@ //! Vt-d fault interrupt handling: MSI-delivered, never polled. //! -//! The handler is bounded, allocates nothing and takes no lock; unit and +//! The handler is bounded, allocates nothing; unit and //! function state lives in fixed arrays of atomics, written once before the //! mask comes off. Whatever the stream, the same things happen first: Bus //! Master Enable cleared on the function that faulted, the first record latched @@ -36,8 +36,7 @@ const FSTS_OVERFLOW: u32 = 1 << 0; // One fault recording register's F bit, in the 32-bit word that carries it. const RECORD_FAULT: u32 = 1 << 31; -// Atomics only: the handler takes no lock; each field is written once before -// the unit's mask comes off. +// Atomics only: each field is written once before the unit's mask comes off. struct FaultUnit { // Physical base of the register window; 0 means this slot is unused. // vtd::window refuses a base of 0, so no armed unit can collide with the sentinel. diff --git a/kernel/src/device.rs b/kernel/src/device.rs index 17c7f60d5a0..66935836c20 100644 --- a/kernel/src/device.rs +++ b/kernel/src/device.rs @@ -79,7 +79,20 @@ impl Claim { impl Drop for Claim { fn drop(&mut self) { match self.what { - Claimed::Class(class) => *taken(class).lock() = false, + Claimed::Class(class) => { + // Before the flag goes: a poll the next holder registers is not this claim's to answer. + match class { + DeviceType::HdaAudio | DeviceType::VirtioSound => { + crate::drivers::AUDIO_WATCH.cancel_polls() + } + DeviceType::Keyboard + | DeviceType::Mouse + | DeviceType::Framebuffer + | DeviceType::PciFunction + | DeviceType::Partition => {} + } + *taken(class).lock() = false; + } // Bus mastering off, then the domain, then the pages: `release` // owns that order, and this is where a dying process reaches it. Claimed::PciFunction(slot) => crate::pcidev::release(slot), diff --git a/kernel/src/drivers/hda.rs b/kernel/src/drivers/hda.rs index 135c620c34e..ca6c39b12be 100644 --- a/kernel/src/drivers/hda.rs +++ b/kernel/src/drivers/hda.rs @@ -166,7 +166,7 @@ pub fn isr_complete() { } ISR.timestamp.store(timestamp, Ordering::Relaxed); ISR.mask.fetch_or(mask, Ordering::Release); - crate::irq_ring::isr_publish(crate::irq_ring::IrqSource::Audio, timestamp); + super::AUDIO_WATCH.post_in_place(); crate::preempt::set_need_resched(); } diff --git a/kernel/src/drivers/mod.rs b/kernel/src/drivers/mod.rs index d71937ef0b7..a4f8cf9f890 100644 --- a/kernel/src/drivers/mod.rs +++ b/kernel/src/drivers/mod.rs @@ -23,4 +23,4 @@ pub use crate::mm::DmaPool; /// What an audio read and an audio poll wait on: one for both backends, since /// at most one binds and the claim's reader cannot know which. -pub static AUDIO_WATCH: crate::watch::Watch = crate::watch::Watch::new(); +pub static AUDIO_WATCH: crate::watch::IrqWatch = crate::watch::IrqWatch::new(); diff --git a/kernel/src/drivers/virtio_sound.rs b/kernel/src/drivers/virtio_sound.rs index 00e91ecf2ca..416d9001b5f 100644 --- a/kernel/src/drivers/virtio_sound.rs +++ b/kernel/src/drivers/virtio_sound.rs @@ -108,8 +108,7 @@ pub fn isr_complete() { return; } isr_push_completion(mask, timestamp); - crate::irq_ring::isr_publish(crate::irq_ring::IrqSource::Audio, timestamp); - // Force a scheduler entry on IRQ return so the record becomes wakes now, not at next tick. + super::AUDIO_WATCH.post_in_place(); crate::preempt::set_need_resched(); } diff --git a/kernel/src/inbox/mod.rs b/kernel/src/inbox/mod.rs index d7c69cc7906..8e6ed014fd1 100644 --- a/kernel/src/inbox/mod.rs +++ b/kernel/src/inbox/mod.rs @@ -19,14 +19,18 @@ //! 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. //! -//! **Locks.** A ring's own lock takes nothing under it and is never taken under -//! a watch's: a post fires its polls with its list let go. A ring's own watch +//! **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. 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}; @@ -36,7 +40,7 @@ use crate::process::{self, Pid}; use crate::scheduler; use crate::sync::Lock; use crate::time::{Deadline, Duration}; -use crate::watch::Watch; +use crate::watch::{IrqLock, IrqWatch}; use crate::DirectMap; use toyos_abi::inbox::{ @@ -62,7 +66,12 @@ impl InboxRef { impl Drop for InboxRef { fn drop(&mut self) { - // Taken out under the lock and let go of outside it: the unmap flushes. + // The page is the completions', so it goes only once no post 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 { + 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"); }; @@ -70,7 +79,7 @@ impl Drop for InboxRef { poll.withdraw(); } // `Unmapped`'s drop flushes; the `Arc` drop after it frees the pages. - drop(state.shm.unmap_from(state.owner_pid)); + drop(completions.shm.unmap_from(state.owner_pid)); } } @@ -193,68 +202,80 @@ const MAX_PENDING_WATCHES: usize = 1024; 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>, /// Threads parked in `submit`; never a poll — see the module header. - watch: Watch, + watch: IrqWatch, } struct RingState { + /// The page's address. The page is [`Completions`]'s, which the teardown + /// lets go of only after it has taken this. shm_phys: DirectMap, - /// A ring's page has no lifetime of its own; it goes with the last handle to the ring. - shm: Arc, submission_size: u32, - completion_size: u32, - /// The kernel's own copy of the completion tail, the only one it reads. - completion_tail: u32, /// Polls still armed as of the last registration, which sweeps the rest. pending: Vec>, owner_pid: Pid, } -impl RingState { - // No accessor below returns a Rust reference into this page — the process - // maps it writable, so only atomics or `read_volatile` are sound here. +// No accessor below returns a Rust reference into a ring's page — the process +// maps it writable, so only atomics or `read_volatile` are sound here. - /// One atomic word of one ring header; never `&RingHeader` — see the block above. - fn ring_word(&self, ring_off: u64, field_off: usize) -> &core::sync::atomic::AtomicU32 { - let ptr = self.shm_phys.as_mut_ptr::(); - // SAFETY: offset is in-bounds and 4-aligned within the 2 MiB page; `AtomicU32` is sound over memory the process also writes. - unsafe { - core::sync::atomic::AtomicU32::from_ptr( - ptr.add(ring_off as usize + field_off) as *mut u32, - ) - } +/// One atomic word of one ring header; never `&RingHeader` — see the block above. +fn ring_word(page: &DirectMap, ring_off: u64, field_off: usize) -> &core::sync::atomic::AtomicU32 { + let ptr = page.as_mut_ptr::(); + // SAFETY: offset is in-bounds and 4-aligned within the 2 MiB page, which outlives both of the ring's halves that name it; `AtomicU32` is sound over memory the process also writes. + unsafe { + core::sync::atomic::AtomicU32::from_ptr( + ptr.add(ring_off as usize + field_off) as *mut u32, + ) } +} +impl RingState { fn submission_head(&self) -> &core::sync::atomic::AtomicU32 { - self.ring_word(SUBMISSION_RING_OFF, core::mem::offset_of!(RingHeader, head)) + ring_word(&self.shm_phys, SUBMISSION_RING_OFF, core::mem::offset_of!(RingHeader, head)) } fn submission_tail(&self) -> &core::sync::atomic::AtomicU32 { - self.ring_word(SUBMISSION_RING_OFF, core::mem::offset_of!(RingHeader, tail)) + ring_word(&self.shm_phys, SUBMISSION_RING_OFF, core::mem::offset_of!(RingHeader, tail)) } + /// One submission entry, copied out by value via `read_volatile` — never a `&Submission`. + fn submission_at(&self, index: u32) -> Submission { + let ptr = self.shm_phys.as_mut_ptr::(); + // SAFETY: `index` is masked by `submission_size` (≤256), keeping the read in-bounds and aligned within the page. + unsafe { (ptr.add(SUBMISSIONS_OFF as usize + index as usize * core::mem::size_of::()) as *const Submission).read_volatile() } + } +} + +/// What a poll's completion writes, which a post from an interrupt handler +/// reaches, 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, + page: DirectMap, + completion_size: u32, + /// The kernel's own copy of the completion tail, the only one it reads. + completion_tail: u32, +} + +impl Completions { fn completion_head(&self) -> &core::sync::atomic::AtomicU32 { - self.ring_word(COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, head)) + ring_word(&self.page, COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, head)) } fn completion_tail_word(&self) -> &core::sync::atomic::AtomicU32 { - self.ring_word(COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, tail)) + ring_word(&self.page, COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, tail)) } fn completion_dropped(&self) -> &core::sync::atomic::AtomicU32 { - self.ring_word(COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, dropped)) - } - - /// One submission entry, copied out by value via `read_volatile` — never a `&Submission`. - fn submission_at(&self, index: u32) -> Submission { - let ptr = self.shm_phys.as_mut_ptr::(); - // SAFETY: `index` is masked by `submission_size` (≤256), keeping the read in-bounds and aligned within the page. - unsafe { (ptr.add(SUBMISSIONS_OFF as usize + index as usize * core::mem::size_of::()) as *const Submission).read_volatile() } + ring_word(&self.page, COMPLETION_RING_OFF, core::mem::offset_of!(RingHeader, dropped)) } /// The address of one completion entry — a pointer, never a `&mut` minted from a shared borrow. fn completion_at(&self, index: u32) -> *mut Completion { - let ptr = self.shm_phys.as_mut_ptr::(); + let ptr = self.page.as_mut_ptr::(); // SAFETY: `index` is masked by `completion_size` (≤512), keeping the offset inside the page. unsafe { ptr.add(COMPLETION_RING_OFF as usize + core::mem::size_of::() + index as usize * core::mem::size_of::()) as *mut Completion } } @@ -269,7 +290,7 @@ impl RingState { return; } let idx = tail & (self.completion_size - 1); - // SAFETY: `idx` is masked to ring size; the ring's lock serializes kernel writers. + // SAFETY: `idx` is masked to ring size; the completions' lock serializes kernel writers. unsafe { self.completion_at(idx).write(Completion { token: user_data, result, flags }) }; self.completion_tail = tail.wrapping_add(1); self.completion_tail_word().store(tail.wrapping_add(1), Ordering::Release); @@ -292,15 +313,27 @@ 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.with_state(|state| state.post_completion(user_data, result, 0)); - if posted.is_ok() { - self.watch.post(); + let posted = self.completions.with(|c| { + // Inside the section whatever lock it is, so `handler-post` reds + // on one that leaves interrupts open. + #[cfg(feature = "boot-actuators")] + crate::watch::handler_post::raise_if_staged(); + 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) + } } /// Largest submission ring a process may ask for. @@ -357,18 +390,68 @@ pub fn create(depth: u32) -> Result<(InboxRef, u64), SyscallError> { let inbox = Arc::new(Inbox { state: Lock::new(Some(RingState { shm_phys, - shm, submission_size, - completion_size, - completion_tail: 0, pending: Vec::new(), owner_pid: pid, })), - watch: Watch::new(), + completions: IrqLock::new(Some(Completions { + shm, + page: shm_phys, + completion_size, + completion_tail: 0, + })), + watch: IrqWatch::new(), }); Ok((InboxRef(inbox), shm_vaddr)) } +/// `handler-post`'s ring: the kernel's own, mapped into no process and +/// submitted to by nobody, which polls a watch and completes as a submission +/// does. +#[cfg(feature = "boot-actuators")] +pub(crate) struct Staged(Arc); + +#[cfg(feature = "boot-actuators")] +impl Staged { + pub(crate) fn new() -> Self { + let depth = 2 * toyos_sched::watch::handler_post::HOLDS; + let shm = SharedMemObject::create(crate::mm::PAGE_2M).expect("handler-post: a ring's page"); + let page = shm.phys_before_mapping(); + write_ring_page(page, depth, depth * 2); + Self(Arc::new(Inbox { + state: Lock::new(None), + completions: IrqLock::new(Some(Completions { + shm, + page, + completion_size: depth * 2, + completion_tail: 0, + })), + watch: IrqWatch::new(), + })) + } + + /// A poll of this ring on `watch`, which that watch's next post completes. + pub(crate) fn poll(&self, watch: &IrqWatch) { + let poll = Arc::new(Poll { + inbox: self.0.clone(), + user_data: 0, + handle: RawHandle(0), + state: Once::new(), + }); + watch.add_poll(PollEntry { poll, direction: Readiness { readable: true, writable: false } }); + } + + pub(crate) fn complete(&self) { + self.0.complete(0, 0); + } + + /// Run `f` holding this ring's own watch's list lock, as a registration + /// in `submit` holds it. + pub(crate) fn holding_its_watch(&self, f: impl FnOnce()) { + self.0.watch.holding(f); + } +} + /// Processes submissions and waits for completions; called from the syscall handler. pub fn submit( inbox: &Arc, @@ -391,7 +474,7 @@ pub fn submit( } loop { - let (count, dropped) = inbox.with_state(|s| (s.completion_count(), s.dropped()))?; + let (count, dropped) = inbox.with_completions(|c| (c.completion_count(), c.dropped()))?; if count >= min_complete || min_complete == 0 { return Ok(count); @@ -418,7 +501,7 @@ pub fn submit( 0, WaitClass::Io, deadline, - || inbox.with_state(|s| s.completion_count()).map_or(true, |n| n >= min_complete), + || inbox.with_completions(|c| c.completion_count()).map_or(true, |n| n >= min_complete), ) .is_err() { diff --git a/kernel/src/irq_ring.rs b/kernel/src/irq_ring.rs index 6d2c7216522..a71e5204c1c 100644 --- a/kernel/src/irq_ring.rs +++ b/kernel/src/irq_ring.rs @@ -10,19 +10,23 @@ use crate::arch::percpu; use crate::scheduler::MAX_CPUS; /// Interrupt sources that drive scheduling; exhaustive, so a new variant requires updating every `match`. +/// The discriminant is `trace::Kind::IrqDrain`'s top byte, pinned for the reason `Kind` is. #[derive(Clone, Copy, PartialEq, Eq, Debug)] pub enum IrqSource { - Audio, - /// Every vector a claimed PCI function delivers on, coalesced into one - /// slot: the record only says a pass is owed, and `pcidev` keeps the - /// per-slot flag that says whose. - UserDev, - Xhci, - I8042, + Xhci = 2, + I8042 = 3, } impl IrqSource { - pub const COUNT: usize = 4; + pub const COUNT: usize = 2; + + /// This source's record in a CPU's slots. + const fn slot(self) -> usize { + match self { + Self::Xhci => 0, + Self::I8042 => 1, + } + } } /// 64-byte aligned so two CPUs' slots never share a cache line. @@ -36,7 +40,7 @@ static SLOTS: [CpuSlots; MAX_CPUS] = pub fn isr_publish(source: IrqSource, timestamp_nanos: u64) { // MSI-X vectors are configured after clock calibration, so a real IRQ never stamps 0. assert!(timestamp_nanos != 0, "irq_ring: zero IRQ timestamp"); - let slot = &SLOTS[percpu::cpu_id() as usize].0[source as usize]; + let slot = &SLOTS[percpu::cpu_id() as usize].0[source.slot()]; // ISRs run with IF=0, so this load-then-store can't interleave with a same-CPU `take`. if slot.load(Ordering::Relaxed) == 0 { slot.store(timestamp_nanos, Ordering::Relaxed); @@ -45,7 +49,7 @@ pub fn isr_publish(source: IrqSource, timestamp_nanos: u64) { /// Consumes the current CPU's pending record for `source`, returning its IRQ-time timestamp. pub fn take(source: IrqSource) -> Option { - let slot = &SLOTS[percpu::cpu_id() as usize].0[source as usize]; + let slot = &SLOTS[percpu::cpu_id() as usize].0[source.slot()]; // Atomic swap: an interrupting ISR sees either the old record or the cleared slot, never a torn value. match slot.swap(0, Ordering::Relaxed) { 0 => None, @@ -59,7 +63,7 @@ pub fn take(source: IrqSource) -> Option { /// True if `source` has an undrained record on this CPU; non-consuming, unlike [`take`]. pub fn pending(source: IrqSource) -> bool { - SLOTS[percpu::cpu_id() as usize].0[source as usize].load(Ordering::Relaxed) != 0 + SLOTS[percpu::cpu_id() as usize].0[source.slot()].load(Ordering::Relaxed) != 0 } /// True if any IRQ record is undrained on the current CPU; non-consuming. diff --git a/kernel/src/object/ops.rs b/kernel/src/object/ops.rs index 991100e4b2b..1d4265044b9 100644 --- a/kernel/src/object/ops.rs +++ b/kernel/src/object/ops.rs @@ -17,7 +17,8 @@ use crate::time::Deadline; use crate::pipe::{self, PipeId}; use crate::process::PipeMap; use crate::user_ptr::{UserBytes, UserBytesMut}; -use crate::watch::Watch; +use crate::inbox::PollEntry; +use crate::watch::{IrqWatch, Watch}; use crate::{device as device_registry, keyboard, mouse}; use super::device::DeviceClaim; @@ -242,14 +243,24 @@ pub fn pipe_write(object: &KObjectRef) -> Option<(PipeId, WaitClass)> { pub enum WatchRef { Static(&'static Watch), Shared(Arc), + /// A device's, which its interrupt handler posts. + Irq(&'static IrqWatch), } -impl core::ops::Deref for WatchRef { - type Target = Watch; - fn deref(&self) -> &Watch { +impl WatchRef { + pub(crate) fn add_poll(&self, entry: PollEntry) { match self { - Self::Static(watch) => watch, - Self::Shared(watch) => watch, + Self::Static(watch) => watch.add_poll(entry), + Self::Shared(watch) => watch.add_poll(entry), + Self::Irq(watch) => watch.add_poll(entry), + } + } + + pub fn cancel_polls(&self) { + match self { + Self::Static(watch) => watch.cancel_polls(), + Self::Shared(watch) => watch.cancel_polls(), + Self::Irq(watch) => watch.cancel_polls(), } } } @@ -266,10 +277,10 @@ pub fn read_watch(object: &KObjectRef) -> Option { device_registry::DeviceType::Keyboard => Some(WatchRef::Static(&keyboard::WATCH)), device_registry::DeviceType::Mouse => Some(WatchRef::Static(&mouse::WATCH)), device_registry::DeviceType::PciFunction => { - d.pci_slot().map(|slot| WatchRef::Static(crate::pcidev::watch(slot))) + d.pci_slot().map(|slot| WatchRef::Irq(crate::pcidev::watch(slot))) } device_registry::DeviceType::HdaAudio | device_registry::DeviceType::VirtioSound => { - Some(WatchRef::Static(&crate::drivers::AUDIO_WATCH)) + Some(WatchRef::Irq(&crate::drivers::AUDIO_WATCH)) } device_registry::DeviceType::Framebuffer => None, // A partition answers its description and has nothing to wait for. diff --git a/kernel/src/pcidev/mod.rs b/kernel/src/pcidev/mod.rs index 5868e669a32..35a4e3afaec 100644 --- a/kernel/src/pcidev/mod.rs +++ b/kernel/src/pcidev/mod.rs @@ -117,7 +117,7 @@ use crate::mm::policy::{CachePolicy, MmioPolicy}; use crate::mm::{align_2m, DirectMap, Mmio, PAGE_2M}; use crate::object::shm::{Region, SharedMemObject}; use crate::sync::Lock; -use crate::watch::Watch; +use crate::watch::IrqWatch; /// How many functions this machine can hand out at once. /// @@ -280,7 +280,7 @@ static BOUND: [Lock>; MAX_FUNCTIONS] = /// What a claimed function's poll waits on, one per slot: two processes each driving a /// function must not learn when the other's device is busy. -static WATCHES: [Watch; MAX_FUNCTIONS] = [const { Watch::new() }; MAX_FUNCTIONS]; +static WATCHES: [IrqWatch; MAX_FUNCTIONS] = [const { IrqWatch::new() }; MAX_FUNCTIONS]; /// Every function this machine enumerated, and the two windows a BAR may be /// moved into. @@ -1727,7 +1727,11 @@ pub fn take_record(slot: usize) -> Result, SyscallError> if IRQ[slot].faulted() { return Err(SyscallError::Io); } - Ok(IRQ[slot].take().map(|count| DeviceIrqRecord { count })) + let taken = IRQ[slot].take(); + if taken.is_some() && IRQ[slot].take_unannounced() { + log!("pcidev: slot {slot} took its first message on vector {:#x}", VECTORS[slot]); + } + Ok(taken.map(|count| DeviceIrqRecord { count })) } /// Whether a read of the claim answers at once: a message is waiting, or the @@ -1736,46 +1740,26 @@ pub fn has_irq(slot: usize) -> bool { IRQ[slot].armed() || IRQ[slot].faulted() } -/// Records one message. Called from the vector's ISR, so it takes no lock and -/// allocates nothing; `record.rs` owns the counting, and `kernel-loom` models -/// it against a concurrent reader. +/// Records one message and posts the claim's watch. Called from the vector's +/// handler, so it allocates nothing; `record.rs` owns the counting, and +/// `kernel-loom` models it against a concurrent reader. pub fn isr(slot: usize) { IRQ[slot].took(); -} - -/// Turn every message taken since the last pass into a wake. -/// -/// On the scheduler pass rather than in the ISR, like every other device in -/// this kernel: a wake takes the inbox lock and an ISR may not. -pub fn drain_pending() { - for (slot, irq) in IRQ.iter().enumerate() { - if !irq.take_pending() { - continue; - } - // A fault's wake is no message. - if !irq.faulted() && irq.take_unannounced() { - log!( - "pcidev: slot {slot} took its first message on vector {:#x}", - VECTORS[slot] - ); - } - WATCHES[slot].post(); - } + WATCHES[slot].post_in_place(); } /// The unit refused this function an access. /// -/// Called from the fault handler, which takes no lock: every call the claim -/// answers refuses from here on, its interrupt read included, and this CPU's -/// next scheduler pass wakes whoever waits on the claim to read that refusal — -/// the pass a message earns, posted the way its ISR posts it. +/// Called from the fault handler: every call the claim answers refuses from +/// here on, its interrupt read included, and the post wakes whoever waits on +/// the claim to read that refusal, as a message's does. pub fn note_fault(slot: usize) { IRQ[slot].fault(); - crate::irq_ring::isr_publish(crate::irq_ring::IrqSource::UserDev, crate::clock::nanos_since_boot()); + WATCHES[slot].post_in_place(); crate::preempt::set_need_resched(); } /// The watch of the function a claim holds at `slot`. -pub fn watch(slot: usize) -> &'static Watch { +pub fn watch(slot: usize) -> &'static IrqWatch { &WATCHES[slot] } diff --git a/kernel/src/pcidev/record.rs b/kernel/src/pcidev/record.rs index 51f27139a2d..72ff55f33b7 100644 --- a/kernel/src/pcidev/record.rs +++ b/kernel/src/pcidev/record.rs @@ -6,20 +6,16 @@ //! not an ordering, and no guest test in this suite lands on it. //! //! **Two parties race**: the ISR, on whichever CPU the unit routed the message -//! to, and the holder reading its record through a syscall on any CPU. The -//! scheduler pass that turns a message into a wake runs on the ISR's own CPU -//! after it, so those two do not interleave. +//! to, and the holder reading its record through a syscall on any CPU. //! -//! **The invariant is that every message is counted exactly once, and owes -//! exactly one wake.** Both are read-modify-writes and neither is a load -//! followed by a store: a reader that loaded a count and then cleared it drops -//! every message the ISR recorded in between, and a driver that misses one -//! waits for a device that has already spoken. No ordering carries anything -//! across these words — each is the whole of what it says — so the orderings -//! here are `Relaxed` and the model is about the interleaving, **but for one -//! edge**: a fault arms the same wake a message does, and the pass that takes -//! that wake has to read the fault, so `pending` is released by [`Interrupt::fault`] -//! and acquired by [`Interrupt::take_pending`]. +//! **The invariant is that every message is counted exactly once.** The count +//! is a read-modify-write on both sides and never a load followed by a store: +//! a reader that loaded a count and then cleared it drops every message the ISR +//! recorded in between, and a driver that misses one waits for a device that +//! has already spoken. No ordering carries anything across these words — each +//! is the whole of what it says — so the orderings here are `Relaxed` and the +//! model is about the interleaving; the holder reads them after the wake the +//! claim's watch post owes it, which orders them. #[cfg(not(feature = "loom"))] use core::sync::atomic::{AtomicBool, AtomicU32, Ordering}; @@ -32,16 +28,13 @@ use loom::sync::atomic::{AtomicBool, AtomicU32, Ordering}; /// control for and `kernel-loom` is the model of. const ORDER: Ordering = Ordering::Relaxed; -/// The negative control: the two read-modify-writes become a load and a store, +/// The negative control: the read-modify-writes become a load and a store, /// which is the whole of what this record's design is. Never on in a kernel /// build. #[cfg(feature = "device-irq-lossy")] macro_rules! take_word { - ($word:expr, $empty:expr) => { - take_word!($word, $empty, ORDER) - }; - ($word:expr, $empty:expr, $order:expr) => {{ - let held = $word.load($order); + ($word:expr, $empty:expr) => {{ + let held = $word.load(ORDER); $word.store($empty, ORDER); held }}; @@ -49,10 +42,7 @@ macro_rules! take_word { #[cfg(not(feature = "device-irq-lossy"))] macro_rules! take_word { ($word:expr, $empty:expr) => { - take_word!($word, $empty, ORDER) - }; - ($word:expr, $empty:expr, $order:expr) => { - $word.swap($empty, $order) + $word.swap($empty, ORDER) }; } @@ -72,12 +62,10 @@ macro_rules! bump { /// What the ISR writes and the claim reads back. /// -/// Atomics only: the handler takes no lock and allocates nothing. +/// Atomics only: the handler allocates nothing. pub struct Interrupt { /// Messages since the holder's last read. count: AtomicU32, - /// Set by the ISR, cleared by the scheduler pass that turns it into a wake. - pending: AtomicBool, /// The unit refused this function an access. Every call the claim answers /// refuses from here on: its bus mastering is gone, so a driver that kept /// going would be driving nothing. @@ -97,7 +85,6 @@ impl Interrupt { pub const fn new() -> Self { Self { count: AtomicU32::new(0), - pending: AtomicBool::new(false), faulted: AtomicBool::new(false), unannounced: AtomicBool::new(true), } @@ -109,7 +96,6 @@ impl Interrupt { pub fn new() -> Self { Self { count: AtomicU32::new(0), - pending: AtomicBool::new(false), faulted: AtomicBool::new(false), unannounced: AtomicBool::new(true), } @@ -122,12 +108,6 @@ impl Interrupt { /// two of these, and what it took plus what is left has to be what arrived. pub fn took(&self) { bump!(self.count); - // `swap` and not a store: [`Self::fault`] releases through this same - // word, and a plain write landing after that release in `pending`'s - // modification order ends the release sequence there — the pass that - // later takes the fault's wake would then synchronize with nothing. - // An RMW extends the sequence instead, whichever order it lands in. - self.pending.swap(true, ORDER); } /// The messages since the last read, or `None` for none. @@ -147,36 +127,19 @@ impl Interrupt { self.count.load(ORDER) != 0 } - /// Whether a wake is owed, taken at most once per message. Answers `true` - /// for the pass that owes it and `false` for every pass after. - /// - /// `swap` for the same reason as [`Self::take`]: two passes that both - /// loaded `true` would both wake one message's watchers. `Acquire`, so - /// the pass that takes a fault's wake reads [`Self::faulted`] set. - pub fn take_pending(&self) -> bool { - take_word!(self.pending, false, Ordering::Acquire) - } - /// Whether this is the first message this slot has taken. Answers `true` /// once per claim and `false` ever after, so a caller may log on it. /// - /// `swap` for [`Self::take_pending`]'s reason: two passes that both loaded - /// `true` would both announce one message. + /// `swap` for [`Self::take`]'s reason: two reads that both loaded `true` + /// would both announce one message. pub fn take_unannounced(&self) -> bool { take_word!(self.unannounced, false) } - /// The unit refused this function an access. Called from the fault handler, - /// which takes no lock: every call the claim answers refuses from here on, - /// and a wake is owed as for a message, because a holder waiting on the - /// claim would otherwise wait for a function that can no longer speak. - /// - /// The wake is a `swap` and not a store: this one races a pass on another - /// CPU, and loom 0.7 lets a plain store be lost to a concurrent `swap`, which - /// C11 forbids, so a store here is a wake the model cannot show is owed. + /// The unit refused this function an access. Called from the fault handler: + /// every call the claim answers refuses from here on. pub fn fault(&self) { self.faulted.store(true, ORDER); - self.pending.swap(true, Ordering::Release); } pub fn faulted(&self) -> bool { @@ -187,7 +150,6 @@ impl Interrupt { /// up. The holder is not running at either point. pub fn clear(&self) { self.count.store(0, ORDER); - self.pending.store(false, ORDER); self.faulted.store(false, ORDER); self.unannounced.store(true, ORDER); } diff --git a/kernel/src/sched/driver.rs b/kernel/src/sched/driver.rs index d917a8c5ff9..1c4cdf8f562 100644 --- a/kernel/src/sched/driver.rs +++ b/kernel/src/sched/driver.rs @@ -671,18 +671,6 @@ fn drain_irqs(entered: super::dump::Entered) { super::dump::serve_if_owed(); // Repaints the panel if whoever owns the screen has drawn over the report. crate::drivers::panic_console::hold_report(); - - if crate::irq_ring::take(crate::irq_ring::IrqSource::UserDev).is_some() { - // Which claim it was is the per-slot flag `pcidev` keeps; the record - // here says only that a pass is owed, so one function's interrupt does - // not wake every user driver in the machine. - crate::pcidev::drain_pending(); - } - if crate::irq_ring::take(crate::irq_ring::IrqSource::Audio).is_some() { - // Both backends share one watch, so a second would need the parking side - // to know which driver bound, which it doesn't. - crate::drivers::AUDIO_WATCH.post(); - } } /// Leave the current stack for this CPU's idle stack and never come back. @@ -715,6 +703,10 @@ extern "C" fn idle_loop() -> ! { if crate::drivers::panic_console::probe_due() { panic!("metal-panic-probe: a fatal report over a desktop that owns the screen"); } + #[cfg(feature = "boot-actuators")] + if crate::actuator::handler_post() { + crate::watch::handler_post::run(); + } crate::scheduler::log_health(); crate::scheduler::reap_finished(); // `pass` below covers this too; here as well so a CPU that diff --git a/kernel/src/sched/payload.rs b/kernel/src/sched/payload.rs index 523e2a6ede0..34d737df6bb 100644 --- a/kernel/src/sched/payload.rs +++ b/kernel/src/sched/payload.rs @@ -8,7 +8,7 @@ use core::sync::atomic::{AtomicU32, AtomicU64, Ordering}; use toyos_sched::fair::{FairShare, ShareState}; use toyos_sched::hw::Nanos; use toyos_sched::msg::Msg; -use toyos_sched::sync::LeafLock; +use toyos_sched::sync::CellLock; use toyos_sched::task::{SchedPayload, TaskAccounting, TaskShared, WaitClass}; use toyos_sched::park::WaitTicket; @@ -28,7 +28,7 @@ impl KernelLock { } } -impl LeafLock for KernelLock { +impl CellLock for KernelLock { fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { f(&mut self.0.lock()) } diff --git a/kernel/src/watch.rs b/kernel/src/watch.rs index f8bf0aee262..a58e8e8a42f 100644 --- a/kernel/src/watch.rs +++ b/kernel/src/watch.rs @@ -9,9 +9,10 @@ //! record of what was posted — a waiter re-reads the object, never the post. //! //! **A post allocates nothing and may be made under any lock but a poll -//! ring's** (`crate::inbox`), so a driver posts from its scheduler-pass drain -//! and never from its ISR, which publishes to `irq_ring` and nothing else. -//! Registration allocates, in the syscall that registers. +//! ring's** (`crate::inbox`). A watch an interrupt handler posts is an +//! [`IrqWatch`]: its list sits behind an [`IrqLock`], and its post is made in +//! place and frees nothing. Every other watch's list lock leaves interrupts +//! open, and no handler takes it. Registration allocates, in the syscall that registers. //! //! A [`Watch`] is a borrowed reference for the whole of a wait: [`Armed`] //! holds it, so an object cannot be freed under a thread waiting on it, and @@ -20,6 +21,7 @@ use alloc::sync::Arc; use toyos_sched::hw::Nanos; +use toyos_sched::sync::CellLock; use toyos_sched::task::{Refused, WaitClass, WakeCause, WakeReason}; use toyos_sched::watch::{Poster, Waiters}; @@ -32,14 +34,66 @@ use crate::sched::payload::{KMsg, KShared, KernelLock, TaskHandle}; use crate::scheduler::Parkable; use crate::time::Deadline; -type Inner = toyos_sched::watch::Watch>>; +type List = Waiters; -/// What an object holds to be waitable. -pub struct Watch(Inner); +/// What an object holds to be waitable, its list behind `L`. +pub struct Waitable>(toyos_sched::watch::Watch); + +/// A watch no interrupt handler reaches. +pub type Watch = Waitable>; + +/// 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`] +/// 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. +pub struct IrqLock(masked::Masked); + +impl IrqLock { + pub const fn new(value: T) -> Self { + Self(masked::Masked::new(value)) + } +} + +impl CellLock for IrqLock { + fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { + // Raised first, so the lock's own release never reaches depth zero, and + // a pass, with interrupts masked. + preempt_off(|_| { + let irq = crate::arch::IrqGuard::close(); + let mut held = self.0.lock(&irq); + #[cfg(feature = "boot-actuators")] + handler_post::raise_if_staged(); + f(&mut held) + }) + } +} + +mod masked { + use crate::arch::IrqGuard; + use crate::sync::{Lock, LockGuard}; + + /// A lock taken only through a borrow of a closed [`IrqGuard`]: never + /// before the mask, and never held past it. + pub struct Masked(Lock); + + impl Masked { + pub const fn new(value: T) -> Self { + Self(Lock::new(value)) + } + + pub fn lock<'a>(&'a self, _closed: &'a IrqGuard) -> LockGuard<'a, T> { + self.0.lock() + } + } +} impl Watch { pub const fn new() -> Self { - Self(Inner::new(KernelLock::new(Waiters::new()))) + Self(toyos_sched::watch::Watch::new(KernelLock::new(Waiters::new()))) } /// Something about the object changed: wake every thread waiting on it and @@ -82,23 +136,51 @@ impl Watch { ) }) } +} + +impl IrqWatch { + pub const fn new() -> Self { + Self(toyos_sched::watch::Watch::new(IrqLock::new(Waiters::new()))) + } + /// Something about the object changed: wake every thread waiting on it and + /// complete 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) { + #[cfg(feature = "boot-actuators")] + handler_post::note_post(self); + preempt_off(|p| { + let env = Poster { cpus: cpus(), kicker: &HW, preempt: p }; + self.0.post_in_place(WakeCause::new(WakeReason::Woken), &env); + }); + } +} + +impl> Waitable { /// A poll ring's entry, from `inbox`'s registration and nowhere else. pub(crate) fn add_poll(&self, entry: PollEntry) { self.0.add_ring(entry); } /// Answer every poll registered here as gone: the source it watched ended. + /// A thread's, since it frees what it answered. pub fn cancel_polls(&self) { self.0.cancel_rings(); } + + /// `handler-post`'s stand where a registration holds the list lock. + #[cfg(feature = "boot-actuators")] + pub(crate) fn holding(&self, f: impl FnOnce()) { + self.0.holding(f); + } } /// A thread's registration on one watch, held across its wait and ended by /// its drop. #[must_use = "a registration must outlive the park it was made for"] -pub struct Armed<'a> { - watch: &'a Watch, +pub struct Armed<'a, L: CellLock = KernelLock> { + watch: &'a Waitable, shared: Arc, task: Arc, /// Wait class for the blocked-time breakdown; the park carries no subject @@ -106,7 +188,7 @@ pub struct Armed<'a> { class: WaitClass, } -impl Drop for Armed<'_> { +impl> Drop for Armed<'_, L> { fn drop(&mut self) { self.watch.0.unregister(&self.shared); } @@ -114,7 +196,11 @@ impl Drop for Armed<'_> { /// Register the running task on `watch`; `None` when there is no current task. /// Call before reading the condition the wait is for. -pub fn arm(watch: &Watch, token: u64, class: WaitClass) -> Option> { +pub fn arm>( + watch: &Waitable, + token: u64, + class: WaitClass, +) -> Option> { let task = crate::sched::driver::current_handle()?; let shared = crate::sched::driver::current_shared()?; watch.0.register(&shared, token); @@ -144,7 +230,11 @@ fn not_revocable() -> ! { /// this thread is cancelled. A return is not an answer — the caller re-reads /// its condition. #[track_caller] -pub fn wait(p: &Parkable, armed: &Armed<'_>, deadline: Deadline) -> Result<(), Cancelled> { +pub fn wait>( + p: &Parkable, + armed: &Armed<'_, L>, + deadline: Deadline, +) -> Result<(), Cancelled> { match wait_inner(p, armed, deadline, Cancel::Answers) { Ok(()) => Ok(()), Err(Ended::Cancelled) => Err(Cancelled(())), @@ -154,7 +244,7 @@ pub fn wait(p: &Parkable, armed: &Armed<'_>, deadline: Deadline) -> Result<(), C /// The same as [`wait`], for a wait a kill may not end. #[track_caller] -pub fn wait_uncancellable(p: &Parkable, armed: &Armed<'_>, deadline: Deadline) { +pub fn wait_uncancellable>(p: &Parkable, armed: &Armed<'_, L>, deadline: Deadline) { match wait_inner(p, armed, deadline, Cancel::Ignores) { Ok(()) => {} Err(Ended::Cancelled) => unreachable!("an uncancellable wait never reports a cancel"), @@ -168,9 +258,9 @@ pub fn wait_uncancellable(p: &Parkable, armed: &Armed<'_>, deadline: Deadline) { /// be waited on again safely, and the caller reads its own condition again to /// tell them apart. #[track_caller] -pub fn wait_until( +pub fn wait_until>( p: &Parkable, - watch: &Watch, + watch: &Waitable, token: u64, class: WaitClass, deadline: Deadline, @@ -210,7 +300,7 @@ mod window { use toyos_sched::task::WaitClass; use toyos_sched::watch::window::{HELD, STEP}; - use super::Armed; + use super::{Armed, CellLock, List}; use crate::time::{Budget, Deadline, Duration}; /// A pipe wait nothing posts — a reader whose writer is idle — ends its @@ -223,7 +313,7 @@ mod window { static POSTED: AtomicU64 = AtomicU64::new(0); static LAPSED: AtomicU64 = AtomicU64::new(0); - pub(super) fn hold(armed: &Armed<'_>) { + pub(super) fn hold>(armed: &Armed<'_, L>) { if !crate::actuator::watch_window() || armed.class != WaitClass::Pipe { return; } @@ -247,6 +337,119 @@ mod window { } } +/// `handler-post`: the last claim slot's vector, which no claim holds, raised +/// on this CPU while it holds preemption off, posts that slot's watch from the +/// handler before any pass can run. Raised inside a post of the watch itself, +/// the handler's post follows once the outer one lets go; raised inside a +/// completion written into a ring that polls the watch, or inside that ring's +/// own watch's list lock, the handler's post completes that poll once the +/// section lets go. A hold counts the posts of the watch made on its CPU, and +/// no other watch's, and lapses at its budget. One run of [`HOLDS`] holds per +/// arm, on whichever idle loop reaches it first with interrupts open; its +/// verdict is one [`Verdict`] line. +#[cfg(feature = "boot-actuators")] +pub mod handler_post { + use core::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering::Relaxed}; + + use toyos_sched::watch::handler_post::{Verdict, HOLDS}; + + use crate::pcidev::{MAX_FUNCTIONS, VECTORS}; + use crate::time::{Budget, Deadline, Duration}; + + const SLOT: usize = MAX_FUNCTIONS - 1; + + const WINDOW: Budget = Budget::of( + Duration::from_secs(1), + "the hold is counted as lapsed, and the verdict line says so", + ); + + const NOBODY: u32 = u32::MAX; + static RAN: AtomicBool = AtomicBool::new(false); + /// The CPU whose next interrupts-off section raises the vector inside itself. + static RAISE_INSIDE: AtomicU32 = AtomicU32::new(NOBODY); + static HOLDING: AtomicU32 = AtomicU32::new(NOBODY); + static POSTS: AtomicU64 = AtomicU64::new(0); + + pub fn run() { + // A hold is a CPU that takes interrupts while it holds. + if RAN.load(Relaxed) || !crate::arch::cpu::interrupts_enabled() || RAN.swap(true, Relaxed) { + return; + } + // A held slot's driver would take counts its device never raised. + assert!(crate::pcidev::held_at(SLOT).is_none(), "handler-post: claim slot {SLOT} is held"); + let me = crate::arch::percpu::cpu_id(); + let claim = crate::pcidev::watch(SLOT); + let ring = crate::inbox::Staged::new(); + let verdict = crate::sched::driver::preempt_off(|_| { + HOLDING.store(me, Relaxed); + // The outer post is one of the two a hold waits for. + let in_a_list = holds(2, || { + RAISE_INSIDE.store(me, Relaxed); + claim.post_in_place(); + }); + let in_a_ring = holds(1, || { + ring.poll(claim); + RAISE_INSIDE.store(me, Relaxed); + ring.complete(); + }); + // Where a submitter's registration holds it. + let in_a_rings_watch = holds(1, || { + ring.poll(claim); + ring.holding_its_watch(raise); + }); + HOLDING.store(NOBODY, Relaxed); + Verdict { in_a_list, in_a_ring, in_a_rings_watch } + }); + crate::log!("{verdict}"); + } + + /// [`HOLDS`] holds of `stage`, answering how many saw `owed` posts of the + /// claim's watch before their budget. + fn holds(owed: u64, stage: impl Fn()) -> u32 { + let mut posted = 0; + for _ in 0..HOLDS { + let before = POSTS.load(Relaxed); + stage(); + let deadline = Deadline::at(crate::clock::now() + WINDOW.duration()); + loop { + if POSTS.load(Relaxed) >= before + owed { + posted += 1; + break; + } + if deadline.reached(crate::clock::now()) { + break; + } + core::hint::spin_loop(); + } + } + posted + } + + fn raise() { + crate::arch::irqchip::send_self(VECTORS[SLOT]); + } + + /// From inside an interrupts-off section: the staged CPU's next one raises + /// the vector while it holds. + pub fn raise_if_staged() { + let staged = RAISE_INSIDE.load(Relaxed); + if staged != NOBODY && staged == crate::arch::percpu::cpu_id() { + RAISE_INSIDE.store(NOBODY, Relaxed); + raise(); + } + } + + pub fn note_post(watch: &super::IrqWatch) { + let holding = HOLDING.load(Relaxed); + if holding != NOBODY + && holding == crate::arch::percpu::cpu_id() + && core::ptr::eq(watch, crate::pcidev::watch(SLOT)) + { + POSTS.fetch_add(1, Relaxed); + } + } +} + /// Register, then park until `ready()` holds, for a wait a kill may not end /// and no deadline bounds. #[track_caller] @@ -267,9 +470,9 @@ pub fn wait_uncancellable_until(p: &Parkable, watch: &Watch, token: u64, ready: } #[track_caller] -fn wait_inner( +fn wait_inner>( _p: &Parkable, - armed: &Armed<'_>, + armed: &Armed<'_, L>, deadline: Deadline, cancel: Cancel, ) -> Result<(), Ended> { diff --git a/src/ci.rs b/src/ci.rs index 73c9084e524..08e533522d4 100644 --- a/src/ci.rs +++ b/src/ci.rs @@ -268,7 +268,6 @@ pub(crate) const CONTROLS: &[Control] = &[ ]), red(KERNEL_LOOM, "device-irq-lossy", Some("device_irq"), &[ "every_message_is_counted_once ... FAILED", - "one_message_is_one_wake ... FAILED", ]), red(KERNEL_LOOM, "dump-report-relaxed", Some("dump_request"), &[ "a_request_filed_during_a_report_is_reported ... FAILED", @@ -290,6 +289,7 @@ pub(crate) const CONTROLS: &[Control] = &[ // with. A double panic, so the verdict is the first one's message. red(SCHED_LOOM, "commit-ignores-notify", Some("loom_watch"), &[ "parked with the condition true and no wake owed: the post was lost", + "parked with both completions written and no wake owed: a ring's post was lost", ]), // The notify's flagged arm answering off a load: a second post reads the // word from before the waiter consumed the first flag. @@ -304,8 +304,16 @@ pub(crate) const CONTROLS: &[Control] = &[ // ring entry. red(SCHED_LOOM, "poll-fire-load-store", Some("loom_watch"), &[ "a_poll_registered_racing_a_post_completes_exactly_once ... FAILED", + "a_poll_registered_racing_a_post_in_place_completes_exactly_once ... FAILED", "a_poll_on_two_watches_racing_both_posts_completes_exactly_once ... FAILED", ]), + // The ring models' lost-completion half: the producer posts before it + // stores the readiness its registrant rechecks. + red(SCHED_LOOM, "fault-posted-before-it-is-set", Some("loom_watch"), &[ + "a poll over a ready object was completed by neither", + "a_poll_registered_racing_a_post_completes_exactly_once ... FAILED", + "a_poll_registered_racing_a_post_in_place_completes_exactly_once ... FAILED", + ]), // Reproduces an open defect // (`issues/kernel/steal-probe-node-dies-with-its-victim.md`) rather than // proving a lie is caught, and goes with its fix. diff --git a/tests/toyos.rs b/tests/toyos.rs index d9134cf074e..f031b536ad4 100644 --- a/tests/toyos.rs +++ b/tests/toyos.rs @@ -1260,6 +1260,11 @@ const MACHINE_TESTS: &[(&str, Sched)] = &[ // waiter between reading its condition and parking, so the peer's post lands where // only the notified bit carries it to the commit. ("blocking_read_window", Sched::Parallel), + // A claim slot's vector posts its watch from the handler: raised on a CPU + // holding preemption off, inside a post of that watch, inside a completion + // into a ring polling it or inside that ring's own watch, it posts once + // that section lets go and before any pass. One boot; the verdict is counts. + ("handler_post_without_a_pass", Sched::Parallel), // A sibling's munmap and mmap staged between a typed copy's translation // and its store (`copy-meets-a-remap`): the store never reaches the region // mapped after it. @@ -12165,6 +12170,19 @@ fn run_machine_test( "smp_failed_ap_leaves_no_hole" => { smp_failed_ap_leaves_no_hole(test_config, c_bins, rust_bins) } + "handler_post_without_a_pass" => { + let mut qemu = QemuInstance::boot_with_options( + test_config, + c_bins, + rust_bins, + BootOptions { kernel_params: &["handler-post"], ..Default::default() }, + ); + let mut log = qemu.boot_log().to_string(); + if !log.lines().any(handler_post_said) { + log += &qemu.drain_until(Duration::from_secs(30), handler_post_said); + } + handler_post(&log) + } "input_merge" => { // The check runs in the kernel and panics on mismatch, so a // failure arrives as a dead boot; the marker is the only proof it @@ -14601,6 +14619,26 @@ fn sysret_ss(log: &str) -> Result<(), String> { Ok(()) } +use toyos_sched::watch::handler_post::{Verdict as HandlerPost, SAID as HANDLER_POST_SAID}; + +fn handler_post_said(line: &str) -> bool { + line.contains(HANDLER_POST_SAID) +} + +/// Every hold was posted into by the handler of the vector raised inside it. +fn handler_post(log: &str) -> Result<(), String> { + let Some(said) = log.lines().find(|line| handler_post_said(line)) else { + return Err(format!("`handler-post` never said its verdict:\n{log}")); + }; + if !said.contains(&HandlerPost::GREEN.to_string()) { + return Err(format!( + "a hold lapsed with no handler's post in it — the wake waited for a pass:\n{said}\n{log}" + )); + } + eprintln!(" [handler-post] {}", said.trim()); + Ok(()) +} + /// The input core merged what it was handed. /// /// Text in, a verdict out: every line it reads is a kernel record, so the diff --git a/toyos-sched/loom/Cargo.toml b/toyos-sched/loom/Cargo.toml index e6a978c80a9..1c9dcfe234c 100644 --- a/toyos-sched/loom/Cargo.toml +++ b/toyos-sched/loom/Cargo.toml @@ -101,6 +101,15 @@ gate-fence-off = [] # cargo test -p toyos-sched-loom --release --features poll-fire-load-store --test loom_watch # poll-fire-load-store = [] +# The control for `poll_racing`'s lost-completion half: its producer posts, then +# stores the readiness, the reverse of the order `pcidev::note_fault` keeps, so +# both `*_racing_a_post*_completes_exactly_once` models must red with a poll +# completed by neither: +# +# cargo test -p toyos-sched-loom --features fault-posted-before-it-is-set --test loom_watch +# +# `toyos-sched` declares no twin, because this decides nothing in `../src`. +fault-posted-before-it-is-set = [] [dependencies] loom = "0.7" diff --git a/toyos-sched/loom/src/model.rs b/toyos-sched/loom/src/model.rs index 4970c1bfe7f..f58573bb9fd 100644 --- a/toyos-sched/loom/src/model.rs +++ b/toyos-sched/loom/src/model.rs @@ -1,5 +1,5 @@ //! Scaffolding shared by the loom models: the message type, the modelled -//! preempt count, the leaf lock and the kick recorder. +//! preempt count, the cell lock and the kick recorder. use loom::sync::atomic::{AtomicUsize, Ordering}; use loom::sync::{Mutex, MutexGuard}; @@ -8,7 +8,7 @@ use crate::hw::{CpuId, Kicker}; use crate::mailbox::{PreemptGuard, SchedMsg}; use crate::sync::Arc; use crate::task::{TaskKey, TaskShared, WakeCause, WakeReason}; -use crate::sync::LeafLock; +use crate::sync::CellLock; use crate::watch::Waiters; pub const CPU0: CpuId = CpuId(0); @@ -113,7 +113,7 @@ pub struct RemoteGuard; #[allow(unsafe_code)] unsafe impl PreemptGuard for RemoteGuard {} -/// `LeafLock` over loom's mutex, so the watch models exercise the real +/// `CellLock` over loom's mutex, so the watch models exercise the real /// critical sections. pub struct LoomLock(Mutex); @@ -123,7 +123,7 @@ impl LoomLock { } } -impl LeafLock for LoomLock { +impl CellLock for LoomLock { fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { f(&mut self.0.lock().unwrap()) } diff --git a/toyos-sched/loom/tests/loom_watch.rs b/toyos-sched/loom/tests/loom_watch.rs index 2b09968f570..6c57e61b7f0 100644 --- a/toyos-sched/loom/tests/loom_watch.rs +++ b/toyos-sched/loom/tests/loom_watch.rs @@ -18,13 +18,16 @@ //! cargo test -p toyos-sched-loom --features commit-ignores-notify --test loom_watch //! ``` //! -//! Three more controls: `notify-flag-load-only` lets a post that finds its bits +//! Four more controls: `notify-flag-load-only` lets a post that finds its bits //! already set answer off a load, and //! `a_second_post_is_not_lost_to_a_flag_the_waiter_consumed` must red; //! `gate-fence-off` removes the [`Gate`]'s two fences, and -//! `a_transition_racing_an_opening_gate_is_never_missed` must red; and +//! `a_transition_racing_an_opening_gate_is_never_missed` must red; //! `poll-fire-load-store`, the kernel's own control for the poll's one-shot -//! answer, which the ring models below compile, must red both poll models here. +//! answer, which the ring models below compile, must red every +//! `*_completes_exactly_once` model here; and `fault-posted-before-it-is-set` +//! posts before the readiness is stored, and both `a_poll_registered_racing_*` +//! 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 @@ -49,6 +52,7 @@ use toyos_sched_loom::park::{prepare, Cancel, Commit, CurrentTask}; use toyos_sched_loom::task::{ Claim, Refused, TaskKey, TaskShared, TaskState, WaitClass, WakeCause, WakeReason, }; +use toyos_sched_loom::sync::CellLock; use toyos_sched_loom::watch::{Fire, Gate, Poster, Ring, Waiters, Watch}; #[path = "../../../kernel/src/inbox/once.rs"] @@ -109,6 +113,15 @@ impl World { self.watch.post(WakeCause::new(WakeReason::Woken), &env); } + fn post_in_place(&self) { + let env = Poster { + cpus: &self.cpus, + kicker: &self.kicks, + preempt: &RemoteGuard, + }; + self.watch.post_in_place(WakeCause::new(WakeReason::Woken), &env); + } + fn post_one(&self, token: u64) -> usize { let env = Poster { cpus: &self.cpus, @@ -274,53 +287,89 @@ fn a_bounded_post_racing_a_timeout_reaches_a_live_waiter() { /// and never by both, which is a completion the process did not ask for. #[test] fn a_poll_registered_racing_a_post_completes_exactly_once() { - model(|| { - let (world, _rx) = world(); - let ready = Arc::new(AtomicBool::new(false)); - let poll = Entry::new(); + model(|| poll_racing(World::post)); +} - let registrant = { - let world = world.clone(); - let ready = ready.clone(); - let poll = poll.clone(); - loom::thread::spawn(move || { - world.watch.add_ring(poll.clone()); - if ready.load(Ordering::Acquire) { - poll.fire(Fire::Ready); - } - }) - }; - let producer = loom::thread::spawn(move || { - ready.store(true, Ordering::Release); - world.post(); - }); - registrant.join().unwrap(); - producer.join().unwrap(); +/// The same, against the post an interrupt handler makes. +#[test] +fn a_poll_registered_racing_a_post_in_place_completes_exactly_once() { + model(|| poll_racing(World::post_in_place)); +} - assert_eq!(poll.posts(), 1, "a poll over a ready object completes once"); - }); +/// `ready` is `Relaxed` on both sides, as a claim's `faulted` is: the list lock +/// alone orders it against the registrant's recheck. +fn poll_racing(post: fn(&World)) { + let (world, _rx) = world(); + let ready = Arc::new(AtomicBool::new(false)); + let poll = Entry::new(); + + let registrant = { + let world = world.clone(); + let ready = ready.clone(); + let poll = poll.clone(); + loom::thread::spawn(move || { + world.watch.add_ring(poll.clone()); + if ready.load(Ordering::Relaxed) { + poll.fire(Fire::Ready); + } + }) + }; + let producer = { + let world = world.clone(); + loom::thread::spawn(move || { + // `fault-posted-before-it-is-set` is the control: posted first, the + // readiness can land after both the post and the recheck. + if cfg!(feature = "fault-posted-before-it-is-set") { + post(&world); + ready.store(true, Ordering::Relaxed); + } else { + ready.store(true, Ordering::Relaxed); + post(&world); + } + }) + }; + registrant.join().unwrap(); + producer.join().unwrap(); + + // While `world` lives: its watch's drop answers a live entry as gone. + assert_ne!(poll.posts(), 0, "a poll over a ready object was completed by neither"); + assert_eq!(poll.posts(), 1, "a poll over a ready object completes once"); + drop(world); } /// The object's end racing its readiness: the poll is answered once, as ready -/// or as gone, and whichever answered it no longer holds it. +/// or as gone. #[test] fn an_end_racing_a_post_answers_a_poll_once() { - model(|| { - let (world, _rx) = world(); - let poll = Entry::new(); - world.watch.add_ring(poll.clone()); + model(|| end_racing(|w| w.watch.cancel_rings(), World::post)); +} - let ender = { - let world = world.clone(); - loom::thread::spawn(move || world.watch.cancel_rings()) - }; - let poster = loom::thread::spawn(move || world.post()); - ender.join().unwrap(); - poster.join().unwrap(); +/// The same, for a thread's end racing a handler's post, which is made in +/// place. +#[test] +fn an_end_racing_a_post_in_place_answers_a_poll_once() { + model(|| end_racing(|w| w.watch.cancel_rings(), World::post_in_place)); +} - assert_eq!(poll.posts(), 1); - assert!(!poll.live()); - }); +fn end_racing(end: fn(&World), post: fn(&World)) { + let (world, _rx) = world(); + let poll = Entry::new(); + world.watch.add_ring(poll.clone()); + + let ender = { + let world = world.clone(); + loom::thread::spawn(move || end(&world)) + }; + let poster = { + let world = world.clone(); + loom::thread::spawn(move || post(&world)) + }; + ender.join().unwrap(); + poster.join().unwrap(); + + assert_eq!(poll.posts(), 1); + assert!(!poll.live()); + drop(world); } /// One waiter, two producers, each storing its own condition and then posting. @@ -594,3 +643,118 @@ 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. +struct PollRing { + cpus: CpuHandles, + kicks: Kicks, + written: LoomLock, + parked: RingWatch, +} + +/// One of that ring's polls, registered on one device's watch. +struct RingPoll { + ring: Arc, + state: once::Once, +} + +struct RingEntry(Arc); + +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); + } + } + + fn live(&self) -> bool { + self.0.state.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. +#[test] +fn two_posts_through_one_rings_lock_lose_no_wake() { + 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 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") + }) + }; + 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 parked = waiting.join().unwrap(); + for poster in posters { + poster.join().unwrap(); + } + let guard = PreemptModel::new(); + let msgs = drain(&mut rx, &guard); + + if parked { + assert_eq!( + msgs, + [Msg::Wake(TaskKey(1), WakeReason::Woken)], + "parked with both completions written and no wake owed: a ring's post was lost", + ); + } else { + assert!(msgs.is_empty(), "a submitter that never parked is owed nothing: {msgs:?}"); + } + ring.parked.unregister(&submitter); + }); +} diff --git a/toyos-sched/sim/src/payload.rs b/toyos-sched/sim/src/payload.rs index b0676b36c7d..c67b3871412 100644 --- a/toyos-sched/sim/src/payload.rs +++ b/toyos-sched/sim/src/payload.rs @@ -11,7 +11,7 @@ use std::sync::{Arc, Mutex}; use toyos_sched::fair::ShareState; use toyos_sched::mailbox::PreemptGuard; -use toyos_sched::sync::LeafLock; +use toyos_sched::sync::CellLock; use toyos_sched::task::{SchedPayload, TaskKey}; use toyos_sched::watch::{Fire, Ring, Waiters}; @@ -29,7 +29,7 @@ impl StdLock { } } -impl LeafLock for StdLock { +impl CellLock for StdLock { fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { f(&mut self.0.lock().expect("the simulator never poisons a lock")) } diff --git a/toyos-sched/src/cpu.rs b/toyos-sched/src/cpu.rs index 0cba829f61c..7db1c4d9e5b 100644 --- a/toyos-sched/src/cpu.rs +++ b/toyos-sched/src/cpu.rs @@ -2438,13 +2438,13 @@ mod tests { use crate::fair::{FairShare, ShareState}; use crate::hw::{Kicker, Machine}; use crate::mailbox::{mailbox, NoPreempt}; - use crate::sync::LeafLock; + use crate::sync::CellLock; use crate::task::{RtState, TaskAccounting, TaskBuilder}; use std::sync::Mutex; struct TestLock(Mutex); - impl LeafLock for TestLock { + impl CellLock for TestLock { fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { f(&mut self.0.lock().expect("a test never poisons a lock")) } diff --git a/toyos-sched/src/fair.rs b/toyos-sched/src/fair.rs index 95aa1256a3b..cbcd4603002 100644 --- a/toyos-sched/src/fair.rs +++ b/toyos-sched/src/fair.rs @@ -5,7 +5,7 @@ use core::num::NonZeroU32; use core::sync::atomic::{AtomicU64, Ordering}; -use crate::sync::LeafLock; +use crate::sync::CellLock; /// Clamp for stored lag at the Runnable→NonRunnable transition: how far /// behind (entitled catch-up) or ahead (throttled on wake) of the frontier a @@ -206,13 +206,13 @@ fn vrt_from_lag(frontier: u64, lag: i64) -> u64 { } /// One process's fair-share pot, reached through any thread that owns it. The -/// cell is supplied by the environment for the reason stated on [`LeafLock`]: +/// cell is supplied by the environment for the reason stated on [`CellLock`]: /// the kernel's is a word-sized spin, the simulator's a mutex. pub struct FairShare { state: L, } -impl> FairShare { +impl> FairShare { pub fn new(state: L) -> Self { Self { state } } diff --git a/toyos-sched/src/sync.rs b/toyos-sched/src/sync.rs index b0e1a982716..1a08ca17168 100644 --- a/toyos-sched/src/sync.rs +++ b/toyos-sched/src/sync.rs @@ -20,13 +20,12 @@ pub use loom::sync::Arc; /// Interior mutability for a small shared cell, supplied by the environment: /// a watch's waiter list and the per-process fair share. The kernel's -/// implementor is a few-instruction, IRQ-off leaf lock that acquires nothing -/// beneath it and is never held across a pass or a switch; the simulator and -/// the loom models supply their own. +/// implementor is never held across a pass or a switch; the simulator and the +/// loom models supply their own. /// /// It lives here rather than in one of its users because the core crate may /// not implement a lock itself — that would need `unsafe`, which only /// `mailbox.rs` is allowed to write. -pub trait LeafLock: Sync { +pub trait CellLock: Sync { fn with(&self, f: impl FnOnce(&mut T) -> R) -> R; } diff --git a/toyos-sched/src/task.rs b/toyos-sched/src/task.rs index 24f99799059..1c490d24cde 100644 --- a/toyos-sched/src/task.rs +++ b/toyos-sched/src/task.rs @@ -13,7 +13,7 @@ use crate::fair::{FairShare, ShareState, QUANTUM_NS}; use crate::hw::{CpuId, Nanos}; use crate::mailbox::MailboxNode; use crate::msg::Msg; -use crate::sync::{Arc, AtomicBool, AtomicU64, LeafLock, Ordering}; +use crate::sync::{Arc, AtomicBool, AtomicU64, CellLock, Ordering}; use crate::park::CommittedTicket; /// Monotonic, never reused. Stale messages keyed by `TaskKey` are provably @@ -31,8 +31,8 @@ pub trait SchedPayload: Sized + Send + 'static { /// The cell the per-process [`FairShare`] lives in. Supplied by the /// environment because the core crate may not implement a lock itself - /// (see [`LeafLock`]). - type ShareLock: LeafLock + Send; + /// (see [`CellLock`]). + type ShareLock: CellLock + Send; } /// Shorthand for the share type a payload implies. diff --git a/toyos-sched/src/watch.rs b/toyos-sched/src/watch.rs index 96867851337..6bd0a6f5b5e 100644 --- a/toyos-sched/src/watch.rs +++ b/toyos-sched/src/watch.rs @@ -15,16 +15,22 @@ //! waiting side does between its registration and its park can be skipped: //! `park::prepare` consumes the flag, and the commit consumes the claim. //! -//! **A post allocates nothing.** Every ring entry is one-shot, so a post takes -//! them all out of the list and fires them with the list lock let go; what it -//! frees is those entries. A registration sweeps the entries a withdrawal left -//! behind, so the list never holds more than the polls live at its last post or -//! registration, plus that one. +//! **Nothing allocates or frees under the list lock**, because an interrupt +//! handler's post can interrupt the allocator's holder while another CPU holds +//! the list waiting for the allocator. A registration takes the entries that +//! can no longer fire out [`FEW`] at a time, grows a full list in a buffer +//! allocated with the lock let go, and drops both with it let go. +//! [`Watch::post`] takes every ring entry out and fires and frees them with the +//! lock let go; [`Watch::post_in_place`], for a handler, which may not free at +//! all, fires them where they stand, and later registrations sweep the dead +//! out, since an entry is one-shot. //! -//! **Lock order.** The list lock is a leaf the environment supplies: -//! [`Ring::fire`] is always called with it let go, and so is the drop of every -//! entry the list lets go of, because an entry's last reference may own another -//! watch. So a post may be made under any lock but a ring's own. +//! **Lock order.** A post in place fires its rings under the list lock, so +//! beneath it are each ring's own lock and the watch that ring's submitters +//! park on, which holds threads and no ring and so nests nothing. The drop of +//! every entry the list lets go of runs with it let go, because an entry's last +//! reference may own another watch. So a post may be made under any lock but a +//! ring's own. use alloc::vec::Vec; use core::marker::PhantomData; @@ -33,7 +39,7 @@ use crate::cpu::CpuHandles; use crate::hw::Kicker; use crate::mailbox::{PreemptGuard, SchedMsg}; use crate::park::{notify, revoke}; -use crate::sync::{fence, Arc, AtomicU32, LeafLock, Ordering}; +use crate::sync::{fence, Arc, AtomicU32, CellLock, Ordering}; use crate::task::{TaskShared, WakeCause}; /// Why a ring entry is posted. @@ -50,14 +56,16 @@ pub enum Fire { pub trait Ring { /// Post this poll's completion. One-shot across every watch the poll is /// registered on: an entry that already fired, or whose poll was withdrawn, - /// posts nothing. Called with no watch's list lock held; may take only its - /// ring's own lock and post only the watch its ring's submitters park on. + /// posts 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. fn live(&self) -> bool; } -/// Everything a watch holds, behind the environment's leaf lock. +/// Everything a watch holds. pub struct Waiters { /// In registration order, so a bounded post reaches the longest waiter /// first. @@ -97,12 +105,16 @@ pub struct Poster<'a, M, K, P> { pub preempt: &'a P, } -pub struct Watch>> { +/// How many ring entries one section takes out of the list to drop with the +/// lock let go, on the stack and so allocating nothing. +pub const FEW: usize = 4; + +pub struct Watch>> { list: L, _msg: PhantomData (M, R)>, } -impl>> Watch { +impl>> Watch { pub const fn new(list: L) -> Self { Self { list, @@ -111,7 +123,7 @@ impl>> Watch { } } -impl>> Watch { +impl>> Watch { /// Register the running task for one wait. After this, and before its /// condition is read: that order is the whole lost-wake argument. A post /// that reached this task during an earlier wait is forgotten here, before @@ -119,14 +131,7 @@ impl>> Watch { pub fn register(&self, task: &Arc>, token: u64) { assert!(task.set_waiting(), "a task waits on at most one watch"); task.forget_posts(); - let dead = self.list.with(|w| { - w.threads.push(Waiter { - task: task.clone(), - token, - }); - sweep(&mut w.rings) - }); - drop(dead); + self.admit(Waiter { task: task.clone(), token }, |w| &mut w.threads); } /// End one wait. Idempotent against [`Self::revoke`], which may have taken @@ -144,28 +149,82 @@ impl>> Watch { /// after this returns and fires the entry itself if the object is already /// ready — the ring's half of the same order a thread keeps. pub fn add_ring(&self, entry: R) { - let dead = self.list.with(|w| { - let dead = sweep(&mut w.rings); - w.rings.push(entry); - dead - }); - drop(dead); + self.admit(entry, |w| &mut w.rings); } /// Something changed: wake every registered thread, and fire and let go of - /// every ring entry. + /// every ring entry, with the list lock let go. pub fn post(&self, cause: WakeCause, env: &Poster<'_, M, K, P>) { - let fired = self.list.with(|w| { + let mut few: [Option; FEW] = [const { None }; FEW]; + let many = self.list.with(|w| { for waiter in &w.threads { notify(&waiter.task, cause, env.cpus, env.kicker, env.preempt); } - core::mem::take(&mut w.rings) + // [`FEW`] leave the list's buffer to the next registration. + if w.rings.len() > FEW { + return core::mem::take(&mut w.rings); + } + // Indexed, so an entry with no slot panics rather than being + // dropped unfired. + for (at, ring) in w.rings.drain(..).enumerate() { + few[at] = Some(ring); + } + Vec::new() }); - for ring in &fired { + for ring in few.iter().flatten().chain(&many) { ring.fire(Fire::Ready); } } + /// [`Self::post`] for a context that may not free, an interrupt handler: + /// every ring entry is fired where it stands, under the list lock, so the + /// cost is every registration on the watch with that lock held. + pub fn post_in_place( + &self, + cause: WakeCause, + env: &Poster<'_, M, K, P>, + ) { + self.list.with(|w| { + for waiter in &w.threads { + notify(&waiter.task, cause, env.cpus, env.kicker, env.preempt); + } + for ring in &w.rings { + ring.fire(Fire::Ready); + } + }); + } + + /// Put `item` on the list `list` picks, taking out up to [`FEW`] ring + /// entries that can no longer fire first. One section while the list has + /// room; a full one grows in a buffer allocated with the lock let go, and + /// what a section took out is dropped with it let go. + fn admit(&self, item: T, list: fn(&mut Waiters) -> &mut Vec) { + let mut item = item; + let mut room = Vec::new(); + loop { + let mut dead: [Option; FEW] = [const { None }; FEW]; + let full = self.list.with(|w| { + take_dead(&mut w.rings, &mut dead); + let v = list(w); + // Only a buffer that holds the whole list and `item`: another + // registration may have outgrown it since it was sized. + if v.len() == v.capacity() && room.capacity() > v.len() { + room.append(v); + core::mem::swap(v, &mut room); + } + if v.len() < v.capacity() { + v.push(item); + return None; + } + Some((v.capacity(), item)) + }); + drop(dead); + let Some((seen, back)) = full else { return }; + item = back; + room = Vec::with_capacity((seen * 2).max(FEW)); + } + } + /// Wake at most `limit` threads registered with `token`, in registration /// order, and answer how many this post woke. A thread that was not parked /// is flagged and spends nothing: its own commit rechecks, and the post @@ -233,6 +292,12 @@ impl>> Watch { } } + /// Run `f` holding the list lock, where a registration holds it, touching + /// nothing on the list: an actuator's way to raise an interrupt there. + pub fn holding(&self, f: impl FnOnce() -> U) -> U { + self.list.with(|_| f()) + } + /// Registered threads, for a model or a report. pub fn threads(&self) -> usize { self.list.with(|w| w.threads.len()) @@ -244,6 +309,22 @@ impl>> Watch { } } +/// Take up to [`FEW`] entries that can no longer fire out of `rings` into +/// `dead`, and answer how many. +fn take_dead(rings: &mut Vec, dead: &mut [Option; FEW]) -> usize { + let (mut at, mut taken) = (0, 0); + // Order is nothing to a ring entry: every post fires them all. + while at < rings.len() && taken < FEW { + if rings[at].live() { + at += 1; + } else { + dead[taken] = Some(rings.swap_remove(at)); + taken += 1; + } + } + taken +} + /// A word in front of a watch that its posters read to learn whether a post is /// owed at all, for a waiter whose condition is a sweep of words the posters /// write: the machine's stop, whose posters are every thread's park, band and @@ -297,15 +378,9 @@ fn gate_fence() { } } -/// Take out the entries that can no longer fire, for the caller to drop once -/// the list lock is let go. -fn sweep(rings: &mut Vec) -> Vec { - rings.extract_if(.., |r| !r.live()).collect() -} - /// An object that ends with polls still registered answers them: dropping a /// watch fires every live ring entry as [`Fire::Gone`]. -impl>> Drop for Watch { +impl>> Drop for Watch { fn drop(&mut self) { let ended = self.list.with(|w| core::mem::take(&mut w.rings)); for ring in &ended { @@ -326,6 +401,45 @@ pub mod window { pub const STEP: u64 = 64; } +/// The `handler-post` actuator's line: the kernel writes it once, and the +/// harness compares it whole against [`Verdict::GREEN`]. +/// +/// [`Verdict::GREEN`]: handler_post::Verdict::GREEN +pub mod handler_post { + use core::fmt; + + /// The line's first word, which the harness waits for. + pub const SAID: &str = "handler-post:"; + /// Holds staged per arm. + pub const HOLDS: u32 = 4; + + /// The holds a handler's post ended, per arm: its vector raised inside a + /// watch's list lock, inside a ring's completions, and inside the list + /// lock of the watch that ring's own submitters park on. + #[derive(Clone, Copy)] + pub struct Verdict { + pub in_a_list: u32, + pub in_a_ring: u32, + pub in_a_rings_watch: u32, + } + + impl Verdict { + pub const GREEN: Self = + Self { in_a_list: HOLDS, in_a_ring: HOLDS, in_a_rings_watch: HOLDS }; + } + + impl fmt::Display for Verdict { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + f, + "{SAID} of {HOLDS} holds per arm, a handler posted into {} inside a list lock, \ + {} inside a ring's completions and {} inside a ring's own watch", + self.in_a_list, self.in_a_ring, self.in_a_rings_watch, + ) + } + } +} + #[cfg(test)] mod tests { extern crate std; @@ -356,7 +470,7 @@ mod tests { } struct StdLock(Mutex); - impl LeafLock for StdLock { + impl CellLock for StdLock { fn with(&self, f: impl FnOnce(&mut T) -> U) -> U { f(&mut self.0.lock().unwrap()) } @@ -607,7 +721,7 @@ mod tests { } struct SharedLock(std::sync::Arc>>); - impl LeafLock> for SharedLock { + impl CellLock> for SharedLock { fn with(&self, f: impl FnOnce(&mut Waiters) -> U) -> U { f(&mut self.0.lock().unwrap()) } @@ -644,19 +758,49 @@ mod tests { assert_eq!(w.live_rings(), 1); } + /// A post in place fires where the entry stands and lets go of nothing: + /// the entry it fired is still the list's until a registration sweeps it. + #[test] + fn a_post_in_place_fires_and_drops_nothing() { + let (handles, _rx) = cpus(); + let env = Poster { cpus: &handles, kicker: &NoKick, preempt: &NoPreempt }; + let w = watch(); + let poll = Arc::new(Poll::default()); + w.add_ring(poll.clone()); + w.post_in_place(woken(), &env); + assert_eq!(poll.posts.load(Ordering::Acquire), 1); + assert_eq!(Arc::strong_count(&poll), 2, "the post let go of the entry it fired"); + let t = task(1); + w.register(&t, 0); + assert_eq!(Arc::strong_count(&poll), 1, "the registration did not sweep it"); + w.unregister(&t); + } + + /// One entry past the [`FEW`] a post takes out on its stack, so a cancel + /// or a drop bounded at that many is caught. #[test] fn cancel_and_drop_answer_every_live_poll_as_gone() { + let polls = || -> Vec<_> { (0..=FEW).map(|_| Arc::new(Poll::default())).collect() }; let w = watch(); - let cancelled = Arc::new(Poll::default()); - w.add_ring(cancelled.clone()); + let cancelled = polls(); + for poll in &cancelled { + w.add_ring(poll.clone()); + } w.cancel_rings(); - assert_eq!(cancelled.state.load(Ordering::Acquire), 2); + for (at, poll) in cancelled.iter().enumerate() { + assert_eq!(poll.state.load(Ordering::Acquire), 2, "poll {at} was not answered as gone"); + assert_eq!(Arc::strong_count(poll), 1, "the cancel kept poll {at}"); + } - let dropped = Arc::new(Poll::default()); - w.add_ring(dropped.clone()); + let dropped = polls(); + for poll in &dropped { + w.add_ring(poll.clone()); + } drop(w); - assert_eq!(dropped.state.load(Ordering::Acquire), 2); - assert_eq!(dropped.posts.load(Ordering::Acquire), 1); + for (at, poll) in dropped.iter().enumerate() { + assert_eq!(poll.state.load(Ordering::Acquire), 2, "dropped poll {at} was not answered as gone"); + assert_eq!(poll.posts.load(Ordering::Acquire), 1); + } } #[test] diff --git a/toyos-sched/tests/watch_list.rs b/toyos-sched/tests/watch_list.rs new file mode 100644 index 00000000000..63c60e7208e --- /dev/null +++ b/toyos-sched/tests/watch_list.rs @@ -0,0 +1,314 @@ +//! What makes a post legal in an interrupt handler, which interrupts whoever +//! holds the allocator's lock: nothing allocates or frees under a watch's list +//! lock, and a post in place frees nothing at all. +//! +//! The allocator below counts every allocation and free this thread makes +//! while it holds a list lock, and every one it makes inside a post in place. +//! The list lock counts its sections, and can run a stage between two of them. + +use std::alloc::{GlobalAlloc, Layout, System}; +use std::cell::{Cell, RefCell}; +use std::sync::atomic::{AtomicU32, Ordering::{AcqRel, Acquire, Relaxed}}; +use std::sync::{Arc, Mutex}; + +use toyos_sched::cpu::{CpuHandle, CpuHandles}; +use toyos_sched::hw::{CpuId, Kicker}; +use toyos_sched::mailbox::{mailbox, PreemptGuard, SchedMsg}; +use toyos_sched::park::{prepare, Cancel, Commit, CurrentTask}; +use toyos_sched::sync::CellLock; +use toyos_sched::task::{TaskKey, TaskShared, TaskState, WaitClass, WakeCause, WakeReason}; +use toyos_sched::watch::{Fire, Poster, Ring, Waiters, Watch, FEW}; + +thread_local! { + static HELD: Cell = const { Cell::new(false) }; + static POSTING: Cell = const { Cell::new(false) }; + static UNDER_LOCK: Cell = const { Cell::new(0) }; + static IN_POST: Cell = const { Cell::new(0) }; + static SECTIONS: Cell = const { Cell::new(0) }; + /// Sections left to start before [`BETWEEN`] runs, ahead of the last one. + static COUNTDOWN: Cell = const { Cell::new(0) }; + static BETWEEN: RefCell>> = const { RefCell::new(None) }; +} + +struct Counting; + +impl Counting { + fn note(&self) { + // `try_with`: a thread allocates while its locals are torn down. + let count = |flag: &'static std::thread::LocalKey>, + counter: &'static std::thread::LocalKey>| { + if flag.try_with(Cell::get).unwrap_or(false) { + let _ = counter.try_with(|n| n.set(n.get() + 1)); + } + }; + count(&HELD, &UNDER_LOCK); + count(&POSTING, &IN_POST); + } +} + +// SAFETY: every method forwards to `System` unchanged after a count. +unsafe impl GlobalAlloc for Counting { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + self.note(); + // SAFETY: the caller's contract, passed on. + unsafe { System.alloc(layout) } + } + + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + self.note(); + // SAFETY: the caller's contract, passed on. + unsafe { System.dealloc(ptr, layout) } + } + + unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, size: usize) -> *mut u8 { + self.note(); + // SAFETY: the caller's contract, passed on. + unsafe { System.realloc(ptr, layout, size) } + } +} + +#[global_allocator] +static ALLOCATOR: Counting = Counting; + +/// The list lock, marking this thread as holding it. +struct Watched(Mutex); + +impl CellLock for Watched { + fn with(&self, f: impl FnOnce(&mut T) -> R) -> R { + // Before the lock is taken, so a stage lands between two sections. + let due = COUNTDOWN.get(); + if due > 0 { + COUNTDOWN.set(due - 1); + if due == 1 { + let between = BETWEEN.take().expect("a countdown is staged with its stage"); + between(); + } + } + SECTIONS.set(SECTIONS.get() + 1); + let mut guard = self.0.lock().unwrap(); + let was = HELD.replace(true); + let out = f(&mut guard); + HELD.set(was); + out + } +} + +/// Run `between` with no list lock held, just before the `nth` section from +/// here takes its lock. +fn stage(nth: usize, between: impl FnOnce() + 'static) { + BETWEEN.set(Some(Box::new(between))); + COUNTDOWN.set(nth); +} + +fn sections(f: impl FnOnce()) -> usize { + let before = SECTIONS.get(); + f(); + SECTIONS.get() - before +} + +#[derive(Debug)] +enum Msg { + Wake, + Retire, +} + +impl SchedMsg for Msg { + fn wake(_key: TaskKey, _cause: WakeCause) -> Self { + Msg::Wake + } + fn retire(_shared: Arc>) -> Self { + Msg::Retire + } +} + +struct NoPreempt; + +// SAFETY: one thread and no scheduler, so nothing can deschedule a push. +unsafe impl PreemptGuard for NoPreempt {} + +struct NoKick; + +impl Kicker for NoKick { + fn kick(&self, _target: CpuId) {} +} + +/// 0 armed, 1 fired, 2 withdrawn. +#[derive(Default)] +struct Poll(AtomicU32); + +struct Entry(Arc); + +impl Ring for Entry { + fn fire(&self, _how: Fire) { + let _ = self.0 .0.compare_exchange(0, 1, AcqRel, Acquire); + } + fn live(&self) -> bool { + self.0 .0.load(Acquire) == 0 + } +} + +type TestWatch = Watch>>; + +const C0: CpuId = CpuId(0); + +fn watch() -> TestWatch { + Watch::new(Watched(Mutex::new(Waiters::new()))) +} + +fn task(key: u64) -> Arc> { + Arc::new(TaskShared::new(TaskKey(key), TaskState::Running(C0))) +} + +fn clean(phase: &str) { + assert_eq!(UNDER_LOCK.get(), 0, "{phase} allocated or freed under the list lock"); + assert_eq!(IN_POST.get(), 0, "{phase}: a post in place allocated or freed"); +} + +#[test] +fn nothing_allocates_or_frees_under_the_list_lock_and_a_post_in_place_frees_nothing() { + let (tx, mut rx) = mailbox::(); + let cpus = CpuHandles::new(vec![CpuHandle::new(C0, tx)]); + let env = Poster { cpus: &cpus, kicker: &NoKick, preempt: &NoPreempt }; + let w = watch(); + + // Nine of each grows both lists from nothing three times over. + let tasks: Vec<_> = (0..9).map(task).collect(); + for (token, t) in tasks.iter().enumerate() { + w.register(t, token as u64); + let ticket = prepare(&CurrentTask::new(t, C0), Cancel::Answers, WaitClass::Other) + .expect("nothing has posted yet"); + assert!(matches!(ticket.commit(), Commit::Parked(_))); + } + let polls: Vec<_> = (0..9).map(|_| Arc::new(Poll::default())).collect(); + for poll in &polls { + w.add_ring(Entry(poll.clone())); + } + clean("registering"); + + POSTING.set(true); + w.post_in_place(WakeCause::new(WakeReason::Woken), &env); + POSTING.set(false); + clean("posting in place"); + assert!(polls.iter().all(|p| p.0.load(Acquire) == 1), "a post in place fired every entry"); + + // Every entry is dead now, and a withdrawn one joins them: these three + // registrations sweep all ten, more than one section takes out. + let withdrawn = Arc::new(Poll::default()); + let first = Entry(withdrawn.clone()); + assert_eq!(sections(|| w.add_ring(first)), 1, "a re-arm after a post in place took the lock twice"); + withdrawn.0.store(2, Relaxed); + let fresh: Vec<_> = (0..2).map(|_| Arc::new(Poll::default())).collect(); + for poll in &fresh { + w.add_ring(Entry(poll.clone())); + } + clean("sweeping"); + assert!(polls.iter().all(|p| Arc::strong_count(p) == 1), "the sweep let go of every fired entry"); + assert_eq!(Arc::strong_count(&withdrawn), 1, "the sweep let go of the withdrawn entry"); + + // A thread's post frees what it fired, with the lock let go, and leaves + // the list's buffer to the next registration. + let later = Arc::new(Poll::default()); + w.add_ring(Entry(later.clone())); + w.post(WakeCause::new(WakeReason::Woken), &env); + clean("posting"); + assert_eq!(Arc::strong_count(&later), 1, "the post let go of the entry it fired"); + let rearmed = Entry(Arc::new(Poll::default())); + assert_eq!(sections(|| w.add_ring(rearmed)), 1, "a re-arm after a post regrew the list"); + + assert_eq!(w.post_n(3, 1, WakeCause::new(WakeReason::Woken), &env), 0, "every waiter is claimed"); + assert_eq!(w.revoke(|token| token % 2 == 0, WakeCause::new(WakeReason::Woken), &env), 5); + for t in &tasks { + w.unregister(t); + } + w.cancel_rings(); + clean("ending"); + + while rx.pop(&NoPreempt).is_some() {} + drop(w); + clean("dropping"); +} + +/// A live entry that is the last owner of a heap cell, so its drop frees, and +/// that counts its fires on a word the test keeps. +struct Owned { + fired: Arc, + _cell: Box, +} + +impl Ring for Owned { + fn fire(&self, _how: Fire) { + self.fired.fetch_add(1, AcqRel); + } + fn live(&self) -> bool { + self.fired.load(Acquire) == 0 + } +} + +type OwnedWatch = Watch>>; + +fn owned(fired: &Arc) -> Owned { + Owned { fired: fired.clone(), _cell: Box::new(0) } +} + +/// A thread's post of `n` live entries: every entry is fired once, and freed +/// with the list lock let go. +fn post_live(n: usize) -> OwnedWatch { + let (tx, _rx) = mailbox::(); + let cpus = CpuHandles::new(vec![CpuHandle::new(C0, tx)]); + let env = Poster { cpus: &cpus, kicker: &NoKick, preempt: &NoPreempt }; + let w: OwnedWatch = Watch::new(Watched(Mutex::new(Waiters::new()))); + let fired: Vec<_> = (0..n).map(|_| Arc::new(AtomicU32::new(0))).collect(); + for word in &fired { + w.add_ring(owned(word)); + } + clean("registering"); + w.post(WakeCause::new(WakeReason::Woken), &env); + clean("posting"); + for (at, word) in fired.iter().enumerate() { + assert_eq!(word.load(Acquire), 1, "entry {at} of {n} fired other than once"); + assert_eq!(Arc::strong_count(word), 1, "entry {at} of {n} was not let go of"); + } + w +} + +/// One entry past the [`FEW`] a post takes out on its stack. +#[test] +fn a_post_of_one_entry_past_its_stack_fires_each_once_and_frees_none_under_the_lock() { + post_live(FEW + 1); +} + +/// Exactly the [`FEW`] a post takes out on its stack: the list's buffer is +/// left to the re-arm. +#[test] +fn a_re_arm_after_a_post_of_few_entries_is_one_section() { + let w = post_live(FEW); + let rearmed = owned(&Arc::new(AtomicU32::new(0))); + assert_eq!(sections(|| w.add_ring(rearmed)), 1, "a re-arm after a post regrew the list"); +} + +/// A registration that finds the list full sizes a bigger buffer with the lock +/// let go. Registrations that fill the list past that buffer before it comes +/// back leave it too small to take the list, and moving the list into it then +/// would grow it under the lock. +#[test] +fn a_list_outgrowing_the_buffer_sized_for_it_is_not_moved_into_it() { + let w = Arc::new(watch()); + let tasks: Vec<_> = (0..17).map(task).collect(); + // Four fill the list's first buffer, so the fifth grows it. + for t in &tasks[..4] { + w.register(t, 0); + } + // Twelve more fill it to sixteen, the buffer the fifth sized holds eight. + let (others, rest) = (w.clone(), tasks[5..].to_vec()); + stage(2, move || { + for t in &rest { + others.register(t, 0); + } + }); + w.register(&tasks[4], 0); + clean("growing past a buffer sized before"); + assert_eq!(w.threads(), 17, "a registration was lost"); + for t in &tasks { + w.unregister(t); + } +}