moirai_executor/hybrid/
control.rs1use 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 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 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}