Expand description
§Moirai - Weaving the Threads of Fate
Moirai is a high-performance hybrid concurrency library for Rust that seamlessly blends asynchronous and parallel execution models. Named after the Greek Fates who controlled the threads of life, Moirai weaves together the best principles from async task scheduling and parallel work-stealing into a unified framework.
§Core Design Principles
Moirai follows elite programming practices:
- SOLID: Single responsibility, open/closed, Liskov substitution, interface segregation, dependency inversion
- CUPID: Composable, Unix philosophy, predictable, idiomatic, domain-centric
- GRASP: Information expert, creator, controller, low coupling, high cohesion
- ACID: Atomicity, consistency, isolation, durability in task execution
§Features
- Zero-cost abstractions: All abstractions compile away to optimal code
- Hybrid execution: Seamlessly mix async and parallel tasks
- Work-stealing scheduler: Intelligent load balancing across CPU cores
- Memory safety: Leverage Rust’s ownership system for safe concurrency
- High performance: Sub-microsecond task scheduling overhead
- NUMA awareness: Optimize for modern multi-socket systems
- Rich iterator combinators: Parallel and async iterator processing
- IPC: Inter-process communication (optional)
- Metrics: Performance monitoring (optional)
- Distributed transport feature gates: Optional transport and iterator helpers without a facade-level remote-closure API
§Performance Characteristics
- Task scheduling overhead: < 1μs per task
- Memory efficiency: Zero-copy task passing where possible
- Scalability: Linear scaling up to CPU core count
- SIMD optimization: 4-8x performance improvement for vectorizable workloads
- NUMA awareness: Reduced memory latency on multi-socket systems
§Safety Guarantees
- Memory safety: All operations are memory-safe by construction
- Data race freedom: Rust’s ownership system prevents data races
- Deadlock prevention: Lock-free data structures where possible
- Resource cleanup: Automatic resource cleanup on task completion
- Error handling: Comprehensive error types with recovery mechanisms
§Quick Start Example
use moirai::Moirai;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
// Create a new runtime with optimal configuration
let runtime = Moirai::builder()
.worker_threads(4)
.build()?;
// CPU-bound parallel computation
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = counter.clone();
let parallel_handle = runtime.spawn_fn(move || {
// Simulate intensive computation
for i in 0..1000 {
counter_clone.fetch_add(i % 100, Ordering::Relaxed);
}
counter_clone.load(Ordering::Relaxed)
});
// Another parallel task
let critical_handle = runtime.spawn_fn(move || "critical task executed");
// Tasks execute concurrently with optimal scheduling
let parallel_result = parallel_handle.join().unwrap().unwrap();
let critical_result = critical_handle.join().unwrap().unwrap();
println!("Parallel result: {}", parallel_result);
println!("Critical result: {}", critical_result);
// Graceful shutdown with resource cleanup
runtime.shutdown();§Advanced Usage Patterns
§Task Chaining and Composition
use moirai::Moirai;
let runtime = Moirai::new()?;
// Chain tasks with dependencies using regular closures
let handle1 = runtime.spawn_fn(|| 42);
let result1 = handle1.join().unwrap().unwrap();
let handle2 = runtime.spawn_fn(move || result1 * 2);
let result2 = handle2.join().unwrap().unwrap();
let handle3 = runtime.spawn_fn(move || result2 + 10);
let result = handle3.join().unwrap().unwrap();
assert_eq!(result, 94); // (42 * 2) + 10§Distributed Boundary
use moirai::Moirai;
let runtime = Moirai::builder()
.worker_threads(2)
.build()?;
// Execute task locally through the verified scheduler facade.
let handle = runtime.spawn_fn(move || "computed locally");
let result = handle.join().unwrap().unwrap();
println!("Result: {}", result);
// Cross-machine execution uses fixed-format capability tokens; arbitrary
// remote closure execution is intentionally outside the public facade.§Migration Guide
§From std::thread
// Before: std::thread
let handle = std::thread::spawn(|| {
expensive_computation()
});
let result = handle.join().unwrap();
// After: Moirai
let runtime = moirai::Moirai::new()?;
let handle = runtime.spawn_fn(|| {
expensive_computation()
});
let result = handle.join().unwrap().unwrap();§From Tokio
// Before: std::thread (since tokio requires async context)
let handle = std::thread::spawn(|| {
async_operation()
});
let result = handle.join().unwrap();
// After: Moirai
let runtime = moirai::Moirai::new()?;
let handle = runtime.spawn_fn(|| {
async_operation()
});
let result = handle.join().unwrap().unwrap();§From Rayon
let data = vec![1, 2, 3, 4, 5];
// Before: Sequential processing
let result: Vec<_> = data.iter()
.map(|x| expensive_transform(x))
.collect();
// After: Moirai parallel processing
let runtime = moirai::Moirai::new()?;
let handles: Vec<_> = data.iter()
.map(|&x| runtime.spawn_fn(move || expensive_transform(&x)))
.collect();
let result: Result<Vec<_>, _> = handles.into_iter()
.map(|h| h.join().unwrap())
.collect();Modules§
- channel
- Unified high-performance channel implementations for Moirai.
- melinoe_
ext - Parallel partitioning drivers for branded Melinoe cell slices.
- ops
- Synchronous data-parallel operators and free functions. Synchronous data-parallel operators over the unified scheduler.
- prelude
- Convenience functions for common operations.
- timeout
- Deadline wrappers over futures.
Structs§
- Adaptive
- Run in parallel only for inputs at or above
ADAPTIVE_PARALLEL_THRESHOLD. - Adaptive
With Threshold - Run in parallel only for inputs at or above the custom threshold
N. - Archived
Message - Transport message bytes plus a typed archive view contract.
- Archived
Universal Receiver - Universal receiver for zero-copy archive views.
- Archived
Universal Sender - Universal sender for archive-serializable values.
- Async
Context - Async execution context for I/O-bound work
- Async
Executor - Native async executor with access to the PAL I/O reactor.
- Async
Handle - A handle to an async task that can be awaited.
- Atomic
Counter - A thread-safe atomic counter with increment and decrement operations.
- Barrier
- A barrier enables multiple threads to synchronize the beginning of some computation.
- Base
Task - Base implementation for common task patterns to reduce redundancy.
- Blocking
Result Wait - Zero-sized blocking wait policy: spins up to
MAX_SPIN_ATTEMPTSthen parks. - Blocking
Task - Potentially blocking task marker.
- Chained
- A task that chains two operations together.
- Closure
- A simple closure-based task implementation.
- Condvar
- A Condition Variable
- Contextual
Task - A task with an explicit context.
- Executor
Config - Configuration settings for executor behavior and performance characteristics.
- File
- Async file handle supporting stateful streams and positioned reads.
- File
Open Options - Declarative open-mode configuration for
File::open_with. - Group
- A collection of related tasks that can be executed as a group.
- Hybrid
Config - Configuration for hybrid execution strategy
- Hybrid
Context - Hybrid context that adapts between parallel and async execution
- Hybrid
Executor - Main hybrid executor that coordinates sync, async, and blocking tasks.
- InMemory
Transport - In-memory transport for local communication
- IpcTransport
- Shared-memory same-machine IPC transport (Unix/Windows only).
Shared-memory same-machine IPC transport (Unix/Windows only).
Shared-memory IPC transport. Holds one
SharedQueuehandle per segment name, created lazily on first use. - Mapped
- A task that maps the output of another task.
- Moirai
- The main Moirai runtime that provides a unified interface for hybrid concurrency.
- Moirai
Builder - Builder for configuring the Moirai runtime.
- Moirai
Compat - Exposes a
tokio::ioreader, writer, or buffered reader through the Moirai I/O traits. - Moirai
Iterator - Main iterator type that adapts to different execution contexts.
- Mutex
- A mutual exclusion primitive useful for protecting shared data
- ParMut
- A mutable parallel view of a slice bound to execution policy
P. - ParRef
- A read-only parallel view of a slice bound to execution policy
P. - Parallel
- Always run in parallel on the shared work-stealing pool.
- Parallel
Context - Parallel execution context for CPU-bound work
- Parameterized
- A task that accepts parameters for customized execution.
- Performance
History - Performance history for adaptive execution decisions
- Range
ParIter - Range parallel iterator.
- Remote
Address - Remote address for cross-machine communication
- RwLock
- A reader-writer lock
- Scheduler
Id - A unique identifier for a scheduler instance.
- Scheduler
Scope - Borrowing scope for scheduler jobs that must complete before the scope exits.
- Scope
- Borrowing scope for spawning parallel sub-tasks that may capture non-
'staticreferences. - Sequential
- Always run sequentially (single thread, no scheduling).
- Spawner
- A task that can spawn other tasks during its execution.
- Task
Builder - Builder for creating and configuring tasks.
- Task
Context - Task execution context and metadata.
- Task
Future - A future adapter that executes a
Taskon first poll. - Task
Handle - A handle to a task that may be running on another thread.
- TaskId
- A unique identifier for tasks in the Moirai runtime.
- Task
Result Sender - Single-producer completion endpoint for a task result.
- TcpListener
- Native async TCP listener with connection management
- TcpStream
- Native async TCP stream with statistics tracking
- Timeout
- Timeout wrapper for futures with comprehensive cancellation
- Tokio
Compat - Exposes a Moirai reader, writer, or buffered reader through the
tokio::iotraits. - Transport
Manager - Transport manager that routes messages to appropriate transport
- VecPar
Iter - Parallel iterator over a vector.
- VecRef
ParIter - Parallel iterator over vector references.
- Work
Bytes - Run in parallel only for an operation that moves at least
Nbytes.
Enums§
- Address
- Address for identifying communication endpoints
- Chunk
Buffers Error - Failure to partition a fixed set of mutable buffers into matching chunks.
- Execution
Context - Concrete execution context enum that wraps different strategy implementations This approach ensures type safety while avoiding dyn-compatibility issues
- Execution
Strategy - Execution strategy selected by the hybrid context.
- Executor
Error - Errors that can occur during executor operations.
- Idle
Hook Registration Error - Failure returned when a worker idle hook cannot be admitted.
- Placement
Failure - Why a worker could not be confined to a logical processor.
- Priority
- Priority levels for task scheduling.
- Scheduler
Error - Errors that can occur during scheduler operations.
- Task
Error - Errors that can occur during task operations.
- Task
Error Kind - Specific kinds of task execution errors.
- Transport
Error - Error types for channel operations
- Worker
Placement - Whether scheduler workers are confined to logical processors.
Constants§
- ADAPTIVE_
PARALLEL_ THRESHOLD - Element count at or above which
Adaptivechooses parallel execution. - MAX_
IDLE_ HOOKS - Maximum number of process-wide worker idle hooks.
- UNIT_
TASK_ BYTES - Bytes of work one scheduled task carries.
Traits§
- Archive
Serialize - Writes a value into transport-owned archive bytes.
- Archive
View - Validates archive bytes and returns a typed borrowed view.
- Async
BufRead - Read bytes from a buffered source asynchronously.
- Async
Iterator - Core async iterator trait for async/await compatible iteration
- Async
Parallel Iterator - Parallel async iterator for CPU+async hybrid workloads
- Async
Read - Read bytes asynchronously.
- Async
Read Ext - Extension methods for types implementing
AsyncRead. - Async
Write - Write bytes asynchronously.
- Async
Write Ext - Extension methods for types implementing
AsyncWrite. - Execution
Base - Base trait for all execution contexts
- Execution
Policy - Compile-time strategy selector for the data-parallel operations in this crate.
- Executor
- Combined executor trait with all capabilities.
- Executor
Control - Provides control operations for executor lifecycle management.
- Indexed
Parallel Iterator - Exact-size boundary for Moirai’s bounded Rayon-style indexed source subset.
- Into
Async Iterator - Trait for converting types into async iterators
- Into
Parallel Iterator - Extension trait for collections to create parallel iterators.
- Into
Parallel RefIterator - Extension trait for collection references to create parallel iterators.
- Parallel
Extend - Trait for collections that can be extended in parallel.
- Parallel
Iterator - Core parallel iterator trait for Moirai’s Rayon-style non-indexed subset.
- Parallel
Slice - Extension trait providing an adaptive parallel view over
&[T]. - Parallel
Slice Mut - Extension trait providing an adaptive mutable parallel view over
&mut [T]. - Result
Wait Policy - Compile-time wait policy for task result handoff.
- Task
- The core trait for executable tasks in the Moirai runtime.
- TaskExt
- Extension methods for tasks.
- Task
Spawner - Core task spawning capabilities.
Functions§
- async_
range - Async range iterator.
- block_
on - Block on a future using the global runtime.
- enumerate_
mut_ with - Apply
f(index, &mut element)to every element ofdatain place, scheduled by policyP. - enumerate_
with - Apply
f(index, &element)to every element ofdata, scheduled by policyP. - fold_
reduce_ with - Parallel fold-reduce over the index domain
0..len, scheduled by policyP. - for_
each_ chunk_ buffers_ mut_ enumerated_ with - Apply
f(index, chunks)to matching chunks from a fixed set of distinct mutable buffers, scheduled by policyP. - for_
each_ chunk_ mut_ enumerated_ with - Like
for_each_chunk_mut_withbut also passes the zero-based chunk index tof(synchronous equivalent ofdata.par_chunks_mut(chunk_size).enumerate().for_each(f)). - for_
each_ chunk_ mut_ with - Apply
fto each consecutivechunk_size-element mutable chunk ofdatain parallel, scheduled by policyP. The final chunk may be shorter. - for_
each_ chunk_ mut_ with_ state - Apply
f(state, chunk)to each consecutive mutable chunk, creating one reusable state value per scheduled worker shard. - for_
each_ chunk_ pair_ mut_ enumerated_ with - Apply
f(index, a_chunk, b_chunk)to pairedchunk_size-element mutable chunks of two distinct buffers in parallel, scheduled by policyP. - for_
each_ chunk_ quad_ mut_ enumerated_ with - Apply
f(index, a_chunk, b_chunk, c_chunk, d_chunk)to four distinct mutable buffers chunked identically, scheduled by policyP. - for_
each_ chunk_ triple_ mut_ enumerated_ with - Apply
f(index, a_chunk, b_chunk, c_chunk)to three distinct mutable buffers chunked identically, scheduled by policyP. - for_
each_ index_ with - Apply
fto every index in0..lenin parallel, scheduled by policyP. - for_
each_ mut_ with - Apply
fto every element ofdatain place, scheduled by policyP. - for_
each_ unit_ task_ many_ mut_ with crate::for_each_unit_task_mut_withoverKbuffers of one type, split in lockstep.- for_
each_ unit_ task_ mut_ with - Apply
f(state, first_unit, units)to consecutive runs of whole units ofdata, each run sized to aboutcrate::UNIT_TASK_BYTESof work. - for_
each_ unit_ task_ pair_ mut_ with - Apply
f(state, first_unit, a_units, b_units)to aligned runs of whole units of two mutable buffers, each run sized to aboutcrate::UNIT_TASK_BYTESof work. - for_
each_ unit_ task_ range_ with - Apply
f(state, first_unit, units)to consecutive runs of whole units of a pass that addresses its own data, each run sized to aboutcrate::UNIT_TASK_BYTESof work. - for_
each_ unit_ task_ triple_ mut_ with - Apply
f(state, first_unit, a_units, b_units, c_units)to aligned runs of whole units of three mutable buffers, each run sized to aboutcrate::UNIT_TASK_BYTESof work. - for_
each_ with - Apply
fto every element ofdata, scheduled by policyP. - global
- Get or initialize the global Moirai runtime.
- initialize
- Initialize the global runtime and its Melinoe partition bridge.
- join
- Adaptive Rayon-style two-closure join.
- join_
with - Run two closures to completion and return both results.
- map_
collect_ index_ with - Parallel map over the index domain
0..len, collecting into aVec<R>in order, scheduled by policyP. - map_
collect_ mut_ with - Map each element of
datain place withf(index, &mut element), collecting each returned value into aVec<R>in order, scheduled by policyP. - map_
collect_ with - Map each element of
datawithf, collecting into aVec<R>in order, scheduled by policyP. - map_
reduce_ with - Map-reduce over
data, scheduled by policyP. - moirai_
iter - Convenience function to create a Moirai iterator.
- moirai_
iter_ async - Create an async iterator.
- moirai_
iter_ hybrid - Create a hybrid iterator.
- moirai_
iter_ parallel - Create a parallel iterator.
- par_
partition_ for_ each - Split
cellsinto disjoint shards ofchunk_sizeand runfon each in parallel. - par_
partition_ for_ each_ with_ policy - Split
cellsinto disjoint shards ofchunk_sizeand runfon each, parallelizing only whenPpermits it for this region size. - par_
partition_ map - Split
cellsinto disjoint shards ofchunk_size, runfon each in parallel, and collect the per-shard results into aVec<R>in partition order. - par_
partition_ map_ with_ policy - Split
cellsinto disjoint shards ofchunk_size, runfon each, and collect the per-shard results in partition order — parallelizing only whenPpermits it for this region size. - par_
range - Parallel range iterator for Moirai’s Rayon-style non-indexed subset.
- reduce_
index_ with - Parallel reduction over the index domain
0..len, scheduled by policyP. - register_
idle_ hook - Registers
hookto run on every worker thread before it parks for work. - run_
idle_ hooks - Runs the hooks registered for the calling worker.
- scope
- Create a borrowing scope for parallel sub-tasks.
- sleep
- Create a delay future that completes after the specified duration
- spawn_
async - Spawn an async task on the global runtime.
- spawn_
fn - Spawn a parallel task on the global runtime.
- timeout
- Timeout wrapper for futures with comprehensive cancellation
- units_
per_ task - Units one task carries when each unit moves
unit_bytes, never fewer than one: a unit at or aboveUNIT_TASK_BYTESis a task on its own.
Type Aliases§
- Executor
Result - A result type for executor operations.
- Idle
Hook - A worker idle hook: a plain function pointer called on the worker thread immediately before it parks for work.
- Moirai
Scope - Completion-only borrowing scope for jobs submitted to the unified scheduler.
- Scheduler
Result - A result type for scheduler operations.
- Task
Result - A result type for task operations.
- Transport
Result - Result type for transport operations