use crate::{MoiraiBuilder, MoiraiScope};
#[cfg(feature = "metrics")]
use moirai_core::executor::Executor;
use moirai_core::{
error::*,
executor::{ExecutorControl, TaskSpawner},
Priority, Task, TaskBuilder, TaskHandle,
};
use moirai_executor::{BlockingTask, HybridExecutor};
use std::{future::Future, sync::Arc, time::Duration};
#[derive(Clone)]
pub struct Moirai {
pub(crate) executor: Arc<HybridExecutor>,
}
impl Moirai {
pub fn new() -> ExecutorResult<Self> {
Self::builder().build()
}
#[must_use]
pub fn builder() -> MoiraiBuilder {
MoiraiBuilder::new()
}
pub fn spawn<T>(&self, task: T) -> TaskHandle<T::Output>
where
T: Task,
{
self.executor.spawn(task).expect("Failed to spawn task")
}
pub fn spawn_fn<F, R>(&self, func: F) -> TaskHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
self.executor
.spawn_blocking(func)
.expect("Failed to spawn blocking task")
}
pub fn spawn_async<F>(&self, future: F) -> TaskHandle<F::Output>
where
F: Future + Send + 'static,
F::Output: Send + 'static,
{
self.executor
.spawn_async(future)
.expect("Failed to spawn async task")
}
pub fn spawn_blocking<F, R>(&self, func: F) -> TaskHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
self.executor
.spawn_blocking(func)
.expect("Failed to spawn blocking task")
}
pub fn spawn_detached<F>(&self, func: F)
where
F: FnOnce() + Send + 'static,
{
self.executor
.spawn_detached(func)
.expect("Failed to spawn detached task");
}
pub fn scope<'scope, F>(&'scope self, body: F) -> ExecutorResult<()>
where
F: FnOnce(&MoiraiScope<'scope>) -> ExecutorResult<()>,
{
self.executor.scope::<BlockingTask, _>(body)
}
pub fn for_each_indexed<'scope, F>(&'scope self, count: usize, task: F) -> ExecutorResult<()>
where
F: Fn(usize) + Send + Sync + 'scope,
{
self.executor
.for_each_indexed::<BlockingTask, _>(count, task)
}
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,
{
self.executor
.map_reduce_indexed::<BlockingTask, _, _, _>(count, identity, map, reduce)
}
pub fn spawn_with_priority<T>(&self, task: T, priority: Priority) -> TaskHandle<T::Output>
where
T: Task,
{
self.executor
.spawn_with_priority(task, priority, None)
.expect("Failed to spawn task with priority")
}
pub fn spawn_fn_with_priority<F, R>(&self, f: F, priority: Priority) -> TaskHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
let task = TaskBuilder::new().build(f);
self.spawn_with_priority(task, priority)
}
pub fn block_on<F>(&self, future: F) -> F::Output
where
F: Future,
{
self.executor.block_on(future)
}
#[must_use]
pub fn try_run(&self) -> bool {
self.executor.try_run()
}
#[must_use]
pub fn has_work(&self) -> bool {
self.executor.has_work()
}
pub fn join(&self) -> ExecutorResult<()> {
self.executor.join()
}
pub fn shutdown(&self) {
self.executor.shutdown();
}
pub fn shutdown_timeout(&self, timeout: Duration) {
self.executor.shutdown_timeout(timeout);
}
#[must_use]
pub fn is_shutting_down(&self) -> bool {
self.executor.is_shutting_down()
}
#[must_use]
pub fn worker_count(&self) -> usize {
self.executor.worker_count()
}
#[must_use]
pub fn load(&self) -> usize {
self.executor.load()
}
#[cfg(feature = "metrics")]
#[must_use]
pub fn stats(&self) -> moirai_core::executor::ExecutorStats {
self.executor.stats()
}
#[must_use]
pub fn channel<T: Send + 'static>(
&self,
) -> (
moirai_core::channel::MpmcSender<T>,
moirai_core::channel::MpmcReceiver<T>,
) {
moirai_core::channel::unbounded()
}
#[must_use]
pub fn bounded_channel<T: Send + 'static>(
&self,
capacity: usize,
) -> (
moirai_core::channel::MpmcSender<T>,
moirai_core::channel::MpmcReceiver<T>,
) {
moirai_core::channel::mpmc(capacity)
}
#[cfg(feature = "gpu")]
pub async fn create_gpu_context(&self) -> Result<moirai_gpu::GpuContext, moirai_gpu::GpuError> {
moirai_gpu::GpuContext::new().await
}
#[cfg(feature = "gpu")]
pub async fn create_gpu_context_with_preferences(
&self,
preferences: moirai_gpu::DevicePreferences,
) -> Result<moirai_gpu::GpuContext, moirai_gpu::GpuError> {
moirai_gpu::GpuContext::with_preferences(preferences).await
}
#[cfg(feature = "gpu")]
pub fn spawn_gpu<T>(
&self,
gpu_context: &moirai_gpu::GpuContext,
task: T,
) -> moirai_gpu::GpuTaskFuture<T::Output>
where
T: moirai_gpu::GpuTask + Send + 'static,
T::Output: Send + 'static,
{
gpu_context.spawn_gpu_task(task)
}
}
impl Default for Moirai {
fn default() -> Self {
Self::new().expect("Failed to create default Moirai runtime")
}
}