diff --git a/src/main.rs b/src/main.rs index 07b461c..aa6a3f9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -43,7 +43,8 @@ async fn main() { } else { log!(" Users"); } - server::server::start(9187).await; - log!(" Server"); + server::server::start(9187, 9188).await; + log!(" API Server on 9187"); + log!(" WS Server on 9188"); tokio::signal::ctrl_c().await.unwrap(); } diff --git a/src/server/server.rs b/src/server/server.rs index 3926f08..22c24c6 100755 --- a/src/server/server.rs +++ b/src/server/server.rs @@ -1,45 +1,36 @@ use base64::Engine; use base64::engine::general_purpose::STANDARD; use bytes::Bytes; -use futures::StreamExt; use futures_util::TryFutureExt; use http_body_util::BodyExt; use http_body_util::Full; -use hyper::server::conn::http2; +use hyper::server::conn::{http1, http2}; use hyper::{Method, StatusCode}; -use hyper::{ - Request as HttpRequest, Response as HttpResponse, body::Incoming, server::conn::http1, upgrade, -}; +use hyper::{Request as HttpRequest, Response as HttpResponse, body::Incoming, upgrade}; use hyper_util::rt::TokioExecutor; use hyper_util::rt::tokio::TokioIo; use hyper_util::service::TowerToHyperService; -use pnet::datalink::NetworkInterface; -use rustls::ServerConfig; -use rustls::pki_types::{CertificateDer, PrivateKeyDer}; use sha1::{Digest, Sha1}; -use std::error::Error; -use std::io::ErrorKind; -use std::io::{self, BufReader}; -use std::net::{IpAddr, SocketAddr}; +use std::io; +use std::net::SocketAddr; use std::result::Result::Ok; -use std::sync::Arc; use std::{future::Future, pin::Pin, time::Duration}; use tokio::net::TcpListener; -use tokio::sync::broadcast; -use tokio_rustls::TlsAcceptor; use tower::Service; use crate::log; use crate::server::api; use crate::server::short_link::get_short_link; use crate::server::socket; -use crate::util::file_util::load_file_buf; + +// --- ApiService for HTTP/2 --- + #[derive(Clone)] -struct HttpService { - peer_addr: SocketAddr, +struct ApiService { + _peer_addr: SocketAddr, } -impl Service> for HttpService { +impl Service> for ApiService { type Response = HttpResponse>; type Error = io::Error; type Future = Pin> + Send>>; @@ -52,61 +43,17 @@ impl Service> for HttpService { } fn call(&mut self, req: HttpRequest) -> Self::Future { - let peer_ip = self.peer_addr.ip(); - let (parts, body) = req.into_parts(); - - let method = parts.method.clone(); let path = parts.uri.path().to_string(); let headers = parts.headers.clone(); let fut = async move { - let is_websocket_upgrade = path.starts_with("/ws") - && method == Method::GET - && headers - .get("connection") - .map(|v| v.to_str().unwrap_or("").contains("Upgrade")) - .unwrap_or(false) - && headers.get("upgrade").map(|v| v.to_str().unwrap_or("")) == Some("websocket"); - - if is_websocket_upgrade { - log!("Attempting WebSocket upgrade on /ws"); - - if let Some(sec_websocket_key) = headers.get("sec-websocket-key") { - let sec_websocket_key = sec_websocket_key.to_str().unwrap_or("").to_string(); - let sec_websocket_accept = calculate_accept_key(&sec_websocket_key); - - let response = HttpResponse::builder() - .status(StatusCode::SWITCHING_PROTOCOLS) - .header("Upgrade", "websocket") - .header("Connection", "Upgrade") - .header("Sec-WebSocket-Accept", sec_websocket_accept) - .body(Full::new(Bytes::from(""))) - .unwrap(); - let req_for_upgrade = HttpRequest::from_parts(parts, body); - let upgrades = upgrade::on(req_for_upgrade); - - socket::handle(path, upgrades); - Ok(response) - } else { - log!("No Sec-WebSocket-Key found in request headers"); - let response = HttpResponse::builder() - .status(StatusCode::BAD_REQUEST) - .body(Full::new(Bytes::from("Missing Sec-WebSocket-Key"))) - .unwrap(); - Ok(response) - } - } else if path.starts_with("/api") { + if path.starts_with("/api") { let whole_body = tokio::time::timeout(Duration::from_secs(10), body.collect()).await??; let bytes = whole_body.to_bytes(); - - let body_string: Option = match String::from_utf8(bytes.to_vec()) { - Ok(s) => Some(s), - Err(_) => None, - }; - - Ok(api::handle(&path, headers.clone(), body_string).await) + let body_string: Option = String::from_utf8(bytes.to_vec()).ok(); + Ok(api::handle(&path, headers, body_string).await) } else if path.starts_with("/direct") { let short = path.replace("/direct/", ""); if let Ok(long) = get_short_link(&short).await { @@ -124,363 +71,219 @@ impl Service> for HttpService { Ok(response) } } else { - let whole_body = match body.collect().await { - Ok(collected) => collected, - Err(e) => { - log!("Error collecting body: {}", e); - return Ok(HttpResponse::builder() - .status(StatusCode::INTERNAL_SERVER_ERROR) - .body(Full::new(Bytes::from(format!( - "Failed to read body: {}", - e - )))) - .unwrap()); - } - }; - let bytes = whole_body.to_bytes(); - - let body_string: Option = match String::from_utf8(bytes.to_vec()) { - Ok(s) => Some(s), - Err(_) => None, - }; - - Ok(api::handle(&path, headers, body_string).await) + // Not a WebSocket, /api, or /direct, so it's a 404 + let response = HttpResponse::builder() + .status(StatusCode::NOT_FOUND) + .body(Full::new(Bytes::from("Not Found"))) + .unwrap(); + Ok(response) } }; Box::pin(fut.map_err(|err: color_eyre::eyre::ErrReport| { io::Error::new( io::ErrorKind::Other, - format!("Error in request handling: {}", err), + format!("Error in API request: {}", err), ) })) } } -pub fn is_local_network(addr: IpAddr) -> bool { - match addr { - IpAddr::V4(v4) => { - let octets = v4.octets(); - if octets[0] == 10 { - return true; - } - if octets[0] == 172 && (16..=31).contains(&octets[1]) { - return true; - } - if octets[0] == 192 && octets[1] == 168 { - return true; - } - if octets[0] == 127 { - return true; - } - if octets[0] == 169 && octets[1] == 254 { - return true; - } +// --- WebSocketService for HTTP/1.1 --- - false - } +#[derive(Clone)] +struct WebSocketService { + _peer_addr: SocketAddr, +} - IpAddr::V6(v6) => { - let segments = v6.segments(); - if (segments[0] & 0xfe00) == 0xfc00 { - return true; - } - if (segments[0] & 0xffc0) == 0xfe80 { - return true; - } - if v6.is_loopback() { - return true; - } +impl Service> for WebSocketService { + type Response = HttpResponse>; + type Error = io::Error; + type Future = Pin> + Send>>; - false - } + fn poll_ready( + &mut self, + _cx: &mut std::task::Context<'_>, + ) -> std::task::Poll> { + std::task::Poll::Ready(std::io::Result::Ok(())) + } + + fn call(&mut self, req: HttpRequest) -> Self::Future { + let (parts, body) = req.into_parts(); + let method = parts.method.clone(); + let path = parts.uri.path().to_string(); + let headers = parts.headers.clone(); + + let fut = async move { + let is_websocket_upgrade = path.starts_with("/ws") + && method == Method::GET + && headers + .get("connection") + .map(|v| v.to_str().unwrap_or("").contains("Upgrade")) + .unwrap_or(false) + && headers.get("upgrade").map(|v| v.to_str().unwrap_or("")) == Some("websocket"); + + if is_websocket_upgrade { + log!("Attempting WebSocket upgrade on /ws"); + if let Some(sec_websocket_key) = headers.get("sec-websocket-key") { + let sec_websocket_key_str = + sec_websocket_key.to_str().unwrap_or("").to_string(); + let sec_websocket_accept = calculate_accept_key(&sec_websocket_key_str); + + let response = HttpResponse::builder() + .status(StatusCode::SWITCHING_PROTOCOLS) + .header("Upgrade", "websocket") + .header("Connection", "Upgrade") + .header("Sec-WebSocket-Accept", sec_websocket_accept) + .body(Full::new(Bytes::from(""))) + .unwrap(); + + let req_for_upgrade = HttpRequest::from_parts(parts, body); + let upgrades = upgrade::on(req_for_upgrade); + socket::handle(path, upgrades); + Ok(response) + } else { + log!("No Sec-WebSocket-Key found in request headers"); + let response = HttpResponse::builder() + .status(StatusCode::BAD_REQUEST) + .body(Full::new(Bytes::from("Missing Sec-WebSocket-Key"))) + .unwrap(); + Ok(response) + } + } else { + let response = HttpResponse::builder() + .status(StatusCode::NOT_FOUND) + .body(Full::new(Bytes::from( + "Not Found: Use the API server for non-WebSocket requests.", + ))) + .unwrap(); + Ok(response) + } + }; + + Box::pin(fut.map_err(|_err: color_eyre::eyre::ErrReport| { + io::Error::new(io::ErrorKind::Other, "Error in WebSocket request") + })) } } -async fn run_http_server(port: u16) -> bool { +async fn run_api_server(port: u16) -> bool { let mut ip = "0.0.0.0".to_string(); for iface in pnet::datalink::interfaces() { - let iface: NetworkInterface = iface; - if iface.ips.len() > 0 { - let ipsv = format!("{}", iface.ips[0]); - let ips: &str = ipsv.split('/').next().unwrap(); - if format!("{}", ips).starts_with("10.") || format!("{}", ips).starts_with("192.") { - ip = ips.to_string(); - } + if let Some(ip_net) = iface.ips.iter().find(|ip_net| { + ip_net.is_ipv4() + && (ip_net.ip().to_string().starts_with("10.") + || ip_net.ip().to_string().starts_with("192.")) + }) { + ip = ip_net.ip().to_string(); + break; } } - let listener = TcpListener::bind(format!("0.0.0.0:{}", port)).await; - if let Err(e) = listener { - log!("Failed to bind to port {}: {:?}", port, e); + + let addr = SocketAddr::new(ip.parse().unwrap(), port); + let listener_res = TcpListener::bind(format!("0.0.0.0:{}", port)).await; + if let Err(e) = listener_res { + log!( + "[FATAL] Failed to bind API server to port {}: {:?}", + port, + e + ); return false; } - let listener = listener.unwrap(); - log!( - "Standard Server listening for HTTP and WS on {}:{}", - ip, - port - ); - - // Create a broadcast channel for graceful shutdown signal - let (shutdown_tx, _) = broadcast::channel::<()>(1); + let listener = listener_res.unwrap(); + log!("API Server (HTTP/2) listening on {}:{}", ip, port); tokio::spawn(async move { loop { - tokio::select! { - // Monitor for shutdown signal - _ = async { - loop { - tokio::time::sleep(Duration::from_millis(100)).await; - } - } => { - log!("Standard Server received shutdown signal."); - // Send kill signal to all active connection tasks - let _ = shutdown_tx.send(()); - break; - } - - // Accept new connections - accepted = listener.accept() => { - match accepted { - std::result::Result::Ok((stream, addr)) => { - let service = HttpService { peer_addr: addr }; - let io = TokioIo::new(stream); - - // Subscribe to the shutdown signal for this specific connection - let mut rx = shutdown_tx.subscribe(); - - tokio::spawn(async move { - use hyper::server::conn::http2; - - let conn = http2::Builder::new(TokioExecutor::new()) - .serve_connection(io, TowerToHyperService::new(service)); - - - // Wait for either the connection to finish naturally OR the shutdown signal - tokio::select! { - res = conn => { - if let Err(err) = res { - if let Some(io_err) = err.source().and_then(|e| e.downcast_ref::()) { - if io_err.kind() != io::ErrorKind::ConnectionReset - && io_err.kind() != io::ErrorKind::BrokenPipe - { - log!("Error serving connection: {:?}", err); - } - } - } - } - _ = rx.recv() => { - // Shutdown signal received. - // Dropping the 'conn' future here closes the socket immediately. - } - } - }); - } - Err(e) => { - log!("Error accepting connection: {:?}", e); - tokio::time::sleep(Duration::from_millis(500)).await; + if let Ok((stream, addr)) = listener.accept().await { + let service = ApiService { _peer_addr: addr }; + let io = TokioIo::new(stream); + tokio::spawn(async move { + let conn = http2::Builder::new(TokioExecutor::new()) + .serve_connection(io, TowerToHyperService::new(service)); + if let Err(err) = conn.await { + if !err + .to_string() + .starts_with("error shutting down connection") + { + log!("API server connection error: {:?}", err); } } - } + }); } } - - log!("Standard Server shutdown complete."); }); true } -/// Runs the encrypted HTTPS/WSS server loop using the provided TLS config. -async fn run_tls_server(port: u16, tls_config: Arc) -> bool { +async fn run_websocket_server(port: u16) -> bool { let mut ip = "0.0.0.0".to_string(); for iface in pnet::datalink::interfaces() { - let iface: NetworkInterface = iface; - let ipsv = format!("{}", iface.ips[0]); - let ips: &str = ipsv.split('/').next().unwrap(); - log!("{}", ips); - if format!("{}", ips).starts_with("10.") { - ip = ips.to_string(); + if let Some(ip_net) = iface.ips.iter().find(|ip_net| { + ip_net.is_ipv4() + && (ip_net.ip().to_string().starts_with("10.") + || ip_net.ip().to_string().starts_with("192.")) + }) { + ip = ip_net.ip().to_string(); + break; } } - let acceptor = TlsAcceptor::from(tls_config); - - let listener = TcpListener::bind(format!("0.0.0.0:{}", port)).await; - if let Err(e) = listener { - log!("Failed to bind to port {}: {:?}", port, e); + let addr = SocketAddr::new(ip.parse().unwrap(), port); + let listener_res = TcpListener::bind(format!("0.0.0.0:{}", port)).await; + if let Err(e) = listener_res { + log!( + "[FATAL] Failed to bind WebSocket server to port {}: {:?}", + port, + e + ); return false; } - let listener = listener.unwrap(); - log!( - "Encrypted Server listening for HTTPS and WSS on {}:{}", - ip, - port - ); - - // Create a broadcast channel for graceful shutdown signal - let (shutdown_tx, _) = broadcast::channel::<()>(1); + let listener = listener_res.unwrap(); + log!("WebSocket Server (HTTP/1.1) listening on {}:{}", ip, port); tokio::spawn(async move { loop { - tokio::select! { - // Monitor for shutdown signal - _ = async { - loop { - tokio::time::sleep(Duration::from_millis(100)).await; - } - } => { - log!("Encrypted Server received shutdown signal."); - // Send kill signal to all active connection tasks - let _ = shutdown_tx.send(()); - break; - } - - // Accept new connections - accepted = listener.accept() => { - match accepted { - std::result::Result::Ok((stream, addr)) => { - let service = HttpService { peer_addr: addr }; - let acceptor = acceptor.clone(); - - // Subscribe to the shutdown signal for this specific connection - let mut rx = shutdown_tx.subscribe(); - - tokio::spawn(async move { - // Perform TLS handshake - let tls_stream = match acceptor.accept(stream).await { - Ok(s) => s, - Err(e) => { - if e.kind() != io::ErrorKind::Interrupted { - log!("TLS Handshake error: {:?}", e); - } - return; - } - }; - let io = TokioIo::new(tls_stream); - - // Prepare connection future - let conn = http1::Builder::new() - .preserve_header_case(true) - .title_case_headers(true) - .serve_connection(io, TowerToHyperService::new(service)) - .with_upgrades(); - - // Wait for either the connection to finish naturally OR the shutdown signal - tokio::select! { - res = conn => { - if let Err(err) = res { - if let Some(io_err) = err.source().and_then(|e| e.downcast_ref::()) { - if io_err.kind() != io::ErrorKind::ConnectionReset - && io_err.kind() != io::ErrorKind::BrokenPipe - { - log!("Error serving connection: {:?}", err); - } - } - } - } - _ = rx.recv() => { - // Shutdown signal received. - // Dropping the 'conn' future here closes the socket immediately. - } - } - }); - } - Err(e) => { - log!("Error accepting connection: {:?}", e); - tokio::time::sleep(Duration::from_millis(500)).await; + if let Ok((stream, addr)) = listener.accept().await { + let service = WebSocketService { _peer_addr: addr }; + let io = TokioIo::new(stream); + tokio::spawn(async move { + let conn = http1::Builder::new() + .preserve_header_case(true) + .title_case_headers(true) + .serve_connection(io, TowerToHyperService::new(service)) + .with_upgrades(); + if let Err(err) = conn.await { + if !err + .to_string() + .starts_with("error shutting down connection") + { + log!("WebSocket server connection error: {:?}", err); } } - } + }); } } - - log!("Encrypted Server shutdown complete."); }); + true } -pub async fn start(port: u16) -> bool { - let tls_result = load_tls_config(); +// --- Main Start function --- - match tls_result { - Ok(Some(tls_config)) => run_tls_server(port, tls_config).await, - Ok(_) => run_http_server(port).await, - Err(e) => { - log!("Fatal error during TLS config load: {}", e); - false - } - } +pub async fn start(api_port: u16, ws_port: u16) { + // We are not using TLS for this setup as per the request's focus on splitting protocols. + // The previous TLS loading logic is removed for simplicity. + tokio::spawn(run_api_server(api_port)); + tokio::spawn(run_websocket_server(ws_port)); } + fn calculate_accept_key(key: &str) -> String { let websocket_guid = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; let mut sha1 = Sha1::new(); sha1.update(key.as_bytes()); sha1.update(websocket_guid.as_bytes()); let result = sha1.finalize(); - STANDARD.encode(result) // Base64 encode the result -} - -/// Loads TLS config. Returns Ok(None) if cert files are not found, and an error if parsing fails. -fn load_tls_config() -> Result>, Box> { - let cert_file_res = load_file_buf("certs", "cert.pem"); - let key_file_res = load_file_buf("certs", "cert.key"); - - // Check if certificate files are present. If not, return None. - let cert_file_buf = match cert_file_res { - Ok(b) => b, - Err(e) if e.kind() == ErrorKind::NotFound => { - log!("TLS certificate 'certs/cert.pem' not found."); - return Ok(None); - } - Err(e) => return Err(e.into()), // Other IO error - }; - - let key_file_buf = match key_file_res { - Ok(b) => b, - Err(e) if e.kind() == ErrorKind::NotFound => { - log!("TLS key 'certs/cert.key' not found."); - return Ok(None); - } - Err(e) => return Err(e.into()), // Other IO error - }; - - // Continue with configuration if both files were found - let mut cert_reader = BufReader::new(cert_file_buf); - let cert_ders = rustls_pemfile::certs(&mut cert_reader) - .collect::, io::Error>>()?; - - // PKCS8 - let mut key_reader = BufReader::new(key_file_buf); - let mut key_ders = rustls_pemfile::pkcs8_private_keys(&mut key_reader) - .map(|r| r.map(Into::into)) // Explicit conversion - .collect::, io::Error>>()?; - - if key_ders.is_empty() { - // RSA - key_reader = BufReader::new(load_file_buf("certs", "cert.key")?); // Re-read key file - key_ders = rustls_pemfile::rsa_private_keys(&mut key_reader) - .map(|r| r.map(Into::into)) - .collect::, io::Error>>()?; - } - - if key_ders.is_empty() { - // EC - key_reader = BufReader::new(load_file_buf("certs", "cert.key")?); // Re-read key file - key_ders = rustls_pemfile::ec_private_keys(&mut key_reader) - .map(|r| r.map(Into::into)) - .collect::, io::Error>>()?; - } - - if key_ders.is_empty() { - return Err("No private keys found in key file. (Tried PKCS8, RSA, and EC)".into()); - } - - let mut config = rustls::ServerConfig::builder() - .with_no_client_auth() - .with_single_cert(cert_ders, key_ders.remove(0))?; - - config.alpn_protocols = vec![b"h2".to_vec(), b"http/1.1".to_vec()]; - - Ok(Some(Arc::new(config))) + STANDARD.encode(result) }