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