Skip to main content

s3_wire/operation/multipart/
managed.rs

1use std::collections::BTreeMap;
2use std::fmt;
3use std::path::{Path, PathBuf};
4use std::time::Duration;
5
6use bytes::Bytes;
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, 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    pub user_metadata: BTreeMap<String, String>,
147    /// Checksum algorithm applied to uploaded parts.
148    pub checksum_algorithm: Option<ChecksumAlgorithm>,
149    /// How S3 should derive the completed object's checksum from its parts.
150    pub checksum_type: Option<ChecksumType>,
151}
152
153impl CreateMultipartUploadRequest {
154    /// Constructs a request without optional metadata.
155    pub fn new(key: ObjectKey) -> Self {
156        Self {
157            key,
158            content_type: None,
159            user_metadata: BTreeMap::new(),
160            checksum_algorithm: None,
161            checksum_type: None,
162        }
163    }
164}
165
166/// Result of initiating a multipart upload.
167#[derive(Clone, Debug, Eq, PartialEq)]
168pub struct CreateMultipartUploadOutput {
169    /// Bucket reported by the service.
170    pub bucket: Option<String>,
171    /// Object key reported by the service.
172    pub key: ObjectKey,
173    /// Opaque upload identifier required by later operations.
174    upload_id: UploadId,
175    /// Checksum algorithm selected by the service.
176    pub checksum_algorithm: Option<String>,
177    /// Checksum aggregation selected by S3.
178    pub checksum_type: Option<ChecksumType>,
179    /// Service request identifiers.
180    pub request_ids: RequestIds,
181}
182
183impl CreateMultipartUploadOutput {
184    pub(crate) fn new(
185        bucket: Option<String>,
186        key: ObjectKey,
187        upload_id: UploadId,
188        checksum_algorithm: Option<String>,
189        checksum_type: Option<ChecksumType>,
190    ) -> Self {
191        Self {
192            bucket,
193            key,
194            upload_id,
195            checksum_algorithm,
196            checksum_type,
197            request_ids: RequestIds::default(),
198        }
199    }
200
201    /// Returns the validated upload identifier required by later operations.
202    pub const fn upload_id(&self) -> &UploadId {
203        &self.upload_id
204    }
205}
206
207/// Explicit state needed to resume or clean up an in-progress multipart upload.
208#[derive(Clone, Debug, Eq, PartialEq)]
209pub struct MultipartUpload {
210    key: ObjectKey,
211    upload_id: UploadId,
212    /// Parts that have completed successfully, in ascending order.
213    completed_parts: Vec<CompletedPart>,
214}
215
216pub(crate) enum MultipartUploadSource {
217    Bytes(Bytes),
218    File(PathBuf),
219}
220
221/// Request for a bounded, automatically cleaned-up multipart upload.
222///
223/// In-memory sources retain the caller's complete byte buffer for the duration
224/// of the operation. File sources retain at most the derived in-flight part
225/// byte bound in memory; their immutable snapshot is disk-backed.
226pub struct ManagedMultipartUploadRequest {
227    pub(crate) key: ObjectKey,
228    pub(crate) source: MultipartUploadSource,
229    pub(crate) content_type: Option<String>,
230    pub(crate) user_metadata: BTreeMap<String, String>,
231    pub(crate) options: MultipartOptions,
232    pub(crate) checksum_algorithm: Option<ChecksumAlgorithm>,
233}
234
235impl ManagedMultipartUploadRequest {
236    /// Constructs a request with a replayable in-memory body.
237    pub fn from_bytes(key: ObjectKey, bytes: impl Into<Bytes>) -> Self {
238        Self {
239            key,
240            source: MultipartUploadSource::Bytes(bytes.into()),
241            content_type: None,
242            user_metadata: BTreeMap::new(),
243            options: MultipartOptions::default(),
244            checksum_algorithm: None,
245        }
246    }
247
248    /// Constructs a request with a replayable regular-file body.
249    ///
250    /// The file is copied into a private disk-backed snapshot before S3 creates
251    /// the multipart upload, so later changes to the source path have no effect.
252    pub fn from_path(key: ObjectKey, path: impl AsRef<Path>) -> Self {
253        Self {
254            key,
255            source: MultipartUploadSource::File(path.as_ref().to_owned()),
256            content_type: None,
257            user_metadata: BTreeMap::new(),
258            options: MultipartOptions::default(),
259            checksum_algorithm: None,
260        }
261    }
262
263    /// Sets the object's media type.
264    pub fn with_content_type(mut self, content_type: impl Into<String>) -> Self {
265        self.content_type = Some(content_type.into());
266        self
267    }
268
269    /// Adds user-defined object metadata.
270    pub fn with_metadata(mut self, name: impl Into<String>, value: impl Into<String>) -> Self {
271        self.user_metadata.insert(name.into(), value.into());
272        self
273    }
274
275    /// Applies validated resource and time bounds to this upload.
276    pub fn with_options(mut self, options: MultipartOptions) -> Self {
277        self.options = options;
278        self
279    }
280
281    /// Calculates and sends this checksum for every uploaded part.
282    ///
283    /// # Errors
284    ///
285    /// Returns an error when local calculation is unavailable for the selected
286    /// algorithm. SHA-1 may still be supplied through primitive multipart APIs.
287    pub fn with_checksum_algorithm(
288        mut self,
289        algorithm: ChecksumAlgorithm,
290    ) -> Result<Self, crate::operation::ChecksumCalculationError> {
291        if algorithm == ChecksumAlgorithm::Sha1 {
292            return Err(crate::operation::ChecksumCalculationError { algorithm });
293        }
294        self.checksum_algorithm = Some(algorithm);
295        Ok(self)
296    }
297
298    /// Returns this upload's resource and time bounds.
299    pub const fn options(&self) -> MultipartOptions {
300        self.options
301    }
302
303    /// Returns the destination object key.
304    pub const fn key(&self) -> &ObjectKey {
305        &self.key
306    }
307}
308
309impl fmt::Debug for ManagedMultipartUploadRequest {
310    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
311        formatter
312            .debug_struct("ManagedMultipartUploadRequest")
313            .field("key", &self.key)
314            .field(
315                "source",
316                &match self.source {
317                    MultipartUploadSource::Bytes(_) => "bytes",
318                    MultipartUploadSource::File(_) => "file",
319                },
320            )
321            .field("content_type", &self.content_type)
322            .field("options", &self.options)
323            .field("checksum_algorithm", &self.checksum_algorithm)
324            .field(
325                "user_metadata_names",
326                &self.user_metadata.keys().collect::<Vec<_>>(),
327            )
328            .finish()
329    }
330}
331
332impl MultipartUpload {
333    /// Starts tracking a newly created upload.
334    pub fn new(key: ObjectKey, upload_id: UploadId) -> Self {
335        Self {
336            key,
337            upload_id,
338            completed_parts: Vec::new(),
339        }
340    }
341
342    /// Records a completed part while retaining canonical ascending order.
343    pub fn record_part(&mut self, part: CompletedPart) -> Result<(), MultipartError> {
344        match self
345            .completed_parts
346            .binary_search_by_key(&part.part_number(), CompletedPart::part_number)
347        {
348            Ok(_) => Err(MultipartError::DuplicatePart(part.part_number())),
349            Err(index) => {
350                self.completed_parts.insert(index, part);
351                Ok(())
352            }
353        }
354    }
355
356    /// Returns completed parts in validated ascending order.
357    pub fn completed_parts(&self) -> &[CompletedPart] {
358        &self.completed_parts
359    }
360
361    /// Returns the destination object key.
362    pub const fn key(&self) -> &ObjectKey {
363        &self.key
364    }
365
366    /// Returns the non-empty service-issued upload identifier.
367    pub const fn upload_id(&self) -> &UploadId {
368        &self.upload_id
369    }
370}