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