Skip to main content

s3_wire/client/object/
operations.rs

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    /// Uploads one object to the configured bucket.
31    ///
32    /// # Errors
33    ///
34    /// Returns an error when request preparation, signing, transport, or response parsing fails.
35    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    /// Downloads one object as a bounded, cancellation-safe response stream.
79    ///
80    /// # Errors
81    ///
82    /// Returns an error when request validation, signing, transport, or header parsing fails.
83    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    /// Retrieves one object's metadata without downloading its body.
128    ///
129    /// # Errors
130    ///
131    /// Returns an error when signing, transport, or response parsing fails.
132    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    /// Deletes one object or object version.
160    ///
161    /// # Errors
162    ///
163    /// Returns an error when signing, transport, or response parsing fails.
164    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    /// Deletes a validated batch of objects in one request.
192    ///
193    /// # Errors
194    ///
195    /// Returns an error when XML serialization, signing, transport, or response parsing fails.
196    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    /// Copies an object entirely within S3.
232    ///
233    /// # Errors
234    ///
235    /// Returns an error when request validation, signing, transport, or response parsing fails.
236    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            &copy_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}