Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 26 additions & 5 deletions DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.

Expand Down
177 changes: 173 additions & 4 deletions crates/turnloop-contract/src/extended.rs
Original file line number Diff line number Diff line change
Expand Up @@ -835,10 +835,9 @@ pub fn handoff_accept_exactly_once<B: Backend>() {
}));
}
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();
Expand Down Expand Up @@ -1004,3 +1003,173 @@ pub fn kernel_accept_exactly_once<B: Backend>() {
}
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<B: Backend>() {
const CEILING: usize = 8;
let mut l = Driver::<B>::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<B: Backend>() {
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::<B>::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<Handle> = (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");
}
Loading