oxedyne/fe2o3/fe2o3_net/src/http/msg.rs
58.1 KiB, 173 runs
created by r1870400018:583, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | #[cfg(feature = "async")] |
| 2 | use crate::conc::AsyncReadIterator; |
| 3 | use crate::{ |
| 4 | constant, |
| 5 | http::{ |
| 6 | fields::{ |
| 7 | ConnectionType, |
| 8 | Cookie, |
| 9 | HeaderFields, |
| 10 | HeaderFieldValue, |
| 11 | HeaderFieldCategory, |
| 12 | HeaderName, |
| 13 | }, |
| 14 | header::{ |
| 15 | HttpHeader, |
| 16 | HttpHeadline, |
| 17 | HttpMethod, |
| 18 | HttpVersion, |
| 19 | }, |
| 20 | status::HttpStatus, |
| 21 | }, |
| 22 | media::{ |
| 23 | ContentTypeValue, |
| 24 | MEDIA_PLAIN_TEXT, |
| 25 | }, |
| 26 | }; |
| 27 | |
| 28 | use oxedyne_fe2o3_core::prelude::*; |
| 29 | |
| 30 | use std::{ |
| 31 | str::FromStr, |
| 32 | path::PathBuf, |
| 33 | time::Duration, |
| 34 | }; |
| 35 | |
| 36 | // `Future`/`Pin` are the async reader's alone; the tokio traits drive every |
| 37 | // stream read and write in this module, all of it behind the `async` feature. |
| 38 | #[cfg(feature = "async")] |
| 39 | use std::{ |
| 40 | future::Future, |
| 41 | pin::Pin, |
| 42 | }; |
| 43 | |
| 44 | #[cfg(feature = "async")] |
| 45 | use tokio::{ |
| 46 | io::{ |
| 47 | AsyncRead, |
| 48 | AsyncReadExt, |
| 49 | AsyncSeekExt, |
| 50 | AsyncWriteExt, |
| 51 | }, |
| 52 | }; |
| 53 | |
| 54 | /// Caller-supplied bounds applied to an HTTP read. Each field is |
| 55 | /// optional so a caller can tune only the bounds it cares about, |
| 56 | /// and the absence of a `ReadLimits` argument (`None`) restores the |
| 57 | /// unbounded pre-feature behaviour used by tests and outbound HTTP |
| 58 | /// clients that trust their own peers. |
| 59 | /// |
| 60 | /// Applied at the server accept path to harden the connection |
| 61 | /// against three common cheap-resource-exhaustion attacks: |
| 62 | /// |
| 63 | /// - Oversized headers (memory blow through a small number of |
| 64 | /// connections) are bounded by `max_header_bytes`. |
| 65 | /// - Oversized bodies (same, per request) are bounded by |
| 66 | /// `max_body_bytes` which is checked against the `Content-Length` |
| 67 | /// header up front and enforced during the body read. |
| 68 | /// - Slowloris-style trickle attacks (pin a connection for hours |
| 69 | /// by sending one byte at a time) are bounded by |
| 70 | /// `header_read_timeout` which wraps every header-phase stream |
| 71 | /// read in a `tokio::time::timeout`. |
| 72 | #[derive(Clone, Debug, Default)] |
| 73 | pub struct ReadLimits { |
| 74 | /// Maximum total bytes accepted for the request / response header |
| 75 | /// block before the reader aborts with `TooBig`. The count covers |
| 76 | /// every byte read up to the terminating `CRLF CRLF`. |
| 77 | pub max_header_bytes: Option<usize>, |
| 78 | /// Maximum bytes permitted in the message body. Checked against |
| 79 | /// `Content-Length` first (rejecting oversize requests before a |
| 80 | /// single byte is read) and enforced as an upper bound during the |
| 81 | /// body read loop. |
| 82 | pub max_body_bytes: Option<usize>, |
| 83 | /// Upper bound on the wall-clock duration between entering the |
| 84 | /// header read loop and the `CRLF CRLF` terminator arriving. A |
| 85 | /// slow client exceeding this limit has its connection dropped |
| 86 | /// with a `TimedOut` error. |
| 87 | pub header_read_timeout: Option<Duration>, |
| 88 | } |
| 89 | |
| 90 | impl ReadLimits { |
| 91 | /// Shortcut for constructing a fully permissive limits struct |
| 92 | /// with explicit `None` for every field. Mostly a readability |
| 93 | /// alias for `ReadLimits::default()`. |
| 94 | pub fn unbounded() -> Self { |
| 95 | Self::default() |
| 96 | } |
| 97 | } |
| 98 | |
| 99 | /// A body that is read off a file as it is written, rather than gathered into |
| 100 | /// memory first. |
| 101 | /// |
| 102 | /// A two hour recording is a couple of gigabytes, and a player seeking through it |
| 103 | /// asks for a megabyte at a time. Reading the file to answer such a request costs |
| 104 | /// two thousand times what the answer weighs, and a handful of concurrent viewers |
| 105 | /// then costs the machine its memory. So the window is named here and copied |
| 106 | /// straight from the file to the socket when the message is written. |
| 107 | #[derive(Clone, Debug)] |
| 108 | pub struct FileWindow { |
| 109 | /// The file to read from. |
| 110 | pub path: PathBuf, |
| 111 | /// Offset of the first byte sent. |
| 112 | pub start: u64, |
| 113 | /// How many bytes to send. |
| 114 | pub len: u64, |
| 115 | } |
| 116 | |
| 117 | impl FileWindow { |
| 118 | |
| 119 | /// Name a window of a file as a message body. |
| 120 | pub fn new(path: PathBuf, start: u64, len: u64) -> Self { |
| 121 | Self { |
| 122 | path, |
| 123 | start, |
| 124 | len, |
| 125 | } |
| 126 | } |
| 127 | |
| 128 | /// Read the window into memory. |
| 129 | /// |
| 130 | /// The opposite of what naming a window is for, and asked only when the |
| 131 | /// bytes have to be transformed before they go out -- a body about to be |
| 132 | /// encoded cannot be copied straight from the file, because the file is no |
| 133 | /// longer what is being sent. The eligibility rules that lead here exclude |
| 134 | /// the large media types, so what is read is markup, script or a module. |
| 135 | #[cfg(feature = "async")] |
| 136 | pub async fn read(&self) -> Outcome<Vec<u8>> { |
| 137 | let mut file = match tokio::fs::File::open(&self.path).await { |
| 138 | Ok(f) => f, |
| 139 | Err(e) => return Err(err!(e, |
| 140 | "Opening {:?} to read bytes {} to {} of it.", |
| 141 | self.path, self.start, self.start + self.len; |
| 142 | IO, File, Read)), |
| 143 | }; |
| 144 | if self.start > 0 { |
| 145 | let result = file.seek(std::io::SeekFrom::Start(self.start)).await; |
| 146 | res!(result, IO, File, Seek); |
| 147 | } |
| 148 | let mut buf = vec![0u8; self.len as usize]; |
| 149 | let result = file.read_exact(&mut buf).await; |
| 150 | res!(result, IO, File, Read); |
| 151 | Ok(buf) |
| 152 | } |
| 153 | |
| 154 | /// Copy the window from the file to the given sink, in chunks of |
| 155 | /// `CHUNK_SIZE`. |
| 156 | /// |
| 157 | /// Exactly `len` bytes are written, or the write fails. A file that shrank |
| 158 | /// between the length being measured and the body being sent would otherwise |
| 159 | /// leave the connection short of the `Content-Length` already promised, which |
| 160 | /// desynchronises every message after it on a kept-alive connection -- so a |
| 161 | /// short read is an error, and the caller drops the connection. |
| 162 | #[cfg(feature = "async")] |
| 163 | pub async fn write_to<const CHUNK_SIZE: usize, W: AsyncWriteExt + Unpin>( |
| 164 | &self, |
| 165 | sink: &mut W, |
| 166 | ) |
| 167 | -> Outcome<()> |
| 168 | { |
| 169 | let mut file = match tokio::fs::File::open(&self.path).await { |
| 170 | Ok(f) => f, |
| 171 | Err(e) => return Err(err!(e, |
| 172 | "Opening {:?} to send bytes {} to {} of it.", |
| 173 | self.path, self.start, self.start + self.len; |
| 174 | IO, File, Read)), |
| 175 | }; |
| 176 | if self.start > 0 { |
| 177 | let result = file.seek(std::io::SeekFrom::Start(self.start)).await; |
| 178 | res!(result, IO, File, Seek); |
| 179 | } |
| 180 | let mut left = self.len; |
| 181 | let mut buf = vec![0u8; CHUNK_SIZE]; |
| 182 | while left > 0 { |
| 183 | let want = std::cmp::min(left as usize, CHUNK_SIZE); |
| 184 | let n = match file.read(&mut buf[..want]).await { |
| 185 | Ok(0) => return Err(err!( |
| 186 | "{:?} ended {} bytes short of the {} promised from offset {}.", |
| 187 | self.path, left, self.len, self.start; |
| 188 | IO, File, Read, Size)), |
| 189 | Ok(n) => n, |
| 190 | Err(e) => return Err(err!(e, |
| 191 | "Reading {:?} at offset {}.", self.path, self.start + (self.len - left); |
| 192 | IO, File, Read)), |
| 193 | }; |
| 194 | let result = sink.write_all(&buf[..n]).await; |
| 195 | res!(result, IO, Network, Wire, Write); |
| 196 | left -= n as u64; |
| 197 | } |
| 198 | Ok(()) |
| 199 | } |
| 200 | } |
| 201 | |
| 202 | #[derive(Debug, Default)] |
| 203 | pub struct HttpMessage { |
| 204 | pub header: HttpHeader, |
| 205 | pub body: Vec<u8>, |
| 206 | /// Send the headers and stop. |
| 207 | /// |
| 208 | /// A `HEAD` asks what a `GET` would answer without the answer itself, and RFC 9110 9.3.2 |
| 209 | /// requires the header fields to be the ones the `GET` would carry -- `Content-Length` |
| 210 | /// included, describing a body that is deliberately not sent. So the body is built as |
| 211 | /// usual, counted as usual, and withheld at the wire. |
| 212 | pub head_only: bool, |
| 213 | /// A body held on disk rather than in `body`, sent straight from the file. |
| 214 | /// |
| 215 | /// When set this *is* the body: `body` is not sent, and `Content-Length` is |
| 216 | /// the window's length. See [`FileWindow`]. |
| 217 | pub file: Option<FileWindow>, |
| 218 | } |
| 219 | |
| 220 | impl HttpMessage { |
| 221 | |
| 222 | pub fn new_response(status: HttpStatus) -> Self { |
| 223 | Self { |
| 224 | header: HttpHeader { |
| 225 | version: HttpVersion::Http1_1, |
| 226 | headline: HttpHeadline::Response { status }, |
| 227 | fields: HeaderFields::default(), |
| 228 | }, |
| 229 | body: Vec::new(), |
| 230 | head_only: false, |
| 231 | file: None, |
| 232 | } |
| 233 | } |
| 234 | |
| 235 | pub fn respond_with_text<S: AsRef<str>>(status: HttpStatus, txt: S) -> Self { |
| 236 | let mut fields = HeaderFields::default(); |
| 237 | fields.insert( |
| 238 | HeaderName::ContentType, |
| 239 | HeaderFieldValue::ContentType(MEDIA_PLAIN_TEXT), |
| 240 | Some(HeaderFieldCategory::Entity as u16), |
| 241 | ); |
| 242 | Self { |
| 243 | header: HttpHeader { |
| 244 | version: HttpVersion::Http1_1, |
| 245 | headline: HttpHeadline::Response { status }, |
| 246 | fields, |
| 247 | }, |
| 248 | body: txt.as_ref().as_bytes().to_vec(), |
| 249 | head_only: false, |
| 250 | file: None, |
| 251 | } |
| 252 | } |
| 253 | |
| 254 | pub fn ok_respond_with_text<S: AsRef<str>>(txt: S) -> Self { |
| 255 | Self::respond_with_text(HttpStatus::OK, txt) |
| 256 | } |
| 257 | |
| 258 | /// Set the message status if the message is a response, otherwise do nothing. |
| 259 | pub fn with_status(mut self, status: HttpStatus) -> Self { |
| 260 | match self.header.headline { |
| 261 | HttpHeadline::Response { .. } => |
| 262 | self.header.headline = HttpHeadline::Response { status }, |
| 263 | _ => () |
| 264 | } |
| 265 | self |
| 266 | } |
| 267 | |
| 268 | pub fn with_field( |
| 269 | mut self, |
| 270 | nam: HeaderName, |
| 271 | val: HeaderFieldValue, |
| 272 | ) |
| 273 | -> Self |
| 274 | { |
| 275 | self.header.fields.insert(nam, val, None); |
| 276 | self |
| 277 | } |
| 278 | |
| 279 | pub fn with_field_with_order( |
| 280 | mut self, |
| 281 | nam: HeaderName, |
| 282 | val: HeaderFieldValue, |
| 283 | ord: Option<u16>, |
| 284 | ) |
| 285 | -> Self |
| 286 | { |
| 287 | self.header.fields.insert(nam, val, ord); |
| 288 | self |
| 289 | } |
| 290 | |
| 291 | /// Set the message status, returning an error if the message is not a response. |
| 292 | pub fn set_response_status(&mut self, status: HttpStatus) -> Outcome<()> { |
| 293 | match &mut self.header.headline { |
| 294 | HttpHeadline::Response { status: old_status } => { |
| 295 | *old_status = status; |
| 296 | Ok(()) |
| 297 | }, |
| 298 | _ => Err(err!("HTTP message is not a response."; Invalid, Mismatch)), |
| 299 | } |
| 300 | } |
| 301 | |
| 302 | pub fn set_response_code(&mut self, code: &str) -> Outcome<()> { |
| 303 | match &mut self.header.headline { |
| 304 | HttpHeadline::Response { status } => { |
| 305 | *status = res!(HttpStatus::from_str(code)); |
| 306 | Ok(()) |
| 307 | }, |
| 308 | _ => Err(err!("HTTP message is not a response."; Invalid, Mismatch)), |
| 309 | } |
| 310 | } |
| 311 | |
| 312 | pub fn set_cookie(mut self, cookie: Cookie) -> Self { |
| 313 | self.header.fields.insert( |
| 314 | HeaderName::SetCookie, |
| 315 | HeaderFieldValue::SetCookie(cookie), |
| 316 | None, |
| 317 | ); |
| 318 | self |
| 319 | } |
| 320 | |
| 321 | pub fn insert( |
| 322 | &mut self, |
| 323 | nam: HeaderName, |
| 324 | val: HeaderFieldValue, |
| 325 | ord: Option<u16>, |
| 326 | ) |
| 327 | -> bool |
| 328 | { |
| 329 | self.header.fields.insert(nam, val, ord) |
| 330 | } |
| 331 | |
| 332 | pub fn log(&self, log_level: LogLevel) { |
| 333 | if log_level >= LogLevel::Trace { |
| 334 | //Header. |
| 335 | let is_request = match &self.header.headline { |
| 336 | HttpHeadline::Request { method, loc } => { |
| 337 | trace!("HTTP Request Header:"); |
| 338 | trace!(" {} {} {}", method, loc, self.header.version); |
| 339 | true |
| 340 | }, |
| 341 | HttpHeadline::Response { status } => { |
| 342 | trace!("HTTP Response Header:"); |
| 343 | trace!(" {} {} {}", self.header.version, status, status.desc()); |
| 344 | false |
| 345 | }, |
| 346 | }; |
| 347 | for (k, header_field_values) in self.header.fields.iter() { |
| 348 | for header_field_value in header_field_values { |
| 349 | let prefix = match header_field_value { |
| 350 | HeaderFieldValue::Generic(_) => " ", |
| 351 | _ => " > ", |
| 352 | }; |
| 353 | trace!("{}{}: {}", prefix, k, header_field_value); |
| 354 | } |
| 355 | } |
| 356 | |
| 357 | // Body. A body on disk is never read to be logged: the point of |
| 358 | // naming it as a window is that nothing gathers it into memory. |
| 359 | if let Some(window) = &self.file { |
| 360 | trace!("HTTP Body: bytes {} to {} of {:?}, {} bytes, sent from the file.", |
| 361 | window.start, window.start + window.len, window.path, window.len); |
| 362 | return; |
| 363 | } |
| 364 | |
| 365 | let body_is_text = match self.body_is_text() { |
| 366 | Some(true) => true, |
| 367 | Some(false) => false, |
| 368 | None => { |
| 369 | trace!("HTTP message has no body."); |
| 370 | return; |
| 371 | }, |
| 372 | }; |
| 373 | const LIM: usize = constant::HTTP_BODY_BYTES_MAX_VIEW; |
| 374 | let (display_all_bytes, bytes_report) = if self.body.len() < LIM { |
| 375 | (true, fmt!("[{} bytes, all displayed]", self.body.len())) |
| 376 | } else { |
| 377 | (false, fmt!("[{} bytes, only {} displayed]", self.body.len(), LIM)) |
| 378 | }; |
| 379 | match is_request { |
| 380 | true => trace!("HTTP Request Body {}:", bytes_report), |
| 381 | false => trace!("HTTP Response Body {}:", bytes_report), |
| 382 | } |
| 383 | // Text dump of body, if necessary. |
| 384 | if body_is_text { |
| 385 | if display_all_bytes { |
| 386 | trace!("\n{}", String::from_utf8_lossy(&self.body[..])); |
| 387 | } else { |
| 388 | trace!("\n{{START}}{}{{END}}", String::from_utf8_lossy(&self.body[..LIM])); |
| 389 | } |
| 390 | return; |
| 391 | } |
| 392 | // Binary dump of body. |
| 393 | let lines = if display_all_bytes { |
| 394 | dump!(" {:02x}", &self.body[..], 16) |
| 395 | } else { |
| 396 | dump!(" {:02x}", &self.body[..LIM], 16) |
| 397 | }; |
| 398 | for line in lines { |
| 399 | trace!(" {}", line); |
| 400 | } |
| 401 | } |
| 402 | } |
| 403 | |
| 404 | #[cfg(feature = "async")] |
| 405 | pub async fn read< |
| 406 | 'a, |
| 407 | const HEADER_CHUNK_SIZE: usize, |
| 408 | const BODY_CHUNK_SIZE: usize, |
| 409 | R: AsyncRead + Unpin, |
| 410 | >( |
| 411 | mut stream: Pin<&mut R>, |
| 412 | remnant: &Vec<u8>, |
| 413 | is_request: Option<bool>, |
| 414 | limits: Option<&ReadLimits>, |
| 415 | ) |
| 416 | -> Outcome<(Option<Self>, Vec<u8>)> |
| 417 | { |
| 418 | trace!("Entered HttpMessage::read"); |
| 419 | let result = HttpHeader::read::<HEADER_CHUNK_SIZE, _>( |
| 420 | stream.as_mut(), |
| 421 | remnant, |
| 422 | is_request, |
| 423 | limits, |
| 424 | ).await; |
| 425 | |
| 426 | match result { |
| 427 | Ok(Some((header, mut remnant, content_length))) => { |
| 428 | trace!("remnant size = {}, content_length = {}", remnant.len(), content_length); |
| 429 | |
| 430 | // Reject the request up front if its declared content |
| 431 | // length exceeds the caller's bound. Cheaper than |
| 432 | // accumulating bytes and erroring halfway through, and |
| 433 | // surfaces the error before any body bytes are read. |
| 434 | if let Some(lim) = limits.and_then(|l| l.max_body_bytes) { |
| 435 | if content_length > lim { |
| 436 | return Err(err!( |
| 437 | "HTTP body of {} bytes exceeds the \ |
| 438 | configured limit of {}.", content_length, lim; |
| 439 | IO, Network, Input, TooBig)); |
| 440 | } |
| 441 | } |
| 442 | |
| 443 | let mut msg = HttpMessage::default(); |
| 444 | msg.header = header; |
| 445 | |
| 446 | // A chunked body carries its own lengths on the wire, and |
| 447 | // says so instead of declaring a `Content-Length`. Much of |
| 448 | // the web answers this way, so a client that cannot decode |
| 449 | // it reads every such page as empty. |
| 450 | if is_chunked(&msg.header.fields) { |
| 451 | let (body, rest) = res!(read_chunked::<BODY_CHUNK_SIZE, _>( |
| 452 | stream.as_mut(), |
| 453 | remnant, |
| 454 | limits, |
| 455 | ).await); |
| 456 | msg.body = body; |
| 457 | // The body now sits decoded in `msg.body`, so the message is |
| 458 | // no longer chunked and must stop saying that it is. A proxy |
| 459 | // re-emitting it would otherwise send `Transfer-Encoding: |
| 460 | // chunked` over a body that is not chunked, and, once |
| 461 | // `write_all` added the length, both framing fields at once -- |
| 462 | // the combination RFC 9112 §6.1 forbids and §6.3 calls a |
| 463 | // smuggling signal. |
| 464 | msg.header.fields.remove(&HeaderName::TransferEncoding); |
| 465 | return Ok((Some(msg), rest)); |
| 466 | } |
| 467 | |
| 468 | if content_length > 0 { |
| 469 | // Reserve for the body the peer declares, but only as far as |
| 470 | // `HTTP_BODY_RESERVE_MAX`: a caller that set no `max_body_bytes` -- an |
| 471 | // outbound client, or a test -- would otherwise hand a stranger's header |
| 472 | // straight to the allocator, and a `Content-Length` of a terabyte costs |
| 473 | // nothing to write. The vector still grows to whatever genuinely arrives; all |
| 474 | // that is given up is one reservation on a body that large. |
| 475 | let mut body = Vec::with_capacity( |
| 476 | std::cmp::min(content_length, constant::HTTP_BODY_RESERVE_MAX), |
| 477 | ); |
| 478 | body.extend_from_slice(&remnant); |
| 479 | let mut bytes_read = body.len(); |
| 480 | |
| 481 | while bytes_read < content_length { |
| 482 | let mut chunk = [0; BODY_CHUNK_SIZE]; |
| 483 | let result = stream.as_mut().read(&mut chunk).await; |
| 484 | match result { |
| 485 | Ok(0) => { |
| 486 | // A reply to `HEAD` states the length of a body |
| 487 | // it will never send, so the stream ends with |
| 488 | // nothing received -- and that is complete, not |
| 489 | // a failure, or a client could never make a HEAD |
| 490 | // request at all. A response that delivered some |
| 491 | // of a declared body and then broke off is a |
| 492 | // genuine truncation, and still an error: an |
| 493 | // empty body is the tell that distinguishes the |
| 494 | // two, since HEAD yields no bytes and a cut-off |
| 495 | // GET yields the part that arrived. A truncated |
| 496 | // *request* is a broken client, dropped as before. |
| 497 | if is_request == Some(false) && body.is_empty() { |
| 498 | msg.body = body; |
| 499 | return Ok((Some(msg), Vec::new())); |
| 500 | } |
| 501 | warn!("UnexpectedEof treated as connection closure."); |
| 502 | return Ok((None, body.to_vec())); |
| 503 | } |
| 504 | Ok(n) => { |
| 505 | body.extend_from_slice(&chunk[..n]); |
| 506 | bytes_read += n; |
| 507 | // Defence-in-depth: the upfront |
| 508 | // Content-Length check already |
| 509 | // rejected oversize requests, but a |
| 510 | // buggy or malicious remote could |
| 511 | // keep streaming beyond the declared |
| 512 | // length. Stop as soon as we have |
| 513 | // enough to satisfy the header. |
| 514 | if let Some(lim) = |
| 515 | limits.and_then(|l| l.max_body_bytes) |
| 516 | { |
| 517 | if bytes_read > lim { |
| 518 | return Err(err!( |
| 519 | "HTTP body overflowed \ |
| 520 | the configured limit of {} \ |
| 521 | bytes during read.", lim; |
| 522 | IO, Network, Input, TooBig)); |
| 523 | } |
| 524 | } |
| 525 | } |
| 526 | Err(e) => return Err(e.into()), |
| 527 | } |
| 528 | } |
| 529 | |
| 530 | remnant = if bytes_read > content_length { |
| 531 | body[content_length..].to_vec() |
| 532 | } else { |
| 533 | Vec::new() |
| 534 | }; |
| 535 | |
| 536 | msg.body = body[..content_length].to_vec(); |
| 537 | Ok((Some(msg), remnant)) |
| 538 | } else { |
| 539 | Ok((Some(msg), remnant)) |
| 540 | } |
| 541 | } |
| 542 | Ok(None) => Ok((None, Vec::new())), |
| 543 | Err(e) => Err(e), |
| 544 | } |
| 545 | } |
| 546 | |
| 547 | #[cfg(feature = "async")] |
| 548 | pub async fn write_all< |
| 549 | R: AsyncWriteExt + Unpin, |
| 550 | >( |
| 551 | mut self, |
| 552 | stream: &mut R, |
| 553 | ) |
| 554 | -> Outcome<()> |
| 555 | { |
| 556 | // `HeaderFields::insert` refuses a `Content-Length` while a |
| 557 | // `Transfer-Encoding` stands, so a message that really is chunked keeps |
| 558 | // its own framing and does not acquire a second, contradictory one. |
| 559 | if !self.is_bodiless_status() { |
| 560 | let _ = self.insert( |
| 561 | HeaderName::ContentLength, |
| 562 | HeaderFieldValue::ContentLength(self.body_len()), |
| 563 | Some(HeaderFieldCategory::Entity as u16), |
| 564 | ); |
| 565 | } |
| 566 | self.log(log_get_level!()); |
| 567 | let result = stream.write_all(&self.header.as_vec()).await; |
| 568 | res!(result); |
| 569 | // The answer to a `HEAD` is the header the matching `GET` would have sent, and nothing |
| 570 | // after it. `Content-Length` above still counts the body that was built, which is what |
| 571 | // the field is for -- a caller asking how big a thing is deserves the true answer. |
| 572 | if !self.head_only { |
| 573 | match &self.file { |
| 574 | // A body on disk goes to the socket a chunk at a time and is |
| 575 | // never gathered into memory, which is the whole point of |
| 576 | // naming it as a window rather than reading it. |
| 577 | Some(window) => res!(window.write_to::< |
| 578 | { constant::HTTP_FILE_BODY_CHUNK_SIZE }, |
| 579 | _, |
| 580 | >(stream).await), |
| 581 | None => { |
| 582 | let result = stream.write_all(&self.body).await; |
| 583 | res!(result); |
| 584 | } |
| 585 | } |
| 586 | } |
| 587 | let result = stream.flush().await; |
| 588 | res!(result); |
| 589 | Ok(()) |
| 590 | } |
| 591 | |
| 592 | /// How many bytes the body comes to, whether it is held in memory or named |
| 593 | /// as a window of a file. |
| 594 | /// |
| 595 | /// This is the `Content-Length` the message will carry, and a `HEAD` answer |
| 596 | /// states it too -- a caller asking how big a thing is deserves the true |
| 597 | /// answer. |
| 598 | pub fn body_len(&self) -> usize { |
| 599 | match &self.file { |
| 600 | Some(window) => window.len as usize, |
| 601 | None => self.body.len(), |
| 602 | } |
| 603 | } |
| 604 | |
| 605 | /// Is this a response whose status forbids a `Content-Length` of zero? |
| 606 | /// |
| 607 | /// RFC 9110 §8.6 forbids the field on a `1xx` or `204`, and on a `304` it |
| 608 | /// would have to state the length of the representation not sent, which a |
| 609 | /// zero here would misstate. |
| 610 | pub fn is_bodiless_status(&self) -> bool { |
| 611 | match &self.header.headline { |
| 612 | HttpHeadline::Response { status } => { |
| 613 | let code = *status as u16; |
| 614 | code < 200 || code == 204 || code == 304 |
| 615 | } |
| 616 | HttpHeadline::Request { .. } => false, |
| 617 | } |
| 618 | } |
| 619 | |
| 620 | /// Send a window of a file as the body, rather than a buffer of bytes. |
| 621 | pub fn with_file_window(mut self, window: FileWindow) -> Self { |
| 622 | self.file = Some(window); |
| 623 | self.body = Vec::new(); |
| 624 | self |
| 625 | } |
| 626 | |
| 627 | /// Answer with the headers and no body, as a `HEAD` asks. |
| 628 | /// |
| 629 | /// `Content-Length` still counts the body that was built -- withholding it is the point of |
| 630 | /// the method, and lying about its size would defeat the reason anyone asks. |
| 631 | pub fn head_only(mut self) -> Self { |
| 632 | self.head_only = true; |
| 633 | self |
| 634 | } |
| 635 | |
| 636 | pub fn body_text(&mut self, txt: &str) { |
| 637 | self.body = txt.as_bytes().to_vec() |
| 638 | } |
| 639 | |
| 640 | pub fn with_body(mut self, byts: Vec<u8>) -> Self { |
| 641 | self.body = byts; |
| 642 | self |
| 643 | } |
| 644 | |
| 645 | pub fn body_as_string(&self) -> std::borrow::Cow<'_, str> { |
| 646 | String::from_utf8_lossy(&self.body[..]) |
| 647 | } |
| 648 | |
| 649 | pub fn body_is_text(&self) -> Option<bool> { |
| 650 | if self.body.len() > 0 { |
| 651 | if let Some(hfv) = self.header.fields.get_one(&HeaderName::ContentType) { |
| 652 | if let HeaderFieldValue::ContentType(ContentTypeValue::MediaType((mt, _))) = hfv { |
| 653 | Some(mt.is_text()) |
| 654 | } else { |
| 655 | None |
| 656 | } |
| 657 | } else { |
| 658 | Some(false) |
| 659 | } |
| 660 | } else { |
| 661 | None |
| 662 | } |
| 663 | } |
| 664 | |
| 665 | pub fn get_connection_close(&self) -> bool { |
| 666 | if let Some(hfv) = self.header.get_a_field_value(&HeaderName::Connection) { |
| 667 | if let HeaderFieldValue::Connection(Some(ct), _) = hfv { |
| 668 | if let ConnectionType::Close = ct { |
| 669 | return true; |
| 670 | } |
| 671 | } |
| 672 | } |
| 673 | false |
| 674 | } |
| 675 | |
| 676 | pub fn set_connection_close(&mut self, close: bool) -> bool { |
| 677 | self.insert( |
| 678 | HeaderName::Connection, |
| 679 | HeaderFieldValue::Connection(Some(ConnectionType::new(close)), Vec::new()), |
| 680 | Some(HeaderFieldCategory::General as u16), |
| 681 | ) |
| 682 | } |
| 683 | |
| 684 | pub fn has_websocket_headers(&self) -> bool { |
| 685 | |
| 686 | let connection_upgrade = match self.header.get_the_field_value(&HeaderName::Connection) { |
| 687 | Ok(HeaderFieldValue::Connection(ct_opt, list)) => match ct_opt { |
| 688 | Some(ConnectionType::KeepAlive) | None => { |
| 689 | list.iter().any(|s| s == "upgrade") |
| 690 | } |
| 691 | _ => false, |
| 692 | }, |
| 693 | _ => false, |
| 694 | }; |
| 695 | |
| 696 | let upgrade_websocket = match self.header.get_the_field_value(&HeaderName::Upgrade) { |
| 697 | Ok(HeaderFieldValue::Upgrade(list)) => { |
| 698 | list.iter().any(|s| s == "websocket") |
| 699 | }, |
| 700 | _ => false, |
| 701 | }; |
| 702 | |
| 703 | trace!("connection_upgrade = {}", connection_upgrade); |
| 704 | trace!("upgrade_websocket = {}", upgrade_websocket); |
| 705 | |
| 706 | if connection_upgrade && upgrade_websocket { |
| 707 | true |
| 708 | } else { |
| 709 | false |
| 710 | } |
| 711 | } |
| 712 | |
| 713 | pub fn is_websocket_upgrade(&self) -> bool { |
| 714 | |
| 715 | let is_websocket_request = match &self.header.headline { |
| 716 | HttpHeadline::Request { method, .. } => { |
| 717 | match method { |
| 718 | HttpMethod::GET => { |
| 719 | true |
| 720 | } |
| 721 | _ => false, |
| 722 | } |
| 723 | } |
| 724 | _ => false, |
| 725 | }; |
| 726 | |
| 727 | let key_present = match self.header.get_the_field_value(&HeaderName::SecWebSocketKey) { |
| 728 | Ok(HeaderFieldValue::SecWebSocketKey(key)) => { |
| 729 | key.len() == 24 && |
| 730 | key.chars().all(|c| |
| 731 | c.is_ascii_alphanumeric() |
| 732 | || c == '+' |
| 733 | || c == '/' |
| 734 | || c == '=' |
| 735 | ) |
| 736 | }, |
| 737 | _ => false, |
| 738 | }; |
| 739 | |
| 740 | trace!("is_websocket_request = {}", is_websocket_request); |
| 741 | trace!("key_present = {}", key_present); |
| 742 | |
| 743 | if is_websocket_request && self.has_websocket_headers() && key_present { |
| 744 | true |
| 745 | } else { |
| 746 | false |
| 747 | } |
| 748 | } |
| 749 | |
| 750 | /// Is this the `101` response that completes a websocket handshake whose key derives |
| 751 | /// `expected_accept_key`? |
| 752 | pub fn is_websocket_handshake(&self, expected_accept_key: &str) -> bool { |
| 753 | self.check_websocket_handshake(expected_accept_key).is_ok() |
| 754 | } |
| 755 | |
| 756 | /// Checks that this response completes a websocket handshake, in the order a client needs to |
| 757 | /// hear about failure: the status first, since a refusal carries no accept key at all, then the |
| 758 | /// upgrade fields, then `Sec-WebSocket-Accept` against `expected_accept_key` (see |
| 759 | /// [`crate::ws::core::accept_key`]). |
| 760 | pub fn check_websocket_handshake(&self, expected_accept_key: &str) -> Outcome<()> { |
| 761 | match self.header.headline { |
| 762 | HttpHeadline::Response { status: HttpStatus::SwitchingProtocols } => (), |
| 763 | HttpHeadline::Response { status } => return Err(err!( |
| 764 | "The server answered the websocket upgrade with status {}, not 101.", |
| 765 | status; |
| 766 | IO, Network, Wire, Invalid, Input)), |
| 767 | _ => return Err(err!( |
| 768 | "Expected a response to the websocket upgrade request, found a request."; |
| 769 | IO, Network, Wire, Invalid, Input)), |
| 770 | } |
| 771 | if !self.has_websocket_headers() { |
| 772 | return Err(err!( |
| 773 | "The server's 101 response lacks 'Connection: Upgrade' or 'Upgrade: websocket'."; |
| 774 | IO, Network, Wire, Invalid, Input, Missing)); |
| 775 | } |
| 776 | match self.header.get_the_field_value(&HeaderName::SecWebSocketAccept) { |
| 777 | Ok(HeaderFieldValue::SecWebSocketKey(key)) if key == expected_accept_key => Ok(()), |
| 778 | Ok(HeaderFieldValue::SecWebSocketKey(key)) => Err(err!( |
| 779 | "The server's Sec-WebSocket-Accept key '{}' is not '{}', the value derived from \ |
| 780 | the key this client sent.", key, expected_accept_key; |
| 781 | IO, Network, Wire, Invalid, Input, Mismatch)), |
| 782 | _ => Err(err!( |
| 783 | "The server's 101 response has no Sec-WebSocket-Accept key."; |
| 784 | IO, Network, Wire, Invalid, Input, Missing)), |
| 785 | } |
| 786 | } |
| 787 | } |
| 788 | |
| 789 | |
| 790 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 791 | // │ CHUNKED TRANSFER ENCODING │ |
| 792 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 793 | |
| 794 | /// Whether the message says its body arrives in chunks (RFC 9112 §7.1). |
| 795 | /// |
| 796 | /// `chunked` is the last coding applied, so a value of `gzip, chunked` is |
| 797 | /// chunked too, and the test is for its presence in the list. |
| 798 | pub fn is_chunked(fields: &HeaderFields) -> bool { |
| 799 | match fields.get_list(&HeaderName::TransferEncoding) { |
| 800 | Some(list) => list.iter().any(|v| fmt!("{}", v).to_lowercase().contains("chunked")), |
| 801 | None => false, |
| 802 | } |
| 803 | } |
| 804 | |
| 805 | /// Read a chunked body to its terminating zero-length chunk, returning the |
| 806 | /// decoded bytes and whatever the peer sent after them. |
| 807 | /// |
| 808 | /// Each chunk is a hex length, optional extensions after a `;`, a CRLF, that |
| 809 | /// many bytes, and another CRLF. A zero length ends the body, after which the |
| 810 | /// peer may send trailer fields, which are consumed and discarded. |
| 811 | /// |
| 812 | /// `limits.max_body_bytes` bounds the *decoded* size: a chunked body declares |
| 813 | /// no total up front, so the only way to bound it is to stop reading when it |
| 814 | /// gets too big, which is what happens here. |
| 815 | #[cfg(feature = "async")] |
| 816 | async fn read_chunked< |
| 817 | const BODY_CHUNK_SIZE: usize, |
| 818 | R: AsyncRead + Unpin, |
| 819 | >( |
| 820 | mut stream: Pin<&mut R>, |
| 821 | remnant: Vec<u8>, |
| 822 | limits: Option<&ReadLimits>, |
| 823 | ) |
| 824 | -> Outcome<(Vec<u8>, Vec<u8>)> |
| 825 | { |
| 826 | let max = limits.and_then(|l| l.max_body_bytes); |
| 827 | let max_trailer = limits.and_then(|l| l.max_header_bytes); // Same bound as the header block. |
| 828 | let mut raw: Vec<u8> = remnant; // Undecoded bytes not yet consumed. |
| 829 | let mut out: Vec<u8> = Vec::new(); // The decoded body. |
| 830 | |
| 831 | loop { |
| 832 | // The chunk size line. Capped independently of any configured limits -- |
| 833 | // see `HTTP_CHUNK_LINE_MAX`. |
| 834 | let line = match res!(take_line::<BODY_CHUNK_SIZE, _>( |
| 835 | stream.as_mut(), &mut raw, Some(constant::HTTP_CHUNK_LINE_MAX), "chunk-size", |
| 836 | ).await) { |
| 837 | Some(l) => l, |
| 838 | None => return Err(err!( |
| 839 | "The peer closed the connection in the middle of a chunked \ |
| 840 | body, before its terminating chunk."; |
| 841 | IO, Network, Wire, Read, Missing)), |
| 842 | }; |
| 843 | let size_txt = match line.split(';').next() { |
| 844 | Some(s) => s.trim().to_string(), |
| 845 | None => String::new(), |
| 846 | }; |
| 847 | let size = match usize::from_str_radix(&size_txt, 16) { |
| 848 | Ok(n) => n, |
| 849 | Err(e) => return Err(err!(e, |
| 850 | "A chunked body declared a chunk size of {:?}, which is not a \ |
| 851 | hexadecimal length.", size_txt; |
| 852 | IO, Network, Wire, Read, Invalid)), |
| 853 | }; |
| 854 | |
| 855 | // The last chunk is empty, and any trailer fields follow it up to a |
| 856 | // blank line. They are read so the stream is left where the next |
| 857 | // message begins, and are then discarded. |
| 858 | if size == 0 { |
| 859 | loop { |
| 860 | match res!(take_line::<BODY_CHUNK_SIZE, _>( |
| 861 | stream.as_mut(), &mut raw, max_trailer, "trailer", |
| 862 | ).await) { |
| 863 | Some(l) if l.is_empty() => break, |
| 864 | Some(_) => continue, |
| 865 | // A peer that hangs up rather than closing off its |
| 866 | // trailers has still sent us the whole body. |
| 867 | None => return Ok((out, Vec::new())), |
| 868 | } |
| 869 | } |
| 870 | return Ok((out, raw)); |
| 871 | } |
| 872 | |
| 873 | if let Some(lim) = max { |
| 874 | if out.len().saturating_add(size) > lim { |
| 875 | return Err(err!( |
| 876 | "A chunked HTTP body overflowed the configured limit of {} \ |
| 877 | bytes.", lim; |
| 878 | IO, Network, Input, TooBig)); |
| 879 | } |
| 880 | } |
| 881 | |
| 882 | // The chunk data, and the CRLF that closes it. The `+ 2` is a checked |
| 883 | // add so a hostile chunk-size line, `ffffffffffffffff` and the like, |
| 884 | // cannot overflow `usize` into a panic or a wrapped-round short slice. |
| 885 | // A caller that set `max_body_bytes` never reaches here with a size |
| 886 | // that large -- the limit above trips first -- but a caller that |
| 887 | // trusts its peer and set no limit must still not be crashable by one. |
| 888 | let need = match size.checked_add(2) { |
| 889 | Some(n) => n, |
| 890 | None => return Err(err!( |
| 891 | "A chunked body declared a chunk of {} bytes, which is too \ |
| 892 | large to be a real length.", size; |
| 893 | IO, Network, Wire, Read, Invalid)), |
| 894 | }; |
| 895 | while raw.len() < need { |
| 896 | if !res!(fill::<BODY_CHUNK_SIZE, _>(stream.as_mut(), &mut raw).await) { |
| 897 | return Err(err!( |
| 898 | "The peer closed the connection {} bytes into a chunk it \ |
| 899 | said was {} bytes long.", raw.len(), size; |
| 900 | IO, Network, Wire, Read, Missing)); |
| 901 | } |
| 902 | } |
| 903 | out.extend_from_slice(&raw[..size]); |
| 904 | raw.drain(..need); |
| 905 | } |
| 906 | } |
| 907 | |
| 908 | /// Take the next CRLF-terminated line off the buffer, reading more from the |
| 909 | /// stream until there is one. `None` means the peer closed first. |
| 910 | /// |
| 911 | /// `max_len` rejects a line -- `TooBig` -- once `raw` passes it with no CRLF |
| 912 | /// found; `None` leaves the line unbounded. `what` names the line kind (e.g. |
| 913 | /// "chunk-size", "trailer") in that error. |
| 914 | /// |
| 915 | /// Only the bytes not yet searched, plus one byte behind them for a CRLF |
| 916 | /// split across a read boundary, are scanned on each pass -- rescanning the |
| 917 | /// whole buffer from the front on every fill is quadratic in a peer that |
| 918 | /// trickles one byte at a time. |
| 919 | #[cfg(feature = "async")] |
| 920 | async fn take_line< |
| 921 | const BODY_CHUNK_SIZE: usize, |
| 922 | R: AsyncRead + Unpin, |
| 923 | >( |
| 924 | mut stream: Pin<&mut R>, |
| 925 | raw: &mut Vec<u8>, |
| 926 | max_len: Option<usize>, |
| 927 | what: &str, |
| 928 | ) |
| 929 | -> Outcome<Option<String>> |
| 930 | { |
| 931 | let mut scanned = 0; // How much of `raw` has already turned up no CRLF. |
| 932 | loop { |
| 933 | if let Some(i) = raw[scanned..].windows(2).position(|w| w == b"\r\n") { |
| 934 | let i = scanned + i; |
| 935 | let line = String::from_utf8_lossy(&raw[..i]).to_string(); |
| 936 | raw.drain(..i + 2); |
| 937 | return Ok(Some(line)); |
| 938 | } |
| 939 | if let Some(lim) = max_len { |
| 940 | if raw.len() > lim { |
| 941 | return Err(err!( |
| 942 | "A chunked HTTP {} line ran to {} bytes with no \ |
| 943 | terminating CRLF, over the {}-byte limit.", |
| 944 | what, raw.len(), lim; |
| 945 | IO, Network, Input, TooBig)); |
| 946 | } |
| 947 | } |
| 948 | scanned = raw.len().saturating_sub(1); |
| 949 | if !res!(fill::<BODY_CHUNK_SIZE, _>(stream.as_mut(), raw).await) { |
| 950 | return Ok(None); |
| 951 | } |
| 952 | } |
| 953 | } |
| 954 | |
| 955 | /// Read one more chunk of bytes onto the buffer. `false` means the peer closed. |
| 956 | #[cfg(feature = "async")] |
| 957 | async fn fill< |
| 958 | const BODY_CHUNK_SIZE: usize, |
| 959 | R: AsyncRead + Unpin, |
| 960 | >( |
| 961 | mut stream: Pin<&mut R>, |
| 962 | raw: &mut Vec<u8>, |
| 963 | ) |
| 964 | -> Outcome<bool> |
| 965 | { |
| 966 | let mut buf = [0u8; BODY_CHUNK_SIZE]; |
| 967 | match stream.as_mut().read(&mut buf).await { |
| 968 | Ok(0) => Ok(false), |
| 969 | Ok(n) => { |
| 970 | raw.extend_from_slice(&buf[..n]); |
| 971 | Ok(true) |
| 972 | } |
| 973 | Err(e) => Err(err!(e, |
| 974 | "Reading a chunked HTTP body."; IO, Network, Wire, Read)), |
| 975 | } |
| 976 | } |
| 977 | |
| 978 | |
| 979 | #[cfg(all(test, feature = "async"))] |
| 980 | mod body_tests { |
| 981 | use super::*; |
| 982 | |
| 983 | /// Read one message off a buffer of wire bytes, as a client reading a |
| 984 | /// response does. |
| 985 | fn read_reply(wire: &str, is_request: Option<bool>) -> Outcome<Option<HttpMessage>> { |
| 986 | let bytes = wire.as_bytes(); |
| 987 | let mut stream = std::io::Cursor::new(bytes); |
| 988 | let rt = res!(tokio::runtime::Runtime::new()); |
| 989 | let (msg, _rest) = res!(rt.block_on(HttpMessage::read::< |
| 990 | { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE }, |
| 991 | { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE }, |
| 992 | _, |
| 993 | >( |
| 994 | Pin::new(&mut stream), |
| 995 | &Vec::new(), |
| 996 | is_request, |
| 997 | None, |
| 998 | ))); |
| 999 | Ok(msg) |
| 1000 | } |
| 1001 | |
| 1002 | /// Like `read_reply`, but with caller-chosen `ReadLimits` rather than none. |
| 1003 | fn read_reply_with_limits( |
| 1004 | wire: &str, |
| 1005 | is_request: Option<bool>, |
| 1006 | limits: &ReadLimits, |
| 1007 | ) -> Outcome<Option<HttpMessage>> { |
| 1008 | let bytes = wire.as_bytes(); |
| 1009 | let mut stream = std::io::Cursor::new(bytes); |
| 1010 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1011 | let (msg, _rest) = res!(rt.block_on(HttpMessage::read::< |
| 1012 | { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE }, |
| 1013 | { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE }, |
| 1014 | _, |
| 1015 | >( |
| 1016 | Pin::new(&mut stream), |
| 1017 | &Vec::new(), |
| 1018 | is_request, |
| 1019 | Some(limits), |
| 1020 | ))); |
| 1021 | Ok(msg) |
| 1022 | } |
| 1023 | |
| 1024 | /// Write a response to a buffer, as the server writing to a socket does. |
| 1025 | fn on_the_wire_bytes(msg: HttpMessage) -> Outcome<Vec<u8>> { |
| 1026 | let mut out: Vec<u8> = Vec::new(); |
| 1027 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1028 | res!(rt.block_on(msg.write_all(&mut out))); |
| 1029 | Ok(out) |
| 1030 | } |
| 1031 | |
| 1032 | fn on_the_wire(msg: HttpMessage) -> Outcome<String> { |
| 1033 | Ok(String::from_utf8_lossy(&res!(on_the_wire_bytes(msg))).to_string()) |
| 1034 | } |
| 1035 | |
| 1036 | /// A `HEAD` answer is the headers the matching `GET` would have sent, and nothing after them |
| 1037 | /// -- `Content-Length` included, still counting the body that was withheld. RFC 9110 9.3.2. |
| 1038 | #[test] |
| 1039 | fn test_a_head_answer_carries_the_length_and_no_body() -> Outcome<()> { |
| 1040 | let body = "<!doctype html><p>Twenty-nine bytes</p>"; |
| 1041 | let get = HttpMessage::respond_with_text(HttpStatus::OK, body); |
| 1042 | let head = HttpMessage::respond_with_text(HttpStatus::OK, body).head_only(); |
| 1043 | |
| 1044 | let get_wire = res!(on_the_wire(get)); |
| 1045 | let head_wire = res!(on_the_wire(head)); |
| 1046 | |
| 1047 | let n = fmt!("content-length: {}", body.len()); |
| 1048 | assert!(get_wire.to_lowercase().contains(&n), "the GET lost its length: {}", get_wire); |
| 1049 | // The same length, describing a body deliberately not sent -- a caller asking how big a |
| 1050 | // thing is deserves the true answer, which is the whole reason to ask with a HEAD. |
| 1051 | assert!(head_wire.to_lowercase().contains(&n), "the HEAD lost its length: {}", head_wire); |
| 1052 | assert!(!head_wire.contains(body), "the HEAD sent the body: {}", head_wire); |
| 1053 | // Headers to the blank line and not a byte past it. |
| 1054 | assert!(head_wire.ends_with("\r\n\r\n"), "the HEAD wrote past its headers: {:?}", head_wire); |
| 1055 | Ok(()) |
| 1056 | } |
| 1057 | |
| 1058 | /// A file that stands in for a large one: a counting pattern, so a window cut |
| 1059 | /// from the wrong offset is obvious rather than plausible. |
| 1060 | fn counting_file(name: &str, len: usize) -> Outcome<std::path::PathBuf> { |
| 1061 | let path = std::env::temp_dir().join(fmt!("fe2o3-window-{}-{}", std::process::id(), name)); |
| 1062 | let bytes: Vec<u8> = (0..len).map(|i| (i % 251) as u8).collect(); |
| 1063 | res!(std::fs::write(&path, &bytes), IO, File); |
| 1064 | Ok(path) |
| 1065 | } |
| 1066 | |
| 1067 | /// The body is a window of the file, and the message carries the window's |
| 1068 | /// length rather than the file's. |
| 1069 | #[test] |
| 1070 | fn test_a_file_window_writes_only_its_own_bytes() -> Outcome<()> { |
| 1071 | let path = res!(counting_file("window", 10_000)); |
| 1072 | let msg = HttpMessage::new_response(HttpStatus::PartialContent) |
| 1073 | .with_file_window(FileWindow::new(path.clone(), 4_000, 1_000)); |
| 1074 | assert_eq!(msg.body_len(), 1_000); |
| 1075 | |
| 1076 | let wire = res!(on_the_wire_bytes(msg)); |
| 1077 | let split = match wire.windows(4).position(|w| w == b"\r\n\r\n") { |
| 1078 | Some(i) => i + 4, |
| 1079 | None => return Err(err!("The response had no header terminator."; Test, Missing)), |
| 1080 | }; |
| 1081 | let head = String::from_utf8_lossy(&wire[..split]).to_lowercase(); |
| 1082 | assert!(head.contains("content-length: 1000"), "unexpected head: {}", head); |
| 1083 | |
| 1084 | let body = &wire[split..]; |
| 1085 | let expected: Vec<u8> = (4_000..5_000).map(|i: usize| (i % 251) as u8).collect(); |
| 1086 | assert_eq!(body.len(), 1_000, "the window was the wrong length"); |
| 1087 | assert_eq!(body, &expected[..], "the window came from the wrong offset"); |
| 1088 | |
| 1089 | let _ = std::fs::remove_file(&path); |
| 1090 | Ok(()) |
| 1091 | } |
| 1092 | |
| 1093 | /// A `HEAD` over a file window states the window's length and sends none of it, |
| 1094 | /// so a caller learns the size without the server reading the file at all. |
| 1095 | #[test] |
| 1096 | fn test_a_head_over_a_file_window_sends_no_bytes() -> Outcome<()> { |
| 1097 | let path = res!(counting_file("head", 10_000)); |
| 1098 | let msg = HttpMessage::new_response(HttpStatus::OK) |
| 1099 | .with_file_window(FileWindow::new(path.clone(), 0, 10_000)) |
| 1100 | .head_only(); |
| 1101 | let wire = res!(on_the_wire_bytes(msg)); |
| 1102 | let text = String::from_utf8_lossy(&wire).to_lowercase(); |
| 1103 | assert!(text.contains("content-length: 10000"), "unexpected head: {}", text); |
| 1104 | assert!(text.ends_with("\r\n\r\n"), "the HEAD wrote past its headers"); |
| 1105 | let _ = std::fs::remove_file(&path); |
| 1106 | Ok(()) |
| 1107 | } |
| 1108 | |
| 1109 | /// A window that runs past the end of the file is an error, not a short body: |
| 1110 | /// a body shorter than the `Content-Length` already promised desynchronises |
| 1111 | /// every message after it on a kept-alive connection. |
| 1112 | #[test] |
| 1113 | fn test_a_window_past_the_end_of_the_file_is_an_error() -> Outcome<()> { |
| 1114 | let path = res!(counting_file("short", 100)); |
| 1115 | let msg = HttpMessage::new_response(HttpStatus::PartialContent) |
| 1116 | .with_file_window(FileWindow::new(path.clone(), 0, 1_000)); |
| 1117 | assert!(on_the_wire_bytes(msg).is_err()); |
| 1118 | let _ = std::fs::remove_file(&path); |
| 1119 | Ok(()) |
| 1120 | } |
| 1121 | |
| 1122 | #[test] |
| 1123 | fn test_a_chunked_response_body_is_decoded() -> Outcome<()> { |
| 1124 | // Much of the web answers this way, and a client that cannot decode |
| 1125 | // it reads every such page as empty. |
| 1126 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1127 | Content-Type: text/html\r\n\ |
| 1128 | Transfer-Encoding: chunked\r\n\ |
| 1129 | \r\n\ |
| 1130 | 9\r\n<!doctype\r\n\ |
| 1131 | 6\r\n html>\r\n\ |
| 1132 | 0\r\n\r\n"; |
| 1133 | let msg = match res!(read_reply(wire, Some(false))) { |
| 1134 | Some(m) => m, |
| 1135 | None => return Err(err!( |
| 1136 | "A chunked response was read as no response at all."; |
| 1137 | Test, Missing)), |
| 1138 | }; |
| 1139 | assert_eq!(msg.body, b"<!doctype html>"); |
| 1140 | Ok(()) |
| 1141 | } |
| 1142 | |
| 1143 | #[test] |
| 1144 | fn test_a_chunk_may_carry_extensions_and_be_followed_by_trailers() -> Outcome<()> { |
| 1145 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1146 | Transfer-Encoding: chunked\r\n\ |
| 1147 | \r\n\ |
| 1148 | 5;name=value\r\nhello\r\n\ |
| 1149 | 0\r\n\ |
| 1150 | X-Checksum: 1234\r\n\ |
| 1151 | \r\n"; |
| 1152 | let msg = match res!(read_reply(wire, Some(false))) { |
| 1153 | Some(m) => m, |
| 1154 | None => return Err(err!( |
| 1155 | "A chunked response with trailers was read as no response."; |
| 1156 | Test, Missing)), |
| 1157 | }; |
| 1158 | assert_eq!(msg.body, b"hello"); |
| 1159 | Ok(()) |
| 1160 | } |
| 1161 | |
| 1162 | #[test] |
| 1163 | fn test_a_reply_to_head_states_a_length_it_never_sends() -> Outcome<()> { |
| 1164 | // A `HEAD` reply carries the `Content-Length` the body *would* have |
| 1165 | // had, and then no body. Waiting for those bytes waits forever, so a |
| 1166 | // client that treats the close as a failure cannot make the request |
| 1167 | // at all. |
| 1168 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1169 | Content-Type: text/html\r\n\ |
| 1170 | Content-Length: 5120\r\n\ |
| 1171 | \r\n"; |
| 1172 | let msg = match res!(read_reply(wire, Some(false))) { |
| 1173 | Some(m) => m, |
| 1174 | None => return Err(err!( |
| 1175 | "A reply to HEAD was read as no reply at all."; |
| 1176 | Test, Missing)), |
| 1177 | }; |
| 1178 | assert!(msg.body.is_empty()); |
| 1179 | Ok(()) |
| 1180 | } |
| 1181 | |
| 1182 | #[test] |
| 1183 | fn test_a_truncated_request_is_still_dropped() -> Outcome<()> { |
| 1184 | // The tolerance above is for responses only. A request whose body |
| 1185 | // stops short is a broken client, and is dropped as it always was. |
| 1186 | let wire = "POST /x HTTP/1.1\r\n\ |
| 1187 | Host: example.test\r\n\ |
| 1188 | Content-Length: 100\r\n\ |
| 1189 | \r\n\ |
| 1190 | short"; |
| 1191 | assert!(res!(read_reply(wire, None)).is_none()); |
| 1192 | assert!(res!(read_reply(wire, Some(true))).is_none()); |
| 1193 | Ok(()) |
| 1194 | } |
| 1195 | |
| 1196 | #[test] |
| 1197 | fn test_a_truncated_response_body_is_dropped_not_kept() -> Outcome<()> { |
| 1198 | // A response that delivered part of a declared body and then broke off |
| 1199 | // is a genuine truncation. The HEAD tolerance keys on an *empty* body, |
| 1200 | // so this partial one must not be waved through as complete. |
| 1201 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1202 | Content-Type: text/html\r\n\ |
| 1203 | Content-Length: 1000\r\n\ |
| 1204 | \r\n\ |
| 1205 | only these bytes arrived"; |
| 1206 | assert!(res!(read_reply(wire, Some(false))).is_none()); |
| 1207 | Ok(()) |
| 1208 | } |
| 1209 | |
| 1210 | /// A proxy reads a chunked response, and writes it on. The body is decoded |
| 1211 | /// by then, so the message must stop claiming to be chunked: emitting |
| 1212 | /// `Transfer-Encoding: chunked` and a `Content-Length` together is forbidden |
| 1213 | /// by RFC 9112 §6.1, and is the request-smuggling signal of §6.3, which is |
| 1214 | /// why browsers and proxies treat such a response as hostile. |
| 1215 | #[test] |
| 1216 | fn test_a_dechunked_body_is_not_re_emitted_as_chunked() -> Outcome<()> { |
| 1217 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1218 | Content-Type: text/html\r\n\ |
| 1219 | Transfer-Encoding: chunked\r\n\ |
| 1220 | \r\n\ |
| 1221 | 9\r\n<!doctype\r\n\ |
| 1222 | 6\r\n html>\r\n\ |
| 1223 | 0\r\n\r\n"; |
| 1224 | let msg = match res!(read_reply(wire, Some(false))) { |
| 1225 | Some(m) => m, |
| 1226 | None => return Err(err!( |
| 1227 | "A chunked response was read as no response at all."; |
| 1228 | Test, Missing)), |
| 1229 | }; |
| 1230 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1231 | let mut out: Vec<u8> = Vec::new(); |
| 1232 | res!(rt.block_on(msg.write_all(&mut out))); |
| 1233 | assert_eq!( |
| 1234 | out, |
| 1235 | b"HTTP/1.1 200 OK\r\n\ |
| 1236 | content-type: text/html\r\n\ |
| 1237 | content-length: 15\r\n\ |
| 1238 | \r\n\ |
| 1239 | <!doctype html>".to_vec(), |
| 1240 | ); |
| 1241 | Ok(()) |
| 1242 | } |
| 1243 | |
| 1244 | #[test] |
| 1245 | fn test_a_chunk_size_that_would_overflow_is_an_error_not_a_panic() -> Outcome<()> { |
| 1246 | // A hostile chunk-size line of `usize::MAX` once overflowed the `+ 2` |
| 1247 | // that reads past the chunk's own CRLF, panicking the reader. It is now |
| 1248 | // rejected as the impossible length it is. |
| 1249 | let wire = "HTTP/1.1 200 OK\r\n\ |
| 1250 | Transfer-Encoding: chunked\r\n\ |
| 1251 | \r\n\ |
| 1252 | ffffffffffffffff\r\n"; |
| 1253 | assert!(read_reply(wire, Some(false)).is_err()); |
| 1254 | Ok(()) |
| 1255 | } |
| 1256 | |
| 1257 | /// A chunk-size line is a hex length and maybe an extension: a few dozen |
| 1258 | /// bytes. A peer that keeps extending it with no CRLF is not describing a |
| 1259 | /// real chunk, it is growing the reader's buffer without limit, and gets |
| 1260 | /// `TooBig` well short of the whole message ever landing in memory. |
| 1261 | #[test] |
| 1262 | fn test_an_endless_chunk_size_line_is_capped() -> Outcome<()> { |
| 1263 | let junk = "a".repeat(constant::HTTP_CHUNK_LINE_MAX + 4_000); |
| 1264 | let wire = fmt!("HTTP/1.1 200 OK\r\n\ |
| 1265 | Transfer-Encoding: chunked\r\n\ |
| 1266 | \r\n\ |
| 1267 | {}", junk); |
| 1268 | match read_reply(&wire, Some(false)) { |
| 1269 | Err(e) => assert!(e.tags().contains(&ErrTag::TooBig), |
| 1270 | "an endless chunk-size line raised the wrong error: {}", e), |
| 1271 | Ok(_) => return Err(err!( |
| 1272 | "An endless chunk-size line was read as if it had a CRLF."; |
| 1273 | Test, Unreachable)), |
| 1274 | } |
| 1275 | Ok(()) |
| 1276 | } |
| 1277 | |
| 1278 | /// A trailer is one more header-like field, so it is bounded the same way |
| 1279 | /// the header block is: by `max_header_bytes`. Without the cap this loops |
| 1280 | /// forever on the never-terminated line above `read_chunked`'s zero chunk. |
| 1281 | #[test] |
| 1282 | fn test_an_endless_trailer_line_is_capped() -> Outcome<()> { |
| 1283 | // Comfortably above `HTTP_DEFAULT_HEADER_CHUNK_SIZE` (1,500), so the header |
| 1284 | // phase's own bulk-read check never trips on a read that happens to reach |
| 1285 | // past the header terminator into this body; comfortably below the junk |
| 1286 | // length below, so it is this trailer cap, not that one, doing the work. |
| 1287 | let limits = ReadLimits { max_header_bytes: Some(2_000), ..Default::default() }; |
| 1288 | let junk = "a".repeat(5_000); |
| 1289 | let wire = fmt!("HTTP/1.1 200 OK\r\n\ |
| 1290 | Transfer-Encoding: chunked\r\n\ |
| 1291 | \r\n\ |
| 1292 | 0\r\n\ |
| 1293 | {}", junk); |
| 1294 | match read_reply_with_limits(&wire, Some(false), &limits) { |
| 1295 | Err(e) => assert!(e.tags().contains(&ErrTag::TooBig), |
| 1296 | "an endless trailer line raised the wrong error: {}", e), |
| 1297 | Ok(_) => return Err(err!( |
| 1298 | "An endless trailer line was read as if it had a CRLF."; |
| 1299 | Test, Unreachable)), |
| 1300 | } |
| 1301 | Ok(()) |
| 1302 | } |
| 1303 | } |
| 1304 | |
| 1305 | |
| 1306 | #[cfg(feature = "async")] |
| 1307 | pub struct HttpMessageReader< |
| 1308 | 'a, |
| 1309 | const HEADER_CHUNK_SIZE: usize, |
| 1310 | const BODY_CHUNK_SIZE: usize, |
| 1311 | R: AsyncRead + Unpin + Send, |
| 1312 | > { |
| 1313 | stream: Pin<&'a mut R>, |
| 1314 | buffer: Vec<u8>, |
| 1315 | limits: Option<ReadLimits>, |
| 1316 | } |
| 1317 | |
| 1318 | #[cfg(feature = "async")] |
| 1319 | impl< |
| 1320 | 'a, |
| 1321 | const HEADER_CHUNK_SIZE: usize, |
| 1322 | const BODY_CHUNK_SIZE: usize, |
| 1323 | R: AsyncRead + Unpin + Send, |
| 1324 | > |
| 1325 | HttpMessageReader<'a, HEADER_CHUNK_SIZE, BODY_CHUNK_SIZE, R> |
| 1326 | { |
| 1327 | /// Build a reader with no read limits applied. Suitable for |
| 1328 | /// outbound HTTP clients, tests, and anywhere the peer is trusted. |
| 1329 | pub fn new( |
| 1330 | stream: Pin<&'a mut R>, |
| 1331 | ) |
| 1332 | -> Self |
| 1333 | { |
| 1334 | Self { |
| 1335 | stream, |
| 1336 | buffer: Vec::new(), |
| 1337 | limits: None, |
| 1338 | } |
| 1339 | } |
| 1340 | |
| 1341 | /// Build a reader that enforces the supplied limits on every |
| 1342 | /// successive read. Use from HTTPS accept paths where the peer |
| 1343 | /// is untrusted. |
| 1344 | pub fn with_limits( |
| 1345 | stream: Pin<&'a mut R>, |
| 1346 | limits: ReadLimits, |
| 1347 | ) |
| 1348 | -> Self |
| 1349 | { |
| 1350 | Self { |
| 1351 | stream, |
| 1352 | buffer: Vec::new(), |
| 1353 | limits: Some(limits), |
| 1354 | } |
| 1355 | } |
| 1356 | |
| 1357 | /// The bytes taken off the stream past the end of the last message. |
| 1358 | /// |
| 1359 | /// A reader in a `keep-alive` loop needs no accessor: the remnant is the |
| 1360 | /// front of the next message and the next [`AsyncReadIterator::next`] |
| 1361 | /// consumes it. A caller that stops reading HTTP and takes the raw socket |
| 1362 | /// does need one -- a protocol upgrade, where the `101` is the last HTTP on |
| 1363 | /// the connection and everything after it belongs to another protocol. |
| 1364 | /// |
| 1365 | /// Without this the remnant is dropped with the reader. That is invisible |
| 1366 | /// almost always, because a client waits for the `101` before it frames |
| 1367 | /// anything, and it is a truncated stream on the occasion the client |
| 1368 | /// pipelines its first frame or the network delivers both in one segment. |
| 1369 | /// Occasional silent truncation in a byte pipe is the shape of bug nobody |
| 1370 | /// can report usefully, so the bytes are made reachable rather than left to |
| 1371 | /// each caller to notice. |
| 1372 | pub fn remnant(&self) -> &[u8] { |
| 1373 | &self.buffer |
| 1374 | } |
| 1375 | |
| 1376 | /// The remnant, taken. For a caller that is finished with the reader, which |
| 1377 | /// is the usual case at an upgrade: the socket outlives the reader and the |
| 1378 | /// bytes have to go with the socket. |
| 1379 | pub fn into_remnant(self) -> Vec<u8> { |
| 1380 | self.buffer |
| 1381 | } |
| 1382 | } |
| 1383 | |
| 1384 | #[cfg(feature = "async")] |
| 1385 | impl< |
| 1386 | 'a, |
| 1387 | const HEADER_CHUNK_SIZE: usize, |
| 1388 | const BODY_CHUNK_SIZE: usize, |
| 1389 | R: AsyncRead + Unpin + Send, |
| 1390 | > |
| 1391 | AsyncReadIterator for HttpMessageReader<'a, HEADER_CHUNK_SIZE, BODY_CHUNK_SIZE, R> |
| 1392 | { |
| 1393 | type Item = Outcome<HttpMessage>; |
| 1394 | |
| 1395 | fn next<'b>(&'b mut self) -> Pin<Box<dyn Future<Output = Option<Self::Item>> + Send + 'b>> { |
| 1396 | let mut stream = self.stream.as_mut(); |
| 1397 | let buffer = &mut self.buffer; |
| 1398 | let limits = self.limits.as_ref(); |
| 1399 | |
| 1400 | Box::pin(async move { |
| 1401 | let result = HttpMessage::read::<HEADER_CHUNK_SIZE, BODY_CHUNK_SIZE, _>( |
| 1402 | stream.as_mut(), |
| 1403 | buffer, |
| 1404 | None, |
| 1405 | limits, |
| 1406 | ) |
| 1407 | .await; |
| 1408 | |
| 1409 | match result { |
| 1410 | Ok((Some(message), remnant)) => { |
| 1411 | *buffer = remnant; |
| 1412 | trace!("Remnant = {} bytes", buffer.len()); |
| 1413 | Some(Ok(message)) |
| 1414 | } |
| 1415 | Ok((None, _)) => None, |
| 1416 | Err(e) => Some(Err(e)), |
| 1417 | } |
| 1418 | }) |
| 1419 | } |
| 1420 | } |
| 1421 | |
| 1422 | |
| 1423 | #[cfg(all(test, feature = "async"))] |
| 1424 | mod reader_tests { |
| 1425 | use super::*; |
| 1426 | |
| 1427 | use crate::conc::AsyncReadIterator; |
| 1428 | |
| 1429 | /// A WebSocket upgrade with the client's first frame already behind it in the |
| 1430 | /// same buffer, which is what a pipelining client or one obliging TCP segment |
| 1431 | /// delivers. |
| 1432 | fn upgrade_then_frame() -> (Vec<u8>, Vec<u8>) { |
| 1433 | let head = "GET /api/mail/tunnel?host=imap.example.com&port=993 HTTP/1.1\r\n\ |
| 1434 | Host: example.com\r\n\ |
| 1435 | Upgrade: websocket\r\n\ |
| 1436 | Connection: Upgrade\r\n\ |
| 1437 | Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\n\ |
| 1438 | Sec-WebSocket-Version: 13\r\n\r\n"; |
| 1439 | // A masked binary frame of four bytes, as a client must send. |
| 1440 | let frame: Vec<u8> = vec![ |
| 1441 | 0x82, 0x84, 0x01, 0x02, 0x03, 0x04, |
| 1442 | 0x16, 0x03, 0x04, 0x05, |
| 1443 | ]; |
| 1444 | let mut wire = head.as_bytes().to_vec(); |
| 1445 | wire.extend_from_slice(&frame); |
| 1446 | (wire, frame) |
| 1447 | } |
| 1448 | |
| 1449 | /// The bytes past the header block survive the read, so a caller taking the |
| 1450 | /// raw socket at an upgrade can go on where the reader stopped. |
| 1451 | /// |
| 1452 | /// The reader fills a header-sized chunk in one go, so those bytes have |
| 1453 | /// already left the stream by the time the request is parsed. Dropped, they |
| 1454 | /// are the front of the next protocol's first message -- and a stream missing |
| 1455 | /// its first frame does not fail, it decodes to nonsense. |
| 1456 | #[test] |
| 1457 | fn test_the_bytes_after_an_upgrade_survive_the_read() -> Outcome<()> { |
| 1458 | let (wire, frame) = upgrade_then_frame(); |
| 1459 | let mut stream = std::io::Cursor::new(wire); |
| 1460 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1461 | let mut reader: HttpMessageReader< |
| 1462 | '_, |
| 1463 | { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE }, |
| 1464 | { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE }, |
| 1465 | _, |
| 1466 | > = HttpMessageReader::new(Pin::new(&mut stream)); |
| 1467 | |
| 1468 | let msg = match rt.block_on(reader.next()) { |
| 1469 | Some(Ok(m)) => m, |
| 1470 | Some(Err(e)) => return Err(err!(e, "The upgrade request did not parse."; Test)), |
| 1471 | None => return Err(err!("The reader found no request at all."; Test, Missing)), |
| 1472 | }; |
| 1473 | assert!(msg.is_websocket_upgrade(), "the request was not read as an upgrade"); |
| 1474 | assert_eq!(reader.remnant(), &frame[..], |
| 1475 | "the frame behind the header block was dropped by the reader"); |
| 1476 | assert_eq!(reader.into_remnant(), frame, |
| 1477 | "the taken remnant differs from the borrowed one"); |
| 1478 | Ok(()) |
| 1479 | } |
| 1480 | |
| 1481 | /// A request with nothing behind it leaves an empty remnant, not a stray |
| 1482 | /// byte: the accessor must not turn the usual case into a corrupted stream. |
| 1483 | #[test] |
| 1484 | fn test_a_request_with_nothing_behind_it_leaves_nothing() -> Outcome<()> { |
| 1485 | let (wire, _) = upgrade_then_frame(); |
| 1486 | let head_end = match wire.windows(4).position(|w| w == b"\r\n\r\n") { |
| 1487 | Some(i) => i + 4, |
| 1488 | None => return Err(err!("The fixture has no header terminator."; Test, Missing)), |
| 1489 | }; |
| 1490 | let mut stream = std::io::Cursor::new(wire[..head_end].to_vec()); |
| 1491 | let rt = res!(tokio::runtime::Runtime::new()); |
| 1492 | let mut reader: HttpMessageReader< |
| 1493 | '_, |
| 1494 | { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE }, |
| 1495 | { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE }, |
| 1496 | _, |
| 1497 | > = HttpMessageReader::new(Pin::new(&mut stream)); |
| 1498 | match rt.block_on(reader.next()) { |
| 1499 | Some(Ok(_)) => (), |
| 1500 | Some(Err(e)) => return Err(err!(e, "The upgrade request did not parse."; Test)), |
| 1501 | None => return Err(err!("The reader found no request."; Test, Missing)), |
| 1502 | } |
| 1503 | assert!(reader.remnant().is_empty(), |
| 1504 | "a request with nothing behind it left {} byte(s) of remnant", |
| 1505 | reader.remnant().len()); |
| 1506 | Ok(()) |
| 1507 | } |
| 1508 | } |