ddsketch/canonical/store/
collapsing_highest.rs

1use datadog_protos::sketches::Store as ProtoStore;
2
3use super::{validate_proto_count, Store};
4use crate::canonical::error::ProtoConversionError;
5
6/// A dense store that collapses highest-indexed bins when capacity is exceeded.
7///
8/// This store maintains a maximum number of bins. When adding a new index would exceed this limit, the highest-indexed
9/// bins are collapsed (merged into the next highest bin), sacrificing accuracy for higher quantiles to preserve
10/// accuracy for lower quantiles.
11///
12/// Use this store when:
13/// - You need bounded memory usage
14/// - Lower quantiles (for example, p1, p5) are more important than higher quantiles
15/// - You're tracking metrics where the minimum values matter most
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct CollapsingHighestDenseStore {
18    /// The bin counts, stored contiguously.
19    bins: Vec<u64>,
20
21    /// The count stored in bins[0] corresponds to this index.
22    offset: i32,
23
24    /// Maximum number of bins to maintain.
25    max_num_bins: usize,
26
27    /// Total count across all bins.
28    count: u64,
29
30    /// Whether collapsing has occurred (accuracy may be compromised for high quantiles).
31    is_collapsed: bool,
32}
33
34impl CollapsingHighestDenseStore {
35    /// Creates an empty `CollapsingHighestDenseStore` with the given maximum number of bins.
36    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    /// Returns `true` if this store has collapsed bins.
48    ///
49    /// If true, accuracy guarantees may not hold for higher quantiles.
50    pub fn is_collapsed(&self) -> bool {
51        self.is_collapsed
52    }
53
54    /// Ensures the store can accommodate the given index, growing and collapsing if necessary.
55    ///
56    /// On return, `index` is always representable: it either falls inside `[offset, offset + bins.len())`, or it sits
57    /// above the window, in which case it belongs to the collapsed highest bin. [`Self::bin_index`] relies on this.
58    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            // Need to append bins - but first check if we need to collapse
67            let new_len = (index - self.offset + 1) as usize;
68
69            if new_len > self.max_num_bins {
70                // The index sits above the widest window the store can represent. Widen the window upwards to the cap
71                // first, so the collapsed counts land on the highest index that's still representable, then let
72                // `bin_index` fold `index` into the top bin.
73                //
74                // Widening matters for accuracy: without it the collapsed ceiling is wherever the window happened to
75                // already end, which makes the result depend on insertion order. Ascending input would pin every
76                // observation to the first index seen instead of spreading it across the available bins. The reference
77                // implementation places the ceiling at `min_index + max_num_bins - 1`, which is what widening to the
78                // cap achieves here.
79                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            // Need to prepend bins
90            let new_len = self.bins.len() + (self.offset - index) as usize;
91
92            if new_len > self.max_num_bins {
93                // The window can't stretch far enough down to cover `index`, so slide its top down to the highest
94                // index that keeps `index` in range, collapsing everything above that point.
95                self.collapse_above(index + self.max_num_bins as i32 - 1);
96            }
97
98            // Now prepend. `collapse_above` has already freed exactly as many bins as this needs, so the full
99            // distance always fits within the capacity.
100            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    /// Collapses every bin above `new_ceiling` into the bin at `new_ceiling`, making it the new highest bin.
125    ///
126    /// Taking the new highest index rather than a number of bins to shift by is what keeps this correct when
127    /// `new_ceiling` lands below the window entirely: that case folds the whole store into a single bin, instead of
128    /// clamping the shift and leaving the window short of where it needs to be.
129    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        // How many bins, counting up from the bottom of the window, sit at or below the new ceiling.
137        let keep = (new_ceiling as i64 - self.offset as i64 + 1).max(0) as usize;
138
139        if keep == 0 {
140            // Every bin sits above the new highest index, so the whole store folds into one bin.
141            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    /// Returns the index into the bins array for the given logical index.
153    ///
154    /// Indices above the window map to the highest bin, which is where collapsed counts accumulate. [`Self::grow`]
155    /// guarantees no index arrives here from below the window, and that the store holds at least one bin.
156    #[inline]
157    fn bin_index(&self, index: i32) -> usize {
158        if index >= self.offset + self.bins.len() as i32 {
159            // Index is above our range, map to highest bin
160            return self.bins.len() - 1;
161        }
162
163        debug_assert!(index >= self.offset, "grow should have made `index` representable");
164
165        // Clamping keeps the bins consistent with `count` even if that invariant is ever broken: the observation
166        // loses accuracy instead of being dropped outright.
167        (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        // Process each bin from the other store
239        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        // Use contiguous encoding for dense store
286        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    /// Creates a collapsing highest dense store with a default of 2048 bins.
295    fn default() -> Self {
296        Self::new(2048)
297    }
298}
299
300#[cfg(test)]
301mod tests {
302    use super::*;
303
304    // Shared `Store` trait conformance suite. The default 2048-bin capacity is far larger than the small index
305    // ranges these cases use, so no collapsing occurs and the standard `Store` contract holds.
306    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        // Add bins 5-9
325        for i in 5..10 {
326            store.add(i, 1);
327        }
328        assert!(!store.is_collapsed());
329
330        // Adding index 0 should trigger collapse of highest bins
331        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        // Add bins 0-4
343        for i in 0..5 {
344            store.add(i, 1);
345        }
346        assert!(!store.is_collapsed());
347
348        // Adding index 10 should trigger collapse since it would need more than 5 bins
349        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        // Adding a lower index should trigger collapse
363        store.add(-1, 1);
364
365        assert!(store.is_collapsed());
366        assert_eq!(store.total_count(), 4);
367
368        // All counts should still be accounted for
369        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}