Skip to main content

moirai_executor/schedule/
seam.rs

1//! The scheduler seam — the dependency-inversion boundary between the executor
2//! façade ([`HybridExecutor`](crate::hybrid::HybridExecutor)) and a concrete
3//! work-stealing runtime.
4//!
5//! `HybridExecutor` depends on these role traits, not on the concrete
6//! [`ThreadScheduler`], so an alternative runtime (for example a single-threaded
7//! scheduler for `wasm32` targets) can be substituted as `HybridExecutor<S>`
8//! without touching the façade. The contract is split per the
9//! interface-segregation principle, so a substitute implements only the roles it
10//! supports. Every method is statically dispatched (generic over the work class
11//! and closure), so the seam compiles to the same code as a direct concrete
12//! call — it is zero-cost.
13//!
14//! The richer borrowing `scope` API is intentionally *not* part of the seam: its
15//! signature exposes a concrete [`SchedulerScope`](crate::schedule::SchedulerScope)
16//! borrow handle, so it stays an inherent capability of `ThreadScheduler` rather
17//! than forcing every substitute to reproduce that machinery.
18
19use moirai_core::{Priority, error::ExecutorResult};
20
21use crate::schedule::class::WorkClass;
22use crate::schedule::runtime::{ScheduleMetrics, ThreadScheduler};
23
24/// Submit type-erased work to a scheduler.
25pub trait WorkSubmit: Send + Sync + 'static {
26    /// Schedule `task` as a job of work class `C` at `priority`, optionally
27    /// biased toward the worker named by `locality_hint`.
28    ///
29    /// # Errors
30    /// Returns [`ExecutorError::ShuttingDown`](moirai_core::error::ExecutorError::ShuttingDown)
31    /// if the scheduler is draining.
32    fn schedule<C, F>(
33        &self,
34        priority: Priority,
35        locality_hint: Option<usize>,
36        task: F,
37    ) -> ExecutorResult<()>
38    where
39        C: WorkClass,
40        F: FnOnce(usize) + Send + 'static;
41}
42
43/// Introspect and control a scheduler's lifecycle.
44pub trait SchedulerControl {
45    /// Number of queued-but-not-yet-running jobs.
46    fn pending_tasks(&self) -> usize;
47    /// Number of jobs currently executing.
48    fn active_workers(&self) -> usize;
49    /// Number of worker threads.
50    fn worker_count(&self) -> usize;
51    /// Whether any work is queued or running.
52    fn has_work(&self) -> bool;
53    /// Block until the scheduler is quiescent (no queued or running work).
54    ///
55    /// # Errors
56    /// Returns an [`ExecutorError`](moirai_core::error::ExecutorError) if the
57    /// scheduler cannot reach quiescence.
58    fn join(&self) -> ExecutorResult<()>;
59    /// Drain queued work and stop the worker sets.
60    ///
61    /// Exactly one non-worker caller joins the worker sets. Other external
62    /// callers wait for that join. Scheduler workers close the blocking lane
63    /// and return before the election; when no external caller remains, worker
64    /// ownership releases the scheduler after every accepted task drains.
65    fn shutdown(&self);
66    /// Snapshot the scheduler metrics.
67    fn metrics(&self) -> ScheduleMetrics;
68}
69
70/// Indexed data-parallel fan-out without per-item result storage.
71///
72/// Calls partition non-empty domains across the available worker-plus-caller
73/// lanes. Operation-level execution policies own profitability thresholds
74/// before invoking this scheduler seam.
75pub trait DataParallel {
76    /// Apply `task` to every index in `0..count`, completing before return.
77    ///
78    /// # Errors
79    /// Returns an [`ExecutorError`](moirai_core::error::ExecutorError) if the
80    /// scheduler is draining or a chunk panics.
81    fn for_each_indexed<C, F>(
82        &self,
83        priority: Priority,
84        locality_hint: Option<usize>,
85        count: usize,
86        task: F,
87    ) -> ExecutorResult<()>
88    where
89        C: WorkClass,
90        F: Fn(usize) + Send + Sync;
91
92    /// Map every index in `0..count` and reduce the results with `identity` as
93    /// the neutral element of `reduce`.
94    ///
95    /// # Errors
96    /// Returns an [`ExecutorError`](moirai_core::error::ExecutorError) if the
97    /// scheduler is draining or a chunk panics.
98    fn map_reduce_indexed<C, T, Map, Reduce>(
99        &self,
100        priority: Priority,
101        locality_hint: Option<usize>,
102        count: usize,
103        identity: T,
104        map: Map,
105        reduce: Reduce,
106    ) -> ExecutorResult<T>
107    where
108        C: WorkClass,
109        T: Send + Clone,
110        Map: Fn(usize) -> T + Send + Sync,
111        Reduce: Fn(T, T) -> T + Send + Sync;
112}
113
114/// The full runtime-scheduler contract that `HybridExecutor` depends on.
115///
116/// `Clone` is required because the async lane clones the scheduler handle into
117/// each future's waker. This is a marker over the role traits, implemented for
118/// every type that satisfies all of them.
119pub trait WorkScheduler: WorkSubmit + SchedulerControl + DataParallel + Clone {}
120
121impl<S> WorkScheduler for S where S: WorkSubmit + SchedulerControl + DataParallel + Clone {}
122
123// ── ThreadScheduler: the canonical implementation ──────────────────────────
124//
125// Each role method forwards to the same-named inherent method via the
126// `Type::method` form, which resolves to the inherent method (preferred over a
127// trait method of the same name), so there is no recursion. The seam therefore
128// adds no behavior — it only re-exposes the existing surface as a substitutable
129// contract.
130
131impl<const BLOCKING_QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> WorkSubmit
132    for ThreadScheduler<BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>
133{
134    fn schedule<C, F>(
135        &self,
136        priority: Priority,
137        locality_hint: Option<usize>,
138        task: F,
139    ) -> ExecutorResult<()>
140    where
141        C: WorkClass,
142        F: FnOnce(usize) + Send + 'static,
143    {
144        ThreadScheduler::schedule::<C, F>(self, priority, locality_hint, task)
145    }
146}
147
148impl<const BLOCKING_QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> SchedulerControl
149    for ThreadScheduler<BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>
150{
151    fn pending_tasks(&self) -> usize {
152        ThreadScheduler::pending_tasks(self)
153    }
154    fn active_workers(&self) -> usize {
155        ThreadScheduler::active_workers(self)
156    }
157    fn worker_count(&self) -> usize {
158        ThreadScheduler::worker_count(self)
159    }
160    fn has_work(&self) -> bool {
161        ThreadScheduler::has_work(self)
162    }
163    fn join(&self) -> ExecutorResult<()> {
164        ThreadScheduler::join(self)
165    }
166    fn shutdown(&self) {
167        ThreadScheduler::shutdown(self)
168    }
169    fn metrics(&self) -> ScheduleMetrics {
170        ThreadScheduler::metrics(self)
171    }
172}
173
174impl<const BLOCKING_QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> DataParallel
175    for ThreadScheduler<BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>
176{
177    fn for_each_indexed<C, F>(
178        &self,
179        priority: Priority,
180        locality_hint: Option<usize>,
181        count: usize,
182        task: F,
183    ) -> ExecutorResult<()>
184    where
185        C: WorkClass,
186        F: Fn(usize) + Send + Sync,
187    {
188        ThreadScheduler::for_each_indexed::<C, F>(self, priority, locality_hint, count, task)
189    }
190
191    fn map_reduce_indexed<C, T, Map, Reduce>(
192        &self,
193        priority: Priority,
194        locality_hint: Option<usize>,
195        count: usize,
196        identity: T,
197        map: Map,
198        reduce: Reduce,
199    ) -> ExecutorResult<T>
200    where
201        C: WorkClass,
202        T: Send + Clone,
203        Map: Fn(usize) -> T + Send + Sync,
204        Reduce: Fn(T, T) -> T + Send + Sync,
205    {
206        ThreadScheduler::map_reduce_indexed::<C, T, Map, Reduce>(
207            self,
208            priority,
209            locality_hint,
210            count,
211            identity,
212            map,
213            reduce,
214        )
215    }
216}