moirai_executor/schedule/runtime/scheduler/scope.rs
1//! SchedulerScope implementation.
2
3use std::{
4 marker::PhantomData,
5 mem,
6 panic::{AssertUnwindSafe, catch_unwind},
7 ptr::NonNull,
8 sync::atomic::Ordering,
9};
10
11use moirai_core::{
12 Priority,
13 error::{ExecutorError, ExecutorResult},
14};
15
16use super::super::super::{class::WorkClass, job::ScheduledJob};
17use super::super::scope_state::{SchedulerScopeState, ScopedTaskCompletion};
18use super::super::types::{SchedulerScope, ThreadScheduler, get_current_worker_id};
19use super::super::worker::{execute_job, lock_mutex, next_shared_job};
20
21impl<'scope, C, const BLOCKING_QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize>
22 SchedulerScope<'scope, C, BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>
23where
24 C: WorkClass,
25{
26 /// Spawn a job into this scope.
27 ///
28 /// The job may borrow values that outlive the scope call. Scoped jobs are
29 /// coalesced into worker-sized scheduler batches and complete before
30 /// `ThreadScheduler::scope` returns. Jobs are not guaranteed to start while
31 /// the scope body is still registering work.
32 ///
33 /// The `usize` the job receives identifies the lane running it. Worker
34 /// lanes are `0..worker_count()`; a job the admission queue turned away
35 /// runs on the calling lane, identified as `worker_count()`. It is a lane
36 /// identity, not an index into the worker set.
37 pub fn spawn<F>(&self, task: F) -> ExecutorResult<()>
38 where
39 F: FnOnce(usize) + Send + 'scope,
40 {
41 self.state().register_task();
42 let completion = ScopedTaskCompletion::new(self.state());
43 let complete = move |succeeded: bool| completion.finish(succeeded);
44
45 // SAFETY: `ThreadScheduler::scope` waits for every scheduled scoped
46 // job and drops unscheduled buffered jobs before borrowed scope data
47 // can expire.
48 let job = unsafe { ScheduledJob::new_scoped_with_completion(task, complete) };
49 self.jobs.borrow_mut().push(job);
50 Ok(())
51 }
52
53 /// Schedule all jobs currently buffered in this scope.
54 ///
55 /// `ThreadScheduler::scope` calls this before waiting, so most callers do
56 /// not need to invoke it directly. It is exposed for two-lane fork/join
57 /// shapes where one branch should enter the scheduler before the caller
58 /// executes the second branch locally. The scope still waits for every
59 /// flushed job before returning, so borrowed data cannot escape.
60 pub fn flush(&self) -> ExecutorResult<()> {
61 let jobs = mem::take(&mut *self.jobs.borrow_mut());
62 if jobs.is_empty() {
63 return Ok(());
64 }
65
66 if jobs.len() == 1 {
67 let job = jobs
68 .into_iter()
69 .next()
70 .expect("single scoped job must exist");
71 return self.schedule_single(job);
72 }
73
74 let worker_count = self.scheduler.worker_count();
75 let chunk_count = jobs.len().min(worker_count.max(1));
76 let chunk_size = jobs.len().div_ceil(chunk_count);
77 let spread_start = self
78 .locality_hint
79 .is_none()
80 .then(|| self.scheduler.select_worker::<C>(self.priority, None));
81 let mut pending_jobs = jobs.into_iter();
82
83 for chunk_index in 0..chunk_count {
84 let mut chunk = Vec::with_capacity(chunk_size);
85 for _ in 0..chunk_size {
86 if let Some(job) = pending_jobs.next() {
87 chunk.push(job);
88 }
89 }
90
91 if chunk.is_empty() {
92 break;
93 }
94
95 // Select the unhinted batch base once, then distribute physical
96 // batches across distinct workers. Re-running selection after each
97 // admission lets a fast first batch change pending/active state and
98 // route later batches back to its occupied lane, defeating the
99 // worker-sized coalescing contract and deadlocking saturated joins.
100 let locality_hint = self.locality_hint.or_else(|| {
101 spread_start.map(|start| start.wrapping_add(chunk_index) % worker_count)
102 });
103 self.schedule_chunk(chunk, locality_hint)?;
104 }
105
106 Ok(())
107 }
108
109 fn schedule_single(&self, job: ScheduledJob) -> ExecutorResult<()> {
110 self.schedule_job(job, self.locality_hint)
111 }
112
113 fn schedule_job(&self, job: ScheduledJob, locality_hint: Option<usize>) -> ExecutorResult<()> {
114 let mut job = Some(job);
115 let admitted = self
116 .scheduler
117 .admit_job::<C>(self.priority, locality_hint, &mut job);
118 self.run_if_refused(admitted, job)
119 }
120
121 fn schedule_chunk(
122 &self,
123 jobs: Vec<ScheduledJob>,
124 locality_hint: Option<usize>,
125 ) -> ExecutorResult<()> {
126 let scoped_job = move |worker_id| {
127 for job in jobs {
128 let _ = job.execute(worker_id);
129 }
130 };
131
132 // Safety: `ThreadScheduler::scope` waits for every scheduled scoped job
133 // and drops unscheduled buffered jobs before borrowed scope data can
134 // expire. A refused job runs below, inside the same scope, so it
135 // observes the same live borrows.
136 let job = unsafe { ScheduledJob::new_scoped(scoped_job) };
137 self.schedule_job(job, locality_hint)
138 }
139
140 /// Run a job the scheduler refused on the calling lane.
141 ///
142 /// A scope promises its caller that every spawned job runs before the scope
143 /// returns. Dropping a job the admission queue rejected breaks that promise
144 /// silently: the caller blocks until the scope joins and then continues as
145 /// though the work happened. Running it here keeps the promise at the cost
146 /// of the parallelism that job would have had — the same trade
147 /// `for_each_indexed` already makes with a rejected chunk, counted by the
148 /// same `admission_caller_runs` surface.
149 ///
150 /// Shutdown is not backpressure and is not absorbed: a scheduler that is
151 /// going away refuses the work, and the error reaches the caller.
152 fn run_if_refused(
153 &self,
154 admitted: ExecutorResult<()>,
155 refused: Option<ScheduledJob>,
156 ) -> ExecutorResult<()> {
157 match (admitted, refused) {
158 (Ok(()), _) => Ok(()),
159 (Err(ExecutorError::ResourceExhausted(_)), Some(job)) => {
160 self.scheduler.record_admission_caller_run();
161 // `execute` contains its own unwind boundary, so a panicking
162 // job marks its completion token failed exactly as it would on
163 // a worker instead of unwinding through the scope body.
164 let _ = job.execute(self.scheduler.caller_lane_id());
165 Ok(())
166 }
167 (Err(error), _) => Err(error),
168 }
169 }
170
171 fn state(&self) -> &SchedulerScopeState {
172 // Safety: `ThreadScheduler::scope` creates this pointer from a local
173 // state value and waits for every scheduled scoped job before returning.
174 unsafe { self.state.as_ref() }
175 }
176}
177
178/// Busy-spin iterations a worker-thread scope waiter performs after exhausting
179/// runnable work before it parks on the scope condvar. The waiter only reaches
180/// this path when its remaining scoped jobs are actively executing on other
181/// workers (nothing left to steal), so a short spin absorbs the common
182/// finish-imminently case without an OS park round-trip; the timed park below
183/// then bounds idle-CPU while `complete_task` provides the real wakeup.
184const SCOPE_HELP_SPIN_LIMIT: usize = 64;
185
186impl<const BLOCKING_QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize>
187 ThreadScheduler<BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>
188{
189 /// Run a borrowing job scope on the scheduler and wait for all spawned jobs.
190 ///
191 /// This is the scheduler-equivalent of a scoped fan-out. It avoids per-task
192 /// result storage when the caller only needs completion, while preserving the
193 /// invariant that borrowed data cannot outlive the scope.
194 ///
195 /// Jobs may borrow data that outlives this call:
196 ///
197 /// ```
198 /// use moirai_core::Priority;
199 /// use moirai_executor::{SyncTask, ThreadScheduler};
200 ///
201 /// let scheduler = ThreadScheduler::new(2, "scope-doc").unwrap();
202 /// let total = std::sync::atomic::AtomicU64::new(0);
203 /// let values = [1_u64, 2, 3];
204 /// scheduler
205 /// .scope::<SyncTask, _>(Priority::Normal, None, |scope| {
206 /// for value in &values {
207 /// scope.spawn(|_| {
208 /// total.fetch_add(*value, std::sync::atomic::Ordering::Relaxed);
209 /// })?;
210 /// }
211 /// Ok(())
212 /// })
213 /// .unwrap();
214 /// assert_eq!(total.into_inner(), 6);
215 /// ```
216 ///
217 /// A job cannot borrow a value local to the body. The body returns and drops
218 /// it before the buffered job is scheduled:
219 ///
220 /// ```compile_fail,E0597
221 /// use moirai_core::Priority;
222 /// use moirai_executor::{SyncTask, ThreadScheduler};
223 ///
224 /// let scheduler = ThreadScheduler::new(2, "scope-doc").unwrap();
225 /// scheduler
226 /// .scope::<SyncTask, _>(Priority::Normal, None, |scope| {
227 /// let local = vec![7_u8; 64];
228 /// let borrowed: &Vec<u8> = &local;
229 /// scope.spawn(move |_| assert_eq!(borrowed[0], 7))
230 /// })
231 /// .unwrap();
232 /// ```
233 pub fn scope<'scope, C, F>(
234 &'scope self,
235 priority: Priority,
236 locality_hint: Option<usize>,
237 body: F,
238 ) -> ExecutorResult<()>
239 where
240 C: WorkClass,
241 F: FnOnce(
242 &SchedulerScope<'scope, C, BLOCKING_QUEUE_CAPACITY, SPIN_LIMIT>,
243 ) -> ExecutorResult<()>,
244 {
245 if self.inner.shutdown.load(Ordering::Acquire) {
246 return Err(ExecutorError::ShuttingDown);
247 }
248
249 let state = SchedulerScopeState::new();
250 let scope = SchedulerScope {
251 scheduler: self,
252 state: NonNull::from(&state),
253 priority,
254 locality_hint,
255 jobs: std::cell::RefCell::new(Vec::new()),
256 _scope: PhantomData,
257 _class: PhantomData,
258 };
259
260 let body_result = catch_unwind(AssertUnwindSafe(|| body(&scope)));
261 // `flush` may enqueue lifetime-erased borrowing jobs before an internal
262 // unwind. Catch it so the scope state remains live through the drain.
263 let flush_result = catch_unwind(AssertUnwindSafe(|| scope.flush()));
264 // A panic here would unwind `scope` while buffered or running jobs still
265 // hold `state` and the caller's borrows, so it cannot be allowed to
266 // escape: the process aborts, as `std::thread::scope` does.
267 if catch_unwind(AssertUnwindSafe(|| self.drain_scope(&state))).is_err() {
268 std::process::abort();
269 }
270
271 match body_result {
272 Err(payload) => std::panic::resume_unwind(payload),
273 Ok(body_result) => match flush_result {
274 Err(payload) => std::panic::resume_unwind(payload),
275 Ok(flush_result) => match body_result {
276 Ok(()) if state.has_panicked() => Err(ExecutorError::SpawnFailed(
277 moirai_core::error::TaskError::Panicked,
278 )),
279 Ok(()) => flush_result.and_then(|()| state.unrun_job_result()),
280 Err(error) => Err(error),
281 },
282 },
283 }
284 }
285
286 /// Wait for every job registered on `state` to complete.
287 ///
288 /// If the caller is itself a scheduler worker, it participates in work
289 /// stealing instead of parking: a worker that blocks inside `scope` while
290 /// its nested scoped jobs sit unrun would otherwise remove itself from the
291 /// pool and deadlock the fork-join (provably so on a single-worker pool, and
292 /// a source of use-after-free on the scope's stack-owned state under
293 /// concurrent nesting). Running its own queue via `next_job` keeps the pool
294 /// making progress, so nesting is deadlock-free and the scope state stays
295 /// live until every borrowing job has completed. `next_job(worker_id)` only
296 /// pops this worker's own deque and steals into it, so the aliasing rules of
297 /// the single-owner Chase–Lev deques are preserved.
298 ///
299 /// A non-worker caller parks (`SchedulerScopeState::wait`): the worker pool
300 /// drains its scoped jobs, so it never starves anything by blocking. A
301 /// caller that helped from its own lane crashed consumers (moirai-iter's
302 /// nested iteration on CI, kwavers' 3-D FFT at 32³ and above); until the
303 /// race is found the join waits as before.
304 pub(super) fn drain_scope(&self, state: &SchedulerScopeState) {
305 let Some(worker_id) = self.own_worker_id() else {
306 state.wait();
307 return;
308 };
309
310 let inner = &self.inner;
311 let mut idle_spins = 0usize;
312 loop {
313 if state.pending_tasks.load(Ordering::Acquire) == 0 {
314 state.wait();
315 return;
316 }
317
318 if let Some(job) = next_shared_job(inner, worker_id) {
319 execute_job(inner, worker_id, job);
320 idle_spins = 0;
321 continue;
322 }
323
324 // Scope still pending but nothing runnable: the remaining scoped jobs
325 // are executing on other workers. Spin briefly, then park on the
326 // scope condvar with a timeout so `complete_task` wakes us while we
327 // still periodically re-probe for freshly stealable work.
328 if idle_spins < SCOPE_HELP_SPIN_LIMIT {
329 idle_spins += 1;
330 core::hint::spin_loop();
331 continue;
332 }
333 idle_spins = 0;
334
335 let guard = lock_mutex(&state.wait_lock);
336 if state.pending_tasks.load(Ordering::Acquire) != 0 {
337 let _ = state
338 .wait_signal
339 .wait_timeout(guard, std::time::Duration::from_micros(50))
340 .unwrap_or_else(|poisoned| poisoned.into_inner());
341 }
342 }
343 }
344
345 /// Index of the calling thread among this scheduler's own workers.
346 ///
347 /// The worker-id thread cache is process-wide, so it names whichever
348 /// scheduler owns the thread. Indexing this scheduler's worker table with
349 /// a foreign id panics or drains the wrong deques, so membership is
350 /// confirmed against the thread each worker registers when it starts.
351 fn own_worker_id(&self) -> Option<usize> {
352 let id = get_current_worker_id()?;
353 let registered = self.inner.workers.get(id)?.thread.get()?;
354 (registered.id() == std::thread::current().id()).then_some(id)
355 }
356}