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#[derive(Copy, Clone, Debug)]
27pub struct Reactor;
28
29pub(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}