Skip to main content

moirai_executor/
lib.rs

1//! # Hybrid Executor Implementation
2//!
3//! This crate provides a high-performance hybrid executor that combines
4//! synchronous, asynchronous, and blocking execution on **one** unified
5//! scheduler facade. Synchronous and async-ready work use the compute
6//! work-stealing pool; potentially blocking work uses a lazily initialized,
7//! bounded lane owned by that scheduler.
8//!
9//! ## Architecture Overview
10//!
11//! - **Static Work-Class Routing**: sync, async, and blocking jobs are routed
12//!   by zero-sized work-class markers; blocking admission is isolated from the
13//!   compute worker pool.
14//! - **Priority-Partitioned Queues**: per-worker Chase-Lev deques indexed by
15//!   [`moirai_core::Priority::index`].
16//! - **Zero-Copy Task Passing**: minimal overhead task distribution.
17
18#![cfg_attr(nightly_tls_active, feature(thread_local))]
19#![deny(missing_docs)]
20
21// Module declarations - following SRP and SOC principles
22pub mod hybrid;
23pub mod metrics;
24pub mod registry;
25pub mod schedule;
26pub mod task;
27
28#[cfg(test)]
29mod counting_wake;
30
31// Re-export key types for clean API
32pub use hybrid::HybridExecutor;
33pub use metrics::ExecutorMetrics;
34pub use registry::{RetentionPolicy, TaskRegistry};
35pub use schedule::{
36    AcceleratorCounts, AcceleratorId, AcceleratorKind, AcceleratorRoute, AcceleratorRoutePolicy,
37    AsyncLaneId, AsyncLanesPerProcess, AsyncTask, BlockingTask, HybridRoutePolicy, HybridRouter,
38    IdleHook, IdleHookRegistrationError, MAX_IDLE_HOOKS, ProcessCount, ProcessId, ProcessRoute,
39    RoutePolicy, RouteSummary, RouteTopology, ScheduleMetrics, SchedulerRoute, SchedulerScope,
40    ServerCount, ServerId, ServerRoute, ServerRoutePolicy, SyncTask, ThreadId, ThreadRoute,
41    ThreadRoutePolicy, ThreadScheduler, WorkClass, WorkerCount, register_idle_hook, run_idle_hooks,
42};
43#[cfg(feature = "scheduler-diagnostics")]
44pub use schedule::{
45    ContendedWakeDecision, DiagnosticWakeDecision, EmptyWakeDecision, SaturatedWakeDecision,
46};
47pub use task::TaskMetadata;
48
49/// Block the current thread until `future` resolves.
50///
51/// This is the Moirai-owned synchronous wait primitive for code that only needs
52/// to bridge an async operation into a synchronous boundary. It uses the same
53/// parking waker as [`moirai_core::executor::ExecutorControl::block_on`] without constructing or
54/// touching the process-wide scheduler.
55pub fn block_on<F>(future: F) -> F::Output
56where
57    F: core::future::Future,
58{
59    schedule::wake::block_on_current_thread(future)
60}
61
62/// Main executor builder for creating configured instances
63pub struct ExecutorBuilder {
64    worker_threads: usize,
65    async_threads: usize,
66    worker_placement: moirai_core::executor::WorkerPlacement,
67}
68
69impl ExecutorBuilder {
70    /// Create a new executor builder with default settings
71    pub fn new() -> Self {
72        Self {
73            worker_threads: themis::CpuTopology::detect()
74                .map(|topology| topology.logical_processors())
75                .or_else(|| std::thread::available_parallelism().ok().map(|n| n.get()))
76                .unwrap_or(4)
77                .max(1),
78            async_threads: 4,
79            worker_placement: moirai_core::executor::WorkerPlacement::default(),
80        }
81    }
82
83    /// Set the number of worker threads
84    pub fn worker_threads(mut self, count: usize) -> Self {
85        self.worker_threads = count;
86        self
87    }
88
89    /// Set whether each worker is confined to one logical processor
90    pub fn worker_placement(mut self, placement: moirai_core::executor::WorkerPlacement) -> Self {
91        self.worker_placement = placement;
92        self
93    }
94
95    /// Set the number of async threads
96    pub fn async_threads(mut self, count: usize) -> Self {
97        self.async_threads = count;
98        self
99    }
100
101    /// Build the hybrid executor
102    pub fn build(self) -> Result<HybridExecutor, Box<dyn std::error::Error>> {
103        let config = moirai_core::executor::ExecutorConfig {
104            worker_threads: self.worker_threads,
105            async_threads: self.async_threads,
106            worker_placement: self.worker_placement,
107            ..moirai_core::executor::ExecutorConfig::default()
108        };
109        HybridExecutor::new(config).map_err(|e| Box::new(e) as Box<dyn std::error::Error>)
110    }
111}
112
113impl Default for ExecutorBuilder {
114    fn default() -> Self {
115        Self::new()
116    }
117}
118
119/// Address-carrying wrapper that lets the bridge move a type-erased context
120/// pointer into `Send` task closures.
121///
122/// The pointee is opaque here — Melinoe owns its type and guarantees it stays
123/// live and unaliased for the whole call — so the only thing this wrapper
124/// asserts is that *moving the address* to another worker is sound. That is
125/// Melinoe's own obligation, discharged in `TaskContext`'s documentation.
126#[derive(Copy, Clone)]
127struct SendContext(*mut ());
128
129// SAFETY: the pointed-to `TaskContext` is documented by Melinoe as valid for
130// the whole executor call and as permitting concurrent field access from
131// distinct tasks. Reading the address on another thread is therefore sound;
132// this wrapper is moved and copied, never dereferenced by this crate.
133unsafe impl Send for SendContext {}
134
135// SAFETY: `Sync` is required because the closure that captures this wrapper is
136// shared across the pool. Sharing the *address* is sound for the same reason as
137// `Send` — the pointee is never accessed through this wrapper, and Melinoe
138// guarantees concurrent access to distinct fields of its context is disjoint.
139unsafe impl Sync for SendContext {}
140
141impl SendContext {
142    /// Recover the erased context pointer.
143    ///
144    /// Reading the field through a method rather than at the field keeps the
145    /// `*mut ()` out of any closure's capture analysis: the closure moves a
146    /// `SendContext` and calls a method, so the raw type never appears in its
147    /// environment.
148    #[inline]
149    fn get(self) -> *mut () {
150        self.0
151    }
152}
153
154/// Moirai's implementation of Melinoe's parallel-executor contract.
155///
156/// Drives a partition's tasks on the shared work-stealing pool, so a branded
157/// partition pays no OS-thread spawn. See [`melinoe::sync::ParallelExecutor`]
158/// for the contract this discharges.
159struct MoiraiExecutor;
160
161// SAFETY: `global().for_each_indexed` owns the complete `0..num_tasks` domain
162// and invokes its closure exactly once per index; it blocks until every
163// scheduled invocation has completed, so the method satisfies both the
164// "every index exactly once" and "no invocation outliving the return"
165// obligations. On scheduler failure it panics *after* joining the scheduled
166// invocations, which is the unwind path the contract permits — Melinoe's
167// `ExecutorDropGuard` contains the omitted slots. The context pointer is
168// forwarded unchanged and never outlives the blocking call.
169unsafe impl melinoe::ParallelExecutor for MoiraiExecutor {
170    unsafe fn run_indexed(num_tasks: usize, task: unsafe fn(usize, *mut ()), context: *mut ()) {
171        // Bind the type-erased context into a `Send + Sync` wrapper before the
172        // closure exists, so the raw pointer is absent from its capture set.
173        // Function pointers already satisfy the pool's `Send + Sync` bound.
174        let context = SendContext(context);
175        let res = global().for_each_indexed::<SyncTask, _>(num_tasks, move |index| {
176            // SAFETY: forwarded from the caller. `for_each_indexed` invokes this
177            // closure exactly once per index and never concurrently for the same
178            // index, so each call addresses a distinct slot of Melinoe's context.
179            unsafe {
180                task(index, context.get());
181            }
182        });
183        if let Err(e) = res {
184            panic!(
185                "Moirai executor failure in Melinoe parallel driver: {:?}",
186                e
187            );
188        }
189    }
190}
191
192fn global_arc() -> &'static std::sync::Arc<HybridExecutor> {
193    static GLOBAL_EXECUTOR: std::sync::OnceLock<std::sync::Arc<HybridExecutor>> =
194        std::sync::OnceLock::new();
195    GLOBAL_EXECUTOR.get_or_init(|| {
196        let executor = std::sync::Arc::new(
197            ExecutorBuilder::new()
198                .build()
199                .expect("initialize global Moirai executor"),
200        );
201        // Register after the pool exists so a re-entrant callback can never
202        // observe a partially initialized scheduler.
203        melinoe::register_parallel_executor::<MoiraiExecutor>();
204        executor
205    })
206}
207
208/// Initialize the shared executor and install its Melinoe partition bridge.
209///
210/// Call this during application startup when code may invoke
211/// `melinoe::sync::partition_*` directly. The function is idempotent: the
212/// scheduler is built once, and the process-global Melinoe slot is refreshed on
213/// each call. Higher-level Moirai partition helpers initialize the same bridge
214/// automatically when they enter the pool path.
215pub fn initialize() {
216    let _ = global_arc();
217    // `clear_parallel_executor` is a supported lifecycle hook for tests and
218    // integrations. Refresh the slot after such a reset without adding an
219    // atomic store to every ordinary `global()` access.
220    melinoe::register_parallel_executor::<MoiraiExecutor>();
221}
222
223/// Borrow the shared, lazily-initialized process-wide executor.
224///
225/// Provides a single default runtime so higher-level crates (e.g.
226/// `moirai-parallel`'s data-parallel primitives) can schedule work without each
227/// constructing — and over-subscribing — their own thread pool. Built once with
228/// the default [`ExecutorBuilder`] configuration on first access.
229///
230/// # Panics
231///
232/// Panics if the executor cannot be initialized, which should not happen under
233/// normal conditions.
234pub fn global() -> &'static HybridExecutor {
235    global_arc()
236}
237
238/// Obtain an owned handle to the shared process-wide executor.
239///
240/// Higher layers (e.g. the `moirai` umbrella's global runtime) wrap this same
241/// `Arc` so that parallel data-parallel work and async tasks run on **one**
242/// unified hybrid scheduler rather than separate thread pools.
243pub fn shared() -> std::sync::Arc<HybridExecutor> {
244    std::sync::Arc::clone(global_arc())
245}