diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index bfb2907..bb95bb6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -41,24 +41,35 @@ jobs: run: cargo fmt --check - name: clippy run: cargo clippy --all-targets --all-features -- -D warnings + - name: ABBA evidence gate tests + run: python3 -m unittest discover -s scripts -p 'test_bench_abba.py' -v test: name: test (real io_uring on runner) runs-on: ubuntu-latest + env: + URING_REQUIRE_DIRECT: "1" steps: - uses: actions/checkout@v7 - name: Install toolchain run: rustup toolchain install stable --profile minimal && rustup default stable - name: Build tests - run: cargo build --tests --all-features + run: cargo build --locked --tests --all-features - name: Test and assert io_uring actually ran run: | set -o pipefail - cargo test --all-features -- --nocapture --test-threads=1 2>&1 | tee test.out + cargo test --locked --all-features -- --nocapture --test-threads=1 2>&1 | tee test.out if grep -q "SKIP " test.out; then echo "::error::a test skipped — io_uring did not run on this runner (vacuous pass)" exit 1 fi + grep -q "DIRECT_OK direct_read_returns_exact_unaligned_ranges" test.out + - name: Benchmark schema and correctness smoke + run: bash scripts/test-benchmark-cli.sh + - name: Instrumented benchmark schema and correctness smoke + env: + BENCH_DIAGNOSTICS: "1" + run: bash scripts/test-benchmark-cli.sh docker-two-legs: name: docker (degradation + real io_uring) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6605b6d..5035bd3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -23,6 +23,24 @@ aims to follow [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## [Unreleased] +### Added + +- Opt-in `diagnostics` feature with per-shard sampled driver-stage histograms, + aggregate snapshots, and measurement-interval deltas. Default builds compile + out the timing fields and sampling work. +- Native O_DIRECT execution gate and byte-exact benchmark CLI smoke coverage. +- Warm-cache ABBA evidence runner with per-leg provenance/resource records, + drift and environment gates, and bounded benchmark process-group cleanup. + +### Changed + +- Benchmark CSV schema v2 obtains headers from the executable and reports + independent setup, workload, and teardown timings, configurable ring depth, + workers, warmup, and instrumentation status. Positional CSV consumers must + migrate to the new header. +- Clarified that driver submission counters count accepted logical reads and + delivery counters count successful channel sends, not caller consumption. + ## [0.2.2] - 2026-09-07 This patch release documents and hardens the public read API while preserving diff --git a/Cargo.toml b/Cargo.toml index b674f24..e91d8ca 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -33,6 +33,9 @@ categories = ["asynchronous", "filesystem"] # probe-drain leak paths deterministically. All gated code is `#[cfg(feature = # "fault-injection")]`, so a default build compiles none of it. fault-injection = [] +# Opt-in sampled timing. No timestamps, histogram storage, or tracing allocations +# are compiled into the default driver. +diagnostics = [] [dependencies] # tracing for driver diagnostics: unifies with the RustFS tracing pipeline for diff --git a/README.md b/README.md index a583041..d30fc77 100644 --- a/README.md +++ b/README.md @@ -78,8 +78,16 @@ The public API is intentionally small and read-only: `submitted == delivered + orphan_reclaimed` holds after all completions have been reaped; `in_flight == 0` indicates a clean shutdown. +The optional `diagnostics` feature exposes sampled stage histograms through +`UringDriver::diagnostics()` and `shard_diagnostics()`. It is off by default. +See [the measurement guide](docs/benchmarking.md#sampled-diagnostics) for sampling, +stage overlap, cancellation, and instrumentation-overhead boundaries. + ## Testing +Benchmark configuration, CSV schema, timing boundaries, and performance gates +are documented in [the benchmarking guide](docs/benchmarking.md). + Linux only; on other hosts `cargo check` builds the empty stub. ```bash diff --git a/bench-concurrent-pread.sh b/bench-concurrent-pread.sh index 355bd52..b75db48 100755 --- a/bench-concurrent-pread.sh +++ b/bench-concurrent-pread.sh @@ -19,7 +19,7 @@ # Sweeps strategy × read_size × concurrency. Caches are dropped before every # timed run so reads hit the device (a large file + random offsets means the # page cache would otherwise skew results run-to-run). Needs root for -# drop_caches; the bench host azure-20780104 runs as root. +# drop_caches; run cold sweeps only on an isolated test host. set -euo pipefail cd "$(dirname "$0")" @@ -28,13 +28,14 @@ REPEAT="${REPEAT:-3}" BIN="${CARGO_TARGET_DIR:-target}/release/examples/concurrent_pread_bench" FILE_SIZE="${FILE_SIZE:-4294967296}" # 4 GiB, >> page cache reuse for random reads -READ_SIZES=(${READ_SIZES:-65536 1048576}) -CONCURRENCIES=(${CONCURRENCIES:-1 8 32 128}) +read -r -a READ_SIZES <<<"${READ_SIZES:-65536 1048576}" +read -r -a CONCURRENCIES <<<"${CONCURRENCIES:-1 8 32 128}" +read -r -a SHARD_COUNTS <<<"${SHARD_COUNTS:-1}" # Bound each cold run's transfer instead of fixing the op count: a cold 1 MiB # read costs ~16x a 64 KiB one, so a fixed op count would make the large-read # legs dominate wall-clock for no extra signal. TOTAL_BYTES="${TOTAL_BYTES:-268435456}" # 256 MiB per run -MIN_OPS="${MIN_OPS:-512}" # enough samples for a p999 +MIN_OPS="${MIN_OPS:-512}" # smoke-sized; not a reliable p999 population MAX_OPS="${MAX_OPS:-4096}" STRATS=(std_open_pread std_cached_pread uring_open_read uring_cached_read) @@ -46,7 +47,9 @@ ops_for() { # read_size -> op count, clamped } mkdir -p "$DIR" -cargo build --release --example concurrent_pread_bench >&2 +build_features=() +if [[ "${BENCH_DIAGNOSTICS:-0}" == 1 ]]; then build_features=(--features diagnostics); fi +cargo build --locked --release --example concurrent_pread_bench "${build_features[@]}" >&2 # Correctness preflight (untimed; output discarded). IOPS cannot distinguish a # strategy that reads the right *number* of bytes from one that reads the wrong @@ -54,11 +57,9 @@ cargo build --release --example concurrent_pread_bench >&2 # which checks each delivered byte against the file's offset-addressable pattern. # A mismatch aborts before any measurement is taken. preflight_verify() { - local vdir="$DIR/verify.$$" vfile strat - rm -rf "$vdir" - mkdir -p "$vdir" - # shellcheck disable=SC2064 - trap "rm -rf '$vdir'" RETURN + local vdir vfile strat + vdir=$(mktemp -d "$DIR/verify.XXXXXX") + trap 'rm -rf -- "$vdir"; trap - RETURN' RETURN vfile="$vdir/verify.bin" for strat in "${STRATS[@]}"; do # Unaligned read size on purpose: exercises the offset bookkeeping. @@ -72,24 +73,38 @@ FILE="$DIR/pread_${FILE_SIZE}.bin" # Create once, untimed, so every cold run below is genuinely cold. "$BIN" std_cached_pread "$FILE" "$FILE_SIZE" 65536 1 1 >/dev/null -echo "cache,strategy,file_size,read_size,concurrency,ops,secs,IOPS,MBps,p50_us,p99_us,p999_us" +HEADER=$("$BIN" --header) +FIELDS=$(awk -F, '{print NF}' <<<"$HEADER") +printf 'cache,%s\n' "$HEADER" for cache in ${CACHES:-cold warm}; do + if [[ "$cache" == cold && "${BENCH_WARMUP_OPS:-0}" != 0 ]]; then + echo "cold sweeps require BENCH_WARMUP_OPS=0" >&2 + exit 1 + fi for read_size in "${READ_SIZES[@]}"; do ops=$(ops_for "$read_size") for conc in "${CONCURRENCIES[@]}"; do for strat in "${STRATS[@]}"; do - for _ in $(seq 1 "$REPEAT"); do - if [ "$cache" = cold ]; then - sync - echo 3 >/proc/sys/vm/drop_caches - else - # Warm takes the device out of the picture, isolating the - # software cost (open, blocking-pool hop, submission). - # On a throughput-throttled disk the cold leg saturates - # and hides exactly the overhead we are pricing. - cat "$FILE" >/dev/null - fi - echo "$cache,$("$BIN" "$strat" "$FILE" "$FILE_SIZE" "$read_size" "$conc" "$ops")" + shards_for_strategy=(1) + if [[ "$strat" == uring_* ]]; then + shards_for_strategy=("${SHARD_COUNTS[@]}") + fi + for shards in "${shards_for_strategy[@]}"; do + for _ in $(seq 1 "$REPEAT"); do + if [ "$cache" = cold ]; then + sync + echo 3 >/proc/sys/vm/drop_caches + else + # Warm takes the device out of the picture, isolating the + # software cost (open, blocking-pool hop, submission). + # On a throughput-throttled disk the cold leg saturates + # and hides exactly the overhead we are pricing. + cat "$FILE" >/dev/null + fi + row=$("$BIN" "$strat" "$FILE" "$FILE_SIZE" "$read_size" "$conc" "$ops" "$shards") + awk -F, -v expected="$FIELDS" 'NF != expected || $1 != 2 || $2 != "measure" {exit 1}' <<<"$row" + printf '%s,%s\n' "$cache" "$row" + done done done done diff --git a/bench-streaming.sh b/bench-streaming.sh index bd4658a..5650e29 100755 --- a/bench-streaming.sh +++ b/bench-streaming.sh @@ -22,7 +22,7 @@ # cold — `drop_caches` before each timed run, so every read hits the device. # # Emits one CSV to stdout. Cold mode needs root (drop_caches); the bench host -# azure-20780104 runs as root. REPEAT>1 runs each config N times so the caller +# must be isolated from other workloads. REPEAT>1 repeats each config so the caller # can take the median. # # BENCH_DIR only ever holds files this script created: the example refuses to @@ -39,12 +39,14 @@ BIN="${CARGO_TARGET_DIR:-target}/release/examples/streaming_bench" STRATEGIES=(std_buffered std_odirect uring_read_at uring_read_at_direct) # Sizes: metadata-ish, mid object, large object. -SIZES=(${SIZES:-65536 16777216 268435456}) -CHUNKS=(${CHUNKS:-131072 1048576}) -QDS=(${QDS:-1 4 16}) +read -r -a SIZES <<<"${SIZES:-65536 16777216 268435456}" +read -r -a CHUNKS <<<"${CHUNKS:-131072 1048576}" +read -r -a QDS <<<"${QDS:-1 4 16}" mkdir -p "$DIR" -cargo build --release --example streaming_bench >&2 +build_features=() +if [[ "${BENCH_DIAGNOSTICS:-0}" == 1 ]]; then build_features=(--features diagnostics); fi +cargo build --locked --release --example streaming_bench "${build_features[@]}" >&2 # Correctness preflight (untimed; its output is discarded). Throughput cannot # distinguish a strategy that reads the right *number* of bytes from one that @@ -56,11 +58,9 @@ cargo build --release --example streaming_bench >&2 # - queue depths 1 and 4, covering the pipelined path's offset bookkeeping. # A mismatch aborts the script before any measurement is taken. preflight_verify() { - local vdir="$DIR/verify.$$" chunk=131072 size strat qd - rm -rf "$vdir" - mkdir -p "$vdir" - # shellcheck disable=SC2064 - trap "rm -rf '$vdir'" RETURN + local vdir chunk=131072 size strat qd + vdir=$(mktemp -d "$DIR/verify.XXXXXX") + trap 'rm -rf -- "$vdir"; trap - RETURN' RETURN for size in $((1048576 + 4096)) $((1048576 + 4097)); do for strat in "${STRATEGIES[@]}"; do for qd in 1 4; do @@ -94,12 +94,21 @@ run() { # strategy size chunk qd cache local file="$DIR/bench_${size}.bin" for _ in $(seq 1 "$REPEAT"); do drop_or_warm "$cache" "$size" "$chunk" - echo "$cache,$("$BIN" "$strat" "$file" "$size" "$chunk" "$qd" "$ALIGN")" + local row + row=$("$BIN" "$strat" "$file" "$size" "$chunk" "$qd" "$ALIGN") + awk -F, -v expected="$FIELDS" 'NF != expected || $1 != 2 || $2 != "measure" {exit 1}' <<<"$row" + printf '%s,%s\n' "$cache" "$row" done } -echo "cache,strategy,size,chunk,qd,align,bytes,secs,MBps,ops" +HEADER=$("$BIN" --header) +FIELDS=$(awk -F, '{print NF}' <<<"$HEADER") +printf 'cache,%s\n' "$HEADER" for cache in warm cold; do + if [[ "$cache" == cold && "${BENCH_WARMUP_RUNS:-0}" != 0 ]]; then + echo "cold sweeps require BENCH_WARMUP_RUNS=0" >&2 + exit 1 + fi for size in "${SIZES[@]}"; do for chunk in "${CHUNKS[@]}"; do run std_buffered "$size" "$chunk" 1 "$cache" diff --git a/docs/benchmarking.md b/docs/benchmarking.md new file mode 100644 index 0000000..6373202 --- /dev/null +++ b/docs/benchmarking.md @@ -0,0 +1,185 @@ +# Read-backend benchmarks + +The examples measure I/O mechanisms, not a complete RustFS `LocalIoBackend` or +S3 request. In particular, the std pread strategies do not model every mmap, +metadata, cache, or reclaim policy used by an application. The streaming example +measures a task-per-chunk completion pipeline, not ordered delivery to a slow +consumer. Compare application backends separately before making rollout choices. + +## Output contract + +Each example prints its CSV header with `--header`. Sweep scripts obtain that +header from the executable rather than maintaining a second copy. Schema version +2 adds `schema_version`, `mode`, explicit setup/teardown times, and configuration +fields. Existing positional CSV consumers must migrate; do not parse v2 rows +using the previous header. `mode=verify` is a correctness run and must never be +included in performance results. Sweeps reject verification rows. + +`secs` measures the workload after preparation and optional warmup, before +driver/runtime teardown. `startup_secs` includes runtime/driver creation, probe, +opening reusable descriptors, generating offsets, and allocating reusable std +buffers. Dataset creation happens before all reported intervals. +`shutdown_secs` includes dropping these prepared resources and the runtime. +Warmup time is excluded from all three fields. + +Task creation/joining, per-operation allocations, error/length checks, result +collection, and buffer disposal remain inside the workload interval. Reported +operation latencies end at the result, before content verification or buffer +drop; total throughput measures the whole loop. Verification changes workload +cost even though its byte checks happen outside an individual latency sample. + +The std streaming strategies reuse a single buffer. io_uring allocates one per +read. These results compare the current implementations including that +difference; they do not isolate syscall overhead or prove buffer-pool benefits. + +## Configuration + +| Variable | Default | Meaning | +| --- | --- | --- | +| `BENCH_WORKERS` | Available parallelism | Tokio workers, 1–1024 | +| `BENCH_RING_ENTRIES` | 128 | Power-of-two SQ depth per shard, 1–32768 | +| `BENCH_WARMUP_OPS` | 0 | Concurrent benchmark operations before measurement | +| `BENCH_WARMUP_RUNS` | 0 | Full streaming passes before measurement | +| `BENCH_VERIFY` | Off | `1` or `true` enables byte-exact verification | +| `BENCH_DIAGNOSTICS` | Off | `1` reports sampled driver-stage histograms to stderr; requires the `diagnostics` feature | +| `SHARD_COUNTS` | 1 | Space-separated shard counts for the concurrent sweep | + +Ring depth is independent of concurrency so saturation can be exercised. +Concurrent positional arguments still accept an optional shard count (1–64). +Std rows report zero ring entries/shards; std streaming also reports zero Tokio +workers because it uses no runtime. `workers` in concurrent std rows remains +nonzero because the tasks dispatch work through Tokio's blocking pool. + +Warmup defaults to zero to avoid relabeling warmed data as cold. Cold sweep +legs require zero warmup. A cache drop describes the initial state only: reads +can repopulate the cache during a run. A warm preload alone does not prove that +the whole working set remains resident. Never clear global page caches on a +shared or production host. + +## Correctness and smoke validation + +On Linux with real io_uring and an O_DIRECT-capable test filesystem: + +```sh +bash scripts/test-benchmark-cli.sh +URING_REQUIRE_DIRECT=1 cargo test --locked --test cancel \ + direct_read_returns_exact_unaligned_ranges -- --nocapture --test-threads=1 +``` + +The smoke checks all eight strategies, schema/field alignment, warmup counts, +small-ring saturation, invalid configuration, and byte-exact unaligned ranges. +It is intentionally small and is not a performance result. The direct test +prints `DIRECT_OK direct_read_returns_exact_unaligned_ranges` only after its +native direct assertions execute. Without `URING_REQUIRE_DIRECT=1`, unsupported +filesystems may still omit the direct assertions; a passing general test suite +alone is not proof of O_DIRECT coverage. + +## Performance acceptance + +### Sampled diagnostics + +Build with `--features diagnostics` to collect one sample per 64 handle +constructions **per shard**, starting with each shard's first handle. Sampling +has a separate sequence per shard so round-robin selection cannot bias every +sample toward shard zero. Invalid requests can consume a sampling position +without recording stages. Deterministic sampling is diagnostic, not an unbiased +estimate for every possible periodic workload. + +`UringDriver::diagnostics()` aggregates shards; +`UringDriver::shard_diagnostics()` preserves shard identity. Each stage has a +count, total nanoseconds, and 64 log2 nanosecond buckets. Bucket zero covers +0–1 ns; bucket i > 0 covers `[2^i, 2^(i+1))`. Counts and sums wrap modulo 2^64. +Snapshots are approximate during concurrent updates; take quiescent snapshots +and use `since()` on the same driver to exclude warmup without resetting it. + +| Stage | Boundary | +| --- | --- | +| `admission` | Handle timing starts to permit acquisition; includes an unpolled saturated handle's inactivity | +| `driver_queue` | Immediately before send through driver intake | +| `preparation` | Intake through initial local SQE backlog insertion | +| `driver_lifetime` | Intake through final read CQE reap; includes preparation, retries and reaper delay | +| `cqe_processing` | Each read CQE's handling and range adjustment, before removal/send | +| `completion_to_poll` | Immediately before result send through caller ready poll, including receiver inactivity | + +These stages overlap; do not sum them or label lifetime as disk latency or +completion-to-poll as Tokio schedule latency. Short reads can generate several +CQE samples. Cancel CQEs themselves are excluded. Abandoned receivers do not +create completion-to-poll samples. Rejection, driver disappearance, and the +bounded-drain leak path can leave some stages unrecorded; do not infer a terminal +CQE from an error delivered by shutdown. + +`BENCH_DIAGNOSTICS=1` makes the examples report **measurement-interval deltas** +to stderr, outside their measured interval. Sweep scripts build with the feature +automatically when this variable is 1. In a manually built binary the variable +does not enable/disable instrumentation; it controls reporting only. CSV +`diagnostics_interval` is 64 whenever the feature was compiled, otherwise zero, +even on std strategies (which do not use the instrumented driver). + +The default build compiles out timing fields, clock reads, histogram storage and +sample allocations. Enabled builds add a per-shard atomic sampling counter per +handle and an Arc/timestamps/histogram updates for sampled operations. Measure +that overhead on target hardware with separate feature-off/feature-on artifacts; +do not assume it is free. Runtime schedule-latency and blocking-pool metrics must +still be correlated in the application, which owns the Tokio runtime. + +### Target-hardware acceptance + +Use release builds and record source revisions, runtime configuration, filesystem, +cache policy, resource allocation, and actual backend execution. Keep the total +operation/byte budget constant when comparing shard counts. Use separate startup +and steady-state results. Do not compare old whole-process timings to v2 `secs`. + +The existing sweeps are exploratory parameter sweeps, not an ABBA acceptance +harness. Run A1/B1/B2/A2 on isolated target hardware, reject baseline drift, and +measure throughput, CPU/op, tail latency, RSS, fallback/errors, and scheduler +signals. The default 512–4096 operations are smoke-sized: p99.9 is not robust +with that sample count. Increase the population and report uncertainty. Closed +loop concurrency also hides overload queueing; evaluate a controlled arrival +rate separately when studying tail latency. + +### ABBA driver runner + +`scripts/bench-abba.py` (Python 3.11+, Linux, GNU time, taskset) runs three or +more A1/B1/B2/A2 rounds against a **pre-created, byte-verified** warm-cache +dataset. Build two release artifacts from the same source with diagnostics off +and on, using separate target directories. Supply explicit binaries, source +revision, CPU affinity, and a reservation description. The tool never stops +services, changes their configuration, clears global caches, or overwrites an +existing result directory. + +```sh +python3 scripts/bench-abba.py \ + --baseline /test/target-off/release/examples/concurrent_pread_bench \ + --candidate /test/target-on/release/examples/concurrent_pread_bench \ + --data-file /test/verified-data.bin --run-dir /test/results/new-run \ + --source-revision COMMIT --reservation-note 'reserved benchmark window' \ + --cpus 0-7 --workers 4 --shards 2 --entries 64 \ + --ops 1000000 --warmup-ops 10000 --rounds 3 --dry-run +``` + +Remove `--dry-run` only after reserving resources. If a CI runner or another +service can introduce load, supply `--require-inactive-unit UNIT`; the runner +refuses to start or continue unless that unit is inactive. Any service stop or +restore is an explicit operator action outside this tool. A reservation note and +process checks are evidence aids, not a substitute for exclusive resources. +Perform A/A calibration first by passing the baseline binary in both positions +and `--candidate-interval 0`; then use the diagnostics candidate and the default +interval of 64. Keep the workload and thresholds fixed between experiments. + +Every leg must run at least five measured seconds by default. If it is too +short, increase operations in a **new** experiment. The tool validates geometry, +feature state, sample count, schema, finite values, derived throughput, process +resource reports, binary identity and dataset metadata. Active build/load/CI +workers, errors, per-leg deadlines, or failed baseline drift invalidate the run. +It stops after the first invalid round rather than expanding the matrix. + +Output contains provenance, each leg's CSV/stderr and whole-process CPU/RSS/ +context-switch report, parsed leg JSON, and a summary. `valid-comparison` means +the evidence passed these gates, **not** that the candidate improved, passed an +overhead budget, or proved a RustFS/S3 benefit. Review signed candidate changes +and resource reports against the application SLO separately. Duration histograms +with too few sampled operations are not reliable tail estimates. + +Runner gate tests: `python3 -m unittest discover -s scripts -p 'test_bench_abba.py'`. + +Tracking and implementation status: [rustfs/backlog#2647](https://github.com/rustfs/backlog/issues/2647). diff --git a/examples/common/mod.rs b/examples/common/mod.rs new file mode 100644 index 0000000..ecc98d4 --- /dev/null +++ b/examples/common/mod.rs @@ -0,0 +1,66 @@ +// Copyright 2024 RustFS Team +// SPDX-License-Identifier: Apache-2.0 + +use std::env; + +pub fn diagnostics_interval() -> u64 { + #[cfg(feature = "diagnostics")] + { + rustfs_uring::DIAGNOSTICS_SAMPLE_INTERVAL + } + #[cfg(not(feature = "diagnostics"))] + { + 0 + } +} + +pub fn setting(name: &str, default: usize, min: usize, max: usize) -> Result { + let value = match env::var(name) { + Ok(raw) => raw.parse::().map_err(|_| format!("{name} must be an integer"))?, + Err(env::VarError::NotPresent) => default, + Err(err) => return Err(format!("{name}: {err}")), + }; + if !(min..=max).contains(&value) { + return Err(format!("{name} must be in {min}..={max}")); + } + Ok(value) +} + +pub fn workers() -> Result { + if env::var("BENCH_DIAGNOSTICS").as_deref() == Ok("1") && !cfg!(feature = "diagnostics") { + return Err("BENCH_DIAGNOSTICS=1 requires building with --features diagnostics".into()); + } + let default = std::thread::available_parallelism().map(usize::from).unwrap_or(1); + setting("BENCH_WORKERS", default, 1, 1024) +} + +#[cfg(feature = "diagnostics")] +pub fn report_diagnostics(snapshot: &rustfs_uring::DiagnosticsSnapshot) { + if env::var("BENCH_DIAGNOSTICS").as_deref() != Ok("1") { + return; + } + for (stage, histogram) in [ + ("admission", &snapshot.admission), + ("driver_queue", &snapshot.driver_queue), + ("preparation", &snapshot.preparation), + ("driver_lifetime", &snapshot.driver_lifetime), + ("cqe_processing", &snapshot.cqe_processing), + ("completion_to_poll", &snapshot.completion_to_poll), + ] { + eprintln!( + "DIAGNOSTICS stage={stage} interval={} count={} total_ns={} buckets={:?}", + rustfs_uring::DIAGNOSTICS_SAMPLE_INTERVAL, + histogram.count, + histogram.total_nanos, + histogram.buckets, + ); + } +} + +pub fn ring_entries(default: usize) -> Result { + let entries = setting("BENCH_RING_ENTRIES", default, 1, 32768)?; + if !entries.is_power_of_two() { + return Err("BENCH_RING_ENTRIES must be a power of two".into()); + } + Ok(entries as u32) +} diff --git a/examples/concurrent_pread_bench.rs b/examples/concurrent_pread_bench.rs index 6ca009e..f5d6697 100644 --- a/examples/concurrent_pread_bench.rs +++ b/examples/concurrent_pread_bench.rs @@ -12,6 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. +#[cfg(target_os = "linux")] +mod common; + #[cfg(target_os = "linux")] #[path = "concurrent_pread_bench/linux.rs"] mod linux; diff --git a/examples/concurrent_pread_bench/linux.rs b/examples/concurrent_pread_bench/linux.rs index 7c3a00b..35d3c59 100644 --- a/examples/concurrent_pread_bench/linux.rs +++ b/examples/concurrent_pread_bench/linux.rs @@ -57,8 +57,9 @@ use std::time::{Duration, Instant}; use rustfs_uring::UringDriver; -/// Keeps `(concurrency * 2).next_power_of_two()` ring entries under the -/// kernel's 32768-entry limit and the arithmetic far from overflow. +const CSV_HEADER: &str = "schema_version,mode,strategy,shards,file_size,read_size,concurrency,ops,secs,IOPS,MBps,p50_us,p99_us,p999_us,startup_secs,shutdown_secs,workers,ring_entries,warmup_ops,diagnostics_interval"; + +/// Bound task fan-out independently of the configurable per-shard ring depth. const MAX_CONCURRENCY: usize = 4096; const MAX_READ_SIZE: usize = 1 << 30; @@ -112,6 +113,9 @@ struct Config { total_ops: usize, shards: usize, verify: bool, + workers: usize, + ring_entries: u32, + warmup_ops: usize, } fn parse_args() -> Result { @@ -135,6 +139,9 @@ fn parse_args() -> Result { total_ops: parse("total_ops", &a[6])? as usize, shards: if a.len() == 8 { parse("shards", &a[7])? as usize } else { 1 }, verify: matches!(std::env::var("BENCH_VERIFY").as_deref(), Ok("1") | Ok("true")), + workers: crate::common::workers()?, + ring_entries: crate::common::ring_entries(128)?, + warmup_ops: crate::common::setting("BENCH_WARMUP_OPS", 0, 0, 10_000_000)?, }; validate(&cfg)?; Ok(cfg) @@ -162,6 +169,9 @@ fn validate(cfg: &Config) -> Result<(), String> { if cfg.total_ops == 0 { return Err("total_ops must be > 0".to_string()); } + if cfg.shards == 0 || cfg.shards > 64 { + return Err("shards must be in 1..=64".to_string()); + } Ok(()) } @@ -261,10 +271,10 @@ fn open_read(path: &str) -> io::Result { /// Deterministic block-aligned offsets, so every strategy reads the exact same /// blocks and the comparison is not confounded by a different access pattern. -fn offsets(cfg: &Config) -> Vec { +fn offsets(cfg: &Config, ops: usize) -> Vec { let blocks = (cfg.file_size - cfg.read_size as u64) / 4096; let mut state = 0x2545_f491_4f6c_dd1du64; - (0..cfg.total_ops) + (0..ops) .map(|_| { state = state.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407); ((state >> 33) % blocks) * 4096 @@ -280,31 +290,51 @@ fn percentile(sorted: &[Duration], p: f64) -> u128 { sorted[idx].as_micros() } -async fn run(cfg: &Config) -> Result, String> { - let offs = Arc::new(offsets(cfg)); - let path = Arc::new(cfg.file.clone()); +struct Prepared { + offsets: Arc>, + warmup_offsets: Arc>, + path: Arc, + cached: Option>, + driver: Option>, +} - // A shared fd is safe for the cached strategies: pread is positional and - // never touches the file description's shared offset. - let cached: Option> = if cfg.strategy.uses_cached_fd() { - Some(Arc::new(open_read(&cfg.file).map_err(|e| format!("open cached fd: {e}"))?)) - } else { - None - }; - let driver = if cfg.strategy.uses_uring() { - let depth = (cfg.concurrency.max(64) * 2).next_power_of_two() as u32; - Some(Arc::new( - UringDriver::probe_and_start_sharded(depth, cfg.shards).map_err(|e| format!("probe io_uring: {e:?}"))?, - )) - } else { - None - }; +impl Prepared { + fn new(cfg: &Config) -> Result { + // Positional reads can share an fd without changing its current offset. + let cached = if cfg.strategy.uses_cached_fd() { + Some(Arc::new(open_read(&cfg.file).map_err(|e| format!("open cached fd: {e}"))?)) + } else { + None + }; + let driver = if cfg.strategy.uses_uring() { + Some(Arc::new( + UringDriver::probe_and_start_sharded(cfg.ring_entries, cfg.shards) + .map_err(|e| format!("probe io_uring: {e:?}"))?, + )) + } else { + None + }; + Ok(Self { + offsets: Arc::new(offsets(cfg, cfg.total_ops)), + warmup_offsets: Arc::new(offsets(cfg, cfg.warmup_ops)), + path: Arc::new(cfg.file.clone()), + cached, + driver, + }) + } +} +async fn run(cfg: &Config, prepared: &Prepared, offs: &Arc>) -> Result, String> { let (strategy, read_size, verify) = (cfg.strategy, cfg.read_size, cfg.verify); - let per_task = cfg.total_ops.div_ceil(cfg.concurrency); + let per_task = offs.len().div_ceil(cfg.concurrency); let mut set = tokio::task::JoinSet::new(); for t in 0..cfg.concurrency { - let (offs, path, cached, driver) = (offs.clone(), path.clone(), cached.clone(), driver.clone()); + let (offs, path, cached, driver) = ( + Arc::clone(offs), + Arc::clone(&prepared.path), + prepared.cached.clone(), + prepared.driver.clone(), + ); let start_idx = (t * per_task).min(offs.len()); let end_idx = (start_idx + per_task).min(offs.len()); set.spawn(async move { @@ -312,7 +342,7 @@ async fn run(cfg: &Config) -> Result, String> { for &off in &offs[start_idx..end_idx] { let t0 = Instant::now(); let bytes: Vec = match strategy { - // Today's StdBackend: a blocking-pool hop that opens then preads. + // Mechanism baseline: one blocking hop for open plus pread. Strategy::StdOpenPread => { let p = path.clone(); tokio::task::spawn_blocking(move || { @@ -337,7 +367,7 @@ async fn run(cfg: &Config) -> Result, String> { .map_err(|e| format!("join: {e}"))? .map_err(|e| format!("pread: {e}"))? } - // Today's UringBackend: still a blocking-pool hop for open+stat. + // Cache-miss shape: blocking open+stat, followed by io_uring. Strategy::UringOpenRead => { let p = path.clone(); let file = tokio::task::spawn_blocking(move || { @@ -373,7 +403,7 @@ async fn run(cfg: &Config) -> Result, String> { }); } - let mut all = Vec::with_capacity(cfg.total_ops); + let mut all = Vec::with_capacity(offs.len()); while let Some(r) = set.join_next().await { all.extend(r.map_err(|e| format!("task panicked: {e}"))??); } @@ -381,6 +411,10 @@ async fn run(cfg: &Config) -> Result, String> { } pub(super) fn main() -> ExitCode { + if std::env::args().nth(1).as_deref() == Some("--header") { + println!("{CSV_HEADER}"); + return ExitCode::SUCCESS; + } let cfg = match parse_args() { Ok(cfg) => cfg, Err(e) => { @@ -393,7 +427,12 @@ pub(super) fn main() -> ExitCode { return ExitCode::FAILURE; } - let rt = match tokio::runtime::Builder::new_multi_thread().enable_all().build() { + let startup = Instant::now(); + let rt = match tokio::runtime::Builder::new_multi_thread() + .worker_threads(cfg.workers) + .enable_all() + .build() + { Ok(rt) => rt, Err(e) => { eprintln!("error: build runtime: {e}"); @@ -401,8 +440,25 @@ pub(super) fn main() -> ExitCode { } }; + let prepared = match Prepared::new(&cfg) { + Ok(prepared) => prepared, + Err(e) => { + eprintln!("error: {e}"); + return ExitCode::FAILURE; + } + }; + let startup_secs = startup.elapsed().as_secs_f64(); + if cfg.warmup_ops != 0 + && let Err(e) = rt.block_on(run(&cfg, &prepared, &prepared.warmup_offsets)) + { + eprintln!("error: warmup: {e}"); + return ExitCode::FAILURE; + } + + #[cfg(feature = "diagnostics")] + let before = prepared.driver.as_ref().map(|driver| driver.diagnostics()); let start = Instant::now(); - let mut lats = match rt.block_on(run(&cfg)) { + let mut lats = match rt.block_on(run(&cfg, &prepared, &prepared.offsets)) { Ok(l) => l, Err(e) => { eprintln!("error: {e}"); @@ -410,6 +466,16 @@ pub(super) fn main() -> ExitCode { } }; let secs = start.elapsed().as_secs_f64(); + #[cfg(feature = "diagnostics")] + if let (Some(driver), Some(before)) = (&prepared.driver, &before) { + crate::common::report_diagnostics(&driver.diagnostics().since(before)); + } + // All operation tasks have joined. Drop the driver and runtime on this + // synchronous thread, outside the steady-state measurement. + let shutdown = Instant::now(); + drop(prepared); + drop(rt); + let shutdown_secs = shutdown.elapsed().as_secs_f64(); if cfg.verify { eprintln!("{}: verified byte-exact across {} reads", cfg.strategy.name(), lats.len()); @@ -419,11 +485,11 @@ pub(super) fn main() -> ExitCode { let ops = lats.len(); let iops = ops as f64 / secs; let mbps = (ops as f64 * cfg.read_size as f64 / (1024.0 * 1024.0)) / secs; - // CSV: strategy,shards,file_size,read_size,concurrency,ops,secs,IOPS,MBps,p50_us,p99_us,p999_us println!( - "{},{},{},{},{},{},{:.6},{:.0},{:.1},{},{},{}", + "2,{},{},{},{},{},{},{},{:.6},{:.0},{:.1},{},{},{},{:.6},{:.6},{},{},{},{}", + if cfg.verify { "verify" } else { "measure" }, cfg.strategy.name(), - cfg.shards, + if cfg.strategy.uses_uring() { cfg.shards } else { 0 }, cfg.file_size, cfg.read_size, cfg.concurrency, @@ -434,6 +500,12 @@ pub(super) fn main() -> ExitCode { percentile(&lats, 0.50), percentile(&lats, 0.99), percentile(&lats, 0.999), + startup_secs, + shutdown_secs, + cfg.workers, + if cfg.strategy.uses_uring() { cfg.ring_entries } else { 0 }, + cfg.warmup_ops, + crate::common::diagnostics_interval(), ); ExitCode::SUCCESS } diff --git a/examples/streaming_bench.rs b/examples/streaming_bench.rs index 903239a..25d6986 100644 --- a/examples/streaming_bench.rs +++ b/examples/streaming_bench.rs @@ -12,6 +12,9 @@ // See the License for the specific language governing permissions and // limitations under the License. +#[cfg(target_os = "linux")] +mod common; + #[cfg(target_os = "linux")] #[path = "streaming_bench/linux.rs"] mod linux; diff --git a/examples/streaming_bench/linux.rs b/examples/streaming_bench/linux.rs index f82d52b..6f09991 100644 --- a/examples/streaming_bench/linux.rs +++ b/examples/streaming_bench/linux.rs @@ -45,7 +45,7 @@ //! not a measurement, and its CSV row must be discarded. use std::fs::{File, OpenOptions}; -use std::io::{self, Read, Write}; +use std::io::{self, Read, Seek, Write}; use std::os::unix::fs::{FileExt, OpenOptionsExt}; use std::process::ExitCode; use std::sync::Arc; @@ -53,9 +53,9 @@ use std::time::Instant; use rustfs_uring::UringDriver; -/// Queue depth ceiling. `probe_and_start` gets `(qd * 2).next_power_of_two()` -/// entries, so this keeps the ring under the kernel's 32768-entry limit and the -/// arithmetic far from overflow. +const CSV_HEADER: &str = "schema_version,mode,strategy,size,chunk,qd,align,bytes,secs,MBps,ops,startup_secs,shutdown_secs,workers,ring_entries,warmup_runs,diagnostics_interval"; + +/// Bound task fan-out independently of the configurable ring depth. const MAX_QD: usize = 4096; /// Smallest logical block size any device reports. const MIN_ALIGN: usize = 512; @@ -107,6 +107,9 @@ struct Config { qd: usize, align: usize, verify: bool, + workers: usize, + ring_entries: u32, + warmup_runs: usize, } fn parse_usize(name: &str, raw: &str) -> Result { @@ -127,6 +130,9 @@ fn parse_args() -> Result { qd: parse_usize("qd", &a[5])?, align: parse_usize("align", &a[6])?, verify: matches!(std::env::var("BENCH_VERIFY").as_deref(), Ok("1") | Ok("true")), + workers: crate::common::workers()?, + ring_entries: crate::common::ring_entries(128)?, + warmup_runs: crate::common::setting("BENCH_WARMUP_RUNS", 0, 0, 1000)?, }; validate(&cfg)?; Ok(cfg) @@ -269,12 +275,12 @@ fn open_read(cfg: &Config, direct: bool) -> io::Result { // --------------------------------------------------------------------------- /// Sequential buffered read: the StdBackend baseline that rides kernel readahead. -fn run_std_buffered(cfg: &Config) -> Result<(usize, usize), String> { - let mut f = open_read(cfg, false).map_err(|e| format!("open: {e}"))?; - let mut buf = vec![0u8; cfg.chunk]; +fn run_std_buffered(cfg: &Config, prepared: &mut Prepared) -> Result<(usize, usize), String> { + let mut f = prepared.file.as_ref(); + let buf = &mut prepared.buf; let (mut total, mut ops) = (0usize, 0usize); loop { - let n = f.read(&mut buf).map_err(|e| format!("read: {e}"))?; + let n = f.read(buf).map_err(|e| format!("read: {e}"))?; if n == 0 { break; } @@ -296,11 +302,11 @@ fn aligned_buf(len: usize, align: usize) -> (Vec, usize) { } /// Sequential O_DIRECT read: page-cache-bypassing baseline. -fn run_std_odirect(cfg: &Config) -> Result<(usize, usize), String> { - let f = open_read(cfg, true).map_err(|e| format!("open O_DIRECT: {e}"))?; +fn run_std_odirect(cfg: &Config, prepared: &mut Prepared) -> Result<(usize, usize), String> { + let f = prepared.file.as_ref(); // `validate` already requires chunk % align == 0; this is the identity. let chunk = cfg.chunk; - let (mut buf, pad) = aligned_buf(chunk, cfg.align); + let (buf, pad) = (&mut prepared.buf, prepared.pad); let (mut total, mut ops, mut off) = (0usize, 0usize, 0u64); while (off as usize) < cfg.size { let n = f @@ -322,16 +328,10 @@ fn run_std_odirect(cfg: &Config) -> Result<(usize, usize), String> { } /// Pipelined io_uring read at depth `qd`. `direct` selects read_at_direct. -async fn run_uring(cfg: &Config, direct: bool) -> Result<(usize, usize), String> { - // `validate` caps qd at MAX_QD, so this cannot overflow u32. - let depth = (cfg.qd * 2).next_power_of_two() as u32; - let driver = Arc::new(UringDriver::probe_and_start(depth).map_err(|e| format!("probe io_uring: {e:?}"))?); - let file = Arc::new(open_read(cfg, direct).map_err(|e| format!("open: {e}"))?); - - let offsets: Vec<(u64, usize)> = (0..cfg.size) - .step_by(cfg.chunk) - .map(|o| (o as u64, cfg.chunk.min(cfg.size - o))) - .collect(); +async fn run_uring(cfg: &Config, prepared: &Prepared, direct: bool) -> Result<(usize, usize), String> { + let driver = prepared.driver.as_ref().expect("uring strategy has a driver"); + let file = &prepared.file; + let offsets = &prepared.offsets; let (mut total, mut ops) = (0usize, 0usize); let mut set = tokio::task::JoinSet::new(); @@ -376,28 +376,114 @@ async fn run_uring(cfg: &Config, direct: bool) -> Result<(usize, usize), String> Ok((total, ops)) } -fn run(cfg: &Config) -> Result<(usize, usize, f64), String> { - let start = Instant::now(); - let (total, ops) = match cfg.strategy { - Strategy::StdBuffered => run_std_buffered(cfg)?, - Strategy::StdODirect => run_std_odirect(cfg)?, +struct Prepared { + file: Arc, + buf: Vec, + pad: usize, + driver: Option>, + offsets: Vec<(u64, usize)>, +} + +struct Measurement { + total: usize, + ops: usize, + secs: f64, + startup_secs: f64, + shutdown_secs: f64, +} + +fn run_once(cfg: &Config, prepared: &mut Prepared, rt: Option<&tokio::runtime::Runtime>) -> Result<(usize, usize), String> { + match cfg.strategy { + Strategy::StdBuffered => run_std_buffered(cfg, prepared), + Strategy::StdODirect => run_std_odirect(cfg, prepared), Strategy::UringReadAt | Strategy::UringReadAtDirect => { - let direct = cfg.strategy == Strategy::UringReadAtDirect; - let rt = tokio::runtime::Builder::new_multi_thread() + rt.expect("uring strategy has a runtime") + .block_on(run_uring(cfg, prepared, cfg.strategy.is_direct())) + } + } +} + +fn run(cfg: &Config) -> Result { + let startup = Instant::now(); + let uses_uring = matches!(cfg.strategy, Strategy::UringReadAt | Strategy::UringReadAtDirect); + let rt = if uses_uring { + Some( + tokio::runtime::Builder::new_multi_thread() + .worker_threads(cfg.workers) .enable_all() .build() - .map_err(|e| format!("runtime: {e}"))?; - rt.block_on(run_uring(cfg, direct))? - } + .map_err(|e| format!("runtime: {e}"))?, + ) + } else { + None + }; + let file = Arc::new(open_read(cfg, cfg.strategy.is_direct()).map_err(|e| format!("open: {e}"))?); + let driver = if uses_uring { + Some(Arc::new( + UringDriver::probe_and_start(cfg.ring_entries).map_err(|e| format!("probe io_uring: {e:?}"))?, + )) + } else { + None }; + let (buf, pad) = match cfg.strategy { + Strategy::StdBuffered => (vec![0u8; cfg.chunk], 0), + Strategy::StdODirect => aligned_buf(cfg.chunk, cfg.align), + _ => (Vec::new(), 0), + }; + let offsets = if uses_uring { + (0..cfg.size) + .step_by(cfg.chunk) + .map(|o| (o as u64, cfg.chunk.min(cfg.size - o))) + .collect() + } else { + Vec::new() + }; + let mut prepared = Prepared { + file, + buf, + pad, + driver, + offsets, + }; + let startup_secs = startup.elapsed().as_secs_f64(); + for _ in 0..cfg.warmup_runs { + let (total, _) = run_once(cfg, &mut prepared, rt.as_ref())?; + if total != cfg.size { + return Err(format!("warmup read {total} of {} bytes", cfg.size)); + } + // Reset the sequential baseline outside the next timed interval. + prepared.file.as_ref().rewind().map_err(|e| format!("rewind: {e}"))?; + } + #[cfg(feature = "diagnostics")] + let before = prepared.driver.as_ref().map(|driver| driver.diagnostics()); + let start = Instant::now(); + let (total, ops) = run_once(cfg, &mut prepared, rt.as_ref())?; let secs = start.elapsed().as_secs_f64(); + #[cfg(feature = "diagnostics")] + if let (Some(driver), Some(before)) = (&prepared.driver, &before) { + crate::common::report_diagnostics(&driver.diagnostics().since(before)); + } + let shutdown = Instant::now(); + drop(prepared); + drop(rt); + let shutdown_secs = shutdown.elapsed().as_secs_f64(); if total != cfg.size { return Err(format!("strategy {} read {total} of {} bytes", cfg.strategy.name(), cfg.size)); } - Ok((total, ops, secs)) + Ok(Measurement { + total, + ops, + secs, + startup_secs, + shutdown_secs, + }) } pub(super) fn main() -> ExitCode { + if std::env::args().nth(1).as_deref() == Some("--header") { + println!("{CSV_HEADER}"); + return ExitCode::SUCCESS; + } let cfg = match parse_args().and_then(|cfg| ensure_file(&cfg.file, cfg.size).map(|()| cfg)) { Ok(cfg) => cfg, Err(e) => { @@ -406,7 +492,13 @@ pub(super) fn main() -> ExitCode { } }; - let (total, ops, secs) = match run(&cfg) { + let Measurement { + total, + ops, + secs, + startup_secs, + shutdown_secs, + } = match run(&cfg) { Ok(v) => v, Err(e) => { eprintln!("streaming_bench: {e}"); @@ -425,9 +517,10 @@ pub(super) fn main() -> ExitCode { ); } let mbps = (total as f64 / (1024.0 * 1024.0)) / secs; - // CSV: strategy,size,chunk,qd,align,bytes,secs,MBps,ops + let uses_uring = matches!(cfg.strategy, Strategy::UringReadAt | Strategy::UringReadAtDirect); println!( - "{},{},{},{},{},{},{:.6},{:.1},{}", + "2,{},{},{},{},{},{},{},{:.6},{:.1},{},{:.6},{:.6},{},{},{},{}", + if cfg.verify { "verify" } else { "measure" }, cfg.strategy.name(), cfg.size, cfg.chunk, @@ -436,7 +529,13 @@ pub(super) fn main() -> ExitCode { total, secs, mbps, - ops + ops, + startup_secs, + shutdown_secs, + if uses_uring { cfg.workers } else { 0 }, + if uses_uring { cfg.ring_entries } else { 0 }, + cfg.warmup_runs, + crate::common::diagnostics_interval(), ); ExitCode::SUCCESS } diff --git a/scripts/bench-abba.py b/scripts/bench-abba.py new file mode 100644 index 0000000..11e787b --- /dev/null +++ b/scripts/bench-abba.py @@ -0,0 +1,277 @@ +#!/usr/bin/env python3 +# Copyright 2024 RustFS Team +# SPDX-License-Identifier: Apache-2.0 +"""Warm-cache driver ABBA runner. Does not stop services or clear global caches.""" + +import argparse +import csv +import hashlib +import io +import json +import math +import os +from pathlib import Path +import signal +import subprocess +import sys +import time + + +GEOMETRY = ("shards", "file_size", "read_size", "concurrency", "ops", "workers", "ring_entries", "warmup_ops") + + +def parse_row(header, output, expected, interval): + names = next(csv.reader([header])) + if len(names) != len(set(names)): + raise ValueError("duplicate CSV column") + rows = list(csv.reader(io.StringIO(output))) + if len(rows) != 1 or len(rows[0]) != len(names): + raise ValueError("expected exactly one row matching the executable header") + row = dict(zip(names, rows[0])) + required = {*GEOMETRY, "diagnostics_interval", "secs", "IOPS", "MBps", "p50_us", "p99_us", "p999_us", + "startup_secs", "shutdown_secs", "schema_version", "mode", "strategy"} + if not required.issubset(row): + raise ValueError("missing required CSV columns") + if row.get("schema_version") != "2" or row.get("mode") != "measure": + raise ValueError("not a schema-v2 measurement") + if row.get("strategy") != "uring_cached_read": + raise ValueError("unexpected backend") + for key in (*GEOMETRY, "diagnostics_interval"): + wanted = interval if key == "diagnostics_interval" else expected[key] + if int(row[key]) != wanted: + raise ValueError(f"configuration mismatch: {key}") + for key in ("secs", "IOPS", "MBps", "p50_us", "p99_us", "p999_us", "startup_secs", "shutdown_secs"): + value = float(row[key]) + if not math.isfinite(value) or value < 0: + raise ValueError(f"invalid numeric field: {key}") + row[key] = value + if row["secs"] <= 0 or row["IOPS"] <= 0: + raise ValueError("empty measurement") + computed_iops = expected["ops"] / row["secs"] + if abs(row["IOPS"] - computed_iops) > max(1.0, computed_iops * 0.00001): + raise ValueError("IOPS does not match operation count and workload duration") + if abs(row["MBps"] - computed_iops * expected["read_size"] / (1024 * 1024)) > 0.1: + raise ValueError("throughput does not match read geometry") + return row + + +def percent_change(after, before): + if before == 0: + if after == 0: + return 0.0 + raise ValueError("zero baseline prevents percentage attribution") + result = 100.0 * (after / before - 1.0) + if not math.isfinite(result): + raise ValueError("non-finite percentage change") + return result + + +def evaluate_round(rows, throughput_drift, tail_drift): + if [row["leg"] for row in rows] != ["A1", "B1", "B2", "A2"]: + raise ValueError("incomplete or reordered ABBA round") + a1, b1, b2, a2 = [row["measurement"] for row in rows] + drift = { + "iops_pct": percent_change(a2["IOPS"], a1["IOPS"]), + "p99_pct": percent_change(a2["p99_us"], a1["p99_us"]), + } + valid = abs(drift["iops_pct"]) <= throughput_drift and abs(drift["p99_pct"]) <= tail_drift + result = {"valid": valid, "baseline_drift": drift} + # Do not calculate candidate attribution after a failed baseline gate. + if valid: + result["candidate_change_pct"] = { + key: percent_change((b1[key] + b2[key]) / 2, (a1[key] + a2[key]) / 2) + for key in ("IOPS", "p99_us", "p999_us") + } + return result + + +def environment_guard(unit): + if unit: + state = subprocess.run(["systemctl", "is-active", unit], text=True, capture_output=True, check=False) + if state.stdout.strip() != "inactive": + raise RuntimeError("required service is not inactive; no service state was changed") + processes = subprocess.check_output(["ps", "-eo", "comm="], text=True).splitlines() + conflicts = {"rustfs", "warp", "cargo", "rustc", "fio", "samply", "perf", "Runner.Worker"} + if any(name.strip() in conflicts for name in processes): + raise RuntimeError("another build, load generator, service or CI worker is active") + + +def execute(command, env, stdout, stderr, timeout, unit): + with stdout.open("w") as out, stderr.open("w") as err: + process = subprocess.Popen(command, env=env, stdin=subprocess.DEVNULL, stdout=out, stderr=err, start_new_session=True) + deadline = time.monotonic() + timeout + try: + while True: + try: + code = process.wait(timeout=0.5) + break + except subprocess.TimeoutExpired: + if time.monotonic() >= deadline: + raise RuntimeError("benchmark exceeded the per-leg deadline") + environment_guard(unit) + environment_guard(unit) + if code != 0: + raise RuntimeError(f"benchmark failed with exit code {code}; inspect the leg stderr") + except BaseException: + if process.poll() is None: + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + finally: + process.wait() + raise + + +def sha256(path): + with path.open("rb") as file: + return hashlib.file_digest(file, "sha256").hexdigest() + + +def data_identity(path): + stat = path.stat() + return [stat.st_dev, stat.st_ino, stat.st_size, stat.st_mtime_ns, stat.st_ctime_ns] + + +def parse_resources(output): + values = next(csv.reader([output.strip()])) + keys = ("user_seconds", "system_seconds", "wall_seconds", "max_rss_kib", "voluntary_switches", "involuntary_switches") + if len(values) != len(keys): + raise ValueError("incomplete process resource report") + result = dict(zip(keys, map(float, values))) + if any(not math.isfinite(value) or value < 0 for value in result.values()): + raise ValueError("invalid process resource report") + if result["wall_seconds"] <= 0 or result["max_rss_kib"] <= 0: + raise ValueError("empty process resource report") + return result + + +def arguments(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--baseline", type=Path, required=True) + parser.add_argument("--candidate", type=Path, required=True) + parser.add_argument("--data-file", type=Path, required=True, help="pre-created, byte-verified benchmark dataset") + parser.add_argument("--run-dir", type=Path, required=True, help="new output directory; existing paths are refused") + parser.add_argument("--source-revision", required=True) + parser.add_argument("--reservation-note", required=True) + parser.add_argument("--require-inactive-unit") + parser.add_argument("--cpus", default="0-7") + parser.add_argument("--read-size", type=int, default=32768) + parser.add_argument("--concurrency", type=int, default=32) + parser.add_argument("--shards", type=int, default=2) + parser.add_argument("--workers", type=int, default=4) + parser.add_argument("--entries", type=int, default=64) + parser.add_argument("--ops", type=int, default=1_000_000) + parser.add_argument("--warmup-ops", type=int, default=10000) + parser.add_argument("--rounds", type=int, default=3) + parser.add_argument("--candidate-interval", type=int, choices=(0, 64), default=64, + help="0 supports A/A calibration; 64 compares diagnostics on against off") + parser.add_argument("--throughput-drift-pct", type=float, default=3) + parser.add_argument("--p99-drift-pct", type=float, default=5) + parser.add_argument("--timeout", type=float, default=300) + parser.add_argument("--cooldown", type=float, default=1) + parser.add_argument("--min-seconds", type=float, default=5) + parser.add_argument("--dry-run", action="store_true") + args = parser.parse_args() + if not (100000 <= args.ops <= 10_000_000) or args.rounds < 3: + parser.error("acceptance requires 100000..10000000 operations and at least three rounds") + if any(not math.isfinite(value) or value < 0 for value in ( + args.throughput_drift_pct, args.p99_drift_pct, args.timeout, args.cooldown, args.min_seconds + )) or args.timeout == 0 or args.min_seconds == 0: + parser.error("invalid gate or duration") + if not (1 <= args.workers <= 1024 and 1 <= args.shards <= 64 and 1 <= args.concurrency <= 4096): + parser.error("invalid worker/shard/concurrency count") + if not (1 <= args.entries <= 32768) or args.entries & (args.entries - 1): + parser.error("entries must be a power of two in 1..32768") + if not (1 <= args.read_size <= 1 << 30 and 0 <= args.warmup_ops <= 10_000_000): + parser.error("invalid read size or warmup count") + if args.data_file.is_symlink() or not args.data_file.is_file(): + parser.error("data-file must be a regular, pre-created file, not a symlink") + if args.data_file.stat().st_size < args.read_size + 4096: + parser.error("dataset is too small") + return args + + +def main(): + args = arguments() + binaries = {"A": args.baseline.resolve(), "B": args.candidate.resolve()} + headers = {key: subprocess.check_output([str(path), "--header"], text=True).strip() for key, path in binaries.items()} + if headers["A"] != headers["B"]: + raise ValueError("baseline and candidate schema differ") + expected = dict(shards=args.shards, file_size=args.data_file.stat().st_size, read_size=args.read_size, + concurrency=args.concurrency, ops=args.ops, workers=args.workers, + ring_entries=args.entries, warmup_ops=args.warmup_ops) + plan = [(round_id, leg) for round_id in range(1, args.rounds + 1) for leg in ("A1", "B1", "B2", "A2")] + if args.dry_run: + print(json.dumps({"plan": plan, "geometry": expected, "candidate_interval": args.candidate_interval})) + return 0 + environment_guard(args.require_inactive_unit) + args.run_dir.mkdir(mode=0o700) + provenance = {"source_revision": args.source_revision, "reservation": args.reservation_note, + "cache": "warm-preload", "geometry": expected, "cpus": args.cpus, + "data_identity": data_identity(args.data_file), + "binaries": {key: {"path": str(path), "sha256": sha256(path)} for key, path in binaries.items()}, + "gates": {"throughput_drift_pct": args.throughput_drift_pct, "p99_drift_pct": args.p99_drift_pct, + "min_seconds": args.min_seconds}, "candidate_interval": args.candidate_interval} + (args.run_dir / "provenance.json").write_text(json.dumps(provenance, indent=2) + "\n") + summary = {"status": "incomplete", "rounds": []} + try: + current_round = [] + for round_id, leg in plan: + environment_guard(args.require_inactive_unit) + if sha256(binaries[leg[0]]) != provenance["binaries"][leg[0]]["sha256"]: + raise RuntimeError("binary changed during the experiment") + if data_identity(args.data_file) != provenance["data_identity"]: + raise RuntimeError("dataset changed during the experiment") + time.sleep(args.cooldown) + preload_deadline = time.monotonic() + args.timeout + with args.data_file.open("rb") as file: + while file.read(8 << 20): + if time.monotonic() >= preload_deadline: + raise RuntimeError("dataset preload exceeded the deadline") + prefix = args.run_dir / f"round-{round_id}-{leg}" + env = dict(os.environ, BENCH_WORKERS=str(args.workers), BENCH_RING_ENTRIES=str(args.entries), + BENCH_WARMUP_OPS=str(args.warmup_ops), BENCH_VERIFY="0", + BENCH_DIAGNOSTICS="1" if leg[0] == "B" and args.candidate_interval else "0") + command = ["/usr/bin/time", "-f", "%U,%S,%e,%M,%w,%c", "-o", str(prefix) + ".resources.csv", + "taskset", "-c", args.cpus, str(binaries[leg[0]]), "uring_cached_read", str(args.data_file), + str(expected["file_size"]), str(args.read_size), str(args.concurrency), str(args.ops), str(args.shards)] + execute(command, env, Path(str(prefix) + ".csv"), Path(str(prefix) + ".stderr"), + args.timeout, args.require_inactive_unit) + interval = args.candidate_interval if leg[0] == "B" else 0 + row = parse_row(headers[leg[0]], Path(str(prefix) + ".csv").read_text(), expected, interval) + resources = parse_resources(Path(str(prefix) + ".resources.csv").read_text()) + if data_identity(args.data_file) != provenance["data_identity"]: + raise RuntimeError("dataset changed during measurement") + if row["secs"] < args.min_seconds: + raise RuntimeError("measurement too short; increase operations in a new experiment") + if row["secs"] > resources["wall_seconds"] + 0.02: + raise RuntimeError("workload duration exceeds the process resource interval") + record = {"round": round_id, "leg": leg, "measurement": row, "resources": resources} + Path(str(prefix) + ".json").write_text(json.dumps(record, indent=2, allow_nan=False) + "\n") + current_round.append(record) + print(f"round={round_id} leg={leg} IOPS={row['IOPS']:.0f} p99_us={row['p99_us']:.0f}", flush=True) + if leg == "A2": + result = evaluate_round(current_round, args.throughput_drift_pct, args.p99_drift_pct) + summary["rounds"].append(result) + current_round = [] + if not result["valid"]: + raise RuntimeError("baseline drift failed; stopped before expanding the experiment") + summary["status"] = "valid-comparison" + summary["scope"] = "driver-only; resource CSV is whole-process CPU/RSS, not steady-state-only" + return 0 + except (ValueError, RuntimeError, OSError, subprocess.SubprocessError) as error: + summary["status"] = "invalid" + summary["reason"] = str(error) + print(str(error), file=sys.stderr) + return 1 + finally: + (args.run_dir / "summary.json").write_text(json.dumps(summary, indent=2, allow_nan=False) + "\n") + + +if __name__ == "__main__": + try: + sys.exit(main()) + except (ValueError, RuntimeError, OSError, subprocess.SubprocessError) as error: + print(str(error), file=sys.stderr) + sys.exit(1) diff --git a/scripts/test-benchmark-cli.sh b/scripts/test-benchmark-cli.sh new file mode 100644 index 0000000..19ef16c --- /dev/null +++ b/scripts/test-benchmark-cli.sh @@ -0,0 +1,75 @@ +#!/usr/bin/env bash +# Copyright 2024 RustFS Team +# SPDX-License-Identifier: Apache-2.0 +set -euo pipefail +cd "$(dirname "$0")/.." + +if [[ "${BENCH_SKIP_BUILD:-0}" != 1 ]]; then + build_features=() + if [[ "${BENCH_DIAGNOSTICS:-0}" == 1 ]]; then build_features=(--features diagnostics); fi + cargo build --locked --examples "${build_features[@]}" +fi +bin="${BENCH_BIN_DIR:-${CARGO_TARGET_DIR:-target}/debug/examples}" +scratch=$(mktemp -d) +trap 'rm -rf "$scratch"' EXIT + +check_row() { + local header=$1 row=$2 strategy=$3 + awk -F, -v header="$header" -v strategy="$strategy" ' + BEGIN { n = split(header, names, ",") } + { + if (NF != n) exit 1 + for (i = 1; i <= NF; i++) value[names[i]] = $i + if (value["schema_version"] != 2 || value["mode"] != "verify" || value["strategy"] != strategy) exit 1 + if (value["secs"] <= 0 || value["startup_secs"] < 0 || value["shutdown_secs"] < 0) exit 1 + if (strategy ~ /^uring_/ && value["ring_entries"] != 4) exit 1 + if (strategy ~ /^std_/ && value["ring_entries"] != 0) exit 1 + if ("warmup_ops" in value && (value["warmup_ops"] != 8 || value["ops"] != 32 || value["workers"] != 2)) exit 1 + if ("warmup_runs" in value && (value["warmup_runs"] != 1 || value["bytes"] != 1048583)) exit 1 + } + ' <<<"$row" +} + +concurrent="$bin/concurrent_pread_bench" +streaming="$bin/streaming_bench" +header=$("$concurrent" --header) +[[ ",$header," == *,shards,* ]] +for strategy in std_open_pread std_cached_pread uring_open_read uring_cached_read; do + row=$(BENCH_VERIFY=1 BENCH_WORKERS=2 BENCH_RING_ENTRIES=4 BENCH_WARMUP_OPS=8 \ + "$concurrent" "$strategy" "$scratch/random.bin" 1048576 4097 8 32 2) + check_row "$header" "$row" "$strategy" +done + +header=$("$streaming" --header) +for strategy in std_buffered std_odirect uring_read_at uring_read_at_direct; do + row=$(BENCH_VERIFY=1 BENCH_WORKERS=2 BENCH_RING_ENTRIES=4 BENCH_WARMUP_RUNS=1 \ + "$streaming" "$strategy" "$scratch/stream.bin" 1048583 65536 8 4096) + check_row "$header" "$row" "$strategy" +done + +if [[ "${BENCH_DIAGNOSTICS:-0}" == 1 ]]; then + # Warmup samples must be excluded. Each of two shards then sees 128 measured + # reads, giving four final-CQE samples in total rather than six. + BENCH_VERIFY=1 BENCH_WORKERS=2 BENCH_RING_ENTRIES=4 BENCH_WARMUP_OPS=128 \ + "$concurrent" uring_cached_read "$scratch/random.bin" 1048576 4097 8 256 2 \ + >/dev/null 2>"$scratch/diagnostics.log" + awk ' + /^DIAGNOSTICS / { + split($2, stage, "="); split($4, count, "=") + if (stage[2] == "cqe_processing") { if (count[2] < 4) exit 1 } + else if (count[2] != 4) exit 1 + seen++ + } + END { if (seen != 6) exit 1 } + ' "$scratch/diagnostics.log" +fi + +if BENCH_RING_ENTRIES=3 "$concurrent" std_cached_pread "$scratch/random.bin" 1048576 4096 1 1; then + echo "non-power-of-two ring depth was accepted" >&2 + exit 1 +fi +if "$concurrent" uring_cached_read "$scratch/random.bin" 1048576 4096 1 1 0; then + echo "zero shards were accepted" >&2 + exit 1 +fi +echo "benchmark schemas, warmup, saturation and byte-exact strategies passed" diff --git a/scripts/test_bench_abba.py b/scripts/test_bench_abba.py new file mode 100644 index 0000000..5339ae5 --- /dev/null +++ b/scripts/test_bench_abba.py @@ -0,0 +1,83 @@ +# Copyright 2024 RustFS Team +# SPDX-License-Identifier: Apache-2.0 +import importlib.util +import os +from pathlib import Path +import sys +import tempfile +import unittest +from unittest.mock import patch + +SPEC = importlib.util.spec_from_file_location("bench_abba", Path(__file__).with_name("bench-abba.py")) +BENCH = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(BENCH) + + +class Gates(unittest.TestCase): + def round(self, a2_iops=100, a2_p99=20): + return [ + {"leg": leg, "measurement": dict(IOPS=iops, p99_us=p99, p999_us=30)} + for leg, iops, p99 in [("A1", 100, 20), ("B1", 90, 22), ("B2", 90, 22), ("A2", a2_iops, a2_p99)] + ] + + def test_stable_baseline_reports_regression_not_improvement(self): + result = BENCH.evaluate_round(self.round(), 3, 5) + self.assertTrue(result["valid"]) + self.assertAlmostEqual(result["candidate_change_pct"]["IOPS"], -10) + self.assertAlmostEqual(result["candidate_change_pct"]["p99_us"], 10) + + def test_drift_rejects_both_directions_without_attribution(self): + for iops in (90, 110): + result = BENCH.evaluate_round(self.round(a2_iops=iops), 3, 5) + self.assertFalse(result["valid"]) + self.assertNotIn("candidate_change_pct", result) + + def test_tail_drift_alone_invalidates_round(self): + self.assertFalse(BENCH.evaluate_round(self.round(a2_p99=22), 3, 5)["valid"]) + + def test_incomplete_or_reordered_round_is_rejected(self): + for rows in (self.round()[:3], list(reversed(self.round()))): + with self.assertRaises(ValueError): + BENCH.evaluate_round(rows, 3, 5) + + def test_zero_baseline_is_not_a_valid_percentage(self): + with self.assertRaises(ValueError): + BENCH.percent_change(1, 0) + + def test_missing_or_nonfinite_resource_reports_are_rejected(self): + self.assertEqual(BENCH.parse_resources("1,2,5,1000,3,4")["max_rss_kib"], 1000) + for invalid in ("1,2", "1,2,nan,3,4,5", "1,2,3,0,0,0"): + with self.assertRaises(ValueError): + BENCH.parse_resources(invalid) + + def test_verification_and_geometry_drift_are_rejected(self): + expected = dict.fromkeys(BENCH.GEOMETRY, 1) + names = ["schema_version", "mode", "strategy", *BENCH.GEOMETRY, "diagnostics_interval", + "secs", "IOPS", "MBps", "p50_us", "p99_us", "p999_us", "startup_secs", "shutdown_secs"] + values = ["2", "measure", "uring_cached_read", *(["1"] * len(BENCH.GEOMETRY)), "64", *(["1"] * 8)] + values[names.index("MBps")] = "0.0" + header = ",".join(names) + BENCH.parse_row(header, ",".join(values), expected, 64) + with self.assertRaises(ValueError): + BENCH.parse_row("schema_version,mode", "2,measure", expected, 64) + for key, replacement in (("mode", "verify"), ("shards", "2"), ("IOPS", "nan"), + ("IOPS", "100"), ("diagnostics_interval", "0")): + invalid = list(values) + invalid[names.index(key)] = replacement + with self.assertRaises(ValueError): + BENCH.parse_row(header, ",".join(invalid), expected, 64) + + def test_timeout_terminates_the_owned_process_group(self): + with tempfile.TemporaryDirectory() as directory, patch.object(BENCH, "environment_guard"): + out = Path(directory) / "out" + err = Path(directory) / "err" + command = [sys.executable, "-c", "import os,time; print(os.getpid(), flush=True); time.sleep(30)"] + with self.assertRaisesRegex(RuntimeError, "deadline"): + BENCH.execute(command, os.environ.copy(), out, err, 0.01, None) + pid = int(out.read_text().strip()) + with self.assertRaises(ProcessLookupError): + os.kill(pid, 0) + + +if __name__ == "__main__": + unittest.main() diff --git a/src/diagnostics.rs b/src/diagnostics.rs new file mode 100644 index 0000000..8c82834 --- /dev/null +++ b/src/diagnostics.rs @@ -0,0 +1,289 @@ +// Copyright 2024 RustFS Team +// SPDX-License-Identifier: Apache-2.0 + +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Instant; + +/// One in every 64 handle constructions per shard is sampled when `diagnostics` +/// is enabled, starting with that shard's first handle. +pub const DIAGNOSTICS_SAMPLE_INTERVAL: u64 = 64; + +/// Cumulative sampled duration histogram. Concurrent snapshots are approximate. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct LatencyHistogram { + /// Number of recorded durations, including error outcomes reaching this stage. + pub count: u64, + /// Sum of recorded durations in nanoseconds; not an average or percentile. + pub total_nanos: u64, + /// Log2 nanosecond buckets: bucket 0 holds 0..=1 ns, bucket i > 0 holds + /// `[2^i, 2^(i+1))` ns. Bucket 63 includes all larger durations. + pub buckets: [u64; 64], +} + +impl Default for LatencyHistogram { + fn default() -> Self { + Self { + count: 0, + total_nanos: 0, + buckets: [0; 64], + } + } +} + +impl LatencyHistogram { + fn merge(&mut self, other: &Self) { + self.count = self.count.wrapping_add(other.count); + self.total_nanos = self.total_nanos.wrapping_add(other.total_nanos); + for (value, addition) in self.buckets.iter_mut().zip(other.buckets) { + *value = value.wrapping_add(addition); + } + } + + fn since(&self, earlier: &Self) -> Self { + let mut delta = Self { + count: self.count.wrapping_sub(earlier.count), + total_nanos: self.total_nanos.wrapping_sub(earlier.total_nanos), + ..Self::default() + }; + for ((out, current), previous) in delta.buckets.iter_mut().zip(self.buckets).zip(earlier.buckets) { + *out = current.wrapping_sub(previous); + } + delta + } +} + +/// Sampled driver-stage durations, with no object/path labels or exporter locks. +/// +/// Counts can differ: cancellation can stop a request before admission or after +/// completion, and short reads can produce several CQEs per logical operation. +/// The stages overlap as documented and must not be blindly summed. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct DiagnosticsSnapshot { + /// Sample creation during handle construction to permit acquisition, including time before a + /// saturated handle's first poll. Not solely semaphore queue time. + pub admission: LatencyHistogram, + /// Just before message send to driver intake, including channel send cost. + pub driver_queue: LatencyHistogram, + /// Driver intake to placing the initial SQE in the local backlog. + pub preparation: LatencyHistogram, + /// Driver intake to reaping the final read CQE, including preparation, + /// backlog, retries, kernel/device time and reaper delay. Not disk latency. + pub driver_lifetime: LatencyHistogram, + /// Processing time for each sampled read CQE, including range adjustment + /// but excluding result send and pending removal. May overlap lifetime. + pub cqe_processing: LatencyHistogram, + /// Just before result send to the caller's ready poll. This includes + /// receiver inactivity and is not a substitute for runtime schedule latency. + pub completion_to_poll: LatencyHistogram, +} + +impl DiagnosticsSnapshot { + pub(crate) fn merge(&mut self, other: &Self) { + self.admission.merge(&other.admission); + self.driver_queue.merge(&other.driver_queue); + self.preparation.merge(&other.preparation); + self.driver_lifetime.merge(&other.driver_lifetime); + self.cqe_processing.merge(&other.cqe_processing); + self.completion_to_poll.merge(&other.completion_to_poll); + } + + /// Subtract an earlier snapshot from the same driver without resetting + /// counters. Prefer quiescent boundaries; concurrent snapshots need not + /// satisfy cross-field conservation identities. Counters wrap modulo 2^64. + pub fn since(&self, earlier: &Self) -> Self { + Self { + admission: self.admission.since(&earlier.admission), + driver_queue: self.driver_queue.since(&earlier.driver_queue), + preparation: self.preparation.since(&earlier.preparation), + driver_lifetime: self.driver_lifetime.since(&earlier.driver_lifetime), + cqe_processing: self.cqe_processing.since(&earlier.cqe_processing), + completion_to_poll: self.completion_to_poll.since(&earlier.completion_to_poll), + } + } +} + +#[derive(Clone, Copy)] +enum Stage { + Admission, + Queue, + Prepare, + Lifetime, + Cqe, + Resume, +} + +struct Histogram { + count: AtomicU64, + total: AtomicU64, + buckets: [AtomicU64; 64], +} + +impl Default for Histogram { + fn default() -> Self { + Self { + count: AtomicU64::new(0), + total: AtomicU64::new(0), + buckets: std::array::from_fn(|_| AtomicU64::new(0)), + } + } +} + +impl Histogram { + fn record(&self, nanos: u64) { + let bucket = 63 - nanos.max(1).leading_zeros() as usize; + self.buckets[bucket].fetch_add(1, Ordering::Relaxed); + self.total.fetch_add(nanos, Ordering::Relaxed); + self.count.fetch_add(1, Ordering::Relaxed); + } + + fn snapshot(&self) -> LatencyHistogram { + LatencyHistogram { + count: self.count.load(Ordering::Relaxed), + total_nanos: self.total.load(Ordering::Relaxed), + buckets: std::array::from_fn(|i| self.buckets[i].load(Ordering::Relaxed)), + } + } +} + +#[derive(Default)] +pub(crate) struct Diagnostics { + stages: [Histogram; 6], + requests: AtomicU64, +} + +impl Diagnostics { + fn record(&self, stage: Stage, nanos: u64) { + self.stages[stage as usize].record(nanos); + } + + pub(crate) fn snapshot(&self) -> DiagnosticsSnapshot { + DiagnosticsSnapshot { + admission: self.stages[Stage::Admission as usize].snapshot(), + driver_queue: self.stages[Stage::Queue as usize].snapshot(), + preparation: self.stages[Stage::Prepare as usize].snapshot(), + driver_lifetime: self.stages[Stage::Lifetime as usize].snapshot(), + cqe_processing: self.stages[Stage::Cqe as usize].snapshot(), + completion_to_poll: self.stages[Stage::Resume as usize].snapshot(), + } + } +} + +pub(crate) struct Trace { + start: Instant, + diagnostics: Arc, + queued: AtomicU64, + entered: AtomicU64, + sent: AtomicU64, +} + +impl Trace { + pub(crate) fn sample(diagnostics: &Arc) -> Option> { + // A global id modulo 64 would sample only shard zero for power-of-two + // round-robin sharding. Each shard therefore owns its sampling sequence. + let sequence = diagnostics.requests.fetch_add(1, Ordering::Relaxed); + sequence.is_multiple_of(DIAGNOSTICS_SAMPLE_INTERVAL).then(|| { + Arc::new(Self { + start: Instant::now(), + diagnostics: Arc::clone(diagnostics), + queued: AtomicU64::new(0), + entered: AtomicU64::new(0), + sent: AtomicU64::new(0), + }) + }) + } + + pub(crate) fn now(&self) -> u64 { + // Reserve zero for "no result sent". Saturate only at ~584 years. + (self.start.elapsed().as_nanos().min(u128::from(u64::MAX - 1)) as u64) + 1 + } + + pub(crate) fn enqueue(&self) { + let now = self.now(); + self.queued.store(now, Ordering::Relaxed); + self.diagnostics.record(Stage::Admission, now - 1); + } + + pub(crate) fn enter(&self) { + let now = self.now(); + self.entered.store(now, Ordering::Relaxed); + self.diagnostics + .record(Stage::Queue, now.saturating_sub(self.queued.load(Ordering::Relaxed))); + } + + pub(crate) fn prepared(&self) { + self.diagnostics + .record(Stage::Prepare, self.now().saturating_sub(self.entered.load(Ordering::Relaxed))); + } + + pub(crate) fn reaped(&self, started: u64, final_cqe: bool) { + self.diagnostics.record(Stage::Cqe, self.now().saturating_sub(started)); + if final_cqe { + self.diagnostics + .record(Stage::Lifetime, started.saturating_sub(self.entered.load(Ordering::Relaxed))); + } + } + + pub(crate) fn sending(&self) { + self.sent.store(self.now(), Ordering::Release); + } + + pub(crate) fn received(&self) { + let sent = self.sent.load(Ordering::Acquire); + if sent != 0 { + self.diagnostics.record(Stage::Resume, self.now().saturating_sub(sent)); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn histogram_boundaries_and_delta_preserve_samples() { + let histogram = Histogram::default(); + for nanos in [0, 1, 2, 3, 4, 7, 8] { + histogram.record(nanos); + } + let before = histogram.snapshot(); + assert_eq!(&before.buckets[..4], &[2, 2, 2, 1]); + assert_eq!(before.count, 7); + assert_eq!(before.total_nanos, 25); + histogram.record(16); + let delta = histogram.snapshot().since(&before); + assert_eq!(delta.count, 1); + assert_eq!(delta.total_nanos, 16); + assert_eq!(delta.buckets[4], 1); + assert_eq!(delta.buckets.iter().sum::(), 1); + } + + #[test] + fn unsampled_handles_do_not_retain_trace_references() { + let diagnostics = Arc::new(Diagnostics::default()); + let count = (0..128).filter(|_| Trace::sample(&diagnostics).is_some()).count(); + assert_eq!(count, 2); + assert_eq!(Arc::strong_count(&diagnostics), 1); + } + + #[test] + fn unsent_results_do_not_record_receiver_latency() { + let diagnostics = Arc::new(Diagnostics::default()); + let trace = Trace::sample(&diagnostics).expect("first handle is sampled"); + trace.received(); + assert_eq!(diagnostics.snapshot().completion_to_poll.count, 0); + } + + #[test] + fn largest_bucket_and_wrapping_delta_are_defined() { + let histogram = Histogram::default(); + histogram.record(u64::MAX); + let before = histogram.snapshot(); + assert_eq!(before.buckets[63], 1); + histogram.record(2); + let delta = histogram.snapshot().since(&before); + assert_eq!(delta.count, 1); + assert_eq!(delta.total_nanos, 2); + assert_eq!(delta.buckets[1], 1); + } +} diff --git a/src/driver.rs b/src/driver.rs index e013e2d..4072c3f 100644 --- a/src/driver.rs +++ b/src/driver.rs @@ -28,6 +28,9 @@ use std::time::{Duration, Instant}; use io_uring::{IoUring, opcode, types}; +#[cfg(feature = "diagnostics")] +use crate::diagnostics::{Diagnostics, DiagnosticsSnapshot, Trace}; + /// Upper bound on how long shutdown waits for in-flight ops to drain before /// leaking the ring+buffers and exiting (C4, rustfs/backlog#1055). ASYNC_CANCEL /// cannot interrupt an in-execution regular-file read on a D-state/NFS-hung @@ -273,6 +276,8 @@ type AcquireFut = Pin, submitted: AtomicU64, delivered: AtomicU64, orphan_reclaimed: AtomicU64, @@ -287,15 +292,17 @@ struct DriverStats { /// Point-in-time copy of the driver counters. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub struct StatsSnapshot { - /// Read ops handed to the kernel. + /// Logical reads accepted into the pending table, before kernel submission. + /// Short-read resubmissions do not increment this count. pub submitted: u64, - /// CQEs whose result was received by a live caller. + /// Logical results successfully sent to the caller's channel. This does not + /// prove the caller has polled or consumed the result. pub delivered: u64, /// CQEs whose caller had dropped the future: the buffer stayed in the /// pending table the whole time and was reclaimed here, at the CQE. pub orphan_reclaimed: u64, - /// Ops submitted but not yet completed. The kernel may still write into - /// their buffers. + /// Logical reads in the pending table, including queued reads not yet + /// submitted. The kernel may still write into submitted reads' buffers. pub in_flight: u64, /// ASYNC_CANCEL CQEs that reported the target op was canceled (res == 0). pub cancel_succeeded: u64, @@ -325,6 +332,8 @@ pub struct StatsSnapshot { enum Msg { Read { id: u64, + #[cfg(feature = "diagnostics")] + timing: Option>, file: Arc, offset: u64, len: usize, @@ -382,6 +391,8 @@ enum Msg { /// only `buf[pad + head .. pad + head + want]` — alignment padding never /// escapes. struct Pending { + #[cfg(feature = "diagnostics")] + timing: Option>, buf: Vec, file: Arc, done: Option>>>, @@ -473,6 +484,8 @@ enum HandleState { #[must_use = "a read handle must be awaited or explicitly dropped"] pub struct ReadHandle { id: u64, + #[cfg(feature = "diagnostics")] + timing: Option>, rx: oneshot::Receiver>>, tx: mpsc::Sender, finished: bool, @@ -538,10 +551,16 @@ impl Future for ReadHandle { else { unreachable!("state was WaitingPermit") }; + #[cfg(feature = "diagnostics")] + if let Some(timing) = &this.timing { + timing.enqueue(); + } if this .tx .send(Msg::Read { id: this.id, + #[cfg(feature = "diagnostics")] + timing: this.timing.clone(), file, offset, len, @@ -561,6 +580,12 @@ impl Future for ReadHandle { match Pin::new(&mut this.rx).poll(cx) { Poll::Ready(res) => { + #[cfg(feature = "diagnostics")] + if !this.finished + && let Some(timing) = &this.timing + { + timing.received(); + } this.finished = true; Poll::Ready(match res { Ok(inner) => inner, @@ -828,6 +853,8 @@ impl UringDriver { // to a ring whose pending table does not hold the op. The rejection paths // below return an `Inert` handle that never sends, but still need a `tx`. let shard = self.shard(); + #[cfg(feature = "diagnostics")] + let timing = Trace::sample(&shard.stats.diagnostics); // `CURRENT_POSITION` is an internal sentinel used only by // `read_current`; accepting it through a positioned API would silently @@ -840,6 +867,8 @@ impl UringDriver { ))); return ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -862,6 +891,8 @@ impl UringDriver { ))); return ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -883,6 +914,8 @@ impl UringDriver { ))); return ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -917,6 +950,8 @@ impl UringDriver { ))); return ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -935,8 +970,14 @@ impl UringDriver { // no await, and the op is in flight the moment `submit` returns, // exactly as with the previous blocking implementation. Ok(permit) => { + #[cfg(feature = "diagnostics")] + if let Some(timing) = &timing { + timing.enqueue(); + } if let Err(mpsc::SendError(msg)) = shard.tx.send(Msg::Read { id, + #[cfg(feature = "diagnostics")] + timing: timing.clone(), file, offset, len, @@ -954,6 +995,8 @@ impl UringDriver { } return ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -965,6 +1008,8 @@ impl UringDriver { shard.wake_efd.signal(); ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -979,6 +1024,8 @@ impl UringDriver { // handle, which awaits it on its first poll and submits then. Err(TryAcquireError::NoPermits) => ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -998,6 +1045,8 @@ impl UringDriver { let _ = done.send(Err(io::Error::other("uring driver shut down"))); ReadHandle { id, + #[cfg(feature = "diagnostics")] + timing, rx, tx: shard.tx.clone(), finished: false, @@ -1029,6 +1078,25 @@ impl UringDriver { snap } + /// Sampled stage histograms aggregated across shards. Available only with + /// the opt-in `diagnostics` feature; the default driver has no timing fields. + /// See [`DiagnosticsSnapshot`] for stage boundaries and snapshot consistency. + #[cfg(feature = "diagnostics")] + pub fn diagnostics(&self) -> DiagnosticsSnapshot { + let mut snapshot = DiagnosticsSnapshot::default(); + for shard in &self.shards { + snapshot.merge(&shard.stats.diagnostics.snapshot()); + } + snapshot + } + + /// Per-shard sampled histograms in stable shard-index order. Snapshot reads + /// allocate only this output vector and do not acquire driver-thread locks. + #[cfg(feature = "diagnostics")] + pub fn shard_diagnostics(&self) -> Vec { + self.shards.iter().map(|shard| shard.stats.diagnostics.snapshot()).collect() + } + /// Test-only fault injection (rustfs/backlog#1103): poison one driver thread /// so it panics with ops in flight, exercising the `DriverState::Drop` abort /// barrier (C2/#1054). Compiled out entirely unless the `fault-injection` @@ -1470,6 +1538,8 @@ fn drive( match msg { Msg::Read { id, + #[cfg(feature = "diagnostics")] + timing, file, offset, len, @@ -1484,6 +1554,10 @@ fn drive( drop(permit); continue; } + #[cfg(feature = "diagnostics")] + if let Some(timing) = &timing { + timing.enter(); + } // `submit` already validated this geometry. let (kernel_offset, head, region_len) = aligned_geometry(offset, len, align).expect("submit validated the geometry"); @@ -1523,6 +1597,8 @@ fn drive( state.pending.insert( id, Pending { + #[cfg(feature = "diagnostics")] + timing, buf, file, done: Some(done), @@ -1543,6 +1619,10 @@ fn drive( stats.submitted.fetch_add(1, Ordering::SeqCst); stats.in_flight.fetch_add(1, Ordering::SeqCst); state.backlog.push_back(sqe); + #[cfg(feature = "diagnostics")] + if let Some(timing) = state.pending.get(&id).and_then(|p| p.timing.as_ref()) { + timing.prepared(); + } } Msg::Cancel { id } => { if state.pending.contains_key(&id) { @@ -1613,6 +1693,12 @@ fn drive( if !state.pending.contains_key(&ud) { continue; } + #[cfg(feature = "diagnostics")] + let timing = state + .pending + .get(&ud) + .and_then(|p| p.timing.as_ref()) + .map(|trace| (Arc::clone(trace), trace.now())); // Decide the next step while borrowing the entry, then act after // the borrow ends (finish removes it; resubmit re-queues an SQE). @@ -1689,6 +1775,10 @@ fn drive( } }; + #[cfg(feature = "diagnostics")] + if let Some((timing, started)) = &timing { + timing.reaped(*started, matches!(&step, ReapStep::Finish(_))); + } match step { ReapStep::Finish(outcome) => { // Content hygiene (C12, rustfs/backlog#1062): the delivered @@ -1697,6 +1787,10 @@ fn drive( // across requests, this ⊆ [0, res) property MUST be // preserved or a previous tenant's object bytes leak. let mut p = state.pending.remove(&ud).expect("checked above"); + #[cfg(feature = "diagnostics")] + if let Some(timing) = &p.timing { + timing.sending(); + } match p.done.take().expect("done sender set at submit").send(outcome) { Ok(()) => stats.delivered.fetch_add(1, Ordering::SeqCst), // Caller dropped the future: the buffer survived in diff --git a/src/lib.rs b/src/lib.rs index aeb1bd9..d75e314 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -48,5 +48,11 @@ #[cfg(target_os = "linux")] mod driver; +#[cfg(all(target_os = "linux", feature = "diagnostics"))] +mod diagnostics; + +#[cfg(all(target_os = "linux", feature = "diagnostics"))] +pub use diagnostics::{DIAGNOSTICS_SAMPLE_INTERVAL, DiagnosticsSnapshot, LatencyHistogram}; + #[cfg(target_os = "linux")] pub use driver::{ProbeFailure, ReadHandle, StatsSnapshot, UringDriver}; diff --git a/tests/cancel.rs b/tests/cancel.rs index 8a2a03c..bf3716a 100644 --- a/tests/cancel.rs +++ b/tests/cancel.rs @@ -526,6 +526,11 @@ async fn direct_read_returns_exact_unaligned_ranges() { } let Some(file) = open_direct(&path) else { + assert_ne!( + std::env::var("URING_REQUIRE_DIRECT").as_deref(), + Ok("1"), + "O_DIRECT is required by this test lane, but the filesystem rejected it" + ); // Not a skip of the io_uring paths — the suite still exercised them — // only the O_DIRECT assertions cannot run on this filesystem. eprintln!("direct_read_returns_exact_unaligned_ranges: filesystem rejects O_DIRECT, assertions not exercised"); @@ -558,6 +563,7 @@ async fn direct_read_returns_exact_unaligned_ranges() { .expect_err("non-power-of-two alignment must be rejected"); assert_eq!(err.kind(), std::io::ErrorKind::InvalidInput, "unexpected error: {err:?}"); + eprintln!("DIRECT_OK direct_read_returns_exact_unaligned_ranges"); driver.shutdown(); let _ = std::fs::remove_file(&path); } @@ -773,3 +779,97 @@ async fn sharded_cancel_routes_to_the_owning_ring() { drop(pipe_write); driver.shutdown(); } + +#[cfg(feature = "diagnostics")] +#[tokio::test(flavor = "multi_thread")] +async fn diagnostics_sample_stages_without_changing_read_results() { + let Some(driver) = sharded_driver_or_skip("diagnostics_sample_stages_without_changing_read_results", 2) else { + return; + }; + let content = make_content(4096); + let (path, file) = temp_file_with(&content, "diagnostics"); + let before = driver.diagnostics(); + for _ in 0..130 { + let got = driver.read_at(Arc::clone(&file), 1, 512).await.expect("sampled read"); + assert_eq!(got, &content[1..513]); + } + let after = driver.diagnostics(); + let delta = after.since(&before); + for histogram in [ + &delta.admission, + &delta.driver_queue, + &delta.preparation, + &delta.driver_lifetime, + &delta.cqe_processing, + &delta.completion_to_poll, + ] { + // Each shard accepts 65 handles and samples its first and 65th. Global + // id modulo 64 would miss one of the two round-robin shards entirely. + assert_eq!(histogram.count, 4, "{delta:?}"); + assert_eq!(histogram.buckets.iter().sum::(), 4); + } + for shard in driver.shard_diagnostics() { + assert_eq!(shard.driver_lifetime.count, 2); + } + assert_eq!(after.since(&after), rustfs_uring::DiagnosticsSnapshot::default()); + driver.shutdown(); + let _ = std::fs::remove_file(path); +} + +#[cfg(feature = "diagnostics")] +#[tokio::test(flavor = "multi_thread")] +async fn diagnostics_do_not_invent_receiver_samples_for_cancelled_reads() { + let Some(driver) = driver_or_skip("diagnostics_do_not_invent_receiver_samples_for_cancelled_reads") else { + return; + }; + let (read, write) = os_pipe(); + let handle = driver.read_current(read, 64); // sampled id 1 + assert!(wait_until(Duration::from_secs(2), || driver.stats().in_flight == 1).await); + drop(handle); + assert!(wait_until(Duration::from_secs(2), || driver.stats().orphan_reclaimed == 1).await); + let snapshot = driver.diagnostics(); + assert_eq!(snapshot.driver_lifetime.count, 1); + assert_eq!(snapshot.completion_to_poll.count, 0); + assert_eq!(snapshot.cqe_processing.count, 1); + driver.shutdown(); + drop(write); +} + +#[cfg(feature = "diagnostics")] +#[tokio::test(flavor = "multi_thread")] +async fn diagnostics_record_deferred_admission_only_after_a_permit_is_acquired() { + let driver = match UringDriver::probe_and_start(1) { + Ok(driver) => driver, + Err(error) => { + assert!(error.is_expected_restriction(), "{error}"); + eprintln!("SKIP diagnostics_record_deferred_admission_only_after_a_permit_is_acquired: {error}"); + return; + } + }; + let (read, write) = os_pipe(); + let held = driver.read_current(Arc::clone(&read), 64); + assert!(wait_until(Duration::from_secs(2), || driver.stats().in_flight == 1).await); + // Move this shard's sample sequence to its next sampled position without + // acquiring any additional permit or issuing a stream read. + for _ in 0..63 { + let error = driver + .read_at(Arc::clone(&read), u64::MAX, 1) + .await + .expect_err("reserved offset"); + assert_eq!(error.kind(), std::io::ErrorKind::InvalidInput); + } + let content = make_content(512); + let (path, file) = temp_file_with(&content, "diagnostics-deferred"); + let deferred = driver.read_at(file, 0, 512); + assert_eq!(driver.diagnostics().admission.count, 1, "waiting is not admission"); + drop(held); + assert_eq!(deferred.await.expect("permit is returned by cancelled read CQE"), content); + let diagnostics = driver.diagnostics(); + assert_eq!(diagnostics.admission.count, 2); + assert_eq!(diagnostics.driver_queue.count, 2); + assert_eq!(diagnostics.driver_lifetime.count, 2); + assert_eq!(diagnostics.completion_to_poll.count, 1); + driver.shutdown(); + drop(write); + let _ = std::fs::remove_file(path); +}