Skip to main content

truce_core/
tasks.rs

1//! Managed background-task pool.
2//!
3//! A process-global pool of worker threads runs plugin
4//! `BackgroundTask::run` handlers off the audio thread. Each plugin
5//! instance owns a
6//! preallocated, wait-free inbound queue via a [`TaskSpawner`]: the
7//! audio thread (or the editor, or `init`) pushes tasks without
8//! allocating or blocking, and a pool worker drains them. Feedback to
9//! the audio thread stays the plugin's job through shared `#[skip]`
10//! channels - the pool owns only the worker threads and the inbound
11//! queue.
12//!
13//! ## Concurrency
14//!
15//! By default drains are **not** mutually exclusive: the stranding-
16//! avoidance handshake clears a sink's `scheduled` flag before draining,
17//! so a burst that re-arms an instance mid-drain can hand a second idle
18//! worker the same sink - `run` can run concurrently with itself for
19//! one instance. Handlers must therefore be reentrancy-safe: talk to the
20//! audio thread only through lock-free / atomic channels (the reverb
21//! example's MPMC handoff), or guard shared mutable state (the
22//! `AudioTap::drain_with` `try_lock` idiom). A plugin that can't meet that
23//! contract sets `BackgroundTask::SERIALIZED = true`, and the pool then
24//! runs that instance's handler one at a time.
25//!
26//! The pool is shared across every instance in the process (one small
27//! set of threads, not one thread per instance) and initializes lazily
28//! the first time any instance actually schedules a task, so a plugin
29//! that never declares a `BackgroundTask` spawns no threads.
30//!
31//! Because the pool is shared and small (`available_parallelism() - 1`,
32//! as few as one thread), task handlers must stay short and
33//! non-blocking: one plugin that blocks on I/O or a lock stalls every
34//! other instance's background work. Long or blocking work belongs on a
35//! plugin's own thread (`AudioTap::spawn_worker`), not the pool.
36
37// Gated on `not(miri)` to match `pin_current_module`, which is a no-op
38// under Miri (no dynamic loader to pin), so these FFI types aren't used.
39#[cfg(all(any(unix, windows), not(miri)))]
40use std::ffi::c_void;
41#[cfg(all(unix, not(miri)))]
42use std::ffi::{c_char, c_int};
43use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering, fence};
44use std::sync::{Arc, OnceLock};
45use std::thread::{self, Thread};
46use std::time::Duration;
47
48use crossbeam_queue::ArrayQueue;
49
50use crate::snapshot::SnapshotPublisher;
51
52/// Preallocated inbound-queue capacity per instance. Mirrors
53/// `EVENT_LIST_PREALLOC`: a block that schedules more tasks than this
54/// drops the overflow (`try_spawn` returns `Err`) rather than
55/// allocating on the audio thread.
56pub const TASK_QUEUE_PREALLOC: usize = 256;
57
58/// How many instances can have pending work queued in the pool at once.
59/// Sized well past any realistic simultaneous-instance count.
60const INJECTOR_CAP: usize = 4096;
61
62/// How long a worker parks before a defensive re-check. Wakes are
63/// explicit (`unpark` after a push), so this only bounds the worst case
64/// if an `unpark` is ever missed.
65const PARK_TIMEOUT: Duration = Duration::from_secs(1);
66
67/// A drainable instance queue, type-erased so the one pool holds many
68/// task types at once.
69trait Drain: Send + Sync {
70    fn drain(&self);
71}
72
73/// Per-instance inbound queue plus the monomorphized handler. Shared
74/// (`Arc`) between the schedulers (audio thread / editor / init, via
75/// [`TaskSpawner`]) and the pool worker that drains it.
76struct Sink<T: Send + 'static> {
77    /// `try_spawn`: FIFO, every queued task runs.
78    queue: ArrayQueue<T>,
79    /// `spawn_coalescing`: a single slot. `force_push` keeps only the
80    /// newest target, and `drain` runs it at most once, so a burst of
81    /// requests between two drains collapses to one execution instead of
82    /// running one build per intermediate target.
83    coalesced: ArrayQueue<T>,
84    /// Coalesces wake-ups: set when this sink is already queued in the
85    /// injector, so a burst of pushes injects it once.
86    scheduled: AtomicBool,
87    /// Serialized ("one-slot") mode: when set, at most one worker runs
88    /// `run` for this sink at a time. `false` (default) lets a second
89    /// worker drain concurrently for throughput.
90    serialized: bool,
91    /// Exclusive-drain guard for [`Self::serialized`]. A worker that finds
92    /// it already held bows out; the holder's re-check loop in `drain`
93    /// picks up whatever the bower-out was injected for, so nothing is
94    /// stranded. Unused in the concurrent (default) mode.
95    draining: AtomicBool,
96    /// `run(task)` is `move |task| task.run(&params)`, built
97    /// once when the instance registers - never per task.
98    run: Box<dyn Fn(T) + Send + Sync>,
99}
100
101impl<T: Send + 'static> Sink<T> {
102    /// Run one task, catching panics so a bad handler can't kill the
103    /// shared worker (which would strand every other instance's tasks).
104    /// `run`/`task` are effectively unwind-safe: `run` is `&`-borrowed and
105    /// a poisoned task is simply dropped.
106    fn run_one(&self, task: T) {
107        let run = &self.run;
108        let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| run(task)));
109    }
110}
111
112impl<T: Send + 'static> Sink<T> {
113    /// Clear `scheduled`, then run every currently-queued task once. The
114    /// clear-before-drain + `SeqCst` fence is the stranding-avoidance
115    /// handshake: a task pushed mid-drain re-arms the sink (its `arm` swap
116    /// sees `scheduled == false`) instead of being stranded. Release/AcqRel
117    /// don't order a store followed by a load of a *different* location, so
118    /// without the `SeqCst` store + fence a worker could clear the flag, read
119    /// the queue empty, and a concurrent producer could push a task and read
120    /// the flag still `true` - stranding it. Only `SeqCst` forbids that.
121    fn drain_queues(&self) {
122        self.scheduled.store(false, Ordering::SeqCst);
123        fence(Ordering::SeqCst);
124        // The coalesced slot held only the newest target, so a burst of
125        // `spawn_coalescing` calls since the last drain runs once here, not
126        // once per intermediate target.
127        if let Some(task) = self.coalesced.pop() {
128            self.run_one(task);
129        }
130        // FIFO tasks each run.
131        while let Some(task) = self.queue.pop() {
132            self.run_one(task);
133        }
134    }
135}
136
137impl<T: Send + 'static> Drain for Sink<T> {
138    fn drain(&self) {
139        // Concurrent (default) mode: a second worker may drain this sink at
140        // the same time. Handlers must be reentrancy-safe (see
141        // `BackgroundTask::SERIALIZED`).
142        if !self.serialized {
143            self.drain_queues();
144            return;
145        }
146        // Serialized ("one-slot") mode: run the handler for this instance on
147        // at most one worker at a time. A worker that finds the guard held
148        // is inert and returns - it touches nothing, so the `scheduled`
149        // handshake below stays the sole no-stranding signal, exactly as in
150        // the concurrent path.
151        if self.draining.swap(true, Ordering::Acquire) {
152            return;
153        }
154        loop {
155            self.drain_queues();
156            self.draining.store(false, Ordering::Release);
157            fence(Ordering::SeqCst);
158            // Re-check the *scheduled flag*, not the queue: a producer that
159            // armed during the drain set it with a SeqCst swap ordered after
160            // `drain_queues`'s SeqCst clear, so we either observe it here and
161            // re-drain, or its swap saw our clear and injected a fresh drain.
162            // (Reading the queue instead would race - a plain queue load
163            // isn't synchronized with the producer's push, so it could miss a
164            // task a re-injection carried and strand it.) `drain_queues`
165            // clears the flag each pass, so the loop makes progress and can't
166            // spin: at most one extra empty drain after the last arm.
167            if !self.scheduled.load(Ordering::SeqCst) {
168                return;
169            }
170            // Work remains. Re-take the guard and drain again; if another
171            // worker took it first, that worker now owns the remainder.
172            if self.draining.swap(true, Ordering::Acquire) {
173                return;
174            }
175        }
176    }
177}
178
179/// Worker-visible pool state: the injector of ready sinks. Held in an
180/// `Arc` so every worker closure can reach it.
181struct Shared {
182    injector: ArrayQueue<Arc<dyn Drain>>,
183    /// Round-robin cursor for choosing which worker to wake.
184    next: AtomicUsize,
185}
186
187struct Pool {
188    shared: Arc<Shared>,
189    workers: Vec<Thread>,
190}
191
192static POOL: OnceLock<Pool> = OnceLock::new();
193
194fn pool() -> &'static Pool {
195    POOL.get_or_init(|| {
196        let shared = Arc::new(Shared {
197            injector: ArrayQueue::new(INJECTOR_CAP),
198            next: AtomicUsize::new(0),
199        });
200        // One fewer than the core count, floored at one, so the pool
201        // never starves the audio and main threads on a small machine.
202        let n = thread::available_parallelism().map_or(1, |p| p.get().saturating_sub(1).max(1));
203        let mut workers = Vec::with_capacity(n);
204        for _ in 0..n {
205            let shared = Arc::clone(&shared);
206            match thread::Builder::new()
207                .name("truce-task-pool".into())
208                .spawn(move || worker_loop(&shared))
209            {
210                Ok(handle) => workers.push(handle.thread().clone()),
211                // A failed spawn (thread/memory exhaustion) must not panic:
212                // pool init can run behind an `extern "C"` boundary in a
213                // host that doesn't catch unwinds (VST3 / VST2 / AAX / LV2),
214                // where an unwind aborts the whole DAW. Keep whatever
215                // workers spawned; if none did, `schedule` drops tasks
216                // instead of queueing work nothing will drain.
217                Err(e) => {
218                    eprintln!("[truce] task-pool worker spawn failed: {e}");
219                    break;
220                }
221            }
222        }
223        // The workers loop forever with no shutdown path, so pin this
224        // module in memory: a host that `dlclose`s it on last-instance
225        // teardown must not leave a parked worker to wake into unmapped
226        // code. Only once we actually have workers to protect.
227        if !workers.is_empty() {
228            pin_current_module();
229        }
230        Pool { shared, workers }
231    })
232}
233
234/// Pin the module `truce-core` is linked into (the plugin cdylib in a
235/// static build, or the hot-reload shell cdylib) so it is never
236/// unmapped. The pool's workers loop forever - park + drain - with no
237/// shutdown path and no `JoinHandle`s; if the host `dlclose`d the module
238/// on last-instance teardown while a worker was parked, it would wake
239/// with its program counter in unmapped code and take down the DAW
240/// (Bitwig, Cubase, and JUCE hosts all unload). Pinning trades a small
241/// permanent mapping for eliminating that crash class - the standard fix
242/// for a persistent plugin helper thread. Called once, and only when the
243/// pool actually spawned workers.
244#[cfg(all(unix, not(miri)))]
245fn pin_current_module() {
246    // Field names mirror the platform's `Dl_info`; only `dli_fname` is
247    // read, the rest are here for correct C layout (`dladdr` writes all).
248    #[repr(C)]
249    #[allow(clippy::struct_field_names)]
250    struct DlInfo {
251        dli_fname: *const c_char,
252        dli_fbase: *mut c_void,
253        dli_sname: *const c_char,
254        dli_saddr: *mut c_void,
255    }
256    unsafe extern "C" {
257        fn dladdr(addr: *const c_void, info: *mut DlInfo) -> c_int;
258        fn dlopen(filename: *const c_char, flags: c_int) -> *mut c_void;
259    }
260    // RTLD_NOLOAD | RTLD_NODELETE: re-reference the already-loaded object
261    // and mark it non-deletable. The flag values differ between glibc and
262    // Darwin.
263    #[cfg(target_os = "linux")]
264    const FLAGS: c_int = 0x0004 | 0x1000;
265    #[cfg(not(target_os = "linux"))]
266    const FLAGS: c_int = 0x0010 | 0x0080;
267
268    // SAFETY: `dladdr` reads the address of a live function in this module
269    // and fills `info` with loader-owned strings valid for the immediate
270    // `dlopen`. NOLOAD only re-references the already-mapped object; the
271    // returned handle is intentionally leaked so NODELETE persists.
272    unsafe {
273        let mut info: DlInfo = core::mem::zeroed();
274        if dladdr(pin_current_module as *const c_void, &raw mut info) != 0
275            && !info.dli_fname.is_null()
276        {
277            let _ = dlopen(info.dli_fname, FLAGS);
278        }
279    }
280}
281
282#[cfg(all(windows, not(miri)))]
283fn pin_current_module() {
284    unsafe extern "system" {
285        fn GetModuleHandleExW(flags: u32, name: *const u16, module: *mut *mut c_void) -> i32;
286    }
287    // GET_MODULE_HANDLE_EX_FLAG_PIN | ..._FROM_ADDRESS: interpret `name`
288    // as an address inside this module and bump its load count so it
289    // never unloads.
290    const PIN_FROM_ADDRESS: u32 = 0x0000_0001 | 0x0000_0004;
291
292    // SAFETY: `name` is the address of a live function in this module;
293    // `module` receives the pinned handle, which we intentionally leak.
294    unsafe {
295        let mut module: *mut c_void = core::ptr::null_mut();
296        let _ = GetModuleHandleExW(
297            PIN_FROM_ADDRESS,
298            pin_current_module as *const u16,
299            &raw mut module,
300        );
301    }
302}
303
304// No-op where there's no dynamic loader to pin against, and under Miri -
305// which has no dlopen/dlclose (so nothing can unload the module) and
306// doesn't support `dladdr` / `GetModuleHandleExW`.
307#[cfg(any(miri, not(any(unix, windows))))]
308fn pin_current_module() {}
309
310/// Eagerly start the shared pool on the calling thread. The shell calls
311/// this at instantiation (the host/main thread) when a plugin wires a
312/// task spawner, so the worker threads exist before the audio thread ever
313/// schedules. Without it a plugin that first schedules from `process()`
314/// (the "rebuild the filter when a knob moves" pattern, with no startup
315/// work in `init` to warm the pool) would cold-start the threads inside
316/// the audio callback. Idempotent: the pool is a process-global singleton
317/// after the first call.
318pub fn warm_pool() {
319    let _ = pool();
320}
321
322fn worker_loop(shared: &Shared) -> ! {
323    loop {
324        while let Some(sink) = shared.injector.pop() {
325            sink.drain();
326        }
327        // Nothing pending: park. A concurrent push + `unpark` either
328        // beats the park (the token makes this return at once) or wakes
329        // us; `PARK_TIMEOUT` is a belt-and-suspenders re-check.
330        thread::park_timeout(PARK_TIMEOUT);
331    }
332}
333
334/// Enqueue a ready sink and wake a worker. Wait-free: `injector.push`
335/// is lock-free and `unpark` is a bounded, non-blocking wake (the same
336/// primitive the manual worker pattern uses). Returns `false` if the
337/// injector is full so the caller can clear `scheduled` and let a later
338/// `arm` retry, rather than leaving the sink flagged-but-unqueued.
339fn schedule(sink: Arc<dyn Drain>) -> bool {
340    let pool = pool();
341    // No workers (every spawn failed at pool init): drop the task rather
342    // than queue work nothing will ever drain, matching the "queue full ->
343    // drop" policy. The caller clears `scheduled` so a later `arm` retries.
344    if pool.workers.is_empty() {
345        return false;
346    }
347    if pool.shared.injector.push(sink).is_err() {
348        return false;
349    }
350    let i = pool.shared.next.fetch_add(1, Ordering::Relaxed) % pool.workers.len();
351    pool.workers[i].unpark();
352    true
353}
354
355/// A cheap-to-clone handle for scheduling background tasks onto the
356/// shared pool. Held by the shell and handed to the plugin through
357/// [`InitContext`], `ProcessContext`, and the editor's `PluginContext`.
358///
359/// A spawner built with [`Self::new`] may run its handler concurrently
360/// with itself for one instance (see the module's Concurrency section);
361/// [`Self::new_serialized`] runs it one at a time.
362pub struct TaskSpawner<T: Send + 'static> {
363    sink: Arc<Sink<T>>,
364}
365
366impl<T: Send + 'static> Clone for TaskSpawner<T> {
367    fn clone(&self) -> Self {
368        Self {
369            sink: Arc::clone(&self.sink),
370        }
371    }
372}
373
374impl<T: Send + 'static> TaskSpawner<T> {
375    /// Register an instance's handler with the shared pool. `run` is the
376    /// monomorphized `move |task| task.run(&params)`, built once
377    /// by the shell. The pool itself is not started until the first task
378    /// is actually scheduled, so constructing a spawner for a plugin that
379    /// never schedules costs only the (small) inbound queue.
380    ///
381    /// The handler may run concurrently with itself for one instance; use
382    /// [`Self::new_serialized`] for a handler that isn't reentrancy-safe.
383    pub fn new(run: impl Fn(T) + Send + Sync + 'static) -> Self {
384        Self::with_mode(run, false)
385    }
386
387    /// Like [`Self::new`], but the pool runs the handler for a given
388    /// instance one at a time ("one-slot" mode). The shell selects this
389    /// when the plugin's `BackgroundTask::SERIALIZED` is `true`.
390    pub fn new_serialized(run: impl Fn(T) + Send + Sync + 'static) -> Self {
391        Self::with_mode(run, true)
392    }
393
394    fn with_mode(run: impl Fn(T) + Send + Sync + 'static, serialized: bool) -> Self {
395        Self {
396            sink: Arc::new(Sink {
397                queue: ArrayQueue::new(TASK_QUEUE_PREALLOC),
398                coalesced: ArrayQueue::new(1),
399                scheduled: AtomicBool::new(false),
400                serialized,
401                draining: AtomicBool::new(false),
402                run: Box::new(run),
403            }),
404        }
405    }
406
407    /// Enqueue a task, running it on the pool as soon as a worker is
408    /// free. Wait-free. Returns `Err(task)` if the inbound queue is full
409    /// (the audio thread decides what to do - drop, or coalesce via
410    /// [`Self::spawn_coalescing`] - rather than block).
411    ///
412    /// # Errors
413    ///
414    /// Returns the task back when the preallocated inbound queue is full.
415    pub fn try_spawn(&self, task: T) -> Result<(), T> {
416        self.sink.queue.push(task)?;
417        self.arm();
418        Ok(())
419    }
420
421    /// Post a task into the single coalescing slot, replacing any
422    /// still-unrun target. Wait-free, never rejects. Only the newest
423    /// survives and the worker runs it at most once per drain, so a knob
424    /// sweep that outruns the handler collapses to one execution, not one
425    /// build per intermediate target. The displaced target drops on the
426    /// caller (the audio thread on the hot path), so a coalescing task
427    /// type should be cheap to drop - a small `Copy` request, not an
428    /// owned buffer.
429    pub fn spawn_coalescing(&self, task: T) {
430        let _ = self.sink.coalesced.force_push(task);
431        self.arm();
432    }
433
434    /// Inject this sink into the pool if it isn't already queued.
435    fn arm(&self) {
436        // Pairs with the SeqCst store + fence in `Sink::drain`: the caller
437        // pushed the task just before this, and that push must be ordered
438        // before the flag swap below, or the StoreLoad race described in
439        // `drain` strands the task. The fence + SeqCst swap give the total
440        // order that Release/AcqRel can't. Still wait-free (one barrier,
441        // no lock/alloc/syscall), and `arm` runs at most once per block.
442        fence(Ordering::SeqCst);
443        if !self.sink.scheduled.swap(true, Ordering::SeqCst) {
444            let sink: Arc<dyn Drain> = Arc::clone(&self.sink) as Arc<dyn Drain>;
445            if !schedule(sink) {
446                // Injector full: we flagged the sink but couldn't queue it.
447                // Clear the flag so the next `arm` re-attempts injection
448                // instead of skipping on a stale `true`.
449                self.sink.scheduled.store(false, Ordering::SeqCst);
450            }
451        }
452    }
453}
454
455/// One type-erased lane. Each element of [`AnyTaskSpawner`] holds one
456/// `TaskSpawner<T>` for a distinct task type.
457type ErasedLane = Arc<dyn std::any::Any + Send + Sync>;
458
459/// A bundle of type-erased [`TaskSpawner`]s - one lane per declared task
460/// type - so the concrete `ProcessContext` / `InitContext` (whose
461/// signatures are fixed by the leaf trait and can't name the plugin's task
462/// types) can carry every lane and hand back the right typed spawner on
463/// demand via [`Self::downcast`]. Cheap to clone (one `Arc`).
464#[derive(Clone)]
465pub struct AnyTaskSpawner(Arc<[ErasedLane]>);
466
467impl AnyTaskSpawner {
468    /// Erase a single typed spawner into a one-lane bundle.
469    #[must_use]
470    pub fn new<T: Send + 'static>(spawner: &TaskSpawner<T>) -> Self {
471        Self(Arc::from(vec![Arc::new(spawner.clone()) as ErasedLane]))
472    }
473
474    /// Bundle several already-erased lanes (one per task type). The
475    /// `plugin!` macro builds the lanes with [`TaskSpawnerBundle`].
476    #[must_use]
477    pub fn from_lanes(lanes: Vec<ErasedLane>) -> Self {
478        Self(Arc::from(lanes))
479    }
480
481    /// Recover the typed spawner for task type `T`, or `None` if no lane of
482    /// that type was declared. Lanes have distinct types, so at most one
483    /// matches.
484    #[must_use]
485    pub fn downcast<T: Send + 'static>(&self) -> Option<TaskSpawner<T>> {
486        self.0
487            .iter()
488            .find_map(|lane| lane.downcast_ref::<TaskSpawner<T>>().cloned())
489    }
490}
491
492/// Builder the `plugin!` macro uses to collect one lane per declared task
493/// type into an [`AnyTaskSpawner`]. Kept separate so the macro never has to
494/// name the erased-lane type.
495#[derive(Default)]
496pub struct TaskSpawnerBundle(Vec<ErasedLane>);
497
498impl TaskSpawnerBundle {
499    #[must_use]
500    pub fn new() -> Self {
501        Self(Vec::new())
502    }
503
504    /// Add one task type's spawner to the bundle.
505    pub fn push<T: Send + 'static>(&mut self, spawner: TaskSpawner<T>) {
506        self.0.push(Arc::new(spawner) as ErasedLane);
507    }
508
509    /// Finish: `Some` bundle, or `None` when no lanes were added (a plugin
510    /// that declared no tasks), matching the `Option<AnyTaskSpawner>` the
511    /// shell threads through.
512    #[must_use]
513    pub fn into_any(self) -> Option<AnyTaskSpawner> {
514        if self.0.is_empty() {
515            None
516        } else {
517            Some(AnyTaskSpawner::from_lanes(self.0))
518        }
519    }
520}
521
522/// Context handed to `init` so a plugin can schedule startup background
523/// work before the first block. Concrete (not generic over the task
524/// type) because `init`'s signature lives on the leaf trait; recover the
525/// typed spawner with [`Self::tasks`]. Params arrive as the separate
526/// `init` argument.
527pub struct InitContext {
528    tasks: Option<AnyTaskSpawner>,
529    /// Handle for the off-thread snapshot lane (large state save), when
530    /// the shell wired one. `None` in `--shell` hot-reload builds, where
531    /// (like `tasks`) it isn't threaded across the dylib boundary yet -
532    /// so the off-thread snapshot path is a static-build feature.
533    snapshot: Option<SnapshotPublisher>,
534}
535
536impl InitContext {
537    #[must_use]
538    pub fn new(tasks: Option<AnyTaskSpawner>) -> Self {
539        Self {
540            tasks,
541            snapshot: None,
542        }
543    }
544
545    /// Attach the off-thread snapshot publisher (see [`SnapshotPublisher`]).
546    #[must_use]
547    pub fn with_snapshot(mut self, snapshot: SnapshotPublisher) -> Self {
548        self.snapshot = Some(snapshot);
549        self
550    }
551
552    /// The task spawner for task type `T`, or `None` if the plugin declared
553    /// no `tasks:` lane of that type on `plugin!`.
554    #[must_use]
555    pub fn tasks<T: Send + 'static>(&self) -> Option<TaskSpawner<T>> {
556        self.tasks.as_ref().and_then(AnyTaskSpawner::downcast::<T>)
557    }
558
559    /// Handle to publish large custom state off the audio thread. Stash it
560    /// in your DSP state and call `publish` from a background-task handler
561    /// after your state changes. `None` in `--shell` builds; use
562    /// `snapshot_into` for small state, which works everywhere.
563    #[must_use]
564    pub fn snapshot_publisher(&self) -> Option<SnapshotPublisher> {
565        self.snapshot.clone()
566    }
567}
568
569#[cfg(test)]
570mod tests {
571    // `TASK_QUEUE_PREALLOC` is 256, so casting it to `u32` for the loop
572    // bounds is always exact.
573    #![allow(clippy::cast_possible_truncation)]
574
575    use super::*;
576    use crate::snapshot::{SnapshotPublisher, SnapshotSlot};
577    use std::sync::atomic::AtomicU32;
578    use std::sync::{Condvar, Mutex};
579    use std::time::Instant;
580
581    #[test]
582    fn init_context_exposes_snapshot_publisher() {
583        let slot = SnapshotSlot::new();
584        let cx = InitContext::new(None).with_snapshot(SnapshotPublisher::new(&slot));
585        // The plugin captures this in `init` and publishes large state
586        // through it off the audio thread.
587        cx.snapshot_publisher()
588            .expect("publisher present")
589            .publish(vec![1, 2, 3]);
590        assert_eq!(slot.read(), Some(vec![1, 2, 3]));
591        // No snapshot wired (the `--shell` / no-slot case) yields None.
592        assert!(InitContext::new(None).snapshot_publisher().is_none());
593    }
594
595    fn wait_until(deadline: Duration, mut done: impl FnMut() -> bool) -> bool {
596        let start = Instant::now();
597        while start.elapsed() < deadline {
598            if done() {
599                return true;
600            }
601            thread::sleep(Duration::from_millis(1));
602        }
603        done()
604    }
605
606    /// Blocking completion latch: the pool handler bumps it, the test
607    /// blocks until a count is reached. Unlike the wall-clock `wait_until`,
608    /// it has no deadline, so it stays deterministic under Miri, whose
609    /// interpreter can't run background tasks within a real-time budget.
610    /// A `Mutex`/`Condvar` (not an `mpsc::Sender`, which is `!Sync`) keeps
611    /// the handler `Fn + Send + Sync`.
612    #[derive(Default)]
613    struct Latch {
614        ran: Mutex<u32>,
615        woke: Condvar,
616    }
617
618    impl Latch {
619        fn bump(&self) {
620            *self.ran.lock().unwrap() += 1;
621            self.woke.notify_all();
622        }
623
624        fn wait_for(&self, target: u32) {
625            let mut ran = self.ran.lock().unwrap();
626            while *ran < target {
627                ran = self.woke.wait(ran).unwrap();
628            }
629        }
630    }
631
632    #[test]
633    fn warm_pool_starts_workers_and_is_idempotent() {
634        // Warming off the audio thread is what keeps the first
635        // audio-thread schedule from cold-starting the workers inline.
636        warm_pool();
637        warm_pool();
638        assert!(
639            !pool().workers.is_empty(),
640            "warming spawns at least one worker"
641        );
642    }
643
644    #[test]
645    fn pin_current_module_is_safe_and_idempotent() {
646        // Smoke: pinning must never panic or crash. In the test binary
647        // the loader query resolves to the harness executable (never
648        // unloaded anyway); we only assert the FFI is benign and can run
649        // more than once.
650        super::pin_current_module();
651        super::pin_current_module();
652    }
653
654    #[test]
655    fn runs_scheduled_tasks_off_thread() {
656        let latch = Arc::new(Latch::default());
657        let sum = Arc::new(AtomicU32::new(0));
658        let (l, s) = (Arc::clone(&latch), Arc::clone(&sum));
659        let spawner = TaskSpawner::<u32>::new(move |n| {
660            s.fetch_add(n, Ordering::Relaxed);
661            l.bump();
662        });
663
664        for n in 1..=10 {
665            spawner.try_spawn(n).expect("queue has room");
666        }
667
668        // Block until all ten ran. The latch mutex orders every handler's
669        // `sum` write before the read below, so the Relaxed sum is exact.
670        latch.wait_for(10);
671        assert_eq!(sum.load(Ordering::Relaxed), 55, "all ten tasks ran");
672    }
673
674    #[test]
675    fn full_queue_returns_the_task() {
676        // Handler blocks on a gate so the queue can actually fill.
677        let gate = Arc::new(AtomicBool::new(false));
678        let g = Arc::clone(&gate);
679        let spawner = TaskSpawner::<u32>::new(move |_| {
680            while !g.load(Ordering::Acquire) {
681                thread::sleep(Duration::from_millis(1));
682            }
683        });
684
685        // First task is picked up and blocks a worker; fill the rest.
686        let mut rejected = 0u32;
687        for n in 0..(TASK_QUEUE_PREALLOC as u32 + 64) {
688            if spawner.try_spawn(n).is_err() {
689                rejected += 1;
690            }
691        }
692        assert!(rejected > 0, "a full inbound queue rejects further tasks");
693        gate.store(true, Ordering::Release);
694    }
695
696    #[test]
697    fn panicking_task_does_not_kill_the_worker() {
698        let latch = Arc::new(Latch::default());
699        let l = Arc::clone(&latch);
700        let spawner = TaskSpawner::<bool>::new(move |should_panic| {
701            assert!(!should_panic, "intentional panic, caught by the pool");
702            l.bump();
703        });
704        spawner.try_spawn(true).expect("queue has room"); // panics in the handler
705        spawner.try_spawn(false).expect("queue has room"); // must still run
706        // If the panic had killed the worker, the survivor never runs and
707        // this blocks forever - surfaced as a hung test, not a false pass.
708        latch.wait_for(1);
709    }
710
711    #[test]
712    fn coalescing_never_rejects() {
713        let last = Arc::new(AtomicU32::new(0));
714        let l = Arc::clone(&last);
715        let spawner = TaskSpawner::<u32>::new(move |n| {
716            l.store(n, Ordering::Relaxed);
717        });
718        for n in 0..(TASK_QUEUE_PREALLOC as u32 * 4) {
719            spawner.spawn_coalescing(n); // never panics, never blocks
720        }
721        let target = TASK_QUEUE_PREALLOC as u32 * 4 - 1;
722        assert!(
723            wait_until(Duration::from_secs(2), || last.load(Ordering::Relaxed)
724                == target),
725            "the newest task always runs"
726        );
727    }
728
729    #[test]
730    fn serialized_runs_one_at_a_time_and_drops_nothing() {
731        // One-slot mode: the handler must never run concurrently with
732        // itself for this instance, and every FIFO task must still run.
733        const N: u32 = 64;
734        let in_flight = Arc::new(AtomicU32::new(0));
735        let peak = Arc::new(AtomicU32::new(0));
736        let latch = Arc::new(Latch::default());
737        let (inf, pk, l) = (
738            Arc::clone(&in_flight),
739            Arc::clone(&peak),
740            Arc::clone(&latch),
741        );
742        let spawner = TaskSpawner::<u32>::new_serialized(move |_| {
743            let now = inf.fetch_add(1, Ordering::AcqRel) + 1;
744            pk.fetch_max(now, Ordering::AcqRel);
745            // Widen the window so a second worker would overlap if the guard
746            // let it - a bare increment could hide a real race.
747            thread::sleep(Duration::from_millis(1));
748            inf.fetch_sub(1, Ordering::AcqRel);
749            l.bump();
750        });
751
752        // Push across a burst so re-arms land mid-drain: each one re-injects
753        // the sink and tempts an idle worker to pick it up concurrently.
754        for n in 0..N {
755            while spawner.try_spawn(n).is_err() {
756                thread::sleep(Duration::from_millis(1));
757            }
758        }
759
760        // Blocks until all N ran; a stranded task would hang here (a hung
761        // test, not a false pass).
762        latch.wait_for(N);
763        assert_eq!(
764            peak.load(Ordering::Acquire),
765            1,
766            "serialized: at most one handler in flight at a time"
767        );
768    }
769
770    #[test]
771    fn coalescing_collapses_to_the_newest() {
772        let runs = Arc::new(AtomicU32::new(0));
773        let last = Arc::new(AtomicU32::new(0));
774        let (r, l) = (Arc::clone(&runs), Arc::clone(&last));
775        let spawner = TaskSpawner::<u32>::new(move |n| {
776            r.fetch_add(1, Ordering::Relaxed);
777            l.store(n, Ordering::Relaxed);
778        });
779
780        // Fill the coalescing slot repeatedly without arming the pool, so
781        // the burst collapses in the slot rather than racing a worker.
782        // Then drain once and confirm the whole burst ran a single time,
783        // as the newest target.
784        for n in 1..=1000 {
785            let _ = spawner.sink.coalesced.force_push(n);
786        }
787        spawner.sink.drain();
788
789        assert_eq!(runs.load(Ordering::Relaxed), 1, "the burst ran once");
790        assert_eq!(last.load(Ordering::Relaxed), 1000, "and it was the newest");
791    }
792}
793
794// Model-checked proof that the schedule/drain handshake can't strand a
795// task under any thread interleaving. Run with:
796//   cargo test -p truce-core --features loom loom
797//
798// loom can't see into crossbeam's `ArrayQueue`, so this models the
799// protocol directly: a one-slot queue (`item`) plus the `scheduled` flag,
800// driven through the exact SeqCst store / fence / swap sequence that
801// `Sink::drain` and `TaskSpawner::arm` use. Weakening either side to
802// Release/AcqRel (dropping the SeqCst or the fences) makes loom find the
803// stranding interleaving; the version below passes.
804#[cfg(all(test, feature = "loom"))]
805mod loom_tests {
806    use loom::sync::Arc;
807    use loom::sync::atomic::{AtomicBool, Ordering, fence};
808    use loom::thread;
809
810    #[test]
811    fn schedule_drain_never_strands_a_task() {
812        loom::model(|| {
813            // Start with a drain in flight: the sink was scheduled
814            // (`flag == true`) and a worker is about to drain an empty
815            // queue, concurrent with a producer pushing one more task.
816            let flag = Arc::new(AtomicBool::new(true));
817            let item = Arc::new(AtomicBool::new(false));
818
819            let (f, i) = (flag.clone(), item.clone());
820            let worker = thread::spawn(move || {
821                // `Sink::drain`: clear the flag, then check the queue. The
822                // presence check must be a plain load (a `pop` reading
823                // empty) - an RMW would always read the latest value in
824                // modification order and so hide the StoreLoad staleness
825                // this test exists to catch.
826                f.store(false, Ordering::SeqCst);
827                fence(Ordering::SeqCst);
828                if i.load(Ordering::Acquire) {
829                    i.store(false, Ordering::Release); // popped it
830                }
831            });
832
833            // `try_spawn` + `arm`: push the task, then flag the sink.
834            item.store(true, Ordering::Release);
835            fence(Ordering::SeqCst);
836            let was_scheduled = flag.swap(true, Ordering::SeqCst);
837            // `was_scheduled == false` => the producer injects a fresh
838            // drain (it set `flag = true`). `true` => it relies on the
839            // in-flight drain to pick the task up.
840            let _ = was_scheduled;
841
842            worker.join().unwrap();
843
844            // Safe end states: the queue is empty (some drain popped it),
845            // or a drain is still scheduled (`flag == true`) to pick it up.
846            // A pending task with `flag == false` is the stranding bug.
847            let pending = item.load(Ordering::SeqCst);
848            let scheduled = flag.load(Ordering::SeqCst);
849            assert!(
850                !pending || scheduled,
851                "task stranded: pending with scheduled == false"
852            );
853        });
854    }
855
856    // The serialized ("one-slot") path adds a `draining` guard for mutual
857    // exclusion. The no-stranding signal stays the `scheduled` handshake: a
858    // worker that loses the guard is inert, and the winner re-checks the
859    // *flag* (not the queue) after each drain, looping until it reads
860    // `false`. This models that body for two workers plus a producer and
861    // asserts the same invariant. Re-checking `item` instead of the flag
862    // (an unsynchronized queue read) makes loom find the interleaving where
863    // the winner misses the producer's task and it strands.
864    #[test]
865    fn serialized_drain_never_strands_a_task() {
866        loom::model(|| {
867            // A sink already scheduled with one queued task; two workers pop
868            // it (the producer's re-inject can hand it to a second worker),
869            // and the producer pushes one more task concurrently.
870            let scheduled = Arc::new(AtomicBool::new(true));
871            let item = Arc::new(AtomicBool::new(true));
872            let draining = Arc::new(AtomicBool::new(false));
873
874            // One execution of the serialized `Sink::drain` path. The
875            // re-check loop is bounded to two passes - enough for the one
876            // extra task a single producer can push.
877            let worker =
878                |scheduled: Arc<AtomicBool>, item: Arc<AtomicBool>, draining: Arc<AtomicBool>| {
879                    if draining.swap(true, Ordering::Acquire) {
880                        return; // loser: inert
881                    }
882                    for _ in 0..2 {
883                        // drain_queues: clear the flag (SeqCst), fence, pop.
884                        scheduled.store(false, Ordering::SeqCst);
885                        fence(Ordering::SeqCst);
886                        let _ = item.swap(false, Ordering::AcqRel);
887                        draining.store(false, Ordering::Release);
888                        fence(Ordering::SeqCst);
889                        // Re-check the flag, not the queue.
890                        if !scheduled.load(Ordering::SeqCst) {
891                            return;
892                        }
893                        if draining.swap(true, Ordering::Acquire) {
894                            return;
895                        }
896                    }
897                };
898
899            let (s1, i1, d1) = (scheduled.clone(), item.clone(), draining.clone());
900            let w1 = thread::spawn(move || worker(s1, i1, d1));
901            let (s2, i2, d2) = (scheduled.clone(), item.clone(), draining.clone());
902            let w2 = thread::spawn(move || worker(s2, i2, d2));
903
904            // `try_spawn` + `arm`: push the task, then flag the sink.
905            item.store(true, Ordering::Release);
906            fence(Ordering::SeqCst);
907            let _ = scheduled.swap(true, Ordering::SeqCst);
908
909            w1.join().unwrap();
910            w2.join().unwrap();
911
912            // Same invariant: no task left pending unless the flag is still
913            // set for a future drain to pick it up.
914            let pending = item.load(Ordering::SeqCst);
915            let is_scheduled = scheduled.load(Ordering::SeqCst);
916            assert!(
917                !pending || is_scheduled,
918                "serialized task stranded: pending with scheduled == false"
919            );
920        });
921    }
922}