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
impl HybridExecutor
Sourcepub fn new(config: ExecutorConfig) -> Result<HybridExecutor, ExecutorError>
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.
Sourcepub fn scope<'scope, C, F>(&'scope self, body: F) -> Result<(), ExecutorError>
pub fn scope<'scope, C, F>(&'scope self, body: F) -> 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,
impl<S> HybridExecutor<S>where
S: WorkScheduler,
Sourcepub fn config(&self) -> &ExecutorConfig
pub fn config(&self) -> &ExecutorConfig
Get executor configuration.
Sourcepub fn shutdown(&mut self) -> Result<(), ExecutorError>
pub fn shutdown(&mut self) -> Result<(), ExecutorError>
Shutdown the executor gracefully.
Sourcepub fn metrics(&self) -> &ExecutorMetrics
pub fn metrics(&self) -> &ExecutorMetrics
Get executor metrics.
Sourcepub fn submit_task<F>(&self, task: F) -> Result<TaskId, ExecutorError>
pub fn submit_task<F>(&self, task: F) -> Result<TaskId, ExecutorError>
Submit an untyped synchronous job.
Sourcepub fn for_each_indexed<'scope, C, F>(
&'scope self,
count: usize,
task: F,
) -> Result<(), ExecutorError>
pub fn for_each_indexed<'scope, C, F>( &'scope self, count: usize, task: F, ) -> Result<(), ExecutorError>
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.
Sourcepub fn map_reduce_indexed<'scope, C, T, Map, Reduce>(
&'scope self,
count: usize,
identity: T,
map: Map,
reduce: Reduce,
) -> Result<T, ExecutorError>
pub fn map_reduce_indexed<'scope, C, T, Map, Reduce>( &'scope self, count: usize, identity: T, map: Map, reduce: Reduce, ) -> Result<T, ExecutorError>
Run indexed map/reduce in worker-sized chunks on the unified scheduler.
Sourcepub fn active_workers(&self) -> usize
pub fn active_workers(&self) -> usize
Get the number of active workers.
Sourcepub fn total_workers(&self) -> usize
pub fn total_workers(&self) -> usize
Get the total number of workers.
Sourcepub fn pending_tasks(&self) -> usize
pub fn pending_tasks(&self) -> usize
Get pending task count across all workers.
Sourcepub fn join(&self) -> Result<(), ExecutorError>
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,
impl<S> Drop for HybridExecutor<S>where
S: WorkScheduler,
Source§impl<S> Executor for HybridExecutor<S>where
S: WorkScheduler,
impl<S> Executor for HybridExecutor<S>where
S: WorkScheduler,
Source§fn stats(&self) -> ExecutorStats
fn stats(&self) -> ExecutorStats
Source§impl<S> ExecutorControl for HybridExecutor<S>where
S: WorkScheduler,
impl<S> ExecutorControl for HybridExecutor<S>where
S: WorkScheduler,
Source§fn shutdown_timeout(&self, timeout: Duration)
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>::Outputwhere
F: Future,
fn block_on<F>(&self, future: F) -> <F as Future>::Outputwhere
F: Future,
Source§fn is_shutting_down(&self) -> bool
fn is_shutting_down(&self) -> bool
Source§fn worker_count(&self) -> usize
fn worker_count(&self) -> usize
Source§impl<S> TaskManager for HybridExecutor<S>where
S: WorkScheduler,
impl<S> TaskManager for HybridExecutor<S>where
S: WorkScheduler,
Source§fn cancel_task(&self, id: TaskId) -> Result<(), ExecutorError>
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
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>
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>
fn task_status(&self, id: TaskId) -> Option<TaskStatus>
Source§impl<S> TaskSpawner for HybridExecutor<S>where
S: WorkScheduler,
impl<S> TaskSpawner for HybridExecutor<S>where
S: WorkScheduler,
Source§fn spawn<T>(
&self,
task: T,
) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
fn spawn<T>( &self, task: T, ) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
Source§fn spawn_async<F>(
&self,
future: F,
) -> Result<TaskHandle<<F as Future>::Output>, ExecutorError>
fn spawn_async<F>( &self, future: F, ) -> Result<TaskHandle<<F as Future>::Output>, ExecutorError>
Source§fn spawn_blocking<F, R>(&self, func: F) -> Result<TaskHandle<R>, ExecutorError>
fn spawn_blocking<F, R>(&self, func: F) -> Result<TaskHandle<R>, ExecutorError>
Source§fn spawn_detached<F>(&self, func: F) -> Result<(), ExecutorError>
fn spawn_detached<F>(&self, func: F) -> Result<(), ExecutorError>
Source§fn spawn_with_priority<T>(
&self,
task: T,
priority: Priority,
locality_hint: Option<usize>,
) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
fn spawn_with_priority<T>( &self, task: T, priority: Priority, locality_hint: Option<usize>, ) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>
Source§fn spawn_local<T>(
&self,
task: T,
) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>where
T: Task + 'static,
fn spawn_local<T>(
&self,
task: T,
) -> Result<TaskHandle<<T as Task>::Output>, ExecutorError>where
T: Task + 'static,
spawn_local)