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}