pub struct Moirai { /* private fields */ }Expand description
The main Moirai runtime that provides a unified interface for hybrid concurrency.
This is the primary entry point for using Moirai. It provides methods for spawning both async and parallel tasks, managing their execution, and coordinating between different execution models.
§Examples
use moirai::Moirai;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
// Create a new runtime
let runtime = Moirai::new()?;
// Spawn a parallel task
let counter = Arc::new(AtomicU32::new(0));
let counter_clone = counter.clone();
let handle = runtime.spawn_fn(move || {
for _ in 0..1000 {
counter_clone.fetch_add(1, Ordering::Relaxed);
}
counter_clone.load(Ordering::Relaxed)
});
// Spawn an async task
let async_handle = runtime.spawn_async(async {
// Simulate some async work
std::thread::sleep(std::time::Duration::from_millis(10));
"async task completed"
});
// The tasks will execute concurrently
println!("Tasks spawned, runtime is working...");
// Shutdown gracefully
runtime.shutdown();Implementations§
Source§impl Moirai
impl Moirai
Sourcepub fn new() -> ExecutorResult<Self>
pub fn new() -> ExecutorResult<Self>
Create a new Moirai runtime with default configuration.
§Errors
Returns an error if the runtime cannot be initialized.
Sourcepub fn builder() -> MoiraiBuilder
pub fn builder() -> MoiraiBuilder
Create a builder for configuring the Moirai runtime.
Sourcepub fn spawn<T>(&self, task: T) -> TaskHandle<T::Output>where
T: Task,
pub fn spawn<T>(&self, task: T) -> TaskHandle<T::Output>where
T: Task,
Spawn a task for parallel execution.
This is a convenience method for spawning CPU-bound tasks.
§Panics
Panics if the executor fails to spawn the task, which should not happen under normal circumstances unless the runtime is shutting down.
Sourcepub fn spawn_fn<F, R>(&self, func: F) -> TaskHandle<R>
pub fn spawn_fn<F, R>(&self, func: F) -> TaskHandle<R>
Spawn a parallel task using a closure.
The task will be executed on the work-stealing thread pool.
§Panics
Panics if the executor fails to spawn the blocking task, which should not happen under normal circumstances unless the runtime is shutting down.
Sourcepub fn spawn_async<F>(&self, future: F) -> TaskHandle<F::Output>
pub fn spawn_async<F>(&self, future: F) -> TaskHandle<F::Output>
Spawn an async task for execution.
The task will be executed on the async thread pool.
§Panics
Panics if the executor fails to spawn the async task, which should not happen under normal circumstances unless the runtime is shutting down.
Sourcepub fn spawn_blocking<F, R>(&self, func: F) -> TaskHandle<R>
pub fn spawn_blocking<F, R>(&self, func: F) -> TaskHandle<R>
Spawn a blocking task that may block the current thread.
Use this for I/O-bound or blocking operations.
§Panics
Panics if the executor fails to spawn the blocking task, which should not happen under normal circumstances unless the runtime is shutting down.
Sourcepub fn spawn_detached<F>(&self, func: F)
pub fn spawn_detached<F>(&self, func: F)
Spawn a fire-and-forget closure whose result is discarded.
This is the cheapest dispatch path: it returns no handle and skips the
per-task result-slot allocation that spawn_fn
performs, making it the right choice for background work whose output is
not needed (event handlers, logging, cache warming). The task is still
drained on shutdown.
§Panics
Panics if the executor fails to spawn the task, which should not happen unless the runtime is shutting down.
Sourcepub fn scope<'scope, F>(&'scope self, body: F) -> ExecutorResult<()>
pub fn scope<'scope, F>(&'scope self, body: F) -> ExecutorResult<()>
Run a completion-only scoped fan-out on the unified scheduler.
Use this when tasks only need to publish side effects through borrowed synchronization primitives and the caller must wait for all tasks before continuing. Scoped jobs may be coalesced and start after the scope body has finished registering work.
§Errors
Returns an executor error if the runtime is shutting down or if a scoped task panics.
Sourcepub fn for_each_indexed<'scope, F>(
&'scope self,
count: usize,
task: F,
) -> ExecutorResult<()>
pub fn for_each_indexed<'scope, F>( &'scope self, count: usize, task: F, ) -> ExecutorResult<()>
Run indexed work in worker-sized chunks on the unified scheduler.
Use this for CPU-bound data-parallel fan-out where the caller needs
completion, not one task handle per item. Work executes through the
compute-worker pool; potentially blocking work belongs on Self::scope.
The closure may borrow data that lives for the call because all chunks
complete before this method returns.
§Errors
Returns an executor error if the runtime is shutting down or if any chunk panics.
Sourcepub fn map_reduce_indexed<'scope, T, Map, Reduce>(
&'scope self,
count: usize,
identity: T,
map: Map,
reduce: Reduce,
) -> ExecutorResult<T>
pub fn map_reduce_indexed<'scope, T, Map, Reduce>( &'scope self, count: usize, identity: T, map: Map, reduce: Reduce, ) -> ExecutorResult<T>
Run indexed map/reduce in worker-sized chunks on the unified scheduler.
identity must be the neutral element for reduce. Use this for
CPU-bound indexed data-parallel reductions where per-item task handles
are not required. Work executes through the compute-worker pool.
§Errors
Returns an executor error if the runtime is shutting down or if any chunk panics.
Sourcepub fn spawn_with_priority<T>(
&self,
task: T,
priority: Priority,
) -> TaskHandle<T::Output>where
T: Task,
pub fn spawn_with_priority<T>(
&self,
task: T,
priority: Priority,
) -> TaskHandle<T::Output>where
T: Task,
Spawn a task with a specific priority.
Higher priority tasks will be executed before lower priority tasks.
§Panics
Panics if the executor fails to spawn the task with priority, which should not happen under normal circumstances unless the runtime is shutting down.
Sourcepub fn spawn_fn_with_priority<F, R>(
&self,
f: F,
priority: Priority,
) -> TaskHandle<R>
pub fn spawn_fn_with_priority<F, R>( &self, f: F, priority: Priority, ) -> TaskHandle<R>
Spawn a closure with priority as a task (convenience method).
Sourcepub fn block_on<F>(&self, future: F) -> F::Outputwhere
F: Future,
pub fn block_on<F>(&self, future: F) -> F::Outputwhere
F: Future,
Block the current thread until the future completes.
This is useful for running async code from synchronous contexts.
Sourcepub fn try_run(&self) -> bool
pub fn try_run(&self) -> bool
Try to run pending tasks without blocking.
Returns true if any tasks were executed, false if no work was available.
Sourcepub fn join(&self) -> ExecutorResult<()>
pub fn join(&self) -> ExecutorResult<()>
Wait until queued and active runtime work completes without shutting down workers.
Use this as a non-destructive process-fusion barrier when producers have finished submitting a batch and the runtime should process all available work before the caller continues. New tasks submitted after this method observes quiescence belong to a later batch.
§Errors
Returns an executor error if the scheduler join operation fails.
Sourcepub fn shutdown(&self)
pub fn shutdown(&self)
Shutdown the runtime gracefully.
This will wait for all currently running tasks to complete before shutting down the thread pools.
Sourcepub fn shutdown_timeout(&self, timeout: Duration)
pub fn shutdown_timeout(&self, timeout: Duration)
Shutdown the runtime with a timeout.
If tasks don’t complete within the timeout, they will be forcefully terminated.
Sourcepub fn is_shutting_down(&self) -> bool
pub fn is_shutting_down(&self) -> bool
Check if the runtime is shutting down.
Sourcepub fn worker_count(&self) -> usize
pub fn worker_count(&self) -> usize
Get the number of worker threads.
Sourcepub fn channel<T: Send + 'static>(&self) -> (MpmcSender<T>, MpmcReceiver<T>)
pub fn channel<T: Send + 'static>(&self) -> (MpmcSender<T>, MpmcReceiver<T>)
Create a channel for communication, bounded at
DEFAULT_CHANNEL_CAPACITY.
Bounded is the default because an unbounded queue converts a slow
consumer into unbounded memory growth: a full channel blocks its
producer (or returns ChannelError::Full from try_send) instead of
allocating. Use Self::bounded_channel when the right capacity is
known; the unbounded queue remains available as
moirai_core::channel::unbounded, whose documentation states the cost.
Sourcepub fn bounded_channel<T: Send + 'static>(
&self,
capacity: usize,
) -> (MpmcSender<T>, MpmcReceiver<T>)
pub fn bounded_channel<T: Send + 'static>( &self, capacity: usize, ) -> (MpmcSender<T>, MpmcReceiver<T>)
Create a bounded channel with an explicit capacity.