Skip to main content

moirai_executor/schedule/runtime/scheduler/
scope.rs

1//! SchedulerScope implementation.
2
3use std::{
4    marker::PhantomData,
5    mem,
6    panic::{catch_unwind, AssertUnwindSafe},
7};
8
9use moirai_core::error::{ExecutorError, ExecutorResult};
10
11use super::super::super::{class::WorkClass, job::ScheduledJob};
12use super::super::types::{SchedulerScope, SchedulerScopeState, ScopedTaskCompletion};
13
14impl<'scope, C, const QUEUE_CAPACITY: usize, const SPIN_LIMIT: usize>
15    SchedulerScope<'scope, C, QUEUE_CAPACITY, SPIN_LIMIT>
16where
17    C: WorkClass,
18{
19    /// Spawn a job into this scope.
20    ///
21    /// The job may borrow values that outlive the scope call. Scoped jobs are
22    /// coalesced into worker-sized scheduler batches and complete before
23    /// `ThreadScheduler::scope` returns. Jobs are not guaranteed to start while
24    /// the scope body is still registering work.
25    ///
26    /// The `usize` the job receives identifies the lane running it. Worker
27    /// lanes are `0..worker_count()`; a job the admission queue turned away
28    /// runs on the calling lane, identified as `worker_count()`. It is a lane
29    /// identity, not an index into the worker set.
30    pub fn spawn<F>(&self, task: F) -> ExecutorResult<()>
31    where
32        F: FnOnce(usize) + Send + 'scope,
33    {
34        self.state().register_task();
35        let completion = ScopedTaskCompletion {
36            state: self.state,
37            _state: PhantomData,
38        };
39        let scoped_task = move |worker_id| {
40            let _completion = completion;
41            let result = catch_unwind(AssertUnwindSafe(|| task(worker_id)));
42            if result.is_err() {
43                _completion.mark_failed();
44            }
45        };
46
47        // Safety: `ThreadScheduler::scope` waits for every scheduled scoped
48        // job and drops unscheduled buffered jobs before borrowed scope data
49        // can expire.
50        let job = unsafe { ScheduledJob::new_scoped(scoped_task) };
51        self.jobs.borrow_mut().push(job);
52        Ok(())
53    }
54
55    /// Schedule all jobs currently buffered in this scope.
56    ///
57    /// `ThreadScheduler::scope` calls this before waiting, so most callers do
58    /// not need to invoke it directly. It is exposed for two-lane fork/join
59    /// shapes where one branch should enter the scheduler before the caller
60    /// executes the second branch locally. The scope still waits for every
61    /// flushed job before returning, so borrowed data cannot escape.
62    pub fn flush(&self) -> ExecutorResult<()> {
63        let jobs = mem::take(&mut *self.jobs.borrow_mut());
64        if jobs.is_empty() {
65            return Ok(());
66        }
67
68        if jobs.len() == 1 {
69            let job = jobs
70                .into_iter()
71                .next()
72                .expect("single scoped job must exist");
73            return self.schedule_single(job);
74        }
75
76        let worker_count = self.scheduler.worker_count();
77        let chunk_count = jobs.len().min(worker_count.max(1));
78        let chunk_size = jobs.len().div_ceil(chunk_count);
79        let mut pending_jobs = jobs.into_iter();
80
81        for _ in 0..chunk_count {
82            let mut chunk = Vec::with_capacity(chunk_size);
83            for _ in 0..chunk_size {
84                if let Some(job) = pending_jobs.next() {
85                    chunk.push(job);
86                }
87            }
88
89            if chunk.is_empty() {
90                break;
91            }
92
93            self.schedule_chunk(chunk)?;
94        }
95
96        Ok(())
97    }
98
99    fn schedule_single(&self, job: ScheduledJob) -> ExecutorResult<()> {
100        let mut job = Some(job);
101        let admitted = self
102            .scheduler
103            .admit_job::<C>(self.priority, self.locality_hint, &mut job);
104        self.run_if_refused(admitted, job)
105    }
106
107    fn schedule_chunk(&self, jobs: Vec<ScheduledJob>) -> ExecutorResult<()> {
108        let scoped_job = move |worker_id| {
109            for job in jobs {
110                let _ = job.execute(worker_id);
111            }
112        };
113
114        // Safety: `ThreadScheduler::scope` waits for every scheduled scoped job
115        // and drops unscheduled buffered jobs before borrowed scope data can
116        // expire. A refused job runs below, inside the same scope, so it
117        // observes the same live borrows.
118        let job = unsafe { ScheduledJob::new_scoped(scoped_job) };
119        self.schedule_single(job)
120    }
121
122    /// Run a job the scheduler refused on the calling lane.
123    ///
124    /// A scope promises its caller that every spawned job runs before the scope
125    /// returns. Dropping a job the admission queue rejected breaks that promise
126    /// silently: the caller blocks until the scope joins and then continues as
127    /// though the work happened. Running it here keeps the promise at the cost
128    /// of the parallelism that job would have had — the same trade
129    /// `for_each_indexed` already makes with a rejected chunk, counted by the
130    /// same `admission_caller_runs` surface.
131    ///
132    /// Shutdown is not backpressure and is not absorbed: a scheduler that is
133    /// going away refuses the work, and the error reaches the caller.
134    fn run_if_refused(
135        &self,
136        admitted: ExecutorResult<()>,
137        refused: Option<ScheduledJob>,
138    ) -> ExecutorResult<()> {
139        match (admitted, refused) {
140            (Ok(()), _) => Ok(()),
141            (Err(ExecutorError::ResourceExhausted(_)), Some(job)) => {
142                self.scheduler.record_admission_caller_run();
143                // `execute` contains its own unwind boundary, so a panicking
144                // job marks its completion token failed exactly as it would on
145                // a worker instead of unwinding through the scope body.
146                let _ = job.execute(self.scheduler.caller_lane_id());
147                Ok(())
148            }
149            (Err(error), _) => Err(error),
150        }
151    }
152
153    fn state(&self) -> &SchedulerScopeState {
154        // Safety: `ThreadScheduler::scope` creates this pointer from a local
155        // state value and waits for every scheduled scoped job before returning.
156        unsafe { self.state.as_ref() }
157    }
158}