1use std::collections::VecDeque;
20use std::io::{BufRead, BufReader, Write};
21use std::path::{Path, PathBuf};
22use std::sync::{Arc, Condvar, Mutex};
23use std::thread::JoinHandle;
24use std::time::{Duration, Instant};
25
26use crate::protocol::{ConcludedCall, Endpoint, EndpointState, Metrics, Outcome, State};
27use crate::sdk::Host;
28
29#[derive(Clone, Debug)]
31pub enum Attempt {
32 Done { url: String },
34 Skip(String),
36 Retry(String),
38 Fail(String),
40}
41
42#[derive(Clone, Debug)]
43pub struct QueueOptions {
44 pub retry_after: Vec<Duration>,
46 pub threads: usize,
48 pub save_to: Option<PathBuf>,
50 pub capacity: usize,
52 pub noun: &'static str,
54 pub endpoint: Option<String>,
57}
58
59impl Default for QueueOptions {
60 fn default() -> Self {
61 QueueOptions { retry_after: [10, 60, 300, 900].map(Duration::from_secs).to_vec(), threads: 2, save_to: None, capacity: 10_000, noun: "call", endpoint: None }
62 }
63}
64
65impl QueueOptions {
66 pub fn saved_in(data_dir: &Path) -> Self {
68 QueueOptions { save_to: Some(data_dir.join("queue.jsonl")), ..Default::default() }
69 }
70}
71
72struct Item {
73 call: ConcludedCall,
74 tries: usize,
75 due: Instant,
76 last_error: String,
77}
78
79#[derive(Default)]
80struct Inner {
81 ready: VecDeque<Item>,
82 waiting: Vec<Item>,
83 busy: usize,
84 closing: Option<Instant>,
85 reported: usize,
87 last_error: String,
88 stats: Tally,
89}
90
91#[derive(Default)]
93struct Tally {
94 retries: u64,
95 bytes: u64,
96 latency_ms: Option<f64>,
97 last_ok: Option<f64>,
98 last_error: Option<f64>,
99 in_a_row: u32,
101 tried: bool,
102 first_try: Option<f64>,
104 changed: bool,
105 sent: Option<Instant>,
106 sent_state: EndpointState,
107}
108
109const METRICS_MIN: Duration = Duration::from_secs(5);
111const METRICS_MAX: Duration = Duration::from_secs(30);
112const DOWN_AFTER: u32 = 3;
115const DOWN_SILENT: Duration = Duration::from_secs(300);
116
117fn unix_now() -> f64 {
118 std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).map_or(0.0, |d| d.as_secs_f64())
119}
120
121impl Inner {
122 fn endpoint_state(&self) -> EndpointState {
123 let t = &self.stats;
124 if !t.tried {
125 return EndpointState::Unknown;
126 }
127 let since = t.last_ok.or(t.first_try).unwrap_or_else(unix_now);
128 let silent = !self.waiting.is_empty() && unix_now() - since > DOWN_SILENT.as_secs_f64();
129 if t.in_a_row >= DOWN_AFTER || silent {
130 EndpointState::Down
131 } else if t.in_a_row > 0 || !self.waiting.is_empty() {
132 EndpointState::Degraded
133 } else {
134 EndpointState::Up
135 }
136 }
137
138 fn metrics(&self, o: &QueueOptions) -> Metrics {
139 let t = &self.stats;
140 let state = self.endpoint_state();
141 Metrics {
142 queued: Some((self.ready.len() + self.waiting.len()) as u64),
143 retrying: Some(self.waiting.len() as u64),
144 in_flight: Some(self.busy as u64),
145 retries: Some(t.retries),
146 bytes_sent: Some(t.bytes),
147 latency_ms: t.latency_ms.map(|l| (l * 10.0).round() / 10.0),
148 last_ok: t.last_ok,
149 last_error: t.last_error,
150 last_error_text: self.last_error.clone(),
151 endpoints: o
152 .endpoint
153 .as_ref()
154 .map(|name| Endpoint { name: name.clone(), state, latency_ms: t.latency_ms, last_ok: t.last_ok, last_error: if state == EndpointState::Up { String::new() } else { self.last_error.clone() } })
155 .into_iter()
156 .collect(),
157 ..Default::default()
158 }
159 }
160
161 fn report(&mut self, host: &Host, o: &QueueOptions) {
164 let state = self.endpoint_state();
165 let since = self.stats.sent.map_or(Duration::MAX, |t| t.elapsed());
166 let due = (self.stats.changed && since >= METRICS_MIN) || since >= METRICS_MAX || state != self.stats.sent_state;
167 if due {
168 host.metrics(&self.metrics(o));
169 self.stats.sent = Some(Instant::now());
170 self.stats.sent_state = state;
171 self.stats.changed = false;
172 }
173 }
174}
175
176struct Shared {
177 inner: Mutex<Inner>,
178 wake: Condvar,
179}
180
181pub struct CallQueue {
182 shared: Arc<Shared>,
183 workers: Vec<JoinHandle<()>>,
184 host: Host,
185 opts: QueueOptions,
186}
187
188impl CallQueue {
189 pub fn start<F>(host: Host, opts: QueueOptions, work: F) -> CallQueue
192 where
193 F: Fn(&ConcludedCall) -> Attempt + Send + Sync + 'static,
194 {
195 let shared = Arc::new(Shared { inner: Mutex::new(Inner::default()), wake: Condvar::new() });
196 if let Some(path) = &opts.save_to {
197 let saved = load(path);
198 if !saved.is_empty() {
199 host.info(format!("{} {}(s) left from last time", saved.len(), opts.noun));
200 let now = Instant::now();
201 shared.inner.lock().unwrap().ready.extend(saved.into_iter().map(|call| Item { call, tries: 0, due: now, last_error: String::new() }));
202 }
203 let _ = std::fs::remove_file(path);
204 }
205 let work = Arc::new(work);
206 let workers = (0..opts.threads.max(1))
207 .map(|i| {
208 let (sh, h, w, o) = (shared.clone(), host.clone(), work.clone(), opts.clone());
209 std::thread::Builder::new().name(format!("queue-{i}")).spawn(move || worker(&sh, &h, &*w, &o)).expect("thread")
210 })
211 .collect();
212 CallQueue { shared, workers, host, opts }
213 }
214
215 pub fn push(&self, call: ConcludedCall) {
216 let mut g = self.shared.inner.lock().unwrap();
217 if g.closing.is_some() {
218 drop(g);
219 self.host.call_result(&call.path, Outcome::Failed, "recorder stopping", "");
220 return;
221 }
222 if g.ready.len() + g.waiting.len() >= self.opts.capacity {
223 drop(g);
224 self.host.call_result(&call.path, Outcome::Failed, format!("too many {}s waiting", self.opts.noun), "");
225 return;
226 }
227 g.ready.push_back(Item { call, tries: 0, due: Instant::now(), last_error: String::new() });
228 g.stats.changed = true;
229 drop(g);
230 self.shared.wake.notify_one();
231 }
232
233 pub fn len(&self) -> usize {
235 let g = self.shared.inner.lock().unwrap();
236 g.ready.len() + g.waiting.len() + g.busy
237 }
238
239 pub fn is_empty(&self) -> bool {
240 self.len() == 0
241 }
242
243 pub fn shutdown(&mut self, grace: Duration) {
246 if self.workers.is_empty() {
247 return;
248 }
249 let deadline = Instant::now() + grace.mul_f64(0.8);
251 self.shared.inner.lock().unwrap().closing = Some(deadline);
252 self.shared.wake.notify_all();
253 for w in self.workers.drain(..) {
254 let _ = w.join();
255 }
256 let mut g = self.shared.inner.lock().unwrap();
257 let g = &mut *g;
258 let left: Vec<Item> = g.ready.drain(..).chain(g.waiting.drain(..)).collect();
259 if left.is_empty() {
260 return;
261 }
262 match &self.opts.save_to {
263 Some(path) => match save(path, &left) {
264 Ok(()) => self.host.info(format!("{} {}(s) saved for next time", left.len(), self.opts.noun)),
265 Err(e) => {
266 self.host.error(format!("couldn't save {} waiting {}(s): {e}", left.len(), self.opts.noun));
267 left.iter().for_each(|i| self.host.call_result(&i.call.path, Outcome::Failed, "recorder stopped", ""));
268 }
269 },
270 None => left.iter().for_each(|i| self.host.call_result(&i.call.path, Outcome::Failed, "recorder stopped", "")),
271 }
272 }
273}
274
275fn worker(sh: &Shared, host: &Host, work: &(dyn Fn(&ConcludedCall) -> Attempt + Send + Sync), o: &QueueOptions) {
276 let mut g = sh.inner.lock().unwrap();
277 loop {
278 let now = Instant::now();
279 let mut i = 0;
281 while i < g.waiting.len() {
282 if g.waiting[i].due <= now {
283 let it = g.waiting.swap_remove(i);
284 g.ready.push_front(it);
285 } else {
286 i += 1;
287 }
288 }
289 if let Some(deadline) = g.closing {
290 if now >= deadline || g.ready.is_empty() {
292 return;
293 }
294 }
295 let Some(mut it) = g.ready.pop_front() else {
296 g.report(host, o);
297 let next = g.waiting.iter().map(|w| w.due).min();
298 let wait = next.map_or(Duration::from_secs(1), |t| t.saturating_duration_since(now)).min(Duration::from_secs(1));
299 g = sh.wake.wait_timeout(g, wait).unwrap().0;
300 continue;
301 };
302 g.busy += 1;
303 drop(g);
304 let started = Instant::now();
305 let r = work(&it.call);
306 let took_ms = started.elapsed().as_secs_f64() * 1000.0;
307 g = sh.inner.lock().unwrap();
308 g.busy -= 1;
309 {
310 let t = &mut g.stats;
311 t.changed = true;
312 t.first_try.get_or_insert_with(unix_now);
313 match &r {
314 Attempt::Done { .. } => {
315 t.tried = true;
316 t.in_a_row = 0;
317 t.last_ok = Some(unix_now());
318 t.latency_ms = Some(t.latency_ms.map_or(took_ms, |l| l + 0.2 * (took_ms - l)));
319 let f = &it.call.files;
320 t.bytes += f.m4a.as_ref().and_then(|m| std::fs::metadata(m).ok()).or_else(|| std::fs::metadata(&f.wav).ok()).map_or(0, |m| m.len());
321 }
322 Attempt::Retry(_) | Attempt::Fail(_) => {
323 t.tried = true;
324 t.in_a_row += 1;
325 t.last_error = Some(unix_now());
326 t.retries += matches!(r, Attempt::Retry(_)) as u64;
327 }
328 Attempt::Skip(_) => {}
329 }
330 }
331 if let Attempt::Fail(m) = &r {
332 g.last_error = m.clone();
333 }
334 match r {
335 Attempt::Done { url } => host.call_result(&it.call.path, Outcome::Ok, "", url),
336 Attempt::Skip(m) => host.call_result(&it.call.path, Outcome::Skipped, m, ""),
337 Attempt::Fail(m) => host.call_result(&it.call.path, Outcome::Failed, m, ""),
338 Attempt::Retry(m) => match o.retry_after.get(it.tries) {
339 Some(d) => {
340 it.tries += 1;
341 it.due = Instant::now() + *d;
342 it.last_error = m.clone();
343 g.last_error = m;
344 g.waiting.push(it);
345 }
346 None => host.call_result(&it.call.path, Outcome::Failed, format!("gave up after {} tries: {m}", it.tries + 1), ""),
347 },
348 }
349 let n = g.waiting.len();
351 if (n == 0) != (g.reported == 0) || n >= g.reported * 2 && n >= 10 {
352 g.reported = n;
353 if n == 0 {
354 host.status(State::Ok, "");
355 } else {
356 let noun = if n == 1 { o.noun.to_string() } else { format!("{}s", o.noun) };
357 host.status(State::Warning, format!("{n} {noun} waiting to retry: {}", g.last_error));
358 }
359 }
360 g.report(host, o);
361 }
362}
363
364fn load(path: &Path) -> Vec<ConcludedCall> {
365 let Ok(f) = std::fs::File::open(path) else {
366 return Vec::new();
367 };
368 BufReader::new(f)
369 .lines()
370 .map_while(Result::ok)
371 .filter_map(|l| serde_json::from_str::<ConcludedCall>(&l).ok())
372 .filter(|c| c.files.wav.exists() || c.files.m4a.as_ref().is_some_and(|m| m.exists()))
374 .collect()
375}
376
377fn save(path: &Path, items: &[Item]) -> std::io::Result<()> {
378 if let Some(d) = path.parent() {
379 std::fs::create_dir_all(d)?;
380 }
381 let mut f = std::io::BufWriter::new(std::fs::File::create(path)?);
382 for it in items {
383 writeln!(f, "{}", serde_json::to_string(&it.call).unwrap_or_default())?;
384 }
385 f.flush()
386}
387
388#[cfg(test)]
389mod tests {
390 use super::*;
391 use crate::testing;
392 use std::sync::atomic::{AtomicUsize, Ordering};
393
394 fn quick(dir: &Path) -> QueueOptions {
395 QueueOptions { retry_after: vec![Duration::from_millis(20), Duration::from_millis(20)], threads: 1, ..QueueOptions::saved_in(dir) }
396 }
397
398 #[test]
399 fn retries_then_succeeds_and_reports() {
400 let dir = testing::temp_dir("queue");
401 let (host, out) = testing::capture();
402 let tries = Arc::new(AtomicUsize::new(0));
403 let t = tries.clone();
404 let mut q = CallQueue::start(host, quick(&dir), move |_| {
405 if t.fetch_add(1, Ordering::SeqCst) == 0 {
406 Attempt::Retry("down".into())
407 } else {
408 Attempt::Done { url: "u".into() }
409 }
410 });
411 let call = testing::call(&dir, "sys1", 5);
412 q.push(call.clone());
413 let t0 = Instant::now();
414 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
415 std::thread::sleep(Duration::from_millis(5));
416 }
417 q.shutdown(Duration::from_secs(1));
418 let out = out.output();
419 assert_eq!(tries.load(Ordering::SeqCst), 2);
420 assert_eq!(out.results(), vec![(call.path, Outcome::Ok, String::new(), "u".into())]);
421 assert!(matches!(out.status(), Some((State::Ok, _))));
422 }
423
424 #[test]
425 fn metrics_follow_the_service() {
426 let dir = testing::temp_dir("queue");
427 let (host, out) = testing::capture();
428 let tries = Arc::new(AtomicUsize::new(0));
429 let t = tries.clone();
430 let opts = QueueOptions { endpoint: Some("Example Calls".into()), ..quick(&dir) };
431 let mut q = CallQueue::start(host, opts, move |_| {
432 if t.fetch_add(1, Ordering::SeqCst) == 0 {
433 Attempt::Retry("503".into())
434 } else {
435 Attempt::Done { url: String::new() }
436 }
437 });
438 q.push(testing::call(&dir, "sys1", 5));
439 let t0 = Instant::now();
440 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
441 std::thread::sleep(Duration::from_millis(5));
442 }
443 q.shutdown(Duration::from_secs(1));
444 let m = out.output().metrics();
445 let states: Vec<EndpointState> = m.iter().filter_map(|x| x.endpoints.first().map(|e| e.state)).collect();
447 assert!(states.contains(&EndpointState::Degraded), "{states:?}");
448 assert_eq!(states.last(), Some(&EndpointState::Up), "{states:?}");
449 let last = m.last().unwrap();
450 assert_eq!(last.retries, Some(1));
451 assert_eq!(last.queued, Some(0));
452 assert!(last.latency_ms.is_some() && last.last_ok.is_some());
453 assert_eq!(last.endpoints[0].name, "Example Calls");
454 }
455
456 #[test]
457 fn gives_up_after_the_last_retry() {
458 let dir = testing::temp_dir("queue");
459 let (host, out) = testing::capture();
460 let mut q = CallQueue::start(host, quick(&dir), |_| Attempt::Retry("down".into()));
461 q.push(testing::call(&dir, "sys1", 5));
462 let t0 = Instant::now();
463 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
464 std::thread::sleep(Duration::from_millis(5));
465 }
466 q.shutdown(Duration::from_secs(1));
467 let r = out.output().results();
468 assert_eq!(r[0].1, Outcome::Failed);
469 assert!(r[0].2.contains("3 tries"), "{}", r[0].2);
470 }
471
472 #[test]
473 fn saves_what_waits_at_shutdown_and_resumes() {
474 let dir = testing::temp_dir("queue");
475 let (host, _) = testing::capture();
476 let opts = QueueOptions { retry_after: vec![Duration::from_secs(3600)], threads: 1, ..QueueOptions::saved_in(&dir) };
477 let mut q = CallQueue::start(host, opts.clone(), |_| Attempt::Retry("down".into()));
478 q.push(testing::call(&dir, "sys1", 5));
479 std::thread::sleep(Duration::from_millis(50));
480 q.shutdown(Duration::from_secs(1));
481 assert!(dir.join("queue.jsonl").exists());
482 let (host, out) = testing::capture();
484 let mut q = CallQueue::start(host, opts, |_| Attempt::Done { url: String::new() });
485 let t0 = Instant::now();
486 while !q.is_empty() && t0.elapsed() < Duration::from_secs(5) {
487 std::thread::sleep(Duration::from_millis(5));
488 }
489 q.shutdown(Duration::from_secs(1));
490 assert_eq!(out.output().results()[0].1, Outcome::Ok);
491 assert!(!dir.join("queue.jsonl").exists());
492 }
493}