Skip to main content

hive_console_sdk/agent/
buffer.rs

1use std::collections::VecDeque;
2
3use tokio::sync::Mutex;
4
5pub struct Buffer<T> {
6    max_size: usize,
7    queue: Mutex<VecDeque<T>>,
8}
9
10pub enum AddStatus<T> {
11    Full { drained: Vec<T> },
12    Ok,
13}
14
15impl<T> Buffer<T> {
16    pub fn new(max_size: usize) -> Self {
17        Self {
18            queue: Mutex::new(VecDeque::with_capacity(max_size)),
19            max_size,
20        }
21    }
22
23    pub async fn add(&self, item: T) -> AddStatus<T> {
24        let mut queue = self.queue.lock().await;
25        if queue.len() >= self.max_size {
26            let mut drained: Vec<T> = queue.drain(..).collect();
27            drained.push(item);
28            AddStatus::Full { drained }
29        } else {
30            queue.push_back(item);
31            AddStatus::Ok
32        }
33    }
34
35    pub async fn drain(&self) -> Vec<T> {
36        let mut queue = self.queue.lock().await;
37        queue.drain(..).collect()
38    }
39}