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, ¶ms)`, 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, ¶ms)`, 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}