libdd-telemetry 8.0.0

Telemetry client allowing to send data as described in https://docs.datadoghq.com/tracing/configure_data_security/?tab=net#telemetry-collection
Documentation
// Copyright 2021-Present Datadog, Inc. https://www.datadoghq.com/
// SPDX-License-Identifier: Apache-2.0

use std::{collections::VecDeque, hash::Hash};

mod queuehashmap {
    use hashbrown::{hash_table::HashTable, DefaultHashBuilder};
    use std::{
        collections::VecDeque,
        hash::{BuildHasher, Hash},
    };

    #[derive(Debug)]
    pub struct QueueHashMap<K, V> {
        table: HashTable<usize>,
        hash_builder: DefaultHashBuilder,
        items: VecDeque<(K, V)>,
        popped: usize,
    }

    impl<K, V> QueueHashMap<K, V>
    where
        K: PartialEq + Eq + Hash,
    {
        pub fn iter(&self) -> impl Iterator<Item = &(K, V)> {
            self.items.iter()
        }

        pub fn iter_idx(&self) -> impl Iterator<Item = usize> {
            self.popped..(self.popped + self.items.len())
        }

        pub fn len(&self) -> usize {
            self.items.len()
        }

        pub fn is_empty(&self) -> bool {
            self.items.is_empty()
        }

        /// Clear the map, reusing existing allocations
        pub fn clear(&mut self) {
            self.table.clear();
            self.items.clear();
            self.popped = 0;
        }

        // Remove the oldest item in the queue and return it
        pub fn pop_front(&mut self) -> Option<(K, V)> {
            let (k, v) = self.items.pop_front()?;
            let hash = make_hash(&self.hash_builder, &k);
            if let Ok(entry) = self.table.find_entry(hash, |&other| other == self.popped) {
                entry.remove();
            }
            debug_assert!(self.items.len() == self.table.len());
            self.popped += 1;
            Some((k, v))
        }

        pub fn get(&self, k: &K) -> Option<&V> {
            let hash = make_hash(&self.hash_builder, k);
            let idx = self
                .table
                .find(hash, |other| &self.items[other - self.popped].0 == k)?;
            Some(&self.items[*idx - self.popped].1)
        }

        pub fn get_mut(&mut self, k: &K) -> Option<&mut V> {
            let hash = make_hash(&self.hash_builder, k);
            let idx = *self
                .table
                .find(hash, |other| &self.items[other - self.popped].0 == k)?;
            Some(&mut self.items[idx - self.popped].1)
        }

        pub fn get_idx(&self, idx: usize) -> Option<&(K, V)> {
            self.items.get(idx - self.popped)
        }

        pub fn get_mut_or_insert(&mut self, key: K, default: V) -> (&mut V, bool) {
            let hash = make_hash(&self.hash_builder, &key);
            if let Some(idx) = self
                .table
                .find(hash, |other| self.items[other - self.popped].0 == key)
            {
                return (&mut self.items[*idx - self.popped].1, false);
            }
            self.insert_nocheck(hash, key, default);

            #[allow(clippy::unwrap_used)]
            (&mut self.items.back_mut().unwrap().1, true)
        }

        // Insert a new item at the back if the queue if it doesn't yet exist.
        //
        // If the key already exists, replace the previous value
        pub fn insert(&mut self, key: K, value: V) -> (usize, bool) {
            let hash = make_hash(&self.hash_builder, &key);
            if let Some(idx) = self
                .table
                .find(hash, |other| self.items[other - self.popped].0 == key)
            {
                self.items[*idx - self.popped].1 = value;
                (*idx, false)
            } else {
                (self.insert_nocheck(hash, key, value), true)
            }
        }

        /// # Safety
        ///
        /// This function inserts a new item in the store unconditionally
        /// If the item already exists, it's drop implementation will not be called, and memory
        /// might leak
        ///
        /// The hash needs to be precomputed too
        fn insert_nocheck(&mut self, hash: u64, key: K, value: V) -> usize {
            let item_index = self.items.len() + self.popped;

            // Separate set and items since set is mutably borrowed, while items is immutable
            let Self {
                table,
                items,
                popped,
                hash_builder,
                ..
            } = self;
            table.insert_unique(hash, item_index, |i| {
                make_hash(hash_builder, &items[i - *popped].0)
            });
            self.items.push_back((key, value));
            item_index
        }
    }

    impl<K, V> Default for QueueHashMap<K, V> {
        fn default() -> Self {
            Self {
                table: HashTable::new(),
                hash_builder: DefaultHashBuilder::default(),
                items: VecDeque::new(),
                popped: 0,
            }
        }
    }

    fn make_hash<T: Hash>(h: &DefaultHashBuilder, i: &T) -> u64 {
        h.hash_one(i)
    }
}

pub use queuehashmap::QueueHashMap;

/// Stable key to use as hash for stored items.
pub trait Keyed<K> {
    fn key(&self) -> K;
}

impl<T: Clone> Keyed<T> for T {
    fn key(&self) -> T {
        self.clone()
    }
}

#[derive(Debug, Default)]
/// Stores telemetry data item, like dependencies and integrations
///
/// * Bounds the length of the collection it uses to prevent memory leaks
/// * Tries to keep a list of items that it has seen (within max number of items)
/// * Tries to keep a list of items that haven't been sent to datadog yet
/// * Deduplicates items, to make sure we don't send the item twice
pub struct Store<T, K = T> {
    // unflushed and set contain indices into
    unflushed: VecDeque<usize>,
    items: QueueHashMap<K, T>,
    max_items: usize,
}

