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}