Skip to main content

truce_core/
tasks.rs

1//! Managed background-task pool.
2//!
3//! A process-global pool of worker threads runs plugin `run_task`
4//! handlers off the audio thread. Each plugin instance owns a
5//! preallocated, wait-free inbound queue via a [`TaskSpawner`]: the
6//! audio thread (or the editor, or `init`) pushes tasks without
7//! allocating or blocking, and a pool worker drains them. Feedback to
8//! the audio thread stays the plugin's job through shared `#[skip]`
9//! channels - the pool owns only the worker threads and the inbound
10//! queue.
11//!
12//! The pool is shared across every instance in the process (one small
13//! set of threads, not one thread per instance) and initializes lazily
14//! the first time any instance actually schedules a task, so a plugin
15//! that never declares a `BackgroundTask` spawns no threads.
16//!
17//! Because the pool is shared and small (`available_parallelism() - 1`,
18//! as few as one thread), `run_task` handlers must stay short and
19//! non-blocking: one plugin that blocks on I/O or a lock stalls every
20//! other instance's background work. Long or blocking work belongs on a
21//! plugin's own thread (`AudioTap::spawn_worker`), not the pool.
22
23use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering, fence};
24use std::sync::{Arc, OnceLock};
25use std::thread::{self, Thread};
26use std::time::Duration;
27
28use crossbeam_queue::ArrayQueue;
29
30/// Preallocated inbound-queue capacity per instance. Mirrors
31/// `EVENT_LIST_PREALLOC`: a block that schedules more tasks than this
32/// drops the overflow (`try_spawn` returns `Err`) rather than
33/// allocating on the audio thread.
34pub const TASK_QUEUE_PREALLOC: usize = 256;
35
36/// How many instances can have pending work queued in the pool at once.
37/// Sized well past any realistic simultaneous-instance count.
38const INJECTOR_CAP: usize = 4096;
39
40/// How long a worker parks before a defensive re-check. Wakes are
41/// explicit (`unpark` after a push), so this only bounds the worst case
42/// if an `unpark` is ever missed.
43const PARK_TIMEOUT: Duration = Duration::from_secs(1);
44
45/// A drainable instance queue, type-erased so the one pool holds many
46/// task types at once.
47trait Drain: Send + Sync {
48    fn drain(&self);
49}
50
51/// Per-instance inbound queue plus the monomorphized handler. Shared
52/// (`Arc`) between the schedulers (audio thread / editor / init, via
53/// [`TaskSpawner`]) and the pool worker that drains it.
54struct Sink<T: Send + 'static> {
55    /// `try_spawn`: FIFO, every queued task runs.
56    queue: ArrayQueue<T>,
57    /// `spawn_coalescing`: a single slot. `force_push` keeps only the
58    /// newest target, and `drain` runs it at most once, so a burst of
59    /// requests between two drains collapses to one execution instead of
60    /// running one build per intermediate target.
61    coalesced: ArrayQueue<T>,
62    /// Coalesces wake-ups: set when this sink is already queued in the
63    /// injector, so a burst of pushes injects it once.
64    scheduled: AtomicBool,
65    /// `run(task)` is `move |task| L::run_task(task, &params)`, built
66    /// once when the instance registers - never per task.
67    run: Box<dyn Fn(T) + Send + Sync>,
68}
69
70impl<T: Send + 'static> Sink<T> {
71    /// Run one task, catching panics so a bad handler can't kill the
72    /// shared worker (which would strand every other instance's tasks).
73    /// `run`/`task` are effectively unwind-safe: `run` is `&`-borrowed and
74    /// a poisoned task is simply dropped.
75    fn run_one(&self, task: T) {
76        let run = &self.run;
77        let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| run(task)));
78    }
79}
80
81impl<T: Send + 'static> Drain for Sink<T> {
82    fn drain(&self) {
83        // Clear the flag before draining so a task pushed mid-drain
84        // re-arms the sink instead of being stranded. This clear and the
85        // queue reads below form a StoreLoad pair against the producer's
86        // push + `arm` swap; Release/AcqRel don't order a store followed
87        // by a load of a *different* location, so without the SeqCst store
88        // + fence a worker could clear the flag, read the queue empty, and
89        // a concurrent producer could push a task and read the flag still
90        // `true` - stranding that task with nothing scheduled to drain it.
91        // Only SeqCst forbids that reordering. The stranded task self-heals
92        // on the sink's next `arm`, so continuous work (a knob sweep) is
93        // fine, but the last one-shot before an idle period would hang.
94        self.scheduled.store(false, Ordering::SeqCst);
95        fence(Ordering::SeqCst);
96        // The coalesced slot held only the newest target, so a burst of
97        // `spawn_coalescing` calls since the last drain runs once here, not
98        // once per intermediate target.
99        if let Some(task) = self.coalesced.pop() {
100            self.run_one(task);
101        }
102        // FIFO tasks each run.
103        while let Some(task) = self.queue.pop() {
104            self.run_one(task);
105        }
106    }
107}
108
109/// Worker-visible pool state: the injector of ready sinks. Held in an
110/// `Arc` so every worker closure can reach it.
111struct Shared {
112    injector: ArrayQueue<Arc<dyn Drain>>,
113    /// Round-robin cursor for choosing which worker to wake.
114    next: AtomicUsize,
115}
116
117struct Pool {
118    shared: Arc<Shared>,
119    workers: Vec<Thread>,
120}
121
122static POOL: OnceLock<Pool> = OnceLock::new();
123
124fn pool() -> &'static Pool {
125    POOL.get_or_init(|| {
126        let shared = Arc::new(Shared {
127            injector: ArrayQueue::new(INJECTOR_CAP),
128            next: AtomicUsize::new(0),
129        });
130        // One fewer than the core count, floored at one, so the pool
131        // never starves the audio and main threads on a small machine.
132        let n = thread::available_parallelism().map_or(1, |p| p.get().saturating_sub(1).max(1));
133        let mut workers = Vec::with_capacity(n);
134        for _ in 0..n {
135            let shared = Arc::clone(&shared);
136            let handle = thread::Builder::new()
137                .name("truce-task-pool".into())
138                .spawn(move || worker_loop(&shared))
139                .expect("spawn truce task-pool worker");
140            workers.push(handle.thread().clone());
141        }
142        Pool { shared, workers }
143    })
144}
145
146fn worker_loop(shared: &Shared) -> ! {
147    loop {
148        while let Some(sink) = shared.injector.pop() {
149            sink.drain();
150        }
151        // Nothing pending: park. A concurrent push + `unpark` either
152        // beats the park (the token makes this return at once) or wakes
153        // us; `PARK_TIMEOUT` is a belt-and-suspenders re-check.
154        thread::park_timeout(PARK_TIMEOUT);
155    }
156}
157
158/// Enqueue a ready sink and wake a worker. Wait-free: `injector.push`
159/// is lock-free and `unpark` is a bounded, non-blocking wake (the same
160/// primitive the manual worker pattern uses). Returns `false` if the
161/// injector is full so the caller can clear `scheduled` and let a later
162/// `arm` retry, rather than leaving the sink flagged-but-unqueued.
163fn schedule(sink: Arc<dyn Drain>) -> bool {
164    let pool = pool();
165    if pool.shared.injector.push(sink).is_err() {
166        return false;
167    }
168    if !pool.workers.is_empty() {
169        let i = pool.shared.next.fetch_add(1, Ordering::Relaxed) % pool.workers.len();
170        pool.workers[i].unpark();
171    }
172    true
173}
174
175/// A cheap-to-clone handle for scheduling background tasks onto the
176/// shared pool. Held by the shell and handed to the plugin through
177/// [`InitContext`], `ProcessContext`, and the editor's `PluginContext`.
178pub struct TaskSpawner<T: Send + 'static> {
179    sink: Arc<Sink<T>>,
180}
181
182impl<T: Send + 'static> Clone for TaskSpawner<T> {
183    fn clone(&self) -> Self {
184        Self {
185            sink: Arc::clone(&self.sink),
186        }
187    }
188}
189
190impl<T: Send + 'static> TaskSpawner<T> {
191    /// Register an instance's handler with the shared pool. `run` is the
192    /// monomorphized `move |task| L::run_task(task, &params)`, built once
193    /// by the shell. The pool itself is not started until the first task
194    /// is actually scheduled, so constructing a spawner for a plugin that
195    /// never schedules costs only the (small) inbound queue.
196    pub fn new(run: impl Fn(T) + Send + Sync + 'static) -> Self {
197        Self {
198            sink: Arc::new(Sink {
199                queue: ArrayQueue::new(TASK_QUEUE_PREALLOC),
200                coalesced: ArrayQueue::new(1),
201                scheduled: AtomicBool::new(false),
202                run: Box::new(run),
203            }),
204        }
205    }
206
207    /// Enqueue a task, running it on the pool as soon as a worker is
208    /// free. Wait-free. Returns `Err(task)` if the inbound queue is full
209    /// (the audio thread decides what to do - drop, or coalesce via
210    /// [`Self::spawn_coalescing`] - rather than block).
211    ///
212    /// # Errors
213    ///
214    /// Returns the task back when the preallocated inbound queue is full.
215    pub fn try_spawn(&self, task: T) -> Result<(), T> {
216        self.sink.queue.push(task)?;
217        self.arm();
218        Ok(())
219    }
220
221    /// Post a task into the single coalescing slot, replacing any
222    /// still-unrun target. Wait-free, never rejects. Only the newest
223    /// survives and the worker runs it at most once per drain, so a knob
224    /// sweep that outruns the handler collapses to one execution, not one
225    /// build per intermediate target. The displaced target drops on the
226    /// caller (the audio thread on the hot path), so a coalescing task
227    /// type should be cheap to drop - a small `Copy` request, not an
228    /// owned buffer.
229    pub fn spawn_coalescing(&self, task: T) {
230        let _ = self.sink.coalesced.force_push(task);
231        self.arm();
232    }
233
234    /// Inject this sink into the pool if it isn't already queued.
235    fn arm(&self) {
236        // Pairs with the SeqCst store + fence in `Sink::drain`: the caller
237        // pushed the task just before this, and that push must be ordered
238        // before the flag swap below, or the StoreLoad race described in
239        // `drain` strands the task. The fence + SeqCst swap give the total
240        // order that Release/AcqRel can't. Still wait-free (one barrier,
241        // no lock/alloc/syscall), and `arm` runs at most once per block.
242        fence(Ordering::SeqCst);
243        if !self.sink.scheduled.swap(true, Ordering::SeqCst) {
244            let sink: Arc<dyn Drain> = Arc::clone(&self.sink) as Arc<dyn Drain>;
245            if !schedule(sink) {
246                // Injector full: we flagged the sink but couldn't queue it.
247                // Clear the flag so the next `arm` re-attempts injection
248                // instead of skipping on a stale `true`.
249                self.sink.scheduled.store(false, Ordering::SeqCst);
250            }
251        }
252    }
253}
254
255/// A type-erased [`TaskSpawner`], so the concrete `ProcessContext` /
256/// `InitContext` (whose signatures are fixed by the leaf trait and can't
257/// name the plugin's task type) can carry the handle and hand back a
258/// typed spawner on demand via [`Self::downcast`].
259#[derive(Clone)]
260pub struct AnyTaskSpawner(Arc<dyn std::any::Any + Send + Sync>);
261
262impl AnyTaskSpawner {
263    /// Erase a typed spawner. Called by the shell when it wires tasks.
264    #[must_use]
265    pub fn new<T: Send + 'static>(spawner: &TaskSpawner<T>) -> Self {
266        Self(Arc::new(spawner.clone()) as Arc<dyn std::any::Any + Send + Sync>)
267    }
268
269    /// Recover the typed spawner. `None` if this handle was erased from a
270    /// different task type (a caller asking for the wrong `T`).
271    #[must_use]
272    pub fn downcast<T: Send + 'static>(&self) -> Option<TaskSpawner<T>> {
273        self.0.downcast_ref::<TaskSpawner<T>>().cloned()
274    }
275}
276
277/// Context handed to `init` so a plugin can schedule startup background
278/// work before the first block. Concrete (not generic over the task
279/// type) because `init`'s signature lives on the leaf trait; recover the
280/// typed spawner with [`Self::tasks`]. Params arrive as the separate
281/// `init` argument.
282pub struct InitContext {
283    tasks: Option<AnyTaskSpawner>,
284}
285
286impl InitContext {
287    #[must_use]
288    pub fn new(tasks: Option<AnyTaskSpawner>) -> Self {
289        Self { tasks }
290    }
291
292    /// The task spawner for this instance's `BackgroundTasks::Task`, or
293    /// `None` if the plugin wired no `tasks:` on `plugin!`.
294    #[must_use]
295    pub fn tasks<T: Send + 'static>(&self) -> Option<TaskSpawner<T>> {
296        self.tasks.as_ref().and_then(AnyTaskSpawner::downcast::<T>)
297    }
298}
299
300#[cfg(test)]
301mod tests {
302    // `TASK_QUEUE_PREALLOC` is 256, so casting it to `u32` for the loop
303    // bounds is always exact.
304    #![allow(clippy::cast_possible_truncation)]
305
306    use super::*;
307    use std::sync::atomic::AtomicU32;
308    use std::sync::{Condvar, Mutex};
309    use std::time::Instant;
310
311    fn wait_until(deadline: Duration, mut done: impl FnMut() -> bool) -> bool {
312        let start = Instant::now();
313        while start.elapsed() < deadline {
314            if done() {
315                return true;
316            }
317            thread::sleep(Duration::from_millis(1));
318        }
319        done()
320    }
321
322    /// Blocking completion latch: the pool handler bumps it, the test
323    /// blocks until a count is reached. Unlike the wall-clock `wait_until`,
324    /// it has no deadline, so it stays deterministic under Miri, whose
325    /// interpreter can't run background tasks within a real-time budget.
326    /// A `Mutex`/`Condvar` (not an `mpsc::Sender`, which is `!Sync`) keeps
327    /// the handler `Fn + Send + Sync`.
328    #[derive(Default)]
329    struct Latch {
330        ran: Mutex<u32>,
331        woke: Condvar,
332    }
333
334    impl Latch {
335        fn bump(&self) {
336            *self.ran.lock().unwrap() += 1;
337            self.woke.notify_all();
338        }
339
340        fn wait_for(&self, target: u32) {
341            let mut ran = self.ran.lock().unwrap();
342            while *ran < target {
343                ran = self.woke.wait(ran).unwrap();
344            }
345        }
346    }
347
348    #[test]
349    fn runs_scheduled_tasks_off_thread() {
350        let latch = Arc::new(Latch::default());
351        let sum = Arc::new(AtomicU32::new(0));
352        let (l, s) = (Arc::clone(&latch), Arc::clone(&sum));
353        let spawner = TaskSpawner::<u32>::new(move |n| {
354            s.fetch_add(n, Ordering::Relaxed);
355            l.bump();
356        });
357
358        for n in 1..=10 {
359            spawner.try_spawn(n).expect("queue has room");
360        }
361
362        // Block until all ten ran. The latch mutex orders every handler's
363        // `sum` write before the read below, so the Relaxed sum is exact.
364        latch.wait_for(10);
365        assert_eq!(sum.load(Ordering::Relaxed), 55, "all ten tasks ran");
366    }
367
368    #[test]
369    fn full_queue_returns_the_task() {
370        // Handler blocks on a gate so the queue can actually fill.
371        let gate = Arc::new(AtomicBool::new(false));
372        let g = Arc::clone(&gate);
373        let spawner = TaskSpawner::<u32>::new(move |_| {
374            while !g.load(Ordering::Acquire) {
375                thread::sleep(Duration::from_millis(1));
376            }
377        });
378
379        // First task is picked up and blocks a worker; fill the rest.
380        let mut rejected = 0u32;
381        for n in 0..(TASK_QUEUE_PREALLOC as u32 + 64) {
382            if spawner.try_spawn(n).is_err() {
383                rejected += 1;
384            }
385        }
386        assert!(rejected > 0, "a full inbound queue rejects further tasks");
387        gate.store(true, Ordering::Release);
388    }
389
390    #[test]
391    fn panicking_task_does_not_kill_the_worker() {
392        let latch = Arc::new(Latch::default());
393        let l = Arc::clone(&latch);
394        let spawner = TaskSpawner::<bool>::new(move |should_panic| {
395            assert!(!should_panic, "intentional panic, caught by the pool");
396            l.bump();
397        });
398        spawner.try_spawn(true).expect("queue has room"); // panics in the handler
399        spawner.try_spawn(false).expect("queue has room"); // must still run
400        // If the panic had killed the worker, the survivor never runs and
401        // this blocks forever - surfaced as a hung test, not a false pass.
402        latch.wait_for(1);
403    }
404
405    #[test]
406    fn coalescing_never_rejects() {
407        let last = Arc::new(AtomicU32::new(0));
408        let l = Arc::clone(&last);
409        let spawner = TaskSpawner::<u32>::new(move |n| {
410            l.store(n, Ordering::Relaxed);
411        });
412        for n in 0..(TASK_QUEUE_PREALLOC as u32 * 4) {
413            spawner.spawn_coalescing(n); // never panics, never blocks
414        }
415        let target = TASK_QUEUE_PREALLOC as u32 * 4 - 1;
416        assert!(
417            wait_until(Duration::from_secs(2), || last.load(Ordering::Relaxed)
418                == target),
419            "the newest task always runs"
420        );
421    }
422
423    #[test]
424    fn coalescing_collapses_to_the_newest() {
425        let runs = Arc::new(AtomicU32::new(0));
426        let last = Arc::new(AtomicU32::new(0));
427        let (r, l) = (Arc::clone(&runs), Arc::clone(&last));
428        let spawner = TaskSpawner::<u32>::new(move |n| {
429            r.fetch_add(1, Ordering::Relaxed);
430            l.store(n, Ordering::Relaxed);
431        });
432
433        // Fill the coalescing slot repeatedly without arming the pool, so
434        // the burst collapses in the slot rather than racing a worker.
435        // Then drain once and confirm the whole burst ran a single time,
436        // as the newest target.
437        for n in 1..=1000 {
438            let _ = spawner.sink.coalesced.force_push(n);
439        }
440        spawner.sink.drain();
441
442        assert_eq!(runs.load(Ordering::Relaxed), 1, "the burst ran once");
443        assert_eq!(last.load(Ordering::Relaxed), 1000, "and it was the newest");
444    }
445}
446
447// Model-checked proof that the schedule/drain handshake can't strand a
448// task under any thread interleaving. Run with:
449//   cargo test -p truce-core --features loom loom
450//
451// loom can't see into crossbeam's `ArrayQueue`, so this models the
452// protocol directly: a one-slot queue (`item`) plus the `scheduled` flag,
453// driven through the exact SeqCst store / fence / swap sequence that
454// `Sink::drain` and `TaskSpawner::arm` use. Weakening either side to
455// Release/AcqRel (dropping the SeqCst or the fences) makes loom find the
456// stranding interleaving; the version below passes.
457#[cfg(all(test, feature = "loom"))]
458mod loom_tests {
459    use loom::sync::Arc;
460    use loom::sync::atomic::{AtomicBool, Ordering, fence};
461    use loom::thread;
462
463    #[test]
464    fn schedule_drain_never_strands_a_task() {
465        loom::model(|| {
466            // Start with a drain in flight: the sink was scheduled
467            // (`flag == true`) and a worker is about to drain an empty
468            // queue, concurrent with a producer pushing one more task.
469            let flag = Arc::new(AtomicBool::new(true));
470            let item = Arc::new(AtomicBool::new(false));
471
472            let (f, i) = (flag.clone(), item.clone());
473            let worker = thread::spawn(move || {
474                // `Sink::drain`: clear the flag, then check the queue. The
475                // presence check must be a plain load (a `pop` reading
476                // empty) - an RMW would always read the latest value in
477                // modification order and so hide the StoreLoad staleness
478                // this test exists to catch.
479                f.store(false, Ordering::SeqCst);
480                fence(Ordering::SeqCst);
481                if i.load(Ordering::Acquire) {
482                    i.store(false, Ordering::Release); // popped it
483                }
484            });
485
486            // `try_spawn` + `arm`: push the task, then flag the sink.
487            item.store(true, Ordering::Release);
488            fence(Ordering::SeqCst);
489            let was_scheduled = flag.swap(true, Ordering::SeqCst);
490            // `was_scheduled == false` => the producer injects a fresh
491            // drain (it set `flag = true`). `true` => it relies on the
492            // in-flight drain to pick the task up.
493            let _ = was_scheduled;
494
495            worker.join().unwrap();
496
497            // Safe end states: the queue is empty (some drain popped it),
498            // or a drain is still scheduled (`flag == true`) to pick it up.
499            // A pending task with `flag == false` is the stranding bug.
500            let pending = item.load(Ordering::SeqCst);
501            let scheduled = flag.load(Ordering::SeqCst);
502            assert!(
503                !pending || scheduled,
504                "task stranded: pending with scheduled == false"
505            );
506        });
507    }
508}