Skip to main content

Crate moirai

Crate moirai 

Source
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.
AdaptiveWithThreshold
Run in parallel only for inputs at or above the custom threshold N.
ArchivedMessage
Transport message bytes plus a typed archive view contract.
ArchivedUniversalReceiver
Universal receiver for zero-copy archive views.
ArchivedUniversalSender
Universal sender for archive-serializable values.
AsyncContext
Async execution context for I/O-bound work
AsyncExecutor
Native async executor with access to the PAL I/O reactor.
AsyncHandle
A handle to an async task that can be awaited.
AtomicCounter
A thread-safe atomic counter with increment and decrement operations.
Barrier
A barrier enables multiple threads to synchronize the beginning of some computation.
BaseTask
Base implementation for common task patterns to reduce redundancy.
BlockingResultWait
Zero-sized blocking wait policy: spins up to MAX_SPIN_ATTEMPTS then parks.
BlockingTask
Potentially blocking task marker.
Chained
A task that chains two operations together.
Closure
A simple closure-based task implementation.
Condvar
A Condition Variable
ContextualTask
A task with an explicit context.
ExecutorConfig
Configuration settings for executor behavior and performance characteristics.
File
Async file handle supporting stateful streams and positioned reads.
FileOpenOptions
Declarative open-mode configuration for File::open_with.
Group
A collection of related tasks that can be executed as a group.
HybridConfig
Configuration for hybrid execution strategy
HybridContext
Hybrid context that adapts between parallel and async execution
HybridExecutor
Main hybrid executor that coordinates sync, async, and blocking tasks.
InMemoryTransport
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 SharedQueue handle 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.
MoiraiBuilder
Builder for configuring the Moirai runtime.
MoiraiCompat
Exposes a tokio::io reader, writer, or buffered reader through the Moirai I/O traits.
MoiraiIterator
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.
ParallelContext
Parallel execution context for CPU-bound work
Parameterized
A task that accepts parameters for customized execution.
PerformanceHistory
Performance history for adaptive execution decisions
RangeParIter
Range parallel iterator.
RemoteAddress
Remote address for cross-machine communication
RwLock
A reader-writer lock
SchedulerId
A unique identifier for a scheduler instance.
SchedulerScope
Borrowing scope for scheduler jobs that must complete before the scope exits.
Scope
Borrowing scope for spawning parallel sub-tasks that may capture non-'static references.
Sequential
Always run sequentially (single thread, no scheduling).
Spawner
A task that can spawn other tasks during its execution.
TaskBuilder
Builder for creating and configuring tasks.
TaskContext
Task execution context and metadata.
TaskFuture
A future adapter that executes a Task on first poll.
TaskHandle
A handle to a task that may be running on another thread.
TaskId
A unique identifier for tasks in the Moirai runtime.
TaskResultSender
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
TokioCompat
Exposes a Moirai reader, writer, or buffered reader through the tokio::io traits.
TransportManager
Transport manager that routes messages to appropriate transport
VecParIter
Parallel iterator over a vector.
VecRefParIter
Parallel iterator over vector references.
WorkBytes
Run in parallel only for an operation that moves at least N bytes.

Enums§

Address
Address for identifying communication endpoints
ChunkBuffersError
Failure to partition a fixed set of mutable buffers into matching chunks.
ExecutionContext
Concrete execution context enum that wraps different strategy implementations This approach ensures type safety while avoiding dyn-compatibility issues
ExecutionStrategy
Execution strategy selected by the hybrid context.
ExecutorError
Errors that can occur during executor operations.
IdleHookRegistrationError
Failure returned when a worker idle hook cannot be admitted.
PlacementFailure
Why a worker could not be confined to a logical processor.
Priority
Priority levels for task scheduling.
SchedulerError
Errors that can occur during scheduler operations.
TaskError
Errors that can occur during task operations.
TaskErrorKind
Specific kinds of task execution errors.
TransportError
Error types for channel operations
WorkerPlacement
Whether scheduler workers are confined to logical processors.

Constants§

ADAPTIVE_PARALLEL_THRESHOLD
Element count at or above which Adaptive chooses parallel execution.
MAX_IDLE_HOOKS
Maximum number of process-wide worker idle hooks.
UNIT_TASK_BYTES
Bytes of work one scheduled task carries.

Traits§

