| 89 | const READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); |
| 90 | |
| 91 | async fn handle_connection<H: MetadataHandler>( |
| 92 | handler: &H, |
| 93 | mut stream: tokio::net::TcpStream, |
| 94 | ) -> Result<()> { |
| 95 | let mut buf = vec![0u8; MAX_REQUEST_BYTES]; |
| 96 | let mut used = 0; |
| 97 | let deadline = tokio::time::sleep(READ_TIMEOUT); |
| 98 | tokio::pin!(deadline); |
| 99 | loop { |
| 100 | tokio::select! { |
| 101 | result = stream.read(&mut buf[used..]) => { |
| 102 | let n = result.map_err(|e| miette::miette!("{e}"))?; |
| 103 | if n == 0 { |
| 104 | return Ok(()); |
| 105 | } |
| 106 | used += n; |
| 107 | if buf[..used].windows(4).any(|w| w == b"\r\n\r\n") { |
| 108 | break; |
| 109 | } |
| 110 | if used >= buf.len() { |
| 111 | let _ = stream |
| 112 | .write_all(b"HTTP/1.1 413 Request Entity Too Large\r\nContent-Length: 0\r\n\r\n") |
| 113 | .await; |
| 114 | return Ok(()); |
| 115 | } |
| 116 | } |
| 117 | () = &mut deadline => { |
| 118 | return Ok(()); |
| 119 | } |
| 120 | } |
| 121 | } |
| 122 | let request = String::from_utf8_lossy(&buf[..used]); |
| 123 | let request_line = request.split("\r\n").next().unwrap_or(""); |
| 124 | let mut parts = request_line.split_whitespace(); |
| 125 | let method = parts.next().unwrap_or(""); |
| 126 | let path = parts.next().unwrap_or("/"); |
| 127 | |
| 128 | tokio::time::timeout( |
| 129 | READ_TIMEOUT, |
| 130 | handler.handle(method, path, &buf[..used], &mut stream), |
| 131 | ) |
| 132 | .await |
| 133 | .unwrap_or_else(|_| { |
| 134 | debug!(method, path, "metadata handler timed out"); |
| 135 | Ok(()) |
| 136 | }) |
| 137 | } |