Oregami
Repositories/oxedyne/fe2o3

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")]
2use crate::conc::AsyncReadIterator;
3use 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
28use oxedyne_fe2o3_core::prelude::*;
29
30use 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")]
39use std::{
40 future::Future,
41 pin::Pin,
42};
43
44#[cfg(feature = "async")]
45use 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)]
73pub 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
90impl 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)]
108pub 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
117impl 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)]
203pub 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
220impl 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.
798pub 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")]
816async 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")]
920async 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")]
957async 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"))]
980mod 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")]
1307pub 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")]
1319impl<
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")]
1385impl<
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"))]
1424mod 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}