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