1use std::{error::Error, fmt, marker, mem, pin::Pin, rc::Rc, task::Context, task::Poll};
3
4use futures_core::Stream;
5use ntex_bytes::{BytePages, Bytes, BytesMut};
6
7#[derive(Debug, PartialEq, Eq, Copy, Clone)]
8pub enum BodySize {
10 None,
15 Empty,
17 Sized(u64),
19 Stream,
21}
22
23impl BodySize {
24 pub fn is_eof(&self) -> bool {
26 matches!(self, BodySize::None | BodySize::Empty | BodySize::Sized(0))
27 }
28}
29
30pub trait MessageBody: 'static {
32 fn size(&self) -> BodySize;
34
35 fn poll_next_chunk(
37 &mut self,
38 cx: &mut Context<'_>,
39 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>>;
40}
41
42impl MessageBody for () {
43 #[inline]
44 fn size(&self) -> BodySize {
45 BodySize::Empty
46 }
47
48 #[inline]
49 fn poll_next_chunk(
50 &mut self,
51 _: &mut Context<'_>,
52 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
53 Poll::Ready(None)
54 }
55}
56
57impl<T: MessageBody> MessageBody for Box<T> {
58 #[inline]
59 fn size(&self) -> BodySize {
60 self.as_ref().size()
61 }
62
63 #[inline]
64 fn poll_next_chunk(
65 &mut self,
66 cx: &mut Context<'_>,
67 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
68 self.as_mut().poll_next_chunk(cx)
69 }
70}
71
72#[derive(Debug)]
73pub enum ResponseBody<B> {
75 Body(B),
77 Other(Body),
79}
80
81impl ResponseBody<Body> {
82 pub fn into_body<B>(self) -> ResponseBody<B> {
84 match self {
85 ResponseBody::Body(b) | ResponseBody::Other(b) => ResponseBody::Other(b),
86 }
87 }
88}
89
90impl From<ResponseBody<Body>> for Body {
91 fn from(b: ResponseBody<Body>) -> Self {
92 match b {
93 ResponseBody::Body(b) | ResponseBody::Other(b) => b,
94 }
95 }
96}
97
98impl<B> From<Body> for ResponseBody<B> {
99 fn from(b: Body) -> Self {
100 ResponseBody::Other(b)
101 }
102}
103
104impl<B> ResponseBody<B> {
105 #[inline]
106 pub fn new(body: B) -> Self {
108 ResponseBody::Body(body)
109 }
110
111 #[inline]
112 #[must_use]
113 pub fn take_body(&mut self) -> ResponseBody<B> {
115 std::mem::replace(self, ResponseBody::Other(Body::None))
116 }
117}
118
119impl<B: MessageBody> ResponseBody<B> {
120 pub fn as_ref(&self) -> Option<&B> {
122 if let ResponseBody::Body(b) = self {
123 Some(b)
124 } else {
125 None
126 }
127 }
128}
129
130impl<B: MessageBody> MessageBody for ResponseBody<B> {
131 #[inline]
132 fn size(&self) -> BodySize {
133 match self {
134 ResponseBody::Body(body) => body.size(),
135 ResponseBody::Other(body) => body.size(),
136 }
137 }
138
139 #[inline]
140 fn poll_next_chunk(
141 &mut self,
142 cx: &mut Context<'_>,
143 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
144 match self {
145 ResponseBody::Body(body) => body.poll_next_chunk(cx),
146 ResponseBody::Other(body) => body.poll_next_chunk(cx),
147 }
148 }
149}
150
151impl<B: MessageBody + Unpin> Stream for ResponseBody<B> {
152 type Item = Result<Bytes, Rc<dyn Error>>;
153
154 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
155 match self.get_mut() {
156 ResponseBody::Body(body) => body.poll_next_chunk(cx),
157 ResponseBody::Other(body) => body.poll_next_chunk(cx),
158 }
159 }
160}
161
162pub enum Body {
164 None,
168 Empty,
170 Bytes(Bytes),
172 Message(Box<dyn MessageBody>),
174}
175
176impl Body {
177 pub fn from_slice(s: &[u8]) -> Body {
179 Body::Bytes(Bytes::copy_from_slice(s))
180 }
181
182 pub fn from_message<B: MessageBody>(body: B) -> Body {
184 Body::Message(Box::new(body))
185 }
186}
187
188impl MessageBody for Body {
189 #[inline]
190 fn size(&self) -> BodySize {
191 match self {
192 Body::None => BodySize::None,
193 Body::Empty => BodySize::Empty,
194 Body::Bytes(bin) => BodySize::Sized(bin.len() as u64),
195 Body::Message(body) => body.size(),
196 }
197 }
198
199 fn poll_next_chunk(
200 &mut self,
201 cx: &mut Context<'_>,
202 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
203 match self {
204 Body::None | Body::Empty => Poll::Ready(None),
205 Body::Bytes(bin) => {
206 let len = bin.len();
207 if len == 0 {
208 Poll::Ready(None)
209 } else {
210 Poll::Ready(Some(Ok(mem::take(bin))))
211 }
212 }
213 Body::Message(body) => body.poll_next_chunk(cx),
214 }
215 }
216}
217
218impl PartialEq for Body {
219 fn eq(&self, other: &Body) -> bool {
220 match (self, other) {
221 (Body::None, Body::None) | (Body::Empty, Body::Empty) => true,
222 (Body::Bytes(b), Body::Bytes(b2)) => b == b2,
223 _ => false,
224 }
225 }
226}
227
228impl fmt::Debug for Body {
229 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
230 match *self {
231 Body::None => write!(f, "Body::None"),
232 Body::Empty => write!(f, "Body::Empty"),
233 Body::Bytes(ref b) => write!(f, "Body::Bytes({b:?})"),
234 Body::Message(_) => write!(f, "Body::Message(_)"),
235 }
236 }
237}
238
239impl From<&'static str> for Body {
240 fn from(s: &'static str) -> Body {
241 Body::Bytes(Bytes::from_static(s.as_ref()))
242 }
243}
244
245impl From<&'static [u8]> for Body {
246 fn from(s: &'static [u8]) -> Body {
247 Body::Bytes(Bytes::from_static(s))
248 }
249}
250
251impl From<Vec<u8>> for Body {
252 fn from(vec: Vec<u8>) -> Body {
253 Body::Bytes(Bytes::from(vec))
254 }
255}
256
257impl From<String> for Body {
258 fn from(s: String) -> Body {
259 s.into_bytes().into()
260 }
261}
262
263impl<'a> From<&'a String> for Body {
264 fn from(s: &'a String) -> Body {
265 Body::Bytes(Bytes::copy_from_slice(AsRef::<[u8]>::as_ref(&s)))
266 }
267}
268
269impl From<Bytes> for Body {
270 fn from(s: Bytes) -> Body {
271 Body::Bytes(s)
272 }
273}
274
275impl From<BytesMut> for Body {
276 fn from(s: BytesMut) -> Body {
277 Body::Bytes(s.freeze())
278 }
279}
280
281impl From<BytePages> for Body {
282 fn from(pages: BytePages) -> Body {
283 Body::from_message(pages)
284 }
285}
286
287impl<S> From<SizedStream<S>> for Body
288where
289 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin + 'static,
290{
291 fn from(s: SizedStream<S>) -> Body {
292 Body::from_message(s)
293 }
294}
295
296impl<S, E> From<BodyStream<S, E>> for Body
297where
298 S: Stream<Item = Result<Bytes, E>> + Unpin + 'static,
299 E: Error + 'static,
300{
301 fn from(s: BodyStream<S, E>) -> Body {
302 Body::from_message(s)
303 }
304}
305
306impl<S> From<BoxedBodyStream<S>> for Body
307where
308 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin + 'static,
309{
310 fn from(s: BoxedBodyStream<S>) -> Body {
311 Body::from_message(s)
312 }
313}
314
315impl MessageBody for Bytes {
316 fn size(&self) -> BodySize {
317 BodySize::Sized(self.len() as u64)
318 }
319
320 fn poll_next_chunk(
321 &mut self,
322 _: &mut Context<'_>,
323 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
324 if self.is_empty() {
325 Poll::Ready(None)
326 } else {
327 Poll::Ready(Some(Ok(mem::take(self))))
328 }
329 }
330}
331
332impl MessageBody for BytesMut {
333 fn size(&self) -> BodySize {
334 BodySize::Sized(self.len() as u64)
335 }
336
337 fn poll_next_chunk(
338 &mut self,
339 _: &mut Context<'_>,
340 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
341 if self.is_empty() {
342 Poll::Ready(None)
343 } else {
344 Poll::Ready(Some(Ok(mem::take(self).freeze())))
345 }
346 }
347}
348
349impl MessageBody for BytePages {
350 fn size(&self) -> BodySize {
351 BodySize::Sized(self.len() as u64)
352 }
353
354 fn poll_next_chunk(
355 &mut self,
356 _: &mut Context<'_>,
357 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
358 if let Some(page) = self.take() {
359 Poll::Ready(Some(Ok(page.freeze())))
360 } else {
361 Poll::Ready(None)
362 }
363 }
364}
365
366impl MessageBody for &'static str {
367 fn size(&self) -> BodySize {
368 BodySize::Sized(self.len() as u64)
369 }
370
371 fn poll_next_chunk(
372 &mut self,
373 _: &mut Context<'_>,
374 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
375 if self.is_empty() {
376 Poll::Ready(None)
377 } else {
378 Poll::Ready(Some(Ok(Bytes::from_static(mem::take(self).as_ref()))))
379 }
380 }
381}
382
383impl MessageBody for &'static [u8] {
384 fn size(&self) -> BodySize {
385 BodySize::Sized(self.len() as u64)
386 }
387
388 fn poll_next_chunk(
389 &mut self,
390 _: &mut Context<'_>,
391 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
392 if self.is_empty() {
393 Poll::Ready(None)
394 } else {
395 Poll::Ready(Some(Ok(Bytes::from_static(mem::take(self)))))
396 }
397 }
398}
399
400impl MessageBody for Vec<u8> {
401 fn size(&self) -> BodySize {
402 BodySize::Sized(self.len() as u64)
403 }
404
405 fn poll_next_chunk(
406 &mut self,
407 _: &mut Context<'_>,
408 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
409 if self.is_empty() {
410 Poll::Ready(None)
411 } else {
412 Poll::Ready(Some(Ok(Bytes::from(mem::take(self)))))
413 }
414 }
415}
416
417impl MessageBody for String {
418 fn size(&self) -> BodySize {
419 BodySize::Sized(self.len() as u64)
420 }
421
422 fn poll_next_chunk(
423 &mut self,
424 _: &mut Context<'_>,
425 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
426 if self.is_empty() {
427 Poll::Ready(None)
428 } else {
429 Poll::Ready(Some(Ok(Bytes::from(mem::take(self).into_bytes()))))
430 }
431 }
432}
433
434pub struct BodyStream<S, E> {
438 stream: S,
439 _t: marker::PhantomData<E>,
440}
441
442impl<S, E> BodyStream<S, E>
443where
444 S: Stream<Item = Result<Bytes, E>> + Unpin,
445 E: Error,
446{
447 pub fn new(stream: S) -> Self {
448 BodyStream {
449 stream,
450 _t: marker::PhantomData,
451 }
452 }
453}
454
455impl<S, E> fmt::Debug for BodyStream<S, E> {
456 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
457 f.debug_struct("BodyStream")
458 .field("stream", &std::any::type_name::<S>())
459 .field("error", &std::any::type_name::<E>())
460 .finish()
461 }
462}
463
464impl<S, E> MessageBody for BodyStream<S, E>
465where
466 S: Stream<Item = Result<Bytes, E>> + Unpin + 'static,
467 E: Error + 'static,
468{
469 fn size(&self) -> BodySize {
470 BodySize::Stream
471 }
472
473 fn poll_next_chunk(
479 &mut self,
480 cx: &mut Context<'_>,
481 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
482 loop {
483 return Poll::Ready(match Pin::new(&mut self.stream).poll_next(cx) {
484 Poll::Ready(Some(Ok(ref bytes))) if bytes.is_empty() => continue,
485 Poll::Ready(opt) => opt.map(|res| {
486 res.map_err(|e| {
487 let e: Rc<dyn Error> = Rc::new(e);
488 e
489 })
490 }),
491 Poll::Pending => return Poll::Pending,
492 });
493 }
494 }
495}
496
497pub struct BoxedBodyStream<S> {
500 stream: S,
501}
502
503impl<S> BoxedBodyStream<S>
504where
505 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin,
506{
507 pub fn new(stream: S) -> Self {
508 BoxedBodyStream { stream }
509 }
510}
511
512impl<S> fmt::Debug for BoxedBodyStream<S> {
513 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
514 f.debug_struct("BoxedBodyStream")
515 .field("stream", &std::any::type_name::<S>())
516 .finish()
517 }
518}
519
520impl<S> MessageBody for BoxedBodyStream<S>
521where
522 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin + 'static,
523{
524 fn size(&self) -> BodySize {
525 BodySize::Stream
526 }
527
528 fn poll_next_chunk(
534 &mut self,
535 cx: &mut Context<'_>,
536 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
537 loop {
538 return Poll::Ready(match Pin::new(&mut self.stream).poll_next(cx) {
539 Poll::Ready(Some(Ok(ref bytes))) if bytes.is_empty() => continue,
540 Poll::Ready(opt) => opt,
541 Poll::Pending => return Poll::Pending,
542 });
543 }
544 }
545}
546
547pub struct SizedStream<S> {
550 size: u64,
551 stream: S,
552}
553
554impl<S> SizedStream<S>
555where
556 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin,
557{
558 pub fn new(size: u64, stream: S) -> Self {
559 SizedStream { size, stream }
560 }
561}
562
563impl<S> fmt::Debug for SizedStream<S> {
564 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
565 f.debug_struct("SizedStream")
566 .field("size", &self.size)
567 .field("stream", &std::any::type_name::<S>())
568 .finish()
569 }
570}
571
572impl<S> MessageBody for SizedStream<S>
573where
574 S: Stream<Item = Result<Bytes, Rc<dyn Error>>> + Unpin + 'static,
575{
576 fn size(&self) -> BodySize {
577 BodySize::Sized(self.size)
578 }
579
580 fn poll_next_chunk(
586 &mut self,
587 cx: &mut Context<'_>,
588 ) -> Poll<Option<Result<Bytes, Rc<dyn Error>>>> {
589 loop {
590 return Poll::Ready(match Pin::new(&mut self.stream).poll_next(cx) {
591 Poll::Ready(Some(Ok(ref bytes))) if bytes.is_empty() => continue,
592 Poll::Ready(val) => val,
593 Poll::Pending => return Poll::Pending,
594 });
595 }
596 }
597}
598
599#[cfg(test)]
600mod tests {
601 use std::{future::poll_fn, future::ready, io};
602
603 use futures_util::stream;
604 use ntex::util::BufMut;
605
606 use super::*;
607
608 impl Body {
609 pub(crate) fn get_ref(&self) -> &[u8] {
610 if let Body::Bytes(bin) = self { bin } else { panic!() }
611 }
612 }
613
614 #[ntex::test]
615 async fn test_size() {
616 assert_eq!(32, std::mem::size_of::<Body>());
617 assert_eq!(32, std::mem::size_of::<ResponseBody<Bytes>>());
618 assert_eq!(40, std::mem::size_of::<ResponseBody<Body>>());
619 }
620
621 #[ntex::test]
622 async fn test_static_str() {
623 assert_eq!(Body::from("").size(), BodySize::Sized(0));
624 assert_eq!(Body::from("test").size(), BodySize::Sized(4));
625 assert_eq!(Body::from("test").get_ref(), b"test");
626
627 assert_eq!("test".size(), BodySize::Sized(4));
628 assert_eq!(
629 poll_fn(|cx| "test".poll_next_chunk(cx)).await.unwrap().ok(),
630 Some(Bytes::from("test"))
631 );
632 assert_eq!(
633 poll_fn(|cx| "test".poll_next_chunk(cx)).await.unwrap().ok(),
634 Some(Bytes::from("test"))
635 );
636 assert!(poll_fn(|cx| "".poll_next_chunk(cx)).await.is_none());
637 }
638
639 #[ntex::test]
640 async fn test_pages() {
641 let mut pages = BytePages::default();
642 pages.put_slice(b"1111");
643 pages.append(Bytes::from("2222"));
644 pages.append(Bytes::from("3333"));
645
646 let mut body = Body::from(pages);
647 assert_eq!(body.size(), BodySize::Sized(12));
648 assert_eq!(
649 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
650 Some(Bytes::from("111122223333"))
651 );
652 assert!(poll_fn(|cx| body.poll_next_chunk(cx)).await.is_none());
653 }
654
655 #[ntex::test]
656 async fn test_static_bytes() {
657 assert_eq!(Body::from(b"test".as_ref()).size(), BodySize::Sized(4));
658 assert_eq!(Body::from(b"test".as_ref()).get_ref(), b"test");
659 assert_eq!(
660 Body::from_slice(b"test".as_ref()).size(),
661 BodySize::Sized(4)
662 );
663 assert_eq!(Body::from_slice(b"test".as_ref()).get_ref(), b"test");
664
665 assert_eq!((&b"test"[..]).size(), BodySize::Sized(4));
666 assert_eq!(
667 poll_fn(|cx| (&b"test"[..]).poll_next_chunk(cx))
668 .await
669 .unwrap()
670 .ok(),
671 Some(Bytes::from("test"))
672 );
673 assert_eq!((&b"test"[..]).size(), BodySize::Sized(4));
674 assert!(poll_fn(|cx| (&b""[..]).poll_next_chunk(cx)).await.is_none());
675 }
676
677 #[ntex::test]
678 async fn test_vec() {
679 assert_eq!(Body::from(Vec::from("test")).size(), BodySize::Sized(4));
680 assert_eq!(Body::from(Vec::from("test")).get_ref(), b"test");
681
682 assert_eq!(Vec::from("test").size(), BodySize::Sized(4));
683 assert_eq!(
684 poll_fn(|cx| Vec::from("test").poll_next_chunk(cx))
685 .await
686 .unwrap()
687 .ok(),
688 Some(Bytes::from("test"))
689 );
690 assert_eq!(
691 poll_fn(|cx| Vec::from("test").poll_next_chunk(cx))
692 .await
693 .unwrap()
694 .ok(),
695 Some(Bytes::from("test"))
696 );
697 assert!(
698 poll_fn(|cx| Vec::<u8>::new().poll_next_chunk(cx))
699 .await
700 .is_none()
701 );
702 }
703
704 #[ntex::test]
705 async fn test_bytes() {
706 let mut b = Bytes::from("test");
707 assert_eq!(Body::from(b.clone()).size(), BodySize::Sized(4));
708 assert_eq!(Body::from(b.clone()).get_ref(), b"test");
709
710 assert_eq!(b.size(), BodySize::Sized(4));
711 assert_eq!(
712 poll_fn(|cx| b.poll_next_chunk(cx)).await.unwrap().ok(),
713 Some(Bytes::from("test"))
714 );
715 assert!(poll_fn(|cx| b.poll_next_chunk(cx)).await.is_none(),);
716 }
717
718 #[ntex::test]
719 async fn test_bytes_mut() {
720 let mut b = Body::from(BytesMut::from("test"));
721 assert_eq!(b.size(), BodySize::Sized(4));
722 assert_eq!(b.get_ref(), b"test");
723 assert_eq!(
724 poll_fn(|cx| b.poll_next_chunk(cx)).await.unwrap().ok(),
725 Some(Bytes::from("test"))
726 );
727 assert!(poll_fn(|cx| b.poll_next_chunk(cx)).await.is_none(),);
728
729 let mut b = BytesMut::from("test");
730 assert_eq!(b.size(), BodySize::Sized(4));
731 assert_eq!(
732 poll_fn(|cx| b.poll_next_chunk(cx)).await.unwrap().ok(),
733 Some(Bytes::from("test"))
734 );
735 assert_eq!(b.size(), BodySize::Sized(0));
736 assert!(poll_fn(|cx| b.poll_next_chunk(cx)).await.is_none(),);
737 }
738
739 #[ntex::test]
740 async fn test_string() {
741 let mut b = "test".to_owned();
742 assert_eq!(Body::from(b.clone()).size(), BodySize::Sized(4));
743 assert_eq!(Body::from(b.clone()).get_ref(), b"test");
744 assert_eq!(Body::from(&b).size(), BodySize::Sized(4));
745 assert_eq!(Body::from(&b).get_ref(), b"test");
746
747 assert_eq!(b.size(), BodySize::Sized(4));
748 assert_eq!(
749 poll_fn(|cx| b.poll_next_chunk(cx)).await.unwrap().ok(),
750 Some(Bytes::from("test"))
751 );
752 assert!(poll_fn(|cx| b.poll_next_chunk(cx)).await.is_none(),);
753 }
754
755 #[ntex::test]
756 async fn test_unit() {
757 assert_eq!(().size(), BodySize::Empty);
758 assert!(poll_fn(|cx| ().poll_next_chunk(cx)).await.is_none());
759 }
760
761 #[ntex::test]
762 async fn test_box() {
763 let mut val = Box::new(());
764 assert_eq!(val.size(), BodySize::Empty);
765 assert!(poll_fn(|cx| val.poll_next_chunk(cx)).await.is_none());
766 }
767
768 #[ntex::test]
769 #[allow(clippy::eq_op)]
770 async fn test_body_eq() {
771 assert_eq!(Body::None, Body::None);
772 assert_ne!(Body::None, Body::Empty);
773 assert_eq!(Body::Empty, Body::Empty);
774 assert_ne!(Body::Empty, Body::None);
775 assert_eq!(
776 Body::Bytes(Bytes::from_static(b"1")),
777 Body::Bytes(Bytes::from_static(b"1"))
778 );
779 assert_ne!(Body::Bytes(Bytes::from_static(b"1")), Body::None);
780 }
781
782 #[ntex::test]
783 async fn test_body_debug() {
784 assert!(format!("{:?}", Body::None).contains("Body::None"));
785 assert!(format!("{:?}", Body::Empty).contains("Body::Empty"));
786 assert!(format!("{:?}", Body::Bytes(Bytes::from_static(b"1"))).contains('1'));
787 }
788
789 #[ntex::test]
790 async fn body_stream() {
791 let st = BodyStream::new(stream::once(ready(Ok::<_, io::Error>(Bytes::from("1")))));
792 assert!(format!("{st:?}").contains("BodyStream"));
793 let body: Body = st.into();
794 assert!(format!("{body:?}").contains("Body::Message(_)"));
795 assert_ne!(body, Body::None);
796
797 let res = ResponseBody::new(body);
798 assert!(res.as_ref().is_some());
799 }
800
801 #[ntex::test]
802 async fn boxed_body_stream() {
803 let st = BoxedBodyStream::new(stream::once(ready(Ok::<_, Rc<dyn Error>>(Bytes::from(
804 "1",
805 )))));
806 assert!(format!("{st:?}").contains("BoxedBodyStream"));
807 let body: Body = st.into();
808 assert!(format!("{body:?}").contains("Body::Message(_)"));
809 assert_ne!(body, Body::None);
810
811 let res = ResponseBody::new(body);
812 assert!(res.as_ref().is_some());
813 }
814
815 #[ntex::test]
816 async fn body_skips_empty_chunks() {
817 let mut body = BodyStream::new(stream::iter(
818 ["1", "", "2"]
819 .iter()
820 .map(|&v| Ok(Bytes::from(v)) as Result<Bytes, io::Error>),
821 ));
822 assert_eq!(
823 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
824 Some(Bytes::from("1")),
825 );
826 assert_eq!(
827 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
828 Some(Bytes::from("2")),
829 );
830 }
831
832 #[ntex::test]
833 async fn sized_skips_empty_chunks() {
834 let mut body = SizedStream::new(
835 2,
836 stream::iter(["1", "", "2"].iter().map(|&v| Ok(Bytes::from(v)))),
837 );
838 assert!(format!("{body:?}").contains("SizedStream"));
839 assert_eq!(
840 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
841 Some(Bytes::from("1")),
842 );
843 assert_eq!(
844 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
845 Some(Bytes::from("2")),
846 );
847 }
848
849 #[ntex::test]
850 async fn boxed_body_skips_empty_chunks() {
851 let mut body = BoxedBodyStream::new(stream::iter(
852 ["1", "", "2"]
853 .iter()
854 .map(|&v| Ok(Bytes::from(v)) as Result<Bytes, Rc<dyn Error>>),
855 ));
856 assert_eq!(
857 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
858 Some(Bytes::from("1")),
859 );
860 assert_eq!(
861 poll_fn(|cx| body.poll_next_chunk(cx)).await.unwrap().ok(),
862 Some(Bytes::from("2")),
863 );
864 assert!(poll_fn(|cx| body.poll_next_chunk(cx)).await.is_none());
865 }
866}