Skip to main content

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}