From cf9eb2a5d9d2ea2c0333377a67cc9176b9dcb1de Mon Sep 17 00:00:00 2001 From: tison Date: Tue, 15 Sep 2026 15:13:37 +0800 Subject: [PATCH 1/4] refactor: use hashcrew CRC and remove Tokio runtime coupling --- ARCHITECTURE.md | 2 +- CHANGELOG.md | 4 + CONFIGURATION.md | 8 +- Cargo.lock | 14 ++- Cargo.toml | 10 +- README.md | 8 +- cache2/Cargo.toml | 5 +- cache2/ERRORS.md | 2 +- cache2/src/cache.rs | 128 ++++++++++++------------ cache2/src/checksum.rs | 43 ++++++-- cache2/src/config/mod.rs | 4 +- cache2/src/config/runtime.rs | 5 +- cache2/src/error.rs | 3 +- cache2/src/io/engine/mod.rs | 70 +++++++------ cache2/src/io/engine/tests.rs | 61 ++++++----- cache2/src/region/file_backend/mod.rs | 11 +- cache2/src/region/file_backend/tests.rs | 13 +-- cache2/src/region/reader.rs | 8 +- cache2/src/region/runtime/mod.rs | 22 ++-- examples/Cargo.toml | 2 +- examples/src/logforth.rs | 8 +- tests-integration/Cargo.toml | 1 + tests-integration/tests/cache.rs | 47 +++++---- 23 files changed, 267 insertions(+), 212 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 4ee1ee2..83d45de 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -114,7 +114,7 @@ Reads and writes use independent bounded engine pools. Reclaim has separate read ### Memory -The managed-memory limit covers the index mapping, heat bits, L1, append buffers, reclaim buffers, metadata, cache-owned thread stacks, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, Tokio, process overhead, and the kernel page cache. `CacheConfig::new` rejects invalid or insufficient memory budgets before file access; actual allocation can still fail during open. +The managed-memory limit covers the index mapping, heat bits, L1, append buffers, reclaim buffers, metadata, worker thread stacks, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, lifecycle thread stacks, timer infrastructure, the application executor, process overhead, and the kernel page cache. `CacheConfig::new` rejects invalid or insufficient memory budgets before file access; actual allocation can still fail during open. ### Storage path diff --git a/CHANGELOG.md b/CHANGELOG.md index ce53fe9..5a5e820 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### Breaking Changes + +- Cache futures now work with any executor and no longer require a Tokio runtime or timer driver. Replace `Cache::open_with_handle(path, config, handle)` with `Cache::open(path, config)`. Opening and closing use lifecycle threads; read admission, I/O deadlines, and cancellation retain their existing bounds, and closing still continues if its returned future is dropped. + ## v0.4.0 (2026-09-10) ### Breaking Changes diff --git a/CONFIGURATION.md b/CONFIGURATION.md index 7deb233..8bd881e 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -6,7 +6,7 @@ This guide explains the interactions between `StorageLayout` and `RuntimeOptions ## Configuration lifecycle -Fill in `StorageOptions` and `RuntimeOptions` using their public fields. `StorageOptions::build()` returns an immutable `StorageLayout` with checked geometry and disk accounting. `CacheConfig::new(storage, runtime)` checks the complete combination and retains the resolved runtime choices and memory requirements. Neither step opens files, starts workers, or needs Tokio. +Fill in `StorageOptions` and `RuntimeOptions` using their public fields. `StorageOptions::build()` returns an immutable `StorageLayout` with checked geometry and disk accounting. `CacheConfig::new(storage, runtime)` checks the complete combination and retains the resolved runtime choices and memory requirements. Neither step opens files, starts workers, or needs an async runtime. ```rust use cache2::{Cache, CacheConfig, L1EvictionPolicy, RuntimeOptions, StorageOptions}; @@ -39,7 +39,7 @@ let adjusted = CacheConfig::new(config.storage().clone(), RuntimeOptions { })?; ``` -`Cache::open` consumes one configuration and uses the current Tokio runtime when first polled. Use `Cache::open_with_handle(path, config, handle)` for an explicit runtime, which must have time enabled and outlive the cache. Clone the configuration before opening when it will be reused. Each open locks files and acquires its own resources; constructing a configuration reserves none of them. +`Cache::open` consumes one configuration and can be polled by any executor. File setup and recovery use a lifecycle thread; the cache is not bound to the executor that opened it. Clone the configuration before opening when it will be reused. Each open locks files and acquires its own resources; constructing a configuration reserves none of them. | Stage | Error operation | What can fail | |-------|-----------------|---------------| @@ -51,7 +51,7 @@ Recovery still validates persisted metadata against the selected layout, and rea ### Migrating from the builder API -Replace `StaticConfig` with `StorageOptions` and call `build()` once to obtain the layout. Replace `RuntimeConfig` setters with `RuntimeOptions` fields, then construct `CacheConfig::new(layout, options)`. Replace `CacheBuilder::open()` with `Cache::open(path, config)`, passing any explicit Tokio handle through `Cache::open_with_handle`. The standalone `validate()` method is removed; disk-usage queries now belong to `StorageLayout` and return `u64` directly. Use `ReadAdmission::Immediate` for the former zero timeout, or `ReadAdmission::Wait { timeout, max_waiters }` for bounded waiting. +Replace `StaticConfig` with `StorageOptions` and call `build()` once to obtain the layout. Replace `RuntimeConfig` setters with `RuntimeOptions` fields, then construct `CacheConfig::new(layout, options)`. Replace `CacheBuilder::open()` or `Cache::open_with_handle(path, config, handle)` with `Cache::open(path, config)`; no runtime handle is needed. The standalone `validate()` method is removed; disk-usage queries now belong to `StorageLayout` and return `u64` directly. Use `ReadAdmission::Immediate` for the former zero timeout, or `ReadAdmission::Wait { timeout, max_waiters }` for bounded waiting. ## Measure the workload envelope first @@ -87,7 +87,7 @@ This is a floor, not a recommended limit. Concurrent L2 reads allocate alignment `CacheConfig::new` rejects a combination whose fixed footprint cannot fit the managed-memory limit. A configuration that barely meets `minimum_memory_bytes()` can still produce read-memory misses or overload once concurrent transient buffers consume the remaining budget. -The managed-memory limit is not an RSS limit. Allocator metadata, Tokio, the application, mapped-file residency, and the kernel page cache are outside it. Buffered I/O can therefore use substantial kernel memory even when the C² managed-memory gauges remain below their limit. +The managed-memory limit is not an RSS limit. Allocator metadata, lifecycle thread stacks, timer infrastructure, the application executor, mapped-file residency, and the kernel page cache are outside it. Buffered I/O can therefore use substantial kernel memory even when the C² managed-memory gauges remain below their limit. ### High-impact interaction map diff --git a/Cargo.lock b/Cargo.lock index 393454c..d0fc704 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -94,6 +94,7 @@ version = "0.4.0" dependencies = [ "asyncband", "crc-fast", + "futures-timer", "hashcrew", "io-uring", "libc", @@ -251,12 +252,18 @@ dependencies = [ name = "examples" version = "0.0.0" dependencies = [ + "asyncband", "cache2", "log", "logforth", - "tokio", ] +[[package]] +name = "futures-timer" +version = "3.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "af43fadb8a98512d547e37b4e92e0ced13e205c061b87b4623eff01d918d6968" + [[package]] name = "generic-array" version = "0.14.7" @@ -281,9 +288,9 @@ dependencies = [ [[package]] name = "hashcrew" -version = "0.2.0" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5fd3730de410a9d5f2ca41072653e79059e2c6fd9cd87300de4becacb1a93154" +checksum = "f73782af6df9e45939f4e6206b646f3cbab68c2cc5c0c375dab045800cafb018" [[package]] name = "heck" @@ -605,6 +612,7 @@ dependencies = [ name = "tests-integration" version = "0.0.0" dependencies = [ + "asyncband", "cache2", "crc-fast", "tokio", diff --git a/Cargo.toml b/Cargo.toml index c5ec144..092c152 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,11 +28,17 @@ rust-version = "1.98" cache2 = { path = "cache2", version = "0.4.0" } # Crates.io dependencies -asyncband = { version = "0.7.1", features = ["barrier", "semaphore", "watch"] } +asyncband = { version = "0.7.1", features = [ + "barrier", + "oneshot", + "semaphore", + "watch", +] } cargo_metadata = { version = "0.23.1" } clap = { version = "4.6.5", features = ["derive"] } crc-fast = { version = "1.10.0", default-features = false, features = ["std"] } -hashcrew = { version = "0.2.0", features = ["std", "xxhash"] } +futures-timer = { version = "3.0.4" } +hashcrew = { version = "0.3.0", features = ["std", "crc", "xxhash"] } io-uring = { version = "0.7.14" } libc = { version = "0.2.189" } log = { version = "0.4.31", features = ["kv"] } diff --git a/README.md b/README.md index aa78223..e765ef3 100644 --- a/README.md +++ b/README.md @@ -34,7 +34,7 @@ async fn run() -> Result<(), Error> { } ``` -`open` uses the current Tokio runtime. Use `Cache::open_with_handle(path, config, handle)` to bind C² to another runtime; it must have time enabled and outlive the cache. +`Cache::open` and all other asynchronous operations work with any executor. File setup and shutdown use short-lived lifecycle threads; read deadlines use an independent shared timer. No Tokio runtime or timer driver is required. ## Semantics @@ -58,13 +58,13 @@ Public failures are `cache2::Error` values with an actionable `ErrorKind`, the f ### Lifecycle -`close_fast`, drop, and an unclean exit make the next open a cold start. `close_warm` publishes a clean recovery image for a warm start. Both close methods work through `Arc` without `Arc::try_unwrap`: the first close call immediately makes every shared handle inert, then fences accepted persistent work on Tokio's blocking pool. Prefer explicit async close because drop closes the cache synchronously. Retained handles may keep bounded in-memory resources allocated until they are dropped, but no longer admit public operations. +`close_fast`, drop, and an unclean exit make the next open a cold start. `close_warm` publishes a clean recovery image for a warm start. Both close methods work through `Arc` without `Arc::try_unwrap`: the first close call immediately makes every shared handle inert, then fences accepted persistent work on a lifecycle thread. Prefer explicit async close because drop closes the cache synchronously. Retained handles may keep bounded in-memory resources allocated until they are dropped, but no longer admit public operations. ## Configuration `StorageOptions` and `RuntimeOptions` are editable inputs. Build the storage options into a `StorageLayout`, then combine it with runtime options using `CacheConfig::new`. The resulting configuration is immutable and ready for `Cache::open(path, config)`. -Construction requires neither file access nor Tokio. Inspect `config.storage().peak_disk_bytes()` and `config.minimum_memory_bytes()` before opening; both queries reuse computed values. Clone a configuration to reuse it across paths or successive opens. File locks, device support, recovery, and actual allocations are checked when opening each instance. +Construction requires neither file access nor an async runtime. Inspect `config.storage().peak_disk_bytes()` and `config.minimum_memory_bytes()` before opening; both queries reuse computed values. Clone a configuration to reuse it across paths or successive opens. File locks, device support, recovery, and actual allocations are checked when opening each instance. ### Persistent layout @@ -98,7 +98,7 @@ C² supports 64-bit Linux and macOS. Buffered positioned I/O is available on bot C² accepts one data-file path. For multiple homogeneous SSDs, expose RAID0 or an equivalent striped block device below the filesystem. Losing any member discards the complete cache. -The managed-memory limit covers the index, L1, append and reclaim buffers, metadata, cache-owned threads, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, Tokio, process overhead, and the kernel page cache. +The managed-memory limit covers the index, L1, append and reclaim buffers, metadata, worker thread stacks, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, lifecycle thread stacks, timer infrastructure, the application executor, process overhead, and the kernel page cache. The on-disk format is versioned. During 0.x, deployments should expect cold starts across releases and monitor `Cache::startup_mode()`. diff --git a/cache2/Cargo.toml b/cache2/Cargo.toml index 1ce4596..5891bbf 100644 --- a/cache2/Cargo.toml +++ b/cache2/Cargo.toml @@ -32,12 +32,13 @@ rustdoc-args = ["--cfg", "docsrs"] [dependencies] asyncband = { workspace = true } -crc-fast = { workspace = true } +futures-timer = { workspace = true } hashcrew = { workspace = true } log = { workspace = true } -tokio = { workspace = true } [dev-dependencies] +asyncband = { workspace = true, features = ["blocking"] } +crc-fast = { workspace = true } quickcheck = { workspace = true } tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/cache2/ERRORS.md b/cache2/ERRORS.md index c30512c..5032028 100644 --- a/cache2/ERRORS.md +++ b/cache2/ERRORS.md @@ -109,4 +109,4 @@ Use `error.kind() == cache2::ErrorKind::Overloaded` for application policy. Use `StorageOptions::build` reports `BuildStorage`; `CacheConfig::new` reports `BuildConfig`. Correct invalid options or an insufficient managed-memory limit before trying again. Successful results retain the checked layout and resource requirements, so inspecting their disk and memory bounds cannot fail. -`Cache::open` and `Cache::open_with_handle` report `Open`. A valid configuration can still encounter a busy file, unavailable device support, an allocation failure, or an unavailable runtime. Handle these according to `ErrorKind` just like other lifecycle failures. +`Cache::open` reports `Open`. A valid configuration can still encounter a busy file, unavailable device support, an allocation failure, or a worker startup failure. Handle these according to `ErrorKind` just like other lifecycle failures. diff --git a/cache2/src/cache.rs b/cache2/src/cache.rs index 9baffaf..e1430a7 100644 --- a/cache2/src/cache.rs +++ b/cache2/src/cache.rs @@ -33,7 +33,7 @@ use std::time::Instant; use std::time::SystemTime; use std::time::UNIX_EPOCH; -use tokio::task::JoinError; +use asyncband::oneshot; use crate::config::CacheConfig; use crate::config::storage::KEY_HASH_SEED; @@ -106,7 +106,7 @@ impl Value { } } -/// An open cache bound to one Tokio runtime. +/// An open cache whose futures can run on any executor. /// /// Prefer [`Self::close_fast`] or [`Self::close_warm`] for async shutdown. /// Either close method can be called through a shared [`Arc`] and immediately @@ -121,7 +121,6 @@ pub struct Cache { startup: StartupMode, path: PathBuf, logical_disk_peak_bytes: u64, - tokio_handle: tokio::runtime::Handle, } impl fmt::Debug for Cache { @@ -134,65 +133,32 @@ impl fmt::Debug for Cache { } impl Cache { - /// Opens a configured cache using the current Tokio runtime at first poll. - /// The runtime must have time enabled and outlive the cache. File setup and - /// recovery run on its blocking pool. + /// Opens a configured cache without requiring a particular async runtime. + /// File setup and recovery run on a separate lifecycle thread. /// /// # Errors /// /// Returns [`ErrorOperation::Open`] for file locking, recovery, allocation, - /// device support, runtime binding, or worker startup failures. Configuration + /// device support, or worker startup failures. Configuration /// has already been checked by [`CacheConfig::new`]. pub async fn open(path: impl AsRef, config: CacheConfig) -> Result { - let handle = tokio::runtime::Handle::try_current().map_err(|error| { - from_io( - ErrorOperation::Open, - io::Error::new(io::ErrorKind::InvalidInput, error.to_string()), - ) - })?; - Self::open_with_handle(path, config, handle).await - } - - /// Opens on an explicit Tokio runtime, including from a caller without an - /// active runtime. The selected runtime must have time enabled and outlive - /// the cache. Each open independently locks files and acquires resources. - /// - /// # Errors - /// - /// Uses the same [`ErrorOperation::Open`] failures as [`Self::open`]. - pub async fn open_with_handle( - path: impl AsRef, - config: CacheConfig, - tokio_handle: tokio::runtime::Handle, - ) -> Result { let path = path.as_ref().to_path_buf(); - let cache_handle = tokio_handle.clone(); let started = Instant::now(); - let result = tokio_handle - .spawn_blocking(move || Self::open_blocking(path, config, cache_handle, started)) - .await - .map_err(|error| { - from_io( - ErrorOperation::Open, - blocking_task_error("cache open", error), - ) - })?; + let result = spawn_lifecycle("cache2-open", move || { + Self::open_blocking(path, config, started) + }) + .await; public_result(ErrorOperation::Open, result) } - fn open_blocking( - path: PathBuf, - config: CacheConfig, - tokio_handle: tokio::runtime::Handle, - started: Instant, - ) -> io::Result { + fn open_blocking(path: PathBuf, config: CacheConfig, started: Instant) -> io::Result { let capacity_bytes = config.storage().capacity_bytes(); let index_slots = config.storage().index_slots(); let index_bytes = u64::try_from(index_slots) .ok() .and_then(recovery_image_index_len) .unwrap_or(0); - let result = Self::open_blocking_inner(path.clone(), config, tokio_handle); + let result = Self::open_blocking_inner(path.clone(), config); match &result { Ok(cache) => { let startup = cache.startup_mode(); @@ -225,11 +191,7 @@ impl Cache { result } - fn open_blocking_inner( - path: PathBuf, - config: CacheConfig, - tokio_handle: tokio::runtime::Handle, - ) -> io::Result { + fn open_blocking_inner(path: PathBuf, config: CacheConfig) -> io::Result { let format_data = DataSuperblock { generation: 1, cache_uuid: next_persistent_id(), @@ -256,7 +218,6 @@ impl Cache { startup, path, logical_disk_peak_bytes, - tokio_handle, }) } /// Reports whether this open started empty or mapped a clean recovery image. @@ -341,9 +302,7 @@ impl Cache { } public_result( ErrorOperation::Get, - self.data_plane - .get_async(key.as_ref(), &self.tokio_handle) - .await, + self.data_plane.get_async(key.as_ref()).await, ) .map(|value| value.map(|inner| Value { inner })) } @@ -397,7 +356,7 @@ impl Cache { Ok(snapshot) } - /// Stops on Tokio's blocking pool and makes the next open a cold start. + /// Stops on a lifecycle thread and makes the next open a cold start. /// Calling this method immediately rejects new operations through every /// shared handle; unique [`Arc`] ownership is not required. The first close /// call wins and continues even if the returned future is dropped. @@ -411,7 +370,7 @@ impl Cache { self.close(false) } - /// Publishes a clean recovery image on Tokio's blocking pool. Calling this + /// Publishes a clean recovery image on a lifecycle thread. Calling this /// method immediately rejects new operations through every shared handle; /// unique [`Arc`] ownership is not required. Accepted mutations and their /// submitted writes are fenced before the image is frozen. The first close @@ -446,13 +405,12 @@ impl Cache { } else { (ErrorOperation::CloseFast, "fast") }; - let tokio_handle = self.tokio_handle.clone(); let close = (!self.closed.swap(true, Ordering::AcqRel)).then(|| { self.data_plane.start_close(); let owner = Arc::clone(&self.owner); let path = self.path.clone(); let started = Instant::now(); - tokio_handle.spawn_blocking(move || { + spawn_lifecycle("cache2-close", move || { let result = owner .lock() .map_err(|_| cache_lifecycle_poisoned()) @@ -469,10 +427,7 @@ impl Cache { }); async move { let result = match close { - Some(close) => match close.await { - Ok(result) => result, - Err(error) => Err(blocking_task_error("cache close", error)), - }, + Some(close) => close.await, None => Err(cache_closed_error()), }; public_result(operation, result) @@ -556,8 +511,25 @@ fn log_cache_close(path: &Path, mode: &'static str, elapsed: Duration, result: & } } -fn blocking_task_error(operation: &'static str, error: JoinError) -> io::Error { - io::Error::other(format!("{operation} task failed: {error}")) +fn spawn_lifecycle( + name: &'static str, + operation: impl FnOnce() -> io::Result + Send + 'static, +) -> impl Future> + Send { + let (sender, receiver) = oneshot::channel(); + // Starting before the returned future is polled preserves close's eager, + // cancellation-independent contract. There is at most one open and one + // close task per cache; request-path operations never spawn threads here. + let thread = std::thread::Builder::new() + .name(name.to_owned()) + .spawn(move || { + let _ = sender.send(operation()); + }); + async move { + drop(thread?); + receiver + .await + .map_err(|_| io::Error::other(format!("{name} task panicked")))? + } } fn sidecar_path(path: &Path, suffix: &str) -> PathBuf { @@ -590,6 +562,8 @@ fn next_persistent_id() -> PersistentId { #[cfg(test)] mod tests { + use asyncband::blocking::FutureExt as _; + use super::*; fn assert_send_sync() {} @@ -599,4 +573,30 @@ mod tests { assert_send_sync::(); assert_send_sync::(); } + + #[test] + fn lifecycle_work_continues_when_its_unpolled_future_is_dropped() { + let (release, wait) = std::sync::mpsc::channel(); + let (finished, completion) = std::sync::mpsc::channel(); + let work = spawn_lifecycle("cache2-test-close", move || { + wait.recv().unwrap(); + finished.send(()).unwrap(); + Ok(()) + }); + drop(work); + release.send(()).unwrap(); + completion.recv_timeout(Duration::from_secs(2)).unwrap(); + } + + #[test] + fn lifecycle_panics_are_reported_as_io_errors() { + let error = spawn_lifecycle("cache2-test-panic", || -> io::Result<()> { + panic!("injected lifecycle failure") + }) + .wait_timeout(Duration::from_secs(2)) + .expect("panicking lifecycle task did not wake its caller") + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::Other); + assert!(error.to_string().contains("cache2-test-panic")); + } } diff --git a/cache2/src/checksum.rs b/cache2/src/checksum.rs index 1bbb1d0..bd2ae88 100644 --- a/cache2/src/checksum.rs +++ b/cache2/src/checksum.rs @@ -18,9 +18,8 @@ //! retains a portable software fallback. This wrapper keeps the cache's codec //! API and checksum values independent of that implementation detail. -use crc_fast::CrcAlgorithm; -use crc_fast::Digest; -use crc_fast::crc32_iscsi; +use hashcrew::crc::Crc32Iscsi; +use hashcrew::crc::crc32_iscsi; /// Computes the standard CRC32C checksum of `bytes`. pub fn crc32c(bytes: &[u8]) -> u32 { @@ -30,13 +29,13 @@ pub fn crc32c(bytes: &[u8]) -> u32 { /// Incremental CRC32C state, useful for checksum a key and value without first joining them in a /// temporary allocation. pub struct Crc32c { - digest: Digest, + digest: Crc32Iscsi, } impl Crc32c { pub fn new() -> Self { Self { - digest: Digest::new(CrcAlgorithm::Crc32Iscsi), + digest: Crc32Iscsi::new(), } } @@ -45,7 +44,7 @@ impl Crc32c { } pub fn finish(self) -> u32 { - self.digest.finalize() as u32 + self.digest.digest() } } @@ -57,6 +56,7 @@ impl Default for Crc32c { #[cfg(test)] mod tests { + use crate::checksum::Crc32c; use crate::checksum::crc32c; #[test] @@ -64,4 +64,35 @@ mod tests { assert_eq!(crc32c(b"123456789"), 0xe306_9283); assert_eq!(crc32c(b""), 0); } + + #[test] + fn fragmented_checksums_match_the_previous_implementation() { + let bytes: Vec<_> = (0..65_544).map(|index| (index * 37) as u8).collect(); + for offset in [0, 1, 7] { + for len in [0, 1, 44, 48, 4092, 4096, 65_537] { + let input = &bytes[offset..offset + len]; + let expected = crc_fast::crc32_iscsi(input); + assert_eq!(crc32c(input), expected); + for split in [0, len.min(44), len.min(56), len / 2, len] { + let mut checksum = Crc32c::new(); + checksum.update(&input[..split]); + checksum.update(&[]); + checksum.update(&input[split..]); + assert_eq!(checksum.finish(), expected, "len={len}, split={split}"); + } + } + } + + // Record headers, index pages, and recovery pages zero their checksum + // field without concatenating the surrounding slices. + for (len, checksum_offset) in [(48, 44), (4096, 56), (4096, 4092)] { + let mut page = bytes[..len].to_vec(); + page[checksum_offset..checksum_offset + 4].fill(0); + let mut checksum = Crc32c::new(); + checksum.update(&page[..checksum_offset]); + checksum.update(&[0; 4]); + checksum.update(&page[checksum_offset + 4..]); + assert_eq!(checksum.finish(), crc_fast::crc32_iscsi(&page)); + } + } } diff --git a/cache2/src/config/mod.rs b/cache2/src/config/mod.rs index 6d3cd3a..96a512a 100644 --- a/cache2/src/config/mod.rs +++ b/cache2/src/config/mod.rs @@ -24,7 +24,7 @@ pub mod storage; /// /// Construction checks runtime settings against the storage layout and managed /// memory limit. It performs bounded calculations without opening files, starting -/// workers, or requiring Tokio. Resource queries reuse the computed results. +/// workers, or requiring an async runtime. Resource queries reuse the computed results. /// Clone a configuration to reuse it across paths or successive opens; each open /// acquires its own resources and can still fail on I/O or allocation. #[derive(Clone, Debug)] @@ -61,7 +61,7 @@ impl CacheConfig { /// /// Created by [`StorageOptions::build`](crate::StorageOptions::build). Changing the geometry or /// index size changes the disk identity, so an incompatible recovery image opens empty. -/// Layout construction neither reserves disk space nor requires a Tokio runtime. +/// Layout construction neither reserves disk space nor requires an async runtime. #[derive(Clone, Debug, Eq, PartialEq)] pub struct StorageLayout { geometry: DataGeometry, diff --git a/cache2/src/config/runtime.rs b/cache2/src/config/runtime.rs index ce309af..41feb1a 100644 --- a/cache2/src/config/runtime.rs +++ b/cache2/src/config/runtime.rs @@ -420,8 +420,9 @@ pub struct RuntimeOptions { /// Bounded L1 eviction policy. Defaults to CLOCK; S3-FIFO adds ghost metadata. pub l1_eviction_policy: L1EvictionPolicy, /// Aggregate cache-managed memory limit, defaulting to 1 GiB. Covers index - /// mappings, L1, buffers, metadata, queues, and cache threads. Allocator - /// overhead, Tokio, application memory, and the kernel page cache are outside it. + /// mappings, L1, buffers, metadata, queues, and worker stacks. Allocator + /// overhead, lifecycle thread stacks, timer infrastructure, application memory, + /// and the kernel page cache are outside it. pub managed_memory_limit_bytes: usize, /// Independently locked L1 shards, from 1 through 65536 (default 32). Powers /// of two give the cheapest routing; more shards require more metadata. diff --git a/cache2/src/error.rs b/cache2/src/error.rs index f8a7641..27c517e 100644 --- a/cache2/src/error.rs +++ b/cache2/src/error.rs @@ -74,8 +74,7 @@ pub enum ErrorOperation { BuildConfig, /// [`StorageOptions::build`](crate::StorageOptions::build). BuildStorage, - /// [`Cache::open`](crate::Cache::open) or - /// [`Cache::open_with_handle`](crate::Cache::open_with_handle). + /// [`Cache::open`](crate::Cache::open). Open, /// [`Cache::put`](crate::Cache::put). Put, diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index 07ad2d8..08ddacd 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -20,8 +20,10 @@ use std::fmt; use std::future::Future; +use std::future::poll_fn; use std::io; use std::pin::Pin; +use std::pin::pin; use std::sync::Arc; use std::sync::Condvar; use std::sync::Mutex; @@ -43,6 +45,7 @@ use std::time::Instant; use asyncband::semaphore::OwnedSemaphorePermit; use asyncband::semaphore::Semaphore; +use futures_timer::Delay; use crate::IoEngineConfig; #[cfg(unix)] @@ -788,26 +791,19 @@ impl BoundedIoRequest { pub async fn wait_async( self, engine: Arc, - tokio_handle: &tokio::runtime::Handle, ) -> Result { let mut request = AsyncRequestGuard::new(self.request, engine); - let deadline = tokio::time::Instant::from_std(self.deadline); - let completion = { - let _entered = tokio_handle.enter(); - tokio::time::timeout_at(deadline, request.request_mut()) - } - .await; + let completion = timeout_at(self.deadline, request.request_mut()).await; if let Ok(completion) = completion { request.disarm(); return Ok(completion); } let cancel_error = request.cancel().err(); - let completion = { - let _entered = tokio_handle.enter(); - tokio::time::timeout(self.cancel_grace, request.request_mut()) - } - .await; + let grace_deadline = Instant::now() + .checked_add(self.cancel_grace) + .unwrap_or_else(Instant::now); + let completion = timeout_at(grace_deadline, request.request_mut()).await; match completion { Ok(completion) => { request.disarm(); @@ -1081,23 +1077,16 @@ impl ReadSlotAdmission { Ok(permit) } - async fn acquire_until( - &self, - deadline: Instant, - tokio_handle: &tokio::runtime::Handle, - ) -> io::Result { + async fn acquire_until(&self, deadline: Instant) -> io::Result { self.ensure_open()?; let acquire = Arc::clone(&self.slots).acquire_owned(1); - { - let _entered = tokio_handle.enter(); - tokio::time::timeout_at(tokio::time::Instant::from_std(deadline), acquire) - } - .await - .map_err(|_| io::Error::new(io::ErrorKind::TimedOut, "L2 read wait deadline expired")) - .and_then(|permit| { - self.ensure_open()?; - Ok(permit) - }) + timeout_at(deadline, acquire) + .await + .map_err(|_| io::Error::new(io::ErrorKind::TimedOut, "L2 read wait deadline expired")) + .and_then(|permit| { + self.ensure_open()?; + Ok(permit) + }) } fn close(&self) { @@ -1107,6 +1096,25 @@ impl ReadSlotAdmission { } } +async fn timeout_at(deadline: Instant, future: F) -> Result { + let mut future = pin!(future); + let mut timer = None; + poll_fn(|context| { + // A ready completion wins even at the deadline. Delay registration is + // unnecessary when I/O or admission completed before the first poll. + if let Poll::Ready(output) = future.as_mut().poll(context) { + return Poll::Ready(Ok(output)); + } + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Poll::Ready(Err(())); + } + let timer = timer.get_or_insert_with(|| Delay::new(remaining)); + Pin::new(timer).poll(context).map(|()| Err(())) + }) + .await +} + struct ReadWaiterGuard<'a> { admission: &'a ReadSlotAdmission, } @@ -1118,11 +1126,7 @@ impl Drop for ReadWaiterGuard<'_> { } impl ReadSlotWaiter { - pub async fn reserve_until( - self, - deadline: Instant, - tokio_handle: &tokio::runtime::Handle, - ) -> io::Result { + pub async fn reserve_until(self, deadline: Instant) -> io::Result { let admission = self .shared .read_slot_admission @@ -1130,7 +1134,7 @@ impl ReadSlotWaiter { .ok_or_else(|| io::Error::other("async read admission is disabled"))?; let _waiter = admission.register_waiter(); self.shared.ensure_accepting()?; - let permit = admission.acquire_until(deadline, tokio_handle).await?; + let permit = admission.acquire_until(deadline).await?; self.shared.try_reserve_read_slot(Some(permit)) } } diff --git a/cache2/src/io/engine/tests.rs b/cache2/src/io/engine/tests.rs index a0fa67c..a97e08d 100644 --- a/cache2/src/io/engine/tests.rs +++ b/cache2/src/io/engine/tests.rs @@ -21,6 +21,8 @@ use std::sync::atomic::AtomicU64; use std::sync::mpsc; use std::time::Duration; +use asyncband::blocking::FutureExt as _; + use super::*; use crate::config::runtime::PosixIoConfig; use crate::io::backend::FileBackend; @@ -54,11 +56,8 @@ async fn spawn_registered_read_slot_waiter( expected_waiters: usize, ) -> tokio::task::JoinHandle> { let slot_waiter = engine.read_slot_waiter(); - let waiter = tokio::spawn(async move { - slot_waiter - .reserve_until(Instant::now() + timeout, &tokio::runtime::Handle::current()) - .await - }); + let waiter = + tokio::spawn(async move { slot_waiter.reserve_until(Instant::now() + timeout).await }); wait_for_registered_read_waiters(engine, expected_waiters).await; waiter } @@ -369,8 +368,8 @@ fn posix_engine_reports_progress_before_a_terminal_short_io_error() { engine.shutdown().unwrap(); } -#[tokio::test] -async fn async_request_is_woken_by_driver_completion() { +#[test] +fn async_request_is_woken_by_driver_completion_without_a_runtime() { let file = TestFile::new(); file.file().set_len(4096).unwrap(); let engine: Arc = Arc::new(BackendIoEngine::new(file.backend(), 2).unwrap()); @@ -382,8 +381,9 @@ async fn async_request_is_woken_by_driver_completion() { .unwrap(); let completion = request - .wait_async(Arc::clone(&engine), &tokio::runtime::Handle::current()) - .await + .wait_async(Arc::clone(&engine)) + .wait_timeout(Duration::from_secs(2)) + .expect("driver completion did not wake the executor") .unwrap(); assert!(matches!(completion.status, CompletionStatus::Completed)); @@ -402,11 +402,7 @@ async fn dropping_async_wait_requests_bounded_cancellation() { ) .unwrap(); let waiter_engine = Arc::clone(&engine); - let waiter = tokio::spawn(async move { - request - .wait_async(waiter_engine, &tokio::runtime::Handle::current()) - .await - }); + let waiter = tokio::spawn(async move { request.wait_async(waiter_engine).await }); tokio::task::yield_now().await; assert!(backend.wait_for_entered(1)); @@ -431,11 +427,7 @@ async fn read_slot_waits_for_cancelled_request_to_release_physical_capacity() { ) .unwrap(); let request_engine = Arc::clone(&engine); - let request_waiter = tokio::spawn(async move { - request - .wait_async(request_engine, &tokio::runtime::Handle::current()) - .await - }); + let request_waiter = tokio::spawn(async move { request.wait_async(request_engine).await }); tokio::task::yield_now().await; assert!(backend.wait_for_entered(1)); @@ -445,8 +437,7 @@ async fn read_slot_waits_for_cancelled_request_to_release_physical_capacity() { let slot_waiter = engine.read_slot_waiter(); let deadline = Instant::now() + Duration::from_secs(1); - let tokio_handle = tokio::runtime::Handle::current(); - let mut reservation = Box::pin(slot_waiter.reserve_until(deadline, &tokio_handle)); + let mut reservation = Box::pin(slot_waiter.reserve_until(deadline)); assert!( tokio::time::timeout(Duration::from_millis(20), reservation.as_mut()) .await @@ -585,8 +576,8 @@ async fn cancelled_queue_head_passes_priority_to_next_read() { engine.shutdown().unwrap(); } -#[tokio::test] -async fn async_read_deadline_keeps_other_slots_available() { +#[test] +fn async_read_deadline_keeps_other_slots_available_without_a_runtime() { let backend = Arc::new(BlockingBackend::default()); let engine: Arc = Arc::new(BackendIoEngine::new(backend.clone(), 2).unwrap()); let resources = resources(); @@ -600,8 +591,9 @@ async fn async_read_deadline_keeps_other_slots_available() { assert!(backend.wait_for_entered(1)); let timeout = request - .wait_async(Arc::clone(&engine), &tokio::runtime::Handle::current()) - .await + .wait_async(Arc::clone(&engine)) + .wait_timeout(Duration::from_secs(2)) + .expect("I/O deadline did not wake the executor") .unwrap_err(); let (error, buffer) = timeout.into_buffer(); assert_eq!(error.kind(), io::ErrorKind::TimedOut); @@ -613,6 +605,25 @@ async fn async_read_deadline_keeps_other_slots_available() { engine.shutdown().unwrap(); } +#[test] +fn read_slot_deadline_without_a_runtime_does_not_leak_capacity() { + let file = TestFile::new(); + let engine = BackendIoEngine::new_with_read_wait(file.backend(), 1).unwrap(); + let held = engine.try_reserve_read().unwrap(); + let reservation = engine + .read_slot_waiter() + .reserve_until(Instant::now() + Duration::from_millis(20)) + .wait_timeout(Duration::from_secs(2)) + .expect("admission deadline did not wake the executor"); + match reservation { + Ok(_) => panic!("a read was admitted without physical capacity"), + Err(error) => assert_eq!(error.kind(), io::ErrorKind::TimedOut), + } + drop(held); + drop(engine.try_reserve_read().unwrap()); + engine.shutdown().unwrap(); +} + #[cfg(unix)] #[test] fn posix_engine_routes_only_aligned_record_io_to_direct() { diff --git a/cache2/src/region/file_backend/mod.rs b/cache2/src/region/file_backend/mod.rs index 93eed41..d6590fd 100644 --- a/cache2/src/region/file_backend/mod.rs +++ b/cache2/src/region/file_backend/mod.rs @@ -262,15 +262,8 @@ impl RegionStore> { } #[cfg(test)] - async fn get_value_async( - &self, - key: &[u8], - tokio_handle: &tokio::runtime::Handle, - ) -> io::Result> { - self.runtime()? - .data_plane()? - .get_async(key, tokio_handle) - .await + async fn get_value_async(&self, key: &[u8]) -> io::Result> { + self.runtime()?.data_plane()?.get_async(key).await } #[cfg(test)] diff --git a/cache2/src/region/file_backend/tests.rs b/cache2/src/region/file_backend/tests.rs index d14e2ee..e7e854a 100644 --- a/cache2/src/region/file_backend/tests.rs +++ b/cache2/src/region/file_backend/tests.rs @@ -438,7 +438,7 @@ fn configured_read_wait_is_bounded_and_cancel_safe() { let plane = store.data_plane_handle().unwrap(); let mut slots: Vec<_> = (0..2).map(|_| plane.reserve_read_slot_for_test()).collect(); tokio_runtime.block_on(async { - let mut cancelled = Box::pin(store.get_value_async(b"queued-read", tokio_runtime.handle())); + let mut cancelled = Box::pin(store.get_value_async(b"queued-read")); assert_pending( cancelled.as_mut(), "saturated read must enter the wait queue", @@ -447,7 +447,7 @@ fn configured_read_wait_is_bounded_and_cancel_safe() { drop(cancelled); }); let value = tokio_runtime.block_on(async { - let mut waiting = Box::pin(store.get_value_async(b"queued-read", tokio_runtime.handle())); + let mut waiting = Box::pin(store.get_value_async(b"queued-read")); assert_pending( waiting.as_mut(), "a cancelled read must release its wait-queue permit", @@ -461,16 +461,13 @@ fn configured_read_wait_is_bounded_and_cancel_safe() { drop(value); let blocked: Vec<_> = (0..2).map(|_| plane.reserve_read_slot_for_test()).collect(); tokio_runtime.block_on(async { - let mut waiting = Box::pin(store.get_value_async(b"queued-read", tokio_runtime.handle())); + let mut waiting = Box::pin(store.get_value_async(b"queued-read")); assert_pending( waiting.as_mut(), "saturated read must enter the bounded wait queue", ) .await; - let queue_full = match store - .get_value_async(b"queued-read", tokio_runtime.handle()) - .await - { + let queue_full = match store.get_value_async(b"queued-read").await { Err(error) => error, Ok(_) => panic!("a second saturated read must not enter a full wait queue"), }; @@ -521,7 +518,7 @@ fn queued_l2_read_does_not_pin_warm_close() { let slot = plane.reserve_read_slot_for_test(); tokio_runtime.block_on(async { - let mut waiting = Box::pin(plane.get_async(b"queued-close", tokio_runtime.handle())); + let mut waiting = Box::pin(plane.get_async(b"queued-close")); assert_pending(waiting.as_mut(), "saturated read must enter the wait queue").await; store.close_warm().unwrap(); diff --git a/cache2/src/region/reader.rs b/cache2/src/region/reader.rs index 792aa77..d6d5e84 100644 --- a/cache2/src/region/reader.rs +++ b/cache2/src/region/reader.rs @@ -98,17 +98,13 @@ impl PendingRead { Self::finish(plan, request_id, completion) } - pub async fn wait_async( - self, - engine: Arc, - tokio_handle: &tokio::runtime::Handle, - ) -> ReadCompletion { + pub async fn wait_async(self, engine: Arc) -> ReadCompletion { let Self { plan, request_id, request, } = self; - let completion = request.wait_async(engine, tokio_handle).await; + let completion = request.wait_async(engine).await; Self::finish(plan, request_id, completion) } diff --git a/cache2/src/region/runtime/mod.rs b/cache2/src/region/runtime/mod.rs index 182155c..2dc7a54 100644 --- a/cache2/src/region/runtime/mod.rs +++ b/cache2/src/region/runtime/mod.rs @@ -358,7 +358,7 @@ impl PendingGet { } } - async fn wait_async(self, tokio_handle: &tokio::runtime::Handle) -> CompletedGet { + async fn wait_async(self) -> CompletedGet { let Self { engine, read, @@ -366,7 +366,7 @@ impl PendingGet { hash, } = self; CompletedGet { - read: read.wait_async(engine, tokio_handle).await, + read: read.wait_async(engine).await, read_token, hash, } @@ -374,7 +374,7 @@ impl PendingGet { } impl WaitingGet { - async fn reserve_async(self, tokio_handle: &tokio::runtime::Handle) -> io::Result { + async fn reserve_async(self) -> io::Result { let Self { engine, slot_waiter, @@ -384,7 +384,7 @@ impl WaitingGet { deadline, waiter_permit, } = self; - let slot = slot_waiter.reserve_until(deadline, tokio_handle).await?; + let slot = slot_waiter.reserve_until(deadline).await?; drop(waiter_permit); Ok(ReservedGet { engine, @@ -913,19 +913,13 @@ impl RegionDataPlane { } } - pub async fn get_async( - &self, - key: &[u8], - tokio_handle: &tokio::runtime::Handle, - ) -> io::Result> { + pub async fn get_async(&self, key: &[u8]) -> io::Result> { match self.prepare_get(key)? { PreparedGet::Complete(value) => Ok(value), - PreparedGet::Pending(pending) => { - self.finish_get(pending.wait_async(tokio_handle).await, key) - } + PreparedGet::Pending(pending) => self.finish_get(pending.wait_async().await, key), PreparedGet::Waiting(waiting) => { let wait_started = self.config.statistics.then(Instant::now); - let reserved = waiting.reserve_async(tokio_handle).await; + let reserved = waiting.reserve_async().await; if let Some(wait_started) = wait_started { self.metrics.record_read_wait(wait_started.elapsed()); } @@ -933,7 +927,7 @@ impl RegionDataPlane { let Some(pending) = self.submit_reserved_get(reserved)? else { return Ok(None); }; - self.finish_get(pending.wait_async(tokio_handle).await, key) + self.finish_get(pending.wait_async().await, key) } } } diff --git a/examples/Cargo.toml b/examples/Cargo.toml index 729cf58..b223a6a 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -21,10 +21,10 @@ license.workspace = true rust-version.workspace = true [dependencies] +asyncband = { workspace = true, features = ["blocking"] } cache2 = { workspace = true } log = { workspace = true } logforth = { workspace = true } -tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } [[example]] name = "logforth" diff --git a/examples/src/logforth.rs b/examples/src/logforth.rs index 29bc62d..a286e42 100644 --- a/examples/src/logforth.rs +++ b/examples/src/logforth.rs @@ -15,6 +15,7 @@ use std::env; use std::io; +use asyncband::blocking::FutureExt as _; use cache2::Cache; use cache2::CacheConfig; use cache2::RuntimeOptions; @@ -24,8 +25,7 @@ use logforth::bridge::log::LogBridge; use logforth::filter::rustlog::RustLogFilterBuilder; use logforth::layout::JsonLayout; -#[tokio::main(flavor = "multi_thread")] -async fn main() -> io::Result<()> { +fn main() -> io::Result<()> { init_logforth(); let path = env::args_os().nth(1).ok_or_else(|| { @@ -40,8 +40,8 @@ async fn main() -> io::Result<()> { } .build()?; let config = CacheConfig::new(storage, RuntimeOptions::default())?; - let cache = Cache::open(path, config).await?; - cache.close_warm().await?; + let cache = Cache::open(path, config).block_on()?; + cache.close_warm().block_on()?; Ok(()) } diff --git a/tests-integration/Cargo.toml b/tests-integration/Cargo.toml index d20b38a..8650ad6 100644 --- a/tests-integration/Cargo.toml +++ b/tests-integration/Cargo.toml @@ -25,6 +25,7 @@ default = [] io-uring = ["cache2/io-uring"] [dev-dependencies] +asyncband = { workspace = true, features = ["blocking"] } cache2 = { workspace = true } crc-fast = { workspace = true } tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/tests-integration/tests/cache.rs b/tests-integration/tests/cache.rs index 039f99f..e0b4409 100644 --- a/tests-integration/tests/cache.rs +++ b/tests-integration/tests/cache.rs @@ -31,6 +31,7 @@ use std::sync::atomic::Ordering; use std::time::Duration; use std::time::Instant; +use asyncband::blocking::FutureExt as _; use cache2::Cache; use cache2::CacheConfig; use cache2::CacheHealth; @@ -206,16 +207,8 @@ async fn completed_reclaim_snapshot(cache: &Cache) -> DetailedCacheSnapshot { } #[test] -fn explicit_tokio_handle_works_from_a_runtime_without_time_enabled() { - let files = TestCache::new("explicit-tokio-handle"); - let cache_runtime = tokio::runtime::Builder::new_multi_thread() - .worker_threads(2) - .enable_time() - .build() - .unwrap(); - let caller_runtime = tokio::runtime::Builder::new_current_thread() - .build() - .unwrap(); +fn cache_works_without_an_async_runtime() { + let files = TestCache::new("without-runtime"); let config = test_config(1); let minimum_memory_bytes = config.minimum_memory_bytes(); let config = CacheConfig::new( @@ -228,11 +221,8 @@ fn explicit_tokio_handle_works_from_a_runtime_without_time_enabled() { .unwrap(); files.assert_absent(); - caller_runtime.block_on(async { - let cache = - Cache::open_with_handle(&files.data, config.clone(), cache_runtime.handle().clone()) - .await - .unwrap(); + async { + let cache = Cache::open(&files.data, config.clone()).await.unwrap(); assert_eq!( cache.snapshot().unwrap().managed_memory_limit_bytes, minimum_memory_bytes @@ -241,16 +231,35 @@ fn explicit_tokio_handle_works_from_a_runtime_without_time_enabled() { cache.drain().await.unwrap(); cache.close_warm().await.unwrap(); - let reopened = Cache::open_with_handle(&files.data, config, cache_runtime.handle().clone()) - .await - .unwrap(); + let reopened = Cache::open(&files.data, config).await.unwrap(); assert_eq!(reopened.startup_mode(), StartupMode::Warm); let value = reopened.get("key").await.unwrap().unwrap(); assert_eq!(value.tier(), CacheTier::L2); assert_eq!(value.as_ref(), b"value"); drop(value); reopened.close_fast().await.unwrap(); - }); + } + .block_on(); +} + +#[test] +fn cache_can_outlive_its_opening_executor() { + let files = TestCache::new("outlive-executor"); + let runtime = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let cache = runtime + .block_on(Cache::open(&files.data, test_config(1))) + .unwrap(); + drop(runtime); + + cache.put_l2("key", "value").unwrap(); + cache.drain().block_on().unwrap(); + assert_eq!( + cache.get("key").block_on().unwrap().unwrap().as_ref(), + b"value" + ); + cache.close_warm().block_on().unwrap(); } #[tokio::test] From 4aaa66e8a2dc8b93fb0fb25648350ef5a7e3ef73 Mon Sep 17 00:00:00 2001 From: tison Date: Tue, 15 Sep 2026 15:22:25 +0800 Subject: [PATCH 2/4] perf: reduce runtime-independent timer overhead --- Cargo.lock | 124 ++++++++++++++++++++++++++++++++++-- Cargo.toml | 2 +- cache2/Cargo.toml | 2 +- cache2/src/cache.rs | 3 +- cache2/src/io/engine/mod.rs | 15 ++--- 5 files changed, 128 insertions(+), 18 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d0fc704..af84b0d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -58,12 +58,36 @@ version = "1.0.104" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "330a5ed07fa54e4702c9d6c4174f74427fc0ef6e214bbd677ae50a5099946470" +[[package]] +name = "async-io" +version = "2.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "456b8a8feb6f42d237746d4b3e9a178494627745c3c56c6ea55d92ba50d026fc" +dependencies = [ + "autocfg", + "cfg-if", + "concurrent-queue", + "futures-io", + "futures-lite", + "parking", + "polling", + "rustix", + "slab", + "windows-sys", +] + [[package]] name = "asyncband" version = "0.7.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2e52766975a4f080528a898235c51e82e65df9db713419067b5368040eeb5659" +[[package]] +name = "autocfg" +version = "1.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" + [[package]] name = "benchmarks" version = "0.0.0" @@ -92,9 +116,9 @@ checksum = "b588b76d00fde79687d7646a9b5bdf3cc0f655e0bbd080335a95d7e96f3587da" name = "cache2" version = "0.4.0" dependencies = [ + "async-io", "asyncband", "crc-fast", - "futures-timer", "hashcrew", "io-uring", "libc", @@ -188,6 +212,15 @@ version = "1.0.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" +[[package]] +name = "concurrent-queue" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crc-fast" version = "1.10.0" @@ -198,6 +231,12 @@ dependencies = [ "spin", ] +[[package]] +name = "crossbeam-utils" +version = "0.8.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a31eee39dddec8330830986fcd7625edb5a24ec90ea038215273bbc3adb08ac6" + [[package]] name = "crypto-common" version = "0.1.7" @@ -248,6 +287,16 @@ dependencies = [ "crypto-common", ] +[[package]] +name = "errno" +version = "0.3.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" +dependencies = [ + "libc", + "windows-sys", +] + [[package]] name = "examples" version = "0.0.0" @@ -259,10 +308,26 @@ dependencies = [ ] [[package]] -name = "futures-timer" -version = "3.0.4" +name = "futures-core" +version = "0.3.34" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "92d699e522242e69e3003b94ecc1f960f3a5e015aa7c5d7486e65ad01dd94f5e" + +[[package]] +name = "futures-io" +version = "0.3.34" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "af43fadb8a98512d547e37b4e92e0ced13e205c061b87b4623eff01d918d6968" +checksum = "53c0fa8157de1303bfffdaa1cc2a673bfffb60102f76b0ef4441659124373fed" + +[[package]] +name = "futures-lite" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f78e10609fe0e0b3f4157ffab1876319b5b0db102a2c60dc4626306dc46b44ad" +dependencies = [ + "futures-core", + "pin-project-lite", +] [[package]] name = "generic-array" @@ -298,6 +363,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hermit-abi" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e17592d60ebacc7d5e169f4663c5f84f9161cc90328abcfe8456f41e4dfcb284" + [[package]] name = "io-uring" version = "0.7.14" @@ -380,6 +451,12 @@ version = "0.2.189" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3eaf3ede3fee6db1a4c2ee091bf8a8b4dccdc6d17f656fb07896ee72867612f2" +[[package]] +name = "linux-raw-sys" +version = "0.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" + [[package]] name = "log" version = "0.4.31" @@ -451,12 +528,32 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" +[[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + [[package]] name = "pin-project-lite" version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" +[[package]] +name = "polling" +version = "3.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5d0e4f59085d47d8241c88ead0f274e8a0cb551f3625263c05eb8dd897c34218" +dependencies = [ + "cfg-if", + "concurrent-queue", + "hermit-abi", + "pin-project-lite", + "rustix", + "windows-sys", +] + [[package]] name = "portable-atomic" version = "1.15.0" @@ -521,6 +618,19 @@ version = "0.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" +[[package]] +name = "rustix" +version = "1.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" +dependencies = [ + "bitflags 2.13.1", + "errno", + "libc", + "linux-raw-sys", + "windows-sys", +] + [[package]] name = "semver" version = "1.0.28" @@ -574,6 +684,12 @@ dependencies = [ "zmij", ] +[[package]] +name = "slab" +version = "0.4.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" + [[package]] name = "spin" version = "0.10.1" diff --git a/Cargo.toml b/Cargo.toml index 092c152..d0473b3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,6 +28,7 @@ rust-version = "1.98" cache2 = { path = "cache2", version = "0.4.0" } # Crates.io dependencies +async-io = { version = "2.6.0" } asyncband = { version = "0.7.1", features = [ "barrier", "oneshot", @@ -37,7 +38,6 @@ asyncband = { version = "0.7.1", features = [ cargo_metadata = { version = "0.23.1" } clap = { version = "4.6.5", features = ["derive"] } crc-fast = { version = "1.10.0", default-features = false, features = ["std"] } -futures-timer = { version = "3.0.4" } hashcrew = { version = "0.3.0", features = ["std", "crc", "xxhash"] } io-uring = { version = "0.7.14" } libc = { version = "0.2.189" } diff --git a/cache2/Cargo.toml b/cache2/Cargo.toml index 5891bbf..d80410b 100644 --- a/cache2/Cargo.toml +++ b/cache2/Cargo.toml @@ -31,8 +31,8 @@ all-features = true rustdoc-args = ["--cfg", "docsrs"] [dependencies] +async-io = { workspace = true } asyncband = { workspace = true } -futures-timer = { workspace = true } hashcrew = { workspace = true } log = { workspace = true } diff --git a/cache2/src/cache.rs b/cache2/src/cache.rs index e1430a7..e273f8d 100644 --- a/cache2/src/cache.rs +++ b/cache2/src/cache.rs @@ -593,8 +593,7 @@ mod tests { let error = spawn_lifecycle("cache2-test-panic", || -> io::Result<()> { panic!("injected lifecycle failure") }) - .wait_timeout(Duration::from_secs(2)) - .expect("panicking lifecycle task did not wake its caller") + .block_on() .unwrap_err(); assert_eq!(error.kind(), io::ErrorKind::Other); assert!(error.to_string().contains("cache2-test-panic")); diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index 08ddacd..f7a04aa 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -43,9 +43,9 @@ use std::thread::JoinHandle; use std::time::Duration; use std::time::Instant; +use async_io::Timer; use asyncband::semaphore::OwnedSemaphorePermit; use asyncband::semaphore::Semaphore; -use futures_timer::Delay; use crate::IoEngineConfig; #[cfg(unix)] @@ -1098,19 +1098,14 @@ impl ReadSlotAdmission { async fn timeout_at(deadline: Instant, future: F) -> Result { let mut future = pin!(future); - let mut timer = None; + let mut timer = Timer::at(deadline); poll_fn(|context| { - // A ready completion wins even at the deadline. Delay registration is - // unnecessary when I/O or admission completed before the first poll. + // A ready completion wins even at the deadline. The timer registers + // only when polled, after I/O or admission has returned Pending. if let Poll::Ready(output) = future.as_mut().poll(context) { return Poll::Ready(Ok(output)); } - let remaining = deadline.saturating_duration_since(Instant::now()); - if remaining.is_zero() { - return Poll::Ready(Err(())); - } - let timer = timer.get_or_insert_with(|| Delay::new(remaining)); - Pin::new(timer).poll(context).map(|()| Err(())) + Pin::new(&mut timer).poll(context).map(|_| Err(())) }) .await } From 3f875e33cd32affcee5a44ca40aca67359120f85 Mon Sep 17 00:00:00 2001 From: tison Date: Tue, 15 Sep 2026 15:28:47 +0800 Subject: [PATCH 3/4] fix: make lifecycle completion visible to ThreadSanitizer --- Cargo.toml | 2 +- cache2/src/cache.rs | 9 ++++++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index d0473b3..90d1ad2 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -31,7 +31,7 @@ cache2 = { path = "cache2", version = "0.4.0" } async-io = { version = "2.6.0" } asyncband = { version = "0.7.1", features = [ "barrier", - "oneshot", + "mpsc", "semaphore", "watch", ] } diff --git a/cache2/src/cache.rs b/cache2/src/cache.rs index e273f8d..e7db636 100644 --- a/cache2/src/cache.rs +++ b/cache2/src/cache.rs @@ -33,7 +33,7 @@ use std::time::Instant; use std::time::SystemTime; use std::time::UNIX_EPOCH; -use asyncband::oneshot; +use asyncband::mpsc; use crate::config::CacheConfig; use crate::config::storage::KEY_HASH_SEED; @@ -515,18 +515,21 @@ fn spawn_lifecycle( name: &'static str, operation: impl FnOnce() -> io::Result + Send + 'static, ) -> impl Future> + Send { - let (sender, receiver) = oneshot::channel(); + // AsyncBand 0.7's oneshot uses atomic fences unsupported by ThreadSanitizer. + // A bounded mpsc channel keeps the lifecycle handoff observable to the checker. + let (sender, mut receiver) = mpsc::bounded(1); // Starting before the returned future is polled preserves close's eager, // cancellation-independent contract. There is at most one open and one // close task per cache; request-path operations never spawn threads here. let thread = std::thread::Builder::new() .name(name.to_owned()) .spawn(move || { - let _ = sender.send(operation()); + let _ = sender.try_send(operation()); }); async move { drop(thread?); receiver + .recv() .await .map_err(|_| io::Error::other(format!("{name} task panicked")))? } From 45638839d56b07def69d3fa7e2056d74cdc6082b Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 18 Sep 2026 21:29:47 +0800 Subject: [PATCH 4/4] refactor: retain only the hashcrew CRC upgrade Integrate the current main branch and keep the CRC32C implementation and compatibility tests. Defer the async runtime changes while preserving the existing PR history. --- Cargo.lock | 4 ++-- Cargo.toml | 2 +- cache2/Cargo.toml | 2 +- cache2/src/checksum.rs | 43 ++++++++++++++++++++++++++++++++++++------ 4 files changed, 41 insertions(+), 10 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 6c334b1..e3534d2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -281,9 +281,9 @@ dependencies = [ [[package]] name = "hashcrew" -version = "0.2.0" +version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5fd3730de410a9d5f2ca41072653e79059e2c6fd9cd87300de4becacb1a93154" +checksum = "f73782af6df9e45939f4e6206b646f3cbab68c2cc5c0c375dab045800cafb018" [[package]] name = "heck" diff --git a/Cargo.toml b/Cargo.toml index 9542fd2..644b5ae 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -32,7 +32,7 @@ asyncband = { version = "0.7.1", features = ["barrier", "semaphore", "watch"] } cargo_metadata = { version = "0.23.1" } clap = { version = "4.6.5", features = ["derive"] } crc-fast = { version = "1.10.0", default-features = false, features = ["std"] } -hashcrew = { version = "0.2.0", features = ["std", "xxhash"] } +hashcrew = { version = "0.3.0", features = ["std", "crc", "xxhash"] } io-uring = { version = "0.7.14" } libc = { version = "0.2.189" } log = { version = "0.4.31", features = ["kv"] } diff --git a/cache2/Cargo.toml b/cache2/Cargo.toml index 6245cfe..bcbbbad 100644 --- a/cache2/Cargo.toml +++ b/cache2/Cargo.toml @@ -32,12 +32,12 @@ rustdoc-args = ["--cfg", "docsrs"] [dependencies] asyncband = { workspace = true } -crc-fast = { workspace = true } hashcrew = { workspace = true } log = { workspace = true } tokio = { workspace = true } [dev-dependencies] +crc-fast = { workspace = true } quickcheck = { workspace = true } tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } diff --git a/cache2/src/checksum.rs b/cache2/src/checksum.rs index 1bbb1d0..bd2ae88 100644 --- a/cache2/src/checksum.rs +++ b/cache2/src/checksum.rs @@ -18,9 +18,8 @@ //! retains a portable software fallback. This wrapper keeps the cache's codec //! API and checksum values independent of that implementation detail. -use crc_fast::CrcAlgorithm; -use crc_fast::Digest; -use crc_fast::crc32_iscsi; +use hashcrew::crc::Crc32Iscsi; +use hashcrew::crc::crc32_iscsi; /// Computes the standard CRC32C checksum of `bytes`. pub fn crc32c(bytes: &[u8]) -> u32 { @@ -30,13 +29,13 @@ pub fn crc32c(bytes: &[u8]) -> u32 { /// Incremental CRC32C state, useful for checksum a key and value without first joining them in a /// temporary allocation. pub struct Crc32c { - digest: Digest, + digest: Crc32Iscsi, } impl Crc32c { pub fn new() -> Self { Self { - digest: Digest::new(CrcAlgorithm::Crc32Iscsi), + digest: Crc32Iscsi::new(), } } @@ -45,7 +44,7 @@ impl Crc32c { } pub fn finish(self) -> u32 { - self.digest.finalize() as u32 + self.digest.digest() } } @@ -57,6 +56,7 @@ impl Default for Crc32c { #[cfg(test)] mod tests { + use crate::checksum::Crc32c; use crate::checksum::crc32c; #[test] @@ -64,4 +64,35 @@ mod tests { assert_eq!(crc32c(b"123456789"), 0xe306_9283); assert_eq!(crc32c(b""), 0); } + + #[test] + fn fragmented_checksums_match_the_previous_implementation() { + let bytes: Vec<_> = (0..65_544).map(|index| (index * 37) as u8).collect(); + for offset in [0, 1, 7] { + for len in [0, 1, 44, 48, 4092, 4096, 65_537] { + let input = &bytes[offset..offset + len]; + let expected = crc_fast::crc32_iscsi(input); + assert_eq!(crc32c(input), expected); + for split in [0, len.min(44), len.min(56), len / 2, len] { + let mut checksum = Crc32c::new(); + checksum.update(&input[..split]); + checksum.update(&[]); + checksum.update(&input[split..]); + assert_eq!(checksum.finish(), expected, "len={len}, split={split}"); + } + } + } + + // Record headers, index pages, and recovery pages zero their checksum + // field without concatenating the surrounding slices. + for (len, checksum_offset) in [(48, 44), (4096, 56), (4096, 4092)] { + let mut page = bytes[..len].to_vec(); + page[checksum_offset..checksum_offset + 4].fill(0); + let mut checksum = Crc32c::new(); + checksum.update(&page[..checksum_offset]); + checksum.update(&[0; 4]); + checksum.update(&page[checksum_offset + 4..]); + assert_eq!(checksum.finish(), crc_fast::crc32_iscsi(&page)); + } + } }