Skip to main content

Moirai

Struct Moirai 

Source
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

Source

pub fn new() -> ExecutorResult<Self>

Create a new Moirai runtime with default configuration.

§Errors

Returns an error if the runtime cannot be initialized.

Source

pub fn builder() -> MoiraiBuilder

Create a builder for configuring the Moirai runtime.

Source

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.

Source

pub fn spawn_fn<F, R>(&self, func: F) -> TaskHandle<R>
where F: FnOnce() -> R + Send + 'static, R: Send + 'static,

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.

Source

pub fn spawn_async<F>(&self, future: F) -> TaskHandle<F::Output>
where F: Future + Send + 'static, F::Output: Send + 'static,

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.

Source

pub fn spawn_blocking<F, R>(&self, func: F) -> TaskHandle<R>
where F: FnOnce() -> R + Send + 'static, R: Send + 'static,

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.

Source

pub fn spawn_detached<F>(&self, func: F)
where F: FnOnce() + Send + 'static,

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.

Source

pub fn scope<'scope, F>(&'scope self, body: F) -> ExecutorResult<()>
where F: FnOnce(&MoiraiScope<'scope>) -> 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.

Source

pub fn for_each_indexed<'scope, F>( &'scope self, count: usize, task: F, ) -> ExecutorResult<()>
where F: Fn(usize) + Send + Sync + 'scope,

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.

Source

pub fn map_reduce_indexed<'scope, T, Map, Reduce>( &'scope self, count: usize, identity: T, map: Map, reduce: Reduce, ) -> ExecutorResult<T>
where T: Send + Clone + 'scope, Map: Fn(usize) -> T + Send + Sync + 'scope, Reduce: Fn(T, T) -> T + Send + Sync + 'scope,

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.

Source

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.

Source

pub fn spawn_fn_with_priority<F, R>( &self, f: F, priority: Priority, ) -> TaskHandle<R>
where F: FnOnce() -> R + Send + 'static, R: Send + 'static,

Spawn a closure with priority as a task (convenience method).

Source

pub fn block_on<F>(&self, future: F) -> F::Output
where F: Future,

Block the current thread until the future completes.

This is useful for running async code from synchronous contexts.

Source

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.

Source

pub fn has_work(&self) -> bool

Returns true when queued or active runtime work exists.

Source

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.

Source

pub fn shutdown(&self)

Shutdown the runtime gracefully.

This will wait for all currently running tasks to complete before shutting down the thread pools.

Source

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.

Source

pub fn is_shutting_down(&self) -> bool

Check if the runtime is shutting down.

Source

pub fn worker_count(&self) -> usize

Get the number of worker threads.

Source

pub fn load(&self) -> usize

Get the current load (number of pending tasks).

Source

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.

Source

pub fn bounded_channel<T: Send + 'static>( &self, capacity: usize, ) -> (MpmcSender<T>, MpmcReceiver<T>)

Create a bounded channel with an explicit capacity.

Trait Implementations§

Source§

impl Clone for Moirai

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Default for Moirai

Source§

fn default() -> Self

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.