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, StatusCode};
6use http_body_util::BodyExt;
7use md5::{Digest as _, Md5};
8
9use super::super::S3Client;
10use super::super::request::{protocol_error, service_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 CopyObjectOutput, CopyObjectRequest, DeleteObjectOutput, DeleteObjectRequest,
20 DeleteObjectsOutput, DeleteObjectsRequest, GetObjectOutput, GetObjectRequest, HeadObjectOutput,
21 HeadObjectRequest, PutObjectOutput, PutObjectRequest,
22};
23use crate::protocol::{
24 CopyObjectResponse, parse_copy_object, parse_delete_objects, serialize_delete_objects,
25};
26use crate::stream::{ByteStream, ResponseStream};
27
28impl S3Client {
29 pub async fn put_object(
35 &self,
36 request: PutObjectRequest<ByteStream>,
37 ) -> Result<PutObjectOutput, S3Error> {
38 let PutObjectRequest {
39 key,
40 body,
41 content_type,
42 user_metadata,
43 conditions,
44 checksum_algorithm,
45 } = request;
46 let prepared = body.prepare().await?;
47 let mut headers = HeaderMap::new();
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.object_operation_target(Some(key.as_str()))?;
53 let deadline = self.deadline();
54 let response = self
55 .send_signed(
56 Method::PUT,
57 target,
58 &[],
59 headers,
60 Some(&prepared),
61 &deadline,
62 )
63 .await?;
64 let response_headers = response.headers().clone();
65 self.collect_response(
66 response,
67 self.inner.config.max_xml_response_size(),
68 &deadline,
69 )
70 .await?;
71 Ok(PutObjectOutput {
72 e_tag: response_header(&response_headers, ETAG.as_str())?,
73 version_id: response_header(&response_headers, "x-amz-version-id")?,
74 checksum: parse_checksum(&response_headers)?,
75 request_ids: parse_request_ids(&response_headers)?,
76 })
77 }
78
79 pub async fn get_object(
85 &self,
86 request: GetObjectRequest,
87 ) -> Result<GetObjectOutput<ResponseStream>, S3Error> {
88 let mut headers = HeaderMap::new();
89 insert_conditions(&mut headers, &request.conditions, "")?;
90 insert_header(
91 &mut headers,
92 HeaderName::from_static("x-amz-checksum-mode"),
93 "ENABLED",
94 )?;
95 if let Some(range) = request.range {
96 insert_header(&mut headers, RANGE, &range.to_header_value())?;
97 }
98 let query = optional_query("versionId", request.version_id.as_deref());
99 let target = self.object_operation_target(Some(request.key.as_str()))?;
100 let deadline = self.deadline();
101 let response = self
102 .send_signed(Method::GET, target, &query, headers, None, &deadline)
103 .await?;
104 let metadata = parse_object_metadata(response.headers())?;
105 let content_range = parse_content_range(response.headers())?;
106 let expected_sha256 =
107 verified_download_sha256(response.headers())?.filter(|_| content_range.is_none());
108 let expected_length = response
109 .headers()
110 .get(CONTENT_LENGTH)
111 .map(parse_u64_value)
112 .transpose()?;
113 let stream = response
114 .into_body()
115 .into_data_stream()
116 .map_err(crate::transport::classify_response_body_error);
117 let body = ResponseStream::with_deadline(
118 stream,
119 expected_length,
120 expected_sha256,
121 self.inner.config.idle_body_timeout(),
122 deadline.instant(),
123 );
124 Ok(GetObjectOutput {
125 metadata,
126 body,
127 content_range,
128 })
129 }
130
131 pub async fn head_object(
137 &self,
138 request: HeadObjectRequest,
139 ) -> Result<HeadObjectOutput, S3Error> {
140 let mut headers = HeaderMap::new();
141 insert_conditions(&mut headers, &request.conditions, "")?;
142 insert_header(
143 &mut headers,
144 HeaderName::from_static("x-amz-checksum-mode"),
145 "ENABLED",
146 )?;
147 let query = optional_query("versionId", request.version_id.as_deref());
148 let target = self.object_operation_target(Some(request.key.as_str()))?;
149 let deadline = self.deadline();
150 let response = self
151 .send_signed(Method::HEAD, target, &query, headers, None, &deadline)
152 .await?;
153 let metadata = parse_object_metadata(response.headers())?;
154 self.collect_response(
155 response,
156 self.inner.config.max_xml_response_size(),
157 &deadline,
158 )
159 .await?;
160 Ok(metadata)
161 }
162
163 pub async fn delete_object(
169 &self,
170 request: DeleteObjectRequest,
171 ) -> Result<DeleteObjectOutput, S3Error> {
172 let mut headers = HeaderMap::new();
173 insert_optional_header(&mut headers, IF_MATCH, request.if_match.as_deref())?;
174 let query = optional_query("versionId", request.version_id.as_deref());
175 let target = self.object_operation_target(Some(request.key.as_str()))?;
176 let deadline = self.deadline();
177 let response = self
178 .send_signed(Method::DELETE, target, &query, headers, None, &deadline)
179 .await?;
180 let response_headers = response.headers().clone();
181 self.collect_response(
182 response,
183 self.inner.config.max_xml_response_size(),
184 &deadline,
185 )
186 .await?;
187 Ok(DeleteObjectOutput {
188 delete_marker: parse_bool_header(&response_headers, "x-amz-delete-marker")?
189 .unwrap_or(false),
190 version_id: response_header(&response_headers, "x-amz-version-id")?,
191 request_ids: parse_request_ids(&response_headers)?,
192 })
193 }
194
195 pub async fn delete_objects(
201 &self,
202 request: DeleteObjectsRequest,
203 ) -> Result<DeleteObjectsOutput, S3Error> {
204 let maximum = self.inner.config.max_xml_response_size();
205 let xml = serialize_delete_objects(&request, maximum).map_err(protocol_error)?;
206 let digest = Md5::digest(&xml);
207 let content_md5 = BASE64_STANDARD.encode(digest);
208 let prepared = ByteStream::from_bytes(xml).prepare().await?;
209 let mut headers = HeaderMap::new();
210 insert_header(&mut headers, CONTENT_TYPE, "application/xml")?;
211 insert_header(
212 &mut headers,
213 HeaderName::from_static("content-md5"),
214 &content_md5,
215 )?;
216 let target = self.object_operation_target(None)?;
217 let deadline = self.deadline();
218 let response = self
219 .send_signed(
220 Method::POST,
221 target,
222 &[("delete".to_owned(), String::new())],
223 headers,
224 Some(&prepared),
225 &deadline,
226 )
227 .await?;
228 let response_headers = response.headers().clone();
229 let body = self.collect_response(response, maximum, &deadline).await?;
230 let mut output = parse_delete_objects(&body, maximum).map_err(protocol_error)?;
231 output.request_ids = parse_request_ids(&response_headers)?;
232 Ok(output)
233 }
234
235 pub async fn copy_object(
241 &self,
242 request: CopyObjectRequest,
243 ) -> Result<CopyObjectOutput, S3Error> {
244 let mut headers = HeaderMap::new();
245 insert_header(
246 &mut headers,
247 HeaderName::from_static("x-amz-copy-source"),
248 ©_source_header(&request.source),
249 )?;
250 insert_conditions(
251 &mut headers,
252 &request.source_conditions,
253 "x-amz-copy-source-",
254 )?;
255 let replaces_metadata = request.content_type.is_some() || request.user_metadata.is_some();
256 insert_optional_header(&mut headers, CONTENT_TYPE, request.content_type.as_deref())?;
257 if let Some(metadata) = &request.user_metadata {
258 insert_user_metadata(&mut headers, metadata)?;
259 }
260 if replaces_metadata {
261 insert_header(
262 &mut headers,
263 HeaderName::from_static("x-amz-metadata-directive"),
264 "REPLACE",
265 )?;
266 }
267 let target = self.object_operation_target(Some(request.destination.as_str()))?;
268 let deadline = self.deadline();
269 let response = self
270 .send_signed(Method::PUT, target, &[], headers, None, &deadline)
271 .await?;
272 let response_headers = response.headers().clone();
273 let maximum = self.inner.config.max_xml_response_size();
274 let body = self.collect_response(response, maximum, &deadline).await?;
275 let mut output = match parse_copy_object(&body, maximum).map_err(protocol_error)? {
276 CopyObjectResponse::Complete(output) => output,
277 CopyObjectResponse::EmbeddedError(parsed) => {
278 return Err(service_error(
279 StatusCode::OK,
280 &response_headers,
281 Some(parsed),
282 ));
283 }
284 };
285 output.version_id = response_header(&response_headers, "x-amz-version-id")?;
286 output.request_ids = parse_request_ids(&response_headers)?;
287 merge_checksum(&mut output.checksum, parse_checksum(&response_headers)?);
288 Ok(output)
289 }
290}