diff --git a/crates/asap_otel_proto/proto/sketchlib_delta/hll_delta.proto b/crates/asap_otel_proto/proto/sketchlib_delta/hll_delta.proto index de314b1c4..6e49f9c86 100644 --- a/crates/asap_otel_proto/proto/sketchlib_delta/hll_delta.proto +++ b/crates/asap_otel_proto/proto/sketchlib_delta/hll_delta.proto @@ -1,23 +1,17 @@ // hll_delta.proto — Delta wire format for HLL. // -// Mirrors `sketchlib-go/proto/hll/hll.proto` HLLDelta + -// HLLRegisterUpdate byte-for-byte. Vendored for the same reason as -// ddsketch_delta.proto — `asap_sketchlib` upstream doesn't yet -// generate delta types. HLL uses max semantics: applying a delta -// is `registers[idx] = max(registers[idx], value)` for each update. +// The increased registers are varint-packed as (index_delta, value) pairs in +// ascending index order (index_delta from the previous index, starting at 0) — +// the same layout the full sparse register state uses. HLL applies a delta with +// max semantics: registers[idx] = max(registers[idx], value). // -// No threshold is meaningful for HLL — any skipped register update -// causes permanent underestimation. +// No threshold is meaningful for HLL — any skipped register update causes +// permanent underestimation. syntax = "proto3"; package sketchlib.v1; message HLLDelta { - repeated HLLRegisterUpdate updates = 1; -} - -message HLLRegisterUpdate { - uint32 index = 1; // register index, 0 … 2^precision-1 - uint32 value = 2; // new register value (only sent when > snapshot) + bytes packed_updates = 1; } diff --git a/data_plane/src/drivers/ingest/otel.rs b/data_plane/src/drivers/ingest/otel.rs index 2ed7af8c7..291c48941 100644 --- a/data_plane/src/drivers/ingest/otel.rs +++ b/data_plane/src/drivers/ingest/otel.rs @@ -2663,7 +2663,7 @@ mod dispatcher_tests { #[test] fn apply_modified_otlp_delta_bytes_hll_round_trip() { - use asap_otel_proto::sketchlib::v1::{HllDelta as PbDelta, HllRegisterUpdate}; + use asap_otel_proto::sketchlib::v1::HllDelta as PbDelta; use prost::Message; let mut acc: Box = @@ -2674,11 +2674,9 @@ mod dispatcher_tests { .inner .registers = vec![1, 5, 3, 7]; + // Packed (index_delta, value) blob for updates {0:4, 2:6}. let bytes = PbDelta { - updates: vec![ - HllRegisterUpdate { index: 0, value: 4 }, - HllRegisterUpdate { index: 2, value: 6 }, - ], + packed_updates: vec![0, 4, 2, 6], } .encode_to_vec(); 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 d83294857..91f7c955b 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 @@ -208,6 +208,12 @@ impl CountMinSketchAccumulator { cells, l1: pb.l1, l2: pb.l2, + // The Go-side CountMinDelta proto now carries an hh_keys field + // (heavy-hitter candidates), mirrored on asap_sketchlib's + // CountMinSketchDelta. The vendored Rust proto bindings here don't + // decode it yet, and CountMin has no TopK to rebuild, so pass an + // empty set — same handling as CountSketch's hh_keys. + hh_keys: Vec::new(), }; self.inner .apply_delta(&delta) @@ -890,7 +896,7 @@ mod tests { let result = CountMinSketchAccumulator::from_sketchlib_proto_bytes(&bytes); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("zero dims")); + assert!(result.unwrap_err().to_string().contains("degenerate dims")); } #[test] 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 e1ad89811..4426a07f0 100644 --- a/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/count_sketch_accumulator.rs @@ -495,7 +495,7 @@ mod tests { let bytes = state.encode_to_vec(); let result = CountSketchAccumulator::from_sketchlib_proto_bytes(&bytes); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("zero dims")); + assert!(result.unwrap_err().to_string().contains("degenerate dims")); } #[test] 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 d6d2a30dc..3561605c1 100644 --- a/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs +++ b/data_plane/src/precompute_engine/operators/hll_sketch_accumulator.rs @@ -12,7 +12,7 @@ //! store round-trip works end-to-end without that richer query surface. use crate::storage_engines::types::{AggregateCore, AggregationType, KeyByLabelValues, SerializableToSink}; -use asap_sketchlib::{HllSketch, HllSketchDelta, HllVariant, MessagePackCodec}; +use asap_sketchlib::{HllSketch, HllVariant, MessagePackCodec}; use serde_json::Value; use std::collections::HashMap; @@ -121,19 +121,11 @@ impl HllSketchAccumulator { &mut self, buffer: &[u8], ) -> Result<(), Box> { - use asap_otel_proto::sketchlib::v1::HllDelta as PbDelta; - use prost::Message; - - let pb = PbDelta::decode(buffer).map_err(|e| format!("decode HLLDelta: {e}"))?; - - let updates = pb - .updates - .into_iter() - .map(|u| (u.index, u.value as u8)) - .collect(); - let delta = HllSketchDelta { updates }; + // The HLLDelta wire format is a varint-packed (index_delta, value) blob; + // decode + apply (register-wise max) via the shared sketch library so + // the unpacking stays a single source of truth. self.inner - .apply_delta(&delta) + .apply_delta_bytes(buffer) .map_err(|e| format!("apply HLLDelta: {e}"))?; Ok(()) } @@ -468,17 +460,16 @@ mod tests { #[test] fn test_apply_proto_delta_bytes_round_trip() { - use asap_otel_proto::sketchlib::v1::{HllDelta as PbDelta, HllRegisterUpdate}; + use asap_otel_proto::sketchlib::v1::HllDelta as PbDelta; use prost::Message; let mut acc = HllSketchAccumulator::new(HllVariant::Regular, 2); acc.inner.registers = vec![1, 5, 3, 7]; + // Packed (index_delta, value) blob for updates {0:4, 2:6}: + // varint(0),varint(4),varint(2),varint(6). let delta_bytes = PbDelta { - updates: vec![ - HllRegisterUpdate { index: 0, value: 4 }, - HllRegisterUpdate { index: 2, value: 6 }, - ], + packed_updates: vec![0, 4, 2, 6], } .encode_to_vec(); 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 3e6dc9e85..b5a0d525f 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 @@ -412,22 +412,12 @@ fn hll_from_proto(buffer: &[u8]) -> Result { )) } -/// Apply a proto-encoded `HllDelta` frame onto the HLL register vector -/// — mirrors `HllSketchAccumulator::apply_proto_delta_bytes` (sparse -/// `(index, value)` updates, `register = max(register, value)`). +/// Apply a proto-encoded `HllDelta` frame onto the HLL register vector — the +/// delta is a varint-packed (index_delta, value) blob; decode + apply +/// (register-wise max) via the shared sketch library so the unpacking stays a +/// single source of truth. fn apply_hll_proto_delta(sk: &mut HllSketch, buffer: &[u8]) -> Result<(), String> { - use asap_otel_proto::sketchlib::v1::HllDelta as PbDelta; - use asap_sketchlib::HllSketchDelta; - use prost::Message; - - let pb = PbDelta::decode(buffer).map_err(|e| format!("decode HLLDelta: {e}"))?; - let updates = pb - .updates - .into_iter() - .map(|u| (u.index, u.value as u8)) - .collect(); - let delta = HllSketchDelta { updates }; - sk.apply_delta(&delta) + sk.apply_delta_bytes(buffer) .map_err(|e| format!("apply HLLDelta: {e}"))?; Ok(()) }