Skip to main content

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;