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
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
//! The parallel-executor seam: how Melinoe hands a partition's tasks to an
//! external scheduler without depending on it.
//!
//! Melinoe is a foundation crate with no dependencies and a `#![no_std]`
//! posture, so it cannot name a scheduler. It instead defines the *contract* a
//! scheduler must satisfy ([`ParallelExecutor`]) and a process-global slot
//! holding one implementation ([`register_parallel_executor`]). A scheduler
//! implements the trait and registers itself; Melinoe's partition drivers then
//! run their shards on it.
//!
//! # Why a trait and not a function pointer
//!
//! The seam used to be a bare
//! `unsafe fn(usize, unsafe fn(usize, *mut ()), *mut ())` stored in an
//! `AtomicPtr<()>`. That shape worked, but it leaked cost into every
//! implementor:
//!
//! * The contract ("invoke every index exactly once, and return only when the
//! last invocation has completed") could not be stated anywhere a
//! implementor would read it — it lived in a doc comment on an unrelated
//! constructor.
//! * `Send` was inexpressible. An implementor that moves its context pointer
//! onto worker threads had to launder it through `usize` to satisfy the
//! compiler, which is exactly the kind of cast that defeats the check.
//! * Melinoe's own `MaybeUninit` out-buffer trusts the contract absolutely. A
//! caller could install any function of the right *arity*, and the machine
//! that contains the fallout ([`super::driver_core`]'s drop guard) exists
//! because that trust cannot be verified.
//!
//! Expressing the seam as a trait does not make the contract verifiable — that
//! is not possible across a crate boundary — but it moves it to one place an
//! implementor cannot miss. The entry point is an associated function rather
//! than a receiver method: the global slot stores no scheduler value, so an
//! implementation cannot accidentally observe fabricated or uninitialized
//! receiver storage.
//!
//! # The zero-alloc global slot
//!
//! A global slot cannot hold a generic type, and `dyn Trait` needs `alloc`,
//! which is an optional feature here. The slot therefore holds a *monomorphized
//! shim*: [`register`] generates one trampoline per implementing type with the
//! same ABI the old design used, and stores its address. The ABI is an
//! implementation detail of this module now, not the public interface.
//!
//! # Registration is process-global and order-sensitive
//!
//! There is one slot for the whole process, and it is read afresh on every
//! partition call. Two consequences follow.
//!
//! **Register before partitioning.** The scoped-thread fallback costs roughly
//! 30 µs per shard — it spawns `min(parts, len) − 1` OS threads per call and
//! joins them — against a pool dispatch that spawns nothing. The gap is large
//! enough that a fine-grained workload can be an order of magnitude slower on
//! the fallback. Because registration is lazily performed by the scheduling
//! layer (a scheduler typically registers from its own first-access
//! initializer), a program whose first partition call happens *before* it
//! touches that scheduler will silently take the fallback path for that call.
//! Registering, or touching the scheduler, at startup avoids this.
//!
//! **A later registration does not retroactively change a call in flight.**
//! Each call loads the slot once; calls already running keep the driver they
//! started with.
use ;
/// A scheduler that can run a partition's independent tasks.
///
/// This is the contract a scheduler honours to drive Melinoe's partitioning.
/// Implement its associated entry point, then pass the type to
/// [`register_parallel_executor`]. The global slot stores only a
/// monomorphized function pointer; no scheduler value or receiver is created.
///
/// # Contract
///
/// An implementation of [`run_indexed`](ParallelExecutor::run_indexed) **must**:
///
/// 1. Invoke `task(index, context)` for **every** `index` in `0..num_tasks`
/// **exactly once**.
/// 2. Return only after the last invocation has completed — no invocation may
/// still be running, or able to start, once `run_indexed` returns.
/// 3. On unwind, either complete every remaining index or unwind through the
/// caller; it must not return normally having skipped one.
///
/// Melinoe depends on (1) and (2) for soundness, not merely correctness: it
/// hands each task a pointer into a `MaybeUninit` out-buffer and reconstructs
/// a fully-initialized `Vec` on return, so a skipped index reads uninitialized
/// memory and a torn return aliases it. Violating the contract is undefined
/// behaviour, which is why [`run_indexed`](ParallelExecutor::run_indexed) is
/// `unsafe` to *invoke* and why implementations are expected to document how
/// they discharge it.
///
/// # `Send` and the context pointer
///
/// `context` is a raw pointer because the tasks must remain type-erased across
/// the global slot. An implementation that moves it to worker threads must do
/// so soundly; where the pointer was laundered through `usize` to pass the
/// compiler, prefer a `Send` wrapper type so the obligation is visible.
///
/// # Safety
///
/// `run_indexed` is unsafe to invoke because its *caller* must uphold the
/// mirror-image obligations: `context` must be a live pointer of the type the
/// task function expects, valid for the whole call, and the task function must
/// expect it. Implementors of the trait do not choose those; Melinoe does.
///
/// # Example
///
/// A sequential stand-in, useful in tests:
///
/// ```
/// # use melinoe::sync::ParallelExecutor;
/// struct Sequentially;
///
/// // SAFETY: the loop runs every index in `0..num_tasks` exactly once and
/// // returns only after the last invocation, as the contract requires.
/// unsafe impl ParallelExecutor for Sequentially {
/// unsafe fn run_indexed(
/// num_tasks: usize,
/// task: unsafe fn(usize, *mut ()),
/// context: *mut (),
/// ) {
/// for index in 0..num_tasks {
/// // SAFETY: forwarded from the caller; this implementation invokes
/// // each index exactly once with the caller's context.
/// unsafe { task(index, context) };
/// }
/// }
/// }
/// ```
pub unsafe
/// The ABI of the registered shim: type-erased, monomorphized per implementor
/// type by [`register`].
type ExecutorFn = unsafe fn;
/// A validated, process-global parallel executor.
///
/// Construct one with [`Executor::new`] and hand it to
/// [`register_parallel_executor`] when an integration needs to hold the
/// function-pointer capability before registration.
;
const _: = assert!;
static PARALLEL_EXECUTOR: = new;
/// Register a global parallel executor to run partition tasks.
///
/// If registered, Melinoe's partition drivers execute their shards on `E`
/// instead of spawning raw OS threads via `std::thread::scope`.
///
/// See the module documentation for why registration is process-global and
/// order-sensitive.
///
/// # Example
///
/// ```
/// # use melinoe::sync::{ParallelExecutor, register_parallel_executor};
/// struct Sequentially;
///
/// // SAFETY: every index in `0..num_tasks` is invoked exactly once, and the
/// // call returns only after the last invocation completes.
/// unsafe impl ParallelExecutor for Sequentially {
/// unsafe fn run_indexed(
/// num_tasks: usize,
/// task: unsafe fn(usize, *mut ()),
/// context: *mut (),
/// ) {
/// for index in 0..num_tasks {
/// // SAFETY: forwarded from the caller.
/// unsafe { task(index, context) };
/// }
/// }
/// }
///
/// register_parallel_executor::<Sequentially>();
/// ```
/// Clear the registered parallel executor, restoring the default scoped-thread
/// partition driver.
///
/// This is primarily a lifecycle and test-isolation hook for integrations that
/// install a process-global scheduler temporarily. Existing partition calls
/// that have already loaded the executor continue under that call's chosen
/// driver; later calls use the default path.
pub