Skip to main content

lean_ctx/core/a2a/
dlq.rs

1use std::collections::BTreeMap;
2use std::sync::{Arc, Mutex};
3
4use chrono::{DateTime, Utc};
5use serde::{Deserialize, Serialize};
6
7use crate::core::agents::AgentRegistry;
8use crate::core::ocla::types::{OclaError, OclaResult};
9
10const MAX_ENTRIES: usize = 1_000;
11
12#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
13pub struct DeadLetter {
14    pub id: String,
15    pub original_message: String,
16    pub target_agent: String,
17    pub error: String,
18    pub attempts: u8,
19    pub first_failed_at: String,
20    pub last_failed_at: String,
21}
22
23#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
24pub struct DlqStats {
25    pub total: usize,
26    pub oldest_age_seconds: u64,
27    pub by_target_agent: BTreeMap<String, usize>,
28}
29
30#[derive(Clone, Default)]
31pub struct DeadLetterQueue {
32    entries: Arc<Mutex<Vec<DeadLetter>>>,
33}
34
35impl DeadLetterQueue {
36    #[must_use]
37    pub fn new() -> Self {
38        Self::default()
39    }
40
41    pub fn enqueue(&self, letter: DeadLetter) {
42        let mut entries = self
43            .entries
44            .lock()
45            .unwrap_or_else(std::sync::PoisonError::into_inner);
46        if entries.len() == MAX_ENTRIES {
47            entries.remove(0);
48        }
49        entries.push(letter);
50    }
51
52    pub fn dequeue(&self, id: &str) -> Option<DeadLetter> {
53        let mut entries = self
54            .entries
55            .lock()
56            .unwrap_or_else(std::sync::PoisonError::into_inner);
57        let position = entries.iter().position(|letter| letter.id == id)?;
58        Some(entries.remove(position))
59    }
60
61    #[must_use]
62    pub fn peek_all(&self) -> Vec<DeadLetter> {
63        self.entries
64            .lock()
65            .unwrap_or_else(std::sync::PoisonError::into_inner)
66            .clone()
67    }
68
69    pub fn retry(&self, id: &str) -> OclaResult<()> {
70        let Some(letter) = self
71            .entries
72            .lock()
73            .unwrap_or_else(std::sync::PoisonError::into_inner)
74            .iter()
75            .find(|letter| letter.id == id)
76            .cloned()
77        else {
78            return Err(OclaError::InvalidRequest(format!(
79                "dead letter not found: {id}"
80            )));
81        };
82
83        match resend(&letter) {
84            Ok(()) => {
85                let _ = self.dequeue(id);
86                Ok(())
87            }
88            Err(error) => {
89                let mut entries = self
90                    .entries
91                    .lock()
92                    .unwrap_or_else(std::sync::PoisonError::into_inner);
93                if let Some(current) = entries.iter_mut().find(|current| current.id == id) {
94                    current.attempts = current.attempts.saturating_add(1);
95                    current.last_failed_at = Utc::now().to_rfc3339();
96                }
97                Err(error)
98            }
99        }
100    }
101
102    #[must_use]
103    pub fn stats(&self) -> DlqStats {
104        let entries = self
105            .entries
106            .lock()
107            .unwrap_or_else(std::sync::PoisonError::into_inner);
108        let now = Utc::now();
109        let mut oldest_age_seconds = 0;
110        let mut by_target_agent = BTreeMap::new();
111
112        for letter in entries.iter() {
113            *by_target_agent
114                .entry(letter.target_agent.clone())
115                .or_insert(0) += 1;
116            if let Ok(failed_at) = DateTime::parse_from_rfc3339(&letter.first_failed_at) {
117                let age =
118                    u64::try_from((now - failed_at.with_timezone(&Utc)).num_seconds()).unwrap_or(0);
119                oldest_age_seconds = oldest_age_seconds.max(age);
120            }
121        }
122
123        DlqStats {
124            total: entries.len(),
125            oldest_age_seconds,
126            by_target_agent,
127        }
128    }
129}
130
131fn resend(letter: &DeadLetter) -> OclaResult<()> {
132    if letter.original_message.trim().is_empty() || letter.target_agent.trim().is_empty() {
133        return Err(OclaError::InvalidRequest(
134            "dead letter message and target_agent are required".to_string(),
135        ));
136    }
137
138    AgentRegistry::mutate_locked(|registry| {
139        registry.post_message(
140            "dead-letter-queue",
141            Some(&letter.target_agent),
142            "retry",
143            &letter.original_message,
144        );
145    })
146    .map(|_| ())
147    .map_err(|error| OclaError::InvalidRequest(format!("dead letter retry failed: {error}")))
148}
149
150#[cfg(test)]
151mod tests {
152    use super::*;
153
154    fn letter(id: &str, target: &str, first_failed_at: &str) -> DeadLetter {
155        DeadLetter {
156            id: id.to_string(),
157            original_message: format!("message-{id}"),
158            target_agent: target.to_string(),
159            error: "delivery failed".to_string(),
160            attempts: 1,
161            first_failed_at: first_failed_at.to_string(),
162            last_failed_at: first_failed_at.to_string(),
163        }
164    }
165
166    #[test]
167    fn enqueue_dequeue_and_peek_preserve_entries() {
168        let queue = DeadLetterQueue::new();
169        let item = letter("one", "agent-a", "2026-01-01T00:00:00Z");
170        queue.enqueue(item.clone());
171
172        assert_eq!(queue.peek_all(), vec![item.clone()]);
173        assert_eq!(queue.dequeue("one"), Some(item));
174        assert!(queue.peek_all().is_empty());
175        assert!(queue.dequeue("missing").is_none());
176    }
177
178    #[test]
179    fn enqueue_evicts_oldest_entry_at_capacity() {
180        let queue = DeadLetterQueue::new();
181        for index in 0..=MAX_ENTRIES {
182            queue.enqueue(letter(
183                &index.to_string(),
184                "agent-a",
185                "2026-01-01T00:00:00Z",
186            ));
187        }
188
189        let entries = queue.peek_all();
190        assert_eq!(entries.len(), MAX_ENTRIES);
191        assert_eq!(entries.first().map(|entry| entry.id.as_str()), Some("1"));
192        assert_eq!(entries.last().map(|entry| entry.id.as_str()), Some("1000"));
193    }
194
195    #[test]
196    fn stats_report_age_and_target_counts() {
197        let queue = DeadLetterQueue::new();
198        queue.enqueue(letter("one", "agent-a", "2020-01-01T00:00:00Z"));
199        queue.enqueue(letter("two", "agent-a", "2026-01-01T00:00:00Z"));
200        queue.enqueue(letter("three", "agent-b", "not-a-timestamp"));
201
202        let stats = queue.stats();
203        assert_eq!(stats.total, 3);
204        assert!(stats.oldest_age_seconds > 0);
205        assert_eq!(stats.by_target_agent.get("agent-a"), Some(&2));
206        assert_eq!(stats.by_target_agent.get("agent-b"), Some(&1));
207    }
208
209    #[test]
210    fn retry_invalid_message_keeps_letter_and_records_attempt() {
211        let queue = DeadLetterQueue::new();
212        let mut item = letter("one", "agent-a", "2026-01-01T00:00:00Z");
213        item.original_message.clear();
214        queue.enqueue(item);
215
216        assert!(queue.retry("one").is_err());
217        let retained = queue.peek_all();
218        assert_eq!(retained[0].attempts, 2);
219        assert!(!retained[0].last_failed_at.is_empty());
220    }
221
222    #[test]
223    fn retry_resends_and_removes_letter() {
224        let _isolated = crate::core::data_dir::isolated_data_dir();
225        let queue = DeadLetterQueue::new();
226        queue.enqueue(letter("one", "agent-a", "2026-01-01T00:00:00Z"));
227
228        queue.retry("one").expect("retry succeeds");
229        assert!(queue.peek_all().is_empty());
230
231        let registry = AgentRegistry::load().expect("registry persisted");
232        assert_eq!(registry.scratchpad[0].to_agent.as_deref(), Some("agent-a"));
233        assert_eq!(registry.scratchpad[0].message, "message-one");
234    }
235}