From 73fbf09854a7a2ac6f471bf6741fa1d58e419e17 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 26 May 2026 08:30:38 -0600 Subject: [PATCH] chore(ddsketch): consume DDSketch wire format without metric scalars The DDSketch wire format dropped its DataPoint-level METRIC scalars (count/sum/min/max) in ProjectASAP/sketchlib-go#243 / asap_sketchlib#57. Update the backend to consume the trimmed format: - proto: DDSketchDelta keeps only `buckets` (1); tags 2-7 reserved. - DDSketchAccumulator: - from_sketchlib_proto_bytes / sketch_reducer / delta_apply now call DdSketch::from_raw(alpha, store_counts, store_offset) (3-arg) and derive count from the bucket store via total_count(). - apply_proto_delta_bytes applies bucket deltas only; count is recomputed from the merged buckets. - serialize_to_json / edge_runtime_adapter drop the removed scalars. - query_statistic STRICT policy: - Quantile -> sketch-derived (unchanged). - Count -> derived from buckets (sum of store counts). - Sum / Min / Max -> return the unavailable-statistic error; these move to controller-provisioned exact Sum / MinMax aggregations. Default aux_stats() is empty for DDSketch so the query path falls through to query_statistic and propagates the error gracefully. - tests: assert Quantile + Count still work and that Sum/Min/Max return the unavailable-statistic error (not a panic / 0); fixtures rebuilt from buckets only. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../sketchlib_delta/ddsketch_delta.proto | 27 +-- data_plane/benches/sketch_db.rs | 4 - data_plane/examples/sketch_db_diag.rs | 4 - data_plane/src/drivers/ingest/otel.rs | 15 +- .../src/precompute_engine/ingest_handler.rs | 26 +-- .../operators/dd_sketch_accumulator.rs | 197 +++++++++++------- .../operators/edge_runtime_adapter.rs | 26 ++- data_plane/src/precompute_engine/worker.rs | 11 +- .../sketch_db/query/delta_apply.rs | 8 +- .../sketch_db/query/sketch_reducer.rs | 8 +- .../storage_engines/sketch_db/query/tests.rs | 4 - ...e2e_controller_plans_and_backend_serves.rs | 22 +- .../tests/e2e_modified_otlp_sketch_path.rs | 26 +-- .../edge_runtime_consumes_precompute_rs.rs | 13 +- 14 files changed, 207 insertions(+), 184 deletions(-) diff --git a/crates/asap_otel_proto/proto/sketchlib_delta/ddsketch_delta.proto b/crates/asap_otel_proto/proto/sketchlib_delta/ddsketch_delta.proto index f0ca10241..f2aef92fc 100644 --- a/crates/asap_otel_proto/proto/sketchlib_delta/ddsketch_delta.proto +++ b/crates/asap_otel_proto/proto/sketchlib_delta/ddsketch_delta.proto @@ -7,25 +7,28 @@ // `asap_sketchlib` picks these up upstream this vendoring can be removed. // // See `sketches/DDSketch/delta.go::ApplyDelta` for the canonical -// merge semantics — additive bucket counts, additive count + sum, -// min (can only decrease) + max (can only increase) as lossless -// scalars when changed. +// merge semantics — additive bucket counts only. The DataPoint-level +// METRIC scalars (count/sum/min/max) were dropped from the wire format +// (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57): the total count is +// recoverable by summing the bucket counts and the remaining aggregates +// are carried by controller-provisioned exact aggregations. syntax = "proto3"; package sketchlib.v1; // DDSketchDelta carries only the buckets that changed above threshold. -// count/sum are additive deltas; min/max are lossless scalars transmitted -// only when they changed. +// Bucket counts are additive. The former count/sum (tags 2-3) and +// min/max scalars (tags 4-7) were removed from the wire format +// (ProjectASAP/sketchlib-go#243): the count is recoverable by summing +// the merged bucket counts; Sum/Min/Max move to exact aggregations. message DDSketchDelta { - repeated DDSketchBucketDelta buckets = 1; - int64 d_count = 2; - double d_sum = 3; - double new_min = 4; - double new_max = 5; - bool min_changed = 6; - bool max_changed = 7; + repeated DDSketchBucketDelta buckets = 1; + + // Tags 2-3 previously carried the additive count/sum deltas and tags + // 4-7 the lossless min/max scalars + their changed flags. Dropped from + // the wire format (ProjectASAP/sketchlib-go#243). + reserved 2, 3, 4, 5, 6, 7; } // DDSketchBucketDelta is the delta for one log-scale bucket. diff --git a/data_plane/benches/sketch_db.rs b/data_plane/benches/sketch_db.rs index f2c62e32f..ed495b232 100644 --- a/data_plane/benches/sketch_db.rs +++ b/data_plane/benches/sketch_db.rs @@ -52,10 +52,6 @@ fn encode_ddsketch(values: &[f64], alpha: f64) -> Vec { alpha: sk.alpha, store_counts: sk.store_counts.clone(), store_offset: sk.store_offset, - count: sk.count, - sum: sk.sum, - min: sk.min, - max: sk.max, }; SketchEnvelope { sketch_state: Some(sketch_envelope::SketchState::Ddsketch(state)), diff --git a/data_plane/examples/sketch_db_diag.rs b/data_plane/examples/sketch_db_diag.rs index 46ebe6761..50a718073 100644 --- a/data_plane/examples/sketch_db_diag.rs +++ b/data_plane/examples/sketch_db_diag.rs @@ -38,10 +38,6 @@ fn ddsketch_payload() -> Vec { alpha: sk.alpha, store_counts: sk.store_counts.clone(), store_offset: sk.store_offset, - count: sk.count, - sum: sk.sum, - min: sk.min, - max: sk.max, }; SketchEnvelope { sketch_state: Some(sketch_envelope::SketchState::Ddsketch(state)), diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index cf54ef28f..77dc6390d 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -2580,9 +2580,11 @@ mod dispatcher_tests { // Base sketch represents the last full snapshot the agent sent. let mut acc: Box = Box::new(DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0, 6, 12.0, 1.0, 3.0), + inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0), }); + // The wire delta now carries only bucket deltas (tags 2-7 + // reserved post ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). let bytes = PbDelta { buckets: vec![ DdSketchBucketDelta { @@ -2594,12 +2596,6 @@ mod dispatcher_tests { d_count: 20, }, ], - d_count: 30, - d_sum: 70.0, - new_min: 0.5, - new_max: 5.0, - min_changed: true, - max_changed: true, } .encode_to_vec(); @@ -2613,9 +2609,8 @@ mod dispatcher_tests { let dd = acc.as_any().downcast_ref::().unwrap(); assert_eq!(dd.inner.store_counts, vec![11, 2, 23]); - assert_eq!(dd.inner.count, 36); - assert_eq!(dd.inner.min, 0.5); - assert_eq!(dd.inner.max, 5.0); + // `count` recomputed from the merged buckets: 11 + 2 + 23 = 36. + assert_eq!(dd.inner.total_count(), 36); } #[test] diff --git a/data_plane/src/precompute_engine/ingest_handler.rs b/data_plane/src/precompute_engine/ingest_handler.rs index 51603d166..5a4097dd2 100644 --- a/data_plane/src/precompute_engine/ingest_handler.rs +++ b/data_plane/src/precompute_engine/ingest_handler.rs @@ -184,13 +184,15 @@ mod tests { // so we're not racing any earlier tests. let series_key = "__name__=latency_ms,inst=a"; let base = DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0, 6, 12.0, 1.0, 3.0), + inner: DdSketch::from_raw(0.01, vec![1, 2, 3], 0), }; state .sketch_snapshots .insert(series_key.to_string(), Box::new(base.clone())); - // First delta adds to bucket 0 and bucket 2. + // First delta adds to bucket 0 and bucket 2. The wire delta now + // carries only bucket deltas (the count/sum/min/max scalar fields + // were dropped, ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). let d1 = PbDelta { buckets: vec![ DdSketchBucketDelta { @@ -202,11 +204,6 @@ mod tests { d_count: 20, }, ], - d_count: 30, - d_sum: 70.0, - new_max: 5.0, - max_changed: true, - ..Default::default() } .encode_to_vec(); let mut acc1 = state @@ -227,11 +224,6 @@ mod tests { index: 1, d_count: 5, }], - d_count: 5, - d_sum: 10.0, - new_max: 6.0, - max_changed: true, - ..Default::default() } .encode_to_vec(); let mut acc2 = state @@ -246,12 +238,10 @@ mod tests { // Base [1,2,3] + d1 [+10 on 0, +20 on 2] = [11,2,23]; // + d2 [+5 on 1] = [11,7,23]. assert_eq!(final_dd.inner.store_counts, vec![11, 7, 23]); - // Counts add: 6 + 30 + 5 = 41. - assert_eq!(final_dd.inner.count, 41); - // Sum: 12 + 70 + 10 = 92. - assert_eq!(final_dd.inner.sum, 92.0); - // Max updated to 6.0 via d2's max_changed flag. - assert_eq!(final_dd.inner.max, 6.0); + // `count` recomputed from the merged buckets: 11 + 7 + 23 = 41. + // (sum/min/max were dropped from the wire format, + // ProjectASAP/sketchlib-go#243 / asap_sketchlib#57.) + assert_eq!(final_dd.inner.total_count(), 41); drop(state); let _ = drain.await; 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 f11b8c20c..3044c446b 100644 --- a/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/dd_sketch_accumulator.rs @@ -6,10 +6,13 @@ //! MessagePack for the sink, and decode from the sketchlib //! `DDSketchState` proto. //! -//! Query semantics (quantile estimation via log-bucket indices) are -//! intentionally deferred — the wire format carries the bucket counts, -//! offset, and aggregates losslessly, so the merge + store round-trip -//! works end-to-end without that richer query surface. +//! Query semantics follow the STRICT policy after the DataPoint-level +//! METRIC scalars were dropped from the wire format +//! (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57): the sketch serves +//! Quantile (log-bucket estimation) and Count (sum of bucket counts). +//! Sum/Min/Max are no longer derivable from the wire bytes and are +//! served by controller-provisioned exact aggregations — `query_statistic` +//! returns the unavailable-statistic error for them. use crate::storage_engines::types::{AggregateCore, AggregationType, KeyByLabelValues, SerializableToSink}; use asap_sketchlib::{DdSketch, DdSketchDelta, MessagePackCodec}; @@ -80,15 +83,13 @@ impl DDSketchAccumulator { ) .into()); } - let inner = DdSketch::from_raw( - state.alpha, - state.store_counts.clone(), - state.store_offset, - state.count, - state.sum, - state.min, - state.max, - ); + // The DataPoint-level METRIC scalars (count/sum/min/max) were + // dropped from `DDSketchState` (ProjectASAP/sketchlib-go#243 / + // asap_sketchlib#57). Reconstruct from the bucket store only: + // `DdSketch::from_raw` now takes just (alpha, store_counts, + // store_offset) and recovers `count` by summing the bucket + // counts via `total_count()`. + let inner = DdSketch::from_raw(state.alpha, state.store_counts.clone(), state.store_offset); Ok(Self { inner }) } @@ -109,6 +110,10 @@ impl DDSketchAccumulator { let pb = PbDelta::decode(buffer).map_err(|e| format!("decode DDSketchDelta: {e}"))?; + // The delta no longer carries d_count/d_sum/min/max + // (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). Apply the + // bucket deltas only; `DdSketch` recomputes its total count from + // the merged bucket counts (`total_count()`). let buckets = pb .buckets .into_iter() @@ -116,12 +121,7 @@ impl DDSketchAccumulator { .collect(); let delta = DdSketchDelta { buckets, - d_count: pb.d_count, - d_sum: pb.d_sum, - min_changed: pb.min_changed, - new_min: pb.new_min, - max_changed: pb.max_changed, - new_max: pb.new_max, + ..Default::default() }; self.inner.apply_delta(&delta); Ok(()) @@ -130,14 +130,14 @@ impl DDSketchAccumulator { impl SerializableToSink for DDSketchAccumulator { fn serialize_to_json(&self) -> Value { + // The DataPoint-level scalars (sum/min/max) are no longer carried + // by `DdSketch` (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). + // `count` is the bucket-derived total via `total_count()`. serde_json::json!({ "alpha": self.inner.alpha, "store_offset": self.inner.store_offset, "bucket_count": self.inner.store_counts.len(), - "count": self.inner.count, - "sum": self.inner.sum, - "min": self.inner.min, - "max": self.inner.max, + "count": self.inner.total_count(), }) } @@ -219,14 +219,30 @@ impl AggregateCore for DDSketchAccumulator { "DDSketchAccumulator: quantile() returned None (sketch empty?)".into() }) } - Statistic::Sum => Ok(self.inner.sum), - Statistic::Count => Ok(self.inner.count as f64), - Statistic::Min => Ok(self.inner.min), - Statistic::Max => Ok(self.inner.max), + // Count is exact, derived by summing the bucket store + // counts — the only DataPoint-level scalar that survives + // the wire-format trim (ProjectASAP/sketchlib-go#243 / + // asap_sketchlib#57). + Statistic::Count => Ok(self.inner.total_count() as f64), + // STRICT policy: the Sum/Min/Max scalars were removed from + // the DDSketch wire format. They are now served by the + // controller-provisioned exact aggregations (an exact `Sum` + // and an exact `MinMax`), NOT estimated from the buckets. + // Surface the unavailable-statistic error so the query path + // routes to those aggregations instead of returning a wrong + // (0 / panicked) value. + Statistic::Sum => Err( + "DDSketchAccumulator: Sum not available from DDSketch wire format \ + (ProjectASAP/sketchlib-go#243); use an exact Sum aggregation" + .into(), + ), + Statistic::Min | Statistic::Max => Err(format!( + "DDSketchAccumulator: {statistic:?} not available from DDSketch wire format \ + (ProjectASAP/sketchlib-go#243); use an exact MinMax aggregation", + ) + .into()), other => Err(format!( - "DDSketchAccumulator: statistic {:?} not supported (only Quantile / Sum / \ - Count / Min / Max)", - other, + "DDSketchAccumulator: statistic {other:?} not supported (only Quantile / Count)", ) .into()), } @@ -237,45 +253,35 @@ impl AggregateCore for DDSketchAccumulator { mod tests { use super::*; - fn encode_state( - alpha: f64, - store_counts: Vec, - store_offset: i32, - count: u64, - sum: f64, - min: f64, - max: f64, - ) -> Vec { + // The DataPoint-level METRIC scalars (count/sum/min/max) were dropped + // from `DdSketchState` (ProjectASAP/sketchlib-go#243 / + // asap_sketchlib#57); the proto now carries only + // `alpha`/`store_counts`/`store_offset`. + fn encode_state(alpha: f64, store_counts: Vec, store_offset: i32) -> Vec { use asap_sketchlib::proto::sketchlib::DdSketchState; use prost::Message; let state = DdSketchState { alpha, store_counts, store_offset, - count, - sum, - min, - max, }; state.encode_to_vec() } #[test] fn test_from_sketchlib_proto_bytes_round_trip() { - let bytes = encode_state(0.01, vec![1, 2, 3, 4], -2, 10, 50.0, 1.0, 4.0); + let bytes = encode_state(0.01, vec![1, 2, 3, 4], -2); let acc = DDSketchAccumulator::from_sketchlib_proto_bytes(&bytes).expect("decode ok"); assert_eq!(acc.inner.alpha, 0.01); assert_eq!(acc.inner.store_counts, vec![1, 2, 3, 4]); assert_eq!(acc.inner.store_offset, -2); - assert_eq!(acc.inner.count, 10); - assert_eq!(acc.inner.sum, 50.0); - assert_eq!(acc.inner.min, 1.0); - assert_eq!(acc.inner.max, 4.0); + // `count` is recovered by summing the bucket store counts. + assert_eq!(acc.inner.total_count(), 10); } #[test] fn test_from_sketchlib_proto_bytes_rejects_invalid_alpha() { - let bytes = encode_state(0.0, vec![1], 0, 1, 1.0, 1.0, 1.0); + let bytes = encode_state(0.0, vec![1], 0); let result = DDSketchAccumulator::from_sketchlib_proto_bytes(&bytes); assert!(result.is_err()); assert!(result.unwrap_err().to_string().contains("alpha")); @@ -293,10 +299,6 @@ mod tests { alpha: 0.01, store_counts: vec![1, 2, 3, 4], store_offset: -2, - count: 10, - sum: 50.0, - min: 1.0, - max: 4.0, }; let env = SketchEnvelope { sketch_state: Some(sketch_envelope::SketchState::Ddsketch(state)), @@ -307,7 +309,7 @@ mod tests { let acc = DDSketchAccumulator::from_sketchlib_proto_bytes(&bytes) .expect("envelope-wrapped decode should succeed"); assert_eq!(acc.inner.alpha, 0.01); - assert_eq!(acc.inner.count, 10); + assert_eq!(acc.inner.total_count(), 10); } #[test] @@ -328,10 +330,10 @@ mod tests { #[test] fn test_aggregate_core_merge_aligns_buckets() { let a = DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![1, 1, 1], -1, 3, 3.0, 1.0, 3.0), + inner: DdSketch::from_raw(0.01, vec![1, 1, 1], -1), }; let b = DDSketchAccumulator { - inner: DdSketch::from_raw(0.01, vec![10, 10, 10], 0, 30, 30.0, 1.0, 3.0), + inner: DdSketch::from_raw(0.01, vec![10, 10, 10], 0), }; let merged_box = a.merge_with(&b).expect("merge ok"); let merged = merged_box @@ -340,7 +342,7 @@ mod tests { .expect("downcast ok"); assert_eq!(merged.inner.store_counts, vec![1, 11, 11, 10]); assert_eq!(merged.inner.store_offset, -1); - assert_eq!(merged.inner.count, 33); + assert_eq!(merged.inner.total_count(), 33); } #[test] @@ -353,14 +355,14 @@ mod tests { #[test] fn test_from_msgpack_bytes_round_trip() { - let original = DdSketch::from_raw(0.01, vec![5, 10, 15, 20], -2, 50, 150.0, 0.25, 8.0); + let original = DdSketch::from_raw(0.01, vec![5, 10, 15, 20], -2); let bytes = original.to_msgpack().unwrap(); let acc = DDSketchAccumulator::from_msgpack_bytes(&bytes).expect("decode ok"); assert_eq!(acc.inner.alpha, 0.01); assert_eq!(acc.inner.store_counts, vec![5, 10, 15, 20]); assert_eq!(acc.inner.store_offset, -2); - assert_eq!(acc.inner.count, 50); - assert_eq!(acc.inner.sum, 150.0); + // `count` is recovered by summing the bucket store counts. + assert_eq!(acc.inner.total_count(), 50); } #[test] @@ -375,8 +377,10 @@ mod tests { use prost::Message; let mut acc = DDSketchAccumulator::new(0.01); - acc.inner = DdSketch::from_raw(0.01, vec![1, 2, 3], 0, 6, 12.0, 1.0, 3.0); + acc.inner = DdSketch::from_raw(0.01, vec![1, 2, 3], 0); + // The wire delta now carries only bucket deltas (tags 2-7 + // reserved); `DdSketchBucketDelta` has just `index` + `d_count`. let bytes = PbDelta { buckets: vec![ DdSketchBucketDelta { @@ -388,21 +392,13 @@ mod tests { d_count: 20, }, ], - d_count: 30, - d_sum: 70.0, - new_min: 0.5, - new_max: 5.0, - min_changed: true, - max_changed: true, } .encode_to_vec(); acc.apply_proto_delta_bytes(&bytes).expect("apply ok"); assert_eq!(acc.inner.store_counts, vec![11, 2, 23]); - assert_eq!(acc.inner.count, 36); - assert_eq!(acc.inner.sum, 82.0); - assert_eq!(acc.inner.min, 0.5); - assert_eq!(acc.inner.max, 5.0); + // `count` recomputed from the merged buckets: 11 + 2 + 23 = 36. + assert_eq!(acc.inner.total_count(), 36); } #[test] @@ -410,4 +406,63 @@ mod tests { let mut acc = DDSketchAccumulator::new(0.01); assert!(acc.apply_proto_delta_bytes(b"not valid proto").is_err()); } + + // ----- query_statistic STRICT policy ----- + // + // After the DataPoint-level METRIC scalars were dropped from the + // DDSketch wire format (ProjectASAP/sketchlib-go#243 / + // asap_sketchlib#57), DDSketch serves only quantiles and Count. + // Sum/Min/Max move to controller-provisioned exact aggregations and + // MUST surface the unavailable-statistic error (never a panic / 0). + + fn sample_accumulator() -> DDSketchAccumulator { + // Build the in-memory sketch from bucket counts only — no scalars. + DDSketchAccumulator { + inner: DdSketch::from_raw(0.01, vec![1, 2, 3, 4], -2), + } + } + + #[test] + fn test_query_statistic_quantile_is_sketch_derived() { + use promql_utilities::query_logics::enums::Statistic; + let acc = sample_accumulator(); + let mut kwargs = HashMap::new(); + kwargs.insert("quantile".to_string(), "0.5".to_string()); + let v = acc + .query_statistic(Statistic::Quantile, &None, &kwargs) + .expect("quantile should be served from the sketch buckets"); + assert!( + v.is_finite() && v > 0.0, + "quantile estimate should be positive finite, got {v}" + ); + } + + #[test] + fn test_query_statistic_count_is_bucket_derived() { + use promql_utilities::query_logics::enums::Statistic; + let acc = sample_accumulator(); + let v = acc + .query_statistic(Statistic::Count, &None, &HashMap::new()) + .expect("count should be derivable from the bucket store"); + // 1 + 2 + 3 + 4 = 10. + assert_eq!(v, 10.0); + } + + #[test] + fn test_query_statistic_sum_min_max_return_unavailable_error() { + use promql_utilities::query_logics::enums::Statistic; + let acc = sample_accumulator(); + for stat in [Statistic::Sum, Statistic::Min, Statistic::Max] { + let result = acc.query_statistic(stat, &None, &HashMap::new()); + assert!( + result.is_err(), + "{stat:?} must return the unavailable-statistic error (not a panic / 0)" + ); + let msg = result.unwrap_err().to_string(); + assert!( + msg.contains("not available"), + "{stat:?} error should explain the statistic is unavailable, got: {msg}" + ); + } + } } diff --git a/data_plane/src/precompute_engine/operators/edge_runtime_adapter.rs b/data_plane/src/precompute_engine/operators/edge_runtime_adapter.rs index e10d113d0..d55b4b825 100644 --- a/data_plane/src/precompute_engine/operators/edge_runtime_adapter.rs +++ b/data_plane/src/precompute_engine/operators/edge_runtime_adapter.rs @@ -215,7 +215,7 @@ pub fn snapshot_ddsketch_via_runtime( // We can't move-construct the wrapper from a non-empty `DdSketch`, // but `Sketch::apply_delta` against the existing snapshot bytes // is equivalent. - if sk.count > 0 { + if sk.total_count() > 0 { // Re-encode the source's state into the canonical envelope // shape that asap-precompute-rs's wrapper recognizes, then // round-trip through `apply_delta`. Mirrors the agent runtime's @@ -249,14 +249,9 @@ pub fn encode_ddsketch_envelope(sk: &asap_sketchlib::DdSketch) -> Vec { alpha: sk.wire_alpha(), store_counts: sk.store_counts.clone(), store_offset: sk.store_offset, - count: sk.count, - sum: sk.sum, - min: if sk.count == 0 { f64::INFINITY } else { sk.min }, - max: if sk.count == 0 { - f64::NEG_INFINITY - } else { - sk.max - }, + // The DataPoint-level scalars (count/sum/min/max) were dropped from + // `DDSketchState` (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57); + // the bucket counts carry all reconstructable state. }; let env = ProtoEnvelope { format_version: 1, @@ -291,14 +286,14 @@ pub fn merge_ddsketches_via_runtime( .into()); } let mut wrapper_a = DDSketchWrapper::new(a.alpha); - if a.count > 0 { + if a.total_count() > 0 { let bridge = encode_ddsketch_envelope(a); wrapper_a .apply_delta(&bridge) .map_err(|e| format!("merge_ddsketches_via_runtime/a: {e}"))?; } let mut wrapper_b = DDSketchWrapper::new(b.alpha); - if b.count > 0 { + if b.total_count() > 0 { let bridge = encode_ddsketch_envelope(b); wrapper_b .apply_delta(&bridge) @@ -329,7 +324,10 @@ mod tests { let state = unwrap_envelope_state(&bytes).expect("decode ok"); match state { Some(SketchState::Ddsketch(s)) => { - assert!(s.count > 0); + // `count` was dropped from `DdSketchState` + // (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57); + // a non-empty sketch carries it in the bucket store. + assert!(s.store_counts.iter().sum::() > 0); assert!(s.alpha > 0.0 && s.alpha < 1.0); } other => panic!("expected DDSketch state, got {other:?}"), @@ -354,7 +352,7 @@ mod tests { ReconstructedSketch::DdSketch(d) => d, ReconstructedSketch::Kll { .. } => panic!("got KLL, expected DDSketch"), }; - assert_eq!(dd.count, 100); + assert_eq!(dd.total_count(), 100); let re_encoded = encode_ddsketch_envelope(&dd); assert_eq!( re_encoded, original_bytes, @@ -387,6 +385,6 @@ mod tests { _ => panic!(), }; let merged = merge_ddsketches_via_runtime(&a_inner, &b_inner).expect("merge ok"); - assert_eq!(merged.count, 20); + assert_eq!(merged.total_count(), 20); } } diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 670347cab..2d6dea3ed 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -2378,7 +2378,8 @@ aggregations: .downcast_ref::() .expect("must downcast back to DDSketchAccumulator"); assert_eq!( - dd.inner.count, 30, + dd.inner.total_count(), + 30, "all 10 first-batch sketches must merge into the persisted output (3 values × 10)" ); } @@ -2470,12 +2471,13 @@ aggregations: .as_ref() .and_then(|k| k.labels.first().cloned()) .unwrap_or_default(); - let expected_count = match zone.as_str() { + let expected_count: u64 = match zone.as_str() { "us-east" => 3, "us-west" => 2, other => panic!("unexpected zone {other}")}; assert_eq!( - dd.inner.count, expected_count, + dd.inner.total_count(), + expected_count, "zone {zone} must roll up exactly {expected_count} per-tuple sketches" ); } @@ -2619,7 +2621,8 @@ aggregations: .downcast_ref::() .expect("must downcast back to DDSketchAccumulator"); assert_eq!( - dd.inner.count, 10, + dd.inner.total_count(), + 10, "all 10 frozen-time sketches must merge into the single emitted output" ); diff --git a/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs b/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs index c1f7de834..3e6dc9e85 100644 --- a/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs +++ b/data_plane/src/storage_engines/sketch_db/query/delta_apply.rs @@ -323,14 +323,14 @@ fn dd_from_proto(buffer: &[u8]) -> Result { state.alpha )); } + // The DataPoint-level scalars (count/sum/min/max) were dropped from + // `DDSketchState` (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57); + // `DdSketch::from_raw` takes only (alpha, store_counts, store_offset) + // and recovers `count` from the bucket store. Ok(DdSketch::from_raw( state.alpha, state.store_counts.clone(), state.store_offset, - state.count, - state.sum, - state.min, - state.max, )) } diff --git a/data_plane/src/storage_engines/sketch_db/query/sketch_reducer.rs b/data_plane/src/storage_engines/sketch_db/query/sketch_reducer.rs index 6593c9747..c6407e66c 100644 --- a/data_plane/src/storage_engines/sketch_db/query/sketch_reducer.rs +++ b/data_plane/src/storage_engines/sketch_db/query/sketch_reducer.rs @@ -1266,14 +1266,14 @@ fn DdSketch_from_sketchlib_proto_bytes(buffer: &[u8]) -> Result Vec { alpha: sk.alpha, store_counts: sk.store_counts.clone(), store_offset: sk.store_offset, - count: sk.count, - sum: sk.sum, - min: sk.min, - max: sk.max, }; let env = SketchEnvelope { sketch_state: Some(sketch_envelope::SketchState::Ddsketch(state)), diff --git a/data_plane/tests/e2e_controller_plans_and_backend_serves.rs b/data_plane/tests/e2e_controller_plans_and_backend_serves.rs index b5d79475c..9bfb1de70 100644 --- a/data_plane/tests/e2e_controller_plans_and_backend_serves.rs +++ b/data_plane/tests/e2e_controller_plans_and_backend_serves.rs @@ -337,24 +337,14 @@ async fn start_full_stack(otlp_http_port: u16, otlp_grpc_port: u16) -> FullStack } } -/// Build a `DdSketchState` proto from raw values. -fn build_dd_sketch_state( - alpha: f64, - store_counts: Vec, - store_offset: i32, - count: u64, - sum: f64, - min: f64, - max: f64, -) -> DdSketchState { +/// Build a `DdSketchState` proto from raw values. The DataPoint-level +/// scalars (count/sum/min/max) were dropped from the wire format +/// (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). +fn build_dd_sketch_state(alpha: f64, store_counts: Vec, store_offset: i32) -> DdSketchState { DdSketchState { alpha, store_counts, store_offset, - count, - sum, - min, - max, } } @@ -923,7 +913,7 @@ async fn controller_plan_to_query_full_roundtrip_ddsketch() { // test uses so we know it's representable. let alpha = 0.01; let store_counts = vec![5u64, 10, 15, 20]; - let dd_state = build_dd_sketch_state(alpha, store_counts, -1, 50, 150.0, 0.25, 8.0); + let dd_state = build_dd_sketch_state(alpha, store_counts, -1); let sketch_bytes = dd_state.encode_to_vec(); // ── 3. POST the sketch DP via OTLP HTTP ──────────────────────────── @@ -962,7 +952,7 @@ async fn controller_plan_to_query_full_roundtrip_ddsketch() { // 1-second tumbling window W. This DP at `now - 1s` is at least // 2 seconds past the start of W, so it moves the engine's // watermark past W's close boundary and triggers the flush. - let watermark_state = build_dd_sketch_state(alpha, Vec::new(), 0, 0, 0.0, 0.0, 0.0); + let watermark_state = build_dd_sketch_state(alpha, Vec::new(), 0); let watermark_req = build_dd_sketch_export( "http_latency_ms", &[("service", "e2e-test")], diff --git a/data_plane/tests/e2e_modified_otlp_sketch_path.rs b/data_plane/tests/e2e_modified_otlp_sketch_path.rs index 1030b733b..c29644068 100644 --- a/data_plane/tests/e2e_modified_otlp_sketch_path.rs +++ b/data_plane/tests/e2e_modified_otlp_sketch_path.rs @@ -751,23 +751,13 @@ fn make_dd_sketch_agg_config( ) } -fn build_dd_sketch_state( - alpha: f64, - store_counts: Vec, - store_offset: i32, - count: u64, - sum: f64, - min: f64, - max: f64, -) -> DdSketchState { +fn build_dd_sketch_state(alpha: f64, store_counts: Vec, store_offset: i32) -> DdSketchState { + // The DataPoint-level scalars (count/sum/min/max) were dropped from + // `DDSketchState` (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). DdSketchState { alpha, store_counts, store_offset, - count, - sum, - min, - max, } } @@ -861,14 +851,14 @@ async fn e2e_dd_sketch_modified_otlp_path() { tokio::time::sleep(tokio::time::Duration::from_millis(400)).await; let store_counts = vec![5u64, 10, 15, 20]; - let dd_state = build_dd_sketch_state(alpha, store_counts.clone(), -1, 50, 150.0, 0.25, 8.0); + let dd_state = build_dd_sketch_state(alpha, store_counts.clone(), -1); let sketch_bytes = dd_state.encode_to_vec(); let client = reqwest::Client::new(); let req = build_dd_sketch_export_request(metric_name, service_label, 100_000_000, sketch_bytes); post_otlp_http(&client, otlp_http_port, req).await; - let watermark_state = build_dd_sketch_state(alpha, Vec::new(), 0, 0, 0.0, 0.0, 0.0); + let watermark_state = build_dd_sketch_state(alpha, Vec::new(), 0); let watermark_req = build_dd_sketch_export_request( metric_name, service_label, @@ -896,8 +886,10 @@ async fn e2e_dd_sketch_modified_otlp_path() { assert_eq!(dd_acc.inner.store_counts, store_counts); assert_eq!(dd_acc.inner.store_offset, -1); - assert_eq!(dd_acc.inner.count, 50); - assert_eq!(dd_acc.inner.sum, 150.0); + // `count` is recovered from the bucket store (5 + 10 + 15 + 20 = 50); + // sum/min/max were dropped from the wire format + // (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). + assert_eq!(dd_acc.inner.total_count(), 50); assert!((dd_acc.inner.alpha - alpha).abs() < f64::EPSILON); } diff --git a/data_plane/tests/edge_runtime_consumes_precompute_rs.rs b/data_plane/tests/edge_runtime_consumes_precompute_rs.rs index 9d71194a5..e045cbf11 100644 --- a/data_plane/tests/edge_runtime_consumes_precompute_rs.rs +++ b/data_plane/tests/edge_runtime_consumes_precompute_rs.rs @@ -63,7 +63,13 @@ fn ddsketch_envelope_round_trip_through_backend_adapter() { ReconstructedSketch::DdSketch(d) => d, _ => panic!("expected DDSketch reconstruction"), }; - assert_eq!(dd.count, 200, "count preserved through runtime adapter"); + // `count` is recovered from the bucket store now that the scalar was + // dropped (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57). + assert_eq!( + dd.total_count(), + 200, + "count preserved through runtime adapter" + ); let re_encoded = encode_ddsketch_envelope(&dd); assert_eq!( @@ -91,7 +97,10 @@ fn ddsketch_envelope_structural_assertions() { .expect("state"); match state { SketchState::Ddsketch(s) => { - assert_eq!(s.count, 3, "structural count"); + // `count` was dropped from `DdSketchState` + // (ProjectASAP/sketchlib-go#243 / asap_sketchlib#57); it is + // recovered by summing the bucket store counts. + assert_eq!(s.store_counts.iter().sum::(), 3, "structural count"); assert!( s.alpha > 0.0 && s.alpha < 1.0, "alpha within (0,1): got {}",