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}