From 3e1be1966bdb72c220157ea6a0ff785ff5a003f3 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 26 May 2026 09:47:22 -0600 Subject: [PATCH] feat(ingest): per-window delta-base rotation for sketch deltas MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The OTLP delta-apply path reconstructed per-series sketch state by merging each window's delta onto a never-reset running base (state(N) = state(N-1) ⊕ delta(N)). Because the edge tumbling window resets per-series sketch state every window, each window's delta is that window's marginal against an empty base — so accumulating forever over-counts across windows (and for HLL the register-wise max base holds the all-time max). See docs/delta-baseline-contract.md §3. This adds per-window base rotation on the consumer side: - Cache entry now carries the window start it was built for. The per-series snapshot cache value changes from a bare boxed accumulator to SnapshotCacheEntry { core, window_start }. - The delta-apply branch detects a per-series window boundary by comparing the data point's start_time_unix_nano against the cached base's window_start. On a new window it resets the cached base to empty (AggregateCore::reset_to_empty) BEFORE applying the new window's delta, so within a window deltas accumulate and at a new window the base starts fresh — state(N) is window N only. - Full frames (PROTO_FULL) keep REPLACE semantics and set the stored window_start. - reset_to_empty is sketch-agnostic: implemented for the additive delta-capable families (DDSketch, CMS, CountSketch, HLL), each preserving its shape/config (dims / relative accuracy / register width); default no-op elsewhere. KLL never deltas. - The existing "no base yet → drop the delta" guard is unchanged. Adds an ingest test that feeds full(win1) → delta(win1) → delta(win2): within win1 the delta accumulates onto the full frame, and at the win2 boundary the base is rotated so the reconstructed state is win2 only, not win1 + win2. Co-Authored-By: Claude Opus 4.7 (1M context) --- data_plane/src/drivers/ingest/otel.rs | 246 +++++++++++++++++- .../src/precompute_engine/ingest_handler.rs | 47 +++- .../operators/count_min_sketch_accumulator.rs | 7 + .../operators/count_sketch_accumulator.rs | 7 + .../operators/dd_sketch_accumulator.rs | 7 + .../operators/hll_sketch_accumulator.rs | 9 + .../src/storage_engines/types/traits.rs | 18 ++ 7 files changed, 325 insertions(+), 16 deletions(-) diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 77dc6390d..2ed7af8c7 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -1287,14 +1287,38 @@ async fn route_modified_otlp_sketches_to_precompute( // one in-flight delta per (metric, labels) so // the next full snapshot replaces the current // cache entry cleanly. + // + // Per-window base rotation + // (`docs/delta-baseline-contract.md` §3): the edge + // tumbling window resets per-series sketch state + // every window, so each window's delta is that + // window's marginal against an empty base. The + // backend therefore must NOT accumulate forever + // (`state(N) = state(N-1) ⊕ delta(N)`), which would + // over-count across windows. Instead we detect a + // window boundary per series — a change in the data + // point's window start (`start_time_unix_nano`) + // versus the `window_start` stored with the cached + // base — and reset the cached base to empty before + // applying the new window's delta. Within a window + // deltas still accumulate; at a new window the base + // starts fresh, so the reconstructed `state(N)` is + // window N only. Sketch-agnostic: the reset is the + // additive families' (DDSketch / CMS / CountSketch / + // HLL) `AggregateCore::reset_to_empty`; KLL never + // deltas. Full frames keep REPLACE semantics and set + // the stored `window_start`. let accumulator: Box = if dp.encoding == ENCODING_PROTO_DELTA || dp.encoding == ENCODING_MSGPACK_DELTA { - let Some(base) = ingest_state + let Some((mut merged, base_window_start)) = ingest_state .sketch_snapshots .get(&series_key) - .map(|e| e.clone_boxed_core()) + .map(|e| (e.core.clone_boxed_core(), e.window_start)) else { + // No base yet → drop the delta (agent must + // resend the next full frame). Unchanged + // guard. decoded_failed += 1; debug!( "OTLP delta-sketch arrived before any base \ @@ -1305,7 +1329,23 @@ async fn route_modified_otlp_sketches_to_precompute( ); continue; }; - let mut merged = base; + // Window boundary: the incoming delta opens a new + // window for this series. Rotate the base to empty + // so the new window starts fresh (state == this + // window only), keeping the sketch's shape/config + // intact for the additive apply below. + if dp.start_time_unix_nano != base_window_start { + debug!( + "OTLP delta-sketch window boundary (metric={}, \ + series_key={}, prev_window_start={}, \ + new_window_start={}); rotating per-series base", + metric.name, + series_key, + base_window_start, + dp.start_time_unix_nano + ); + merged.reset_to_empty(); + } if let Err(e) = apply_modified_otlp_delta_bytes( dp.kind, dp.encoding, @@ -1326,16 +1366,24 @@ async fn route_modified_otlp_sketches_to_precompute( ); continue; } - ingest_state - .sketch_snapshots - .insert(series_key.clone(), merged.clone_boxed_core()); + ingest_state.sketch_snapshots.insert( + series_key.clone(), + crate::precompute_engine::ingest_handler::SnapshotCacheEntry { + core: merged.clone_boxed_core(), + window_start: dp.start_time_unix_nano, + }, + ); merged } else { match decode_modified_otlp_sketch_bytes(dp.kind, dp.encoding, &dp.sketch) { Ok(acc) => { - ingest_state - .sketch_snapshots - .insert(series_key.clone(), acc.clone_boxed_core()); + ingest_state.sketch_snapshots.insert( + series_key.clone(), + crate::precompute_engine::ingest_handler::SnapshotCacheEntry { + core: acc.clone_boxed_core(), + window_start: dp.start_time_unix_nano, + }, + ); acc } Err(e) => { @@ -2894,6 +2942,186 @@ mod sid_resolution_tests { let _ = drain.await; } + /// Per-window base rotation (`docs/delta-baseline-contract.md` §3): + /// the backend must NOT accumulate deltas across windows. For one + /// series, a full frame opens window 1, a delta in window 1 (same + /// `start_time_unix_nano`) accumulates onto it, then a delta in + /// window 2 (a NEW `start_time_unix_nano`) must reset the cached base + /// to empty first — so the reconstructed state is window 2's delta + /// only, NOT window1 + window2. + /// + /// Uses DDSketch (an additive family) so accumulation vs. reset is + /// directly observable on the bucket counts. + #[tokio::test] + async fn delta_apply_rotates_per_series_base_at_window_boundary() { + use crate::precompute_engine::operators::DDSketchAccumulator; + use asap_otel_proto::sketchlib::v1::{DdSketchBucketDelta, DdSketchDelta as PbDelta}; + use asap_sketchlib::proto::sketchlib::DdSketchState; + use prost::Message; + + let (state, drain) = make_state().await; + + const WIN1_START: u64 = 1_000_000; + const WIN2_START: u64 = 2_000_000; + + // Build a DDSketch DataPoint with explicit encoding / window-start + // / payload so we can stage a full frame then per-window deltas. + let make_dp = |start: u64, ts: u64, encoding: i32, sketch: Vec| DdSketchDataPoint { + attributes: vec![kv("zone", "z0")], + start_time_unix_nano: start, + time_unix_nano: ts, + sketch, + encoding, + exemplars: Vec::new(), + flags: 0, + series_id: 0, + }; + + // The cache key (series_key) is derived from the canonical metric + // name + attrs; recompute it the same way the ingest loop does so + // we can read the reconstructed base back out. + let mut attrs = HashMap::new(); + attrs.insert("zone".to_string(), "z0".to_string()); + let series_key = format_series_key( + canonical_sketch_metric_name("http_latency_ms", SketchKind::DdSketch), + &attrs, + ); + + // ── Window 1: full frame. Base buckets [10, 0, 5]. ── + let full_w1 = DdSketchState { + alpha: 0.01, + store_counts: vec![10, 0, 5], + store_offset: 0, + } + .encode_to_vec(); + route_modified_otlp_sketches_to_precompute( + &build_request( + "http_latency_ms", + make_dp(WIN1_START, 11_000_000, 1, full_w1), + ), + &state, + ) + .await; + + // ── Window 1: delta (SAME window_start). Adds +3 to bucket 0, + // +7 to bucket 1. Within the window this accumulates onto the + // full frame → [13, 7, 5]. ── + let delta_w1 = PbDelta { + buckets: vec![ + DdSketchBucketDelta { + index: 0, + d_count: 3, + }, + DdSketchBucketDelta { + index: 1, + d_count: 7, + }, + ], + } + .encode_to_vec(); + route_modified_otlp_sketches_to_precompute( + &build_request( + "http_latency_ms", + make_dp(WIN1_START, 12_000_000, 2, delta_w1), + ), + &state, + ) + .await; + + { + let entry = state + .sketch_snapshots + .get(&series_key) + .expect("base cached after full + delta in window 1"); + let dd = entry + .core + .as_any() + .downcast_ref::() + .expect("DDSketch base"); + assert_eq!( + dd.inner.store_counts, + vec![13, 7, 5], + "within window 1 the delta accumulates onto the full frame" + ); + assert_eq!( + entry.window_start, WIN1_START, + "cached window_start tracks window 1" + ); + } + + // ── Window 2: delta with a NEW window_start. Adds +20 to bucket + // 2. With per-window base rotation the base is reset to empty + // BEFORE this delta is applied, so the reconstructed state is + // window 2 ONLY: count 20 — NOT window1 + window2 (count 45). + let delta_w2 = PbDelta { + buckets: vec![DdSketchBucketDelta { + index: 2, + d_count: 20, + }], + } + .encode_to_vec(); + route_modified_otlp_sketches_to_precompute( + &build_request( + "http_latency_ms", + make_dp(WIN2_START, 21_000_000, 2, delta_w2), + ), + &state, + ) + .await; + + { + let entry = state + .sketch_snapshots + .get(&series_key) + .expect("base still cached after window 2 delta"); + let dd = entry + .core + .as_any() + .downcast_ref::() + .expect("DDSketch base"); + // The base was rotated to empty before the window-2 delta, so + // it holds window 2 ONLY: total count 20 (the +20 on bucket 2), + // NOT window1 + window2 (which would be 13 + 7 + 5 + 20 = 45). + // Asserting on `total_count` keeps the check independent of the + // empty-sketch's store offset/layout (a fresh sketch re-bases + // its store offset around the first touched bucket). + assert_eq!( + dd.inner.total_count(), + 20, + "new window_start rotates the base to empty: state == window 2 only \ + (count 20), NOT the all-time accumulation (45)" + ); + // Only bucket 2 carries mass; the window-1 buckets (0 and 1) + // were dropped by the rotation. + let bucket_count = |abs_idx: i32| -> u64 { + let i = abs_idx - dd.inner.store_offset; + if i >= 0 && (i as usize) < dd.inner.store_counts.len() { + dd.inner.store_counts[i as usize] + } else { + 0 + } + }; + assert_eq!(bucket_count(2), 20, "window 2's +20 lands on bucket 2"); + assert_eq!( + bucket_count(0), + 0, + "window 1's bucket 0 mass was rotated away" + ); + assert_eq!( + bucket_count(1), + 0, + "window 1's bucket 1 mass was rotated away" + ); + assert_eq!( + entry.window_start, WIN2_START, + "cached window_start advanced to window 2" + ); + } + + drop(state); + let _ = drain.await; + } + #[tokio::test] async fn second_emit_with_cached_sid_and_no_attrs_hits_same_instance() { // Round-trip: emit DP with attrs → cache the assignment → diff --git a/data_plane/src/precompute_engine/ingest_handler.rs b/data_plane/src/precompute_engine/ingest_handler.rs index 5a4097dd2..7217414ab 100644 --- a/data_plane/src/precompute_engine/ingest_handler.rs +++ b/data_plane/src/precompute_engine/ingest_handler.rs @@ -4,6 +4,25 @@ use crate::precompute_engine::worker::parse_labels_from_series_key; use asap_types::aggregation_config::AggregationConfig; use std::sync::Arc; +/// One per-series entry in the delta-reconstitution snapshot cache. +/// +/// Carries the reconstructed accumulator base **plus** the start of the +/// tumbling window that base belongs to. The ingest path uses +/// `window_start` to drive per-window base rotation: when a delta frame +/// arrives whose data-point window start differs from the cached +/// `window_start`, the cached `core` is reset to empty before the new +/// window's delta is applied, so the reconstructed state reflects that +/// window only rather than an all-time accumulation across windows (see +/// `docs/delta-baseline-contract.md` §3). +pub struct SnapshotCacheEntry { + /// Reconstructed per-series accumulator base. + pub core: Box, + /// `start_time_unix_nano` of the window this base was built for. + /// Full frames set it from their own data point; delta frames + /// compare against it to detect a window boundary. + pub window_start: u64, +} + /// Shared state for the ingest path. /// /// Holds the worker router plus the aggregation configs needed for group-key @@ -43,7 +62,11 @@ pub struct IngestState { /// follow-up will add TTL-based eviction keyed by last-seen /// timestamp so long-running deployments don't leak memory /// on retired series. - pub sketch_snapshots: dashmap::DashMap>, + /// + /// The value is a [`SnapshotCacheEntry`] — the reconstructed base + /// plus the window start it belongs to — so the delta-apply path can + /// rotate (reset) the base at a per-series window boundary. + pub sketch_snapshots: dashmap::DashMap, /// Phase 4 — centralized series_id resolver. Shared across the OTLP /// receive path (sid resolution + `unknown_series_ids` population) and /// the `ResolveSeriesIDs` RPC (eager batch resolution from the agent's @@ -186,9 +209,13 @@ mod tests { let base = DDSketchAccumulator { inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0), }; - state - .sketch_snapshots - .insert(series_key.to_string(), Box::new(base.clone())); + state.sketch_snapshots.insert( + series_key.to_string(), + SnapshotCacheEntry { + core: Box::new(base.clone()), + window_start: 0, + }, + ); // First delta adds to bucket 0 and bucket 2. The wire delta now // carries only bucket deltas (the count/sum/min/max scalar fields @@ -210,12 +237,17 @@ mod tests { .sketch_snapshots .get(series_key) .unwrap() + .core .clone_boxed_core(); apply_modified_otlp_delta_bytes(SketchKind::DdSketch, ENCODING_PROTO_DELTA, &mut acc1, &d1) .expect("apply first delta"); - state - .sketch_snapshots - .insert(series_key.to_string(), acc1.clone_boxed_core()); + state.sketch_snapshots.insert( + series_key.to_string(), + SnapshotCacheEntry { + core: acc1.clone_boxed_core(), + window_start: 0, + }, + ); // Second delta — picks up on top of the first, proving the // cache refresh is transitive. @@ -230,6 +262,7 @@ mod tests { .sketch_snapshots .get(series_key) .unwrap() + .core .clone_boxed_core(); apply_modified_otlp_delta_bytes(SketchKind::DdSketch, ENCODING_PROTO_DELTA, &mut acc2, &d2) .expect("apply second delta"); diff --git a/data_plane/src/precompute_engine/operators/count_min_sketch_accumulator.rs b/data_plane/src/precompute_engine/operators/count_min_sketch_accumulator.rs index 692f29b6f..d83294857 100644 --- a/data_plane/src/precompute_engine/operators/count_min_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/count_min_sketch_accumulator.rs @@ -393,6 +393,13 @@ impl AggregateCore for CountMinSketchAccumulator { "CountMinSketchAccumulator" } + /// Per-window base rotation: rebuild an empty counter matrix with + /// the same (rows, cols) so the next window's additive cell deltas + /// align to the identical hash geometry. + fn reset_to_empty(&mut self) { + self.inner = CountMinSketch::new(self.inner.rows(), self.inner.cols()); + } + fn as_any(&self) -> &dyn std::any::Any { self } diff --git a/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs b/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs index cca852c0b..e1ad89811 100644 --- a/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs @@ -232,6 +232,13 @@ impl AggregateCore for CountSketchAccumulator { "CountSketchAccumulator" } + /// Per-window base rotation: rebuild an empty signed-counter matrix + /// with the same (rows, cols) so the next window's additive cell + /// deltas align to the identical hash geometry. + fn reset_to_empty(&mut self) { + self.inner = CountSketch::new(self.inner.rows, self.inner.cols); + } + fn as_any(&self) -> &dyn std::any::Any { self } diff --git a/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs b/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs index 3044c446b..a28b03972 100644 --- a/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs @@ -155,6 +155,13 @@ impl AggregateCore for DDSketchAccumulator { "DDSketchAccumulator" } + /// Per-window base rotation: drop all bucket counts but keep the + /// relative-accuracy parameter so the next window's bucket deltas + /// index into the same log-bucket layout. + fn reset_to_empty(&mut self) { + self.inner = DdSketch::new(self.inner.alpha); + } + fn as_any(&self) -> &dyn std::any::Any { self } diff --git a/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs b/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs index f35b44fda..d6d2a30dc 100644 --- a/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs @@ -165,6 +165,15 @@ impl AggregateCore for HllSketchAccumulator { "HllSketchAccumulator" } + /// Per-window base rotation: zero the registers but keep the variant + /// and precision. Critical for HLL — its register-wise `max` merge + /// has no inverse, so a never-reset base accumulates the all-time-max + /// across windows (`docs/delta-baseline-contract.md` §1.5); rotating + /// to an empty register array makes per-window cardinality correct. + fn reset_to_empty(&mut self) { + self.inner = HllSketch::new(self.inner.variant, self.inner.precision); + } + fn as_any(&self) -> &dyn std::any::Any { self } diff --git a/data_plane/src/storage_engines/types/traits.rs b/data_plane/src/storage_engines/types/traits.rs index 6aea12344..846c45b24 100644 --- a/data_plane/src/storage_engines/types/traits.rs +++ b/data_plane/src/storage_engines/types/traits.rs @@ -83,6 +83,24 @@ pub trait AggregateCore: SerializableToSink + Send + Sync { fn aux_stats(&self) -> AuxStats { AuxStats::empty() } + + /// Reset the sketch state to empty **in place**, preserving its + /// shape / configuration (dimensions, relative accuracy, register + /// width, …) so a subsequent delta-apply lands on a clean, + /// same-shape base. + /// + /// Used by the OTLP ingest path's per-window base rotation: when a + /// delta frame opens a new tumbling window for a series, the cached + /// base is reset here before the new window's delta is applied, so + /// the reconstructed state reflects that window only rather than an + /// all-time accumulation across windows (see + /// `docs/delta-baseline-contract.md` §3). + /// + /// The default is a no-op: only the delta-capable, additive families + /// (DDSketch, CMS, CountSketch, HLL) ever reach the rotation path and + /// override this. KLL never deltas, and the non-sketch accumulators + /// are never cached as a delta base. + fn reset_to_empty(&mut self) {} } /// Four typed auxiliary scalars tracked alongside every sketch entry: