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, 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    /// Uploads one object to the configured bucket.
30    ///
31    /// # Errors
32    ///
33    /// Returns an error when request preparation, signing, transport, or response parsing fails.
34    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    /// Downloads one object as a bounded, cancellation-safe response stream.
80    ///
81    /// # Errors
82    ///
83    /// Returns an error when request validation, signing, transport, or header parsing fails.
84    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    /// Retrieves one object's metadata without downloading its body.
132    ///
133    /// # Errors
134    ///
135    /// Returns an error when signing, transport, or response parsing fails.
136    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    /// Deletes one object or object version.
164    ///
165    /// # Errors
166    ///
167    /// Returns an error when signing, transport, or response parsing fails.
168    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    /// Deletes a validated batch of objects in one request.
196    ///
197    /// # Errors
198    ///
199    /// Returns an error when XML serialization, signing, transport, or response parsing fails.
200    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    /// Copies an object entirely within S3.
236    ///
237    /// # Errors
238    ///
239    /// Returns an error when request validation, signing, transport, or response parsing fails.
240    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            &copy_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}