ddsketch/canonical/store/
collapsing_lowest.rs1use datadog_protos::sketches::Store as ProtoStore;
2
3use super::{validate_proto_count, Store};
4use crate::canonical::error::ProtoConversionError;
5
6#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct CollapsingLowestDenseStore {
18 bins: Vec<u64>,
20
21 offset: i32,
23
24 max_num_bins: usize,
26
27 count: u64,
29
30 is_collapsed: bool,
32}
33
34impl CollapsingLowestDenseStore {
35 pub fn new(max_num_bins: usize) -> Self {
37 assert!(max_num_bins >= 1, "max_num_bins must be at least 1");
38 Self {
39 bins: Vec::new(),
40 offset: 0,
41 max_num_bins,
42 count: 0,
43 is_collapsed: false,
44 }
45 }
46
47 pub fn is_collapsed(&self) -> bool {
51 self.is_collapsed
52 }
53
54 fn grow(&mut self, index: i32) {
59 if self.bins.is_empty() {
60 self.bins.push(0);
61 self.offset = index;
62 return;
63 }
64
65 if index < self.offset {
66 let num_prepend = (self.offset - index) as usize;
68 let new_len = self.bins.len() + num_prepend;
69
70 if new_len > self.max_num_bins {
71 let widen_by = self.max_num_bins - self.bins.len();
81 if widen_by > 0 {
82 let mut widened = vec![0u64; self.max_num_bins];
83 widened[widen_by..].copy_from_slice(&self.bins);
84 self.bins = widened;
85 self.offset -= widen_by as i32;
86 }
87
88 self.is_collapsed = true;
89 return;
90 }
91
92 let mut new_bins = vec![0u64; new_len];
93 new_bins[num_prepend..].copy_from_slice(&self.bins);
94 self.bins = new_bins;
95 self.offset = index;
96 } else if index >= self.offset + self.bins.len() as i32 {
97 let new_len = (index - self.offset + 1) as usize;
99
100 if new_len > self.max_num_bins {
101 self.collapse_below(index - self.max_num_bins as i32 + 1);
104 }
105
106 let target_len = ((index - self.offset + 1) as usize).min(self.max_num_bins);
108 if target_len > self.bins.len() {
109 self.bins.resize(target_len, 0);
110 }
111 }
112
113 debug_assert!(
114 index < self.offset + self.bins.len() as i32,
115 "grow must leave `index` representable: index={}, offset={}, len={}",
116 index,
117 self.offset,
118 self.bins.len()
119 );
120 debug_assert!(
121 self.bins.len() <= self.max_num_bins,
122 "grow must respect the bin cap: len={}, max_num_bins={}",
123 self.bins.len(),
124 self.max_num_bins
125 );
126 }
127
128 fn collapse_below(&mut self, new_offset: i32) {
134 debug_assert!(
135 new_offset > self.offset,
136 "collapse_below only slides the window upwards"
137 );
138
139 if self.bins.is_empty() {
140 return;
141 }
142
143 self.is_collapsed = true;
144
145 let n = (new_offset - self.offset) as usize;
146 if n >= self.bins.len() {
147 let collapsed_count: u64 = self.bins.iter().sum();
149 self.bins.clear();
150 self.bins.push(collapsed_count);
151 } else {
152 let collapsed_count: u64 = self.bins[..n].iter().sum();
153 self.bins[n] = self.bins[n].saturating_add(collapsed_count);
154 self.bins.drain(..n);
155 }
156
157 self.offset = new_offset;
158 }
159
160 #[inline]
165 fn bin_index(&self, index: i32) -> usize {
166 if index < self.offset {
167 return 0;
169 }
170
171 let idx = (index - self.offset) as usize;
172 debug_assert!(idx < self.bins.len(), "grow should have made `index` representable");
173
174 idx.min(self.bins.len() - 1)
177 }
178}
179
180impl Store for CollapsingLowestDenseStore {
181 fn add(&mut self, index: i32, count: u64) {
182 if count == 0 {
183 return;
184 }
185
186 self.grow(index);
187
188 let bin_idx = self.bin_index(index);
189 self.bins[bin_idx] = self.bins[bin_idx].saturating_add(count);
190 self.count = self.count.saturating_add(count);
191 }
192
193 fn total_count(&self) -> u64 {
194 self.count
195 }
196
197 fn min_index(&self) -> Option<i32> {
198 if self.bins.is_empty() {
199 return None;
200 }
201
202 for (i, &count) in self.bins.iter().enumerate() {
203 if count > 0 {
204 return Some(self.offset + i as i32);
205 }
206 }
207 None
208 }
209
210 fn max_index(&self) -> Option<i32> {
211 if self.bins.is_empty() {
212 return None;
213 }
214
215 for (i, &count) in self.bins.iter().enumerate().rev() {
216 if count > 0 {
217 return Some(self.offset + i as i32);
218 }
219 }
220 None
221 }
222
223 fn key_at_rank(&self, rank: u64) -> Option<i32> {
224 if rank >= self.count {
225 return None;
226 }
227
228 let mut cumulative = 0u64;
229 for (i, &count) in self.bins.iter().enumerate() {
230 cumulative += count;
231 if cumulative > rank {
232 return Some(self.offset + i as i32);
233 }
234 }
235 None
236 }
237
238 fn merge(&mut self, other: &Self) {
239 if other.bins.is_empty() {
240 return;
241 }
242
243 if other.is_collapsed {
244 self.is_collapsed = true;
245 }
246
247 for (i, &count) in other.bins.iter().enumerate() {
249 if count > 0 {
250 let index = other.offset + i as i32;
251 self.add(index, count);
252 }
253 }
254
255 }
259
260 fn is_empty(&self) -> bool {
261 self.count == 0
262 }
263
264 fn clear(&mut self) {
265 self.bins.clear();
266 self.offset = 0;
267 self.count = 0;
268 self.is_collapsed = false;
269 }
270
271 fn merge_from_proto(&mut self, proto: &ProtoStore) -> Result<(), ProtoConversionError> {
272 for (&index, &count) in &proto.binCounts {
274 let count = validate_proto_count(index, count)?;
275 if count > 0 {
276 self.add(index, count);
277 }
278 }
279
280 let offset = proto.contiguousBinIndexOffset;
282 for (i, &count) in proto.contiguousBinCounts.iter().enumerate() {
283 let index = offset + i as i32;
284 let count = validate_proto_count(index, count)?;
285 if count > 0 {
286 self.add(index, count);
287 }
288 }
289
290 Ok(())
291 }
292
293 fn to_proto(&self) -> ProtoStore {
294 let mut proto = ProtoStore::new();
295
296 if self.bins.is_empty() {
297 return proto;
298 }
299
300 proto.contiguousBinIndexOffset = self.offset;
302 proto.contiguousBinCounts = self.bins.iter().map(|&c| c as f64).collect();
303
304 proto
305 }
306}
307
308impl Default for CollapsingLowestDenseStore {
309 fn default() -> Self {
311 Self::new(2048)
312 }
313}
314
315#[cfg(test)]
316mod tests {
317 use super::*;
318
319 crate::canonical::store::store_conformance_tests!(CollapsingLowestDenseStore);
322
323 #[test]
324 fn within_limit_does_not_collapse() {
325 let mut store = CollapsingLowestDenseStore::new(10);
326 for i in 0..10 {
327 store.add(i, 1);
328 }
329
330 assert_eq!(store.total_count(), 10);
331 assert!(!store.is_collapsed());
332 assert_eq!(store.bins.len(), 10);
333 }
334
335 #[test]
336 fn collapse_on_high_index() {
337 let mut store = CollapsingLowestDenseStore::new(5);
338
339 for i in 0..5 {
341 store.add(i, 1);
342 }
343 assert!(!store.is_collapsed());
344
345 store.add(5, 1);
347
348 assert!(store.is_collapsed());
349 assert_eq!(store.total_count(), 6);
350 assert!(store.bins.len() <= 5);
351 }
352
353 #[test]
354 fn collapse_on_low_index() {
355 let mut store = CollapsingLowestDenseStore::new(5);
356
357 for i in 5..10 {
359 store.add(i, 1);
360 }
361 assert!(!store.is_collapsed());
362
363 store.add(0, 1);
365
366 assert!(store.is_collapsed());
367 assert_eq!(store.total_count(), 6);
368 }
369
370 #[test]
371 fn key_at_rank_after_collapse() {
372 let mut store = CollapsingLowestDenseStore::new(3);
373
374 store.add(0, 1);
375 store.add(1, 1);
376 store.add(2, 1);
377 store.add(3, 1); assert!(store.is_collapsed());
380 assert_eq!(store.total_count(), 4);
381
382 assert!(store.key_at_rank(0).is_some());
384 assert!(store.key_at_rank(3).is_some());
385 assert!(store.key_at_rank(4).is_none());
386 }
387
388 #[test]
389 fn merge_respects_collapse() {
390 let mut store1 = CollapsingLowestDenseStore::new(5);
391 store1.add(0, 1);
392
393 let mut store2 = CollapsingLowestDenseStore::new(5);
394 for i in 0..10 {
395 store2.add(i, 1);
396 }
397
398 assert!(store2.is_collapsed());
399
400 store1.merge(&store2);
401
402 assert!(store1.is_collapsed());
403 }
404}