//! Streaming, cancellation, size-bound, and timeout semantics of the HTTP //! engine against a real loopback socket (`SMH-SPEC-SPEC0001`, network transport). use smith::http::{ HttpBodyEnd, HttpBodyItem, HttpExchange, HttpExecutor, HttpMethod, HttpResponse, }; use smith::tool::CancelHandle; use smith_harness::http::{MAX_BODY_BYTES, UreqHttpExecutor}; use std::io::{Read, Write}; use std::net::{TcpListener, TcpStream}; use std::sync::mpsc::{self, Receiver, RecvTimeoutError, SyncSender}; use std::time::{Duration, Instant}; const EVENTS: [&str; 3] = ["data: one\n\n", "data: two\n\n", "data: three\n\n"]; /// Upper bound on any single wait for the fixture or the engine. const DEADLINE: Duration = Duration::from_secs(10); /// Accept one connection, read the request head, hand the socket to /// `serve`, then drain until the client closes so unread request bytes /// never turn the close into a reset that discards in-flight response data. fn spawn_server(serve: impl FnOnce(&mut TcpStream) + Send + 'static) -> u16 { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let port = listener.local_addr().unwrap().port(); std::thread::spawn(move || { let (mut socket, _) = listener.accept().unwrap(); let mut buffer = [0u8; 8192]; let _ = socket.read(&mut buffer); serve(&mut socket); let _ = socket.shutdown(std::net::Shutdown::Write); let _ = std::io::copy(&mut socket, &mut std::io::sink()); }); port } /// Chunked SSE: head plus first event at once; each later event only once /// the test sends on the returned gate, so the client must stream what it has. fn serve_events() -> (SyncSender<()>, impl FnOnce(&mut TcpStream) + Send + 'static) { let (gate, next) = mpsc::sync_channel::<()>(EVENTS.len()); let serve = move |socket: &mut TcpStream| { let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\n\ Transfer-Encoding: chunked\r\nConnection: close\r\n\r\n"; let _ = socket.write_all(head.as_bytes()); for (index, event) in EVENTS.iter().enumerate() { // synchronize: the test releases each later event. if index > 0 && next.recv_timeout(DEADLINE).is_err() { return; } let _ = write!(socket, "{:x}\r\n{event}\r\n", event.len()); let _ = socket.flush(); } let _ = socket.write_all(b"0\r\n\r\n"); }; (gate, serve) } /// Hold a fixture until the test is done with it, bounded by `DEADLINE`. fn stall(release: &Receiver<()>) { // synchronize: the stall lasts until the test releases it. let _ = release.recv_timeout(DEADLINE); } /// The next body item, or a failure once `DEADLINE` passes. fn next_item(response: &HttpResponse) -> HttpBodyItem { let item = response.body.recv_timeout(DEADLINE); assert!(item.is_ok(), "body stalled: {item:?}"); item.unwrap() } fn chunk(event: &str) -> HttpBodyItem { HttpBodyItem::Chunk(smith::http::HttpChunk { bytes: event.as_bytes().to_vec(), }) } /// Every chunk, then how the body ended. fn drain(response: HttpResponse) -> (Vec>, Option) { let mut chunks = Vec::new(); for item in response.body { match item { HttpBodyItem::Chunk(chunk) => chunks.push(chunk.bytes), HttpBodyItem::End(end) => return (chunks, Some(end)), } } (chunks, None) } fn post(executor: &dyn HttpExecutor, port: u16, cancel: CancelHandle) -> HttpResponse { let exchange = HttpExchange { method: HttpMethod::Post, url: format!("http://127.0.0.1:{port}/v1/chat/completions"), headers: Vec::new(), body: b"{}".to_vec(), }; executor.execute(exchange, cancel).unwrap() } #[test] fn chunks_arrive_incrementally() { let executor = UreqHttpExecutor::new(); let (gate, serve) = serve_events(); let port = spawn_server(serve); let response = post(&executor, port, CancelHandle::new()); assert_eq!(response.head.status, 200); // Each event arrives before the server may send the next: a buffering // client would stall here instead. for (index, event) in EVENTS.iter().enumerate() { if index > 0 { gate.send(()).unwrap(); } assert_eq!(next_item(&response), chunk(event)); } assert_eq!( next_item(&response), HttpBodyItem::End(HttpBodyEnd::Complete) ); } #[test] fn cancel_after_first_chunk_stops_the_reader() { let executor = UreqHttpExecutor::new(); let (gate, serve) = serve_events(); let port = spawn_server(serve); let cancel = CancelHandle::new(); let response = post(&executor, port, cancel.clone()); assert_eq!(next_item(&response), chunk(EVENTS[0])); cancel.cancel(); // The reader delivers nothing read after cancellation and ends the body // once its in-flight read returns, which the next event forces. gate.send(()).unwrap(); assert_eq!( next_item(&response), HttpBodyItem::End(HttpBodyEnd::Cancelled) ); assert_eq!( response.body.recv_timeout(DEADLINE), Err(RecvTimeoutError::Disconnected) ); } #[test] fn body_over_bound_is_cut_off() { let oversized = MAX_BODY_BYTES + 1024 * 1024; let executor = UreqHttpExecutor::new(); let port = spawn_server(move |socket| { let head = format!("HTTP/1.1 200 OK\r\nContent-Length: {oversized}\r\n\r\n"); let _ = socket.write_all(head.as_bytes()); let block = vec![b'x'; 64 * 1024]; let mut sent = 0; while sent < oversized && socket.write_all(&block).is_ok() { sent += block.len() as u64; } }); let response = post(&executor, port, CancelHandle::new()); let (chunks, end) = drain(response); let received: u64 = chunks.iter().map(|b| b.len() as u64).sum(); assert_eq!( end, Some(HttpBodyEnd::Truncated { limit: MAX_BODY_BYTES }) ); assert!(received <= MAX_BODY_BYTES, "received {received}"); assert!(received > MAX_BODY_BYTES - 64 * 1024, "received {received}"); } #[test] fn ureq_gap_timeout_ends_a_stalled_body() { let gap = Duration::from_millis(300); let executor = UreqHttpExecutor::with_timeouts(DEADLINE, gap); let (release, stalled) = mpsc::sync_channel::<()>(1); let port = spawn_server(move |socket| { let head = "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n"; let _ = socket.write_all(head.as_bytes()); // Keep streaming for longer than `gap` in total, then stall once. for _ in 0..3 { let _ = socket.write_all(b"2\r\nok\r\n"); // keep-with-deadline: fixture pacing below `gap`; sleeps last at least this long. std::thread::sleep(gap / 3); } let _ = socket.write_all(b"2\r\nok\r\n"); stall(&stalled); let _ = socket.write_all(b"5\r\nlate!\r\n0\r\n\r\n"); }); let start = Instant::now(); let response = post(&executor, port, CancelHandle::new()); let (chunks, end) = drain(response); let ended = start.elapsed(); let _ = release.send(()); let body: Vec = chunks.into_iter().flatten().collect(); assert_eq!(body, b"okokokok"); assert!( matches!(end, Some(HttpBodyEnd::Failed { .. })), "stalled body ended with {end:?}" ); // keep-with-deadline: time is the contract; the stream outlived `gap` // in total, and the stall failed long before its release. assert!(ended >= gap * 2, "ended at {ended:?}"); assert!(ended < DEADLINE / 2, "ended at {ended:?}"); } #[test] fn ureq_head_timeout_fails_the_exchange() { let head = Duration::from_millis(300); let executor = UreqHttpExecutor::with_timeouts(head, DEADLINE); let (release, stalled) = mpsc::sync_channel::<()>(1); let port = spawn_server(move |_| stall(&stalled)); let exchange = HttpExchange { method: HttpMethod::Post, url: format!("http://127.0.0.1:{port}/"), headers: Vec::new(), body: Vec::new(), }; let start = Instant::now(); let result = executor.execute(exchange, CancelHandle::new()); let failed = start.elapsed(); let _ = release.send(()); assert!(result.is_err()); // keep-with-deadline: time is the contract; failed at `head`, long before the stall's release. assert!(failed < DEADLINE / 2, "failed at {failed:?}"); }