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