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

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
High-performance async file handle with native implementation
FileOpenOptions
Configuration for file operations
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 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
Wrapper providing Moirai’s native I/O traits compatibility for Tokio types.
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
Wrapper providing Tokio’s I/O traits compatibility.
TransportManager
Transport manager that routes messages to appropriate transport
VecParIter
Parallel iterator over a vector.
VecRefParIter
Parallel iterator over vector references.

Enums§

Address
Address for identifying communication endpoints
ExecutionContext
Concrete execution context enum that wraps different strategy implementations This approach ensures type safety while avoiding dyn-compatibility issues
ExecutionStrategy
ExecutorError
Errors that can occur during executor operations.
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

Constants§

ADAPTIVE_PARALLEL_THRESHOLD
Element count at or above which Adaptive chooses parallel execution.

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_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_with
Apply f to every element of data, scheduled by policy P.
global
Get or initialize the global Moirai runtime.
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_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_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.
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

Type Aliases§

ExecutorResult
A result type for executor operations.
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