1use crate::error_page;
5use crate::headers::Headers;
6use crate::method::Method;
7use crate::panic;
8use crate::request::Request;
9use crate::response::Response;
10use crate::router::Router;
11use crate::status::Status;
12use crate::url;
13use rustlavel_core::{Context, Error, Result};
14use std::net::SocketAddr;
15use std::sync::Arc;
16use std::sync::atomic::{AtomicUsize, Ordering};
17use std::time::{Duration, Instant};
18use tokio::io::{AsyncReadExt, AsyncWriteExt, BufWriter};
19use tokio::net::{TcpListener, TcpStream};
20
21#[derive(Debug, Clone)]
23pub struct Limits {
24 pub max_header_bytes: usize,
25 pub max_body_bytes: usize,
26 pub keep_alive_timeout: Duration,
28 pub header_timeout: Duration,
30}
31
32impl Default for Limits {
33 fn default() -> Self {
34 Limits {
35 max_header_bytes: 64 * 1024,
36 max_body_bytes: 10 * 1024 * 1024,
37 keep_alive_timeout: Duration::from_secs(15),
38 header_timeout: Duration::from_secs(10),
39 }
40 }
41}
42
43fn port_attempts(config: &rustlavel_core::Config) -> u16 {
54 let default = if config.is_production() { 1 } else { 10 };
55 config.int("server.port_attempts", default).clamp(1, 1000) as u16
56}
57
58fn check_framing(headers: &Headers) -> Result<()> {
70 let lengths = headers.get_all("content-length");
82 if lengths.len() > 1 && lengths.iter().any(|value| value != &lengths[0]) {
83 return Err(Error::Protocol(
84 "more than one Content-Length, and they disagree".into(),
85 ));
86 }
87 if !lengths.is_empty() && headers.get("transfer-encoding").is_some() {
88 return Err(Error::Protocol(
89 "both Transfer-Encoding and Content-Length: a message may say where it ends \
90 once, not twice"
91 .into(),
92 ));
93 }
94 let encodings = headers.get_all("transfer-encoding");
99 if !encodings.is_empty() {
100 let listed: Vec<&str> = encodings
101 .iter()
102 .flat_map(|value| value.split(','))
103 .map(str::trim)
104 .filter(|value| !value.is_empty())
105 .collect();
106 if listed.last() != Some(&"chunked")
107 || listed.iter().filter(|value| **value == "chunked").count() != 1
108 {
109 return Err(Error::Protocol(format!(
110 "unsupported Transfer-Encoding: {}",
111 listed.join(", ")
112 )));
113 }
114 }
115 Ok(())
116}
117
118async fn bind_walking(addr: &str, attempts: u16) -> Result<TcpListener> {
135 let Some((host, first)) = addr
138 .rsplit_once(':')
139 .and_then(|(host, port)| port.parse::<u16>().ok().map(|port| (host, port)))
140 .filter(|(_, port)| *port != 0)
141 else {
142 return TcpListener::bind(addr).await.map_err(Error::Io);
143 };
144
145 let mut tried = first;
146 for offset in 0..attempts {
147 let Some(port) = first.checked_add(offset) else { break };
148 tried = port;
149 match TcpListener::bind(format!("{host}:{port}")).await {
150 Ok(listener) => {
151 if offset > 0 {
152 rustlavel_core::warn!(
153 "port {first} is in use, so this is serving on {port} instead"
154 );
155 }
156 return Ok(listener);
157 }
158 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => continue,
159 Err(e) => return Err(Error::Io(e)),
160 }
161 }
162
163 Err(Error::msg(if first == tried {
164 format!(
165 "port {first} is already in use. Something else is listening on it — `lsof -i :{first}` \
166 says what — so stop that, or set SERVER_PORT to a free port."
167 )
168 } else {
169 format!(
170 "every port from {first} to {tried} is already in use. Stop whatever is holding them \
171 — `lsof -i :{first}` names the first — or set SERVER_PORT to a free one."
172 )
173 }))
174}
175
176#[cfg(unix)]
186async fn stop_requested() {
187 use tokio::signal::unix::{SignalKind, signal};
188
189 let mut terminate = match signal(SignalKind::terminate()) {
190 Ok(stream) => stream,
191 Err(error) => {
194 rustlavel_core::warn!("cannot listen for SIGTERM: {error}");
195 let _ = tokio::signal::ctrl_c().await;
196 return;
197 }
198 };
199
200 tokio::select! {
201 _ = tokio::signal::ctrl_c() => {}
202 _ = terminate.recv() => {}
203 }
204}
205
206#[cfg(not(unix))]
207async fn stop_requested() {
208 let _ = tokio::signal::ctrl_c().await;
209}
210
211pub type OnShutdown = Box<dyn FnOnce() -> crate::handler::BoxFuture<()> + Send + Sync>;
217
218pub struct Server {
219 router: Arc<Router>,
220 context: Context,
221 limits: Limits,
222 on_shutdown: Vec<OnShutdown>,
230}
231
232impl Server {
233 pub fn new(mut router: Router, context: Context) -> Self {
234 router.finalize();
235 let limits = Limits {
236 max_body_bytes: context.config().int("server.max_body_bytes", 10 * 1024 * 1024) as usize,
237 ..Limits::default()
238 };
239 Server { router: Arc::new(router), context, limits, on_shutdown: Vec::new() }
240 }
241
242 pub fn on_shutdown(mut self, work: OnShutdown) -> Self {
249 self.on_shutdown.push(work);
250 self
251 }
252
253 pub fn limits(mut self, limits: Limits) -> Self {
254 self.limits = limits;
255 self
256 }
257
258 pub async fn listen(mut self, addr: impl Into<String>) -> Result<()> {
264 let addr = addr.into();
265 let listener = bind_walking(&addr, port_attempts(self.context.config())).await?;
266 let local = listener.local_addr().map_err(Error::Io)?;
267
268 panic::install_hook();
269 error_page::set_debug(self.context.config().debug());
270
271 rustlavel_core::info!("Rustlavel serving on http://{local}");
272 rustlavel_core::info!("Press Ctrl-C to stop");
273
274 let in_flight = Arc::new(AtomicUsize::new(0));
275 let mut on_shutdown = std::mem::take(&mut self.on_shutdown);
278 let shared = Arc::new(self);
279
280 loop {
281 let accepted = tokio::select! {
282 result = listener.accept() => result,
283 _ = stop_requested() => break,
284 };
285
286 let (stream, peer) = match accepted {
287 Ok(pair) => pair,
288 Err(e) => {
291 rustlavel_core::warn!("accept failed: {e}");
292 continue;
293 }
294 };
295
296 let server = Arc::clone(&shared);
297 let counter = Arc::clone(&in_flight);
298 counter.fetch_add(1, Ordering::SeqCst);
299 tokio::spawn(async move {
300 if let Err(e) = server.serve_connection(stream, peer).await {
301 rustlavel_core::debug!("connection closed: {e}");
302 }
303 counter.fetch_sub(1, Ordering::SeqCst);
304 });
305 }
306
307 rustlavel_core::info!("Shutting down, waiting for in-flight requests…");
308 let deadline = Instant::now() + Duration::from_secs(10);
309 while in_flight.load(Ordering::SeqCst) > 0 && Instant::now() < deadline {
310 tokio::time::sleep(Duration::from_millis(25)).await;
311 }
312 for work in on_shutdown.drain(..) {
315 work().await;
316 }
317
318 rustlavel_core::info!("Goodbye.");
319 Ok(())
320 }
321
322 pub(crate) async fn serve_connection(&self, stream: TcpStream, peer: SocketAddr) -> Result<()> {
323 let _ = stream.set_nodelay(true);
325 let (mut reader, writer) = stream.into_split();
326 let mut writer = BufWriter::new(writer);
327 let mut buffer: Vec<u8> = Vec::with_capacity(2048);
328
329 loop {
330 let head = match self.read_head(&mut reader, &mut buffer).await? {
331 Some(head) => head,
332 None => return Ok(()),
334 };
335
336 let (mut request, keep_alive) = match self.parse(&head, &mut reader, &mut buffer, peer).await {
337 Ok(parsed) => parsed,
338 Err(error) => {
339 let response = Response::new(Status::BAD_REQUEST).with_text(error.to_string());
340 writer.write_all(&response.to_bytes(true)).await.map_err(Error::Io)?;
341 writer.flush().await.map_err(Error::Io)?;
342 return Ok(());
343 }
344 };
345
346 request.context = self.context.clone();
347 let is_head = request.method() == Method::Head;
348 let mut response = self.dispatch(request).await;
349
350 if response.upgrades() {
353 let head = response.to_bytes(false);
360 let upgrade = response.take_upgrade().expect("checked just above");
361 writer.write_all(&head).await.map_err(Error::Io)?;
362 writer.flush().await.map_err(Error::Io)?;
363
364 let upgraded = crate::upgrade::Upgraded {
365 reader: Box::new(reader),
366 writer: Box::new(writer),
367 buffered: std::mem::take(&mut buffer),
370 };
371 upgrade.run(upgraded).await;
372 return Ok(());
373 }
374
375 if !keep_alive {
376 response.headers.set("connection", "close");
377 }
378 writer.write_all(&response.to_bytes(!is_head)).await.map_err(Error::Io)?;
379 writer.flush().await.map_err(Error::Io)?;
380
381 if !keep_alive {
382 return Ok(());
383 }
384 }
385 }
386
387 async fn dispatch(&self, request: Request) -> Response {
390 let started = Instant::now();
391 let method = request.method();
392 let path = request.path().to_string();
393
394 let response = self.router.dispatch(request).await;
397 let elapsed = started.elapsed();
398
399 if rustlavel_core::log::enabled(rustlavel_core::log::Level::Debug) {
400 rustlavel_core::debug!(
401 "{method} {path} → {} ({:.1}ms)",
402 response.status.code(),
403 elapsed.as_secs_f64() * 1000.0
404 );
405 }
406
407 response
408 }
409
410 async fn read_head(
412 &self,
413 reader: &mut tokio::net::tcp::OwnedReadHalf,
414 buffer: &mut Vec<u8>,
415 ) -> Result<Option<Vec<u8>>> {
416 let mut timeout = self.limits.keep_alive_timeout;
419
420 loop {
421 if let Some(end) = find_head_end(buffer) {
422 let head = buffer[..end].to_vec();
423 buffer.drain(..end);
424 return Ok(Some(head));
425 }
426 if buffer.len() > self.limits.max_header_bytes {
427 return Err(Error::Protocol("request headers are too large".into()));
428 }
429
430 let mut chunk = [0u8; 4096];
431 let read = match tokio::time::timeout(timeout, reader.read(&mut chunk)).await {
432 Ok(Ok(0)) if buffer.is_empty() => return Ok(None),
433 Ok(Ok(0)) => return Err(Error::Protocol("connection closed mid-request".into())),
434 Ok(Ok(n)) => n,
435 Ok(Err(e)) => return Err(Error::Io(e)),
436 Err(_) if buffer.is_empty() => return Ok(None),
437 Err(_) => return Err(Error::Protocol("timed out reading request headers".into())),
438 };
439 buffer.extend_from_slice(&chunk[..read]);
440 timeout = self.limits.header_timeout;
441 }
442 }
443
444 async fn parse(
445 &self,
446 head: &[u8],
447 reader: &mut tokio::net::tcp::OwnedReadHalf,
448 buffer: &mut Vec<u8>,
449 peer: SocketAddr,
450 ) -> Result<(Request, bool)> {
451 let text = std::str::from_utf8(head).map_err(|_| Error::Protocol("headers are not UTF-8".into()))?;
452 let mut lines = text.split("\r\n");
453
454 let request_line = lines.next().ok_or_else(|| Error::Protocol("empty request".into()))?;
455 let mut parts = request_line.split(' ');
456 let method = parts
457 .next()
458 .and_then(Method::parse)
459 .ok_or_else(|| Error::Protocol("unsupported method".into()))?;
460 let target = parts.next().ok_or_else(|| Error::Protocol("missing request target".into()))?;
461 let version = parts.next().unwrap_or("HTTP/1.1");
462
463 let mut headers = Headers::new();
464 for line in lines {
465 if line.is_empty() {
466 continue;
467 }
468 let (name, value) = line
469 .split_once(':')
470 .ok_or_else(|| Error::Protocol(format!("malformed header line: {line}")))?;
471
472 if name.ends_with(' ') || name.ends_with('\t') {
478 return Err(Error::Protocol(
479 "a header name may not be followed by whitespace before the colon".into(),
480 ));
481 }
482 headers.append(name.trim(), value.trim());
483 }
484
485 check_framing(&headers)?;
486
487 let target = match target.find("://") {
489 Some(scheme_end) => match target[scheme_end + 3..].find('/') {
490 Some(path_start) => &target[scheme_end + 3 + path_start..],
491 None => "/",
492 },
493 None => target,
494 };
495
496 let body = self.read_body(&headers, reader, buffer).await?;
497
498 let keep_alive = match headers.get("connection") {
499 Some(value) if value.eq_ignore_ascii_case("close") => false,
500 Some(value) if value.eq_ignore_ascii_case("keep-alive") => true,
501 _ => version != "HTTP/1.0",
502 };
503
504 let (path, query) = url::split_target(target);
505 let mut request = Request::new(method, target);
506 request.path = url::decode(path);
507 request.query = url::parse_query(query);
508 request.headers = headers;
509 request.peer = Some(peer);
510 Ok((request.with_body(body), keep_alive))
511 }
512
513 async fn read_body(
514 &self,
515 headers: &Headers,
516 reader: &mut tokio::net::tcp::OwnedReadHalf,
517 buffer: &mut Vec<u8>,
518 ) -> Result<Vec<u8>> {
519 if headers.get("transfer-encoding").is_some() {
523 return self.read_chunked_body(reader, buffer).await;
524 }
525
526 let Some(length) = headers.content_length() else {
527 return Ok(Vec::new());
528 };
529 if length > self.limits.max_body_bytes {
530 return Err(Error::Protocol("request body is too large".into()));
531 }
532
533 while buffer.len() < length {
534 let mut chunk = vec![0u8; (length - buffer.len()).min(64 * 1024)];
535 let read = tokio::time::timeout(self.limits.header_timeout, reader.read(&mut chunk))
536 .await
537 .map_err(|_| Error::Protocol("timed out reading request body".into()))?
538 .map_err(Error::Io)?;
539 if read == 0 {
540 return Err(Error::Protocol("request body ended early".into()));
541 }
542 buffer.extend_from_slice(&chunk[..read]);
543 }
544
545 Ok(buffer.drain(..length).collect())
546 }
547
548 async fn read_chunked_body(
549 &self,
550 reader: &mut tokio::net::tcp::OwnedReadHalf,
551 buffer: &mut Vec<u8>,
552 ) -> Result<Vec<u8>> {
553 let mut body = Vec::new();
554
555 loop {
556 let line_end = loop {
558 if let Some(at) = find_crlf(buffer) {
559 break at;
560 }
561 if !fill(reader, buffer, self.limits.header_timeout).await? {
562 return Err(Error::Protocol("chunked body ended early".into()));
563 }
564 };
565
566 let header: Vec<u8> = buffer.drain(..line_end + 2).collect();
567 let size_text = String::from_utf8_lossy(&header[..line_end]);
568 let size = usize::from_str_radix(size_text.split(';').next().unwrap_or("").trim(), 16)
569 .map_err(|_| Error::Protocol("invalid chunk size".into()))?;
570
571 if size == 0 {
572 loop {
575 let end = loop {
576 if let Some(at) = find_crlf(buffer) {
577 break at;
578 }
579 if !fill(reader, buffer, self.limits.header_timeout).await? {
580 return Ok(body);
581 }
582 };
583 buffer.drain(..end + 2);
584 if end == 0 {
585 return Ok(body);
586 }
587 }
588 }
589
590 if body.len() + size > self.limits.max_body_bytes {
591 return Err(Error::Protocol("request body is too large".into()));
592 }
593
594 while buffer.len() < size + 2 {
595 if !fill(reader, buffer, self.limits.header_timeout).await? {
596 return Err(Error::Protocol("chunked body ended early".into()));
597 }
598 }
599 body.extend(buffer.drain(..size));
600 buffer.drain(..2);
601 }
602 }
603}
604
605async fn fill(
606 reader: &mut tokio::net::tcp::OwnedReadHalf,
607 buffer: &mut Vec<u8>,
608 timeout: Duration,
609) -> Result<bool> {
610 let mut chunk = [0u8; 4096];
611 let read = tokio::time::timeout(timeout, reader.read(&mut chunk))
612 .await
613 .map_err(|_| Error::Protocol("timed out reading request body".into()))?
614 .map_err(Error::Io)?;
615 buffer.extend_from_slice(&chunk[..read]);
616 Ok(read > 0)
617}
618
619fn find_head_end(buffer: &[u8]) -> Option<usize> {
621 buffer.windows(4).position(|w| w == b"\r\n\r\n").map(|at| at + 4)
622}
623
624fn find_crlf(buffer: &[u8]) -> Option<usize> {
625 buffer.windows(2).position(|w| w == b"\r\n")
626}
627
628#[cfg(test)]
629mod tests {
630
631
632 fn framing_of(raw: &str) -> Result<()> {
635 let mut headers = Headers::new();
636 for line in raw.split("\r\n").skip(1) {
637 if line.is_empty() {
638 break;
639 }
640 let (name, value) = line.split_once(':').expect("a header line");
641 if name.ends_with(' ') || name.ends_with('\t') {
642 return Err(Error::Protocol("whitespace before the colon".into()));
643 }
644 headers.append(name.trim(), value.trim());
645 }
646 check_framing(&headers)
647 }
648
649 #[test]
658 fn a_message_that_says_where_it_ends_twice_is_refused() {
659 let ambiguous = [
660 (
661 "two lengths that disagree",
662 "POST / HTTP/1.1\r\nhost: x\r\ncontent-length: 6\r\ncontent-length: 0\r\n\r\nsmuggl",
663 ),
664 (
665 "a length and a chunked encoding",
666 "POST / HTTP/1.1\r\nhost: x\r\ncontent-length: 6\r\ntransfer-encoding: chunked\r\n\r\n0\r\n\r\n",
667 ),
668 (
669 "an encoding that is not chunked last",
670 "POST / HTTP/1.1\r\nhost: x\r\ntransfer-encoding: chunked, identity\r\n\r\n0\r\n\r\n",
671 ),
672 (
673 "whitespace before the colon",
674 "POST / HTTP/1.1\r\nhost: x\r\ncontent-length : 6\r\n\r\nsmuggl",
675 ),
676 ];
677
678 for (what, raw) in ambiguous {
679 assert!(
680 framing_of(raw).is_err(),
681 "{what}: accepted a message with two answers for where its body ends"
682 );
683 }
684 }
685
686 #[test]
689 fn an_unambiguous_message_still_parses() {
690 for raw in [
691 "GET / HTTP/1.1\r\nhost: x\r\n\r\n",
692 "POST / HTTP/1.1\r\nhost: x\r\ncontent-length: 3\r\n\r\nabc",
693 "POST / HTTP/1.1\r\nhost: x\r\ntransfer-encoding: chunked\r\n\r\n0\r\n\r\n",
694 "POST / HTTP/1.1\r\nhost: x\r\ncontent-length: 3\r\ncontent-length: 3\r\n\r\nabc",
696 ] {
697 assert!(framing_of(raw).is_ok(), "refused an ordinary message: {raw:?}");
698 }
699 }
700
701 #[test]
704 fn production_gets_one_attempt_and_development_gets_more() {
705 use rustlavel_core::Config;
706
707 let production = Config::with_defaults();
708 production.set("app.env", "production");
709 assert_eq!(port_attempts(&production), 1);
710
711 let local = Config::with_defaults();
712 local.set("app.env", "local");
713 assert!(port_attempts(&local) > 1);
714 }
715
716 #[test]
718 fn the_setting_overrides_the_environment_both_ways() {
719 use rustlavel_core::Config;
720
721 let production = Config::with_defaults();
722 production.set("app.env", "production");
723 production.set("server.port_attempts", "5");
724 assert_eq!(port_attempts(&production), 5);
725
726 let local = Config::with_defaults();
727 local.set("app.env", "local");
728 local.set("server.port_attempts", "1");
729 assert_eq!(port_attempts(&local), 1);
730 }
731
732 #[tokio::test]
735 async fn a_taken_port_moves_to_the_next_one() {
736 let held = TcpListener::bind("127.0.0.1:0").await.unwrap();
737 let taken = held.local_addr().unwrap().port();
738
739 let listener = bind_walking(&format!("127.0.0.1:{taken}"), 10).await.unwrap();
740 assert_ne!(listener.local_addr().unwrap().port(), taken);
741 assert!(listener.local_addr().unwrap().port() > taken);
742 }
743
744 #[tokio::test]
748 async fn a_single_attempt_does_not_move() {
749 let held = TcpListener::bind("127.0.0.1:0").await.unwrap();
750 let taken = held.local_addr().unwrap().port();
751
752 let error = bind_walking(&format!("127.0.0.1:{taken}"), 1).await.unwrap_err();
753 let message = error.to_string();
754 assert!(message.contains(&taken.to_string()), "{message}");
755 assert!(message.contains("lsof"), "the message has to say what to do: {message}");
756 }
757
758 #[tokio::test]
761 async fn port_zero_is_left_to_the_operating_system() {
762 let listener = bind_walking("127.0.0.1:0", 10).await.unwrap();
763 assert_ne!(listener.local_addr().unwrap().port(), 0);
764 }
765
766 #[tokio::test]
769 async fn exhausting_the_range_says_what_it_tried() {
770 let first = TcpListener::bind("127.0.0.1:0").await.unwrap();
772 let start = first.local_addr().unwrap().port();
773 let mut held = vec![first];
774 for offset in 1..3u16 {
775 if let Ok(listener) = TcpListener::bind(format!("127.0.0.1:{}", start + offset)).await {
777 held.push(listener);
778 }
779 }
780
781 let error = bind_walking(&format!("127.0.0.1:{start}"), 3).await;
782 if let Err(error) = error {
783 let message = error.to_string();
784 assert!(message.contains(&start.to_string()), "{message}");
785 assert!(message.contains(&(start + 2).to_string()), "{message}");
786 }
787 }
788 use super::*;
789
790 #[test]
791 fn finds_the_end_of_a_header_block() {
792 assert_eq!(find_head_end(b"GET / HTTP/1.1\r\n\r\nbody"), Some(18));
793 assert_eq!(find_head_end(b"GET / HTTP/1.1\r\n"), None);
794 }
795
796 #[tokio::test]
797 async fn parses_a_request_with_a_body() {
798 let server = Server::new(Router::new(), Context::default());
799 let head = b"POST /users?page=2 HTTP/1.1\r\nHost: localhost\r\nContent-Type: application/json\r\nContent-Length: 14\r\n\r\n";
800
801 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
802 let addr = listener.local_addr().unwrap();
803 tokio::spawn(async move {
804 let (mut stream, _) = listener.accept().await.unwrap();
805 stream.write_all(br#"{"name":"ada"}"#).await.unwrap();
806 });
807 let stream = TcpStream::connect(addr).await.unwrap();
808 let (mut reader, _writer) = stream.into_split();
809
810 let mut buffer = Vec::new();
811 let (mut request, keep_alive) =
812 server.parse(head, &mut reader, &mut buffer, addr).await.unwrap();
813
814 assert_eq!(request.method(), Method::Post);
815 assert_eq!(request.path(), "/users");
816 assert_eq!(request.query("page"), Some("2"));
817 assert_eq!(request.header("host"), Some("localhost"));
818 assert_eq!(request.input("name").as_deref(), Some("ada"));
819 assert!(keep_alive);
820 }
821
822 #[tokio::test]
823 async fn http_1_0_closes_by_default() {
824 let server = Server::new(Router::new(), Context::default());
825 let head = b"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n";
826
827 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
828 let addr = listener.local_addr().unwrap();
829 tokio::spawn(async move {
830 let _ = listener.accept().await;
831 });
832 let (mut reader, _w) = TcpStream::connect(addr).await.unwrap().into_split();
833
834 let mut buffer = Vec::new();
835 let (_request, keep_alive) =
836 server.parse(head, &mut reader, &mut buffer, addr).await.unwrap();
837
838 assert!(!keep_alive);
839 }
840
841 #[tokio::test]
842 async fn rejects_a_body_larger_than_the_limit() {
843 let mut server = Server::new(Router::new(), Context::default());
844 server.limits.max_body_bytes = 8;
845 let head = b"POST / HTTP/1.1\r\nContent-Length: 9999\r\n\r\n";
846
847 let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
848 let addr = listener.local_addr().unwrap();
849 tokio::spawn(async move {
850 let _ = listener.accept().await;
851 });
852 let (mut reader, _w) = TcpStream::connect(addr).await.unwrap().into_split();
853
854 let mut buffer = Vec::new();
855 let error = server.parse(head, &mut reader, &mut buffer, addr).await.unwrap_err();
856
857 assert!(error.to_string().contains("too large"));
858 }
859}