hive_console_sdk/agent/
buffer.rs1use 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}