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::{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    /// 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        } = 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    /// Downloads one object as a bounded, cancellation-safe response stream.
78    ///
79    /// # Errors
80    ///
81    /// Returns an error when request validation, signing, transport, or header parsing fails.
82    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    /// Retrieves one object's metadata without downloading its body.
127    ///
128    /// # Errors
129    ///
130    /// Returns an error when signing, transport, or response parsing fails.
131    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    /// Deletes one object or object version.
159    ///
160    /// # Errors
161    ///
162    /// Returns an error when signing, transport, or response parsing fails.
163    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    /// Deletes a validated batch of objects in one request.
191    ///
192    /// # Errors
193    ///
194    /// Returns an error when XML serialization, signing, transport, or response parsing fails.
195    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    /// Copies an object entirely within S3.
231    ///
232    /// # Errors
233    ///
234    /// Returns an error when request validation, signing, transport, or response parsing fails.
235    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            &copy_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}