diff --git a/DESIGN.md b/DESIGN.md index 2454b19..19c7134 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -596,8 +596,22 @@ 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. + +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 @@ -607,9 +621,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 11274fd..83842b4 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(); @@ -1004,3 +1003,173 @@ 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()); +} + +/// 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 e16bfbc..eec0027 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -231,6 +231,14 @@ 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 multishot_accept_parked_after_its_read_fired() { + super::multishot_accept_parked_after_its_read_fired::(); + } + #[test] fn cancellation_close() { cancel_close_ordering::(); } @@ -287,6 +295,14 @@ 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 rebuild_keeps_in_flight_pool_completions() { + rebuild_never_drops_pool_completions::(); + } + #[test] fn reuse_port_share_binds_twice() { reuse_port_share::(); } @@ -1413,6 +1429,232 @@ 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); +} + +/// 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 0f91c77..2bdee7e 100644 --- a/crates/turnloop-contract/tests/wasi.rs +++ b/crates/turnloop-contract/tests/wasi.rs @@ -231,6 +231,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 1923c9a..b8f65eb 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -34,6 +34,10 @@ contract!( paged_growth_preserves_handles, 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, pooled_lease_backpressure, io_and_posts_progress_with_repeating_timers 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..82b5843 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,48 @@ 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); + // 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. + 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 +352,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 +381,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 +411,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 +638,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 +735,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 +1214,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/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 7290b9d..276996c 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,12 +50,19 @@ 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. 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, @@ -101,7 +108,17 @@ 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 + /// below `PoolConfig::max_undelivered`, which sizes the ring (issue #88). + undelivered: usize, config: Config, _local: PhantomData>, } @@ -115,6 +132,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 +147,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 @@ -170,11 +196,57 @@ impl Driver { refs: 0, outstanding: 0, reserved_handles: 0, + multishot_armed: 0, + multishot_reserved: 0, native_pending: 0, + undelivered: 0, config, _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 @@ -259,7 +331,9 @@ impl Driver { external_wait: false, fs: false, reserved_handles: 0, + multishot: false, native, + port: false, previous, next: None, }) @@ -291,13 +365,41 @@ 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); 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; } + if op.port { + self.undelivered -= 1; + } if let Some(previous) = op.previous { self.ops.get_mut(previous.key).expect("previous").next = op.next; } @@ -856,6 +958,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 { @@ -863,7 +966,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, @@ -1072,34 +1182,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 @@ -1122,8 +1251,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) { @@ -1148,6 +1279,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)?; @@ -1175,6 +1309,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), @@ -1369,7 +1506,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); @@ -1482,23 +1621,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), } @@ -1557,7 +1692,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; @@ -1953,6 +2088,7 @@ mod clock_contract { fn poll( &mut self, timeout: Option, + _budget: Budget, events: &mut Vec>, ) -> Result { assert_eq!(timeout, Some(Duration::ZERO)); @@ -2358,3 +2494,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