Skip to main content

karyon_core/async_runtime/
executor.rs

1use std::{future::Future, sync::Arc};
2
3#[cfg(feature = "smol")]
4use std::{
5    future::pending,
6    num::NonZeroUsize,
7    panic::{catch_unwind, AssertUnwindSafe},
8    thread,
9};
10
11#[cfg(feature = "smol")]
12use log::error;
13
14use once_cell::sync::OnceCell;
15
16#[cfg(feature = "smol")]
17pub use smol::Executor as SmolEx;
18
19#[cfg(feature = "tokio")]
20pub use tokio::runtime::Runtime;
21
22use super::Task;
23
24/// A handle to a multi-threaded async executor.
25///
26/// karyon does not run the executor for you. The caller owns the runtime and
27/// is responsible for driving it. On `tokio` this is just a `Runtime`. On
28/// `smol` you must spin up worker threads yourself, e.g. via
29/// [`easy_parallel`](https://docs.rs/easy-parallel):
30///
31/// ```ignore
32/// use std::sync::Arc;
33/// use async_channel::unbounded;
34/// use easy_parallel::Parallel;
35/// use smol::{future, Executor as SmolEx};
36///
37/// let ex = Arc::new(SmolEx::new());
38/// let (signal, shutdown) = unbounded::<()>();
39///
40/// let num_threads = std::thread::available_parallelism()
41///     .map(|n| n.get())
42///     .unwrap_or(1);
43///
44/// Parallel::new()
45///     .each(0..num_threads, |_| future::block_on(ex.run(shutdown.recv())))
46///     .finish(|| future::block_on(async {
47///         // your async main here
48///         drop(signal);
49///     }));
50/// ```
51///
52/// Pass `ex.clone().into()` to karyon APIs that take an [`Executor`].
53#[derive(Clone)]
54pub struct Executor {
55    #[cfg(feature = "smol")]
56    inner: Arc<SmolEx<'static>>,
57    #[cfg(feature = "tokio")]
58    inner: Arc<Runtime>,
59}
60
61impl Executor {
62    pub fn spawn<T: Send + 'static>(
63        &self,
64        future: impl Future<Output = T> + Send + 'static,
65    ) -> Task<T> {
66        self.inner.spawn(future).into()
67    }
68
69    #[cfg(feature = "tokio")]
70    pub fn handle(&self) -> &tokio::runtime::Handle {
71        self.inner.handle()
72    }
73}
74
75static GLOBAL_EXECUTOR: OnceCell<Executor> = OnceCell::new();
76
77/// Returns the process-wide global executor. Multi-threaded on both
78/// runtimes: `smol` runs one worker thread per core, `tokio` uses its
79/// own worker threads.
80///
81/// Note: this is the convenience path. Resource-conscious users should
82/// create and drive their own executor and pass it to karyon APIs; no
83/// worker threads are spawned until the first call to this function.
84pub fn global_executor() -> Executor {
85    #[cfg(feature = "smol")]
86    fn init_executor() -> Executor {
87        let ex = Arc::new(smol::Executor::new());
88        let num_threads = thread::available_parallelism()
89            .map(NonZeroUsize::get)
90            .unwrap_or(1);
91        for i in 0..num_threads {
92            let ex = ex.clone();
93            thread::Builder::new()
94                .name(format!("smol-executor-{i}"))
95                .spawn(move || loop {
96                    // A panicking task unwinds out of block_on; log it
97                    // and keep the worker alive.
98                    let run = AssertUnwindSafe(|| smol::block_on(ex.run(pending::<()>())));
99                    if catch_unwind(run).is_err() {
100                        error!("global executor worker recovered from a task panic");
101                    }
102                })
103                .expect("cannot spawn executor thread");
104        }
105        // Prevent spawning another thread by running the process driver on this
106        // executor. see https://github.com/smol-rs/smol/blob/master/src/spawn.rs
107        ex.spawn(async_process::driver()).detach();
108        Executor { inner: ex }
109    }
110
111    #[cfg(feature = "tokio")]
112    fn init_executor() -> Executor {
113        let ex = Arc::new(tokio::runtime::Runtime::new().expect("cannot build tokio runtime"));
114        Executor { inner: ex }
115    }
116
117    GLOBAL_EXECUTOR.get_or_init(init_executor).clone()
118}
119
120#[cfg(feature = "smol")]
121impl From<Arc<smol::Executor<'static>>> for Executor {
122    fn from(ex: Arc<smol::Executor<'static>>) -> Executor {
123        Executor { inner: ex }
124    }
125}
126
127#[cfg(feature = "tokio")]
128impl From<Arc<tokio::runtime::Runtime>> for Executor {
129    fn from(rt: Arc<tokio::runtime::Runtime>) -> Executor {
130        Executor { inner: rt }
131    }
132}
133
134#[cfg(feature = "smol")]
135impl From<smol::Executor<'static>> for Executor {
136    fn from(ex: smol::Executor<'static>) -> Executor {
137        Executor {
138            inner: Arc::new(ex),
139        }
140    }
141}
142
143#[cfg(feature = "tokio")]
144impl From<tokio::runtime::Runtime> for Executor {
145    fn from(rt: tokio::runtime::Runtime) -> Executor {
146        Executor {
147            inner: Arc::new(rt),
148        }
149    }
150}