diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 910c28b..db41900 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -80,7 +80,7 @@ This layout converts small foreground mutations into ordered, batched Region wri 5. Return after bounded in-memory admission. 6. The shard worker writes sealed batches and publishes their L2 mappings only after write completion. -Full staging or short-path contention returns structured `ErrorKind::Overloaded`. Admission is shard-local. Success means accepted staging; `put_l2` becomes visible when publication completes. `drain` fences all mutations accepted before its operation barrier and waits for their Region writes and L2 publication, but does not issue the recovery durability syncs. +Full staging or short-path contention returns structured `ErrorKind::Overloaded`. Staging admission is shard-local; optional adaptive fill budgets are shared across shards. Success means accepted staging; `put_l2` becomes visible when publication completes. `drain` fences all mutations accepted before its operation barrier and waits for their Region writes and L2 publication, but does not issue the recovery durability syncs. ### `get` @@ -114,6 +114,14 @@ Reads and writes use independent bounded engine pools. Reclaim has separate read Each lane uses one concrete `IoEngine` for admission, submission, cancellation, statistics, and shutdown. Driver-specific constructors start POSIX workers or an io_uring driver behind the same bounded command and completion protocol. Callers share the engine through `Arc`; its final owner joins the workers. Submitted requests retain their buffers and capacity until actual completion, independently of the caller's wait deadline. +### Pre-timeout fill pressure + +`io::fill_control` owns optional background observations and adaptive fill admission. `Disabled` retains the existing path. Enabled modes preallocate a fixed observation table sized by append and reclaim worker counts. Background workers register before engine admission and retain their observations through completion validation and publication. A full table skips that request and counts `dropped_observations` instead of failing I/O. Workers checkpoint at one quarter of the I/O deadline, capped at 500 ms, and pause new fills without entering timeout recovery. Aggregate validated throughput provides an approximate drain estimate for snapshots and does not drive pause decisions. Pressure is `Healthy` or `Paused` and does not change terminal health or prove a device fault. + +Adaptive foreground `put`/`put_l2` load pause-holder state and otherwise only compete for staging. Byte and record ceilings are instance-wide, shared by every shard worker; they pace non-essential background flush. Urgent, drain, and rotation flushes always proceed. A flush that cannot take a span refunds the consumed budget and retries. `Observe` does not delay flush or reject fills and counts pause refusals as `would_reject`; it still checkpoints waits so pause is observable. Adaptive rejects new fills immediately while paused and skips optional reinsertion. Pause is released when the slow I/O completes, not after later publication. Reads, deletes, and essential reclaim bypass fill budgets. Close stops admission independently of outstanding I/O. + +See [adaptive fill admission](CONFIGURATION.md#adaptive-fill-admission) for rate ceilings, pause conditions, bounded bursts, and tuning limits. + ### Memory The managed-memory limit covers the index mapping, heat bits, L1, append buffers, reclaim buffers, metadata, cache-owned thread stacks, recovery scratch, and transient reads. Total deployment memory additionally includes allocator metadata, Tokio, process overhead, and the kernel page cache. `CacheConfig::new` rejects invalid or insufficient memory budgets before file access; actual allocation can still fail during open. diff --git a/BENCHMARK.md b/BENCHMARK.md index 924cdc9..1399504 100644 --- a/BENCHMARK.md +++ b/BENCHMARK.md @@ -114,6 +114,8 @@ Repeat sizes in `CACHE_SOAK_VALUE_BYTES` to weight a production distribution. Us | io_uring | `_IO_URING__SQPOLL_MS` | absent (disabled) | | io_uring | `_IO_URING__SQPOLL_CPU` | absent (unpinned) | +Optional fill admission uses `_FILL_CONTROL=disabled|observe|adaptive` under the same three prefixes. Enabled modes require explicit `_FILL_BYTES_PER_SECOND` and `_FILL_OPERATIONS_PER_SECOND` ceilings. The harness prints the effective options and adds a `type=fill_control` report beside cache records. Compare the modes in alternating order with identical ceilings and traffic; record rejected and hypothetical fills alongside completed throughput. High ceilings help measure instrumentation overhead, while lower ceilings and injected storage stalls exercise admission behavior. Buffered macOS results do not qualify Linux NVMe or cgroup throttling. + Select the backend with `_IO_ENGINE=posix|io-uring`. POSIX workers bound concurrent operations. io_uring ring count and aggregate in-flight limit are independent; changing one does not rewrite the other. Ring count must not exceed the in-flight limit. IOPOLL requires `_IO_MODE=direct`; SQPOLL CPU requires an idle timeout. Only the selected backend's settings are read. Each harness prints the resulting `IoEngineOptions`, and buffer estimates use its actual concurrency. The request benchmark's default read-wait capacity follows the selected read pool's in-flight limit. `CACHE_BENCH_STATS` is now `CACHE_BENCH_ACTIVITY_COUNTERS`. Machine-readable reports use `version=2`, renaming the cache record field `statistics_enabled` to `activity_counters_enabled`; the counter population is unchanged. diff --git a/CHANGELOG.md b/CHANGELOG.md index 90c54a1..42a7966 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ - All optional statistics now use `RuntimeOptions::stats`: replace `runtime.statistics` with `runtime.stats.activity_counters` and `CacheSnapshot::statistics_enabled` with `activity_counters_enabled`. Activity counters, terminal request counters, and latency collection remain independent and disabled by default; counter semantics and the on-disk format are unchanged. See [the configuration migration guide](CONFIGURATION.md#migrating-from-05). - The benchmarking `RegionIndexTurnoverReport` field `config` is renamed to `options`, matching the `RegionIndexTurnoverOptions` input it carries. +### Features + +- Optional `RuntimeOptions::fill_control` observes background request age and estimated drain time before timeout. `Observe` reports pause refusals as `would_reject`. `Adaptive` pauses new `put`/`put_l2` fills when a background worker checkpoints an old I/O (one quarter of the deadline, capped at 500 ms) or timeout recovery is active, releases checkpoint pause when that I/O completes, paces non-essential background flush by encoded bytes and record count, and suppresses optional reinsertion while paused. Foreground admission loads pause-holder state and otherwise only competes for staging. A full observation table skips that request instead of failing I/O. Reads, deletes, accepted writes, and essential reclaim retain their paths. `CacheSnapshot::fill_control` exposes pressure and accounting independently of statistics. Controller storage is included in managed memory. Disabled by default; see [configuration](CONFIGURATION.md#adaptive-fill-admission) for ceilings and measurement limits. + ### Improvements - Background write and reclaim timeouts now enter a reversible `CacheHealth::Recovering` state: new cache fills return overload while existing reads and deletes remain available. Original requests keep their bounded buffers and Regions and are never resubmitted; all affected work must complete validation and publication before fills resume. `RuntimeOptions::io_recovery_timeout` defaults to `None` for recovery until completion or close; use `Some(duration)` to bound recovery or `Some(Duration::ZERO)` for immediate cancellation. Recovery checks at fixed one-second intervals. Close interrupts recovery and preserves the existing unfenced-write safeguards; drain may wait indefinitely. Actual I/O errors and invalid completions remain terminal. diff --git a/CONFIGURATION.md b/CONFIGURATION.md index bd44d76..47d203f 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -104,6 +104,7 @@ The managed-memory limit is not an RSS limit. Allocator metadata, Tokio, the app | Enable or lengthen read waiting | Trades immediate misses for bounded wait | Adds the configured wait capacity, can raise p99, and returns overload when queue, memory, or deadline is exhausted | | Increase write execution capacity | Adds write I/O concurrency | Adds queue pressure and device contention; POSIX also adds worker stacks and neither backend adds foreground staging space | | Increase reclaim concurrency | Recycles more Regions concurrently | Adds one Region buffer and logical worker per request slot and can compete for device and append-staging capacity | +| Enable Adaptive fill control | Pauses new fills; paces flush | Ceilings are instance-wide; they do not change reads, deletes, or `CacheHealth` on pre-timeout pause | | Increase L1 capacity | Retains more reusable values | Reduces L2 demand but consumes retained bytes and fixed metadata from the same managed-memory limit | | Increase L1 shards | Reduces shard-local contention | Adds fixed metadata and controls; too many small shards reduce useful capacity efficiency | | Lower write flush threshold | Requests earlier partial publication | Increases write operation count and reduces batching without increasing staging capacity | @@ -379,6 +380,9 @@ Use `Cache::snapshot()` for regular telemetry and `Cache::detailed_snapshot()` f | `reinsert_budget_skipped` rises | Hot live bytes exceed the fixed reclaim allowance | Treat retention as best effort; change capacity/workload geometry rather than worker count | | High `l1_bypasses` with useful candidates | L1 contention, slots, or byte pressure | Inspect L1 occupancy, retained bytes, shards, capacity, and oversized values | | Managed-memory peak approaches the limit | Fixed or transient memory pressure | Rebalance index, L1, staging, I/O topology, and read headroom | +| `fill_control.pressure` is `Paused` while `health` is `Running` | Pre-timeout slow background I/O | Expected Adaptive load-shed; correlate with device queueing before treating it as a fault | +| Pause refusals with little device wait | Checkpoint fired without device saturation | Confirm I/O timeouts and host scheduling; `FillLimits` do not drive pause | +| `write_rejections` while fill pressure is `Healthy` | Staging or paced-flush pressure | Distinguish from pause; change ceilings only after measuring sustainable fill | | Unexpected cold start after configuration change | Static identity or append-shard rebind failed | Verify Region/index geometry and available Free Regions; cache loss is safe | Always correlate cache counters with device latency, physical IOPS, filesystem behavior, process RSS, and authoritative-backend load. Cache throughput alone can reward configurations that merely turn work into fast misses. @@ -393,8 +397,9 @@ Always correlate cache counters with device latency, physical IOPS, filesystem b 6. Tune read execution and wait semantics using hit rate, overload, and p99. 7. Tune append shards, write execution, and flush threshold using acceptance, publication, and device counters. 8. Increase reclaim concurrency only if one reclaimer cannot maintain Free Regions. -9. Re-run after selecting Buffered versus Direct or POSIX versus io_uring; backend topology values are not interchangeable. -10. Validate cold and warm opens, `drain`, overload behavior, managed-memory peak, and final value correctness before deployment. +9. Enable fill control after write and reclaim topology is stable: `Observe` with chosen instance-wide ceilings, then `Adaptive` with the same values. +10. Re-run after selecting Buffered versus Direct or POSIX versus io_uring; backend topology values are not interchangeable. +11. Validate cold and warm opens, `drain`, overload behavior, managed-memory peak, and final value correctness before deployment. Change one resource family at a time and alternate baseline and candidate runs on the same host. Use multiple fresh-cache samples for turnover and burst tests. For storage qualification, use a dataset larger than host RAM and follow [Validation](BENCHMARK.md). @@ -414,6 +419,8 @@ Change one resource family at a time and alternate baseline and candidate runs o | Maximum waiting reads | 1 through 65,536 in `Wait`; defaults to the read execution capacity | | L1 shards | 1 through 65,536 | | Write flush threshold | 4 KiB multiple from 4 KiB through 4 MiB | +| Fill bytes per second | 640 through 1 TiB/s when fill control is enabled; instance-wide | +| Fill records per second | 10 through 4,294,967,295 when fill control is enabled; instance-wide | | Managed-memory limit | Nonzero, at least L1 capacity, and large enough for the validated fixed footprint | Storage and configuration construction enforce these bounds before file access. Open checks the selected filesystem, device, and runtime environment. @@ -425,3 +432,28 @@ Set `options.reclaim_io_timeout = Duration::from_secs(30)` after constructing `R ### Background I/O recovery `options.io_recovery_timeout = None` is the default: background admission/completion timeouts enter `CacheHealth::Recovering` and wait until the original operation completes or close interrupts recovery. Set `Some(Duration::from_secs(300))` for a five-minute additional budget, or `Some(Duration::ZERO)` for immediate cancellation. Finite combined deadlines must be representable. Recovery uses fixed one-second checks, without an exponential backoff or repeated submission of issued I/O. Admission retries are allowed only when the engine returns the unsubmitted operation. New `put` and `put_l2` calls return overload while any background operation is recovering; reads and deletes keep their normal behavior. Already accepted work retains its bounded resources. Fills resume only after every affected operation passes validation and publication. Drain can wait indefinitely in the default mode; close interrupts recovery within a polling interval, then applies bounded cancellation. Normal close work still has its ordinary deadlines. Actual I/O errors, invalid completions, and finite-budget exhaustion retain terminal failure safeguards. This is not an automatic reopen of a failed instance. + +### Adaptive fill admission + +`RuntimeOptions::fill_control` defaults to `FillControlOptions::Disabled`. `Observe` records pause pressure and hypothetical rejections while retaining ordinary admission and flush cadence. `Adaptive` rejects new `put` and `put_l2` fills with `ErrorKind::Overloaded` when paused, and paces non-essential background flush by the configured ceilings. Neither mode changes foreground reads, deletes, accepted flushes, essential reclaim, or the existing timeout recovery policy. Adaptive also suppresses optional hot-record reinsertion while paused. Pressure is separate from `CacheHealth`: a pre-timeout pause leaves a healthy cache `Running`. + +```rust +use cache2::{FillControlOptions, FillLimits, RuntimeOptions}; + +let mut options = RuntimeOptions::default(); +// Instance-wide mixed-write starting point; not a device sequential-write spec. +let limits = FillLimits::new(256 * 1024 * 1024, 20_000); +options.fill_control = FillControlOptions::Observe(limits); +// After evaluating pressure and would_reject, enforce the same policy: +options.fill_control = FillControlOptions::Adaptive(limits); +``` + +The two ceilings are instance-wide: every append shard shares one token bucket. They constrain logical encoded fill bytes (charged in 64-byte units) and fill records from the worker's staging snapshot, not device bandwidth or physical IOPS. `put` and `put_l2` do not consume the bucket; Adaptive rejects those calls only while paused. The sealed span may grow by later encodes before the lock. Valid ceilings are 640 bytes/s through 1 TiB/s and 10 through 4,294,967,295 records/s. Adaptive refill adds elapsed credit up to 100 ms of the configured ceilings and at least one maximum-size record; intervals too short to mint a unit or record do not move the refill clock. An idle controller never accumulates more than that bounded burst. A flush larger than that burst still proceeds on remaining credit, and a flush that cannot take a span refunds only what it consumed. Urgent, drain, and rotation flushes always proceed; other flushes wait for budget and back-pressure staging into ordinary write overload. Observe does not delay flush and counts `would_reject` only for pause. + +Background workers checkpoint outstanding I/O at one quarter of the normal deadline, capped at 500 ms, and pause new Adaptive fills without entering timeout recovery. Checkpoint pause follows live holder counts and lifts when that I/O completes, even if later validation or reclaim scanning is still running. Observe uses the same checkpoints so it can report `would_reject`. Estimated drain time is reported for snapshots from validated throughput since open and does not pause admission. These are conservative pressure signals: very short deadlines or a descheduled worker can still reach timeout first. Idle time alone never looks like a stall. Pause and the rate bucket are independent: a busy device is shed by checkpoint or timeout recovery, not by copying sequential-write specifications into `FillLimits`. Normal timeout recovery keeps its admission fence until all affected work is validated and published. + +Start with `Observe` and the same ceilings you plan to enforce. Inspect `CacheSnapshot::fill_control` (`pressure`, `would_reject`, `oldest_operation_ns`) together with device queueing, then switch to `Adaptive` without changing the numbers. The example `256 MiB/s` and `20_000` records/s is a conservative mixed-write starting point, not an NVMe datasheet copy. Both dimensions must hold: if typical encoded records are about 1 KiB, a 256 MiB/s byte ceiling needs on the order of 200_000–500_000 records/s or the record limit fires first. To pause ingest on stalls while leaving healthy flush uncapped, raise both ceilings until they no longer clip measured fill (for example 1–2 GiB/s with a matching record rate). To keep headroom for reads and reclaim on a shared NVMe, set the byte ceiling around one third to one half of measured sustainable cache write; four write workers and a 4 MiB flush threshold often land in the 256–512 MiB/s range. Do not use the legal minima (640 bytes/s, 10 records/s) as NVMe defaults. Buffered macOS results do not qualify Linux NVMe; use [Validation](BENCHMARK.md) on the target device. + +`CacheSnapshot::fill_control` reports pressure, enforcement, configured or paused rates, rejections, hypothetical pause refusals, dropped observations, outstanding bytes/operations, oldest age, and estimated drain time independently of `RuntimeOptions::stats`. Outstanding bytes exclude unflushed staging. A zero drain estimate means no estimate is available when work is pending. A full observation table skips that request rather than failing I/O. These measurements cannot distinguish device throttling from CPU scheduling, kernel queueing, or slow validation; correlate them with host metrics before diagnosing hardware. `cache_fill_pressure_changed` logs state transitions under `cache2::health`. + +Enabled modes preallocate one observation per append/reclaim worker, included in `CacheConfig::minimum_memory_bytes()`. Background observations use atomics over this bounded table. Foreground fill admission loads pause-holder and recovery counters without allocation, clocks, or locks from the controller. Reinsertion suppression reads atomics. Close stops admission independently of outstanding I/O. Disabled mode creates no controller or observations. diff --git a/README.md b/README.md index aa03ad3..736cd63 100644 --- a/README.md +++ b/README.md @@ -83,6 +83,7 @@ See the [configuration guide](CONFIGURATION.md#configuration-lifecycle) for exam | Memory | `managed_memory_limit_bytes` | 1 GiB across cache-managed allocations. | | I/O mode | `io_mode` | Buffered I/O. | | Metrics | `stats: StatsOptions` | Health/resource gauges always available; activity, request, and latency collection opt in. | +| Fill control | `fill_control: FillControlOptions` | Disabled. `Observe` reports pause; `Adaptive` rejects new fills. Ceilings are cache-wide. | Changing the append-shard count rebinds recovered Active Regions during a warm open. Growth uses available Free Regions; when there are not enough, the disposable cache safely starts empty. @@ -106,7 +107,7 @@ The on-disk format is versioned. During 0.x, deployments should expect cold star ### Metrics -`Cache::snapshot()` provides lock-free health and resource gauges. Setting `RuntimeOptions::stats.activity_counters` to `true` adds cumulative cache and I/O counters. `Cache::detailed_snapshot()` samples L1, index, write-buffer pressure, and Region metadata for periodic diagnostics. +`Cache::snapshot()` provides health and resource gauges using atomics. Setting `RuntimeOptions::stats.activity_counters` to `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. @@ -131,6 +132,12 @@ RUST_LOG=cache2=info cargo run --package examples --example logforth -- /tmp/cac `cache_opened` reports the index backing, mapping extent, validation mode, and whether warm mutations use copy-on-write. `cache_recovery_cold` records why a clean image was rejected or why private mapping fell back to a cold start. `cache_miss_only` records the first terminal index-validation or I/O failure. +### Reclaim read deadline + +Set `RuntimeOptions::reclaim_io_timeout` to change the normal background reclaim deadline, for example `Duration::from_secs(30)` (default five seconds). Background write and reclaim timeouts enter `CacheHealth::Recovering`: new fills return overload while reads and deletes remain available. `RuntimeOptions::io_recovery_timeout` defaults to `None`, allowing recovery until completion or close. Use `Some(Duration::from_secs(300))` to limit the additional wait, or `Some(Duration::ZERO)` for immediate cancellation. Original requests retain their resources and are never resubmitted; fills resume after validation and publication of all affected work. Close interrupts recovery; drain may wait indefinitely. Actual I/O errors and invalid completions still fail the cache. + +Optional [adaptive fill admission](CONFIGURATION.md#adaptive-fill-admission) detects slow background progress before timeout. `Adaptive` pauses new fills immediately and paces non-essential flush with instance-wide byte and record ceilings. Start with `FillControlOptions::Observe` to inspect pause pressure and hypothetical rejections, then use `Adaptive` to enforce the same limits. It is disabled by default; reads, deletes, accepted writes, and essential reclaim retain their existing paths. + ## Development C² requires Rust 1.98.0. @@ -155,6 +162,3 @@ The root workspace keeps the publishable crate, integration tests, benchmarks, e Licensed under the [Apache License, Version 2.0](LICENSE). -### Reclaim read deadline - -Set `RuntimeOptions::reclaim_io_timeout` to change the normal background reclaim deadline, for example `Duration::from_secs(30)` (default five seconds). Background write and reclaim timeouts enter `CacheHealth::Recovering`: new fills return overload while reads and deletes remain available. `RuntimeOptions::io_recovery_timeout` defaults to `None`, allowing recovery until completion or close. Use `Some(Duration::from_secs(300))` to limit the additional wait, or `Some(Duration::ZERO)` for immediate cancellation. Original requests retain their resources and are never resubmitted; fills resume after validation and publication of all affected work. Close interrupts recovery; drain may wait indefinitely. Actual I/O errors and invalid completions still fail the cache. diff --git a/benchmarks/cache/main.rs b/benchmarks/cache/main.rs index 34b90fb..852e9c8 100644 --- a/benchmarks/cache/main.rs +++ b/benchmarks/cache/main.rs @@ -80,6 +80,7 @@ struct BenchConfig { write_clients: usize, clients: usize, io_engine: IoEngineOptions, + fill_control: cache2::FillControlOptions, io_mode: IoMode, l1_eviction_policy: L1EvictionPolicy, stats: cache2::StatsOptions, @@ -229,6 +230,7 @@ impl BenchConfig { write_clients, clients, io_engine, + fill_control: benchmarks::config::fill_control_from_env("CACHE_BENCH")?, io_mode, l1_eviction_policy, stats, @@ -246,6 +248,8 @@ impl BenchConfig { fn runtime_options(&self) -> RuntimeOptions { let mut options = RuntimeOptions::default(); options.io_engine = self.io_engine; + options.fill_control = self.fill_control; + println!("fill_control={:?}", self.fill_control); options.io_mode = self.io_mode; options.append_shards = self.append_shards; options.l1_capacity_bytes = self.l1_capacity_bytes; diff --git a/benchmarks/cache_soak/main.rs b/benchmarks/cache_soak/main.rs index 0d6db57..48e9ff8 100644 --- a/benchmarks/cache_soak/main.rs +++ b/benchmarks/cache_soak/main.rs @@ -89,6 +89,7 @@ struct SoakConfig { require_path_coverage: bool, require_reinsert_coverage: bool, io_engine: IoEngineOptions, + fill_control: cache2::FillControlOptions, io_mode: IoMode, l1_eviction_policy: L1EvictionPolicy, directory: PathBuf, @@ -193,6 +194,7 @@ impl SoakConfig { require_path_coverage, require_reinsert_coverage, io_engine, + fill_control: benchmarks::config::fill_control_from_env("CACHE_SOAK")?, io_mode, l1_eviction_policy, directory, @@ -209,6 +211,8 @@ impl SoakConfig { fn runtime_options(&self) -> RuntimeOptions { let mut options = RuntimeOptions::default(); options.io_engine = self.io_engine; + options.fill_control = self.fill_control; + println!("fill_control={:?}", self.fill_control); options.io_mode = self.io_mode; options.append_shards = self.append_shards; options.l1_capacity_bytes = self.l1_capacity_bytes; diff --git a/benchmarks/mixed_workloads/main.rs b/benchmarks/mixed_workloads/main.rs index ffd6e43..8b3c437 100644 --- a/benchmarks/mixed_workloads/main.rs +++ b/benchmarks/mixed_workloads/main.rs @@ -217,6 +217,7 @@ struct HarnessOptions { latency_sample_interval: usize, seed: u64, io_engine: IoEngineOptions, + fill_control: cache2::FillControlOptions, io_mode: IoMode, l1_eviction_policy: L1EvictionPolicy, directory: PathBuf, @@ -268,6 +269,7 @@ impl HarnessOptions { latency_sample_interval, seed, io_engine, + fill_control: benchmarks::config::fill_control_from_env("CACHE_WORKLOAD")?, io_mode, l1_eviction_policy, directory, @@ -346,6 +348,7 @@ impl HarnessOptions { latency_sample_interval: self.latency_sample_interval, seed: self.seed, io_engine: self.io_engine, + fill_control: self.fill_control, io_mode: self.io_mode, l1_eviction_policy: self.l1_eviction_policy, directory: self.directory.clone(), @@ -367,6 +370,7 @@ struct ScenarioConfig { latency_sample_interval: usize, seed: u64, io_engine: IoEngineOptions, + fill_control: cache2::FillControlOptions, io_mode: IoMode, l1_eviction_policy: L1EvictionPolicy, directory: PathBuf, @@ -383,6 +387,8 @@ impl ScenarioConfig { fn runtime_options(&self) -> RuntimeOptions { let mut options = RuntimeOptions::default(); options.io_engine = self.io_engine; + options.fill_control = self.fill_control; + println!("fill_control={:?}", self.fill_control); options.io_mode = self.io_mode; options.append_shards = self.append_shards; options.l1_capacity_bytes = self.l1_capacity_bytes; diff --git a/benchmarks/src/config.rs b/benchmarks/src/config.rs index 6560070..67daa01 100644 --- a/benchmarks/src/config.rs +++ b/benchmarks/src/config.rs @@ -59,6 +59,25 @@ pub fn io_engine_from_env(prefix: &str) -> io::Result { } } +/// Reads optional fill-control mode and its explicit logical rate ceilings. +pub fn fill_control_from_env(prefix: &str) -> io::Result { + let name = format!("{prefix}_FILL_CONTROL"); + let mode = setting::(&name)?.unwrap_or_else(|| "disabled".into()); + if mode == "disabled" { + return Ok(cache2::FillControlOptions::Disabled); + } + let bytes = setting(&format!("{prefix}_FILL_BYTES_PER_SECOND"))? + .ok_or_else(|| invalid("enabled fill control requires FILL_BYTES_PER_SECOND"))?; + let operations = setting(&format!("{prefix}_FILL_OPERATIONS_PER_SECOND"))? + .ok_or_else(|| invalid("enabled fill control requires FILL_OPERATIONS_PER_SECOND"))?; + let options = cache2::FillLimits::new(bytes, operations); + match mode.as_str() { + "observe" => Ok(cache2::FillControlOptions::Observe(options)), + "adaptive" => Ok(cache2::FillControlOptions::Adaptive(options)), + _ => Err(invalid(format!("unsupported {name}: {mode}"))), + } +} + /// Maximum active reads across the selected backend's pool. pub fn read_max_in_flight(options: IoEngineOptions) -> usize { match options { diff --git a/benchmarks/src/report.rs b/benchmarks/src/report.rs index e96ac2c..3df5609 100644 --- a/benchmarks/src/report.rs +++ b/benchmarks/src/report.rs @@ -479,6 +479,24 @@ pub fn emit_cache_report( detailed.region.physical_record_count, format_bytes(cache.logical_disk_peak_bytes as f64), ); + let fill = cache.fill_control; + println!( + "report version=2 type=fill_control benchmark={} scenario={} phase={} pressure={:?} enforcing={} bytes_per_second={} records_per_second={} rejections={} would_reject={} dropped_observations={} outstanding_operations={} outstanding_bytes={} oldest_operation_ns={} estimated_drain_ns={}", + benchmark, + scenario, + phase, + fill.pressure, + fill.enforcing, + fill.bytes_per_second, + fill.records_per_second, + fill.rejections, + fill.would_reject, + fill.dropped_observations, + fill.outstanding_operations, + fill.outstanding_bytes, + fill.oldest_operation_ns, + fill.estimated_drain_ns, + ); println!( "report version=2 type=cache benchmark={} scenario={} phase={} health={:?} activity_counters_enabled={} puts={} deletes={} written_bytes={} served_bytes={} l1_hits={} l1_misses={} l2_hits={} l2_misses={} l2_read_memory_misses={} l2_read_busy_misses={} l2_read_overloads={} l2_read_wait_ns={} promotions={} l1_evictions={} l1_bypasses={} write_rejections={} io_failures={} rotations={} reclaimed_regions={} reclaim_bytes={} reclaim_records={} reinsert_records={} reinsert_bytes={} reinsert_skipped={} reinsert_budget_skipped={}", benchmark, diff --git a/cache2/ERRORS.md b/cache2/ERRORS.md index c18ab88..c7567fe 100644 --- a/cache2/ERRORS.md +++ b/cache2/ERRORS.md @@ -32,6 +32,8 @@ fn cache_value(cache: &Cache, key: &[u8], value: &[u8]) -> Result<(), Error> { Do not retry without a bound. C² deliberately exposes pressure instead of building unbounded queues. With the default immediate-read policy, read-pool or buffer pressure is `Ok(None)`. When read waiting is enabled, queue saturation, buffer pressure, and deadline expiry are `ErrorKind::Overloaded`. +With `FillControlOptions::Adaptive`, new fills also return `Overloaded` when the controller pauses. Staging back-pressure from paced flush uses the ordinary write-overload path. Inspect `CacheSnapshot::fill_control` to distinguish controller pause from other admission pressure. Pre-timeout pressure does not change `CacheHealth::Running` or identify a hardware failure. `Observe` counts `would_reject` for pause refusals only. Reads, deletes, accepted writes, and essential reclaim retain their ordinary paths. + ## Classifications | `ErrorKind` | Meaning | Usual response | diff --git a/cache2/src/cache.rs b/cache2/src/cache.rs index f138406..c91a124 100644 --- a/cache2/src/cache.rs +++ b/cache2/src/cache.rs @@ -396,7 +396,7 @@ impl Cache { public_result(ErrorOperation::Drain, self.data_plane.drain_async().await) } - /// Returns a lock-free operational snapshot. Activity and I/O counters are + /// Returns an operational snapshot using atomics. Activity and I/O counters are /// cumulative for this open and are populated only when /// `RuntimeOptions::stats.activity_counters` is enabled; health and resource gauges are /// always available. diff --git a/cache2/src/config/runtime.rs b/cache2/src/config/runtime.rs index 56a0d90..390ed74 100644 --- a/cache2/src/config/runtime.rs +++ b/cache2/src/config/runtime.rs @@ -353,6 +353,44 @@ pub enum ReadAdmission { }, } +/// Optional pre-timeout pressure observation and fill admission control. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub enum FillControlOptions { + /// No controller, timer, or additional request accounting. + #[default] + Disabled, + /// Report pressure and hypothetical rejections without changing admission. + Observe(FillLimits), + /// Reject new fills under pressure. Accepted writes and essential reclaim continue. + Adaptive(FillLimits), +} + +/// Logical fill-rate ceilings shared by [`FillControlOptions::Observe`] and +/// [`FillControlOptions::Adaptive`]. They are instance-wide, not per worker, +/// and are not device bandwidth or IOPS guarantees. Foreground `put` does not +/// consume them; Adaptive uses them to pace non-essential background flush. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct FillLimits { + /// Maximum encoded fill bytes per second, from 640 through 1 TiB/s. + /// Instance-wide across all shard workers. + pub max_bytes_per_second: u64, + /// Maximum fill records per second, from 10 through 4,294,967,295. + /// Instance-wide across all shard workers. + pub max_records_per_second: u32, +} + +impl FillLimits { + /// Creates unchecked rate ceilings. [`CacheConfig::new`] validates them. + pub const fn new(max_bytes_per_second: u64, max_records_per_second: u32) -> Self { + Self { + max_bytes_per_second, + max_records_per_second, + } + } +} + /// Process-local resource choices, checked together by [`CacheConfig::new`]. /// /// These values may change across opens. Warm recovery rebinds append shards @@ -387,6 +425,10 @@ pub struct RuntimeOptions { /// can wait indefinitely with `None`. Actual I/O errors and invalid /// completions still fail the instance. pub io_recovery_timeout: Option, + /// Optional pre-timeout fill pressure control. Disabled by default. + /// Enabled modes reserve bounded worker observations; Adaptive paces + /// background flush and pauses new fills when outstanding work is old. + pub fill_control: FillControlOptions, /// Hash-routed append paths, from 1 through 256 (default 4). Each needs one /// Active Region, two Region-sized buffers, and a worker. The layout also needs a /// spare Region. @@ -421,6 +463,7 @@ impl Default for RuntimeOptions { read_admission: ReadAdmission::Immediate, reclaim_io_timeout: Duration::from_secs(5), io_recovery_timeout: None, + fill_control: FillControlOptions::Disabled, append_shards: DEFAULT_APPEND_SHARDS, l1_capacity_bytes: DEFAULT_L1_CAPACITY_BYTES, l1_eviction_policy: L1EvictionPolicy::Clock, @@ -496,6 +539,11 @@ impl CacheConfig { let index_slots = storage.index_slots; runtime.resolve()?; let stats_bytes = Recorder::allocation_bytes(runtime.stats)?; + let fill_bytes = crate::io::fill_control::FillController::allocation_bytes( + runtime.fill_control, + runtime.append_shards as usize + + IoPoolTopology::reclaim(runtime.io_engine).max_in_flight(), + )?; 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", @@ -511,6 +559,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)) + .and_then(|bytes| bytes.checked_add(fill_bytes)) .ok_or_else(|| invalid_config("fixed memory requirements overflow"))?; let (reserved_memory_bytes, minimum_memory_bytes) = runtime.memory_requirements(geometry, fixed_bytes)?; @@ -802,6 +851,51 @@ mod tests { use crate::ErrorKind; use crate::StorageOptions; + #[test] + fn fill_control_validates_ceilings_and_accounts_observation_memory() { + let storage = StorageOptions::new(1024 * 1024 * 1024).build().unwrap(); + let base = CacheConfig::new(storage.clone(), RuntimeOptions::default()).unwrap(); + assert_eq!(base.runtime().fill_control, FillControlOptions::Disabled); + for (bytes, operations) in [(639, 100), ((1 << 40) + 1, 100), (640, 9)] { + let options = RuntimeOptions { + fill_control: FillControlOptions::Adaptive(FillLimits::new(bytes, operations)), + ..RuntimeOptions::default() + }; + assert_eq!( + CacheConfig::new(storage.clone(), options) + .unwrap_err() + .kind(), + ErrorKind::InvalidInput + ); + } + CacheConfig::new( + storage.clone(), + RuntimeOptions { + fill_control: FillControlOptions::Adaptive(FillLimits::new(640, u32::MAX)), + ..RuntimeOptions::default() + }, + ) + .unwrap(); + let mut minimum = None; + for mode in [FillControlOptions::Observe, FillControlOptions::Adaptive] { + let config = CacheConfig::new( + storage.clone(), + RuntimeOptions { + fill_control: mode(FillLimits::new(64_000, 100)), + ..RuntimeOptions::default() + }, + ) + .unwrap(); + let extra = config.minimum_memory_bytes() - base.minimum_memory_bytes(); + assert!(extra > 0); + assert!(extra < CACHE_THREAD_STACK_BYTES); + if let Some(previous) = minimum { + assert_eq!(extra, previous); + } + minimum = Some(extra); + } + } + #[test] fn recovery_timeout_allows_zero_and_rejects_overflow() { let storage = crate::StorageOptions::new(1024 * 1024 * 1024) diff --git a/cache2/src/io/engine/mod.rs b/cache2/src/io/engine/mod.rs index 6eff683..f0570cb 100644 --- a/cache2/src/io/engine/mod.rs +++ b/cache2/src/io/engine/mod.rs @@ -777,13 +777,38 @@ impl BoundedIoRequest { ) -> Result { let original = self.deadline; loop { - self.request = match self.request.wait_until(self.deadline) { - Ok(completion) => return Ok(completion), + let cap = if Instant::now() < original { + recovery.wait_cap(original) + } else { + self.deadline + }; + self.request = match self.request.wait_until(cap) { + Ok(completion) => { + recovery.clear_slow(); + return Ok(completion); + } Err(request) => request, }; + if Instant::now() < original { + recovery.note_slow(); + continue; + } match recovery.next_deadline(original) { Some(deadline) => self.deadline = deadline, - None => return self.wait(engine), + None => { + return match self.wait(engine) { + Ok(completion) => { + recovery.clear_slow(); + Ok(completion) + } + Err(exceeded) => { + if exceeded.completion.is_some() { + recovery.clear_slow(); + } + Err(exceeded) + } + }; + } } } } @@ -958,10 +983,14 @@ pub fn submit_background_io( timeout: Duration, recovery: &mut RecoveryAttempt<'_>, ) -> Result { + let bytes = match &operation { + IoOperation::Read { buffer, .. } | IoOperation::Write { buffer, .. } => buffer.len() as u64, + }; + recovery.start(bytes, timeout); let original = Instant::now() .checked_add(timeout) .unwrap_or_else(Instant::now); - let mut deadline = original; + let mut deadline = recovery.wait_cap(original); loop { match submit_cache_io_until(engine, operation, deadline, CACHE_IO_CANCEL_GRACE) { Ok(mut request) => { @@ -969,6 +998,12 @@ pub fn submit_background_io( return Ok(request); } Err(error) if error.error.kind() == io::ErrorKind::TimedOut => { + if Instant::now() < original { + recovery.note_slow(); + deadline = original; + operation = error.operation; + continue; + } let Some(next) = recovery.next_deadline(original) else { return Err(error); }; diff --git a/cache2/src/io/engine/recovery.rs b/cache2/src/io/engine/recovery.rs index 23ea450..f505b8f 100644 --- a/cache2/src/io/engine/recovery.rs +++ b/cache2/src/io/engine/recovery.rs @@ -14,17 +14,22 @@ //! Reversible background timeout recovery, separate from the health latch. +use std::sync::Arc; use std::sync::atomic::AtomicBool; use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::time::Duration; use std::time::Instant; +use crate::io::fill_control::FillController; +use crate::io::fill_control::Observation; + /// Shared across background workers. Resource ownership remains with each worker. pub struct BackgroundRecovery { timeout: Option, pending: AtomicUsize, stopped: AtomicBool, + pub fill: Option>, } impl BackgroundRecovery { @@ -33,6 +38,14 @@ impl BackgroundRecovery { timeout, pending: AtomicUsize::new(0), stopped: AtomicBool::new(false), + fill: None, + } + } + + pub fn with_fill(timeout: Option, fill: Option>) -> Self { + Self { + fill, + ..Self::new(timeout) } } @@ -40,6 +53,10 @@ impl BackgroundRecovery { RecoveryAttempt { recovery: self, entered: false, + slow: false, + started: None, + io_timeout: Duration::ZERO, + observation: None, } } @@ -49,6 +66,9 @@ impl BackgroundRecovery { pub fn stop(&self) { self.stopped.store(true, Ordering::Release); + if let Some(fill) = &self.fill { + fill.stop(); + } } } @@ -57,9 +77,52 @@ impl BackgroundRecovery { pub struct RecoveryAttempt<'a> { recovery: &'a BackgroundRecovery, entered: bool, + slow: bool, + started: Option, + io_timeout: Duration, + observation: Option>, } impl RecoveryAttempt<'_> { + pub fn start(&mut self, bytes: u64, timeout: Duration) { + self.started = Some(Instant::now()); + self.io_timeout = timeout; + if let Some(fill) = &self.recovery.fill { + self.observation = fill.observe(bytes, timeout); + } + } + + /// First wait bound: a fill checkpoint before the normal I/O timeout. + pub fn wait_cap(&self, original: Instant) -> Instant { + if self.recovery.fill.is_none() || self.slow { + return original; + } + let Some(start) = self.started else { + return original; + }; + FillController::checkpoint(start, self.io_timeout).min(original) + } + + pub fn note_slow(&mut self) { + if self.slow { + return; + } + self.slow = true; + if let Some(observation) = &mut self.observation { + observation.note_slow(); + } + } + + pub fn clear_slow(&mut self) { + if !self.slow { + return; + } + self.slow = false; + if let Some(observation) = &mut self.observation { + observation.clear_slow(); + } + } + /// Poll completion/admission at fixed one-second intervals, without extending /// a configured total budget. Shutdown also terminates unlimited recovery. pub fn next_deadline(&mut self, original: Instant) -> Option { @@ -80,6 +143,9 @@ impl RecoveryAttempt<'_> { }; if !self.entered { self.entered = true; + if let Some(fill) = &self.recovery.fill { + fill.set_recovering(true); + } if self.recovery.pending.fetch_add(1, Ordering::AcqRel) == 0 { log::warn!(target: "cache2::health", event = "cache_io_recovery_started"; "background I/O timed out; pausing cache fills while retaining owned requests"); @@ -89,10 +155,19 @@ impl RecoveryAttempt<'_> { } /// Called only after operation-result validation and publication succeed. - pub fn finish(self) { - if self.entered && self.recovery.pending.fetch_sub(1, Ordering::AcqRel) == 1 { - log::info!(target: "cache2::health", event = "cache_io_recovery_completed"; + pub fn finish(mut self) { + if let Some(observation) = self.observation.take() { + observation.finish(); + } + if self.entered { + let last = self.recovery.pending.fetch_sub(1, Ordering::AcqRel) == 1; + if let Some(fill) = &self.recovery.fill { + fill.set_recovering(false); + } + if last { + log::info!(target: "cache2::health", event = "cache_io_recovery_completed"; "all timed-out background operations recovered and passed validation"); + } } } } diff --git a/cache2/src/io/engine/tests.rs b/cache2/src/io/engine/tests.rs index 69dc461..81b5c9f 100644 --- a/cache2/src/io/engine/tests.rs +++ b/cache2/src/io/engine/tests.rs @@ -1391,3 +1391,71 @@ fn shutdown_interrupts_unlimited_recovery_without_releasing_pending_write() { }); engine.shutdown().unwrap(); } + +#[test] +fn adaptive_pressure_pauses_before_real_engine_timeout_and_resumes_after_io_completes() { + use crate::FillControlOptions; + use crate::FillLimits; + use crate::FillPressure; + use crate::io::fill_control::FillController; + for enforcing in [false, true] { + let settings = FillLimits::new(1024 * 1024, 1000); + let mode = if enforcing { + FillControlOptions::Adaptive(settings) + } else { + FillControlOptions::Observe(settings) + }; + let control = FillController::new(mode, 1, 4096).unwrap().unwrap(); + let recovery = BackgroundRecovery::with_fill(None, Some(Arc::clone(&control))); + let backend = Arc::new(BlockingBackend::default()); + let engine = IoEngine::for_test(backend.clone(), 1).unwrap(); + let memory = managed_memory(); + std::thread::scope(|scope| { + let (returned_tx, returned_rx) = mpsc::channel(); + let (validate_tx, validate_rx) = mpsc::channel(); + let recovery = &recovery; + let engine = &engine; + let memory = &memory; + scope.spawn(move || { + let mut attempt = recovery.attempt(); + let request = submit_background_io( + engine, + IoOperation::write(WritePoint::Record, write_buffer(memory, &[5; 4096]), 0), + Duration::from_secs(4), + &mut attempt, + ) + .unwrap(); + let completion = request.wait_with_recovery(engine, &mut attempt).unwrap(); + returned_tx.send(completion).unwrap(); + validate_rx.recv().unwrap(); + attempt.finish(); + }); + assert!(backend.wait_for_entered(1)); + let deadline = Instant::now() + Duration::from_secs(2); + while control.snapshot().pressure != FillPressure::Paused && Instant::now() < deadline { + std::thread::sleep(Duration::from_millis(10)); + } + let paused = control.snapshot(); + let recovered = recovery.is_recovering(); + let admission = control.try_admit(); + backend.release(); + let completion = returned_rx.recv_timeout(Duration::from_secs(2)).unwrap(); + assert_eq!(paused.pressure, FillPressure::Paused); + assert!(!recovered, "must react before the normal I/O timeout"); + assert_eq!(admission, !enforcing); + assert_eq!(paused.outstanding_bytes, 4096); + assert!(completion.into_io_result().0.is_ok()); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + validate_tx.send(()).unwrap(); + }); + let deadline = Instant::now() + Duration::from_secs(2); + while control.snapshot().pressure == FillPressure::Paused && Instant::now() < deadline { + std::thread::sleep(Duration::from_millis(10)); + } + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + assert_eq!(lock_unpoisoned(&backend.state).entered, 1); + engine.shutdown().unwrap(); + } +} diff --git a/cache2/src/io/fill_control.rs b/cache2/src/io/fill_control.rs new file mode 100644 index 0000000..dde7541 --- /dev/null +++ b/cache2/src/io/fill_control.rs @@ -0,0 +1,832 @@ +// 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. + +//! Bounded pre-timeout observation and worker-paced fill admission. + +use std::io; +use std::sync::Arc; +use std::sync::atomic::AtomicBool; +use std::sync::atomic::AtomicU8; +use std::sync::atomic::AtomicU64; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; +use std::time::Duration; +use std::time::Instant; + +use crate::FillControlOptions; +use crate::FillControlSnapshot; +use crate::FillLimits; +use crate::FillPressure; + +const UNIT: u64 = 64; +const MAX_CAS_ATTEMPTS: usize = 4; +const MIN_BYTES_PER_SECOND: u64 = 640; +const MAX_BYTES_PER_SECOND: u64 = 1 << 40; +const MIN_RECORDS_PER_SECOND: u32 = 10; +const MAX_RECORDS_PER_SECOND: u32 = u32::MAX; +const STALL_CAP: Duration = Duration::from_millis(500); +const TICKS_PER_SECOND: u64 = 10; + +/// Wake shard workers this often when a non-essential flush is waiting for budget. +pub const FLUSH_RETRY: Duration = Duration::from_millis(100); + +/// Credit taken by one non-essential flush. Zero for Observe and essential work. +#[derive(Clone, Copy)] +pub struct FlushCharge { + pub(crate) ops: u64, + pub(crate) units: u64, +} + +fn packed_ops(value: u64) -> u64 { + value >> 32 +} +fn packed_units(value: u64) -> u64 { + value & u64::from(u32::MAX) +} +fn pack_credit(ops: u64, units: u64) -> u64 { + (ops << 32) | units +} + +fn encode_pressure(pressure: FillPressure) -> u8 { + match pressure { + FillPressure::Disabled => 0, + FillPressure::Healthy => 1, + FillPressure::Paused => 2, + } +} + +fn decode_pressure(value: u8) -> FillPressure { + match value { + 1 => FillPressure::Healthy, + 2 => FillPressure::Paused, + _ => FillPressure::Disabled, + } +} + +fn nanos(duration: Duration) -> u64 { + duration.as_nanos().min(u128::from(u64::MAX)) as u64 +} + +struct Slot { + start_ns: AtomicU64, + bytes: AtomicU64, +} + +pub struct FillController { + options: FillLimits, + enforcing: bool, + max_units: u64, + max_ops: u64, + origin: Instant, + credit: AtomicU64, + last_refill_ns: AtomicU64, + pressure: AtomicU8, + stopped: AtomicBool, + recovering: AtomicUsize, + pause_holders: AtomicUsize, + rejections: AtomicU64, + would_reject: AtomicU64, + dropped_observations: AtomicU64, + completed_bytes: AtomicU64, + completed_ops: AtomicU64, + slots: Box<[Slot]>, +} + +impl FillController { + pub fn allocation_bytes(options: FillControlOptions, slots: usize) -> io::Result { + let settings = match options { + FillControlOptions::Disabled => return Ok(0), + FillControlOptions::Observe(settings) | FillControlOptions::Adaptive(settings) => { + settings + } + }; + if !(MIN_BYTES_PER_SECOND..=MAX_BYTES_PER_SECOND).contains(&settings.max_bytes_per_second) + || !(MIN_RECORDS_PER_SECOND..=MAX_RECORDS_PER_SECOND) + .contains(&settings.max_records_per_second) + { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "fill rate ceilings are out of range", + )); + } + slots + .checked_mul(size_of::()) + .and_then(|bytes| bytes.checked_add(size_of::() + 256)) + .ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidInput, + "fill controller memory overflow", + ) + }) + } + + pub fn new( + options: FillControlOptions, + slots: usize, + max_record: u64, + ) -> io::Result>> { + Self::allocation_bytes(options, slots)?; + let (settings, enforcing) = match options { + FillControlOptions::Disabled => return Ok(None), + FillControlOptions::Observe(settings) => (settings, false), + FillControlOptions::Adaptive(settings) => (settings, true), + }; + let max_units = (settings.max_bytes_per_second / TICKS_PER_SECOND) + .max(max_record) + .div_ceil(UNIT); + if max_units > u64::from(u32::MAX) { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "fill burst exceeds budget representation", + )); + } + let max_ops = u64::from(settings.max_records_per_second / TICKS_PER_SECOND as u32).max(1); + let mut table = Vec::new(); + table.try_reserve_exact(slots).map_err(|_| { + io::Error::new( + io::ErrorKind::OutOfMemory, + "cannot allocate fill observations", + ) + })?; + table.resize_with(slots, || Slot { + start_ns: AtomicU64::new(0), + bytes: AtomicU64::new(0), + }); + Ok(Some(Arc::new(Self { + options: settings, + enforcing, + max_units, + max_ops, + origin: Instant::now(), + credit: AtomicU64::new(pack_credit(max_ops, max_units)), + last_refill_ns: AtomicU64::new(0), + pressure: AtomicU8::new(encode_pressure(FillPressure::Healthy)), + stopped: AtomicBool::new(false), + recovering: AtomicUsize::new(0), + pause_holders: AtomicUsize::new(0), + rejections: AtomicU64::new(0), + would_reject: AtomicU64::new(0), + dropped_observations: AtomicU64::new(0), + completed_bytes: AtomicU64::new(0), + completed_ops: AtomicU64::new(0), + slots: table.into_boxed_slice(), + }))) + } + + pub fn checkpoint(start: Instant, timeout: Duration) -> Instant { + start + (timeout / 4).min(STALL_CAP) + } + + fn elapsed_ns(&self) -> u64 { + nanos(self.origin.elapsed()) + } + + pub fn stop(&self) { + self.stopped.store(true, Ordering::Release); + self.publish(); + } + + pub fn set_recovering(&self, recovering: bool) { + if recovering { + self.recovering.fetch_add(1, Ordering::AcqRel); + } else { + loop { + let holds = self.recovering.load(Ordering::Acquire); + if holds == 0 + || self + .recovering + .compare_exchange(holds, holds - 1, Ordering::AcqRel, Ordering::Relaxed) + .is_ok() + { + break; + } + } + } + self.publish(); + } + + fn hold_pause(&self) { + self.pause_holders.fetch_add(1, Ordering::AcqRel); + self.publish(); + } + + fn release_pause(&self) { + self.pause_holders.fetch_sub(1, Ordering::AcqRel); + self.publish(); + } + + fn publish(&self) { + let paused = self.stopped.load(Ordering::Acquire) + || self.recovering.load(Ordering::Acquire) != 0 + || self.pause_holders.load(Ordering::Acquire) != 0; + let next = if paused { + FillPressure::Paused + } else { + FillPressure::Healthy + }; + let previous = + decode_pressure(self.pressure.swap(encode_pressure(next), Ordering::Release)); + if previous != next && previous != FillPressure::Disabled { + log::info!(target: "cache2::health", event = "cache_fill_pressure_changed", pressure:? = next; + "cache fill admission pressure changed"); + } + } + + fn paused(&self) -> bool { + self.stopped.load(Ordering::Acquire) + || self.recovering.load(Ordering::Acquire) != 0 + || self.pause_holders.load(Ordering::Acquire) != 0 + } + + pub fn suppress_reinsertion(&self) -> bool { + self.enforcing && self.paused() + } + + pub fn snapshot(&self) -> FillControlSnapshot { + let pressure = if self.paused() { + FillPressure::Paused + } else { + FillPressure::Healthy + }; + let now = self.elapsed_ns(); + let mut outstanding_operations = 0_u64; + let mut outstanding_bytes = 0_u64; + let mut oldest = 0_u64; + for slot in self.slots.iter() { + let start = slot.start_ns.load(Ordering::Acquire); + if start == 0 { + continue; + } + outstanding_operations += 1; + outstanding_bytes = + outstanding_bytes.saturating_add(slot.bytes.load(Ordering::Relaxed)); + oldest = oldest.max(now.saturating_sub(start)); + } + let completed = self.completed_bytes.load(Ordering::Relaxed); + let drain = if outstanding_operations != 0 && completed != 0 && now != 0 { + (outstanding_bytes as f64 / (completed as f64 / (now as f64 / 1e9)) * 1e9) + .min(u64::MAX as f64) as u64 + } else { + 0 + }; + let (byte_rate, record_rate) = if pressure == FillPressure::Paused { + (0, 0) + } else { + ( + self.options.max_bytes_per_second, + self.options.max_records_per_second, + ) + }; + FillControlSnapshot { + pressure, + enforcing: self.enforcing, + bytes_per_second: byte_rate, + records_per_second: record_rate, + rejections: self.rejections.load(Ordering::Relaxed), + would_reject: self.would_reject.load(Ordering::Relaxed), + dropped_observations: self.dropped_observations.load(Ordering::Relaxed), + outstanding_operations, + outstanding_bytes, + oldest_operation_ns: oldest, + estimated_drain_ns: drain, + } + } + + /// Foreground admission: Adaptive rejects only while paused. + pub fn try_admit(&self) -> bool { + if !self.enforcing { + if self.paused() { + self.would_reject.fetch_add(1, Ordering::Relaxed); + } + return true; + } + if self.paused() { + self.rejections.fetch_add(1, Ordering::Relaxed); + false + } else { + true + } + } + + /// Background flush pacing. Essential flushes always proceed. + pub fn try_flush(&self, bytes: u64, records: u32, essential: bool) -> Option { + if essential || !self.enforcing { + return Some(FlushCharge { ops: 0, units: 0 }); + } + self.refill(); + let raw_units = bytes.div_ceil(UNIT).max(1); + let raw_ops = u64::from(records.max(1)); + let units = raw_units.min(self.max_units); + let ops = raw_ops.min(self.max_ops); + let byte_oversized = raw_units >= self.max_units; + let record_oversized = raw_ops >= self.max_ops; + let mut value = self.credit.load(Ordering::Relaxed); + for _ in 0..MAX_CAS_ATTEMPTS { + let have_units = packed_units(value); + let have_ops = packed_ops(value); + if have_units == 0 || have_ops == 0 { + return None; + } + if !byte_oversized && have_units < units { + return None; + } + if !record_oversized && have_ops < ops { + return None; + } + let take_units = if byte_oversized { have_units } else { units }; + let take_ops = if record_oversized { have_ops } else { ops }; + let next = pack_credit(have_ops - take_ops, have_units - take_units); + match self + .credit + .compare_exchange(value, next, Ordering::AcqRel, Ordering::Relaxed) + { + Ok(_) => { + return Some(FlushCharge { + ops: take_ops, + units: take_units, + }); + } + Err(current) => value = current, + } + } + None + } + + pub fn refund_flush(&self, charge: FlushCharge) { + if !self.enforcing || (charge.ops == 0 && charge.units == 0) { + return; + } + self.add_credit(charge.ops, charge.units); + } + + fn refill(&self) { + let now = self.elapsed_ns(); + let last = self.last_refill_ns.load(Ordering::Relaxed); + let Some(dt) = now.checked_sub(last).filter(|dt| *dt != 0) else { + return; + }; + let add_units = self + .options + .max_bytes_per_second + .saturating_mul(dt) + .checked_div(1_000_000_000) + .unwrap_or(0) + / UNIT; + let add_ops = u64::from(self.options.max_records_per_second) + .saturating_mul(dt) + .checked_div(1_000_000_000) + .unwrap_or(0); + if add_units == 0 && add_ops == 0 { + return; + } + if self + .last_refill_ns + .compare_exchange(last, now, Ordering::Relaxed, Ordering::Relaxed) + .is_err() + { + return; + } + self.add_credit(add_ops, add_units); + } + + fn add_credit(&self, ops: u64, units: u64) { + if ops == 0 && units == 0 { + return; + } + let mut value = self.credit.load(Ordering::Relaxed); + loop { + let next = pack_credit( + (packed_ops(value) + ops).min(self.max_ops), + (packed_units(value) + units).min(self.max_units), + ); + match self + .credit + .compare_exchange(value, next, Ordering::Relaxed, Ordering::Relaxed) + { + Ok(_) => return, + Err(current) => value = current, + } + } + } + + pub fn observe(&self, bytes: u64, _timeout: Duration) -> Option> { + let start = self.elapsed_ns().max(1); + for (index, slot) in self.slots.iter().enumerate() { + if slot + .start_ns + .compare_exchange(0, start, Ordering::AcqRel, Ordering::Acquire) + .is_ok() + { + slot.bytes.store(bytes, Ordering::Relaxed); + return Some(Observation { + control: self, + index, + succeeded: false, + slow: false, + }); + } + } + self.dropped_observations.fetch_add(1, Ordering::Relaxed); + None + } +} + +/// One worker-owned observation. Failed work is removed, never marked successful. +pub struct Observation<'a> { + control: &'a FillController, + index: usize, + succeeded: bool, + slow: bool, +} +impl Observation<'_> { + pub fn note_slow(&mut self) { + if self.slow { + return; + } + self.slow = true; + self.control.hold_pause(); + } + + pub fn clear_slow(&mut self) { + if !self.slow { + return; + } + self.slow = false; + self.control.release_pause(); + } + + pub fn finish(mut self) { + self.succeeded = true; + } +} +impl Drop for Observation<'_> { + fn drop(&mut self) { + let slot = &self.control.slots[self.index]; + let bytes = slot.bytes.load(Ordering::Relaxed); + slot.start_ns.store(0, Ordering::Release); + if self.succeeded { + self.control.completed_ops.fetch_add(1, Ordering::Relaxed); + self.control + .completed_bytes + .fetch_add(bytes, Ordering::Relaxed); + } + if self.slow { + self.control.release_pause(); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn control(enforcing: bool) -> Arc { + let settings = FillLimits::new(64_000, 100); + FillController::new( + if enforcing { + FillControlOptions::Adaptive(settings) + } else { + FillControlOptions::Observe(settings) + }, + 4, + 4096, + ) + .unwrap() + .unwrap() + } + + fn restore_burst(control: &FillController) { + control.credit.store( + pack_credit(control.max_ops, control.max_units), + Ordering::Release, + ); + } + + fn freeze_refill(control: &FillController) { + control + .last_refill_ns + .store(control.elapsed_ns(), Ordering::Relaxed); + } + + #[test] + fn foreground_admit_ignores_flush_budget() { + let control = control(true); + for _ in 0..10 { + assert!(control.try_flush(64, 1, false).is_some()); + } + assert!(control.try_flush(64, 1, false).is_none()); + assert!(control.try_admit()); + } + + #[test] + fn record_budget_paces_nonessential_flush() { + let control = control(true); + for _ in 0..10 { + assert!(control.try_flush(64, 1, false).is_some()); + } + assert!(control.try_flush(64, 1, false).is_none()); + assert!(control.try_flush(64, 1, true).is_some()); + restore_burst(&control); + assert!(control.try_flush(64, 1, false).is_some()); + } + + #[test] + fn byte_budget_paces_nonessential_flush() { + let control = control(true); + assert!(control.try_flush(4096, 1, false).is_some()); + assert!(control.try_flush(4096, 1, false).is_none()); + assert!(control.try_flush(64, 1, false).is_some()); + restore_burst(&control); + assert!(control.try_flush(4096, 1, false).is_some()); + } + + #[test] + fn old_request_pauses_while_other_requests_progress() { + let control = control(true); + let mut old = control.observe(4096, Duration::from_secs(2)).unwrap(); + let fast = control.observe(64, Duration::from_secs(2)).unwrap(); + fast.finish(); + old.note_slow(); + assert_eq!(control.snapshot().pressure, FillPressure::Paused); + assert!(!control.try_admit()); + drop(old); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + } + + #[test] + fn observe_never_rejects_or_suppresses_reinsertion() { + let control = control(false); + let mut old = control.observe(4096, Duration::from_secs(2)).unwrap(); + old.note_slow(); + assert_eq!(control.snapshot().pressure, FillPressure::Paused); + for _ in 0..20 { + assert!(control.try_admit()); + } + assert_eq!(control.snapshot().would_reject, 20); + assert_eq!(control.snapshot().rejections, 0); + assert!(!control.suppress_reinsertion()); + } + + #[test] + fn observe_would_reject_counts_pause_not_budget() { + let control = control(false); + for _ in 0..20 { + assert!(control.try_flush(4096, 1, false).is_some()); + assert!(control.try_admit()); + } + assert_eq!(control.snapshot().would_reject, 0); + let mut old = control.observe(4096, Duration::from_secs(2)).unwrap(); + old.note_slow(); + assert_eq!(control.snapshot().pressure, FillPressure::Paused); + assert!(control.try_admit()); + assert_eq!(control.snapshot().would_reject, 1); + assert_eq!(control.snapshot().rejections, 0); + } + + #[test] + fn idle_time_does_not_inflate_credit() { + let control = control(true); + control.last_refill_ns.store(0, Ordering::Relaxed); + control.refill(); + let credit = control.credit.load(Ordering::Relaxed); + assert_eq!(packed_ops(credit), control.max_ops); + assert_eq!(packed_units(credit), control.max_units); + } + + #[test] + fn large_flush_remains_eligible_and_stop_cannot_reopen_admission() { + let control = FillController::new( + FillControlOptions::Adaptive(FillLimits::new(640, 10)), + 1, + 4096, + ) + .unwrap() + .unwrap(); + assert!(control.try_flush(4096, 1, false).is_some()); + restore_burst(&control); + assert!(control.try_flush(4096, 1, false).is_some()); + control.stop(); + assert!(!control.try_admit()); + } + + #[test] + fn recovering_fence_does_not_need_outstanding_io() { + let control = control(true); + control.set_recovering(true); + assert_eq!(control.snapshot().pressure, FillPressure::Paused); + assert!(!control.try_admit()); + control.set_recovering(false); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + } + + #[test] + fn overlapping_recovery_holds_resume_when_all_release() { + let control = control(true); + control.set_recovering(true); + control.set_recovering(true); + assert!(!control.try_admit()); + control.set_recovering(false); + assert_eq!(control.snapshot().pressure, FillPressure::Paused); + assert!(!control.try_admit()); + control.set_recovering(false); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + } + + #[test] + fn failed_observations_release_capacity_without_reporting_success() { + let control = FillController::new( + FillControlOptions::Observe(FillLimits::new(640, 10)), + 1, + 4096, + ) + .unwrap() + .unwrap(); + let pending = control.observe(4096, Duration::from_secs(5)).unwrap(); + assert!(control.observe(64, Duration::from_secs(5)).is_none()); + assert_eq!(control.snapshot().dropped_observations, 1); + drop(pending); + let pending = control.observe(4096, Duration::from_secs(5)).unwrap(); + drop(pending); + assert_eq!(control.completed_ops.load(Ordering::Relaxed), 0); + assert_eq!(control.completed_bytes.load(Ordering::Relaxed), 0); + assert_eq!(control.slots[0].start_ns.load(Ordering::Acquire), 0); + } + + #[test] + fn concurrent_flush_never_exceeds_shared_budget() { + let control = control(true); + freeze_refill(&control); + let accepted = AtomicU64::new(0); + std::thread::scope(|scope| { + for _ in 0..8 { + let control = &control; + let accepted = &accepted; + scope.spawn(move || { + for _ in 0..100 { + if control.try_flush(640, 1, false).is_some() { + accepted.fetch_add(1, Ordering::Relaxed); + } + } + }); + } + }); + assert_eq!(accepted.load(Ordering::Relaxed), 10); + } + + #[test] + fn io_completion_releases_pause_before_observation_drop() { + let control = control(true); + let mut pending = control.observe(4096, Duration::from_secs(2)).unwrap(); + pending.note_slow(); + assert!(!control.try_admit()); + pending.clear_slow(); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + pending.finish(); + } + + #[test] + fn concurrent_pause_holders_resume_when_all_release() { + let control = control(true); + std::thread::scope(|scope| { + for _ in 0..8 { + let control = &control; + scope.spawn(move || { + let mut pending = control.observe(64, Duration::from_secs(2)).unwrap(); + pending.note_slow(); + assert!(!control.try_admit()); + drop(pending); + }); + } + }); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert_eq!(control.pause_holders.load(Ordering::Acquire), 0); + assert!(control.try_admit()); + } + + #[test] + fn overlapping_hold_and_release_cannot_stick_paused() { + let control = control(true); + for _ in 0..1_000 { + let mut first = control.observe(64, Duration::from_secs(2)).unwrap(); + let mut second = control.observe(64, Duration::from_secs(2)).unwrap(); + first.note_slow(); + second.note_slow(); + drop(first); + drop(second); + assert_eq!(control.pause_holders.load(Ordering::Acquire), 0); + assert_eq!(control.snapshot().pressure, FillPressure::Healthy); + assert!(control.try_admit()); + } + } + + #[test] + fn burst_sized_flush_proceeds_on_remaining_credit() { + let control = control(true); + freeze_refill(&control); + assert!(control.try_flush(64, 1, false).is_some()); + assert!(control.try_flush(1_000_000, 1, false).is_some()); + restore_burst(&control); + freeze_refill(&control); + assert!(control.try_flush(64, 1, false).is_some()); + assert!(control.try_flush(64, 1_000_000, false).is_some()); + } + + #[test] + fn refund_restores_nonessential_budget() { + let control = control(true); + freeze_refill(&control); + let mut last = FlushCharge { ops: 0, units: 0 }; + for _ in 0..10 { + last = control.try_flush(64, 1, false).unwrap(); + } + assert!(control.try_flush(64, 1, false).is_none()); + control.refund_flush(last); + assert!(control.try_flush(64, 1, false).is_some()); + } + + #[test] + fn sub_tick_refills_do_not_burn_elapsed_time() { + let control = FillController::new( + FillControlOptions::Adaptive(FillLimits::new(640, 10)), + 1, + 4096, + ) + .unwrap() + .unwrap(); + freeze_refill(&control); + while control.try_flush(64, 1, false).is_some() {} + let one_ms_ago = control.elapsed_ns().saturating_sub(1_000_000); + control.last_refill_ns.store(one_ms_ago, Ordering::Relaxed); + for _ in 0..32 { + control.refill(); + } + assert_eq!(control.last_refill_ns.load(Ordering::Relaxed), one_ms_ago); + std::thread::sleep(Duration::from_millis(110)); + control.refill(); + assert!(control.try_flush(64, 1, false).is_some()); + } + + #[test] + fn concurrent_refunds_are_not_dropped() { + let control = FillController::new( + FillControlOptions::Adaptive(FillLimits::new(64_000, 10_000)), + 1, + 4096, + ) + .unwrap() + .unwrap(); + freeze_refill(&control); + while control.try_flush(64, 1, false).is_some() {} + let before = control.credit.load(Ordering::Relaxed); + std::thread::scope(|scope| { + for _ in 0..32 { + let control = &control; + scope.spawn(move || { + control.refund_flush(FlushCharge { ops: 1, units: 1 }); + }); + } + }); + let credit = control.credit.load(Ordering::Relaxed); + assert_eq!( + packed_ops(credit), + (packed_ops(before) + 32).min(control.max_ops) + ); + assert_eq!( + packed_units(credit), + (packed_units(before) + 32).min(control.max_units) + ); + } + + #[test] + fn oversized_refund_restores_only_consumed_remainder() { + let control = control(true); + freeze_refill(&control); + for _ in 0..9 { + assert!(control.try_flush(64, 1, false).is_some()); + } + let charge = control.try_flush(1_000_000, 1, false).unwrap(); + assert_eq!(charge.ops, 1); + assert_eq!(charge.units, control.max_units - 9); + control.refund_flush(charge); + let credit = control.credit.load(Ordering::Relaxed); + assert_eq!(packed_ops(credit), 1); + assert_eq!(packed_units(credit), charge.units); + assert!(control.try_flush(64, 1, false).is_some()); + assert!(control.try_flush(64, 1, false).is_none()); + } +} diff --git a/cache2/src/io/mod.rs b/cache2/src/io/mod.rs index 9086108..7bda564 100644 --- a/cache2/src/io/mod.rs +++ b/cache2/src/io/mod.rs @@ -16,3 +16,4 @@ pub mod backend; pub mod engine; +pub mod fill_control; diff --git a/cache2/src/lib.rs b/cache2/src/lib.rs index 8c73851..0d8b567 100644 --- a/cache2/src/lib.rs +++ b/cache2/src/lib.rs @@ -35,6 +35,8 @@ pub use self::cache::Value; mod config; pub use self::config::CacheConfig; pub use self::config::StorageLayout; +pub use self::config::runtime::FillControlOptions; +pub use self::config::runtime::FillLimits; pub use self::config::runtime::IoEngineOptions; pub use self::config::runtime::IoMode; pub use self::config::runtime::IoUringOptions; @@ -56,6 +58,8 @@ pub use self::snapshot::CacheL1Snapshot; pub use self::snapshot::CacheReclaimSnapshot; pub use self::snapshot::CacheSnapshot; pub use self::snapshot::DetailedCacheSnapshot; +pub use self::snapshot::FillControlSnapshot; +pub use self::snapshot::FillPressure; pub use self::snapshot::RegionSnapshot; pub use self::snapshot::StartupMode; diff --git a/cache2/src/region/runtime/metrics.rs b/cache2/src/region/runtime/metrics.rs index cbe0cf2..2c03f2d 100644 --- a/cache2/src/region/runtime/metrics.rs +++ b/cache2/src/region/runtime/metrics.rs @@ -215,6 +215,7 @@ impl RuntimeMetrics { CacheSnapshot { metrics_epoch: self.metrics_epoch, health, + fill_control: crate::snapshot::FillControlSnapshot::default(), activity_counters_enabled, puts, deletes, diff --git a/cache2/src/region/runtime/mod.rs b/cache2/src/region/runtime/mod.rs index 6bd22c6..bc98686 100644 --- a/cache2/src/region/runtime/mod.rs +++ b/cache2/src/region/runtime/mod.rs @@ -61,6 +61,9 @@ use crate::io::engine::ReadSlotWaiter; use crate::io::engine::build_file_engine; use crate::io::engine::recovery::BackgroundRecovery; use crate::io::engine::submit_background_io; +use crate::io::fill_control::FLUSH_RETRY; +use crate::io::fill_control::FillController; +use crate::io::fill_control::FlushCharge; use crate::managed_memory::BufferLease; use crate::managed_memory::CACHE_THREAD_STACK_BYTES; use crate::managed_memory::ManagedMemory; @@ -839,6 +842,14 @@ impl RegionDataPlane { return Err(write_overload_error()); } }; + if let Some(fill) = &running.recovery.fill + && !fill.try_admit() + { + if running.activity_counters { + running.metrics.record_write_rejection(); + } + return Err(write_overload_error()); + } let staged = self.core.try_stage_value( &running.staging, shard_id, @@ -1296,6 +1307,9 @@ impl RegionDataPlane { { snapshot.health = crate::snapshot::CacheHealth::Recovering; } + if let Some(fill) = &running.recovery.fill { + snapshot.fill_control = fill.snapshot(); + } snapshot.io = aggregate_io_stats( &running.read_engines, &running.write_engines, @@ -1497,6 +1511,11 @@ fn start_running( io::Error::new(io::ErrorKind::OutOfMemory, "cannot allocate shard controls") })?; shards.resize_with(shard_count, || Arc::new(ShardControl::new())); + let fill = FillController::new( + runtime.fill_control, + shard_count + reclaim_worker_count, + data.geometry.region_size, + )?; let shared = Arc::new(RunningShared { core, read_engines, @@ -1506,7 +1525,7 @@ fn start_running( reclaim_engines, reclaim_control: ReclaimControl::new(), reclaim_io_timeout: runtime.reclaim_io_timeout, - recovery: BackgroundRecovery::new(runtime.io_recovery_timeout), + recovery: BackgroundRecovery::with_fill(runtime.io_recovery_timeout, fill), managed_memory, metrics, memory, @@ -1816,7 +1835,12 @@ fn reclaim_worker_result( // Keep one completion boundary per source Region while each // reclaimer rotates through a disjoint subset of append shards. let reinsert_shard = reinsert_shards.take(); - let preserve_hot = shared.core.reclaim_can_reinsert()?; + let preserve_hot = shared.core.reclaim_can_reinsert()? + && !shared + .recovery + .fill + .as_ref() + .is_some_and(|fill| fill.suppress_reinsertion()); let reinsert_operation = if preserve_hot { shared.operations.try_enter() } else { @@ -1896,14 +1920,32 @@ fn shard_worker_result( )?); } if force_flush || fill.bytes >= shared.write_flush_threshold_bytes { - let engine = shared.write_engine_for(shard_id as u64); - shared.core.flush_staging_shard( - &shared.staging, - engine.as_ref(), - shard_id, - &shared.recovery, - )?; - deadline = None; + let essential = flags & (WAKE_URGENT | WAKE_ROTATE) != 0 || draining; + let bytes = fill.bytes as u64; + let records = u32::try_from(fill.records).unwrap_or(u32::MAX); + let charge = match &shared.recovery.fill { + None => Some(FlushCharge { ops: 0, units: 0 }), + Some(control) => control.try_flush(bytes, records, essential), + }; + if let Some(charge) = charge { + let engine = shared.write_engine_for(shard_id as u64); + match shared.core.flush_staging_shard( + &shared.staging, + engine.as_ref(), + shard_id, + &shared.recovery, + )? { + Some(_) => deadline = None, + None => { + if let Some(control) = &shared.recovery.fill { + control.refund_flush(charge); + } + deadline = Some(Instant::now() + STAGING_RETRY_DELAY); + } + } + } else { + deadline = Some(Instant::now() + FLUSH_RETRY); + } } } Ok(None) => { diff --git a/cache2/src/region/runtime/shutdown_tests.rs b/cache2/src/region/runtime/shutdown_tests.rs index d0310a0..c47801b 100644 --- a/cache2/src/region/runtime/shutdown_tests.rs +++ b/cache2/src/region/runtime/shutdown_tests.rs @@ -272,3 +272,79 @@ fn recovery_rejects_fills_preserves_reads_and_waits_for_all_workers() { store.close_fast().unwrap(); std::fs::remove_dir_all(root).unwrap(); } + +#[test] +fn adaptive_pressure_preserves_reads_and_deletes_and_resumes_fills() { + use crate::region::file_backend::FileRegionBackend; + use crate::region::file_backend::RegionFiles; + use crate::region::recovery::PersistentId; + use crate::region::store::RegionStore; + use crate::snapshot::CacheHealth; + let root = env::temp_dir().join(format!("cache2-fill-admission-{}", std::process::id())); + std::fs::create_dir_all(&root).unwrap(); + let files = RegionFiles::new(root.join("data"), root.join("state"), root.join("image")); + let data = DataSuperblock { + generation: 1, + cache_uuid: PersistentId::from_bytes([1; 16]).unwrap(), + data_identity: PersistentId::from_bytes([2; 16]).unwrap(), + geometry: DataGeometry { + data_file_len: DataGeometry::expected_file_len(4096, 4).unwrap(), + region_size: 4096, + region_count: 4, + }, + hash_seed: 3, + storage_fingerprint: 4, + }; + let config = RuntimeOptions { + fill_control: crate::FillControlOptions::Adaptive(crate::FillLimits::new(1_048_576, 1000)), + append_shards: 1, + l1_capacity_bytes: 0, + ..RuntimeOptions::default() + }; + let mut store = RegionStore::open( + 8, + FileRegionBackend::for_test_with_options(files, data, 8, config), + ) + .unwrap(); + let plane = store.data_plane_handle().unwrap(); + plane.put(b"existing", b"value").unwrap(); + plane.drain().unwrap(); + let fill = plane.shared.recovery.fill.as_ref().unwrap(); + fill.set_recovering(true); + let deadline = Instant::now() + Duration::from_secs(2); + while fill.snapshot().pressure != crate::FillPressure::Paused { + assert!(Instant::now() < deadline, "controller did not pause fills"); + std::thread::sleep(Duration::from_millis(1)); + } + let snapshot = plane.snapshot().unwrap(); + assert_eq!(snapshot.health, CacheHealth::Running); + assert_eq!(snapshot.fill_control.pressure, crate::FillPressure::Paused); + assert_eq!( + plane.put(b"new", b"value").unwrap_err().kind(), + io::ErrorKind::WouldBlock + ); + assert_eq!( + plane.put_l2(b"new", b"value").unwrap_err().kind(), + io::ErrorKind::WouldBlock + ); + assert_eq!(plane.get(b"existing").unwrap().unwrap().value(), b"value"); + plane.delete(b"existing").unwrap(); + assert!(plane.get(b"existing").unwrap().is_none()); + assert!(plane.put(b"new", b"value").is_err()); + fill.set_recovering(false); + let deadline = Instant::now() + Duration::from_secs(2); + loop { + match plane.put(b"new", b"value") { + Ok(_) => break, + Err(error) => { + assert_eq!(error.kind(), io::ErrorKind::WouldBlock); + assert!(Instant::now() < deadline, "fills did not resume"); + std::thread::sleep(Duration::from_millis(1)); + } + } + } + plane.drain().unwrap(); + assert_eq!(plane.get(b"new").unwrap().unwrap().value(), b"value"); + store.close_fast().unwrap(); + std::fs::remove_dir_all(root).unwrap(); +} diff --git a/cache2/src/snapshot.rs b/cache2/src/snapshot.rs index 28bd540..8b5e8eb 100644 --- a/cache2/src/snapshot.rs +++ b/cache2/src/snapshot.rs @@ -39,9 +39,50 @@ pub enum CacheHealth { Failed, } -/// Lock-free point-in-time operational counters and cache-owned resource -/// accounting. Counters are process-local and reset on every open. Concurrent -/// updates may appear across fields at slightly different instants. +/// Admission pressure, independent of terminal cache health. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub enum FillPressure { + /// Pressure observation is disabled. + #[default] + Disabled, + /// Configured rate ceilings apply. + Healthy, + /// New fills are paused while outstanding work is old or stalled. + Paused, +} + +/// Always available when fill control is enabled, independent of statistics. +#[non_exhaustive] +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub struct FillControlSnapshot { + /// Current measured pressure; Observe mode does not enforce it. + pub pressure: FillPressure, + /// Whether the controller enforces its admission decisions. + pub enforcing: bool, + /// Configured encoded-byte ceiling when healthy; zero while paused. + pub bytes_per_second: u64, + /// Configured fill-record ceiling when healthy; zero while paused. + pub records_per_second: u32, + /// Fills rejected by Adaptive while paused. + pub rejections: u64, + /// Observe-mode pause refusals. + pub would_reject: u64, + /// Background observations skipped because the table was full. + pub dropped_observations: u64, + /// Background operations awaiting completion or validation. + pub outstanding_operations: u64, + /// Bytes held by those background operations; excludes unflushed staging. + pub outstanding_bytes: u64, + /// Age of the oldest background operation. + pub oldest_operation_ns: u64, + /// Estimated drain time using validated background throughput since open. + pub estimated_drain_ns: u64, +} + +/// Point-in-time operational counters and cache-owned resource accounting. +/// Sampling uses atomics. Counters are process-local and reset on every open. +/// Concurrent updates may appear across fields at slightly different instants. #[non_exhaustive] #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct CacheSnapshot { @@ -51,6 +92,8 @@ pub struct CacheSnapshot { pub metrics_epoch: u64, /// Current cache availability. pub health: CacheHealth, + /// Pre-timeout fill pressure and controller accounting. + pub fill_control: FillControlSnapshot, /// Whether optional cumulative activity and I/O counters are enabled. pub activity_counters_enabled: bool, /// Accepted `put` and `put_l2` operations. @@ -245,7 +288,7 @@ pub struct RegionSnapshot { #[non_exhaustive] #[derive(Clone, Debug, Eq, PartialEq)] pub struct DetailedCacheSnapshot { - /// Lock-free summary sampled for this diagnostic. + /// Operational summary sampled for this diagnostic. pub summary: CacheSnapshot, /// Mutations rejected specifically because an append buffer needed progress. pub write_buffer_rejections: u64,