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