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