Skip to main content

rpi_agent/
queue.rs

1//! Mirrors the queue portion of `packages/agent/src/agent.ts`
2//! (`PendingMessageQueue`) + the `QueueMode` type from
3//! `packages/agent/src/types.ts`.
4//!
5//! The agent loop has two drain points:
6//! - **steering**: fired mid-run after a turn's tool batch settles (before the
7//!   next LLM call), to inject guidance the app pushed while the model was
8//!   working. Drained via [`PendingMessageQueue::try_drain`].
9//! - **follow-up**: fired when the loop would otherwise stop (no tool calls +
10//!   no steering), to keep the conversation going with queued user messages.
11//!
12//! Both reuse the same [`PendingMessageQueue`] type with a [`QueueMode`]
13//! controlling how many messages a single `try_drain` releases.
14
15use std::collections::VecDeque;
16
17use crate::message::AgentMessage;
18use crate::types::QueueMode;
19
20/// A FIFO of `AgentMessage`s waiting to be injected at a drain point. Mirrors
21/// TS `PendingMessageQueue`. `Clone` so `Agent` can keep a steering instance
22/// and a follow-up instance independently.
23#[derive(Debug, Clone, Default)]
24pub struct PendingMessageQueue {
25    mode: QueueMode,
26    pending: VecDeque<AgentMessage>,
27}
28
29impl PendingMessageQueue {
30    pub fn new(mode: QueueMode) -> Self {
31        Self {
32            mode,
33            pending: VecDeque::new(),
34        }
35    }
36
37    pub fn mode(&self) -> QueueMode {
38        self.mode
39    }
40
41    /// Append a message to the tail. Mirrors TS `enqueue`.
42    pub fn enqueue(&mut self, message: AgentMessage) {
43        self.pending.push_back(message);
44    }
45
46    pub fn is_empty(&self) -> bool {
47        self.pending.is_empty()
48    }
49
50    pub fn len(&self) -> usize {
51        self.pending.len()
52    }
53
54    /// Drain queued messages per [`QueueMode`]:
55    /// - `All`: return and remove every queued message, in enqueue order.
56    /// - `OneAtATime`: return and remove only the oldest, leaving the rest.
57    ///
58    /// Returns an empty `Vec` when nothing is queued (the loop treats an empty
59    /// drain as "no injection"). Mirrors TS `tryDrain`.
60    pub fn try_drain(&mut self) -> Vec<AgentMessage> {
61        if self.pending.is_empty() {
62            return Vec::new();
63        }
64        match self.mode {
65            QueueMode::All => self.pending.drain(..).collect(),
66            QueueMode::OneAtATime => self
67                .pending
68                .pop_front()
69                .map(|m| vec![m])
70                .unwrap_or_default(),
71        }
72    }
73}
74
75#[cfg(test)]
76mod tests {
77    use super::*;
78    use rpi_ai::types::UserMessage;
79    use AgentMessage;
80
81    fn user(s: &str) -> AgentMessage {
82        AgentMessage::User(UserMessage::new(s, 0))
83    }
84
85    #[test]
86    fn drain_all_empties_queue() {
87        let mut q = PendingMessageQueue::new(QueueMode::All);
88        q.enqueue(user("a"));
89        q.enqueue(user("b"));
90        let drained = q.try_drain();
91        assert_eq!(drained.len(), 2);
92        assert_eq!(drained[0].role().as_str(), "user");
93        assert!(q.is_empty());
94        assert!(q.try_drain().is_empty());
95    }
96
97    #[test]
98    fn drain_one_at_a_time_keeps_tail() {
99        let mut q = PendingMessageQueue::new(QueueMode::OneAtATime);
100        q.enqueue(user("a"));
101        q.enqueue(user("b"));
102        let first = q.try_drain();
103        assert_eq!(first.len(), 1);
104        assert_eq!(q.len(), 1);
105        let second = q.try_drain();
106        assert_eq!(second.len(), 1);
107        assert!(q.is_empty());
108    }
109
110    #[test]
111    fn drain_empty_returns_empty() {
112        let mut q = PendingMessageQueue::new(QueueMode::All);
113        assert!(q.try_drain().is_empty());
114    }
115}