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::{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}