ddsketch/canonical/store/
collapsing_highest.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 CollapsingHighestDenseStore {
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 CollapsingHighestDenseStore {
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 + self.bins.len() as i32 {
66 let new_len = (index - self.offset + 1) as usize;
68
69 if new_len > self.max_num_bins {
70 if self.bins.len() < self.max_num_bins {
80 self.bins.resize(self.max_num_bins, 0);
81 }
82
83 self.is_collapsed = true;
84 return;
85 }
86
87 self.bins.resize(new_len, 0);
88 } else if index < self.offset {
89 let new_len = self.bins.len() + (self.offset - index) as usize;
91
92 if new_len > self.max_num_bins {
93 self.collapse_above(index + self.max_num_bins as i32 - 1);
96 }
97
98 let num_prepend = (self.offset - index) as usize;
101 if num_prepend > 0 {
102 let mut new_bins = vec![0u64; num_prepend + self.bins.len()];
103 new_bins[num_prepend..].copy_from_slice(&self.bins);
104 self.bins = new_bins;
105 self.offset = index;
106 }
107 }
108
109 debug_assert!(
110 index >= self.offset,
111 "grow must leave `index` representable: index={}, offset={}, len={}",
112 index,
113 self.offset,
114 self.bins.len()
115 );
116 debug_assert!(
117 self.bins.len() <= self.max_num_bins,
118 "grow must respect the bin cap: len={}, max_num_bins={}",
119 self.bins.len(),
120 self.max_num_bins
121 );
122 }
123
124 fn collapse_above(&mut self, new_ceiling: i32) {
130 if self.bins.is_empty() {
131 return;
132 }
133
134 self.is_collapsed = true;
135
136 let keep = (new_ceiling as i64 - self.offset as i64 + 1).max(0) as usize;
138
139 if keep == 0 {
140 let collapsed_count: u64 = self.bins.iter().sum();
142 self.bins.clear();
143 self.bins.push(collapsed_count);
144 self.offset = new_ceiling;
145 } else if keep < self.bins.len() {
146 let collapsed_count: u64 = self.bins[keep..].iter().sum();
147 self.bins[keep - 1] = self.bins[keep - 1].saturating_add(collapsed_count);
148 self.bins.truncate(keep);
149 }
150 }
151
152 #[inline]
157 fn bin_index(&self, index: i32) -> usize {
158 if index >= self.offset + self.bins.len() as i32 {
159 return self.bins.len() - 1;
161 }
162
163 debug_assert!(index >= self.offset, "grow should have made `index` representable");
164
165 (index.max(self.offset) - self.offset) as usize
168 }
169}
170
171impl Store for CollapsingHighestDenseStore {
172 fn add(&mut self, index: i32, count: u64) {
173 if count == 0 {
174 return;
175 }
176
177 self.grow(index);
178
179 let bin_idx = self.bin_index(index);
180 self.bins[bin_idx] = self.bins[bin_idx].saturating_add(count);
181 self.count = self.count.saturating_add(count);
182 }
183
184 fn total_count(&self) -> u64 {
185 self.count
186 }
187
188 fn min_index(&self) -> Option<i32> {
189 if self.bins.is_empty() {
190 return None;
191 }
192
193 for (i, &count) in self.bins.iter().enumerate() {
194 if count > 0 {
195 return Some(self.offset + i as i32);
196 }
197 }
198 None
199 }
200
201 fn max_index(&self) -> Option<i32> {
202 if self.bins.is_empty() {
203 return None;
204 }
205
206 for (i, &count) in self.bins.iter().enumerate().rev() {
207 if count > 0 {
208 return Some(self.offset + i as i32);
209 }
210 }
211 None
212 }
213
214 fn key_at_rank(&self, rank: u64) -> Option<i32> {
215 if rank >= self.count {
216 return None;
217 }
218
219 let mut cumulative = 0u64;
220 for (i, &count) in self.bins.iter().enumerate() {
221 cumulative += count;
222 if cumulative > rank {
223 return Some(self.offset + i as i32);
224 }
225 }
226 None
227 }
228
229 fn merge(&mut self, other: &Self) {
230 if other.bins.is_empty() {
231 return;
232 }
233
234 if other.is_collapsed {
235 self.is_collapsed = true;
236 }
237
238 for (i, &count) in other.bins.iter().enumerate() {
240 if count > 0 {
241 let index = other.offset + i as i32;
242 self.add(index, count);
243 }
244 }
245 }
246
247 fn is_empty(&self) -> bool {
248 self.count == 0
249 }
250
251 fn clear(&mut self) {
252 self.bins.clear();
253 self.offset = 0;
254 self.count = 0;
255 self.is_collapsed = false;
256 }
257
258 fn merge_from_proto(&mut self, proto: &ProtoStore) -> Result<(), ProtoConversionError> {
259 for (&index, &count) in &proto.binCounts {
260 let count = validate_proto_count(index, count)?;
261 if count > 0 {
262 self.add(index, count);
263 }
264 }
265
266 let offset = proto.contiguousBinIndexOffset;
267 for (i, &count) in proto.contiguousBinCounts.iter().enumerate() {
268 let index = offset + i as i32;
269 let count = validate_proto_count(index, count)?;
270 if count > 0 {
271 self.add(index, count);
272 }
273 }
274
275 Ok(())
276 }
277
278 fn to_proto(&self) -> ProtoStore {
279 let mut proto = ProtoStore::new();
280
281 if self.bins.is_empty() {
282 return proto;
283 }
284
285 proto.contiguousBinIndexOffset = self.offset;
287 proto.contiguousBinCounts = self.bins.iter().map(|&c| c as f64).collect();
288
289 proto
290 }
291}
292
293impl Default for CollapsingHighestDenseStore {
294 fn default() -> Self {
296 Self::new(2048)
297 }
298}
299
300#[cfg(test)]
301mod tests {
302 use super::*;
303
304 crate::canonical::store::store_conformance_tests!(CollapsingHighestDenseStore);
307
308 #[test]
309 fn within_limit_does_not_collapse() {
310 let mut store = CollapsingHighestDenseStore::new(10);
311 for i in 0..10 {
312 store.add(i, 1);
313 }
314
315 assert_eq!(store.total_count(), 10);
316 assert!(!store.is_collapsed());
317 assert_eq!(store.bins.len(), 10);
318 }
319
320 #[test]
321 fn collapse_on_low_index() {
322 let mut store = CollapsingHighestDenseStore::new(5);
323
324 for i in 5..10 {
326 store.add(i, 1);
327 }
328 assert!(!store.is_collapsed());
329
330 store.add(0, 1);
332
333 assert!(store.is_collapsed());
334 assert_eq!(store.total_count(), 6);
335 assert!(store.bins.len() <= 5);
336 }
337
338 #[test]
339 fn collapse_on_high_index() {
340 let mut store = CollapsingHighestDenseStore::new(5);
341
342 for i in 0..5 {
344 store.add(i, 1);
345 }
346 assert!(!store.is_collapsed());
347
348 store.add(10, 1);
350
351 assert!(store.is_collapsed());
352 assert_eq!(store.total_count(), 6);
353 }
354
355 #[test]
356 fn key_at_rank_after_collapse() {
357 let mut store = CollapsingHighestDenseStore::new(3);
358
359 store.add(0, 1);
360 store.add(1, 1);
361 store.add(2, 1);
362 store.add(-1, 1);
364
365 assert!(store.is_collapsed());
366 assert_eq!(store.total_count(), 4);
367
368 assert!(store.key_at_rank(0).is_some());
370 assert!(store.key_at_rank(3).is_some());
371 assert!(store.key_at_rank(4).is_none());
372 }
373
374 #[test]
375 fn merge_respects_collapse() {
376 let mut store1 = CollapsingHighestDenseStore::new(5);
377 store1.add(0, 1);
378
379 let mut store2 = CollapsingHighestDenseStore::new(5);
380 for i in 0..10 {
381 store2.add(i, 1);
382 }
383
384 assert!(store2.is_collapsed());
385
386 store1.merge(&store2);
387
388 assert!(store1.is_collapsed());
389 }
390}