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