Skip to main content

hyper/proto/h1/
io.rs

1use std::cmp;
2use std::fmt;
3use std::io::{self, IoSlice};
4use std::pin::Pin;
5use std::task::{Context, Poll};
6
7use crate::rt::{Read, ReadBuf, Write};
8use bytes::{Buf, BufMut, Bytes, BytesMut};
9use futures_core::ready;
10
11use super::{Http1Transaction, ParseContext, ParsedMessage};
12use crate::common::buf::BufList;
13
14/// The initial buffer size allocated before trying to read from IO.
15pub(crate) const INIT_BUFFER_SIZE: usize = 8192;
16
17/// The minimum value that can be set to max buffer size.
18pub(crate) const MINIMUM_MAX_BUFFER_SIZE: usize = INIT_BUFFER_SIZE;
19
20/// The default maximum read buffer size. If the buffer gets this big and
21/// a message is still not complete, a `TooLarge` error is triggered.
22// Note: if this changes, update server::conn::Http::max_buf_size docs.
23pub(crate) const DEFAULT_MAX_BUFFER_SIZE: usize = 8192 + 4096 * 100;
24
25/// The maximum number of distinct `Buf`s to hold in a list before requiring
26/// a flush. Only affects when the buffer strategy is to queue buffers.
27///
28/// Note that a flush can happen before reaching the maximum. This simply
29/// forces a flush if the queue gets this big.
30const MAX_BUF_LIST_BUFFERS: usize = 16;
31
32pub(crate) struct Buffered<T, B> {
33    flush_pipeline: bool,
34    io: T,
35    partial_len: Option<usize>,
36    read_blocked: bool,
37    read_buf: BytesMut,
38    read_buf_strategy: ReadStrategy,
39    write_buf: WriteBuf<B>,
40}
41
42impl<T, B> fmt::Debug for Buffered<T, B>
43where
44    B: Buf,
45{
46    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
47        f.debug_struct("Buffered")
48            .field("read_buf", &self.read_buf)
49            .field("write_buf", &self.write_buf)
50            .finish()
51    }
52}
53
54impl<T, B> Buffered<T, B>
55where
56    T: Read + Write + Unpin,
57    B: Buf,
58{
59    pub(crate) fn new(io: T) -> Buffered<T, B> {
60        let strategy = if io.is_write_vectored() {
61            WriteStrategy::Queue
62        } else {
63            WriteStrategy::Flatten
64        };
65        let write_buf = WriteBuf::new(strategy);
66        Buffered {
67            flush_pipeline: false,
68            io,
69            partial_len: None,
70            read_blocked: false,
71            read_buf: BytesMut::with_capacity(0),
72            read_buf_strategy: ReadStrategy::default(),
73            write_buf,
74        }
75    }
76
77    #[cfg(feature = "server")]
78    pub(crate) fn set_flush_pipeline(&mut self, enabled: bool) {
79        debug_assert!(!self.write_buf.has_remaining());
80        self.flush_pipeline = enabled;
81        if enabled {
82            self.set_write_strategy_flatten();
83        }
84    }
85
86    pub(crate) fn set_max_buf_size(&mut self, max: usize) {
87        assert!(
88            max >= MINIMUM_MAX_BUFFER_SIZE,
89            "The max_buf_size cannot be smaller than {MINIMUM_MAX_BUFFER_SIZE}.",
90        );
91        self.read_buf_strategy = ReadStrategy::with_max(max);
92        self.write_buf.max_buf_size = max;
93    }
94
95    #[cfg(feature = "client")]
96    pub(crate) fn set_read_buf_exact_size(&mut self, sz: usize) {
97        self.read_buf_strategy = ReadStrategy::Exact(sz);
98    }
99
100    pub(crate) fn set_write_strategy_flatten(&mut self) {
101        // this should always be called only at construction time,
102        // so this assert is here to catch myself
103        debug_assert_eq!(self.write_buf.queue.bufs_cnt(), 0);
104        self.write_buf.set_strategy(WriteStrategy::Flatten);
105    }
106
107    pub(crate) fn set_write_strategy_queue(&mut self) {
108        // this should always be called only at construction time,
109        // so this assert is here to catch myself
110        debug_assert_eq!(self.write_buf.queue.bufs_cnt(), 0);
111        self.write_buf.set_strategy(WriteStrategy::Queue);
112    }
113
114    pub(crate) fn read_buf(&self) -> &[u8] {
115        self.read_buf.as_ref()
116    }
117
118    #[cfg(test)]
119    #[cfg(feature = "nightly")]
120    pub(super) fn read_buf_mut(&mut self) -> &mut BytesMut {
121        &mut self.read_buf
122    }
123
124    /// Return the "allocated" available space, not the potential space
125    /// that could be allocated in the future.
126    fn read_buf_remaining_mut(&self) -> usize {
127        self.read_buf.capacity() - self.read_buf.len()
128    }
129
130    /// Return whether we can append to the headers buffer.
131    ///
132    /// Reasons we can't:
133    /// - The write buf is in queue mode, and some of the past body is still
134    ///   needing to be flushed.
135    pub(crate) fn can_headers_buf(&self) -> bool {
136        !self.write_buf.queue.has_remaining()
137    }
138
139    pub(crate) fn headers_buf(&mut self) -> &mut Vec<u8> {
140        let buf = self.write_buf.headers_mut();
141        &mut buf.bytes
142    }
143
144    pub(super) fn write_buf(&mut self) -> &mut WriteBuf<B> {
145        &mut self.write_buf
146    }
147
148    pub(crate) fn buffer<BB: Buf + Into<B>>(&mut self, buf: BB) {
149        self.write_buf.buffer(buf);
150    }
151
152    pub(crate) fn can_buffer(&self) -> bool {
153        self.flush_pipeline || self.write_buf.can_buffer()
154    }
155
156    pub(crate) fn consume_leading_lines(&mut self) {
157        if !self.read_buf.is_empty() {
158            let mut i = 0;
159            while i < self.read_buf.len() {
160                match self.read_buf[i] {
161                    b'\r' | b'\n' => i += 1,
162                    _ => break,
163                }
164            }
165            self.read_buf.advance(i);
166        }
167    }
168
169    pub(super) fn parse<S>(
170        &mut self,
171        cx: &mut Context<'_>,
172        parse_ctx: ParseContext<'_>,
173    ) -> Poll<crate::Result<ParsedMessage<S::Incoming>>>
174    where
175        S: Http1Transaction,
176    {
177        loop {
178            if let Some(msg) = super::role::parse_headers::<S>(
179                &mut self.read_buf,
180                self.partial_len,
181                ParseContext {
182                    cached_headers: parse_ctx.cached_headers,
183                    req_method: parse_ctx.req_method,
184                    h1_parser_config: parse_ctx.h1_parser_config.clone(),
185                    h1_max_headers: parse_ctx.h1_max_headers,
186                    preserve_header_case: parse_ctx.preserve_header_case,
187                    #[cfg(feature = "ffi")]
188                    preserve_header_order: parse_ctx.preserve_header_order,
189                    h09_responses: parse_ctx.h09_responses,
190                    #[cfg(feature = "client")]
191                    on_informational: parse_ctx.on_informational,
192                },
193            )? {
194                debug!("parsed {} headers", msg.head.headers.len());
195                self.partial_len = None;
196                return Poll::Ready(Ok(msg));
197            } else {
198                let max = self.read_buf_strategy.max();
199                let curr_len = self.read_buf.len();
200                if curr_len >= max {
201                    debug!("max_buf_size ({}) reached, closing", max);
202                    return Poll::Ready(Err(crate::Error::new_too_large()));
203                }
204                if curr_len > 0 {
205                    trace!("partial headers; {} bytes so far", curr_len);
206                    self.partial_len = Some(curr_len);
207                } else {
208                    // 1xx gobled some bytes
209                    self.partial_len = None;
210                }
211            }
212            if ready!(self.poll_read_from_io(cx)).map_err(crate::Error::new_io)? == 0 {
213                trace!("parse eof");
214                return Poll::Ready(Err(crate::Error::new_incomplete()));
215            }
216        }
217    }
218
219    pub(crate) fn poll_read_from_io(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<usize>> {
220        self.read_blocked = false;
221        // Get the next amount to allocate, but make sure we don't go over
222        // the max read buf size configured.
223        let next = cmp::min(
224            self.read_buf_strategy.next(),
225            self.read_buf_strategy
226                .max()
227                .saturating_sub(self.read_buf.len()),
228        );
229        if self.read_buf_remaining_mut() < next {
230            self.read_buf.reserve(next);
231        }
232
233        // SAFETY: ReadBuf and poll_read promise not to set any uninitialized
234        // bytes onto `dst`.
235        let dst = unsafe { self.read_buf.chunk_mut().as_uninit_slice_mut() };
236        let mut buf = ReadBuf::uninit(dst);
237        match Pin::new(&mut self.io).poll_read(cx, buf.unfilled()) {
238            Poll::Ready(Ok(_)) => {
239                let n = buf.filled().len();
240                trace!("received {} bytes", n);
241                // Safety: we just read that many bytes into the
242                // uninitialized part of the buffer, so this is okay.
243                // @tokio pls give me back `poll_read_buf` thanks
244                unsafe {
245                    self.read_buf.advance_mut(n);
246                }
247                self.read_buf_strategy.record(n);
248                Poll::Ready(Ok(n))
249            }
250            Poll::Pending => {
251                self.read_blocked = true;
252                Poll::Pending
253            }
254            Poll::Ready(Err(e)) => Poll::Ready(Err(e)),
255        }
256    }
257
258    pub(crate) fn into_inner(self) -> (T, Bytes) {
259        (self.io, self.read_buf.freeze())
260    }
261
262    pub(crate) fn is_read_blocked(&self) -> bool {
263        self.read_blocked
264    }
265
266    pub(crate) fn poll_flush(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
267        if self.flush_pipeline && !self.read_buf.is_empty() {
268            Poll::Ready(Ok(()))
269        } else if self.write_buf.remaining() == 0 {
270            Pin::new(&mut self.io).poll_flush(cx)
271        } else {
272            if let WriteStrategy::Flatten = self.write_buf.strategy {
273                return self.poll_flush_flattened(cx);
274            }
275
276            const MAX_WRITEV_BUFS: usize = 64;
277            loop {
278                let n = {
279                    let mut iovs = [IoSlice::new(&[]); MAX_WRITEV_BUFS];
280                    let len = self.write_buf.chunks_vectored(&mut iovs);
281                    ready!(Pin::new(&mut self.io).poll_write_vectored(cx, &iovs[..len]))?
282                };
283                // TODO(eliza): we have to do this manually because
284                // `poll_write_buf` doesn't exist in Tokio 0.3 yet...when
285                // `poll_write_buf` comes back, the manual advance will need to leave!
286                self.write_buf.advance(n);
287                debug!("flushed {} bytes", n);
288                if self.write_buf.remaining() == 0 {
289                    break;
290                } else if n == 0 {
291                    trace!(
292                        "write returned zero, but {} bytes remaining",
293                        self.write_buf.remaining()
294                    );
295                    return Poll::Ready(Err(io::ErrorKind::WriteZero.into()));
296                }
297            }
298            Pin::new(&mut self.io).poll_flush(cx)
299        }
300    }
301
302    /// Specialized version of `flush` when strategy is Flatten.
303    ///
304    /// Since all buffered bytes are flattened into the single headers buffer,
305    /// that skips some bookkeeping around using multiple buffers.
306    fn poll_flush_flattened(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
307        loop {
308            let n = ready!(Pin::new(&mut self.io).poll_write(cx, self.write_buf.headers.chunk()))?;
309            debug!("flushed {} bytes", n);
310            self.write_buf.headers.advance(n);
311            if self.write_buf.headers.remaining() == 0 {
312                self.write_buf.headers.reset();
313                break;
314            } else if n == 0 {
315                trace!(
316                    "write returned zero, but {} bytes remaining",
317                    self.write_buf.remaining()
318                );
319                return Poll::Ready(Err(io::ErrorKind::WriteZero.into()));
320            }
321        }
322        Pin::new(&mut self.io).poll_flush(cx)
323    }
324
325    pub(crate) fn poll_shutdown(&mut self, cx: &mut Context<'_>) -> Poll<io::Result<()>> {
326        ready!(self.poll_flush(cx))?;
327        Pin::new(&mut self.io).poll_shutdown(cx)
328    }
329
330    #[cfg(test)]
331    fn flush(&mut self) -> impl std::future::Future<Output = io::Result<()>> + '_ {
332        futures_util::future::poll_fn(move |cx| self.poll_flush(cx))
333    }
334}
335
336// The `B` is a `Buf`, we never project a pin to it
337impl<T: Unpin, B> Unpin for Buffered<T, B> {}
338
339// TODO: This trait is old... at least rename to PollBytes or something...
340pub(crate) trait MemRead {
341    fn read_mem(&mut self, cx: &mut Context<'_>, len: usize) -> Poll<io::Result<Bytes>>;
342}
343
344impl<T, B> MemRead for Buffered<T, B>
345where
346    T: Read + Write + Unpin,
347    B: Buf,
348{
349    fn read_mem(&mut self, cx: &mut Context<'_>, len: usize) -> Poll<io::Result<Bytes>> {
350        if !self.read_buf.is_empty() {
351            let n = std::cmp::min(len, self.read_buf.len());
352            Poll::Ready(Ok(self.read_buf.split_to(n).freeze()))
353        } else {
354            let n = ready!(self.poll_read_from_io(cx))?;
355            Poll::Ready(Ok(self.read_buf.split_to(::std::cmp::min(len, n)).freeze()))
356        }
357    }
358}
359
360#[derive(Clone, Copy, Debug)]
361enum ReadStrategy {
362    Adaptive {
363        decrease_now: bool,
364        next: usize,
365        max: usize,
366    },
367    #[cfg(feature = "client")]
368    Exact(usize),
369}
370
371impl ReadStrategy {
372    fn with_max(max: usize) -> ReadStrategy {
373        ReadStrategy::Adaptive {
374            decrease_now: false,
375            next: INIT_BUFFER_SIZE,
376            max,
377        }
378    }
379
380    fn next(&self) -> usize {
381        match *self {
382            ReadStrategy::Adaptive { next, .. } => next,
383            #[cfg(feature = "client")]
384            ReadStrategy::Exact(exact) => exact,
385        }
386    }
387
388    fn max(&self) -> usize {
389        match *self {
390            ReadStrategy::Adaptive { max, .. } => max,
391            #[cfg(feature = "client")]
392            ReadStrategy::Exact(exact) => exact,
393        }
394    }
395
396    fn record(&mut self, bytes_read: usize) {
397        match *self {
398            ReadStrategy::Adaptive {
399                ref mut decrease_now,
400                ref mut next,
401                max,
402                ..
403            } => {
404                if bytes_read >= *next {
405                    *next = cmp::min(incr_power_of_two(*next), max);
406                    *decrease_now = false;
407                } else {
408                    let decr_to = prev_power_of_two(*next);
409                    if bytes_read < decr_to {
410                        if *decrease_now {
411                            *next = cmp::max(decr_to, INIT_BUFFER_SIZE);
412                            *decrease_now = false;
413                        } else {
414                            // Decreasing is a two "record" process.
415                            *decrease_now = true;
416                        }
417                    } else {
418                        // A read within the current range should cancel
419                        // a potential decrease, since we just saw proof
420                        // that we still need this size.
421                        *decrease_now = false;
422                    }
423                }
424            }
425            #[cfg(feature = "client")]
426            ReadStrategy::Exact(_) => (),
427        }
428    }
429}
430
431fn incr_power_of_two(n: usize) -> usize {
432    n.saturating_mul(2)
433}
434
435fn prev_power_of_two(n: usize) -> usize {
436    // Only way this shift can underflow is if n is less than 4.
437    // (Which would means `usize::MAX >> 64` and underflowed!)
438    debug_assert!(n >= 4);
439    (usize::MAX >> (n.leading_zeros() + 2)) + 1
440}
441
442impl Default for ReadStrategy {
443    fn default() -> ReadStrategy {
444        ReadStrategy::with_max(DEFAULT_MAX_BUFFER_SIZE)
445    }
446}
447
448#[derive(Clone)]
449pub(crate) struct Cursor<T> {
450    bytes: T,
451    pos: usize,
452}
453
454impl<T: AsRef<[u8]>> Cursor<T> {
455    #[inline]
456    pub(crate) fn new(bytes: T) -> Cursor<T> {
457        Cursor { bytes, pos: 0 }
458    }
459}
460
461impl Cursor<Vec<u8>> {
462    /// If we've advanced the position a bit in this cursor, and wish to
463    /// extend the underlying vector, we may wish to unshift the "read" bytes
464    /// off, and move everything else over.
465    fn maybe_unshift(&mut self, additional: usize) {
466        if self.pos == 0 {
467            // nothing to do
468            return;
469        }
470
471        if self.bytes.capacity() - self.bytes.len() >= additional {
472            // there's room!
473            return;
474        }
475
476        self.bytes.drain(0..self.pos);
477        self.pos = 0;
478    }
479
480    fn reset(&mut self) {
481        self.pos = 0;
482        self.bytes.clear();
483    }
484}
485
486impl<T: AsRef<[u8]>> fmt::Debug for Cursor<T> {
487    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
488        f.debug_struct("Cursor")
489            .field("pos", &self.pos)
490            .field("len", &self.bytes.as_ref().len())
491            .finish()
492    }
493}
494
495impl<T: AsRef<[u8]>> Buf for Cursor<T> {
496    #[inline]
497    fn remaining(&self) -> usize {
498        self.bytes.as_ref().len() - self.pos
499    }
500
501    #[inline]
502    fn chunk(&self) -> &[u8] {
503        &self.bytes.as_ref()[self.pos..]
504    }
505
506    #[inline]
507    fn advance(&mut self, cnt: usize) {
508        debug_assert!(self.pos + cnt <= self.bytes.as_ref().len());
509        self.pos += cnt;
510    }
511}
512
513// an internal buffer to collect writes before flushes
514pub(super) struct WriteBuf<B> {
515    /// Re-usable buffer that holds message headers.
516    headers: Cursor<Vec<u8>>,
517    max_buf_size: usize,
518    /// Deque of user buffers if strategy is `Queue`.
519    queue: BufList<B>,
520    strategy: WriteStrategy,
521}
522
523impl<B: Buf> WriteBuf<B> {
524    fn new(strategy: WriteStrategy) -> WriteBuf<B> {
525        WriteBuf {
526            headers: Cursor::new(Vec::with_capacity(INIT_BUFFER_SIZE)),
527            max_buf_size: DEFAULT_MAX_BUFFER_SIZE,
528            queue: BufList::new(),
529            strategy,
530        }
531    }
532}
533
534impl<B> WriteBuf<B>
535where
536    B: Buf,
537{
538    fn set_strategy(&mut self, strategy: WriteStrategy) {
539        self.strategy = strategy;
540    }
541
542    pub(super) fn buffer<BB: Buf + Into<B>>(&mut self, mut buf: BB) {
543        debug_assert!(buf.has_remaining());
544        match self.strategy {
545            WriteStrategy::Flatten => {
546                let head = self.headers_mut();
547
548                head.maybe_unshift(buf.remaining());
549                trace!(
550                    self.len = head.remaining(),
551                    buf.len = buf.remaining(),
552                    "buffer.flatten"
553                );
554                //perf: This is a little faster than <Vec as BufMut>>::put,
555                //but accomplishes the same result.
556                loop {
557                    let adv = {
558                        let slice = buf.chunk();
559                        if slice.is_empty() {
560                            return;
561                        }
562                        head.bytes.extend_from_slice(slice);
563                        slice.len()
564                    };
565                    buf.advance(adv);
566                }
567            }
568            WriteStrategy::Queue => {
569                trace!(
570                    self.len = self.remaining(),
571                    buf.len = buf.remaining(),
572                    "buffer.queue"
573                );
574                self.queue.push(buf.into());
575            }
576        }
577    }
578
579    fn can_buffer(&self) -> bool {
580        match self.strategy {
581            WriteStrategy::Flatten => self.remaining() < self.max_buf_size,
582            WriteStrategy::Queue => {
583                self.queue.bufs_cnt() < MAX_BUF_LIST_BUFFERS && self.remaining() < self.max_buf_size
584            }
585        }
586    }
587
588    fn headers_mut(&mut self) -> &mut Cursor<Vec<u8>> {
589        debug_assert!(!self.queue.has_remaining());
590        &mut self.headers
591    }
592}
593
594impl<B: Buf> fmt::Debug for WriteBuf<B> {
595    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
596        f.debug_struct("WriteBuf")
597            .field("remaining", &self.remaining())
598            .field("strategy", &self.strategy)
599            .finish()
600    }
601}
602
603impl<B: Buf> Buf for WriteBuf<B> {
604    #[inline]
605    fn remaining(&self) -> usize {
606        self.headers.remaining() + self.queue.remaining()
607    }
608
609    #[inline]
610    fn chunk(&self) -> &[u8] {
611        let headers = self.headers.chunk();
612        if !headers.is_empty() {
613            headers
614        } else {
615            self.queue.chunk()
616        }
617    }
618
619    #[inline]
620    fn advance(&mut self, cnt: usize) {
621        let hrem = self.headers.remaining();
622
623        match hrem.cmp(&cnt) {
624            cmp::Ordering::Equal => self.headers.reset(),
625            cmp::Ordering::Greater => self.headers.advance(cnt),
626            cmp::Ordering::Less => {
627                let qcnt = cnt - hrem;
628                self.headers.reset();
629                self.queue.advance(qcnt);
630            }
631        }
632    }
633
634    #[inline]
635    fn chunks_vectored<'t>(&'t self, dst: &mut [IoSlice<'t>]) -> usize {
636        let n = self.headers.chunks_vectored(dst);
637        self.queue.chunks_vectored(&mut dst[n..]) + n
638    }
639}
640
641#[derive(Debug)]
642enum WriteStrategy {
643    Flatten,
644    Queue,
645}
646
647#[cfg(test)]
648mod tests {
649    use super::*;
650    use crate::common::io::Compat;
651    use std::time::Duration;
652
653    use tokio_test::io::Builder as Mock;
654
655    // #[cfg(feature = "nightly")]
656    // use test::Bencher;
657
658    /*
659    impl<T: Read> MemRead for AsyncIo<T> {
660        fn read_mem(&mut self, len: usize) -> Poll<Bytes, io::Error> {
661            let mut v = vec![0; len];
662            let n = try_nb!(self.read(v.as_mut_slice()));
663            Ok(Async::Ready(BytesMut::from(&v[..n]).freeze()))
664        }
665    }
666    */
667
668    #[tokio::test]
669    #[ignore]
670    async fn iobuf_write_empty_slice() {
671        // TODO(eliza): can i have writev back pls T_T
672        // // First, let's just check that the Mock would normally return an
673        // // error on an unexpected write, even if the buffer is empty...
674        // let mut mock = Mock::new().build();
675        // futures_util::future::poll_fn(|cx| {
676        //     Pin::new(&mut mock).poll_write_buf(cx, &mut Cursor::new(&[]))
677        // })
678        // .await
679        // .expect_err("should be a broken pipe");
680
681        // // underlying io will return the logic error upon write,
682        // // so we are testing that the io_buf does not trigger a write
683        // // when there is nothing to flush
684        // let mock = Mock::new().build();
685        // let mut io_buf = Buffered::<_, Cursor<Vec<u8>>>::new(mock);
686        // io_buf.flush().await.expect("should short-circuit flush");
687    }
688
689    #[cfg(not(miri))]
690    #[tokio::test]
691    async fn parse_reads_until_blocked() {
692        use crate::proto::h1::ClientTransaction;
693
694        let _ = pretty_env_logger::try_init();
695        let mock = Mock::new()
696            // Split over multiple reads will read all of it
697            .read(b"HTTP/1.1 200 OK\r\n")
698            .read(b"Server: hyper\r\n")
699            // missing last line ending
700            .wait(Duration::from_secs(1))
701            .build();
702
703        let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(Compat::new(mock));
704
705        // We expect a `parse` to be not ready, and so can't await it directly.
706        // Rather, this `poll_fn` will wrap the `Poll` result.
707        futures_util::future::poll_fn(|cx| {
708            let parse_ctx = ParseContext {
709                cached_headers: &mut None,
710                req_method: &mut None,
711                h1_parser_config: Default::default(),
712                h1_max_headers: None,
713                preserve_header_case: false,
714                #[cfg(feature = "ffi")]
715                preserve_header_order: false,
716                h09_responses: false,
717                #[cfg(feature = "client")]
718                on_informational: &mut None,
719            };
720            assert!(buffered
721                .parse::<ClientTransaction>(cx, parse_ctx)
722                .is_pending());
723            Poll::Ready(())
724        })
725        .await;
726
727        assert_eq!(
728            buffered.read_buf,
729            b"HTTP/1.1 200 OK\r\nServer: hyper\r\n"[..]
730        );
731    }
732
733    #[test]
734    fn read_strategy_adaptive_increments() {
735        let mut strategy = ReadStrategy::default();
736        assert_eq!(strategy.next(), 8192);
737
738        // Grows if record == next
739        strategy.record(8192);
740        assert_eq!(strategy.next(), 16384);
741
742        strategy.record(16384);
743        assert_eq!(strategy.next(), 32768);
744
745        // Enormous records still increment at same rate
746        strategy.record(usize::MAX);
747        assert_eq!(strategy.next(), 65536);
748
749        let max = strategy.max();
750        while strategy.next() < max {
751            strategy.record(max);
752        }
753
754        assert_eq!(strategy.next(), max, "never goes over max");
755        strategy.record(max + 1);
756        assert_eq!(strategy.next(), max, "never goes over max");
757    }
758
759    #[test]
760    fn read_strategy_adaptive_decrements() {
761        let mut strategy = ReadStrategy::default();
762        strategy.record(8192);
763        assert_eq!(strategy.next(), 16384);
764
765        strategy.record(1);
766        assert_eq!(
767            strategy.next(),
768            16384,
769            "first smaller record doesn't decrement yet"
770        );
771        strategy.record(8192);
772        assert_eq!(strategy.next(), 16384, "record was with range");
773
774        strategy.record(1);
775        assert_eq!(
776            strategy.next(),
777            16384,
778            "in-range record should make this the 'first' again"
779        );
780
781        strategy.record(1);
782        assert_eq!(strategy.next(), 8192, "second smaller record decrements");
783
784        strategy.record(1);
785        assert_eq!(strategy.next(), 8192, "first doesn't decrement");
786        strategy.record(1);
787        assert_eq!(strategy.next(), 8192, "doesn't decrement under minimum");
788    }
789
790    #[test]
791    fn read_strategy_adaptive_stays_the_same() {
792        let mut strategy = ReadStrategy::default();
793        strategy.record(8192);
794        assert_eq!(strategy.next(), 16384);
795
796        strategy.record(8193);
797        assert_eq!(
798            strategy.next(),
799            16384,
800            "first smaller record doesn't decrement yet"
801        );
802
803        strategy.record(8193);
804        assert_eq!(
805            strategy.next(),
806            16384,
807            "with current step does not decrement"
808        );
809    }
810
811    #[test]
812    fn read_strategy_adaptive_max_fuzz() {
813        fn fuzz(max: usize) {
814            let mut strategy = ReadStrategy::with_max(max);
815            while strategy.next() < max {
816                strategy.record(usize::MAX);
817            }
818            let mut next = strategy.next();
819            while next > 8192 {
820                strategy.record(1);
821                strategy.record(1);
822                next = strategy.next();
823                assert!(
824                    next.is_power_of_two(),
825                    "decrement should be powers of two: {} (max = {})",
826                    next,
827                    max,
828                );
829            }
830        }
831
832        let mut max = 8192;
833        while max < usize::MAX {
834            fuzz(max);
835            max = (max / 2).saturating_mul(3);
836        }
837        fuzz(usize::MAX);
838    }
839
840    #[test]
841    #[should_panic]
842    #[cfg(debug_assertions)] // needs to trigger a debug_assert
843    fn write_buf_requires_non_empty_bufs() {
844        let mock = Mock::new().build();
845        let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(Compat::new(mock));
846
847        buffered.buffer(Cursor::new(Vec::new()));
848    }
849
850    /*
851    TODO: needs tokio_test::io to allow configure write_buf calls
852    #[test]
853    fn write_buf_queue() {
854        let _ = pretty_env_logger::try_init();
855
856        let mock = AsyncIo::new_buf(vec![], 1024);
857        let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(mock);
858
859
860        buffered.headers_buf().extend(b"hello ");
861        buffered.buffer(Cursor::new(b"world, ".to_vec()));
862        buffered.buffer(Cursor::new(b"it's ".to_vec()));
863        buffered.buffer(Cursor::new(b"hyper!".to_vec()));
864        assert_eq!(buffered.write_buf.queue.bufs_cnt(), 3);
865        buffered.flush().unwrap();
866
867        assert_eq!(buffered.io, b"hello world, it's hyper!");
868        assert_eq!(buffered.io.num_writes(), 1);
869        assert_eq!(buffered.write_buf.queue.bufs_cnt(), 0);
870    }
871    */
872
873    #[cfg(not(miri))]
874    #[tokio::test]
875    async fn write_buf_flatten() {
876        let _ = pretty_env_logger::try_init();
877
878        let mock = Mock::new().write(b"hello world, it's hyper!").build();
879
880        let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(Compat::new(mock));
881        buffered.write_buf.set_strategy(WriteStrategy::Flatten);
882
883        buffered.headers_buf().extend(b"hello ");
884        buffered.buffer(Cursor::new(b"world, ".to_vec()));
885        buffered.buffer(Cursor::new(b"it's ".to_vec()));
886        buffered.buffer(Cursor::new(b"hyper!".to_vec()));
887        assert_eq!(buffered.write_buf.queue.bufs_cnt(), 0);
888
889        buffered.flush().await.expect("flush");
890    }
891
892    #[test]
893    fn write_buf_flatten_partially_flushed() {
894        let _ = pretty_env_logger::try_init();
895
896        let b = |s: &str| Cursor::new(s.as_bytes().to_vec());
897
898        let mut write_buf = WriteBuf::<Cursor<Vec<u8>>>::new(WriteStrategy::Flatten);
899
900        write_buf.buffer(b("hello "));
901        write_buf.buffer(b("world, "));
902
903        assert_eq!(write_buf.chunk(), b"hello world, ");
904
905        // advance most of the way, but not all
906        write_buf.advance(11);
907
908        assert_eq!(write_buf.chunk(), b", ");
909        assert_eq!(write_buf.headers.pos, 11);
910        assert_eq!(write_buf.headers.bytes.capacity(), INIT_BUFFER_SIZE);
911
912        // there's still room in the headers buffer, so just push on the end
913        write_buf.buffer(b("it's hyper!"));
914
915        assert_eq!(write_buf.chunk(), b", it's hyper!");
916        assert_eq!(write_buf.headers.pos, 11);
917
918        let rem1 = write_buf.remaining();
919        let cap = write_buf.headers.bytes.capacity();
920
921        // but when this would go over capacity, don't copy the old bytes
922        write_buf.buffer(Cursor::new(vec![b'X'; cap]));
923        assert_eq!(write_buf.remaining(), cap + rem1);
924        assert_eq!(write_buf.headers.pos, 0);
925    }
926
927    #[cfg(not(miri))]
928    #[tokio::test]
929    async fn write_buf_queue_disable_auto() {
930        let _ = pretty_env_logger::try_init();
931
932        let mock = Mock::new()
933            .write(b"hello ")
934            .write(b"world, ")
935            .write(b"it's ")
936            .write(b"hyper!")
937            .build();
938
939        let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(Compat::new(mock));
940        buffered.write_buf.set_strategy(WriteStrategy::Queue);
941
942        // we have 4 buffers, and vec IO disabled, but explicitly said
943        // don't try to auto detect (via setting strategy above)
944
945        buffered.headers_buf().extend(b"hello ");
946        buffered.buffer(Cursor::new(b"world, ".to_vec()));
947        buffered.buffer(Cursor::new(b"it's ".to_vec()));
948        buffered.buffer(Cursor::new(b"hyper!".to_vec()));
949        assert_eq!(buffered.write_buf.queue.bufs_cnt(), 3);
950
951        buffered.flush().await.expect("flush");
952
953        assert_eq!(buffered.write_buf.queue.bufs_cnt(), 0);
954    }
955
956    // #[cfg(feature = "nightly")]
957    // #[bench]
958    // fn bench_write_buf_flatten_buffer_chunk(b: &mut Bencher) {
959    //     let s = "Hello, World!";
960    //     b.bytes = s.len() as u64;
961
962    //     let mut write_buf = WriteBuf::<bytes::Bytes>::new();
963    //     write_buf.set_strategy(WriteStrategy::Flatten);
964    //     b.iter(|| {
965    //         let chunk = bytes::Bytes::from(s);
966    //         write_buf.buffer(chunk);
967    //         ::test::black_box(&write_buf);
968    //         write_buf.headers.bytes.clear();
969    //     })
970    // }
971}