Skip to main content

libdd_telemetry/worker/
store.rs

1// Copyright 2021-Present Datadog, Inc. https://www.datadoghq.com/
2// SPDX-License-Identifier: Apache-2.0
3
4use std::{collections::VecDeque, hash::Hash};
5
6mod queuehashmap {
7    use hashbrown::{hash_table::HashTable, DefaultHashBuilder};
8    use std::{
9        collections::VecDeque,
10        hash::{BuildHasher, Hash},
11    };
12
13    #[derive(Debug)]
14    pub struct QueueHashMap<K, V> {
15        table: HashTable<usize>,
16        hash_builder: DefaultHashBuilder,
17        items: VecDeque<(K, V)>,
18        popped: usize,
19    }
20
21    impl<K, V> QueueHashMap<K, V>
22    where
23        K: PartialEq + Eq + Hash,
24    {
25        pub fn iter(&self) -> impl Iterator<Item = &(K, V)> {
26            self.items.iter()
27        }
28
29        pub fn iter_idx(&self) -> impl Iterator<Item = usize> {
30            self.popped..(self.popped + self.items.len())
31        }
32
33        pub fn len(&self) -> usize {
34            self.items.len()
35        }
36
37        pub fn is_empty(&self) -> bool {
38            self.items.is_empty()
39        }
40
41        /// Clear the map, reusing existing allocations
42        pub fn clear(&mut self) {
43            self.table.clear();
44            self.items.clear();
45            self.popped = 0;
46        }
47
48        // Remove the oldest item in the queue and return it
49        pub fn pop_front(&mut self) -> Option<(K, V)> {
50            let (k, v) = self.items.pop_front()?;
51            let hash = make_hash(&self.hash_builder, &k);
52            if let Ok(entry) = self.table.find_entry(hash, |&other| other == self.popped) {
53                entry.remove();
54            }
55            debug_assert!(self.items.len() == self.table.len());
56            self.popped += 1;
57            Some((k, v))
58        }
59
60        pub fn get(&self, k: &K) -> Option<&V> {
61            let hash = make_hash(&self.hash_builder, k);
62            let idx = self
63                .table
64                .find(hash, |other| &self.items[other - self.popped].0 == k)?;
65            Some(&self.items[*idx - self.popped].1)
66        }
67
68        pub fn get_mut(&mut self, k: &K) -> Option<&mut V> {
69            let hash = make_hash(&self.hash_builder, k);
70            let idx = *self
71                .table
72                .find(hash, |other| &self.items[other - self.popped].0 == k)?;
73            Some(&mut self.items[idx - self.popped].1)
74        }
75
76        pub fn get_idx(&self, idx: usize) -> Option<&(K, V)> {
77            self.items.get(idx - self.popped)
78        }
79
80        pub fn get_mut_or_insert(&mut self, key: K, default: V) -> (&mut V, bool) {
81            let hash = make_hash(&self.hash_builder, &key);
82            if let Some(idx) = self
83                .table
84                .find(hash, |other| self.items[other - self.popped].0 == key)
85            {
86                return (&mut self.items[*idx - self.popped].1, false);
87            }
88            self.insert_nocheck(hash, key, default);
89
90            #[allow(clippy::unwrap_used)]
91            (&mut self.items.back_mut().unwrap().1, true)
92        }
93
94        // Insert a new item at the back if the queue if it doesn't yet exist.
95        //
96        // If the key already exists, replace the previous value
97        pub fn insert(&mut self, key: K, value: V) -> (usize, bool) {
98            let hash = make_hash(&self.hash_builder, &key);
99            if let Some(idx) = self
100                .table
101                .find(hash, |other| self.items[other - self.popped].0 == key)
102            {
103                self.items[*idx - self.popped].1 = value;
104                (*idx, false)
105            } else {
106                (self.insert_nocheck(hash, key, value), true)
107            }
108        }
109
110        /// # Safety
111        ///
112        /// This function inserts a new item in the store unconditionally
113        /// If the item already exists, it's drop implementation will not be called, and memory
114        /// might leak
115        ///
116        /// The hash needs to be precomputed too
117        fn insert_nocheck(&mut self, hash: u64, key: K, value: V) -> usize {
118            let item_index = self.items.len() + self.popped;
119
120            // Separate set and items since set is mutably borrowed, while items is immutable
121            let Self {
122                table,
123                items,
124                popped,
125                hash_builder,
126                ..
127            } = self;
128            table.insert_unique(hash, item_index, |i| {
129                make_hash(hash_builder, &items[i - *popped].0)
130            });
131            self.items.push_back((key, value));
132            item_index
133        }
134    }
135
136    impl<K, V> Default for QueueHashMap<K, V> {
137        fn default() -> Self {
138            Self {
139                table: HashTable::new(),
140                hash_builder: DefaultHashBuilder::default(),
141                items: VecDeque::new(),
142                popped: 0,
143            }
144        }
145    }
146
147    fn make_hash<T: Hash>(h: &DefaultHashBuilder, i: &T) -> u64 {
148        h.hash_one(i)
149    }
150}
151
152pub use queuehashmap::QueueHashMap;
153
154/// Stable key to use as hash for stored items.
155pub trait Keyed<K> {
156    fn key(&self) -> K;
157}
158
159impl<T: Clone> Keyed<T> for T {
160    fn key(&self) -> T {
161        self.clone()
162    }
163}
164
165#[derive(Debug, Default)]
166/// Stores telemetry data item, like dependencies and integrations
167///
168/// * Bounds the length of the collection it uses to prevent memory leaks
169/// * Tries to keep a list of items that it has seen (within max number of items)
170/// * Tries to keep a list of items that haven't been sent to datadog yet
171/// * Deduplicates items, to make sure we don't send the item twice
172pub struct Store<T, K = T> {
173    // unflushed and set contain indices into
174    unflushed: VecDeque<usize>,
175    items: QueueHashMap<K, T>,
176    max_items: usize,
177}
178
179impl<T, K> Store<T, K>
180where
181    T: Keyed<K> + PartialEq,
182    K: PartialEq + Eq + Hash + Clone,
183{
184    pub fn new(max_items: usize) -> Self {
185        Self {
186            unflushed: VecDeque::new(),
187            items: QueueHashMap::default(),
188            max_items,
189        }
190    }
191
192    pub fn insert(&mut self, item: T) {
193        let key = item.key();
194        if let Some(existing) = self.items.get(&key) {
195            if *existing == item {
196                // Exact duplicate: already stored, and already sent or queued.
197                return;
198            }
199            // The entry was refreshed rather than duplicated (e.g. SCA attaching CVE reachability
200            // metadata to a dependency that was already reported). Re-queue it to actually deliver.
201            let (idx, _) = self.items.insert(key, item);
202            if !self.unflushed.contains(&idx) {
203                self.unflushed.push_back(idx);
204            }
205            return;
206        }
207        if self.items.len() == self.max_items {
208            self.items.pop_front();
209        }
210        let (idx, _) = self.items.insert(key, item);
211        if self.unflushed.len() == self.max_items {
212            self.unflushed.pop_front();
213        }
214        self.unflushed.push_back(idx);
215    }
216
217    // Reinsert all already flushed items in the flush queue
218    pub fn unflush_stored(&mut self) {
219        self.unflushed.clear();
220        for i in self.items.iter_idx() {
221            self.unflushed.push_back(i);
222        }
223    }
224
225    // Remove the first `count` items in the queue
226    pub fn removed_flushed(&mut self, count: usize) {
227        for _ in 0..count {
228            self.unflushed.pop_front();
229        }
230    }
231
232    pub fn flush_not_empty(&self) -> bool {
233        !self.unflushed.is_empty()
234    }
235
236    pub fn unflushed(&self) -> impl Iterator<Item = &T> {
237        self.unflushed
238            .iter()
239            .flat_map(|i| Some(&self.items.get_idx(*i)?.1))
240    }
241
242    pub fn len_unflushed(&self) -> usize {
243        self.unflushed.len()
244    }
245
246    pub fn len_stored(&self) -> usize {
247        self.items.len()
248    }
249
250    /// Discard pending unflushed items and clear stored dedupe history.
251    pub fn clear(&mut self) {
252        self.unflushed.clear();
253        self.items.clear();
254    }
255}
256
257impl<T, K> Extend<T> for Store<T, K>
258where
259    T: Keyed<K> + PartialEq,
260    K: PartialEq + Eq + Hash + Clone,
261{
262    fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
263        for i in iter {
264            self.insert(i)
265        }
266    }
267}
268
269#[cfg(test)]
270mod tests {
271    use super::*;
272
273    #[test]
274    fn test_smoke_insert() {
275        let mut store = Store::new(10);
276        store.insert("hello");
277        store.insert("world");
278        store.insert("world");
279
280        assert_eq!(store.unflushed.len(), 2);
281        assert_eq!(store.items.len(), 2);
282        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"hello", &"world"]);
283
284        store.removed_flushed(1);
285        assert_eq!(store.items.len(), 2);
286        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"world"]);
287
288        store.removed_flushed(1);
289        assert_eq!(store.items.len(), 2);
290        assert!(store.unflushed().next().is_none());
291
292        store.insert("hello");
293        assert!(store.unflushed().next().is_none());
294    }
295
296    /// A keyed item whose value can change while its key stays the same, mirroring a Dependency
297    /// that gets SCA reachability metadata attached after it was first reported.
298    #[derive(Debug, PartialEq, Clone)]
299    struct Keyish(&'static str, u32);
300
301    impl Keyed<&'static str> for Keyish {
302        fn key(&self) -> &'static str {
303            self.0
304        }
305    }
306
307    #[test]
308    fn test_updated_item_is_requeued_for_flush() {
309        let mut store: Store<Keyish, &'static str> = Store::new(10);
310        store.insert(Keyish("requests", 0));
311        store.removed_flushed(1);
312        assert!(store.unflushed().next().is_none(), "initial send drains");
313
314        // Same key, unchanged value: deduplicated, nothing to re-send.
315        store.insert(Keyish("requests", 0));
316        assert!(
317            store.unflushed().next().is_none(),
318            "identical re-insert must not re-queue"
319        );
320
321        // Same key, updated value: must be re-queued so the update is actually delivered.
322        store.insert(Keyish("requests", 1));
323        assert_eq!(
324            store.items.len(),
325            1,
326            "update refreshes in place, no duplicate entry"
327        );
328        assert_eq!(
329            store.unflushed().collect::<Vec<_>>(),
330            &[&Keyish("requests", 1)]
331        );
332
333        // Updating again while still queued must not enqueue it twice.
334        store.insert(Keyish("requests", 2));
335        assert_eq!(
336            store.unflushed().collect::<Vec<_>>(),
337            &[&Keyish("requests", 2)]
338        );
339    }
340
341    #[test]
342    fn test_insert_spill() {
343        let mut store = Store::new(5);
344        for i in 2..15 {
345            store.insert(i);
346        }
347        assert_eq!(store.unflushed.len(), 5);
348        assert_eq!(store.items.len(), 5);
349
350        assert_eq!(
351            store.unflushed().collect::<Vec<_>>(),
352            &[&10, &11, &12, &13, &14]
353        )
354    }
355
356    #[test]
357    fn test_insert_spill_no_unflush() {
358        let mut store = Store::new(5);
359        for i in 2..7 {
360            store.insert(i);
361        }
362        assert_eq!(store.unflushed.len(), 5);
363
364        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
365        store.removed_flushed(4);
366
367        for i in 7..10 {
368            store.insert(i);
369        }
370
371        assert_eq!(store.unflushed.len(), 4);
372        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&6, &7, &8, &9]);
373    }
374
375    #[test]
376    fn test_unflush_stored() {
377        let mut store = Store::new(5);
378        for i in 2..7 {
379            store.insert(i);
380        }
381        assert_eq!(store.unflushed.len(), 5);
382
383        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
384        store.unflush_stored();
385        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
386    }
387}