agent_base/engine/runtime/
message_queue.rs1use std::collections::VecDeque;
2use std::sync::Mutex;
3
4#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
6pub enum QueueMode {
7 #[default]
9 All,
10 OneAtATime,
12}
13
14pub struct MessageQueue {
24 steering: Mutex<VecDeque<String>>,
25 follow_up: Mutex<VecDeque<String>>,
26 mode: Mutex<QueueMode>,
27}
28
29impl MessageQueue {
30 pub fn new() -> Self {
31 Self {
32 steering: Mutex::new(VecDeque::new()),
33 follow_up: Mutex::new(VecDeque::new()),
34 mode: Mutex::new(QueueMode::default()),
35 }
36 }
37
38 #[allow(dead_code)]
39 pub fn with_mode(mode: QueueMode) -> Self {
40 Self {
41 steering: Mutex::new(VecDeque::new()),
42 follow_up: Mutex::new(VecDeque::new()),
43 mode: Mutex::new(mode),
44 }
45 }
46
47 pub fn set_mode(&self, mode: QueueMode) {
49 *self.mode.lock().unwrap_or_else(|e| e.into_inner()) = mode;
50 }
51
52 #[allow(dead_code)]
54 pub fn mode(&self) -> QueueMode {
55 *self.mode.lock().unwrap_or_else(|e| e.into_inner())
56 }
57
58 pub fn steer(&self, message: String) {
60 self.steering
61 .lock()
62 .unwrap_or_else(|e| e.into_inner())
63 .push_back(message);
64 }
65
66 pub fn follow_up(&self, message: String) {
68 self.follow_up
69 .lock()
70 .unwrap_or_else(|e| e.into_inner())
71 .push_back(message);
72 }
73
74 pub fn drain_steering(&self) -> Vec<String> {
76 let mode = *self.mode.lock().unwrap_or_else(|e| e.into_inner());
77 let mut queue = self.steering.lock().unwrap_or_else(|e| e.into_inner());
78 match mode {
79 QueueMode::All => queue.drain(..).collect(),
80 QueueMode::OneAtATime => queue.pop_front().into_iter().collect(),
81 }
82 }
83
84 pub fn drain_follow_up(&self) -> Vec<String> {
86 let mode = *self.mode.lock().unwrap_or_else(|e| e.into_inner());
87 let mut queue = self.follow_up.lock().unwrap_or_else(|e| e.into_inner());
88 match mode {
89 QueueMode::All => queue.drain(..).collect(),
90 QueueMode::OneAtATime => queue.pop_front().into_iter().collect(),
91 }
92 }
93
94 #[allow(dead_code)]
96 pub fn has_steering(&self) -> bool {
97 !self
98 .steering
99 .lock()
100 .unwrap_or_else(|e| e.into_inner())
101 .is_empty()
102 }
103}
104
105#[cfg(test)]
106mod tests {
107 use super::*;
108
109 #[test]
112 fn new_uses_default_mode() {
113 let mq = MessageQueue::new();
114 assert_eq!(mq.mode(), QueueMode::All);
115 }
116
117 #[test]
118 fn with_mode_sets_correct_mode() {
119 let mq = MessageQueue::with_mode(QueueMode::OneAtATime);
120 assert_eq!(mq.mode(), QueueMode::OneAtATime);
121 }
122
123 #[test]
124 fn set_mode_changes_at_runtime() {
125 let mq = MessageQueue::new();
126 assert_eq!(mq.mode(), QueueMode::All);
127 mq.set_mode(QueueMode::OneAtATime);
128 assert_eq!(mq.mode(), QueueMode::OneAtATime);
129 }
130
131 #[test]
134 fn steer_and_drain_all() {
135 let mq = MessageQueue::new();
136 mq.steer("msg1".into());
137 mq.steer("msg2".into());
138 mq.steer("msg3".into());
139
140 let drained = mq.drain_steering();
141 assert_eq!(drained, vec!["msg1", "msg2", "msg3"]);
142 }
143
144 #[test]
145 fn steer_and_drain_one_at_a_time() {
146 let mq = MessageQueue::with_mode(QueueMode::OneAtATime);
147 mq.steer("msg1".into());
148 mq.steer("msg2".into());
149 mq.steer("msg3".into());
150
151 assert_eq!(mq.drain_steering(), vec!["msg1"]);
152 assert_eq!(mq.drain_steering(), vec!["msg2"]);
153 assert_eq!(mq.drain_steering(), vec!["msg3"]);
154 assert!(mq.drain_steering().is_empty());
155 }
156
157 #[test]
158 fn drain_empty_steering_returns_empty() {
159 let mq = MessageQueue::new();
160 assert!(mq.drain_steering().is_empty());
161 }
162
163 #[test]
164 fn drain_empty_steering_one_at_a_time_returns_empty() {
165 let mq = MessageQueue::with_mode(QueueMode::OneAtATime);
166 assert!(mq.drain_steering().is_empty());
167 }
168
169 #[test]
170 fn has_steering_reports_correctly() {
171 let mq = MessageQueue::new();
172 assert!(!mq.has_steering());
173 mq.steer("msg".into());
174 assert!(mq.has_steering());
175 mq.drain_steering();
176 assert!(!mq.has_steering());
177 }
178
179 #[test]
182 fn follow_up_and_drain_all() {
183 let mq = MessageQueue::new();
184 mq.follow_up("f1".into());
185 mq.follow_up("f2".into());
186
187 assert_eq!(mq.drain_follow_up(), vec!["f1", "f2"]);
188 }
189
190 #[test]
191 fn follow_up_and_drain_one_at_a_time() {
192 let mq = MessageQueue::with_mode(QueueMode::OneAtATime);
193 mq.follow_up("f1".into());
194 mq.follow_up("f2".into());
195
196 assert_eq!(mq.drain_follow_up(), vec!["f1"]);
197 assert_eq!(mq.drain_follow_up(), vec!["f2"]);
198 assert!(mq.drain_follow_up().is_empty());
199 }
200
201 #[test]
202 fn drain_empty_follow_up_returns_empty() {
203 let mq = MessageQueue::new();
204 assert!(mq.drain_follow_up().is_empty());
205 }
206
207 #[test]
210 fn steering_and_follow_up_are_independent() {
211 let mq = MessageQueue::new();
212 mq.steer("s1".into());
213 mq.follow_up("f1".into());
214
215 assert_eq!(mq.drain_steering(), vec!["s1"]);
217 assert_eq!(mq.drain_follow_up(), vec!["f1"]);
218 }
219
220 #[test]
221 fn mode_switch_mid_queue() {
222 let mq = MessageQueue::new(); mq.steer("s1".into());
224 mq.steer("s2".into());
225
226 assert_eq!(mq.drain_steering(), vec!["s1", "s2"]);
228
229 mq.set_mode(QueueMode::OneAtATime);
231 mq.steer("s3".into());
232 mq.steer("s4".into());
233 assert_eq!(mq.drain_steering(), vec!["s3"]);
234 assert_eq!(mq.drain_steering(), vec!["s4"]);
235 }
236
237 #[test]
240 fn message_queue_is_send_and_sync() {
241 fn assert_send_sync<T: Send + Sync>() {}
242 assert_send_sync::<MessageQueue>();
243 }
244}