gapirs_common/api/
client.rs1use 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 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 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, 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}