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
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,9 @@ aims to follow [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

### Added

- Per-shard `shard_iowq_setup` startup snapshots of the existing best-effort
io-wq worker-limit registration, with requested values, previous limits on
success, and error classification on failure.
- Opt-in `SharedReadBudget` whole-driver quota reservations, retained through
deferred and kernel-owned reads, including bounded-drain leaks. Independent
drivers keep independent admission and shutdown; default constructors are unchanged.
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ assert_eq!(snapshot.delivered + snapshot.orphan_reclaimed, snapshot.submitted);
- `with_shard_policy(ShardPolicy::CapacityAware)` — opt-in capacity-aware routing for positioned reads. Constructors keep `ShardPolicy::RoundRobin` by default.
- `request_shutdown()` — close admission and request cancellation/drain without joining; `is_finished()` reports advisory thread completion, not a clean drain.
- `shutdown_async()` — with the default-off `tokio-runtime` feature, transfer consuming cleanup to Tokio's blocking pool at method call time. See [shutdown ownership and runtime boundaries](docs/shutdown.md).
- `shard_iowq_setup()` — inspect each shard's best-effort startup registration request and result. Success returns the *previous* kernel limits; see [io-wq startup reporting](docs/iowq-startup.md).

### Shard selection

Expand Down
35 changes: 35 additions & 0 deletions docs/iowq-startup.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
# Io-wq registration at driver startup

Each io_uring shard makes one best-effort `IORING_REGISTER_IOWQ_MAX_WORKERS`
request after the ring and its completion eventfd are ready. The request remains
`[16, 0]`: sixteen bounded workers per NUMA node, and zero to leave the
unbounded-worker limit unchanged. A kernel that rejects this registration still
allows the driver to start and serve reads.

`UringDriver::shard_iowq_setup()` returns one immutable `IoWqSetup` per shard in
shard order. It copies startup records and does not call the kernel or add work
to the read path. Every record contains the requested values and an
`IoWqRegistration` result:

- `Registered { previous_limits }` means the request succeeded. The
`io-uring` registration call overwrites its input array with the limits from
**before** this request. These values are historical; they are not a query of
the resulting limit or active worker count.
- `Failed { error_kind, raw_os_error }` records why registration failed. Some
errors have no errno. A failed best-effort setting is distinct from the
real-read startup probe and does not automatically disable io_uring or poison
an application's unsupported-disk cache.

This snapshot can explain whether an attempted bounded-worker setting was
accepted on a particular ring. It does not establish a process-wide thread cap,
current kernel worker count, per-disk CPU share, or performance benefit. Shards
and other applications may share kernel io-wq resources; the effective limit
can also differ from the requested value. Operational validation must combine
this startup record with runtime process/kernel observations.

Unit tests inject success with a rewritten previous-limit array and failure
with and without errno. A native test creates two shards, reads byte-exact data,
and checks that their startup records remain stable. It accepts either
registration outcome because kernel support and policy vary; it does not force
an io-wq worker to run. See [acceptance status](optimization-status.md) for the
latest CI result.
26 changes: 20 additions & 6 deletions docs/optimization-status.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,14 @@ integration are separate gates; a checked code item does not close the roadmap.
| --- | --- | --- |
| 1.1–1.4: benchmark schema, timing, diagnostics, direct positive gate | Merged in [#15](https://github.com/rustfs/uring/pull/15), CI passed | Diagnostics overhead and real LocalIoBackend baseline |
| 1.5: ABBA evidence tooling | Explicit same-binary calibration implemented/reviewed; 20 gate/cleanup tests pass | Stable native calibration, valid comparison and application integration |
| 2.1–2.4: completion recovery, direct integrity, byte admission, guarantee boundaries | Implemented and independently reviewed; native regression suite passed | PR merge; application-wide limits remain separate |
| 2.4 API follow-up: shutdown control and runtime adapter | Implemented/reviewed; native CI #55 passed at `a942ac1` | PR merge; no hard cleanup deadline or performance claim |
| 3.1–3.3: bounded turns, explicit batch notifications, cancellation efficiency | Implemented and independently reviewed; native regression suite passed | CPU/syscall and tail-latency comparison; PR merge |
| 4.1: capacity-aware routing | Implemented and independently reviewed, opt-in; native tests passed | Controlled slow-shard/mixed-load performance; PR merge |
| 2.1–2.4: completion recovery, direct integrity, byte admission, guarantee boundaries | Implemented/reviewed; native regression suite passed; #16 merged | Application-wide limits remain separate |
| 2.4 API follow-up: shutdown control and runtime adapter | Implemented/reviewed; native CI #55 passed; #16 merged | No hard cleanup deadline or performance claim |
| 3.1–3.3: bounded turns, explicit batch notifications, cancellation efficiency | Implemented/reviewed; native regression suite passed; #16 merged | CPU/syscall and tail-latency comparison |
| 4.1: capacity-aware routing | Implemented/reviewed, opt-in; native tests passed; #16 merged | Controlled slow-shard/mixed-load performance |
| 4.2: system-wide budgets and probe offload | Application probe offload, driver-thread budget and logical chunk limit implemented/reviewed in draft PRs #8072/#8074/#8076 | Full CI/merge; physical byte/result/io-wq budgets and [dependency wiring](rustfs-integration.md) remain open |
| 4.2 pool prerequisite | Library whole-driver shared quota implemented/reviewed; native CI #57 passed at `b09b966` | Application dependency/fallback/result ownership wiring and merge |
| 5: owned buffers and direct FD cache | Dual-mode exact invalidation implemented/reviewed in application draft PR #8075; direct caching and owned buffers not enabled | Full CI/merge, profile evidence, lease/pool lifetime and complete invalidation integration |
| 4.2 pool prerequisite | Library whole-driver shared quota implemented/reviewed; native CI #57 passed; #16 merged | Application dependency/fallback/result ownership wiring |
| 4.2 io-wq registration evidence | Per-shard startup query reviewed; [PR #17 CI](https://github.com/rustfs/uring/actions/runs/35808804770) passed at `193ab21` | PR merge; actual worker count and process-wide budget remain unproven |
| 5: owned buffers and direct FD cache | Dual-mode exact invalidation in application draft PR #8075 has passing CI; direct caching and owned buffers not enabled | PR merge, profile evidence, lease/pool lifetime and complete invalidation integration |
| 6: ordered streaming prefetch | [Example-only contract experiment](ordered-prefetch.md) implemented/reviewed; 11 portable tests and native CLI CI #54 passed | Production consumer contract, bitrot/S3 and performance evidence |
| 7: advanced ring/runtime modes | Not enabled or implemented | Earlier gates, capability/fallback and isolated benefit evidence |

Expand All @@ -38,6 +39,18 @@ final-head CI and merge status separately.

## Follow-up evidence and examples

The [io-wq startup record](iowq-startup.md) captures the existing registration
result without changing its requested `[16, 0]` values, fallback, or read path.
It records kernel-returned *previous* limits on success and error kind/errno on
failure. Deterministic injected tests and a two-shard native test are added;
two independent reviews passed. [CI at `193ab21`](https://github.com/rustfs/uring/actions/runs/35808804770)
passed all three jobs: 147 native all-feature tests with no skips, the positive
O_DIRECT marker, ordered-prefetch and both benchmark CLI smokes, and restricted/unrestricted
Docker coverage. The new native test read exact bytes from two shards and
observed stable startup records; it did not force io-wq worker execution. The
snapshot does not measure effective worker counts or establish a process-wide
cap. A later documentation-only head has its own CI status.

`3931124` implements `SharedReadBudget` whole-driver reservations, with public
tests in `7e94f36`. Seven deterministic private tests and nine public/native
tests cover concurrent accounting, large `usize` quotas, message/deferred
Expand Down Expand Up @@ -136,6 +149,7 @@ establish stable calibration without loosening gates after observing results.
- [Public read and admission contracts](../README.md)
- [Shutdown ownership and Tokio adapter](shutdown.md)
- [Shared whole-driver read-budget reservations](shared-read-budget.md)
- [Io-wq startup registration evidence](iowq-startup.md)
- [Completion recovery and direct integrity](fault-recovery.md)
- [Cancellation and eventfd behavior](cancellation-efficiency.md)
- [Driver turn fairness](driver-fairness.md)
Expand Down
1 change: 1 addition & 0 deletions src/admission_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -187,6 +187,7 @@ fn fake_driver(limits: ReadLimits) -> (UringDriver, mpsc::Receiver<Msg>) {
tx,
handle: None,
stats: Arc::new(DriverStats::default()),
iowq_setup: mock_iowq_setup(),
sem: count,
wake_efd: Arc::new(EventFd::new().unwrap()),
}],
Expand Down
1 change: 1 addition & 0 deletions src/batch_read_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ fn driver(shard_count: usize, capacity: usize, bytes: Option<usize>) -> (UringDr
tx,
handle: None,
stats: Arc::new(DriverStats::default()),
iowq_setup: mock_iowq_setup(),
sem,
wake_efd: Arc::new(EventFd::new().unwrap()),
}
Expand Down
120 changes: 118 additions & 2 deletions src/driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,109 @@ const IDLE_HEARTBEAT: Duration = Duration::from_secs(1);
/// default.
const IOWQ_MAX_BOUNDED_WORKERS: u32 = 16;

/// Startup result of one ring's best-effort io-wq worker-limit registration.
///
/// Registration success reports the limits that were in effect *before* the
/// request. It does not report the resulting limit or the number of workers.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct IoWqSetup {
/// Requested `[bounded, unbounded]` workers per NUMA node. Zero leaves a
/// worker class unchanged; the driver currently requests `[16, 0]`.
pub requested: [u32; 2],
/// The registration result, captured during this shard's startup.
pub result: IoWqRegistration,
}

/// Outcome of one ring's best-effort io-wq worker-limit registration.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IoWqRegistration {
/// The kernel accepted the request and returned the previous limits.
Registered {
/// `[bounded, unbounded]` limits before this registration call.
previous_limits: [u32; 2],
},
/// Registration failed; the driver continued to start normally.
Failed {
/// Standard I/O error category, including errors without an errno.
error_kind: io::ErrorKind,
/// Kernel errno when one was provided.
raw_os_error: Option<i32>,
},
}

fn register_iowq_setup(register: impl FnOnce(&mut [u32; 2]) -> io::Result<()>) -> IoWqSetup {
let requested = [IOWQ_MAX_BOUNDED_WORKERS, 0];
let mut limits = requested;
let result = match register(&mut limits) {
Ok(()) => IoWqRegistration::Registered { previous_limits: limits },
Err(error) => IoWqRegistration::Failed {
error_kind: error.kind(),
raw_os_error: error.raw_os_error(),
},
};
IoWqSetup { requested, result }
}

#[cfg(test)]
fn mock_iowq_setup() -> IoWqSetup {
register_iowq_setup(|_| Err(io::Error::other("mock shard has no io-wq ring")))
}

#[cfg(test)]
mod iowq_setup_tests {
use super::*;

#[test]
fn successful_registration_preserves_the_request_and_labels_kernel_output_as_previous() {
let mut calls = 0;
let setup = register_iowq_setup(|limits| {
calls += 1;
assert_eq!(*limits, [16, 0]);
*limits = [37, 5];
Ok(())
});
assert_eq!(calls, 1);
assert_eq!(setup.requested, [16, 0]);
assert_eq!(
setup.result,
IoWqRegistration::Registered {
previous_limits: [37, 5]
}
);
}

#[test]
fn failed_registration_records_errno_without_claiming_previous_limits() {
for errno in [libc::EINVAL, libc::EOPNOTSUPP, libc::EPERM] {
let setup = register_iowq_setup(|limits| {
assert_eq!(*limits, [16, 0]);
*limits = [91, 92]; // A failed call cannot supply valid previous limits.
Err(io::Error::from_raw_os_error(errno))
});
assert_eq!(setup.requested, [16, 0]);
assert_eq!(
setup.result,
IoWqRegistration::Failed {
error_kind: io::Error::from_raw_os_error(errno).kind(),
raw_os_error: Some(errno),
}
);
}
}

#[test]
fn failed_registration_without_errno_keeps_its_error_kind() {
let setup = register_iowq_setup(|_| Err(io::Error::other("registration unavailable")));
assert_eq!(
setup.result,
IoWqRegistration::Failed {
error_kind: io::ErrorKind::Other,
raw_os_error: None,
}
);
}
}

/// Consecutive non-transient `ring.submit()` failures the driver tolerates
/// before it stops retrying silently and shuts the shard down, so callers get a
/// driver-gone error and fall back to the std backend instead of stalling
Expand Down Expand Up @@ -980,6 +1083,7 @@ struct Shard {
tx: mpsc::Sender<Msg>,
handle: Option<JoinHandle<()>>,
stats: Arc<DriverStats>,
iowq_setup: IoWqSetup,
/// Backpressure permits (one per allowed in-flight op on this ring). Closed
/// when the driver thread exits so any waiting `ReadHandle` resolves with a
/// driver-gone error instead of hanging (rustfs/backlog#1102).
Expand Down Expand Up @@ -1263,8 +1367,7 @@ impl UringDriver {
// TasksMax/RLIMIT_NPROC (rustfs/backlog#1169). Best-effort: 0 leaves the
// unbounded pool unchanged, and a kernel without this op (< 5.15) keeps
// the default — neither is fatal to a working ring.
let mut iowq_max = [IOWQ_MAX_BOUNDED_WORKERS, 0u32];
let _ = ring.submitter().register_iowq_max_workers(&mut iowq_max);
let iowq_setup = register_iowq_setup(|limits| ring.submitter().register_iowq_max_workers(limits));

let wake_efd = Arc::new(EventFd::new().map_err(ProbeFailure::Setup)?);
let thread_wake = Arc::clone(&wake_efd);
Expand Down Expand Up @@ -1305,6 +1408,7 @@ impl UringDriver {
tx,
handle: Some(handle),
stats,
iowq_setup,
sem,
wake_efd,
})
Expand Down Expand Up @@ -1651,6 +1755,16 @@ impl UringDriver {
}
}

/// Io-wq registration results captured when the shards started, in shard
/// order. Querying only copies stored data; it does not call the kernel.
///
/// A successful registration reports the *previous* limits. Neither that
/// value nor the requested value proves the resulting worker count or a
/// process-wide thread cap. A failed registration does not disable reads.
pub fn shard_iowq_setup(&self) -> Vec<IoWqSetup> {
self.shards.iter().map(|shard| shard.iowq_setup).collect()
}

/// Counters summed across every shard. The conservation identities the
/// cancel-safety tests assert (`submitted == delivered + orphan_reclaimed`,
/// `in_flight == 0` after a clean drain) hold per shard, so they hold for
Expand Down Expand Up @@ -2625,6 +2739,7 @@ mod shared_budget_reservation_tests {
tx,
handle: None,
stats: Arc::new(DriverStats::default()),
iowq_setup: mock_iowq_setup(),
sem,
wake_efd: Arc::new(EventFd::new().expect("mock wake fd")),
}],
Expand Down Expand Up @@ -2827,6 +2942,7 @@ mod shutdown_request_tests {
tx,
handle: None,
stats: Arc::new(DriverStats::default()),
iowq_setup: mock_iowq_setup(),
sem,
wake_efd: Arc::new(EventFd::new().expect("mock wake fd")),
}
Expand Down
3 changes: 2 additions & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -56,5 +56,6 @@ pub use diagnostics::{DIAGNOSTICS_SAMPLE_INTERVAL, DiagnosticsSnapshot, LatencyH

#[cfg(target_os = "linux")]
pub use driver::{
MAX_BATCH_READS, ProbeFailure, ReadHandle, ReadLimits, ReadRequest, ShardPolicy, SharedReadBudget, StatsSnapshot, UringDriver,
IoWqRegistration, IoWqSetup, MAX_BATCH_READS, ProbeFailure, ReadHandle, ReadLimits, ReadRequest, ShardPolicy,
SharedReadBudget, StatsSnapshot, UringDriver,
};
1 change: 1 addition & 0 deletions src/shard_policy_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ fn driver(limits: ReadLimits) -> (UringDriver, Vec<mpsc::Receiver<Msg>>) {
tx,
handle: None,
stats: Arc::new(DriverStats::default()),
iowq_setup: mock_iowq_setup(),
sem,
wake_efd: Arc::new(EventFd::new().unwrap()),
}
Expand Down
Loading
Loading