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}