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_idx(&self, idx: usize) -> Option<&(K, V)> {
69            self.items.get(idx - self.popped)
70        }
71
72        pub fn get_mut_or_insert(&mut self, key: K, default: V) -> (&mut V, bool) {
73            let hash = make_hash(&self.hash_builder, &key);
74            if let Some(idx) = self
75                .table
76                .find(hash, |other| self.items[other - self.popped].0 == key)
77            {
78                return (&mut self.items[*idx - self.popped].1, false);
79            }
80            self.insert_nocheck(hash, key, default);
81
82            #[allow(clippy::unwrap_used)]
83            (&mut self.items.back_mut().unwrap().1, true)
84        }
85
86        // Insert a new item at the back if the queue if it doesn't yet exist.
87        //
88        // If the key already exists, replace the previous value
89        pub fn insert(&mut self, key: K, value: V) -> (usize, bool) {
90            let hash = make_hash(&self.hash_builder, &key);
91            if let Some(idx) = self
92                .table
93                .find(hash, |other| self.items[other - self.popped].0 == key)
94            {
95                self.items[*idx - self.popped].1 = value;
96                (*idx, false)
97            } else {
98                (self.insert_nocheck(hash, key, value), true)
99            }
100        }
101
102        /// # Safety
103        ///
104        /// This function inserts a new item in the store unconditionally
105        /// If the item already exists, it's drop implementation will not be called, and memory
106        /// might leak
107        ///
108        /// The hash needs to be precomputed too
109        fn insert_nocheck(&mut self, hash: u64, key: K, value: V) -> usize {
110            let item_index = self.items.len() + self.popped;
111
112            // Separate set and items since set is mutably borrowed, while items is immutable
113            let Self {
114                table,
115                items,
116                popped,
117                hash_builder,
118                ..
119            } = self;
120            table.insert_unique(hash, item_index, |i| {
121                make_hash(hash_builder, &items[i - *popped].0)
122            });
123            self.items.push_back((key, value));
124            item_index
125        }
126    }
127
128    impl<K, V> Default for QueueHashMap<K, V> {
129        fn default() -> Self {
130            Self {
131                table: HashTable::new(),
132                hash_builder: DefaultHashBuilder::default(),
133                items: VecDeque::new(),
134                popped: 0,
135            }
136        }
137    }
138
139    fn make_hash<T: Hash>(h: &DefaultHashBuilder, i: &T) -> u64 {
140        h.hash_one(i)
141    }
142}
143
144pub use queuehashmap::QueueHashMap;
145
146#[derive(Debug, Default)]
147/// Stores telemetry data item, like dependencies and integrations
148///
149/// * Bounds the length of the collection it uses to prevent memory leaks
150/// * Tries to keep a list of items that it has seen (within max number of items)
151/// * Tries to keep a list of items that haven't been sent to datadog yet
152/// * Deduplicates items, to make sure we don't send the item twice
153pub struct Store<T> {
154    // unflushed and set contain indices into
155    unflushed: VecDeque<usize>,
156    items: QueueHashMap<T, ()>,
157    max_items: usize,
158}
159
160impl<T> Store<T>
161where
162    T: PartialEq + Eq + Hash,
163{
164    pub fn new(max_items: usize) -> Self {
165        Self {
166            unflushed: VecDeque::new(),
167            items: QueueHashMap::default(),
168            max_items,
169        }
170    }
171
172    pub fn insert(&mut self, item: T) {
173        if self.items.get(&item).is_some() {
174            return;
175        }
176        if self.items.len() == self.max_items {
177            self.items.pop_front();
178        }
179        let (idx, _) = self.items.insert(item, ());
180        if self.unflushed.len() == self.max_items {
181            self.unflushed.pop_front();
182        }
183        self.unflushed.push_back(idx);
184    }
185
186    // Reinsert all already flushed items in the flush queue
187    pub fn unflush_stored(&mut self) {
188        self.unflushed.clear();
189        for i in self.items.iter_idx() {
190            self.unflushed.push_back(i);
191        }
192    }
193
194    // Remove the first `count` items in the queue
195    pub fn removed_flushed(&mut self, count: usize) {
196        for _ in 0..count {
197            self.unflushed.pop_front();
198        }
199    }
200
201    pub fn flush_not_empty(&self) -> bool {
202        !self.unflushed.is_empty()
203    }
204
205    pub fn unflushed(&self) -> impl Iterator<Item = &T> {
206        self.unflushed
207            .iter()
208            .flat_map(|i| Some(&self.items.get_idx(*i)?.0))
209    }
210
211    pub fn len_unflushed(&self) -> usize {
212        self.unflushed.len()
213    }
214
215    pub fn len_stored(&self) -> usize {
216        self.items.len()
217    }
218
219    /// Discard pending unflushed items and clear stored dedupe history.
220    pub fn clear(&mut self) {
221        self.unflushed.clear();
222        self.items.clear();
223    }
224}
225
226impl<T> Extend<T> for Store<T>
227where
228    T: PartialEq + Eq + Hash,
229{
230    fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
231        for i in iter {
232            self.insert(i)
233        }
234    }
235}
236
237#[cfg(test)]
238mod tests {
239    use super::*;
240
241    #[test]
242    fn test_smoke_insert() {
243        let mut store = Store::new(10);
244        store.insert("hello");
245        store.insert("world");
246        store.insert("world");
247
248        assert_eq!(store.unflushed.len(), 2);
249        assert_eq!(store.items.len(), 2);
250        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"hello", &"world"]);
251
252        store.removed_flushed(1);
253        assert_eq!(store.items.len(), 2);
254        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"world"]);
255
256        store.removed_flushed(1);
257        assert_eq!(store.items.len(), 2);
258        assert!(store.unflushed().next().is_none());
259
260        store.insert("hello");
261        assert!(store.unflushed().next().is_none());
262    }
263
264    #[test]
265    fn test_insert_spill() {
266        let mut store = Store::new(5);
267        for i in 2..15 {
268            store.insert(i);
269        }
270        assert_eq!(store.unflushed.len(), 5);
271        assert_eq!(store.items.len(), 5);
272
273        assert_eq!(
274            store.unflushed().collect::<Vec<_>>(),
275            &[&10, &11, &12, &13, &14]
276        )
277    }
278
279    #[test]
280    fn test_insert_spill_no_unflush() {
281        let mut store = Store::new(5);
282        for i in 2..7 {
283            store.insert(i);
284        }
285        assert_eq!(store.unflushed.len(), 5);
286
287        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
288        store.removed_flushed(4);
289
290        for i in 7..10 {
291            store.insert(i);
292        }
293
294        assert_eq!(store.unflushed.len(), 4);
295        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&6, &7, &8, &9]);
296    }
297
298    #[test]
299    fn test_unflush_stored() {
300        let mut store = Store::new(5);
301        for i in 2..7 {
302            store.insert(i);
303        }
304        assert_eq!(store.unflushed.len(), 5);
305
306        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
307        store.unflush_stored();
308        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
309    }
310}