s3_wire/operation/multipart/
managed.rs1use 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#[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 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 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 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 pub const fn part_size(self) -> u64 {
97 self.part_size
98 }
99
100 pub const fn concurrency(self) -> usize {
102 self.concurrency
103 }
104
105 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 pub const fn transfer_timeout(self) -> Duration {
112 self.transfer_timeout
113 }
114
115 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#[derive(Clone, Debug, Eq, PartialEq)]
140pub struct CreateMultipartUploadRequest {
141 pub key: ObjectKey,
143 pub content_type: Option<String>,
145 pub user_metadata: BTreeMap<String, String>,
147 pub checksum_algorithm: Option<ChecksumAlgorithm>,
149 pub checksum_type: Option<ChecksumType>,
151}
152
153impl CreateMultipartUploadRequest {
154 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#[derive(Clone, Debug, Eq, PartialEq)]
168pub struct CreateMultipartUploadOutput {
169 pub bucket: Option<String>,
171 pub key: ObjectKey,
173 upload_id: UploadId,
175 pub checksum_algorithm: Option<String>,
177 pub checksum_type: Option<ChecksumType>,
179 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 pub const fn upload_id(&self) -> &UploadId {
203 &self.upload_id
204 }
205}
206
207#[derive(Clone, Debug, Eq, PartialEq)]
209pub struct MultipartUpload {
210 key: ObjectKey,
211 upload_id: UploadId,
212 completed_parts: Vec<CompletedPart>,
214}
215
216pub(crate) enum MultipartUploadSource {
217 Bytes(Bytes),
218 File(PathBuf),
219}
220
221pub 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 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 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 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 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 pub fn with_options(mut self, options: MultipartOptions) -> Self {
277 self.options = options;
278 self
279 }
280
281 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 pub const fn options(&self) -> MultipartOptions {
300 self.options
301 }
302
303 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 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 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 pub fn completed_parts(&self) -> &[CompletedPart] {
358 &self.completed_parts
359 }
360
361 pub const fn key(&self) -> &ObjectKey {
363 &self.key
364 }
365
366 pub const fn upload_id(&self) -> &UploadId {
368 &self.upload_id
369 }
370}