ArchiveSerialize
Writes a value into transport-owned archive bytes.
ArchiveView
Validates archive bytes and returns a typed borrowed view.
AsyncBufRead
Read bytes from a buffered source asynchronously.
AsyncIterator
Core async iterator trait for async/await compatible iteration
AsyncParallelIterator
Parallel async iterator for CPU+async hybrid workloads
AsyncRead
Read bytes asynchronously.
AsyncReadExt
Extension methods for types implementing AsyncRead.
AsyncWrite
Write bytes asynchronously.
AsyncWriteExt
Extension methods for types implementing AsyncWrite.
ExecutionBase
Base trait for all execution contexts
ExecutionPolicy
Compile-time strategy selector for the data-parallel operations in this crate.
Executor
Combined executor trait with all capabilities.
ExecutorControl
Provides control operations for executor lifecycle management.
IndexedParallelIterator
Exact-size boundary for Moirai’s bounded Rayon-style indexed source subset.
IntoAsyncIterator
Trait for converting types into async iterators
IntoParallelIterator
Extension trait for collections to create parallel iterators.
IntoParallelRefIterator
Extension trait for collection references to create parallel iterators.
ParallelExtend
Trait for collections that can be extended in parallel.
ParallelIterator
Core parallel iterator trait for Moirai’s Rayon-style non-indexed subset.
ParallelSlice
Extension trait providing an adaptive parallel view over &[T].
ParallelSliceMut
Extension trait providing an adaptive mutable parallel view over &mut [T].
ResultWaitPolicy
Compile-time wait policy for task result handoff.
Task
The core trait for executable tasks in the Moirai runtime.
TaskExt
Extension methods for tasks.
TaskSpawner
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 of data in place, scheduled by policy P.
enumerate_with
Apply f(index, &element) to every element of data, scheduled by policy P.
fold_reduce_with
Parallel fold-reduce over the index domain 0..len, scheduled by policy P.
for_each_chunk_buffers_mut_enumerated_with
Apply f(index, chunks) to matching chunks from a fixed set of distinct mutable buffers, scheduled by policy P.
for_each_chunk_mut_enumerated_with
Like for_each_chunk_mut_with but also passes the zero-based chunk index to f (synchronous equivalent of data.par_chunks_mut(chunk_size).enumerate().for_each(f)).
for_each_chunk_mut_with
Apply f to each consecutive chunk_size-element mutable chunk of data in parallel, scheduled by policy P. 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 paired chunk_size-element mutable chunks of two distinct buffers in parallel, scheduled by policy P.
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 policy P.
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 policy P.
for_each_index_with
Apply f to every index in 0..len in parallel, scheduled by policy P.
for_each_mut_with
Apply f to every element of data in place, scheduled by policy P.
for_each_unit_task_many_mut_with
crate::for_each_unit_task_mut_with over K buffers of one type, split in lockstep.
for_each_unit_task_mut_with
Apply f(state, first_unit, units) to consecutive runs of whole units of data, each run sized to about crate::UNIT_TASK_BYTES of 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 about crate::UNIT_TASK_BYTES of 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 about crate::UNIT_TASK_BYTES of 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 about crate::UNIT_TASK_BYTES of work.
for_each_with
Apply f to every element of data, scheduled by policy P.
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 a Vec<R> in order, scheduled by policy P.
map_collect_mut_with
Map each element of data in place with f(index, &mut element), collecting each returned value into a Vec<R> in order, scheduled by policy P.
map_collect_with
Map each element of data with f, collecting into a Vec<R> in order, scheduled by policy P.
map_reduce_with
Map-reduce over data, scheduled by policy P.
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 cells into disjoint shards of chunk_size and run f on each in parallel.
par_partition_for_each_with_policy
Split cells into disjoint shards of chunk_size and run f on each, parallelizing only when P permits it for this region size.
par_partition_map
Split cells into disjoint shards of chunk_size, run f on each in parallel, and collect the per-shard results into a Vec<R> in partition order.
par_partition_map_with_policy
Split cells into disjoint shards of chunk_size, run f on each, and collect the per-shard results in partition order — parallelizing only when P permits 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 policy P.
register_idle_hook
Registers hook to 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 above UNIT_TASK_BYTES is a task on its own.

Type Aliases§

ExecutorResult
A result type for executor operations.
IdleHook
A worker idle hook: a plain function pointer called on the worker thread immediately before it parks for work.
MoiraiScope
Completion-only borrowing scope for jobs submitted to the unified scheduler.
SchedulerResult
A result type for scheduler operations.
TaskResult
A result type for task operations.
TransportResult
Result type for transport operations