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
20// Module declarations - following SRP and SOC principles
21pub mod hybrid;
22pub mod metrics;
23pub mod registry;
24pub mod schedule;
25pub mod task;
26
27// Re-export key types for clean API
28pub use hybrid::HybridExecutor;
29pub use metrics::ExecutorMetrics;
30pub use registry::TaskRegistry;
31pub use schedule::{
32    AcceleratorCounts, AcceleratorId, AcceleratorKind, AcceleratorRoute, AcceleratorRoutePolicy,
33    AsyncLaneId, AsyncLanesPerProcess, AsyncTask, BlockingTask, HybridRoutePolicy, HybridRouter,
34    ProcessCount, ProcessId, ProcessRoute, RoutePolicy, RouteSummary, RouteTopology,
35    ScheduleMetrics, SchedulerRoute, SchedulerScope, ServerCount, ServerId, ServerRoute,
36    ServerRoutePolicy, SyncTask, ThreadId, ThreadRoute, ThreadRoutePolicy, ThreadScheduler,
37    WorkClass, WorkerCount,
38};
39#[cfg(feature = "scheduler-diagnostics")]
40pub use schedule::{
41    ContendedWakeDecision, DiagnosticWakeDecision, EmptyWakeDecision, SaturatedWakeDecision,
42};
43pub use task::TaskMetadata;
44
45/// Block the current thread until `future` resolves.
46///
47/// This is the Moirai-owned synchronous wait primitive for code that only needs
48/// to bridge an async operation into a synchronous boundary. It uses the same
49/// parking waker as [`moirai_core::executor::ExecutorControl::block_on`] without constructing or
50/// touching the process-wide scheduler.
51pub fn block_on<F>(future: F) -> F::Output
52where
53    F: core::future::Future,
54{
55    schedule::wake::block_on_current_thread(future)
56}
57
58/// Main executor builder for creating configured instances
59pub struct ExecutorBuilder {
60    worker_threads: usize,
61    async_threads: usize,
62}
63
64impl ExecutorBuilder {
65    /// Create a new executor builder with default settings
66    pub fn new() -> Self {
67        Self {
68            worker_threads: std::thread::available_parallelism()
69                .map(|n| n.get())
70                .unwrap_or(4),
71            async_threads: 4,
72        }
73    }
74
75    /// Set the number of worker threads
76    pub fn worker_threads(mut self, count: usize) -> Self {
77        self.worker_threads = count;
78        self
79    }
80
81    /// Set the number of async threads
82    pub fn async_threads(mut self, count: usize) -> Self {
83        self.async_threads = count;
84        self
85    }
86
87    /// Build the hybrid executor
88    pub fn build(self) -> Result<HybridExecutor, Box<dyn std::error::Error>> {
89        let config = moirai_core::executor::ExecutorConfig {
90            worker_threads: self.worker_threads,
91            async_threads: self.async_threads,
92            ..moirai_core::executor::ExecutorConfig::default()
93        };
94        HybridExecutor::new(config).map_err(|e| Box::new(e) as Box<dyn std::error::Error>)
95    }
96}
97
98impl Default for ExecutorBuilder {
99    fn default() -> Self {
100        Self::new()
101    }
102}
103
104/// Address-carrying wrapper that lets the melinoe bridge move a raw data
105/// pointer into `Send` task closures; safety is owed by the bridge caller.
106#[derive(Copy, Clone)]
107struct SendPtr(usize);
108
109unsafe fn melinoe_executor_bridge(
110    num_tasks: usize,
111    task_fn: unsafe fn(usize, *mut ()),
112    data: *mut (),
113) {
114    let data_ptr = SendPtr(data as usize);
115    let res = global().for_each_indexed::<SyncTask, _>(num_tasks, move |index| {
116        let p = data_ptr;
117        // SAFETY: task_fn is called concurrently on separate indices.
118        unsafe {
119            task_fn(index, p.0 as *mut ());
120        }
121    });
122    if let Err(e) = res {
123        panic!(
124            "Moirai executor failure in Melinoe parallel driver: {:?}",
125            e
126        );
127    }
128}
129
130// SAFETY: on success, `for_each_indexed` owns the complete `0..num_tasks`
131// domain and invokes its closure once per index. On scheduler failure, the
132// bridge panics after `for_each_indexed` has joined every scheduled invocation;
133// Melinoe's unwind guard handles omitted slots. The unchanged context pointer
134// never outlives the blocking scheduler call.
135const MELINOE_EXECUTOR: melinoe::ParallelExecutor =
136    unsafe { melinoe::ParallelExecutor::new(melinoe_executor_bridge) };
137
138fn global_arc() -> &'static std::sync::Arc<HybridExecutor> {
139    static GLOBAL_EXECUTOR: std::sync::OnceLock<std::sync::Arc<HybridExecutor>> =
140        std::sync::OnceLock::new();
141    GLOBAL_EXECUTOR.get_or_init(|| {
142        let exec = std::sync::Arc::new(
143            ExecutorBuilder::new()
144                .build()
145                .expect("initialize global Moirai executor"),
146        );
147        // Register the global parallel executor in melinoe.
148        melinoe::register_parallel_executor(MELINOE_EXECUTOR);
149        exec
150    })
151}
152
153/// Borrow the shared, lazily-initialized process-wide executor.
154///
155/// Provides a single default runtime so higher-level crates (e.g.
156/// `moirai-parallel`'s data-parallel primitives) can schedule work without each
157/// constructing — and over-subscribing — their own thread pool. Built once with
158/// the default [`ExecutorBuilder`] configuration on first access.
159///
160/// # Panics
161///
162/// Panics if the executor cannot be initialized, which should not happen under
163/// normal conditions.
164pub fn global() -> &'static HybridExecutor {
165    global_arc()
166}
167
168/// Obtain an owned handle to the shared process-wide executor.
169///
170/// Higher layers (e.g. the `moirai` umbrella's global runtime) wrap this same
171/// `Arc` so that parallel data-parallel work and async tasks run on **one**
172/// unified hybrid scheduler rather than separate thread pools.
173pub fn shared() -> std::sync::Arc<HybridExecutor> {
174    std::sync::Arc::clone(global_arc())
175}