Skip to main content

moirai/
runtime.rs

1use crate::{MoiraiBuilder, MoiraiScope};
2#[cfg(feature = "metrics")]
3use moirai_core::executor::Executor;
4use moirai_core::{
5    Priority, Task, TaskBuilder, TaskHandle,
6    error::*,
7    executor::{ExecutorControl, TaskSpawner},
8};
9use moirai_executor::{BlockingTask, HybridExecutor, SyncTask};
10use std::{future::Future, sync::Arc, time::Duration};
11
12/// The main Moirai runtime that provides a unified interface for hybrid concurrency.
13///
14/// This is the primary entry point for using Moirai. It provides methods for spawning
15/// both async and parallel tasks, managing their execution, and coordinating between
16/// different execution models.
17///
18/// # Examples
19///
20/// ```
21/// use moirai::Moirai;
22/// use std::sync::atomic::{AtomicU32, Ordering};
23/// use std::sync::Arc;
24///
25/// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
26/// // Create a new runtime
27/// let runtime = Moirai::new()?;
28///
29/// // Spawn a parallel task
30/// let counter = Arc::new(AtomicU32::new(0));
31/// let counter_clone = counter.clone();
32/// let handle = runtime.spawn_fn(move || {
33///     for _ in 0..1000 {
34///         counter_clone.fetch_add(1, Ordering::Relaxed);
35///     }
36///     counter_clone.load(Ordering::Relaxed)
37/// });
38///
39/// // Spawn an async task
40/// let async_handle = runtime.spawn_async(async {
41///     // Simulate some async work
42///     std::thread::sleep(std::time::Duration::from_millis(10));
43///     "async task completed"
44/// });
45///
46/// // The tasks will execute concurrently
47/// println!("Tasks spawned, runtime is working...");
48///
49/// // Shutdown gracefully
50/// runtime.shutdown();
51/// # Ok(())
52/// # }
53/// ```
54#[derive(Clone)]
55pub struct Moirai {
56    pub(crate) executor: Arc<HybridExecutor>,
57}
58
59impl Moirai {
60    /// Create a new Moirai runtime with default configuration.
61    ///
62    /// # Errors
63    ///
64    /// Returns an error if the runtime cannot be initialized.
65    pub fn new() -> ExecutorResult<Self> {
66        Self::builder().build()
67    }
68
69    /// Create a builder for configuring the Moirai runtime.
70    #[must_use]
71    pub fn builder() -> MoiraiBuilder {
72        MoiraiBuilder::new()
73    }
74
75    fn spawn_or_panic<T>(
76        &self,
77        result: ExecutorResult<TaskHandle<T>>,
78        message: &'static str,
79    ) -> TaskHandle<T> {
80        result.expect(message)
81    }
82
83    fn dispatch_or_panic(&self, result: ExecutorResult<()>, message: &'static str) {
84        result.expect(message);
85    }
86
87    /// Spawn a task for parallel execution.
88    ///
89    /// This is a convenience method for spawning CPU-bound tasks.
90    ///
91    /// # Panics
92    ///
93    /// Panics if the executor fails to spawn the task, which should not happen
94    /// under normal circumstances unless the runtime is shutting down.
95    pub fn spawn<T>(&self, task: T) -> TaskHandle<T::Output>
96    where
97        T: Task,
98    {
99        self.spawn_or_panic(self.executor.spawn(task), "Failed to spawn task")
100    }
101
102    /// Spawn a parallel task using a closure.
103    ///
104    /// The task will be executed on the work-stealing thread pool.
105    ///
106    /// # Panics
107    ///
108    /// Panics if the executor fails to spawn the blocking task, which should not happen
109    /// under normal circumstances unless the runtime is shutting down.
110    pub fn spawn_fn<F, R>(&self, func: F) -> TaskHandle<R>
111    where
112        F: FnOnce() -> R + Send + 'static,
113        R: Send + 'static,
114    {
115        self.spawn_or_panic(
116            self.executor.spawn_blocking(func),
117            "Failed to spawn blocking task",
118        )
119    }
120
121    /// Spawn an async task for execution.
122    ///
123    /// The task will be executed on the async thread pool.
124    ///
125    /// # Panics
126    ///
127    /// Panics if the executor fails to spawn the async task, which should not happen
128    /// under normal circumstances unless the runtime is shutting down.
129    pub fn spawn_async<F>(&self, future: F) -> TaskHandle<F::Output>
130    where
131        F: Future + Send + 'static,
132        F::Output: Send + 'static,
133    {
134        self.spawn_or_panic(
135            self.executor.spawn_async(future),
136            "Failed to spawn async task",
137        )
138    }
139
140    /// Spawn a blocking task that may block the current thread.
141    ///
142    /// Use this for I/O-bound or blocking operations.
143    ///
144    /// # Panics
145    ///
146    /// Panics if the executor fails to spawn the blocking task, which should not happen
147    /// under normal circumstances unless the runtime is shutting down.
148    pub fn spawn_blocking<F, R>(&self, func: F) -> TaskHandle<R>
149    where
150        F: FnOnce() -> R + Send + 'static,
151        R: Send + 'static,
152    {
153        self.spawn_or_panic(
154            self.executor.spawn_blocking(func),
155            "Failed to spawn blocking task",
156        )
157    }
158
159    /// Spawn a fire-and-forget closure whose result is discarded.
160    ///
161    /// This is the cheapest dispatch path: it returns no handle and skips the
162    /// per-task result-slot allocation that [`spawn_fn`](Self::spawn_fn)
163    /// performs, making it the right choice for background work whose output is
164    /// not needed (event handlers, logging, cache warming). The task is still
165    /// drained on [`shutdown`](Self::shutdown).
166    ///
167    /// # Panics
168    ///
169    /// Panics if the executor fails to spawn the task, which should not happen
170    /// unless the runtime is shutting down.
171    pub fn spawn_detached<F>(&self, func: F)
172    where
173        F: FnOnce() + Send + 'static,
174    {
175        self.dispatch_or_panic(
176            self.executor.spawn_detached(func),
177            "Failed to spawn detached task",
178        );
179    }
180
181    /// Run a completion-only scoped fan-out on the unified scheduler.
182    ///
183    /// Use this when tasks only need to publish side effects through borrowed
184    /// synchronization primitives and the caller must wait for all tasks before
185    /// continuing. Scoped jobs may be coalesced and start after the scope body
186    /// has finished registering work.
187    ///
188    /// # Errors
189    ///
190    /// Returns an executor error if the runtime is shutting down or if a scoped
191    /// task panics.
192    pub fn scope<'scope, F>(&'scope self, body: F) -> ExecutorResult<()>
193    where
194        F: FnOnce(&MoiraiScope<'scope>) -> ExecutorResult<()>,
195    {
196        self.executor.scope::<BlockingTask, _>(body)
197    }
198
199    /// Run indexed work in worker-sized chunks on the unified scheduler.
200    ///
201    /// Use this for CPU-bound data-parallel fan-out where the caller needs
202    /// completion, not one task handle per item. Work executes through the
203    /// compute-worker pool; potentially blocking work belongs on [`Self::scope`].
204    /// The closure may borrow data that lives for the call because all chunks
205    /// complete before this method returns.
206    ///
207    /// # Errors
208    ///
209    /// Returns an executor error if the runtime is shutting down or if any
210    /// chunk panics.
211    pub fn for_each_indexed<'scope, F>(&'scope self, count: usize, task: F) -> ExecutorResult<()>
212    where
213        F: Fn(usize) + Send + Sync + 'scope,
214    {
215        self.executor.for_each_indexed::<SyncTask, _>(count, task)
216    }
217
218    /// Run indexed map/reduce in worker-sized chunks on the unified scheduler.
219    ///
220    /// `identity` must be the neutral element for `reduce`. Use this for
221    /// CPU-bound indexed data-parallel reductions where per-item task handles
222    /// are not required. Work executes through the compute-worker pool.
223    ///
224    /// # Errors
225    ///
226    /// Returns an executor error if the runtime is shutting down or if any
227    /// chunk panics.
228    pub fn map_reduce_indexed<'scope, T, Map, Reduce>(
229        &'scope self,
230        count: usize,
231        identity: T,
232        map: Map,
233        reduce: Reduce,
234    ) -> ExecutorResult<T>
235    where
236        T: Send + Clone + 'scope,
237        Map: Fn(usize) -> T + Send + Sync + 'scope,
238        Reduce: Fn(T, T) -> T + Send + Sync + 'scope,
239    {
240        self.executor
241            .map_reduce_indexed::<SyncTask, _, _, _>(count, identity, map, reduce)
242    }
243
244    /// Spawn a task with a specific priority.
245    ///
246    /// Higher priority tasks will be executed before lower priority tasks.
247    ///
248    /// # Panics
249    ///
250    /// Panics if the executor fails to spawn the task with priority, which should not happen
251    /// under normal circumstances unless the runtime is shutting down.
252    pub fn spawn_with_priority<T>(&self, task: T, priority: Priority) -> TaskHandle<T::Output>
253    where
254        T: Task,
255    {
256        self.spawn_or_panic(
257            self.executor.spawn_with_priority(task, priority, None),
258            "Failed to spawn task with priority",
259        )
260    }
261
262    /// Spawn a closure with priority as a task (convenience method).
263    pub fn spawn_fn_with_priority<F, R>(&self, f: F, priority: Priority) -> TaskHandle<R>
264    where
265        F: FnOnce() -> R + Send + 'static,
266        R: Send + 'static,
267    {
268        // Let the executor handle ID assignment and priority
269        let task = TaskBuilder::new().build(f);
270        self.spawn_with_priority(task, priority)
271    }
272
273    /// Block the current thread until the future completes.
274    ///
275    /// This is useful for running async code from synchronous contexts.
276    pub fn block_on<F>(&self, future: F) -> F::Output
277    where
278        F: Future,
279    {
280        self.executor.block_on(future)
281    }
282
283    /// Try to run pending tasks without blocking.
284    ///
285    /// Returns `true` if any tasks were executed, `false` if no work was available.
286    #[must_use]
287    pub fn try_run(&self) -> bool {
288        self.executor.try_run()
289    }
290
291    /// Returns true when queued or active runtime work exists.
292    #[must_use]
293    pub fn has_work(&self) -> bool {
294        self.executor.has_work()
295    }
296
297    /// Wait until queued and active runtime work completes without shutting down workers.
298    ///
299    /// Use this as a non-destructive process-fusion barrier when producers have
300    /// finished submitting a batch and the runtime should process all available
301    /// work before the caller continues. New tasks submitted after this method
302    /// observes quiescence belong to a later batch.
303    ///
304    /// # Errors
305    ///
306    /// Returns an executor error if the scheduler join operation fails.
307    pub fn join(&self) -> ExecutorResult<()> {
308        self.executor.join()
309    }
310
311    /// Shutdown the runtime gracefully.
312    ///
313    /// This will wait for all currently running tasks to complete before
314    /// shutting down the thread pools.
315    pub fn shutdown(&self) {
316        self.executor.shutdown();
317    }
318
319    /// Shutdown the runtime with a timeout.
320    ///
321    /// If tasks don't complete within the timeout, they will be forcefully
322    /// terminated.
323    pub fn shutdown_timeout(&self, timeout: Duration) {
324        self.executor.shutdown_timeout(timeout);
325    }
326
327    /// Check if the runtime is shutting down.
328    #[must_use]
329    pub fn is_shutting_down(&self) -> bool {
330        self.executor.is_shutting_down()
331    }
332
333    /// Get the number of worker threads.
334    #[must_use]
335    pub fn worker_count(&self) -> usize {
336        self.executor.worker_count()
337    }
338
339    /// Get the current load (number of pending tasks).
340    #[must_use]
341    pub fn load(&self) -> usize {
342        self.executor.load()
343    }
344
345    /// Get runtime statistics.
346    #[cfg(feature = "metrics")]
347    #[must_use]
348    pub fn stats(&self) -> moirai_core::executor::ExecutorStats {
349        self.executor.stats()
350    }
351
352    /// Create a channel for communication, bounded at
353    /// [`DEFAULT_CHANNEL_CAPACITY`].
354    ///
355    /// Bounded is the default because an unbounded queue converts a slow
356    /// consumer into unbounded memory growth: a full channel blocks its
357    /// producer (or returns [`ChannelError::Full`] from `try_send`) instead of
358    /// allocating. Use [`Self::bounded_channel`] when the right capacity is
359    /// known; the unbounded queue remains available as
360    /// `moirai_core::channel::unbounded`, whose documentation states the cost.
361    ///
362    /// [`DEFAULT_CHANNEL_CAPACITY`]: moirai_core::channel::DEFAULT_CHANNEL_CAPACITY
363    /// [`ChannelError::Full`]: moirai_core::channel::ChannelError::Full
364    #[must_use]
365    pub fn channel<T: Send + 'static>(
366        &self,
367    ) -> (
368        moirai_core::channel::MpmcSender<T>,
369        moirai_core::channel::MpmcReceiver<T>,
370    ) {
371        moirai_core::channel::mpmc(moirai_core::channel::DEFAULT_CHANNEL_CAPACITY)
372    }
373
374    /// Create a bounded channel with an explicit capacity.
375    #[must_use]
376    pub fn bounded_channel<T: Send + 'static>(
377        &self,
378        capacity: usize,
379    ) -> (
380        moirai_core::channel::MpmcSender<T>,
381        moirai_core::channel::MpmcReceiver<T>,
382    ) {
383        moirai_core::channel::mpmc(capacity)
384    }
385}
386
387impl Default for Moirai {
388    fn default() -> Self {
389        Self::new().expect("Failed to create default Moirai runtime")
390    }
391}