Skip to main content

gapirs_common/api/
client.rs

1use crate::api::client::media::multipart::{MultipartBody, MultipartState};
2use crate::api::client::media::{MediaUpload, MediaUploadStream};
3use crate::errors::ApiCallError;
4use crate::ConnectRequirements;
5use bytes::{Buf, Bytes};
6use derive_more::Debug;
7use futures_core::Stream;
8use futures_util::AsyncRead;
9use http_body_util::combinators::BoxBody;
10use http_body_util::{BodyExt, Full, StreamBody};
11use hyper::body::{Body, Frame};
12use media::{AsyncMediaUpload, AsyncMediaUploadStream};
13use serde::{Deserialize, Serialize};
14use std::error::Error;
15use std::io::{Read, Seek};
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19use url::Url;
20
21#[derive(Default, Debug)]
22pub enum OutgoingBodyContent<const N: usize = { 1024 * 10 }> {
23    #[default]
24    Empty,
25    Json(String),
26    Stream(MediaUpload<dyn MediaUploadStream>),
27    AsyncStream(AsyncMediaUpload<dyn AsyncMediaUploadStream>),
28    Multipart(MultipartBody),
29}
30impl OutgoingBodyContent {
31    pub fn from_body_and_media(
32        body: Option<impl Into<String>>,
33        media: Option<MediaUpload<dyn MediaUploadStream>>,
34    ) -> Self {
35        match (body, media) {
36            (Some(body), Some(media)) => {
37                OutgoingBodyContent::Multipart(MultipartBody::new(media, body))
38            }
39            (None, Some(media)) => OutgoingBodyContent::Stream(media),
40            (Some(body), None) => OutgoingBodyContent::Json(body.into()),
41            (None, None) => OutgoingBodyContent::Empty,
42        }
43    }
44    pub fn get_content_type(&self) -> Option<String> {
45        Some(match self {
46            OutgoingBodyContent::Empty => return None,
47            OutgoingBodyContent::Json(_) => "application/json".into(),
48            OutgoingBodyContent::Stream(x) => x.mime_type.as_str().into(),
49            OutgoingBodyContent::AsyncStream(x) => x.mime_type.as_str().into(),
50            OutgoingBodyContent::Multipart(multipart) => {
51                format!("multipart/related; boundary={}", multipart.boundary)
52            }
53        })
54    }
55
56    pub fn get_length(&self) -> u64 {
57        match self {
58            OutgoingBodyContent::Empty => 0,
59            OutgoingBodyContent::Json(string) => Bytes::from(string.to_string()).len() as u64,
60            OutgoingBodyContent::Stream(x) => x.length,
61            OutgoingBodyContent::AsyncStream(x) => x.length,
62            OutgoingBodyContent::Multipart(x) => {
63                x.get_body_before_media().len() as u64
64                    + x.media_body.length
65                    + x.get_body_after_media().len() as u64
66            }
67        }
68    }
69
70    fn poll_stream(
71        cx: &mut Context,
72        stream: &mut MediaUpload<dyn MediaUploadStream>,
73    ) -> Poll<Option<Result<Bytes, Box<dyn Error + Send + Sync + 'static>>>> {
74        stream.body.as_mut().poll_next(cx)
75    }
76}
77impl<const ASYNC_STREAM_BUF_SIZE: usize> OutgoingBodyContent<ASYNC_STREAM_BUF_SIZE> {
78    fn poll_stream_async(
79        cx: &mut Context,
80        stream: &mut AsyncMediaUpload<dyn AsyncMediaUploadStream>,
81    ) -> Poll<Option<Result<Bytes, Box<dyn Error + Send + Sync + 'static>>>> {
82        let mut buf = [0u8; ASYNC_STREAM_BUF_SIZE];
83        let x = Pin::new(&mut stream.body).poll_read(cx, &mut buf)?;
84        match x {
85            Poll::Ready(x) => {
86                let x = Bytes::copy_from_slice(&buf[0..x]);
87                Poll::Ready(Some(Ok(x)))
88            }
89            Poll::Pending => Poll::Pending,
90        }
91    }
92}
93
94type PollResult<T, E> = Poll<Option<Result<T, Box<E>>>>;
95fn bytes_to_frame<E: ?Sized>(data: PollResult<Bytes, E>) -> PollResult<Frame<Bytes>, E> {
96    // println!("returning bytes: {data:?}");
97    match data {
98        Poll::Ready(Some(Ok(data))) => {
99            if data.is_empty() {
100                Poll::Ready(None)
101            } else {
102                Poll::Ready(Some(Ok(Frame::data(data))))
103            }
104        }
105        Poll::Ready(Some(Err(err))) => Poll::Ready(Some(Err(err))),
106        Poll::Ready(None) => Poll::Ready(None),
107        Poll::Pending => Poll::Pending,
108    }
109}
110pub mod media;
111
112impl Body for OutgoingBodyContent {
113    type Data = Bytes;
114    type Error = Box<dyn Error + Send + Sync + 'static>;
115
116    fn poll_frame(
117        self: Pin<&mut Self>,
118        cx: &mut Context<'_>,
119    ) -> Poll<Option<Result<Frame<Self::Data>, Self::Error>>> {
120        match &*self {
121            OutgoingBodyContent::Empty => Poll::Ready(None),
122            OutgoingBodyContent::Json(value) => {
123                let json_str = value.to_string();
124                *self.get_mut() = OutgoingBodyContent::Empty;
125                Poll::Ready(Some(Ok(Frame::data(Bytes::from(json_str)))))
126            }
127            OutgoingBodyContent::Multipart(_) => {
128                use media::multipart::MultipartState;
129                let multipart = self.get_mut();
130                let multipart = match multipart {
131                    OutgoingBodyContent::Multipart(multipart) => multipart,
132                    _ => unreachable!(),
133                };
134                // dbg!("Polling multipart", &multipart);
135                match multipart.state {
136                    MultipartState::NotStarted => {
137                        multipart.state = MultipartState::Polling;
138                        let body_string = multipart.get_body_before_media();
139                        Poll::Ready(Some(Ok(Frame::data(Bytes::from(body_string)))))
140                    }
141                    MultipartState::Polling => {
142                        let poll = Self::poll_stream(cx, &mut multipart.media_body);
143                        match poll {
144                            Poll::Ready(None) => {
145                                multipart.state = MultipartState::Done;
146                                let string = multipart.get_body_after_media();
147                                Poll::Ready(Some(Ok(Frame::data(Bytes::from(string)))))
148                            }
149                            Poll::Ready(Some(Ok(x))) => Poll::Ready(Some(Ok(Frame::data(x)))),
150                            _ => bytes_to_frame(poll),
151                        }
152                    }
153                    MultipartState::Done => Poll::Ready(None),
154                }
155            }
156            OutgoingBodyContent::Stream(_) => {
157                let stream = self.get_mut();
158                let stream = match stream {
159                    OutgoingBodyContent::Stream(stream) => stream,
160                    _ => unreachable!(),
161                };
162
163                bytes_to_frame(Self::poll_stream(cx, stream))
164            }
165            OutgoingBodyContent::AsyncStream(_) => {
166                let stream = self.get_mut();
167                let stream = match stream {
168                    OutgoingBodyContent::AsyncStream(stream) => stream,
169                    _ => unreachable!(),
170                };
171
172                bytes_to_frame(Self::poll_stream_async(cx, stream))
173            }
174        }
175    }
176
177    fn is_end_stream(&self) -> bool {
178        match self {
179            OutgoingBodyContent::Empty => true,
180            OutgoingBodyContent::Json(_) => false, // we will turn it into empty after first poll
181            OutgoingBodyContent::Stream(_) => false,
182            OutgoingBodyContent::AsyncStream(_) => false,
183            OutgoingBodyContent::Multipart(_) => false,
184        }
185    }
186
187    fn size_hint(&self) -> http_body::SizeHint {
188        http_body::SizeHint::default()
189    }
190}