Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 7 additions & 13 deletions crates/asap_otel_proto/proto/sketchlib_delta/hll_delta.proto
Original file line number Diff line number Diff line change
@@ -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;
}
8 changes: 3 additions & 5 deletions data_plane/src/drivers/ingest/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn AggregateCore> =
Expand All @@ -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();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -121,19 +121,11 @@ impl HllSketchAccumulator {
&mut self,
buffer: &[u8],
) -> Result<(), Box<dyn std::error::Error>> {
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(())
}
Expand Down Expand Up @@ -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();

Expand Down
20 changes: 5 additions & 15 deletions data_plane/src/storage_engines/sketch_db/query/delta_apply.rs
Original file line number Diff line number Diff line change
Expand Up @@ -412,22 +412,12 @@ fn hll_from_proto(buffer: &[u8]) -> Result<HllSketch, String> {
))
}

/// 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(())
}