1use base64::Engine as _;
2use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
3use futures_util::TryStreamExt;
4use http::header::{CONTENT_LENGTH, CONTENT_TYPE, ETAG, IF_MATCH, RANGE};
5use http::{HeaderMap, HeaderName, Method};
6use http_body_util::BodyExt;
7use md5::{Digest as _, Md5};
8
9use super::super::S3Client;
10use super::super::request::protocol_error;
11use super::headers::{
12 copy_source_header, insert_conditions, insert_header, insert_optional_header,
13 insert_upload_checksum, insert_user_metadata, merge_checksum, optional_query,
14 parse_bool_header, parse_checksum, parse_content_range, parse_object_metadata,
15 parse_request_ids, parse_u64_value, response_header, verified_download_sha256,
16};
17use crate::error::S3Error;
18use crate::operation::{
19 CopyMetadataDirective, CopyObjectOutput, CopyObjectRequest, DeleteObjectOutput,
20 DeleteObjectRequest, DeleteObjectsOutput, DeleteObjectsRequest, GetObjectOutput,
21 GetObjectRequest, HeadObjectOutput, HeadObjectRequest, PutObjectOutput, PutObjectRequest,
22};
23use crate::protocol::{
24 CopyObjectResponse, ParsedS3Error, ProtocolError, parse_copy_object, parse_delete_objects,
25 serialize_delete_objects,
26};
27use crate::stream::{ByteStream, ResponseStream};
28
29impl S3Client {
30 pub async fn put_object(&self, request: PutObjectRequest) -> Result<PutObjectOutput, S3Error> {
36 let deadline = self.deadline();
37 let PutObjectRequest {
38 key,
39 body,
40 content_type,
41 user_metadata,
42 conditions,
43 checksum_algorithm,
44 } = request;
45 let prepared = deadline.prepare_body(body).await?;
46 let mut headers = HeaderMap::new();
47 insert_optional_header(&mut headers, CONTENT_TYPE, content_type.as_deref())?;
48 insert_conditions(&mut headers, &conditions, "")?;
49 insert_user_metadata(&mut headers, &user_metadata)?;
50 insert_upload_checksum(&mut headers, checksum_algorithm, &prepared)?;
51 let target = self.operation_target(Some(key.as_str()))?;
52 let response = self
53 .send_signed(
54 Method::PUT,
55 target,
56 &[],
57 headers,
58 Some(&prepared),
59 &deadline,
60 )
61 .await?;
62 let response_headers = response.headers().clone();
63 self.collect_response(
64 response,
65 self.inner.config.max_xml_response_size(),
66 &deadline,
67 )
68 .await?;
69 Ok(PutObjectOutput {
70 e_tag: response_header(&response_headers, ETAG.as_str())?,
71 version_id: response_header(&response_headers, "x-amz-version-id")?,
72 checksum: parse_checksum(&response_headers)?,
73 request_ids: parse_request_ids(&response_headers)?,
74 })
75 }
76
77 pub async fn get_object(&self, request: GetObjectRequest) -> Result<GetObjectOutput, S3Error> {
83 let mut headers = HeaderMap::new();
84 insert_conditions(&mut headers, &request.conditions, "")?;
85 insert_header(
86 &mut headers,
87 HeaderName::from_static("x-amz-checksum-mode"),
88 "ENABLED",
89 )?;
90 if let Some(range) = request.range {
91 insert_header(&mut headers, RANGE, &range.to_header_value())?;
92 }
93 let query = optional_query("versionId", request.version_id.as_deref());
94 let target = self.operation_target(Some(request.key.as_str()))?;
95 let deadline = self.deadline();
96 let response = self
97 .send_signed(Method::GET, target, &query, headers, None, &deadline)
98 .await?;
99 let metadata = parse_object_metadata(response.headers())?;
100 let content_range = parse_content_range(response.headers())?;
101 let expected_sha256 =
102 verified_download_sha256(response.headers())?.filter(|_| content_range.is_none());
103 let expected_length = response
104 .headers()
105 .get(CONTENT_LENGTH)
106 .map(parse_u64_value)
107 .transpose()?;
108 let stream = response
109 .into_body()
110 .into_data_stream()
111 .map_err(crate::transport::classify_response_body_error);
112 let body = ResponseStream::with_deadline(
113 stream,
114 expected_length,
115 expected_sha256,
116 self.inner.config.idle_body_timeout(),
117 deadline.instant(),
118 );
119 Ok(GetObjectOutput {
120 metadata,
121 body,
122 content_range,
123 })
124 }
125
126 pub async fn head_object(
132 &self,
133 request: HeadObjectRequest,
134 ) -> Result<HeadObjectOutput, S3Error> {
135 let mut headers = HeaderMap::new();
136 insert_conditions(&mut headers, &request.conditions, "")?;
137 insert_header(
138 &mut headers,
139 HeaderName::from_static("x-amz-checksum-mode"),
140 "ENABLED",
141 )?;
142 let query = optional_query("versionId", request.version_id.as_deref());
143 let target = self.operation_target(Some(request.key.as_str()))?;
144 let deadline = self.deadline();
145 let response = self
146 .send_signed(Method::HEAD, target, &query, headers, None, &deadline)
147 .await?;
148 let metadata = parse_object_metadata(response.headers())?;
149 self.collect_response(
150 response,
151 self.inner.config.max_xml_response_size(),
152 &deadline,
153 )
154 .await?;
155 Ok(metadata)
156 }
157
158 pub async fn delete_object(
164 &self,
165 request: DeleteObjectRequest,
166 ) -> Result<DeleteObjectOutput, S3Error> {
167 let mut headers = HeaderMap::new();
168 insert_optional_header(&mut headers, IF_MATCH, request.if_match.as_deref())?;
169 let query = optional_query("versionId", request.version_id.as_deref());
170 let target = self.operation_target(Some(request.key.as_str()))?;
171 let deadline = self.deadline();
172 let response = self
173 .send_signed(Method::DELETE, target, &query, headers, None, &deadline)
174 .await?;
175 let response_headers = response.headers().clone();
176 self.collect_response(
177 response,
178 self.inner.config.max_xml_response_size(),
179 &deadline,
180 )
181 .await?;
182 Ok(DeleteObjectOutput {
183 delete_marker: parse_bool_header(&response_headers, "x-amz-delete-marker")?
184 .unwrap_or(false),
185 version_id: response_header(&response_headers, "x-amz-version-id")?,
186 request_ids: parse_request_ids(&response_headers)?,
187 })
188 }
189
190 pub async fn delete_objects(
196 &self,
197 request: DeleteObjectsRequest,
198 ) -> Result<DeleteObjectsOutput, S3Error> {
199 let deadline = self.deadline();
200 let maximum = self.inner.config.max_xml_response_size();
201 let xml = serialize_delete_objects(&request, maximum).map_err(protocol_error)?;
202 let digest = Md5::digest(&xml);
203 let content_md5 = BASE64_STANDARD.encode(digest);
204 let prepared = deadline.prepare_body(ByteStream::from_bytes(xml)).await?;
205 let mut headers = HeaderMap::new();
206 insert_header(&mut headers, CONTENT_TYPE, "application/xml")?;
207 insert_header(
208 &mut headers,
209 HeaderName::from_static("content-md5"),
210 &content_md5,
211 )?;
212 let target = self.operation_target(None)?;
213 let response = self
214 .send_signed(
215 Method::POST,
216 target,
217 &[("delete".to_owned(), String::new())],
218 headers,
219 Some(&prepared),
220 &deadline,
221 )
222 .await?;
223 let response_headers = response.headers().clone();
224 let body = self.collect_response(response, maximum, &deadline).await?;
225 let mut output = parse_delete_objects(&body, maximum).map_err(protocol_error)?;
226 output.request_ids = parse_request_ids(&response_headers)?;
227 Ok(output)
228 }
229
230 pub async fn copy_object(
236 &self,
237 request: CopyObjectRequest,
238 ) -> Result<CopyObjectOutput, S3Error> {
239 let mut headers = HeaderMap::new();
240 insert_header(
241 &mut headers,
242 HeaderName::from_static("x-amz-copy-source"),
243 ©_source_header(&request.source),
244 )?;
245 insert_conditions(
246 &mut headers,
247 &request.source_conditions,
248 "x-amz-copy-source-",
249 )?;
250 match &request.metadata {
251 CopyMetadataDirective::Copy => {}
252 CopyMetadataDirective::Replace {
253 content_type,
254 user_metadata,
255 } => {
256 insert_optional_header(&mut headers, CONTENT_TYPE, content_type.as_deref())?;
257 insert_user_metadata(&mut headers, user_metadata)?;
258 insert_header(
259 &mut headers,
260 HeaderName::from_static("x-amz-metadata-directive"),
261 "REPLACE",
262 )?;
263 }
264 }
265 let target = self.operation_target(Some(request.destination.as_str()))?;
266 let deadline = self.deadline();
267 let maximum = self.inner.config.max_xml_response_size();
268 let response = self
269 .send_signed_collected_xml(
270 Method::PUT,
271 target,
272 &[],
273 headers,
274 None,
275 &deadline,
276 maximum,
277 copy_embedded_error,
278 )
279 .await?;
280 let response_headers = response.headers;
281 let mut output = match parse_copy_object(&response.body, maximum).map_err(protocol_error)? {
282 CopyObjectResponse::Complete(output) => output,
283 CopyObjectResponse::EmbeddedError(_) => {
284 unreachable!("embedded copy errors are handled by request execution")
285 }
286 };
287 output.version_id = response_header(&response_headers, "x-amz-version-id")?;
288 output.request_ids = parse_request_ids(&response_headers)?;
289 merge_checksum(&mut output.checksum, parse_checksum(&response_headers)?);
290 Ok(output)
291 }
292}
293
294fn copy_embedded_error(
295 body: &[u8],
296 maximum: usize,
297) -> Result<Option<ParsedS3Error>, ProtocolError> {
298 match parse_copy_object(body, maximum)? {
299 CopyObjectResponse::Complete(_) => Ok(None),
300 CopyObjectResponse::EmbeddedError(error) => Ok(Some(error)),
301 }
302}