Skip to main content

HybridExecutor

Struct HybridExecutor 

Source
pub struct HybridExecutor<S = ThreadScheduler>
where S: WorkScheduler,
{ /* private fields */ }
Expand description

Main hybrid executor that coordinates sync, async, and blocking tasks.

Generic over the work-stealing runtime S behind the WorkScheduler seam; the default ThreadScheduler backs the production runtime, while the parameter lets a substitute (e.g. a single-threaded wasm32 scheduler) be plugged in without touching this façade.

Implementations§

Source§

impl HybridExecutor

Source

pub fn new(config: ExecutorConfig) -> Result<HybridExecutor, ExecutorError>

Create a new hybrid executor with the given configuration.

§Errors

Returns ExecutorError::InvalidConfiguration when the configured global admission bound cannot supply at least two slots per worker, ExecutorError::InvalidLocalQueueInitialCapacity when local capacity cannot normalize or form the required deque allocation layouts, or propagates scheduler construction failures.

Source

pub fn scope<'scope, C, F>(&'scope self, body: F) -> Result<(), ExecutorError>
where C: WorkClass, F: FnOnce(&SchedulerScope<'scope, C>) -> Result<(), ExecutorError>,

Run a scoped fan-out directly on the unified scheduler.

This path is for completion-only work that does not require per-task result handles or lifecycle metadata. It preserves borrowing semantics by waiting for all spawned jobs before returning.

scope is inherent to the default ThreadScheduler backing because its signature exposes a concrete SchedulerScope borrow handle, which is outside the substitutable WorkScheduler seam.

Source§

impl<S> HybridExecutor<S>
where S: WorkScheduler,

Source

pub fn config(&self) -> &ExecutorConfig

Get executor configuration.

Source

pub fn shutdown(&mut self) -> Result<(), ExecutorError>

Shutdown the executor gracefully.

Source

pub fn metrics(&self) -> &ExecutorMetrics

Get executor metrics.

Source

pub fn submit_task<F>(&self, task: F) -> Result<TaskId, ExecutorError>
where F: FnOnce() + Send + 'static,

Submit an untyped synchronous job.

Source

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

Run indexed work in worker-sized chunks on the unified scheduler.

This path avoids per-item task handles and lifecycle metadata when the caller only needs completion for a bounded index domain.

Source

pub fn map_reduce_indexed<'scope, C, T, Map, Reduce>( &'scope self, count: usize, identity: T, map: Map, reduce: Reduce, ) -> Result<T, ExecutorError>
where C: WorkClass, 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.

Source

pub fn active_workers(&self) -> usize

Get the number of active workers.

Source

pub fn total_workers(&self) -> usize

Get the total number of workers.

Source

pub fn pending_tasks(&self) -> usize

Get pending task count across all workers.

Source

pub fn has_work(&self) -> bool

Returns true when queued or active scheduler work exists.

Source

pub fn join(&self) -> Result<(), ExecutorError>

Wait until queued and active scheduler work completes without shutting down workers.

Trait Implementations§

Source§

impl<S> Drop for HybridExecutor<S>
where S: WorkScheduler,

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more
Source§

impl<S> Executor for HybridExecutor<S>
where S: WorkScheduler,

Source§

fn stats(&self) -> ExecutorStats

Get comprehensive executor statistics. Read more
Source§

impl<S> ExecutorControl for HybridExecutor<S>
where S: WorkScheduler,

Source§

fn shutdown_timeout(&self, timeout: Duration)

Graceful shutdown bounded by timeout for the caller.

The drain runs on a helper thread; this call returns once the drain completes or timeout elapses, whichever comes first. If the deadline lapses, workers keep draining in the background and the (idempotent) scheduler shutdown is re-joined on executor drop.

Source§

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

Block the current thread until the future completes. Read more
Source§

fn try_run(&self) -> bool

Attempt to run tasks without blocking. Read more
Source§

fn shutdown(&self)

Shutdown the executor gracefully. Read more
Source§

fn is_shutting_down(&self) -> bool

Check if the executor is shutting down. Read more
Source§

fn worker_count(&self) -> usize

Get the number of worker threads. Read more
Source§

fn load(&self) -> usize

Get the current load (number of pending tasks). Read more
Source§

impl<S> TaskManager for HybridExecutor<S>
where S: WorkScheduler,

Source§

fn cancel_task(&self, id: TaskId) -> Result<(), ExecutorError>

Cooperative cancellation.

Contract: a queued task that has not started skips its body when a worker dequeues it — the task completes with TaskError::Cancelled and status Cancelled. A task that already started is not preempted; it runs to completion and this call still returns Ok(()) (the request is recorded but has no effect). Cancelling an already-completed task is a no-op Ok(()). An unknown task ID is an error.

Source§

fn wait_for_task( &self, id: TaskId, timeout: Option<Duration>, ) -> impl Future<Output = Result<(), ExecutorError>> + Send

Event-driven completion wait.

The future registers a completion waker with the task registry and returns Pending; mark_completed/mark_cancelled wake it exactly (no polling loop and no thread-blocking sleep). The deadline is checked on every poll — at creation, on each completion wake, and on any external poll after expiry. No in-scope timer exists (the async timer lives in moirai-async), so if the task never completes, observing the expiry requires the caller to poll after the deadline (e.g. a timeout-aware runtime); a completion wake always resolves promptly.

Source§

fn task_stats(&self, id: TaskId) -> Option<TaskStats>

Statistics limited to what the executor actually tracks.

priority is the value recorded at spawn; timing fields come from the registry’s lifecycle timestamps.

Source§

fn task_status(&self, id: TaskId) -> Option<TaskStatus>

Get the current status of a task. Read more
Source§

impl<S> TaskSpawner for HybridExecutor<S>
where S: WorkScheduler,

Source§

fn spawn<T>( &self, task: T, ) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
where T: Task + Send + 'static, <T as Task>::Output: Send + 'static,

Spawns a new task for execution. Read more
Source§

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

Spawns an asynchronous task (Future) for execution. Read more
Source§

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

Spawns a blocking task that may perform I/O or CPU-intensive work. Read more
Source§

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

Spawns a fire-and-forget task whose result is discarded. Read more
Source§

fn spawn_with_priority<T>( &self, task: T, priority: Priority, locality_hint: Option<usize>, ) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
where T: Task + Send + 'static, <T as Task>::Output: Send + 'static,

Spawns a task with specific priority and scheduling hints. Read more
Source§

fn spawn_local<T>( &self, task: T, ) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
where T: Task + 'static,

Spawn a task on the current thread’s local queue for better locality (inspired by Tokio’s spawn_local)

Auto Trait Implementations§

§

impl<S> Freeze for HybridExecutor<S>
where S: Freeze,

§

impl<S> RefUnwindSafe for HybridExecutor<S>
where S: RefUnwindSafe,

§

impl<S> Send for HybridExecutor<S>

§

impl<S> Sync for HybridExecutor<S>

§

impl<S> Unpin for HybridExecutor<S>
where S: Unpin,

§

impl<S> UnsafeUnpin for HybridExecutor<S>
where S: UnsafeUnpin,

§

impl<S> UnwindSafe for HybridExecutor<S>
where S: UnwindSafe,

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> 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, 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.