From 193ab21f28c161bbf3791b5f8efd12c41bcb7ffa Mon Sep 17 00:00:00 2001 From: Hauser Date: Wed, 23 Sep 2026 10:01:48 +0800 Subject: [PATCH 1/2] feat(driver): expose per-shard io-wq registration results Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- CHANGELOG.md | 3 + README.md | 1 + docs/iowq-startup.md | 35 +++++++++++ docs/optimization-status.md | 21 +++++-- src/admission_tests.rs | 1 + src/batch_read_tests.rs | 1 + src/driver.rs | 120 +++++++++++++++++++++++++++++++++++- src/lib.rs | 3 +- src/shard_policy_tests.rs | 1 + tests/iowq_setup.rs | 91 +++++++++++++++++++++++++++ 10 files changed, 268 insertions(+), 9 deletions(-) create mode 100644 docs/iowq-startup.md create mode 100644 tests/iowq_setup.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index abdf337..6add655 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/README.md b/README.md index 03dadae..b411960 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/docs/iowq-startup.md b/docs/iowq-startup.md new file mode 100644 index 0000000..48f89e2 --- /dev/null +++ b/docs/iowq-startup.md @@ -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. diff --git a/docs/optimization-status.md b/docs/optimization-status.md index ff25684..78bb95c 100644 --- a/docs/optimization-status.md +++ b/docs/optimization-status.md @@ -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, deterministic and native tests implemented; two independent reviews passed | New-head native CI; 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 | @@ -38,6 +39,13 @@ 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; a +new-head native CI result remains pending; two independent reviews passed. The snapshot +does not measure effective worker counts or establish a process-wide cap. + `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 @@ -136,6 +144,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) diff --git a/src/admission_tests.rs b/src/admission_tests.rs index a984cac..b7e76e0 100644 --- a/src/admission_tests.rs +++ b/src/admission_tests.rs @@ -187,6 +187,7 @@ fn fake_driver(limits: ReadLimits) -> (UringDriver, mpsc::Receiver) { tx, handle: None, stats: Arc::new(DriverStats::default()), + iowq_setup: mock_iowq_setup(), sem: count, wake_efd: Arc::new(EventFd::new().unwrap()), }], diff --git a/src/batch_read_tests.rs b/src/batch_read_tests.rs index 0796c2e..4cb1eb8 100644 --- a/src/batch_read_tests.rs +++ b/src/batch_read_tests.rs @@ -20,6 +20,7 @@ fn driver(shard_count: usize, capacity: usize, bytes: Option) -> (UringDr tx, handle: None, stats: Arc::new(DriverStats::default()), + iowq_setup: mock_iowq_setup(), sem, wake_efd: Arc::new(EventFd::new().unwrap()), } diff --git a/src/driver.rs b/src/driver.rs index f1f4f5c..d700cb0 100644 --- a/src/driver.rs +++ b/src/driver.rs @@ -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, + }, +} + +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 @@ -980,6 +1083,7 @@ struct Shard { tx: mpsc::Sender, handle: Option>, stats: Arc, + 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). @@ -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); @@ -1305,6 +1408,7 @@ impl UringDriver { tx, handle: Some(handle), stats, + iowq_setup, sem, wake_efd, }) @@ -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 { + 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 @@ -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")), }], @@ -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")), } diff --git a/src/lib.rs b/src/lib.rs index 6b965e1..ca47b5b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -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, }; diff --git a/src/shard_policy_tests.rs b/src/shard_policy_tests.rs index 689981c..830e55e 100644 --- a/src/shard_policy_tests.rs +++ b/src/shard_policy_tests.rs @@ -20,6 +20,7 @@ fn driver(limits: ReadLimits) -> (UringDriver, Vec>) { tx, handle: None, stats: Arc::new(DriverStats::default()), + iowq_setup: mock_iowq_setup(), sem, wake_efd: Arc::new(EventFd::new().unwrap()), } diff --git a/tests/iowq_setup.rs b/tests/iowq_setup.rs new file mode 100644 index 0000000..7335869 --- /dev/null +++ b/tests/iowq_setup.rs @@ -0,0 +1,91 @@ +// Copyright 2024 RustFS Team +// SPDX-License-Identifier: Apache-2.0 + +//! Native Linux io-wq startup reporting. Registration failure remains a valid +//! best-effort outcome; this test does not force that failure or prove a worker +//! limit was enforced by a particular I/O workload. +#![cfg(target_os = "linux")] + +use std::fs::File; +use std::io; +use std::sync::Arc; +use std::time::Duration; + +use rustfs_uring::{IoWqRegistration, IoWqSetup, UringDriver}; + +fn assert_requested_limits(setup: &IoWqSetup) { + assert_eq!(setup.requested, [16, 0], "requested limits must not be replaced by kernel output"); + match &setup.result { + IoWqRegistration::Registered { .. } => { + // The kernel returns PREVIOUS limits. Neither equality with [16, 0] + // nor inequality is portable, and they are not a current-limit query. + } + IoWqRegistration::Failed { + error_kind, + raw_os_error, + } => { + if let Some(errno) = raw_os_error { + assert_eq!(*error_kind, io::Error::from_raw_os_error(*errno).kind()); + } + } + } +} + +#[tokio::test(flavor = "current_thread")] +async fn each_shard_reports_startup_iowq_result_without_changing_read_correctness() { + let driver = match UringDriver::probe_and_start_sharded(8, 2) { + Ok(driver) => driver, + Err(error) => { + assert!(error.is_expected_restriction(), "unexpected io_uring probe failure: {error}"); + eprintln!("SKIP each_shard_reports_startup_iowq_result_without_changing_read_correctness: {error}"); + return; + } + }; + let before = driver.shard_iowq_setup(); + assert_eq!(before.len(), 2, "one setup result is required for every started shard"); + for setup in &before { + assert_requested_limits(setup); + } + + // Default round-robin routes the two reads to distinct owning shards. This + // uses real kernel reads but does not assert that /dev/zero used io-wq. + let file = Arc::new(File::open("/dev/zero").expect("open deterministic read fixture")); + for offset in [0, 17] { + let bytes = tokio::time::timeout(Duration::from_secs(5), driver.read_at(Arc::clone(&file), offset, 32)) + .await + .expect("ordinary fixture read should complete") + .expect("best-effort io-wq registration must not disable working reads"); + assert_eq!(bytes, [0; 32]); + } + + let after = driver.shard_iowq_setup(); + assert_eq!(after.len(), 2); + for (before, after) in before.iter().zip(&after) { + assert_requested_limits(after); + match (&before.result, &after.result) { + ( + IoWqRegistration::Registered { previous_limits: before }, + IoWqRegistration::Registered { previous_limits: after }, + ) => assert_eq!(before, after, "previous-limit startup snapshot must remain historical"), + ( + IoWqRegistration::Failed { + error_kind: before_kind, + raw_os_error: before_errno, + }, + IoWqRegistration::Failed { + error_kind: after_kind, + raw_os_error: after_errno, + }, + ) => { + assert_eq!(before_kind, after_kind); + assert_eq!(before_errno, after_errno); + } + _ => panic!("startup registration outcome changed after reads"), + } + } + let stats = driver.shutdown(); + assert_eq!(stats.submitted, 2); + assert_eq!(stats.delivered, 2); + assert_eq!(stats.orphan_reclaimed, 0); + assert_eq!(stats.in_flight, 0); +} From 2b793582f6ff2e3eadeeffb33e0bba72aab87217 Mon Sep 17 00:00:00 2001 From: Hauser Date: Wed, 23 Sep 2026 10:07:34 +0800 Subject: [PATCH 2/2] docs(driver): record io-wq native verification Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- docs/optimization-status.md | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/docs/optimization-status.md b/docs/optimization-status.md index 78bb95c..1c55cfc 100644 --- a/docs/optimization-status.md +++ b/docs/optimization-status.md @@ -14,7 +14,7 @@ integration are separate gates; a checked code item does not close the roadmap. | 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; #16 merged | Application dependency/fallback/result ownership wiring | -| 4.2 io-wq registration evidence | Per-shard startup query, deterministic and native tests implemented; two independent reviews passed | New-head native CI; actual worker count and process-wide budget remain unproven | +| 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 | @@ -42,9 +42,14 @@ final-head CI and merge status separately. 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; a -new-head native CI result remains pending; two independent reviews passed. The snapshot -does not measure effective worker counts or establish a process-wide cap. +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