Skip to main content

trunk_recorder_plugin/
queue.rs

1//! A queue of calls worked off the event thread — what an uploader needs:
2//! retries with backoff, a warning status while calls are waiting, results
3//! and [`Metrics`] (queue depth, upload time, how the service is answering)
4//! reported for you, and calls still waiting at shutdown saved to the data
5//! folder and picked up at the next start.
6//!
7//! ```ignore
8//! let queue = CallQueue::start(host.clone(), QueueOptions::saved_in(&setup.data_dir), move |call| {
9//!     match upload(call) {
10//!         Ok(url) => Attempt::Done { url },
11//!         Err(e) if e.is_temporary() => Attempt::Retry(e.to_string()),
12//!         Err(e) => Attempt::Fail(e.to_string()),
13//!     }
14//! });
15//! // in call_concluded:   queue.push(call);
16//! // in shutdown:         self.queue.shutdown(grace);
17//! ```
18
19use 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/// What became of one try at a call.
30#[derive(Clone, Debug)]
31pub enum Attempt {
32    /// Done; `url`: where it can be found now (or empty).
33    Done { url: String },
34    /// Not for this plugin (reported as skipped).
35    Skip(String),
36    /// Failed for now (the service is down, the network's out): try again later.
37    Retry(String),
38    /// Failed for good (it was refused): report it and move on.
39    Fail(String),
40}
41
42#[derive(Clone, Debug)]
43pub struct QueueOptions {
44    /// How long to wait before each retry; after the last, the call fails.
45    pub retry_after: Vec<Duration>,
46    /// Calls worked on at once.
47    pub threads: usize,
48    /// Where to save calls still waiting at shutdown (None: they fail).
49    pub save_to: Option<PathBuf>,
50    /// Calls waiting beyond this are refused (reported as failed).
51    pub capacity: usize,
52    /// A name for the status line, e.g. "upload" → "3 uploads waiting to retry".
53    pub noun: &'static str,
54    /// The service the work goes to ("Broadcastify Calls"), for the
55    /// dashboard's endpoint status. None: the plugin's own name is shown.
56    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    /// The defaults, saving waiting calls to `queue.jsonl` in `data_dir`.
67    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    /// The waiting count last reported in a status.
86    reported: usize,
87    last_error: String,
88    stats: Tally,
89}
90
91/// What the queue has done, for [`Metrics`].
92#[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    /// Failed tries in a row (0 after a success).
100    in_a_row: u32,
101    tried: bool,
102    /// Unix s of the first try (the clock for "silent" before any success).
103    first_try: Option<f64>,
104    changed: bool,
105    sent: Option<Instant>,
106    sent_state: EndpointState,
107}
108
109/// Metrics go out at most this often, and at least this often while something changes.
110const METRICS_MIN: Duration = Duration::from_secs(5);
111const METRICS_MAX: Duration = Duration::from_secs(30);
112/// The service counts as down after this many failed tries in a row, or
113/// this long without a success while calls wait.
114const 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    /// Send metrics when something changed and it's been a while (or the
162    /// service's state changed: at once).
163    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    /// Start `opts.threads` workers running `work` on each call pushed — and
190    /// on calls saved by the last shutdown, first.
191    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    /// Calls queued or being worked on.
234    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    /// Work what's ready until `grace` is nearly up; save the rest (or fail it).
244    /// Call it from [`crate::Plugin::shutdown`]; calls pushed after it are refused.
245    pub fn shutdown(&mut self, grace: Duration) {
246        if self.workers.is_empty() {
247            return;
248        }
249        // Leave a little of the grace period for saving.
250        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        // Retries that are due go to the front: they've waited longest.
280        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            // Stopping: work what's ready (retries not yet due are saved) until the deadline.
291            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        // Say when calls start or stop waiting on retries (not on every change).
350        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        // (The user may have deleted the call meanwhile.)
373        .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        // The service went degraded on the retry (reported at once), then up.
446        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        // Next start: it's worked first.
483        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}