diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 4ee1ee2..f488f11 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -120,6 +120,12 @@ The managed-memory limit covers the index mapping, heat bits, L1, append buffers C² owns one logical data path. Multi-device deployments stripe below the filesystem with RAID0 or an equivalent layer. Request routing, recovery identity, and descriptor count stay independent of device topology. +## Statistics + +Existing health/resource snapshots and optional legacy activity counters retain their populations. Additional `RuntimeOptions::stats` recorders use cache-owned, preallocated atomic stripes whose allocation is validated and included in managed memory. Public operations contribute one exclusive terminal outcome; original L2 hits remain L2 even after promotion. L1 hit, L2 lookup and mutation timing are independently disabled, full or randomly sampled. L2 clocks start immediately after L1 miss, before index lookup and admission, so full L2 timing does not require clocks on L1 hits; full I/O timing reuses engine timestamps and preserves separate read/write/reclaim populations. Fixed scalar thread-local routing and sampling state cannot retain a cache instance or grow a per-cache registry. Recording adds no allocation, queue, recorder lock or worker. + +Statistics snapshots load cumulative counters without resetting them or scanning metadata. Histogram bucket counts and duration sums may reflect slightly different instants during concurrent updates; quiescent snapshots are exact. Snapshot collection runs on the caller, with bounded output determined by the fixed operation/outcome set and histogram layout. Applications own all metric conversion and transport. Sampled observations remain explicitly distinct from full-population metrics, and invalid overflowed distributions are flagged for applications to omit. Additional recorder control objects fit the fixed runtime control reservation; variable stripe storage is separately charged. + ## Recovery and failures ### Persistent artifacts diff --git a/BENCHMARK.md b/BENCHMARK.md index 4e0158c..f99b6f5 100644 --- a/BENCHMARK.md +++ b/BENCHMARK.md @@ -43,6 +43,8 @@ The main controls are grouped below. See `benchmarks/cache/main.rs` for defaults | Measurement | `CACHE_BENCH_READ_LATENCY_SAMPLE_INTERVAL` (default 16; zero disables), `CACHE_BENCH_STATS` (default false; true adds cache/I/O/resource records) | | Gates | `CACHE_BENCH_MIN_PUT_OPS`, `CACHE_BENCH_MIN_RESIDENT_L1_OPS`, `CACHE_BENCH_MIN_L2_OPS`, `CACHE_BENCH_MAX_WARM_CLOSE_MS` | +The additional stats implementation is controlled independently by `CACHE_BENCH_REQUEST_STATS` (complete request counters, default false), `CACHE_BENCH_L1_LATENCY_SAMPLE_INTERVAL`, `CACHE_BENCH_L2_LATENCY_SAMPLE_INTERVAL`, and `CACHE_BENCH_MUTATION_LATENCY_SAMPLE_INTERVAL` (independent: 0 off, 1 full, greater values sample with that mean interval; default 0), `CACHE_BENCH_IO_LATENCY` (full engine latency, default false), and `CACHE_BENCH_STATS_SHARDS` (default 16). These instrument the library itself. `CACHE_BENCH_READ_LATENCY_SAMPLE_INTERVAL` remains the independent benchmark observer and should be held constant across comparisons. The effective additional settings are printed with each run. Compare disabled, counters-only, full and sampled modes on the same workload, keeping actual hit/overload populations in view. + `CACHE_BENCH_STATS=true` adds cache accounting on the measured request path. Use the same setting for baseline and candidate runs. For device measurements, use a data set larger than host RAM and no larger than half of L2 capacity. Run baseline and candidate in alternating order at least five times, compare medians, and retain every sample. Throughput does not replace correctness, overload, memory, or latency checks. diff --git a/CHANGELOG.md b/CHANGELOG.md index ce53fe9..1364dd3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,14 @@ ## Unreleased +### Bug Fixes + +- Read I/O duration accounting now includes buffer preparation and scheduling between slot reservation and submission, matching the documented reservation-to-completion interval. + +### Features + +- Added `RuntimeOptions::stats` and `Cache::stats_snapshot()` for complete public request outcomes, independently configured full or sampled L1-hit, L2-lookup and mutation latency, and full read/write/reclaim I/O latency. L2 timing starts after L1 miss, allowing full L2 collection without clocks on L1 hits; structured durations identify their scope and collection mode. Recorders have bounded managed-memory accounting; snapshots preserve cumulative values and distinguish disabled families and sampled observations. The structured snapshots are independent of monitoring SDKs; applications own metric conversion and export. The existing `statistics` switch retains its behavior; fully specified runtime-option literals must add `stats` or use `..RuntimeOptions::default()`. + ## v0.4.0 (2026-09-10) ### Breaking Changes diff --git a/CONFIGURATION.md b/CONFIGURATION.md index 7deb233..016e4bb 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -295,6 +295,14 @@ IOPOLL is an additional explicit per-pool opt-in through `IoUringPoolConfig::wit Health and managed-resource gauges are always available. Enable `RuntimeOptions::statistics` while tuning to obtain cumulative request, index, L1, and I/O counters. Enabled counters add relaxed atomic work on active paths, so measure the overhead before leaving all activity statistics enabled in a latency-critical deployment. +`RuntimeOptions::stats` adds independent controls: `request_counters` records all terminal public results; `l1_latency`, `l2_latency`, and `mutation_latency` independently select `LatencyMode::Off`, `Full`, or `Sampled { interval }`; `io_latency` records every engine completion by role; and `shards` selects a power of two from 1 through 64 (default 16). Defaults disable all additional recorders and preserve the legacy `statistics` switch. Fully specified `RuntimeOptions` literals need `stats: StatsOptions::default()` or a struct-update default. + +Request counters distinguish original L1/L2 hits, misses, accepted mutations, overload, invalid input, unavailable mutations, other errors and cancelled gets. Never-polled futures contribute nothing; cancelled durations are partial lifetimes. L1 timing covers successful L1 lookups from the first poll. L2 timing begins immediately after L1 miss and covers index lookup, admission waiting, I/O and promotion, excluding the initial L1 lookup; early misses before L1 lookup contribute only to request counters. L1 miss discards the L1 timer and makes an independent L2 sampling decision. Put timing ends at acceptance, not publication or durability. I/O timing reuses engine timestamps and covers slot reservation through terminal completion, including engine queueing rather than just device service. An I/O histogram may complete after its caller has cancelled. + +For low overhead start with counters and, if needed, full L2 lookup and I/O latency. Keep L1 latency off or sampled independently. Enabling only L2 timing starts no clock on L1 hits. Enable full timing for a selected population when every observed latency matters, or explicitly choose sampling after measuring overhead and per-series sample volume. Combine only matching timing scopes and sample intervals; each structured request row carries its own `latency_scope` and `latency_mode`. Applications choose metric names and attributes and own conversion, timestamps, scheduling and transport. Sampled bucket counts and sums are never inflated to full request volume. Use a full histogram with an aligned bucket boundary for exact observed SLO-threshold counts. No sampling setting guarantees detecting isolated long-tail events. + +`stats_snapshot()` reports disabled families explicitly, preserves cumulative values across concurrent readers, and resets on reopen using the existing `metrics_epoch`. It does not scan L1/index/Region metadata or start background workers. Each histogram includes a validity flag; arithmetic overflow invalidates that series until reopen and applications must omit its distribution when exporting. Detailed occupancy remains an explicit `detailed_snapshot()` operation. Only enabled latency populations allocate histogram stripes. The fixed histogram boundaries are documented by `LATENCY_BUCKET_UPPER_BOUNDS_NS`; counts above the largest boundary remain in an unbounded bucket with their true duration sum. + ## Goal-oriented profiles ### Balanced starting point diff --git a/README.md b/README.md index aa78223..a95a8fa 100644 --- a/README.md +++ b/README.md @@ -108,6 +108,8 @@ The on-disk format is versioned. During 0.x, deployments should expect cold star `Cache::snapshot()` provides lock-free health and resource gauges. `RuntimeOptions { statistics: true, .. }` adds cumulative cache and I/O counters. `Cache::detailed_snapshot()` samples L1, index, write-buffer pressure, and Region metadata for periodic diagnostics. +`RuntimeOptions::stats` independently enables complete public request outcomes, L1-hit, L2-lookup and mutation latency (each `Off`, `Full`, or `Sampled`), and full I/O latency by read/write/reclaim role. `Cache::stats_snapshot()` combines these with the existing summary without metadata scans. Structured request rows include their timing scope and collection mode. Applications own metric conversion, timestamps, scheduling and transport. Run `cargo run --example stats -- ` for an example. Full timing avoids sampling work; sampled histograms retain actual sample counts and cannot guarantee observation of rare tail events. Recorder storage is preallocated, bounded and charged to managed memory. + C² exposes snapshots for integration with the application's metrics SDK. An OpenTelemetry or Prometheus adapter can export: - get outcomes from `l1_hits`, `l2_hits`, `l2_misses`, and `l2_read_overloads`; diff --git a/benchmarks/cache/main.rs b/benchmarks/cache/main.rs index 1328956..7b7d664 100644 --- a/benchmarks/cache/main.rs +++ b/benchmarks/cache/main.rs @@ -82,6 +82,7 @@ struct BenchConfig { io_mode: IoMode, l1_eviction_policy: L1EvictionPolicy, statistics_enabled: bool, + stats: cache2::StatsOptions, directory: PathBuf, } @@ -154,6 +155,23 @@ impl BenchConfig { } }; let statistics_enabled = env_bool("CACHE_BENCH_STATS", false)?; + let latency = |name| -> io::Result { + Ok(match env_u32(name, 0)? { + 0 => cache2::LatencyMode::Off, + 1 => cache2::LatencyMode::Full, + interval => cache2::LatencyMode::Sampled { + interval: std::num::NonZeroU32::new(interval).expect("positive interval"), + }, + }) + }; + let stats = cache2::StatsOptions { + request_counters: env_bool("CACHE_BENCH_REQUEST_STATS", false)?, + l1_latency: latency("CACHE_BENCH_L1_LATENCY_SAMPLE_INTERVAL")?, + l2_latency: latency("CACHE_BENCH_L2_LATENCY_SAMPLE_INTERVAL")?, + mutation_latency: latency("CACHE_BENCH_MUTATION_LATENCY_SAMPLE_INTERVAL")?, + io_latency: env_bool("CACHE_BENCH_IO_LATENCY", false)?, + shards: env_usize("CACHE_BENCH_STATS_SHARDS", 16)?, + }; let directory = env::var_os("CACHE_BENCH_DIR") .map(PathBuf::from) .unwrap_or_else(env::temp_dir); @@ -264,6 +282,7 @@ impl BenchConfig { io_mode, l1_eviction_policy, statistics_enabled, + stats, directory, }) } @@ -285,6 +304,7 @@ impl BenchConfig { l1_eviction_policy: self.l1_eviction_policy, managed_memory_limit_bytes: self.managed_memory_limit_bytes, statistics: self.statistics_enabled, + stats: self.stats, read_admission: if self.read_io_wait_timeout.is_zero() { ReadAdmission::Immediate } else { @@ -411,6 +431,7 @@ async fn run(config: BenchConfig) -> io::Result<()> { }; println!("C² cache benchmark"); + println!("additional_stats={:?}", config.stats); println!( "entries={} index_slots={} index_load={:.1}% resident_entries={} hot_entries={} hot_read_interval={} value={} B data={:.1} MiB memory={:.1} MiB initial_l1={:.1} MiB managed_memory_limit={:.1} MiB append_shards={} read_workers={} read_wait_capacity={} read_wait_timeout_us={} read_latency_sample_interval={} write_workers={} reclaim_workers={} write_clients={} read_clients={} l2_clients={} l1_entry_eligible={} l1_eviction={:?} engine={:?} mode={:?} statistics={}", config.entries, diff --git a/cache2/src/cache.rs b/cache2/src/cache.rs index 9baffaf..aded359 100644 --- a/cache2/src/cache.rs +++ b/cache2/src/cache.rs @@ -56,17 +56,16 @@ use crate::snapshot::CacheSnapshot; use crate::snapshot::DetailedCacheSnapshot; use crate::snapshot::StartupMode; -/// Storage tier that backs a returned [`Value`]. +/// Storage tier that served a lookup. /// -/// This describes the value's backing after the lookup completes, not -/// necessarily the tier where the lookup began. A successful L2 promotion is -/// therefore backed by L1. +/// An L2 value promoted before return still reports L2, even though its owned +/// bytes are backed by L1. This agrees with the hit counters. #[non_exhaustive] #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub enum CacheTier { /// Process-local in-memory storage. L1, - /// A transient buffer read from the Region store. + /// A validated Region lookup, including a value promoted to L1. L2, } @@ -93,10 +92,10 @@ impl AsRef<[u8]> for Value { } impl Value { - /// Returns the storage tier backing this value. + /// Returns the storage tier that served this lookup. /// - /// A Region hit successfully promoted before return reports - /// [`CacheTier::L1`]; hit counters still record it as an L2 hit. + /// A Region hit successfully promoted before return still reports + /// [`CacheTier::L2`], matching the hit counters. pub const fn tier(&self) -> CacheTier { if self.inner.is_l1() { CacheTier::L1 @@ -116,6 +115,8 @@ impl Value { pub struct Cache { // Keep public reads off the write-mutated admission counter's cache line. closed: AtomicBool, + read_statistics: bool, + mutation_statistics: bool, data_plane: RegionDataPlane, owner: Arc>>>, startup: StartupMode, @@ -245,12 +246,20 @@ impl Cache { ); let index_slots = config.storage().index_slots(); let logical_disk_peak_bytes = config.storage().peak_disk_bytes(); + let stats = config.runtime().stats; + let read_statistics = stats.request_counters + || stats.l1_latency != crate::LatencyMode::Off + || stats.l2_latency != crate::LatencyMode::Off; + let mutation_statistics = + stats.request_counters || stats.mutation_latency != crate::LatencyMode::Off; let backend = FileRegionBackend::new(files, format_data, config); let store = RegionStore::open(index_slots, backend)?; let startup = store.startup(); let data_plane = store.data_plane_handle()?; Ok(Cache { closed: AtomicBool::new(false), + read_statistics, + mutation_statistics, data_plane, owner: Arc::new(Mutex::new(store)), startup, @@ -277,11 +286,9 @@ impl Cache { /// close starts. Runtime and device failures use their corresponding structured /// classifications. pub fn put(&self, key: impl AsRef<[u8]>, value: impl AsRef<[u8]>) -> Result { - self.ensure_open(ErrorOperation::Put)?; - public_result( - ErrorOperation::Put, - self.data_plane.put(key.as_ref(), value.as_ref()), - ) + self.mutate(ErrorOperation::Put, crate::RequestOperation::Put, || { + self.data_plane.put(key.as_ref(), value.as_ref()) + }) } /// Stages a value directly for L2 and returns its monotonic mutation @@ -297,10 +304,10 @@ impl Cache { /// [`Self::put`], including unavailable after close starts, with /// [`ErrorOperation::PutL2`](crate::ErrorOperation::PutL2) as its context. pub fn put_l2(&self, key: impl AsRef<[u8]>, value: impl AsRef<[u8]>) -> Result { - self.ensure_open(ErrorOperation::PutL2)?; - public_result( + self.mutate( ErrorOperation::PutL2, - self.data_plane.put_l2(key.as_ref(), value.as_ref()), + crate::RequestOperation::PutL2, + || self.data_plane.put_l2(key.as_ref(), value.as_ref()), ) } @@ -316,8 +323,11 @@ impl Cache { /// busy. Returns [`ErrorKind::Unavailable`](crate::ErrorKind::Unavailable) after close starts. /// Runtime and device failures remain explicit. pub fn delete(&self, key: impl AsRef<[u8]>) -> Result { - self.ensure_open(ErrorOperation::Delete)?; - public_result(ErrorOperation::Delete, self.data_plane.delete(key.as_ref())) + self.mutate( + ErrorOperation::Delete, + crate::RequestOperation::Delete, + || self.data_plane.delete(key.as_ref()), + ) } /// Looks up a value in L1 and then L2. @@ -336,16 +346,42 @@ impl Cache { /// open transition reads to misses instead of surfacing an application /// error. pub async fn get(&self, key: impl AsRef<[u8]> + Send) -> Result, Error> { - if self.is_closed() { - return Ok(None); + // Keep the disabled arm identical to the uninstrumented read path: no + // recorder access, guard, clock, TLS or terminal-result classification. + if !self.read_statistics { + if self.is_closed() { + return Ok(None); + } + return public_result( + ErrorOperation::Get, + self.data_plane + .get_async(key.as_ref(), &self.tokio_handle, None) + .await, + ) + .map(|value| value.map(|inner| Value { inner })); } - public_result( - ErrorOperation::Get, - self.data_plane - .get_async(key.as_ref(), &self.tokio_handle) - .await, - ) - .map(|value| value.map(|inner| Value { inner })) + let mut guard = self + .data_plane + .stats_recorder() + .begin(crate::RequestOperation::Get); + let result = if self.is_closed() { + Ok(None) + } else { + public_result( + ErrorOperation::Get, + self.data_plane + .get_async(key.as_ref(), &self.tokio_handle, Some(&mut guard)) + .await, + ) + .map(|value| value.map(|inner| Value { inner })) + }; + let outcome = match &result { + Ok(Some(value)) if value.inner.is_l1() => crate::RequestOutcome::L1Hit, + Ok(Some(_)) => crate::RequestOutcome::L2Hit, + _ => crate::RequestOutcome::Miss, + }; + guard.finish(&result, outcome); + result } /// Waits until all accepted mutations have completed their Region writes @@ -377,6 +413,20 @@ impl Cache { Ok(snapshot) } + /// Returns cumulative legacy statistics, complete optional request outcomes, + /// and optional request/I/O latency histograms without scanning metadata. + /// + /// Repeated or concurrent readers do not reset counts. Duration populations + /// and enablement are explicit in [`crate::CacheStatsSnapshot`]. Snapshot + /// allocations are bounded by the configured stripes and fixed series set. + /// + /// # Errors + /// + /// Uses the same availability and runtime failures as [`Self::snapshot`]. + pub fn stats_snapshot(&self) -> Result { + Ok(self.data_plane.stats_recorder().snapshot(self.snapshot()?)) + } + /// Samples L1, index, write-buffer rejection, I/O, and Region state in /// addition to the regular cache summary. This periodic diagnostic briefly /// reads every configured L1 shard and index partition and scans Region @@ -427,6 +477,25 @@ impl Cache { self.close(true) } + #[inline] + fn mutate( + &self, + operation: ErrorOperation, + stats_operation: crate::RequestOperation, + mutation: impl FnOnce() -> io::Result, + ) -> Result { + let guard = self + .mutation_statistics + .then(|| self.data_plane.stats_recorder().begin(stats_operation)); + let result = self + .ensure_open(operation) + .and_then(|()| public_result(operation, mutation())); + if let Some(guard) = guard { + guard.finish(&result, crate::RequestOutcome::Accepted); + } + result + } + #[inline(always)] fn ensure_open(&self, operation: ErrorOperation) -> Result<(), Error> { if self.is_closed() { diff --git a/cache2/src/config/runtime.rs b/cache2/src/config/runtime.rs index ce309af..2bb1abe 100644 --- a/cache2/src/config/runtime.rs +++ b/cache2/src/config/runtime.rs @@ -433,6 +433,9 @@ pub struct RuntimeOptions { /// Enable cumulative request, cache, and I/O counters. Defaults to false; /// health and managed-resource gauges remain available. pub statistics: bool, + /// Additional request accounting and latency distributions. Independent of + /// `statistics`; defaults disable all additional recorders. + pub stats: crate::StatsOptions, } impl Default for RuntimeOptions { @@ -448,6 +451,7 @@ impl Default for RuntimeOptions { l1_shards: DEFAULT_L1_SHARDS, write_flush_threshold_bytes: MAX_WRITE_FLUSH_THRESHOLD_BYTES, statistics: false, + stats: crate::StatsOptions::default(), } } } @@ -506,6 +510,7 @@ impl CacheConfig { let geometry = storage.geometry; let index_slots = storage.index_slots; runtime.resolve()?; + let stats_bytes = crate::stats::recording::Recorder::allocation_bytes(runtime.stats)?; if geometry.region_count <= runtime.append_shards { return Err(invalid_config( "append shards require valid geometry with one Active Region each plus one spare Region", @@ -520,6 +525,7 @@ impl CacheConfig { )?; let fixed_bytes = runtime_fixed_memory_bytes(index_slots, geometry.region_count)? .checked_add(l1_metadata_bytes) + .and_then(|bytes| bytes.checked_add(stats_bytes)) .ok_or_else(|| invalid_config("fixed memory requirements overflow"))?; let (reserved_memory_bytes, minimum_memory_bytes) = runtime.memory_requirements(geometry, fixed_bytes)?; diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index 07ad2d8..9d2f1a1 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -944,6 +944,8 @@ fn submit_cache_io_until( } pub trait IoEngine: Send + Sync { + /// Installed once during construction, before any requests are admitted. + fn set_latency_recorder(&self, recorder: crate::stats::recording::IoTiming); fn try_reserve_read(&self) -> io::Result; fn read_slot_waiter(&self) -> ReadSlotWaiter; fn submit_reserved_read( @@ -1009,6 +1011,7 @@ struct IoSlot { /// Dropping it before submission releases the slot immediately. pub struct ReadSlot { slot: IoSlot, + reserved_at: Option, } /// An async reservation handle backed by the engine's physical slot state. @@ -1165,6 +1168,7 @@ const fn active_write_slots(state: u64) -> usize { } struct RuntimeShared { + latency: std::sync::OnceLock, max_in_flight: usize, statistics_enabled: bool, accepting: AtomicBool, @@ -1196,6 +1200,7 @@ enum SlotWaitError { impl RuntimeShared { fn new(max_in_flight: usize, statistics_enabled: bool, read_wait_enabled: bool) -> Self { Self { + latency: std::sync::OnceLock::new(), max_in_flight, statistics_enabled, accepting: AtomicBool::new(true), @@ -1276,7 +1281,9 @@ impl RuntimeShared { .try_reserve_slot(false) .ok_or_else(|| io::Error::new(io::ErrorKind::WouldBlock, "no I/O slot is available"))?; slot.read_permit = permit; - Ok(ReadSlot { slot }) + let reserved_at = + (self.statistics_enabled || self.latency.get().is_some()).then(Instant::now); + Ok(ReadSlot { slot, reserved_at }) } fn stop_accepting_slots(&self) { @@ -1475,8 +1482,19 @@ impl RuntimeShared { self.requests_failed.fetch_add(1, Ordering::Relaxed); } } - if let Some(submitted_at) = submitted_at { - add_duration_ns(&self.request_time_ns, submitted_at.elapsed()); + } + if let Some(submitted_at) = submitted_at { + let elapsed = submitted_at.elapsed(); + if self.statistics_enabled { + add_duration_ns(&self.request_time_ns, elapsed); + } + if let Some(recorder) = self.latency.get() { + let outcome = match &status { + CompletionStatus::Completed => crate::IoOutcome::Completed, + CompletionStatus::Cancelled => crate::IoOutcome::Cancelled, + CompletionStatus::Failed(_) => crate::IoOutcome::Failed, + }; + recorder.record(elapsed, outcome); } } drop(slot); @@ -1654,8 +1672,9 @@ impl RuntimeInner { if let Err(error) = operation.validate() { return Err(SubmitError { error, operation }); } - let request_started = self.shared.statistics_enabled.then(Instant::now); - self.submit_with_slot(operation, slot.slot, request_started, true) + // Include buffer preparation and scheduling after the reservation. + // Dropped, unsubmitted reservations still produce no observation. + self.submit_with_slot(operation, slot.slot, slot.reserved_at, true) } #[cfg(test)] @@ -1737,7 +1756,9 @@ impl RuntimeInner { add_duration_ns(&self.shared.slot_wait_ns, slot_wait_started.elapsed()); } - let request_started = self.shared.statistics_enabled.then(Instant::now); + let request_started = (self.shared.statistics_enabled + || self.shared.latency.get().is_some()) + .then(Instant::now); self.submit_with_slot(operation, slot, request_started, nonblocking) } diff --git a/cache2/src/io/engine/posix.rs b/cache2/src/io/engine/posix.rs index eec286e..afc2315 100644 --- a/cache2/src/io/engine/posix.rs +++ b/cache2/src/io/engine/posix.rs @@ -174,6 +174,13 @@ impl BackendIoEngine { } impl IoEngine for BackendIoEngine { + fn set_latency_recorder(&self, recorder: crate::stats::recording::IoTiming) { + assert!( + self.inner.shared.latency.set(recorder).is_ok(), + "I/O recorder installed twice" + ); + } + fn try_reserve_read(&self) -> io::Result { self.inner.try_reserve_read() } diff --git a/cache2/src/io/engine/tests.rs b/cache2/src/io/engine/tests.rs index a0fa67c..e276e7c 100644 --- a/cache2/src/io/engine/tests.rs +++ b/cache2/src/io/engine/tests.rs @@ -395,6 +395,14 @@ async fn async_request_is_woken_by_driver_completion() { async fn dropping_async_wait_requests_bounded_cancellation() { let backend = Arc::new(BlockingBackend::default()); let engine: Arc = Arc::new(BackendIoEngine::new(backend.clone(), 1).unwrap()); + let recorder = Arc::new( + crate::stats::recording::Recorder::new(crate::StatsOptions { + io_latency: true, + ..crate::StatsOptions::default() + }) + .unwrap(), + ); + engine.set_latency_recorder(recorder.io_timing(crate::IoRole::Read).unwrap()); let resources = resources(); let request = submit_cache_io( engine.as_ref(), @@ -412,9 +420,82 @@ async fn dropping_async_wait_requests_bounded_cancellation() { waiter.abort(); assert!(waiter.await.unwrap_err().is_cancelled()); + assert!( + recorder + .io_snapshot() + .iter() + .all(|row| row.latency.count == 0) + ); backend.release(); engine.shutdown().unwrap(); assert_eq!(engine.in_flight(), 0); + assert_eq!( + recorder + .io_snapshot() + .iter() + .map(|row| row.latency.count) + .sum::(), + 1 + ); +} + +#[tokio::test] +async fn reserved_read_latency_includes_time_before_submission() { + for statistics_enabled in [false, true] { + let file = TestFile::new(); + file.file().set_len(4096).unwrap(); + let engine = BackendIoEngine::new_with_workers_and_statistics( + file.backend(), + 1, + 1, + statistics_enabled, + true, + ) + .unwrap(); + let recorder = Arc::new( + crate::stats::recording::Recorder::new(crate::StatsOptions { + io_latency: true, + ..crate::StatsOptions::default() + }) + .unwrap(), + ); + engine.set_latency_recorder(recorder.io_timing(crate::IoRole::Read).unwrap()); + drop(engine.try_reserve_read().unwrap()); + let slot = engine.try_reserve_read().unwrap(); + let reserved_at = slot.reserved_at.unwrap(); + tokio::time::sleep(Duration::from_millis(20)).await; + let before_submit = reserved_at.elapsed(); + assert!( + recorder + .io_snapshot() + .iter() + .all(|row| row.latency.count == 0) + ); + let completion = submit_cache_read( + &engine, + slot, + IoOperation::read(read_buffer(&resources(), 4096), 0), + ) + .unwrap() + .wait(&engine) + .unwrap(); + assert!(matches!(completion.status, CompletionStatus::Completed)); + engine.shutdown().unwrap(); + let snapshots = recorder.io_snapshot(); + let completed = snapshots + .iter() + .find(|row| { + row.role == crate::IoRole::Read && row.outcome == crate::IoOutcome::Completed + }) + .unwrap(); + assert_eq!(completed.latency.count, 1); + assert!(completed.latency.sum_ns >= before_submit.as_nanos()); + if statistics_enabled { + assert!( + u128::from(engine.stats().requests.request_time_ns) >= before_submit.as_nanos() + ); + } + } } #[tokio::test] @@ -1041,3 +1122,34 @@ fn quarantined_completion_does_not_return_a_potentially_live_buffer() { drop(shared); assert_eq!(resources.managed_memory_snapshot().current_bytes, 0); } + +#[test] +fn io_histograms_include_failures_when_legacy_statistics_are_disabled() { + let file = TestFile::new(); + let engine = + BackendIoEngine::new_with_workers_and_statistics(file.backend(), 1, 1, false, false) + .unwrap(); + let recorder = Arc::new( + crate::stats::recording::Recorder::new(crate::StatsOptions { + io_latency: true, + ..crate::StatsOptions::default() + }) + .unwrap(), + ); + engine.set_latency_recorder(recorder.io_timing(crate::IoRole::Read).unwrap()); + let resources = resources(); + let completion = engine + .read_exact_at(read_buffer(&resources, 4096), 0) + .unwrap() + .wait(); + assert!(matches!(completion.status, CompletionStatus::Failed(_))); + assert_eq!(engine.stats(), EngineIoSnapshot::default()); + let snapshot = recorder.io_snapshot(); + let failed = snapshot + .iter() + .find(|row| row.role == crate::IoRole::Read && row.outcome == crate::IoOutcome::Failed) + .unwrap(); + assert_eq!(failed.latency.count, 1); + assert!(failed.latency.valid); + engine.shutdown().unwrap(); +} diff --git a/cache2/src/io/engine/uring.rs b/cache2/src/io/engine/uring.rs index 02e5426..06f6147 100644 --- a/cache2/src/io/engine/uring.rs +++ b/cache2/src/io/engine/uring.rs @@ -245,6 +245,13 @@ impl UringIoEngine { } impl IoEngine for UringIoEngine { + fn set_latency_recorder(&self, recorder: crate::stats::recording::IoTiming) { + assert!( + self.inner.shared.latency.set(recorder).is_ok(), + "I/O recorder installed twice" + ); + } + fn try_reserve_read(&self) -> io::Result { self.inner.try_reserve_read() } diff --git a/cache2/src/lib.rs b/cache2/src/lib.rs index da6a711..7d09e89 100644 --- a/cache2/src/lib.rs +++ b/cache2/src/lib.rs @@ -70,3 +70,17 @@ mod resources; mod fixtures; #[cfg(test)] mod property_tests; + +mod stats; +pub use self::stats::CacheStatsSnapshot; +pub use self::stats::IoLatencySnapshot; +pub use self::stats::IoOutcome; +pub use self::stats::IoRole; +pub use self::stats::LATENCY_BUCKET_UPPER_BOUNDS_NS; +pub use self::stats::LatencyMode; +pub use self::stats::LatencySnapshot; +pub use self::stats::RequestLatencyScope; +pub use self::stats::RequestOperation; +pub use self::stats::RequestOutcome; +pub use self::stats::RequestStatsSnapshot; +pub use self::stats::StatsOptions; diff --git a/cache2/src/region/file_backend/mod.rs b/cache2/src/region/file_backend/mod.rs index 93eed41..9439a15 100644 --- a/cache2/src/region/file_backend/mod.rs +++ b/cache2/src/region/file_backend/mod.rs @@ -269,7 +269,7 @@ impl RegionStore> { ) -> io::Result> { self.runtime()? .data_plane()? - .get_async(key, tokio_handle) + .get_async(key, tokio_handle, None) .await } diff --git a/cache2/src/region/file_backend/tests.rs b/cache2/src/region/file_backend/tests.rs index d14e2ee..a939bfd 100644 --- a/cache2/src/region/file_backend/tests.rs +++ b/cache2/src/region/file_backend/tests.rs @@ -521,7 +521,7 @@ fn queued_l2_read_does_not_pin_warm_close() { let slot = plane.reserve_read_slot_for_test(); tokio_runtime.block_on(async { - let mut waiting = Box::pin(plane.get_async(b"queued-close", tokio_runtime.handle())); + let mut waiting = Box::pin(plane.get_async(b"queued-close", tokio_runtime.handle(), None)); assert_pending(waiting.as_mut(), "saturated read must enter the wait queue").await; store.close_warm().unwrap(); diff --git a/cache2/src/region/runtime/metrics.rs b/cache2/src/region/runtime/metrics.rs index 1d4adfc..578c997 100644 --- a/cache2/src/region/runtime/metrics.rs +++ b/cache2/src/region/runtime/metrics.rs @@ -34,6 +34,7 @@ static NEXT_METRICS_EPOCH: AtomicU64 = AtomicU64::new(1); pub struct RuntimeMetrics { metrics_epoch: u64, + pub stats: std::sync::Arc, pub lifecycle: AtomicU8, activity: Box<[ActivityMetrics]>, l2_read_overloads: AtomicU64, @@ -86,7 +87,7 @@ impl ActivityMetrics { } impl RuntimeMetrics { - pub fn new(shard_count: usize) -> io::Result { + pub fn new(shard_count: usize, stats: crate::StatsOptions) -> io::Result { let mut activity = Vec::new(); activity.try_reserve_exact(shard_count).map_err(|_| { io::Error::new( @@ -97,6 +98,7 @@ impl RuntimeMetrics { activity.resize_with(shard_count, ActivityMetrics::new); Ok(Self { metrics_epoch: NEXT_METRICS_EPOCH.fetch_add(1, Ordering::Relaxed), + stats: std::sync::Arc::new(crate::stats::recording::Recorder::new(stats)?), lifecycle: AtomicU8::new(LIFECYCLE_RUNNING), activity: activity.into_boxed_slice(), l2_read_overloads: AtomicU64::new(0), diff --git a/cache2/src/region/runtime/mod.rs b/cache2/src/region/runtime/mod.rs index 182155c..382fcfe 100644 --- a/cache2/src/region/runtime/mod.rs +++ b/cache2/src/region/runtime/mod.rs @@ -736,6 +736,10 @@ impl ShardControl { } impl RegionDataPlane { + pub fn stats_recorder(&self) -> &crate::stats::recording::Recorder { + &self.metrics.stats + } + pub fn new( core: Arc, data: DataSuperblock, @@ -758,7 +762,7 @@ impl RegionDataPlane { let config = configuration.runtime().clone(); core.configure_reclaim_workers(IoPoolTopology::reclaim(config.io_engine).max_in_flight)?; core.set_index_statistics_enabled(config.statistics); - let metrics = Arc::new(RuntimeMetrics::new(core.shard_count())?); + let metrics = Arc::new(RuntimeMetrics::new(core.shard_count(), config.stats)?); let operations = Arc::new(MutationGate::new()); let running = start_running( Arc::clone(&core), @@ -904,7 +908,7 @@ impl RegionDataPlane { #[cfg(test)] pub fn get(&self, key: &[u8]) -> io::Result> { - match self.prepare_get(key)? { + match self.prepare_get(key, None)? { PreparedGet::Complete(value) => Ok(value), PreparedGet::Pending(pending) => self.finish_get(pending.wait(), key), PreparedGet::Waiting(_) => Err(io::Error::other( @@ -917,8 +921,9 @@ impl RegionDataPlane { &self, key: &[u8], tokio_handle: &tokio::runtime::Handle, + guard: Option<&mut crate::stats::recording::RequestGuard<'_>>, ) -> io::Result> { - match self.prepare_get(key)? { + match self.prepare_get(key, guard)? { PreparedGet::Complete(value) => Ok(value), PreparedGet::Pending(pending) => { self.finish_get(pending.wait_async(tokio_handle).await, key) @@ -990,7 +995,11 @@ impl RegionDataPlane { } } - fn prepare_get(&self, key: &[u8]) -> io::Result { + fn prepare_get( + &self, + key: &[u8], + guard: Option<&mut crate::stats::recording::RequestGuard<'_>>, + ) -> io::Result { if key.len() > MAX_KEY_SIZE { if self.config.statistics { let activity = self.metrics.activity(0); @@ -1023,6 +1032,9 @@ impl RegionDataPlane { return Ok(PreparedGet::Complete(Some(HybridValueRead::L1(value)))); } MemoryLookup::Miss(token) => { + if let Some(guard) = guard { + guard.enter_l2(); + } if let Some(activity) = activity { RuntimeMetrics::increment(&activity.l1_misses); } @@ -1439,6 +1451,17 @@ fn start_running( IoPoolTopology::reclaim(config.io_engine), false, )?; + for (engines, role) in [ + (&read_engines, crate::IoRole::Read), + (&write_engines, crate::IoRole::Write), + (&reclaim_engines, crate::IoRole::Reclaim), + ] { + for engine in engines.iter() { + if let Some(timing) = metrics.stats.io_timing(role) { + engine.set_latency_recorder(timing); + } + } + } let mut shards = Vec::new(); shards.try_reserve_exact(shard_count).map_err(|_| { io::Error::new(io::ErrorKind::OutOfMemory, "cannot allocate shard controls") @@ -2589,7 +2612,7 @@ mod tests { #[test] fn read_resource_misses_remain_separately_observable() { - let metrics = RuntimeMetrics::new(1).unwrap(); + let metrics = RuntimeMetrics::new(1, crate::StatsOptions::default()).unwrap(); let activity = metrics.activity(0); RuntimeMetrics::add(&activity.l2_misses, 2); RuntimeMetrics::increment(&activity.l2_read_memory_misses); diff --git a/cache2/src/region/runtime/shutdown_tests.rs b/cache2/src/region/runtime/shutdown_tests.rs index c8c0a50..7311fd2 100644 --- a/cache2/src/region/runtime/shutdown_tests.rs +++ b/cache2/src/region/runtime/shutdown_tests.rs @@ -104,6 +104,9 @@ struct RacingEngine { } impl IoEngine for RacingEngine { + fn set_latency_recorder(&self, recorder: crate::stats::recording::IoTiming) { + self.inner.set_latency_recorder(recorder); + } fn try_reserve_read(&self) -> io::Result { self.inner.try_reserve_read() } diff --git a/cache2/src/stats/mod.rs b/cache2/src/stats/mod.rs new file mode 100644 index 0000000..9bcd9fe --- /dev/null +++ b/cache2/src/stats/mod.rs @@ -0,0 +1,342 @@ +// Copyright 2026 ScopeDB, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Optional request accounting and cumulative latency distributions. +//! +//! The legacy `RuntimeOptions::statistics` switch retains its existing behavior. +//! These options add public request outcomes and timing independently. Sampled +//! durations describe observed requests only; they are not full-population SLOs. + +use std::num::NonZeroU32; + +use crate::snapshot::CacheSnapshot; + +pub(crate) mod recording; + +/// Duration collection for one timing scope, selected before its outcome is known. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub enum LatencyMode { + /// No clocks or histogram updates for this scope. + #[default] + Off, + /// Observe every terminal operation in this scope; no sampling RNG work. + Full, + /// Observe each request with probability approximately `1 / interval`. + /// Counts and sums retain sampled semantics. One is equivalent to full. + Sampled { + /// Mean number of requests per observation. + interval: NonZeroU32, + }, +} + +/// Additional statistics allocated once per open, independent of legacy counters. +/// +/// Defaults allocate no counter or histogram stripes. Enable `request_counters` for complete +/// terminal accounting even when durations are sampled. All enabled storage is +/// included in the managed-memory plan. Configuration cannot change while open. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct StatsOptions { + /// Count every terminal get, put, put_l2 and delete by exclusive result. + pub request_counters: bool, + /// Time L1 hits from the first poll to return. L1 misses discard this timer. + pub l1_latency: LatencyMode, + /// Time reads from L1 miss through L2 lookup, waiting, I/O and promotion. + /// Excludes the initial L1 lookup. Early misses before L1 lookup are not timed. + pub l2_latency: LatencyMode, + /// Time put, put_l2 and delete from entry to their terminal result. + pub mutation_latency: LatencyMode, + /// Fully observe engine requests from slot reservation to terminal completion. + /// Split by foreground read, append write and reclaim read; not device latency. + pub io_latency: bool, + /// Fixed recorder stripes. Must be a power of two in 1..=64; defaults to 16. + /// Colliding threads use relaxed atomic updates, never a recorder lock. + pub shards: usize, +} + +impl Default for StatsOptions { + fn default() -> Self { + Self { + request_counters: false, + l1_latency: LatencyMode::Off, + l2_latency: LatencyMode::Off, + mutation_latency: LatencyMode::Off, + io_latency: false, + shards: 16, + } + } +} + +impl StatsOptions { + pub(crate) fn requests_enabled(self) -> bool { + self.request_counters + || self.l1_latency != LatencyMode::Off + || self.l2_latency != LatencyMode::Off + || self.mutation_latency != LatencyMode::Off + } + + pub(crate) fn latency_mode(self, scope: RequestLatencyScope) -> LatencyMode { + match scope { + RequestLatencyScope::L1Hit => self.l1_latency, + RequestLatencyScope::L2Lookup => self.l2_latency, + RequestLatencyScope::Mutation => self.mutation_latency, + } + } +} + +/// Public operation whose terminal outcome is recorded. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RequestOperation { + /// Lookup, starting at first poll; an unpolled future is not a request. + Get, + /// Acceptance into staging with best-effort L1 admission. + Put, + /// Acceptance into staging without L1 admission. + PutL2, + /// Bounded index deletion and best-effort L1 cleanup. + Delete, +} + +impl RequestOperation { + /// Stable low-cardinality metric label. + pub const fn as_str(self) -> &'static str { + match self { + Self::Get => "get", + Self::Put => "put", + Self::PutL2 => "put_l2", + Self::Delete => "delete", + } + } +} + +/// Mutually exclusive terminal result; internal events may overlap these counts. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RequestOutcome { + /// A value found in L1 at lookup, not an L2 value subsequently promoted. + L1Hit, + /// A validated L2 value, including values promoted to L1 before return. + L2Hit, + /// No value returned, including fail-open and post-close gets. + Miss, + /// Mutation accepted; does not imply flush or durability. + Accepted, + /// Explicit bounded admission or deadline rejection. + Overloaded, + /// Invalid public mutation input. + InvalidInput, + /// Mutation attempted after shutdown or loss of availability. + Unavailable, + /// Another explicit public error; safe fail-open reads remain misses. + Error, + /// A polled get future dropped before returning; duration is partial lifetime. + Cancelled, +} + +impl RequestOutcome { + /// Stable low-cardinality metric label. + pub const fn as_str(self) -> &'static str { + match self { + Self::L1Hit => "l1_hit", + Self::L2Hit => "l2_hit", + Self::Miss => "miss", + Self::Accepted => "accepted", + Self::Overloaded => "overloaded", + Self::InvalidInput => "invalid_input", + Self::Unavailable => "unavailable", + Self::Error => "error", + Self::Cancelled => "cancelled", + } + } +} + +/// Engine ownership, independent of buffered/direct file-operation accounting. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum IoRole { + /// Foreground L2 lookups. + Read, + /// Foreground and reinsertion append batches. + Write, + /// Source Region reads by reclaim workers. + Reclaim, +} + +impl IoRole { + /// Stable metric label. + pub const fn as_str(self) -> &'static str { + match self { + Self::Read => "read", + Self::Write => "write", + Self::Reclaim => "reclaim", + } + } +} + +/// Terminal engine result, distinct from caller cancellation or timeout. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum IoOutcome { + /// Operation completed successfully. + Completed, + /// Engine completed the operation through cancellation. + Cancelled, + /// Engine reported a failure. + Failed, +} + +impl IoOutcome { + /// Stable metric label. + pub const fn as_str(self) -> &'static str { + match self { + Self::Completed => "completed", + Self::Cancelled => "cancelled", + Self::Failed => "failed", + } + } +} + +/// Version-one explicit duration bucket boundaries, in nanoseconds, inclusive. +/// +/// A final unbounded bucket follows these finite bounds. The layout includes +/// common SLO thresholds and is shared by all instances. Quantiles have bucket +/// resolution, not an exact-value or relative-error guarantee. +pub const LATENCY_BUCKET_UPPER_BOUNDS_NS: &[u64] = &[ + 0, + 50, + 100, + 250, + 500, + 1_000, + 2_500, + 5_000, + 10_000, + 25_000, + 50_000, + 100_000, + 250_000, + 500_000, + 1_000_000, + 2_500_000, + 5_000_000, + 10_000_000, + 25_000_000, + 50_000_000, + 100_000_000, + 250_000_000, + 500_000_000, + 1_000_000_000, + 2_500_000_000, + 5_000_000_000, + 10_000_000_000, + 30_000_000_000, + 60_000_000_000, +]; + +/// Cumulative observations for one population since open. +/// +/// Snapshots never reset or consume counts. Concurrent updates can appear at +/// slightly different times in buckets and sum; quiescent snapshots are exact. +/// The sum covers actual observed durations, including overflow-bucket values. +#[non_exhaustive] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct LatencySnapshot { + /// Noncumulative bucket counts: one per finite bound plus the unbounded bucket. + pub bucket_counts: Box<[u64]>, + /// Sum of these bucket counts, not an estimate of full request volume. + pub count: u64, + /// Sum of observed durations in nanoseconds, widened when stripes are merged. + pub sum_ns: u128, + /// False after recorder arithmetic overflow. Do not export invalid distributions. + pub valid: bool, +} + +/// Interval measured by a request duration distribution. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum RequestLatencyScope { + /// First public get poll through return of an L1 hit. + L1Hit, + /// L1 miss through terminal result, including index lookup, waiting and promotion. + /// This is not the full public get duration. + L2Lookup, + /// Public put, put_l2 or delete entry through its terminal result. + Mutation, +} + +impl RequestLatencyScope { + pub(crate) fn for_request(operation: RequestOperation, outcome: RequestOutcome) -> Self { + match operation { + RequestOperation::Get if outcome == RequestOutcome::L1Hit => Self::L1Hit, + RequestOperation::Get => Self::L2Lookup, + _ => Self::Mutation, + } + } +} + +/// One allowed public-operation/result combination. +#[non_exhaustive] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RequestStatsSnapshot { + /// Public operation. + pub operation: RequestOperation, + /// Exclusive terminal result. + pub outcome: RequestOutcome, + /// Full event count, or None when request counters are disabled. + pub count: Option, + /// Interval covered by this row's optional duration distribution. + pub latency_scope: RequestLatencyScope, + /// Collection mode for this row, independent of other tiers and mutations. + pub latency_mode: LatencyMode, + /// Full or sampled duration distribution, or None when timing is disabled. + /// L1 hits cover the public lookup; other get rows start at L1 miss. + /// Early misses before L1 lookup contribute only to the request counter. + pub latency: Option, +} + +/// Fully observed engine durations for one role/result combination. +#[non_exhaustive] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct IoLatencySnapshot { + /// Engine ownership. + pub role: IoRole, + /// Engine terminal result. + pub outcome: IoOutcome, + /// Full duration observations, irrespective of public-request sampling. + pub latency: LatencySnapshot, +} + +/// Cumulative monitoring view without metadata scans or recorder locks. +/// +/// Legacy activity availability is given by summary.statistics_enabled; new +/// families use explicit options and optional values. The application owns +/// collection scheduling, timestamp assignment and export to its monitoring SDK. Empty vectors +/// indicate a disabled family, not a zero-event population. Returned allocations belong to +/// the caller; cache-owned recorder storage is fixed and charged at open. +#[non_exhaustive] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct CacheStatsSnapshot { + /// Existing health, resource, activity, I/O and reclaim metrics and reset epoch. + pub summary: CacheSnapshot, + /// Exact collection settings for this open. + pub options: StatsOptions, + /// Bytes reserved for counter/histogram stripes, excluding this snapshot. + /// Recorder controls are covered by the fixed runtime control reservation. + pub recorder_bytes: usize, + /// Finite operation/result series, independent of user keys and error strings. + pub requests: Vec, + /// Full I/O duration distributions; empty when I/O latency is disabled. + pub io_latency: Vec, +} diff --git a/cache2/src/stats/recording.rs b/cache2/src/stats/recording.rs new file mode 100644 index 0000000..a785869 --- /dev/null +++ b/cache2/src/stats/recording.rs @@ -0,0 +1,578 @@ +// Copyright 2026 ScopeDB, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::cell::Cell; +use std::io; +use std::sync::Arc; +use std::sync::atomic::AtomicBool; +use std::sync::atomic::AtomicU64; +use std::sync::atomic::Ordering::Relaxed; +use std::time::Duration; +use std::time::Instant; + +use super::*; +use crate::Error; +use crate::ErrorKind; + +const BUCKETS: usize = LATENCY_BUCKET_UPPER_BOUNDS_NS.len() + 1; +const REQUEST_SERIES: &[(RequestOperation, RequestOutcome)] = &[ + (RequestOperation::Get, RequestOutcome::L1Hit), + (RequestOperation::Get, RequestOutcome::L2Hit), + (RequestOperation::Get, RequestOutcome::Miss), + (RequestOperation::Get, RequestOutcome::Overloaded), + (RequestOperation::Get, RequestOutcome::Error), + (RequestOperation::Get, RequestOutcome::Cancelled), + (RequestOperation::Put, RequestOutcome::Accepted), + (RequestOperation::Put, RequestOutcome::Overloaded), + (RequestOperation::Put, RequestOutcome::InvalidInput), + (RequestOperation::Put, RequestOutcome::Unavailable), + (RequestOperation::Put, RequestOutcome::Error), + (RequestOperation::PutL2, RequestOutcome::Accepted), + (RequestOperation::PutL2, RequestOutcome::Overloaded), + (RequestOperation::PutL2, RequestOutcome::InvalidInput), + (RequestOperation::PutL2, RequestOutcome::Unavailable), + (RequestOperation::PutL2, RequestOutcome::Error), + (RequestOperation::Delete, RequestOutcome::Accepted), + (RequestOperation::Delete, RequestOutcome::Overloaded), + (RequestOperation::Delete, RequestOutcome::InvalidInput), + (RequestOperation::Delete, RequestOutcome::Unavailable), + (RequestOperation::Delete, RequestOutcome::Error), +]; +const REQUEST_INDEX: [[usize; 9]; 4] = { + let mut indices = [[usize::MAX; 9]; 4]; + let mut index = 0; + while index < REQUEST_SERIES.len() { + let (operation, outcome) = REQUEST_SERIES[index]; + indices[operation as usize][outcome as usize] = index; + index += 1; + } + indices +}; +const IO_ROLES: [IoRole; 3] = [IoRole::Read, IoRole::Write, IoRole::Reclaim]; +const IO_OUTCOMES: [IoOutcome; 3] = [ + IoOutcome::Completed, + IoOutcome::Cancelled, + IoOutcome::Failed, +]; +const IO_SERIES: usize = IO_ROLES.len() * IO_OUTCOMES.len(); + +// Keep independent writers' counters off the same cache line. Histograms group +// all outcomes of a stripe, rather than routing hot keys to the same recorder. +#[repr(align(128))] +struct Counters([AtomicU64; REQUEST_SERIES.len()]); + +#[repr(align(128))] +struct Histogram { + buckets: [AtomicU64; BUCKETS], + sum_ns: AtomicU64, + invalid: AtomicBool, +} + +impl Histogram { + fn new() -> Self { + Self { + buckets: std::array::from_fn(|_| AtomicU64::new(0)), + sum_ns: AtomicU64::new(0), + invalid: AtomicBool::new(false), + } + } + + fn record(&self, duration: Duration) { + let Ok(ns) = u64::try_from(duration.as_nanos()) else { + self.invalid.store(true, Relaxed); + return; + }; + let bucket = LATENCY_BUCKET_UPPER_BOUNDS_NS.partition_point(|bound| *bound < ns); + let old_count = self.buckets[bucket].fetch_add(1, Relaxed); + let old_sum = self.sum_ns.fetch_add(ns, Relaxed); + if old_count == u64::MAX || old_sum.checked_add(ns).is_none() { + self.invalid.store(true, Relaxed); + } + } +} + +static NEXT_THREAD: AtomicU64 = AtomicU64::new(1); +thread_local! { + // A bounded scalar per thread, not a map of cache instances. A different + // sampling rate uses the same uniform random stream without reinitializing. + static THREAD: Cell<(u64, u64)> = { + let id = NEXT_THREAD.fetch_add(1, Relaxed); + Cell::new((id, id.wrapping_mul(0x9e3779b97f4a7c15).max(1))) + }; +} + +fn thread_sample(mode: LatencyMode) -> (usize, bool) { + THREAD.with(|cell| { + let (id, mut state) = cell.get(); + let sampled = match mode { + LatencyMode::Off => false, + LatencyMode::Full => true, + LatencyMode::Sampled { interval } if interval.get() == 1 => true, + LatencyMode::Sampled { interval } => { + state ^= state >> 12; + state ^= state << 25; + state ^= state >> 27; + cell.set((id, state)); + state.wrapping_mul(0x2545f4914f6cdd1d) <= u64::MAX / u64::from(interval.get()) + } + }; + (id as usize, sampled) + }) +} + +pub struct Recorder { + options: StatsOptions, + counters: Box<[Counters]>, + request_histograms: Box<[Histogram]>, + histogram_indices: [usize; REQUEST_SERIES.len()], + histogram_series: usize, + io_histograms: Box<[Histogram]>, +} + +impl Recorder { + pub fn allocation_bytes(options: StatsOptions) -> io::Result { + if !options.shards.is_power_of_two() || options.shards > 64 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "stats shards must be a power of two in 1..=64", + )); + } + let counters = usize::from(options.request_counters) * size_of::(); + let request = REQUEST_SERIES + .iter() + .filter(|&&(operation, outcome)| { + options.latency_mode(RequestLatencyScope::for_request(operation, outcome)) + != LatencyMode::Off + }) + .count(); + let io = usize::from(options.io_latency) * IO_SERIES; + Ok(options.shards * (counters + (request + io) * size_of::())) + } + + pub fn new(options: StatsOptions) -> io::Result { + Self::allocation_bytes(options)?; + let counters = allocate( + usize::from(options.request_counters) * options.shards, + || Counters(std::array::from_fn(|_| AtomicU64::new(0))), + )?; + let mut histogram_indices = [usize::MAX; REQUEST_SERIES.len()]; + let mut histogram_series = 0; + for (index, &(operation, outcome)) in REQUEST_SERIES.iter().enumerate() { + if options.latency_mode(RequestLatencyScope::for_request(operation, outcome)) + != LatencyMode::Off + { + histogram_indices[index] = histogram_series; + histogram_series += 1; + } + } + Ok(Self { + options, + counters, + request_histograms: allocate(options.shards * histogram_series, Histogram::new)?, + histogram_indices, + histogram_series, + io_histograms: allocate( + usize::from(options.io_latency) * options.shards * IO_SERIES, + Histogram::new, + )?, + }) + } + + #[inline] + pub fn begin(&self, operation: RequestOperation) -> RequestGuard<'_> { + if !self.options.requests_enabled() { + return RequestGuard { + recorder: None, + operation, + stripe: 0, + timer: None, + }; + } + let scope = if operation == RequestOperation::Get { + RequestLatencyScope::L1Hit + } else { + RequestLatencyScope::Mutation + }; + let mode = self.options.latency_mode(scope); + let (thread, sampled) = if mode == LatencyMode::Off && !self.options.request_counters { + (0, false) + } else { + thread_sample(mode) + }; + RequestGuard { + recorder: Some(self), + operation, + stripe: thread & (self.options.shards - 1), + timer: sampled.then(|| (scope, Instant::now())), + } + } + + pub fn io_timing(self: &Arc, role: IoRole) -> Option { + self.options.io_latency.then(|| IoTiming { + recorder: Arc::clone(self), + role, + }) + } + + pub fn snapshot(&self, summary: CacheSnapshot) -> CacheStatsSnapshot { + let mut requests = Vec::new(); + if self.options.requests_enabled() { + for (index, &(operation, outcome)) in REQUEST_SERIES.iter().enumerate() { + let scope = RequestLatencyScope::for_request(operation, outcome); + requests.push(RequestStatsSnapshot { + operation, + outcome, + count: self.options.request_counters.then(|| { + self.counters.iter().fold(0u64, |sum, stripe| { + sum.saturating_add(stripe.0[index].load(Relaxed)) + }) + }), + latency_scope: scope, + latency_mode: self.options.latency_mode(scope), + latency: (self.histogram_indices[index] != usize::MAX).then(|| { + histogram_snapshot( + &self.request_histograms, + self.histogram_series, + self.histogram_indices[index], + ) + }), + }); + } + } + CacheStatsSnapshot { + summary, + options: self.options, + recorder_bytes: self.counters.len() * size_of::() + + (self.request_histograms.len() + self.io_histograms.len()) + * size_of::(), + requests, + io_latency: self.io_snapshot(), + } + } + pub fn io_snapshot(&self) -> Vec { + let mut io_latency = Vec::new(); + if self.options.io_latency { + for role in IO_ROLES { + for outcome in IO_OUTCOMES { + io_latency.push(IoLatencySnapshot { + role, + outcome, + latency: histogram_snapshot( + &self.io_histograms, + IO_SERIES, + io_index(role, outcome), + ), + }); + } + } + } + io_latency + } +} + +fn allocate(count: usize, init: impl FnMut() -> T) -> io::Result> { + let mut items = Vec::new(); + items.try_reserve_exact(count).map_err(|_| { + io::Error::new( + io::ErrorKind::OutOfMemory, + "cannot allocate statistics recorders", + ) + })?; + items.resize_with(count, init); + Ok(items.into_boxed_slice()) +} + +fn histogram_snapshot(histograms: &[Histogram], series: usize, index: usize) -> LatencySnapshot { + let mut snapshot = LatencySnapshot { + bucket_counts: vec![0u64; BUCKETS].into_boxed_slice(), + count: 0, + sum_ns: 0, + valid: true, + }; + for histogram in histograms.iter().skip(index).step_by(series) { + for (sum, bucket) in snapshot.bucket_counts.iter_mut().zip(&histogram.buckets) { + match sum.checked_add(bucket.load(Relaxed)) { + Some(value) => *sum = value, + None => { + snapshot.valid = false; + *sum = u64::MAX; + } + } + } + snapshot.sum_ns += u128::from(histogram.sum_ns.load(Relaxed)); + snapshot.valid &= !histogram.invalid.load(Relaxed); + } + for count in &snapshot.bucket_counts { + match snapshot.count.checked_add(*count) { + Some(value) => snapshot.count = value, + None => { + snapshot.valid = false; + snapshot.count = u64::MAX; + } + } + } + snapshot +} + +pub struct RequestGuard<'a> { + recorder: Option<&'a Recorder>, + operation: RequestOperation, + stripe: usize, + timer: Option<(RequestLatencyScope, Instant)>, +} + +impl RequestGuard<'_> { + pub fn enter_l2(&mut self) { + let Some(recorder) = self.recorder else { + return; + }; + let (thread, sampled) = thread_sample(recorder.options.l2_latency); + self.stripe = thread & (recorder.options.shards - 1); + self.timer = sampled.then(|| (RequestLatencyScope::L2Lookup, Instant::now())); + } + + #[inline] + pub fn finish(mut self, result: &Result, success: RequestOutcome) { + if self.recorder.is_none() { + return; + } + let outcome = match result { + Ok(_) => success, + Err(error) => match error.kind() { + ErrorKind::Overloaded => RequestOutcome::Overloaded, + ErrorKind::InvalidInput if self.operation != RequestOperation::Get => { + RequestOutcome::InvalidInput + } + ErrorKind::Unavailable if self.operation != RequestOperation::Get => { + RequestOutcome::Unavailable + } + _ => RequestOutcome::Error, + }, + }; + self.record(outcome); + } + + fn record(&mut self, outcome: RequestOutcome) { + let Some(recorder) = self.recorder.take() else { + return; + }; + let index = REQUEST_INDEX[self.operation as usize][outcome as usize]; + // Stop the clock before recorder work; caller-side observers measure its + // overhead separately. Count once at the same terminal boundary. + let scope = RequestLatencyScope::for_request(self.operation, outcome); + let elapsed = self + .timer + .filter(|(timed_scope, _)| *timed_scope == scope) + .map(|(_, start)| start.elapsed()); + if recorder.options.request_counters { + recorder.counters[self.stripe].0[index].fetch_add(1, Relaxed); + } + if let Some(elapsed) = elapsed { + recorder.request_histograms + [self.stripe * recorder.histogram_series + recorder.histogram_indices[index]] + .record(elapsed); + } + } +} + +impl Drop for RequestGuard<'_> { + fn drop(&mut self) { + if self.operation == RequestOperation::Get && !std::thread::panicking() { + self.record(RequestOutcome::Cancelled); + } + } +} + +#[derive(Clone)] +pub struct IoTiming { + recorder: Arc, + role: IoRole, +} + +impl IoTiming { + pub fn record(&self, duration: Duration, outcome: IoOutcome) { + let index = io_index(self.role, outcome); + let stripe = thread_sample(LatencyMode::Off).0 & (self.recorder.options.shards - 1); + self.recorder.io_histograms[stripe * IO_SERIES + index].record(duration); + } +} + +fn io_index(role: IoRole, outcome: IoOutcome) -> usize { + role as usize * IO_OUTCOMES.len() + outcome as usize +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn histogram_boundaries_sum_overflow_and_invalidity() { + let histogram = Histogram::new(); + let mut expected_sum = 0u128; + for &bound in LATENCY_BUCKET_UPPER_BOUNDS_NS { + histogram.record(Duration::from_nanos(bound)); + expected_sum += u128::from(bound); + } + histogram.record(Duration::from_nanos(60_000_000_001)); + expected_sum += 60_000_000_001; + let snapshot = histogram_snapshot(std::slice::from_ref(&histogram), 1, 0); + assert!(snapshot.valid); + assert_eq!(snapshot.bucket_counts.as_ref(), &[1; BUCKETS]); + assert_eq!(snapshot.count, BUCKETS as u64); + assert_eq!(snapshot.sum_ns, expected_sum); + histogram.sum_ns.store(u64::MAX, Relaxed); + histogram.record(Duration::from_nanos(1)); + assert!(!histogram_snapshot(&[histogram], 1, 0).valid); + } + + #[test] + fn guards_count_once_and_cancellation_is_a_separate_population() { + let recorder = Recorder::new(StatsOptions { + request_counters: true, + l1_latency: LatencyMode::Full, + l2_latency: LatencyMode::Full, + mutation_latency: LatencyMode::Full, + shards: 1, + ..StatsOptions::default() + }) + .unwrap(); + let mut hit = recorder.begin(RequestOperation::Get); + hit.enter_l2(); + hit.finish(&Ok::<_, Error>(()), RequestOutcome::L2Hit); + let mut cancelled = recorder.begin(RequestOperation::Get); + cancelled.enter_l2(); + drop(cancelled); + assert_eq!(recorder.counters[0].0[1].load(Relaxed), 1); + assert_eq!(recorder.counters[0].0[5].load(Relaxed), 1); + assert_eq!( + histogram_snapshot(&recorder.request_histograms, REQUEST_SERIES.len(), 1).count, + 1 + ); + assert_eq!( + histogram_snapshot(&recorder.request_histograms, REQUEST_SERIES.len(), 5).count, + 1 + ); + assert_eq!( + recorder.counters[0] + .0 + .iter() + .map(|count| count.load(Relaxed)) + .sum::(), + 2 + ); + } + + #[test] + fn l2_timing_starts_at_the_l1_miss() { + let recorder = Recorder::new(StatsOptions { + l2_latency: LatencyMode::Full, + shards: 1, + ..StatsOptions::default() + }) + .unwrap(); + let mut guard = recorder.begin(RequestOperation::Get); + assert!( + guard.timer.is_none(), + "L2-only timing must not start an L1 clock" + ); + let transition = Instant::now(); + guard.enter_l2(); + assert!(guard.timer.unwrap().1 >= transition); + drop(guard); + let recorder = Recorder::new(StatsOptions { + l1_latency: LatencyMode::Full, + l2_latency: LatencyMode::Full, + ..StatsOptions::default() + }) + .unwrap(); + let mut guard = recorder.begin(RequestOperation::Get); + guard.timer = Some(( + RequestLatencyScope::L1Hit, + Instant::now() - Duration::from_secs(3600), + )); + let transition = Instant::now(); + guard.enter_l2(); + assert!( + guard.timer.unwrap().1 >= transition, + "L2 must discard the L1 timer" + ); + guard.finish(&Ok::<_, Error>(()), RequestOutcome::Miss); + } + + #[test] + fn sampling_is_independent_of_callers_switching_rates() { + THREAD.with(|thread| thread.set((1, 12345))); + let mut sampled = [0usize; 2]; + for _ in 0..100_000 { + for (index, interval) in [16, 64].into_iter().enumerate() { + sampled[index] += usize::from( + thread_sample(LatencyMode::Sampled { + interval: NonZeroU32::new(interval).unwrap(), + }) + .1, + ); + } + } + assert!((5900..6600).contains(&sampled[0]), "{sampled:?}"); + assert!((1400..1750).contains(&sampled[1]), "{sampled:?}"); + } + + #[test] + fn concurrent_readers_do_not_consume_or_lose_observations() { + let recorder = Recorder::new(StatsOptions { + request_counters: true, + l1_latency: LatencyMode::Full, + l2_latency: LatencyMode::Full, + mutation_latency: LatencyMode::Full, + shards: 1, + ..StatsOptions::default() + }) + .unwrap(); + std::thread::scope(|scope| { + for _ in 0..8 { + let recorder = &recorder; + scope.spawn(move || { + for _ in 0..10_000 { + recorder + .begin(RequestOperation::Get) + .finish(&Ok::<_, Error>(()), RequestOutcome::L1Hit); + } + }); + } + for _ in 0..2 { + let recorder = &recorder; + scope.spawn(move || { + let mut previous = 0; + for _ in 0..1000 { + let snapshot = histogram_snapshot( + &recorder.request_histograms, + REQUEST_SERIES.len(), + 0, + ); + assert!(snapshot.valid); + assert!(snapshot.count >= previous); + previous = snapshot.count; + } + }); + } + }); + assert_eq!(recorder.counters[0].0[0].load(Relaxed), 80_000); + let first = histogram_snapshot(&recorder.request_histograms, REQUEST_SERIES.len(), 0); + assert_eq!(first.count, 80_000); + assert_eq!( + first, + histogram_snapshot(&recorder.request_histograms, REQUEST_SERIES.len(), 0) + ); + } +} diff --git a/examples/Cargo.toml b/examples/Cargo.toml index 729cf58..b3f297d 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -30,5 +30,9 @@ tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } name = "logforth" path = "src/logforth.rs" +[[example]] +name = "stats" +path = "src/stats.rs" + [lints] workspace = true diff --git a/examples/src/stats.rs b/examples/src/stats.rs new file mode 100644 index 0000000..09b1e5f --- /dev/null +++ b/examples/src/stats.rs @@ -0,0 +1,58 @@ +// Copyright 2026 ScopeDB, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::env; +use std::num::NonZeroU32; + +use cache2::Cache; +use cache2::CacheConfig; +use cache2::LatencyMode; +use cache2::RuntimeOptions; +use cache2::StatsOptions; +use cache2::StorageOptions; + +#[tokio::main] +async fn main() -> Result<(), Box> { + let path = env::args_os() + .nth(1) + .ok_or("usage: stats ")?; + let storage = StorageOptions { + region_size_bytes: 1024 * 1024, + ..StorageOptions::new(8 * 1024 * 1024) + } + .build()?; + let runtime = RuntimeOptions { + statistics: true, + stats: StatsOptions { + request_counters: true, + l1_latency: LatencyMode::Sampled { + interval: NonZeroU32::new(64).unwrap(), + }, + l2_latency: LatencyMode::Full, + io_latency: true, + ..StatsOptions::default() + }, + ..RuntimeOptions::default() + }; + let cache = Cache::open(path, CacheConfig::new(storage, runtime)?).await?; + cache.put("example", "value")?; + let _value = cache.get("example").await?; + cache.drain().await?; + // Pass this owned, cumulative snapshot to the application's metrics adapter. + // It carries collection modes/scopes; bucket bounds and sums use nanoseconds. + let stats = cache.stats_snapshot()?; + println!("{stats:#?}"); + cache.close_fast().await?; + Ok(()) +} diff --git a/tests-integration/tests/cache.rs b/tests-integration/tests/cache.rs index 039f99f..adbff87 100644 --- a/tests-integration/tests/cache.rs +++ b/tests-integration/tests/cache.rs @@ -568,7 +568,7 @@ async fn runtime_options_can_change_across_a_warm_reopen() { cache.drain().await.unwrap(); cache.close_warm().await.unwrap(); - let retuned = RuntimeOptions { + let reopened_options = RuntimeOptions { io_engine: IoEngineConfig::Posix(PosixIoConfig::new(7, 2, 2)), l1_capacity_bytes: 2 * 1024 * 1024, l1_eviction_policy: L1EvictionPolicy::S3Fifo, @@ -583,7 +583,7 @@ async fn runtime_options_can_change_across_a_warm_reopen() { }; let reopened = Cache::open( &files.data, - CacheConfig::new(test_storage(), retuned).unwrap(), + CacheConfig::new(test_storage(), reopened_options).unwrap(), ) .await .unwrap(); @@ -1292,3 +1292,197 @@ async fn unsupported_cache_format_versions_cold_start_empty() { reopened.close_fast().await.unwrap(); } } + +#[tokio::test] +async fn structured_stats_track_original_hit_tier_and_full_distributions() { + use cache2::IoOutcome; + use cache2::IoRole; + use cache2::LatencyMode; + use cache2::RequestOperation; + use cache2::RequestOutcome; + use cache2::StatsOptions; + + let files = TestCache::new("structured-stats"); + let mut options = test_runtime_options(1, 2); + options.statistics = false; + options.stats = StatsOptions { + request_counters: true, + l1_latency: LatencyMode::Full, + l2_latency: LatencyMode::Full, + mutation_latency: LatencyMode::Full, + io_latency: true, + shards: 2, + }; + let base = CacheConfig::new(test_storage(), test_runtime_options(1, 2)).unwrap(); + let config = CacheConfig::new(test_storage(), options).unwrap(); + let cache = Cache::open(&files.data, config.clone()).await.unwrap(); + let before = cache.stats_snapshot().unwrap(); + assert_eq!( + config.minimum_memory_bytes() - base.minimum_memory_bytes(), + before.recorder_bytes + ); + assert!(before.summary.managed_memory_bytes >= before.recorder_bytes); + // Future construction must not count or start timing. + drop(cache.get(b"never-polled")); + cache.put_l2(b"l2", b"value").unwrap(); + cache.drain().await.unwrap(); + let value = cache.get(b"l2").await.unwrap().unwrap(); + assert_eq!(value.as_ref(), b"value"); + assert_eq!(value.tier(), CacheTier::L2); + assert!(cache.get(b"l2").await.unwrap().is_some()); + assert!(cache.get(b"missing").await.unwrap().is_none()); + assert_eq!( + cache.put(vec![0; 4097], b"value").unwrap_err().kind(), + ErrorKind::InvalidInput + ); + cache.delete(b"l2").unwrap(); + let stats = cache.stats_snapshot().unwrap(); + assert!(!stats.summary.statistics_enabled); + assert_eq!(stats.summary.l2_hits, 0); + for (operation, outcome) in [ + (RequestOperation::PutL2, RequestOutcome::Accepted), + (RequestOperation::Get, RequestOutcome::L2Hit), + (RequestOperation::Get, RequestOutcome::L1Hit), + (RequestOperation::Get, RequestOutcome::Miss), + (RequestOperation::Put, RequestOutcome::InvalidInput), + (RequestOperation::Delete, RequestOutcome::Accepted), + ] { + let row = stats + .requests + .iter() + .find(|row| row.operation == operation && row.outcome == outcome) + .unwrap(); + assert_eq!(row.count, Some(1), "{row:?}"); + assert_eq!(row.latency.as_ref().unwrap().count, 1); + } + assert_eq!( + stats + .requests + .iter() + .map(|row| row.count.unwrap()) + .sum::(), + 6 + ); + let read = stats + .io_latency + .iter() + .find(|row| row.role == IoRole::Read && row.outcome == IoOutcome::Completed) + .unwrap(); + assert_eq!(read.latency.count, 1); + let write = stats + .io_latency + .iter() + .find(|row| row.role == IoRole::Write && row.outcome == IoOutcome::Completed) + .unwrap(); + assert!(write.latency.count >= 1); + let epoch = stats.summary.metrics_epoch; + cache.close_warm().await.unwrap(); + assert!(cache.stats_snapshot().is_err()); + drop(cache); + let reopened = Cache::open(&files.data, config).await.unwrap(); + let reset = reopened.stats_snapshot().unwrap(); + assert_ne!(reset.summary.metrics_epoch, epoch); + assert!(reset.requests.iter().all(|row| row.count == Some(0))); + reopened.close_fast().await.unwrap(); +} + +#[tokio::test] +async fn disabled_stats_do_not_allocate_recorders() { + let files = TestCache::new("disabled-stats"); + let cache = Cache::open(&files.data, test_config(1)).await.unwrap(); + let stats = cache.stats_snapshot().unwrap(); + assert!(stats.requests.is_empty()); + assert!(stats.io_latency.is_empty()); + assert_eq!(stats.recorder_bytes, 0); + cache.close_fast().await.unwrap(); +} + +#[tokio::test] +async fn tier_latency_modes_are_independent_and_describe_their_scope() { + use std::num::NonZeroU32; + + use cache2::LatencyMode; + use cache2::RequestOperation; + use cache2::RequestOutcome; + use cache2::StatsOptions; + + let sampled = LatencyMode::Sampled { + interval: NonZeroU32::new(64).unwrap(), + }; + for (l1_latency, l2_latency) in [ + (LatencyMode::Off, LatencyMode::Full), + (sampled, LatencyMode::Full), + (LatencyMode::Full, LatencyMode::Off), + (LatencyMode::Full, sampled), + ] { + let files = TestCache::new("tier-latency"); + let mut options = test_runtime_options(1, 2); + options.stats = StatsOptions { + request_counters: true, + l1_latency, + l2_latency, + ..StatsOptions::default() + }; + let cache = Cache::open( + &files.data, + CacheConfig::new(test_storage(), options).unwrap(), + ) + .await + .unwrap(); + cache.put_l2(b"key", b"value").unwrap(); + cache.drain().await.unwrap(); + assert_eq!( + cache.get(b"key").await.unwrap().unwrap().tier(), + CacheTier::L2 + ); + for _ in 0..4096 { + assert_eq!( + cache.get(b"key").await.unwrap().unwrap().tier(), + CacheTier::L1 + ); + assert!(cache.get(b"absent").await.unwrap().is_none()); + } + // Rejected before the L1 lookup: count the miss without inventing an L2 duration. + assert!(cache.get(vec![0; 4097]).await.unwrap().is_none()); + let stats = cache.stats_snapshot().unwrap(); + for (outcome, mode, count, scope) in [ + ( + RequestOutcome::L1Hit, + l1_latency, + 4096, + cache2::RequestLatencyScope::L1Hit, + ), + ( + RequestOutcome::Miss, + l2_latency, + 4097, + cache2::RequestLatencyScope::L2Lookup, + ), + ] { + let row = stats + .requests + .iter() + .find(|row| row.operation == RequestOperation::Get && row.outcome == outcome) + .unwrap(); + assert_eq!(row.count, Some(count)); + assert_eq!(row.latency_mode, mode); + assert_eq!(row.latency_scope, scope); + match mode { + LatencyMode::Off => assert!(row.latency.is_none()), + LatencyMode::Full => assert_eq!(row.latency.as_ref().unwrap().count, 4096), + LatencyMode::Sampled { .. } => { + assert!((10..140).contains(&row.latency.as_ref().unwrap().count)) + } + _ => unreachable!(), + } + } + assert!( + stats + .requests + .iter() + .filter(|row| row.operation != RequestOperation::Get) + .all(|row| row.latency.is_none()) + ); + cache.close_fast().await.unwrap(); + } +} diff --git a/tests-integration/tests/config.rs b/tests-integration/tests/config.rs index 3f38db1..8f2f0d9 100644 --- a/tests-integration/tests/config.rs +++ b/tests-integration/tests/config.rs @@ -284,3 +284,40 @@ fn invalid_runtime_options_are_rejected_when_building_configuration() { assert_eq!(error.operation(), ErrorOperation::BuildConfig, "{case}"); } } + +#[test] +fn stats_storage_is_bounded_and_included_in_the_memory_floor() { + let options = RuntimeOptions { + append_shards: 2, + ..RuntimeOptions::default() + }; + let base = CacheConfig::new(test_storage(), options.clone()).unwrap(); + for shards in [0, 3, 65, usize::MAX] { + let mut invalid = options.clone(); + invalid.stats.shards = shards; + assert_eq!( + CacheConfig::new(test_storage(), invalid) + .unwrap_err() + .kind(), + ErrorKind::InvalidInput + ); + } + let mut enabled = options; + enabled.stats = cache2::StatsOptions { + request_counters: true, + l1_latency: cache2::LatencyMode::Full, + l2_latency: cache2::LatencyMode::Full, + mutation_latency: cache2::LatencyMode::Full, + io_latency: true, + shards: 64, + }; + let configured = CacheConfig::new(test_storage(), enabled.clone()).unwrap(); + assert!(configured.minimum_memory_bytes() > base.minimum_memory_bytes()); + enabled.managed_memory_limit_bytes = base.minimum_memory_bytes(); + assert_eq!( + CacheConfig::new(test_storage(), enabled) + .unwrap_err() + .kind(), + ErrorKind::InvalidInput + ); +}