Skip to main content

moirai/
lib.rs

1//! # Moirai - Weaving the Threads of Fate
2//!
3//! Moirai is a high-performance hybrid concurrency library for Rust that seamlessly
4//! blends asynchronous and parallel execution models. Named after the Greek Fates
5//! who controlled the threads of life, Moirai weaves together the best principles
6//! from async task scheduling and parallel work-stealing into a unified framework.
7//!
8//! ## Core Design Principles
9//!
10//! Moirai follows elite programming practices:
11//! - **SOLID**: Single responsibility, open/closed, Liskov substitution, interface segregation, dependency inversion
12//! - **CUPID**: Composable, Unix philosophy, predictable, idiomatic, domain-centric
13//! - **GRASP**: Information expert, creator, controller, low coupling, high cohesion
14//! - **ACID**: Atomicity, consistency, isolation, durability in task execution
15//!
16//! ## Features
17//!
18//! - **Zero-cost abstractions**: All abstractions compile away to optimal code
19//! - **Hybrid execution**: Seamlessly mix async and parallel tasks
20//! - **Work-stealing scheduler**: Intelligent load balancing across CPU cores
21//! - **Memory safety**: Leverage Rust's ownership system for safe concurrency
22//! - **High performance**: Sub-microsecond task scheduling overhead
23//! - **NUMA awareness**: Optimize for modern multi-socket systems
24//! - **Rich iterator combinators**: Parallel and async iterator processing
25//! - **IPC**: Inter-process communication (optional)
26//! - **Metrics**: Performance monitoring (optional)
27//! - **Distributed transport feature gates**: Optional transport and iterator helpers without a
28//!   facade-level remote-closure API
29//!
30//! ## Performance Characteristics
31//!
32//! - **Task scheduling overhead**: < 1μs per task
33//! - **Memory efficiency**: Zero-copy task passing where possible
34//! - **Scalability**: Linear scaling up to CPU core count
35//! - **SIMD optimization**: 4-8x performance improvement for vectorizable workloads
36//! - **NUMA awareness**: Reduced memory latency on multi-socket systems
37//!
38//! ## Safety Guarantees
39//!
40//! - **Memory safety**: All operations are memory-safe by construction
41//! - **Data race freedom**: Rust's ownership system prevents data races
42//! - **Deadlock prevention**: Lock-free data structures where possible
43//! - **Resource cleanup**: Automatic resource cleanup on task completion
44//! - **Error handling**: Comprehensive error types with recovery mechanisms
45//!
46//! ## Quick Start Example
47//!
48//! ```rust
49//! use moirai::Moirai;
50//! use std::sync::atomic::{AtomicU32, Ordering};
51//! use std::sync::Arc;
52//!
53//! # fn example() -> Result<(), Box<dyn std::error::Error>> {
54//! // Create a new runtime with optimal configuration
55//! let runtime = Moirai::builder()
56//!     .worker_threads(4)
57//!     .build()?;
58//!
59//! // CPU-bound parallel computation
60//! let counter = Arc::new(AtomicU32::new(0));
61//! let counter_clone = counter.clone();
62//! let parallel_handle = runtime.spawn_fn(move || {
63//!     // Simulate intensive computation
64//!     for i in 0..1000 {
65//!         counter_clone.fetch_add(i % 100, Ordering::Relaxed);
66//!     }
67//!     counter_clone.load(Ordering::Relaxed)
68//! });
69//!
70//! // Another parallel task
71//! let critical_handle = runtime.spawn_fn(move || "critical task executed");
72//!
73//! // Tasks execute concurrently with optimal scheduling
74//! let parallel_result = parallel_handle.join().unwrap().unwrap();
75//! let critical_result = critical_handle.join().unwrap().unwrap();
76//!
77//! println!("Parallel result: {}", parallel_result);
78//! println!("Critical result: {}", critical_result);
79//!
80//! // Graceful shutdown with resource cleanup
81//! runtime.shutdown();
82//! # Ok(())
83//! # }
84//! ```
85//!
86//! ## Advanced Usage Patterns
87//!
88//! ### Task Chaining and Composition
89//!
90//! ```rust
91//! use moirai::Moirai;
92//!
93//! # fn chaining_example() -> Result<(), Box<dyn std::error::Error>> {
94//! let runtime = Moirai::new()?;
95//!
96//! // Chain tasks with dependencies using regular closures
97//! let handle1 = runtime.spawn_fn(|| 42);
98//! let result1 = handle1.join().unwrap().unwrap();
99//!
100//! let handle2 = runtime.spawn_fn(move || result1 * 2);
101//! let result2 = handle2.join().unwrap().unwrap();
102//!
103//! let handle3 = runtime.spawn_fn(move || result2 + 10);
104//! let result = handle3.join().unwrap().unwrap();
105//!
106//! assert_eq!(result, 94); // (42 * 2) + 10
107//! # Ok(())
108//! # }
109//! ```
110//!
111//! ### Distributed Boundary
112//!
113//! ```rust
114//! use moirai::Moirai;
115//!
116//! # fn boundary_example() -> Result<(), Box<dyn std::error::Error>> {
117//! let runtime = Moirai::builder()
118//!     .worker_threads(2)
119//!     .build()?;
120//!
121//! // Execute task locally through the verified scheduler facade.
122//! let handle = runtime.spawn_fn(move || "computed locally");
123//! let result = handle.join().unwrap().unwrap();
124//! println!("Result: {}", result);
125//!
126//! // Cross-machine execution uses fixed-format capability tokens; arbitrary
127//! // remote closure execution is intentionally outside the public facade.
128//! # Ok(())
129//! # }
130//! ```
131//!
132//! ## Migration Guide
133//!
134//! ### From `std::thread`
135//!
136//! ```rust
137//! # fn expensive_computation() -> i32 { 42 }
138//! # fn example() -> Result<(), Box<dyn std::error::Error>> {
139//! // Before: std::thread
140//! let handle = std::thread::spawn(|| {
141//!     expensive_computation()
142//! });
143//! let result = handle.join().unwrap();
144//!
145//! // After: Moirai
146//! let runtime = moirai::Moirai::new()?;
147//! let handle = runtime.spawn_fn(|| {
148//!     expensive_computation()
149//! });
150//! let result = handle.join().unwrap().unwrap();
151//! # Ok(())
152//! # }
153//! ```
154//!
155//! ### From Tokio
156//!
157//! ```rust
158//! # fn async_operation() -> String { "result".to_string() }
159//! # fn example() -> Result<(), Box<dyn std::error::Error>> {
160//! // Before: std::thread (since tokio requires async context)
161//! let handle = std::thread::spawn(|| {
162//!     async_operation()
163//! });
164//! let result = handle.join().unwrap();
165//!
166//! // After: Moirai
167//! let runtime = moirai::Moirai::new()?;
168//! let handle = runtime.spawn_fn(|| {
169//!     async_operation()
170//! });
171//! let result = handle.join().unwrap().unwrap();
172//! # Ok(())
173//! # }
174//! ```
175//!
176//! ### From Rayon
177//!
178//! ```rust
179//! # fn expensive_transform(x: &i32) -> i32 { x * 2 }
180//! # fn example() -> Result<(), Box<dyn std::error::Error>> {
181//! let data = vec![1, 2, 3, 4, 5];
182//!
183//! // Before: Sequential processing
184//! let result: Vec<_> = data.iter()
185//!     .map(|x| expensive_transform(x))
186//!     .collect();
187//!
188//! // After: Moirai parallel processing
189//! let runtime = moirai::Moirai::new()?;
190//! let handles: Vec<_> = data.iter()
191//!     .map(|&x| runtime.spawn_fn(move || expensive_transform(&x)))
192//!     .collect();
193//! let result: Result<Vec<_>, _> = handles.into_iter()
194//!     .map(|h| h.join().unwrap())
195//!     .collect();
196//! # Ok(())
197//! # }
198//! ```
199
200#![deny(missing_docs)]
201#![deny(unsafe_op_in_unsafe_fn)]
202
203// Re-export core functionality (avoiding ExecutorStats conflict)
204pub use moirai_core::{
205    Priority, Task, TaskContext, TaskHandle, TaskId,
206    error::*,
207    executor::{Executor, ExecutorConfig, ExecutorControl, TaskSpawner, WorkerPlacement},
208    scheduler::*,
209    task::*,
210};
211
212// Re-export executor functionality
213pub use moirai_executor::schedule::{
214    IdleHook, IdleHookRegistrationError, MAX_IDLE_HOOKS, register_idle_hook, run_idle_hooks,
215};
216pub use moirai_executor::{BlockingTask, HybridExecutor, SchedulerScope};
217
218/// Completion-only borrowing scope for jobs submitted to the unified scheduler.
219pub type MoiraiScope<'scope> = SchedulerScope<'scope, BlockingTask>;
220
221// Re-export transport functionality. The typed cross-boundary channel is the
222// rkyv-style archive pair (`ArchivedUniversalSender`/`ArchivedUniversalReceiver`);
223// the old non-functional `Universal*` placeholders were removed.
224/// Shared-memory same-machine IPC transport (Unix/Windows only).
225#[cfg(any(unix, windows))]
226pub use moirai_transport::IpcTransport;
227pub use moirai_transport::{
228    Address, ArchiveSerialize, ArchiveView, ArchivedMessage, ArchivedUniversalReceiver,
229    ArchivedUniversalSender, InMemoryTransport, RemoteAddress, TransportError, TransportManager,
230    TransportResult,
231};
232
233#[cfg(feature = "distributed")]
234mod routed;
235
236#[cfg(feature = "distributed")]
237pub use moirai_executor::schedule::{
238    AsyncTask, HybridRoutePolicy, HybridRouter, RoutePolicy, RouteTopology, SchedulerRoute,
239    ServerRoutePolicy, SyncTask, ThreadRoutePolicy, WorkClass,
240};
241
242#[cfg(feature = "distributed")]
243pub use moirai_transport::{
244    process::{ProcessDropPolicy, ProcessSpec, ProcessWaitPolicy},
245    remote_task::{
246        EchoBytesCapability, IntoRemoteOperation, RemoteCapability, RemoteCapabilityToken,
247        RemoteTaskId, RemoteTaskOperationKind, RemoteTaskOutput, RemoteTaskResult,
248        SumU64Capability,
249    },
250    route::{
251        ProcessEndpoint, RouteAddressBook, RouteNamespace, RouteResolution, RouteService,
252        RoutedProcessTaskError, RoutedProcessTaskOutput, ServerEndpoint,
253    },
254};
255
256#[cfg(feature = "distributed")]
257pub use routed::{FixedRemoteTask, RoutedProcessTarget, RoutedServerTarget};
258
259// Re-export channel functionality from core
260pub use moirai_core::channel;
261
262#[cfg(feature = "network")]
263pub use moirai_transport::TcpTransport;
264
265// Re-export synchronization primitives
266pub use moirai_sync::{AtomicCounter, Barrier, Condvar, Mutex, RwLock};
267
268// Re-export metrics functionality
269#[cfg(feature = "metrics")]
270pub use moirai_metrics::MetricsCollector;
271
272// Re-export async functionality (specific imports to avoid conflicts)
273#[cfg(feature = "async")]
274pub use moirai_async::{
275    File, FileOpenOptions, Timeout,
276    io::{
277        AsyncBufRead, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, MoiraiCompat, TokioCompat,
278    },
279    timer::{sleep, timeout},
280};
281
282#[cfg(all(feature = "async", any(unix, windows)))]
283pub use moirai_async::{TcpListener, TcpStream};
284
285#[cfg(all(feature = "async", not(target_arch = "wasm32")))]
286pub use moirai_async::executor::{AsyncExecutor, AsyncHandle};
287
288// Re-export iterator functionality
289#[cfg(feature = "iter")]
290pub use moirai_iter::{
291    AsyncContext, AsyncIterator, AsyncParallelIterator, ExecutionBase, ExecutionContext,
292    ExecutionStrategy, HybridConfig, HybridContext, IndexedParallelIterator, IntoAsyncIterator,
293    IntoParallelIterator, IntoParallelRefIterator, MoiraiIterator, ParallelContext, ParallelExtend,
294    ParallelIterator, PerformanceHistory, RangeParIter, VecParIter, VecRefParIter, async_range,
295    moirai_iter, moirai_iter_async, moirai_iter_hybrid, moirai_iter_parallel, par_range,
296};
297
298// Re-export GPU functionality
299#[cfg(feature = "gpu")]
300pub use moirai_gpu::prelude::*;
301
302// Synchronous data-parallel primitives (rayon-replacement surface), provided by
303// the `moirai-parallel` domain crate: monomorphized ExecutionPolicy + the
304// adaptive `par_*` helpers.
305#[cfg(feature = "parallel")]
306pub use moirai_parallel::*;
307
308#[cfg(all(feature = "parallel", feature = "melinoe"))]
309pub use moirai_parallel::melinoe_ext::*;
310
311// Submodules
312mod builder;
313mod global;
314/// Convenience functions for common operations.
315///
316/// Common imports for Moirai users.
317pub mod prelude;
318mod runtime;
319
320#[cfg(test)]
321mod tests;
322
323// Facade re-exports
324pub use builder::MoiraiBuilder;
325pub use global::{block_on, global, initialize, spawn_async, spawn_fn};
326pub use runtime::Moirai;