From b656c5de9953cf2bc13a357dba175dc3a88875ae Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:53:47 +0200 Subject: [PATCH 1/4] Give undelivered pool results their own ceiling, and size the ring by it WorkPort's result ring was sized by max_operations, because that was the only bound on what can sit in it: every pool-delivered operation holds an operation credit until delivery, and push is an assert!. queue_capacity does not bound it (a finished job leaves the queue while its result waits, and long jobs never enter the queue), so the ring could not simply be made smaller. A host that sets max_operations to 32 768 for I/O therefore paid for a 32 768-slot result ring its pool work would never fill. The accounting is now separate, as the correction on #88 asks. PoolConfig gains max_undelivered (default 4096): a per-loop ceiling on operations that deliver through the ring - blocking jobs of both occupancy classes, pool lookups, native typed file requests and external waits. The driver admits each one against it before anything is created, refusing with ResourceLimit exactly like a full queue_capacity, charges the operation, and releases the credit when the operation retires, which is after its result has been popped from the ring. The ring is sized by min(max_undelivered, max_operations), so the assert's invariant now holds against the new counter: the ring's occupancy never exceeds the charged operations, and those never exceed its capacity. A zero ceiling is InvalidInput. max_undelivered is per loop, so it is left out of the process-wide pool configuration check: loops that differ only there share one pool. Tests: a contract scenario submits exactly max_undelivered jobs of both classes without turning, checks that one more blocking job, long job and external wait are each refused with ResourceLimit, lets every result land in the ring (filling it to capacity) and then delivers each exactly once, twice over to prove delivery returns every credit. It fails without the admission check. A unit test checks the ring is sized by the ceiling, never above max_operations, and that a zero ceiling is rejected. Breaking: PoolConfig has a new public field, so struct literals without ..Default::default() no longer compile. Closes #88 --- DESIGN.md | 12 ++- crates/turnloop-contract/src/extended.rs | 84 ++++++++++++++++++++ crates/turnloop-contract/src/lib.rs | 4 + crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/blocking.rs | 53 ++++++++++++- crates/turnloop/src/driver.rs | 95 ++++++++++++++++++++++- crates/turnloop/src/queue.rs | 11 ++- 7 files changed, 251 insertions(+), 9 deletions(-) diff --git a/DESIGN.md b/DESIGN.md index 2454b19..3ad4c33 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -596,8 +596,16 @@ slab (`kernel`, and the `bridges` parallel to it) maps a completion packet's pointer back to an operation index by pointer arithmetic over one allocation; paging it needs a different reverse map, which is its own change. The cross-thread result rings are lock-free and index by a power-of-two mask, where the capacity is -also the backpressure bound. `pooled_buffers` is a separate question, tracked in -#43: it is the remaining per-loop cost that scales with configuration. +also the backpressure bound. The post ring is sized by `post_capacity`; the +result ring that pool workers and helper threads deliver into is sized by +`PoolConfig::max_undelivered` (capped at `max_operations`), not by +`max_operations` itself (#88). Every operation delivering through that ring — +blocking jobs of both classes, pool lookups, typed file requests, external waits +— holds one of its credits from acceptance until its result leaves the ring, and +a submission past the ceiling is refused with `ResourceLimit`, so the ring can +never be full when a worker pushes. `pooled_buffers` is the remaining per-loop +cost that scales with configuration, the other half of #88. + **Reaching the ceiling is backpressure.** An operation whose completion creates a handle — an accept, a handle receive — reserves its handle slot when it is diff --git a/crates/turnloop-contract/src/extended.rs b/crates/turnloop-contract/src/extended.rs index 11274fd..d85f738 100644 --- a/crates/turnloop-contract/src/extended.rs +++ b/crates/turnloop-contract/src/extended.rs @@ -1004,3 +1004,87 @@ pub fn kernel_accept_exactly_once() { } mt_verify(served, &per_loop, "reuse-port"); } + +/// Undelivered pool results have a ceiling of their own, and the loop's result +/// ring is sized by it instead of by `max_operations` (issue #88). +/// +/// A job keeps its credit until its result is delivered, so results really do +/// pile up: here jobs that finish at once are submitted without turning the +/// loop, and every result waits in the ring. The submission that would not fit +/// is refused with `ResourceLimit` before any work exists, whichever kind of +/// pool-delivered operation it is. Without the ceiling a worker would find the +/// ring full, which is an `assert!` on the worker's own thread. +pub fn undelivered_pool_results_are_bounded() { + const CEILING: usize = 8; + let mut l = Driver::::new(Config { + // An I/O ceiling this large no longer sizes the result ring. + max_operations: 32_768, + blocking_pool: PoolConfig { + max_undelivered: CEILING, + ..PoolConfig::default() + }, + ..Config::default() + }) + .expect("loop"); + let condition = WaitCondition::new(0).expect("condition"); + let finished = Arc::new(AtomicUsize::new(0)); + let mut out = Completions::default(); + // Twice: delivery must return every credit it took. + for round in 0..2 { + for i in 0..CEILING { + let finished = finished.clone(); + let occupancy = if i % 2 == 0 { + Occupancy::Bounded + } else { + Occupancy::Long + }; + l.blocking_with( + move |_| { + finished.fetch_add(1, Ordering::AcqRel); + Ok(Payload::U64(i as u64)) + }, + occupancy, + Token(i as u64), + ) + .expect("within the ceiling"); + } + let refused = [ + l.blocking(|| unreachable!("refused"), Token(99)), + l.blocking_with(|_| unreachable!("refused"), Occupancy::Long, Token(99)), + l.external_wait(&condition, 0, None, Token(99)), + ]; + for attempt in refused { + assert_eq!( + attempt.expect_err("at the ceiling").kind, + ErrorKind::ResourceLimit, + "pool-delivered work past the ceiling is refused before it exists" + ); + } + // Every job runs and publishes while the loop is not turned, so the + // ring holds all CEILING results at once: exactly its capacity. + until("every job ran", || { + finished.load(Ordering::Acquire) == (round + 1) * CEILING + }); + until("every result was published", || { + let s = pool_stats(); + s.busy == 0 && s.long_busy == 0 + }); + let mut seen = [false; CEILING]; + let deadline = l.now() + Duration::from_secs(10); + while seen.contains(&false) { + assert!(l.now() < deadline, "results delivered: {seen:?}"); + l.turn(Timeout::Until(deadline), &mut out).expect("turn"); + for c in out.drain() { + let OpResult::Blocking(Payload::U64(v)) = c.result else { + panic!("unexpected {:?}", c.result); + }; + assert_eq!(c.token, Token(v)); + assert!(!seen[v as usize], "result {v} delivered twice"); + seen[v as usize] = true; + } + } + } + l.turn(Timeout::Now, &mut out).expect("nothing left"); + assert!(out.is_empty()); + assert!(!l.alive()); +} diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index ebc2b11..018705a 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -279,6 +279,10 @@ mod native { long_jobs_settle_once_on_cancel_panic_and_shutdown::(); } #[test] + fn undelivered_pool_results_have_a_ceiling() { + undelivered_pool_results_are_bounded::(); + } + #[test] fn reuse_port_share_binds_twice() { reuse_port_share::(); } diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index 7529ab3..2e662c9 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -32,6 +32,7 @@ contract!( paged_growth_preserves_handles, accept_reserves_its_handle_slot, an_armed_accept_keeps_its_slot, + undelivered_pool_results_are_bounded, ready_timer_liveness, pooled_lease_backpressure, io_and_posts_progress_with_repeating_timers diff --git a/crates/turnloop/src/blocking.rs b/crates/turnloop/src/blocking.rs index 40f9bc1..5f266b8 100644 --- a/crates/turnloop/src/blocking.rs +++ b/crates/turnloop/src/blocking.rs @@ -51,6 +51,33 @@ pub struct PoolConfig { /// for peak long occupancy only while it lasts. Reuse within the window /// costs no thread creation. pub long_idle_timeout: Duration, + /// Ceiling on one loop's pool-delivered operations whose result has not + /// been delivered yet (issue #88). Unlike the fields above, this one is + /// **per loop**: it is not part of the process-wide configuration, and loops + /// that differ only here share one pool. + /// + /// Every operation whose result a pool worker or helper thread hands back + /// counts against it from acceptance until its completion is taken off the + /// loop's result ring: [`Driver::blocking`](crate::Driver::blocking) and + /// [`Driver::blocking_with`](crate::Driver::blocking_with) jobs of both + /// occupancy classes, lookups served by the pool, typed filesystem requests + /// on native targets and [`Driver::external_wait`](crate::Driver::external_wait) + /// registrations. A submission that would exceed it is refused with + /// [`ErrorKind::ResourceLimit`] before the work exists, the way a full + /// [`queue_capacity`](Self::queue_capacity) refuses one. + /// + /// The loop's result ring is sized by this ceiling (capped at + /// [`Config::max_operations`](crate::Config::max_operations), which bounds + /// every operation anyway), not by `max_operations` itself. The ring must + /// be able to hold every undelivered result at once — a worker that finds + /// it full has nowhere to put a result — so this is the number that makes + /// the ring safe, and a host that raises `max_operations` for I/O no longer + /// pays for a result ring its pool work will never fill. + /// + /// `queue_capacity` does not bound the same thing: a bounded job leaves the + /// queue when a worker takes it, long jobs never enter it, and a finished + /// job's result can wait in the ring while the queue refills. + pub max_undelivered: usize, } impl Default for PoolConfig { fn default() -> Self { @@ -59,6 +86,19 @@ impl Default for PoolConfig { queue_capacity: 1024, long_threads_max: 512, long_idle_timeout: Duration::from_secs(10), + max_undelivered: 4096, + } + } +} +impl PoolConfig { + /// The fields fixed process-wide by the first submission. `max_undelivered` + /// belongs to each loop and is left out, so it never makes two loops' + /// configurations disagree. + #[cfg(not(target_arch = "wasm32"))] + fn process_wide(self) -> Self { + Self { + max_undelivered: 0, + ..self } } } @@ -203,6 +243,11 @@ impl WorkPort { pub fn is_empty(&self) -> bool { self.queue.is_empty() } + /// Slots in the result ring, for tests of its sizing. + #[cfg(all(test, not(loom), not(target_arch = "wasm32")))] + pub(crate) fn capacity(&self) -> usize { + self.queue.capacity() + } pub fn pop(&self) -> Option { self.queue.pop() } @@ -229,8 +274,10 @@ impl WorkPort { if self.closed.load(Ordering::Acquire) { return; } - // One reserved core operation credit per job remains held until delivery. - // Thus this separate queue (>= max_operations) cannot overflow. + // Every producer here is an operation the driver admitted against + // `PoolConfig::max_undelivered`, and it keeps that credit until its + // result is popped from this ring, which is at least that large. + // Thus the ring cannot overflow; a full ring is a broken invariant. assert!( self.queue.push(result).is_ok(), "blocking completion credit invariant" @@ -453,7 +500,7 @@ mod native { .get_or_init(|| start(config)) .as_ref() .map_err(|&e| e)?; - if pool.config != config { + if pool.config.process_wide() != config.process_wide() { return Err(Error::new(ErrorKind::InvalidInput)); } Ok(pool) diff --git a/crates/turnloop/src/driver.rs b/crates/turnloop/src/driver.rs index 45303f4..582f409 100644 --- a/crates/turnloop/src/driver.rs +++ b/crates/turnloop/src/driver.rs @@ -56,6 +56,9 @@ struct Op { /// backend accepted natively (such as WASI DNS). DESIGN §10 rule 3 keys /// queued-turn discovery on these operations. native: bool, + /// Delivers through the loop's result ring (`WorkPort`), and so holds one + /// of `PoolConfig::max_undelivered` credits until it retires. + port: bool, job_cancel: Option>, previous: Option, next: Option, @@ -102,6 +105,10 @@ pub struct Driver { /// for but not yet delivered. See [`Driver::submit`]. reserved_handles: usize, native_pending: usize, + /// Operations delivering through `work_port` that have not retired: the + /// ring's occupancy can never exceed this, and admission keeps it at or + /// below `PoolConfig::max_undelivered`, which sizes the ring (issue #88). + undelivered: usize, config: Config, _local: PhantomData>, } @@ -115,6 +122,7 @@ impl Driver { || config.events_per_turn == 0 || config.post_capacity == 0 || config.pooled_buffer_size == 0 + || config.blocking_pool.max_undelivered == 0 { return Err(Error::new(ErrorKind::InvalidInput)); } @@ -129,7 +137,15 @@ impl Driver { let notifier = Notifier::new(backend.waker()); backend.set_notifier(notifier.clone()); let poster = Poster::new(config.post_capacity, notifier.clone()); - let work_port = crate::blocking::WorkPort::new(config.max_operations, notifier.clone()); + // Sized by the results that can be undelivered at once, not by every + // operation the loop may hold: see `PoolConfig::max_undelivered`. + let work_port = crate::blocking::WorkPort::new( + config + .blocking_pool + .max_undelivered + .min(config.max_operations), + notifier.clone(), + ); let owner = loop { let current = NEXT_OWNER.load(Ordering::Relaxed); let next = current @@ -171,6 +187,7 @@ impl Driver { outstanding: 0, reserved_handles: 0, native_pending: 0, + undelivered: 0, config, _local: PhantomData, }) @@ -260,6 +277,7 @@ impl Driver { fs: false, reserved_handles: 0, native, + port: false, previous, next: None, }) @@ -291,6 +309,24 @@ impl Driver { key, }) } + /// Refuse a submission whose result would have no room in the result ring. + /// + /// Checked before anything is created, like the pool's own queue limit, so + /// a refused submission owes no completion. See `PoolConfig::max_undelivered`. + fn admit_port(&self) -> Result<()> { + if self.undelivered >= self.config.blocking_pool.max_undelivered { + return Err(Error::new(ErrorKind::ResourceLimit)); + } + Ok(()) + } + /// Charge an admitted operation to the result ring until it retires. + fn charge_port(&mut self, op: OpId) { + let op = self.ops.get_mut(op.key).expect("admitted op"); + debug_assert!(!op.port); + op.port = true; + self.undelivered += 1; + debug_assert!(self.undelivered <= self.config.blocking_pool.max_undelivered); + } fn retire(&mut self, id: OpId) -> Option { let op = self.ops.remove(id.key)?; self.connect_deadlines.cancel(id.key); @@ -298,6 +334,9 @@ impl Driver { if op.native { self.native_pending -= 1; } + if op.port { + self.undelivered -= 1; + } if let Some(previous) = op.previous { self.ops.get_mut(previous.key).expect("previous").next = op.next; } @@ -1109,8 +1148,10 @@ impl Driver { deadline: Option, token: Token, ) -> Result { + self.admit_port()?; let op = self.new_op(None, token)?; self.ops.get_mut(op.key).expect("new wait").external_wait = true; + self.charge_port(op); if let Err(e) = crate::external_wait::submit(op, self.work_port.clone(), condition, expected, deadline) { @@ -1135,6 +1176,9 @@ impl Driver { } else { Kind::Socket }; + if B::FILESYSTEM == Filesystem::Pool { + self.admit_port()?; + } let target = request.handle(); if let Some(h) = target { let r = self.resource(h)?; @@ -1162,6 +1206,9 @@ impl Driver { } }; self.ops.get_mut(op.key).expect("new request").fs = true; + if B::FILESYSTEM == Filesystem::Pool { + self.charge_port(op); + } let accepted = match B::FILESYSTEM { #[cfg(not(target_arch = "wasm32"))] Filesystem::Pool => self.files.submit(op, handle, request), @@ -1356,7 +1403,9 @@ impl Driver { &std::sync::Arc, ) -> Box Result + Send>, { + self.admit_port()?; let op = self.new_op(None, token)?; + self.charge_port(op); let cancel = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); self.ops.get_mut(op.key).expect("new op").job_cancel = Some(cancel.clone()); let f = make(&cancel); @@ -2345,3 +2394,47 @@ impl Driver { Ok(h) } } + +#[cfg(all(test, not(loom), not(target_arch = "wasm32")))] +mod result_ring { + use super::*; + /// The result ring holds every result that may be undelivered at once, and + /// no more: `max_undelivered`, capped by `max_operations` (issue #88). + #[test] + fn is_sized_by_undelivered_results_not_by_operations() { + let pool = PoolConfig { + max_undelivered: 100, + ..PoolConfig::default() + }; + let l = Loop::new(Config { + max_operations: 32_768, + blocking_pool: pool, + ..Config::default() + }) + .expect("loop"); + assert_eq!( + l.work_port.capacity(), + 128, + "a large I/O ceiling is not paid for" + ); + let l = Loop::new(Config { + max_operations: 16, + ..Config::default() + }) + .expect("loop"); + assert_eq!( + l.work_port.capacity(), + 16, + "never more than every operation" + ); + let none = PoolConfig { + max_undelivered: 0, + ..PoolConfig::default() + }; + let refused = Loop::new(Config { + blocking_pool: none, + ..Config::default() + }); + assert_eq!(refused.err().map(|e| e.kind), Some(ErrorKind::InvalidInput)); + } +} diff --git a/crates/turnloop/src/queue.rs b/crates/turnloop/src/queue.rs index f2fabe5..3bedab8 100644 --- a/crates/turnloop/src/queue.rs +++ b/crates/turnloop/src/queue.rs @@ -23,8 +23,9 @@ unsafe impl Send for Queue {} /// `Slot`'s initial value **is** the all-zero bit pattern: `state` starts at 0, /// which is the empty state, and `value` is a `MaybeUninit` for which every /// pattern is valid. Building the slots with `map`/`collect` writes each one, -/// which makes the entire ring resident at construction — 3.15 MiB for a -/// 32 768-slot `WorkPort`, in a process that may never submit a blocking job. +/// which makes the entire ring resident at construction — 3.15 MiB for the +/// 32 768-slot `WorkPort` a large `max_operations` used to give a loop, in a +/// process that may never submit a blocking job. /// /// Asking the allocator for zeroed memory instead gives a ready ring with no /// writes at all, and a zeroed allocation this large is fresh pages the OS @@ -117,6 +118,10 @@ impl Queue { pub fn is_empty(&self) -> bool { self.occupied.load(Ordering::Acquire) == 0 } + #[cfg(all(test, not(loom), not(target_arch = "wasm32")))] + pub fn capacity(&self) -> usize { + self.slots.len() + } } impl Drop for Queue { fn drop(&mut self) { @@ -185,7 +190,7 @@ mod zeroed_ring { /// slots were constructed individually. /// /// The slots come from `alloc_zeroed` so that a large ring costs only the - /// pages it touches — `WorkPort` sizes its ring from `max_operations`, which + /// pages it touches — `WorkPort` sized its ring from `max_operations`, which /// a host legitimately sets to 32 768, and writing every slot made all /// 3.15 MiB of it resident in a process that may never submit a blocking /// job. That is only sound because `Slot`'s empty state IS the all-zero bit From d271c675821d871cebb0841834b506fa22e92bb4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:54:20 +0200 Subject: [PATCH 2/4] Budget a multishot accept's connections to the handle slots left A multishot accept holds one handle reservation, so it was protected for one connection at a time, but one poll can hand it up to events_per_turn connections. Once the loop reached its ceiling mid-batch, the rest were accepted from the kernel and then destroyed, completing with ResourceLimit - the behaviour #75 removed for single-shot accepts. Backend::poll now takes a Budget. Its only field, accepts, is the number of connections multishot accepts may deliver in this poll: every free handle slot that no single-shot accept or handle receive has reserved. Each nonterminal Accepted/PipeAccepted event spends one. That is exactly what the core can house, so a batch can no longer outrun the ceiling. The budget throttles only multishot accepts, per source, as the issue requires for DESIGN section 10 rule 3: reads, writes, timers, posts, single-shot accepts and every other source fill the event vector exactly as before, so the fairness contract is untouched. A multishot accept out of budget takes nothing more from the kernel and is parked in a per-backend throttled list, which is neither runnable work (has_work) nor a reason to end the wait, so a turn at the ceiling blocks for its timeout instead of spinning (rule 4a). The first poll with a positive budget resumes it. Per backend: - epoll/kqueue (shared Unix engine), WASI 0.2, WASI 0.3: the listener's head accept is skipped and parked; readiness stays cached, so the connection stays in the backlog and is accepted when the listener resumes. WASI 0.2 also stops subscribing a parked listener's level-triggered pollable, which would otherwise end every wait at once. - IOCP: a parked accept is not re-armed with AcceptEx, and a connection an already-armed AcceptEx completed stays in the operation until it can be housed; cancellation still reaches a parked operation. - web: has no listener; it takes the parameter and ignores it. In the core, the multishot reservations are pooled rather than owned per operation. With per-operation ownership, two armed listeners could each hold a slot while one of them spent both within the budget, and the second delivery would then find its slot promised to the other listener and fail. A multishot delivery now spends any pooled reservation and the pool is refilled as far as the ceiling allows; the budget (free slots minus single-shot reservations) counts exactly the deliveries that can then succeed. Test: a new contract scenario floods a multishot accept on a loop with room for three connections (events_per_turn is 256), asserts no accept ever completes with an error and no more than three are live, that turns at the ceiling deliver nothing and block (at most 10 turns in 100 ms), and that every one of 12 clients is then accepted as slots are freed - none was destroyed. It fails without the budget (ResourceLimit completions). The existing fairness (timer_backlog_io_and_post_progress, sustained_posts_idle_io) and no-spin (no_spin, quiet_deadline_accounting) contract tests pass unchanged on macOS/kqueue. epoll, IOCP, WASI 0.2/0.3 and web were compile-checked (clippy --target) only on this host. Breaking: Backend::poll has a new parameter (the internal backend contract). Closes #77 --- DESIGN.md | 13 ++- crates/turnloop-contract/src/extended.rs | 7 +- crates/turnloop-contract/src/lib.rs | 116 ++++++++++++++++++++++ crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/backend/iocp/mod.rs | 51 ++++++++-- crates/turnloop/src/backend/iocp/timer.rs | 24 ++++- crates/turnloop/src/backend/mod.rs | 42 +++++++- crates/turnloop/src/backend/unix.rs | 78 +++++++++++++-- crates/turnloop/src/backend/wasi_p2.rs | 60 ++++++++++- crates/turnloop/src/backend/wasi_p3.rs | 56 ++++++++++- crates/turnloop/src/backend/web.rs | 4 +- crates/turnloop/src/driver.rs | 114 ++++++++++++++------- 12 files changed, 496 insertions(+), 70 deletions(-) diff --git a/DESIGN.md b/DESIGN.md index 3ad4c33..52e7334 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -615,9 +615,16 @@ with `ResourceLimit` instead, before the kernel is asked. The pending connection stays in the listener's backlog, which is the queue meant to absorb it. This replaces accepting a connection and then destroying it for want of a slot, which is what a host at its ceiling did before, and is why it surfaced as a connection -refusal rather than a delay. A multishot accept holds one such reservation, so it -is protected for one connection at a time; bounding a whole batch needs a per-turn -native event budget on `Backend::poll`, which is a separate change. +refusal rather than a delay. A multishot accept holds one such reservation (the +reservations of all armed multishot accepts are pooled), so it is protected for +one connection at a time; a whole batch is bounded by the `Budget` the core +passes to `Backend::poll` (#77): the free slots no single-shot accept or handle +receive holds. Each connection a multishot accept delivers spends one; a +multishot accept out of budget takes nothing more from the kernel, so the +connection stays in the backlog, and it is not runnable work either, so a turn at +the ceiling still blocks for its timeout (rule 4a). The budget throttles only +multishot accepts: every other source, in the same poll, fills the output exactly +as before, which is what rule 3's fairness contract depends on. **Rule 3 rationale (tl-i01b, spec-owner decision, 2026-09-15).** The original rule 3 ("none when completions are already queued") forbade even a zero-timeout poll. Backend revision 2 has no separate no-wait discovery primitive: `poll(Duration::ZERO)` is the only way to learn about fresh readiness or completions, and on Unix, after cached readiness reaches `EAGAIN`, only the poller marks a resource ready again. libuv and Node make the same trade: `uv_run` computes `uv_backend_timeout()`, which is zero while pending, idle or closing work exists, and still calls `uv__io_poll` with that zero timeout (see the libuv [loop API](https://docs.libuv.org/en/v1.x/loop.html) and [`src/unix/core.c`](https://github.com/libuv/libuv/blob/v1.x/src/unix/core.c)). The [tl-i01 probes](docs/lanes/tl-i01.md#evidence--specification-decision) showed that skipping discovery breaks the unchanged fairness contract: skipping the native step whenever work was queued failed the timer/I/O/post fairness test after two seconds, and a guard that only drained cached work delivered 64 posts with zero reads and failed the same fairness test. A separate no-wait collection API on all six backends is not justified before measurement. Rule 4a and its raw zero-event accounting are unchanged: an empty discovery poll still counts toward the same no-spin bound. diff --git a/crates/turnloop-contract/src/extended.rs b/crates/turnloop-contract/src/extended.rs index d85f738..f46af3c 100644 --- a/crates/turnloop-contract/src/extended.rs +++ b/crates/turnloop-contract/src/extended.rs @@ -835,10 +835,9 @@ pub fn handoff_accept_exactly_once() { })); } let clients = mt_clients(addr, MT_CONNECTIONS); - // Single-shot accept, re-armed per connection. turnloop#77 is open: a - // multishot accept can outrun the handle ceiling within one turn, and this - // loop deliberately runs at a low ceiling, so depending on multishot here - // would be testing that open issue rather than the handoff. + // Single-shot accept, re-armed per connection, so this tests the handoff + // alone. A multishot accept at a low ceiling is covered separately, by + // `multishot_accept_respects_the_handle_ceiling` (turnloop#77). acceptor.accept(listener, Token(0)).expect("accept"); let mut handed = 0; let mut out = Completions::default(); diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index 018705a..496aa07 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -231,6 +231,10 @@ mod native { an_armed_accept_keeps_its_slot::(); } #[test] + fn multishot_accept_waits_at_the_handle_ceiling() { + multishot_accept_respects_the_handle_ceiling::(); + } + #[test] fn cancellation_close() { cancel_close_ordering::(); } @@ -1189,6 +1193,118 @@ pub fn an_armed_accept_keeps_its_slot() { l.close(spare, Token(6)).expect("close"); } +/// A multishot accept never outruns the handle ceiling within one turn (#77). +/// +/// One poll can bring a multishot accept many connections, up to the event +/// budget. With fewer free handle slots than that, every connection past the +/// ceiling used to be taken from the kernel and then destroyed, surfacing as an +/// accept that completed with `ResourceLimit`. The core now tells the backend +/// how many connections it can house, so the rest wait in the listener's +/// backlog, a turn at the ceiling blocks instead of spinning, and each waiting +/// connection is accepted as soon as a slot is free. +pub fn multishot_accept_respects_the_handle_ceiling() { + // Room for the listener and three connections, far below the event budget. + const ROOM: usize = 3; + const CLIENTS: usize = 12; + assert!(ROOM < Config::default().events_per_turn); + let mut l = Driver::::new(Config { + max_handles: 1 + ROOM, + max_operations: 64, + ..Config::default() + }) + .expect("loop"); + let listener = l + .tcp_listen(localhost(), &ListenOpts::default()) + .expect("listen"); + let addr = l.local_addr(listener).expect("addr"); + let accept = l + .accept_start(listener, Token(1)) + .expect("multishot accept"); + // Flood the listener before the loop turns, so connections are waiting in + // the backlog together and one poll can find all of them at once. + let clients: Vec = (0..CLIENTS) + .map(|_| std::net::TcpStream::connect(addr).expect("client")) + .collect(); + let mut out = Completions::with_capacity(4 * CLIENTS); + let mut live = Vec::new(); + let mut accepted = 0; + let collect = |out: &mut Completions, live: &mut Vec, accepted: &mut usize| { + for c in out.drain() { + match c.result { + OpResult::Accepted { conn, .. } => { + assert_eq!(c.op, Some(accept)); + assert!(!c.terminal, "the multishot accept stays armed"); + live.push(conn); + *accepted += 1; + } + OpResult::Closed => {} + OpResult::Err(e) => panic!("an accepted connection was destroyed: {e:?}"), + other => panic!("unexpected {other:?}"), + } + assert!(live.len() <= ROOM, "more connections than handle slots"); + } + }; + let deadline = Instant::now() + Duration::from_secs(20); + while live.len() < ROOM { + assert!(Instant::now() < deadline, "only {} accepted", live.len()); + l.turn(Timeout::After(Duration::from_millis(50)), &mut out) + .expect("turn"); + collect(&mut out, &mut live, &mut accepted); + } + // At the ceiling the remaining connections wait in the backlog: turns + // deliver nothing, and each blocks for its timeout instead of spinning. + let quiet = Instant::now(); + let mut turns = 0; + while quiet.elapsed() < Duration::from_millis(100) { + let info = l + .turn(Timeout::After(Duration::from_millis(20)), &mut out) + .expect("turn at the ceiling"); + assert!(info.os_waits + info.discovery_polls <= 1); + collect(&mut out, &mut live, &mut accepted); + assert_eq!(accepted, ROOM, "nothing past the ceiling was delivered"); + turns += 1; + } + assert!( + turns <= 10, + "{turns} turns in 100 ms at the ceiling: a spin" + ); + // Every slot that frees up admits the next waiting connection, until each + // client has been accepted exactly once: none was lost along the way. + while accepted < CLIENTS { + assert!( + Instant::now() < deadline, + "only {accepted} of {CLIENTS} accepted" + ); + if live.len() == ROOM { + let conn = live.pop().expect("a connection"); + l.close(conn, Token(2)).expect("close"); + } + l.turn(Timeout::After(Duration::from_millis(50)), &mut out) + .expect("turn"); + collect(&mut out, &mut live, &mut accepted); + } + assert_eq!(accepted, CLIENTS); + assert!(l.stop(accept)); + for conn in live.drain(..) { + l.close(conn, Token(2)).expect("close"); + } + l.close(listener, Token(3)).expect("close listener"); + let deadline = Instant::now() + Duration::from_secs(5); + while l.alive() { + assert!(Instant::now() < deadline, "the loop never drained"); + l.turn(Timeout::After(Duration::from_millis(50)), &mut out) + .expect("drain"); + for c in out.drain() { + assert!( + matches!(c.result, OpResult::Stopped | OpResult::Closed), + "unexpected {:?}", + c.result + ); + } + } + drop(clients); +} + pub fn pooled_lease_backpressure() { let mut l = Driver::::new(Config { pooled_buffers: 1, diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index 2e662c9..9876cfa 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -32,6 +32,7 @@ contract!( paged_growth_preserves_handles, accept_reserves_its_handle_slot, an_armed_accept_keeps_its_slot, + multishot_accept_respects_the_handle_ceiling, undelivered_pool_results_are_bounded, ready_timer_liveness, pooled_lease_backpressure, diff --git a/crates/turnloop/src/backend/iocp/mod.rs b/crates/turnloop/src/backend/iocp/mod.rs index 72203fc..d054d20 100644 --- a/crates/turnloop/src/backend/iocp/mod.rs +++ b/crates/turnloop/src/backend/iocp/mod.rs @@ -16,7 +16,7 @@ mod timer; mod watch; use crate::{ - backend::{Backend, Event, Operation, Outcome, PollInfo, Request}, + backend::{Backend, Budget, Event, Operation, Outcome, PollInfo, Request}, *, }; use integration::EventIntegration; @@ -239,6 +239,9 @@ struct Pending { queued: bool, pool_wait: bool, listener_wait: bool, + /// Listed in `throttled`: a multishot accept held back by the budget. It + /// is not re-armed, and a connection it already completed stays here. + throttled: bool, cancelled: bool, completion: Option>, offset: usize, @@ -280,6 +283,11 @@ pub struct Iocp { /// `OVERLAPPED` pointer back to an operation index needs `kernel` contiguous. bridges: Vec>, ready: VecDeque, + /// Multishot accepts out of budget: not runnable, so they neither count + /// as work nor shorten a wait, until a poll with a positive budget (#77). + throttled: VecDeque, + /// Connections multishot accepts may still deliver in the current poll. + accepts: usize, pool_waiting: VecDeque, pool: BufferPool, port: Arc, @@ -403,6 +411,7 @@ impl Iocp { && !p.queued && !p.waiting && !p.listener_wait + && (!p.throttled || p.cancelled) { p.queued = true; self.ready.push_back(i); @@ -455,6 +464,17 @@ impl Iocp { continue; }; p.queued = false; + if !p.cancelled && self.accepts == 0 && Budget::throttles(&p.request.operation) { + // Out of budget: take nothing more from the kernel, and keep a + // connection that has already completed in `p` until a poll can + // house it. + if !p.throttled { + p.throttled = true; + self.throttled.push_back(i); + } + self.ops[i] = Some(p); + continue; + } let result = if p.cancelled { Ok(Some((Outcome::Cancelled, true))) } else { @@ -476,11 +496,15 @@ impl Iocp { if terminal { self.unlink(i, &p); } - events.push(Event { + let event = Event { op: p.request.op, terminal, result: Ok(result), - }); + }; + if Budget::spends(&event) { + self.accepts -= 1; + } + events.push(event); if terminal { continue; } @@ -1122,6 +1146,8 @@ unsafe impl Backend for Iocp { kernel, bridges, ready: VecDeque::with_capacity(page_reserve(config.max_operations)), + throttled: VecDeque::with_capacity(page_reserve(config.max_operations)), + accepts: usize::MAX, pool_waiting: VecDeque::with_capacity(page_reserve(config.max_operations)), pool, wake: Arc::new(IocpWake { @@ -1483,6 +1509,7 @@ unsafe impl Backend for Iocp { queued: false, pool_wait: false, listener_wait: false, + throttled: false, cancelled: false, completion: None, offset: 0, @@ -1540,8 +1567,18 @@ unsafe impl Backend for Iocp { fn poll( &mut self, timeout: Option, + budget: Budget, events: &mut Vec>, ) -> Result { + self.accepts = budget.accepts; + if self.accepts != 0 { + while let Some(i) = self.throttled.pop_front() { + if let Some(p) = self.ops[i].as_mut().filter(|p| p.throttled) { + p.throttled = false; + self.schedule(i); + } + } + } if let Some(event) = &self.event { event.check()?; } @@ -1698,7 +1735,7 @@ impl Drop for Iocp { self.watches.shutdown(); let mut events = Vec::with_capacity(64); while self.ops.iter().any(Option::is_some) || self.watches.pending() { - if self.poll(None, &mut events).is_err() { + if self.poll(None, Budget::UNLIMITED, &mut events).is_err() { std::process::abort(); } events.clear(); @@ -1775,7 +1812,7 @@ mod tests { backend.port.post(WAKE, 0).expect("wake"); } let info = backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("discovery"); assert_eq!( (info.waits, info.discovery_polls, info.zero_event_waits), @@ -1791,7 +1828,7 @@ mod tests { let result = unsafe { WaitForSingleObject(event as _, 2000) }; assert_eq!(result, WAIT_OBJECT_0, "helper must forward a real packet"); let info = backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("helper drain"); assert_eq!( (info.waits, info.discovery_polls, info.zero_event_waits), @@ -1857,7 +1894,7 @@ mod tests { for (i, op) in self.ops.into_iter().enumerate() { let info = self .backend - .poll(Some(Duration::ZERO), &mut out) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut out) .expect("drain valid packet"); assert_eq!((info.waits, info.discovery_polls), (0, 0)); assert_eq!(out.len(), 1); diff --git a/crates/turnloop/src/backend/iocp/timer.rs b/crates/turnloop/src/backend/iocp/timer.rs index 16b568e..c23d722 100644 --- a/crates/turnloop/src/backend/iocp/timer.rs +++ b/crates/turnloop/src/backend/iocp/timer.rs @@ -268,7 +268,11 @@ mod tests { assert!( fixture .backend - .poll(Some(Duration::from_secs(30)), &mut out) + .poll( + Some(Duration::from_secs(30)), + crate::backend::Budget::UNLIMITED, + &mut out + ) .is_err() ); assert_eq!( @@ -288,7 +292,11 @@ mod tests { backend.port.post(WAKE, 0).expect("wake first wait"); let mut out = Vec::with_capacity(1); let info = backend - .poll(Some(Duration::from_secs(30)), &mut out) + .poll( + Some(Duration::from_secs(30)), + crate::backend::Budget::UNLIMITED, + &mut out, + ) .expect("first wait"); assert_eq!( (info.waits, info.discovery_polls, info.zero_event_waits), @@ -316,7 +324,11 @@ mod tests { }, ); let info = backend - .poll(Some(Duration::from_secs(30)), &mut out) + .poll( + Some(Duration::from_secs(30)), + crate::backend::Budget::UNLIMITED, + &mut out, + ) .expect("drain pending timer"); assert_eq!( (info.waits, info.discovery_polls, info.zero_event_waits), @@ -332,7 +344,11 @@ mod tests { // Reuse is allowed after acknowledgement, and the next generation runs. backend.port.post(WAKE, 0).expect("wake rearmed wait"); backend - .poll(Some(Duration::from_secs(30)), &mut out) + .poll( + Some(Duration::from_secs(30)), + crate::backend::Budget::UNLIMITED, + &mut out, + ) .expect("reuse timer"); assert_eq!(backend.timer.generation, generation + 1); assert_eq!(CANCELS.with(Cell::get), 2); diff --git a/crates/turnloop/src/backend/mod.rs b/crates/turnloop/src/backend/mod.rs index dcb17c3..285d30a 100644 --- a/crates/turnloop/src/backend/mod.rs +++ b/crates/turnloop/src/backend/mod.rs @@ -44,6 +44,18 @@ //! queued work, whatever `has_work` reports (web, whose poll never //! enters the OS, keeps draining cached host work). `None` means an unbounded //! wait. Durations must retain sub-millisecond precision. +//! * `Budget::accepts` bounds the connections multishot accepts deliver in one +//! poll: nonterminal `Accepted`/`PipeAccepted` events (turnloop#77). It is the +//! number of handles the core can still house for them, so each one past it +//! would be accepted from the kernel and then destroyed. It throttles only +//! those operations: every other source, single-shot accepts included, fills +//! the vector exactly as before (DESIGN §10 rule 3 fairness). A multishot +//! accept held back by it takes nothing further from the kernel (a readiness +//! backend leaves the connection in the listener's backlog; a completion +//! backend does not re-arm, and keeps an already completed one in its +//! operation storage), and is neither runnable work for `has_work` nor a +//! reason to end the wait early (rule 4a): the operation is resumed by the +//! first poll whose budget is positive again. //! * Readiness backends execute I/O in poll, cache readiness until EAGAIN, and //! requeue partially processed work fairly. Completion backends drain native //! completions. WASI 0.2 polls pollables; 0.3 drives a waitable set; web drains @@ -267,6 +279,31 @@ pub enum Filesystem { /// No filesystem: `Loop::fs` returns Unsupported. Unsupported, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +/// What the core can house from one `poll` (see the module documentation). +pub struct Budget { + /// Connections multishot accepts may deliver in this poll: the handle + /// slots the core holds for them, which each such event spends. + pub accepts: usize, +} +impl Budget { + /// No limit on any source. For backend-level tests that bypass the core. + pub const UNLIMITED: Self = Self { + accepts: usize::MAX, + }; + /// Whether this operation's deliveries are counted against `accepts`. + pub fn throttles(operation: &Operation) -> bool { + matches!(operation, Operation::Accept { multishot: true }) + } + /// Whether this event spent one of `accepts`. + pub fn spends(event: &Event) -> bool { + !event.terminal + && matches!( + event.result, + Ok(Outcome::Accepted { .. } | Outcome::PipeAccepted(_)) + ) + } +} #[derive(Clone, Copy, Debug, Default)] /// Instrumentation for the single bounded backend wait. pub struct PollInfo { @@ -436,10 +473,13 @@ pub unsafe trait Backend: Sized + 'static { fn cancel(&mut self, op: OpId) -> Result<()>; /// Whether completions or immediately runnable cached operations are available. fn has_work(&self) -> bool; - /// Append bounded completions, performing at most one OS wait with the supplied exact timeout. + /// Append bounded completions, performing at most one OS wait with the + /// supplied exact timeout, and delivering no more multishot connections + /// than `budget` allows. fn poll( &mut self, timeout: Option, + budget: Budget, events: &mut Vec>, ) -> Result; /// Release a quiescent resource after Closed was appended to host output. diff --git a/crates/turnloop/src/backend/unix.rs b/crates/turnloop/src/backend/unix.rs index 125ecf7..be0c62a 100644 --- a/crates/turnloop/src/backend/unix.rs +++ b/crates/turnloop/src/backend/unix.rs @@ -10,7 +10,7 @@ use super::{ }; use crate::slots::{Slots, page_reserve}; use crate::{ - backend::{Backend, Event, Operation, Outcome, PollInfo, Request}, + backend::{Backend, Budget, Event, Operation, Outcome, PollInfo, Request}, *, }; use std::{ @@ -126,6 +126,8 @@ struct Resource { heads: [Option; 2], tails: [Option; 2], queued: bool, + /// Listed in `throttled`: a ready multishot accept held back by the budget. + throttled: bool, } struct Pending { request: Request, @@ -139,6 +141,12 @@ pub struct Unix { resources: Slots, ops: Slots, ready: VecDeque, + /// Listeners whose multishot accept is ready but out of budget. They are + /// not runnable, so they neither count as work nor shorten a wait, and they + /// return to `ready` at the first poll with a positive budget (turnloop#77). + throttled: VecDeque, + /// Connections multishot accepts may still deliver in the current poll. + accepts: usize, cancelled: VecDeque, polled: Vec, pool: BufferPool, @@ -191,6 +199,7 @@ impl Unix { heads: [None; 2], tails: [None; 2], queued: false, + throttled: false, }); Ok(()) } @@ -198,9 +207,34 @@ impl Unix { let Some(r) = self.resources.get_mut(h.index()).and_then(Option::as_mut) else { return; }; - if r.handle == h && !r.queued && (0..2).any(|d| r.ready[d] && r.heads[d].is_some()) { + if r.handle != h || r.queued { + return; + } + let held = self.accepts == 0 && r.ready[0] && holds_throttled(&self.ops, r.heads[0]); + if (0..2).any(|d| r.ready[d] && r.heads[d].is_some() && !(d == 0 && held)) { r.queued = true; self.ready.push_back(h); + } else if held && !r.throttled { + r.throttled = true; + self.throttled.push_back(h); + } + } + /// Start a poll with `budget`: a positive one resumes every held listener. + fn begin(&mut self, budget: Budget) { + self.accepts = budget.accepts; + if self.accepts == 0 { + return; + } + while let Some(h) = self.throttled.pop_front() { + if let Some(r) = self + .resources + .get_mut(h.index()) + .and_then(Option::as_mut) + .filter(|r| r.handle == h) + { + r.throttled = false; + self.schedule(h); + } } } fn unlink(&mut self, h: Handle, i: usize, d: usize) { @@ -255,6 +289,11 @@ impl Unix { let Some(i) = r.heads[d] else { continue; }; + if self.accepts == 0 && holds_throttled(&self.ops, Some(i)) { + // Out of budget: the connection stays in the backlog, and + // `schedule` below parks the listener until a later poll. + continue; + } let p = self.ops[i].as_mut().expect("queued op"); let result = execute(r, p, &self.pool); let event = match result { @@ -276,6 +315,9 @@ impl Unix { }), }; if let Some(e) = event { + if Budget::spends(&e) { + self.accepts -= 1; + } if e.terminal { self.unlink(h, i, d); self.ops[i] = None; @@ -287,6 +329,12 @@ impl Unix { } } } +/// Whether the operation at the head of a direction is one the budget throttles. +fn holds_throttled(ops: &Slots, head: Option) -> bool { + head.and_then(|i| ops.get(i)) + .and_then(Option::as_ref) + .is_some_and(|p| Budget::throttles(&p.request.operation)) +} // SAFETY: all I/O executes synchronously in poll; a terminal event removes its // request, and owned descriptors/requests are dropped without outstanding native // buffer access. Readiness events contain generation keys, never buffer pointers. @@ -302,6 +350,8 @@ unsafe impl Backend for Unix { resources: Slots::new(config.max_handles), ops: Slots::new(config.max_operations), ready: VecDeque::with_capacity(page_reserve(config.max_handles)), + throttled: VecDeque::with_capacity(page_reserve(config.max_handles)), + accepts: usize::MAX, cancelled: VecDeque::with_capacity(page_reserve(config.max_operations)), polled: Vec::with_capacity(config.events_per_turn), files: super::files::Files::new(config, pool.clone()), @@ -757,8 +807,10 @@ unsafe impl Backend for Unix { fn poll( &mut self, timeout: Option, + budget: Budget, events: &mut Vec>, ) -> Result { + self.begin(budget); while events.len() < events.capacity() { let Some(op) = self.cancelled.pop_front() else { break; @@ -817,6 +869,7 @@ unsafe impl Backend for Unix { let _ = self.poller.deregister(fd); } self.ready.retain(|&at| at != h); + self.throttled.retain(|&at| at != h); self.resources[h.index()] = None; } // Closing the final descriptor removes its registration from epoll/kqueue. @@ -837,6 +890,7 @@ unsafe impl Backend for Unix { } // Remove a stale scheduling entry before the slot can be reused. self.ready.retain(|&at| at != h); + self.throttled.retain(|&at| at != h); Ok(self.resources[h.index()] .take() .expect("validated") @@ -1316,7 +1370,7 @@ mod udp_tests { }) .expect("old receive"); backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("arm old receive"); assert!(events.is_empty()); assert_eq!( @@ -1333,7 +1387,7 @@ mod udp_tests { assert!(backend.polled.iter().any(|e| e.key == old.key() && e.read)); backend.cancel(old_op).expect("cancel old receive"); backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("acknowledgement"); assert_eq!(events.len(), 1); assert_eq!(events[0].op, old_op); @@ -1388,14 +1442,18 @@ mod udp_tests { ); backend.release(old); // A stale handle must not release the new fd. backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("arm replacement receive"); assert!(events.is_empty(), "old event batch crossed generations"); // A stale event is not a completion, and must not cause a no-spin // violation after the new socket reports EAGAIN. let at = Instant::now() + Duration::from_millis(2); let info = backend - .poll(Some(Duration::from_millis(2)), &mut events) + .poll( + Some(Duration::from_millis(2)), + Budget::UNLIMITED, + &mut events, + ) .expect("new empty socket waits"); assert!(events.is_empty(), "old datagram crossed socket lifetime"); assert_eq!(info.waits, 1); @@ -1403,7 +1461,7 @@ mod udp_tests { assert!(Instant::now() >= at); assert_eq!(peer.send_to(b"new", endpoint).expect("new packet"), 3); backend - .poll(Some(Duration::from_secs(1)), &mut events) + .poll(Some(Duration::from_secs(1)), Budget::UNLIMITED, &mut events) .expect("new delivery"); assert_eq!(events.len(), 1); let event = events.pop().expect("one completion"); @@ -1422,7 +1480,7 @@ mod udp_tests { other => panic!("unexpected reused-socket result: {other:?}"), } backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("no duplicates"); assert!(events.is_empty()); Ok(()) @@ -1485,7 +1543,7 @@ mod process_races { .expect("exit operation"); let mut events = Vec::with_capacity(4); backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("exit poll"); assert_eq!(events.len(), 1); assert!(matches!( @@ -1497,7 +1555,7 @@ mod process_races { )); events.clear(); backend - .poll(Some(Duration::ZERO), &mut events) + .poll(Some(Duration::ZERO), Budget::UNLIMITED, &mut events) .expect("duplicate check"); assert!(events.is_empty()); let mut code = 0; diff --git a/crates/turnloop/src/backend/wasi_p2.rs b/crates/turnloop/src/backend/wasi_p2.rs index cf0e356..ba2e8b4 100644 --- a/crates/turnloop/src/backend/wasi_p2.rs +++ b/crates/turnloop/src/backend/wasi_p2.rs @@ -5,7 +5,7 @@ mod abi; mod fs; mod sockopt; use crate::{ - backend::{Backend, Event, Filesystem, Operation, Outcome, PollInfo, Request, Wake}, + backend::{Backend, Budget, Event, Filesystem, Operation, Outcome, PollInfo, Request, Wake}, *, }; use std::{ @@ -83,6 +83,8 @@ struct Resource { heads: [Option; 2], tails: [Option; 2], queued: bool, + /// Listed in `throttled`: a ready multishot accept held back by the budget. + throttled: bool, } struct Pending { request: Request, @@ -120,6 +122,11 @@ pub struct WasiP2 { resources: Slots, ops: Slots, ready: VecDeque, + /// Listeners whose multishot accept is ready but out of budget: neither + /// runnable nor subscribed, until a poll with a positive budget (#77). + throttled: VecDeque, + /// Connections multishot accepts may still deliver in the current poll. + accepts: usize, cancelled: VecDeque, handles: Vec, owners: Vec, @@ -227,6 +234,7 @@ impl WasiP2 { heads: [None; 2], tails: [None; 2], queued: false, + throttled: false, }); Ok(()) } @@ -234,9 +242,34 @@ impl WasiP2 { let Some(r) = self.resources.get_mut(h.index()).and_then(Option::as_mut) else { return; }; - if r.handle == h && !r.queued && (0..2).any(|d| r.ready[d] && r.heads[d].is_some()) { + if r.handle != h || r.queued { + return; + } + let held = self.accepts == 0 && r.ready[0] && holds_throttled(&self.ops, r.heads[0]); + if (0..2).any(|d| r.ready[d] && r.heads[d].is_some() && !(d == 0 && held)) { r.queued = true; self.ready.push_back(h); + } else if held && !r.throttled { + r.throttled = true; + self.throttled.push_back(h); + } + } + /// Start a poll with `budget`: a positive one resumes every held listener. + fn begin(&mut self, budget: Budget) { + self.accepts = budget.accepts; + if self.accepts == 0 { + return; + } + while let Some(h) = self.throttled.pop_front() { + if let Some(r) = self + .resources + .get_mut(h.index()) + .and_then(Option::as_mut) + .filter(|r| r.handle == h) + { + r.throttled = false; + self.schedule(h); + } } } fn unlink(&mut self, h: Handle, i: usize, d: usize) { @@ -291,6 +324,10 @@ impl WasiP2 { let Some(i) = r.heads[d] else { continue; }; + if self.accepts == 0 && holds_throttled(&self.ops, Some(i)) { + // Out of budget: `schedule` below parks the listener. + continue; + } let p = self.ops[i].as_mut().expect("queued op"); let result = execute(r, p, &self.pool, &mut self.scratch); let event = match result { @@ -311,6 +348,9 @@ impl WasiP2 { }), }; if let Some(e) = event { + if Budget::spends(&e) { + self.accepts -= 1; + } if e.terminal { self.unlink(h, i, d); self.ops[i] = None; @@ -404,6 +444,8 @@ unsafe impl Backend for WasiP2 { resources: Slots::new(config.max_handles), ops: Slots::new(config.max_operations), ready: VecDeque::with_capacity(page_reserve(config.max_handles)), + throttled: VecDeque::with_capacity(page_reserve(config.max_handles)), + accepts: usize::MAX, cancelled: VecDeque::with_capacity(page_reserve(config.max_operations)), handles: Vec::with_capacity(batch), owners: Vec::with_capacity(batch), @@ -654,8 +696,10 @@ unsafe impl Backend for WasiP2 { fn poll( &mut self, timeout: Option, + budget: Budget, events: &mut Vec>, ) -> Result { + self.begin(budget); while events.len() < events.capacity() { let Some(op) = self.cancelled.pop_front() else { break; @@ -678,7 +722,9 @@ unsafe impl Backend for WasiP2 { self.indices.clear(); for r in self.resources.iter().flatten() { for d in 0..2 { - if r.heads[d].is_none() { + // A held listener is not subscribed: its pollable is level + // triggered, and would end every wait at once (rule 4a). + if r.heads[d].is_none() || (d == 0 && r.throttled) { continue; } let t = &r.transport; @@ -752,6 +798,7 @@ unsafe impl Backend for WasiP2 { self.files.release(h); if self.get(h).is_ok() { self.ready.retain(|&at| at != h); + self.throttled.retain(|&at| at != h); self.resources[h.index()] = None; } } @@ -966,3 +1013,10 @@ fn receive( !multishot, ))) } + +/// Whether the operation at the head of a direction is one the budget throttles. +fn holds_throttled(ops: &Slots, head: Option) -> bool { + head.and_then(|i| ops.get(i)) + .and_then(Option::as_ref) + .is_some_and(|p| Budget::throttles(&p.request.operation)) +} diff --git a/crates/turnloop/src/backend/wasi_p3.rs b/crates/turnloop/src/backend/wasi_p3.rs index 2b6baea..47faa0e 100644 --- a/crates/turnloop/src/backend/wasi_p3.rs +++ b/crates/turnloop/src/backend/wasi_p3.rs @@ -9,7 +9,7 @@ mod return_storage; mod sockopt; mod wait_set; use crate::{ - backend::{Backend, Event, Filesystem, Operation, Outcome, PollInfo, Request, Wake}, + backend::{Backend, Budget, Event, Filesystem, Operation, Outcome, PollInfo, Request, Wake}, *, }; use std::{ @@ -99,6 +99,8 @@ struct Resource { heads: [Option; 2], tails: [Option; 2], queued: bool, + /// Listed in `throttled`: a ready multishot accept held back by the budget. + throttled: bool, } #[derive(Clone, Copy)] enum WaitKind { @@ -136,6 +138,11 @@ pub struct WasiP3 { resources: Slots, ops: Slots, ready: VecDeque, + /// Listeners whose multishot accept is ready but out of budget: neither + /// runnable nor subscribed, until a poll with a positive budget (#77). + throttled: VecDeque, + /// Connections multishot accepts may still deliver in the current poll. + accepts: usize, cancelled: VecDeque, pool: BufferPool, wake: Arc, @@ -241,6 +248,7 @@ impl WasiP3 { heads: [None; 2], tails: [None; 2], queued: false, + throttled: false, }); Ok(()) } @@ -248,9 +256,34 @@ impl WasiP3 { let Some(r) = self.resources.get_mut(h.index()).and_then(Option::as_mut) else { return; }; - if r.handle == h && !r.queued && (0..2).any(|d| r.ready[d] && r.heads[d].is_some()) { + if r.handle != h || r.queued { + return; + } + let held = self.accepts == 0 && r.ready[0] && holds_throttled(&self.ops, r.heads[0]); + if (0..2).any(|d| r.ready[d] && r.heads[d].is_some() && !(d == 0 && held)) { r.queued = true; self.ready.push_back(h); + } else if held && !r.throttled { + r.throttled = true; + self.throttled.push_back(h); + } + } + /// Start a poll with `budget`: a positive one resumes every held listener. + fn begin(&mut self, budget: Budget) { + self.accepts = budget.accepts; + if self.accepts == 0 { + return; + } + while let Some(h) = self.throttled.pop_front() { + if let Some(r) = self + .resources + .get_mut(h.index()) + .and_then(Option::as_mut) + .filter(|r| r.handle == h) + { + r.throttled = false; + self.schedule(h); + } } } fn unlink(&mut self, h: Handle, i: usize, d: usize) { @@ -305,6 +338,10 @@ impl WasiP3 { let Some(i) = r.heads[d] else { continue; }; + if self.accepts == 0 && holds_throttled(&self.ops, Some(i)) { + // Out of budget: `schedule` below parks the listener. + continue; + } let p = self.ops[i].as_mut().expect("queued op"); let result = execute(r, p, &self.pool, &self.wait_set); let event = match result { @@ -330,6 +367,9 @@ impl WasiP3 { }), }; if let Some(e) = event { + if Budget::spends(&e) { + self.accepts -= 1; + } if e.terminal { self.unlink(h, i, d); self.ops[i] = None; @@ -357,6 +397,8 @@ unsafe impl Backend for WasiP3 { resources: Slots::new(config.max_handles), ops: Slots::new(config.max_operations), ready: VecDeque::with_capacity(page_reserve(config.max_handles)), + throttled: VecDeque::with_capacity(page_reserve(config.max_handles)), + accepts: usize::MAX, cancelled: VecDeque::with_capacity(page_reserve(config.max_operations)), files: super::wasi_fs::Files::new(config, pool.clone()), pool, @@ -582,8 +624,10 @@ unsafe impl Backend for WasiP3 { fn poll( &mut self, timeout: Option, + budget: Budget, events: &mut Vec>, ) -> Result { + self.begin(budget); while events.len() < events.capacity() { let Some(op) = self.cancelled.pop_front() else { break; @@ -677,6 +721,7 @@ unsafe impl Backend for WasiP3 { self.files.release(h); if self.get(h).is_ok() { self.ready.retain(|&at| at != h); + self.throttled.retain(|&at| at != h); self.resources[h.index()] = None; } } @@ -1155,3 +1200,10 @@ fn quiesce(p: &mut Pending, set: &WaitSet, r: &mut Resource) { } } } + +/// Whether the operation at the head of a direction is one the budget throttles. +fn holds_throttled(ops: &Slots, head: Option) -> bool { + head.and_then(|i| ops.get(i)) + .and_then(Option::as_ref) + .is_some_and(|p| Budget::throttles(&p.request.operation)) +} diff --git a/crates/turnloop/src/backend/web.rs b/crates/turnloop/src/backend/web.rs index c59c803..21728f4 100644 --- a/crates/turnloop/src/backend/web.rs +++ b/crates/turnloop/src/backend/web.rs @@ -8,7 +8,7 @@ //! sockets are themselves unsupported (DESIGN §7.5). use crate::slots::Slots; use crate::{ - backend::{Backend, Event, Operation, Outcome, PollInfo, Request, Wake}, + backend::{Backend, Budget, Event, Operation, Outcome, PollInfo, Request, Wake}, *, }; use js_sys::{Function, Uint8Array}; @@ -322,6 +322,8 @@ unsafe impl Backend for Web { fn poll( &mut self, timeout: Option, + // No listener exists on the web: nothing here creates a handle. + _budget: Budget, events: &mut Vec>, ) -> Result { if timeout != Some(Duration::ZERO) { diff --git a/crates/turnloop/src/driver.rs b/crates/turnloop/src/driver.rs index 582f409..2892c83 100644 --- a/crates/turnloop/src/driver.rs +++ b/crates/turnloop/src/driver.rs @@ -1,5 +1,5 @@ use crate::{ - backend::{Backend, Event, Filesystem, Operation, Outcome, Request}, + backend::{Backend, Budget, Event, Filesystem, Operation, Outcome, Request}, fs::FsOutput, table::Table, timer::DriverTimerQueue as TimerQueue, @@ -50,8 +50,12 @@ struct Op { /// A typed filesystem request (pool service or backend, per `B::FILESYSTEM`). fs: bool, /// Handle slots this operation holds against the ceiling, released when it - /// retires. Non-zero only for operations whose completion creates a handle. + /// retires. Non-zero only for single-shot operations whose completion + /// creates a handle; a multishot accept draws on `multishot_reserved`. reserved_handles: usize, + /// A multishot accept, counted in `multishot_armed`, whose connections the + /// backend delivers against the poll's `Budget` (turnloop#77). + multishot: bool, /// Counted in `native_pending`: a socket-handle operation or a request the /// backend accepted natively (such as WASI DNS). DESIGN §10 rule 3 keys /// queued-turn discovery on these operations. @@ -104,6 +108,12 @@ pub struct Driver { /// Handle slots promised to accepts and handle receives that have been asked /// for but not yet delivered. See [`Driver::submit`]. reserved_handles: usize, + /// Armed multishot accepts, and the part of `reserved_handles` promised to + /// them: one each, as far as the ceiling allows. The promises are pooled + /// rather than owned per operation, because the poll budget is shared by + /// every multishot accept, and any of them may spend any of these slots. + multishot_armed: usize, + multishot_reserved: usize, native_pending: usize, /// Operations delivering through `work_port` that have not retired: the /// ring's occupancy can never exceed this, and admission keeps it at or @@ -186,6 +196,8 @@ impl Driver { refs: 0, outstanding: 0, reserved_handles: 0, + multishot_armed: 0, + multishot_reserved: 0, native_pending: 0, undelivered: 0, config, @@ -276,6 +288,7 @@ impl Driver { external_wait: false, fs: false, reserved_handles: 0, + multishot: false, native, port: false, previous, @@ -331,6 +344,13 @@ impl Driver { let op = self.ops.remove(id.key)?; self.connect_deadlines.cancel(id.key); self.reserved_handles -= op.reserved_handles; + if op.multishot { + self.multishot_armed -= 1; + if self.multishot_reserved > self.multishot_armed { + self.multishot_reserved -= 1; + self.reserved_handles -= 1; + } + } if op.native { self.native_pending -= 1; } @@ -882,6 +902,7 @@ impl Driver { operation, Operation::Accept { .. } | Operation::RecvHandle )); + let multishot = matches!(operation, Operation::Accept { multishot: true }); // Only a reserving operation is gated: reads, writes and everything else // on an existing handle must still work at the ceiling. if reserve != 0 && self.handles.remaining() < self.reserved_handles + reserve { @@ -889,7 +910,14 @@ impl Driver { } let op = self.new_op(Some(h), token)?; self.reserved_handles += reserve; - self.ops.get_mut(op.key).expect("new op").reserved_handles = reserve; + let new = self.ops.get_mut(op.key).expect("new op"); + if multishot { + new.multishot = true; + self.multishot_armed += 1; + self.multishot_reserved += 1; + } else { + new.reserved_handles = reserve; + } self.assert_reservations(); if let Err(e) = self.backend.submit(Request { op, @@ -1098,34 +1126,53 @@ impl Driver { /// /// Releasing the reservation before creating the handle is what makes it /// real: `new_handle` refuses slots that are still promised, so the only - /// thing that can spend this one is the accept that reserved it. A multishot - /// accept stays armed and reserves again for its next connection; if the - /// loop is at its ceiling it holds nothing, and its next connection is - /// refused the way a fresh submission would be. + /// thing that can spend this one is the accept that reserved it. + /// + /// A multishot accept spends one of the slots pooled for every multishot + /// accept, if any is left, and the pool is then refilled as far as the + /// ceiling allows. The backend delivered this connection within the poll's + /// [`Budget`], which counts exactly the free slots not promised to + /// single-shot operations, so there is room for it either way. fn attach_reserved( &mut self, transport: B::Detached, id: OpId, token: Token, - terminal: bool, ) -> Result { - let held = self - .ops - .get_mut(id.key) - .map_or(0, |op| std::mem::take(&mut op.reserved_handles)); - self.reserved_handles -= held; - let attached = self.attach(transport, token); - if !terminal && held != 0 && self.handles.remaining() > self.reserved_handles { - self.reserved_handles += held; - if let Some(op) = self.ops.get_mut(id.key) { - op.reserved_handles = held; - } else { - self.reserved_handles -= held; + let multishot = self.ops.get(id.key).is_some_and(|op| op.multishot); + let attached = if multishot { + if self.multishot_reserved != 0 { + self.multishot_reserved -= 1; + self.reserved_handles -= 1; } - } + let attached = self.attach(transport, token); + while self.multishot_reserved < self.multishot_armed + && self.handles.remaining() > self.reserved_handles + { + self.multishot_reserved += 1; + self.reserved_handles += 1; + } + attached + } else { + let held = self + .ops + .get_mut(id.key) + .map_or(0, |op| std::mem::take(&mut op.reserved_handles)); + self.reserved_handles -= held; + self.attach(transport, token) + }; self.assert_reservations(); attached } + /// What the backend may deliver in the next poll (turnloop#77): every free + /// handle slot that no single-shot accept or handle receive is holding. A + /// multishot accept's connections are counted against it, so a batch can + /// never outrun the ceiling; nothing else is. + fn budget(&self) -> Budget { + Budget { + accepts: self.handles.remaining() - (self.reserved_handles - self.multishot_reserved), + } + } /// Register an owning transport on this loop; failure drops the rejected transport. /// /// Reports `ResourceLimit` at the handle ceiling, and also when the only @@ -1518,23 +1565,19 @@ impl Driver { Ok(Outcome::Resolved(addresses)) => OpResult::Resolved(addresses), Ok(Outcome::Exited(status)) => OpResult::Exited(status), Ok(Outcome::Signal(signal)) => OpResult::Signal(signal), - Ok(Outcome::PipeAccepted(d)) => { - match self.attach_reserved(d, e.op, op.token, e.terminal) { - Ok(conn) => OpResult::PipeAccepted { conn }, - Err(e) => OpResult::Err(e), - } - } - Ok(Outcome::HandleReceived(d)) => { - match self.attach_reserved(d, e.op, op.token, e.terminal) { - Ok(handle) => OpResult::HandleReceived { handle }, - Err(e) => OpResult::Err(e), - } - } + Ok(Outcome::PipeAccepted(d)) => match self.attach_reserved(d, e.op, op.token) { + Ok(conn) => OpResult::PipeAccepted { conn }, + Err(e) => OpResult::Err(e), + }, + Ok(Outcome::HandleReceived(d)) => match self.attach_reserved(d, e.op, op.token) { + Ok(handle) => OpResult::HandleReceived { handle }, + Err(e) => OpResult::Err(e), + }, Ok(Outcome::HandleSent) => OpResult::HandleSent, Err(e) => OpResult::Err(e), Ok(Outcome::Connected) => OpResult::Connected, Ok(Outcome::Accepted { transport, peer }) => { - match self.attach_reserved(transport, e.op, op.token, e.terminal) { + match self.attach_reserved(transport, e.op, op.token) { Ok(conn) => OpResult::Accepted { conn, peer }, Err(e) => OpResult::Err(e), } @@ -1593,7 +1636,7 @@ impl Driver { if timeout == Some(Duration::ZERO) || queued || notified || !self.notifier.park() { timeout = Some(Duration::ZERO); } - let poll = self.backend.poll(timeout, &mut self.events); + let poll = self.backend.poll(timeout, self.budget(), &mut self.events); self.notifier.running(); let poll = poll?; waits = poll.waits; @@ -1989,6 +2032,7 @@ mod clock_contract { fn poll( &mut self, timeout: Option, + _budget: Budget, events: &mut Vec>, ) -> Result { assert_eq!(timeout, Some(Duration::ZERO)); From 4db32bed5983fe909c5e80ab3539157cbee04594 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:54:31 +0200 Subject: [PATCH 3/4] Add Loop::rebuild, which refuses instead of dropping pool completions A loop's Config is fixed at construction, so raising a profile means building a new loop, and dropping the old one closes its WorkPort: a pool job already in flight runs to the end and its completion is discarded. Perry works around this by always submitting pool work at its network profile. Of the two fixes #43 offers, growing a live loop in place is not feasible without redesigning structures other threads or the kernel hold: the post and result rings are lock-free and written by other threads (their capacity is their mask and their backpressure bound), and the Windows OVERLAPPED slab is contiguous because the kernel holds its addresses; max_operations and max_handles also size storage in every backend. So this takes the second option and makes replacement safe and loud. Driver::rebuild(config) builds a new loop in place, but only when the old one owes nothing a new loop could not deliver: no handle, no operation in flight (pool jobs, pool lookups, typed file requests and external waits included - every one is counted in `outstanding` until its completion is handed to the host), no queued completion, no result in the ring and no queued post. Otherwise it returns WouldBlock and changes nothing: the host turns until the work is delivered and retries, the same contract as detach. A config that Driver::new rejects fails the same way. Handles, notifiers, posters and the integration primitive from before belong to the old loop, as documented. Test: the issue's acceptance scenario, adapted to the option taken. A small profile loop submits a pool job that blocks on a gate; an upgrade to a larger profile is refused with WouldBlock while the job runs, and again while its result is waiting in the ring; the job's completion is then delivered exactly once (and nothing follows it); the upgrade succeeds, the new loop holds 256 handles where the old held 4, still refuses a rebuild while handles are live, runs pool work, and rebuilds again once drained. Closes #43 --- DESIGN.md | 6 ++ crates/turnloop-contract/src/extended.rs | 86 +++++++++++++++++++++++ crates/turnloop-contract/src/lib.rs | 4 ++ crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/driver.rs | 43 ++++++++++++ 5 files changed, 140 insertions(+) diff --git a/DESIGN.md b/DESIGN.md index 52e7334..19c7134 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -606,6 +606,12 @@ a submission past the ceiling is refused with `ResourceLimit`, so the ring can never be full when a worker pushes. `pooled_buffers` is the remaining per-loop cost that scales with configuration, the other half of #88. +A loop's configuration cannot grow in place: the rings above are written by +other threads and the `OVERLAPPED` slab is held by the kernel. `Loop::rebuild` +replaces a loop with one built from a new configuration, and refuses with +`WouldBlock` while the old one still owes anything — a handle, an operation +(pool jobs included), a queued completion or post — so a profile upgrade can +never discard an in-flight pool completion (#43). **Reaching the ceiling is backpressure.** An operation whose completion creates a handle — an accept, a handle receive — reserves its handle slot when it is diff --git a/crates/turnloop-contract/src/extended.rs b/crates/turnloop-contract/src/extended.rs index f46af3c..83842b4 100644 --- a/crates/turnloop-contract/src/extended.rs +++ b/crates/turnloop-contract/src/extended.rs @@ -1087,3 +1087,89 @@ pub fn undelivered_pool_results_are_bounded() { assert!(out.is_empty()); assert!(!l.alive()); } + +/// A loop is rebuilt with a larger profile without losing a pool completion +/// (issue #43). +/// +/// Dropping a loop to rebuild it discards its in-flight pool results: the job +/// runs and its completion goes nowhere. `rebuild` refuses while anything is +/// owed — the job running, then its result waiting in the ring — leaves the +/// loop working, and succeeds once the job has been delivered exactly once. +pub fn rebuild_never_drops_pool_completions() { + let small = Config { + max_handles: 4, + max_operations: 8, + pooled_buffers: 0, + ..Config::default() + }; + let large = Config { + max_handles: 256, + max_operations: 512, + ..Config::default() + }; + let mut l = Driver::::new(small).expect("small profile"); + let gate = Gate::new(); + let entered = Arc::new(AtomicBool::new(false)); + let job = { + let (gate, entered) = (gate.clone(), entered.clone()); + l.blocking( + move || { + entered.store(true, Ordering::Release); + gate.wait(); + Ok(Payload::U64(42)) + }, + Token(1), + ) + .expect("job") + }; + until("the job is running", || entered.load(Ordering::Acquire)); + let refused = l.rebuild(large).expect_err("the job is in flight"); + assert_eq!(refused.kind, ErrorKind::WouldBlock); + gate.release(); + until("the result was published", || { + let s = pool_stats(); + s.busy == 0 && s.queued == 0 + }); + let refused = l.rebuild(large).expect_err("the result is not delivered"); + assert_eq!(refused.kind, ErrorKind::WouldBlock); + let mut out = Completions::default(); + // `settle` also proves that nothing follows the one completion. + assert!(matches!( + settle(&mut l, job, &mut out), + OpResult::Blocking(Payload::U64(42)) + )); + l.rebuild(large).expect("nothing is owed"); + // The larger profile is in effect: far past the small handle ceiling. + let at = l.now() + Duration::from_secs(60); + let timers: Vec = (0..large.max_handles) + .map(|i| { + l.timer(at, None, Token(100 + i as u64)) + .expect("room in the larger profile") + }) + .collect(); + assert_eq!( + l.timer(at, None, Token(9)) + .expect_err("the larger ceiling") + .kind, + ErrorKind::ResourceLimit + ); + // A live handle is owed too, and the rebuilt loop still runs pool work. + assert_eq!( + l.rebuild(small).expect_err("live handles").kind, + ErrorKind::WouldBlock + ); + let job = l.blocking(|| Ok(Payload::U64(7)), Token(2)).expect("job"); + assert!(matches!( + settle(&mut l, job, &mut out), + OpResult::Blocking(Payload::U64(7)) + )); + for h in timers { + l.close(h, Token(3)).expect("close"); + } + let deadline = l.now() + Duration::from_secs(10); + while l.alive() { + assert!(l.now() < deadline, "the timers never closed"); + l.turn(Timeout::Now, &mut out).expect("turn"); + } + l.rebuild(small).expect("drained"); +} diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index 496aa07..d630e19 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -287,6 +287,10 @@ mod native { undelivered_pool_results_are_bounded::(); } #[test] + fn rebuild_keeps_in_flight_pool_completions() { + rebuild_never_drops_pool_completions::(); + } + #[test] fn reuse_port_share_binds_twice() { reuse_port_share::(); } diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index 9876cfa..a8a3d63 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -34,6 +34,7 @@ contract!( an_armed_accept_keeps_its_slot, multishot_accept_respects_the_handle_ceiling, undelivered_pool_results_are_bounded, + rebuild_never_drops_pool_completions, ready_timer_liveness, pooled_lease_backpressure, io_and_posts_progress_with_repeating_timers diff --git a/crates/turnloop/src/driver.rs b/crates/turnloop/src/driver.rs index 2892c83..d1c0280 100644 --- a/crates/turnloop/src/driver.rs +++ b/crates/turnloop/src/driver.rs @@ -204,6 +204,49 @@ impl Driver { _local: PhantomData, }) } + /// Rebuild this loop from a new `config`, in place (issue #43). + /// + /// A loop's configuration is fixed when it is built, so raising a profile — + /// a process that starts doing network I/O, say — means building a new loop. + /// Dropping the old one to do that is a cancellation nobody asked for: a pool + /// job it accepted still runs, and its result is discarded because its loop + /// is gone, so the job's completion never arrives. This is the checked way + /// to replace a loop, which refuses instead of losing anything. + /// + /// The loop is rebuilt only when it owns nothing a new loop could not + /// deliver: no handle, no operation in flight (pool jobs, pool lookups, + /// typed file requests and external waits included), no completion awaiting + /// delivery and no queued post. Otherwise this is `WouldBlock` and nothing + /// changes: turn the loop until that work has been delivered, then retry — + /// the same contract as [`detach`](Self::detach). A `config` that + /// [`Driver::new`] would reject fails the same way, also without changing + /// anything. + /// + /// On success this is a new loop. Handles, operation ids, [`Notifier`]s, + /// [`Poster`]s and an [`integration`](Self::integration) primitive obtained + /// before belong to the old one and no longer reach it: fetch new ones. + /// Quiesce other threads' posters first: a post that lands between the + /// check and the rebuild goes down with the old loop, exactly as it would if + /// the loop were dropped, and later ones are refused with `NotFound`. + /// + /// Growing the old loop in place instead is not offered: its capacities size + /// structures that cannot grow under live use — the lock-free post and + /// result rings that other threads are writing into, and the Windows + /// `OVERLAPPED` slab whose addresses the kernel holds. + pub fn rebuild(&mut self, config: Config) -> Result<()> { + self.assert_owner(); + if self.outstanding != 0 + || self.handles.remaining() != self.config.max_handles + || !self.queued.is_empty() + || !self.poster.is_empty() + || !self.work_port.is_empty() + { + return Err(Error::new(ErrorKind::WouldBlock)); + } + debug_assert_eq!(self.undelivered, 0, "every pool operation is outstanding"); + *self = Self::new(config)?; + Ok(()) + } /// Reserved slots are a subset of free ones. /// /// Every path that moves either side keeps this: a reserving submission From b63d71d70dcb16aa12dfc61294c56fe81e942282 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 20:14:55 +0200 Subject: [PATCH 4/4] wasi 0.3: take a parked accept's fired waitable out of the wait-set Review of #106: when the budget parks a multishot accept whose stream read has already fired, run_ready skips execute(), and execute() -> finish_wait() is the only place that removes the waitable from the wait-set. The accept's outcome is already held in `code`, so the membership serves nothing, and if the host re-delivered the event each wait at the ceiling would return at once: a turn loop that produces nothing (rule 4a). schedule() now removes that waitable when it parks the listener. This runs no outcome side effect and loses no connection: the result stays in `code`, and execute() consumes it when a positive budget resumes the listener, where finish_wait()'s own removal is a harmless repeat (join to 0 is idempotent). A parked accept can only be in one of two states, not yet started (no waitable) or fired (outcome in `code`), because a started, unfired read clears `ready`, so nothing needs to rejoin the set. On wasmtime 46 the stream-read event was delivered once even without this change (checked with instrumentation), so no spin was observed there. The fix keeps the backend from depending on that host behaviour. Test: a new contract scenario reaches this path deliberately. Two listeners share the ceiling; both accepts start waiting on empty backlogs, the second listener's connections take every free slot, and a client then connects to the first, so its read fires with no slot to house the connection. The test asserts turns at the ceiling deliver nothing and block (at most 20 turns in 200 ms, one native wait each), and that freeing one slot delivers exactly that connection. It runs on native and Windows, and on WASI 0.2 and 0.3 together with the existing multishot ceiling test, which tests/wasi.rs did not run before. --- crates/turnloop-contract/src/lib.rs | 118 ++++++++++++++++++++++ crates/turnloop-contract/tests/wasi.rs | 12 +++ crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/backend/wasi_p3.rs | 14 +++ 4 files changed, 145 insertions(+) diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index d630e19..5edbe5d 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -235,6 +235,10 @@ mod native { multishot_accept_respects_the_handle_ceiling::(); } #[test] + fn multishot_accept_parked_after_its_read_fired() { + super::multishot_accept_parked_after_its_read_fired::(); + } + #[test] fn cancellation_close() { cancel_close_ordering::(); } @@ -1309,6 +1313,120 @@ pub fn multishot_accept_respects_the_handle_ceiling() { drop(clients); } +/// A multishot accept whose connection arrives once the budget is spent (#77). +/// +/// Two listeners share the handle ceiling. Both accepts are started while the +/// backlogs are empty, then the second listener's connections take every free +/// slot, so the first listener's accept is already waiting on the backend when +/// its connection arrives with nothing left to house it. That connection must +/// wait without the loop spinning — on WASI 0.3 the fired waitable would +/// otherwise stay in the wait-set — and must be delivered, not lost, as soon as +/// a slot frees up. +pub fn multishot_accept_parked_after_its_read_fired() { + let mut l = Driver::::new(Config { + max_handles: 4, + max_operations: 16, + ..Config::default() + }) + .expect("loop"); + let first = l + .tcp_listen(localhost(), &ListenOpts::default()) + .expect("listen"); + let second = l + .tcp_listen(localhost(), &ListenOpts::default()) + .expect("listen"); + let first_addr = l.local_addr(first).expect("addr"); + let second_addr = l.local_addr(second).expect("addr"); + let first_accept = l.accept_start(first, Token(1)).expect("accept"); + let second_accept = l.accept_start(second, Token(2)).expect("accept"); + let mut out = Completions::with_capacity(16); + // Both accepts start waiting on empty backlogs. + for _ in 0..3 { + l.turn(Timeout::After(Duration::from_millis(5)), &mut out) + .expect("turn"); + assert!(out.is_empty(), "nothing has connected yet"); + } + // The second listener's two connections take both free slots. + let mut clients: Vec = (0..2) + .map(|_| std::net::TcpStream::connect(second_addr).expect("client")) + .collect(); + let mut live = Vec::new(); + let deadline = Instant::now() + Duration::from_secs(10); + while live.len() < 2 { + assert!(Instant::now() < deadline, "only {} accepted", live.len()); + l.turn(Timeout::After(Duration::from_millis(20)), &mut out) + .expect("turn"); + for c in out.drain() { + let OpResult::Accepted { conn, .. } = c.result else { + panic!("unexpected {:?}", c.result); + }; + assert_eq!(c.op, Some(second_accept)); + live.push(conn); + } + } + clients.push(std::net::TcpStream::connect(first_addr).expect("client")); + // At the ceiling the first listener's connection waits: turns deliver + // nothing, and each blocks for its timeout instead of spinning. + let quiet = Instant::now(); + let mut turns = 0; + while quiet.elapsed() < Duration::from_millis(200) { + let info = l + .turn(Timeout::After(Duration::from_millis(20)), &mut out) + .expect("turn at the ceiling"); + assert!(info.os_waits + info.discovery_polls <= 1); + assert!( + out.is_empty(), + "nothing past the ceiling: {:?}", + out[0].result + ); + turns += 1; + } + assert!( + turns <= 20, + "{turns} turns in 200 ms at the ceiling: a spin" + ); + // One freed slot delivers the waiting connection to the first listener. + l.close(live.pop().expect("a connection"), Token(3)) + .expect("close"); + let mut delivered = false; + while !delivered { + assert!(Instant::now() < deadline, "the waiting connection was lost"); + l.turn(Timeout::After(Duration::from_millis(20)), &mut out) + .expect("turn"); + for c in out.drain() { + match c.result { + OpResult::Closed => {} + OpResult::Accepted { conn, .. } => { + assert_eq!(c.op, Some(first_accept)); + assert!(!delivered, "one connection, one delivery"); + live.push(conn); + delivered = true; + } + other => panic!("unexpected {other:?}"), + } + } + } + assert!(l.stop(first_accept)); + assert!(l.stop(second_accept)); + for h in live.drain(..).chain([first, second]) { + l.close(h, Token(4)).expect("close"); + } + let deadline = Instant::now() + Duration::from_secs(5); + while l.alive() { + assert!(Instant::now() < deadline, "the loop never drained"); + l.turn(Timeout::After(Duration::from_millis(20)), &mut out) + .expect("drain"); + for c in out.drain() { + assert!( + matches!(c.result, OpResult::Stopped | OpResult::Closed), + "unexpected {:?}", + c.result + ); + } + } + drop(clients); +} + pub fn pooled_lease_backpressure() { let mut l = Driver::::new(Config { pooled_buffers: 1, diff --git a/crates/turnloop-contract/tests/wasi.rs b/crates/turnloop-contract/tests/wasi.rs index 851b0de..968409d 100644 --- a/crates/turnloop-contract/tests/wasi.rs +++ b/crates/turnloop-contract/tests/wasi.rs @@ -227,6 +227,18 @@ fn nodelay_accept_default_rejects_the_listener() { fn no_spin() { contract::no_spin::(); } +/// turnloop#77: a multishot accept waits at the handle ceiling without +/// destroying connections, and a turn at the ceiling blocks instead of spinning. +#[test] +fn multishot_accept_waits_at_the_handle_ceiling() { + contract::multishot_accept_respects_the_handle_ceiling::(); +} +/// A connection that reaches an already-started accept once the budget is +/// spent waits without a spin (0.3: the fired waitable leaves the wait-set). +#[test] +fn multishot_accept_parked_after_its_read_fired() { + contract::multishot_accept_parked_after_its_read_fired::(); +} /// DESIGN §10.4a and the `PollInfo` contract: the private WASI deadline /// pollable (0.2) / subtask (0.3) is this wait's timeout, not native work. #[test] diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index a8a3d63..7ed9933 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -33,6 +33,7 @@ contract!( accept_reserves_its_handle_slot, an_armed_accept_keeps_its_slot, multishot_accept_respects_the_handle_ceiling, + multishot_accept_parked_after_its_read_fired, undelivered_pool_results_are_bounded, rebuild_never_drops_pool_completions, ready_timer_liveness, diff --git a/crates/turnloop/src/backend/wasi_p3.rs b/crates/turnloop/src/backend/wasi_p3.rs index 47faa0e..82b5843 100644 --- a/crates/turnloop/src/backend/wasi_p3.rs +++ b/crates/turnloop/src/backend/wasi_p3.rs @@ -266,6 +266,20 @@ impl WasiP3 { } else if held && !r.throttled { r.throttled = true; self.throttled.push_back(h); + // A parked accept whose read already fired keeps its outcome in + // `code`, which `execute` consumes once the budget resumes it; only + // `finish_wait` would otherwise take the waitable out of the set. + // Leave it there and every wait while the ceiling holds could return + // it again at once, a spin that produces nothing (rule 4a). Removing + // it runs no side effect and loses no connection; `finish_wait` + // removes it again harmlessly, since `join(_, 0)` is idempotent. + if let Some(p) = r.heads[0].and_then(|i| self.ops[i].as_ref()) + && let Some((waitable, _)) = p.wait + && p.code.is_some() + && waitable != 0 + { + self.wait_set.remove(waitable); + } } } /// Start a poll with `budget`: a positive one resumes every held listener.