ddsketch/canonical/store/
mod.rs

1//! Sketch storage.
2
3use datadog_protos::sketches::Store as ProtoStore;
4
5use super::error::ProtoConversionError;
6
7mod collapsing_highest;
8pub use self::collapsing_highest::CollapsingHighestDenseStore;
9
10mod collapsing_lowest;
11pub use self::collapsing_lowest::CollapsingLowestDenseStore;
12
13mod dense;
14pub use self::dense::DenseStore;
15
16mod sparse;
17pub use self::sparse::SparseStore;
18
19#[cfg(test)]
20mod reference_tests;
21
22/// Storage for sketch observations.
23///
24/// Stores manage holding the counts of mapped values, such that they contain a list of bins and the number of
25/// observations currently counted in each bin.
26pub trait Store: Clone + Send + Sync {
27    /// Adds a count to the bin at the given index.
28    fn add(&mut self, index: i32, count: u64);
29
30    /// Returns the total count across all bins.
31    fn total_count(&self) -> u64;
32
33    /// Returns the minimum index with a non-zero count, or `None` if empty.
34    fn min_index(&self) -> Option<i32>;
35
36    /// Returns the maximum index with a non-zero count, or `None` if empty.
37    fn max_index(&self) -> Option<i32>;
38
39    /// Returns the index of the bin containing the given rank.
40    ///
41    /// The rank is 0-indexed, so rank 0 is the first observation.
42    fn key_at_rank(&self, rank: u64) -> Option<i32>;
43
44    /// Merges another store into this one.
45    fn merge(&mut self, other: &Self);
46
47    /// Returns `true` if the store is empty.
48    fn is_empty(&self) -> bool;
49
50    /// Clears all bins from the store.
51    fn clear(&mut self);
52
53    /// Populates this store from a protobuf `Store`.
54    fn merge_from_proto(&mut self, proto: &ProtoStore) -> Result<(), ProtoConversionError>;
55
56    /// Converts this store to a protobuf `Store`.
57    fn to_proto(&self) -> ProtoStore;
58}
59
60/// Validates and converts a protobuf `f64` count to `u64`.
61///
62/// # Errors
63///
64/// If the count is negative, or has a fractional part, an error is returned.
65pub(crate) fn validate_proto_count(index: i32, count: f64) -> Result<u64, ProtoConversionError> {
66    if count < 0.0 {
67        return Err(ProtoConversionError::NegativeBinCount { index, count });
68    }
69    if count.fract() != 0.0 {
70        return Err(ProtoConversionError::NonIntegerBinCount { index, count });
71    }
72    Ok(count as u64)
73}
74
75/// Generates the shared [`Store`] trait conformance suite for a concrete store type.
76///
77/// Every `Store` implementation shares the same observable contract for adding counts, reporting
78/// totals/min/max, ranking, merging, clearing, and round-tripping through protobuf. Rather than hand-duplicate that
79/// suite across each sibling implementation, each implementation's `tests` module invokes this macro with its
80/// concrete type. Implementation-specific behavior (collapsing, private-field layout, the dense-only bin iterator)
81/// is still covered by inline tests in the respective module.
82///
83/// The store type must implement `Default`. The collapsing stores default to 2048 bins, which is far larger than
84/// any index range these cases use, so no collapsing occurs and the shared assertions hold for every
85/// implementation.
86#[cfg(test)]
87macro_rules! store_conformance_tests {
88    ($store:ty) => {
89        #[test]
90        fn add_single_value_sets_count_and_bounds() {
91            let mut store = <$store>::default();
92            store.add(5, 1);
93
94            assert_eq!(store.total_count(), 1);
95            assert_eq!(store.min_index(), Some(5));
96            assert_eq!(store.max_index(), Some(5));
97        }
98
99        #[test]
100        fn add_repeated_at_same_index_accumulates() {
101            let mut store = <$store>::default();
102            store.add(5, 3);
103            store.add(5, 2);
104
105            assert_eq!(store.total_count(), 5);
106            assert_eq!(store.min_index(), Some(5));
107            assert_eq!(store.max_index(), Some(5));
108        }
109
110        #[test]
111        fn add_distinct_indices_tracks_distribution_and_bounds() {
112            let mut store = <$store>::default();
113            store.add(5, 1);
114            store.add(10, 2);
115            store.add(3, 3);
116
117            assert_eq!(store.total_count(), 6);
118            assert_eq!(store.min_index(), Some(3));
119            assert_eq!(store.max_index(), Some(10));
120            assert_eq!(
121                $crate::canonical::store::conformance::distribution(&store),
122                [(3, 3), (5, 1), (10, 2)].into_iter().collect()
123            );
124        }
125
126        #[test]
127        fn add_with_zero_count_is_a_noop() {
128            let mut store = <$store>::default();
129            store.add(5, 0);
130
131            assert!(store.is_empty());
132            assert_eq!(store.total_count(), 0);
133            assert_eq!(store.min_index(), None);
134        }
135
136        #[test]
137        fn key_at_rank_maps_each_rank_to_its_index() {
138            let mut store = <$store>::default();
139            store.add(5, 3);
140            store.add(10, 2);
141
142            assert_eq!(store.key_at_rank(0), Some(5));
143            assert_eq!(store.key_at_rank(2), Some(5));
144            assert_eq!(store.key_at_rank(3), Some(10));
145            assert_eq!(store.key_at_rank(4), Some(10));
146            assert_eq!(store.key_at_rank(5), None);
147        }
148
149        #[test]
150        fn empty_store_reports_no_min_max_or_rank() {
151            let store = <$store>::default();
152
153            assert!(store.is_empty());
154            assert_eq!(store.total_count(), 0);
155            assert_eq!(store.min_index(), None);
156            assert_eq!(store.max_index(), None);
157            assert_eq!(store.key_at_rank(0), None);
158        }
159
160        #[test]
161        fn merge_combines_distributions() {
162            let mut store1 = <$store>::default();
163            store1.add(5, 2);
164            store1.add(10, 1);
165
166            let mut store2 = <$store>::default();
167            store2.add(5, 1);
168            store2.add(15, 3);
169
170            store1.merge(&store2);
171
172            assert_eq!(store1.total_count(), 7);
173            assert_eq!(store1.min_index(), Some(5));
174            assert_eq!(store1.max_index(), Some(15));
175            assert_eq!(
176                $crate::canonical::store::conformance::distribution(&store1),
177                [(5, 3), (10, 1), (15, 3)].into_iter().collect()
178            );
179        }
180
181        #[test]
182        fn merge_from_empty_store_is_a_noop() {
183            let mut store = <$store>::default();
184            store.add(5, 2);
185
186            let empty = <$store>::default();
187            store.merge(&empty);
188
189            assert_eq!(store.total_count(), 2);
190            assert_eq!(
191                $crate::canonical::store::conformance::distribution(&store),
192                [(5, 2)].into_iter().collect()
193            );
194        }
195
196        #[test]
197        fn clear_resets_to_empty() {
198            let mut store = <$store>::default();
199            store.add(5, 2);
200            store.add(10, 1);
201
202            store.clear();
203
204            assert!(store.is_empty());
205            assert_eq!(store.total_count(), 0);
206            assert_eq!(store.min_index(), None);
207        }
208
209        #[test]
210        fn handles_negative_indices() {
211            let mut store = <$store>::default();
212            store.add(-5, 1);
213            store.add(5, 1);
214
215            assert_eq!(store.total_count(), 2);
216            assert_eq!(store.min_index(), Some(-5));
217            assert_eq!(store.max_index(), Some(5));
218            assert_eq!(
219                $crate::canonical::store::conformance::distribution(&store),
220                [(-5, 1), (5, 1)].into_iter().collect()
221            );
222        }
223
224        #[test]
225        fn handles_widely_scattered_indices() {
226            let mut store = <$store>::default();
227            store.add(-1000, 1);
228            store.add(0, 2);
229            store.add(1000, 3);
230
231            assert_eq!(store.total_count(), 6);
232            assert_eq!(store.min_index(), Some(-1000));
233            assert_eq!(store.max_index(), Some(1000));
234            assert_eq!(
235                $crate::canonical::store::conformance::distribution(&store),
236                [(-1000, 1), (0, 2), (1000, 3)].into_iter().collect()
237            );
238        }
239
240        #[test]
241        fn proto_round_trip_preserves_distribution() {
242            let mut store = <$store>::default();
243            store.add(-3, 4);
244            store.add(5, 2);
245            store.add(10, 1);
246
247            let proto = store.to_proto();
248            let mut restored = <$store>::default();
249            restored
250                .merge_from_proto(&proto)
251                .expect("round-tripped proto must be valid");
252
253            assert_eq!(restored.total_count(), store.total_count());
254            assert_eq!(
255                $crate::canonical::store::conformance::distribution(&restored),
256                $crate::canonical::store::conformance::distribution(&store)
257            );
258        }
259    };
260}
261
262#[cfg(test)]
263pub(crate) use store_conformance_tests;
264
265#[cfg(test)]
266pub(crate) mod conformance {
267    use std::collections::BTreeMap;
268
269    use super::Store;
270
271    /// Reconstructs the full `index -> count` distribution of a store using only the public [`Store`] trait API.
272    ///
273    /// This walks every rank in `[0, total_count)` and tallies the index each rank maps to. Doing this through the
274    /// trait (rather than each implementation's private fields) lets the shared conformance suite assert exact bin
275    /// distributions for every `Store`, including after merges and protobuf round-trips.
276    pub(crate) fn distribution<S: Store>(store: &S) -> BTreeMap<i32, u64> {
277        let mut dist = BTreeMap::new();
278        for rank in 0..store.total_count() {
279            let index = store
280                .key_at_rank(rank)
281                .expect("every rank below total_count must map to an index");
282            *dist.entry(index).or_insert(0) += 1;
283        }
284        dist
285    }
286}