Skip to main content

ntex_net/tokio/
mod.rs

1#![allow(unused_variables)]
2#[cfg(unix)]
3use std::os::unix::net::UnixStream as OsUnixStream;
4use std::{io::Result as IoResult, net, net::SocketAddr, path::PathBuf};
5
6use ntex_io::Io;
7use ntex_service::cfg::SharedCfg;
8
9#[doc(hidden)]
10pub mod internal {
11    pub use tok_io::*;
12}
13
14use crate::channel::{self, Receiver};
15
16mod io;
17
18pub use self::io::TokioIoBoxed;
19
20pub(crate) struct TcpStream(tok_io::net::TcpStream);
21
22#[cfg(unix)]
23pub(crate) struct UnixStream(tok_io::net::UnixStream);
24
25/// ntex reactor based on tokio runtime
26#[derive(Copy, Clone, Debug)]
27pub struct Reactor;
28
29/// Runs the provided future, blocking the current thread until the future
30/// completes.
31pub(crate) fn block_on<F: Future<Output = ()>>(fut: F) {
32    if let Ok(hnd) = tok_io::runtime::Handle::try_current() {
33        log::debug!("Use existing tokio runtime and block on future");
34        hnd.block_on(tok_io::task::LocalSet::new().run_until(fut));
35    } else {
36        log::debug!("Create tokio runtime and block on future");
37
38        let rt = tok_io::runtime::Builder::new_current_thread()
39            .enable_all()
40            .build()
41            .unwrap();
42
43        tok_io::task::LocalSet::new().block_on(&rt, fut);
44    }
45}
46
47impl ntex_rt::Driver for Reactor {
48    fn run(&self, _: &ntex_rt::Runtime) -> std::io::Result<()> {
49        panic!("Not supported")
50    }
51
52    fn handle(&self) -> Box<dyn ntex_rt::Notify> {
53        panic!("Not supported")
54    }
55
56    fn clear(&self) {}
57}
58
59impl crate::Reactor for Reactor {
60    fn tcp_connect(&self, addr: SocketAddr, cfg: SharedCfg) -> Receiver<Io> {
61        let (tx, rx) = channel::create();
62        ntex_rt::spawn(async move {
63            let result = async {
64                let sock = tok_io::net::TcpStream::connect(addr).await?;
65                sock.set_nodelay(true)?;
66                Ok(Io::new(TcpStream(sock), cfg))
67            }
68            .await;
69            let _ = tx.send(result);
70        });
71
72        rx
73    }
74
75    fn unix_connect(&self, addr: PathBuf, cfg: SharedCfg) -> Receiver<Io> {
76        #[cfg(unix)]
77        {
78            let (tx, rx) = channel::create();
79            ntex_rt::spawn(async move {
80                let result = async {
81                    let sock = tok_io::net::UnixStream::connect(addr).await?;
82                    Ok(Io::new(UnixStream(sock), cfg))
83                }
84                .await;
85                let _ = tx.send(result);
86            });
87
88            rx
89        }
90
91        #[cfg(not(unix))]
92        {
93            Receiver::new(Err(std::io::Error::other(
94                "Unix domain sockets are not supported",
95            )))
96        }
97    }
98
99    fn from_tcp_stream(&self, stream: net::TcpStream, cfg: SharedCfg) -> IoResult<Io> {
100        stream.set_nonblocking(true)?;
101        stream.set_nodelay(true)?;
102        Ok(Io::new(
103            TcpStream(tok_io::net::TcpStream::from_std(stream)?),
104            cfg,
105        ))
106    }
107
108    #[cfg(unix)]
109    fn from_unix_stream(&self, stream: OsUnixStream, cfg: SharedCfg) -> IoResult<Io> {
110        stream.set_nonblocking(true)?;
111        Ok(Io::new(
112            UnixStream(tok_io::net::UnixStream::from_std(stream)?),
113            cfg,
114        ))
115    }
116}