tenferro_cpu/domain_executor.rs
1use std::fmt::Debug;
2use std::num::NonZeroUsize;
3use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
4use std::sync::Arc;
5
6use rayon::prelude::*;
7
8/// Inner parallel-region support offered by a CPU domain executor.
9///
10/// # Examples
11///
12/// ```rust
13/// use tenferro_cpu::CpuInnerParallelism;
14///
15/// assert_ne!(CpuInnerParallelism::None, CpuInnerParallelism::Rayon);
16/// ```
17#[derive(Clone, Copy, Debug, Eq, PartialEq)]
18pub enum CpuInnerParallelism {
19 /// The executor cannot host provider-owned inner parallel regions.
20 None,
21 /// The executor can host a Rayon-compatible inner parallel region.
22 Rayon,
23}
24
25/// Re-entry capability of one CPU domain executor.
26///
27/// This describes executor-level same-executor entry only. It never grants
28/// permission for recursive public [`crate::CpuBackend`] entry, which remains a
29/// separate backend contract.
30///
31/// # Examples
32///
33/// ```rust
34/// use tenferro_cpu::CpuExecutorReentrancy;
35///
36/// let policy = CpuExecutorReentrancy::Rejected;
37/// assert_eq!(policy, CpuExecutorReentrancy::Rejected);
38/// ```
39#[derive(Clone, Copy, Debug, Eq, PartialEq)]
40pub enum CpuExecutorReentrancy {
41 /// Nested entry into the same executor is rejected.
42 Rejected,
43 /// The executor supports nested entry into that same executor.
44 SameExecutor,
45}
46
47/// Affinity claim made by a CPU domain executor.
48///
49/// # Examples
50///
51/// ```rust
52/// use tenferro_cpu::CpuExecutorAffinity;
53///
54/// let affinity = CpuExecutorAffinity::CallerDeclaredUnverified;
55/// assert_ne!(affinity, CpuExecutorAffinity::TenferroDomainVerified);
56/// ```
57#[derive(Clone, Copy, Debug, Eq, PartialEq)]
58pub enum CpuExecutorAffinity {
59 /// Tenferro confined every worker to the declared domain CPU set and verified
60 /// the resulting mask.
61 ///
62 /// The claim covers the domain's CPU set, not a distinct CPU per worker. A
63 /// provider that creates its own threads therefore inherits the whole set.
64 TenferroDomainVerified,
65 /// The caller declared worker placement, but tenferro did not verify it.
66 CallerDeclaredUnverified,
67 /// The executor makes no worker-placement claim.
68 None,
69}
70
71/// Ownership of CPU executor shutdown.
72///
73/// # Examples
74///
75/// ```rust
76/// use tenferro_cpu::CpuExecutorShutdown;
77///
78/// assert_ne!(
79/// CpuExecutorShutdown::TenferroOwned,
80/// CpuExecutorShutdown::CallerOwned,
81/// );
82/// ```
83#[derive(Clone, Copy, Debug, Eq, PartialEq)]
84pub enum CpuExecutorShutdown {
85 /// Tenferro owns executor shutdown.
86 TenferroOwned,
87 /// The caller owns executor shutdown and executor lifetime policy.
88 CallerOwned,
89}
90
91/// Immutable construction-time capabilities of a CPU domain executor.
92///
93/// # Examples
94///
95/// ```rust
96/// use std::num::NonZeroUsize;
97/// use tenferro_cpu::{
98/// CpuDomainExecutorCapabilities, CpuExecutorAffinity, CpuExecutorReentrancy,
99/// CpuExecutorShutdown, CpuInnerParallelism,
100/// };
101///
102/// let capabilities = CpuDomainExecutorCapabilities {
103/// worker_count: NonZeroUsize::new(4).unwrap(),
104/// outer_parallelism: true,
105/// inner_parallelism: CpuInnerParallelism::Rayon,
106/// reentrancy: CpuExecutorReentrancy::Rejected,
107/// affinity: CpuExecutorAffinity::TenferroDomainVerified,
108/// shutdown: CpuExecutorShutdown::TenferroOwned,
109/// };
110/// assert_eq!(capabilities.worker_count.get(), 4);
111/// ```
112#[derive(Clone, Copy, Debug, Eq, PartialEq)]
113pub struct CpuDomainExecutorCapabilities {
114 /// Number of workers made available to this domain.
115 pub worker_count: NonZeroUsize,
116 /// Whether indexed outer fork/join submission is supported.
117 pub outer_parallelism: bool,
118 /// Provider-owned inner parallel-region support.
119 pub inner_parallelism: CpuInnerParallelism,
120 /// Same-executor re-entry capability.
121 pub reentrancy: CpuExecutorReentrancy,
122 /// Worker-affinity claim and verification level.
123 pub affinity: CpuExecutorAffinity,
124 /// Executor shutdown owner.
125 pub shutdown: CpuExecutorShutdown,
126}
127
128/// Failure at the CPU executor admission or scheduling boundary.
129///
130/// Operation and provider errors do not belong in this type. Executors use
131/// these variants only for their own admission, scheduling, cancellation, and
132/// panic-bridge failures.
133///
134/// # Examples
135///
136/// ```rust
137/// use tenferro_cpu::CpuDomainExecutorError;
138///
139/// let error = CpuDomainExecutorError::Admission {
140/// message: "domain is busy".to_string(),
141/// };
142/// assert!(matches!(error, CpuDomainExecutorError::Admission { .. }));
143/// ```
144#[derive(Clone, Debug, Eq, PartialEq, thiserror::Error)]
145pub enum CpuDomainExecutorError {
146 /// The executor rejected entry before scheduling work.
147 #[error("CPU domain executor admission failed: {message}")]
148 Admission {
149 /// Executor-owned diagnostic.
150 message: String,
151 },
152 /// The executor could not schedule or complete submitted work.
153 #[error("CPU domain executor scheduling failed: {message}")]
154 Scheduling {
155 /// Executor-owned diagnostic.
156 message: String,
157 },
158 /// The executor cancelled submitted work.
159 #[error("CPU domain executor cancelled work: {message}")]
160 Cancellation {
161 /// Executor-owned diagnostic.
162 message: String,
163 },
164 /// The executor converted a worker panic into a typed failure.
165 #[error("CPU domain executor worker panicked: {message}")]
166 PanicBridge {
167 /// Executor-owned diagnostic.
168 message: String,
169 },
170}
171
172/// One borrowed job installed synchronously into a CPU domain executor.
173///
174/// The executor must finish using the job before [`CpuDomainExecutor::install`]
175/// returns; a borrowed job never escapes that call.
176///
177/// # Examples
178///
179/// ```rust
180/// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJob};
181///
182/// struct Job(bool);
183/// impl ScopedCpuJob for Job {
184/// fn run(&mut self) -> Result<(), CpuDomainExecutorError> {
185/// self.0 = true;
186/// Ok(())
187/// }
188/// }
189/// let mut job = Job(false);
190/// job.run().unwrap();
191/// assert!(job.0);
192/// ```
193pub trait ScopedCpuJob: Send {
194 /// Run this job once on the executor-selected calling context.
195 ///
196 /// # Examples
197 ///
198 /// ```rust
199 /// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJob};
200 ///
201 /// struct Job;
202 /// impl ScopedCpuJob for Job {
203 /// fn run(&mut self) -> Result<(), CpuDomainExecutorError> { Ok(()) }
204 /// }
205 /// assert!(Job.run().is_ok());
206 /// ```
207 ///
208 /// # Errors
209 ///
210 /// Returns [`CpuDomainExecutorError::Admission`],
211 /// [`CpuDomainExecutorError::Scheduling`],
212 /// [`CpuDomainExecutorError::Cancellation`], or
213 /// [`CpuDomainExecutorError::PanicBridge`] when that failure is observed
214 /// while running the job.
215 fn run(&mut self) -> Result<(), CpuDomainExecutorError>;
216}
217
218/// Synchronously submitted indexed jobs for engine-owned outer scheduling.
219///
220/// Implementations expose a borrowed logical range `0..len()` without
221/// allocating a job collection. Every indexed call must finish before
222/// [`CpuDomainExecutor::submit`] returns.
223///
224/// # Examples
225///
226/// ```rust
227/// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJobs};
228///
229/// struct Jobs;
230/// impl ScopedCpuJobs for Jobs {
231/// fn len(&self) -> usize { 2 }
232/// fn run(&self, index: usize) -> Result<(), CpuDomainExecutorError> {
233/// assert!(index < self.len());
234/// Ok(())
235/// }
236/// }
237/// let jobs: &dyn ScopedCpuJobs = &Jobs;
238/// assert_eq!(jobs.len(), 2);
239/// jobs.run(1).unwrap();
240/// ```
241pub trait ScopedCpuJobs: Sync {
242 /// Return the number of indexed jobs in this synchronous submission.
243 ///
244 /// # Examples
245 ///
246 /// ```rust
247 /// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJobs};
248 ///
249 /// struct Jobs;
250 /// impl ScopedCpuJobs for Jobs {
251 /// fn len(&self) -> usize { 3 }
252 /// fn run(&self, _index: usize) -> Result<(), CpuDomainExecutorError> { Ok(()) }
253 /// }
254 /// assert_eq!(Jobs.len(), 3);
255 /// ```
256 fn len(&self) -> usize;
257
258 /// Return whether this submission contains no indexed jobs.
259 ///
260 /// # Examples
261 ///
262 /// ```rust
263 /// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJobs};
264 ///
265 /// struct Jobs;
266 /// impl ScopedCpuJobs for Jobs {
267 /// fn len(&self) -> usize { 0 }
268 /// fn run(&self, _index: usize) -> Result<(), CpuDomainExecutorError> { Ok(()) }
269 /// }
270 /// assert!(Jobs.is_empty());
271 /// ```
272 fn is_empty(&self) -> bool {
273 self.len() == 0
274 }
275
276 /// Run one indexed job synchronously.
277 ///
278 /// # Examples
279 ///
280 /// ```rust
281 /// use tenferro_cpu::{CpuDomainExecutorError, ScopedCpuJobs};
282 ///
283 /// struct Jobs;
284 /// impl ScopedCpuJobs for Jobs {
285 /// fn len(&self) -> usize { 1 }
286 /// fn run(&self, index: usize) -> Result<(), CpuDomainExecutorError> {
287 /// assert_eq!(index, 0);
288 /// Ok(())
289 /// }
290 /// }
291 /// Jobs.run(0).unwrap();
292 /// ```
293 ///
294 /// # Errors
295 ///
296 /// Returns [`CpuDomainExecutorError::Admission`],
297 /// [`CpuDomainExecutorError::Scheduling`],
298 /// [`CpuDomainExecutorError::Cancellation`], or
299 /// [`CpuDomainExecutorError::PanicBridge`] when that failure is observed
300 /// while running the indexed job.
301 fn run(&self, index: usize) -> Result<(), CpuDomainExecutorError>;
302}
303
304/// Object-safe synchronous executor for one CPU resource domain.
305///
306/// `submit` is an indexed fork/join boundary and `install` is one borrowed
307/// provider-owned inner-region entry. Neither method may retain its borrowed
308/// job after returning.
309///
310/// # Examples
311///
312/// ```rust
313/// use std::num::NonZeroUsize;
314/// use tenferro_cpu::{
315/// CpuDomainExecutor, CpuDomainExecutorCapabilities, CpuDomainExecutorError,
316/// CpuExecutorAffinity, CpuExecutorReentrancy, CpuExecutorShutdown,
317/// CpuInnerParallelism, ScopedCpuJob, ScopedCpuJobs,
318/// };
319///
320/// #[derive(Debug)]
321/// struct Inline;
322/// impl CpuDomainExecutor for Inline {
323/// fn capabilities(&self) -> CpuDomainExecutorCapabilities {
324/// CpuDomainExecutorCapabilities {
325/// worker_count: NonZeroUsize::new(1).unwrap(),
326/// outer_parallelism: false,
327/// inner_parallelism: CpuInnerParallelism::None,
328/// reentrancy: CpuExecutorReentrancy::Rejected,
329/// affinity: CpuExecutorAffinity::None,
330/// shutdown: CpuExecutorShutdown::CallerOwned,
331/// }
332/// }
333/// fn submit(&self, _jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError> {
334/// Ok(())
335/// }
336/// fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError> {
337/// job.run()
338/// }
339/// }
340/// let executor: &dyn CpuDomainExecutor = &Inline;
341/// assert_eq!(executor.capabilities().worker_count.get(), 1);
342/// ```
343pub trait CpuDomainExecutor: Debug + Send + Sync + 'static {
344 /// Return immutable construction-time executor capabilities.
345 ///
346 /// # Examples
347 ///
348 /// ```rust
349 /// use std::num::NonZeroUsize;
350 /// use tenferro_cpu::{
351 /// CpuDomainExecutor, CpuDomainExecutorCapabilities, CpuDomainExecutorError,
352 /// CpuExecutorAffinity, CpuExecutorReentrancy, CpuExecutorShutdown,
353 /// CpuInnerParallelism, ScopedCpuJob, ScopedCpuJobs,
354 /// };
355 /// # #[derive(Debug)] struct Inline;
356 /// # impl CpuDomainExecutor for Inline {
357 /// # fn capabilities(&self) -> CpuDomainExecutorCapabilities {
358 /// # CpuDomainExecutorCapabilities { worker_count: NonZeroUsize::new(2).unwrap(),
359 /// # outer_parallelism: true, inner_parallelism: CpuInnerParallelism::None,
360 /// # reentrancy: CpuExecutorReentrancy::Rejected, affinity: CpuExecutorAffinity::None,
361 /// # shutdown: CpuExecutorShutdown::CallerOwned }
362 /// # }
363 /// # fn submit(&self, _jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError> { Ok(()) }
364 /// # fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError> { job.run() }
365 /// # }
366 /// let executor: &dyn CpuDomainExecutor = &Inline;
367 /// assert_eq!(executor.capabilities().worker_count.get(), 2);
368 /// ```
369 fn capabilities(&self) -> CpuDomainExecutorCapabilities;
370
371 /// Submit all indexed jobs as one synchronous fork/join operation.
372 ///
373 /// All `0..jobs.len()` jobs must be complete when this method returns.
374 ///
375 /// # Examples
376 ///
377 /// ```rust
378 /// use std::num::NonZeroUsize;
379 /// use std::sync::atomic::{AtomicUsize, Ordering};
380 /// use tenferro_cpu::{
381 /// CpuDomainExecutor, CpuDomainExecutorCapabilities, CpuDomainExecutorError,
382 /// CpuExecutorAffinity, CpuExecutorReentrancy, CpuExecutorShutdown,
383 /// CpuInnerParallelism, ScopedCpuJob, ScopedCpuJobs,
384 /// };
385 /// # #[derive(Debug)] struct Inline;
386 /// # impl CpuDomainExecutor for Inline {
387 /// # fn capabilities(&self) -> CpuDomainExecutorCapabilities {
388 /// # CpuDomainExecutorCapabilities { worker_count: NonZeroUsize::new(1).unwrap(),
389 /// # outer_parallelism: true, inner_parallelism: CpuInnerParallelism::None,
390 /// # reentrancy: CpuExecutorReentrancy::Rejected, affinity: CpuExecutorAffinity::None,
391 /// # shutdown: CpuExecutorShutdown::CallerOwned }
392 /// # }
393 /// # fn submit(&self, jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError> {
394 /// # for index in 0..jobs.len() { jobs.run(index)?; } Ok(())
395 /// # }
396 /// # fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError> { job.run() }
397 /// # }
398 /// struct Jobs<'a>(&'a AtomicUsize);
399 /// impl ScopedCpuJobs for Jobs<'_> {
400 /// fn len(&self) -> usize { 2 }
401 /// fn run(&self, _index: usize) -> Result<(), CpuDomainExecutorError> {
402 /// self.0.fetch_add(1, Ordering::Relaxed);
403 /// Ok(())
404 /// }
405 /// }
406 /// let count = AtomicUsize::new(0);
407 /// Inline.submit(&Jobs(&count)).unwrap();
408 /// assert_eq!(count.load(Ordering::Relaxed), 2);
409 /// ```
410 ///
411 /// # Errors
412 ///
413 /// Returns [`CpuDomainExecutorError::Admission`],
414 /// [`CpuDomainExecutorError::Scheduling`],
415 /// [`CpuDomainExecutorError::Cancellation`], or
416 /// [`CpuDomainExecutorError::PanicBridge`] for executor-owned failures.
417 fn submit(&self, jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError>;
418
419 /// Enter one synchronous provider-owned inner parallel region.
420 ///
421 /// # Examples
422 ///
423 /// ```rust
424 /// use std::num::NonZeroUsize;
425 /// use tenferro_cpu::{
426 /// CpuDomainExecutor, CpuDomainExecutorCapabilities, CpuDomainExecutorError,
427 /// CpuExecutorAffinity, CpuExecutorReentrancy, CpuExecutorShutdown,
428 /// CpuInnerParallelism, ScopedCpuJob, ScopedCpuJobs,
429 /// };
430 /// # #[derive(Debug)] struct Inline;
431 /// # impl CpuDomainExecutor for Inline {
432 /// # fn capabilities(&self) -> CpuDomainExecutorCapabilities {
433 /// # CpuDomainExecutorCapabilities { worker_count: NonZeroUsize::new(1).unwrap(),
434 /// # outer_parallelism: false, inner_parallelism: CpuInnerParallelism::None,
435 /// # reentrancy: CpuExecutorReentrancy::Rejected, affinity: CpuExecutorAffinity::None,
436 /// # shutdown: CpuExecutorShutdown::CallerOwned }
437 /// # }
438 /// # fn submit(&self, _jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError> { Ok(()) }
439 /// # fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError> { job.run() }
440 /// # }
441 /// struct Job(bool);
442 /// impl ScopedCpuJob for Job {
443 /// fn run(&mut self) -> Result<(), CpuDomainExecutorError> {
444 /// self.0 = true;
445 /// Ok(())
446 /// }
447 /// }
448 /// let mut job = Job(false);
449 /// Inline.install(&mut job).unwrap();
450 /// assert!(job.0);
451 /// ```
452 ///
453 /// # Errors
454 ///
455 /// Returns [`CpuDomainExecutorError::Admission`],
456 /// [`CpuDomainExecutorError::Scheduling`],
457 /// [`CpuDomainExecutorError::Cancellation`], or
458 /// [`CpuDomainExecutorError::PanicBridge`] for executor-owned failures.
459 fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError>;
460
461 /// The Rayon pool this executor runs jobs on, when it has one.
462 ///
463 /// Providers reach it through [`crate::CpuExecutionContext::rayon_pool`],
464 /// which exposes it only to contexts that own an inner parallel region.
465 /// Executors that are not backed by one Rayon pool keep the default
466 /// `None`.
467 ///
468 /// # Examples
469 ///
470 /// ```rust
471 /// use std::sync::Arc;
472 /// use tenferro_cpu::{CpuDomainExecutor, RayonCpuDomainExecutor};
473 /// let pool = Arc::new(rayon::ThreadPoolBuilder::new().num_threads(2).build()?);
474 /// assert!(RayonCpuDomainExecutor::new(pool).rayon_pool().is_some());
475 /// # Ok::<(), rayon::ThreadPoolBuildError>(())
476 /// ```
477 fn rayon_pool(&self) -> Option<&rayon::ThreadPool> {
478 None
479 }
480}
481
482/// Adapter that executes CPU-domain jobs on one caller-owned Rayon pool.
483///
484/// The adapter retains the supplied pool and never creates, reconfigures, or
485/// shuts it down. [`CpuDomainExecutor`] remains the primary injection contract;
486/// this type is only a convenience for Rayon hosts.
487///
488/// # Examples
489///
490/// ```rust
491/// use std::sync::Arc;
492/// use tenferro_cpu::{CpuDomainExecutor, RayonCpuDomainExecutor};
493///
494/// let pool = Arc::new(rayon::ThreadPoolBuilder::new().num_threads(2).build()?);
495/// let executor = RayonCpuDomainExecutor::new(Arc::clone(&pool));
496/// assert_eq!(executor.capabilities().worker_count.get(), 2);
497/// assert_eq!(Arc::strong_count(&pool), 2);
498/// # Ok::<(), rayon::ThreadPoolBuildError>(())
499/// ```
500pub struct RayonCpuDomainExecutor {
501 pool: Arc<rayon::ThreadPool>,
502}
503
504impl std::fmt::Debug for RayonCpuDomainExecutor {
505 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
506 formatter
507 .debug_struct("RayonCpuDomainExecutor")
508 .field("worker_count", &self.pool.current_num_threads())
509 .finish_non_exhaustive()
510 }
511}
512
513impl RayonCpuDomainExecutor {
514 /// Retain one caller-owned Rayon pool as a CPU-domain executor.
515 ///
516 /// # Examples
517 ///
518 /// ```rust
519 /// use std::sync::Arc;
520 /// use tenferro_cpu::{CpuDomainExecutor, RayonCpuDomainExecutor};
521 ///
522 /// let pool = Arc::new(rayon::ThreadPoolBuilder::new().num_threads(2).build()?);
523 /// let executor = RayonCpuDomainExecutor::new(pool);
524 /// assert_eq!(executor.capabilities().worker_count.get(), 2);
525 /// # Ok::<(), rayon::ThreadPoolBuildError>(())
526 /// ```
527 pub fn new(pool: Arc<rayon::ThreadPool>) -> Self {
528 Self { pool }
529 }
530}
531
532impl CpuDomainExecutor for RayonCpuDomainExecutor {
533 fn capabilities(&self) -> CpuDomainExecutorCapabilities {
534 // INVARIANT: Rayon rejects thread pools with zero workers.
535 let worker_count =
536 NonZeroUsize::new(self.pool.current_num_threads()).unwrap_or(NonZeroUsize::MIN);
537 CpuDomainExecutorCapabilities {
538 worker_count,
539 outer_parallelism: worker_count.get() > 1,
540 inner_parallelism: CpuInnerParallelism::Rayon,
541 reentrancy: CpuExecutorReentrancy::SameExecutor,
542 affinity: CpuExecutorAffinity::None,
543 shutdown: CpuExecutorShutdown::CallerOwned,
544 }
545 }
546
547 fn submit(&self, jobs: &dyn ScopedCpuJobs) -> Result<(), CpuDomainExecutorError> {
548 self.pool.install(|| {
549 (0..jobs.len())
550 .into_par_iter()
551 .try_for_each(|index| jobs.run(index))
552 })
553 }
554
555 fn install(&self, job: &mut dyn ScopedCpuJob) -> Result<(), CpuDomainExecutorError> {
556 self.pool.install(|| job.run())
557 }
558
559 fn rayon_pool(&self) -> Option<&rayon::ThreadPool> {
560 Some(&self.pool)
561 }
562}
563
564pub(crate) struct ScopedJob<F, R> {
565 operation: Option<F>,
566 result: Option<R>,
567}
568
569pub(crate) fn scoped_job<F, R>(operation: F) -> ScopedJob<F, R>
570where
571 F: FnOnce() -> R + Send,
572 R: Send,
573{
574 ScopedJob {
575 operation: Some(operation),
576 result: None,
577 }
578}
579
580impl<F, R> ScopedJob<F, R> {
581 fn into_result(self) -> Result<R, CpuDomainExecutorError> {
582 self.result
583 .ok_or_else(|| CpuDomainExecutorError::Scheduling {
584 message: "executor returned success without running the scoped CPU job".to_string(),
585 })
586 }
587}
588
589impl<F, R> ScopedCpuJob for ScopedJob<F, R>
590where
591 F: FnOnce() -> R + Send,
592 R: Send,
593{
594 fn run(&mut self) -> Result<(), CpuDomainExecutorError> {
595 let operation =
596 self.operation
597 .take()
598 .ok_or_else(|| CpuDomainExecutorError::Scheduling {
599 message: "executor attempted to run a scoped CPU job more than once"
600 .to_string(),
601 })?;
602 self.result = Some(operation());
603 Ok(())
604 }
605}
606
607pub(crate) fn install_scoped<F, R>(
608 executor: &dyn CpuDomainExecutor,
609 operation: F,
610) -> Result<R, CpuDomainExecutorError>
611where
612 F: FnOnce() -> R + Send,
613 R: Send,
614{
615 let mut job = scoped_job(operation);
616 executor.install(&mut job)?;
617 job.into_result()
618}
619
620pub(crate) struct IndexedJobs<F> {
621 len: usize,
622 run: F,
623 invalid_index_attempt: InvalidIndexAudit,
624}
625
626const INVALID_INDEX_EMPTY: u8 = 0;
627const INVALID_INDEX_WRITING: u8 = 1;
628const INVALID_INDEX_READY: u8 = 2;
629
630// INVARIANT: the two-phase state publishes every usize value, including
631// usize::MAX, without a sentinel collision. Valid `run` calls never touch this
632// audit, and the post-submit Acquire observes the selected invalid index after
633// its Release publication without locking or allocating.
634struct InvalidIndexAudit {
635 state: AtomicU8,
636 index: AtomicUsize,
637}
638
639impl InvalidIndexAudit {
640 const fn new() -> Self {
641 Self {
642 state: AtomicU8::new(INVALID_INDEX_EMPTY),
643 index: AtomicUsize::new(0),
644 }
645 }
646
647 fn record(&self, index: usize) {
648 if self
649 .state
650 .compare_exchange(
651 INVALID_INDEX_EMPTY,
652 INVALID_INDEX_WRITING,
653 Ordering::AcqRel,
654 Ordering::Acquire,
655 )
656 .is_ok()
657 {
658 self.index.store(index, Ordering::Relaxed);
659 self.state.store(INVALID_INDEX_READY, Ordering::Release);
660 }
661 }
662
663 fn load(&self) -> Option<usize> {
664 (self.state.load(Ordering::Acquire) == INVALID_INDEX_READY)
665 .then(|| self.index.load(Ordering::Relaxed))
666 }
667}
668
669pub(crate) fn indexed_jobs<F>(len: usize, run: F) -> IndexedJobs<F>
670where
671 F: Fn(usize) -> Result<(), CpuDomainExecutorError> + Sync,
672{
673 IndexedJobs {
674 len,
675 run,
676 invalid_index_attempt: InvalidIndexAudit::new(),
677 }
678}
679
680impl<F> IndexedJobs<F> {
681 pub(crate) fn invalid_index_attempt(&self) -> Option<usize> {
682 self.invalid_index_attempt.load()
683 }
684}
685
686impl<F> ScopedCpuJobs for IndexedJobs<F>
687where
688 F: Fn(usize) -> Result<(), CpuDomainExecutorError> + Sync,
689{
690 fn len(&self) -> usize {
691 self.len
692 }
693
694 fn run(&self, index: usize) -> Result<(), CpuDomainExecutorError> {
695 if index >= self.len {
696 self.invalid_index_attempt.record(index);
697 return Err(CpuDomainExecutorError::Scheduling {
698 message: format!(
699 "executor requested scoped CPU job index {index}, but the submission has {} jobs",
700 self.len
701 ),
702 });
703 }
704 (self.run)(index)
705 }
706}
707
708#[cfg(test)]
709mod tests;