Skip to main content

ntex_rt/
rt_tokio.rs

1use std::task::{Context, Poll, ready};
2use std::{fmt, future::Future, pin::Pin};
3
4use async_channel::Sender;
5
6use crate::arbiter::{Arbiter, ArbiterCommand};
7
8#[inline]
9/// Spawn a future on the current thread.
10///
11/// This does not create a new Arbiter or Arbiter address,
12/// it is simply a helper for spawning futures on the current thread.
13///
14/// # Panics
15///
16/// This function panics if ntex system is not running.
17pub fn spawn<F>(f: F) -> JoinHandle<F::Output>
18where
19    F: Future + 'static,
20{
21    let task = tok_io::task::spawn_local(crate::task::wrap(f));
22
23    JoinHandle {
24        task: Some(Either::Task(task)),
25    }
26}
27
28#[derive(Debug, Copy, Clone)]
29pub struct JoinError;
30
31impl fmt::Display for JoinError {
32    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
33        write!(f, "JoinError")
34    }
35}
36
37impl std::error::Error for JoinError {}
38
39#[derive(Debug)]
40enum Either<T> {
41    Task(tok_io::task::JoinHandle<T>),
42    Spawn(oneshot::AsyncReceiver<T>),
43}
44
45#[derive(Debug)]
46pub struct JoinHandle<T> {
47    task: Option<Either<T>>,
48}
49
50impl<T> JoinHandle<T> {
51    /// Cancels the task.
52    pub fn cancel(mut self) {
53        if let Some(Either::Task(fut)) = self.task.take() {
54            fut.abort();
55        }
56    }
57
58    /// Detaches the task to let it keep running in the background.
59    pub fn detach(mut self) {
60        self.task.take();
61    }
62
63    /// Returns true if the current task is finished.
64    pub fn is_finished(&self) -> bool {
65        match &self.task {
66            Some(Either::Task(fut)) => fut.is_finished(),
67            Some(Either::Spawn(fut)) => fut.is_closed(),
68            None => true,
69        }
70    }
71}
72
73impl<T> Drop for JoinHandle<T> {
74    fn drop(&mut self) {
75        self.task.take();
76    }
77}
78
79impl<T> Future for JoinHandle<T> {
80    type Output = Result<T, JoinError>;
81
82    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
83        Poll::Ready(match self.task.as_mut() {
84            Some(Either::Task(fut)) => ready!(Pin::new(fut).poll(cx)).map_err(|_| JoinError),
85            Some(Either::Spawn(fut)) => ready!(Pin::new(fut).poll(cx)).map_err(|_| JoinError),
86            None => Err(JoinError),
87        })
88    }
89}
90
91#[derive(Clone, Debug)]
92/// Handle to the runtime.
93pub struct Handle(Sender<ArbiterCommand>);
94
95impl Handle {
96    pub(crate) fn new(sender: Sender<ArbiterCommand>) -> Self {
97        Self(sender)
98    }
99
100    pub fn current() -> Self {
101        Self(Arbiter::current().0.sender.clone())
102    }
103
104    #[inline]
105    /// Wake up runtime.
106    pub fn notify(&self) {}
107
108    /// Spawns a new asynchronous task, returning a [`JoinHandle`] for it.
109    ///
110    /// Spawning a task enables the task to execute concurrently to other tasks.
111    /// There is no guarantee that a spawned task will execute to completion.
112    pub fn spawn<F>(&self, future: F) -> JoinHandle<F::Output>
113    where
114        F: Future + Send + 'static,
115        F::Output: Send + 'static,
116    {
117        let (tx, rx) = oneshot::async_channel();
118
119        let _ = self
120            .0
121            .try_send(ArbiterCommand::Execute(Box::pin(async move {
122                let result = future.await;
123                let _ = tx.send(result);
124            })));
125
126        JoinHandle {
127            task: Some(Either::Spawn(rx)),
128        }
129    }
130}