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