//! Bounded, loopback-only audit probes for the 2026-09-07 review. //! These ignored tests intentionally confirm the CURRENT vulnerable behavior. //! A passing reproducer is evidence of a finding, not a security regression gate. //! Run: cargo test --locked --test security_audit -- --ignored --nocapture --test-threads=1 use std::{ sync::{ Arc, atomic::{AtomicUsize, Ordering}, }, time::Duration, }; use async_trait::async_trait; use bytes::Bytes; use rustserver::{ Result, fbe::{FbeDecoder, SimpleMessage, SimpleNotify, SimpleRequest, decode_frame}, http::{ HttpClient, HttpRequest, HttpResponse, HttpServer, HttpServerHandler, HttpSession, ParseOutcome, }, protocol::{SimpleProtocolContext, SimpleProtocolHandler, SimpleProtocolServer}, tcp::{TcpOptions, TcpServer}, tls::{TlsServer, TlsServerConfig}, udp::{UdpClient, UdpClientHandler}, }; use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::{TcpListener, TcpStream, UdpSocket}, sync::{Notify, mpsc}, time, }; const DEADLINE: Duration = Duration::from_secs(5); const OBSERVE: Duration = Duration::from_millis(250); #[derive(Default)] struct HttpOk(AtomicUsize); #[async_trait] impl HttpServerHandler for HttpOk { async fn on_request(&self, _: &HttpSession, _: HttpRequest) -> Result> { self.0.fetch_add(1, Ordering::SeqCst); Ok(Some(HttpResponse::ok())) } } #[tokio::test] #[ignore = "audit reproducer: confirms default HTTP pipeline deadlock"] async fn audit_http_pipeline_deadlocks_and_prevents_stop() { let handler = Arc::new(HttpOk::default()); let server = HttpServer::new("127.0.0.1:0", handler.clone()); server.start().await.unwrap(); let mut peer = TcpStream::connect(server.local_addr().await.unwrap()) .await .unwrap(); // 29,700 bytes: fits in the default 64 KiB read chunk, exceeds 1,024 writes. let input = b"GET / HTTP/1.1\r\nHost: x\r\n\r\n".repeat(1_100); peer.write_all(&input).await.unwrap(); let mut response_bytes = Vec::new(); let read = time::timeout(OBSERVE, peer.read_to_end(&mut response_bytes)).await; assert!(read.is_err(), "persistent connection unexpectedly closed"); let processed_requests = handler.0.load(Ordering::SeqCst); assert!( processed_requests > 0 && processed_requests < 1_100, "pipeline did not stall: {processed_requests}" ); let stop = time::timeout(OBSERVE, server.stop()).await; assert!( stop.is_err(), "shutdown unexpectedly escaped callback/queue deadlock" ); println!( "HTTP default pipeline: {processed_requests}/1100 handlers, {} response bytes; stop stalled", response_bytes.len() ); // Runtime teardown cancels the deadlocked tasks; no external listener remains. drop(peer); } #[test] #[ignore = "audit reproducer: confirms unvalidated HTTP serialization"] fn audit_http_crlf_injection() { let request = HttpRequest::get("/first HTTP/1.1\r\nHost: x\r\n\r\nGET /injected"); let bytes = request.to_bytes(); let ParseOutcome::Complete { message, remainder } = HttpRequest::parse(&bytes).unwrap() else { panic!("first injected request was incomplete"); }; assert_eq!(message.target(), "/first"); let ParseOutcome::Complete { message, .. } = HttpRequest::parse(remainder).unwrap() else { panic!("second injected request was incomplete"); }; assert_eq!(message.target(), "/injected"); let response = HttpResponse::ok().with_header("X-User", "name\r\nSet-Cookie: audit=injected"); let bytes = response.to_bytes(); let ParseOutcome::Complete { message, .. } = HttpResponse::parse(&bytes).unwrap() else { panic!("injected response was incomplete"); }; assert_eq!(message.header("Set-Cookie"), Some("audit=injected")); println!("One request serialized as two requests; CRLF in header created Set-Cookie"); } #[cfg(unix)] #[test] #[ignore = "audit reproducer: confirms ancestor symlink escape during refresh"] fn audit_static_refresh_follows_replaced_ancestor() { use rustserver::http::StaticFileCache; use std::{fs, os::unix::fs::symlink}; let base = std::env::temp_dir().join(format!("rustserver-audit-{}", uuid::Uuid::new_v4())); let public = base.join("public"); let nested = public.join("nested"); let outside = base.join("outside"); fs::create_dir_all(&nested).unwrap(); fs::create_dir_all(&outside).unwrap(); fs::write(nested.join("value.txt"), b"public fixture").unwrap(); fs::write(outside.join("value.txt"), b"private fixture outside mount").unwrap(); let cache = StaticFileCache::new(); cache .add_static_content(&public, "/static", Duration::ZERO) .unwrap(); fs::rename(&nested, public.join("original")).unwrap(); symlink(&outside, &nested).unwrap(); cache.watchdog().unwrap(); let response = cache .find(&HttpRequest::get("/static/nested/value.txt")) .unwrap(); assert_eq!(response.body(), b"private fixture outside mount"); fs::remove_dir_all(&base).unwrap(); println!( "Static cache served outside-mount fixture after parent directory was replaced with symlink" ); } #[derive(Default)] struct GatedProtocol { entered: Notify, release: Notify, processed: AtomicUsize, } #[async_trait] impl SimpleProtocolHandler for GatedProtocol { async fn on_connected(&self, _: SimpleProtocolContext) { self.entered.notify_one(); self.release.notified().await; } async fn on_notify(&self, _: SimpleProtocolContext, _: SimpleNotify) { self.processed.fetch_add(1, Ordering::SeqCst); } } #[tokio::test] #[ignore = "audit reproducer: confirms protocol queue consumes while handler is blocked"] async fn audit_protocol_queue_has_no_backpressure() { let handler = Arc::new(GatedProtocol::default()); let server = SimpleProtocolServer::new("127.0.0.1:0").with_handler(handler.clone()); server.start().await.unwrap(); let mut peer = TcpStream::connect(server.local_addr().await.unwrap()) .await .unwrap(); time::timeout(DEADLINE, handler.entered.notified()) .await .unwrap(); let frame = SimpleMessage::Notify(SimpleNotify::new(vec![b'a'; 4_096])) .encode() .unwrap(); let count = 2_048; time::timeout(DEADLINE, async { for _ in 0..count { peer.write_all(&frame).await.unwrap(); } while server.stats().bytes_received() < (frame.len() * count) as u64 { time::sleep(Duration::from_millis(2)).await; } }) .await .unwrap(); assert_eq!(handler.processed.load(Ordering::SeqCst), 0); println!( "Protocol queue: {} received bytes / {count} frames, zero callbacks processed", frame.len() * count ); handler.release.notify_one(); time::timeout(DEADLINE, async { while handler.processed.load(Ordering::SeqCst) < count { time::sleep(Duration::from_millis(2)).await; } }) .await .unwrap(); drop(peer); time::timeout(DEADLINE, server.stop()) .await .unwrap() .unwrap(); } #[tokio::test] #[ignore = "audit reproducer: confirms TCP write does not observe disconnect"] async fn audit_tcp_slow_reader_blocks_shutdown() { let server = TcpServer::new("127.0.0.1:0").with_options(TcpOptions { send_buffer_size: Some(4_096), ..TcpOptions::default() }); server.start().await.unwrap(); let peer = TcpStream::connect(server.local_addr().await.unwrap()) .await .unwrap(); time::timeout(DEADLINE, async { while server.connected_sessions().await == 0 { tokio::task::yield_now().await; } }) .await .unwrap(); let session = server.sessions().await.pop().unwrap(); session.wait_connected().await; session .send(Bytes::from(vec![b'x'; 8 * 1024 * 1024])) .await .unwrap(); time::sleep(OBSERVE).await; assert!(session.stats().bytes_pending() > 0); let mut stop = tokio::spawn(async move { server.stop().await }); assert!(time::timeout(OBSERVE, &mut stop).await.is_err()); println!("TCP stop remained pending until non-reading peer was closed"); drop(peer); time::timeout(DEADLINE, stop) .await .unwrap() .unwrap() .unwrap(); } #[tokio::test] #[ignore = "audit reproducer: confirms HTTP request timeout excludes request write"] async fn audit_http_request_timeout_does_not_cover_write() { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let client = HttpClient::with_tcp_options( listener.local_addr().unwrap().to_string(), TcpOptions { send_buffer_size: Some(4_096), ..TcpOptions::default() }, ); client.connect().await.unwrap(); let (peer, _) = listener.accept().await.unwrap(); let request = HttpRequest::post("/", vec![b'x'; 8 * 1024 * 1024]); let declared = Duration::from_millis(25); let mut pending = Box::pin(client.request_with_timeout(request, declared)); assert!(time::timeout(OBSERVE, &mut pending).await.is_err()); println!("HTTP request configured for 25ms remained pending after 250ms during write"); drop(peer); let _ = time::timeout(DEADLINE, &mut pending).await.unwrap(); if client.is_connected().await { let _ = time::timeout(DEADLINE, client.disconnect()).await.unwrap(); } } struct UdpCollector(mpsc::Sender<(std::net::SocketAddr, Bytes)>); #[async_trait] impl UdpClientHandler for UdpCollector { async fn on_received(&self, _: &UdpClient, peer: std::net::SocketAddr, data: Bytes) { self.0.send((peer, data)).await.unwrap(); } } #[tokio::test] #[ignore = "audit reproducer: confirms unicast client delivers unrelated sender data"] async fn audit_udp_client_accepts_unrelated_sender() { let intended = UdpSocket::bind("127.0.0.1:0").await.unwrap(); let unrelated = UdpSocket::bind("127.0.0.1:0").await.unwrap(); let (tx, mut rx) = mpsc::channel(1); let client = UdpClient::new(intended.local_addr().unwrap().to_string()) .with_handler(Arc::new(UdpCollector(tx))); client.connect().await.unwrap(); let port = client.local_addr().await.unwrap().port(); unrelated .send_to(b"forged fixture", (std::net::Ipv4Addr::LOCALHOST, port)) .await .unwrap(); let (sender, data) = time::timeout(DEADLINE, rx.recv()).await.unwrap().unwrap(); assert_eq!(sender, unrelated.local_addr().unwrap()); assert_eq!(data, b"forged fixture"[..]); client.disconnect().await.unwrap(); println!("Unicast UDP callback received a datagram from a socket other than configured remote"); } #[test] #[ignore = "audit reproducer: confirms bodyless-status response boundary error"] fn audit_http_304_consumes_next_response_as_body() { let next = b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok"; let mut input = format!( "HTTP/1.1 304 Not Modified\r\nContent-Length: {}\r\n\r\n", next.len() ) .into_bytes(); input.extend_from_slice(next); let ParseOutcome::Complete { message, remainder } = HttpResponse::parse(&input).unwrap() else { panic!("response was incomplete"); }; assert_eq!(message.status(), 304); assert_eq!(message.body(), next); assert!(remainder.is_empty()); println!("304 response incorrectly consumed entire following 200 response as its body"); } #[test] #[ignore = "audit reproducer: confirms ambiguous Transfer-Encoding is accepted"] fn audit_http_accepts_transfer_encoding_with_content_length() { let input = b"POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: identity\r\nContent-Length: 4\r\n\r\nbody"; let ParseOutcome::Complete { message, .. } = HttpRequest::parse(input).unwrap() else { panic!("ambiguous request was incomplete"); }; assert_eq!(message.body(), b"body"); println!("Request containing both Transfer-Encoding and Content-Length was accepted"); } #[tokio::test] #[ignore = "audit reproducer: confirms pending TLS handshakes have no timeout or count limit"] async fn audit_tls_silent_handshakes_accumulate() { let directory = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("tools/certificates"); let config = TlsServerConfig::from_pem_files(directory.join("server.crt"), directory.join("server.key")) .unwrap(); let server = TlsServer::new("127.0.0.1:0", config); server.start().await.unwrap(); let address = server.local_addr().await.unwrap(); let count = 128; let mut silent_peers = Vec::with_capacity(count); for _ in 0..count { silent_peers.push(TcpStream::connect(address).await.unwrap()); } time::timeout(DEADLINE, async { while server.connected_sessions().await < count { tokio::task::yield_now().await; } }) .await .unwrap(); assert_eq!(server.connected_sessions().await, count); time::sleep(OBSERVE).await; assert_eq!(server.connected_sessions().await, count); println!("{count} silent TCP peers remained registered as pending TLS sessions"); time::timeout(DEADLINE, server.stop()) .await .unwrap() .unwrap(); drop(silent_peers); } #[test] #[ignore = "bounded deterministic parser mutation check; not a coverage-guided fuzz campaign"] fn audit_parser_mutations_do_not_panic() { let frame = SimpleMessage::Request(SimpleRequest::new(uuid::Uuid::nil(), vec![b'x'; 256])) .encode() .unwrap(); let request = HttpRequest::post("/audit", b"body").to_bytes(); let response = HttpResponse::get(b"body").to_bytes(); let mut state = 0x10fe_04ca_1133_7788_u64; for index in 0..20_000 { state ^= state << 13; state ^= state >> 7; state ^= state << 17; let seed = match index % 3 { 0 => &frame, 1 => &request, _ => &response, }; let mut input = seed.clone(); let position = usize::try_from(state % u64::try_from(input.len()).unwrap()).unwrap(); input[position] = state.to_le_bytes()[3]; if index % 5 == 0 { input.truncate(position); } let _ = decode_frame(&input); let mut decoder = FbeDecoder::new(); let _ = decoder.push(&input); let _ = HttpRequest::parse(&input); let _ = HttpResponse::parse(&input); } println!("20,000 deterministic mutations completed without parser panic"); }