1use std::collections::VecDeque;
19use std::io::{BufRead, BufReader, Write};
20use std::path::{Path, PathBuf};
21use std::sync::{Arc, Condvar, Mutex};
22use std::thread::JoinHandle;
23use std::time::{Duration, Instant};
24
25use crate::protocol::{ConcludedCall, Outcome, State};
26use crate::sdk::Host;
27
28#[derive(Clone, Debug)]
30pub enum Attempt {
31 Done { url: String },
33 Skip(String),
35 Retry(String),
37 Fail(String),
39}
40
41#[derive(Clone, Debug)]
42pub struct QueueOptions {
43 pub retry_after: Vec<Duration>,
45 pub threads: usize,
47 pub save_to: Option<PathBuf>,
49 pub capacity: usize,
51 pub noun: &'static str,
53}
54
55impl Default for QueueOptions {
56 fn default() -> Self {
57 QueueOptions { retry_after: [10, 60, 300, 900].map(Duration::from_secs).to_vec(), threads: 2, save_to: None, capacity: 10_000, noun: "call" }
58 }
59}
60
61impl QueueOptions {
62 pub fn saved_in(data_dir: &Path) -> Self {
64 QueueOptions { save_to: Some(data_dir.join("queue.jsonl")), ..Default::default() }
65 }
66}
67
68struct Item {
69 call: ConcludedCall,
70 tries: usize,
71 due: Instant,
72 last_error: String,
73}
74
75#[derive(Default)]
76struct Inner {
77 ready: VecDeque<Item>,
78 waiting: Vec<Item>,
79 busy: usize,
80 closing: Option<Instant>,
81 reported: usize,
83 last_error: String,
84}
85
86struct Shared {
87 inner: Mutex<Inner>,
88 wake: Condvar,
89}
90
91pub struct CallQueue {
92 shared: Arc<Shared>,
93 workers: Vec<JoinHandle<()>>,
94 host: Host,
95 opts: QueueOptions,
96}
97
98impl CallQueue {
99 pub fn start<F>(host: Host, opts: QueueOptions, work: F) -> CallQueue
102 where
103 F: Fn(&ConcludedCall) -> Attempt + Send + Sync + 'static,
104 {
105 let shared = Arc::new(Shared { inner: Mutex::new(Inner::default()), wake: Condvar::new() });
106 if let Some(path) = &opts.save_to {
107 let saved = load(path);
108 if !saved.is_empty() {
109 host.info(format!("{} {}(s) left from last time", saved.len(), opts.noun));
110 let now = Instant::now();
111 shared.inner.lock().unwrap().ready.extend(saved.into_iter().map(|call| Item { call, tries: 0, due: now, last_error: String::new() }));
112 }
113 let _ = std::fs::remove_file(path);
114 }
115 let work = Arc::new(work);
116 let workers = (0..opts.threads.max(1))
117 .map(|i| {
118 let (sh, h, w, o) = (shared.clone(), host.clone(), work.clone(), opts.clone());
119 std::thread::Builder::new().name(format!("queue-{i}")).spawn(move || worker(&sh, &h, &*w, &o)).expect("thread")
120 })
121 .collect();
122 CallQueue { shared, workers, host, opts }
123 }
124
125 pub fn push(&self, call: ConcludedCall) {
126 let mut g = self.shared.inner.lock().unwrap();
127 if g.closing.is_some() {
128 drop(g);
129 self.host.call_result(&call.path, Outcome::Failed, "recorder stopping", "");
130 return;
131 }
132 if g.ready.len() + g.waiting.len() >= self.opts.capacity {
133 drop(g);
134 self.host.call_result(&call.path, Outcome::Failed, format!("too many {}s waiting", self.opts.noun), "");
135 return;
136 }
137 g.ready.push_back(Item { call, tries: 0, due: Instant::now(), last_error: String::new() });
138 drop(g);
139 self.shared.wake.notify_one();
140 }
141
142 pub fn len(&self) -> usize {
144 let g = self.shared.inner.lock().unwrap();
145 g.ready.len() + g.waiting.len() + g.busy
146 }
147
148 pub fn is_empty(&self) -> bool {
149 self.len() == 0
150 }
151
152 pub fn shutdown(&mut self, grace: Duration) {
155 if self.workers.is_empty() {
156 return;
157 }
158 let deadline = Instant::now() + grace.mul_f64(0.8);
160 self.shared.inner.lock().unwrap().closing = Some(deadline);
161 self.shared.wake.notify_all();
162 for w in self.workers.drain(..) {
163 let _ = w.join();
164 }
165 let mut g = self.shared.inner.lock().unwrap();
166 let g = &mut *g;
167 let left: Vec<Item> = g.ready.drain(..).chain(g.waiting.drain(..)).collect();
168 if left.is_empty() {
169 return;
170 }
171 match &self.opts.save_to {
172 Some(path) => match save(path, &left) {
173 Ok(()) => self.host.info(format!("{} {}(s) saved for next time", left.len(), self.opts.noun)),
174 Err(e) => {
175 self.host.error(format!("couldn't save {} waiting {}(s): {e}", left.len(), self.opts.noun));
176 left.iter().for_each(|i| self.host.call_result(&i.call.path, Outcome::Failed, "recorder stopped", ""));
177 }
178 },
179 None => left.iter().for_each(|i| self.host.call_result(&i.call.path, Outcome::Failed, "recorder stopped", "")),
180 }
181 }
182}
183
184fn worker(sh: &Shared, host: &Host, work: &(dyn Fn(&ConcludedCall) -> Attempt + Send + Sync), o: &QueueOptions) {
185 let mut g = sh.inner.lock().unwrap();
186 loop {
187 let now = Instant::now();
188 let mut i = 0;
190 while i < g.waiting.len() {
191 if g.waiting[i].due <= now {
192 let it = g.waiting.swap_remove(i);
193 g.ready.push_front(it);
194 } else {
195 i += 1;
196 }
197 }
198 if let Some(deadline) = g.closing {
199 if now >= deadline || g.ready.is_empty() {
201 return;
202 }
203 }
204 let Some(mut it) = g.ready.pop_front() else {
205 let next = g.waiting.iter().map(|w| w.due).min();
206 let wait = next.map_or(Duration::from_secs(1), |t| t.saturating_duration_since(now)).min(Duration::from_secs(1));
207 g = sh.wake.wait_timeout(g, wait).unwrap().0;
208 continue;
209 };
210 g.busy += 1;
211 drop(g);
212 let r = work(&it.call);
213 g = sh.inner.lock().unwrap();
214 g.busy -= 1;
215 match r {
216 Attempt::Done { url } => host.call_result(&it.call.path, Outcome::Ok, "", url),
217 Attempt::Skip(m) => host.call_result(&it.call.path, Outcome::Skipped, m, ""),
218 Attempt::Fail(m) => host.call_result(&it.call.path, Outcome::Failed, m, ""),
219 Attempt::Retry(m) => match o.retry_after.get(it.tries) {
220 Some(d) => {
221 it.tries += 1;
222 it.due = Instant::now() + *d;
223 it.last_error = m.clone();
224 g.last_error = m;
225 g.waiting.push(it);
226 }
227 None => host.call_result(&it.call.path, Outcome::Failed, format!("gave up after {} tries: {m}", it.tries + 1), ""),
228 },
229 }
230 let n = g.waiting.len();
232 if (n == 0) != (g.reported == 0) || n >= g.reported * 2 && n >= 10 {
233 g.reported = n;
234 if n == 0 {
235 host.status(State::Ok, "");
236 } else {
237 let noun = if n == 1 { o.noun.to_string() } else { format!("{}s", o.noun) };
238 host.status(State::Warning, format!("{n} {noun} waiting to retry: {}", g.last_error));
239 }
240 }
241 }
242}
243
244fn load(path: &Path) -> Vec<ConcludedCall> {
245 let Ok(f) = std::fs::File::open(path) else {
246 return Vec::new();
247 };
248 BufReader::new(f)
249 .lines()
250 .map_while(Result::ok)
251 .filter_map(|l| serde_json::from_str::<ConcludedCall>(&l).ok())
252 .filter(|c| c.files.wav.exists() || c.files.m4a.as_ref().is_some_and(|m| m.exists()))
254 .collect()
255}
256
257fn save(path: &Path, items: &[Item]) -> std::io::Result<()> {
258 if let Some(d) = path.parent() {
259 std::fs::create_dir_all(d)?;
260 }
261 let mut f = std::io::BufWriter::new(std::fs::File::create(path)?);
262 for it in items {
263 writeln!(f, "{}", serde_json::to_string(&it.call).unwrap_or_default())?;
264 }
265 f.flush()
266}
267
268#[cfg(test)]
269mod tests {
270 use super::*;
271 use crate::testing;
272 use std::sync::atomic::{AtomicUsize, Ordering};
273
274 fn quick(dir: &Path) -> QueueOptions {
275 QueueOptions { retry_after: vec![Duration::from_millis(20), Duration::from_millis(20)], threads: 1, ..QueueOptions::saved_in(dir) }
276 }
277
278 #[test]
279 fn retries_then_succeeds_and_reports() {
280 let dir = testing::temp_dir("queue");
281 let (host, out) = testing::capture();
282 let tries = Arc::new(AtomicUsize::new(0));
283 let t = tries.clone();
284 let mut q = CallQueue::start(host, quick(&dir), move |_| {
285 if t.fetch_add(1, Ordering::SeqCst) == 0 {
286 Attempt::Retry("down".into())
287 } else {
288 Attempt::Done { url: "u".into() }
289 }
290 });
291 let call = testing::call(&dir, "sys1", 5);
292 q.push(call.clone());
293 let t0 = Instant::now();
294 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
295 std::thread::sleep(Duration::from_millis(5));
296 }
297 q.shutdown(Duration::from_secs(1));
298 let out = out.output();
299 assert_eq!(tries.load(Ordering::SeqCst), 2);
300 assert_eq!(out.results(), vec![(call.path, Outcome::Ok, String::new(), "u".into())]);
301 assert!(matches!(out.status(), Some((State::Ok, _))));
302 }
303
304 #[test]
305 fn gives_up_after_the_last_retry() {
306 let dir = testing::temp_dir("queue");
307 let (host, out) = testing::capture();
308 let mut q = CallQueue::start(host, quick(&dir), |_| Attempt::Retry("down".into()));
309 q.push(testing::call(&dir, "sys1", 5));
310 let t0 = Instant::now();
311 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
312 std::thread::sleep(Duration::from_millis(5));
313 }
314 q.shutdown(Duration::from_secs(1));
315 let r = out.output().results();
316 assert_eq!(r[0].1, Outcome::Failed);
317 assert!(r[0].2.contains("3 tries"), "{}", r[0].2);
318 }
319
320 #[test]
321 fn saves_what_waits_at_shutdown_and_resumes() {
322 let dir = testing::temp_dir("queue");
323 let (host, _) = testing::capture();
324 let opts = QueueOptions { retry_after: vec![Duration::from_secs(3600)], threads: 1, ..QueueOptions::saved_in(&dir) };
325 let mut q = CallQueue::start(host, opts.clone(), |_| Attempt::Retry("down".into()));
326 q.push(testing::call(&dir, "sys1", 5));
327 std::thread::sleep(Duration::from_millis(50));
328 q.shutdown(Duration::from_secs(1));
329 assert!(dir.join("queue.jsonl").exists());
330 let (host, out) = testing::capture();
332 let mut q = CallQueue::start(host, opts, |_| Attempt::Done { url: String::new() });
333 let t0 = Instant::now();
334 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
335 std::thread::sleep(Duration::from_millis(5));
336 }
337 q.shutdown(Duration::from_secs(1));
338 assert_eq!(out.output().results()[0].1, Outcome::Ok);
339 assert!(!dir.join("queue.jsonl").exists());
340 }
341}