Skip to main content

scheduler/
scheduler.rs

1mod clock;
2mod executor;
3mod test_scheduler;
4#[cfg(test)]
5mod tests;
6
7pub use clock::*;
8pub use executor::*;
9pub use test_scheduler::*;
10
11static TEST_SCHEDULER_CREATED: std::sync::atomic::AtomicBool =
12    std::sync::atomic::AtomicBool::new(false);
13
14/// Whether this process has created a [`TestScheduler`].
15///
16/// Test harnesses create one before the test body runs, and applications never
17/// do. Code with no executor at hand, such as expensive data-structure invariant
18/// checks, can use this to run only in test processes. Code that has an
19/// executor should prefer [`BackgroundExecutor::is_test`], which is correct per
20/// executor.
21pub fn test_scheduler_created() -> bool {
22    TEST_SCHEDULER_CREATED.load(std::sync::atomic::Ordering::Relaxed)
23}
24
25use async_task::Runnable;
26use futures::channel::oneshot;
27use std::{
28    any::Any,
29    future::Future,
30    panic::Location,
31    pin::Pin,
32    sync::Arc,
33    task::{Context, Poll},
34    time::Duration,
35};
36
37/// Task priority for background tasks.
38///
39/// Higher priority tasks are more likely to be scheduled before lower priority tasks,
40/// but this is not a strict guarantee - the scheduler may interleave tasks of different
41/// priorities to prevent starvation.
42#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Hash)]
43#[repr(u8)]
44pub enum Priority {
45    /// Realtime priority
46    ///
47    /// Spawning a task with this priority will spin it off on a separate thread dedicated just to that task. Only use for audio.
48    RealtimeAudio,
49    /// High priority - use for tasks critical to user experience/responsiveness.
50    High,
51    /// Medium priority - suitable for most use cases.
52    #[default]
53    Medium,
54    /// Low priority - use for background work that can be deprioritized.
55    Low,
56}
57
58impl Priority {
59    /// Returns the relative probability weight for this priority level.
60    /// Used by schedulers to determine task selection probability.
61    pub const fn weight(self) -> u32 {
62        match self {
63            Priority::High => 60,
64            Priority::Medium => 30,
65            Priority::Low => 10,
66            // realtime priorities are not considered for probability scheduling
67            Priority::RealtimeAudio => 0,
68        }
69    }
70}
71
72#[derive(Clone, Copy, Debug)]
73pub struct SpawnTime(pub Instant);
74
75/// Metadata attached to runnables for debugging and profiling.
76#[derive(Clone, Debug)]
77pub struct RunnableMeta {
78    /// The source location where the task was spawned.
79    pub location: &'static Location<'static>,
80    /// The moment the task was spawned.
81    pub spawned: SpawnTime,
82}
83
84impl RunnableMeta {
85    #[track_caller]
86    pub fn new_with_callers_location() -> Self {
87        Self {
88            location: core::panic::Location::caller(),
89            spawned: SpawnTime(Instant::now()),
90        }
91    }
92}
93
94pub trait Scheduler: Send + Sync {
95    /// Block until the given future completes or timeout occurs.
96    ///
97    /// Returns `true` if the future completed, `false` if it timed out.
98    /// The future is passed as a pinned mutable reference so the caller
99    /// retains ownership and can continue polling or return it on timeout.
100    #[cfg(not(target_family = "wasm"))]
101    fn block(
102        &self,
103        session_id: Option<SessionId>,
104        future: Pin<&mut dyn Future<Output = ()>>,
105        timeout: Option<Duration>,
106    ) -> bool;
107
108    /// Schedule a runnable on the local (session-pinned) queue for `session_id`.
109    /// Runnables scheduled here run in order on whichever thread drains the
110    /// session — the main thread for ordinary sessions, or a dedicated OS
111    /// thread for sessions created via `spawn_dedicated_thread`.
112    fn schedule_local(&self, session_id: SessionId, runnable: Runnable<RunnableMeta>);
113
114    /// Schedule a background task with the given priority.
115    fn schedule_background_with_priority(
116        &self,
117        runnable: Runnable<RunnableMeta>,
118        priority: Priority,
119    );
120
121    /// Spawn a closure on a dedicated realtime thread for audio processing.
122    fn spawn_realtime(&self, f: Box<dyn FnOnce() + Send>);
123
124    /// Schedule a background task with default (medium) priority.
125    fn schedule_background(&self, runnable: Runnable<RunnableMeta>) {
126        self.schedule_background_with_priority(runnable, Priority::default());
127    }
128
129    #[track_caller]
130    fn timer(&self, timeout: Duration) -> Timer;
131    fn clock(&self) -> Arc<dyn Clock>;
132
133    /// Spawn a closure on a fresh session pinned to its own [`LocalExecutor`].
134    ///
135    /// `PlatformScheduler` runs the closure on a new OS thread (see
136    /// [`spawn_dedicated_thread`]). `TestScheduler` runs it on the test
137    /// scheduler's loop alongside everything else so determinism under
138    /// `TestScheduler::many` is preserved.
139    ///
140    /// This is the dyn-safe entry point: the closure's output is type-erased
141    /// as `Box<dyn Any + Send + Sync>` so the trait stays object-safe.
142    /// Callers typically reach for the type-safe wrappers on
143    /// [`LocalExecutor::spawn_dedicated`] and
144    /// [`BackgroundExecutor::spawn_dedicated`], which compose this method
145    /// with [`Task::downcast`] to recover the closure's concrete return type.
146    fn spawn_dedicated(
147        self: Arc<Self>,
148        f: Box<
149            dyn FnOnce(
150                    LocalExecutor,
151                )
152                    -> Pin<Box<dyn Future<Output = Box<dyn Any + Send + Sync>> + 'static>>
153                + Send
154                + 'static,
155        >,
156    ) -> Task<Box<dyn Any + Send + Sync>>;
157
158    fn as_test(&self) -> Option<&TestScheduler> {
159        None
160    }
161}
162
163/// Spawn work on a fresh OS thread that's exclusive to the returned task and
164/// anything spawned on the executor it provides. Blocking syscalls inside that
165/// work don't disturb any other executor in the process.
166///
167/// `f` is called on the dedicated thread with a [`LocalExecutor`] pinned
168/// to it. The future `f` returns may freely be `!Send`. The returned `Task`
169/// resolves to that future's output: dropping it cancels the root, but
170/// detached children keep running until they finish. The thread shuts down
171/// once the executor and every task on it are gone.
172///
173/// This function never blocks: the returned task starts as an asynchronous
174/// rendezvous that resolves to the root task's handle once the dedicated
175/// thread has spawned it, and behaves like that handle from then on. That
176/// makes it safe to call from threads that must not block, such as the web
177/// main thread. On wasm targets the thread is a web worker and requires the
178/// `wasm-threads` cargo feature (and a shared-memory build); without it this
179/// panics. Spawning a worker is far more expensive than an OS thread —
180/// module instantiation plus a fresh function table — so on the web treat a
181/// dedicated session as a long-lived place to send work, not a per-job
182/// convenience.
183///
184/// The caller is responsible for supplying a `session_id` that's distinct from
185/// every other live session on `scheduler`. Concrete schedulers typically wrap
186/// this in an inherent method that allocates the id from their own counter.
187pub fn spawn_dedicated_thread<F, Fut>(
188    session_id: SessionId,
189    scheduler: Arc<dyn Scheduler>,
190    f: F,
191) -> Task<Fut::Output>
192where
193    F: FnOnce(LocalExecutor) -> Fut + Send + 'static,
194    Fut: Future + 'static,
195    Fut::Output: Send + 'static,
196{
197    let (task, delivery) = Task::rendezvous();
198    let thread_body = move || {
199        let (runnable_sender, runnable_receiver) = flume::unbounded::<Runnable<RunnableMeta>>();
200        let dispatch = move |runnable: Runnable<RunnableMeta>| {
201            let _ = runnable_sender.send(runnable);
202        };
203        let executor = LocalExecutor::new(session_id, scheduler, dispatch);
204        let root_task = executor.spawn(f(executor.clone()));
205        // If the caller already dropped or detached the rendezvous task,
206        // delivery applies that disposition to the root task here.
207        delivery.deliver(root_task);
208        // After this drop, every strong reference to the runnable sender
209        // lives inside a spawned task or a user-held executor clone. The
210        // recv loop exits once all of those are gone.
211        drop(executor);
212
213        while let Ok(runnable) = runnable_receiver.recv() {
214            runnable.run();
215        }
216    };
217    spawn_dedicated_os_thread(session_id, thread_body);
218    task
219}
220
221fn spawn_dedicated_os_thread(session_id: SessionId, thread_body: impl FnOnce() + Send + 'static) {
222    let thread_name = format!("spawn_dedicated session {:?}", session_id);
223    #[cfg(not(target_family = "wasm"))]
224    std::thread::Builder::new()
225        .name(thread_name)
226        .spawn(thread_body)
227        .expect("failed to spawn dedicated thread");
228    #[cfg(all(target_family = "wasm", feature = "wasm-threads"))]
229    wasm_thread::Builder::new()
230        .name(thread_name)
231        .spawn(thread_body)
232        .expect("failed to spawn dedicated thread");
233    #[cfg(all(target_family = "wasm", not(feature = "wasm-threads")))]
234    {
235        let _ = (thread_name, thread_body);
236        panic!("spawn_dedicated on wasm requires the scheduler crate's `wasm-threads` feature");
237    }
238}
239
240#[derive(Copy, Clone, Debug, Eq, PartialEq, Ord, PartialOrd, Hash)]
241pub struct SessionId(u16);
242
243impl SessionId {
244    pub fn new(id: u16) -> Self {
245        SessionId(id)
246    }
247}
248
249pub struct Timer(oneshot::Receiver<()>);
250
251impl Timer {
252    pub fn new(rx: oneshot::Receiver<()>) -> Self {
253        Timer(rx)
254    }
255}
256
257impl Future for Timer {
258    type Output = ();
259
260    fn poll(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll<()> {
261        match Pin::new(&mut self.0).poll(cx) {
262            Poll::Ready(_) => Poll::Ready(()),
263            Poll::Pending => Poll::Pending,
264        }
265    }
266}