Skip to main content

s3_wire/client/multipart/
operations.rs

1use http::header::{CONTENT_TYPE, ETAG};
2use http::{HeaderMap, HeaderValue, Method};
3
4use super::headers::{
5    checksum_headers, create_headers, optional_header, request_ids, required_header,
6    response_checksums,
7};
8use super::query::{create_query, list_query, upload_part_query, upload_query};
9use crate::client::S3Client;
10use crate::client::request::OperationDeadline;
11use crate::error::S3Error;
12use crate::operation::{
13    AbortMultipartUploadRequest, CompleteMultipartUploadOutput, CompleteMultipartUploadRequest,
14    CreateMultipartUploadOutput, CreateMultipartUploadRequest, ListMultipartUploadsOutput,
15    ListMultipartUploadsRequest, RequestIds, UploadPartOutput, UploadPartRequest,
16};
17use crate::protocol::{
18    CompleteMultipartResponse, ParsedS3Error, ProtocolError, parse_complete_multipart_upload,
19    parse_create_multipart_upload, parse_list_multipart_uploads,
20    serialize_complete_multipart_upload,
21};
22use crate::stream::ByteStream;
23
24impl S3Client {
25    /// Initiates a multipart upload for an object.
26    ///
27    /// # Errors
28    ///
29    /// Returns an error for invalid headers, transport or service failures, or
30    /// a malformed or oversized service response.
31    pub async fn create_multipart_upload(
32        &self,
33        request: CreateMultipartUploadRequest,
34    ) -> Result<CreateMultipartUploadOutput, S3Error> {
35        let deadline = self.deadline();
36        self.create_multipart_upload_with_deadline(request, &deadline)
37            .await
38    }
39
40    pub(in crate::client) async fn create_multipart_upload_with_deadline(
41        &self,
42        request: CreateMultipartUploadRequest,
43        deadline: &OperationDeadline,
44    ) -> Result<CreateMultipartUploadOutput, S3Error> {
45        let target = self.operation_target(Some(request.key.as_str()))?;
46        let headers = create_headers(&request)?;
47        let response = self
48            .send_signed(
49                Method::POST,
50                target,
51                &create_query(),
52                headers,
53                None,
54                deadline,
55            )
56            .await?;
57        let request_ids = request_ids(response.headers());
58        let body = self
59            .collect_response(response, self.config().max_xml_response_size(), deadline)
60            .await?;
61        let mut output =
62            parse_create_multipart_upload(&body, self.config().max_xml_response_size())
63                .map_err(crate::client::request::protocol_error)?;
64        output.request_ids = request_ids;
65        Ok(output)
66    }
67
68    /// Uploads one part of an in-progress multipart upload.
69    ///
70    /// # Errors
71    ///
72    /// Returns an error for an invalid checksum, body preparation failure,
73    /// transport or service failure, or malformed response headers.
74    pub async fn upload_part(
75        &self,
76        request: UploadPartRequest,
77    ) -> Result<UploadPartOutput, S3Error> {
78        let deadline = self.deadline();
79        self.upload_part_with_deadline(request, &deadline).await
80    }
81
82    pub(in crate::client) async fn upload_part_with_deadline(
83        &self,
84        request: UploadPartRequest,
85        deadline: &OperationDeadline,
86    ) -> Result<UploadPartOutput, S3Error> {
87        let target = self.operation_target(Some(request.key().as_str()))?;
88        let query = upload_part_query(request.part_number().get(), request.upload_id());
89        let headers = checksum_headers(request.checksum())?;
90        let part_number = request.part_number();
91        let body = deadline.prepare_body(request.into_body()).await?;
92        let response = self
93            .send_signed(Method::PUT, target, &query, headers, Some(&body), deadline)
94            .await?;
95        let headers = response.headers().clone();
96        self.drain_success_response(response, deadline).await?;
97        let e_tag = required_header(
98            &headers,
99            ETAG.as_str(),
100            "upload-part response has no usable ETag",
101        )?;
102        let checksum = response_checksums(&headers)?;
103        let request_ids = request_ids(&headers);
104        Ok(UploadPartOutput {
105            part_number,
106            e_tag,
107            checksum,
108            request_ids,
109        })
110    }
111
112    /// Completes an in-progress multipart upload.
113    ///
114    /// # Errors
115    ///
116    /// Returns an error when the completion document cannot be serialized, the
117    /// request fails, or S3 returns a malformed, oversized, or embedded error.
118    pub async fn complete_multipart_upload(
119        &self,
120        request: CompleteMultipartUploadRequest,
121    ) -> Result<CompleteMultipartUploadOutput, S3Error> {
122        let deadline = self.deadline();
123        self.complete_multipart_upload_with_deadline(request, &deadline)
124            .await
125    }
126
127    pub(in crate::client) async fn complete_multipart_upload_with_deadline(
128        &self,
129        request: CompleteMultipartUploadRequest,
130        deadline: &OperationDeadline,
131    ) -> Result<CompleteMultipartUploadOutput, S3Error> {
132        let target = self.operation_target(Some(request.key.as_str()))?;
133        let query = upload_query(request.upload_id());
134        let document =
135            serialize_complete_multipart_upload(&request, self.config().max_xml_response_size())
136                .map_err(crate::client::request::protocol_error)?;
137        let body = deadline.prepare_body(ByteStream::from(document)).await?;
138        let mut headers = HeaderMap::new();
139        headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/xml"));
140        let maximum = self.config().max_xml_response_size();
141        let response = self
142            .send_signed_collected_xml(
143                Method::POST,
144                target,
145                &query,
146                headers,
147                Some(&body),
148                deadline,
149                maximum,
150                complete_multipart_embedded_error,
151            )
152            .await?;
153        let response_headers = response.headers;
154        let request_ids = request_ids(&response_headers);
155        let version_id = optional_header(&response_headers, "x-amz-version-id")?;
156        match parse_complete_multipart_upload(&response.body, maximum)
157            .map_err(crate::client::request::protocol_error)?
158        {
159            CompleteMultipartResponse::Complete(mut output) => {
160                output.version_id = version_id;
161                output.request_ids = request_ids;
162                Ok(output)
163            }
164            CompleteMultipartResponse::EmbeddedError(_) => {
165                unreachable!("embedded completion errors are handled by request execution")
166            }
167        }
168    }
169
170    /// Aborts an in-progress multipart upload.
171    ///
172    /// # Errors
173    ///
174    /// Returns an error when target construction, signing, transport, or the
175    /// S3 service rejects the abort request.
176    pub async fn abort_multipart_upload(
177        &self,
178        request: AbortMultipartUploadRequest,
179    ) -> Result<RequestIds, S3Error> {
180        let deadline = self.deadline();
181        self.abort_multipart_upload_with_deadline(request, &deadline)
182            .await
183    }
184
185    pub(in crate::client) async fn abort_multipart_upload_with_deadline(
186        &self,
187        request: AbortMultipartUploadRequest,
188        deadline: &OperationDeadline,
189    ) -> Result<RequestIds, S3Error> {
190        let target = self.operation_target(Some(request.key().as_str()))?;
191        let query = upload_query(request.upload_id());
192        let response = self
193            .send_signed(
194                Method::DELETE,
195                target,
196                &query,
197                HeaderMap::new(),
198                None,
199                deadline,
200            )
201            .await?;
202        let request_ids = request_ids(response.headers());
203        self.drain_success_response(response, deadline).await?;
204        Ok(request_ids)
205    }
206
207    /// Retrieves one bounded page of in-progress multipart uploads.
208    ///
209    /// # Errors
210    ///
211    /// Returns an error for inconsistent pagination markers, a failed request,
212    /// or a malformed or oversized listing response.
213    pub async fn list_multipart_uploads(
214        &self,
215        request: ListMultipartUploadsRequest,
216    ) -> Result<ListMultipartUploadsOutput, S3Error> {
217        let deadline = self.deadline();
218        self.list_multipart_uploads_with_deadline(&request, &deadline)
219            .await
220    }
221
222    pub(in crate::client) async fn list_multipart_uploads_with_deadline(
223        &self,
224        request: &ListMultipartUploadsRequest,
225        deadline: &OperationDeadline,
226    ) -> Result<ListMultipartUploadsOutput, S3Error> {
227        let target = self.operation_target(None)?;
228        let query = list_query(request)?;
229        let response = self
230            .send_signed(
231                Method::GET,
232                target,
233                &query,
234                HeaderMap::new(),
235                None,
236                deadline,
237            )
238            .await?;
239        let request_ids = request_ids(response.headers());
240        let body = self
241            .collect_response(response, self.config().max_xml_response_size(), deadline)
242            .await?;
243        let mut output = parse_list_multipart_uploads(&body, self.config().max_xml_response_size())
244            .map_err(crate::client::request::protocol_error)?;
245        output.request_ids = request_ids;
246        Ok(output)
247    }
248}
249
250fn complete_multipart_embedded_error(
251    body: &[u8],
252    maximum: usize,
253) -> Result<Option<ParsedS3Error>, ProtocolError> {
254    match parse_complete_multipart_upload(body, maximum)? {
255        CompleteMultipartResponse::Complete(_) => Ok(None),
256        CompleteMultipartResponse::EmbeddedError(error) => Ok(Some(error)),
257    }
258}