diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 95369582..edbb3034 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -851,7 +851,7 @@ async fn spawn_memory_diagnostics( let series_count = sketch_index.series_len(); let approx_bytes = sketch_index.approx_memory_bytes(); info!( - "[MEMORY_DIAG] SketchStore: {} instance(s), {} sid(s) with state, {:.2} KB approx sealed bytes", + "[MEMORY_DIAG] SketchStore: {} instance(s), {} sid(s) with state, {:.2} KB approx in-memory bytes (hot current_epoch + sealed)", instance_count, series_count, approx_bytes as f64 / 1024.0, diff --git a/data_plane/src/storage_engines/sketch_db/index/epoch_columnar.rs b/data_plane/src/storage_engines/sketch_db/index/epoch_columnar.rs index fe168088..8e613afc 100644 --- a/data_plane/src/storage_engines/sketch_db/index/epoch_columnar.rs +++ b/data_plane/src/storage_engines/sketch_db/index/epoch_columnar.rs @@ -158,6 +158,20 @@ impl

MutableEpoch

{ self.windows_set.len() } + /// Iterate every `(window, label_id, &payload)` entry in insertion + /// order. Used by the persistence layer's memory accounting so the + /// un-sealed hot epoch's footprint is visible (mirrors + /// [`SealedEpoch::entries`]). + pub fn iter_entries( + &self, + ) -> impl Iterator { + self.windows_col + .iter() + .zip(self.label_ids_col.iter()) + .zip(self.payloads_col.iter()) + .map(|((w, id), p)| (*w, *id, p)) + } + pub fn len(&self) -> usize { self.windows_col.len() } @@ -394,6 +408,49 @@ impl

MutableEpoch

{ dropped } + /// Split off every entry whose window-END is at or before + /// `cutoff_end`, returning them as a fresh `MutableEpoch` (the + /// retained, more-recent entries stay in `self`). Mirrors + /// [`Self::evict_window_ends_before`] but PRESERVES the aged entries + /// (returned) instead of dropping them, so a time-driven seal can + /// turn the aged-but-un-sealed tail of `current_epoch` into a sealed + /// epoch the persistence flusher can make durable. + /// + /// Returns `None` when nothing is old enough to split (so the caller + /// can skip the seal+rotate entirely). O(N) — one rebuild pass, same + /// shape as `evict_window_ends_before`. + pub fn split_window_ends_before(&mut self, cutoff_end: u64) -> Option> { + // O(1) skip: the earliest window-START is already past the + // cutoff, so no window can END at/before it either. + match self.min_start { + Some(min_s) if min_s > cutoff_end => return None, + None => return None, + _ => {} + } + let old_windows = std::mem::take(&mut self.windows_col); + let old_ids = std::mem::take(&mut self.label_ids_col); + let old_payloads = std::mem::take(&mut self.payloads_col); + self.windows_set.clear(); + self.window_to_ids = None; + self.last_window = None; + self.min_start = None; + self.max_end = None; + + let mut aged: MutableEpoch

= MutableEpoch::new(); + for ((w, id), p) in old_windows.into_iter().zip(old_ids).zip(old_payloads) { + if w.1 <= cutoff_end { + aged.insert(w, id, p); + } else { + self.insert(w, id, p); + } + } + if aged.is_empty() { + None + } else { + Some(aged) + } + } + /// Remove all entries whose window is in `windows`. /// Mirrors the legacy `SketchStore` CircularBuffer /// cleanup contract. O(N) — rebuilds columns in one pass. @@ -869,6 +926,37 @@ impl SidStoreData { } } + /// Time-driven seal for the persistence tier: roll every window in + /// `current_epoch` whose END is at or before `cutoff_end` into a + /// freshly-sealed epoch, leaving the more-recent windows in + /// `current_epoch`. Returns the number of distinct windows sealed. + /// + /// The count-driven [`Self::maybe_rotate_epoch`] cadence only seals + /// once `current_epoch` accumulates `seal_window_count` DISTINCT + /// windows. A slow or stalled series never reaches that threshold, so + /// its aged windows sit un-sealed in `current_epoch` forever — and the + /// flusher's hot-window phase only ever flushes SEALED epochs, so they + /// are never made durable (the live `parts/`-stays-empty bug). This + /// method, called by the flusher each tick with `cutoff_end = now - + /// hot_window`, guarantees that any window older than the hot window + /// becomes sealed (and therefore flushable) regardless of cadence. + /// + /// No-op when persistence is disabled (the in-memory-only path bounds + /// memory via retention-drop, not flush). + pub fn seal_aged_windows(&mut self, cutoff_end: u64) -> usize { + if !self.persistence_enabled { + return 0; + } + let Some(aged) = self.current_epoch.split_window_ends_before(cutoff_end) else { + return 0; + }; + let sealed_windows = aged.distinct_windows(); + let sealed = SealedEpoch::from_mutable(aged); + self.sealed_epochs.insert(self.current_epoch_id, sealed); + self.current_epoch_id += 1; + sealed_windows + } + fn maybe_rotate_epoch(&mut self) { // The effective rotation threshold is the smaller of the // test-only `epoch_capacity` and the durable-tier @@ -1205,6 +1293,64 @@ mod tests { ); } + #[test] + fn split_window_ends_before_partitions_aged_tail() { + let mut e = MutableEpoch::::new(); + e.insert((0, 100), 1, 1); + e.insert((100, 200), 2, 2); + e.insert((150, 250), 3, 3); // straddles 200 (ends after) → retained + e.insert((200, 300), 4, 4); + + // Cutoff 200: aged = windows ending <= 200 → (0,100) & (100,200). + let aged = e.split_window_ends_before(200).expect("aged tail exists"); + let mut aged_ids: Vec<_> = aged.iter_entries().map(|(_, id, _)| id).collect(); + aged_ids.sort(); + assert_eq!(aged_ids, vec![1, 2]); + // Retained: (150,250) & (200,300). + let mut kept_ids: Vec<_> = e.iter_entries().map(|(_, id, _)| id).collect(); + kept_ids.sort(); + assert_eq!(kept_ids, vec![3, 4]); + assert_eq!(e.min_start(), Some(150)); + assert_eq!(e.max_end(), Some(300)); + } + + #[test] + fn split_window_ends_before_noop_when_all_recent() { + let mut e = MutableEpoch::::new(); + e.insert((1000, 1100), 1, 1); + assert!(e.split_window_ends_before(500).is_none()); + assert_eq!(e.distinct_windows(), 1); + } + + #[test] + fn seal_aged_windows_rolls_unsealed_tail_under_persistence() { + // The time-driven seal: aged windows below the count-cadence must + // still seal so the flusher can make them durable. + let mut s = SidStoreData::::new(); + s.seal_window_count = Some(10); // high cadence → never count-seals + s.persistence_enabled = true; + for i in 0..3u32 { + let start = (i as u64) * 30_000; + s.insert((start, start + 30_000), "series".into(), i); + } + // Below cadence → nothing sealed yet. + assert_eq!(s.sealed_epochs.len(), 0); + // Seal everything ending at/before 90_000 (all 3 windows). + let sealed = s.seal_aged_windows(90_000); + assert_eq!(sealed, 3, "all three aged windows should seal"); + assert_eq!(s.sealed_epochs.len(), 1); + assert_eq!(s.current_epoch.distinct_windows(), 0); + } + + #[test] + fn seal_aged_windows_noop_when_persistence_disabled() { + let mut s = SidStoreData::::new(); + s.persistence_enabled = false; + s.insert((0, 30_000), "series".into(), 1); + assert_eq!(s.seal_aged_windows(60_000), 0); + assert_eq!(s.sealed_epochs.len(), 0); + } + #[test] fn sealed_epochs_survive_for_flush_under_persistence() { // Sealed epochs accumulate (pending flush) and are NOT dropped by diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index 0f27b580..6820ce9a 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -73,6 +73,51 @@ fn tag_to_encoding(tag: u8) -> SketchEncoding { } } +/// Reconstruct an exact-aggregation accumulator from its on-disk +/// `(type_name, bytes)` pair so the durable tier can serve the +/// exact-agg query path (`query_exact_agg_range` / `sum by (...)`) after +/// flush+evict. Covers the deterministic scalar accumulators the live +/// marquee `sum by (zone)` path uses; the sketch-backed accumulator forms +/// (DDSketch/KLL/HLL/CountSketch — registered as `AggKind::Sketch`) are +/// served as opaque bytes via [`SketchStore::query_range`] and are NOT +/// reconstructed here. Returns `None` for an unrecognized `type_name` +/// (the caller skips the disk entry rather than fabricating a wrong +/// payload) — see the remaining-follow-up note in the PR. +fn reconstruct_exact_agg( + type_name: &str, + bytes: &[u8], +) -> Option> { + use crate::precompute_engine::operators::{ + IncreaseAccumulator, MinMaxAccumulator, MultipleIncreaseAccumulator, + MultipleSumAccumulator, SumAccumulator, + }; + use crate::storage_engines::types::AggregateCore; + match type_name { + "SumAccumulator" => SumAccumulator::deserialize_from_bytes(bytes) + .ok() + .map(|a| Box::new(a) as Box), + "IncreaseAccumulator" => IncreaseAccumulator::deserialize_from_bytes(bytes) + .ok() + .map(|a| Box::new(a) as Box), + "MinMaxAccumulator" => MinMaxAccumulator::deserialize_from_bytes(bytes) + .ok() + .map(|a| Box::new(a) as Box), + "MultipleSumAccumulator" => MultipleSumAccumulator::deserialize_from_bytes(bytes) + .ok() + .map(|a| Box::new(a) as Box), + "MultipleIncreaseAccumulator" => MultipleIncreaseAccumulator::deserialize_from_bytes(bytes) + .ok() + .map(|a| Box::new(a) as Box), + // `MultipleMinMaxAccumulator` needs an external `sub_type` + // (min/max) not recorded in the part, and the sketch-backed + // accumulator forms have no generic byte factory — both are left + // to the deferred exact-agg/sketch precompute read-back work (see + // PR follow-up note). They are still served from memory; only the + // evicted-to-disk portion is skipped for these types. + _ => None, + } +} + /// Joint helper shared by [`SketchStore::ingest_precompute_for_agg_config`] /// and [`SketchStore::ingest_precompute_with_sid`] — folds the /// grouping-label values on `output` against the @@ -843,32 +888,29 @@ impl SketchStore { BTreeMap, BTreeMap>, )> { - let store = match self.series.get(&sid) { - Some(s) => s.clone(), - None => return Vec::new(), - }; - let guard = store.write().unwrap(); - let mut by_label_id: HashMap< - LabelValuesId, + // Key by the resolved label MAP (not `LabelValuesId`) so the + // in-memory tier and the durable disk tier — which carry + // independent intern spaces — union by label identity. Mirrors + // `query_range`. Absent in-memory store is NOT an early return: + // under persistence the windows may have been flushed-then-evicted + // (or recovered from disk on restart), so the disk union below + // still runs. + let mut by_label_map: HashMap< + BTreeMap, BTreeMap>, > = HashMap::new(); - let mut buf: Vec<(TimestampRange, LabelValuesId, &AggPayload)> = Vec::new(); - guard - .current_epoch - .range_query_into(start_unix_ms, end_unix_ms, &mut buf); - for (win, label_id, payload) in &buf { - if let Some(p) = payload.as_exact_agg() { - by_label_id - .entry(*label_id) - .or_default() - .insert(win.1 as i64, Arc::from(p.clone_boxed_core())); - } - } - buf.clear(); + if let Some(store) = self.series.get(&sid).map(|s| s.clone()) { + let guard = store.write().unwrap(); + let mut by_label_id: HashMap< + LabelValuesId, + BTreeMap>, + > = HashMap::new(); - for sealed in guard.sealed_epochs.values() { - sealed.range_query_into(start_unix_ms, end_unix_ms, &mut buf); + let mut buf: Vec<(TimestampRange, LabelValuesId, &AggPayload)> = Vec::new(); + guard + .current_epoch + .range_query_into(start_unix_ms, end_unix_ms, &mut buf); for (win, label_id, payload) in &buf { if let Some(p) = payload.as_exact_agg() { by_label_id @@ -878,15 +920,98 @@ impl SketchStore { } } buf.clear(); + + for sealed in guard.sealed_epochs.values() { + sealed.range_query_into(start_unix_ms, end_unix_ms, &mut buf); + for (win, label_id, payload) in &buf { + if let Some(p) = payload.as_exact_agg() { + by_label_id + .entry(*label_id) + .or_default() + .insert(win.1 as i64, Arc::from(p.clone_boxed_core())); + } + } + buf.clear(); + } + + for (label_id, samples) in by_label_id { + let label_values_map = + guard.intern.resolve(label_id).cloned().unwrap_or_default(); + by_label_map.entry(label_values_map).or_default().extend(samples); + } + drop(guard); } - by_label_id - .into_iter() - .map(|(label_id, samples)| { - let label_values_map = guard.intern.resolve(label_id).cloned().unwrap_or_default(); - (label_values_map, samples) - }) - .collect() + // Union the durable disk tier for the flushed-then-evicted portion + // of the range. In-memory wins on a window-end collision. + self.union_disk_exact_agg_into(sid, start_unix_ms, end_unix_ms, &mut by_label_map); + + by_label_map.into_iter().collect() + } + + /// Union the durable disk tier's exact-aggregation entries into + /// `by_label_map` for `[start, end)`. No-op when persistence is off. + /// Disk entries are reconstructed via [`reconstruct_exact_agg`]; an + /// unrecognized accumulator type is skipped (sketch-backed forms are + /// served by `query_range`, not here — see PR follow-up note). + /// Containment scan (`start_ts >= start && end_ts <= end`) matches the + /// in-memory exact-agg `range_query_into`. + fn union_disk_exact_agg_into( + &self, + sid: u64, + start_unix_ms: u64, + end_unix_ms: u64, + by_label_map: &mut HashMap< + BTreeMap, + BTreeMap>, + >, + ) { + let handle = { + let g = self.persistence_read.read().unwrap(); + match g.as_ref() { + Some(h) => Arc::clone(h), + None => return, + } + }; + let Some(keys) = self.sid_group_by_keys(sid) else { + return; + }; + let parts = handle + .manifest + .live_parts_overlapping(start_unix_ms, end_unix_ms); + for pe in &parts { + let reader = match handle.part_cache.get_or_load(pe.part_id) { + Ok(r) => r, + Err(e) => { + tracing::warn!(part_id = pe.part_id, error = %e, "exact-agg disk read: open part failed"); + continue; + } + }; + for rec in reader.index_records() { + if rec.agg_id != sid { + continue; + } + // Containment, matching the in-memory exact-agg scan. + if !(rec.start_ts >= start_unix_ms && rec.end_ts <= end_unix_ms) { + continue; + } + let Ok(entry) = reader.load_entry(&rec) else { + continue; + }; + let Some(acc) = + reconstruct_exact_agg(&entry.sketch_type_name, &entry.sketch_bytes) + else { + continue; + }; + let label_map = Self::rebuild_label_map(&keys, &entry.label); + by_label_map + .entry(label_map) + .or_default() + // In-memory wins — only fill window-ends disk uniquely owns. + .entry(rec.end_ts as i64) + .or_insert_with(|| Arc::from(acc)); + } + } } /// Actual coverage bounds `(min_window_start_ms, max_window_end_ms)` @@ -906,43 +1031,76 @@ impl SketchStore { start_unix_ms: u64, end_unix_ms: u64, ) -> Option<(u64, u64)> { - let store = self.series.get(&sid)?.clone(); - let guard = store.read().unwrap(); let mut min_start: u64 = u64::MAX; let mut max_end: u64 = 0; let mut any = false; - let mut buf: Vec<(TimestampRange, LabelValuesId, &AggPayload)> = Vec::new(); - guard - .current_epoch - .range_query_into(start_unix_ms, end_unix_ms, &mut buf); - for (win, _label_id, payload) in &buf { - if payload.as_exact_agg().is_some() { - any = true; - if win.0 < min_start { - min_start = win.0; + // In-memory tier (absent store is not an early return — disk may + // still cover the range after flush+evict / restart). + if let Some(store) = self.series.get(&sid).map(|s| s.clone()) { + let guard = store.read().unwrap(); + let mut buf: Vec<(TimestampRange, LabelValuesId, &AggPayload)> = Vec::new(); + guard + .current_epoch + .range_query_into(start_unix_ms, end_unix_ms, &mut buf); + for (win, _label_id, payload) in &buf { + if payload.as_exact_agg().is_some() { + any = true; + min_start = min_start.min(win.0); + max_end = max_end.max(win.1); } - if win.1 > max_end { - max_end = win.1; + } + buf.clear(); + + for sealed in guard.sealed_epochs.values() { + sealed.range_query_into(start_unix_ms, end_unix_ms, &mut buf); + for (win, _label_id, payload) in &buf { + if payload.as_exact_agg().is_some() { + any = true; + min_start = min_start.min(win.0); + max_end = max_end.max(win.1); + } } + buf.clear(); } } - buf.clear(); - for sealed in guard.sealed_epochs.values() { - sealed.range_query_into(start_unix_ms, end_unix_ms, &mut buf); - for (win, _label_id, payload) in &buf { - if payload.as_exact_agg().is_some() { - any = true; - if win.0 < min_start { - min_start = win.0; + // Durable disk tier — same coverage-aware divisor must see the + // flushed-then-evicted windows, else `rate` over-divides by the + // nominal `[r]` once the recent data ages onto disk. + if let Some(handle) = { + let g = self.persistence_read.read().unwrap(); + g.as_ref().map(Arc::clone) + } { + let parts = handle + .manifest + .live_parts_overlapping(start_unix_ms, end_unix_ms); + for pe in &parts { + let Ok(reader) = handle.part_cache.get_or_load(pe.part_id) else { + continue; + }; + for rec in reader.index_records() { + if rec.agg_id != sid { + continue; + } + if !(rec.start_ts >= start_unix_ms && rec.end_ts <= end_unix_ms) { + continue; } - if win.1 > max_end { - max_end = win.1; + // Only count entries that reconstruct as exact-agg + // (skip sketch-backed disk entries under this sid). + let Ok(entry) = reader.load_entry(&rec) else { + continue; + }; + if reconstruct_exact_agg(&entry.sketch_type_name, &entry.sketch_bytes) + .is_none() + { + continue; } + any = true; + min_start = min_start.min(rec.start_ts); + max_end = max_end.max(rec.end_ts); } } - buf.clear(); } if any { @@ -1696,12 +1854,35 @@ impl crate::storage_engines::sketch_db::index::persistence::EpochSource for Sket data.sealed_epochs.remove(&epoch_id); } + fn seal_aged_epochs(&self, cutoff_end: u64) -> usize { + let mut sealed = 0usize; + for entry in self.series.iter() { + let Ok(mut data) = entry.value().write() else { + continue; + }; + sealed += data.seal_aged_windows(cutoff_end); + } + sealed + } + fn approx_memory_bytes(&self) -> usize { + // Account for BOTH `current_epoch` (hot, un-sealed) AND sealed + // epochs. Counting sealed-only under-reports the true footprint + // (the live "0.00–0.12 KB approx sealed bytes" diagnostic) and, + // worse, blinds the flusher's memory-pressure trigger to the bulk + // of memory — which under persistence (retention-drop disabled) + // lives in `current_epoch` until the time-driven seal rolls it + // over. The hot-window seal (`seal_aged_epochs`) handles the + // common case; this keeps the memory-pressure backstop honest for + // a burst that outruns the hot window. let mut total = 0usize; for entry in self.series.iter() { let Ok(data) = entry.value().read() else { continue; }; + for (_, _, payload) in data.current_epoch.iter_entries() { + total += payload.approx_bytes(); + } for epoch in data.sealed_epochs.values() { for (_, _, payload) in &epoch.entries { total += payload.approx_bytes(); @@ -2701,6 +2882,200 @@ mod tests { // No persistence handle → seal cadence disabled → no sealing. assert!(idx.persistence_read.read().unwrap().is_none()); } + + // ── LIVE-scenario regression tests (fix/sketch-durable-live) ──────── + // + // The unit tests above use `hot_window_ms: Some(0)` + epoch-1970 + // timestamps, which force-flush everything immediately. The LIVE run + // (`--persistence-seal-window-count=4 --persistence-hot-window-secs=120`) + // ingests panes stamped at WALL-CLOCK ms and uses a 120s hot window, + // and exposed three bugs these helpers must reproduce. + + fn now_ms_wall() -> u64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as u64 + } + + /// Mirror of the live config: seal every 4 windows, 120s hot window, + /// memory limit high (so the flush is driven by the hot-window + /// watermark, exactly as in the live run that produced empty parts/). + fn live_cfg(disk_path: std::path::PathBuf) -> SketchStorePersistenceConfig { + SketchStorePersistenceConfig { + memory_limit_bytes: 2048 * 1024 * 1024, + memory_low_watermark_bytes: 2048 * 1024 * 1024 * 8 / 10, + hard_cap_bytes: 2048 * 1024 * 1024 * 125 / 100, + hot_window_ms: Some(120_000), // live: --persistence-hot-window-secs=120 + delete_older_than_ms: None, + flush_interval: std::time::Duration::from_millis(20), + disk_path, + part_cache_bytes: 1 << 20, + seal_window_count: 4, // live: --persistence-seal-window-count=4 + } + } + + /// BUG #1 + #3 (most severe): in the LIVE run the freshest windows of + /// every series sit UN-SEALED in `current_epoch` — sealing only fires + /// once `current_epoch` reaches the cadence (4 distinct windows). The + /// flusher's hot-window phase ONLY ever considers SEALED epochs, so + /// any window that ages past the 120s hot window while still in + /// `current_epoch` (because the series stopped/slowed before hitting + /// cadence) is NEVER flushed. That is exactly what produced the empty + /// `parts/` + 0-byte manifest log on node2 after 13 min of ingest: + /// data flowed (the resolver WAL grew) but nothing was ever made + /// durable, so a `docker restart` recovered `live=0` and the post- + /// restart query returned "No result". + /// + /// This test ingests aged panes that DON'T reach cadence-4, so they + /// stay un-sealed, then asserts the flusher still makes them durable + /// and they survive a restart. On origin/main nothing flushes. + #[test] + fn live_aged_unsealed_panes_flush_and_survive_restart() { + let tmp = tempfile::TempDir::new().unwrap(); + let disk = tmp.path().to_path_buf(); + // Panes ending 10 minutes ago → comfortably behind the 120s hot + // window the moment they're ingested. + let base = now_ms_wall().saturating_sub(10 * 60 * 1000); + { + let idx = Arc::new(SketchStore::new()); + idx.register(meta_with_host_key(7001)); + let p = idx.start_persistence(live_cfg(disk.clone())).unwrap(); + // Only 3 distinct 30s panes — BELOW the 4-window seal cadence, + // so they never rotate into sealed_epochs and (on origin/main) + // the flusher's sealed-only hot-window scan never sees them. + for i in 0..3u64 { + let s = base + i * 30_000; + idx.append_sample(7001, lv_host("a"), (s, s + 30_000), sample((i + 1) as u8)); + } + // These aged windows MUST become durable parts even though the + // cadence was never reached. On origin/main this never happens + // → empty parts/, matching the live failure. + let flushed = wait_until( + || !p.manifest.live_parts().is_empty(), + std::time::Duration::from_secs(5), + ); + assert!( + flushed, + "LIVE BUG #1: aged un-sealed panes never flushed to disk \ + (parts={}, sealed={}, sealed_bytes={})", + p.manifest.live_parts().len(), + idx.list_sealed_epochs_len(), + idx.approx_memory_bytes(), + ); + let mut p = p; + p.shutdown(); + } + + // "Restart" on the SAME dir — recovery must reload the parts. + let idx2 = Arc::new(SketchStore::new()); + idx2.register(meta_with_host_key(7001)); + let p2 = idx2.start_persistence(live_cfg(disk.clone())).unwrap(); + assert!( + !p2.manifest.live_parts().is_empty(), + "LIVE BUG #1: recovery found 0 live parts after restart" + ); + let series = idx2.query_range(7001, base, base + 90_000); + assert_eq!(series.len(), 1, "recovered data not queryable after restart"); + assert!( + !series[0].samples.is_empty(), + "restart query returned No result — flushed data lost" + ); + drop(p2); + } + + /// BUG #2: after flush+evict, an exact-agg (`sum by (zone)` shape) + /// range query must still resolve from disk. On origin/main + /// `query_exact_agg_range` reads ONLY the in-memory current+sealed + /// epochs — it never unions disk parts — so once the windows are + /// flushed-then-evicted the query returns empty ("No result"). + #[test] + fn live_exact_agg_resolves_from_disk_after_evict() { + use crate::storage_engines::types::AggregationType; + let tmp = tempfile::TempDir::new().unwrap(); + let idx = Arc::new(SketchStore::new()); + // Register an ExactAgg(Sum) sid keyed by `zone`. + let mut m = meta(8001); + m.metric_name = "http_requests_total".into(); + m.group_by_keys = ["zone".to_string()].into_iter().collect(); + m.agg_kind = AggKind::ExactAgg { + agg_type: AggregationType::Sum, + parameters_canonical: String::new(), + spatial_filter_canonical: String::new(), + }; + idx.register(m); + let p = idx.start_persistence(durable_cfg(tmp.path().to_path_buf())).unwrap(); + + let lv_zone = |v: &str| { + let mut x = BTreeMap::new(); + x.insert("zone".to_string(), v.to_string()); + x + }; + for i in 0..10u64 { + let s = i * 30_000; + idx.append_precompute( + 8001, + lv_zone("z0"), + (s, s + 30_000), + Box::new(crate::precompute_engine::operators::SumAccumulator::with_sum( + (i + 1) as f64, + )), + ); + } + assert!( + wait_until( + || idx.approx_memory_bytes() == 0 && idx.list_sealed_epochs_len() == 0, + std::time::Duration::from_secs(5), + ), + "exact-agg windows never fully evicted" + ); + // Query the EVICTED portion [0, 150_000) — must come back from disk. + let series = idx.query_exact_agg_range(8001, 0, 150_000); + assert!( + !series.is_empty(), + "LIVE BUG #2: exact-agg query returned No result after flush+evict \ + (disk read-back missing)" + ); + let (_label, samples) = &series[0]; + assert!( + samples.contains_key(&30_000), + "LIVE BUG #2: evicted exact-agg window (0,30000) missing from disk read-back" + ); + // The coverage-bounds helper (rate divisor) must also see disk. + let cov = idx.exact_agg_coverage_bounds(8001, 0, 150_000); + assert!( + cov.is_some(), + "LIVE BUG #2: exact_agg_coverage_bounds blind to disk after evict" + ); + drop(p); + } + + /// BUG #3: the memory diagnostic + the flusher's memory-pressure + /// trigger must account for `current_epoch`, not just sealed epochs. + /// On origin/main `approx_memory_bytes()` sums ONLY sealed epochs, so + /// a store holding megabytes of un-sealed `current_epoch` data reports + /// ~0 bytes (the live "0.00–0.12 KB approx sealed bytes" under-report) + /// and the flusher's `mem > memory_limit` trigger never fires. + #[test] + fn live_total_memory_accounts_for_current_epoch() { + let idx = SketchStore::new(); + idx.register(meta_with_host_key(9001)); + // No persistence → no sealing → all data sits in current_epoch. + for i in 0..20u64 { + let s = i * 30_000; + idx.append_sample(9001, lv_host("a"), (s, s + 30_000), sample((i + 1) as u8)); + } + assert_eq!( + idx.list_sealed_epochs_len(), + 0, + "precondition: nothing sealed (no persistence)" + ); + assert!( + idx.approx_memory_bytes() > 0, + "LIVE BUG #3: approx_memory_bytes() reports 0 while current_epoch holds 20 \ + windows — the diagnostic under-reports and the flusher's memory trigger is blind" + ); + } } // 2026-05 reorg: generic epoch-partitioned columnar storage lives diff --git a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs index 728f9521..6963168f 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/flusher.rs @@ -262,6 +262,22 @@ struct TickStats { fn run_tick(shared: &Arc, source: &S) -> PersistResult { let mut stats = TickStats::default(); let now = now_ms(); + let cfg = &shared.cfg; + + // ---- Phase 0: time-driven seal of aged un-sealed windows ---- + // The count-driven seal cadence (`seal_window_count`) only fires once + // a series' `current_epoch` reaches N distinct windows. A slow or + // stalled series never gets there, so its aged windows sit un-sealed + // in `current_epoch` — and the hot-window phase below only flushes + // SEALED epochs, so those windows would never be made durable. Roll + // any window older than the hot window into a sealed epoch first so + // it becomes flushable this same tick. (Without `hot_window_ms` the + // durable tier is purely memory-pressure driven and phase 1 below + // handles eviction.) + if let Some(hot) = cfg.hot_window_ms { + let cutoff = now.saturating_sub(hot); + source.seal_aged_epochs(cutoff); + } // ---- Collect candidates across phase 1 and phase 2 ---- let mut all = source.list_sealed_epochs(); @@ -269,7 +285,6 @@ fn run_tick(shared: &Arc, source: &S) -> PersistR all.sort_by_key(|r| r.end_ts); let mem = source.approx_memory_bytes(); - let cfg = &shared.cfg; let need_memory_pressure = mem > cfg.memory_limit_bytes; let low_water_goal = cfg.memory_low_watermark_bytes; @@ -494,6 +509,86 @@ mod tests { } } + /// A fake source whose epochs start UN-SEALED — they only become + /// visible to `list_sealed_epochs` once `seal_aged_epochs` rolls them + /// over. Models the live `current_epoch` → `sealed_epochs` transition + /// the flusher's phase-0 time-seal drives. + struct LazySealSource { + unsealed: StdMutex>, + sealed: StdMutex>, + memory: AtomicU64, + } + + impl LazySealSource { + fn new(snapshots: Vec) -> Self { + let mut map = HashMap::new(); + let mut total = 0u64; + for s in snapshots { + total += s.approx_bytes as u64; + map.insert((s.agg_id, s.epoch_id), s); + } + Self { + unsealed: StdMutex::new(map), + sealed: StdMutex::new(HashMap::new()), + memory: AtomicU64::new(total), + } + } + } + + impl EpochSource for LazySealSource { + fn list_sealed_epochs(&self) -> Vec { + self.sealed + .lock() + .unwrap() + .values() + .map(|s| SealedEpochRef { + agg_id: s.agg_id, + epoch_id: s.epoch_id, + end_ts: s.max_ts, + approx_bytes: s.approx_bytes, + }) + .collect() + } + + fn seal_aged_epochs(&self, cutoff_end: u64) -> usize { + let mut un = self.unsealed.lock().unwrap(); + let mut sealed = self.sealed.lock().unwrap(); + let aged: Vec<(u64, u64)> = un + .iter() + .filter(|(_, s)| s.max_ts <= cutoff_end) + .map(|(k, _)| *k) + .collect(); + let mut n = 0; + for k in aged { + if let Some(s) = un.remove(&k) { + sealed.insert(k, s); + n += 1; + } + } + n + } + + fn snapshot_sealed_epoch( + &self, + agg_id: u64, + epoch_id: u64, + ) -> PersistResult> { + Ok(self.sealed.lock().unwrap().get(&(agg_id, epoch_id)).cloned()) + } + + fn evict_sealed_epoch(&self, agg_id: u64, epoch_id: u64) { + let mut map = self.sealed.lock().unwrap(); + if let Some(removed) = map.remove(&(agg_id, epoch_id)) { + self.memory + .fetch_sub(removed.approx_bytes as u64, Ordering::Relaxed); + } + } + + fn approx_memory_bytes(&self) -> usize { + self.memory.load(Ordering::Relaxed) as usize + } + } + fn snap(agg_id: u64, epoch_id: u64, min_ts: u64, max_ts: u64, approx: usize) -> EpochSnapshot { EpochSnapshot { agg_id, @@ -584,6 +679,40 @@ mod tests { assert!(!manifest.live_parts().is_empty()); } + #[test] + fn hot_window_time_seals_unsealed_aged_epochs_then_flushes() { + // LIVE BUG #1 at the flusher level: an epoch that is still + // UN-SEALED (below the count cadence) but aged past the hot window + // must be time-sealed by phase 0 and then flushed in the SAME + // tick. Without the `seal_aged_epochs` hook the flusher's + // sealed-only scan never sees it → empty parts/ (the live bug). + let tmp = TempDir::new().unwrap(); + let now = now_ms(); + // Aged 10 min; never sealed by the source on its own. + let snaps = vec![snap( + 1, + 1, + now.saturating_sub(600_000), + now.saturating_sub(595_000), + 100, + )]; + let source = Arc::new(LazySealSource::new(snaps)); + let manifest = Arc::new(Manifest::init(tmp.path()).unwrap()); + let mut cfg = test_cfg(tmp.path().to_path_buf(), 100_000); // under mem limit + cfg.hot_window_ms = Some(120_000); // 120s hot window, like live + let mut handle = FlusherHandle::start(cfg, manifest.clone(), source.clone()).unwrap(); + + let deadline = Instant::now() + Duration::from_secs(3); + while manifest.live_parts().is_empty() && Instant::now() < deadline { + thread::sleep(Duration::from_millis(20)); + } + handle.shutdown(); + assert!( + !manifest.live_parts().is_empty(), + "time-driven seal did not make the aged un-sealed epoch durable" + ); + } + #[test] fn t2_sweep_deletes_old_parts() { let tmp = TempDir::new().unwrap(); diff --git a/data_plane/src/storage_engines/sketch_db/persistence/source.rs b/data_plane/src/storage_engines/sketch_db/persistence/source.rs index c1c4da85..edb304fd 100644 --- a/data_plane/src/storage_engines/sketch_db/persistence/source.rs +++ b/data_plane/src/storage_engines/sketch_db/persistence/source.rs @@ -99,6 +99,23 @@ pub struct EpochSnapshotEntry { pub trait EpochSource: Send + Sync { fn list_sealed_epochs(&self) -> Vec; + /// Time-driven seal: roll every un-sealed `current_epoch` window + /// whose END is at or before `cutoff_end` into a sealed epoch so the + /// flusher can make it durable. Called by the flusher each tick with + /// `cutoff_end = now - hot_window` BEFORE [`list_sealed_epochs`]. + /// + /// The count-driven seal cadence (`seal_window_count`) only fires once + /// `current_epoch` reaches N distinct windows; a slow/stalled series + /// never reaches it, leaving aged windows un-sealed — and the flusher + /// only flushes SEALED epochs, so those windows are never made durable + /// (the live `parts/`-stays-empty bug). This hook closes that gap. + /// + /// Default impl is a no-op so existing/test sources need not implement + /// it. Returns the number of distinct windows newly sealed. + fn seal_aged_epochs(&self, _cutoff_end: u64) -> usize { + 0 + } + fn snapshot_sealed_epoch( &self, agg_id: u64,