Skip to main content

moirai_executor/hybrid/
control.rs

1use std::sync::atomic::Ordering;
2
3#[cfg(feature = "metrics")]
4use moirai_core::executor::ExecutorStats;
5use moirai_core::executor::{Executor, ExecutorControl};
6
7use super::HybridExecutor;
8use crate::schedule::WorkScheduler;
9
10impl<S: WorkScheduler> ExecutorControl for HybridExecutor<S> {
11    fn block_on<F>(&self, future: F) -> F::Output
12    where
13        F: core::future::Future,
14    {
15        crate::schedule::wake::block_on_current_thread(future)
16    }
17
18    fn try_run(&self) -> bool {
19        self.refresh_scheduler_metrics();
20        self.scheduler.has_work()
21    }
22
23    fn shutdown(&self) {
24        self.shutdown_signal.store(true, Ordering::Release);
25        self.scheduler.shutdown();
26    }
27
28    /// Graceful shutdown bounded by `timeout` for the *caller*.
29    ///
30    /// The drain runs on a helper thread; this call returns once the drain
31    /// completes or `timeout` elapses, whichever comes first. If the deadline
32    /// lapses, workers keep draining in the background and the (idempotent)
33    /// scheduler shutdown is re-joined on executor drop.
34    fn shutdown_timeout(&self, timeout: core::time::Duration) {
35        self.shutdown_signal.store(true, Ordering::Release);
36
37        let scheduler = self.scheduler.clone();
38        let (done_sender, done_receiver) = std::sync::mpsc::sync_channel::<()>(1);
39        std::thread::spawn(move || {
40            scheduler.shutdown();
41            let _ = done_sender.send(());
42        });
43
44        // A timeout here is the expected bounded outcome, not a failure to
45        // mask: the drain continues in the background by contract (above).
46        let _ = done_receiver.recv_timeout(timeout);
47    }
48
49    fn is_shutting_down(&self) -> bool {
50        self.shutdown_signal.load(Ordering::Acquire)
51    }
52
53    fn worker_count(&self) -> usize {
54        self.scheduler.worker_count()
55    }
56
57    fn load(&self) -> usize {
58        self.scheduler.pending_tasks()
59    }
60}
61
62impl<S: WorkScheduler> Executor for HybridExecutor<S> {
63    #[cfg(feature = "metrics")]
64    fn stats(&self) -> ExecutorStats {
65        self.refresh_scheduler_metrics();
66        ExecutorStats {
67            tasks_executed: self.metrics.tasks_completed.load(Ordering::Acquire),
68            tasks_queued: self.scheduler.pending_tasks(),
69            avg_execution_time_ns: self.metrics.average_task_duration().as_nanos() as u64,
70        }
71    }
72}