Skip to main content

aws_smithy_types/
body.rs

1/*
2 * Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
3 * SPDX-License-Identifier: Apache-2.0
4 */
5
6//! Types for representing the body of an HTTP request or response
7
8use bytes::Bytes;
9use pin_project_lite::pin_project;
10use std::collections::VecDeque;
11use std::error::Error as StdError;
12use std::fmt::{self, Debug, Formatter};
13use std::future::poll_fn;
14use std::pin::Pin;
15use std::sync::Arc;
16use std::task::{Context, Poll};
17
18/// This module is named after the `http-body` version number since we anticipate
19/// needing to provide equivalent functionality for 1.x of that crate in the future.
20/// The name has a suffix `_x` to avoid name collision with a third-party `http-body-0-4`.
21#[cfg(feature = "http-body-0-4-x")]
22pub mod http_body_0_4_x;
23#[cfg(feature = "http-body-1-x")]
24pub mod http_body_1_x;
25
26/// A generic, boxed error that's `Send` and `Sync`
27pub type Error = Box<dyn StdError + Send + Sync>;
28
29// Converts an `http` 0.2.x trailer `HeaderMap` into an `http` 1.x `HeaderMap`. Only needed on the
30// legacy http-body 0.4.x trailer code path inside `poll_next_trailers`, which exists only when the
31// `http-body-0-4-x` feature is enabled. The `http-body-1-x` path (including `rt-tokio`, which now
32// rides the 1.x body path) never needs the `http` 0.2.x crate.
33#[cfg(feature = "http-body-0-4-x")]
34fn convert_trailers_0x_1x(input: http::HeaderMap) -> http_1x::HeaderMap {
35    let mut map = http_1x::HeaderMap::with_capacity(input.capacity());
36    let mut mem: Option<http::HeaderName> = None;
37    for (k, v) in input.into_iter() {
38        let name = k.or_else(|| mem.clone()).unwrap();
39        map.append(
40            http_1x::HeaderName::from_bytes(name.as_str().as_bytes()).expect("already validated"),
41            http_1x::HeaderValue::from_bytes(v.as_bytes()).expect("already validated"),
42        );
43        mem = Some(name);
44    }
45    map
46}
47
48// Converts an `http` 1.x `HeaderMap` into an `http` 0.2.x `HeaderMap`. Shared by the http-body
49// 0.4.x adapters in `body/http_body_0_4_x.rs` and `body/http_body_1_x.rs`, so it lives here to
50// avoid duplication. Only needed on the legacy 0.4.x path (gated on `http-body-0-4-x`).
51#[cfg(feature = "http-body-0-4-x")]
52pub(crate) fn convert_headers_1x_0x(input: http_1x::HeaderMap) -> http::HeaderMap {
53    let mut map = http::HeaderMap::with_capacity(input.capacity());
54    let mut mem: Option<http_1x::HeaderName> = None;
55    for (k, v) in input.into_iter() {
56        let name = k.or_else(|| mem.clone()).unwrap();
57        map.append(
58            http::HeaderName::from_bytes(name.as_str().as_bytes()).expect("already validated"),
59            http::HeaderValue::from_bytes(v.as_bytes()).expect("already validated"),
60        );
61        mem = Some(name);
62    }
63    map
64}
65
66pin_project! {
67    /// SdkBody type
68    ///
69    /// This is the Body used for dispatching all HTTP Requests.
70    /// For handling responses, the type of the body will be controlled
71    /// by the HTTP stack.
72    ///
73    pub struct SdkBody {
74        #[pin]
75        inner: Inner,
76        // An optional function to recreate the inner body
77        //
78        // In the event of retry, this function will be called to generate a new body. See
79        // [`try_clone()`](SdkBody::try_clone)
80        rebuild: Option<Arc<dyn (Fn() -> Inner) + Send + Sync>>,
81        bytes_contents: Option<Bytes>,
82        // Here the optionality indicates whether we have started streaming trailers, and the
83        // VecDeque serves as a buffer for trailer frames that are polled by poll_next instead
84        // of poll_next_trailers
85        trailers: Option<VecDeque<http_1x::HeaderMap>>,
86    }
87}
88
89impl Debug for SdkBody {
90    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
91        f.debug_struct("SdkBody")
92            .field("inner", &self.inner)
93            .field("retryable", &self.rebuild.is_some())
94            .finish()
95    }
96}
97
98/// A boxed generic HTTP body that, when consumed, will result in [`Bytes`] or an [`Error`].
99#[allow(dead_code)]
100enum BoxBody {
101    // The legacy http-body 0.4.x box body. Only exists when the `http-body-0-4-x` feature is
102    // enabled; the `http-body-1-x` path (including `rt-tokio`) does not pull in the `http` 0.2.x /
103    // `http-body` 0.4.x crates.
104    #[cfg(feature = "http-body-0-4-x")]
105    HttpBody04(#[allow(dead_code)] http_body_0_4::combinators::BoxBody<Bytes, Error>),
106
107    #[cfg(feature = "http-body-1-x")]
108    HttpBody1(#[allow(dead_code)] http_body_util::combinators::BoxBody<Bytes, Error>),
109}
110
111pin_project! {
112    #[project = InnerProj]
113    enum Inner {
114        // An in-memory body
115        Once {
116            inner: Option<Bytes>
117        },
118        // A streaming body
119        Dyn {
120            #[pin]
121            inner: BoxBody,
122        },
123
124        /// When a streaming body is transferred out to a stream parser, the body is replaced with
125        /// `Taken`. This will return an Error when polled. Attempting to read data out of a `Taken`
126        /// Body is a bug.
127        Taken,
128    }
129}
130
131impl Debug for Inner {
132    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
133        match &self {
134            Inner::Once { inner: once } => f.debug_tuple("Once").field(once).finish(),
135            Inner::Dyn { .. } => write!(f, "BoxBody"),
136            Inner::Taken => f.debug_tuple("Taken").finish(),
137        }
138    }
139}
140
141impl SdkBody {
142    /// Construct an explicitly retryable SDK body
143    ///
144    /// _Note: This is probably not what you want_
145    ///
146    /// All bodies constructed from in-memory data (`String`, `Vec<u8>`, `Bytes`, etc.) will be
147    /// retryable out of the box. If you want to read data from a file, you should use
148    /// [`ByteStream::from_path`](crate::byte_stream::ByteStream::from_path). This function
149    /// is only necessary when you need to enable retries for your own streaming container.
150    pub fn retryable(f: impl Fn() -> SdkBody + Send + Sync + 'static) -> Self {
151        let initial = f();
152        SdkBody {
153            inner: initial.inner,
154            rebuild: Some(Arc::new(move || f().inner)),
155            bytes_contents: initial.bytes_contents,
156            trailers: None,
157        }
158    }
159
160    /// When an SdkBody is read, the inner data must be consumed. In order to do this, the SdkBody
161    /// is swapped with a "taken" body. This "taken" body cannot be read but aids in debugging.
162    pub fn taken() -> Self {
163        Self {
164            inner: Inner::Taken,
165            rebuild: None,
166            bytes_contents: None,
167            trailers: None,
168        }
169    }
170
171    /// Create an empty SdkBody for requests and responses that don't transfer any data in the body.
172    pub fn empty() -> Self {
173        Self {
174            inner: Inner::Once { inner: None },
175            rebuild: Some(Arc::new(|| Inner::Once { inner: None })),
176            bytes_contents: Some(Bytes::new()),
177            trailers: None,
178        }
179    }
180
181    pub(crate) async fn next(&mut self) -> Option<Result<Bytes, Error>> {
182        let mut me = Pin::new(self);
183        poll_fn(|cx| me.as_mut().poll_next(cx)).await
184    }
185
186    pub(crate) fn poll_next(
187        self: Pin<&mut Self>,
188        #[allow(unused)] cx: &mut Context<'_>,
189    ) -> Poll<Option<Result<Bytes, Error>>> {
190        let this = self.project();
191        match this.inner.project() {
192            InnerProj::Once { ref mut inner } => {
193                let data = inner.take();
194                match data {
195                    Some(bytes) if bytes.is_empty() => Poll::Ready(None),
196                    Some(bytes) => Poll::Ready(Some(Ok(bytes))),
197                    None => Poll::Ready(None),
198                }
199            }
200            InnerProj::Dyn { inner: body } => match body.get_mut() {
201                #[cfg(feature = "http-body-0-4-x")]
202                BoxBody::HttpBody04(box_body) => {
203                    use http_body_0_4::Body;
204                    Pin::new(box_body).poll_data(cx)
205                }
206                #[cfg(feature = "http-body-1-x")]
207                BoxBody::HttpBody1(box_body) => {
208                    // If this is polled after the trailers have been cached end early
209                    if this.trailers.is_some() {
210                        return Poll::Ready(None);
211                    }
212                    use http_body_1_0::Body;
213                    let maybe_data = Pin::new(box_body).poll_frame(cx);
214                    match maybe_data {
215                        Poll::Ready(Some(Ok(frame))) => {
216                            if frame.is_data() {
217                                Poll::Ready(Some(Ok(frame
218                                    .into_data()
219                                    .expect("Confirmed data frame"))))
220                            } else if frame.is_trailers() {
221                                let trailers =
222                                    frame.into_trailers().expect("Confirmed trailer frame");
223                                // Buffer the trailers for the trailer poll
224                                this.trailers.get_or_insert_with(VecDeque::new).push_back(trailers);
225
226                                Poll::Ready(None)
227                            } else {
228                                unreachable!("Frame must be either data or trailers");
229                            }
230                        }
231                        Poll::Ready(Some(Err(err))) => Poll::Ready(Some(Err(err))),
232                        Poll::Ready(None) => Poll::Ready(None),
233                        Poll::Pending => Poll::Pending,
234                    }
235                }
236                #[allow(unreachable_patterns)]
237                _ => unreachable!(
238                    "enabling `http-body-0-4-x` or `http-body-1-x` is the only way to create the `Dyn` variant"
239                ),
240            },
241            InnerProj::Taken => {
242                Poll::Ready(Some(Err("A `Taken` body should never be polled".into())))
243            }
244        }
245    }
246
247    #[allow(dead_code)]
248    #[cfg(feature = "http-body-0-4-x")]
249    pub(crate) fn from_body_0_4_internal<T, E>(body: T) -> Self
250    where
251        T: http_body_0_4::Body<Data = Bytes, Error = E> + Send + Sync + 'static,
252        E: Into<Error> + 'static,
253    {
254        Self {
255            inner: Inner::Dyn {
256                inner: BoxBody::HttpBody04(http_body_0_4::combinators::BoxBody::new(
257                    body.map_err(Into::into),
258                )),
259            },
260            rebuild: None,
261            bytes_contents: None,
262            trailers: None,
263        }
264    }
265
266    #[cfg(feature = "http-body-1-x")]
267    pub(crate) fn from_body_1_x_internal<T, E>(body: T) -> Self
268    where
269        T: http_body_1_0::Body<Data = Bytes, Error = E> + Send + Sync + 'static,
270        E: Into<Error> + 'static,
271    {
272        use http_body_util::BodyExt;
273        Self {
274            inner: Inner::Dyn {
275                inner: BoxBody::HttpBody1(http_body_util::combinators::BoxBody::new(
276                    body.map_err(Into::into),
277                )),
278            },
279            rebuild: None,
280            bytes_contents: None,
281            trailers: None,
282        }
283    }
284
285    #[cfg(any(feature = "http-body-0-4-x", feature = "http-body-1-x",))]
286    pub(crate) fn poll_next_trailers(
287        self: Pin<&mut Self>,
288        cx: &mut Context<'_>,
289    ) -> Poll<Result<Option<http_1x::HeaderMap<http_1x::HeaderValue>>, Error>> {
290        let this = self.project();
291        match this.inner.project() {
292            InnerProj::Once { .. } => Poll::Ready(Ok(None)),
293            InnerProj::Dyn { inner } => match inner.get_mut() {
294                #[cfg(feature = "http-body-0-4-x")]
295                BoxBody::HttpBody04(box_body) => {
296                    use http_body_0_4::Body;
297                    let polled = Pin::new(box_body).poll_trailers(cx);
298
299                    match polled {
300                        Poll::Ready(Ok(maybe_trailers)) => {
301                            let http_1x_trailers = maybe_trailers.map(convert_trailers_0x_1x);
302                            Poll::Ready(Ok(http_1x_trailers))
303                        }
304                        Poll::Ready(Err(err)) => Poll::Ready(Err(err)),
305                        Poll::Pending => Poll::Pending,
306                    }
307                }
308                #[cfg(feature = "http-body-1-x")]
309                BoxBody::HttpBody1(box_body) => {
310                    use http_body_1_0::Body;
311                    // Return the cached trailers without polling
312                    if let Some(trailer_buf) = this.trailers {
313                        if let Some(next_trailer) = trailer_buf.pop_front() {
314                            return Poll::Ready(Ok(Some(next_trailer)));
315                        }
316                    }
317
318                    let polled = Pin::new(box_body).poll_frame(cx);
319                    match polled {
320                        Poll::Ready(Some(Ok(maybe_trailers))) => {
321                            if maybe_trailers.is_data() {
322                                Poll::Ready(Err("Trailers polled while body still has data".into()))
323                            } else {
324                                let trailers = maybe_trailers
325                                    .into_trailers()
326                                    .expect("Frame must be trailers because it is not data");
327                                Poll::Ready(Ok(Some(trailers)))
328                            }
329                        }
330                        Poll::Ready(None) => Poll::Ready(Ok(None)),
331                        Poll::Ready(Some(Err(err))) => Poll::Ready(Err(err)),
332                        Poll::Pending => Poll::Pending,
333                    }
334                }
335            },
336            InnerProj::Taken => Poll::Ready(Err(
337                "A `Taken` body should never be polled for trailers".into(),
338            )),
339        }
340    }
341
342    /// If possible, return a reference to this body as `&[u8]`
343    ///
344    /// If this SdkBody is NOT streaming, this will return the byte slab
345    /// If this SdkBody is streaming, this will return `None`
346    pub fn bytes(&self) -> Option<&[u8]> {
347        match &self.bytes_contents {
348            Some(b) => Some(b),
349            None => None,
350        }
351    }
352
353    /// Attempt to clone this SdkBody. This will fail if the inner data is not cloneable, such as when
354    /// it is a single-use stream that can't be recreated.
355    pub fn try_clone(&self) -> Option<Self> {
356        self.rebuild.as_ref().map(|rebuild| {
357            let next = rebuild();
358            Self {
359                inner: next,
360                rebuild: self.rebuild.clone(),
361                bytes_contents: self.bytes_contents.clone(),
362                trailers: self.trailers.clone(),
363            }
364        })
365    }
366
367    /// Return `true` if this SdkBody is streaming, `false` if it is in-memory.
368    pub fn is_streaming(&self) -> bool {
369        matches!(self.inner, Inner::Dyn { .. })
370    }
371
372    /// Return the length, in bytes, of this SdkBody. If this returns `None`, then the body does not
373    /// have a known length.
374    pub fn content_length(&self) -> Option<u64> {
375        match self.bounds_on_remaining_length() {
376            (lo, Some(hi)) if lo == hi => Some(lo),
377            _ => None,
378        }
379    }
380
381    #[allow(dead_code)] // used by a feature-gated `http-body`'s trait method
382    pub(crate) fn is_end_stream(&self) -> bool {
383        match &self.inner {
384            Inner::Once { inner: None } => true,
385            Inner::Once { inner: Some(bytes) } => bytes.is_empty(),
386            Inner::Dyn { inner: box_body } => match box_body {
387                #[cfg(feature = "http-body-0-4-x")]
388                BoxBody::HttpBody04(box_body) => {
389                    use http_body_0_4::Body;
390                    box_body.is_end_stream()
391                }
392                #[cfg(feature = "http-body-1-x")]
393                BoxBody::HttpBody1(box_body) => {
394                    use http_body_1_0::Body;
395                    box_body.is_end_stream()
396                }
397                #[allow(unreachable_patterns)]
398                _ => unreachable!(
399                    "enabling `http-body-0-4-x` or `http-body-1-x` is the only way to create the `Dyn` variant"
400                ),
401            },
402            Inner::Taken => true,
403        }
404    }
405
406    pub(crate) fn bounds_on_remaining_length(&self) -> (u64, Option<u64>) {
407        match &self.inner {
408            Inner::Once { inner: None } => (0, Some(0)),
409            Inner::Once { inner: Some(bytes) } => {
410                let len = bytes.len() as u64;
411                (len, Some(len))
412            }
413            Inner::Dyn { inner: box_body } => match box_body {
414                #[cfg(feature = "http-body-0-4-x")]
415                BoxBody::HttpBody04(box_body) => {
416                    use http_body_0_4::Body;
417                    let hint = box_body.size_hint();
418                    (hint.lower(), hint.upper())
419                }
420                #[cfg(feature = "http-body-1-x")]
421                BoxBody::HttpBody1(box_body) => {
422                    use http_body_1_0::Body;
423                    let hint = box_body.size_hint();
424                    (hint.lower(), hint.upper())
425                }
426                #[allow(unreachable_patterns)]
427                _ => unreachable!(
428                    "enabling `http-body-0-4-x` or `http-body-1-x` is the only way to create the `Dyn` variant"
429                ),
430            },
431            Inner::Taken => (0, Some(0)),
432        }
433    }
434
435    /// Given a function to modify an `SdkBody`, run that function against this `SdkBody` before
436    /// returning the result.
437    pub fn map(self, f: impl Fn(SdkBody) -> SdkBody + Sync + Send + 'static) -> SdkBody {
438        if self.rebuild.is_some() {
439            SdkBody::retryable(move || f(self.try_clone().unwrap()))
440        } else {
441            f(self)
442        }
443    }
444
445    /// Update this `SdkBody` with `map`. **This function MUST NOT alter the data of the body.**
446    ///
447    /// This function is useful for adding metadata like progress tracking to an [`SdkBody`] that
448    /// does not alter the actual byte data. If your mapper alters the contents of the body, use [`SdkBody::map`]
449    /// instead.
450    pub fn map_preserve_contents(
451        self,
452        f: impl Fn(SdkBody) -> SdkBody + Sync + Send + 'static,
453    ) -> SdkBody {
454        let contents = self.bytes_contents.clone();
455        let mut out = if self.rebuild.is_some() {
456            SdkBody::retryable(move || f(self.try_clone().unwrap()))
457        } else {
458            f(self)
459        };
460        out.bytes_contents = contents;
461        out
462    }
463}
464
465impl From<&str> for SdkBody {
466    fn from(s: &str) -> Self {
467        Self::from(s.as_bytes())
468    }
469}
470
471impl From<Bytes> for SdkBody {
472    fn from(bytes: Bytes) -> Self {
473        let b = bytes.clone();
474        SdkBody {
475            inner: Inner::Once {
476                inner: Some(bytes.clone()),
477            },
478            rebuild: Some(Arc::new(move || Inner::Once {
479                inner: Some(bytes.clone()),
480            })),
481            bytes_contents: Some(b),
482            trailers: None,
483        }
484    }
485}
486
487impl From<Vec<u8>> for SdkBody {
488    fn from(data: Vec<u8>) -> Self {
489        Self::from(Bytes::from(data))
490    }
491}
492
493impl From<String> for SdkBody {
494    fn from(s: String) -> Self {
495        Self::from(s.into_bytes())
496    }
497}
498
499impl From<&[u8]> for SdkBody {
500    fn from(data: &[u8]) -> Self {
501        Self::from(Bytes::copy_from_slice(data))
502    }
503}
504
505#[cfg(test)]
506mod test {
507    use crate::body::SdkBody;
508    use std::pin::Pin;
509
510    #[test]
511    fn valid_size_hint() {
512        assert_eq!(SdkBody::from("hello").content_length(), Some(5));
513        assert_eq!(SdkBody::from("").content_length(), Some(0));
514    }
515
516    #[allow(clippy::bool_assert_comparison)]
517    #[test]
518    fn valid_eos() {
519        assert_eq!(SdkBody::from("hello").is_end_stream(), false);
520        assert_eq!(SdkBody::from("").is_end_stream(), true);
521    }
522
523    #[tokio::test]
524    async fn http_body_consumes_data() {
525        let mut body = SdkBody::from("hello!");
526        let mut body = Pin::new(&mut body);
527        assert!(!body.is_end_stream());
528        let data = body.next().await;
529        assert!(data.is_some());
530        let data = body.next().await;
531        assert!(data.is_none());
532        assert!(body.is_end_stream());
533    }
534
535    #[tokio::test]
536    async fn empty_body_returns_none() {
537        // Its important to avoid sending empty chunks of data to avoid H2 data frame problems
538        let mut body = SdkBody::from("");
539        let mut body = Pin::new(&mut body);
540        let data = body.next().await;
541        assert!(data.is_none());
542    }
543
544    #[test]
545    fn sdkbody_debug_once() {
546        let body = SdkBody::from("123");
547        assert!(format!("{body:?}").contains("Once"));
548    }
549
550    #[test]
551    fn sdk_body_is_send() {
552        fn is_send<T: Send>() {}
553        is_send::<SdkBody>()
554    }
555}