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
14pub(crate) const INIT_BUFFER_SIZE: usize = 8192;
16
17pub(crate) const MINIMUM_MAX_BUFFER_SIZE: usize = INIT_BUFFER_SIZE;
19
20pub(crate) const DEFAULT_MAX_BUFFER_SIZE: usize = 8192 + 4096 * 100;
24
25const 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 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 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 fn read_buf_remaining_mut(&self) -> usize {
127 self.read_buf.capacity() - self.read_buf.len()
128 }
129
130 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 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 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 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 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 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 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
336impl<T: Unpin, B> Unpin for Buffered<T, B> {}
338
339pub(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 *decrease_now = true;
416 }
417 } else {
418 *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 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 fn maybe_unshift(&mut self, additional: usize) {
466 if self.pos == 0 {
467 return;
469 }
470
471 if self.bytes.capacity() - self.bytes.len() >= additional {
472 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
513pub(super) struct WriteBuf<B> {
515 headers: Cursor<Vec<u8>>,
517 max_buf_size: usize,
518 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 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 #[tokio::test]
669 #[ignore]
670 async fn iobuf_write_empty_slice() {
671 }
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 .read(b"HTTP/1.1 200 OK\r\n")
698 .read(b"Server: hyper\r\n")
699 .wait(Duration::from_secs(1))
701 .build();
702
703 let mut buffered = Buffered::<_, Cursor<Vec<u8>>>::new(Compat::new(mock));
704
705 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 strategy.record(8192);
740 assert_eq!(strategy.next(), 16384);
741
742 strategy.record(16384);
743 assert_eq!(strategy.next(), 32768);
744
745 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)] 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 #[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 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 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 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 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 }