vtcode_commons/interjection/
events.rs1use 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}