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//! reported for you, and calls still waiting at shutdown saved to the data
4//! folder and picked up at the next start.
5//!
6//! ```ignore
7//! let queue = CallQueue::start(host.clone(), QueueOptions::saved_in(&setup.data_dir), move |call| {
8//!     match upload(call) {
9//!         Ok(url) => Attempt::Done { url },
10//!         Err(e) if e.is_temporary() => Attempt::Retry(e.to_string()),
11//!         Err(e) => Attempt::Fail(e.to_string()),
12//!     }
13//! });
14//! // in call_concluded:   queue.push(call);
15//! // in shutdown:         self.queue.shutdown(grace);
16//! ```
17
18use 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/// What became of one try at a call.
29#[derive(Clone, Debug)]
30pub enum Attempt {
31    /// Done; `url`: where it can be found now (or empty).
32    Done { url: String },
33    /// Not for this plugin (reported as skipped).
34    Skip(String),
35    /// Failed for now (the service is down, the network's out): try again later.
36    Retry(String),
37    /// Failed for good (it was refused): report it and move on.
38    Fail(String),
39}
40
41#[derive(Clone, Debug)]
42pub struct QueueOptions {
43    /// How long to wait before each retry; after the last, the call fails.
44    pub retry_after: Vec<Duration>,
45    /// Calls worked on at once.
46    pub threads: usize,
47    /// Where to save calls still waiting at shutdown (None: they fail).
48    pub save_to: Option<PathBuf>,
49    /// Calls waiting beyond this are refused (reported as failed).
50    pub capacity: usize,
51    /// A name for the status line, e.g. "upload" → "3 uploads waiting to retry".
52    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    /// The defaults, saving waiting calls to `queue.jsonl` in `data_dir`.
63    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    /// The waiting count last reported in a status.
82    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    /// Start `opts.threads` workers running `work` on each call pushed — and
100    /// on calls saved by the last shutdown, first.
101    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    /// Calls queued or being worked on.
143    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    /// Work what's ready until `grace` is nearly up; save the rest (or fail it).
153    /// Call it from [`crate::Plugin::shutdown`]; calls pushed after it are refused.
154    pub fn shutdown(&mut self, grace: Duration) {
155        if self.workers.is_empty() {
156            return;
157        }
158        // Leave a little of the grace period for saving.
159        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        // Retries that are due go to the front: they've waited longest.
189        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            // Stopping: work what's ready (retries not yet due are saved) until the deadline.
200            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        // Say when calls start or stop waiting on retries (not on every change).
231        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        // (The user may have deleted the call meanwhile.)
253        .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        // Next start: it's worked first.
331        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}