From db4ea0ff2e70b3ce580ed0e5bdff8f01877744e7 Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 21 Sep 2026 21:34:43 +0800 Subject: [PATCH 1/4] refactor: bind buffer ownership to memory reservations Store aligned allocations directly in BufferLease and reuse MemoryReservation for individual buffers and aggregate staging charges. Field drop order frees the allocation before releasing its charge; allocation errors release the reservation automatically. Remove the optional allocation, BufferOwner enum, custom lease destructor, and redundant deallocation state while preserving alignment and memory-limit behavior. Validation: managed-memory and append-staging unit tests passed. --- ARCHITECTURE.md | 2 + cache2/src/managed_memory.rs | 149 +++++++++++------------------------ cache2/src/region/staging.rs | 8 +- 3 files changed, 54 insertions(+), 105 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index ccdffa1..2cb7ced 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -152,6 +152,8 @@ See [adaptive fill admission](CONFIGURATION.md#adaptive-fill-admission) for rate 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. +`MemoryReservation` holds a charge until its allocation is released. `BufferLease` owns an aligned allocation directly and drops it before returning an optional individual charge; append staging instead keeps one aggregate reservation for its fixed buffers and record arrays. A failed buffer allocation releases its reservation automatically. + ### Storage path C² owns one logical data path. Multi-device deployments stripe below the filesystem with RAID0 or an equivalent layer. Request routing, recovery identity, and descriptor count stay independent of device topology. diff --git a/cache2/src/managed_memory.rs b/cache2/src/managed_memory.rs index 3b3dab9..b100930 100644 --- a/cache2/src/managed_memory.rs +++ b/cache2/src/managed_memory.rs @@ -58,15 +58,14 @@ pub struct ManagedMemory { memory: Arc, } -/// A fixed runtime allocation charged to the same hard memory limit as the -/// request pools. The owner keeps this guard for exactly as long as the -/// associated bounded structure exists. -pub struct RuntimeMemoryReservation { +/// A charge against the cache-wide memory limit. Keep this guard until the +/// associated allocation has been released. +pub struct MemoryReservation { memory: Arc, bytes: usize, } -impl Drop for RuntimeMemoryReservation { +impl Drop for MemoryReservation { fn drop(&mut self) { self.memory.release(self.bytes); } @@ -87,14 +86,11 @@ impl ManagedMemory { Ok(Self { memory }) } - pub fn reserve_runtime_memory( - &self, - bytes: usize, - ) -> Result { + pub fn reserve(&self, bytes: usize) -> Result { if !self.memory.try_reserve(bytes) { return Err(ManagedMemoryError::Allocation); } - Ok(RuntimeMemoryReservation { + Ok(MemoryReservation { memory: Arc::clone(&self.memory), bytes, }) @@ -104,7 +100,16 @@ impl ManagedMemory { /// cache-wide hard memory limit. The caller maps failure to either a /// fail-open miss or an explicit bounded-wait overload. pub fn try_read_buffer(&self, length: usize) -> Option { - BufferLease::try_standalone(length, Arc::clone(&self.memory)) + let capacity = align_up(length, BUFFER_ALIGNMENT)?; + if capacity == 0 || capacity > isize::MAX as usize { + return None; + } + let reservation = self.reserve(capacity).ok()?; + let buffer = AlignedBuffer::try_new(capacity)?; + Some(BufferLease { + buffer, + _reservation: Some(reservation), + }) } pub fn snapshot(&self) -> ManagedMemorySnapshot { @@ -125,13 +130,10 @@ pub struct ManagedMemorySnapshot { } pub struct BufferLease { - owner: BufferOwner, - buffer: Option, -} - -enum BufferOwner { - Fixed, - Standalone { memory: Arc }, + // Fields drop in declaration order: free the allocation before returning + // its budget. Fixed staging buffers use their owner's aggregate charge. + buffer: AlignedBuffer, + _reservation: Option, } impl BufferLease { @@ -141,44 +143,17 @@ impl BufferLease { "fixed buffer size must be a non-zero 4096-byte multiple", )); } - let ptr = allocate_buffer(length).ok_or(ManagedMemoryError::Allocation)?; - let mut buffer = AlignedBuffer { - ptr, - capacity: length, - initialized: 0, - }; + let mut buffer = AlignedBuffer::try_new(length).ok_or(ManagedMemoryError::Allocation)?; buffer.prepare_zeroed(length); Ok(Self { - owner: BufferOwner::Fixed, - buffer: Some(buffer), - }) - } - - fn try_standalone(length: usize, memory: Arc) -> Option { - let maximum = align_up(length, BUFFER_ALIGNMENT)?; - if maximum == 0 || maximum > isize::MAX as usize { - return None; - } - if !memory.try_reserve(maximum) { - return None; - } - let Some(ptr) = allocate_buffer(maximum) else { - memory.release(maximum); - return None; - }; - Some(Self { - owner: BufferOwner::Standalone { memory }, - buffer: Some(AlignedBuffer { - ptr, - capacity: maximum, - initialized: 0, - }), + buffer, + _reservation: None, }) } #[cfg(test)] pub fn prepare(&mut self, length: usize) -> Result<&mut [u8], ()> { - let buffer = self.buffer.as_mut().expect("buffer lease owns a buffer"); + let buffer = &mut self.buffer; if length > buffer.capacity { return Err(()); } @@ -193,7 +168,7 @@ impl BufferLease { /// Fresh capacity is zeroed before it is exposed as initialized bytes. #[cfg(test)] fn grow_preserving(&mut self, length: usize) -> Result<&mut [u8], ()> { - let buffer = self.buffer.as_mut().expect("buffer lease owns a buffer"); + let buffer = &mut self.buffer; if length > buffer.capacity { return Err(()); } @@ -202,7 +177,7 @@ impl BufferLease { } pub fn prepared(&self, length: usize) -> Result<&[u8], ()> { - let buffer = self.buffer.as_ref().ok_or(())?; + let buffer = &self.buffer; if length > buffer.initialized { return Err(()); } @@ -212,7 +187,7 @@ impl BufferLease { } pub fn prepared_mut(&mut self, length: usize) -> Result<&mut [u8], ()> { - let buffer = self.buffer.as_mut().ok_or(())?; + let buffer = &mut self.buffer; if length > buffer.initialized { return Err(()); } @@ -220,13 +195,11 @@ impl BufferLease { } pub fn has_capacity(&self, length: usize) -> bool { - self.buffer - .as_ref() - .is_some_and(|buffer| length <= buffer.capacity) + length <= self.buffer.capacity } pub fn read_target(&self, length: usize) -> Result<*mut u8, ()> { - let buffer = self.buffer.as_ref().ok_or(())?; + let buffer = &self.buffer; if length > buffer.capacity { return Err(()); } @@ -234,7 +207,7 @@ impl BufferLease { } pub fn mark_initialized(&mut self, length: usize) -> Result<(), ()> { - let buffer = self.buffer.as_mut().ok_or(())?; + let buffer = &mut self.buffer; if length > buffer.capacity { return Err(()); } @@ -244,26 +217,7 @@ impl BufferLease { #[cfg(test)] fn address(&self) -> usize { - self.buffer - .as_ref() - .expect("buffer lease owns a buffer") - .ptr - .as_ptr() as usize - } -} - -impl Drop for BufferLease { - fn drop(&mut self) { - if let Some(buffer) = self.buffer.take() { - match &self.owner { - BufferOwner::Fixed => drop(buffer), - BufferOwner::Standalone { memory, .. } => { - let capacity = buffer.capacity; - drop(buffer); - memory.release(capacity); - } - } - } + self.buffer.ptr.as_ptr() as usize } } @@ -278,6 +232,18 @@ struct AlignedBuffer { unsafe impl Send for AlignedBuffer {} impl AlignedBuffer { + fn try_new(capacity: usize) -> Option { + let layout = Layout::from_size_align(capacity, BUFFER_ALIGNMENT).ok()?; + // SAFETY: both constructors validate that capacity is non-zero; the + // layout has valid power-of-two alignment and fits in isize. + let ptr = NonNull::new(unsafe { alloc(layout) })?; + Some(Self { + ptr, + capacity, + initialized: 0, + }) + } + fn prepare_zeroed(&mut self, length: usize) { debug_assert!(length <= self.capacity); // SAFETY: the allocation is valid for `capacity` bytes and this value @@ -309,34 +275,15 @@ impl AlignedBuffer { // mutable borrow is exclusive, and `length <= initialized`. unsafe { slice::from_raw_parts_mut(self.ptr.as_ptr(), length) } } - - fn deallocate(&mut self) { - if self.capacity == 0 { - return; - } - deallocate_buffer(self.ptr, self.capacity); - self.ptr = NonNull::dangling(); - self.capacity = 0; - self.initialized = 0; - } -} - -fn allocate_buffer(capacity: usize) -> Option> { - let layout = Layout::from_size_align(capacity, BUFFER_ALIGNMENT).ok()?; - // SAFETY: `layout` has non-zero size and valid power-of-two alignment. - NonNull::new(unsafe { alloc(layout) }) -} - -fn deallocate_buffer(pointer: NonNull, capacity: usize) { - let layout = Layout::from_size_align(capacity, BUFFER_ALIGNMENT) - .expect("stored aligned-buffer layout is valid"); - // SAFETY: `pointer` was allocated with this exact layout and is owned here. - unsafe { dealloc(pointer.as_ptr(), layout) }; } impl Drop for AlignedBuffer { fn drop(&mut self) { - self.deallocate(); + let layout = Layout::from_size_align(self.capacity, BUFFER_ALIGNMENT) + .expect("stored aligned-buffer layout is valid"); + // SAFETY: the pointer was allocated with this exact layout and this + // buffer owns it until drop. + unsafe { dealloc(self.ptr.as_ptr(), layout) }; } } diff --git a/cache2/src/region/staging.rs b/cache2/src/region/staging.rs index bc229e6..c39f188 100644 --- a/cache2/src/region/staging.rs +++ b/cache2/src/region/staging.rs @@ -27,7 +27,7 @@ use crate::managed_memory::BUFFER_ALIGNMENT; use crate::managed_memory::BufferLease; use crate::managed_memory::ManagedMemory; use crate::managed_memory::ManagedMemoryError; -use crate::managed_memory::RuntimeMemoryReservation; +use crate::managed_memory::MemoryReservation; use crate::region::index::packed::IndexEntry; use crate::region::index::packed::MAX_RECORD_LEN; use crate::region::index::packed::PackedLocation; @@ -252,7 +252,7 @@ pub struct AppendStaging { shards: Vec, chunk_bytes: usize, region_size: u64, - _memory: RuntimeMemoryReservation, + _reservation: MemoryReservation, } impl AppendStaging { @@ -301,7 +301,7 @@ impl AppendStaging { .ok_or(ManagedMemoryError::Allocation)?; // Keep the aggregate reservation alive so eager buffers and record // vectors participate in the hard memory limit. - let memory = managed_memory.reserve_runtime_memory(reserved)?; + let reservation = managed_memory.reserve(reserved)?; let mut shards = Vec::new(); shards @@ -329,7 +329,7 @@ impl AppendStaging { shards, chunk_bytes, region_size, - _memory: memory, + _reservation: reservation, }) } From 0887617186843de6a2188e9d0648cedeaac36aa4 Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 21 Sep 2026 21:34:52 +0800 Subject: [PATCH 2/4] test: exercise buffer initialization through production interfaces Remove test-only preparation and growth methods from buffer ownership types. Initialize engine test buffers through the same target and completion interfaces as production, and use prepared ranges for fixed staging buffers. Replace growth-only coverage with initialized-prefix, reuse, zero-initialization, and out-of-bounds access checks. Validation: cargo x check, cargo x test including extended tests, cargo x lint, and Linux all-targets workspace compilation with all features passed. --- cache2/src/io/engine/tests.rs | 6 ++- cache2/src/managed_memory.rs | 87 +++++++++++++---------------------- cache2/src/region/appender.rs | 4 +- 3 files changed, 39 insertions(+), 58 deletions(-) diff --git a/cache2/src/io/engine/tests.rs b/cache2/src/io/engine/tests.rs index cfce4f8..bd7ecd9 100644 --- a/cache2/src/io/engine/tests.rs +++ b/cache2/src/io/engine/tests.rs @@ -199,7 +199,11 @@ fn read_buffer(managed_memory: &Arc, length: usize) -> IoBuffer { fn write_buffer(managed_memory: &Arc, bytes: &[u8]) -> IoBuffer { let mut lease = managed_memory.try_read_buffer(bytes.len()).unwrap(); - lease.prepare(bytes.len()).unwrap().copy_from_slice(bytes); + let target = lease.read_target(bytes.len()).unwrap(); + // SAFETY: the exclusively owned target fits the source, and the allocations + // do not overlap. Publish the initialized range only after copying it. + unsafe { target.copy_from_nonoverlapping(bytes.as_ptr(), bytes.len()) }; + lease.mark_initialized(bytes.len()).unwrap(); IoBuffer::for_write(lease, bytes.len()).unwrap() } diff --git a/cache2/src/managed_memory.rs b/cache2/src/managed_memory.rs index b100930..99c4757 100644 --- a/cache2/src/managed_memory.rs +++ b/cache2/src/managed_memory.rs @@ -151,31 +151,6 @@ impl BufferLease { }) } - #[cfg(test)] - pub fn prepare(&mut self, length: usize) -> Result<&mut [u8], ()> { - let buffer = &mut self.buffer; - if length > buffer.capacity { - return Err(()); - } - // Callers encode complete records. Clearing here also fixes padding and - // prevents bytes from a prior key/value escaping into a later write. - buffer.prepare_zeroed(length); - Ok(buffer.prefix_mut(length)) - } - - /// Grow the leased buffer without clearing bytes already in the buffer. - /// - /// Fresh capacity is zeroed before it is exposed as initialized bytes. - #[cfg(test)] - fn grow_preserving(&mut self, length: usize) -> Result<&mut [u8], ()> { - let buffer = &mut self.buffer; - if length > buffer.capacity { - return Err(()); - } - buffer.zero_uninitialized_through(length); - Ok(buffer.prefix_mut(length)) - } - pub fn prepared(&self, length: usize) -> Result<&[u8], ()> { let buffer = &self.buffer; if length > buffer.initialized { @@ -253,21 +228,6 @@ impl AlignedBuffer { self.initialized = self.initialized.max(length); } - #[cfg(test)] - fn zero_uninitialized_through(&mut self, length: usize) { - debug_assert!(length <= self.capacity); - if length > self.initialized { - // SAFETY: the uninitialized tail is inside the owned allocation. - unsafe { - self.ptr - .as_ptr() - .add(self.initialized) - .write_bytes(0, length - self.initialized); - } - self.initialized = length; - } - } - fn prefix_mut(&mut self, length: usize) -> &mut [u8] { debug_assert!(length <= self.capacity); debug_assert!(length <= self.initialized); @@ -383,25 +343,42 @@ mod tests { } #[test] - fn preserving_growth_keeps_prefix_and_zeroes_fresh_capacity() { - let mut buffer = BufferLease::try_fixed(3 * BUFFER_ALIGNMENT).unwrap(); - - let first = buffer.grow_preserving(BUFFER_ALIGNMENT).unwrap(); - assert!(first.iter().all(|byte| *byte == 0)); - first.fill(0x5a); - - let grown = buffer.grow_preserving(2 * BUFFER_ALIGNMENT).unwrap(); - assert!(grown[..BUFFER_ALIGNMENT].iter().all(|byte| *byte == 0x5a)); - assert!(grown[BUFFER_ALIGNMENT..].iter().all(|byte| *byte == 0)); - assert_eq!(buffer.address() % BUFFER_ALIGNMENT, 0); + fn read_buffers_expose_only_the_initialized_prefix() { + let managed_memory = ManagedMemory::try_new(limits()).unwrap(); + let mut buffer = managed_memory.try_read_buffer(5000).unwrap(); + let bytes = b"completed read prefix"; + let target = buffer.read_target(5000).unwrap(); + // SAFETY: simulate a completed read into the exclusively owned target; + // the source and destination do not overlap and the bytes fit. + unsafe { target.copy_from_nonoverlapping(bytes.as_ptr(), bytes.len()) }; + + assert!(buffer.prepared(bytes.len()).is_err()); + buffer.mark_initialized(bytes.len()).unwrap(); + assert_eq!(buffer.prepared(bytes.len()).unwrap(), bytes); + assert!(buffer.prepared(bytes.len() + 1).is_err()); + assert!(buffer.prepared_mut(bytes.len() + 1).is_err()); + + // A shorter read can reuse the buffer without invalidating its prefix. + buffer.mark_initialized(1).unwrap(); + assert_eq!(buffer.prepared(bytes.len()).unwrap(), bytes); + assert!(buffer.mark_initialized(2 * BUFFER_ALIGNMENT + 1).is_err()); + assert!(buffer.prepared(bytes.len() + 1).is_err()); } #[test] - fn failed_preserving_growth_keeps_the_existing_buffer() { - let mut buffer = BufferLease::try_fixed(2 * BUFFER_ALIGNMENT).unwrap(); - buffer.grow_preserving(BUFFER_ALIGNMENT).unwrap().fill(0xa5); + fn fixed_buffers_are_zeroed_and_reject_out_of_bounds_access() { + let mut buffer = BufferLease::try_fixed(BUFFER_ALIGNMENT).unwrap(); + assert!( + buffer + .prepared(BUFFER_ALIGNMENT) + .unwrap() + .iter() + .all(|byte| *byte == 0) + ); + buffer.prepared_mut(BUFFER_ALIGNMENT).unwrap().fill(0xa5); - assert!(buffer.grow_preserving(3 * BUFFER_ALIGNMENT).is_err()); + assert!(buffer.read_target(BUFFER_ALIGNMENT + 1).is_err()); + assert!(buffer.prepared_mut(BUFFER_ALIGNMENT + 1).is_err()); assert!( buffer .prepared(BUFFER_ALIGNMENT) diff --git a/cache2/src/region/appender.rs b/cache2/src/region/appender.rs index 53f9091..f34ed75 100644 --- a/cache2/src/region/appender.rs +++ b/cache2/src/region/appender.rs @@ -305,7 +305,7 @@ mod tests { }); let engine = IoEngine::for_test(io.clone(), 1).unwrap(); let mut lease = BufferLease::try_fixed(4096).unwrap(); - lease.prepare(4096).unwrap().fill(0x5a); + lease.prepared_mut(4096).unwrap().fill(0x5a); let absolute = DATA_REGION_AREA_OFFSET + geometry().region_size; let io_recovery = IoRecovery::new(Some(Duration::from_secs(5))); let mut attempt = BackgroundIoAttempt::new(&io_recovery, None); @@ -337,7 +337,7 @@ mod tests { let io = Arc::new(RecordingIo::default()); let engine = IoEngine::for_test(io.clone(), 1).unwrap(); let mut lease = BufferLease::try_fixed(4096).unwrap(); - lease.prepare(4096).unwrap().fill(0x5a); + lease.prepared_mut(4096).unwrap().fill(0x5a); let buffer = IoBuffer::for_write(lease, 4096).unwrap(); let absolute = DATA_REGION_AREA_OFFSET + geometry().region_size; From fdc14486f2c98c3408a31447ed5c7062c32cd626 Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 21 Sep 2026 21:54:17 +0800 Subject: [PATCH 3/4] refactor: replace underscore bindings with lint expectations Give retained ownership fields, scope guards, and required trait parameters ordinary names with narrowly scoped expect attributes. Preserve drop order and guard lifetimes, including write-lock poisoning tests, and rename actively used fields without suppressions. Remove unused return-value bindings and obsolete benchmark parameters. Document the naming convention and the distinction between retained guards and discarded values in CONTRIBUTING.md. Validation: cargo x check, cargo x test including extended tests, cargo x lint, and Linux all-targets workspace compilation with all features passed. --- CONTRIBUTING.md | 2 ++ benchmarks/cache/main.rs | 19 ++++--------- cache2/src/cache/runtime/mod.rs | 24 ++++++++++++---- cache2/src/io/engine/mod.rs | 24 +++++++++++----- cache2/src/io/engine/tests.rs | 39 ++++++++++++++++++++++---- cache2/src/io/engine/uring.rs | 6 +++- cache2/src/io/file.rs | 22 ++++++++++++--- cache2/src/managed_memory.rs | 10 +++++-- cache2/src/memory/mod.rs | 12 +++++--- cache2/src/region/appender.rs | 6 +++- cache2/src/region/index/storage/mod.rs | 19 +++++++++---- cache2/src/region/mod.rs | 3 +- cache2/src/region/persistence/tests.rs | 9 ++++-- cache2/src/region/reader.rs | 7 ++++- cache2/src/region/staging.rs | 8 ++++-- examples/src/stats.rs | 6 +++- 16 files changed, 158 insertions(+), 58 deletions(-) diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 7643719..c135551 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -33,6 +33,8 @@ Follow the surrounding code and the design constraints in [ARCHITECTURE.md](ARCH Name types by their concrete domain role. Use `Options` for editable inputs awaiting validation and `Config` for validated or resolved configuration. Keep names such as `Layout` for geometry and use specific domain nouns such as `Desc`, `Candidate`, `Victims`, or `Inspection` for runtime data. Avoid `Plan` as a type-name suffix; computing data before using it does not by itself make that data a plan. Keep related fields, local variables, and functions consistent with the type's role. +Do not prefix Rust identifiers with an underscore to suppress unused-code warnings. Remove unnecessary bindings and private parameters. For required but unread fields, guards, or trait parameters, use ordinary names with a narrowly scoped `#[expect(dead_code)]` or `#[expect(unused_variables)]`; explain non-obvious ownership or lifetime requirements with `reason`. Preserve guard lifetimes and drop order. Use the wildcard pattern `_` only when intentionally discarding a value without binding it. + Declare restricted visibility at module boundaries and use `pub` for items in those modules' APIs. Keep items private when only their defining module and its descendants need them. For items reachable through public modules or re-exported public types, reserve `pub` for intentional public API and use narrower visibility for internal callers. ## Documentation diff --git a/benchmarks/cache/main.rs b/benchmarks/cache/main.rs index 98bce7e..9cf057b 100644 --- a/benchmarks/cache/main.rs +++ b/benchmarks/cache/main.rs @@ -408,7 +408,6 @@ async fn run(config: BenchConfig) -> io::Result<()> { .transpose()?; report( "put_drain", - "put + drain", "write", config.write_clients, &write.measurement, @@ -437,7 +436,7 @@ async fn run(config: BenchConfig) -> io::Result<()> { cache.close_warm().await?; drop(cache); let warm_close = started.elapsed(); - report_latency("warm_close", "warm close", warm_close); + report_latency("warm_close", warm_close); let cache = Arc::new(Cache::open(files.data(), cache_config.clone()).await?); if cache.startup_mode() != StartupMode::Warm { @@ -459,7 +458,6 @@ async fn run(config: BenchConfig) -> io::Result<()> { if l1_entry_eligible { report( "l2_promote", - "L2 get + promote", "read", config.l2_clients(), &l2_read.measurement, @@ -470,7 +468,6 @@ async fn run(config: BenchConfig) -> io::Result<()> { } else { report( "l2_read", - "L2 get", "read", config.l2_clients(), &l2_read.measurement, @@ -485,7 +482,7 @@ async fn run(config: BenchConfig) -> io::Result<()> { l2_read.measurement.elapsed, &l2_read.primary, ); - report_read_latency("l2_latency", "L2 get latency", &l2_read.latency); + report_read_latency("l2_latency", &l2_read.latency); l2_read.measurement } else { let _ = concurrent_writes( @@ -522,7 +519,6 @@ async fn run(config: BenchConfig) -> io::Result<()> { .await?; report( "l2_hot_scan", - "L2 cold scan", "read", config.l2_clients(), &cold_scan.measurement, @@ -536,7 +532,7 @@ async fn run(config: BenchConfig) -> io::Result<()> { cold_scan.measurement.elapsed, &cold_scan.primary, ); - report_read_latency("l2_scan_latency", "L2 scan latency", &cold_scan.latency); + report_read_latency("l2_scan_latency", &cold_scan.latency); report_tiers( "hot_during_scan", "hot during scan", @@ -668,7 +664,6 @@ async fn run(config: BenchConfig) -> io::Result<()> { .await?; report( "resident_l1", - "resident L1 get", "read", config.clients, &resident.measurement, @@ -676,7 +671,7 @@ async fn run(config: BenchConfig) -> io::Result<()> { Some(&resident.latency), config.read_latency_sample_interval, ); - report_read_latency("resident_l1_latency", "L1 get latency", &resident.latency); + report_read_latency("resident_l1_latency", &resident.latency); if let Some(before) = before { let after = cache.snapshot()?; println!( @@ -961,10 +956,8 @@ fn verify_value(ordinal: usize, value: &[u8]) -> io::Result<()> { Ok(()) } -#[allow(clippy::too_many_arguments)] fn report( phase: &str, - _name: &str, operation: &str, workers: usize, measurement: &Measurement, @@ -1032,7 +1025,7 @@ fn report_tiers(phase: &str, name: &str, elapsed: Duration, tiers: &TierCounts) ); } -fn report_latency(phase: &str, _name: &str, elapsed: Duration) { +fn report_latency(phase: &str, elapsed: Duration) { JobReport::new("cache", None, phase, "control", elapsed, 1).emit(); println!( "result phase={phase} elapsed_ns={} operations=0 bytes=0 ops_per_sec=0.000 mib_per_sec=0.000 checksum=0000000000000000", @@ -1040,7 +1033,7 @@ fn report_latency(phase: &str, _name: &str, elapsed: Duration) { ); } -fn report_read_latency(phase: &str, _name: &str, latency: &LatencyHistogram) { +fn report_read_latency(phase: &str, latency: &LatencyHistogram) { let summary = latency.summary(); if summary.samples == 0 { return; diff --git a/cache2/src/cache/runtime/mod.rs b/cache2/src/cache/runtime/mod.rs index ad068a0..42b26b6 100644 --- a/cache2/src/cache/runtime/mod.rs +++ b/cache2/src/cache/runtime/mod.rs @@ -883,12 +883,12 @@ impl CacheRuntime { current_bytes, } => { if ADMIT_L1 { - let _published = state.memory.publish(hash, key, value, seqno); + state.memory.publish(hash, key, value, seqno); } else { // Prevent an older exact-key L1 value from indefinitely // shadowing the prefetched L2 record. Contention remains a // valid best-effort stale outcome. - let _removed = state.memory.delete(hash, key, seqno); + state.memory.delete(hash, key, seqno); } if should_wake_write( previous_bytes, @@ -940,7 +940,7 @@ impl CacheRuntime { } return Err(write_overload_error()); }; - let _removed = state.memory.delete(hash, key, seqno); + state.memory.delete(hash, key, seqno); if let Some(activity) = activity { RuntimeMetrics::increment(&activity.deletes); } @@ -1268,7 +1268,11 @@ impl CacheRuntime { pub fn drain(&self) -> io::Result<()> { let operations = self.operations.begin_drain()?; operations.wait()?; - let _draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); + #[expect( + unused_variables, + reason = "Restore lifecycle state when draining finishes." + )] + let draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); let state = &self.state; drain_shards(state, false) } @@ -1276,7 +1280,11 @@ impl CacheRuntime { pub async fn drain_async(&self) -> io::Result<()> { let operations = self.operations.begin_drain()?; operations.wait_async().await; - let _draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); + #[expect( + unused_variables, + reason = "Restore lifecycle state on completion or cancellation." + )] + let draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); let state = &self.state; drain_shards_async(state, false).await } @@ -1361,7 +1369,11 @@ impl CacheRuntime { .get(shard_id) .expect("test shard exists"); let result = panic::catch_unwind(AssertUnwindSafe(|| { - let _state = shard.state.lock().unwrap(); + #[expect( + unused_variables, + reason = "Keep the guard alive so unwinding poisons the lock." + )] + let state = shard.state.lock().unwrap(); panic!("poison shard gate"); })); assert!(result.is_err()); diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index e8ea5d3..15b97ba 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -847,7 +847,8 @@ impl BoundedIoRequest { let mut request = AsyncRequestGuard::new(self.request, engine); let deadline = tokio::time::Instant::from_std(self.deadline); let completion = { - let _entered = tokio_handle.enter(); + #[expect(unused_variables)] + let entered = tokio_handle.enter(); tokio::time::timeout_at(deadline, request.request_mut()) } .await; @@ -858,7 +859,8 @@ impl BoundedIoRequest { let cancel_error = request.cancel().err(); let completion = { - let _entered = tokio_handle.enter(); + #[expect(unused_variables)] + let entered = tokio_handle.enter(); tokio::time::timeout(self.cancel_grace, request.request_mut()) } .await; @@ -1140,7 +1142,8 @@ impl ReadSlotAdmission { self.ensure_open()?; let acquire = Arc::clone(&self.slots).acquire_owned(1); { - let _entered = tokio_handle.enter(); + #[expect(unused_variables)] + let entered = tokio_handle.enter(); tokio::time::timeout_at(tokio::time::Instant::from_std(deadline), acquire) } .await @@ -1179,7 +1182,11 @@ impl ReadSlotWaiter { .read_slot_admission .as_ref() .ok_or_else(|| io::Error::other("async read admission is disabled"))?; - let _waiter = admission.register_waiter(); + #[expect( + unused_variables, + reason = "Unregister the waiter on completion or cancellation." + )] + let waiter = admission.register_waiter(); self.state.ensure_accepting()?; let permit = admission.acquire_until(deadline, tokio_handle).await?; self.state.try_reserve_read_slot(Some(permit)) @@ -1189,7 +1196,8 @@ impl ReadSlotWaiter { impl Drop for IoSlot { fn drop(&mut self) { { - let _slot = lock_unpoisoned(&self.state.slot_lock); + #[expect(unused_variables)] + let slot = lock_unpoisoned(&self.state.slot_lock); self.state .slot_state .fetch_sub(slot_delta(self.write), Ordering::AcqRel); @@ -1336,7 +1344,8 @@ impl EngineState { fn stop_accepting_slots(&self) { { - let _slot = lock_unpoisoned(&self.slot_lock); + #[expect(unused_variables)] + let slot = lock_unpoisoned(&self.slot_lock); self.accepting.store(false, Ordering::Release); self.slot_available.notify_all(); } @@ -1406,7 +1415,8 @@ impl EngineState { /// Taking the same mutex used around the check-and-wait transition makes /// `cancelled.store(true, Release); wake_slot_waiters()` lossless. fn wake_slot_waiters(&self) { - let _slot = lock_unpoisoned(&self.slot_lock); + #[expect(unused_variables)] + let slot = lock_unpoisoned(&self.slot_lock); self.slot_available.notify_all(); } diff --git a/cache2/src/io/engine/tests.rs b/cache2/src/io/engine/tests.rs index bd7ecd9..5c6d912 100644 --- a/cache2/src/io/engine/tests.rs +++ b/cache2/src/io/engine/tests.rs @@ -118,13 +118,22 @@ impl BlockingIo { } impl PositionedIo for BlockingIo { - fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result { + fn read_at( + &self, + buffer: &mut [u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { self.enter_and_wait(); buffer.fill(0); Ok(buffer.len()) } - fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + buffer: &[u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { self.enter_and_wait(); Ok(buffer.len()) } @@ -149,7 +158,11 @@ impl PanicOnceIo { } impl PositionedIo for PanicOnceIo { - fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result { + fn read_at( + &self, + buffer: &mut [u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { if self.panic_next_read.swap(false, Ordering::AcqRel) { panic!("injected io panic"); } @@ -157,13 +170,22 @@ impl PositionedIo for PanicOnceIo { Ok(buffer.len()) } - fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + buffer: &[u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { Ok(buffer.len()) } } impl PositionedIo for ShortThenErrorIo { - fn read_at(&self, buffer: &mut [u8], _offset: u64) -> io::Result { + fn read_at( + &self, + buffer: &mut [u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { if self.read_calls.fetch_add(1, Ordering::Relaxed) == 0 { let transferred = 3.min(buffer.len()); buffer[..transferred].fill(0x5a); @@ -173,7 +195,12 @@ impl PositionedIo for ShortThenErrorIo { } } - fn write_at(&self, _point: WritePoint, buffer: &[u8], _offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + buffer: &[u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { if self.write_calls.fetch_add(1, Ordering::Relaxed) == 0 { Ok(3.min(buffer.len())) } else { diff --git a/cache2/src/io/engine/uring.rs b/cache2/src/io/engine/uring.rs index b037fcb..2eb98b4 100644 --- a/cache2/src/io/engine/uring.rs +++ b/cache2/src/io/engine/uring.rs @@ -1020,7 +1020,11 @@ mod tests { let depth = 2; let state = Arc::new(EngineState::new(depth, false, false)); let (commands, receiver) = mpsc::sync_channel(depth * 2 + 1); - let (_wake_sender, wake_receiver) = UnixStream::pair().unwrap(); + #[expect( + unused_variables, + reason = "Keep the peer open while exercising the driver." + )] + let (wake_sender, wake_receiver) = UnixStream::pair().unwrap(); let producer = Arc::new(CancelledCommandProducer { state: Arc::clone(&state), commands, diff --git a/cache2/src/io/file.rs b/cache2/src/io/file.rs index f2a5f7b..4cb2b58 100644 --- a/cache2/src/io/file.rs +++ b/cache2/src/io/file.rs @@ -523,7 +523,7 @@ impl StorageFile for CacheFile { } } - fn sync(&self, _point: SyncPoint, mode: SyncMode) -> io::Result<()> { + fn sync(&self, #[expect(unused_variables)] point: SyncPoint, mode: SyncMode) -> io::Result<()> { match mode { SyncMode::Data => self.file.sync_data(), SyncMode::All => self.file.sync_all(), @@ -572,7 +572,12 @@ impl PositionedIo for CacheFile { self.file.read_at(buffer, offset) } - fn write_at(&self, _point: WritePoint, buffer: &[u8], offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + buffer: &[u8], + offset: u64, + ) -> io::Result { self.file.write_at(buffer, offset) } } @@ -874,7 +879,11 @@ mod tests { } impl PositionedIo for InterruptedIo { - fn read_at(&self, _buffer: &mut [u8], _offset: u64) -> io::Result { + fn read_at( + &self, + #[expect(unused_variables)] buffer: &mut [u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { self.calls.fetch_add(1, Ordering::Relaxed); Err(io::Error::new( io::ErrorKind::Interrupted, @@ -882,7 +891,12 @@ mod tests { )) } - fn write_at(&self, _point: WritePoint, _buffer: &[u8], _offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + #[expect(unused_variables)] buffer: &[u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { self.calls.fetch_add(1, Ordering::Relaxed); Err(io::Error::new( io::ErrorKind::Interrupted, diff --git a/cache2/src/managed_memory.rs b/cache2/src/managed_memory.rs index 99c4757..da037f3 100644 --- a/cache2/src/managed_memory.rs +++ b/cache2/src/managed_memory.rs @@ -108,7 +108,7 @@ impl ManagedMemory { let buffer = AlignedBuffer::try_new(capacity)?; Some(BufferLease { buffer, - _reservation: Some(reservation), + reservation: Some(reservation), }) } @@ -133,7 +133,11 @@ pub struct BufferLease { // Fields drop in declaration order: free the allocation before returning // its budget. Fixed staging buffers use their owner's aggregate charge. buffer: AlignedBuffer, - _reservation: Option, + #[expect( + dead_code, + reason = "Returns the memory charge after the allocation is dropped." + )] + reservation: Option, } impl BufferLease { @@ -147,7 +151,7 @@ impl BufferLease { buffer.prepare_zeroed(length); Ok(Self { buffer, - _reservation: None, + reservation: None, }) } diff --git a/cache2/src/memory/mod.rs b/cache2/src/memory/mod.rs index 037e868..5e664e8 100644 --- a/cache2/src/memory/mod.rs +++ b/cache2/src/memory/mod.rs @@ -167,7 +167,7 @@ struct MemoryCharge { struct MemoryValueInner { bytes: Box<[u8]>, key_length: usize, - _charge: MemoryCharge, + charge: MemoryCharge, } #[derive(Clone)] @@ -197,7 +197,7 @@ impl MemoryValue { Ok(Self(Arc::new(MemoryValueInner { bytes: bytes.into_boxed_slice(), key_length: key.len(), - _charge: charge, + charge, }))) } @@ -225,7 +225,7 @@ impl MemoryValue { fn disarm_exclusive_charge(&mut self) -> usize { let inner = Arc::get_mut(&mut self.0).expect("exclusive resident value gained an owner"); - inner._charge.disarm() + inner.charge.disarm() } } @@ -501,7 +501,11 @@ impl MemoryShard { .charged_bytes() ); } - let _entry = self.slots[index].take().expect("resident slot disappeared"); + #[expect( + unused_variables, + reason = "Release the value after slot bookkeeping is complete." + )] + let entry = self.slots[index].take().expect("resident slot disappeared"); self.resident_bytes = self.resident_bytes.saturating_sub(resident_bytes); debug_assert!(self.free_count < self.slots.len()); let packed = u32::try_from(index).expect("memory slot index exceeds u32"); diff --git a/cache2/src/region/appender.rs b/cache2/src/region/appender.rs index f34ed75..d7d898e 100644 --- a/cache2/src/region/appender.rs +++ b/cache2/src/region/appender.rs @@ -262,7 +262,11 @@ mod tests { } impl PositionedIo for RecordingIo { - fn read_at(&self, _buffer: &mut [u8], _offset: u64) -> io::Result { + fn read_at( + &self, + #[expect(unused_variables)] buffer: &mut [u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { Err(io::Error::new(io::ErrorKind::Unsupported, "read unused")) } diff --git a/cache2/src/region/index/storage/mod.rs b/cache2/src/region/index/storage/mod.rs index 6a4e059..7b0bcaa 100644 --- a/cache2/src/region/index/storage/mod.rs +++ b/cache2/src/region/index/storage/mod.rs @@ -641,7 +641,7 @@ impl IndexStorage { let data_pointer = backing.as_mut_ptr(); let page_states = allocate_page_states(layout.page_count, PAGE_STATE_VALID)?; let core = Arc::new(IndexStorageCore { - _backing: UnsafeCell::new(backing), + backing: UnsafeCell::new(backing), data_pointer, data_offset: 0, slot_count, @@ -706,7 +706,7 @@ impl IndexStorage { let data_pointer = backing.as_mut_ptr(); let page_states = allocate_page_states(layout.page_count, PAGE_STATE_UNCHECKED)?; let core = Arc::new(IndexStorageCore { - _backing: UnsafeCell::new(backing), + backing: UnsafeCell::new(backing), data_pointer, data_offset, slot_count, @@ -920,7 +920,11 @@ impl IndexStorage { } struct IndexStorageCore { - _backing: UnsafeCell, + #[expect( + dead_code, + reason = "Owns the mapping accessed through the cached data pointer." + )] + backing: UnsafeCell, data_pointer: *mut u8, data_offset: usize, slot_count: usize, @@ -1146,7 +1150,7 @@ impl IndexStorageCore { fn data_ptr(&self) -> *const u8 { // SAFETY: `data_pointer` addresses the stable allocation owned by - // `_backing`, and `data_offset + image_len` was checked against that + // `backing`, and `data_offset + image_len` was checked against that // allocation during construction. unsafe { self.data_pointer.cast_const().add(self.data_offset) } } @@ -1538,7 +1542,12 @@ impl PartitionedIndexStorage { pub fn poison_hash_partition_for_test(&self, hash: u64) { let partition = index_partition_for(hash, self.partitions.len()); let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - let _guard = self.partitions[partition].write().unwrap(); + #[expect( + unused_variables, + clippy::readonly_write_lock, + reason = "Unwinding must hold a write guard to poison the partition." + )] + let guard = self.partitions[partition].write().unwrap(); panic!("poison index partition for test"); })); assert!(result.is_err()); diff --git a/cache2/src/region/mod.rs b/cache2/src/region/mod.rs index 2bb9c1c..b064479 100644 --- a/cache2/src/region/mod.rs +++ b/cache2/src/region/mod.rs @@ -881,7 +881,8 @@ impl RegionStore { )); } let payload = RecordPayload::new(key, value); - let _shard_mutation = match self.append_gate(shard_id)?.mutation.try_lock() { + #[expect(unused_variables)] + let shard_mutation = match self.append_gate(shard_id)?.mutation.try_lock() { Ok(guard) => guard, Err(TryLockError::WouldBlock) => return Ok(RegionStageValue::NeedsProgress), Err(TryLockError::Poisoned(_)) => { diff --git a/cache2/src/region/persistence/tests.rs b/cache2/src/region/persistence/tests.rs index 3836ad6..65b9d99 100644 --- a/cache2/src/region/persistence/tests.rs +++ b/cache2/src/region/persistence/tests.rs @@ -369,7 +369,11 @@ fn run_crash_child(case: &str, paths: RegionPaths) -> ! { let data = data_path_superblock(); match case { "open" => { - let _store = CacheSession::for_test(paths, data, 4096).unwrap(); + #[expect( + unused_variables, + reason = "Keep the session open until the process is killed." + )] + let session = CacheSession::for_test(paths, data, 4096).unwrap(); kill_process(); } "write" | "drain" => { @@ -913,8 +917,7 @@ fn completed_owned_span_publishes_index_without_a_steady_state_sync() { outcome => panic!("unexpected staging outcome: {outcome:?}"), } } - let (first_key, first_hash, _first_seqno) = - first.expect("4 MiB span must contain target-size records"); + let (first_key, first_hash, _) = first.expect("4 MiB span must contain target-size records"); let (last_key, last_hash, last_seqno) = last.expect("4 MiB span must retain its final record"); assert!(staged_records > 240); assert_eq!(regions.lookup_snapshot(first_hash).unwrap(), None); diff --git a/cache2/src/region/reader.rs b/cache2/src/region/reader.rs index 4245027..6000b60 100644 --- a/cache2/src/region/reader.rs +++ b/cache2/src/region/reader.rs @@ -346,7 +346,12 @@ mod tests { Ok(buffer.len()) } - fn write_at(&self, _point: WritePoint, _buffer: &[u8], _offset: u64) -> io::Result { + fn write_at( + &self, + #[expect(unused_variables)] point: WritePoint, + #[expect(unused_variables)] buffer: &[u8], + #[expect(unused_variables)] offset: u64, + ) -> io::Result { Err(io::Error::new(io::ErrorKind::Unsupported, "write unused")) } } diff --git a/cache2/src/region/staging.rs b/cache2/src/region/staging.rs index c39f188..95faaf0 100644 --- a/cache2/src/region/staging.rs +++ b/cache2/src/region/staging.rs @@ -252,7 +252,11 @@ pub struct AppendStaging { shards: Vec, chunk_bytes: usize, region_size: u64, - _reservation: MemoryReservation, + #[expect( + dead_code, + reason = "Returns the aggregate charge when staging is dropped." + )] + reservation: MemoryReservation, } impl AppendStaging { @@ -329,7 +333,7 @@ impl AppendStaging { shards, chunk_bytes, region_size, - _reservation: reservation, + reservation, }) } diff --git a/examples/src/stats.rs b/examples/src/stats.rs index b4c7659..7a999b6 100644 --- a/examples/src/stats.rs +++ b/examples/src/stats.rs @@ -39,7 +39,11 @@ async fn main() -> Result<(), Box> { runtime.stats.io_latency = true; let cache = Cache::open(path, CacheConfig::new(storage, runtime)?).await?; cache.put("example", "value")?; - let _value = cache.get("example").await?; + #[expect( + unused_variables, + reason = "Retain the value through snapshot collection and close." + )] + let value = cache.get("example").await?; cache.drain().await?; // Pass this owned, cumulative snapshot to the application's metrics adapter. // It carries collection modes/scopes; bucket bounds and sums use nanoseconds. From d6f4b879043be6eb07839eb3cc9e75f0856ef35b Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 21 Sep 2026 22:20:51 +0800 Subject: [PATCH 4/4] refactor: keep conventional local guard bindings Use underscore-prefixed local bindings for guards and resources retained until scope exit, removing redundant lint expectations without changing drop timing. Keep ordinary field names and precise dead-code expectations for unread ownership fields. Document this distinction in CONTRIBUTING.md. Unused trait parameters and unnecessary bindings retain their existing cleanup. Validation: cargo x check, cargo x test (including extended tests), cargo x lint, and Linux workspace compilation with all targets and features passed. --- CONTRIBUTING.md | 2 +- cache2/src/cache/runtime/mod.rs | 18 +++--------------- cache2/src/io/engine/mod.rs | 24 +++++++----------------- cache2/src/io/engine/uring.rs | 6 +----- cache2/src/memory/mod.rs | 6 +----- cache2/src/region/index/storage/mod.rs | 7 +------ cache2/src/region/mod.rs | 3 +-- cache2/src/region/persistence/tests.rs | 6 +----- examples/src/stats.rs | 6 +----- 9 files changed, 17 insertions(+), 61 deletions(-) diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index c135551..ef442ed 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -33,7 +33,7 @@ Follow the surrounding code and the design constraints in [ARCHITECTURE.md](ARCH Name types by their concrete domain role. Use `Options` for editable inputs awaiting validation and `Config` for validated or resolved configuration. Keep names such as `Layout` for geometry and use specific domain nouns such as `Desc`, `Candidate`, `Victims`, or `Inspection` for runtime data. Avoid `Plan` as a type-name suffix; computing data before using it does not by itself make that data a plan. Keep related fields, local variables, and functions consistent with the type's role. -Do not prefix Rust identifiers with an underscore to suppress unused-code warnings. Remove unnecessary bindings and private parameters. For required but unread fields, guards, or trait parameters, use ordinary names with a narrowly scoped `#[expect(dead_code)]` or `#[expect(unused_variables)]`; explain non-obvious ownership or lifetime requirements with `reason`. Preserve guard lifetimes and drop order. Use the wildcard pattern `_` only when intentionally discarding a value without binding it. +Use ordinary names for fields; annotate required but unread ownership fields with a narrowly scoped `#[expect(dead_code)]` and explain their role with `reason`. Local guards and other resources retained until scope exit may use conventional underscore-prefixed bindings such as `_guard` without a lint attribute. Preserve their lifetimes and drop order; the wildcard pattern `_` does not retain a value. Remove unnecessary bindings and private parameters, and use a narrowly scoped `#[expect(unused_variables)]` for required but unused trait parameters. Declare restricted visibility at module boundaries and use `pub` for items in those modules' APIs. Keep items private when only their defining module and its descendants need them. For items reachable through public modules or re-exported public types, reserve `pub` for intentional public API and use narrower visibility for internal callers. diff --git a/cache2/src/cache/runtime/mod.rs b/cache2/src/cache/runtime/mod.rs index 42b26b6..6c1d4cf 100644 --- a/cache2/src/cache/runtime/mod.rs +++ b/cache2/src/cache/runtime/mod.rs @@ -1268,11 +1268,7 @@ impl CacheRuntime { pub fn drain(&self) -> io::Result<()> { let operations = self.operations.begin_drain()?; operations.wait()?; - #[expect( - unused_variables, - reason = "Restore lifecycle state when draining finishes." - )] - let draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); + let _draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); let state = &self.state; drain_shards(state, false) } @@ -1280,11 +1276,7 @@ impl CacheRuntime { pub async fn drain_async(&self) -> io::Result<()> { let operations = self.operations.begin_drain()?; operations.wait_async().await; - #[expect( - unused_variables, - reason = "Restore lifecycle state on completion or cancellation." - )] - let draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); + let _draining = LifecycleDrainingGuard::enter(&self.metrics.lifecycle, &self.operations); let state = &self.state; drain_shards_async(state, false).await } @@ -1369,11 +1361,7 @@ impl CacheRuntime { .get(shard_id) .expect("test shard exists"); let result = panic::catch_unwind(AssertUnwindSafe(|| { - #[expect( - unused_variables, - reason = "Keep the guard alive so unwinding poisons the lock." - )] - let state = shard.state.lock().unwrap(); + let _state = shard.state.lock().unwrap(); panic!("poison shard gate"); })); assert!(result.is_err()); diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index 15b97ba..e8ea5d3 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -847,8 +847,7 @@ impl BoundedIoRequest { let mut request = AsyncRequestGuard::new(self.request, engine); let deadline = tokio::time::Instant::from_std(self.deadline); let completion = { - #[expect(unused_variables)] - let entered = tokio_handle.enter(); + let _entered = tokio_handle.enter(); tokio::time::timeout_at(deadline, request.request_mut()) } .await; @@ -859,8 +858,7 @@ impl BoundedIoRequest { let cancel_error = request.cancel().err(); let completion = { - #[expect(unused_variables)] - let entered = tokio_handle.enter(); + let _entered = tokio_handle.enter(); tokio::time::timeout(self.cancel_grace, request.request_mut()) } .await; @@ -1142,8 +1140,7 @@ impl ReadSlotAdmission { self.ensure_open()?; let acquire = Arc::clone(&self.slots).acquire_owned(1); { - #[expect(unused_variables)] - let entered = tokio_handle.enter(); + let _entered = tokio_handle.enter(); tokio::time::timeout_at(tokio::time::Instant::from_std(deadline), acquire) } .await @@ -1182,11 +1179,7 @@ impl ReadSlotWaiter { .read_slot_admission .as_ref() .ok_or_else(|| io::Error::other("async read admission is disabled"))?; - #[expect( - unused_variables, - reason = "Unregister the waiter on completion or cancellation." - )] - let waiter = admission.register_waiter(); + let _waiter = admission.register_waiter(); self.state.ensure_accepting()?; let permit = admission.acquire_until(deadline, tokio_handle).await?; self.state.try_reserve_read_slot(Some(permit)) @@ -1196,8 +1189,7 @@ impl ReadSlotWaiter { impl Drop for IoSlot { fn drop(&mut self) { { - #[expect(unused_variables)] - let slot = lock_unpoisoned(&self.state.slot_lock); + let _slot = lock_unpoisoned(&self.state.slot_lock); self.state .slot_state .fetch_sub(slot_delta(self.write), Ordering::AcqRel); @@ -1344,8 +1336,7 @@ impl EngineState { fn stop_accepting_slots(&self) { { - #[expect(unused_variables)] - let slot = lock_unpoisoned(&self.slot_lock); + let _slot = lock_unpoisoned(&self.slot_lock); self.accepting.store(false, Ordering::Release); self.slot_available.notify_all(); } @@ -1415,8 +1406,7 @@ impl EngineState { /// Taking the same mutex used around the check-and-wait transition makes /// `cancelled.store(true, Release); wake_slot_waiters()` lossless. fn wake_slot_waiters(&self) { - #[expect(unused_variables)] - let slot = lock_unpoisoned(&self.slot_lock); + let _slot = lock_unpoisoned(&self.slot_lock); self.slot_available.notify_all(); } diff --git a/cache2/src/io/engine/uring.rs b/cache2/src/io/engine/uring.rs index 2eb98b4..b037fcb 100644 --- a/cache2/src/io/engine/uring.rs +++ b/cache2/src/io/engine/uring.rs @@ -1020,11 +1020,7 @@ mod tests { let depth = 2; let state = Arc::new(EngineState::new(depth, false, false)); let (commands, receiver) = mpsc::sync_channel(depth * 2 + 1); - #[expect( - unused_variables, - reason = "Keep the peer open while exercising the driver." - )] - let (wake_sender, wake_receiver) = UnixStream::pair().unwrap(); + let (_wake_sender, wake_receiver) = UnixStream::pair().unwrap(); let producer = Arc::new(CancelledCommandProducer { state: Arc::clone(&state), commands, diff --git a/cache2/src/memory/mod.rs b/cache2/src/memory/mod.rs index 5e664e8..1eaf1e1 100644 --- a/cache2/src/memory/mod.rs +++ b/cache2/src/memory/mod.rs @@ -501,11 +501,7 @@ impl MemoryShard { .charged_bytes() ); } - #[expect( - unused_variables, - reason = "Release the value after slot bookkeeping is complete." - )] - let entry = self.slots[index].take().expect("resident slot disappeared"); + let _entry = self.slots[index].take().expect("resident slot disappeared"); self.resident_bytes = self.resident_bytes.saturating_sub(resident_bytes); debug_assert!(self.free_count < self.slots.len()); let packed = u32::try_from(index).expect("memory slot index exceeds u32"); diff --git a/cache2/src/region/index/storage/mod.rs b/cache2/src/region/index/storage/mod.rs index 7b0bcaa..6a51531 100644 --- a/cache2/src/region/index/storage/mod.rs +++ b/cache2/src/region/index/storage/mod.rs @@ -1542,12 +1542,7 @@ impl PartitionedIndexStorage { pub fn poison_hash_partition_for_test(&self, hash: u64) { let partition = index_partition_for(hash, self.partitions.len()); let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { - #[expect( - unused_variables, - clippy::readonly_write_lock, - reason = "Unwinding must hold a write guard to poison the partition." - )] - let guard = self.partitions[partition].write().unwrap(); + let _guard = self.partitions[partition].write().unwrap(); panic!("poison index partition for test"); })); assert!(result.is_err()); diff --git a/cache2/src/region/mod.rs b/cache2/src/region/mod.rs index b064479..2bb9c1c 100644 --- a/cache2/src/region/mod.rs +++ b/cache2/src/region/mod.rs @@ -881,8 +881,7 @@ impl RegionStore { )); } let payload = RecordPayload::new(key, value); - #[expect(unused_variables)] - let shard_mutation = match self.append_gate(shard_id)?.mutation.try_lock() { + let _shard_mutation = match self.append_gate(shard_id)?.mutation.try_lock() { Ok(guard) => guard, Err(TryLockError::WouldBlock) => return Ok(RegionStageValue::NeedsProgress), Err(TryLockError::Poisoned(_)) => { diff --git a/cache2/src/region/persistence/tests.rs b/cache2/src/region/persistence/tests.rs index 65b9d99..733e2f0 100644 --- a/cache2/src/region/persistence/tests.rs +++ b/cache2/src/region/persistence/tests.rs @@ -369,11 +369,7 @@ fn run_crash_child(case: &str, paths: RegionPaths) -> ! { let data = data_path_superblock(); match case { "open" => { - #[expect( - unused_variables, - reason = "Keep the session open until the process is killed." - )] - let session = CacheSession::for_test(paths, data, 4096).unwrap(); + let _session = CacheSession::for_test(paths, data, 4096).unwrap(); kill_process(); } "write" | "drain" => { diff --git a/examples/src/stats.rs b/examples/src/stats.rs index 7a999b6..b4c7659 100644 --- a/examples/src/stats.rs +++ b/examples/src/stats.rs @@ -39,11 +39,7 @@ async fn main() -> Result<(), Box> { runtime.stats.io_latency = true; let cache = Cache::open(path, CacheConfig::new(storage, runtime)?).await?; cache.put("example", "value")?; - #[expect( - unused_variables, - reason = "Retain the value through snapshot collection and close." - )] - let value = cache.get("example").await?; + let _value = cache.get("example").await?; cache.drain().await?; // Pass this owned, cumulative snapshot to the application's metrics adapter. // It carries collection modes/scopes; bucket bounds and sums use nanoseconds.