ddsketch/canonical/store/
collapsing_lowest.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 lowest-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 lowest-indexed
9/// bins are collapsed (merged into the next lowest bin), sacrificing accuracy for lower quantiles to preserve accuracy
10/// for higher quantiles.
11///
12/// Use this store when:
13/// - You need bounded memory usage
14/// - Higher quantiles (for example, p95, p99) are more important than lower quantiles
15/// - You're tracking latencies or other metrics where the tail matters most
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct CollapsingLowestDenseStore {
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 low quantiles).
31    is_collapsed: bool,
32}
33
34impl CollapsingLowestDenseStore {
35    /// Creates an empty `CollapsingLowestDenseStore` 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 lower 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 falls
57    /// below `offset`, in which case it belongs to the collapsed lowest 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 {
66            // Need to prepend bins - but first check if we need to collapse
67            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                // The index sits below the widest window the store can represent. Widen the window downwards to the
72                // cap first, so the collapsed counts land on the lowest index that's still representable, then let
73                // `bin_index` fold `index` into `bins[0]`.
74                //
75                // Widening matters for accuracy: without it the collapsed floor is wherever the window happened to
76                // already start, which makes the result depend on insertion order. Descending input would pin every
77                // observation to the first index seen instead of spreading it across the available bins. The
78                // reference implementation places the floor at `max_index - max_num_bins + 1`, which is what widening
79                // to the cap achieves here.
80                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            // Need to append bins
98            let new_len = (index - self.offset + 1) as usize;
99
100            if new_len > self.max_num_bins {
101                // The window can't stretch far enough to cover `index`, so slide its bottom up to the lowest index
102                // that keeps `index` in range, collapsing everything below that point.
103                self.collapse_below(index - self.max_num_bins as i32 + 1);
104            }
105
106            // Now append
107            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    /// Collapses every bin below `new_offset` into the bin at `new_offset`, making it the new lowest bin.
129    ///
130    /// Taking the new lowest index rather than a number of bins to shift by is what keeps this correct when
131    /// `new_offset` lands above the window entirely: that case folds the whole store into a single bin, instead of
132    /// clamping the shift and leaving the window short of where it needs to be.
133    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            // Every bin sits below the new lowest index, so the whole store folds into one bin.
148            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    /// Returns the index into the bins array for the given logical index.
161    ///
162    /// Indices below the window map to the lowest bin, which is where collapsed counts accumulate. [`Self::grow`]
163    /// guarantees no index arrives here from above the window, and that the store holds at least one bin.
164    #[inline]
165    fn bin_index(&self, index: i32) -> usize {
166        if index < self.offset {
167            // Index is below our range, map to lowest bin
168            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        // Clamping keeps the bins consistent with `count` even if that invariant is ever broken: the observation
175        // loses accuracy instead of being dropped outright.
176        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        // Process each bin from the other store
248        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        // Adjust count since add() already incremented it
256        // We need to subtract the other.count we added and use the saturating add
257        // Actually, the adds already handle this correctly
258    }
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        // Process sparse binCounts
273        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        // Process contiguous bins
281        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        // Use contiguous encoding for dense store
301        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    /// Creates a collapsing lowest dense store with a default of 2048 bins.
310    fn default() -> Self {
311        Self::new(2048)
312    }
313}
314
315#[cfg(test)]
316mod tests {
317    use super::*;
318
319    // Shared `Store` trait conformance suite. The default 2048-bin capacity is far larger than the small index
320    // ranges these cases use, so no collapsing occurs and the standard `Store` contract holds.
321    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        // Add bins 0-4
340        for i in 0..5 {
341            store.add(i, 1);
342        }
343        assert!(!store.is_collapsed());
344
345        // Adding index 5 should trigger collapse
346        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        // Add bins 5-9
358        for i in 5..10 {
359            store.add(i, 1);
360        }
361        assert!(!store.is_collapsed());
362
363        // Adding index 0 should trigger collapse since it would need 10 bins
364        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); // This should trigger collapse
378
379        assert!(store.is_collapsed());
380        assert_eq!(store.total_count(), 4);
381
382        // All counts should still be accounted for
383        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}