use tokio::io::{AsyncWriteExt, AsyncReadExt};
use std::time::Duration;
use tokio::net::TcpStream;
use tokio::time::timeout;
use crate::mc_text::ServerStatus;
use crate::packets::{ClientHandshake, ServerQueryResponse, StatusQuery};
use anyhow::{anyhow, Result};
use tokio::net::lookup_host;
use tokio_socks::tcp::Socks5Stream;
fn is_domain(addr: &str) -> bool {
addr.parse::<std::net::IpAddr>().is_err()
}
pub struct Connection<T> {
pub is_initialized: bool,
pub stream: Option<T>,
pub timeout: Option<u64>,
pub proxy_addr: Option<(String, u16)>,
pub addr: (String, u16),
}
impl Connection<TcpStream> {
pub fn new(addr: (String, u16)) -> Self {
Self {
stream: None,
timeout: None,
is_initialized: true,
proxy_addr: None,
addr,
}
}
pub async fn connect(&mut self) -> Result<Self> {
let _timeout = self.timeout.unwrap_or(8000);
#[cfg(not(feature = "resolve"))]
{
let addr = self.addr.clone();
if is_domain(&addr.0) {
return Err(anyhow!(r#"Enable feature "resolve" to enable domain resolving"#));
}
match &self.proxy_addr {
None => {
let stream = timeout(Duration::from_millis(_timeout), TcpStream::connect(addr.clone())).await??;
Ok(Self {
stream: Some(stream),
is_initialized: true,
timeout: self.timeout.clone(),
proxy_addr: self.proxy_addr.clone(),
addr: self.addr.clone(),
})
}
Some(proxy_addr) => {
let stream = timeout(
Duration::from_millis(_timeout),
Socks5Stream::connect(
(proxy_addr.0.as_str(), proxy_addr.1),
(addr.0.as_str(), addr.1)
)
).await??;
Ok(Self {
stream: Some(stream.into_inner()),
is_initialized: true,
timeout: self.timeout.clone(),
proxy_addr: self.proxy_addr.clone(),
addr: self.addr.clone(),
})
}
}
}
#[cfg(feature = "resolve")]
{
match &self.proxy_addr {
Some(proxy_addr) => {
let stream = timeout(
Duration::from_millis(_timeout),
Socks5Stream::connect(
(proxy_addr.0.as_str(), proxy_addr.1),
(self.addr.0.as_str(), self.addr.1),
)
).await??;
Ok(Self {
stream: Some(stream.into_inner()),
is_initialized: true,
timeout: self.timeout.clone(),
proxy_addr: self.proxy_addr.clone(),
addr: self.addr.clone(),
})
}
None => {
let host_port = format!("{}:{}", self.addr.0, self.addr.1);
let mut addrs = lookup_host(host_port).await?;
if let Some(sock_addr) = addrs.next() {
let stream = timeout(Duration::from_millis(_timeout), TcpStream::connect(sock_addr)).await??;
Ok(Self {
stream: Some(stream),
is_initialized: true,
timeout: self.timeout.clone(),
proxy_addr: None,
addr: self.addr.clone(),
})
} else {
Err(anyhow!("Could not resolve address: {}", self.addr.0))
}
}
}
}
}
pub fn timeout(mut self, timeout: u64) -> Result<Self> {
if !self.is_initialized {
return Err(anyhow!("using: Connection::new((addr, port)).timeout(u64)"));
}
self.timeout = Some(timeout);
Ok(self)
}
pub fn proxy_socks5(mut self, proxy_addr: (String, u16)) -> Result<Self> {
if !self.is_initialized {
return Err(anyhow!("using: Connection::new((ip, port)).proxy((ip, port))"));
}
self.proxy_addr = Some(proxy_addr);
Ok(self)
}
pub async fn send_handshake(&mut self) -> Result<()> {
let stream = match &mut self.stream {
Some(s) => s,
None => return Err(anyhow!("TCPstream is None. Maybe you forgot to .connect() ?")),
};
let ip = self.addr.0.clone();
let port = self.addr.1;
let handshake = ClientHandshake::new(ip, port);
let bytes = handshake.to_bytes();
timeout(
Duration::from_millis(self.timeout.unwrap_or(9000)),
stream.write_all(bytes.as_slice())
).await??;
Ok(())
}
async fn __send_query_packet(&mut self) -> Result<()> {
let query = StatusQuery::new();
let bytes = query.to_bytes();
let stream = match &mut self.stream {
Some(s) => s,
None => return Err(anyhow!("TCPstream is None. Maybe you forgot to .connect()?")),
};
stream.write_all(bytes.as_slice()).await?;
Ok(())
}
async fn __read_status_packet(&mut self) -> Result<ServerQueryResponse> {
let mut buf = [0u8; 10_000];
let stream = match &mut self.stream {
Some(s) => s,
None => return Err(anyhow!("TCPstream is None. Maybe you forgot to .connect()?")),
};
let n = stream.read(&mut buf).await?;
let status_packet = ServerQueryResponse::from(&buf[..n]).await;
Ok(status_packet)
}
pub async fn get_status(&mut self) -> Result<ServerStatus> {
self.__send_query_packet().await?;
let _status = self.__read_status_packet().await?;
Ok(_status.parse_status()?)
}
pub async fn ping(&mut self) -> Result<ServerStatus> {
self.send_handshake().await?;
self.__send_query_packet().await?;
let status = self.__read_status_packet().await?;
status.parse_status()
}
}