1use async_trait::async_trait;
7use aws_config::stalled_stream_protection::StalledStreamProtectionConfig;
8use aws_config::{BehaviorVersion, Region};
9use aws_credential_types::Credentials;
10use aws_sdk_s3::config::{RequestChecksumCalculation, ResponseChecksumValidation};
11use aws_sdk_s3::error::ProvideErrorMetadata as _;
12use aws_sdk_s3::Client;
13use std::fmt;
14use tracing::warn;
15
16mod google_cloud_storage;
17
18use google_cloud_storage::{GoogleCloudStorageXml, GoogleUploadSource};
19
20use super::runtime::CloudRuntime;
21use super::s3_common::{
22 apply_prefix, is_not_found_code, normalize_prefix, strip_listed_key_prefix,
23};
24use super::{
25 combine_cleanup_failure, range_header, BlobBody, CloudAccessOutcome, CloudAccessState,
26 CloudHome, CloudHomeError, CloudHomeJoinInfo, ExactSlotStorage, MultipartUpload, RevokeOutcome,
27 UploadControl,
28};
29use coven_foundation::id_provider::{IdRef, UuidProvider};
30use coven_protocol::objects::{ObjectSlot, StorageBackendFailure};
31
32#[derive(Clone)]
34pub struct S3CloudHome {
35 runtime: CloudRuntime,
36 client: Client,
37 sts_client: Option<aws_sdk_sts::Client>,
38 bucket: String,
39 region: String,
40 endpoint: Option<String>,
41 access_key: String,
42 secret_key: String,
43 key_prefix: Option<String>,
44 google_xml: Option<GoogleCloudStorageXml>,
45 clock: coven_foundation::clock::ClockRef,
46 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
47 ids: IdRef,
48}
49
50#[derive(Clone, Debug, PartialEq, Eq)]
51struct S3ExactMetadata {
52 size: u64,
53 sha256: String,
54}
55
56fn sha256_base64(hash: coven_protocol::store_commit::ObjectHash) -> String {
57 use base64::Engine as _;
58 base64::engine::general_purpose::STANDARD.encode(hash.as_bytes())
59}
60
61fn sha256_bytes_base64(bytes: &[u8]) -> String {
62 use base64::Engine as _;
63 use sha2::{Digest, Sha256};
64 base64::engine::general_purpose::STANDARD.encode(Sha256::digest(bytes))
65}
66
67fn create_only_put_failed(
68 error: &aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::put_object::PutObjectError>,
69) -> bool {
70 use aws_sdk_s3::error::ProvideErrorMetadata;
71 let status = match error {
72 aws_sdk_s3::error::SdkError::ServiceError(service) => Some(service.raw().status().as_u16()),
73 _ => None,
74 };
75 status == Some(412)
76 || matches!(
77 error.code(),
78 Some("PreconditionFailed" | "ConditionalRequestConflict")
79 )
80}
81
82fn checksum_put_failed(
83 error: &aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::put_object::PutObjectError>,
84) -> bool {
85 use aws_sdk_s3::error::ProvideErrorMetadata;
86 matches!(
87 error.code(),
88 Some("BadDigest" | "InvalidDigest" | "XAmzContentSHA256Mismatch")
89 )
90}
91
92enum S3CreateOnlyPutError {
93 AlreadyExists(String),
94 ChecksumRejected(CloudHomeError),
95 Other(CloudHomeError),
96}
97
98impl S3CreateOnlyPutError {
99 fn into_cloud_error(self) -> CloudHomeError {
100 match self {
101 Self::AlreadyExists(key) => CloudHomeError::AlreadyExists(key),
102 Self::ChecksumRejected(error) => error,
103 Self::Other(error) => error,
104 }
105 }
106}
107
108impl S3CloudHome {
109 #[allow(clippy::too_many_arguments)]
110 async fn new(
111 runtime: CloudRuntime,
112 bucket: String,
113 region: String,
114 endpoint: Option<String>,
115 access_key: String,
116 secret_key: String,
117 key_prefix: Option<String>,
118 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
119 clock: coven_foundation::clock::ClockRef,
120 ) -> Result<Self, CloudHomeError> {
121 for (name, value) in [
122 ("bucket", bucket.as_str()),
123 ("region", region.as_str()),
124 ("access key", access_key.as_str()),
125 ("secret key", secret_key.as_str()),
126 ] {
127 if value.trim().is_empty() {
128 return Err(CloudHomeError::Configuration(format!(
129 "S3 {name} must not be empty"
130 )));
131 }
132 }
133 let endpoint = endpoint
134 .map(|endpoint| {
135 coven_protocol::provider::canonical_custom_s3_origin(&endpoint).map_err(|error| {
136 CloudHomeError::configuration("validate custom S3 endpoint", error)
137 })
138 })
139 .transpose()?;
140 let google_xml = GoogleCloudStorageXml::for_endpoint(endpoint.as_deref())?;
141 let credentials =
142 Credentials::new(&access_key, &secret_key, None, None, "coven-cloud-home");
143
144 let http_client = smithy_transport_reqwest::ReqwestHttpClient::new();
149
150 let mut builder = aws_config::defaults(BehaviorVersion::latest())
151 .region(Region::new(region.clone()))
152 .credentials_provider(credentials)
153 .http_client(http_client)
154 .stalled_stream_protection(
163 StalledStreamProtectionConfig::enabled()
164 .grace_period(std::time::Duration::from_secs(60))
165 .build(),
166 );
167
168 if let Some(ref ep) = endpoint {
169 builder = builder.endpoint_url(ep);
170 }
171
172 let aws_config = builder.load().await;
173 let s3_builder = aws_sdk_s3::config::Builder::from(&aws_config)
174 .force_path_style(true)
175 .request_checksum_calculation(RequestChecksumCalculation::WhenRequired)
195 .response_checksum_validation(ResponseChecksumValidation::WhenRequired);
196 let s3_config = s3_builder.build();
197 let client = Client::from_conf(s3_config);
198 let sts_client = endpoint
199 .is_none()
200 .then(|| aws_sdk_sts::Client::new(&aws_config));
201
202 Ok(S3CloudHome {
203 runtime,
204 client,
205 sts_client,
206 bucket,
207 region,
208 endpoint,
209 access_key,
210 secret_key,
211 key_prefix: normalize_prefix(key_prefix),
214 google_xml,
215 clock,
216 exact_upload_verification,
217 ids: std::sync::Arc::new(UuidProvider),
218 })
219 }
220
221 fn full_key(&self, key: &str) -> String {
223 apply_prefix(self.key_prefix.as_deref(), key)
224 }
225
226 async fn open_multipart_sink(
227 &self,
228 key: &str,
229 completion: MultipartCompletion,
230 exact_sha256: Option<String>,
231 ) -> Result<Box<S3PartSink>, CloudHomeError> {
232 let full = self.full_key(key);
233 let uses_checksum = exact_sha256.is_some();
234 let upload_id = {
235 let key = key.to_string();
236 let full = full.clone();
237 let client = self.client.clone();
238 let bucket = self.bucket.clone();
239 self.runtime
240 .run_cloud(move || async move {
241 let mut request = client.create_multipart_upload().bucket(&bucket).key(&full);
242 if uses_checksum {
243 request = request
244 .checksum_algorithm(aws_sdk_s3::types::ChecksumAlgorithm::Sha256)
245 .checksum_type(aws_sdk_s3::types::ChecksumType::FullObject);
246 }
247 let create = request.send().await.map_err(|error| {
248 s3_operation_error(format!("multipart create {key}"), error)
249 })?;
250 create
251 .upload_id()
252 .ok_or_else(|| {
253 CloudHomeError::Transport(format!(
254 "multipart create {key}: no upload id returned"
255 ))
256 })
257 .map(str::to_string)
258 })
259 .await?
260 };
261 let (commands, receiver) = tokio::sync::mpsc::channel(1);
262 let owner = S3MultipartOwner {
263 client: self.client.clone(),
264 bucket: self.bucket.clone(),
265 key: full,
266 logical_key: key.to_string(),
267 upload_id,
268 completed: Vec::new(),
269 next_part_number: 1,
270 completion,
271 exact_sha256,
272 };
273 Ok(Box::new(S3PartSink {
274 commands: Some(commands),
275 owner: Some(
276 self.runtime
277 .spawn(move || owner.run(receiver))
278 .map_err(|error| {
279 CloudHomeError::transport("start S3 multipart owner", error)
280 })?,
281 ),
282 }))
283 }
284
285 async fn put_create_only_raw(
286 &self,
287 key: &str,
288 data: Vec<u8>,
289 checksum_sha256: Option<String>,
290 control: UploadControl,
291 ) -> Result<(), S3CreateOnlyPutError> {
292 let full = self.full_key(key);
293 let logical_key = key.to_string();
294 let client = self.client.clone();
295 let bucket = self.bucket.clone();
296 let google_xml = self.google_xml.clone();
297 let endpoint = self.endpoint.clone();
298 let region = self.region.clone();
299 let access_key = self.access_key.clone();
300 let secret_key = self.secret_key.clone();
301 let now = self.clock.now();
302 self.runtime
303 .run(move || async move {
304 if let Some(google_xml) = google_xml {
305 let endpoint = endpoint.ok_or_else(|| {
306 S3CreateOnlyPutError::Other(CloudHomeError::Configuration(
307 "Google Cloud Storage XML endpoint is absent".to_string(),
308 ))
309 })?;
310 let size = data.len() as u64;
311 let payload_hash = hex::encode(
312 coven_protocol::store_commit::ObjectHash::digest(&data).as_bytes(),
313 );
314 return google_xml
315 .create_only(
316 &endpoint,
317 &bucket,
318 ®ion,
319 &access_key,
320 &secret_key,
321 &full,
322 GoogleUploadSource::Bytes(data),
323 size,
324 &payload_hash,
325 now,
326 control,
327 )
328 .await
329 .map_err(|error| match error {
330 CloudHomeError::AlreadyExists(_) => {
331 S3CreateOnlyPutError::AlreadyExists(logical_key)
332 }
333 error => S3CreateOnlyPutError::Other(error),
334 });
335 }
336 let content_length = i64::try_from(data.len()).map_err(|_| {
337 S3CreateOnlyPutError::Other(CloudHomeError::Transport(format!(
338 "object {logical_key} exceeds S3's content-length range"
339 )))
340 })?;
341 let body =
342 reqwest::Body::wrap_stream(control.stream_part(bytes::Bytes::from(data), 0));
343 let body = aws_sdk_s3::primitives::ByteStream::new(
344 aws_sdk_s3::primitives::SdkBody::from_body_1_x(body),
345 );
346 let mut request = client
347 .put_object()
348 .bucket(&bucket)
349 .key(&full)
350 .content_length(content_length)
351 .body(body);
352 if let Some(checksum) = checksum_sha256 {
353 request = request.checksum_sha256(checksum);
354 }
355 let result = request.if_none_match("*").send().await;
356 result.map_err(|error| {
357 if create_only_put_failed(&error) {
358 S3CreateOnlyPutError::AlreadyExists(logical_key.clone())
359 } else if checksum_put_failed(&error) {
360 S3CreateOnlyPutError::ChecksumRejected(CloudHomeError::transport(
361 format!("S3 rejected the SHA-256 request checksum for {logical_key}"),
362 error,
363 ))
364 } else {
365 S3CreateOnlyPutError::Other(put_object_error(&logical_key, error))
366 }
367 })?;
368 Ok(())
369 })
370 .await
371 .map_err(|error| {
372 S3CreateOnlyPutError::Other(CloudHomeError::transport(
373 "run S3 create-only operation",
374 error,
375 ))
376 })?
377 }
378
379 async fn put_google_exact_create_only(
380 &self,
381 key: &str,
382 source: GoogleUploadSource,
383 size: u64,
384 payload_hash: String,
385 control: UploadControl,
386 ) -> Result<(), CloudHomeError> {
387 let full = self.full_key(key);
388 let google_xml = self.google_xml.clone().ok_or_else(|| {
389 CloudHomeError::Configuration(
390 "Google Cloud Storage exact creator is absent".to_string(),
391 )
392 })?;
393 let endpoint = self.endpoint.clone().ok_or_else(|| {
394 CloudHomeError::Configuration("Google Cloud Storage endpoint is absent".to_string())
395 })?;
396 let bucket = self.bucket.clone();
397 let region = self.region.clone();
398 let access_key = self.access_key.clone();
399 let secret_key = self.secret_key.clone();
400 let now = self.clock.now();
401 self.runtime
402 .run(move || async move {
403 google_xml
404 .create_only(
405 &endpoint,
406 &bucket,
407 ®ion,
408 &access_key,
409 &secret_key,
410 &full,
411 source,
412 size,
413 &payload_hash,
414 now,
415 control,
416 )
417 .await
418 })
419 .await
420 .map_err(|error| CloudHomeError::transport("run S3 create-only operation", error))?
421 }
422
423 async fn put_create_only(
424 &self,
425 key: &str,
426 data: Vec<u8>,
427 checksum_sha256: Option<String>,
428 ) -> Result<(), CloudHomeError> {
429 self.put_create_only_raw(
430 key,
431 data,
432 checksum_sha256,
433 UploadControl::running(super::no_progress()),
434 )
435 .await
436 .map_err(S3CreateOnlyPutError::into_cloud_error)
437 }
438
439 async fn append_create_only(
440 &self,
441 key: &str,
442 body: BlobBody,
443 exact_sha256: Option<String>,
444 control: &UploadControl,
445 ) -> Result<(), CloudHomeError> {
446 if body.len() <= self.multipart_threshold() {
447 let data = body.collect().await?;
448 return self
449 .put_create_only_raw(key, data, exact_sha256, control.clone())
450 .await
451 .map_err(S3CreateOnlyPutError::into_cloud_error);
452 }
453 let sink = self
454 .open_multipart_sink(key, MultipartCompletion::CreateOnly, exact_sha256)
455 .await?;
456 MultipartUpload::new(key, body, sink, control).run().await
457 }
458
459 async fn create_at_slot(
460 &self,
461 slot: &ObjectSlot,
462 body: BlobBody,
463 exact_sha256: Option<String>,
464 control: &UploadControl,
465 ) -> Result<(), CloudHomeError> {
466 slot.require_logical_key_for("S3")
467 .map_err(CloudHomeError::from)?;
468 self.append_create_only(slot.logical_key(), body, exact_sha256, control)
469 .await
470 }
471
472 async fn read_exact_to_file(
473 &self,
474 slot: &ObjectSlot,
475 destination: &std::path::Path,
476 progress: super::DownloadProgress,
477 ) -> Result<(), super::CloudFileReadError> {
478 slot.require_logical_key_for("S3")
479 .map_err(|error| super::CloudFileReadError::Source(CloudHomeError::from(error)))?;
480 let full = self.full_key(slot.logical_key());
481 let key = slot.logical_key().to_string();
482 let client = self.client.clone();
483 let bucket = self.bucket.clone();
484 let destination = destination.to_path_buf();
485 self.runtime
486 .run_file_read(move || async move {
487 let response = client
488 .get_object()
489 .bucket(&bucket)
490 .key(&full)
491 .send()
492 .await
493 .map_err(|error| get_object_error(&key, error))?;
494 let stream = futures_util::stream::unfold(
495 (response.body, key),
496 |(mut body, key)| async move {
497 body.next().await.map(|result| {
498 let result = result.map_err(|error| {
499 body_read_error("read appended body", &key, error)
500 });
501 (result, (body, key))
502 })
503 },
504 );
505 super::write_cloud_object_stream(&destination, Box::pin(stream), progress).await?;
506 Ok::<(), super::CloudFileReadError>(())
507 })
508 .await
509 }
510
511 async fn exact_metadata(&self, slot: &ObjectSlot) -> Result<S3ExactMetadata, CloudHomeError> {
512 slot.require_logical_key_for("S3")?;
513 let full = self.full_key(slot.logical_key());
514 let key = slot.logical_key().to_string();
515 let client = self.client.clone();
516 let bucket = self.bucket.clone();
517 self.runtime
518 .run_cloud(move || async move {
519 use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError};
520 let response = client
521 .head_object()
522 .bucket(&bucket)
523 .key(&full)
524 .checksum_mode(aws_sdk_s3::types::ChecksumMode::Enabled)
525 .send()
526 .await
527 .map_err(|error| {
528 let status = match &error {
529 SdkError::ServiceError(service) => {
530 Some(service.raw().status().as_u16())
531 }
532 _ => None,
533 };
534 if is_not_found_code(error.code()) || status == Some(404) {
535 CloudHomeError::NotFound(key.clone())
536 } else {
537 s3_operation_error(format!("head exact S3 object {key}"), error)
538 }
539 })?;
540 let size = response
541 .content_length()
542 .and_then(|size| u64::try_from(size).ok())
543 .ok_or_else(|| {
544 CloudHomeError::Transport(format!(
545 "head exact S3 object {key}: missing content length"
546 ))
547 })?;
548 let sha256 = response
549 .checksum_sha256()
550 .filter(|checksum| !checksum.is_empty())
551 .ok_or_else(|| {
552 CloudHomeError::Transport(format!(
553 "head exact S3 object {key}: missing SHA-256 checksum"
554 ))
555 })?
556 .to_string();
557 Ok(S3ExactMetadata { size, sha256 })
558 })
559 .await
560 }
561
562 async fn verify_exact_upload(
563 &self,
564 upload: &super::ExactUpload<'_>,
565 created_response_was_observed: bool,
566 ) -> Result<(), CloudHomeError> {
567 use coven_foundation::config::ExactUploadVerification;
568
569 if self.google_xml.is_some() {
570 if created_response_was_observed {
571 return Ok(());
572 }
573 let bytes = self.read_at(upload.object().slot()).await?;
574 return upload.verify_stored_bytes(&bytes);
575 }
576
577 match self.exact_upload_verification {
578 ExactUploadVerification::UploadChecksum if created_response_was_observed => Ok(()),
579 ExactUploadVerification::UploadChecksum | ExactUploadVerification::MetadataHash => {
580 let metadata = self.exact_metadata(upload.object().slot()).await?;
581 if metadata.size != upload.object().stored_size()
582 || metadata.sha256 != sha256_base64(upload.object().stored_hash())
583 {
584 return Err(CloudHomeError::SlotCollision(
585 upload.object().slot().logical_key().to_string(),
586 ));
587 }
588 Ok(())
589 }
590 ExactUploadVerification::Readback => {
591 let bytes = self.read_at(upload.object().slot()).await?;
592 upload.verify_stored_bytes(&bytes)
593 }
594 ExactUploadVerification::Unchecked => {
595 super::exact_upload::accept_unchecked_create_response(
596 created_response_was_observed,
597 upload.object(),
598 )
599 }
600 }
601 }
602
603 async fn probe_exact_slots(&self) -> Result<(), CloudHomeError> {
607 use coven_foundation::config::ExactUploadVerification;
608
609 let suffix = self.ids.new_id();
610 let key = format!("__coven_probe__/exact-{suffix}");
611 let bad_key = format!("__coven_probe__/bad-checksum-{suffix}");
612 let bytes = b"coven exact-slot checksum probe".to_vec();
613 let checksum = sha256_bytes_base64(&bytes);
614 let sends_checksum = self.google_xml.is_none()
615 && matches!(
616 self.exact_upload_verification,
617 ExactUploadVerification::UploadChecksum | ExactUploadVerification::MetadataHash
618 );
619 let operation = async {
620 self.put_create_only(
621 &key,
622 bytes.clone(),
623 sends_checksum.then(|| checksum.clone()),
624 )
625 .await?;
626
627 match self
628 .put_create_only(
629 &key,
630 bytes.clone(),
631 sends_checksum.then(|| checksum.clone()),
632 )
633 .await
634 {
635 Err(CloudHomeError::AlreadyExists(_)) => {}
636 Ok(()) => {
637 return Err(CloudHomeError::Configuration(
638 "S3 endpoint did not enforce atomic exact-slot creation".to_string(),
639 ));
640 }
641 Err(error) => return Err(error),
642 }
643
644 if self.read(&key).await? != bytes {
645 return Err(CloudHomeError::Configuration(
646 "S3 exact-slot readback returned different bytes".to_string(),
647 ));
648 }
649 let listed = self.list("__coven_probe__/").await?;
650 if !listed.iter().any(|listed_key| listed_key == &key) {
651 return Err(CloudHomeError::Configuration(
652 "S3 listing did not return the exact-slot probe object".to_string(),
653 ));
654 }
655
656 if self.google_xml.is_none() {
657 match self.exact_upload_verification {
658 ExactUploadVerification::UploadChecksum => {
659 let wrong = sha256_bytes_base64(b"different bytes");
660 match self
661 .put_create_only_raw(
662 &bad_key,
663 bytes.clone(),
664 Some(wrong),
665 UploadControl::running(super::no_progress()),
666 )
667 .await
668 {
669 Err(S3CreateOnlyPutError::ChecksumRejected(_)) => {}
670 Ok(()) => {
671 return Err(CloudHomeError::Configuration(
672 "S3 endpoint accepted an object whose SHA-256 request checksum was wrong"
673 .to_string(),
674 ));
675 }
676 Err(error) => return Err(error.into_cloud_error()),
677 }
678 }
679 ExactUploadVerification::MetadataHash => {
680 let slot = ObjectSlot::logical(key.clone())?;
681 let metadata = self.exact_metadata(&slot).await?;
682 if metadata.size != bytes.len() as u64 || metadata.sha256 != checksum {
683 return Err(CloudHomeError::Configuration(
684 "S3 endpoint did not return the uploaded SHA-256 through HeadObject"
685 .to_string(),
686 ));
687 }
688 }
689 ExactUploadVerification::Readback | ExactUploadVerification::Unchecked => {}
690 }
691 }
692 Ok(())
693 }
694 .await;
695
696 let cleanup_key = self.delete(&key).await;
697 let cleanup_bad = self.delete(&bad_key).await;
698 let cleanup = match (cleanup_key, cleanup_bad) {
699 (Ok(()), Ok(())) => Ok(()),
700 (Err(error), Ok(())) | (Ok(()), Err(error)) => Err(error),
701 (Err(first), Err(second)) => Err(CloudHomeError::CleanupFailed {
702 operation: Box::new(first),
703 cleanup: Box::new(second),
704 }),
705 };
706 match (operation, cleanup) {
707 (Ok(()), cleanup) => cleanup,
708 (Err(operation), Ok(())) => Err(operation),
709 (Err(operation), Err(cleanup)) => Err(combine_cleanup_failure(operation, Err(cleanup))),
710 }
711 }
712
713 #[cfg(test)]
714 async fn provision_test_bucket(&self) {
715 let client = self.client.clone();
716 let bucket = self.bucket.clone();
717 self.runtime
718 .run_cloud(move || async move {
719 client
720 .create_bucket()
721 .bucket(&bucket)
722 .send()
723 .await
724 .map_err(|error| {
725 CloudHomeError::transport(format!("create test S3 bucket {bucket}"), error)
726 })?;
727 Ok(())
728 })
729 .await
730 .expect("create test bucket");
731 }
732}
733
734#[allow(clippy::too_many_arguments)]
735pub(crate) async fn open_cloud_home(
736 runtime: CloudRuntime,
737 bucket: String,
738 region: String,
739 endpoint: Option<String>,
740 access_key: String,
741 secret_key: String,
742 key_prefix: Option<String>,
743 exact_upload_verification: coven_foundation::config::ExactUploadVerification,
744 clock: coven_foundation::clock::ClockRef,
745) -> Result<S3CloudHome, CloudHomeError> {
746 let home_runtime = runtime.clone();
747 runtime
748 .run_cloud(move || {
749 S3CloudHome::new(
750 home_runtime,
751 bucket,
752 region,
753 endpoint,
754 access_key,
755 secret_key,
756 key_prefix,
757 exact_upload_verification,
758 clock,
759 )
760 })
761 .await
762}
763
764#[derive(Clone, Copy, PartialEq, Eq)]
770enum MultipartCompletion {
771 Mutable,
772 CreateOnly,
773}
774
775struct S3PartSink {
776 commands: Option<tokio::sync::mpsc::Sender<S3MultipartCommand>>,
777 owner: Option<tokio::task::JoinHandle<Result<(), CloudHomeError>>>,
778}
779
780enum S3MultipartCommand {
781 SendPart {
782 part: bytes::Bytes,
783 offset: u64,
784 control: UploadControl,
785 response: tokio::sync::oneshot::Sender<Result<(), CloudHomeError>>,
786 },
787 Abort,
788 Finish,
789}
790
791struct S3MultipartOwner {
792 client: Client,
793 bucket: String,
794 key: String,
796 logical_key: String,
797 upload_id: String,
798 completed: Vec<aws_sdk_s3::types::CompletedPart>,
799 next_part_number: i32,
800 completion: MultipartCompletion,
801 exact_sha256: Option<String>,
802}
803
804impl S3MultipartOwner {
805 async fn run(
806 mut self,
807 mut commands: tokio::sync::mpsc::Receiver<S3MultipartCommand>,
808 ) -> Result<(), CloudHomeError> {
809 while let Some(command) = commands.recv().await {
810 match command {
811 S3MultipartCommand::SendPart {
812 part,
813 offset,
814 control,
815 response,
816 } => {
817 let result = self.send_part(part, offset, control).await;
818 if response.send(result).is_err() {
819 return self.abort().await;
820 }
821 }
822 S3MultipartCommand::Abort => return self.abort().await,
823 S3MultipartCommand::Finish => return self.finish().await,
824 }
825 }
826 self.abort().await
827 }
828
829 async fn send_part(
830 &mut self,
831 part: bytes::Bytes,
832 offset: u64,
833 control: UploadControl,
834 ) -> Result<(), CloudHomeError> {
835 let part_number = self.next_part_number;
836 self.next_part_number += 1;
837 let part_sha256 = self
838 .exact_sha256
839 .as_ref()
840 .map(|_| sha256_bytes_base64(&part));
841 let content_length = i64::try_from(part.len()).map_err(|_| {
842 CloudHomeError::Transport(format!(
843 "multipart part {part_number} for {} exceeds S3's content-length range",
844 self.key
845 ))
846 })?;
847 let request_body = reqwest::Body::wrap_stream(control.stream_part(part, offset));
848 let body = aws_sdk_s3::primitives::ByteStream::new(
849 aws_sdk_s3::primitives::SdkBody::from_body_1_x(request_body),
850 );
851 let mut request = self
852 .client
853 .upload_part()
854 .bucket(&self.bucket)
855 .key(&self.key)
856 .upload_id(&self.upload_id)
857 .part_number(part_number)
858 .content_length(content_length)
859 .body(body);
860 if let Some(checksum) = part_sha256.as_ref() {
861 request = request.checksum_sha256(checksum);
862 }
863 let uploaded = request.send().await.map_err(|error| {
864 s3_operation_error(
865 format!("upload multipart part {part_number} for {}", self.key),
866 error,
867 )
868 })?;
869 let mut completed = aws_sdk_s3::types::CompletedPart::builder()
870 .part_number(part_number)
871 .set_e_tag(uploaded.e_tag().map(str::to_string));
872 if let Some(checksum) = part_sha256 {
873 completed = completed.checksum_sha256(checksum);
874 }
875 self.completed.push(completed.build());
876 Ok(())
877 }
878
879 async fn abort(&mut self) -> Result<(), CloudHomeError> {
880 self.client
881 .abort_multipart_upload()
882 .bucket(&self.bucket)
883 .key(&self.key)
884 .upload_id(&self.upload_id)
885 .send()
886 .await
887 .map_err(|error| s3_operation_error(format!("abort multipart {}", self.key), error))?;
888 Ok(())
889 }
890
891 async fn finish(&mut self) -> Result<(), CloudHomeError> {
892 let completed_upload = aws_sdk_s3::types::CompletedMultipartUpload::builder()
893 .set_parts(Some(std::mem::take(&mut self.completed)))
894 .build();
895 let request = self
896 .client
897 .complete_multipart_upload()
898 .bucket(&self.bucket)
899 .key(&self.key)
900 .upload_id(&self.upload_id)
901 .multipart_upload(completed_upload);
902 let mut request = match self.completion {
903 MultipartCompletion::Mutable => request,
904 MultipartCompletion::CreateOnly => request.if_none_match("*"),
905 };
906 if let Some(checksum) = self.exact_sha256.as_ref() {
907 request = request
908 .checksum_sha256(checksum)
909 .checksum_type(aws_sdk_s3::types::ChecksumType::FullObject);
910 }
911 let operation = request.send().await.map(|_| ()).map_err(|error| {
912 use aws_sdk_s3::error::ProvideErrorMetadata;
913 if self.completion == MultipartCompletion::CreateOnly
914 && matches!(
915 error.code(),
916 Some("PreconditionFailed" | "ConditionalRequestConflict")
917 )
918 {
919 CloudHomeError::AlreadyExists(self.logical_key.clone())
920 } else {
921 s3_operation_error(format!("complete multipart {}", self.key), error)
922 }
923 });
924 match operation {
925 Ok(()) => Ok(()),
926 Err(operation) => {
927 let cleanup = self.abort().await;
928 Err(combine_cleanup_failure(operation, cleanup))
929 }
930 }
931 }
932}
933
934impl Drop for S3PartSink {
935 fn drop(&mut self) {
936 self.commands.take();
937 }
938}
939
940impl S3PartSink {
941 async fn settle(&mut self, command: S3MultipartCommand) -> Result<(), CloudHomeError> {
942 let commands = self.commands.take().ok_or_else(|| {
943 CloudHomeError::Transport("S3 multipart upload is already settled".to_string())
944 })?;
945 let send_result = commands.send(command).await;
946 drop(commands);
947 let owner = self
948 .owner
949 .take()
950 .ok_or_else(|| CloudHomeError::Transport("S3 multipart owner is absent".to_string()))?;
951 let result = owner.await.map_err(|error| {
952 CloudHomeError::transport("S3 multipart owner task failed".to_string(), error)
953 })?;
954 match (send_result, result) {
955 (Ok(()), result) => result,
956 (Err(_), Err(error)) => Err(error),
957 (Err(_), Ok(())) => Err(CloudHomeError::Transport(
958 "S3 multipart owner stopped before receiving its terminal command".to_string(),
959 )),
960 }
961 }
962}
963
964mod provider_identity;
965use provider_identity::*;
966
967#[async_trait]
968impl super::PartSink for S3PartSink {
969 fn part_size(&self) -> usize {
970 MULTIPART_PART_SIZE
971 }
972
973 async fn send_part(
974 &mut self,
975 part: bytes::Bytes,
976 offset: u64,
977 _is_last: bool,
978 control: &UploadControl,
979 ) -> Result<(), CloudHomeError> {
980 let commands = self.commands.as_ref().ok_or_else(|| {
981 CloudHomeError::Transport("S3 multipart upload is already settled".to_string())
982 })?;
983 let (response, result) = tokio::sync::oneshot::channel();
984 commands
985 .send(S3MultipartCommand::SendPart {
986 part,
987 offset,
988 control: control.clone(),
989 response,
990 })
991 .await
992 .map_err(|_| {
993 CloudHomeError::Transport(
994 "S3 multipart owner stopped before part upload".to_string(),
995 )
996 })?;
997 result.await.map_err(|_| {
998 CloudHomeError::Transport("S3 multipart owner stopped during part upload".to_string())
999 })?
1000 }
1001
1002 async fn abort(&mut self) -> Result<(), CloudHomeError> {
1003 if self.commands.is_none() {
1004 return Ok(());
1005 }
1006 self.settle(S3MultipartCommand::Abort).await
1007 }
1008
1009 async fn finish(mut self: Box<Self>) -> Result<(), CloudHomeError> {
1010 self.settle(S3MultipartCommand::Finish).await
1011 }
1012}
1013
1014const MULTIPART_THRESHOLD: usize = 8 * 1024 * 1024;
1018
1019const MULTIPART_PART_SIZE: usize = 8 * 1024 * 1024;
1023
1024fn body_read_error<E>(context: &str, key: &str, err: E) -> CloudHomeError
1025where
1026 E: std::error::Error + Send + Sync + 'static,
1027{
1028 CloudHomeError::transport(format!("{context} for {key}"), err)
1029}
1030
1031fn s3_backend_failure(code: Option<&str>, status: Option<u16>) -> StorageBackendFailure {
1032 fn from_status(status: Option<u16>) -> StorageBackendFailure {
1033 match status {
1034 Some(401) => StorageBackendFailure::Authentication,
1035 Some(403) => StorageBackendFailure::PermissionDenied,
1036 Some(404) => StorageBackendFailure::ContainerNotFound,
1037 Some(429 | 500..=599) | None => StorageBackendFailure::Transport,
1038 Some(_) => StorageBackendFailure::Configuration,
1039 }
1040 }
1041
1042 match code {
1043 Some(
1044 "InvalidAccessKeyId"
1045 | "InvalidAccessKey"
1046 | "InvalidClientTokenId"
1047 | "SignatureDoesNotMatch"
1048 | "IncompleteSignature"
1049 | "MissingAuthenticationToken"
1050 | "UnrecognizedClientException"
1051 | "InvalidToken"
1052 | "ExpiredToken"
1053 | "TokenRefreshRequired",
1054 ) => StorageBackendFailure::Authentication,
1055 Some("AccessDenied" | "AllAccessDisabled" | "AccountProblem") => {
1056 StorageBackendFailure::PermissionDenied
1057 }
1058 Some("NoSuchBucket") => StorageBackendFailure::ContainerNotFound,
1059 Some(
1060 "PermanentRedirect"
1061 | "AuthorizationHeaderMalformed"
1062 | "IncorrectEndpoint"
1063 | "IllegalLocationConstraintException",
1064 ) => StorageBackendFailure::RegionMismatch,
1065 Some("OverQuota" | "QuotaExceeded" | "InsufficientStorage") => {
1066 StorageBackendFailure::QuotaExceeded
1067 }
1068 Some(
1069 "InternalError"
1070 | "RequestTimeout"
1071 | "RequestTimeoutException"
1072 | "ServiceUnavailable"
1073 | "SlowDown"
1074 | "Throttling"
1075 | "ThrottlingException"
1076 | "RequestLimitExceeded",
1077 ) => StorageBackendFailure::Transport,
1078 Some(_) | None => from_status(status),
1079 }
1080}
1081
1082fn s3_operation_error<E>(
1083 operation: impl Into<String>,
1084 error: aws_sdk_s3::error::SdkError<E>,
1085) -> CloudHomeError
1086where
1087 E: aws_sdk_s3::error::ProvideErrorMetadata + std::error::Error + Send + Sync + 'static,
1088{
1089 let status = match &error {
1090 aws_sdk_s3::error::SdkError::ServiceError(service) => Some(service.raw().status().as_u16()),
1091 _ => None,
1092 };
1093 let kind = s3_backend_failure(error.code(), status);
1094 CloudHomeError::backend(kind, operation, S3SdkError(error))
1095}
1096
1097#[derive(Debug)]
1098struct S3SdkError<E>(aws_sdk_s3::error::SdkError<E>);
1099
1100impl<E> fmt::Display for S3SdkError<E>
1101where
1102 E: std::error::Error + 'static,
1103{
1104 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
1105 aws_sdk_s3::error::DisplayErrorContext(&self.0).fmt(formatter)
1106 }
1107}
1108
1109impl<E> std::error::Error for S3SdkError<E>
1110where
1111 E: std::error::Error + Send + Sync + 'static,
1112{
1113 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
1114 Some(&self.0)
1115 }
1116}
1117
1118fn get_object_error(
1125 key: &str,
1126 err: aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::get_object::GetObjectError>,
1127) -> CloudHomeError {
1128 use aws_sdk_s3::error::ProvideErrorMetadata;
1129 match err.code() {
1130 Some("NoSuchKey") => CloudHomeError::NotFound(key.to_string()),
1131 Some(code) => s3_operation_error(
1132 match err.message() {
1133 Some(msg) => format!("get {key}: S3 {code}: {msg}"),
1134 None => format!("get {key}: S3 {code} (no message provided)"),
1135 },
1136 err,
1137 ),
1138 None => s3_operation_error(format!("get {key}"), err),
1141 }
1142}
1143
1144fn put_object_error(
1156 key: &str,
1157 err: aws_sdk_s3::error::SdkError<aws_sdk_s3::operation::put_object::PutObjectError>,
1158) -> CloudHomeError {
1159 use aws_sdk_s3::error::ProvideErrorMetadata;
1160 match err.code() {
1161 Some(code) => s3_operation_error(
1162 match err.message() {
1163 Some(msg) => format!("put {key}: S3 {code}: {msg}"),
1164 None => format!("put {key}: S3 {code} (no message provided)"),
1165 },
1166 err,
1167 ),
1168 None => s3_operation_error(format!("put {key}"), err),
1169 }
1170}
1171
1172#[async_trait]
1173impl CloudHome for S3CloudHome {
1174 async fn probe(&self) -> Result<(), CloudHomeError> {
1175 self.probe_exact_slots().await
1176 }
1177
1178 async fn put_object(&self, key: &str, data: Vec<u8>) -> Result<(), CloudHomeError> {
1179 let full = self.full_key(key);
1180 let key = key.to_string();
1181 let client = self.client.clone();
1182 let bucket = self.bucket.clone();
1183 self.runtime
1184 .run_cloud(move || async move {
1185 client
1186 .put_object()
1187 .bucket(&bucket)
1188 .key(&full)
1189 .body(data.into())
1190 .send()
1191 .await
1192 .map_err(|e| put_object_error(&key, e))?;
1193 Ok(())
1194 })
1195 .await
1196 }
1197
1198 async fn open_multipart<'a>(
1199 &'a self,
1200 key: &str,
1201 _total_len: u64,
1202 ) -> Result<super::BoxPartSink<'a>, CloudHomeError> {
1203 Ok(self
1204 .open_multipart_sink(key, MultipartCompletion::Mutable, None)
1205 .await?)
1206 }
1207
1208 fn multipart_threshold(&self) -> u64 {
1209 MULTIPART_THRESHOLD as u64
1210 }
1211
1212 async fn read(&self, key: &str) -> Result<Vec<u8>, CloudHomeError> {
1213 let full = self.full_key(key);
1214 let key = key.to_string();
1215 let client = self.client.clone();
1216 let bucket = self.bucket.clone();
1217 self.runtime
1220 .run_cloud(move || async move {
1221 let resp = client
1222 .get_object()
1223 .bucket(&bucket)
1224 .key(&full)
1225 .send()
1226 .await
1227 .map_err(|e| get_object_error(&key, e))?;
1228
1229 let bytes = resp
1230 .body
1231 .collect()
1232 .await
1233 .map_err(|e| body_read_error("read body", &key, e))?
1234 .into_bytes()
1235 .to_vec();
1236
1237 Ok(bytes)
1238 })
1239 .await
1240 }
1241
1242 async fn read_range(&self, key: &str, start: u64, end: u64) -> Result<Vec<u8>, CloudHomeError> {
1243 let full = self.full_key(key);
1244 let range = range_header(start, end);
1245 let key = key.to_string();
1246 let client = self.client.clone();
1247 let bucket = self.bucket.clone();
1248 self.runtime
1249 .run_cloud(move || async move {
1250 let resp = client
1251 .get_object()
1252 .bucket(&bucket)
1253 .key(&full)
1254 .range(range)
1255 .send()
1256 .await
1257 .map_err(|e| get_object_error(&key, e))?;
1258
1259 let bytes = resp
1260 .body
1261 .collect()
1262 .await
1263 .map_err(|e| body_read_error("read range body", &key, e))?
1264 .into_bytes()
1265 .to_vec();
1266
1267 let expected = end - start;
1274 if bytes.len() as u64 != expected {
1275 return Err(CloudHomeError::Transport(format!(
1276 "read range {key}: expected {expected} bytes for range {start}..{end}, \
1277 got {} — the provider likely ignored Range and returned the whole object",
1278 bytes.len()
1279 )));
1280 }
1281
1282 Ok(bytes)
1283 })
1284 .await
1285 }
1286
1287 async fn list(&self, prefix: &str) -> Result<Vec<String>, CloudHomeError> {
1288 let full_prefix = self.full_key(prefix);
1289 let key_prefix = self.key_prefix.clone();
1290 let prefix = prefix.to_string();
1291 let client = self.client.clone();
1292 let bucket = self.bucket.clone();
1293 self.runtime
1296 .run_cloud(move || async move {
1297 let mut keys = Vec::new();
1298 let mut continuation_token: Option<String> = None;
1299
1300 loop {
1301 let mut req = client
1302 .list_objects_v2()
1303 .bucket(&bucket)
1304 .prefix(&full_prefix);
1305
1306 if let Some(token) = continuation_token.take() {
1307 req = req.continuation_token(token);
1308 }
1309
1310 let resp = req
1311 .send()
1312 .await
1313 .map_err(|error| s3_operation_error(format!("list {prefix}"), error))?;
1314
1315 for obj in resp.contents() {
1316 let Some(key) = obj.key() else {
1317 warn!("list {prefix}: S3 returned an object with no key; skipping it");
1318 continue;
1319 };
1320 let Some(stripped) =
1321 strip_listed_key_prefix(key_prefix.as_deref(), &full_prefix, key)
1322 else {
1323 warn!(
1324 "list {prefix}: key {key} is outside the configured S3 prefix {:?}; \
1325 skipping it",
1326 key_prefix
1327 );
1328 continue;
1329 };
1330 keys.push(stripped.to_string());
1331 }
1332
1333 if resp.is_truncated() == Some(true) {
1334 let token = resp.next_continuation_token().ok_or_else(|| {
1335 CloudHomeError::Transport(format!(
1336 "list {prefix}: S3 truncated but returned no continuation token"
1337 ))
1338 })?;
1339 continuation_token = Some(token.to_string());
1340 } else {
1341 break;
1342 }
1343 }
1344
1345 Ok(keys)
1346 })
1347 .await
1348 }
1349
1350 async fn delete(&self, key: &str) -> Result<(), CloudHomeError> {
1351 let full = self.full_key(key);
1352 let key = key.to_string();
1353 let client = self.client.clone();
1354 let bucket = self.bucket.clone();
1355 self.runtime
1356 .run_cloud(move || async move {
1357 use aws_sdk_s3::error::ProvideErrorMetadata;
1358 if let Err(e) = client
1359 .delete_object()
1360 .bucket(&bucket)
1361 .key(&full)
1362 .send()
1363 .await
1364 {
1365 if !is_not_found_code(e.code()) {
1371 return Err(s3_operation_error(format!("delete {key}"), e));
1372 }
1373 }
1374 Ok(())
1375 })
1376 .await
1377 }
1378
1379 async fn exists(&self, key: &str) -> Result<bool, CloudHomeError> {
1380 let full = self.full_key(key);
1381 let key = key.to_string();
1382 let client = self.client.clone();
1383 let bucket = self.bucket.clone();
1384 self.runtime
1385 .run_cloud(move || async move {
1386 use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError};
1387 match client.head_object().bucket(&bucket).key(&full).send().await {
1388 Ok(_) => Ok(true),
1389 Err(e) => {
1392 let status = match &e {
1393 SdkError::ServiceError(svc) => Some(svc.raw().status().as_u16()),
1394 _ => None,
1395 };
1396 if is_not_found_code(e.code()) || status == Some(404) {
1397 Ok(false)
1398 } else {
1399 Err(s3_operation_error(format!("head {key}"), e))
1400 }
1401 }
1402 }
1403 })
1404 .await
1405 }
1406
1407 async fn set_access(
1408 &self,
1409 desired: CloudAccessState,
1410 ) -> Result<CloudAccessOutcome, CloudHomeError> {
1411 Ok(match desired {
1412 CloudAccessState::Present { .. } => {
1413 CloudAccessOutcome::Present(CloudHomeJoinInfo::S3 {
1414 bucket: self.bucket.clone(),
1415 region: self.region.clone(),
1416 endpoint: self.endpoint.clone(),
1417 access_key: self.access_key.clone(),
1418 secret_key: self.secret_key.clone(),
1419 key_prefix: self.key_prefix.clone(),
1420 })
1421 }
1422 CloudAccessState::Absent { .. } => {
1423 CloudAccessOutcome::Absent(RevokeOutcome::Unsupported)
1424 }
1425 })
1426 }
1427}
1428
1429mod exact;
1430
1431#[cfg(test)]
1432mod tests;