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