Skip to main content

agent_base/engine/runtime/
message_queue.rs

1use std::collections::VecDeque;
2use std::sync::Mutex;
3
4/// Controls how queued messages are drained each turn.
5#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
6pub enum QueueMode {
7    /// Drain all queued messages at once.
8    #[default]
9    All,
10    /// Drain one message per iteration (the rest wait for the next turn).
11    OneAtATime,
12}
13
14// (no separate impl Default needed — derived above)
15
16/// Dual-queue message system for steering and follow-up messages.
17///
18/// Inspired by Pi's `steeringQueue` / `followUpQueue` pattern:
19/// - **Steering**: messages injected mid-run, processed at the next turn.
20/// - **Follow-up**: messages processed after the agent stops naturally.
21///
22/// `QueueMode` controls how many messages are drained per turn.
23pub 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    /// Set the drain mode for both queues at runtime.
48    pub fn set_mode(&self, mode: QueueMode) {
49        *self.mode.lock().unwrap_or_else(|e| e.into_inner()) = mode;
50    }
51
52    /// Get the current drain mode.
53    #[allow(dead_code)]
54    pub fn mode(&self) -> QueueMode {
55        *self.mode.lock().unwrap_or_else(|e| e.into_inner())
56    }
57
58    /// Push a steering message — will be processed at the start of the next turn.
59    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    /// Push a follow-up message — will be processed after the agent stops.
67    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    /// Drain steering messages according to the current `QueueMode`.
75    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    /// Drain follow-up messages according to the current `QueueMode`.
85    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    /// Check if the steering queue has any pending messages.
95    #[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    // ── construction ──
110
111    #[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    // ── steering queue ──
132
133    #[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    // ── follow-up queue ──
180
181    #[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    // ── independence between queues ──
208
209    #[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        // Draining steering does not affect follow-up
216        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(); // All mode
223        mq.steer("s1".into());
224        mq.steer("s2".into());
225
226        // Drain all first
227        assert_eq!(mq.drain_steering(), vec!["s1", "s2"]);
228
229        // Switch to OneAtATime
230        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    // ── thread safety (basic) ──
238
239    #[test]
240    fn message_queue_is_send_and_sync() {
241        fn assert_send_sync<T: Send + Sync>() {}
242        assert_send_sync::<MessageQueue>();
243    }
244}