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}