1use 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#[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
26pub type Error = Box<dyn StdError + Send + Sync>;
28
29#[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#[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 pub struct SdkBody {
74 #[pin]
75 inner: Inner,
76 rebuild: Option<Arc<dyn (Fn() -> Inner) + Send + Sync>>,
81 bytes_contents: Option<Bytes>,
82 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#[allow(dead_code)]
100enum BoxBody {
101 #[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 Once {
116 inner: Option<Bytes>
117 },
118 Dyn {
120 #[pin]
121 inner: BoxBody,
122 },
123
124 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 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 pub fn taken() -> Self {
163 Self {
164 inner: Inner::Taken,
165 rebuild: None,
166 bytes_contents: None,
167 trailers: None,
168 }
169 }
170
171 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.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 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 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 pub fn bytes(&self) -> Option<&[u8]> {
347 match &self.bytes_contents {
348 Some(b) => Some(b),
349 None => None,
350 }
351 }
352
353 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 pub fn is_streaming(&self) -> bool {
369 matches!(self.inner, Inner::Dyn { .. })
370 }
371
372 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)] 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 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 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 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}