Skip to main content

vtcode_commons/interjection/
events.rs

1use std::sync::{Arc, Mutex, MutexGuard};
2
3#[derive(Debug)]
4pub struct EventQueue<E> {
5    events: Arc<Mutex<Vec<E>>>,
6}
7
8impl<E> Clone for EventQueue<E> {
9    fn clone(&self) -> Self {
10        Self { events: Arc::clone(&self.events) }
11    }
12}
13
14impl<E> Default for EventQueue<E> {
15    fn default() -> Self {
16        Self::new()
17    }
18}
19
20impl<E> EventQueue<E> {
21    pub fn new() -> Self {
22        Self { events: Arc::new(Mutex::new(Vec::new())) }
23    }
24
25    pub fn push(&self, event: E) {
26        self.lock().push(event);
27    }
28
29    pub fn push_capped(&self, event: E, max: usize) {
30        let mut q = self.lock();
31        q.push(event);
32        if q.len() > max {
33            let excess = q.len() - max;
34            q.drain(..excess);
35        }
36    }
37
38    pub fn len(&self) -> usize {
39        self.lock().len()
40    }
41
42    pub fn is_empty(&self) -> bool {
43        self.lock().is_empty()
44    }
45
46    pub fn drain_matching(&self, take: impl Fn(&E) -> bool) -> Vec<E> {
47        let mut q = self.lock();
48        let (matched, kept): (Vec<E>, Vec<E>) = std::mem::take(&mut *q).into_iter().partition(|e| take(e));
49        *q = kept;
50        matched
51    }
52
53    pub fn drain_all(&self) -> Vec<E> {
54        std::mem::take(&mut *self.lock())
55    }
56
57    pub fn clear(&self) {
58        self.lock().clear();
59    }
60
61    fn lock(&self) -> MutexGuard<'_, Vec<E>> {
62        self.events.lock().unwrap_or_else(|e| e.into_inner())
63    }
64}
65
66impl<E: Clone> EventQueue<E> {
67    pub fn snapshot(&self) -> Vec<E> {
68        self.lock().clone()
69    }
70}
71
72#[cfg(test)]
73mod tests {
74    use super::*;
75
76    #[test]
77    fn push_and_len() {
78        let q: EventQueue<u32> = EventQueue::new();
79        assert!(q.is_empty());
80        q.push(1);
81        q.push(2);
82        assert_eq!(q.len(), 2);
83    }
84
85    #[test]
86    fn clones_share_one_queue() {
87        let q: EventQueue<u32> = EventQueue::new();
88        let q2 = q.clone();
89        q.push(7);
90        assert_eq!(q2.len(), 1);
91    }
92
93    #[test]
94    fn push_capped_drops_oldest() {
95        let q: EventQueue<u32> = EventQueue::new();
96        for i in 0..5 {
97            q.push_capped(i, 3);
98        }
99        assert_eq!(q.drain_matching(|_| true), vec![2, 3, 4]);
100        assert!(q.is_empty());
101    }
102
103    #[test]
104    fn drain_matching_returns_matched_retains_rest_fifo() {
105        let q: EventQueue<u32> = EventQueue::new();
106        for i in 0..6 {
107            q.push(i);
108        }
109        let evens = q.drain_matching(|n| n % 2 == 0);
110        assert_eq!(evens, vec![0, 2, 4]);
111        assert_eq!(q.drain_matching(|_| true), vec![1, 3, 5]);
112    }
113
114    #[test]
115    fn push_capped_under_limit_keeps_all() {
116        let q: EventQueue<u32> = EventQueue::new();
117        q.push_capped(1, 5);
118        q.push_capped(2, 5);
119        assert_eq!(q.drain_matching(|_| true), vec![1, 2]);
120    }
121
122    #[test]
123    fn drain_matching_none_match_retains_all() {
124        let q: EventQueue<u32> = EventQueue::new();
125        q.push(1);
126        q.push(2);
127        assert!(q.drain_matching(|n| *n > 10).is_empty());
128        assert_eq!(q.len(), 2);
129    }
130
131    #[test]
132    fn drain_matching_on_empty_is_empty() {
133        let q: EventQueue<u32> = EventQueue::new();
134        assert!(q.drain_matching(|_| true).is_empty());
135    }
136
137    #[test]
138    fn drain_all_empties_in_fifo_order() {
139        let q: EventQueue<u32> = EventQueue::new();
140        q.push(1);
141        q.push(2);
142        assert_eq!(q.drain_all(), vec![1, 2]);
143        assert!(q.is_empty());
144    }
145
146    #[test]
147    fn clear_discards_all() {
148        let q: EventQueue<u32> = EventQueue::new();
149        q.push(1);
150        q.clear();
151        assert!(q.is_empty());
152    }
153
154    #[test]
155    fn snapshot_reads_without_draining() {
156        let q: EventQueue<u32> = EventQueue::new();
157        q.push(9);
158        assert_eq!(q.snapshot(), vec![9]);
159        assert_eq!(q.len(), 1);
160    }
161}