From 568e2d9a1588c684ddd6510d698e5db02f181465 Mon Sep 17 00:00:00 2001 From: Tin Dang Date: Mon, 20 Jul 2026 11:24:21 +0700 Subject: [PATCH] =?UTF-8?q?refactor(storage):=20one=20spill=20pipeline=20?= =?UTF-8?q?=E2=80=94=20sync=20eviction=20routes=20through=20the=20durable?= =?UTF-8?q?=20batch=20writer=20(W2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Moon had two spill families: the async/batch path (SpillThread flush_buffer: shared multi-entry files, one manifest commit per file) and a sync per-victim path (kv_spill::spill_to_datafile: ONE single-entry file + ONE shard-thread-blocking durable manifest commit PER KEY) used by the background eviction tick, the memory-pressure cascade's sync fallback, and db_quota. The two evolved independently and that seam already bit three times: the task #34 plain-drop misclassification, the task #45 is_string gate, and the #139 spill-test blindness were all divergence bugs between these copies. This PR makes the batch machinery the ONE spill implementation: - select_victim (policy dispatch) and build_spill_payload (victim serialization) exist once — previously copy-pasted at all three entry points. - evict_batch_durable_no_aof is renamed evict_batch_durable and now backs every SpillContext sync path: the sync wrapper reclaims its whole deficit as batches; evict_one_with_spill's spill arm is a deficit-1 delegation. Write-then-durable-then-drop ordering, D1 remove-before-cold-index, and fail-closed retention on I/O error are inherited, not re-implemented. - Red-first: sync-evicting 40 keys now produces 1 shared .mpf file (was 40 files + 40 blocking manifest fsync round-trips). - Fail-closed serialization (red-first): all three sites previously did serialize_collection(..).unwrap_or_default() — a corrupt value was "spilled" as an EMPTY body and evicted (silent durable data loss on reload). Unserializable victims are now retained hot, skipped, and loudly logged, on the sync, batch, and async paths. - evict_batch_durable now calls record_eviction() per removed key (the batch path previously under-counted the eviction metric). - kv_spill::spill_to_datafile is no longer a production eviction path; it remains the single-entry primitive test fixtures use (same page layout the batch salvage path emits). Gates: lib suite 4421 green (macOS + Linux VM release-fast); eviction pin tests (spill-I/O-failure OOM retention, non-string durable spill, db_index restamp) all green; offload/cold integration suites + crash matrix on Linux VM; clippy -D warnings on default and tokio,jemalloc feature sets. Stacked on refactor/w1-value-codec (PR #396). author: Tin Dang --- CHANGELOG.md | 30 +++ src/config.rs | 2 +- src/storage/db.rs | 2 +- src/storage/eviction.rs | 465 +++++++++++++++++++-------------- src/storage/tiered/kv_spill.rs | 22 +- 5 files changed, 313 insertions(+), 208 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 93326c7c..c74b7ef1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,7 +6,37 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Performance +- **Sync eviction spill now batches like the async path (W2).** The + tick-driven sync spill paths (`run_eviction`, the memory-pressure cascade's + sync fallback) previously wrote **one single-entry `.mpf` file and did one + shard-thread-blocking durable manifest commit per victim key**; they now + route through the same `flush_buffer` batch writer the background + `SpillThread` and the no-AOF durable path use — shared multi-entry files + (up to 256 keys), one manifest commit per file. Evicting N keys under + memory pressure costs ~N/256 manifest fsync round-trips instead of N. + ### Changed +- **One spill pipeline (W2).** Victim policy dispatch (`select_victim`) and + victim serialization (`build_spill_payload`) now exist once instead of + being copy-pasted at all three spill entry points; the sync + `SpillContext` path delegates to the one durable batch spiller + (`evict_batch_durable`, formerly `evict_batch_durable_no_aof` — no longer + no-AOF-only). The sync/async divergence class behind the task #34 + plain-drop misclassification, the task #45 `is_string` gate, and the #139 + test blindness is structurally gone. `kv_spill::spill_to_datafile` is no + longer a production eviction path (kept as the single-entry primitive test + fixtures use). + +### Fixed +- **Unserializable eviction victims are retained, fail-closed.** All three + spill sites previously did `serialize_collection(..).unwrap_or_default()`: + a victim whose value failed to serialize (in-memory corruption) was + "spilled" as an EMPTY value body and evicted — silent durable data loss on + reload. Such victims are now retained hot and skipped, loudly logged. +- **Batch-durable evictions now count in the eviction metric.** The durable + batch path (`--appendonly no` cascade) removed keys without calling + `record_eviction()`; single-victim paths did. Both now record. - **One value codec for RDB snapshots and KV disk-offload spill (W1).** The ~150-line value-body encoder/decoder that existed twice — `rdb::write_entry`'s value section and `kv_serde::serialize_collection` — kept bit-compatible only diff --git a/src/config.rs b/src/config.rs index f661b3b7..32fd18d7 100644 --- a/src/config.rs +++ b/src/config.rs @@ -988,7 +988,7 @@ impl ServerConfig { /// (`--appendonly yes`) nor RDB (`--save`) durability is on (GCP /// benchmark finding, 2026-07-10). /// - /// The durable-spill eviction path (`evict_batch_durable_no_aof`) needs + /// The durable-spill eviction path (`evict_batch_durable`) needs /// a `ShardManifest`, which today is only threaded through the /// tick-driven memory-pressure cascade /// (`shard::persistence_tick::handle_memory_pressure`) — itself gated on diff --git a/src/storage/db.rs b/src/storage/db.rs index 473bec6f..79ca5726 100644 --- a/src/storage/db.rs +++ b/src/storage/db.rs @@ -3163,7 +3163,7 @@ mod tests { // (.planning/reviews/storage-audit-2026-07-12-kv.md) // // Disk-offload spills Hash/List/Set/ZSet/Stream via - // `evict_one_async_spill` / `evict_batch_durable_no_aof` (production + // `evict_one_async_spill` / `evict_batch_durable` (production // paths, `kv_serde::serialize_collection` handles every type), but the // type-specific accessors used to consult only `self.data`, never the // `ColdIndex` — so a spilled collection read as absent, and a write diff --git a/src/storage/eviction.rs b/src/storage/eviction.rs index bce27a76..e00c7d10 100644 --- a/src/storage/eviction.rs +++ b/src/storage/eviction.rs @@ -449,7 +449,27 @@ pub fn try_evict_if_needed_with_spill_and_total_budget_reporting( return Err(oom_error()); } let before = db.estimated_memory(); - if !evict_one_with_spill(db, config, &policy, spill.as_deref_mut(), on_plain_drop) { + // W2: with a SpillContext, reclaim the whole deficit as durable + // batches (shared multi-entry files, ONE manifest commit per file) + // instead of one single-entry file + one blocking manifest fsync + // round-trip per victim. Without one, plain-drop one victim at a + // time (unchanged). + let progressed = match spill.as_deref_mut() { + Some(ctx) => { + evict_batch_durable( + db, + config, + &policy, + ctx.shard_dir, + ctx.next_file_id, + ctx.manifest, + ctx.db_index, + current_total.saturating_sub(budget), + ) > 0 + } + None => evict_one_with_spill(db, config, &policy, None, on_plain_drop), + }; + if !progressed { return Err(oom_error()); } let after = db.estimated_memory(); @@ -658,7 +678,7 @@ pub fn try_evict_if_needed_async_spill_budget_reporting( /// Under `--appendonly no` there is no second copy of the value anywhere /// once it leaves RAM, so the fast path would be unconditional, unrecoverable /// data loss on a crash. This function instead batches victims into -/// [`evict_batch_durable_no_aof`], which performs a **fully synchronous, +/// [`evict_batch_durable`], which performs a **fully synchronous, /// durable spill (pwrite + fsync + manifest commit) before dropping any hot /// value**, provided a `ShardManifest` is available (today only the /// tick-driven memory-pressure cascade threads one through — see @@ -733,7 +753,7 @@ pub fn try_evict_if_needed_async_spill_with_total_budget_reporting( if config.appendonly != "yes" { let Some(manifest) = manifest else { // No manifest reachable from this call site and no AOF backstop: - // durable spill (evict_batch_durable_no_aof, below) is impossible + // durable spill (evict_batch_durable, below) is impossible // here. That does NOT mean the cap goes unenforced, though: an // evicting policy (AllKeys*/Volatile*) must still honor Redis // semantics and reclaim the budget by DROPPING victims with no @@ -764,7 +784,7 @@ pub fn try_evict_if_needed_async_spill_with_total_budget_reporting( } let before = db.estimated_memory(); let deficit = current_total.saturating_sub(budget); - let reclaimed_keys = evict_batch_durable_no_aof( + let reclaimed_keys = evict_batch_durable( db, config, &policy, @@ -808,7 +828,7 @@ pub fn try_evict_if_needed_async_spill_with_total_budget_reporting( } /// Maximum victims collected into a single durable batch by -/// [`evict_batch_durable_no_aof`], mirroring `SpillThread`'s own +/// [`evict_batch_durable`], mirroring `SpillThread`'s own /// `FLUSH_ENTRY_CAP` (256) so the no-AOF-backstop path's on-disk batch size /// matches the async path's. const NO_AOF_BATCH_CAP: usize = 256; @@ -821,7 +841,70 @@ const NO_AOF_BATCH_CAP: usize = 256; /// only guards a pathological near-empty or heavily-duplicate keyspace. const NO_AOF_BATCH_STALL_LIMIT: usize = 16; -/// Batch-durable spill for the no-AOF-backstop case (`--appendonly no`). +/// Pick one eviction victim according to `policy` (W2: the single +/// implementation of the policy dispatch — previously copy-pasted in +/// `evict_batch_durable`, `evict_one_async_spill`, and +/// `evict_one_with_spill`). +fn select_victim( + db: &Database, + config: &RuntimeConfig, + policy: &EvictionPolicy, +) -> Option { + match policy { + EvictionPolicy::NoEviction => None, + EvictionPolicy::AllKeysLru => find_victim_lru(db, config.maxmemory_samples, false), + EvictionPolicy::AllKeysLfu => { + find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, false) + } + EvictionPolicy::AllKeysRandom => find_victim_random(db, false), + EvictionPolicy::VolatileLru => find_victim_lru(db, config.maxmemory_samples, true), + EvictionPolicy::VolatileLfu => { + find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, true) + } + EvictionPolicy::VolatileRandom => find_victim_random(db, true), + EvictionPolicy::VolatileTtl => find_victim_volatile_ttl(db, config.maxmemory_samples), + } +} + +/// Serialize a victim entry's value for spill (W2: the single implementation +/// — previously copy-pasted at all three spill entry points). +/// +/// Returns `(value_type, value_bytes, flags, ttl_ms)`, or `None` when the +/// value cannot be faithfully serialized (in-memory corruption, e.g. an +/// unparseable sorted-set listpack score). **`None` is fail-closed**: the +/// caller must retain the hot value and skip the victim — the pre-W2 +/// `unwrap_or_default()` at every call site durably spilled an EMPTY value +/// body and evicted the key (silent data loss on reload). +fn build_spill_payload( + entry: &crate::storage::Entry, +) -> Option<(ValueType, Bytes, u8, Option)> { + let val_ref = entry.as_redis_value(); + let (value_type, value_bytes): (ValueType, Bytes) = match val_ref { + RedisValueRef::String(s) => (ValueType::String, Bytes::copy_from_slice(s)), + ref other => { + let vt = kv_spill::value_type_of(other); + // `None` = corrupt value (serialize_collection logs it loudly). + let buf = kv_serde::serialize_collection(other)?; + (vt, Bytes::from(buf)) + } + }; + let mut flags: u8 = 0; + let ttl_ms = if entry.has_expiry() { + flags |= entry_flags::HAS_TTL; + Some(entry.expires_at_ms(0)) + } else { + None + }; + Some((value_type, value_bytes, flags, ttl_ms)) +} + +/// Batch-durable synchronous spill (W2: the ONE sync spill implementation). +/// +/// Originally built for the no-AOF-backstop case (`--appendonly no`); since +/// W2 it also backs every `SpillContext` sync path (`evict_one_with_spill`, +/// the tick-driven `run_eviction`, the memory-pressure cascade's sync +/// fallback), which previously wrote one single-entry file + one blocking +/// manifest commit per victim through `kv_spill::spill_to_datafile`. /// /// Collects up to [`NO_AOF_BATCH_CAP`] distinct victims (or until the /// estimated bytes staged cover `deficit`, whichever comes first) WITHOUT @@ -838,7 +921,7 @@ const NO_AOF_BATCH_STALL_LIMIT: usize = 16; /// Returns the number of keys durably evicted (0 if nothing could be /// evicted -- the caller surfaces OOM). #[allow(clippy::too_many_arguments)] -fn evict_batch_durable_no_aof( +fn evict_batch_durable( db: &mut Database, config: &RuntimeConfig, policy: &EvictionPolicy, @@ -854,21 +937,7 @@ fn evict_batch_durable_no_aof( let mut stall = 0usize; while buffer.len() < NO_AOF_BATCH_CAP && staged_bytes < deficit { - let victim = match policy { - EvictionPolicy::NoEviction => None, - EvictionPolicy::AllKeysLru => find_victim_lru(db, config.maxmemory_samples, false), - EvictionPolicy::AllKeysLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, false) - } - EvictionPolicy::AllKeysRandom => find_victim_random(db, false), - EvictionPolicy::VolatileLru => find_victim_lru(db, config.maxmemory_samples, true), - EvictionPolicy::VolatileLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, true) - } - EvictionPolicy::VolatileRandom => find_victim_random(db, true), - EvictionPolicy::VolatileTtl => find_victim_volatile_ttl(db, config.maxmemory_samples), - }; - let key = match victim { + let key = match select_victim(db, config, policy) { Some(k) => k, None => break, }; @@ -887,22 +956,16 @@ fn evict_batch_durable_no_aof( // already gone, so keep sampling for more victims. continue; }; - let val_ref = entry.as_redis_value(); - let collection_buf: Vec; - let (value_type, value_bytes): (ValueType, &[u8]) = match val_ref { - RedisValueRef::String(s) => (ValueType::String, s), - ref other => { - let vt = kv_spill::value_type_of(other); - collection_buf = kv_serde::serialize_collection(other).unwrap_or_default(); - (vt, collection_buf.as_slice()) - } - }; - let mut flags: u8 = 0; - let ttl_ms = if entry.has_expiry() { - flags |= entry_flags::HAS_TTL; - Some(entry.expires_at_ms(0)) - } else { - None + // Fail-closed: an unserializable (corrupt) value is retained hot and + // skipped — its `seen` entry stops the sampler from re-picking it + // into this batch, and staging continues with other victims. + let Some((value_type, value_bytes, flags, ttl_ms)) = build_spill_payload(entry) else { + warn!( + key = %String::from_utf8_lossy(key.as_bytes()), + "kv_spill: victim value cannot be serialized; retaining hot \ + value and skipping it (fail-closed, no empty-body spill)" + ); + continue; }; staged_bytes += key.as_bytes().len() + value_bytes.len(); @@ -911,7 +974,7 @@ fn evict_batch_durable_no_aof( buffer.push(SpillRequest { key: Bytes::copy_from_slice(key.as_bytes()), db_index, - value_bytes: Bytes::copy_from_slice(value_bytes), + value_bytes, value_type, flags, ttl_ms, @@ -957,6 +1020,7 @@ fn evict_batch_durable_no_aof( // so the ColdIndex insert MUST come AFTER remove, never before, // or `remove` wipes out the very entry being added. db.remove(&entry.key); + crate::admin::metrics_setup::record_eviction(); if let Some(ref mut ci) = db.cold_index { ci.insert( entry.key, @@ -979,7 +1043,7 @@ fn evict_batch_durable_no_aof( /// /// AOF-backstopped fast path only (see the durability-ordering doc comment /// on [`try_evict_if_needed_async_spill_with_total_budget`], which routes to -/// [`evict_batch_durable_no_aof`] instead when `--appendonly no`). Extracts +/// [`evict_batch_durable`] instead when `--appendonly no`). Extracts /// the entry, removes it from DashTable (immediate RAM relief), then sends a /// `SpillRequest` to the background thread for pwrite. fn evict_one_async_spill( @@ -991,23 +1055,8 @@ fn evict_one_async_spill( next_file_id: &mut u64, db_index: usize, ) -> bool { - // Find victim key using same policy logic as sync path - let victim = match policy { - EvictionPolicy::NoEviction => None, - EvictionPolicy::AllKeysLru => find_victim_lru(db, config.maxmemory_samples, false), - EvictionPolicy::AllKeysLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, false) - } - EvictionPolicy::AllKeysRandom => find_victim_random(db, false), - EvictionPolicy::VolatileLru => find_victim_lru(db, config.maxmemory_samples, true), - EvictionPolicy::VolatileLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, true) - } - EvictionPolicy::VolatileRandom => find_victim_random(db, true), - EvictionPolicy::VolatileTtl => find_victim_volatile_ttl(db, config.maxmemory_samples), - }; - - let key = match victim { + // Find victim key using the shared policy dispatch (same as sync path) + let key = match select_victim(db, config, policy) { Some(k) => k, None => return false, }; @@ -1015,26 +1064,16 @@ fn evict_one_async_spill( // Build SpillRequest from the entry BEFORE removing it from DashTable. // This is CPU work only -- no I/O on the event loop. if let Some(entry) = db.data().get(key.as_bytes()) { - let val_ref = entry.as_redis_value(); - - // Determine value_type and serialize value bytes - let collection_buf: Vec; - let (value_type, value_bytes): (ValueType, &[u8]) = match val_ref { - RedisValueRef::String(s) => (ValueType::String, s), - ref other => { - let vt = kv_spill::value_type_of(other); - collection_buf = kv_serde::serialize_collection(other).unwrap_or_default(); - (vt, collection_buf.as_slice()) - } - }; - - // Determine flags and TTL - let mut flags: u8 = 0; - let ttl_ms = if entry.has_expiry() { - flags |= entry_flags::HAS_TTL; - Some(entry.expires_at_ms(0)) - } else { - None + // Fail-closed: a value that cannot be faithfully serialized must not + // be evicted (pre-W2 this spilled an EMPTY body and dropped the key). + // Bail; the caller surfaces OOM rather than losing data. + let Some((value_type, value_bytes, flags, ttl_ms)) = build_spill_payload(entry) else { + warn!( + key = %String::from_utf8_lossy(key.as_bytes()), + "kv_spill: async-spill victim cannot be serialized; retaining \ + hot value (fail-closed, no empty-body spill)" + ); + return false; }; let file_id = *next_file_id; @@ -1043,7 +1082,7 @@ fn evict_one_async_spill( let req = SpillRequest { key: Bytes::copy_from_slice(key.as_bytes()), db_index, - value_bytes: Bytes::copy_from_slice(value_bytes), + value_bytes, value_type, flags, ttl_ms, @@ -1100,128 +1139,34 @@ pub(crate) fn evict_one_with_spill( spill: Option<&mut SpillContext<'_>>, on_plain_drop: &mut dyn FnMut(&[u8]), ) -> bool { - // Find victim key using policy-specific sampling - let victim = match policy { - EvictionPolicy::NoEviction => None, - EvictionPolicy::AllKeysLru => find_victim_lru(db, config.maxmemory_samples, false), - EvictionPolicy::AllKeysLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, false) - } - EvictionPolicy::AllKeysRandom => find_victim_random(db, false), - EvictionPolicy::VolatileLru => find_victim_lru(db, config.maxmemory_samples, true), - EvictionPolicy::VolatileLfu => { - find_victim_lfu(db, config.maxmemory_samples, config.lfu_decay_time, true) - } - EvictionPolicy::VolatileRandom => find_victim_random(db, true), - EvictionPolicy::VolatileTtl => find_victim_volatile_ttl(db, config.maxmemory_samples), - }; + // W2: the SpillContext arm delegates to the ONE durable batch spiller + // (deficit 1 byte stages exactly one victim). Everything the per-victim + // path used to hand-roll — write-then-durable-then-drop ordering, the + // task #34/#45 every-type spill guarantees, fail-closed retention on + // I/O error, D1 remove-before-cold-index ordering — is inherited from + // `evict_batch_durable` instead of being a second implementation that + // can drift (the drift is exactly what caused #34 defect 1 and #45). + if let Some(ctx) = spill { + return evict_batch_durable( + db, + config, + policy, + ctx.shard_dir, + ctx.next_file_id, + ctx.manifest, + ctx.db_index, + 1, + ) > 0; + } - let key = match victim { + // Plain-drop path (no spill context): pick one victim and delete it. + let key = match select_victim(db, config, policy) { Some(k) => k, None => return false, }; - - // Spill to disk before removing, if context provided. `spilled` tracks - // whether a durable cold copy was ACTUALLY written for this specific - // victim -- not merely whether a `SpillContext` was passed in. Task #34 - // review (defect 1): the previous `is_plain_drop = spill.is_none()` - // snapshot, taken before this block ran, mis-classified any non-string - // victim (Hash/List/Set/ZSet/Stream) under a live `SpillContext` as - // "spilled" even though the `is_string` gate below never touches them -- - // they fell through to the unconditional `db.remove` with no cold copy - // AND no plain-drop report: silent data loss with no AOF/replication DEL - // record. Tying the flag to the branch that actually performs the write - // makes the two cases (no `SpillContext` vs. a `SpillContext` that - // couldn't spill this value type) equivalent by construction. - // task #45: `kv_spill::spill_to_datafile` has fully supported collection - // types (Hash/List/Set/ZSet/Stream, serialized via `kv_serde`) since the - // async write-path eviction gate started using it -- but this synchronous - // path (background tick `timers::run_eviction`, the memory-pressure - // cascade's sync-spill fallback, and `db_quota`) still gated the call - // behind an `is_string` check left over from before that support existed. - // A Hash/List/Set/ZSet/Stream victim picked here while a `SpillContext` - // WAS present therefore never got a byte written to `shard_dir`, yet fell - // through to the same unconditional `db.remove` a genuine plain-drop - // takes -- but WITHOUT the `on_plain_drop` accounting (the old - // `is_plain_drop = spill.is_none()` snapshot read `false` even though - // nothing durable exists anywhere for that key). Spilling every type here - // closes that gap; `spill_to_datafile`'s own type match is exhaustive - // (compile error if a new `RedisValueRef` variant is ever added without a - // `ValueType` mapping). - // - // `file_id`/`ttl_for_cold` are captured from the immutable `entry` - // borrow BEFORE `db.remove` (which also clears any pre-existing cold - // copy of this key, D1: DEL-before-spill ordering) and the index insert - // into `db.cold_index` (a different field of `db`, taken as `&mut` only - // AFTER the `entry` borrow has ended) -- same split-borrow shape as - // `evict_batch_durable_no_aof`'s `db.remove` -> `ci.insert` ordering - // above. - let mut spilled = false; - let mut spilled_file_id = 0u64; - let mut spilled_ttl_ms: Option = None; - let mut spilled_value_type = crate::persistence::kv_page::ValueType::String; - if let Some(ctx) = spill { - if let Some(entry) = db.data().get(key.as_bytes()) { - let file_id = *ctx.next_file_id; - let ttl_ms = if entry.has_expiry() { - Some(entry.expires_at_ms(0)) - } else { - None - }; - if let Err(e) = kv_spill::spill_to_datafile( - ctx.shard_dir, - file_id, - key.as_bytes(), - entry, - ctx.db_index, - ctx.manifest, - None, - ) { - warn!( - key = %String::from_utf8_lossy(key.as_bytes()), - error = %e, - "kv_spill: I/O error during spill; retaining hot value, \ - aborting this eviction attempt (no silent drop)" - ); - // Spill failed: the value has no durable copy anywhere. - // Do NOT evict -- keep the hot value resident and let the - // caller retry (or surface OOM) instead of losing data. - return false; - } - *ctx.next_file_id += 1; - spilled = true; - spilled_file_id = file_id; - spilled_ttl_ms = ttl_ms; - spilled_value_type = kv_spill::value_type_of(&entry.as_redis_value()); - } - } - db.remove(key.as_bytes()); crate::admin::metrics_setup::record_eviction(); - if spilled { - // Register the cold location so `Database::exists`/`get`/the - // type-specific accessors' `promote_cold_if_present` can find this - // key again -- `spill_to_datafile` itself only inserts into a - // `ColdIndex` it is explicitly handed (`None` here, split-borrow - // constraint above), so the caller owns the insert. Single-page spill - // only (no overflow-chain support in this path, same limitation - // `spill_to_datafile`'s single-file layout has), so `page_idx`/ - // `slot_idx` are always the entry's sole slot: 0/0. - if let Some(ref mut ci) = db.cold_index { - ci.insert( - Bytes::copy_from_slice(key.as_bytes()), - crate::storage::tiered::cold_index::ColdLocation { - file_id: spilled_file_id, - page_idx: 0, - slot_idx: 0, - ttl_ms: spilled_ttl_ms, - value_type: spilled_value_type, - }, - ); - } - } else { - on_plain_drop(key.as_bytes()); - } + on_plain_drop(key.as_bytes()); true } @@ -2030,7 +1975,7 @@ mod tests { /// RED→GREEN (GCP benchmark finding, 2026-07-10 -- policy-aware /// fail-close): when no `ShardManifest` is reachable (e.g. the inline /// per-connection write-path eviction gate) AND there is no AOF - /// backstop, durable spill (`evict_batch_durable_no_aof`) is impossible. + /// backstop, durable spill (`evict_batch_durable`) is impossible. /// That must NOT mean the cap goes unenforced for an evicting policy, /// though -- Redis semantics require `allkeys-lru` (and friends) to /// honor `--maxmemory` by DROPPING victims (no tiering) when no durable @@ -2374,6 +2319,136 @@ mod tests { ); } + /// RED→GREEN (W2): the sync-spill eviction path must batch victims into + /// shared multi-entry files (one manifest commit per file), not one + /// single-entry file + one blocking manifest commit per key. Before W2, + /// evicting 40 keys through `try_evict_if_needed_with_spill` produced 40 + /// `heap-*.mpf` files and 40 shard-thread-blocking manifest fsync + /// round-trips; the async path had already solved this (FLUSH_ENTRY_CAP + /// batching) — the sync path re-derived its own per-key layout, the exact + /// divergence class behind tasks #34/#45/#139. + #[test] + fn sync_spill_batches_victims_into_shared_files() { + let tmp = tempfile::tempdir().unwrap(); + let shard_dir = tmp.path(); + let manifest_path = tmp.path().join("shard.manifest"); + let mut manifest = ShardManifest::create(&manifest_path).unwrap(); + let mut next_file_id = 1u64; + + let mut db = Database::new(); + db.cold_index = Some(crate::storage::tiered::cold_index::ColdIndex::new()); + for i in 0..40u8 { + db.set_string( + Bytes::copy_from_slice(format!("batch_key_{i:02}").as_bytes()), + Bytes::from_static(b"batch_value_payload"), + ); + } + assert_eq!(db.len(), 40); + + let config = make_config(1, "allkeys-lru"); // budget 1 byte: evict all + let mut ctx = SpillContext { + shard_dir, + manifest: &mut manifest, + next_file_id: &mut next_file_id, + db_index: 0, + }; + + let result = try_evict_if_needed_with_spill(&mut db, &config, Some(&mut ctx)); + assert!(result.is_ok(), "eviction must succeed: {result:?}"); + assert_eq!(db.len(), 0, "all victims must be evicted from RAM"); + + // Every key must remain cold-readable. + let ci = db.cold_index.as_ref().unwrap(); + for i in 0..40u8 { + let key = format!("batch_key_{i:02}"); + assert!( + ci.lookup(key.as_bytes()).is_some(), + "{key} must be registered in ColdIndex after sync spill" + ); + } + + // The batching claim itself: 40 small same-db victims must share + // files (the async path packs ~256/file; assert a generous bound so + // the test pins "batched", not an exact packing constant). + let file_count = std::fs::read_dir(shard_dir.join("data")) + .unwrap() + .flatten() + .filter(|e| e.file_name().to_string_lossy().ends_with(".mpf")) + .count(); + assert!( + file_count <= 4, + "sync spill must batch victims into shared files: \ + 40 keys produced {file_count} .mpf files (pre-W2: 40)" + ); + } + + /// RED→GREEN (W2): a victim whose value cannot be faithfully serialized + /// (in-memory corruption, e.g. a sorted-set listpack score that does not + /// parse — see `ValueCodecError::CorruptScore`) must be retained hot, + /// fail-closed. Before W2 all three spill sites did + /// `serialize_collection(..).unwrap_or_default()`: the victim was + /// "spilled" as an EMPTY value body and evicted — silent, durable data + /// loss (reload of the empty body yields a decode miss). + #[test] + fn sync_spill_unserializable_victim_retained_fail_closed() { + let tmp = tempfile::tempdir().unwrap(); + let shard_dir = tmp.path(); + let manifest_path = tmp.path().join("shard.manifest"); + let mut manifest = ShardManifest::create(&manifest_path).unwrap(); + let mut next_file_id = 1u64; + + let mut db = Database::new(); + db.cold_index = Some(crate::storage::tiered::cold_index::ColdIndex::new()); + // Inject a corrupt SortedSetListpack (unparseable score) the way a + // memory-corruption event would leave it — unreachable via ZADD. + let mut lp = crate::storage::listpack::Listpack::new(); + lp.push_back(b"member"); + lp.push_back(b"not-a-number"); + let mut entry = Entry::new_string(Bytes::new()); + entry.value = crate::storage::compact_value::CompactValue::from_redis_value( + crate::storage::RedisValue::SortedSetListpack(lp), + ); + db.insert_for_load(Bytes::from_static(b"corrupt_zset"), entry); + db.recalculate_memory(); + assert_eq!(db.len(), 1); + + let config = make_config(1, "allkeys-lru"); + let mut ctx = SpillContext { + shard_dir, + manifest: &mut manifest, + next_file_id: &mut next_file_id, + db_index: 0, + }; + + let mut reported: Vec> = Vec::new(); + let evicted = evict_one_with_spill( + &mut db, + &config, + &EvictionPolicy::AllKeysLru, + Some(&mut ctx), + &mut |key| reported.push(key.to_vec()), + ); + + assert!( + !evicted, + "an unserializable victim must NOT be evicted (fail-closed): \ + pre-W2 it was spilled as an EMPTY value body and dropped from RAM" + ); + assert_eq!( + db.len(), + 1, + "the corrupt value must be retained hot — no durable copy exists" + ); + assert!( + db.cold_index + .as_ref() + .unwrap() + .lookup(b"corrupt_zset") + .is_none(), + "no ColdIndex entry may point at an empty/absent spill body" + ); + } + #[test] fn test_evict_without_spill_unchanged() { let mut db = Database::new(); diff --git a/src/storage/tiered/kv_spill.rs b/src/storage/tiered/kv_spill.rs index e3eadadd..3bb327aa 100644 --- a/src/storage/tiered/kv_spill.rs +++ b/src/storage/tiered/kv_spill.rs @@ -112,19 +112,19 @@ pub fn write_kv_spill_pages( /// `ColdLocation::value_type` (#364). pub use crate::storage::value_codec::value_type_of; -/// Spill a single evicted KV entry to a DataFile on disk. +/// Spill a single KV entry to its own DataFile on disk. /// -/// Creates a single-page `.mpf` file at `{shard_dir}/data/heap-{file_id:06}.mpf`, -/// writes a `KvLeafPage` containing the entry, and registers the file in the -/// shard manifest. +/// Creates a `.mpf` file at `{shard_dir}/data/heap-{file_id:06}.mpf` (leaf +/// page + overflow chain for oversized values), writes the entry, and +/// registers the file in the shard manifest with a blocking durable commit. /// -/// String entries are fully supported. Non-string types (hash, list, set, zset, -/// stream) are skipped with a warning -- overflow serialization is future work. -/// -/// If the entry does not fit in a single 4KB page, it is skipped (oversized -/// entries require overflow pages, also future work). -/// -/// Returns `Ok(())` on success, skip, or best-effort failure logging. +/// **No longer a production eviction path (W2).** Every eviction spill — +/// sync and async — now routes through `spill_thread::flush_buffer` batches +/// (shared multi-entry files, one manifest commit per file). This helper +/// remains as the single-entry primitive: test fixtures use it to place one +/// key cold deterministically (same `build_kv_spill_pages` + +/// `write_kv_spill_pages` layout the batch salvage path emits, so the format +/// stays production-exercised via `spill_single_entry`). /// /// `db_index` is the logical database the entry was evicted FROM — it is /// stamped into the manifest `FileEntry` so `rebuild_from_manifest_per_db`