From d4f56f16d4d67a609e061a24e54c7f9a24ce2aa2 Mon Sep 17 00:00:00 2001 From: Yehor Smoliakov Date: Fri, 25 Sep 2026 10:10:55 +0000 Subject: [PATCH 1/3] Evaluate nightly allocator API for event loop --- .github/workflows/anyio.yml | 4 +- .github/workflows/tests.yml | 40 ++- .github/workflows/vibeio-native.yml | 4 +- .github/workflows/wheels.yml | 8 +- docs/development.md | 15 +- docs/hotpath-lab.md | 13 +- rust-toolchain.toml | 2 +- scripts/hotpath_lab.py | 46 +++- src/lib.rs | 1 + src/transport/stream/server_core.rs | 7 + src/vibeio/allocator.rs | 371 ++++++++++++++++++++++++++++ src/vibeio/lib.rs | 2 + tests/tooling/test_hotpath_lab.py | 25 +- tools/vibeio-check/src/lib.rs | 1 + 14 files changed, 504 insertions(+), 35 deletions(-) create mode 100644 src/vibeio/allocator.rs diff --git a/.github/workflows/anyio.yml b/.github/workflows/anyio.yml index 2098d07..295f447 100644 --- a/.github/workflows/anyio.yml +++ b/.github/workflows/anyio.yml @@ -69,9 +69,9 @@ jobs: cache-suffix: anyio-${{ matrix.os_name }}-py${{ matrix.python_version }} - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Cache Cargo registry uses: Swatinem/rust-cache@v2 diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index e055949..6b99c4b 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -48,9 +48,9 @@ jobs: python-version: "3.13" - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Cache Cargo registry uses: Swatinem/rust-cache@v2 @@ -80,9 +80,9 @@ jobs: python-version: "3.13" - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" components: clippy - name: Cache Cargo registry @@ -91,6 +91,30 @@ jobs: - name: Run Clippy run: cargo clippy --all-targets --all-features --locked -- -D warnings + allocator-miri: + name: allocator miri (strict provenance) + runs-on: ubuntu-latest + timeout-minutes: 15 + steps: + - name: Check out repository + uses: actions/checkout@v6 + with: + persist-credentials: false + + - name: Set up Rust + uses: dtolnay/rust-toolchain@master + with: + toolchain: "nightly-2026-09-25" + components: miri,rust-src + + - name: Cache Cargo registry + uses: Swatinem/rust-cache@v2 + + - name: Run allocator tests under Miri + env: + MIRIFLAGS: -Zmiri-strict-provenance + run: cargo miri test --manifest-path tools/vibeio-check/Cargo.toml allocator::native::tests + type-check: name: pyright runs-on: ubuntu-latest @@ -165,9 +189,9 @@ jobs: cache-suffix: anyio - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Cache Cargo registry uses: Swatinem/rust-cache@v2 @@ -283,9 +307,9 @@ jobs: cache-suffix: ${{ matrix.os_name }}-py${{ matrix.python_version }} - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Cache Cargo registry uses: Swatinem/rust-cache@v2 diff --git a/.github/workflows/vibeio-native.yml b/.github/workflows/vibeio-native.yml index 743cc5d..0636c6a 100644 --- a/.github/workflows/vibeio-native.yml +++ b/.github/workflows/vibeio-native.yml @@ -47,9 +47,9 @@ jobs: persist-credentials: false - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" components: clippy, rustfmt - name: Cache runtime harness build diff --git a/.github/workflows/wheels.yml b/.github/workflows/wheels.yml index a5b4a15..c46c27a 100644 --- a/.github/workflows/wheels.yml +++ b/.github/workflows/wheels.yml @@ -64,9 +64,9 @@ jobs: uses: astral-sh/setup-uv@08807647e7069bb48b6ef5acd8ec9567f424441b - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Install LLVM tools for PGO if: github.event_name == 'workflow_dispatch' && inputs.pgo @@ -114,9 +114,9 @@ jobs: uses: astral-sh/setup-uv@08807647e7069bb48b6ef5acd8ec9567f424441b - name: Set up Rust - uses: dtolnay/rust-toolchain@stable + uses: dtolnay/rust-toolchain@master with: - toolchain: "1.98.1" + toolchain: "nightly-2026-09-25" - name: Cache Cargo registry and target uses: Swatinem/rust-cache@v2 diff --git a/docs/development.md b/docs/development.md index a8d6ef8..e5d188a 100644 --- a/docs/development.md +++ b/docs/development.md @@ -8,7 +8,7 @@ Local development uses Python 3.14.7, pinned in `.python-version`. Install that interpreter before running the `uv` commands below. This development pin does not change the package's Python 3.10+ support or the multi-version test matrix. -Local builds and build/test CI use Rust `1.98.1`, pinned in +Local builds and build/test CI use Rust `nightly-2026-09-25`, pinned in [`rust-toolchain.toml`](https://github.com/RustedBytes/rsloop/blob/master/rust-toolchain.toml). Rustup selects it automatically inside this repository. LLVM tools remain optional for PGO builds. @@ -19,6 +19,19 @@ Quick Rust check: cargo check ``` +The nightly pin supplies the allocator API merged in +[rust-lang/rust#156882](https://github.com/rust-lang/rust/pull/156882). Rsloop +contains a tested internal prototype for a bounded, `System`-backed recycler: +only blocks up to 4 KiB with alignment up to 64 bytes are cached, no bin retains +more than 32 blocks, and total retained memory is capped at 256 KiB. Its first +ready-queue, timer, task, and transport rollout was not enabled: the balanced +holdout gate found no statistically supported primary speedup and did find +workload regressions. Stream payload buffers therefore keep their existing +purpose-built pools, and production collections still use the global allocator. +The Python and public Rust APIs do not expose allocator selection. The remaining +`allocator_ext` feature gate can be removed once the prototype is retired or its +allocator-aware collections stabilize. + Build the extension and install it into the current environment: ```bash diff --git a/docs/hotpath-lab.md b/docs/hotpath-lab.md index 37d6056..08e483e 100644 --- a/docs/hotpath-lab.md +++ b/docs/hotpath-lab.md @@ -79,7 +79,7 @@ This harness does not apply model output or create commits automatically. process performs one full warmup and one measured workload. - `samples.jsonl` is flushed after every process. It includes binary identity, elapsed time, process CPU time, peak RSS, and per-process latency summaries. - Warmups are retained but excluded from inference. Peak RSS includes warmups; + Warmups are retained but excluded from timing inference. Peak RSS includes warmups; process CPU time includes event-loop setup/teardown beyond the workload timer. `harness/` retains the exact runner and benchmark sources for audit (restore their usual repository layout when rerunning them). @@ -87,11 +87,14 @@ This harness does not apply model output or create commits automatically. individual requests. The workload-family intervals use a Bonferroni-adjusted 95% confidence target. Negative percentage changes mean faster. Twelve paired blocks are a minimum, not a promise of adequate statistical power. -- The default gate requires the primary interval to be entirely below -1%, and - **every** workload's interval upper bound to be at most +3%. Any lower bound - above +3% rejects the candidate. Short runs (<0.25 s), fewer than 12 blocks, +- The default balanced gate requires the primary timing interval to be entirely + below -1%, and **every** workload's timing and peak-RSS interval upper bound + to be at most +3%. Any lower bound above +3% rejects the candidate. Short runs + (<0.25 s), fewer than 12 blocks, or uncertain bounds cannot pass. The bulk workload is 64–128 MiB per connection, not the old very short 2 MiB transfer. +- Peak RSS uses the same paired-process bootstrap as timing and participates in + the balanced gate; per-request latency and process CPU remain diagnostic. - Holdout changes concurrency, payloads, task batch size, and work volume. `tls_http` and `websocket_messages` can be added with `--workloads` for entirely untrained scenario families. Default profile training uses the six core @@ -105,7 +108,7 @@ affinity/power controls where appropriate and record them externally. This harness does not pin CPUs on Windows or control thermal state. Child timeouts fail the experiment; incomplete results are retained, not promoted. -`performance_gate_passed` is **not automatic acceptance**. Require both test +`balanced_gate_passed` is **not automatic acceptance**. Require both test logs, inspect diagnostic latency/RSS regressions, and run a fresh confirmation experiment. Rust tests use a debug build of the preserved source; Python tests exercise the exact release extension. The current repository's Python tests are diff --git a/rust-toolchain.toml b/rust-toolchain.toml index e630049..bc11642 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,4 +1,4 @@ [toolchain] -channel = "1.98.1" +channel = "nightly-2026-09-25" profile = "minimal" components = ["clippy", "rustfmt"] diff --git a/scripts/hotpath_lab.py b/scripts/hotpath_lab.py index dd59e70..e931daa 100644 --- a/scripts/hotpath_lab.py +++ b/scripts/hotpath_lab.py @@ -250,6 +250,9 @@ def build(args) -> None: outputs.mkdir() for suffix in ("ll", "s", "pdb"): matches = list((release / "deps").glob(f"rsloop*.{suffix}")) + # Nightly places extra emit artifacts for the generated crate source + # beside that source in OUT_DIR instead of release/deps. + matches.extend((release / "build" / "rsloop").glob(f"*/out/rsloop*.{suffix}")) if suffix != "pdb" and len(matches) != 1: raise ValueError( f"Expected exactly one .{suffix} compiler output, found {matches}" @@ -473,20 +476,30 @@ def paired_estimate( def performance_decision( estimates: dict, + rss_estimates: dict, *, primary: str, minimum_gain: float, regression_budget: float, + rss_regression_budget: float, reliable: bool, ) -> str: if not reliable: return "insufficient_measurement" if any(value["ci_pct"][0] > regression_budget for value in estimates.values()): return "reject_regression" + if any( + value["ci_pct"][0] > rss_regression_budget + for value in rss_estimates.values() + ): + return "reject_rss_regression" if estimates[primary]["ci_pct"][1] < -minimum_gain and all( value["ci_pct"][1] <= regression_budget for value in estimates.values() + ) and all( + value["ci_pct"][1] <= rss_regression_budget + for value in rss_estimates.values() ): - return "performance_gate_passed" + return "balanced_gate_passed" return "inconclusive" @@ -528,7 +541,7 @@ def base_flags(manifest): "candidate": str(candidate), "artifact_manifests": manifests, "primary": args.primary, - "metric": "seconds", + "metrics": ["seconds", "peak_rss_bytes"], "suite": args.suite, "seed": args.seed, "blocks": args.blocks, @@ -536,6 +549,7 @@ def base_flags(manifest): "min_seconds": args.min_seconds, "minimum_gain_pct": args.minimum_gain, "regression_budget_pct": args.regression_budget, + "rss_regression_budget_pct": args.rss_regression_budget, "orders": orders, "configurations": {name: workload_config(name, args.suite) for name in names}, "runner_sha256": sha256(Path(__file__)), @@ -592,6 +606,7 @@ def base_flags(manifest): raise ValueError("Harness changed during timing; results cannot be promoted") # Bonferroni-adjust intervals across the predeclared family of workload metrics. estimates = {} + rss_estimates = {} reliable = args.blocks >= 12 for index, name in enumerate(names): values = { @@ -602,29 +617,51 @@ def base_flags(manifest): estimates[name] = paired_estimate( baseline=values["baseline"], candidate=values["candidate"], - alpha=0.05 / len(names), + alpha=0.05 / (2 * len(names)), seed=args.seed + index, ) estimates[name]["min_observed_seconds"] = min( values["baseline"] + values["candidate"] ) + rss_values = { + label: [ + p["samples"][label]["result"]["peak_rss_bytes"] + for p in pairs[name] + ] + for label in ("baseline", "candidate") + } + rss_estimates[name] = paired_estimate( + baseline=rss_values["baseline"], + candidate=rss_values["candidate"], + alpha=0.05 / (2 * len(names)), + seed=args.seed + len(names) + index, + ) + rss_estimates[name]["baseline_median_bytes"] = rss_estimates[name].pop( + "baseline_median_seconds" + ) + rss_estimates[name]["candidate_median_bytes"] = rss_estimates[name].pop( + "candidate_median_seconds" + ) decision = performance_decision( estimates, + rss_estimates, primary=args.primary, minimum_gain=args.minimum_gain, regression_budget=args.regression_budget, + rss_regression_budget=args.rss_regression_budget, reliable=reliable, ) report = { "schema": SCHEMA, "decision": decision, "estimates": estimates, + "rss_estimates": rss_estimates, "correctness": "NOT certified by timing; run the test subcommand on both artifacts", "notes": [ "Negative changes are faster; positive changes are slower.", "Percentile bootstrap resamples paired fresh processes, not request latencies.", "Intervals are approximate; host load, serial correlation, and repeated candidate searches can still bias inference.", - "Per-process p95/p99, CPU time, and RSS are diagnostic only, not promotion metrics.", + "Peak RSS is a promotion metric; p95/p99 latency and CPU time remain diagnostic.", ], } write_json(out / "report.json", report) @@ -919,6 +956,7 @@ def main() -> None: p.add_argument("--min-seconds", type=float, default=0.25) p.add_argument("--minimum-gain", type=float, default=1.0) p.add_argument("--regression-budget", type=float, default=3.0) + p.add_argument("--rss-regression-budget", type=float, default=3.0) p.add_argument("--timeout", type=float, default=180.0) p.add_argument("--out", type=Path, required=True) p.set_defaults(func=compare) diff --git a/src/lib.rs b/src/lib.rs index bbe4459..0565d03 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,4 +1,5 @@ #![warn(missing_docs)] +#![cfg_attr(not(kani), feature(allocator_ext))] //! Native extension entry point and Rust interoperability API for `rsloop`. //! diff --git a/src/transport/stream/server_core.rs b/src/transport/stream/server_core.rs index a0f5b73..fc8c388 100644 --- a/src/transport/stream/server_core.rs +++ b/src/transport/stream/server_core.rs @@ -116,6 +116,13 @@ impl ServerCore { return None; } let limit = max_pending_tls_handshakes(); + #[cfg(not(kani))] + let reserved = self.pending_tls_handshakes.try_update( + Ordering::AcqRel, + Ordering::Acquire, + |current| reserve_tls_slot(current, limit, false), + ); + #[cfg(kani)] let reserved = self.pending_tls_handshakes.fetch_update( Ordering::AcqRel, Ordering::Acquire, diff --git a/src/vibeio/allocator.rs b/src/vibeio/allocator.rs new file mode 100644 index 0000000..9ac9511 --- /dev/null +++ b/src/vibeio/allocator.rs @@ -0,0 +1,371 @@ +//! Experimental bounded allocator evaluated for event-loop collections. + +#[cfg(not(kani))] +mod native { + use std::alloc::{AllocError, Allocator, AllocatorClone, Layout, System}; + use std::ptr::NonNull; + use std::sync::{Arc, Mutex, MutexGuard}; + + const MIN_CACHED_SIZE: usize = 8; + const MAX_CACHED_SIZE: usize = 4 * 1024; + const MAX_CACHED_ALIGN: usize = 64; + const MAX_BLOCKS_PER_BIN: usize = 32; + const MAX_RETAINED_BYTES: usize = 256 * 1024; + const SIZE_BIN_COUNT: usize = 10; + const ALIGN_BIN_COUNT: usize = 7; + const BIN_COUNT: usize = SIZE_BIN_COUNT * ALIGN_BIN_COUNT; + + #[derive(Clone, Copy)] + struct CachedBlock { + pointer: NonNull, + layout: Layout, + } + + // SAFETY: a cached block has been invalidated by `deallocate`; its bytes + // are inaccessible until the pool transfers ownership from under its mutex. + unsafe impl Send for CachedBlock {} + + struct PoolState { + bins: Vec>, + retained_bytes: usize, + #[cfg(test)] + hits: usize, + #[cfg(test)] + misses: usize, + } + + impl PoolState { + fn new() -> Self { + let bins = (0..BIN_COUNT) + .map(|_| Vec::with_capacity(MAX_BLOCKS_PER_BIN)) + .collect(); + Self { + bins, + retained_bytes: 0, + #[cfg(test)] + hits: 0, + #[cfg(test)] + misses: 0, + } + } + } + + struct PoolInner { + state: Mutex, + } + + impl PoolInner { + #[inline] + fn lock(&self) -> MutexGuard<'_, PoolState> { + match self.state.lock() { + Ok(state) => state, + Err(poisoned) => poisoned.into_inner(), + } + } + } + + impl Drop for PoolInner { + fn drop(&mut self) { + let state = match self.state.get_mut() { + Ok(state) => state, + Err(poisoned) => poisoned.into_inner(), + }; + for bin in &mut state.bins { + for block in bin.drain(..) { + // SAFETY: cached blocks were allocated from `System` with this + // exact layout and were invalidated before entering the cache. + unsafe { + System.deallocate(block.pointer, block.layout); + } + } + } + } + } + + /// Cloneable allocator handle shared by one event loop and its runtimes. + #[derive(Clone)] + pub(crate) struct RuntimeAllocator { + inner: Arc, + } + + impl RuntimeAllocator { + pub(crate) fn new() -> Self { + Self { + inner: Arc::new(PoolInner { + state: Mutex::new(PoolState::new()), + }), + } + } + + #[inline] + fn cached_layout(layout: Layout) -> Option<(usize, Layout)> { + if layout.size() == 0 + || layout.size() > MAX_CACHED_SIZE + || layout.align() > MAX_CACHED_ALIGN + { + return None; + } + let size = layout.size().max(MIN_CACHED_SIZE).next_power_of_two(); + let size_index = + size.trailing_zeros() as usize - MIN_CACHED_SIZE.trailing_zeros() as usize; + let align_index = layout.align().trailing_zeros() as usize; + let index = size_index * ALIGN_BIN_COUNT + align_index; + // SAFETY: `size` is non-zero and `layout.align()` is a power of two. + let actual = unsafe { Layout::from_size_align_unchecked(size, layout.align()) }; + Some((index, actual)) + } + + #[inline] + fn actual_layout(layout: Layout) -> Layout { + Self::cached_layout(layout).map_or(layout, |(_, actual)| actual) + } + + #[cfg(test)] + fn stats(&self) -> (usize, usize, usize) { + let state = self.inner.lock(); + (state.hits, state.misses, state.retained_bytes) + } + } + + // SAFETY: live blocks remain owned by `System`; the mutex-protected cache + // contains only invalidated blocks, preserves their exact normalized + // layouts, and never touches them until a new allocation takes ownership. + unsafe impl Allocator for RuntimeAllocator { + fn allocate(&self, layout: Layout) -> Result, AllocError> { + let Some((index, actual)) = Self::cached_layout(layout) else { + return System.allocate(layout); + }; + + { + let mut state = self.inner.lock(); + if let Some(block) = state.bins[index].pop() { + state.retained_bytes = state.retained_bytes.saturating_sub(block.layout.size()); + #[cfg(test)] + { + state.hits += 1; + } + // SAFETY: the cache owns this non-null `System` allocation. + return Ok(NonNull::slice_from_raw_parts(block.pointer, actual.size())); + } + #[cfg(test)] + { + state.misses += 1; + } + } + + let block = System.allocate(actual)?; + // SAFETY: a successful allocation always returns a non-null data pointer. + let ptr = unsafe { NonNull::new_unchecked(block.as_ptr().cast::()) }; + Ok(NonNull::slice_from_raw_parts(ptr, actual.size())) + } + + unsafe fn deallocate(&self, ptr: NonNull, layout: Layout) { + let Some((index, actual)) = Self::cached_layout(layout) else { + // SAFETY: upheld by this method's caller. + unsafe { System.deallocate(ptr, layout) }; + return; + }; + + let mut state = self.inner.lock(); + if state.bins[index].len() < MAX_BLOCKS_PER_BIN + && state.retained_bytes.saturating_add(actual.size()) <= MAX_RETAINED_BYTES + { + state.bins[index].push(CachedBlock { + pointer: ptr, + layout: actual, + }); + state.retained_bytes += actual.size(); + return; + } + drop(state); + // SAFETY: the block was allocated with the normalized layout above. + unsafe { System.deallocate(ptr, actual) }; + } + + unsafe fn grow( + &self, + ptr: NonNull, + old_layout: Layout, + new_layout: Layout, + ) -> Result, AllocError> { + let old_actual = Self::actual_layout(old_layout); + let new_actual = Self::actual_layout(new_layout); + // SAFETY: every block returned by this allocator is System-backed + // with its normalized layout, and the caller upholds grow's rules. + unsafe { System.grow(ptr, old_actual, new_actual) } + } + + unsafe fn grow_zeroed( + &self, + ptr: NonNull, + old_layout: Layout, + new_layout: Layout, + ) -> Result, AllocError> { + let old_actual = Self::actual_layout(old_layout); + let new_actual = Self::actual_layout(new_layout); + // SAFETY: every block returned by this allocator is System-backed + // with its normalized layout, and the caller upholds grow's rules. + unsafe { System.grow_zeroed(ptr, old_actual, new_actual) } + } + + unsafe fn shrink( + &self, + ptr: NonNull, + old_layout: Layout, + new_layout: Layout, + ) -> Result, AllocError> { + let old_actual = Self::actual_layout(old_layout); + let new_actual = Self::actual_layout(new_layout); + // SAFETY: every block returned by this allocator is System-backed + // with its normalized layout, and the caller upholds shrink's rules. + unsafe { System.shrink(ptr, old_actual, new_actual) } + } + } + + // SAFETY: clones share the same pool, so either clone can free an allocation + // made by the other without invalidating allocations when a clone drops. + unsafe impl AllocatorClone for RuntimeAllocator {} + + #[cfg(test)] + mod tests { + use super::*; + + #[test] + fn cached_block_is_reused_and_bounded() { + let allocator = RuntimeAllocator::new(); + let layout = Layout::from_size_align(24, 8).unwrap(); + let first = allocator.allocate(layout).unwrap(); + let first_ptr = first.as_ptr().cast::(); + // SAFETY: `first` was allocated above with `layout`. + unsafe { allocator.deallocate(NonNull::new_unchecked(first_ptr), layout) }; + assert_eq!(allocator.stats().2, 32); + + let second = allocator.allocate(layout).unwrap(); + assert_eq!(second.as_ptr().cast::(), first_ptr); + assert_eq!(allocator.stats(), (1, 1, 0)); + // SAFETY: `second` is live and uses the same layout. + unsafe { + allocator.deallocate(NonNull::new_unchecked(second.as_ptr().cast::()), layout); + } + } + + #[test] + fn zeroed_allocation_clears_reused_memory() { + let allocator = RuntimeAllocator::new(); + let layout = Layout::from_size_align(16, 8).unwrap(); + let block = allocator.allocate(layout).unwrap(); + let ptr = block.as_ptr().cast::(); + // SAFETY: the allocation contains at least 16 writable bytes. + unsafe { + ptr.write_bytes(0xA5, 16); + allocator.deallocate(NonNull::new_unchecked(ptr), layout); + } + let zeroed = allocator.allocate_zeroed(layout).unwrap(); + let ptr = zeroed.as_ptr().cast::(); + // SAFETY: the returned allocation contains at least 16 initialized bytes. + let bytes = unsafe { std::slice::from_raw_parts(ptr, 16) }; + assert!(bytes.iter().all(|byte| *byte == 0)); + // SAFETY: `zeroed` is live and uses `layout`. + unsafe { allocator.deallocate(NonNull::new_unchecked(ptr), layout) }; + } + + #[test] + fn large_allocations_bypass_the_cache() { + let allocator = RuntimeAllocator::new(); + let layout = Layout::from_size_align(MAX_CACHED_SIZE + 1, 8).unwrap(); + let block = allocator.allocate(layout).unwrap(); + // SAFETY: `block` was allocated above with `layout`. + unsafe { + allocator.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); + } + assert_eq!(allocator.stats(), (0, 0, 0)); + } + + #[test] + fn clones_share_cached_blocks_across_threads() { + let allocator = RuntimeAllocator::new(); + let other = allocator.clone(); + let layout = Layout::from_size_align(64, 16).unwrap(); + let address = std::thread::spawn(move || { + let block = other.allocate(layout).unwrap(); + let address = block.as_ptr().cast::() as usize; + // SAFETY: the block was allocated and is deallocated on this thread. + unsafe { + other.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); + } + address + }) + .join() + .unwrap(); + assert_eq!(allocator.stats().2, 64); + let reused = allocator.allocate(layout).unwrap(); + assert_eq!(reused.as_ptr().cast::() as usize, address); + // SAFETY: `reused` is live and was allocated with `layout`. + unsafe { + allocator.deallocate(NonNull::new_unchecked(reused.as_ptr().cast::()), layout); + } + } + + #[test] + fn bin_limit_bounds_retained_memory() { + let allocator = RuntimeAllocator::new(); + let layout = Layout::from_size_align(MAX_CACHED_SIZE, 8).unwrap(); + let blocks = (0..(MAX_BLOCKS_PER_BIN + 8)) + .map(|_| allocator.allocate(layout).unwrap()) + .collect::>(); + for block in blocks { + // SAFETY: every block is live and was allocated with `layout`. + unsafe { + allocator + .deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); + } + } + assert_eq!(allocator.stats().2, MAX_BLOCKS_PER_BIN * MAX_CACHED_SIZE); + } + + #[test] + fn grow_preserves_initialized_bytes() { + let allocator = RuntimeAllocator::new(); + let old_layout = Layout::from_size_align(24, 8).unwrap(); + let new_layout = Layout::from_size_align(80, 8).unwrap(); + let block = allocator.allocate(old_layout).unwrap(); + let pointer = block.as_ptr().cast::(); + // SAFETY: the old allocation has 24 writable bytes and is live for grow. + let grown = unsafe { + for offset in 0..24 { + pointer.add(offset).write(offset as u8); + } + allocator + .grow(NonNull::new_unchecked(pointer), old_layout, new_layout) + .unwrap() + }; + let pointer = grown.as_ptr().cast::(); + // SAFETY: grow preserves the first 24 initialized bytes. + for offset in 0..24 { + // SAFETY: grow preserves the first 24 initialized bytes. + assert_eq!(unsafe { pointer.add(offset).read() }, offset as u8); + } + // SAFETY: `grown` is live and uses `new_layout`. + unsafe { + allocator.deallocate(NonNull::new_unchecked(pointer), new_layout); + } + } + + #[test] + fn poisoned_metadata_lock_is_recovered() { + let allocator = RuntimeAllocator::new(); + let poisoned = allocator.clone(); + let _ = std::panic::catch_unwind(move || { + let _guard = poisoned.inner.lock(); + panic!("poison allocator metadata for the recovery test"); + }); + let layout = Layout::from_size_align(32, 8).unwrap(); + let block = allocator.allocate(layout).unwrap(); + // SAFETY: `block` is live and uses `layout`. + unsafe { + allocator.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); + } + assert_eq!(allocator.stats().2, 32); + } + } +} diff --git a/src/vibeio/lib.rs b/src/vibeio/lib.rs index 53dceda..ed2fe30 100644 --- a/src/vibeio/lib.rs +++ b/src/vibeio/lib.rs @@ -51,6 +51,8 @@ //! - `splice` - enables splice support (Linux). //! - `blocking-default` - enables the default blocking thread pool. +#[cfg(test)] +pub(crate) mod allocator; pub mod blocking; mod builder; mod driver; diff --git a/tests/tooling/test_hotpath_lab.py b/tests/tooling/test_hotpath_lab.py index 21cab97..1a24a16 100644 --- a/tests/tooling/test_hotpath_lab.py +++ b/tests/tooling/test_hotpath_lab.py @@ -48,23 +48,27 @@ def test_invalid_samples_are_rejected(baseline, candidate): @pytest.mark.parametrize( - "primary,guard,reliable,expected", + "primary,guard,rss,reliable,expected", [ - ((-4, -2), (-1, 2), True, "performance_gate_passed"), - ((-4, 0.1), (-1, 2), True, "inconclusive"), - ((-4, -2), (-1, 4), True, "inconclusive"), - ((-4, -2), (4, 5), True, "reject_regression"), - ((-4, -2), (-1, 2), False, "insufficient_measurement"), + ((-4, -2), (-1, 2), (-2, 2), True, "balanced_gate_passed"), + ((-4, 0.1), (-1, 2), (-2, 2), True, "inconclusive"), + ((-4, -2), (-1, 4), (-2, 2), True, "inconclusive"), + ((-4, -2), (4, 5), (-2, 2), True, "reject_regression"), + ((-4, -2), (-1, 2), (4, 5), True, "reject_rss_regression"), + ((-4, -2), (-1, 2), (-2, 4), True, "inconclusive"), + ((-4, -2), (-1, 2), (-2, 2), False, "insufficient_measurement"), ], ) -def test_gate_requires_gain_and_bounded_regression(primary, guard, reliable, expected): +def test_gate_requires_gain_and_bounded_regression(primary, guard, rss, reliable, expected): estimates = {"tcp": {"ci_pct": primary}, "tasks": {"ci_pct": guard}} assert ( lab.performance_decision( estimates, + {"tcp": {"ci_pct": rss}, "tasks": {"ci_pct": rss}}, primary="tcp", minimum_gain=1.0, regression_budget=3.0, + rss_regression_budget=3.0, reliable=reliable, ) == expected @@ -266,7 +270,11 @@ def test_compare_archives_harness_and_refuses_to_promote_short_experiment( "flags": [], } monkeypatch.setattr(lab, "load_artifact", lambda path: manifest) - monkeypatch.setattr(lab, "run_sample", lambda *args: {"result": {"seconds": 1.0}}) + monkeypatch.setattr( + lab, + "run_sample", + lambda *args: {"result": {"seconds": 1.0, "peak_rss_bytes": 1_000_000}}, + ) args = SimpleNamespace( baseline=tmp_path / "baseline", candidate=tmp_path / "candidate", @@ -280,6 +288,7 @@ def test_compare_archives_harness_and_refuses_to_promote_short_experiment( min_seconds=0.25, minimum_gain=1.0, regression_budget=3.0, + rss_regression_budget=3.0, timeout=10.0, ) lab.compare(args) diff --git a/tools/vibeio-check/src/lib.rs b/tools/vibeio-check/src/lib.rs index 69874f3..9a48fd4 100644 --- a/tools/vibeio-check/src/lib.rs +++ b/tools/vibeio-check/src/lib.rs @@ -1,5 +1,6 @@ //! Build/test harness for rsloop's embedded runtime; no copied implementation. #![doc = include_str!("../EXAMPLES.md")] +#![feature(allocator_ext)] #[path = "../../../src/vibeio/lib.rs"] pub mod vibeio; From a8b840472e80e71faa9bd0a124d78335245c03d8 Mon Sep 17 00:00:00 2001 From: Yehor Smoliakov Date: Fri, 25 Sep 2026 10:52:09 +0000 Subject: [PATCH 2/3] Add opt-in scheduler batch reuse with stabilized allocator API Replace the shared allocator prototype with a bounded owner-thread cache. Validate allocator and drain safety with Miri, and measure randomized paired runtime and application workloads. Keep the feature disabled by default: native idle turns improve about 15%, but TCP workloads regress. --- .github/workflows/tests.yml | 2 +- Cargo.toml | 2 + docs/allocator-batch-results.md | 146 ++++++++ docs/development.md | 57 ++- scripts/compare_runtime_turns.py | 86 +++++ src/lib.rs | 1 - src/vibeio/allocator.rs | 371 ------------------- src/vibeio/batch_allocator.rs | 355 ++++++++++++++++++ src/vibeio/executor.rs | 52 ++- src/vibeio/lib.rs | 5 +- tools/vibeio-check/Cargo.toml | 1 + tools/vibeio-check/examples/runtime_turns.rs | 52 +++ tools/vibeio-check/src/lib.rs | 1 - 13 files changed, 740 insertions(+), 391 deletions(-) create mode 100644 docs/allocator-batch-results.md create mode 100644 scripts/compare_runtime_turns.py delete mode 100644 src/vibeio/allocator.rs create mode 100644 src/vibeio/batch_allocator.rs create mode 100644 tools/vibeio-check/examples/runtime_turns.rs diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 6b99c4b..a367419 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -113,7 +113,7 @@ jobs: - name: Run allocator tests under Miri env: MIRIFLAGS: -Zmiri-strict-provenance - run: cargo miri test --manifest-path tools/vibeio-check/Cargo.toml allocator::native::tests + run: cargo miri test --manifest-path tools/vibeio-check/Cargo.toml --features scheduler-batch-cache batch_ type-check: name: pyright diff --git a/Cargo.toml b/Cargo.toml index d85c0b5..972f009 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -60,6 +60,8 @@ windows-sys = { version = "0.61", features = ["Wdk_Foundation", "Wdk_Storage_Fil [features] default = [] +# Opt-in: faster scheduler turns, but measured TCP workloads can regress. +scheduler-batch-cache = [] blocking-default = ["dep:rusty_pool"] # Embedded-runtime modules are opt-in; production wheels retain the default # networking/timer build unless a feature is explicitly requested. diff --git a/docs/allocator-batch-results.md b/docs/allocator-batch-results.md new file mode 100644 index 0000000..277d8f5 --- /dev/null +++ b/docs/allocator-batch-results.md @@ -0,0 +1,146 @@ +# Scheduler batch allocation reuse + +Measured on 2026-09-25. This change uses the stabilized allocator API available +in the pinned `nightly-2026-09-25` compiler, without an `allocator_ext` feature +gate. It does **not** claim compatibility with an older stable-channel compiler. + +The cache is **opt-in** through `scheduler-batch-cache`, not enabled in default +builds or wheels. Its scheduler-turn improvement comes with a confirmed TCP +slowdown, so the general-purpose event loop retains its original allocation +and drain path. + +## What changes + +`Runtime::block_on` previously allocated and freed its ready-task vector on +every entry. With the feature enabled, the runtime owns a single-slot, exact-layout +`System` allocation cache, used through `Allocator` and `Vec::with_capacity_in`. +Successive calls reuse that storage; batches within a call keep the vector's +capacity. Tests verify one system allocation across repeated runtime turns, +including root-future unwinding and reuse afterward. + +The usual retained allocation is 2 KiB on 64-bit systems, with a hard 16 KiB +ceiling per runtime. Runtime destruction frees it. This reduces allocator +traffic, not the live size of tasks or a guaranteed amount of process RSS. +No mutex, allocator `Arc` clone, size-class search, or public configuration is +introduced. The cache sits outside `RuntimeInner` to preserve the hot shared +scheduler state's layout. + +Allocator-aware `Vec::drain` is not in the stabilized subset. A private full +drain uses stable vector primitives, retains the allocation, preserves FIFO +order, and drops unconsumed tasks on unwind. Strict-provenance Miri tests cover +partial consumption, destructor panics, zero-sized elements, forgotten drains, +alignment, growth, zeroing, disjoint simultaneous allocations, and runtime +panic cleanup. Kani's older compiler retains an ordinary-vector fallback; +its proofs do not verify the custom allocator. + +The previous shared size-class allocator remains rejected. Early versions of +this change also regressed busy turns when replacing the vector each batch or +using optional task slots. Neither implementation is retained. + +## Native scheduler measurement + +Host: Intel Core i9-9900K, Linux x86-64; Rust 1.100.0-nightly +(`f7575a9da`, LLVM 23.1.1). Identical `runtime_turns.rs` source and lockfile were +built with the same compiler and release profile, against `d4f56f1` and the +new implementation with `--features scheduler-batch-cache`. No other builds +or tests ran during timing. CPU affinity +was unrestricted; these are local results, not a cross-platform guarantee. + +Each comparison used 24 randomized, balanced baseline/candidate process pairs, +10,000 warm-up turns per process, and paired bootstrap 95% confidence intervals. +Negative percentages mean less elapsed time. + +| Workload | Turns per sample | Baseline median | Candidate median | Change (95% interval) | +| --- | ---: | ---: | ---: | ---: | +| Idle scheduler entry/exit | 7,000,000 | 0.76260 s | 0.64408 s | -15.41% [-15.87%, -14.88%] | +| 64 continuously ready tasks | 500,000 | 1.12497 s | 1.12506 s | +0.30% [-0.48%, +1.14%] | + +The idle-turn gain passed the declared 1% timing threshold; the busy result is +inconclusive and within the 3% regression budget at this confidence level. +These measurements isolate scheduler turns, not Python or network throughput. + +Reproduce with `scripts/compare_runtime_turns.py` as described in +[Development](development.md). Use `--iterations 7000000 --tasks 0 --seed 999` +for idle turns and `--iterations 500000 --tasks 64 --seed 1000` for busy turns. +Both use `--blocks 24` and separate output directories. + +Local raw plans, binary hashes, paired samples, and reports are in +`target/allocator-batch-optin-idle` and `target/allocator-batch-optin-active`. +These ignored directories are local experiment artifacts, not checked-in data. + +## Python application workloads + +The six-workload holdout used CPython 3.14.7, 24 paired process blocks, +seed 997, and `tcp_streams` as the predeclared primary. Compiler versions, +non-PGO flags, benchmark source, and artifact hashes were checked by the +hot-path harness. Each interval below uses 99.583% confidence (Bonferroni +adjustment across six timing and six RSS metrics). + +| Workload | Time change (interval) | Peak RSS change (interval) | +| --- | ---: | ---: | +| Callbacks | -0.14% [-0.70%, +0.36%] | +0.00% [-0.01%, +0.02%] | +| Tasks | -0.64% [-1.50%, +0.14%] | -0.19% [-0.41%, +0.01%] | +| TCP streams | +2.10% [+0.64%, +3.56%] | +0.00% [-0.23%, +0.23%] | +| HTTP keep-alive | -0.30% [-6.67%, +3.78%] | +0.01% [-0.28%, +0.32%] | +| Mixed streams | -0.71% [-2.98%, +1.33%] | +0.03% [-0.29%, +0.38%] | +| Bulk transfer | -0.43% [-2.01%, +0.82%] | -0.00% [-0.02%, +0.01%] | + +The application-wide gate is **inconclusive**, not passed: the primary did +not improve, TCP showed a small slowdown, and the TCP/HTTP upper timing bounds +exceed the 3% budget. Peak RSS stayed well within budget, but there is no +established application-wide memory reduction or throughput improvement. + +An independent 40-pair TCP confirmation (seed 998) found **+2.75%** elapsed time, +with a 97.5% interval of **[+1.79%, +3.85%]**. RSS changed +0.04% +[-0.14%, +0.24%]. This confirmed the tradeoff and is why the cache is opt-in. +Raw confirmation data is in `target/allocator-batch-final-tcp-confirmation`. + +The exact baseline is `d4f56f1`; the local cache-enabled measurement snapshot is +`a87092ec214b6daa9288818109f58cb6034846d3`. It predates the final feature guard: +these application results evaluate the cache integration, not the final default +build (which uses the original path). Artifact directories +are `target/allocator-batch-baseline-final` and `target/allocator-batch-final`; +the complete plan, samples, manifests and report are under +`target/allocator-batch-final-holdout`. + +The following command records the cache-enabled snapshot comparison. To repeat +the experiment on later revisions, use a cache-enabled candidate build; an +ordinary default build no longer enables the cache. Use fresh output directories. + +```bash +.venv/bin/python scripts/hotpath_lab.py compare \ + --baseline target/allocator-batch-baseline-final \ + --candidate target/allocator-batch-final \ + --primary tcp_streams --suite holdout --blocks 24 --seed 997 \ + --out target/allocator-batch-final-holdout +``` + +## Correctness checks + +- Root Rust suites: 307 passed with default features; 407 with all features. +- Embedded runtime all-feature suite: 283 passed; 16 doctests passed. +- Strict-provenance Miri: 10 allocator/drain/runtime regression tests passed. +- Kani: all 31 `merge_` harnesses passed, using the ordinary-vector fallback. +- Clippy: both crates, all targets/features, warnings denied. +- Python: 192 passed and 2 skipped on **each** isolated extension; tooling, + stress, and slow-network markers were excluded (73 deselected). +- Final default development install: the same 192 Python tests passed, with + 2 skipped and 73 deselected. +- Repository tooling: 67 passed. The new comparison runner also passes Ruff. + +Relevant commands: + +```bash +python3 scripts/run_rust_tests.py --all-features +cargo test --manifest-path tools/vibeio-check/Cargo.toml --all-features --locked +MIRIFLAGS=-Zmiri-strict-provenance cargo miri test \ + --manifest-path tools/vibeio-check/Cargo.toml --features scheduler-batch-cache batch_ +cargo kani --harness merge_ -j 2 --output-format terse +PYTEST_ADDOPTS="-m 'not tooling and not stress and not slow_network'" \ + .venv/bin/python scripts/hotpath_lab.py test \ + --artifact target/allocator-batch-final --python-child +``` + +The Python command was also run against `target/allocator-batch-baseline-final`. +These checks do not replace cross-platform CI or the excluded stress/network +suites. diff --git a/docs/development.md b/docs/development.md index e5d188a..8d9530f 100644 --- a/docs/development.md +++ b/docs/development.md @@ -20,17 +20,52 @@ cargo check ``` The nightly pin supplies the allocator API merged in -[rust-lang/rust#156882](https://github.com/rust-lang/rust/pull/156882). Rsloop -contains a tested internal prototype for a bounded, `System`-backed recycler: -only blocks up to 4 KiB with alignment up to 64 bytes are cached, no bin retains -more than 32 blocks, and total retained memory is capped at 256 KiB. Its first -ready-queue, timer, task, and transport rollout was not enabled: the balanced -holdout gate found no statistically supported primary speedup and did find -workload regressions. Stream payload buffers therefore keep their existing -purpose-built pools, and production collections still use the global allocator. -The Python and public Rust APIs do not expose allocator selection. The remaining -`allocator_ext` feature gate can be removed once the prototype is retired or its -allocator-aware collections stabilize. +[rust-lang/rust#156882](https://github.com/rust-lang/rust/pull/156882). The opt-in +`scheduler-batch-cache` Cargo feature uses the stabilized `Allocator` and +`Vec::with_capacity_in` APIs for scheduler batch storage. With it enabled, +each embedded runtime caches one exact-layout `System` allocation +between `block_on` calls: normally 2 KiB on 64-bit platforms, with a hard 16 KiB +retention ceiling. The cache belongs to the runtime's owner thread and needs no +mutex or reference-counted allocator handles. Simultaneous allocations remain +disjoint, all task values are dropped before storage is recycled, and cached +storage is freed when the runtime is destroyed. + +A full-batch drain retains the vector's capacity within each `block_on` call; +it uses stabilized vector primitives because allocator-aware `Vec::drain` is +still outside the stabilized subset. Unconsumed elements are dropped on early +exit or unwind. The cache lives outside the shared scheduler state to preserve +that hot structure's layout. + +This replaces the rejected shared, size-class recycler experiment. Ready queues, +timers, task futures, transport queues, and stream payload pools keep their +existing allocation strategies. There is no `allocator_ext` feature gate or +runtime allocator-selection API. Default builds retain ordinary `Vec` allocation +and draining: the cache improved native scheduler turns but regressed TCP +workloads, so it is not enabled in production wheels. Kani's older compiler also +uses ordinary `Vec` for +scheduler proofs; strict-provenance Miri tests exercise the actual allocator. + +For an isolated measurement of scheduler entry/exit costs, build +`tools/vibeio-check/examples/runtime_turns.rs` against both revisions with the +same toolchain, lockfile, profile, and benchmark source, then run: + +```bash +cargo build --manifest-path tools/vibeio-check/Cargo.toml --example runtime_turns --release --locked --features scheduler-batch-cache +.venv/bin/python scripts/compare_runtime_turns.py \ + --baseline /path/to/baseline/runtime_turns \ + --candidate tools/vibeio-check/target/release/examples/runtime_turns \ + --out target/runtime-turns-comparison +``` + +The runner records binary hashes and randomized paired process order before +measurement. This benchmark isolates runtime turns; use the +[hot-path workload gate](hotpath-lab.md) separately to evaluate application +timing and peak RSS. +See the [measured results and limitations](allocator-batch-results.md). + +To evaluate the cache with your own Python workload, build explicitly with +`uv run --with maturin maturin develop --release --features scheduler-batch-cache`. +Rebuild without that feature to restore the default path. Build the extension and install it into the current environment: diff --git a/scripts/compare_runtime_turns.py b/scripts/compare_runtime_turns.py new file mode 100644 index 0000000..3b1196b --- /dev/null +++ b/scripts/compare_runtime_turns.py @@ -0,0 +1,86 @@ +"""Compare identical runtime_turns benchmark sources built against two revisions.""" + +from __future__ import annotations + +import argparse +import json +import os +import platform +import random +import subprocess +from pathlib import Path + +from hotpath_lab import balanced_orders, paired_estimate, sha256, write_json + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--baseline", type=Path, required=True) + parser.add_argument("--candidate", type=Path, required=True) + parser.add_argument("--out", type=Path, required=True) + parser.add_argument("--blocks", type=int, default=24) + parser.add_argument("--iterations", type=int, default=3_000_000) + parser.add_argument("--seed", type=int, default=123) + parser.add_argument("--tasks", type=int, default=0) + args = parser.parse_args() + if args.iterations <= 0 or args.tasks < 0: + parser.error("--iterations must be positive and --tasks nonnegative") + binaries = { + name: getattr(args, name).resolve() for name in ("baseline", "candidate") + } + orders = balanced_orders(args.blocks, random.Random(args.seed)) + args.out.mkdir(parents=True, exist_ok=False) + hashes = {name: sha256(path) for name, path in binaries.items()} + write_json( + args.out / "plan.json", + { + "binaries": {name: str(path) for name, path in binaries.items()}, + "sha256": hashes, + "orders": orders, + "iterations": args.iterations, + "tasks": args.tasks, + "seed": args.seed, + "minimum_gain_pct": 1, + "minimum_sample_seconds": 0.25, + "platform": platform.platform(), + "affinity": sorted(os.sched_getaffinity(0)) + if hasattr(os, "sched_getaffinity") + else None, + "runner_sha256": sha256(Path(__file__)), + }, + ) + values = {name: [] for name in binaries} + with (args.out / "samples.jsonl").open("x") as samples: + for block, order in enumerate(orders): + for name in order: + result = json.loads( + subprocess.check_output( + [str(binaries[name]), str(args.iterations), str(args.tasks)], + text=True, + timeout=120, + ) + ) + assert result["iterations"] == args.iterations + assert result["tasks"] == args.tasks + values[name].append(result["seconds"]) + samples.write( + json.dumps({"block": block, "label": name, **result}) + "\n" + ) + samples.flush() + print(f"Completed paired block {block + 1}/{args.blocks}", flush=True) + if hashes != {name: sha256(path) for name, path in binaries.items()}: + raise ValueError("Benchmark binaries changed during measurement") + report = paired_estimate(**values, seed=args.seed) + report["min_observed_seconds"] = min(values["baseline"] + values["candidate"]) + reliable = args.blocks >= 12 and report["min_observed_seconds"] >= 0.25 + report["decision"] = ( + "timing_gate_passed" + if reliable and report["ci_pct"][1] < -1 + else "inconclusive" + ) + write_json(args.out / "report.json", report) + print(json.dumps(report, indent=2)) + + +if __name__ == "__main__": + main() diff --git a/src/lib.rs b/src/lib.rs index 0565d03..bbe4459 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,5 +1,4 @@ #![warn(missing_docs)] -#![cfg_attr(not(kani), feature(allocator_ext))] //! Native extension entry point and Rust interoperability API for `rsloop`. //! diff --git a/src/vibeio/allocator.rs b/src/vibeio/allocator.rs deleted file mode 100644 index 9ac9511..0000000 --- a/src/vibeio/allocator.rs +++ /dev/null @@ -1,371 +0,0 @@ -//! Experimental bounded allocator evaluated for event-loop collections. - -#[cfg(not(kani))] -mod native { - use std::alloc::{AllocError, Allocator, AllocatorClone, Layout, System}; - use std::ptr::NonNull; - use std::sync::{Arc, Mutex, MutexGuard}; - - const MIN_CACHED_SIZE: usize = 8; - const MAX_CACHED_SIZE: usize = 4 * 1024; - const MAX_CACHED_ALIGN: usize = 64; - const MAX_BLOCKS_PER_BIN: usize = 32; - const MAX_RETAINED_BYTES: usize = 256 * 1024; - const SIZE_BIN_COUNT: usize = 10; - const ALIGN_BIN_COUNT: usize = 7; - const BIN_COUNT: usize = SIZE_BIN_COUNT * ALIGN_BIN_COUNT; - - #[derive(Clone, Copy)] - struct CachedBlock { - pointer: NonNull, - layout: Layout, - } - - // SAFETY: a cached block has been invalidated by `deallocate`; its bytes - // are inaccessible until the pool transfers ownership from under its mutex. - unsafe impl Send for CachedBlock {} - - struct PoolState { - bins: Vec>, - retained_bytes: usize, - #[cfg(test)] - hits: usize, - #[cfg(test)] - misses: usize, - } - - impl PoolState { - fn new() -> Self { - let bins = (0..BIN_COUNT) - .map(|_| Vec::with_capacity(MAX_BLOCKS_PER_BIN)) - .collect(); - Self { - bins, - retained_bytes: 0, - #[cfg(test)] - hits: 0, - #[cfg(test)] - misses: 0, - } - } - } - - struct PoolInner { - state: Mutex, - } - - impl PoolInner { - #[inline] - fn lock(&self) -> MutexGuard<'_, PoolState> { - match self.state.lock() { - Ok(state) => state, - Err(poisoned) => poisoned.into_inner(), - } - } - } - - impl Drop for PoolInner { - fn drop(&mut self) { - let state = match self.state.get_mut() { - Ok(state) => state, - Err(poisoned) => poisoned.into_inner(), - }; - for bin in &mut state.bins { - for block in bin.drain(..) { - // SAFETY: cached blocks were allocated from `System` with this - // exact layout and were invalidated before entering the cache. - unsafe { - System.deallocate(block.pointer, block.layout); - } - } - } - } - } - - /// Cloneable allocator handle shared by one event loop and its runtimes. - #[derive(Clone)] - pub(crate) struct RuntimeAllocator { - inner: Arc, - } - - impl RuntimeAllocator { - pub(crate) fn new() -> Self { - Self { - inner: Arc::new(PoolInner { - state: Mutex::new(PoolState::new()), - }), - } - } - - #[inline] - fn cached_layout(layout: Layout) -> Option<(usize, Layout)> { - if layout.size() == 0 - || layout.size() > MAX_CACHED_SIZE - || layout.align() > MAX_CACHED_ALIGN - { - return None; - } - let size = layout.size().max(MIN_CACHED_SIZE).next_power_of_two(); - let size_index = - size.trailing_zeros() as usize - MIN_CACHED_SIZE.trailing_zeros() as usize; - let align_index = layout.align().trailing_zeros() as usize; - let index = size_index * ALIGN_BIN_COUNT + align_index; - // SAFETY: `size` is non-zero and `layout.align()` is a power of two. - let actual = unsafe { Layout::from_size_align_unchecked(size, layout.align()) }; - Some((index, actual)) - } - - #[inline] - fn actual_layout(layout: Layout) -> Layout { - Self::cached_layout(layout).map_or(layout, |(_, actual)| actual) - } - - #[cfg(test)] - fn stats(&self) -> (usize, usize, usize) { - let state = self.inner.lock(); - (state.hits, state.misses, state.retained_bytes) - } - } - - // SAFETY: live blocks remain owned by `System`; the mutex-protected cache - // contains only invalidated blocks, preserves their exact normalized - // layouts, and never touches them until a new allocation takes ownership. - unsafe impl Allocator for RuntimeAllocator { - fn allocate(&self, layout: Layout) -> Result, AllocError> { - let Some((index, actual)) = Self::cached_layout(layout) else { - return System.allocate(layout); - }; - - { - let mut state = self.inner.lock(); - if let Some(block) = state.bins[index].pop() { - state.retained_bytes = state.retained_bytes.saturating_sub(block.layout.size()); - #[cfg(test)] - { - state.hits += 1; - } - // SAFETY: the cache owns this non-null `System` allocation. - return Ok(NonNull::slice_from_raw_parts(block.pointer, actual.size())); - } - #[cfg(test)] - { - state.misses += 1; - } - } - - let block = System.allocate(actual)?; - // SAFETY: a successful allocation always returns a non-null data pointer. - let ptr = unsafe { NonNull::new_unchecked(block.as_ptr().cast::()) }; - Ok(NonNull::slice_from_raw_parts(ptr, actual.size())) - } - - unsafe fn deallocate(&self, ptr: NonNull, layout: Layout) { - let Some((index, actual)) = Self::cached_layout(layout) else { - // SAFETY: upheld by this method's caller. - unsafe { System.deallocate(ptr, layout) }; - return; - }; - - let mut state = self.inner.lock(); - if state.bins[index].len() < MAX_BLOCKS_PER_BIN - && state.retained_bytes.saturating_add(actual.size()) <= MAX_RETAINED_BYTES - { - state.bins[index].push(CachedBlock { - pointer: ptr, - layout: actual, - }); - state.retained_bytes += actual.size(); - return; - } - drop(state); - // SAFETY: the block was allocated with the normalized layout above. - unsafe { System.deallocate(ptr, actual) }; - } - - unsafe fn grow( - &self, - ptr: NonNull, - old_layout: Layout, - new_layout: Layout, - ) -> Result, AllocError> { - let old_actual = Self::actual_layout(old_layout); - let new_actual = Self::actual_layout(new_layout); - // SAFETY: every block returned by this allocator is System-backed - // with its normalized layout, and the caller upholds grow's rules. - unsafe { System.grow(ptr, old_actual, new_actual) } - } - - unsafe fn grow_zeroed( - &self, - ptr: NonNull, - old_layout: Layout, - new_layout: Layout, - ) -> Result, AllocError> { - let old_actual = Self::actual_layout(old_layout); - let new_actual = Self::actual_layout(new_layout); - // SAFETY: every block returned by this allocator is System-backed - // with its normalized layout, and the caller upholds grow's rules. - unsafe { System.grow_zeroed(ptr, old_actual, new_actual) } - } - - unsafe fn shrink( - &self, - ptr: NonNull, - old_layout: Layout, - new_layout: Layout, - ) -> Result, AllocError> { - let old_actual = Self::actual_layout(old_layout); - let new_actual = Self::actual_layout(new_layout); - // SAFETY: every block returned by this allocator is System-backed - // with its normalized layout, and the caller upholds shrink's rules. - unsafe { System.shrink(ptr, old_actual, new_actual) } - } - } - - // SAFETY: clones share the same pool, so either clone can free an allocation - // made by the other without invalidating allocations when a clone drops. - unsafe impl AllocatorClone for RuntimeAllocator {} - - #[cfg(test)] - mod tests { - use super::*; - - #[test] - fn cached_block_is_reused_and_bounded() { - let allocator = RuntimeAllocator::new(); - let layout = Layout::from_size_align(24, 8).unwrap(); - let first = allocator.allocate(layout).unwrap(); - let first_ptr = first.as_ptr().cast::(); - // SAFETY: `first` was allocated above with `layout`. - unsafe { allocator.deallocate(NonNull::new_unchecked(first_ptr), layout) }; - assert_eq!(allocator.stats().2, 32); - - let second = allocator.allocate(layout).unwrap(); - assert_eq!(second.as_ptr().cast::(), first_ptr); - assert_eq!(allocator.stats(), (1, 1, 0)); - // SAFETY: `second` is live and uses the same layout. - unsafe { - allocator.deallocate(NonNull::new_unchecked(second.as_ptr().cast::()), layout); - } - } - - #[test] - fn zeroed_allocation_clears_reused_memory() { - let allocator = RuntimeAllocator::new(); - let layout = Layout::from_size_align(16, 8).unwrap(); - let block = allocator.allocate(layout).unwrap(); - let ptr = block.as_ptr().cast::(); - // SAFETY: the allocation contains at least 16 writable bytes. - unsafe { - ptr.write_bytes(0xA5, 16); - allocator.deallocate(NonNull::new_unchecked(ptr), layout); - } - let zeroed = allocator.allocate_zeroed(layout).unwrap(); - let ptr = zeroed.as_ptr().cast::(); - // SAFETY: the returned allocation contains at least 16 initialized bytes. - let bytes = unsafe { std::slice::from_raw_parts(ptr, 16) }; - assert!(bytes.iter().all(|byte| *byte == 0)); - // SAFETY: `zeroed` is live and uses `layout`. - unsafe { allocator.deallocate(NonNull::new_unchecked(ptr), layout) }; - } - - #[test] - fn large_allocations_bypass_the_cache() { - let allocator = RuntimeAllocator::new(); - let layout = Layout::from_size_align(MAX_CACHED_SIZE + 1, 8).unwrap(); - let block = allocator.allocate(layout).unwrap(); - // SAFETY: `block` was allocated above with `layout`. - unsafe { - allocator.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); - } - assert_eq!(allocator.stats(), (0, 0, 0)); - } - - #[test] - fn clones_share_cached_blocks_across_threads() { - let allocator = RuntimeAllocator::new(); - let other = allocator.clone(); - let layout = Layout::from_size_align(64, 16).unwrap(); - let address = std::thread::spawn(move || { - let block = other.allocate(layout).unwrap(); - let address = block.as_ptr().cast::() as usize; - // SAFETY: the block was allocated and is deallocated on this thread. - unsafe { - other.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); - } - address - }) - .join() - .unwrap(); - assert_eq!(allocator.stats().2, 64); - let reused = allocator.allocate(layout).unwrap(); - assert_eq!(reused.as_ptr().cast::() as usize, address); - // SAFETY: `reused` is live and was allocated with `layout`. - unsafe { - allocator.deallocate(NonNull::new_unchecked(reused.as_ptr().cast::()), layout); - } - } - - #[test] - fn bin_limit_bounds_retained_memory() { - let allocator = RuntimeAllocator::new(); - let layout = Layout::from_size_align(MAX_CACHED_SIZE, 8).unwrap(); - let blocks = (0..(MAX_BLOCKS_PER_BIN + 8)) - .map(|_| allocator.allocate(layout).unwrap()) - .collect::>(); - for block in blocks { - // SAFETY: every block is live and was allocated with `layout`. - unsafe { - allocator - .deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); - } - } - assert_eq!(allocator.stats().2, MAX_BLOCKS_PER_BIN * MAX_CACHED_SIZE); - } - - #[test] - fn grow_preserves_initialized_bytes() { - let allocator = RuntimeAllocator::new(); - let old_layout = Layout::from_size_align(24, 8).unwrap(); - let new_layout = Layout::from_size_align(80, 8).unwrap(); - let block = allocator.allocate(old_layout).unwrap(); - let pointer = block.as_ptr().cast::(); - // SAFETY: the old allocation has 24 writable bytes and is live for grow. - let grown = unsafe { - for offset in 0..24 { - pointer.add(offset).write(offset as u8); - } - allocator - .grow(NonNull::new_unchecked(pointer), old_layout, new_layout) - .unwrap() - }; - let pointer = grown.as_ptr().cast::(); - // SAFETY: grow preserves the first 24 initialized bytes. - for offset in 0..24 { - // SAFETY: grow preserves the first 24 initialized bytes. - assert_eq!(unsafe { pointer.add(offset).read() }, offset as u8); - } - // SAFETY: `grown` is live and uses `new_layout`. - unsafe { - allocator.deallocate(NonNull::new_unchecked(pointer), new_layout); - } - } - - #[test] - fn poisoned_metadata_lock_is_recovered() { - let allocator = RuntimeAllocator::new(); - let poisoned = allocator.clone(); - let _ = std::panic::catch_unwind(move || { - let _guard = poisoned.inner.lock(); - panic!("poison allocator metadata for the recovery test"); - }); - let layout = Layout::from_size_align(32, 8).unwrap(); - let block = allocator.allocate(layout).unwrap(); - // SAFETY: `block` is live and uses `layout`. - unsafe { - allocator.deallocate(NonNull::new_unchecked(block.as_ptr().cast::()), layout); - } - assert_eq!(allocator.stats().2, 32); - } - } -} diff --git a/src/vibeio/batch_allocator.rs b/src/vibeio/batch_allocator.rs new file mode 100644 index 0000000..c2bb471 --- /dev/null +++ b/src/vibeio/batch_allocator.rs @@ -0,0 +1,355 @@ +//! One cached allocation for successive scheduler batches on an owner thread. + +#[cfg(all(feature = "scheduler-batch-cache", not(kani)))] +mod native { + use std::alloc::{AllocError, Allocator, Layout, System}; + use std::cell::Cell; + use std::ptr::NonNull; + + #[derive(Clone, Copy)] + struct Block { + ptr: NonNull, + layout: Layout, + } + + /// Not Sync: this cache belongs to the runtime's owner thread. + #[derive(Default)] + pub(crate) struct BatchAllocator { + cached: Cell>, + #[cfg(test)] + system_allocations: Cell, + } + + #[cfg(test)] + impl BatchAllocator { + pub(crate) fn system_allocations(&self) -> usize { + self.system_allocations.get() + } + } + + const MAX_RETAINED: usize = 16 * 1024; + + // SAFETY: all blocks are System allocations with exact layouts. Only + // deallocated blocks enter the cache; live blocks remain disjoint and are + // never inspected. Cell operations cannot unwind or invoke user code. + unsafe impl Allocator for BatchAllocator { + #[inline] + fn allocate(&self, layout: Layout) -> Result, AllocError> { + if let Some(block) = self.cached.get() + && block.layout == layout + { + self.cached.set(None); + return Ok(NonNull::slice_from_raw_parts(block.ptr, layout.size())); + } + #[cfg(test)] + self.system_allocations + .set(self.system_allocations.get().wrapping_add(1)); + System.allocate(layout) + } + + #[inline] + unsafe fn deallocate(&self, ptr: NonNull, layout: Layout) { + if layout.size() != 0 && layout.size() <= MAX_RETAINED && self.cached.get().is_none() { + self.cached.set(Some(Block { ptr, layout })); + } else { + // SAFETY: the caller supplies a live System-backed block with + // its exact layout; zero-sized blocks also delegate to System. + unsafe { System.deallocate(ptr, layout) }; + } + } + } + + impl Drop for BatchAllocator { + fn drop(&mut self) { + if let Some(block) = self.cached.take() { + // SAFETY: the cached block was returned by its previous owner; + // outstanding allocations are never freed by allocator drop. + unsafe { System.deallocate(block.ptr, block.layout) }; + } + } + } + + pub(crate) type Batch<'a, T> = Vec; + + #[inline] + pub(crate) fn batch(allocator: &BatchAllocator, capacity: usize) -> Batch<'_, T> { + Vec::with_capacity_in(capacity, allocator) + } + + #[cfg(test)] + mod tests { + use super::*; + use crate::vibeio::batch_allocator::drain_batch; + + #[test] + fn partial_drain_drops_remaining_elements_in_order_and_keeps_capacity() { + let allocator = BatchAllocator::default(); + let drops = std::cell::RefCell::new(Vec::new()); + struct Tracked<'a>(usize, &'a std::cell::RefCell>); + impl Drop for Tracked<'_> { + fn drop(&mut self) { + self.1.borrow_mut().push(self.0); + } + } + let mut values = batch(&allocator, 8); + values.extend((0..4).map(|i| Tracked(i, &drops))); + let address = values.as_ptr(); + let mut drain = drain_batch(&mut values); + assert_eq!(drain.len(), 4); + drop(drain.next().unwrap()); + assert_eq!(drain.size_hint(), (3, Some(3))); + drop(drain); + assert_eq!(*drops.borrow(), [0, 1, 2, 3]); + assert!(values.is_empty()); + assert_eq!(values.capacity(), 8); + values.push(Tracked(4, &drops)); + assert_eq!(values.as_ptr(), address); + drop(values); + assert_eq!(*drops.borrow(), [0, 1, 2, 3, 4]); + assert_eq!(allocator.system_allocations(), 1); + } + + #[test] + fn drain_cleanup_continues_after_element_destructor_panics() { + let allocator = BatchAllocator::default(); + let drops = Cell::new(0); + struct Tracked<'a>(bool, &'a Cell); + impl Drop for Tracked<'_> { + fn drop(&mut self) { + self.1.set(self.1.get() + 1); + assert!(!self.0, "element destructor panic"); + } + } + let mut values = batch(&allocator, 4); + values.extend([ + Tracked(false, &drops), + Tracked(true, &drops), + Tracked(false, &drops), + ]); + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let mut drain = drain_batch(&mut values); + drop(drain.next()); + drop(drain); + })); + assert!(result.is_err()); + assert_eq!(drops.get(), 3); + assert!(values.is_empty()); + drop(values); + assert_eq!(drops.get(), 3); + } + + #[test] + fn drain_handles_zero_sized_elements_empty_batches_and_forgetting() { + // Thread-local state avoids interference from parallel tests while + // allowing a genuinely zero-sized element with a destructor. + thread_local! { static DROPS: Cell = const { Cell::new(0) }; } + struct Zst; + impl Drop for Zst { + fn drop(&mut self) { + DROPS.with(|drops| drops.set(drops.get() + 1)); + } + } + DROPS.with(|drops| drops.set(0)); + let allocator = BatchAllocator::default(); + let mut values = batch(&allocator, 4); + values.extend([Zst, Zst, Zst]); + let mut drain = drain_batch(&mut values); + drop(drain.next()); + drop(drain); + DROPS.with(|drops| assert_eq!(drops.get(), 3)); + assert!(drain_batch(&mut values).next().is_none()); + values.push(Zst); + std::mem::forget(drain_batch(&mut values)); + assert!(values.is_empty()); + values.push(Zst); + drop(values); + DROPS.with(|drops| assert_eq!(drops.get(), 4)); + + let mut numbers = batch(&allocator, 4); + numbers.extend([1, 2, 3]); + std::mem::forget(drain_batch(&mut numbers)); + numbers.push(42); + assert_eq!(&numbers[..], &[42]); + } + + #[test] + fn successive_batches_reuse_storage_but_nested_batches_are_disjoint() { + let allocator = BatchAllocator::default(); + let mut first = batch::(&allocator, 256); + first.extend(0..256); + let address = first.as_ptr(); + let mut nested = batch::(&allocator, 256); + nested.push(42); + assert_ne!(address, nested.as_ptr()); + drop(first); + let reused = batch::(&allocator, 256); + assert_eq!(address, reused.as_ptr()); + assert_eq!(allocator.system_allocations(), 2); + assert!(reused.is_empty()); + assert_eq!(nested[0], 42); + } + + #[test] + fn growth_alignment_zeroing_and_retention_bound() { + let allocator = BatchAllocator::default(); + let mut values = batch::(&allocator, 1); + values.extend(0..4096); + assert!(values.iter().copied().eq(0..4096)); + drop(values); + assert!(allocator.cached.get().unwrap().layout.size() <= MAX_RETAINED); + let allocator = BatchAllocator::default(); + let layout = Layout::from_size_align(32, 64).unwrap(); + let block = allocator.allocate_zeroed(layout).unwrap(); + let ptr = block.cast::(); + assert_eq!(ptr.as_ptr().addr() % 64, 0); + // SAFETY: allocate_zeroed initialized all 32 bytes of this block. + unsafe { + assert!( + std::slice::from_raw_parts(ptr.as_ptr(), 32) + .iter() + .all(|b| *b == 0) + ); + allocator.deallocate(ptr, layout); + } + assert_eq!(allocator.cached.get().unwrap().layout, layout); + let empty = batch::<()>(&allocator, 256); + drop(empty); + } + + #[test] + fn unwinding_drops_elements_before_recycling_storage() { + let allocator = BatchAllocator::default(); + let drops = Cell::new(0); + struct Tracked<'a>(&'a Cell); + impl Drop for Tracked<'_> { + fn drop(&mut self) { + self.0.set(self.0.get() + 1); + } + } + let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let mut values = batch(&allocator, 256); + values.push(Tracked(&drops)); + panic!("exercise batch unwinding"); + })); + assert_eq!(drops.get(), 1); + assert!(allocator.cached.get().is_some()); + } + + #[test] + fn reused_storage_is_zeroed_and_mismatched_layouts_do_not_consume_cache() { + let allocator = BatchAllocator::default(); + let layout = Layout::from_size_align(32, 8).unwrap(); + let block = allocator.allocate(layout).unwrap().cast::(); + // SAFETY: this is a live 32-byte allocation made with layout. + unsafe { + block.as_ptr().write_bytes(0xA5, 32); + allocator.deallocate(block, layout); + } + let different = Layout::from_size_align(32, 64).unwrap(); + let aligned = allocator.allocate(different).unwrap().cast::(); + assert_eq!(allocator.cached.get().unwrap().ptr, block); + let zeroed = allocator.allocate_zeroed(layout).unwrap().cast::(); + assert_eq!(zeroed, block); + // SAFETY: both blocks remain live with the respective exact layouts; + // allocate_zeroed initialized every byte read below. + unsafe { + assert!( + std::slice::from_raw_parts(zeroed.as_ptr(), 32) + .iter() + .all(|b| *b == 0) + ); + allocator.deallocate(zeroed, layout); + allocator.deallocate(aligned, different); + } + } + + #[test] + fn zero_and_oversized_allocations_are_never_retained() { + let allocator = BatchAllocator::default(); + for size in [0, MAX_RETAINED + 1] { + let layout = Layout::from_size_align(size, 8).unwrap(); + let block = allocator.allocate(layout).unwrap().cast::(); + // SAFETY: this live allocation was returned for layout above. + unsafe { allocator.deallocate(block, layout) }; + assert!(allocator.cached.get().is_none()); + } + } + } +} + +#[cfg(all(feature = "scheduler-batch-cache", not(kani)))] +pub(crate) use native::{Batch, BatchAllocator, batch}; + +#[cfg(all(feature = "scheduler-batch-cache", not(kani)))] +pub(crate) use draining::drain_batch; + +#[cfg(all(feature = "scheduler-batch-cache", not(kani)))] +mod draining { + use super::Batch; + + /// Full-range drain using stabilized Vec primitives. The allocation stays with + /// its Vec, and the iterator owns the removed elements until they are yielded. + pub(crate) struct BatchDrain<'a, T> { + remaining: &'a mut [T], + } + + pub(crate) fn drain_batch<'a, T>(batch: &'a mut Batch<'_, T>) -> BatchDrain<'a, T> { + let len = batch.len(); + // SAFETY: the old len elements are initialized. Setting len to zero transfers + // their drop responsibility to the iterator, whose borrow prevents the Vec + // from moving/freeing its buffer until the iterator is dropped or forgotten. + unsafe { + batch.set_len(0); + BatchDrain { + remaining: std::slice::from_raw_parts_mut(batch.as_mut_ptr(), len), + } + } + } + + impl Iterator for BatchDrain<'_, T> { + type Item = T; + + #[inline] + fn next(&mut self) -> Option { + let (first, remaining) = std::mem::take(&mut self.remaining).split_first_mut()?; + self.remaining = remaining; + // SAFETY: first is initialized and has been removed from the iterator's + // remaining slice. Neither the iterator nor Vec will drop it again. + Some(unsafe { std::ptr::read(first) }) + } + + fn size_hint(&self) -> (usize, Option) { + (self.remaining.len(), Some(self.remaining.len())) + } + } + + impl ExactSizeIterator for BatchDrain<'_, T> {} + + impl Drop for BatchDrain<'_, T> { + fn drop(&mut self) { + // SAFETY: only initialized, unyielded elements remain. Slice drop glue + // also drops later elements if one destructor panics. + unsafe { std::ptr::drop_in_place(self.remaining) }; + } + } +} + +// Default builds preserve the original allocation/drain path. Kani's compiler +// also predates stabilization, so its scheduler proofs use this implementation. +#[cfg(any(kani, not(feature = "scheduler-batch-cache")))] +mod ordinary { + #[derive(Default)] + pub(crate) struct BatchAllocator {} + pub(crate) type Batch<'a, T> = Vec; + #[inline] + pub(crate) fn batch(_: &BatchAllocator, capacity: usize) -> Batch<'_, T> { + Vec::with_capacity(capacity) + } + #[inline] + pub(crate) fn drain_batch<'a, T>(batch: &'a mut Batch<'_, T>) -> std::vec::Drain<'a, T> { + batch.drain(..) + } +} + +#[cfg(any(kani, not(feature = "scheduler-batch-cache")))] +pub(crate) use ordinary::{Batch, BatchAllocator, batch, drain_batch}; diff --git a/src/vibeio/executor.rs b/src/vibeio/executor.rs index 4fe2430..d4f9191 100644 --- a/src/vibeio/executor.rs +++ b/src/vibeio/executor.rs @@ -28,6 +28,7 @@ use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::task::{Context, Poll, Wake, Waker}; +use crate::vibeio::batch_allocator::{Batch, BatchAllocator, batch, drain_batch}; use crossbeam_queue::SegQueue; use slab::Slab; @@ -455,6 +456,7 @@ pub(crate) struct RuntimeInner { /// See "Spawning and joining tasks" in `tools/vibeio-check/EXAMPLES.md`. pub struct Runtime { inner: Option>, + batch_allocator: BatchAllocator, } impl RuntimeInner { @@ -513,7 +515,7 @@ impl RuntimeInner { /// Drain ready tasks into the given batch. #[inline] - fn drain_ready(&self, batch: &mut Vec>, mut budget: usize) { + fn drain_ready(&self, batch: &mut Batch<'_, Rc>, mut budget: usize) { if budget != 0 { let slab = self.token_to_task.borrow(); while budget != 0 { @@ -611,6 +613,7 @@ impl Runtime { interrupt_pending: Arc::new(AtomicBool::new(false)), }); Runtime { + batch_allocator: BatchAllocator::default(), inner: Some(Rc::new(RuntimeInner { queue: ready_queue, next_task: Rc::new(RefCell::new(None)), @@ -701,7 +704,7 @@ impl Runtime { Arc::clone(&inner.remote_wake.interrupt_pending), ); let root_waker = root_notify.waker(); - let mut batch = Vec::with_capacity(inner.task_batch_size); + let mut batch = batch(&self.batch_allocator, inner.task_batch_size); loop { if root_notify.take_ready() { @@ -778,7 +781,7 @@ impl Runtime { } } - for task in batch.drain(..) { + for task in drain_batch(&mut batch) { let mut future_slot = task.future.borrow_mut(); if let Some(mut future) = future_slot.take() { drop(future_slot); @@ -845,6 +848,47 @@ impl Drop for Runtime { #[cfg(test)] mod tests { + #[cfg(all(feature = "scheduler-batch-cache", not(kani)))] + #[test] + fn repeated_runtime_turns_allocate_batch_storage_only_once() { + let runtime = super::Runtime::new(crate::vibeio::driver::AnyDriver::new_mock()); + for _ in 0..32 { + runtime.poll_once(); + let task = runtime.spawn(async { 77 }); + assert_eq!(runtime.block_on(task), 77); + } + let allocator = &runtime.batch_allocator; + assert_eq!(allocator.system_allocations(), 1); + + let panic = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + runtime.block_on(async { panic!("root future panic") }); + })); + assert!(panic.is_err()); + assert_eq!(runtime.block_on(async { 42 }), 42); + assert_eq!(allocator.system_allocations(), 1); + } + + #[cfg(all(feature = "scheduler-batch-cache", not(kani)))] + #[test] + fn batch_poll_panic_releases_unpolled_task_owners() { + let runtime = super::Runtime::new(crate::vibeio::driver::AnyDriver::new_mock()); + let first = runtime.spawn(async { panic!("task poll panic") }); + let second = runtime.spawn(std::future::pending::<()>()); + let second_task = second.state.borrow().task.clone(); + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + runtime.block_on(std::future::pending::<()>()); + })); + assert!(result.is_err()); + // Only the slab may retain the task after the abandoned batch unwinds. + assert_eq!(second_task.strong_count(), 1); + second.cancel(); + assert_eq!(runtime.block_on(async { 42 }), 42); + assert_eq!(runtime.batch_allocator.system_allocations(), 1); + drop(first); + drop(runtime); + assert!(second_task.upgrade().is_none()); + } + #[cfg(feature = "process")] #[test] fn concurrent_reaper_requests_share_one_channel() { @@ -881,7 +925,7 @@ mod tests { .map(|handle| handle.state.borrow().task.upgrade().unwrap()) .collect(); let inner = runtime.inner.as_ref().unwrap(); - let mut batch = Vec::new(); + let mut batch = batch(&runtime.batch_allocator, 0); inner.drain_ready(&mut batch, 0); assert!(batch.is_empty()); assert_eq!(inner.queue.borrow().len(), 3); diff --git a/src/vibeio/lib.rs b/src/vibeio/lib.rs index ed2fe30..8c3090f 100644 --- a/src/vibeio/lib.rs +++ b/src/vibeio/lib.rs @@ -49,10 +49,11 @@ //! - `pipe` - enables pipe support. //! - `stdio` - enables standard I/O support. //! - `splice` - enables splice support (Linux). +//! - `scheduler-batch-cache` - opts into scheduler batch allocation reuse; +//! benchmark your workload, since TCP workloads can regress. //! - `blocking-default` - enables the default blocking thread pool. -#[cfg(test)] -pub(crate) mod allocator; +mod batch_allocator; pub mod blocking; mod builder; mod driver; diff --git a/tools/vibeio-check/Cargo.toml b/tools/vibeio-check/Cargo.toml index 26cd5bf..d12ac54 100644 --- a/tools/vibeio-check/Cargo.toml +++ b/tools/vibeio-check/Cargo.toml @@ -35,6 +35,7 @@ windows-sys = { version = "0.61", features = ["Wdk_Foundation", "Wdk_Storage_Fil # Compile the actual embedded sources without PyO3/TLS native build dependencies. [features] default = [] +scheduler-batch-cache = [] blocking-default = ["dep:rusty_pool"] fs = ["dep:async-std"] pipe = [] diff --git a/tools/vibeio-check/examples/runtime_turns.rs b/tools/vibeio-check/examples/runtime_turns.rs new file mode 100644 index 0000000..8d5a22c --- /dev/null +++ b/tools/vibeio-check/examples/runtime_turns.rs @@ -0,0 +1,52 @@ +//! Scheduler entry/exit benchmark, without Python or socket traffic noise. +//! Run the same source against both revisions, in alternating fresh processes. + +use rsloop_vibeio_check::vibeio::RuntimeBuilder; +use std::hint::black_box; +use std::task::Poll; +use std::time::Instant; + +fn main() { + let mut args = std::env::args().skip(1); + let iterations: usize = args + .next() + .unwrap_or_else(|| "3000000".into()) + .parse() + .expect("positive iteration count"); + let tasks: usize = args + .next() + .unwrap_or_else(|| "0".into()) + .parse() + .expect("number of ready tasks"); + assert!(iterations > 0); + let runtime = RuntimeBuilder::new().rsloop_profile().build().unwrap(); + let _tasks: Vec<_> = (0..tasks) + .map(|_| { + runtime.spawn(std::future::poll_fn(|cx| { + cx.waker().wake_by_ref(); + Poll::<()>::Pending + })) + }) + .collect(); + for (measured, count) in [(false, 10_000), (true, iterations)] { + let started = Instant::now(); + for _ in 0..count { + let mut polled = false; + black_box(&runtime).block_on(std::future::poll_fn(move |cx| { + if polled { + Poll::Ready(()) + } else { + polled = true; + cx.waker().wake_by_ref(); + Poll::Pending + } + })); + } + if measured { + println!( + "{{\"iterations\":{count},\"tasks\":{tasks},\"seconds\":{}}}", + started.elapsed().as_secs_f64() + ); + } + } +} diff --git a/tools/vibeio-check/src/lib.rs b/tools/vibeio-check/src/lib.rs index 9a48fd4..69874f3 100644 --- a/tools/vibeio-check/src/lib.rs +++ b/tools/vibeio-check/src/lib.rs @@ -1,6 +1,5 @@ //! Build/test harness for rsloop's embedded runtime; no copied implementation. #![doc = include_str!("../EXAMPLES.md")] -#![feature(allocator_ext)] #[path = "../../../src/vibeio/lib.rs"] pub mod vibeio; From 5e4d070ce54ffd65d24b837246692cd612b5107d Mon Sep 17 00:00:00 2001 From: Yehor Smoliakov Date: Fri, 25 Sep 2026 11:25:25 +0000 Subject: [PATCH 3/3] Optimize rsloop callback and task binding hot paths Make the PyLoop wrapper immutable and use stack-backed vectorcall arguments for ordinary task options. Preserve factory and eager-start compatibility paths. Add regression coverage and document the passing seven-workload timing/RSS holdout and Python 3.10, 3.14, and free-threaded checks. --- benches/compare_event_loops.py | 32 +++++++ docs/development.md | 3 + docs/rsloop-bindings-performance.md | 121 ++++++++++++++++++++++++ scripts/hotpath_lab.py | 8 +- src/bindings/loop_api.rs | 4 +- src/bindings/loop_api/fast_callbacks.rs | 4 +- src/bindings/loop_api/tasks.rs | 47 ++++++--- tests/test_fast_callbacks.py | 81 ++++++++++++++++ tests/tooling/test_benchmark_loops.py | 25 +++++ tests/tooling/test_hotpath_lab.py | 9 ++ 10 files changed, 319 insertions(+), 15 deletions(-) create mode 100644 docs/rsloop-bindings-performance.md diff --git a/benches/compare_event_loops.py b/benches/compare_event_loops.py index e7ef3a5..46fd43c 100644 --- a/benches/compare_event_loops.py +++ b/benches/compare_event_loops.py @@ -3,6 +3,7 @@ import argparse import asyncio +import contextvars import ctypes import gc import importlib @@ -373,6 +374,37 @@ async def tiny_task() -> None: ) +async def bench_task_options( + loop_name: str, iterations: int, batch_size: int +) -> ChildResult: + """Exercise named, explicit-context Task construction (Python 3.11+).""" + + async def tiny_task() -> None: + await asyncio.sleep(0) + + loop = asyncio.get_running_loop() + context = contextvars.copy_context() + start = time.perf_counter() + remaining = iterations + while remaining > 0: + current_batch = min(batch_size, remaining) + tasks = [ + loop.create_task(tiny_task(), name="benchmark-task", context=context) + for _ in range(current_batch) + ] + await asyncio.gather(*tasks) + remaining -= current_batch + return ChildResult( + loop_name, + "task_options", + time.perf_counter() - start, + iterations, + baseline_rss_bytes=0, + peak_rss_bytes=0, + peak_rss_delta_bytes=0, + ) + + async def maybe_wait_closed(writer: asyncio.StreamWriter) -> None: wait_closed = getattr(writer, "wait_closed", None) if wait_closed is None: diff --git a/docs/development.md b/docs/development.md index 8d9530f..279d6f5 100644 --- a/docs/development.md +++ b/docs/development.md @@ -298,6 +298,9 @@ For isolated Rust/LLVM optimization experiments, runtime-only PGO training, LLM review packets, and randomized paired comparisons, see the [hot-path experiment harness](hotpath-lab.md). +For rsloop's callback/task binding optimizations, their Rust API implications, +and application-level measurements, see [binding hot paths](rsloop-bindings-performance.md). + ## Current state of the project This is still an alpha-stage project. diff --git a/docs/rsloop-bindings-performance.md b/docs/rsloop-bindings-performance.md new file mode 100644 index 0000000..1fb967a --- /dev/null +++ b/docs/rsloop-bindings-performance.md @@ -0,0 +1,121 @@ +# Rsloop binding hot paths + +This experiment targets rsloop itself, not the embedded vibeio scheduler. +The allocator-cache feature remains disabled in both benchmark builds. + +## Findings and implementation + +Code inspection covered callback scheduling and dispatch, task construction, +the loop-thread ready queue, and stream delivery. Several transport costs are +already addressed by pooled reads, ownership transfer, and cached protocol +callbacks. This change does not alter transport buffering, scheduler fairness, +or I/O routing. + +Two avoidable costs remained in rsloop's Python bindings: + +1. `PyLoop` stores only an `Arc` but used PyO3's mutable-object borrow + bookkeeping on calls. It is now a frozen Rust wrapper; `LoopCore` retains + its existing locks and atomics. Hot callback and task paths use shared + `get()` access. This follows PyO3's + [frozen-class model](https://pyo3.rs/main/class#frozen-classes-opting-out-of-interior-mutability). +2. Named/context task construction created a keyword dictionary, copied it, + and queried Python's running-loop function. Ordinary options now use the + existing cached vectorcall keyword tuples and an explicit loop argument. + A fixed four-pointer stack array also replaces the temporary Rust vector. + Custom task factories and explicit `eager_start` retain their existing path. + +Python subclass attributes, weak references, context capture, and reentrant +calls remain supported. The Rust API becomes shared-only for the wrapped +`PyLoop`: use `get()` or `borrow()`, not `borrow_mut()`, and change loop state +through `LoopCore`. This is a Rust source-compatibility consideration, not a +freeze of Python subclass dictionaries or of loop state. + +## Measurement protocol + +Host: Intel Core i9-9900K, Linux x86-64, CPython 3.14.7, pinned Rust +`nightly-2026-09-25`. Baseline: `a8b8404`. The local candidate snapshot +`be0170024921` contains the production binding changes; later edits add +benchmark coverage, compatibility-test guards, and this report. + +The extension artifacts are immutable, built with the same compiler and +release flags. Each paired block runs baseline and candidate in fresh +processes, with balanced randomized order and an unmeasured warm-up. Builds +and tests do not run during measurement. The primary workload is +`task_options`; every included workload also has a 3% elapsed-time and peak-RSS +regression budget. Confidence intervals are paired bootstrap intervals with +Bonferroni adjustment across the timing/RSS family. + +The new opt-in `task_options` benchmark creates explicitly named, +explicit-context tasks, each yielding once. It reuses one context for all +tasks. It is separate from ordinary `tasks`, and does not replace or alter +the six default benchmark workloads. It requires Python 3.11+. + +Training used 12 pairs, seed 2001, and callbacks/tasks/task-options/TCP. Its +balanced gate passed: callback time changed -2.98%, ordinary tasks -3.12%, +and task options -10.10%; the TCP interval overlapped zero. + +## Holdout results + +The unchanged candidate passed the **balanced timing/RSS gate** over 24 paired +blocks (seed 2002). Holdout changed task batch sizes, work volumes, payloads, +and connection counts. Negative percentages mean less elapsed time or RSS. +Intervals below use 99.643% confidence per metric, targeting 95% across all +14 timing/RSS metrics. + +| Workload | Elapsed-time change (interval) | Peak RSS change (interval) | +| --- | ---: | ---: | +| Callbacks | -3.26% [-5.65%, -0.69%] | +0.00% [-0.02%, +0.02%] | +| Ordinary tasks | -2.62% [-3.46%, -1.74%] | -0.08% [-0.35%, +0.19%] | +| Named/context tasks | -10.36% [-11.29%, -9.37%] | +0.34% [+0.19%, +0.51%] | +| TCP streams | +0.15% [-1.09%, +1.42%] | -0.15% [-0.44%, +0.18%] | +| HTTP keep-alive | -1.71% [-3.60%, +0.13%] | -0.05% [-0.49%, +0.35%] | +| Mixed streams | -1.76% [-3.34%, -0.21%] | -0.05% [-0.60%, +0.53%] | +| Bulk transfer | -0.39% [-2.03%, +1.07%] | -0.01% [-0.02%, +0.01%] | + +This establishes a gain on the measured Python scheduling workloads, not a +universal network speedup or a process-memory reduction. In particular, the +task-options workload had a small RSS increase despite fewer transient +allocations. Every timing and RSS interval upper bound stayed below +3%. + +## Correctness and compatibility + +- All-feature Rust tests: 407 passed; all-target/all-feature Clippy passed. +- CPython 3.14.7 development install: 209 passed, 2 skipped. +- CPython 3.14.0 free-threaded: 210 passed, 1 skipped, including parallel loops + and verification that importing the extension does not re-enable the GIL. +- CPython 3.10.18: 194 passed, 17 skipped. Explicit Task-context cases are + skipped because that stdlib option requires Python 3.11+. +- Repository tooling: 69 passed, including the new workload-selection and + benchmark-execution tests. + +Python commands selected `not tooling and not stress and not slow_network` +(75 deselected in the final suite). This is Linux validation, not a substitute +for the cross-platform CI matrix or the excluded stress/network suites. +The new regressions cover all name/context combinations, before and during +a running loop, with debug mode on and off, plus subclass attributes, weak +references, and reentrant scheduling. Existing factory/eager-start tests also +pass. No additional unsafe Rust was introduced; `LoopCore` synchronization +is unchanged. + +## Reproduction + +Reproduction (use fresh output directories): + +```bash +.venv/bin/python scripts/hotpath_lab.py build \ + --revision a8b8404 --out target/rsloop-bindings-baseline +# Build a committed snapshot containing this change as the candidate. +.venv/bin/python scripts/hotpath_lab.py build \ + --revision --out target/rsloop-bindings-candidate +.venv/bin/python scripts/hotpath_lab.py compare \ + --baseline target/rsloop-bindings-baseline \ + --candidate target/rsloop-bindings-candidate \ + --workloads callbacks,tasks,task_options,tcp_streams,http_keepalive,mixed_streams,bulk_transfer \ + --primary task_options --suite holdout --blocks 24 --seed 2002 \ + --out target/rsloop-bindings-holdout +``` + +Local source snapshots, extension hashes, execution plans, samples, and reports +are under `target/rsloop-bindings-*`. These ignored local artifacts are not +checked-in benchmark data. Results are host/workload-specific, not a guarantee +of equal gains in other applications. diff --git a/scripts/hotpath_lab.py b/scripts/hotpath_lab.py index e931daa..4e310ca 100644 --- a/scripts/hotpath_lab.py +++ b/scripts/hotpath_lab.py @@ -37,6 +37,7 @@ "bulk_transfer", "tls_http", "websocket_messages", + "task_options", ) DEFAULT_WORKLOADS = ",".join(WORKLOADS[:6]) @@ -342,7 +343,7 @@ def child(args) -> None: ): argv.extend(["--" + key.replace("_", "-"), str(config[key])]) matrix_args = None - if name not in WORKLOADS[:3]: + if name not in (*WORKLOADS[:3], "task_options"): sys.argv = argv matrix_args = matrix.parse_args() matrix.validate_args(matrix_args) @@ -360,6 +361,11 @@ def child(args) -> None: "rsloop", config["tasks"], config["task_batch_size"] ) result = asdict(small.run_with_loop("rsloop", coro)) + elif name == "task_options": + coro = small.bench_task_options( + "rsloop", config["tasks"], config["task_batch_size"] + ) + result = asdict(small.run_with_loop("rsloop", coro)) elif name == "tcp_streams": coro = small.bench_tcp_streams( "rsloop", config["tcp_roundtrips"], config["payload_size"] diff --git a/src/bindings/loop_api.rs b/src/bindings/loop_api.rs index 2eac811..e4f45be 100644 --- a/src/bindings/loop_api.rs +++ b/src/bindings/loop_api.rs @@ -49,8 +49,10 @@ use crate::engine::{CallbackKind, LoopCore, LoopCoreError}; /// to `Duration`/`Instant` (issue #48). const MAX_TIMER_DELAY_SECS: f64 = 100.0 * 365.0 * 24.0 * 60.0 * 60.0; -#[pyclass(subclass, module = "rsloop._loop", weakref)] +#[pyclass(subclass, module = "rsloop._loop", weakref, frozen)] /// Python-visible event loop; scheduling and lifecycle state live in `LoopCore`. +/// The Rust shell is immutable: `LoopCore` synchronizes its own state, so Python +/// method calls need no additional PyO3 borrow bookkeeping on this wrapper. pub struct PyLoop { /// Shared scheduling, lifecycle, and runtime state for this loop. pub core: Arc, diff --git a/src/bindings/loop_api/fast_callbacks.rs b/src/bindings/loop_api/fast_callbacks.rs index a81ff7b..ddaeced 100644 --- a/src/bindings/loop_api/fast_callbacks.rs +++ b/src/bindings/loop_api/fast_callbacks.rs @@ -59,7 +59,7 @@ unsafe fn schedule( ), }; let handle = - slf.borrow() + slf.get() .core .schedule_callback_args(py, kind, callback, callback_args, None)?; return Ok(handle.into_ptr()); @@ -125,7 +125,7 @@ unsafe fn schedule( _ => CallbackArgs::Many(PyTuple::new(py, (1..positional).map(value))?.unbind()), }; let handle = - slf.borrow() + slf.get() .core .schedule_callback_args(py, kind, callback, callback_args, context)?; Ok(handle.into_ptr()) diff --git a/src/bindings/loop_api/tasks.rs b/src/bindings/loop_api/tasks.rs index 75e9984..4d6f0e6 100644 --- a/src/bindings/loop_api/tasks.rs +++ b/src/bindings/loop_api/tasks.rs @@ -71,7 +71,7 @@ pub(crate) fn try_fast_create_future( let Ok(pyloop) = loop_obj.bind(py).cast_exact::() else { return Ok(None); }; - if !pyloop.borrow().core.on_runtime_thread() { + if !pyloop.get().core.on_runtime_thread() { return Ok(None); } create_asyncio_future_for_running_loop(py).map(Some) @@ -85,7 +85,7 @@ pub(crate) fn try_fast_create_task( let Ok(pyloop) = loop_obj.bind(py).cast_exact::() else { return Ok(None); }; - let core = &pyloop.borrow().core; + let core = &pyloop.get().core; if !core.on_runtime_thread() || core.has_task_factory() { return Ok(None); } @@ -103,14 +103,21 @@ fn create_asyncio_task_for_loop( { let name = name.as_ref(); let context = context.as_ref(); - let mut args = Vec::with_capacity(4); - args.push(coro.as_ptr()); - args.push(loop_obj.as_ptr()); + // Vectorcall borrows at most four pointers for this synchronous call; + // constructor keyword options do not need a heap allocation. + let mut args = [ + coro.as_ptr(), + loop_obj.as_ptr(), + std::ptr::null_mut(), + std::ptr::null_mut(), + ]; + let mut next = 2; if let Some(name) = name { - args.push(name.as_ptr()); + args[next] = name.as_ptr(); + next += 1; } if let Some(context) = context { - args.push(context.as_ptr()); + args[next] = context.as_ptr(); } let cls = asyncio_task_cls(py)?.as_ptr(); @@ -172,7 +179,7 @@ fn trim_task_source_traceback(py: Python<'_>, task: &Py) -> PyResult<()> } pub(super) fn create_future(slf: Py, py: Python<'_>) -> PyResult> { - if slf.borrow(py).core.on_runtime_thread() { + if slf.get().core.on_runtime_thread() { return create_asyncio_future_for_running_loop(py); } @@ -208,14 +215,13 @@ pub(super) fn create_task( .as_ref() .is_none_or(|kwargs| kwargs.bind(py).is_empty()); if bare { - let loop_ref = slf.borrow(py); + let loop_ref = slf.get(); if !loop_ref.core.has_task_factory() && loop_ref.core.on_runtime_thread() { - drop(loop_ref); return create_asyncio_task_for_running_loop(py, slf.bind(py).as_any(), coro); } } - let core = Arc::clone(&slf.borrow(py).core); + let core = Arc::clone(&slf.get().core); let task_kwarg_support = asyncio_task_kwarg_support(py)?; let extra_kwargs = kwargs .as_ref() @@ -253,6 +259,25 @@ pub(super) fn create_task( ))); } + // Ordinary Task options fit the cached vectorcall keyword layouts. Avoid + // constructing a kwargs dict and then copying it in the generic path. + // An explicit loop is correct both inside and outside a running loop and + // also avoids a Python _get_running_loop call here. Custom factories and + // eager_start retain their existing compatibility path below. + if task_factory.is_none() && eager_start.is_none() { + let created = create_asyncio_task_for_loop( + py, + &loop_obj, + coro, + name.filter(|_| task_kwarg_support.name), + context.filter(|_| task_kwarg_support.context), + )?; + if core.get_debug() { + trim_task_source_traceback(py, &created)?; + } + return Ok(created); + } + let task_kwargs = if has_kwargs || task_factory.is_some() { let task_kwargs = PyDict::new(py); if let Some(kwargs_in) = kwargs.as_ref() { diff --git a/tests/test_fast_callbacks.py b/tests/test_fast_callbacks.py index 903dff1..c2d74a1 100644 --- a/tests/test_fast_callbacks.py +++ b/tests/test_fast_callbacks.py @@ -1,7 +1,9 @@ from __future__ import annotations +import asyncio import contextvars import gc +import sys import threading import weakref from typing import Any, cast @@ -11,6 +13,85 @@ class TestFastCallback: + def test_loop_shell_preserves_subclass_state_and_reentrant_calls(self): + class SubLoop(rsloop.Loop): + pass + + loop = SubLoop() + reference = weakref.ref(loop) + loop.events = [] + + def callback(active_loop): + active_loop.events.append(active_loop.get_debug()) + active_loop.set_debug(True) + active_loop.call_soon(active_loop.events.append, active_loop.get_debug()) + active_loop.call_soon(active_loop.stop) + + try: + loop.call_soon(callback, loop) + loop.run_forever() + assert loop.events == [False, True] + assert reference() is loop + finally: + loop.close() + del loop + gc.collect() + assert reference() is None + + @pytest.mark.parametrize("running", [False, True]) + @pytest.mark.parametrize("named", [False, True]) + @pytest.mark.parametrize("debug", [False, True]) + @pytest.mark.parametrize( + "explicit_context", + [ + False, + pytest.param( + True, + marks=pytest.mark.skipif( + sys.version_info < (3, 11), + reason="asyncio.Task context requires Python 3.11+", + ), + ), + ], + ) + def test_task_constructor_option_layouts( + self, running, named, debug, explicit_context + ): + variable = contextvars.ContextVar("task_option_probe", default="ambient") + context = contextvars.Context() + context.run(variable.set, "explicit") + loop = rsloop.new_event_loop() + loop.set_debug(debug) + + async def probe(): + return variable.get(), asyncio.get_running_loop() + + def create(): + options = {} + if named: + options["name"] = "named-task" + if explicit_context: + options["context"] = context + task = loop.create_task(probe(), **options) + assert task.get_loop() is loop + if named: + assert task.get_name() == "named-task" + if debug: + assert task._source_traceback + return task + + async def while_running(): + return await create() + + try: + value, owning_loop = loop.run_until_complete( + while_running() if running else create() + ) + assert value == ("explicit" if explicit_context else "ambient") + assert owning_loop is loop + finally: + loop.close() + def test_arguments_keywords_and_context_capture(self): loop = rsloop.new_event_loop() events = [] diff --git a/tests/tooling/test_benchmark_loops.py b/tests/tooling/test_benchmark_loops.py index 2a854f2..0336b58 100644 --- a/tests/tooling/test_benchmark_loops.py +++ b/tests/tooling/test_benchmark_loops.py @@ -1,5 +1,6 @@ """Loop selection tests without optional native benchmark dependencies.""" +import asyncio import json import sys from pathlib import Path @@ -15,6 +16,30 @@ class TestBenchmarkLoop: + @pytest.mark.skipif(sys.version_info < (3, 11), reason="Task context needs 3.11+") + def test_task_options_workload_exercises_every_task(self, monkeypatch): + options = [] + + async def exercise(): + loop = asyncio.get_running_loop() + create_task = loop.create_task + + def capture(coro, **kwargs): + options.append(kwargs) + return create_task(coro, **kwargs) + + with monkeypatch.context() as patch: + patch.setattr(loop, "create_task", capture) + return await comparison.bench_task_options("asyncio", 5, 2) + + result = asyncio.run(exercise()) + assert result.operations == 5 + assert result.workload == "task_options" + assert result.seconds > 0 + assert len(options) == 5 + assert all(option["name"] == "benchmark-task" for option in options) + assert all(option["context"] is options[0]["context"] for option in options) + def test_python315_profiler_command_writes_native_thread_flamegraph( self, tmp_path, monkeypatch, mocker ): diff --git a/tests/tooling/test_hotpath_lab.py b/tests/tooling/test_hotpath_lab.py index 1a24a16..ea9e05d 100644 --- a/tests/tooling/test_hotpath_lab.py +++ b/tests/tooling/test_hotpath_lab.py @@ -29,6 +29,15 @@ def test_orders_are_balanced_and_deterministic(): lab.balanced_orders(7, random.Random(0)) +def test_task_options_is_opt_in_and_has_a_distinct_holdout_shape(): + assert lab.select_workloads("task_options") == ["task_options"] + assert "task_options" not in lab.DEFAULT_WORKLOADS.split(",") + training = lab.workload_config("task_options", "training") + holdout = lab.workload_config("task_options", "holdout") + assert training["tasks"] != holdout["tasks"] + assert training["task_batch_size"] != holdout["task_batch_size"] + + def test_paired_estimate_uses_ratios_not_unpaired_medians(): baseline = [1.0, 10.0, 100.0, 1000.0] candidate = [value * 0.9 for value in baseline]