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}