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
27 changes: 15 additions & 12 deletions crates/asap_otel_proto/proto/sketchlib_delta/ddsketch_delta.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 0 additions & 4 deletions data_plane/benches/sketch_db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,10 +52,6 @@ fn encode_ddsketch(values: &[f64], alpha: f64) -> Vec<u8> {
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)),
Expand Down
4 changes: 0 additions & 4 deletions data_plane/examples/sketch_db_diag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,10 +38,6 @@ fn ddsketch_payload() -> Vec<u8> {
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)),
Expand Down
15 changes: 5 additions & 10 deletions data_plane/src/drivers/ingest/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2580,9 +2580,11 @@ mod dispatcher_tests {

// Base sketch represents the last full snapshot the agent sent.
let mut acc: Box<dyn AggregateCore> = 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 {
Expand All @@ -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();

Expand All @@ -2613,9 +2609,8 @@ mod dispatcher_tests {

let dd = acc.as_any().downcast_ref::<DDSketchAccumulator>().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]
Expand Down
26 changes: 8 additions & 18 deletions data_plane/src/precompute_engine/ingest_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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;
Expand Down
Loading