From 3bfac6848a50c0f8d12e0d52c46544b57692e785 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 11:19:41 +0100 Subject: [PATCH 01/10] Expand optimization benchmarks --- benches/range_cache.rs | 163 ++++++++++++++++++++++++++++++++--------- 1 file changed, 128 insertions(+), 35 deletions(-) diff --git a/benches/range_cache.rs b/benches/range_cache.rs index cf5ab93..f7702e4 100644 --- a/benches/range_cache.rs +++ b/benches/range_cache.rs @@ -13,14 +13,20 @@ use range_cache::{CacheCapacity, RangeCache}; fn full_hits(criterion: &mut Criterion) { let mut group = criterion.benchmark_group("full_hit"); for size in [1_024, 65_536, 1_048_576] { - let cache = RangeCache::new(CacheCapacity::Unbounded); - cache - .insert(0_u8, 0..size, Bytes::from(vec![1; size])) - .expect("benchmark insert"); - group.bench_with_input( - BenchmarkId::new("cached_bytes", size), - &size, - |bencher, &size| { + for (policy, capacity) in [ + ("unbounded_cached_bytes", CacheCapacity::Unbounded), + ( + "bounded_cached_bytes", + CacheCapacity::Bounded( + NonZeroUsize::new(size).expect("benchmark capacity is non-zero"), + ), + ), + ] { + let cache = RangeCache::new(capacity); + cache + .insert(0_u8, 0..size, Bytes::from(vec![1; size])) + .expect("benchmark insert"); + group.bench_with_input(BenchmarkId::new(policy, size), &size, |bencher, &size| { bencher.iter(|| { black_box( cache @@ -28,8 +34,8 @@ fn full_hits(criterion: &mut Criterion) { .expect("benchmark range"), ) }); - }, - ); + }); + } } group.finish(); } @@ -161,6 +167,59 @@ fn insertion(criterion: &mut Criterion) { overlap_group.finish(); } +fn sparse_insertion(criterion: &mut Criterion) { + let mut group = criterion.benchmark_group("sparse_insertion"); + for ranges in [8, 64, 512] { + group.bench_with_input( + BenchmarkId::new("beginning_resident_ranges", ranges), + &ranges, + |bencher, &ranges| { + bencher.iter_batched_ref( + || { + ( + fragmented_cache_from(ranges, 32), + Some(Bytes::from_static(&[2; 16])), + ) + }, + |(cache, payload)| { + let Some(payload) = payload.take() else { + panic!("benchmark payload is available"); + }; + black_box( + cache + .insert(0_u8, 0..16, payload) + .expect("benchmark insertion"), + ) + }, + BatchSize::SmallInput, + ); + }, + ); + group.bench_with_input( + BenchmarkId::new("end_resident_ranges", ranges), + &ranges, + |bencher, &ranges| { + let start = ranges * 32; + bencher.iter_batched_ref( + || (fragmented_cache(ranges), Some(Bytes::from_static(&[2; 16]))), + |(cache, payload)| { + let Some(payload) = payload.take() else { + panic!("benchmark payload is available"); + }; + black_box( + cache + .insert(0_u8, start..start + 16, payload) + .expect("benchmark insertion"), + ) + }, + BatchSize::SmallInput, + ); + }, + ); + } + group.finish(); +} + fn eviction(criterion: &mut Criterion) { let mut group = criterion.benchmark_group("eviction"); for ranges in [8, 64, 512] { @@ -199,39 +258,53 @@ fn concurrent_hits(criterion: &mut Criterion) { let mut group = criterion.benchmark_group("concurrent_hit"); group.throughput(Throughput::Elements(1)); for workers in [1, 2, 4, 8] { - let cache = RangeCache::new(CacheCapacity::Unbounded); - for key in 0..workers { - cache - .insert(key, 0..4_096, Bytes::from(vec![1; 4_096])) - .expect("benchmark insert"); - } + for (policy, capacity) in [ + ("unbounded", CacheCapacity::Unbounded), + ( + "bounded", + CacheCapacity::Bounded( + NonZeroUsize::new(workers * 4_096).expect("benchmark capacity is non-zero"), + ), + ), + ] { + let cache = RangeCache::new(capacity); + for key in 0..workers { + cache + .insert(key, 0..4_096, Bytes::from(vec![1; 4_096])) + .expect("benchmark insert"); + } - group.bench_with_input( - BenchmarkId::new("shared_key_workers", workers), - &workers, - |bencher, &workers| { - bencher.iter_custom(|iterations| { - concurrent_hit_duration(&cache, workers, iterations, true) - }); - }, - ); - group.bench_with_input( - BenchmarkId::new("independent_key_workers", workers), - &workers, - |bencher, &workers| { - bencher.iter_custom(|iterations| { - concurrent_hit_duration(&cache, workers, iterations, false) - }); - }, - ); + group.bench_with_input( + BenchmarkId::new(format!("{policy}_shared_key_workers"), workers), + &workers, + |bencher, &workers| { + bencher.iter_custom(|iterations| { + concurrent_hit_duration(&cache, workers, iterations, true) + }); + }, + ); + group.bench_with_input( + BenchmarkId::new(format!("{policy}_independent_key_workers"), workers), + &workers, + |bencher, &workers| { + bencher.iter_custom(|iterations| { + concurrent_hit_duration(&cache, workers, iterations, false) + }); + }, + ); + } } group.finish(); } fn fragmented_cache(ranges: usize) -> RangeCache { + fragmented_cache_from(ranges, 0) +} + +fn fragmented_cache_from(ranges: usize, offset: usize) -> RangeCache { let cache = RangeCache::new(CacheCapacity::Unbounded); for range in 0..ranges { - let start = range * 32; + let start = offset + range * 32; cache .insert(0, start..start + 16, Bytes::from_static(&[1; 16])) .expect("benchmark insert"); @@ -420,6 +493,25 @@ mod asynchronous { BatchSize::SmallInput, ); }); + group.bench_function("partial_single_gap_4096_bytes", |bencher| { + bencher.iter_batched_ref( + || { + let cache = RangeCache::new(CacheCapacity::Unbounded); + cache + .insert(0, 0..2_048, data.slice(0..2_048)) + .expect("benchmark insert"); + reader(Arc::new(Source::immediate(data.clone())), cache, 1) + }, + |reader| { + black_box( + runtime + .block_on(reader.read(&0, 0..4_096)) + .expect("benchmark partial read"), + ) + }, + BatchSize::SmallInput, + ); + }); group.bench_function("warm_4096_bytes", |bencher| { bencher.iter(|| { black_box( @@ -536,6 +628,7 @@ criterion_group!( misses, gap_calculation, insertion, + sparse_insertion, eviction, concurrent_hits, async_benchmarks From e47f39e5bd631f5f0a702f3650a682ce6df715d2 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 11:53:43 +0100 Subject: [PATCH 02/10] Optimize cache eviction metadata --- src/cache.rs | 153 +++++++++++++++++++++++++++++++++------------- tests/property.rs | 121 +++++++++++++++++++++++++++++++++++- 2 files changed, 232 insertions(+), 42 deletions(-) diff --git a/src/cache.rs b/src/cache.rs index 0384b84..7ee282c 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -1,7 +1,7 @@ //! Core cache types. use std::{ - collections::{BTreeMap, BTreeSet}, + collections::BTreeMap, num::NonZeroUsize, ops::{Bound, Range}, sync::Arc, @@ -69,7 +69,7 @@ pub struct CacheSnapshot { struct CacheBlock { end: usize, bytes: Bytes, - last_access: u64, + last_access: Option, } #[derive(Default)] @@ -82,63 +82,94 @@ struct Statistics { evictions: u64, } +enum EvictionPolicy { + Unbounded, + Bounded { + lru: BTreeMap, + next_access: u64, + }, +} + struct State { ranges: BTreeMap>, - lru: BTreeSet<(u64, K, usize)>, + eviction: EvictionPolicy, resident_bytes: usize, - next_access: u64, + resident_ranges: usize, statistics: Statistics, } -impl Default for State { - fn default() -> Self { +impl State { + fn new(capacity: CacheCapacity) -> Self { Self { ranges: BTreeMap::new(), - lru: BTreeSet::new(), + eviction: match capacity { + CacheCapacity::Bounded(_) => EvictionPolicy::Bounded { + lru: BTreeMap::new(), + next_access: 0, + }, + CacheCapacity::Unbounded => EvictionPolicy::Unbounded, + }, resident_bytes: 0, - next_access: 0, + resident_ranges: 0, statistics: Statistics::default(), } } } impl State { - fn take_access(&mut self) -> u64 { - let access = self.next_access; - self.next_access = self - .next_access + fn take_access(&mut self) -> Option { + let EvictionPolicy::Bounded { next_access, .. } = &mut self.eviction else { + return None; + }; + let access = *next_access; + *next_access = next_access .checked_add(1) .expect("range cache LRU clock exhausted"); - access + Some(access) } fn touch(&mut self, key: &K, start: usize) { + if matches!(self.eviction, EvictionPolicy::Unbounded) { + return; + } + let Some(previous) = self .ranges .get(key) .and_then(|ranges| ranges.get(&start)) - .map(|block| block.last_access) + .and_then(|block| block.last_access) else { - return; + panic!("bounded resident range must have an access value"); }; + let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { + unreachable!("unbounded policy returned before touching recency"); + }; + let Some((resident_key, resident_start)) = lru.remove(&previous) else { + panic!("resident range must have an LRU entry"); + }; assert!( - self.lru.remove(&(previous, key.clone(), start)), - "resident range must have an LRU entry" + &resident_key == key && resident_start == start, + "LRU entry must identify the resident range" ); - let next = self.take_access(); + let next = self + .take_access() + .expect("bounded cache provides access values"); self.ranges .get_mut(key) .and_then(|ranges| ranges.get_mut(&start)) .expect("touched range remains resident") - .last_access = next; + .last_access = Some(next); + let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { + unreachable!("bounded policy remains bounded"); + }; assert!( - self.lru.insert((next, key.clone(), start)), + lru.insert(next, (resident_key, start)).is_none(), "LRU access value is unique" ); } - fn remove(&mut self, key: &K, start: usize) -> Option { + fn take_block(&mut self, key: &K, start: usize) -> Option { let (block, key_is_empty) = { let ranges = self.ranges.get_mut(key)?; let block = ranges.remove(&start)?; @@ -148,17 +179,55 @@ impl State { if key_is_empty { self.ranges.remove(key); } - assert!( - self.lru.remove(&(block.last_access, key.clone(), start)), - "removed range must have an LRU entry" - ); self.resident_bytes = self .resident_bytes .checked_sub(block.bytes.len()) .expect("resident byte accounting cannot underflow"); + self.resident_ranges = self + .resident_ranges + .checked_sub(1) + .expect("resident range accounting cannot underflow"); Some(block) } + fn remove(&mut self, key: &K, start: usize) -> Option { + let block = self.take_block(key, start)?; + match (&mut self.eviction, block.last_access) { + (EvictionPolicy::Unbounded, None) => {} + (EvictionPolicy::Bounded { lru, .. }, Some(access)) => { + let Some((resident_key, resident_start)) = lru.remove(&access) else { + panic!("removed range must have an LRU entry"); + }; + assert!( + &resident_key == key && resident_start == start, + "LRU entry must identify the removed range" + ); + } + (EvictionPolicy::Unbounded, Some(_)) | (EvictionPolicy::Bounded { .. }, None) => { + panic!("range access metadata must match the eviction policy"); + } + } + Some(block) + } + + fn evict_oldest(&mut self) -> CacheBlock { + let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { + panic!("only bounded caches evict"); + }; + let Some((access, (key, start))) = lru.pop_first() else { + panic!("resident bytes require an LRU entry"); + }; + let block = self + .take_block(&key, start) + .expect("LRU range remains resident"); + assert_eq!( + block.last_access, + Some(access), + "evicted range access value matches its LRU entry" + ); + block + } + fn covered(&self, key: &K, requested: &Range) -> Vec<(usize, Range)> { let Some(ranges) = self.ranges.get(key) else { return Vec::new(); @@ -209,7 +278,7 @@ impl RangeCache { #[must_use] pub fn new(capacity: CacheCapacity) -> Self { Self { - inner: Arc::new(Mutex::new(State::default())), + inner: Arc::new(Mutex::new(State::new(capacity))), capacity, } } @@ -381,10 +450,15 @@ impl RangeCache { } let access = state.take_access(); - assert!( - state.lru.insert((access, key.clone(), merged_start)), - "new range has a unique LRU entry" - ); + if let Some(access) = access { + let EvictionPolicy::Bounded { lru, .. } = &mut state.eviction else { + unreachable!("access values belong to bounded caches"); + }; + assert!( + lru.insert(access, (key.clone(), merged_start)).is_none(), + "new range has a unique LRU entry" + ); + } let previous = state.ranges.entry(key).or_default().insert( merged_start, CacheBlock { @@ -395,18 +469,12 @@ impl RangeCache { ); assert!(previous.is_none(), "merged range start must be vacant"); state.resident_bytes += merged_length; + state.resident_ranges += 1; state.statistics.insertions += 1; if let CacheCapacity::Bounded(capacity) = self.capacity { while state.resident_bytes > capacity.get() { - let (_, oldest_key, oldest_start) = state - .lru - .first() - .cloned() - .expect("resident bytes require an LRU entry"); - state - .remove(&oldest_key, oldest_start) - .expect("LRU range remains resident"); + let _ = state.evict_oldest(); state.statistics.evictions += 1; } } @@ -444,12 +512,15 @@ impl RangeCache { pub fn clear(&self) -> Invalidation { let mut state = self.inner.lock(); let invalidation = Invalidation { - ranges: state.lru.len(), + ranges: state.resident_ranges, bytes: state.resident_bytes, }; state.ranges.clear(); - state.lru.clear(); + if let EvictionPolicy::Bounded { lru, .. } = &mut state.eviction { + lru.clear(); + } state.resident_bytes = 0; + state.resident_ranges = 0; invalidation } @@ -461,7 +532,7 @@ impl RangeCache { capacity: self.capacity, resident_bytes: state.resident_bytes, keys: state.ranges.len(), - ranges: state.lru.len(), + ranges: state.resident_ranges, hits: state.statistics.hits, partial_hits: state.statistics.partial_hits, misses: state.statistics.misses, diff --git a/tests/property.rs b/tests/property.rs index d43ea0c..99df34d 100644 --- a/tests/property.rs +++ b/tests/property.rs @@ -1,6 +1,8 @@ +use std::num::NonZeroUsize; + use bytes::Bytes; use proptest::prelude::*; -use range_cache::{CacheCapacity, InsertOutcome, RangeCache}; +use range_cache::{CacheCapacity, InsertOutcome, Invalidation, RangeCache}; const KEY_COUNT: usize = 3; const SOURCE_LENGTH: usize = 32; @@ -13,6 +15,14 @@ enum Operation { Invalidate, } +#[derive(Clone, Copy, Debug)] +enum BoundedOperation { + Insert, + Get, + Invalidate, + Clear, +} + fn operation() -> impl Strategy { ( prop_oneof![ @@ -28,6 +38,18 @@ fn operation() -> impl Strategy { ) } +fn bounded_operation() -> impl Strategy { + ( + prop_oneof![ + 4 => Just(BoundedOperation::Insert), + 4 => Just(BoundedOperation::Get), + 2 => Just(BoundedOperation::Invalidate), + 1 => Just(BoundedOperation::Clear), + ], + 0..5_usize, + ) +} + fn model_missing(bytes: &[Option], start: usize, end: usize) -> Vec> { let mut missing = Vec::new(); let mut cursor = start; @@ -136,4 +158,101 @@ proptest! { prop_assert_eq!(snapshot.ranges, model_range_count(&model)); } } + + #[test] + fn bounded_cache_matches_reference_lru( + operations in prop::collection::vec(bounded_operation(), 1..200), + ) { + let cache = RangeCache::new(CacheCapacity::Bounded( + NonZeroUsize::new(8).expect("test capacity is non-zero"), + )); + let mut access_by_key = [None; 5]; + let mut next_access = 0_u64; + let mut hits = 0_u64; + let mut misses = 0_u64; + let mut insertions = 0_u64; + let mut evictions = 0_u64; + + for (operation, key) in operations { + match operation { + BoundedOperation::Insert => { + let outcome = cache + .insert( + key, + 0..4, + Bytes::from(vec![u8::try_from(key).expect("key fits in u8"); 4]), + ) + .expect("valid insert"); + if access_by_key[key].is_some() { + prop_assert_eq!(outcome, InsertOutcome::AlreadyCovered); + } else { + prop_assert_eq!(outcome, InsertOutcome::Inserted); + access_by_key[key] = Some(next_access); + next_access += 1; + insertions += 1; + + if access_by_key.iter().flatten().count() > 2 { + let (oldest_key, _) = access_by_key + .iter() + .enumerate() + .filter_map(|(key, access)| access.map(|access| (key, access))) + .min_by_key(|(_, access)| *access) + .expect("an over-capacity model has an oldest entry"); + access_by_key[oldest_key] = None; + evictions += 1; + } + } + } + BoundedOperation::Get => { + let expected = access_by_key[key] + .is_some() + .then(|| { + Bytes::from(vec![u8::try_from(key).expect("key fits in u8"); 4]) + }); + prop_assert_eq!(cache.get(&key, 0..4).expect("valid range"), expected); + if access_by_key[key].is_some() { + access_by_key[key] = Some(next_access); + next_access += 1; + hits += 1; + } else { + misses += 1; + } + } + BoundedOperation::Invalidate => { + let expected = if access_by_key[key].take().is_some() { + Invalidation { + ranges: 1, + bytes: 4, + } + } else { + Invalidation::default() + }; + prop_assert_eq!(cache.invalidate(&key), expected); + } + BoundedOperation::Clear => { + let ranges = access_by_key.iter().flatten().count(); + prop_assert_eq!( + cache.clear(), + Invalidation { + ranges, + bytes: ranges * 4, + } + ); + access_by_key.fill(None); + } + } + + let ranges = access_by_key.iter().flatten().count(); + let snapshot = cache.snapshot(); + prop_assert_eq!(snapshot.resident_bytes, ranges * 4); + prop_assert_eq!(snapshot.keys, ranges); + prop_assert_eq!(snapshot.ranges, ranges); + prop_assert_eq!(snapshot.hits, hits); + prop_assert_eq!(snapshot.partial_hits, 0); + prop_assert_eq!(snapshot.misses, misses); + prop_assert_eq!(snapshot.insertions, insertions); + prop_assert_eq!(snapshot.admissions_rejected_too_large, 0); + prop_assert_eq!(snapshot.evictions, evictions); + } + } } From f290db7d3c75029854b0825367fcd45b20e8e3e4 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 12:05:10 +0100 Subject: [PATCH 03/10] Bound sparse range traversal --- src/cache.rs | 327 +++++++++++++++++++++++++--------------------- tests/property.rs | 113 +++++++++++++++- 2 files changed, 288 insertions(+), 152 deletions(-) diff --git a/src/cache.rs b/src/cache.rs index 7ee282c..9e038a0 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -69,7 +69,7 @@ pub struct CacheSnapshot { struct CacheBlock { end: usize, bytes: Bytes, - last_access: Option, + last_access: u64, } #[derive(Default)] @@ -90,6 +90,44 @@ enum EvictionPolicy { }, } +impl EvictionPolicy { + fn take_access(&mut self) -> Option { + let Self::Bounded { next_access, .. } = self else { + return None; + }; + let access = *next_access; + *next_access = next_access + .checked_add(1) + .expect("range cache LRU clock exhausted"); + Some(access) + } +} + +impl EvictionPolicy { + #[inline] + fn touch(&mut self, key: &K, start: usize, block: &mut CacheBlock) { + let Self::Bounded { lru, next_access } = self else { + return; + }; + let previous = block.last_access; + let Some((resident_key, resident_start)) = lru.remove(&previous) else { + panic!("resident range must have an LRU entry"); + }; + debug_assert!( + &resident_key == key && resident_start == start, + "LRU entry must identify the resident range" + ); + + let access = *next_access; + *next_access = next_access + .checked_add(1) + .expect("range cache LRU clock exhausted"); + block.last_access = access; + let replaced = lru.insert(access, (resident_key, resident_start)); + debug_assert!(replaced.is_none(), "LRU access value is unique"); + } +} + struct State { ranges: BTreeMap>, eviction: EvictionPolicy, @@ -117,56 +155,19 @@ impl State { } impl State { - fn take_access(&mut self) -> Option { - let EvictionPolicy::Bounded { next_access, .. } = &mut self.eviction else { - return None; - }; - let access = *next_access; - *next_access = next_access - .checked_add(1) - .expect("range cache LRU clock exhausted"); - Some(access) - } - + #[inline] fn touch(&mut self, key: &K, start: usize) { if matches!(self.eviction, EvictionPolicy::Unbounded) { return; } - - let Some(previous) = self - .ranges - .get(key) - .and_then(|ranges| ranges.get(&start)) - .and_then(|block| block.last_access) - else { - panic!("bounded resident range must have an access value"); - }; - - let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { - unreachable!("unbounded policy returned before touching recency"); - }; - let Some((resident_key, resident_start)) = lru.remove(&previous) else { - panic!("resident range must have an LRU entry"); - }; - assert!( - &resident_key == key && resident_start == start, - "LRU entry must identify the resident range" - ); - let next = self - .take_access() - .expect("bounded cache provides access values"); - self.ranges + let Self { + ranges, eviction, .. + } = self; + let block = ranges .get_mut(key) .and_then(|ranges| ranges.get_mut(&start)) - .expect("touched range remains resident") - .last_access = Some(next); - let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { - unreachable!("bounded policy remains bounded"); - }; - assert!( - lru.insert(next, (resident_key, start)).is_none(), - "LRU access value is unique" - ); + .expect("touched range remains resident"); + eviction.touch(key, start, block); } fn take_block(&mut self, key: &K, start: usize) -> Option { @@ -192,20 +193,14 @@ impl State { fn remove(&mut self, key: &K, start: usize) -> Option { let block = self.take_block(key, start)?; - match (&mut self.eviction, block.last_access) { - (EvictionPolicy::Unbounded, None) => {} - (EvictionPolicy::Bounded { lru, .. }, Some(access)) => { - let Some((resident_key, resident_start)) = lru.remove(&access) else { - panic!("removed range must have an LRU entry"); - }; - assert!( - &resident_key == key && resident_start == start, - "LRU entry must identify the removed range" - ); - } - (EvictionPolicy::Unbounded, Some(_)) | (EvictionPolicy::Bounded { .. }, None) => { - panic!("range access metadata must match the eviction policy"); - } + if let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction { + let Some((resident_key, resident_start)) = lru.remove(&block.last_access) else { + panic!("removed range must have an LRU entry"); + }; + assert!( + &resident_key == key && resident_start == start, + "LRU entry must identify the removed range" + ); } Some(block) } @@ -221,41 +216,11 @@ impl State { .take_block(&key, start) .expect("LRU range remains resident"); assert_eq!( - block.last_access, - Some(access), + block.last_access, access, "evicted range access value matches its LRU entry" ); block } - - fn covered(&self, key: &K, requested: &Range) -> Vec<(usize, Range)> { - let Some(ranges) = self.ranges.get(key) else { - return Vec::new(); - }; - - let mut covered = Vec::new(); - match ranges.range(..=requested.start).next_back() { - Some((&start, block)) if block.end > requested.start => { - covered.push((start, requested.start..block.end.min(requested.end))); - } - Some(_) | None => {} - } - - covered.extend( - ranges - .range(( - Bound::Excluded(requested.start), - Bound::Excluded(requested.end), - )) - .map(|(&start, block)| { - ( - start, - start.max(requested.start)..block.end.min(requested.end), - ) - }), - ); - covered - } } /// A cloneable, thread-safe sparse cache of byte ranges keyed by `K`. @@ -318,22 +283,39 @@ impl RangeCache { (start, block.bytes.slice(offset..offset + range.len())) }) }); - if let Some((start, bytes)) = hit { state.statistics.hits += 1; state.touch(key, start); return Ok(Some(bytes)); } - let covered = state.covered(key, &range); - if covered.is_empty() { - state.statistics.misses += 1; - } else { - state.statistics.partial_hits += 1; - for (start, _) in covered { - state.touch(key, start); + let mut has_coverage = false; + { + let State { + ranges, eviction, .. + } = &mut *state; + if let Some(ranges) = ranges.get_mut(key) { + if let Some((&start, block)) = ranges + .range_mut(..=range.start) + .next_back() + .filter(|(_, block)| block.end > range.start) + { + has_coverage = true; + eviction.touch(key, start, block); + } + for (&start, block) in + ranges.range_mut((Bound::Excluded(range.start), Bound::Excluded(range.end))) + { + has_coverage = true; + eviction.touch(key, start, block); + } } } + if has_coverage { + state.statistics.partial_hits += 1; + } else { + state.statistics.misses += 1; + } Ok(None) } @@ -353,8 +335,35 @@ impl RangeCache { } let state = self.inner.lock(); - let covered = state.covered(key, &range); - Ok(missing_ranges(&range, &covered)) + let Some(ranges) = state.ranges.get(key) else { + return Ok(vec![range]); + }; + + let mut missing = Vec::new(); + let mut cursor = range.start; + if let Some((_, block)) = ranges + .range(..=range.start) + .next_back() + .filter(|(_, block)| block.end > range.start) + { + cursor = cursor.max(block.end.min(range.end)); + } + + for (&start, block) in + ranges.range((Bound::Excluded(range.start), Bound::Excluded(range.end))) + { + if cursor < start { + missing.push(cursor..start); + } + cursor = cursor.max(block.end.min(range.end)); + if cursor == range.end { + break; + } + } + if cursor < range.end { + missing.push(cursor..range.end); + } + Ok(missing) } /// Inserts `bytes` for `range`, merging adjacent and overlapping ranges. @@ -407,16 +416,22 @@ impl RangeCache { let mut merged_end = range.end; let mut affected = Vec::new(); if let Some(ranges) = state.ranges.get(&key) { - for (&start, block) in ranges { - if block.end < merged_start { - continue; + let following = match ranges.range(..=range.start).next_back() { + Some((&start, block)) if block.end >= range.start => { + merged_start = start; + merged_end = merged_end.max(block.end); + affected.push((start, block.end)); + Bound::Excluded(start) } + Some(_) | None => Bound::Included(range.start), + }; + + for (&start, block) in ranges.range((following, Bound::Unbounded)) { if start > merged_end { break; } - merged_start = merged_start.min(start); merged_end = merged_end.max(block.end); - affected.push((start, block.end, block.bytes.clone())); + affected.push((start, block.end)); } } @@ -434,7 +449,8 @@ impl RangeCache { bytes } else { let mut merged = vec![0; merged_length]; - for (start, end, cached) in &affected { + for &(start, end) in &affected { + let cached = &state.ranges[&key][&start].bytes; let offset = start - merged_start; merged[offset..offset + (end - start)].copy_from_slice(cached); } @@ -443,13 +459,13 @@ impl RangeCache { Bytes::from(merged) }; - for (start, _, _) in affected { + for (start, _) in affected { state .remove(&key, start) .expect("affected range remains resident"); } - let access = state.take_access(); + let access = state.eviction.take_access(); if let Some(access) = access { let EvictionPolicy::Bounded { lru, .. } = &mut state.eviction else { unreachable!("access values belong to bounded caches"); @@ -464,7 +480,7 @@ impl RangeCache { CacheBlock { end: merged_end, bytes: merged_bytes, - last_access: access, + last_access: access.unwrap_or_default(), }, ); assert!(previous.is_none(), "merged range start must be vacant"); @@ -567,32 +583,59 @@ impl RangeCache { return Ok(ReadPlan::Complete(bytes)); } - let covered = state.covered(key, &range); - let missing = missing_ranges(&range, &covered); - let cached = covered - .iter() - .map(|(start, covered_range)| { - let block = state - .ranges - .get(key) - .and_then(|ranges| ranges.get(start)) - .expect("covered range remains resident"); - let offset = covered_range.start - start; - ( - covered_range.clone(), - block.bytes.slice(offset..offset + covered_range.len()), - ) - }) - .collect(); - - if covered.is_empty() { - state.statistics.misses += 1; - } else { - state.statistics.partial_hits += 1; - for (start, _) in covered { - state.touch(key, start); + let mut cached = Vec::new(); + let mut missing = Vec::new(); + let mut has_coverage = false; + let mut cursor = range.start; + { + let State { + ranges, eviction, .. + } = &mut *state; + if let Some(ranges) = ranges.get_mut(key) { + if let Some((&start, block)) = ranges + .range_mut(..=range.start) + .next_back() + .filter(|(_, block)| block.end > range.start) + { + let covered_end = block.end.min(range.end); + let offset = range.start - start; + cached.push(( + range.start..covered_end, + block + .bytes + .slice(offset..offset + covered_end - range.start), + )); + has_coverage = true; + cursor = cursor.max(covered_end); + eviction.touch(key, start, block); + } + + for (&start, block) in + ranges.range_mut((Bound::Excluded(range.start), Bound::Excluded(range.end))) + { + if cursor < start { + missing.push(cursor..start); + } + let covered_end = block.end.min(range.end); + cached.push((start..covered_end, block.bytes.slice(..covered_end - start))); + has_coverage = true; + cursor = cursor.max(covered_end); + eviction.touch(key, start, block); + if cursor == range.end { + break; + } + } } } + if cursor < range.end { + missing.push(cursor..range.end); + } + + if has_coverage { + state.statistics.partial_hits += 1; + } else { + state.statistics.misses += 1; + } Ok(ReadPlan::Fetch { cached, missing }) } } @@ -607,24 +650,6 @@ fn validate_range(range: &Range) -> Result<(), RangeError> { Ok(()) } -fn missing_ranges( - requested: &Range, - covered: &[(usize, Range)], -) -> Vec> { - let mut missing = Vec::new(); - let mut cursor = requested.start; - for (_, range) in covered { - if cursor < range.start { - missing.push(cursor..range.start); - } - cursor = cursor.max(range.end); - } - if cursor < requested.end { - missing.push(cursor..requested.end); - } - missing -} - #[cfg(feature = "async")] pub(crate) enum ReadPlan { Complete(Bytes), diff --git a/tests/property.rs b/tests/property.rs index 99df34d..a9fecd6 100644 --- a/tests/property.rs +++ b/tests/property.rs @@ -1,4 +1,4 @@ -use std::num::NonZeroUsize; +use std::{collections::BTreeMap, num::NonZeroUsize, ops::Range}; use bytes::Bytes; use proptest::prelude::*; @@ -23,6 +23,14 @@ enum BoundedOperation { Clear, } +#[derive(Clone, Copy, Debug)] +enum SparseOperation { + Insert, + Get, + Missing, + Invalidate, +} + fn operation() -> impl Strategy { ( prop_oneof![ @@ -50,6 +58,20 @@ fn bounded_operation() -> impl Strategy { ) } +fn sparse_operation() -> impl Strategy { + ( + prop_oneof![ + 4 => Just(SparseOperation::Insert), + 2 => Just(SparseOperation::Get), + 2 => Just(SparseOperation::Missing), + 1 => Just(SparseOperation::Invalidate), + ], + 0..1_000_000_usize, + 1..65_usize, + any::(), + ) +} + fn model_missing(bytes: &[Option], start: usize, end: usize) -> Vec> { let mut missing = Vec::new(); let mut cursor = start; @@ -82,6 +104,35 @@ fn model_range_count(model: &[[Option; SOURCE_LENGTH]; KEY_COUNT]) -> usize .sum() } +fn sparse_missing(model: &BTreeMap, range: Range) -> Vec> { + let mut missing = Vec::new(); + let mut cursor = range.start; + while cursor < range.end { + if model.contains_key(&cursor) { + cursor += 1; + continue; + } + let start = cursor; + while cursor < range.end && !model.contains_key(&cursor) { + cursor += 1; + } + missing.push(start..cursor); + } + missing +} + +fn sparse_range_count(model: &BTreeMap) -> usize { + model + .keys() + .scan(None, |previous, &offset| { + let starts_range = previous.is_none_or(|previous| previous + 1 != offset); + *previous = Some(offset); + Some(starts_range) + }) + .filter(|starts_range| *starts_range) + .count() +} + proptest! { #![proptest_config(ProptestConfig::with_cases(256))] @@ -255,4 +306,64 @@ proptest! { prop_assert_eq!(snapshot.evictions, evictions); } } + + #[test] + fn sparse_ranges_match_reference_coverage( + operations in prop::collection::vec(sparse_operation(), 1..100), + ) { + let cache = RangeCache::new(CacheCapacity::Unbounded); + let mut model = BTreeMap::new(); + + for (operation, start, length, value) in operations { + let end = start + length; + match operation { + SparseOperation::Insert => { + let already_covered = (start..end).all(|offset| model.contains_key(&offset)); + let outcome = cache + .insert((), start..end, Bytes::from(vec![value; length])) + .expect("valid insert"); + prop_assert_eq!( + outcome, + if already_covered { + InsertOutcome::AlreadyCovered + } else { + InsertOutcome::Inserted + } + ); + if !already_covered { + model.extend((start..end).map(|offset| (offset, value))); + } + } + SparseOperation::Get => { + let expected = (start..end) + .map(|offset| model.get(&offset).copied()) + .collect::>>() + .map(Bytes::from); + prop_assert_eq!(cache.get(&(), start..end).expect("valid range"), expected); + } + SparseOperation::Missing => { + prop_assert_eq!( + cache.missing_ranges(&(), start..end).expect("valid range"), + sparse_missing(&model, start..end) + ); + } + SparseOperation::Invalidate => { + let ranges = sparse_range_count(&model); + prop_assert_eq!( + cache.invalidate(&()), + Invalidation { + ranges, + bytes: model.len(), + } + ); + model.clear(); + } + } + + let snapshot = cache.snapshot(); + prop_assert_eq!(snapshot.resident_bytes, model.len()); + prop_assert_eq!(snapshot.keys, usize::from(!model.is_empty())); + prop_assert_eq!(snapshot.ranges, sparse_range_count(&model)); + } + } } From 5d225d660d35490a505c630bb062cbf48b8539c8 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 12:08:38 +0100 Subject: [PATCH 04/10] Fast-path single-gap reads --- README.md | 40 ++++++++----------------------- src/reader.rs | 56 +++++++++++++++++++++++++------------------ tests/async_reader.rs | 38 +++++++++++++++++++++++++++++ 3 files changed, 81 insertions(+), 53 deletions(-) diff --git a/README.md b/README.md index bfc55ec..bf95fc8 100644 --- a/README.md +++ b/README.md @@ -92,19 +92,16 @@ async fn main() -> Result<(), Box> { convert::Infallible, num::NonZeroUsize, ops::Range, - sync::{Arc, Mutex}, + sync::Arc, }; use bytes::Bytes; use range_cache::{CacheCapacity, CachedReader, RangeCache, RangeReader, ReaderConfig}; - struct ObjectStore { - data: Bytes, - fetched: Mutex>>, - } + struct MemorySource(Bytes); #[async_trait::async_trait] - impl RangeReader for ObjectStore { + impl RangeReader for MemorySource { type Error = Infallible; async fn read_range( @@ -112,42 +109,25 @@ async fn main() -> Result<(), Box> { _key: &String, range: Range, ) -> Result { - self.fetched - .lock() - .expect("fetch log lock is not poisoned") - .push(range.clone()); - Ok(self.data.slice(range)) + Ok(self.0.slice(range)) } } - let source = Arc::new(ObjectStore { - data: Bytes::from_static(b"abcdefghijklmnop"), - fetched: Mutex::new(Vec::new()), - }); + let cache = RangeCache::new(CacheCapacity::Bounded( + NonZeroUsize::new(1024).expect("capacity is non-zero"), + )); + let source = Arc::new(MemorySource(Bytes::from_static(b"abcdefghijklmnop"))); let reader = CachedReader::new( - Arc::clone(&source), - RangeCache::new(CacheCapacity::Bounded( - NonZeroUsize::new(1024).expect("capacity is non-zero"), - )), + source, + cache, ReaderConfig::new(NonZeroUsize::new(4).expect("concurrency is non-zero")), ); let key = String::from("s3://bucket/index"); - assert_eq!( - reader.read(&key, 0..8).await?, - Bytes::from_static(b"abcdefgh"), - ); assert_eq!( reader.read(&key, 4..12).await?, Bytes::from_static(b"efghijkl"), ); - assert_eq!( - *source - .fetched - .lock() - .expect("fetch log lock is not poisoned"), - vec![0..8, 8..12], - ); Ok(()) } diff --git a/src/reader.rs b/src/reader.rs index 23fae79..2a48aa6 100644 --- a/src/reader.rs +++ b/src/reader.rs @@ -91,8 +91,7 @@ impl Default for InFlightRegistry { } struct InFlightEntry { - gate: AsyncMutex<()>, - successful_response: Mutex>, + response: AsyncMutex>, } impl InFlightRegistry { @@ -102,8 +101,7 @@ impl InFlightRegistry { let current = entries.get(&request).and_then(Weak::upgrade); let entry = current.unwrap_or_else(|| { Arc::new(InFlightEntry { - gate: AsyncMutex::new(()), - successful_response: Mutex::new(None), + response: AsyncMutex::new(None), }) }); entries.insert(request.clone(), Arc::downgrade(&entry)); @@ -209,11 +207,24 @@ where /// Panics only if the privately owned fetch semaphore is unexpectedly /// closed or an internal read-plan coverage invariant is violated. pub async fn read(&self, key: &K, range: Range) -> Result> { - let (mut cached, missing) = match self.cache.read_plan(key, range.clone())? { + let (mut cached, mut missing) = match self.cache.read_plan(key, range.clone())? { ReadPlan::Complete(bytes) => return Ok(bytes), ReadPlan::Fetch { cached, missing } => (cached, missing), }; + if missing.len() == 1 { + let gap = missing.pop().expect("one missing gap exists"); + let fetched = self.fetch_gap(key, gap).await?; + if cached.is_empty() && fetched.0 == range { + return Ok(fetched.1); + } + + let position = + cached.partition_point(|(chunk_range, _)| chunk_range.start < fetched.0.start); + cached.insert(position, fetched); + return Ok(reconstruct(range, cached)); + } + let fetched = stream::iter(missing) .map(|gap| self.fetch_gap(key, gap)) .buffer_unordered(self.config.max_fetch_concurrency.get()) @@ -221,21 +232,7 @@ where .await?; cached.extend(fetched); cached.sort_unstable_by_key(|(chunk_range, _)| chunk_range.start); - - if cached.len() == 1 && cached[0].0 == range { - return Ok(cached.pop().expect("single chunk exists").1); - } - - let mut reconstructed = BytesMut::with_capacity(range.len()); - let mut cursor = range.start; - for (chunk_range, bytes) in cached { - assert_eq!(chunk_range.start, cursor, "read plan has no gaps"); - assert_eq!(chunk_range.len(), bytes.len(), "chunk length is exact"); - reconstructed.extend_from_slice(&bytes); - cursor = chunk_range.end; - } - assert_eq!(cursor, range.end, "read plan covers the request"); - Ok(reconstructed.freeze()) + Ok(reconstruct(range, cached)) } async fn fetch_gap( @@ -248,9 +245,9 @@ where start: range.start, end: range.end, }); - let _request_guard = registration.entry.gate.lock().await; + let mut response = registration.entry.response.lock().await; - if let Some(bytes) = registration.entry.successful_response.lock().clone() { + if let Some(bytes) = response.clone() { return Ok((range, bytes)); } @@ -280,11 +277,24 @@ where let _ = self .cache .insert(key.clone(), range.clone(), bytes.clone())?; - *registration.entry.successful_response.lock() = Some(bytes.clone()); + *response = Some(bytes.clone()); Ok((range, bytes)) } } +fn reconstruct(range: Range, chunks: Vec<(Range, Bytes)>) -> Bytes { + let mut reconstructed = BytesMut::with_capacity(range.len()); + let mut cursor = range.start; + for (chunk_range, bytes) in chunks { + assert_eq!(chunk_range.start, cursor, "read plan has no gaps"); + assert_eq!(chunk_range.len(), bytes.len(), "chunk length is exact"); + reconstructed.extend_from_slice(&bytes); + cursor = chunk_range.end; + } + assert_eq!(cursor, range.end, "read plan covers the request"); + reconstructed.freeze() +} + #[async_trait] impl RangeReader for CachedReader where diff --git a/tests/async_reader.rs b/tests/async_reader.rs index 2541b20..3dbfd68 100644 --- a/tests/async_reader.rs +++ b/tests/async_reader.rs @@ -164,6 +164,44 @@ async fn full_hits_skip_the_source_and_partial_hits_fetch_only_gaps() { assert_eq!(source.calls().len(), 1); } +#[tokio::test] +async fn one_gap_partial_reads_insert_the_fetched_chunk_in_order() { + let source = Arc::new(TestSource::new(b"abcdefghijklmnop")); + let cache = RangeCache::new(CacheCapacity::Unbounded); + cache + .insert(String::from("prefix"), 4..8, Bytes::from_static(b"efgh")) + .expect("valid insert"); + cache + .insert(String::from("middle"), 0..2, Bytes::from_static(b"ab")) + .expect("valid insert"); + cache + .insert(String::from("middle"), 6..8, Bytes::from_static(b"gh")) + .expect("valid insert"); + let reader = reader(Arc::clone(&source), cache, 2); + + assert_eq!( + reader + .read(&String::from("prefix"), 0..8) + .await + .expect("prefix gap read"), + Bytes::from_static(b"abcdefgh") + ); + assert_eq!( + reader + .read(&String::from("middle"), 0..8) + .await + .expect("middle gap read"), + Bytes::from_static(b"abcdefgh") + ); + assert_eq!( + source.calls(), + vec![ + (String::from("prefix"), 0..4), + (String::from("middle"), 2..6), + ] + ); +} + #[tokio::test] async fn missing_gaps_obey_the_global_fetch_concurrency_limit() { let source = Arc::new( From 115e1f98d9b0c258f5fcd4efeeaab2c165de00e4 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 12:27:09 +0100 Subject: [PATCH 05/10] Keep bounded LRU entries stable --- src/cache.rs | 56 ++++++++++++++++++++++++++++++--------------------- tests/core.rs | 35 +++++++++++++++++++++++++++++++- 2 files changed, 67 insertions(+), 24 deletions(-) diff --git a/src/cache.rs b/src/cache.rs index 9e038a0..6d5d626 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -85,21 +85,39 @@ struct Statistics { enum EvictionPolicy { Unbounded, Bounded { - lru: BTreeMap, + // Entries are boxed once, then moved between access keys without + // moving a potentially large K through B-tree nodes on every touch. + lru: BTreeMap>>, next_access: u64, }, } -impl EvictionPolicy { - fn take_access(&mut self) -> Option { - let Self::Bounded { next_access, .. } = self else { - return None; +struct LruEntry { + key: K, + start: usize, +} + +impl EvictionPolicy { + fn register(&mut self, key: &K, start: usize) -> u64 { + let Self::Bounded { lru, next_access } = self else { + return 0; }; let access = *next_access; *next_access = next_access .checked_add(1) .expect("range cache LRU clock exhausted"); - Some(access) + assert!( + lru.insert( + access, + Box::new(LruEntry { + key: key.clone(), + start, + }), + ) + .is_none(), + "new range has a unique LRU entry" + ); + access } } @@ -110,11 +128,11 @@ impl EvictionPolicy { return; }; let previous = block.last_access; - let Some((resident_key, resident_start)) = lru.remove(&previous) else { + let Some(resident) = lru.remove(&previous) else { panic!("resident range must have an LRU entry"); }; debug_assert!( - &resident_key == key && resident_start == start, + &resident.key == key && resident.start == start, "LRU entry must identify the resident range" ); @@ -123,7 +141,7 @@ impl EvictionPolicy { .checked_add(1) .expect("range cache LRU clock exhausted"); block.last_access = access; - let replaced = lru.insert(access, (resident_key, resident_start)); + let replaced = lru.insert(access, resident); debug_assert!(replaced.is_none(), "LRU access value is unique"); } } @@ -194,11 +212,11 @@ impl State { fn remove(&mut self, key: &K, start: usize) -> Option { let block = self.take_block(key, start)?; if let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction { - let Some((resident_key, resident_start)) = lru.remove(&block.last_access) else { + let Some(resident) = lru.remove(&block.last_access) else { panic!("removed range must have an LRU entry"); }; assert!( - &resident_key == key && resident_start == start, + &resident.key == key && resident.start == start, "LRU entry must identify the removed range" ); } @@ -209,9 +227,10 @@ impl State { let EvictionPolicy::Bounded { lru, .. } = &mut self.eviction else { panic!("only bounded caches evict"); }; - let Some((access, (key, start))) = lru.pop_first() else { + let Some((access, resident)) = lru.pop_first() else { panic!("resident bytes require an LRU entry"); }; + let LruEntry { key, start } = *resident; let block = self .take_block(&key, start) .expect("LRU range remains resident"); @@ -465,22 +484,13 @@ impl RangeCache { .expect("affected range remains resident"); } - let access = state.eviction.take_access(); - if let Some(access) = access { - let EvictionPolicy::Bounded { lru, .. } = &mut state.eviction else { - unreachable!("access values belong to bounded caches"); - }; - assert!( - lru.insert(access, (key.clone(), merged_start)).is_none(), - "new range has a unique LRU entry" - ); - } + let access = state.eviction.register(&key, merged_start); let previous = state.ranges.entry(key).or_default().insert( merged_start, CacheBlock { end: merged_end, bytes: merged_bytes, - last_access: access.unwrap_or_default(), + last_access: access, }, ); assert!(previous.is_none(), "merged range start must be vacant"); diff --git a/tests/core.rs b/tests/core.rs index 50905c4..deae6e8 100644 --- a/tests/core.rs +++ b/tests/core.rs @@ -1,4 +1,7 @@ -use std::num::NonZeroUsize; +use std::{ + num::NonZeroUsize, + sync::atomic::{AtomicUsize, Ordering}, +}; use bytes::Bytes; use range_cache::{CacheCapacity, InsertOutcome, Invalidation, RangeCache, RangeError}; @@ -9,6 +12,18 @@ fn bounded(bytes: usize) -> RangeCache<&'static str> { )) } +static KEY_CLONES: AtomicUsize = AtomicUsize::new(0); + +#[derive(Debug, Eq, Ord, PartialEq, PartialOrd)] +struct CountedKey(u8); + +impl Clone for CountedKey { + fn clone(&self) -> Self { + KEY_CLONES.fetch_add(1, Ordering::Relaxed); + Self(self.0) + } +} + #[test] fn cache_is_cloneable_send_and_sync() { fn assert_send_sync() {} @@ -278,6 +293,24 @@ fn bounded_cache_never_exceeds_capacity_and_reads_update_lru() { assert_eq!(snapshot.evictions, 1); } +#[test] +fn bounded_hits_reuse_the_resident_lru_key() { + let cache = RangeCache::new(CacheCapacity::Bounded( + NonZeroUsize::new(8).expect("test capacity is non-zero"), + )); + let key = CountedKey(1); + cache + .insert(key.clone(), 0..4, Bytes::from_static(b"data")) + .expect("valid insert"); + KEY_CLONES.store(0, Ordering::Relaxed); + + assert_eq!( + cache.get(&key, 0..4).expect("valid range"), + Some(Bytes::from_static(b"data")) + ); + assert_eq!(KEY_CLONES.load(Ordering::Relaxed), 0); +} + #[test] fn oversized_merge_does_not_mutate_existing_state() { let cache = bounded(8); From ae685c70d85984a54f8e7262328febbc6ac5976e Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 12:55:34 +0100 Subject: [PATCH 06/10] Document cache performance gains --- CHANGELOG.md | 4 ++++ README.md | 22 +++++++++++++--------- 2 files changed, 17 insertions(+), 9 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 87b35c4..b635bc8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,10 @@ project follows [Semantic Versioning](https://semver.org/). - Raised the minimum supported Rust version from 1.85 to 1.86. - Updated canonical repository links after the project ownership transfer. +- Removed recency bookkeeping from unbounded caches and made bounded touches + reuse stable access-keyed LRU entries. +- Bounded sparse insertion and read planning to relevant ordered ranges. +- Added direct one-gap async reads and consolidated in-flight response locking. ## [0.1.0] - 2026-07-21 diff --git a/README.md b/README.md index bf95fc8..466b07a 100644 --- a/README.md +++ b/README.md @@ -144,20 +144,24 @@ latency is reported instead of apparent byte throughput because the returned | Operation | Workload | Median estimate | | --- | --- | ---: | -| Full hit | 32 KiB requested from a 64 KiB cached range | 19.89 ns | -| Gap calculation | 64 resident ranges | 498.41 ns | -| Overlapping insertion | Merge across 64 resident ranges | 3.16 µs | -| Eviction | Insert with 64 resident ranges at capacity | 288.34 ns | -| Concurrent hit | 8 workers sharing one key | 73.21 ns/read | -| Warm read-through | 4 KiB cached read | 85.80 ns | -| Fragmented reconstruction | 64 alternating cached/missing segments | 18.00 µs | -| Coalesced read-through | 32 identical concurrent readers | 11.17 µs | +| Unbounded full hit | 32 KiB requested from a 64 KiB cached range | 10.04 ns | +| Bounded full hit | 32 KiB requested from a 64 KiB cached range | 18.15 ns | +| Gap calculation | 64 resident ranges | 237.17 ns | +| Overlapping insertion | Merge across 64 resident ranges | 2.27 µs | +| Sparse end insertion | 512 resident ranges | 90.25 ns | +| Eviction | Insert with 64 resident ranges at capacity | 180.85 ns | +| Concurrent hit | 8 workers sharing one unbounded key | 47.87 ns/read | +| Cold read-through | One 4 KiB missing range | 282.98 ns | +| Partial read-through | One 2 KiB gap in a 4 KiB read | 840.30 ns | +| Warm read-through | 4 KiB cached read | 76.89 ns | +| Fragmented reconstruction | 64 alternating cached/missing segments | 13.80 µs | +| Coalesced read-through | 32 identical concurrent readers | 4.77 µs | The coalesced 32-reader case performs one 4 KiB source read; issuing those reads directly would perform 32 calls and fetch 128 KiB. Measured with `cargo bench --all-features --bench range_cache -- --noplot` on -commit `4e9e7baa61ede11bf70b0d246954f3187c1328d8` using `rustc 1.97.1` on macOS +commit `115e1f98d9b0c258f5fcd4efeeaab2c165de00e4` using `rustc 1.97.1` on macOS 26.5, an Apple M4 Pro (14 cores), and 24 GiB of memory. These numbers describe that machine and revision; they are not cross-platform performance guarantees. From 07fc6815c8b59f9728ed66fac882c98bcf37cf26 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 13:28:41 +0100 Subject: [PATCH 07/10] Harden cache failure and concurrency tests --- src/cache.rs | 116 +++++++++++++ src/reader.rs | 160 ++++++++++++++++-- tests/async_reader.rs | 368 +++++++++++++++++++++++++++++++++++++++--- tests/core.rs | 143 +++++++++++++++- 4 files changed, 753 insertions(+), 34 deletions(-) diff --git a/src/cache.rs b/src/cache.rs index 6d5d626..e6c843c 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -668,3 +668,119 @@ pub(crate) enum ReadPlan { missing: Vec>, }, } + +#[cfg(test)] +mod tests { + use std::{num::NonZeroUsize, sync::Arc}; + + use bytes::Bytes; + use parking_lot::Mutex; + + use super::{CacheBlock, CacheCapacity, EvictionPolicy, RangeCache, State}; + + fn bounded_state() -> State<&'static str> { + State::new(CacheCapacity::Bounded( + NonZeroUsize::new(8).expect("test capacity is non-zero"), + )) + } + + fn add_untracked_block(state: &mut State<&'static str>) { + state.ranges.entry("key").or_default().insert( + 0, + CacheBlock { + end: 1, + bytes: Bytes::from_static(b"x"), + last_access: 0, + }, + ); + state.resident_bytes = 1; + state.resident_ranges = 1; + } + + #[test] + #[should_panic(expected = "resident range must have an LRU entry")] + fn touching_an_untracked_bounded_range_panics() { + let mut state = bounded_state(); + add_untracked_block(&mut state); + state.touch(&"key", 0); + } + + #[test] + #[should_panic(expected = "removed range must have an LRU entry")] + fn removing_an_untracked_bounded_range_panics() { + let mut state = bounded_state(); + add_untracked_block(&mut state); + let _ = state.remove(&"key", 0); + } + + #[test] + #[should_panic(expected = "only bounded caches evict")] + fn unbounded_state_cannot_evict() { + let mut state = State::::new(CacheCapacity::Unbounded); + let _ = state.evict_oldest(); + } + + #[test] + #[should_panic(expected = "resident bytes require an LRU entry")] + fn empty_bounded_state_cannot_evict() { + let mut state = State::::new(CacheCapacity::Bounded( + NonZeroUsize::new(1).expect("test capacity is non-zero"), + )); + let _ = state.evict_oldest(); + } + + #[test] + #[should_panic(expected = "range cache LRU clock exhausted")] + fn registering_after_lru_clock_exhaustion_panics() { + let mut eviction = EvictionPolicy::Bounded { + lru: Default::default(), + next_access: u64::MAX, + }; + let _ = eviction.register(&"key", 0); + } + + #[test] + fn absent_internal_ranges_are_noops() { + let mut state = State::<&str>::new(CacheCapacity::Unbounded); + let _ = state.take_block(&"missing", 0); + state.ranges.insert("empty", Default::default()); + let _ = state.take_block(&"empty", 0); + let _ = state.remove(&"missing", 0); + assert_eq!(state.resident_bytes, 0); + assert_eq!(state.resident_ranges, 0); + } + + #[test] + fn overlapping_internal_ranges_do_not_create_negative_gaps() { + let mut state = State::new(CacheCapacity::Unbounded); + state.ranges.entry("key").or_default().insert( + 0, + CacheBlock { + end: 4, + bytes: Bytes::from_static(b"abcd"), + last_access: 0, + }, + ); + state.ranges.entry("key").or_default().insert( + 2, + CacheBlock { + end: 6, + bytes: Bytes::from_static(b"cdef"), + last_access: 0, + }, + ); + state.resident_bytes = 8; + state.resident_ranges = 2; + let cache = RangeCache { + inner: Arc::new(Mutex::new(state)), + capacity: CacheCapacity::Unbounded, + }; + + assert_eq!( + cache.missing_ranges(&"key", 0..6).expect("valid range"), + Vec::>::new() + ); + #[cfg(feature = "async")] + let _ = cache.read_plan(&"key", 0..6).expect("valid read plan"); + } +} diff --git a/src/reader.rs b/src/reader.rs index 2a48aa6..4e55783 100644 --- a/src/reader.rs +++ b/src/reader.rs @@ -276,7 +276,8 @@ where let _ = self .cache - .insert(key.clone(), range.clone(), bytes.clone())?; + .insert(key.clone(), range.clone(), bytes.clone()) + .expect("validated source bytes match the requested range"); *response = Some(bytes.clone()); Ok((range, bytes)) } @@ -310,49 +311,80 @@ where #[cfg(test)] mod tests { - use std::{convert::Infallible, num::NonZeroUsize, sync::Arc, time::Duration}; + use std::{ + num::NonZeroUsize, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + }; use async_trait::async_trait; use bytes::Bytes; use tokio::sync::Semaphore; - use super::{CachedReader, RangeReader, ReaderConfig}; + use super::{CachedReader, RangeReader, ReadError, ReaderConfig}; use crate::{CacheCapacity, RangeCache}; - struct BlockingSource { + #[derive(Clone, Copy, Debug, thiserror::Error)] + #[error("controlled source failure")] + struct TestError; + + struct ControlledSource { started: Semaphore, + release: Semaphore, + calls: AtomicUsize, + fail: bool, } #[async_trait] - impl RangeReader for BlockingSource { - type Error = Infallible; + impl RangeReader for ControlledSource { + type Error = TestError; async fn read_range( &self, _key: &String, range: std::ops::Range, ) -> Result { + self.calls.fetch_add(1, Ordering::SeqCst); self.started.add_permits(1); - tokio::time::sleep(Duration::from_secs(60)).await; + self.release + .acquire() + .await + .expect("source release semaphore remains open") + .forget(); + if self.fail { + return Err(TestError); + } Ok(Bytes::from(vec![0; range.len()])) } } + fn source(fail: bool) -> Arc { + Arc::new(ControlledSource { + started: Semaphore::new(0), + release: Semaphore::new(0), + calls: AtomicUsize::new(0), + fail, + }) + } + + async fn read_key( + reader: CachedReader, + ) -> Result> { + reader.read(&String::from("key"), 0..4).await + } + #[tokio::test] async fn cancelled_leader_removes_its_in_flight_registration() { - let source = Arc::new(BlockingSource { - started: Semaphore::new(0), - }); + let source = source(false); let reader = CachedReader::new( Arc::clone(&source), RangeCache::new(CacheCapacity::Unbounded), ReaderConfig::new(NonZeroUsize::new(1).expect("non-zero")), ); let task_reader = reader.clone(); - let task = tokio::spawn(async move { - let key = String::from("key"); - task_reader.read(&key, 0..4).await - }); + let task = tokio::spawn(read_key(task_reader)); source .started .acquire() @@ -365,4 +397,104 @@ mod tests { tokio::task::yield_now().await; assert!(reader.in_flight.entries.lock().is_empty()); } + + #[tokio::test] + async fn completed_requests_remove_their_in_flight_registration() { + let source = source(false); + let reader = CachedReader::new( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + ReaderConfig::new(NonZeroUsize::new(1).expect("non-zero")), + ); + let task_reader = reader.clone(); + let task = tokio::spawn(read_key(task_reader)); + source + .started + .acquire() + .await + .expect("source semaphore remains open") + .forget(); + source.release.add_permits(1); + assert_eq!( + task.await + .expect("task completed") + .expect("source read succeeds"), + Bytes::from_static(&[0; 4]) + ); + assert!(reader.in_flight.entries.lock().is_empty()); + } + + #[tokio::test] + async fn failed_requests_remove_their_in_flight_registration() { + let source = source(true); + let reader = CachedReader::new( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + ReaderConfig::new(NonZeroUsize::new(1).expect("non-zero")), + ); + let task_reader = reader.clone(); + let task = tokio::spawn(read_key(task_reader)); + source + .started + .acquire() + .await + .expect("source semaphore remains open") + .forget(); + source.release.add_permits(1); + assert_eq!( + task.await + .expect("task completed") + .expect_err("source read fails") + .to_string(), + "range source failed: controlled source failure" + ); + assert!(reader.in_flight.entries.lock().is_empty()); + } + + #[tokio::test] + async fn fetch_gap_rechecks_the_cache_before_reading_the_source() { + let source = source(false); + let cache = RangeCache::new(CacheCapacity::Unbounded); + let key = String::from("key"); + cache + .insert(key.clone(), 0..4, Bytes::from_static(b"data")) + .expect("valid insert"); + let reader = CachedReader::new( + Arc::clone(&source), + cache, + ReaderConfig::new(NonZeroUsize::new(1).expect("non-zero")), + ); + + assert_eq!( + reader + .fetch_gap(&key, 0..4) + .await + .expect("cached gap succeeds"), + (0..4, Bytes::from_static(b"data")) + ); + assert_eq!(source.calls.load(Ordering::SeqCst), 0); + assert!(reader.in_flight.entries.lock().is_empty()); + } + + #[tokio::test] + async fn fetch_gap_propagates_range_validation_errors() { + let source = source(false); + let reader = CachedReader::new( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + ReaderConfig::new(NonZeroUsize::new(1).expect("non-zero")), + ); + let reversed = std::ops::Range { start: 4, end: 3 }; + + assert_eq!( + reader + .fetch_gap(&String::from("key"), reversed) + .await + .expect_err("reversed range fails") + .to_string(), + "reversed byte range 4..3" + ); + assert_eq!(source.calls.load(Ordering::SeqCst), 0); + assert!(reader.in_flight.entries.lock().is_empty()); + } } diff --git a/tests/async_reader.rs b/tests/async_reader.rs index 3dbfd68..3944f67 100644 --- a/tests/async_reader.rs +++ b/tests/async_reader.rs @@ -1,13 +1,13 @@ #![cfg(feature = "async")] use std::{ + collections::BTreeMap, num::NonZeroUsize, ops::Range, sync::{ Arc, Mutex, atomic::{AtomicUsize, Ordering}, }, - time::Duration, }; use async_trait::async_trait; @@ -31,9 +31,11 @@ impl Drop for ActiveRead<'_> { struct TestSource { data: Bytes, - delay: Duration, failures_remaining: AtomicUsize, short_reads_remaining: AtomicUsize, + long_reads_remaining: AtomicUsize, + failures_by_range: Mutex)>>, + gates: Mutex>>, calls: Mutex)>>, active: AtomicUsize, max_active: AtomicUsize, @@ -44,9 +46,11 @@ impl TestSource { fn new(data: &'static [u8]) -> Self { Self { data: Bytes::from_static(data), - delay: Duration::ZERO, failures_remaining: AtomicUsize::new(0), short_reads_remaining: AtomicUsize::new(0), + long_reads_remaining: AtomicUsize::new(0), + failures_by_range: Mutex::new(Vec::new()), + gates: Mutex::new(BTreeMap::new()), calls: Mutex::new(Vec::new()), active: AtomicUsize::new(0), max_active: AtomicUsize::new(0), @@ -54,11 +58,6 @@ impl TestSource { } } - fn with_delay(mut self, delay: Duration) -> Self { - self.delay = delay; - self - } - fn with_failures(self, failures: usize) -> Self { self.failures_remaining.store(failures, Ordering::SeqCst); self @@ -70,6 +69,40 @@ impl TestSource { self } + fn with_long_reads(self, long_reads: usize) -> Self { + self.long_reads_remaining + .store(long_reads, Ordering::SeqCst); + self + } + + fn with_failure_for(self, key: &str, range: Range) -> Self { + self.failures_by_range + .lock() + .expect("failure script lock is not poisoned") + .push((String::from(key), range)); + self + } + + fn with_gate(self, key: &str, range: Range) -> Self { + self.gates + .lock() + .expect("gate lock is not poisoned") + .insert( + (String::from(key), range.start, range.end), + Arc::new(tokio::sync::Semaphore::new(0)), + ); + self + } + + fn release(&self, key: &str, range: Range) { + self.gates + .lock() + .expect("gate lock is not poisoned") + .get(&(String::from(key), range.start, range.end)) + .expect("scripted source gate exists") + .add_permits(1); + } + fn calls(&self) -> Vec<(String, Range)> { self.calls .lock() @@ -110,15 +143,43 @@ impl RangeReader for TestSource { self.max_active.fetch_max(active, Ordering::SeqCst); let _active_read = ActiveRead(&self.active); - tokio::time::sleep(self.delay).await; + let gate = self + .gates + .lock() + .expect("gate lock is not poisoned") + .get(&(key.clone(), range.start, range.end)) + .cloned(); + if let Some(gate) = gate { + gate.acquire() + .await + .expect("source gate remains open") + .forget(); + } if Self::take(&self.failures_remaining) { return Err(TestError); } + let mut failures_by_range = self + .failures_by_range + .lock() + .expect("failure script lock is not poisoned"); + if let Some(index) = failures_by_range + .iter() + .position(|failure| failure == &(key.clone(), range.clone())) + { + failures_by_range.remove(index); + return Err(TestError); + } + drop(failures_by_range); let response = self.data.slice(range); if Self::take(&self.short_reads_remaining) && !response.is_empty() { return Ok(response.slice(..response.len() - 1)); } + if Self::take(&self.long_reads_remaining) { + let mut longer = response.to_vec(); + longer.push(0); + return Ok(Bytes::from(longer)); + } Ok(response) } } @@ -202,10 +263,94 @@ async fn one_gap_partial_reads_insert_the_fetched_chunk_in_order() { ); } +#[tokio::test] +async fn multiple_gaps_complete_out_of_order_and_reconstruct_exactly() { + let source = Arc::new( + TestSource::new(b"abcdefghijkl") + .with_gate("key", 2..4) + .with_gate("key", 6..8), + ); + let cache = RangeCache::new(CacheCapacity::Unbounded); + for (range, bytes) in [ + (0..2, Bytes::from_static(b"ab")), + (4..6, Bytes::from_static(b"ef")), + (8..12, Bytes::from_static(b"ijkl")), + ] { + cache + .insert(String::from("key"), range, bytes) + .expect("valid insert"); + } + let reader = reader(Arc::clone(&source), cache, 2); + let task_reader = reader.clone(); + let task = tokio::spawn(async move { task_reader.read(&String::from("key"), 0..12).await }); + + source.wait_for_calls(2).await; + assert_eq!(source.max_active.load(Ordering::SeqCst), 2); + source.release("key", 6..8); + tokio::task::yield_now().await; + source.release("key", 2..4); + + assert_eq!( + task.await + .expect("task completed") + .expect("multi-gap read succeeds"), + Bytes::from_static(b"abcdefghijkl") + ); + assert_eq!( + source.calls(), + vec![(String::from("key"), 2..4), (String::from("key"), 6..8),] + ); +} + +#[tokio::test] +async fn completed_gaps_remain_cached_when_another_gap_fails() { + let source = Arc::new(TestSource::new(b"abcdefghij").with_failure_for("key", 6..8)); + let cache = RangeCache::new(CacheCapacity::Unbounded); + for (range, bytes) in [ + (0..2, Bytes::from_static(b"ab")), + (4..6, Bytes::from_static(b"ef")), + (8..10, Bytes::from_static(b"ij")), + ] { + cache + .insert(String::from("key"), range, bytes) + .expect("valid insert"); + } + let reader = reader(Arc::clone(&source), cache, 2); + let key = String::from("key"); + + assert!(matches!( + reader.read(&key, 0..10).await, + Err(ReadError::Source(TestError)) + )); + assert_eq!( + reader + .cache() + .missing_ranges(&key, 0..10) + .expect("valid range"), + vec![6..8] + ); + assert_eq!( + reader.read(&key, 0..10).await.expect("retry succeeds"), + Bytes::from_static(b"abcdefghij") + ); + assert_eq!( + source.calls(), + vec![ + (String::from("key"), 2..4), + (String::from("key"), 6..8), + (String::from("key"), 6..8), + ] + ); +} + #[tokio::test] async fn missing_gaps_obey_the_global_fetch_concurrency_limit() { let source = Arc::new( - TestSource::new(b"abcdefghijklmnopqrstuvwxyz012345").with_delay(Duration::from_millis(30)), + TestSource::new(b"abcdefghijklmnopqrstuvwxyz012345") + .with_gate("first", 0..8) + .with_gate("second", 8..16) + .with_gate("third", 16..24) + .with_gate("fourth", 24..32), ); let reader = reader( Arc::clone(&source), @@ -223,6 +368,15 @@ async fn missing_gaps_obey_the_global_fetch_concurrency_limit() { let key = String::from(key); tasks.spawn(async move { task_reader.read(&key, range).await }); } + source.wait_for_calls(2).await; + assert_eq!(source.max_active.load(Ordering::SeqCst), 2); + for (key, range) in source.calls() { + source.release(&key, range); + } + source.wait_for_calls(4).await; + for (key, range) in source.calls().into_iter().skip(2) { + source.release(&key, range); + } while let Some(result) = tasks.join_next().await { result.expect("task completed").expect("source read"); } @@ -233,8 +387,7 @@ async fn missing_gaps_obey_the_global_fetch_concurrency_limit() { #[tokio::test] async fn identical_key_and_gap_requests_are_coalesced() { - let source = - Arc::new(TestSource::new(b"abcdefghijklmnop").with_delay(Duration::from_millis(40))); + let source = Arc::new(TestSource::new(b"abcdefghijklmnop").with_gate("key", 0..8)); let reader = reader( Arc::clone(&source), RangeCache::new(CacheCapacity::Unbounded), @@ -250,6 +403,8 @@ async fn identical_key_and_gap_requests_are_coalesced() { .expect("coalesced read") }); } + source.wait_for_calls(1).await; + source.release("key", 0..8); while let Some(result) = tasks.join_next().await { assert_eq!( result.expect("task completed"), @@ -261,8 +416,11 @@ async fn identical_key_and_gap_requests_are_coalesced() { #[tokio::test] async fn merely_overlapping_requests_are_not_coalesced() { - let source = - Arc::new(TestSource::new(b"abcdefghijklmnop").with_delay(Duration::from_millis(40))); + let source = Arc::new( + TestSource::new(b"abcdefghijklmnop") + .with_gate("key", 0..8) + .with_gate("key", 4..12), + ); let reader = reader( Arc::clone(&source), RangeCache::new(CacheCapacity::Unbounded), @@ -282,6 +440,9 @@ async fn merely_overlapping_requests_are_not_coalesced() { .await .expect("second read") }); + source.wait_for_calls(2).await; + source.release("key", 0..8); + source.release("key", 4..12); assert_eq!( first.await.expect("first task"), Bytes::from_static(b"abcdefgh") @@ -293,6 +454,38 @@ async fn merely_overlapping_requests_are_not_coalesced() { assert_eq!(source.calls().len(), 2); } +#[tokio::test] +async fn identical_ranges_for_different_keys_are_not_coalesced() { + let source = Arc::new( + TestSource::new(b"abcdefgh") + .with_gate("first", 0..8) + .with_gate("second", 0..8), + ); + let reader = reader( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + 2, + ); + let first_reader = reader.clone(); + let second_reader = reader.clone(); + let first = tokio::spawn(async move { first_reader.read(&String::from("first"), 0..8).await }); + let second = + tokio::spawn(async move { second_reader.read(&String::from("second"), 0..8).await }); + + source.wait_for_calls(2).await; + source.release("first", 0..8); + source.release("second", 0..8); + assert_eq!( + first.await.expect("first task").expect("first read"), + Bytes::from_static(b"abcdefgh") + ); + assert_eq!( + second.await.expect("second task").expect("second read"), + Bytes::from_static(b"abcdefgh") + ); + assert_eq!(source.calls().len(), 2); +} + #[tokio::test] async fn source_failures_are_not_cached_and_can_be_retried() { let source = Arc::new(TestSource::new(b"abcdefgh").with_failures(1)); @@ -315,6 +508,50 @@ async fn source_failures_are_not_cached_and_can_be_retried() { assert_eq!(source.calls().len(), 2); } +#[tokio::test] +async fn failed_coalesced_leader_allows_a_waiter_to_retry() { + let source = Arc::new( + TestSource::new(b"abcdefgh") + .with_failures(1) + .with_gate("key", 0..8), + ); + let reader = reader( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + 1, + ); + let first_reader = reader.clone(); + let second_reader = reader.clone(); + let first = tokio::spawn(async move { first_reader.read(&String::from("key"), 0..8).await }); + let second = tokio::spawn(async move { second_reader.read(&String::from("key"), 0..8).await }); + + source.wait_for_calls(1).await; + source.release("key", 0..8); + source.wait_for_calls(2).await; + source.release("key", 0..8); + let results = [ + first.await.expect("first task"), + second.await.expect("second task"), + ]; + assert_eq!( + results + .iter() + .filter(|result| matches!(result, Err(ReadError::Source(TestError)))) + .count(), + 1 + ); + assert_eq!( + results + .iter() + .filter(|result| { + matches!(result, Ok(bytes) if bytes == &Bytes::from_static(b"abcdefgh")) + }) + .count(), + 1 + ); + assert_eq!(source.calls().len(), 2); +} + #[tokio::test] async fn short_reads_are_not_cached_and_can_be_retried() { let source = Arc::new(TestSource::new(b"abcdefgh").with_short_reads(1)); @@ -341,6 +578,32 @@ async fn short_reads_are_not_cached_and_can_be_retried() { assert_eq!(source.calls().len(), 2); } +#[tokio::test] +async fn long_reads_are_not_cached_and_can_be_retried() { + let source = Arc::new(TestSource::new(b"abcdefgh").with_long_reads(1)); + let reader = reader( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + 1, + ); + let key = String::from("key"); + + assert!(matches!( + reader.read(&key, 0..8).await, + Err(ReadError::ShortRead { + range, + expected: 8, + actual: 9, + }) if range == (0..8) + )); + assert_eq!(reader.cache().snapshot().resident_bytes, 0); + assert_eq!( + reader.read(&key, 0..8).await.expect("retry succeeds"), + Bytes::from_static(b"abcdefgh") + ); + assert_eq!(source.calls().len(), 2); +} + #[tokio::test] async fn fetched_bytes_are_returned_when_cache_admission_is_too_large() { let source = Arc::new(TestSource::new(b"abcdefgh")); @@ -363,7 +626,7 @@ async fn fetched_bytes_are_returned_when_cache_admission_is_too_large() { #[tokio::test] async fn identical_waiters_share_a_successful_response_that_is_too_large_to_cache() { - let source = Arc::new(TestSource::new(b"abcdefgh").with_delay(Duration::from_millis(40))); + let source = Arc::new(TestSource::new(b"abcdefgh").with_gate("key", 0..8)); let cache = RangeCache::new(CacheCapacity::Bounded( NonZeroUsize::new(4).expect("non-zero"), )); @@ -382,6 +645,8 @@ async fn identical_waiters_share_a_successful_response_that_is_too_large_to_cach .await .expect("shared response") }); + source.wait_for_calls(1).await; + source.release("key", 0..8); assert_eq!( first.await.expect("first task"), @@ -397,7 +662,7 @@ async fn identical_waiters_share_a_successful_response_that_is_too_large_to_cach #[tokio::test] async fn cancelled_leader_allows_an_identical_waiter_to_retry() { - let source = Arc::new(TestSource::new(b"abcdefgh").with_delay(Duration::from_millis(100))); + let source = Arc::new(TestSource::new(b"abcdefgh").with_gate("key", 0..8)); let reader = reader( Arc::clone(&source), RangeCache::new(CacheCapacity::Unbounded), @@ -414,9 +679,11 @@ async fn cancelled_leader_allows_an_identical_waiter_to_retry() { .await .expect("waiter retries") }); - tokio::time::sleep(Duration::from_millis(10)).await; + tokio::task::yield_now().await; leader.abort(); assert!(leader.await.expect_err("leader cancelled").is_cancelled()); + source.wait_for_calls(2).await; + source.release("key", 0..8); assert_eq!( waiter.await.expect("waiter task"), Bytes::from_static(b"abcdefgh") @@ -426,7 +693,7 @@ async fn cancelled_leader_allows_an_identical_waiter_to_retry() { #[tokio::test] async fn cancelled_waiter_does_not_cancel_the_leader() { - let source = Arc::new(TestSource::new(b"abcdefgh").with_delay(Duration::from_millis(80))); + let source = Arc::new(TestSource::new(b"abcdefgh").with_gate("key", 0..8)); let reader = reader( Arc::clone(&source), RangeCache::new(CacheCapacity::Unbounded), @@ -443,9 +710,10 @@ async fn cancelled_waiter_does_not_cancel_the_leader() { let waiter_reader = reader.clone(); let waiter = tokio::spawn(async move { waiter_reader.read(&String::from("key"), 0..8).await }); - tokio::time::sleep(Duration::from_millis(10)).await; + tokio::task::yield_now().await; waiter.abort(); assert!(waiter.await.expect_err("waiter cancelled").is_cancelled()); + source.release("key", 0..8); assert_eq!( leader.await.expect("leader task"), Bytes::from_static(b"abcdefgh") @@ -453,6 +721,51 @@ async fn cancelled_waiter_does_not_cancel_the_leader() { assert_eq!(source.calls().len(), 1); } +#[tokio::test] +async fn cancelling_a_fetch_permit_waiter_does_not_consume_a_permit() { + let source = Arc::new( + TestSource::new(b"abcdefgh") + .with_gate("first", 0..4) + .with_gate("second", 4..8), + ); + let reader = reader( + Arc::clone(&source), + RangeCache::new(CacheCapacity::Unbounded), + 1, + ); + let first_reader = reader.clone(); + let first = tokio::spawn(async move { first_reader.read(&String::from("first"), 0..4).await }); + source.wait_for_calls(1).await; + + let cancelled_reader = reader.clone(); + let cancelled = + tokio::spawn(async move { cancelled_reader.read(&String::from("second"), 4..8).await }); + tokio::task::yield_now().await; + assert_eq!(source.calls().len(), 1); + cancelled.abort(); + assert!( + cancelled + .await + .expect_err("waiting task was cancelled") + .is_cancelled() + ); + + source.release("first", 0..4); + assert_eq!( + first.await.expect("first task").expect("first read"), + Bytes::from_static(b"abcd") + ); + + let retry_reader = reader.clone(); + let retry = tokio::spawn(async move { retry_reader.read(&String::from("second"), 4..8).await }); + source.wait_for_calls(2).await; + source.release("second", 4..8); + assert_eq!( + retry.await.expect("retry task").expect("retry read"), + Bytes::from_static(b"efgh") + ); +} + #[tokio::test] async fn empty_ranges_succeed_without_source_access_and_reversed_ranges_fail() { let source = Arc::new(TestSource::new(b"abcdefgh")); @@ -495,3 +808,20 @@ async fn range_reader_is_object_safe() { Bytes::from_static(b"abcd") ); } + +#[test] +fn reader_accessors_return_constructor_values() { + let source = Arc::new(TestSource::new(b"abcdefgh")); + let capacity = CacheCapacity::Bounded(NonZeroUsize::new(8).expect("test capacity is non-zero")); + let config = ReaderConfig::new(NonZeroUsize::new(3).expect("test concurrency is non-zero")); + let reader = CachedReader::new( + Arc::clone(&source), + RangeCache::::new(capacity), + config, + ); + + assert!(Arc::ptr_eq(reader.source(), &source)); + assert_eq!(reader.cache().capacity(), capacity); + assert_eq!(reader.config(), config); + assert_eq!(config.max_fetch_concurrency().get(), 3); +} diff --git a/tests/core.rs b/tests/core.rs index deae6e8..b005f87 100644 --- a/tests/core.rs +++ b/tests/core.rs @@ -1,6 +1,10 @@ use std::{ num::NonZeroUsize, - sync::atomic::{AtomicUsize, Ordering}, + sync::{ + Arc, Barrier, + atomic::{AtomicUsize, Ordering}, + }, + thread, }; use bytes::Bytes; @@ -40,6 +44,16 @@ fn cache_is_cloneable_send_and_sync() { ); } +#[test] +fn capacity_reports_the_configured_policy() { + let bounded = CacheCapacity::Bounded(NonZeroUsize::new(8).expect("test capacity is non-zero")); + assert_eq!(RangeCache::::new(bounded).capacity(), bounded); + assert_eq!( + RangeCache::::new(CacheCapacity::Unbounded).capacity(), + CacheCapacity::Unbounded + ); +} + #[test] fn empty_ranges_succeed_and_matching_empty_inserts_are_noops() { let cache = bounded(8); @@ -90,6 +104,44 @@ fn reversed_ranges_and_payload_mismatches_are_errors() { actual: 1, }) ); + assert_eq!( + cache.insert("key", std::ops::Range { start: 5, end: 4 }, Bytes::new()), + Err(RangeError::ReversedRange { start: 5, end: 4 }) + ); + assert_eq!( + cache.insert("key", 2..5, Bytes::from_static(b"long")), + Err(RangeError::PayloadLengthMismatch { + range: 2..5, + expected: 3, + actual: 4, + }) + ); + assert_eq!( + RangeError::ReversedRange { start: 5, end: 4 }.to_string(), + "reversed byte range 5..4" + ); +} + +#[test] +fn ranges_near_usize_max_remain_valid() { + let cache = RangeCache::new(CacheCapacity::Unbounded); + let start = usize::MAX - 8; + cache + .insert("key", start..usize::MAX, Bytes::from_static(b"abcdefgh")) + .expect("valid high-offset insert"); + + assert_eq!( + cache + .get(&"key", start + 2..usize::MAX - 2) + .expect("valid high-offset read"), + Some(Bytes::from_static(b"cdef")) + ); + assert_eq!( + cache + .missing_ranges(&"key", start - 2..usize::MAX) + .expect("valid high-offset range"), + vec![start - 2..start] + ); } #[test] @@ -293,6 +345,95 @@ fn bounded_cache_never_exceeds_capacity_and_reads_update_lru() { assert_eq!(snapshot.evictions, 1); } +#[test] +fn bounded_partial_hits_touch_each_covered_range() { + let cache = bounded(8); + cache + .insert("first", 0..2, Bytes::from_static(b"ab")) + .expect("valid insert"); + cache + .insert("first", 4..6, Bytes::from_static(b"ef")) + .expect("valid insert"); + cache + .insert("oldest", 0..4, Bytes::from_static(b"wxyz")) + .expect("valid insert"); + + assert_eq!(cache.get(&"first", 0..6).expect("valid range"), None); + cache + .insert("new", 0..2, Bytes::from_static(b"12")) + .expect("valid insert"); + + assert_eq!(cache.get(&"oldest", 0..4).expect("valid range"), None); + assert_eq!( + cache.get(&"first", 0..2).expect("valid range"), + Some(Bytes::from_static(b"ab")) + ); + assert_eq!( + cache.get(&"first", 4..6).expect("valid range"), + Some(Bytes::from_static(b"ef")) + ); +} + +#[test] +fn merged_range_is_newer_than_ranges_it_replaces() { + let cache = bounded(8); + cache + .insert("oldest", 0..4, Bytes::from_static(b"wxyz")) + .expect("valid insert"); + cache + .insert("merged", 0..2, Bytes::from_static(b"ab")) + .expect("valid insert"); + cache + .insert("merged", 4..6, Bytes::from_static(b"ef")) + .expect("valid insert"); + cache + .insert("merged", 2..4, Bytes::from_static(b"cd")) + .expect("valid insert"); + + assert_eq!(cache.get(&"oldest", 0..4).expect("valid range"), None); + assert_eq!( + cache.get(&"merged", 0..6).expect("valid range"), + Some(Bytes::from_static(b"abcdef")) + ); + let snapshot = cache.snapshot(); + assert_eq!(snapshot.resident_bytes, 6); + assert_eq!(snapshot.ranges, 1); + assert_eq!(snapshot.evictions, 1); +} + +#[test] +fn concurrent_hits_preserve_bytes_and_exact_counters() { + let cache = Arc::new(RangeCache::new(CacheCapacity::Unbounded)); + cache + .insert("key", 0..8, Bytes::from_static(b"abcdefgh")) + .expect("valid insert"); + let workers = if cfg!(miri) { 2 } else { 8 }; + let iterations = if cfg!(miri) { 8 } else { 1_000 }; + let ready = Arc::new(Barrier::new(workers + 1)); + + thread::scope(|scope| { + for _ in 0..workers { + let cache = Arc::clone(&cache); + let ready = Arc::clone(&ready); + scope.spawn(move || { + ready.wait(); + for _ in 0..iterations { + assert_eq!( + cache.get(&"key", 0..8).expect("valid range"), + Some(Bytes::from_static(b"abcdefgh")) + ); + } + }); + } + ready.wait(); + }); + + let snapshot = cache.snapshot(); + assert_eq!(snapshot.hits, u64::try_from(workers * iterations).unwrap()); + assert_eq!(snapshot.resident_bytes, 8); + assert_eq!(snapshot.ranges, 1); +} + #[test] fn bounded_hits_reuse_the_resident_lru_key() { let cache = RangeCache::new(CacheCapacity::Bounded( From b2ec24ef35f2686dbee9270206a35819a9c7df38 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 13:37:48 +0100 Subject: [PATCH 08/10] Strengthen property and fuzz testing --- Cargo.toml | 2 +- fuzz/.gitignore | 3 + fuzz/Cargo.toml | 23 ++ .../corpus/range_cache_state_machine/bridging | 1 + fuzz/corpus/range_cache_state_machine/empty | 1 + .../corpus/range_cache_state_machine/eviction | 1 + .../range_cache_state_machine/high-offset | 1 + .../range_cache_state_machine/invalidation | 1 + .../range_cache_state_machine/mismatched | 1 + .../range_cache_state_machine/overlapping | 1 + .../corpus/range_cache_state_machine/reversed | 1 + .../fuzz_targets/range_cache_state_machine.rs | 350 ++++++++++++++++++ src/cache.rs | 6 +- tests/property.rs | 139 +++++-- 14 files changed, 488 insertions(+), 43 deletions(-) create mode 100644 fuzz/.gitignore create mode 100644 fuzz/Cargo.toml create mode 100644 fuzz/corpus/range_cache_state_machine/bridging create mode 100644 fuzz/corpus/range_cache_state_machine/empty create mode 100644 fuzz/corpus/range_cache_state_machine/eviction create mode 100644 fuzz/corpus/range_cache_state_machine/high-offset create mode 100644 fuzz/corpus/range_cache_state_machine/invalidation create mode 100644 fuzz/corpus/range_cache_state_machine/mismatched create mode 100644 fuzz/corpus/range_cache_state_machine/overlapping create mode 100644 fuzz/corpus/range_cache_state_machine/reversed create mode 100644 fuzz/fuzz_targets/range_cache_state_machine.rs diff --git a/Cargo.toml b/Cargo.toml index 08babcb..c1f5f13 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,7 @@ documentation = "https://docs.rs/range-cache" readme = "README.md" keywords = ["cache", "range", "bytes", "async", "storage"] categories = ["caching", "data-structures", "asynchronous"] -exclude = [".github/"] +exclude = [".github/", "fuzz/"] [features] default = [] diff --git a/fuzz/.gitignore b/fuzz/.gitignore new file mode 100644 index 0000000..5208f22 --- /dev/null +++ b/fuzz/.gitignore @@ -0,0 +1,3 @@ +artifacts/ +coverage/ +target/ diff --git a/fuzz/Cargo.toml b/fuzz/Cargo.toml new file mode 100644 index 0000000..0e5b9a6 --- /dev/null +++ b/fuzz/Cargo.toml @@ -0,0 +1,23 @@ +[package] +name = "range-cache-fuzz" +version = "0.0.0" +publish = false +edition = "2024" + +[package.metadata] +cargo-fuzz = true + +[dependencies] +bytes = "1" +libfuzzer-sys = "0.4" +range-cache = { path = ".." } + +[[bin]] +name = "range_cache_state_machine" +path = "fuzz_targets/range_cache_state_machine.rs" +test = false +doc = false +bench = false + +[workspace] +members = ["."] diff --git a/fuzz/corpus/range_cache_state_machine/bridging b/fuzz/corpus/range_cache_state_machine/bridging new file mode 100644 index 0000000..9d5081e --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/bridging @@ -0,0 +1 @@ + A0!#0aA0%'0bA0#%0c diff --git a/fuzz/corpus/range_cache_state_machine/empty b/fuzz/corpus/range_cache_state_machine/empty new file mode 100644 index 0000000..573541a --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/empty @@ -0,0 +1 @@ +0 diff --git a/fuzz/corpus/range_cache_state_machine/eviction b/fuzz/corpus/range_cache_state_machine/eviction new file mode 100644 index 0000000..71d2bbc --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/eviction @@ -0,0 +1 @@ +AA0!"0aA1"#0bC0!"0x diff --git a/fuzz/corpus/range_cache_state_machine/high-offset b/fuzz/corpus/range_cache_state_machine/high-offset new file mode 100644 index 0000000..ee60ebe --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/high-offset @@ -0,0 +1 @@ +éaaaaaA0!#0x diff --git a/fuzz/corpus/range_cache_state_machine/invalidation b/fuzz/corpus/range_cache_state_machine/invalidation new file mode 100644 index 0000000..0690121 --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/invalidation @@ -0,0 +1 @@ + A0!#0aD0!!!!A1#%0bE0!!!! diff --git a/fuzz/corpus/range_cache_state_machine/mismatched b/fuzz/corpus/range_cache_state_machine/mismatched new file mode 100644 index 0000000..9bc3599 --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/mismatched @@ -0,0 +1 @@ + A0!$1x diff --git a/fuzz/corpus/range_cache_state_machine/overlapping b/fuzz/corpus/range_cache_state_machine/overlapping new file mode 100644 index 0000000..ccb34f0 --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/overlapping @@ -0,0 +1 @@ + A0!%0aA0#'0bB0!'0x diff --git a/fuzz/corpus/range_cache_state_machine/reversed b/fuzz/corpus/range_cache_state_machine/reversed new file mode 100644 index 0000000..4eb7f86 --- /dev/null +++ b/fuzz/corpus/range_cache_state_machine/reversed @@ -0,0 +1 @@ + B0$"0x diff --git a/fuzz/fuzz_targets/range_cache_state_machine.rs b/fuzz/fuzz_targets/range_cache_state_machine.rs new file mode 100644 index 0000000..70bc19e --- /dev/null +++ b/fuzz/fuzz_targets/range_cache_state_machine.rs @@ -0,0 +1,350 @@ +#![no_main] + +use std::{collections::BTreeMap, num::NonZeroUsize, ops::Range}; + +use bytes::Bytes; +use libfuzzer_sys::fuzz_target; +use range_cache::{ + CacheCapacity, CacheSnapshot, InsertOutcome, Invalidation, RangeCache, RangeError, +}; + +const DOMAIN: usize = 32; +const KEY_COUNT: usize = 4; + +#[derive(Clone, Copy)] +struct ModelRange { + end: usize, + last_access: u64, +} + +struct Model { + base: usize, + capacity: CacheCapacity, + bytes: [[Option; DOMAIN]; KEY_COUNT], + ranges: [BTreeMap; KEY_COUNT], + next_access: u64, + hits: u64, + partial_hits: u64, + misses: u64, + insertions: u64, + admissions_rejected_too_large: u64, + evictions: u64, +} + +impl Model { + fn new(base: usize, capacity: CacheCapacity) -> Self { + Self { + base, + capacity, + bytes: [[None; DOMAIN]; KEY_COUNT], + ranges: std::array::from_fn(|_| BTreeMap::new()), + next_access: 0, + hits: 0, + partial_hits: 0, + misses: 0, + insertions: 0, + admissions_rejected_too_large: 0, + evictions: 0, + } + } + + fn absolute(&self, offset: usize) -> usize { + self.base + offset + } + + fn range_error(&self, start: usize, end: usize) -> RangeError { + RangeError::ReversedRange { + start: self.absolute(start), + end: self.absolute(end), + } + } + + fn insert( + &mut self, + key: usize, + start: usize, + end: usize, + payload: &[u8], + ) -> Result { + if start > end { + return Err(self.range_error(start, end)); + } + let expected = end - start; + if payload.len() != expected { + return Err(RangeError::PayloadLengthMismatch { + range: self.absolute(start)..self.absolute(end), + expected, + actual: payload.len(), + }); + } + if start == end { + return Ok(InsertOutcome::AlreadyCovered); + } + + let containing = self.ranges[key] + .range(..=start) + .next_back() + .is_some_and(|(_, block)| end <= block.end); + if containing { + return Ok(InsertOutcome::AlreadyCovered); + } + + let mut merged_start = start; + let mut merged_end = end; + let mut affected = Vec::new(); + for (&block_start, block) in &self.ranges[key] { + if block.end < merged_start { + continue; + } + if block_start > merged_end { + break; + } + merged_start = merged_start.min(block_start); + merged_end = merged_end.max(block.end); + affected.push(block_start); + } + + let merged_len = merged_end - merged_start; + if matches!(self.capacity, CacheCapacity::Bounded(capacity) if merged_len > capacity.get()) + { + self.admissions_rejected_too_large += 1; + return Ok(InsertOutcome::TooLarge); + } + + self.bytes[key][start..end] + .iter_mut() + .zip(payload) + .for_each(|(slot, &byte)| *slot = Some(byte)); + for block_start in affected { + self.ranges[key].remove(&block_start); + } + let access = self.next_access; + self.next_access += 1; + self.ranges[key].insert( + merged_start, + ModelRange { + end: merged_end, + last_access: access, + }, + ); + self.insertions += 1; + + if let CacheCapacity::Bounded(capacity) = self.capacity { + while self.resident_bytes() > capacity.get() { + let (oldest_key, oldest_start, oldest_end) = self + .ranges + .iter() + .enumerate() + .flat_map(|(key, ranges)| { + ranges + .iter() + .map(move |(&start, block)| (key, start, block.end, block.last_access)) + }) + .min_by_key(|(_, _, _, access)| *access) + .map(|(key, start, end, _)| (key, start, end)) + .expect("over-capacity model has a resident range"); + self.ranges[oldest_key].remove(&oldest_start); + self.bytes[oldest_key][oldest_start..oldest_end].fill(None); + self.evictions += 1; + } + } + Ok(InsertOutcome::Inserted) + } + + fn get(&mut self, key: usize, start: usize, end: usize) -> Result>, RangeError> { + if start > end { + return Err(self.range_error(start, end)); + } + if start == end { + self.hits += 1; + return Ok(Some(Vec::new())); + } + + let hit = self.ranges[key] + .range(..=start) + .next_back() + .filter(|(_, block)| end <= block.end) + .map(|(&block_start, _)| block_start); + if let Some(block_start) = hit { + self.touch(key, block_start); + self.hits += 1; + return Ok(Some( + self.bytes[key][start..end] + .iter() + .map(|byte| byte.expect("covered model byte exists")) + .collect(), + )); + } + + let covered = self.ranges[key] + .iter() + .filter(|(block_start, block)| **block_start < end && block.end > start) + .map(|(&block_start, _)| block_start) + .collect::>(); + if covered.is_empty() { + self.misses += 1; + } else { + for block_start in covered { + self.touch(key, block_start); + } + self.partial_hits += 1; + } + Ok(None) + } + + fn missing( + &self, + key: usize, + start: usize, + end: usize, + ) -> Result>, RangeError> { + if start > end { + return Err(self.range_error(start, end)); + } + let mut missing = Vec::new(); + let mut cursor = start; + while cursor < end { + if self.bytes[key][cursor].is_some() { + cursor += 1; + continue; + } + let gap_start = cursor; + while cursor < end && self.bytes[key][cursor].is_none() { + cursor += 1; + } + missing.push(self.absolute(gap_start)..self.absolute(cursor)); + } + Ok(missing) + } + + fn invalidate(&mut self, key: usize) -> Invalidation { + let invalidation = Invalidation { + ranges: self.ranges[key].len(), + bytes: self.bytes[key].iter().flatten().count(), + }; + self.ranges[key].clear(); + self.bytes[key].fill(None); + invalidation + } + + fn clear(&mut self) -> Invalidation { + let invalidation = Invalidation { + ranges: self.ranges.iter().map(BTreeMap::len).sum(), + bytes: self.resident_bytes(), + }; + self.ranges.iter_mut().for_each(BTreeMap::clear); + self.bytes.iter_mut().for_each(|bytes| bytes.fill(None)); + invalidation + } + + fn touch(&mut self, key: usize, start: usize) { + let block = self.ranges[key] + .get_mut(&start) + .expect("touched model range exists"); + block.last_access = self.next_access; + self.next_access += 1; + } + + fn resident_bytes(&self) -> usize { + self.bytes.iter().flatten().flatten().count() + } + + fn snapshot(&self) -> CacheSnapshot { + CacheSnapshot { + capacity: self.capacity, + resident_bytes: self.resident_bytes(), + keys: self + .ranges + .iter() + .filter(|ranges| !ranges.is_empty()) + .count(), + ranges: self.ranges.iter().map(BTreeMap::len).sum(), + hits: self.hits, + partial_hits: self.partial_hits, + misses: self.misses, + insertions: self.insertions, + admissions_rejected_too_large: self.admissions_rejected_too_large, + evictions: self.evictions, + } + } +} + +fuzz_target!(|data: &[u8]| { + let Some((&configuration, operations)) = data.split_first() else { + return; + }; + let capacity = if configuration & 1 == 0 { + CacheCapacity::Unbounded + } else { + CacheCapacity::Bounded( + NonZeroUsize::new(usize::from(configuration >> 1) % DOMAIN + 1) + .expect("derived capacity is non-zero"), + ) + }; + let base = if configuration & 0x80 == 0 { + 0 + } else { + usize::MAX - DOMAIN + }; + let cache = RangeCache::new(capacity); + let mut model = Model::new(base, capacity); + + for operation in operations.chunks_exact(6).take(256) { + let key = usize::from(operation[1]) % KEY_COUNT; + let start = usize::from(operation[2]) % (DOMAIN + 1); + let end = usize::from(operation[3]) % (DOMAIN + 1); + let absolute = model.absolute(start)..model.absolute(end); + + match operation[0] % 5 { + 0 => { + let range_len = end.saturating_sub(start); + let payload_len = match operation[4] % 3 { + 0 => range_len, + 1 => range_len.saturating_sub(1), + _ => (range_len + 1).min(DOMAIN + 1), + }; + let payload = (0..payload_len) + .map(|offset| operation[5].wrapping_add(offset as u8)) + .collect::>(); + assert_eq!( + cache.insert(key as u8, absolute, Bytes::from(payload.clone())), + model.insert(key, start, end, &payload) + ); + } + 1 => { + let expected = model + .get(key, start, end) + .map(|bytes| bytes.map(Bytes::from)); + assert_eq!(cache.get(&(key as u8), absolute), expected); + } + 2 => { + assert_eq!( + cache.missing_ranges(&(key as u8), absolute), + model.missing(key, start, end) + ); + } + 3 => { + assert_eq!(cache.invalidate(&(key as u8)), model.invalidate(key)); + } + 4 => { + assert_eq!(cache.clear(), model.clear()); + } + _ => unreachable!("opcode is reduced modulo five"), + } + + assert_eq!(cache.snapshot(), model.snapshot()); + for checked_key in 0..KEY_COUNT { + assert_eq!( + cache + .missing_ranges( + &(checked_key as u8), + model.absolute(0)..model.absolute(DOMAIN), + ) + .expect("model verification range is valid"), + model + .missing(checked_key, 0, DOMAIN) + .expect("model verification range is valid") + ); + } + } +}); diff --git a/src/cache.rs b/src/cache.rs index e6c843c..cc30e86 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -671,7 +671,7 @@ pub(crate) enum ReadPlan { #[cfg(test)] mod tests { - use std::{num::NonZeroUsize, sync::Arc}; + use std::{collections::BTreeMap, num::NonZeroUsize, sync::Arc}; use bytes::Bytes; use parking_lot::Mutex; @@ -733,7 +733,7 @@ mod tests { #[should_panic(expected = "range cache LRU clock exhausted")] fn registering_after_lru_clock_exhaustion_panics() { let mut eviction = EvictionPolicy::Bounded { - lru: Default::default(), + lru: BTreeMap::default(), next_access: u64::MAX, }; let _ = eviction.register(&"key", 0); @@ -743,7 +743,7 @@ mod tests { fn absent_internal_ranges_are_noops() { let mut state = State::<&str>::new(CacheCapacity::Unbounded); let _ = state.take_block(&"missing", 0); - state.ranges.insert("empty", Default::default()); + state.ranges.insert("empty", BTreeMap::default()); let _ = state.take_block(&"empty", 0); let _ = state.remove(&"missing", 0); assert_eq!(state.resident_bytes, 0); diff --git a/tests/property.rs b/tests/property.rs index a9fecd6..26d7d9c 100644 --- a/tests/property.rs +++ b/tests/property.rs @@ -6,6 +6,8 @@ use range_cache::{CacheCapacity, InsertOutcome, Invalidation, RangeCache}; const KEY_COUNT: usize = 3; const SOURCE_LENGTH: usize = 32; +const BOUNDED_KEY_COUNT: usize = 3; +const SLOT_COUNT: usize = 4; #[derive(Clone, Copy, Debug)] enum Operation { @@ -19,6 +21,7 @@ enum Operation { enum BoundedOperation { Insert, Get, + Missing, Invalidate, Clear, } @@ -46,15 +49,18 @@ fn operation() -> impl Strategy { ) } -fn bounded_operation() -> impl Strategy { +fn bounded_operation() -> impl Strategy { ( prop_oneof![ 4 => Just(BoundedOperation::Insert), 4 => Just(BoundedOperation::Get), + 2 => Just(BoundedOperation::Missing), 2 => Just(BoundedOperation::Invalidate), 1 => Just(BoundedOperation::Clear), ], - 0..5_usize, + 0..BOUNDED_KEY_COUNT, + 0..SLOT_COUNT, + 0..SLOT_COUNT, ) } @@ -217,89 +223,144 @@ proptest! { let cache = RangeCache::new(CacheCapacity::Bounded( NonZeroUsize::new(8).expect("test capacity is non-zero"), )); - let mut access_by_key = [None; 5]; + let mut access_by_range = [[None; SLOT_COUNT]; BOUNDED_KEY_COUNT]; let mut next_access = 0_u64; let mut hits = 0_u64; + let mut partial_hits = 0_u64; let mut misses = 0_u64; let mut insertions = 0_u64; let mut evictions = 0_u64; - for (operation, key) in operations { + for (operation, key, first, second) in operations { + let first_slot = first.min(second); + let last_slot = first.max(second); + let range = first_slot * 4..last_slot * 4 + 2; match operation { BoundedOperation::Insert => { + let slot = first; + let block_start = slot * 4; + let payload = Bytes::from(vec![ + u8::try_from(key).expect("key fits in u8"), + u8::try_from(slot).expect("slot fits in u8"), + ]); let outcome = cache - .insert( - key, - 0..4, - Bytes::from(vec![u8::try_from(key).expect("key fits in u8"); 4]), - ) + .insert(key, block_start..block_start + 2, payload) .expect("valid insert"); - if access_by_key[key].is_some() { + if access_by_range[key][slot].is_some() { prop_assert_eq!(outcome, InsertOutcome::AlreadyCovered); } else { prop_assert_eq!(outcome, InsertOutcome::Inserted); - access_by_key[key] = Some(next_access); + access_by_range[key][slot] = Some(next_access); next_access += 1; insertions += 1; - if access_by_key.iter().flatten().count() > 2 { - let (oldest_key, _) = access_by_key + if access_by_range.iter().flatten().flatten().count() > 4 { + let (oldest_key, oldest_slot, _) = access_by_range .iter() .enumerate() - .filter_map(|(key, access)| access.map(|access| (key, access))) - .min_by_key(|(_, access)| *access) + .flat_map(|(key, ranges)| { + ranges.iter().enumerate().filter_map( + move |(slot, access)| { + access.map(|access| (key, slot, access)) + }, + ) + }) + .min_by_key(|(_, _, access)| *access) .expect("an over-capacity model has an oldest entry"); - access_by_key[oldest_key] = None; + access_by_range[oldest_key][oldest_slot] = None; evictions += 1; } } } BoundedOperation::Get => { - let expected = access_by_key[key] - .is_some() - .then(|| { - Bytes::from(vec![u8::try_from(key).expect("key fits in u8"); 4]) - }); - prop_assert_eq!(cache.get(&key, 0..4).expect("valid range"), expected); - if access_by_key[key].is_some() { - access_by_key[key] = Some(next_access); + let resident_slots = (first_slot..=last_slot) + .filter(|&slot| access_by_range[key][slot].is_some()) + .collect::>(); + let expected = (first_slot == last_slot + && access_by_range[key][first_slot].is_some()) + .then(|| { + Bytes::from(vec![ + u8::try_from(key).expect("key fits in u8"), + u8::try_from(first_slot).expect("slot fits in u8"), + ]) + }); + let is_hit = expected.is_some(); + prop_assert_eq!(cache.get(&key, range).expect("valid range"), expected); + if is_hit { + access_by_range[key][first_slot] = Some(next_access); next_access += 1; hits += 1; - } else { + } else if resident_slots.is_empty() { misses += 1; + } else { + for slot in resident_slots { + access_by_range[key][slot] = Some(next_access); + next_access += 1; + } + partial_hits += 1; } } - BoundedOperation::Invalidate => { - let expected = if access_by_key[key].take().is_some() { - Invalidation { - ranges: 1, - bytes: 4, + BoundedOperation::Missing => { + let mut expected = Vec::new(); + let mut cursor = range.start; + for (slot, access) in access_by_range[key] + .iter() + .enumerate() + .take(last_slot + 1) + .skip(first_slot) + { + let start = slot * 4; + if access.is_none() { + continue; } - } else { - Invalidation::default() + if cursor < start { + expected.push(cursor..start); + } + cursor = start + 2; + } + if cursor < range.end { + expected.push(cursor..range.end); + } + prop_assert_eq!( + cache.missing_ranges(&key, range).expect("valid range"), + expected + ); + } + BoundedOperation::Invalidate => { + let ranges = access_by_range[key].iter().flatten().count(); + let expected = Invalidation { + ranges, + bytes: ranges * 2, }; prop_assert_eq!(cache.invalidate(&key), expected); + access_by_range[key].fill(None); } BoundedOperation::Clear => { - let ranges = access_by_key.iter().flatten().count(); + let ranges = access_by_range.iter().flatten().flatten().count(); prop_assert_eq!( cache.clear(), Invalidation { ranges, - bytes: ranges * 4, + bytes: ranges * 2, } ); - access_by_key.fill(None); + access_by_range.fill([None; SLOT_COUNT]); } } - let ranges = access_by_key.iter().flatten().count(); + let ranges = access_by_range.iter().flatten().flatten().count(); let snapshot = cache.snapshot(); - prop_assert_eq!(snapshot.resident_bytes, ranges * 4); - prop_assert_eq!(snapshot.keys, ranges); + prop_assert_eq!(snapshot.resident_bytes, ranges * 2); + prop_assert_eq!( + snapshot.keys, + access_by_range + .iter() + .filter(|ranges| ranges.iter().any(Option::is_some)) + .count() + ); prop_assert_eq!(snapshot.ranges, ranges); prop_assert_eq!(snapshot.hits, hits); - prop_assert_eq!(snapshot.partial_hits, 0); + prop_assert_eq!(snapshot.partial_hits, partial_hits); prop_assert_eq!(snapshot.misses, misses); prop_assert_eq!(snapshot.insertions, insertions); prop_assert_eq!(snapshot.admissions_rejected_too_large, 0); From cf8ddc4004cc676120436e457a1e06172ad7fcb9 Mon Sep 17 00:00:00 2001 From: xav-db Date: Fri, 24 Jul 2026 13:42:25 +0100 Subject: [PATCH 09/10] Add continuous robustness checks --- .github/workflows/ci.yml | 3 +- .github/workflows/robustness.yml | 76 ++++++++++++++++++++++++++++++ CHANGELOG.md | 3 ++ CONTRIBUTING.md | 45 +++++++++++++++++- scripts/check-source-coverage.mjs | 78 +++++++++++++++++++++++++++++++ 5 files changed, 202 insertions(+), 3 deletions(-) create mode 100644 .github/workflows/robustness.yml create mode 100644 scripts/check-source-coverage.mjs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9d48fc9..7be083f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -56,7 +56,8 @@ jobs: run: rustup toolchain install stable --profile minimal --component llvm-tools-preview - run: rustup override set stable - run: cargo install cargo-llvm-cov --locked - - run: cargo llvm-cov --all-features --workspace --fail-under-lines 95 + - run: cargo llvm-cov --all-features --workspace --json --output-path target/coverage.json + - run: node scripts/check-source-coverage.mjs target/coverage.json package: name: Package diff --git a/.github/workflows/robustness.yml b/.github/workflows/robustness.yml new file mode 100644 index 0000000..00f0fa6 --- /dev/null +++ b/.github/workflows/robustness.yml @@ -0,0 +1,76 @@ +name: Robustness + +on: + pull_request: + workflow_dispatch: + schedule: + - cron: "30 2 * * *" + - cron: "0 3 * * 0" + +permissions: + contents: read + +jobs: + fuzz-pr: + name: Deterministic fuzz smoke + if: github.event_name == 'pull_request' || github.event_name == 'workflow_dispatch' + runs-on: ubuntu-latest + timeout-minutes: 15 + steps: + - uses: actions/checkout@v7 + - name: Install pinned nightly + run: rustup toolchain install nightly-2026-02-17 --profile minimal + - name: Install cargo-fuzz + run: cargo install cargo-fuzz --version 0.13.2 --locked + - name: Run 10,000 state-machine cases + run: >- + cargo +nightly-2026-02-17 fuzz run range_cache_state_machine + fuzz/corpus/range_cache_state_machine -- + -runs=10000 -seed=0 -timeout=5 -max_len=4096 + - name: Upload crash artifacts + if: failure() + uses: actions/upload-artifact@v4 + with: + name: range-cache-fuzz-pr-artifacts + path: fuzz/artifacts + if-no-files-found: ignore + + fuzz-nightly: + name: 15-minute nightly fuzz + if: github.event_name == 'schedule' && github.event.schedule == '30 2 * * *' + runs-on: ubuntu-latest + timeout-minutes: 25 + steps: + - uses: actions/checkout@v7 + - name: Install pinned nightly + run: rustup toolchain install nightly-2026-02-17 --profile minimal + - name: Install cargo-fuzz + run: cargo install cargo-fuzz --version 0.13.2 --locked + - name: Fuzz for 15 minutes + run: >- + cargo +nightly-2026-02-17 fuzz run range_cache_state_machine + fuzz/corpus/range_cache_state_machine -- + -max_total_time=900 -seed=0 -timeout=5 -max_len=4096 + - name: Upload crash artifacts + if: failure() + uses: actions/upload-artifact@v4 + with: + name: range-cache-fuzz-nightly-artifacts + path: fuzz/artifacts + if-no-files-found: ignore + + miri-weekly: + name: Weekly core Miri + if: github.event_name == 'schedule' && github.event.schedule == '0 3 * * 0' + runs-on: ubuntu-latest + timeout-minutes: 30 + steps: + - uses: actions/checkout@v7 + - name: Install pinned nightly and Miri + run: >- + rustup toolchain install nightly-2026-02-17 + --profile minimal --component miri,rust-src + - name: Run core and private-invariant tests + run: >- + cargo +nightly-2026-02-17 miri test + --no-default-features --lib --test core diff --git a/CHANGELOG.md b/CHANGELOG.md index b635bc8..7a3e4e2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,9 @@ project follows [Semantic Versioning](https://semver.org/). - Expanded Criterion microbenchmarks for cache operations, contention, and async read-through behavior. - Complete core and async README examples plus documented benchmark results. +- Deterministic failure, cancellation, concurrency, and private-invariant tests. +- A public-API state-machine fuzz target with checked-in edge-case corpora. +- Source-normalized 100% coverage enforcement, nightly fuzzing, and weekly Miri. ### Changed diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index a6cb759..21d6d0f 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -8,16 +8,57 @@ cargo fmt --all -- --check cargo test --lib --tests --no-default-features cargo test --all-features cargo clippy --all-targets --all-features -- -D warnings +rustfmt --edition 2024 --check $(rg --files -g '*.rs' -g '!target/**' -g '!fuzz/target/**') RUSTDOCFLAGS="-D warnings" cargo test --doc --no-default-features RUSTDOCFLAGS="-D warnings" cargo test --doc --all-features cargo bench --all-features --bench range_cache -- --test cargo publish --dry-run ``` -For coverage, install `cargo-llvm-cov` and run: +## Coverage + +Install `cargo-llvm-cov`, then generate and check the source-normalized report: + +```bash +cargo llvm-cov --all-features --workspace --json --output-path target/coverage.json +node scripts/check-source-coverage.mjs target/coverage.json +``` + +The checker requires 100% source function, line, and region coverage. LLVM emits +generic Rust functions once per test binary and raw summaries can count a source +region as missed in one monomorphization even when another executes it. The +checker collapses only identical source coordinates, using their maximum count; +it does not exclude files, functions, or regions. + +## Fuzzing + +The state-machine target uses only the public core API and checks every +operation against an independent dense reference model. Install the pinned tools +and reproduce the PR run with: + +```bash +rustup toolchain install nightly-2026-02-17 --profile minimal +cargo install cargo-fuzz --version 0.13.2 --locked +cargo +nightly-2026-02-17 fuzz run range_cache_state_machine fuzz/corpus/range_cache_state_machine -- -runs=10000 -seed=0 -timeout=5 -max_len=4096 +``` + +Reproduce and minimize a crash with: + +```bash +cargo +nightly-2026-02-17 fuzz run range_cache_state_machine fuzz/artifacts/range_cache_state_machine/crash-... +cargo +nightly-2026-02-17 fuzz tmin range_cache_state_machine fuzz/artifacts/range_cache_state_machine/crash-... +``` + +Promote the minimized input into `fuzz/corpus/range_cache_state_machine/` and +add a readable deterministic regression test that captures the same failure. + +## Miri + +Run the core and private-invariant tests under the pinned nightly: ```bash -cargo llvm-cov --all-features --workspace --fail-under-lines 95 +rustup toolchain install nightly-2026-02-17 --profile minimal --component miri,rust-src +cargo +nightly-2026-02-17 miri test --no-default-features --lib --test core ``` Use a current stable toolchain for development. Changes must remain compatible diff --git a/scripts/check-source-coverage.mjs b/scripts/check-source-coverage.mjs new file mode 100644 index 0000000..934087d --- /dev/null +++ b/scripts/check-source-coverage.mjs @@ -0,0 +1,78 @@ +import fs from "node:fs"; +import path from "node:path"; + +const [reportPath] = process.argv.slice(2); +if (reportPath === undefined) { + console.error("usage: node scripts/check-source-coverage.mjs "); + process.exit(2); +} + +const report = JSON.parse(fs.readFileSync(reportPath, "utf8")); +const sourceRoot = `${path.resolve("src")}${path.sep}`; +const sourceFiles = new Map(); + +for (const datum of report.data ?? []) { + for (const file of datum.files ?? []) { + const filename = path.resolve(file.filename); + if (filename.startsWith(sourceRoot)) { + sourceFiles.set(filename, file); + } + } +} + +if (sourceFiles.size === 0) { + console.error("coverage report contains no crate source files"); + process.exit(1); +} + +const regions = new Map(); +for (const datum of report.data ?? []) { + for (const fn of datum.functions ?? []) { + for (const region of fn.regions ?? []) { + const filename = path.resolve(fn.filenames[region[5]]); + const isCodeRegion = region[7] === 0; + if (!isCodeRegion || !sourceFiles.has(filename)) { + continue; + } + + const coordinates = region.slice(0, 4).join(":"); + const key = `${filename}:${coordinates}`; + regions.set(key, Math.max(regions.get(key) ?? 0, region[4])); + } + } +} + +if (regions.size === 0) { + console.error("coverage report contains no crate source regions"); + process.exit(1); +} + +const uncoveredRegions = [...regions] + .filter(([, count]) => count === 0) + .map(([region]) => region); +const functions = [...sourceFiles.values()].reduce( + (totals, file) => ({ + count: totals.count + file.summary.functions.count, + covered: totals.covered + file.summary.functions.covered, + }), + { count: 0, covered: 0 }, +); +const lines = [...sourceFiles.values()].reduce( + (count, file) => count + file.summary.lines.count, + 0, +); + +if (functions.covered !== functions.count || uncoveredRegions.length > 0) { + console.error( + `source coverage failed: functions ${functions.covered}/${functions.count}, ` + + `regions ${regions.size - uncoveredRegions.length}/${regions.size}`, + ); + for (const region of uncoveredRegions.slice(0, 20)) { + console.error(`uncovered: ${region}`); + } + process.exit(1); +} + +console.log(`source functions: ${functions.count}/${functions.count} (100.00%)`); +console.log(`source lines: ${lines}/${lines} (100.00%)`); +console.log(`source regions: ${regions.size}/${regions.size} (100.00%)`); From 95787e9144e04e7cdf7ee897cd69101681200fce Mon Sep 17 00:00:00 2001 From: xav-db Date: Wed, 29 Jul 2026 18:27:34 +0100 Subject: [PATCH 10/10] Prepare range-cache 0.1.1 release --- CHANGELOG.md | 5 ++++- Cargo.toml | 2 +- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7a3e4e2..7720602 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,8 @@ project follows [Semantic Versioning](https://semver.org/). ## [Unreleased] +## [0.1.1] - 2026-07-29 + ### Added - Expanded Criterion microbenchmarks for cache operations, contention, and @@ -37,5 +39,6 @@ project follows [Semantic Versioning](https://semver.org/). - Reference-model property tests, Criterion benchmarks, cross-platform CI, and coverage enforcement. -[Unreleased]: https://github.com/xav-db/range-cache/compare/v0.1.0...HEAD +[Unreleased]: https://github.com/xav-db/range-cache/compare/v0.1.1...HEAD +[0.1.1]: https://github.com/xav-db/range-cache/compare/v0.1.0...v0.1.1 [0.1.0]: https://github.com/xav-db/range-cache/releases/tag/v0.1.0 diff --git a/Cargo.toml b/Cargo.toml index c1f5f13..65d5d00 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "range-cache" -version = "0.1.0" +version = "0.1.1" edition = "2024" rust-version = "1.86" authors = ["HelixDB, Inc. "]