impl<T, K> Store<T, K>
where
    T: Keyed<K> + PartialEq,
    K: PartialEq + Eq + Hash + Clone,
{
    pub fn new(max_items: usize) -> Self {
        Self {
            unflushed: VecDeque::new(),
            items: QueueHashMap::default(),
            max_items,
        }
    }

    pub fn insert(&mut self, item: T) {
        let key = item.key();
        if let Some(existing) = self.items.get(&key) {
            if *existing == item {
                // Exact duplicate: already stored, and already sent or queued.
                return;
            }
            // The entry was refreshed rather than duplicated (e.g. SCA attaching CVE reachability
            // metadata to a dependency that was already reported). Re-queue it to actually deliver.
            let (idx, _) = self.items.insert(key, item);
            if !self.unflushed.contains(&idx) {
                self.unflushed.push_back(idx);
            }
            return;
        }
        if self.items.len() == self.max_items {
            self.items.pop_front();
        }
        let (idx, _) = self.items.insert(key, item);
        if self.unflushed.len() == self.max_items {
            self.unflushed.pop_front();
        }
        self.unflushed.push_back(idx);
    }

    // Reinsert all already flushed items in the flush queue
    pub fn unflush_stored(&mut self) {
        self.unflushed.clear();
        for i in self.items.iter_idx() {
            self.unflushed.push_back(i);
        }
    }

    // Remove the first `count` items in the queue
    pub fn removed_flushed(&mut self, count: usize) {
        for _ in 0..count {
            self.unflushed.pop_front();
        }
    }

    pub fn flush_not_empty(&self) -> bool {
        !self.unflushed.is_empty()
    }

    pub fn unflushed(&self) -> impl Iterator<Item = &T> {
        self.unflushed
            .iter()
            .flat_map(|i| Some(&self.items.get_idx(*i)?.1))
    }

    pub fn len_unflushed(&self) -> usize {
        self.unflushed.len()
    }

    pub fn len_stored(&self) -> usize {
        self.items.len()
    }

    /// Discard pending unflushed items and clear stored dedupe history.
    pub fn clear(&mut self) {
        self.unflushed.clear();
        self.items.clear();
    }
}

impl<T, K> Extend<T> for Store<T, K>
where
    T: Keyed<K> + PartialEq,
    K: PartialEq + Eq + Hash + Clone,
{
    fn extend<I: IntoIterator<Item = T>>(&mut self, iter: I) {
        for i in iter {
            self.insert(i)
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn test_smoke_insert() {
        let mut store = Store::new(10);
        store.insert("hello");
        store.insert("world");
        store.insert("world");

        assert_eq!(store.unflushed.len(), 2);
        assert_eq!(store.items.len(), 2);
        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"hello", &"world"]);

        store.removed_flushed(1);
        assert_eq!(store.items.len(), 2);
        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&"world"]);

        store.removed_flushed(1);
        assert_eq!(store.items.len(), 2);
        assert!(store.unflushed().next().is_none());

        store.insert("hello");
        assert!(store.unflushed().next().is_none());
    }

    /// A keyed item whose value can change while its key stays the same, mirroring a Dependency
    /// that gets SCA reachability metadata attached after it was first reported.
    #[derive(Debug, PartialEq, Clone)]
    struct Keyish(&'static str, u32);

    impl Keyed<&'static str> for Keyish {
        fn key(&self) -> &'static str {
            self.0
        }
    }

    #[test]
    fn test_updated_item_is_requeued_for_flush() {
        let mut store: Store<Keyish, &'static str> = Store::new(10);
        store.insert(Keyish("requests", 0));
        store.removed_flushed(1);
        assert!(store.unflushed().next().is_none(), "initial send drains");

        // Same key, unchanged value: deduplicated, nothing to re-send.
        store.insert(Keyish("requests", 0));
        assert!(
            store.unflushed().next().is_none(),
            "identical re-insert must not re-queue"
        );

        // Same key, updated value: must be re-queued so the update is actually delivered.
        store.insert(Keyish("requests", 1));
        assert_eq!(
            store.items.len(),
            1,
            "update refreshes in place, no duplicate entry"
        );
        assert_eq!(
            store.unflushed().collect::<Vec<_>>(),
            &[&Keyish("requests", 1)]
        );

        // Updating again while still queued must not enqueue it twice.
        store.insert(Keyish("requests", 2));
        assert_eq!(
            store.unflushed().collect::<Vec<_>>(),
            &[&Keyish("requests", 2)]
        );
    }

    #[test]
    fn test_insert_spill() {
        let mut store = Store::new(5);
        for i in 2..15 {
            store.insert(i);
        }
        assert_eq!(store.unflushed.len(), 5);
        assert_eq!(store.items.len(), 5);

        assert_eq!(
            store.unflushed().collect::<Vec<_>>(),
            &[&10, &11, &12, &13, &14]
        )
    }

    #[test]
    fn test_insert_spill_no_unflush() {
        let mut store = Store::new(5);
        for i in 2..7 {
            store.insert(i);
        }
        assert_eq!(store.unflushed.len(), 5);

        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
        store.removed_flushed(4);

        for i in 7..10 {
            store.insert(i);
        }

        assert_eq!(store.unflushed.len(), 4);
        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&6, &7, &8, &9]);
    }

    #[test]
    fn test_unflush_stored() {
        let mut store = Store::new(5);
        for i in 2..7 {
            store.insert(i);
        }
        assert_eq!(store.unflushed.len(), 5);

        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
        store.unflush_stored();
        assert_eq!(store.unflushed().collect::<Vec<_>>(), &[&2, &3, &4, &5, &6]);
    }
}