Skip to main content

s3_wire/operation/multipart/
managed.rs

1use std::collections::BTreeMap;
2use std::path::{Path, PathBuf};
3use std::time::Duration;
4
5use bytes::Bytes;
6use http::HeaderMap;
7
8use super::{CompletedPart, MultipartError, UploadId};
9use crate::operation::{ChecksumAlgorithm, ChecksumType, ObjectKey, RequestIds};
10
11const MIN_PART_SIZE: u64 = 5 * 1024 * 1024;
12const MAX_PART_SIZE: u64 = 5 * 1024 * 1024 * 1024;
13const MAX_CONCURRENCY: usize = 64;
14
15/// Validated resource and time bounds for one managed multipart upload.
16#[derive(Clone, Copy, Debug, Eq, PartialEq)]
17pub struct MultipartOptions {
18    part_size: u64,
19    concurrency: usize,
20    transfer_timeout: Duration,
21    cleanup_timeout: Duration,
22}
23
24impl MultipartOptions {
25    /// Creates multipart options with bounded transfer and cleanup defaults.
26    ///
27    /// # Errors
28    ///
29    /// Returns an error unless `part_size` is between 5 MiB and 5 GiB and
30    /// `concurrency` is between 1 and 64.
31    pub fn new(part_size: u64, concurrency: usize) -> Result<Self, crate::error::S3Error> {
32        if !(MIN_PART_SIZE..=MAX_PART_SIZE).contains(&part_size) {
33            return Err(crate::error::S3Error::configuration(
34                "multipart part size must be between 5 MiB and 5 GiB",
35            ));
36        }
37        if !(1..=MAX_CONCURRENCY).contains(&concurrency) {
38            return Err(crate::error::S3Error::configuration(
39                "multipart concurrency must be between 1 and 64",
40            ));
41        }
42        part_size
43            .checked_mul(u64::try_from(concurrency).map_err(|_| {
44                crate::error::S3Error::configuration("multipart concurrency does not fit in u64")
45            })?)
46            .ok_or_else(|| {
47                crate::error::S3Error::configuration("multipart buffered byte bound overflow")
48            })?;
49        Ok(Self {
50            part_size,
51            concurrency,
52            transfer_timeout: Duration::from_secs(5 * 60),
53            cleanup_timeout: Duration::from_secs(30),
54        })
55    }
56
57    /// Sets the deadline for preparation, upload, and completion.
58    ///
59    /// # Errors
60    ///
61    /// Returns an error when `timeout` is zero.
62    pub fn with_transfer_timeout(
63        mut self,
64        timeout: Duration,
65    ) -> Result<Self, crate::error::S3Error> {
66        if timeout.is_zero() {
67            return Err(crate::error::S3Error::configuration(
68                "multipart transfer timeout must be greater than zero",
69            ));
70        }
71        validate_deadline(timeout, "multipart transfer timeout")?;
72        self.transfer_timeout = timeout;
73        Ok(self)
74    }
75
76    /// Sets the separate deadline for quiescing in-flight parts and aborting.
77    ///
78    /// # Errors
79    ///
80    /// Returns an error when `timeout` is zero.
81    pub fn with_cleanup_timeout(
82        mut self,
83        timeout: Duration,
84    ) -> Result<Self, crate::error::S3Error> {
85        if timeout.is_zero() {
86            return Err(crate::error::S3Error::configuration(
87                "multipart cleanup timeout must be greater than zero",
88            ));
89        }
90        validate_deadline(timeout, "multipart cleanup timeout")?;
91        self.cleanup_timeout = timeout;
92        Ok(self)
93    }
94
95    /// Returns the size of each part except the final part.
96    pub const fn part_size(self) -> u64 {
97        self.part_size
98    }
99
100    /// Returns the maximum number of part requests in flight.
101    pub const fn concurrency(self) -> usize {
102        self.concurrency
103    }
104
105    /// Returns the derived maximum bytes retained by in-flight part buffers.
106    pub fn maximum_buffered_bytes(self) -> u64 {
107        self.part_size * u64::try_from(self.concurrency).expect("validated concurrency fits in u64")
108    }
109
110    /// Returns the transfer deadline duration.
111    pub const fn transfer_timeout(self) -> Duration {
112        self.transfer_timeout
113    }
114
115    /// Returns the part-quiescing and abort cleanup deadline duration.
116    pub const fn cleanup_timeout(self) -> Duration {
117        self.cleanup_timeout
118    }
119}
120
121fn validate_deadline(timeout: Duration, name: &str) -> Result<(), crate::error::S3Error> {
122    std::time::Instant::now()
123        .checked_add(timeout)
124        .map(|_| ())
125        .ok_or_else(|| {
126            crate::error::S3Error::configuration(format!(
127                "{name} is too large to represent as a deadline"
128            ))
129        })
130}
131
132impl Default for MultipartOptions {
133    fn default() -> Self {
134        Self::new(8 * 1024 * 1024, 4).expect("default multipart options are valid")
135    }
136}
137
138/// Request to initiate a multipart upload.
139#[derive(Clone, derive_more::Debug, Eq, PartialEq)]
140pub struct CreateMultipartUploadRequest {
141    /// Destination object key.
142    pub key: ObjectKey,
143    /// Optional object media type.
144    pub content_type: Option<String>,
145    /// User-defined object metadata.
146    #[debug("{:?}", self.user_metadata.keys().collect::<Vec<_>>())]
147    pub user_metadata: BTreeMap<String, String>,
148    /// Checksum algorithm applied to uploaded parts.
149    pub checksum_algorithm: Option<ChecksumAlgorithm>,
150    /// How S3 should derive the completed object's checksum from its parts.
151    pub checksum_type: Option<ChecksumType>,
152    /// Additional request headers. Values are signed and repeated values are preserved.
153    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
154    #[debug("{:?}", "<redacted>")]
155    pub headers: HeaderMap,
156}
157
158impl CreateMultipartUploadRequest {
159    /// Constructs a request without optional metadata.
160    pub fn new(key: ObjectKey) -> Self {
161        Self {
162            key,
163            content_type: None,
164            user_metadata: BTreeMap::new(),
165            checksum_algorithm: None,
166            checksum_type: None,
167            headers: HeaderMap::new(),
168        }
169    }
170
171    /// Replaces the request's additional headers.
172    pub fn with_headers(mut self, headers: HeaderMap) -> Self {
173        self.headers = headers;
174        self
175    }
176}
177
178/// Result of initiating a multipart upload.
179#[derive(Clone, Debug, Eq, PartialEq)]
180pub struct CreateMultipartUploadOutput {
181    /// Bucket reported by the service.
182    pub bucket: Option<String>,
183    /// Object key reported by the service.
184    pub key: ObjectKey,
185    /// Opaque upload identifier required by later operations.
186    upload_id: UploadId,
187    /// Checksum algorithm selected by the service.
188    pub checksum_algorithm: Option<String>,
189    /// Checksum aggregation selected by S3.
190    pub checksum_type: Option<ChecksumType>,
191    /// Service request identifiers.
192    pub request_ids: RequestIds,
193}
194
195impl CreateMultipartUploadOutput {
196    pub(crate) fn new(
197        bucket: Option<String>,
198        key: ObjectKey,
199        upload_id: UploadId,
200        checksum_algorithm: Option<String>,
201        checksum_type: Option<ChecksumType>,
202    ) -> Self {
203        Self {
204            bucket,
205            key,
206            upload_id,
207            checksum_algorithm,
208            checksum_type,
209            request_ids: RequestIds::default(),
210        }
211    }
212
213    /// Returns the validated upload identifier required by later operations.
214    pub const fn upload_id(&self) -> &UploadId {
215        &self.upload_id
216    }
217}
218
219/// Explicit state needed to resume or clean up an in-progress multipart upload.
220#[derive(Clone, Debug, Eq, PartialEq)]
221pub struct MultipartUpload {
222    key: ObjectKey,
223    upload_id: UploadId,
224    /// Parts that have completed successfully, in ascending order.
225    completed_parts: Vec<CompletedPart>,
226}
227
228pub(crate) enum MultipartUploadSource {
229    Bytes(Bytes),
230    File(PathBuf),
231}
232
233/// Phase-specific headers for a managed multipart upload.
234///
235/// Content metadata, tagging, ACL, and SSE-KMS headers typically belong on
236/// `create`. SSE-C headers must usually be repeated on `create` and
237/// `upload_part`. Completion checksums and conditions belong on `complete`,
238/// while requester-pays, expected-owner, and abort conditions belong on
239/// `abort`. AWS support varies by bucket type and feature.
240#[derive(Clone, Default, derive_more::Debug, Eq, PartialEq)]
241pub struct ManagedMultipartHeaders {
242    /// Headers sent when creating the upload. Values are signed and repeated values are preserved.
243    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
244    #[debug("{:?}", "<redacted>")]
245    pub create: HeaderMap,
246    /// Headers cloned into every part. Values are signed and repeated values are preserved.
247    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
248    #[debug("{:?}", "<redacted>")]
249    pub upload_part: HeaderMap,
250    /// Completion headers. Values are signed and repeated values are preserved.
251    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
252    #[debug("{:?}", "<redacted>")]
253    pub complete: HeaderMap,
254    /// Abort headers. Values are signed and repeated values are preserved.
255    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
256    #[debug("{:?}", "<redacted>")]
257    pub abort: HeaderMap,
258}
259
260/// Request for a bounded, automatically cleaned-up multipart upload.
261///
262/// In-memory sources retain the caller's complete byte buffer for the duration
263/// of the operation. File sources retain at most the derived in-flight part
264/// byte bound in memory; their immutable snapshot is disk-backed.
265#[derive(derive_more::Debug)]
266pub struct ManagedMultipartUploadRequest {
267    pub(crate) key: ObjectKey,
268    #[debug(
269        "{:?}",
270        match self.source {
271            MultipartUploadSource::Bytes(_) => "bytes",
272            MultipartUploadSource::File(_) => "file",
273        }
274    )]
275    pub(crate) source: MultipartUploadSource,
276    pub(crate) content_type: Option<String>,
277    #[debug("{:?}", self.user_metadata.keys().collect::<Vec<_>>())]
278    pub(crate) user_metadata: BTreeMap<String, String>,
279    pub(crate) options: MultipartOptions,
280    pub(crate) checksum_algorithm: Option<ChecksumAlgorithm>,
281    /// Phase-specific additional headers. Values are signed and repeated values are preserved.
282    /// Generated-name collisions are errors; signing- and transport-owned headers are rejected.
283    pub headers: ManagedMultipartHeaders,
284}
285
286impl ManagedMultipartUploadRequest {
287    /// Constructs a request with a replayable in-memory body.
288    pub fn from_bytes(key: ObjectKey, bytes: impl Into<Bytes>) -> Self {
289        Self {
290            key,
291            source: MultipartUploadSource::Bytes(bytes.into()),
292            content_type: None,
293            user_metadata: BTreeMap::new(),
294            options: MultipartOptions::default(),
295            checksum_algorithm: None,
296            headers: ManagedMultipartHeaders::default(),
297        }
298    }
299
300    /// Constructs a request with a replayable regular-file body.
301    ///
302    /// The file is copied into a private disk-backed snapshot before S3 creates
303    /// the multipart upload, so later changes to the source path have no effect.
304    pub fn from_path(key: ObjectKey, path: impl AsRef<Path>) -> Self {
305        Self {
306            key,
307            source: MultipartUploadSource::File(path.as_ref().to_owned()),
308            content_type: None,
309            user_metadata: BTreeMap::new(),
310            options: MultipartOptions::default(),
311            checksum_algorithm: None,
312            headers: ManagedMultipartHeaders::default(),
313        }
314    }
315
316    /// Replaces the upload's phase-specific additional headers.
317    pub fn with_headers(mut self, headers: ManagedMultipartHeaders) -> Self {
318        self.headers = headers;
319        self
320    }
321
322    /// Sets the object's media type.
323    pub fn with_content_type(mut self, content_type: impl Into<String>) -> Self {
324        self.content_type = Some(content_type.into());
325        self
326    }
327
328    /// Adds user-defined object metadata.
329    pub fn with_metadata(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
330        self.user_metadata.insert(name.into(), value.into());
331        self
332    }
333
334    /// Applies validated resource and time bounds to this upload.
335    pub fn with_options(mut self, options: MultipartOptions) -> Self {
336        self.options = options;
337        self
338    }
339
340    /// Calculates and sends this checksum for every uploaded part.
341    ///
342    /// # Errors
343    ///
344    /// Returns an error when local calculation is unavailable for the selected
345    /// algorithm. SHA-1 may still be supplied through primitive multipart APIs.
346    pub fn with_checksum_algorithm(
347        mut self,
348        algorithm: ChecksumAlgorithm,
349    ) -> Result<Self, crate::operation::ChecksumCalculationError> {
350        if algorithm == ChecksumAlgorithm::Sha1 {
351            return Err(crate::operation::ChecksumCalculationError { algorithm });
352        }
353        self.checksum_algorithm = Some(algorithm);
354        Ok(self)
355    }
356
357    /// Returns this upload's resource and time bounds.
358    pub const fn options(&self) -> MultipartOptions {
359        self.options
360    }
361
362    /// Returns the destination object key.
363    pub const fn key(&self) -> &ObjectKey {
364        &self.key
365    }
366}
367
368impl MultipartUpload {
369    /// Starts tracking a newly created upload.
370    pub fn new(key: ObjectKey, upload_id: UploadId) -> Self {
371        Self {
372            key,
373            upload_id,
374            completed_parts: Vec::new(),
375        }
376    }
377
378    /// Records a completed part while retaining canonical ascending order.
379    pub fn record_part(&mut self, part: CompletedPart) -> Result<(), MultipartError> {
380        match self
381            .completed_parts
382            .binary_search_by_key(&part.part_number(), CompletedPart::part_number)
383        {
384            Ok(_) => Err(MultipartError::DuplicatePart(part.part_number())),
385            Err(index) => {
386                self.completed_parts.insert(index, part);
387                Ok(())
388            }
389        }
390    }
391
392    /// Returns completed parts in validated ascending order.
393    pub fn completed_parts(&self) -> &[CompletedPart] {
394        &self.completed_parts
395    }
396
397    /// Returns the destination object key.
398    pub const fn key(&self) -> &ObjectKey {
399        &self.key
400    }
401
402    /// Returns the non-empty service-issued upload identifier.
403    pub const fn upload_id(&self) -> &UploadId {
404        &self.upload_id
405    }
406}