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::{error::ExecutorResult, Priority};
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 /// Signal drain and stop the workers.
60 fn shutdown(&self);
61 /// Snapshot the scheduler metrics.
62 fn metrics(&self) -> ScheduleMetrics;
63}
64
65/// Indexed data-parallel fan-out without per-item result storage.
66///
67/// Calls partition non-empty domains across the available worker-plus-caller
68/// lanes. Operation-level execution policies own profitability thresholds
69/// before invoking this scheduler seam.
70pub trait DataParallel {
71 /// Apply `task` to every index in `0..count`, completing before return.
72 ///
73 /// # Errors
74 /// Returns an [`ExecutorError`](moirai_core::error::ExecutorError) if the
75 /// scheduler is draining or a chunk panics.
76 fn for_each_indexed<C, F>(
77 &self,
78 priority: Priority,
79 locality_hint: Option<usize>,
80 count: usize,
81 task: F,
82 ) -> ExecutorResult<()>
83 where
84 C: WorkClass,
85 F: Fn(usize) + Send + Sync;
86
87 /// Map every index in `0..count` and reduce the results with `identity` as
88 /// the neutral element of `reduce`.
89 ///
90 /// # Errors
91 /// Returns an [`ExecutorError`](moirai_core::error::ExecutorError) if the
92 /// scheduler is draining or a chunk panics.
93 fn map_reduce_indexed<C, T, Map, Reduce>(
94 &self,
95 priority: Priority,
96 locality_hint: Option<usize>,
97 count: usize,
98 identity: T,
99 map: Map,
100 reduce: Reduce,
101 ) -> ExecutorResult<T>
102 where
103 C: WorkClass,
104 T: Send + Clone,
105 Map: Fn(usize) -> T + Send + Sync,
106 Reduce: Fn(T, T) -> T + Send + Sync;
107}
108
109/// The full runtime-scheduler contract that `HybridExecutor` depends on.
110///
111/// `Clone` is required because the async lane clones the scheduler handle into
112/// each future's waker. This is a marker over the role traits, implemented for
113/// every type that satisfies all of them.
114pub trait WorkScheduler: WorkSubmit + SchedulerControl + DataParallel + Clone {}
115
116impl<S> WorkScheduler for S where S: WorkSubmit + SchedulerControl + DataParallel + Clone {}
117
118// ── ThreadScheduler: the canonical implementation ──────────────────────────
119//
120// Each role method forwards to the same-named inherent method via the
121// `Type::method` form, which resolves to the inherent method (preferred over a
122// trait method of the same name), so there is no recursion. The seam therefore
123// adds no behavior — it only re-exposes the existing surface as a substitutable
124// contract.
125
126impl<const QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> WorkSubmit
127 for ThreadScheduler<QUEUE_CAPACITY, SPIN_LIMIT>
128{
129 fn schedule<C, F>(
130 &self,
131 priority: Priority,
132 locality_hint: Option<usize>,
133 task: F,
134 ) -> ExecutorResult<()>
135 where
136 C: WorkClass,
137 F: FnOnce(usize) + Send + 'static,
138 {
139 ThreadScheduler::schedule::<C, F>(self, priority, locality_hint, task)
140 }
141}
142
143impl<const QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> SchedulerControl
144 for ThreadScheduler<QUEUE_CAPACITY, SPIN_LIMIT>
145{
146 fn pending_tasks(&self) -> usize {
147 ThreadScheduler::pending_tasks(self)
148 }
149 fn active_workers(&self) -> usize {
150 ThreadScheduler::active_workers(self)
151 }
152 fn worker_count(&self) -> usize {
153 ThreadScheduler::worker_count(self)
154 }
155 fn has_work(&self) -> bool {
156 ThreadScheduler::has_work(self)
157 }
158 fn join(&self) -> ExecutorResult<()> {
159 ThreadScheduler::join(self)
160 }
161 fn shutdown(&self) {
162 ThreadScheduler::shutdown(self)
163 }
164 fn metrics(&self) -> ScheduleMetrics {
165 ThreadScheduler::metrics(self)
166 }
167}
168
169impl<const QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize> DataParallel
170 for ThreadScheduler<QUEUE_CAPACITY, SPIN_LIMIT>
171{
172 fn for_each_indexed<C, F>(
173 &self,
174 priority: Priority,
175 locality_hint: Option<usize>,
176 count: usize,
177 task: F,
178 ) -> ExecutorResult<()>
179 where
180 C: WorkClass,
181 F: Fn(usize) + Send + Sync,
182 {
183 ThreadScheduler::for_each_indexed::<C, F>(self, priority, locality_hint, count, task)
184 }
185
186 fn map_reduce_indexed<C, T, Map, Reduce>(
187 &self,
188 priority: Priority,
189 locality_hint: Option<usize>,
190 count: usize,
191 identity: T,
192 map: Map,
193 reduce: Reduce,
194 ) -> ExecutorResult<T>
195 where
196 C: WorkClass,
197 T: Send + Clone,
198 Map: Fn(usize) -> T + Send + Sync,
199 Reduce: Fn(T, T) -> T + Send + Sync,
200 {
201 ThreadScheduler::map_reduce_indexed::<C, T, Map, Reduce>(
202 self,
203 priority,
204 locality_hint,
205 count,
206 identity,
207 map,
208 reduce,
209 )
210 }
211}