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