1use std::num::NonZeroU64;
112use std::sync::Arc;
113use std::sync::atomic::Ordering;
114use std::time::{Duration, SystemTime};
115
116use base64::Engine as _;
117use bytes::Bytes;
118use futures_util::StreamExt;
119use objectstore_types::metadata::Metadata;
120use objectstore_types::range::ByteRange;
121use objectstore_types::resumable::UploadProgress;
122use objectstore_types::time::Timestamp;
123use sentry::{Hub, SentryFutureExt};
124use serde::{Deserialize, Serialize};
125
126use crate::backend::changelog::{Change, ChangeGuard, ChangeLog, ChangeManager, ChangePhase};
127use crate::backend::common::{
128 Backend, DeleteResponse, ExpiryTarget, ExpiryUpdate, GetResponse, HighVolumeBackend,
129 MetadataResponse, MultipartUploadBackend, PutResponse, SetExpiryResponse, TieredGet,
130 TieredMetadata, TieredUpdate, TieredWrite, Tombstone,
131};
132use crate::backend::{HighVolumeStorageConfig, MultipartUploadStorageConfig};
133use crate::error::{Error, ErrorKind, Result, ResultExt as _};
134use crate::id::ObjectId;
135use crate::multipart::{
136 AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
137 ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
138};
139use crate::resumable::{BackendToken, Session};
140use crate::stream::{ClientStream, SizedPeek, counting_stream};
141
142const BACKEND_SIZE_THRESHOLD: usize = 1024 * 1024; const RESUMABLE_UPLOAD_TTL: Duration = Duration::from_hours(5 * 24);
147
148const MULTIPART_COMPLETE_CLEANUP_DELAY: Duration = Duration::from_hours(24);
153
154fn new_long_term_revision(id: &ObjectId) -> ObjectId {
160 ObjectId {
161 context: id.context.clone(),
162 key: format!("{}/{}", id.key, uuid::Uuid::now_v7()),
163 }
164}
165
166#[derive(Debug, Clone, Deserialize, Serialize)]
187pub struct TieredStorageConfig {
188 pub high_volume: HighVolumeStorageConfig,
193 pub long_term: MultipartUploadStorageConfig,
197}
198
199#[derive(Debug)]
250pub struct TieredStorage {
251 inner: Arc<ChangeManager>,
252}
253
254impl TieredStorage {
255 pub fn new(
257 high_volume: Box<dyn HighVolumeBackend>,
258 long_term: Box<dyn MultipartUploadBackend>,
259 changelog: Box<dyn ChangeLog>,
260 ) -> Self {
261 let inner = ChangeManager::new(high_volume, long_term, changelog);
262 let hub = Hub::new_from_top(Hub::current());
263 tokio::spawn(inner.clone().recover().bind_hub(hub));
266 Self { inner }
267 }
268
269 async fn record_change(&self, change: Change) -> Result<ChangeGuard> {
271 self.inner.clone().record(change).await
272 }
273
274 async fn record_assembling(&self, change: Change) -> Result<ChangeGuard> {
277 self.inner.clone().record_assembling(change).await
278 }
279
280 async fn check_upload_marker(&self, revision: &ObjectId) -> Result<()> {
282 if !self
283 .inner
284 .high_volume
285 .has_upload_marker(revision, Timestamp::now())
286 .await?
287 {
288 return Err(ErrorKind::UploadSessionGone.into());
289 }
290 Ok(())
291 }
292
293 async fn delete_upload_marker(&self, revision: &ObjectId) -> Result<()> {
295 if !self
296 .inner
297 .high_volume
298 .delete_upload_marker(revision, Timestamp::now())
299 .await?
300 {
301 return Err(ErrorKind::UploadSessionGone.into());
302 }
303 Ok(())
304 }
305
306 fn backend_type(&self, choice: &BackendChoice) -> &'static str {
308 match choice {
309 BackendChoice::HighVolume => self.inner.high_volume.name(),
310 BackendChoice::LongTerm => self.inner.long_term.name(),
311 }
312 }
313
314 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
319 async fn put_high_volume(
320 &self,
321 id: &ObjectId,
322 metadata: &Metadata,
323 payload: Bytes,
324 access_time: Timestamp,
325 ) -> Result<()> {
326 let tombstone_opt = self
327 .inner
328 .high_volume
329 .put_non_tombstone(id, metadata, payload.clone(), access_time)
330 .await?;
331
332 let Some(Tombstone { target, .. }) = tombstone_opt else {
333 return Ok(());
335 };
336
337 let mut guard = self
339 .record_change(Change {
340 id: id.clone(),
341 new: None,
342 old: Some(target.clone()),
343 cleanup_after: None,
344 })
345 .await?;
346
347 let write = TieredWrite::Object(metadata.clone(), payload);
348 guard.advance(ChangePhase::Written);
349
350 let written = self
351 .inner
352 .high_volume
353 .compare_and_write(id, Some(&target), write, access_time)
354 .await?;
355
356 guard.advance(ChangePhase::compare_and_write(written));
358
359 Ok(())
360 }
361
362 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
367 async fn put_long_term(
368 &self,
369 id: &ObjectId,
370 metadata: &Metadata,
371 stream: ClientStream,
372 access_time: Timestamp,
373 ) -> Result<()> {
374 let current = match self
376 .inner
377 .high_volume
378 .get_tiered_metadata(id, access_time)
379 .await?
380 {
381 TieredMetadata::Tombstone(t) => Some(t.target),
382 _ => None,
383 };
384
385 let new = new_long_term_revision(id);
387 let mut guard = self
388 .record_change(Change {
389 id: id.clone(),
390 new: Some(new.clone()),
391 old: current.clone(),
392 cleanup_after: None,
393 })
394 .await?;
395
396 self.inner
397 .long_term
398 .put_object(&new, metadata, stream, access_time)
399 .await?;
400 guard.advance(ChangePhase::Written);
401
402 let tombstone = Tombstone {
404 target: new.clone(),
405 time_expires: metadata.time_expires,
406 };
407 let written = self
408 .inner
409 .high_volume
410 .compare_and_write(
411 id,
412 current.as_ref(),
413 TieredWrite::Tombstone(tombstone),
414 access_time,
415 )
416 .await?;
417
418 guard.advance(ChangePhase::compare_and_write(written));
420
421 Ok(())
422 }
423}
424
425type LongTermBackendToken = BackendToken;
426
427#[derive(Debug, Serialize, Deserialize)]
428struct TieredResumableToken {
429 inner: LongTermBackendToken,
430 revision: String,
431 time_expires: Option<Timestamp>,
432}
433
434impl TieredResumableToken {
435 fn decode(token: &BackendToken) -> Result<Self> {
436 serde_json::from_str(token).map_err(|_| ErrorKind::UnknownUploadSession.into())
437 }
438
439 fn into_inner_session(self, session: &Session) -> Session {
441 Session {
442 object_id: ObjectId {
443 context: session.object_id.context.clone(),
444 key: self.revision,
445 },
446 backend_token: self.inner,
447 ..session.clone()
448 }
449 }
450}
451
452#[async_trait::async_trait]
453impl Backend for TieredStorage {
454 fn name(&self) -> &'static str {
455 "tiered"
456 }
457
458 fn upload_granularity(&self) -> u64 {
459 self.inner.long_term.upload_granularity()
460 }
461
462 fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
463 Ok(self)
464 }
465
466 #[tracing::instrument(level = "debug", fields(?id, upload_length), skip_all)]
467 async fn create_upload_session(
468 &self,
469 id: &ObjectId,
470 metadata: &Metadata,
471 upload_length: NonZeroU64,
472 ) -> Result<Option<BackendToken>> {
473 if upload_length.get() <= BACKEND_SIZE_THRESHOLD as u64 {
474 return Ok(None);
475 }
476
477 let revision = new_long_term_revision(id);
478 let Some(inner) = self
479 .inner
480 .long_term
481 .create_upload_session(&revision, metadata, upload_length)
482 .await?
483 else {
484 return Ok(None);
485 };
486
487 if let Err(error) = self
488 .inner
489 .high_volume
490 .create_upload_marker(&revision, Timestamp::now() + RESUMABLE_UPLOAD_TTL)
491 .await
492 {
493 let session = Session {
495 object_id: revision,
496 upload_length,
497 backend_token: inner,
498 };
499 if let Err(cleanup_error) = self.inner.long_term.cancel_upload(&session).await {
500 objectstore_log::warn!(
501 !!&cleanup_error,
502 "Failed to cancel upload after marker creation failed"
503 );
504 }
505 return Err(error);
506 }
507 let token = TieredResumableToken {
508 revision: revision.key,
509 inner,
510 time_expires: metadata.time_expires,
511 };
512 Ok(Some(serde_json::to_string(&token).context(
513 ErrorKind::Internal,
514 "encoding tiered resumable session",
515 )?))
516 }
517
518 #[tracing::instrument(level = "debug", fields(?session, offset, content_length), skip_all)]
519 async fn put_chunk(
520 &self,
521 session: &Session,
522 offset: u64,
523 content_length: u64,
524 stream: ClientStream,
525 ) -> Result<UploadProgress> {
526 let tiered = TieredResumableToken::decode(&session.backend_token)?;
527 let time_expires = tiered.time_expires;
528 let inner_session = tiered.into_inner_session(session);
529 let id = &session.object_id;
530 let revision = &inner_session.object_id;
531 let end = offset
532 .checked_add(content_length)
533 .filter(|end| *end <= session.upload_length.get())
534 .ok_or(ErrorKind::ChunkExceedsUploadLength {
535 offset,
536 content_length,
537 upload_length: session.upload_length.get(),
538 })?;
539
540 let granularity = self.upload_granularity();
541 if content_length > 0 && content_length < granularity && end != session.upload_length.get()
542 {
543 return Err(ErrorKind::ChunkTooSmall {
544 chunk_length: content_length,
545 upload_granularity: granularity,
546 }
547 .into());
548 }
549
550 if end != session.upload_length.get() {
552 let progress = self
553 .inner
554 .long_term
555 .put_chunk(&inner_session, offset, content_length, stream)
556 .await?;
557 return match progress {
558 UploadProgress::Incomplete { .. } => Ok(progress),
559 UploadProgress::Complete => Err(ErrorKind::UploadSessionGone.into()),
560 };
561 }
562
563 let (progress, current) = tokio::join!(
565 self.inner.long_term.upload_offset(&inner_session),
566 self.inner
567 .high_volume
568 .get_tiered_metadata(id, Timestamp::now()),
569 );
570 match progress? {
571 UploadProgress::Complete => return Err(ErrorKind::UploadSessionGone.into()),
572 UploadProgress::Incomplete { offset: actual } if actual < offset => {
573 return Err(ErrorKind::UploadOffsetMismatch { offset: actual }.into());
574 }
575 UploadProgress::Incomplete { .. } => {}
576 }
577 let current = match current? {
578 TieredMetadata::Tombstone(t) if t.target == *revision => {
579 return Ok(UploadProgress::Complete);
580 }
581 TieredMetadata::Tombstone(t) => Some(t.target),
582 _ => None,
583 };
584
585 self.delete_upload_marker(revision).await?;
587
588 let mut guard = self
590 .record_change(Change {
591 id: id.clone(),
592 new: Some(revision.clone()),
593 old: current.clone(),
594 cleanup_after: None,
595 })
596 .await?;
597
598 let progress = self
599 .inner
600 .long_term
601 .put_chunk(&inner_session, offset, content_length, stream)
602 .await?;
603 if progress != UploadProgress::Complete {
604 return Err(ErrorKind::UploadSessionGone.into());
605 }
606 guard.advance(ChangePhase::Written);
607
608 let written = self
609 .inner
610 .high_volume
611 .compare_and_write(
612 id,
613 current.as_ref(),
614 TieredWrite::Tombstone(Tombstone {
615 target: revision.clone(),
616 time_expires,
617 }),
618 Timestamp::now(),
619 )
620 .await?;
621 guard.advance(ChangePhase::compare_and_write(written));
622 Ok(UploadProgress::Complete)
623 }
624
625 #[tracing::instrument(level = "debug", fields(?session), skip_all)]
626 async fn upload_offset(&self, session: &Session) -> Result<UploadProgress> {
627 let tiered = TieredResumableToken::decode(&session.backend_token)?;
628 let inner_session = tiered.into_inner_session(session);
629 self.check_upload_marker(&inner_session.object_id).await?;
630 match self.inner.long_term.upload_offset(&inner_session).await? {
631 UploadProgress::Incomplete { offset } => Ok(UploadProgress::Incomplete { offset }),
632 UploadProgress::Complete => Err(ErrorKind::UploadSessionGone.into()),
633 }
634 }
635
636 #[tracing::instrument(level = "debug", fields(?session), skip_all)]
637 async fn cancel_upload(&self, session: &Session) -> Result<()> {
638 let tiered = TieredResumableToken::decode(&session.backend_token)?;
639 let inner_session = tiered.into_inner_session(session);
640 self.delete_upload_marker(&inner_session.object_id).await?;
641 if let Err(error) = self.inner.long_term.cancel_upload(&inner_session).await {
642 objectstore_log::warn!(!!&error, "Failed to cancel upload after deleting marker");
643 }
644 Ok(())
645 }
646
647 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
648 async fn put_object(
649 &self,
650 id: &ObjectId,
651 metadata: &Metadata,
652 stream: ClientStream,
653 access_time: Timestamp,
654 ) -> Result<PutResponse> {
655 let timer = objectstore_metrics::timer!("put.latency", usecase = id.usecase().to_owned());
656 if metadata.origin.is_none() {
657 objectstore_metrics::count!("put.origin_missing", usecase = id.usecase().to_owned());
658 }
659
660 let peeked = SizedPeek::new(stream, BACKEND_SIZE_THRESHOLD).await?;
661 objectstore_metrics::record!(
662 "put.first_chunk.latency" = timer.elapsed(),
663 usecase = id.usecase().to_owned(),
664 complete = if peeked.is_exhausted() { "yes" } else { "no" },
665 );
666
667 let (backend_choice, stored_size) = if peeked.is_exhausted() {
668 let payload = peeked.into_bytes().await?;
669 let payload_len = payload.len() as u64;
670 self.put_high_volume(id, metadata, payload, access_time)
671 .await?;
672 (BackendChoice::HighVolume, payload_len)
673 } else {
674 let (stored_size, stream) = counting_stream(peeked.into_stream());
675 self.put_long_term(id, metadata, stream.boxed(), access_time)
676 .await?;
677 (BackendChoice::LongTerm, stored_size.load(Ordering::Acquire))
678 };
679
680 let backend_ty = self.backend_type(&backend_choice);
681 timer
682 .tag("backend_choice", backend_choice.as_str())
683 .tag("backend_type", backend_ty)
684 .record();
685 objectstore_metrics::record!(
686 "put.size" = stored_size,
687 usecase = id.usecase().to_owned(),
688 backend_choice = backend_choice.as_str(),
689 backend_type = backend_ty,
690 upload_type = "direct",
691 );
692
693 Ok(())
694 }
695
696 #[tracing::instrument(level = "debug", skip(self))]
697 async fn get_object(
698 &self,
699 id: &ObjectId,
700 access_time: Timestamp,
701 range: Option<ByteRange>,
702 ) -> Result<GetResponse> {
703 let timer = objectstore_metrics::timer!(
704 "get.latency.pre-response",
705 usecase = id.usecase().to_owned(),
706 );
707
708 let hv_result = self
709 .inner
710 .high_volume
711 .get_tiered_object(id, access_time, range)
712 .await?;
713 let (result, backend_choice) = match hv_result {
714 TieredGet::NotFound => (None, BackendChoice::HighVolume),
715 TieredGet::Object(metadata, content_range, stream) => (
716 Some((metadata, content_range, stream)),
717 BackendChoice::HighVolume,
718 ),
719 TieredGet::Tombstone(tombstone) => (
720 self.inner
721 .long_term
722 .get_object(&tombstone.target, access_time, range)
723 .await?
724 .map(|(meta, range, stream)| (align_expiry(meta, &tombstone), range, stream)),
725 BackendChoice::LongTerm,
726 ),
727 };
728
729 let backend_type = self.backend_type(&backend_choice);
730 timer
731 .tag("backend_choice", backend_choice.as_str())
732 .tag("backend_type", backend_type)
733 .record();
734
735 if let Some((ref metadata, ref content_range, _)) = result {
736 let size = content_range.map(|cr| cr.len() as usize).or(metadata.size);
737 if let Some(size) = size {
738 objectstore_metrics::record!(
739 "get.size" = size,
740 usecase = id.usecase().to_owned(),
741 backend_choice = backend_choice.as_str(),
742 backend_type = backend_type,
743 );
744 }
745 }
746
747 Ok(result)
748 }
749
750 #[tracing::instrument(level = "debug", skip(self))]
751 async fn get_metadata(
752 &self,
753 id: &ObjectId,
754 access_time: Timestamp,
755 ) -> Result<MetadataResponse> {
756 let timer = objectstore_metrics::timer!("head.latency", usecase = id.usecase().to_owned());
757
758 let hv_result = self
759 .inner
760 .high_volume
761 .get_tiered_metadata(id, access_time)
762 .await?;
763 let (result, backend_choice) = match hv_result {
764 TieredMetadata::NotFound => (None, BackendChoice::HighVolume),
765 TieredMetadata::Object(metadata) => (Some(metadata), BackendChoice::HighVolume),
766 TieredMetadata::Tombstone(tombstone) => (
767 self.inner
768 .long_term
769 .get_metadata(&tombstone.target, access_time)
770 .await?
771 .map(|metadata| align_expiry(metadata, &tombstone)),
772 BackendChoice::LongTerm,
773 ),
774 };
775
776 timer
777 .tag("backend_choice", backend_choice.as_str())
778 .tag("backend_type", self.backend_type(&backend_choice))
779 .record();
780
781 Ok(result)
782 }
783
784 async fn set_expiry(
785 &self,
786 id: &ObjectId,
787 target: ExpiryUpdate,
788 access_time: Timestamp,
789 ) -> Result<SetExpiryResponse> {
790 match self
791 .inner
792 .high_volume
793 .get_tiered_metadata(id, access_time)
794 .await?
795 {
796 TieredMetadata::NotFound => Ok(SetExpiryResponse::NotFound),
797 TieredMetadata::Object(_) => {
798 self.inner
799 .high_volume
800 .compare_and_update(id, None, TieredUpdate::SetExpiry(target), access_time)
801 .await
802 }
803 TieredMetadata::Tombstone(tombstone) => {
804 let deadline = match self
807 .inner
808 .long_term
809 .set_expiry(&tombstone.target, target, access_time)
810 .await?
811 {
812 SetExpiryResponse::Satisfied(deadline) => deadline,
813 outcome @ (SetExpiryResponse::NotFound | SetExpiryResponse::Rejected) => {
814 return Ok(outcome);
815 }
816 };
817
818 let update = TieredUpdate::SetExpiry(ExpiryTarget::At(deadline).into());
820
821 self.inner
827 .high_volume
828 .compare_and_update(id, Some(&tombstone.target), update, access_time)
829 .await
830 }
831 }
832 }
833
834 #[tracing::instrument(level = "debug", skip(self))]
835 async fn delete_object(&self, id: &ObjectId, access_time: Timestamp) -> Result<DeleteResponse> {
836 let timer =
837 objectstore_metrics::timer!("delete.latency", usecase = id.usecase().to_owned());
838
839 let mut backend_choice = BackendChoice::HighVolume;
840
841 if let Some(tombstone) = self
842 .inner
843 .high_volume
844 .delete_non_tombstone(id, access_time)
845 .await?
846 {
847 backend_choice = BackendChoice::LongTerm;
848
849 let mut guard = self
850 .record_change(Change {
851 id: id.clone(),
852 new: None,
853 old: Some(tombstone.target.clone()),
854 cleanup_after: None,
855 })
856 .await?;
857 guard.advance(ChangePhase::Written);
858
859 let deleted = self
861 .inner
862 .high_volume
863 .compare_and_write(
864 id,
865 Some(&tombstone.target),
866 TieredWrite::Delete,
867 access_time,
868 )
869 .await?;
870
871 guard.advance(ChangePhase::compare_and_write(deleted));
873 }
874
875 timer
876 .tag("backend_choice", backend_choice.as_str())
877 .tag("backend_type", self.backend_type(&backend_choice))
878 .record();
879
880 Ok(())
881 }
882
883 async fn join(&self) {
884 self.inner.tracker.close();
885 tokio::join!(
886 self.inner.high_volume.join(),
887 self.inner.long_term.join(),
888 self.inner.tracker.wait()
889 );
890 }
891}
892
893#[derive(Debug)]
894enum BackendChoice {
895 HighVolume,
896 LongTerm,
897}
898
899impl BackendChoice {
900 fn as_str(&self) -> &'static str {
901 match self {
902 BackendChoice::HighVolume => "high-volume",
903 BackendChoice::LongTerm => "long-term",
904 }
905 }
906}
907
908impl std::fmt::Display for BackendChoice {
909 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
910 f.write_str(self.as_str())
911 }
912}
913
914fn effective_expiry(
918 redirect_expiry: Option<Timestamp>,
919 blob_expiry: Option<Timestamp>,
920) -> Option<Timestamp> {
921 match (redirect_expiry, blob_expiry) {
922 (Some(redirect), Some(blob)) => Some(redirect.min(blob)),
923 (Some(expiry), None) | (None, Some(expiry)) => Some(expiry),
924 (None, None) => None,
925 }
926}
927
928fn align_expiry(mut metadata: Metadata, tombstone: &Tombstone) -> Metadata {
933 metadata.time_expires = effective_expiry(tombstone.time_expires, metadata.time_expires);
934 metadata
935}
936
937#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
939struct TieredUploadId {
940 revision: String,
941 upload_id: UploadId,
942}
943
944impl TryInto<UploadId> for TieredUploadId {
945 type Error = Error;
946
947 fn try_into(self) -> Result<UploadId, Self::Error> {
948 let json =
949 serde_json::to_vec(&self).context(ErrorKind::Internal, "encoding tiered upload ID")?;
950 Ok(UploadId::new(
951 base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(json),
952 )?)
953 }
954}
955
956impl TryFrom<&UploadId> for TieredUploadId {
957 type Error = Error;
958
959 fn try_from(value: &UploadId) -> Result<Self, Self::Error> {
960 let json = base64::engine::general_purpose::URL_SAFE_NO_PAD
961 .decode(value.as_bytes())
962 .kind(ErrorKind::InvalidUploadId)?;
963 serde_json::from_slice(&json).kind(ErrorKind::InvalidUploadId)
964 }
965}
966
967#[async_trait::async_trait]
968impl MultipartUploadBackend for TieredStorage {
969 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
970 async fn initiate_multipart(
971 &self,
972 id: &ObjectId,
973 metadata: &Metadata,
974 ) -> Result<InitiateMultipartResponse> {
975 let timer = objectstore_metrics::timer!(
976 "multipart.initiate.latency",
977 usecase = id.usecase().to_owned(),
978 );
979 let physical = new_long_term_revision(id);
980
981 let upload_id = self
982 .inner
983 .long_term
984 .initiate_multipart(&physical, metadata)
985 .await?;
986
987 let id = TieredUploadId {
988 revision: physical.key,
989 upload_id,
990 };
991 let id = id.try_into()?;
992
993 timer.record();
994 Ok(id)
995 }
996
997 #[tracing::instrument(level = "debug", fields(?id, part_number, content_length), skip_all)]
998 async fn upload_part(
999 &self,
1000 id: &ObjectId,
1001 upload_id: &UploadId,
1002 part_number: PartNumber,
1003 content_length: u64,
1004 content_md5: Option<&str>,
1005 body: ClientStream,
1006 ) -> Result<UploadPartResponse> {
1007 let timer = objectstore_metrics::timer!(
1008 "multipart.upload_part.latency",
1009 usecase = id.usecase().to_owned(),
1010 );
1011 let tiered: TieredUploadId = upload_id.try_into()?;
1012
1013 let physical = ObjectId {
1014 context: id.context.clone(),
1015 key: tiered.revision,
1016 };
1017
1018 let etag = self
1019 .inner
1020 .long_term
1021 .upload_part(
1022 &physical,
1023 &tiered.upload_id,
1024 part_number,
1025 content_length,
1026 content_md5,
1027 body,
1028 )
1029 .await?;
1030
1031 timer.record();
1032 objectstore_metrics::record!(
1033 "multipart.upload_part.size" = content_length,
1034 usecase = id.usecase().to_owned(),
1035 );
1036
1037 Ok(etag)
1038 }
1039
1040 #[tracing::instrument(level = "debug", skip(self, upload_id))]
1041 async fn list_parts(
1042 &self,
1043 id: &ObjectId,
1044 upload_id: &UploadId,
1045 max_parts: Option<u32>,
1046 part_number_marker: Option<PartNumber>,
1047 ) -> Result<ListPartsResponse> {
1048 let timer = objectstore_metrics::timer!(
1049 "multipart.list_parts.latency",
1050 usecase = id.usecase().to_owned(),
1051 );
1052 let tiered: TieredUploadId = upload_id.try_into()?;
1053
1054 let physical = ObjectId {
1055 context: id.context.clone(),
1056 key: tiered.revision,
1057 };
1058
1059 let response = self
1060 .inner
1061 .long_term
1062 .list_parts(&physical, &tiered.upload_id, max_parts, part_number_marker)
1063 .await?;
1064
1065 timer.record();
1066 Ok(response)
1067 }
1068
1069 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1070 async fn abort_multipart(
1071 &self,
1072 id: &ObjectId,
1073 upload_id: &UploadId,
1074 ) -> Result<AbortMultipartResponse> {
1075 let timer = objectstore_metrics::timer!(
1076 "multipart.abort.latency",
1077 usecase = id.usecase().to_owned(),
1078 );
1079 let tiered: TieredUploadId = upload_id.try_into()?;
1080
1081 let physical = ObjectId {
1082 context: id.context.clone(),
1083 key: tiered.revision,
1084 };
1085
1086 let () = self
1087 .inner
1088 .long_term
1089 .abort_multipart(&physical, &tiered.upload_id)
1090 .await?;
1091
1092 timer.record();
1093 Ok(())
1094 }
1095
1096 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1097 async fn complete_multipart(
1098 &self,
1099 id: &ObjectId,
1100 upload_id: &UploadId,
1101 parts: Vec<CompletedPart>,
1102 access_time: Timestamp,
1103 ) -> Result<CompleteMultipartResponse> {
1104 let timer = objectstore_metrics::timer!(
1105 "multipart.complete.latency",
1106 usecase = id.usecase().to_owned(),
1107 );
1108 let part_count = parts.len();
1109 let tiered: TieredUploadId = upload_id.try_into()?;
1110
1111 let physical = ObjectId {
1112 context: id.context.clone(),
1113 key: tiered.revision,
1114 };
1115
1116 let current = match self
1118 .inner
1119 .high_volume
1120 .get_tiered_metadata(id, access_time)
1121 .await?
1122 {
1123 TieredMetadata::Tombstone(t) if t.target == physical => {
1125 timer.record();
1126 return Ok(None);
1127 }
1128 TieredMetadata::Tombstone(t) => Some(t.target),
1129 _ => None,
1130 };
1131
1132 let mut guard = self
1135 .record_assembling(Change {
1136 id: id.clone(),
1137 new: Some(physical.clone()),
1138 old: current.clone(),
1139 cleanup_after: Some(SystemTime::now() + MULTIPART_COMPLETE_CLEANUP_DELAY),
1140 })
1141 .await?;
1142
1143 let maybe_complete_multipart_err = match self
1145 .inner
1146 .long_term
1147 .complete_multipart(&physical, &tiered.upload_id, parts, access_time)
1148 .await
1149 {
1150 Ok(error) => {
1153 if error.is_some() {
1154 return Ok(error);
1155 }
1156 None
1157 }
1158 Err(err) => Some(err),
1164 };
1165
1166 let metadata = self
1173 .inner
1174 .long_term
1175 .get_metadata(&physical, access_time)
1176 .await;
1177
1178 let metadata = match (metadata, maybe_complete_multipart_err) {
1179 (Ok(Some(metadata)), _) => metadata,
1181 (Ok(None), Some(err)) => return Err(err),
1183 (Ok(None), None) => {
1186 objectstore_log::error!(
1187 id = ?id,
1188 upload_id = ?upload_id,
1189 physical = ?physical,
1190 "complete_multipart call succeeded on long_term backend, but subsequent get_metadata found no object"
1191 );
1192 return Err(Error::new(
1193 ErrorKind::BackendFailure,
1194 "tiered multipart object missing from long-term storage",
1195 ));
1196 }
1197 (Err(get_metadata_err), maybe_complete_multipart_err) => {
1199 return Err(maybe_complete_multipart_err.unwrap_or(get_metadata_err));
1203 }
1204 };
1205
1206 let tombstone = Tombstone {
1208 target: physical.clone(),
1209 time_expires: metadata.time_expires,
1210 };
1211 let written = self
1212 .inner
1213 .high_volume
1214 .compare_and_write(
1215 id,
1216 current.as_ref(),
1217 TieredWrite::Tombstone(tombstone),
1218 access_time,
1219 )
1220 .await?;
1221
1222 guard.advance(ChangePhase::compare_and_write(written));
1224
1225 timer.record();
1226 objectstore_metrics::record!(
1227 "multipart.complete.part_count" = part_count as u64,
1228 usecase = id.usecase().to_owned(),
1229 );
1230 if let Some(size) = metadata.size {
1231 objectstore_metrics::record!(
1232 "put.size" = size as u64,
1233 usecase = id.usecase().to_owned(),
1234 backend_choice = BackendChoice::LongTerm.as_str(),
1235 backend_type = self.backend_type(&BackendChoice::LongTerm),
1236 upload_type = "multipart",
1237 );
1238 }
1239
1240 Ok(None)
1241 }
1242}
1243
1244#[cfg(test)]
1245mod tests {
1246 use std::num::{NonZeroU32, NonZeroU64};
1247 use std::sync::Mutex as StdMutex;
1248
1249 use futures::lock::Mutex;
1250 use objectstore_types::metadata::{ExpirationPolicy, Metadata};
1251 use objectstore_types::scope::{Scope, Scopes};
1252
1253 use super::*;
1254 use crate::backend::bigtable::{BigTableBackend, BigTableConfig};
1255 use crate::backend::changelog::{InMemoryChangeLog, NoopChangeLog};
1256 use crate::backend::gcs::{GcsBackend, GcsConfig};
1257 use crate::backend::in_memory::InMemoryBackend;
1258 use crate::backend::testing::{Hooks, TestBackend};
1259 use crate::change_stream::ChangeStreamFactory;
1260 use crate::error::Error;
1261 use crate::id::ObjectContext;
1262 use crate::stream::{self, ClientStream};
1263
1264 fn make_context() -> ObjectContext {
1265 ObjectContext {
1266 usecase: "testing".into(),
1267 scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1268 }
1269 }
1270
1271 fn make_id(key: &str) -> ObjectId {
1272 ObjectId::new(make_context(), key.into())
1273 }
1274
1275 fn make_tiered_storage() -> (
1276 TieredStorage,
1277 InMemoryBackend,
1278 InMemoryBackend,
1279 InMemoryChangeLog,
1280 ) {
1281 let hv = InMemoryBackend::new("in-memory-hv");
1282 let lt = InMemoryBackend::new("in-memory-lt");
1283 let changelog = InMemoryChangeLog::default();
1284 let storage = TieredStorage::new(
1285 Box::new(hv.clone()),
1286 Box::new(lt.clone()),
1287 Box::new(changelog.clone()),
1288 );
1289 (storage, hv, lt, changelog)
1290 }
1291
1292 async fn resumable_token(
1293 storage: &TieredStorage,
1294 id: &ObjectId,
1295 metadata: &Metadata,
1296 length: u64,
1297 ) -> Session {
1298 let backend_token = storage
1299 .create_upload_session(id, metadata, NonZeroU64::new(length).unwrap())
1300 .await
1301 .unwrap()
1302 .unwrap();
1303 Session {
1304 object_id: id.clone(),
1305 upload_length: NonZeroU64::new(length).unwrap(),
1306 backend_token,
1307 }
1308 }
1309
1310 fn upload_revision(session: &Session) -> ObjectId {
1311 TieredResumableToken::decode(&session.backend_token)
1312 .unwrap()
1313 .into_inner_session(session)
1314 .object_id
1315 }
1316
1317 #[tokio::test]
1318 async fn resumable_inmemory() -> anyhow::Result<()> {
1319 let (storage, _, _, _) = make_tiered_storage();
1320 let id = make_id("tiered-resumable-inmemory");
1321 let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1];
1322 let token =
1323 resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
1324
1325 let revision = upload_revision(&token);
1326 assert!(
1327 storage
1328 .inner
1329 .high_volume
1330 .has_upload_marker(&revision, Timestamp::now())
1331 .await?
1332 );
1333 assert_eq!(
1334 storage.upload_offset(&token).await?,
1335 UploadProgress::Incomplete { offset: 0 }
1336 );
1337 let split = 256 * 1024;
1338 assert_eq!(
1339 storage
1340 .put_chunk(
1341 &token,
1342 0,
1343 split as u64,
1344 stream::single(payload[..split].to_vec())
1345 )
1346 .await?,
1347 UploadProgress::Incomplete {
1348 offset: split as u64
1349 }
1350 );
1351 assert!(
1352 storage
1353 .inner
1354 .high_volume
1355 .has_upload_marker(&revision, Timestamp::now())
1356 .await?
1357 );
1358 assert_eq!(
1359 storage.upload_offset(&token).await?,
1360 UploadProgress::Incomplete {
1361 offset: split as u64
1362 }
1363 );
1364 assert!(storage.get_metadata(&id, Timestamp::now()).await?.is_none());
1365
1366 assert_eq!(
1368 storage
1369 .put_chunk(
1370 &token,
1371 split as u64,
1372 (payload.len() - split) as u64,
1373 stream::single(payload[split..].to_vec())
1374 )
1375 .await?,
1376 UploadProgress::Complete
1377 );
1378 assert!(
1379 !storage
1380 .inner
1381 .high_volume
1382 .has_upload_marker(&revision, Timestamp::now())
1383 .await?
1384 );
1385 assert_eq!(
1386 storage.upload_offset(&token).await.unwrap_err().kind(),
1387 ErrorKind::UploadSessionGone
1388 );
1389 let (_, _, body) = storage
1390 .get_object(&id, Timestamp::now(), None)
1391 .await?
1392 .unwrap();
1393 assert_eq!(stream::read_to_vec(body).await?, payload);
1394 Ok(())
1395 }
1396
1397 #[tokio::test]
1398 async fn resumable_bigtable_and_gcs() -> anyhow::Result<()> {
1399 let streams = ChangeStreamFactory::default();
1400 let lt = GcsBackend::new(
1401 GcsConfig {
1402 endpoint: Some("http://localhost:8087".into()),
1403 bucket: "test-bucket".into(),
1404 cogs: None,
1405 },
1406 &streams,
1407 )
1408 .await?;
1409 let hv = BigTableBackend::new(
1410 BigTableConfig {
1411 endpoint: Some("localhost:8086".into()),
1412 project_id: "testing".into(),
1413 instance_name: "objectstore".into(),
1414 table_name: "objectstore".into(),
1415 connections: None,
1416 rpc_timeout: Duration::from_secs(2),
1417 cogs: None,
1418 },
1419 &streams,
1420 )
1421 .await?;
1422 let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(NoopChangeLog));
1423 let id = make_id(&format!("tiered-resumable-{}", uuid::Uuid::now_v7()));
1424 let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1];
1425 let token =
1426 resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
1427
1428 let error = storage
1429 .put_chunk(&token, 0, 1, stream::single("a"))
1430 .await
1431 .unwrap_err();
1432 assert_eq!(
1433 error.kind(),
1434 ErrorKind::ChunkTooSmall {
1435 chunk_length: 1,
1436 upload_granularity: 256 * 1024,
1437 }
1438 );
1439
1440 let revision = upload_revision(&token);
1441 assert!(
1442 storage
1443 .inner
1444 .high_volume
1445 .has_upload_marker(&revision, Timestamp::now())
1446 .await?
1447 );
1448 assert_eq!(
1449 storage.upload_offset(&token).await?,
1450 UploadProgress::Incomplete { offset: 0 }
1451 );
1452 let split = 256 * 1024;
1453 assert_eq!(
1454 storage
1455 .put_chunk(
1456 &token,
1457 0,
1458 split as u64,
1459 stream::single(payload[..split].to_vec())
1460 )
1461 .await?,
1462 UploadProgress::Incomplete {
1463 offset: split as u64
1464 }
1465 );
1466 assert!(
1467 storage
1468 .inner
1469 .high_volume
1470 .has_upload_marker(&revision, Timestamp::now())
1471 .await?
1472 );
1473 assert_eq!(
1474 storage.upload_offset(&token).await?,
1475 UploadProgress::Incomplete {
1476 offset: split as u64
1477 }
1478 );
1479 assert!(storage.get_metadata(&id, Timestamp::now()).await?.is_none());
1480
1481 assert_eq!(
1483 storage
1484 .put_chunk(
1485 &token,
1486 split as u64,
1487 (payload.len() - split) as u64,
1488 stream::single(payload[split..].to_vec())
1489 )
1490 .await?,
1491 UploadProgress::Complete
1492 );
1493 assert!(
1494 !storage
1495 .inner
1496 .high_volume
1497 .has_upload_marker(&revision, Timestamp::now())
1498 .await?
1499 );
1500 assert_eq!(
1501 storage.upload_offset(&token).await.unwrap_err().kind(),
1502 ErrorKind::UploadSessionGone
1503 );
1504 let (_, _, body) = storage
1505 .get_object(&id, Timestamp::now(), None)
1506 .await?
1507 .unwrap();
1508 assert_eq!(stream::read_to_vec(body).await?, payload);
1509 Ok(())
1510 }
1511
1512 #[tokio::test]
1513 async fn resumable_invalid_chunks() {
1514 let (storage, _, _, _) = make_tiered_storage();
1515 let id = make_id("resumable-invalid");
1516 let length = BACKEND_SIZE_THRESHOLD as u64 + 1;
1517 let token = resumable_token(&storage, &id, &Metadata::default(), length).await;
1518
1519 assert_eq!(
1521 storage
1522 .put_chunk(&token, u64::MAX, 1, stream::single("x"))
1523 .await
1524 .unwrap_err()
1525 .kind(),
1526 ErrorKind::ChunkExceedsUploadLength {
1527 offset: u64::MAX,
1528 content_length: 1,
1529 upload_length: length
1530 }
1531 );
1532 assert_eq!(
1533 storage
1534 .put_chunk(&token, length - 1, 1, stream::single("x"))
1535 .await
1536 .unwrap_err()
1537 .kind(),
1538 ErrorKind::UploadOffsetMismatch { offset: 0 }
1539 );
1540 }
1541
1542 #[tokio::test]
1543 async fn resumable_cancel() {
1544 let (storage, hv, _, _) = make_tiered_storage();
1545 let id = make_id("resumable-invalid");
1546 let token = resumable_token(
1547 &storage,
1548 &id,
1549 &Metadata::default(),
1550 BACKEND_SIZE_THRESHOLD as u64 + 1,
1551 )
1552 .await;
1553
1554 let revision = upload_revision(&token);
1555 assert!(
1556 hv.has_upload_marker(&revision, Timestamp::now())
1557 .await
1558 .unwrap()
1559 );
1560 storage.cancel_upload(&token).await.unwrap();
1561 assert!(
1562 !hv.has_upload_marker(&revision, Timestamp::now())
1563 .await
1564 .unwrap()
1565 );
1566 assert_eq!(
1567 storage.cancel_upload(&token).await.unwrap_err().kind(),
1568 ErrorKind::UploadSessionGone
1569 );
1570 assert_eq!(
1571 storage.upload_offset(&token).await.unwrap_err().kind(),
1572 ErrorKind::UploadSessionGone
1573 );
1574 assert!(!hv.contains(&id));
1575 }
1576
1577 #[derive(Clone, Debug)]
1578 struct ExpiryHook {
1579 label: &'static str,
1580 events: Arc<StdMutex<Vec<&'static str>>>,
1581 reject: bool,
1582 }
1583
1584 type ExpiryTestStorage = (
1585 TieredStorage,
1586 TestBackend<ExpiryHook>,
1587 TestBackend<ExpiryHook>,
1588 Arc<StdMutex<Vec<&'static str>>>,
1589 );
1590
1591 #[async_trait::async_trait]
1592 impl Hooks for ExpiryHook {
1593 async fn set_expiry(
1594 &self,
1595 inner: &InMemoryBackend,
1596 id: &ObjectId,
1597 target: ExpiryUpdate,
1598 access_time: Timestamp,
1599 ) -> Result<SetExpiryResponse> {
1600 self.events.lock().unwrap().push(self.label);
1601 if self.reject {
1602 Ok(SetExpiryResponse::Rejected)
1603 } else {
1604 inner.set_expiry(id, target, access_time).await
1605 }
1606 }
1607
1608 async fn compare_and_update(
1609 &self,
1610 inner: &InMemoryBackend,
1611 id: &ObjectId,
1612 current: Option<&ObjectId>,
1613 update: TieredUpdate,
1614 access_time: Timestamp,
1615 ) -> Result<SetExpiryResponse> {
1616 self.events.lock().unwrap().push(self.label);
1617 if self.reject {
1618 Ok(SetExpiryResponse::Rejected)
1619 } else {
1620 inner
1621 .compare_and_update(id, current, update, access_time)
1622 .await
1623 }
1624 }
1625 }
1626
1627 fn tiered_with_expiry_hooks(hv_reject: bool, lt_reject: bool) -> ExpiryTestStorage {
1628 let events = Arc::new(StdMutex::new(Vec::new()));
1629 let hv = TestBackend::new(ExpiryHook {
1630 label: "hv",
1631 events: Arc::clone(&events),
1632 reject: hv_reject,
1633 });
1634 let lt = TestBackend::new(ExpiryHook {
1635 label: "lt",
1636 events: Arc::clone(&events),
1637 reject: lt_reject,
1638 });
1639 let storage = TieredStorage::new(
1640 Box::new(hv.clone()),
1641 Box::new(lt.clone()),
1642 Box::new(NoopChangeLog),
1643 );
1644 (storage, hv, lt, events)
1645 }
1646
1647 async fn seed_redirect(
1648 hv: &InMemoryBackend,
1649 lt: &InMemoryBackend,
1650 id: &ObjectId,
1651 target: &ObjectId,
1652 expiry: Timestamp,
1653 ) {
1654 let metadata = Metadata {
1655 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_mins(10)),
1656 time_created: expiry.checked_sub(Duration::from_mins(10)),
1657 time_expires: Some(expiry),
1658 ..Default::default()
1659 };
1660 lt.put_object(
1661 target,
1662 &metadata,
1663 stream::single("payload"),
1664 Timestamp::now(),
1665 )
1666 .await
1667 .unwrap();
1668 hv.compare_and_write(
1669 id,
1670 None,
1671 TieredWrite::Tombstone(Tombstone {
1672 target: target.clone(),
1673 time_expires: metadata.time_expires,
1674 }),
1675 Timestamp::now(),
1676 )
1677 .await
1678 .unwrap();
1679 }
1680
1681 #[tokio::test]
1682 async fn set_expiry() {
1683 let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, false);
1684 let id = make_id("tiered-expiry-order");
1685 let target = new_long_term_revision(&id);
1686 let old_expiry = Timestamp::now() + Duration::from_mins(10);
1687 seed_redirect(&hv.inner, <.inner, &id, &target, old_expiry).await;
1688
1689 let created = old_expiry - Duration::from_mins(10);
1690 let requested = created + Duration::from_hours(1);
1691 let blob_expiry = requested + Duration::from_hours(1);
1692 lt.inner
1693 .put_object(
1694 &target,
1695 &Metadata {
1696 expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
1697 time_created: Some(created),
1698 time_expires: Some(blob_expiry),
1699 ..Default::default()
1700 },
1701 stream::single("payload"),
1702 Timestamp::now(),
1703 )
1704 .await
1705 .unwrap();
1706 let update = ExpiryUpdate {
1707 target: ExpiryTarget::FromCreation(Duration::from_hours(1)),
1708 max: Some(Duration::ZERO),
1709 };
1710 assert_eq!(
1711 storage
1712 .set_expiry(&id, update, created)
1713 .await
1714 .unwrap_err()
1715 .kind(),
1716 ErrorKind::InvalidMetadata,
1717 );
1718 assert_eq!(events.lock().unwrap().as_slice(), &["lt"]);
1719 events.lock().unwrap().clear();
1720 assert_eq!(
1721 hv.inner.get(&id).expect_tombstone().time_expires,
1722 Some(old_expiry)
1723 );
1724 assert_eq!(
1725 storage
1726 .set_expiry(
1727 &id,
1728 ExpiryUpdate {
1729 max: Some(Duration::from_hours(1)),
1730 ..update
1731 },
1732 Timestamp::now(),
1733 )
1734 .await
1735 .unwrap(),
1736 SetExpiryResponse::Satisfied(requested)
1737 );
1738 assert_eq!(events.lock().unwrap().as_slice(), &["lt", "hv"]);
1739 assert_eq!(
1740 lt.inner.get(&target).expect_object().0.time_expires,
1741 Some(blob_expiry)
1742 );
1743 assert_eq!(
1744 hv.inner.get(&id).expect_tombstone().time_expires,
1745 Some(requested)
1746 );
1747 }
1748
1749 #[derive(Clone, Debug)]
1750 struct ReplaceInlineOnLookup {
1751 replacement: Metadata,
1752 replaced: Arc<StdMutex<bool>>,
1753 }
1754
1755 #[async_trait::async_trait]
1756 impl Hooks for ReplaceInlineOnLookup {
1757 async fn get_tiered_metadata(
1758 &self,
1759 inner: &InMemoryBackend,
1760 id: &ObjectId,
1761 access_time: Timestamp,
1762 ) -> Result<TieredMetadata> {
1763 let observed = inner.get_tiered_metadata(id, access_time).await?;
1764 let should_replace = {
1765 let mut replaced = self.replaced.lock().unwrap();
1766 !std::mem::replace(&mut *replaced, true)
1767 };
1768 if should_replace {
1769 inner
1770 .put_object(
1771 id,
1772 &self.replacement,
1773 stream::single("replacement"),
1774 access_time,
1775 )
1776 .await?;
1777 }
1778 Ok(observed)
1779 }
1780 }
1781
1782 #[tokio::test]
1783 async fn inline_creation_target_uses_update_snapshot() {
1784 let access_time = Timestamp::from_unix_secs(1_700_000_000).unwrap();
1785 let old_created = access_time - Duration::from_hours(2);
1786 let new_created = access_time - Duration::from_mins(30);
1787 let old_expiry = access_time + Duration::from_mins(10);
1788 let resolved = new_created + Duration::from_hours(2);
1789 let replacement = Metadata {
1790 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1791 time_created: Some(new_created),
1792 time_expires: Some(old_expiry),
1793 ..Default::default()
1794 };
1795 let hv = TestBackend::new(ReplaceInlineOnLookup {
1796 replacement,
1797 replaced: Arc::new(StdMutex::new(false)),
1798 });
1799 let storage = TieredStorage::new(
1800 Box::new(hv.clone()),
1801 Box::new(InMemoryBackend::new("lt")),
1802 Box::new(NoopChangeLog),
1803 );
1804 let id = make_id("inline-update-snapshot");
1805 hv.inner
1806 .put_object(
1807 &id,
1808 &Metadata {
1809 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1810 time_created: Some(old_created),
1811 time_expires: Some(old_expiry),
1812 ..Default::default()
1813 },
1814 stream::single("original"),
1815 access_time,
1816 )
1817 .await
1818 .unwrap();
1819
1820 assert_eq!(
1821 storage
1822 .set_expiry(
1823 &id,
1824 ExpiryTarget::FromCreation(Duration::from_hours(2)).into(),
1825 access_time,
1826 )
1827 .await
1828 .unwrap(),
1829 SetExpiryResponse::Satisfied(resolved)
1830 );
1831 let (metadata, payload) = hv.inner.get(&id).expect_object();
1832 assert_eq!(
1833 metadata.expiration_policy,
1834 ExpirationPolicy::TimeToLive(Duration::from_hours(2))
1835 );
1836 assert_eq!(metadata.time_created, Some(new_created));
1837 assert_eq!(metadata.time_expires, Some(resolved));
1838 assert_eq!(payload, Bytes::from_static(b"replacement"));
1839 }
1840
1841 #[tokio::test]
1842 async fn expiry_hv_conflict() {
1843 let (storage, hv, lt, events) = tiered_with_expiry_hooks(true, false);
1844 let id = make_id("tiered-hv-failure");
1845 let target = new_long_term_revision(&id);
1846 let old_expiry = Timestamp::now() + Duration::from_mins(10);
1847 seed_redirect(&hv.inner, <.inner, &id, &target, old_expiry).await;
1848
1849 let requested = old_expiry + Duration::from_mins(50);
1850 assert_eq!(
1851 storage
1852 .set_expiry(
1853 &id,
1854 ExpiryTarget::FromCreation(Duration::from_hours(1)).into(),
1855 Timestamp::now(),
1856 )
1857 .await
1858 .unwrap(),
1859 SetExpiryResponse::Rejected
1860 );
1861 assert_eq!(events.lock().unwrap().as_slice(), &["lt", "hv"]);
1862 assert_eq!(
1863 hv.inner.get(&id).expect_tombstone().time_expires,
1864 Some(old_expiry)
1865 );
1866 assert_eq!(
1867 lt.inner.get(&target).expect_object().0.time_expires,
1868 Some(requested)
1869 );
1870 assert_eq!(
1871 storage
1872 .get_metadata(&id, Timestamp::now())
1873 .await
1874 .unwrap()
1875 .unwrap()
1876 .time_expires,
1877 Some(old_expiry)
1878 );
1879 assert_eq!(
1880 lt.inner.get(&target).expect_object().0.expiration_policy,
1881 ExpirationPolicy::TimeToLive(Duration::from_hours(1))
1882 );
1883 }
1884
1885 #[tokio::test]
1886 async fn expiry_lt_conflict() {
1887 let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, true);
1888 let id = make_id("tiered-lt-conflict");
1889 let target = new_long_term_revision(&id);
1890 seed_redirect(
1891 &hv.inner,
1892 <.inner,
1893 &id,
1894 &target,
1895 Timestamp::now() + Duration::from_mins(10),
1896 )
1897 .await;
1898
1899 assert_eq!(
1900 storage
1901 .set_expiry(
1902 &id,
1903 ExpiryTarget::At(Timestamp::now() + Duration::from_hours(1)).into(),
1904 Timestamp::now()
1905 )
1906 .await
1907 .unwrap(),
1908 SetExpiryResponse::Rejected
1909 );
1910 assert_eq!(events.lock().unwrap().as_slice(), &["lt"]);
1911 }
1912
1913 #[tokio::test]
1914 async fn expiry_not_found() {
1915 for missing_blob in [false, true] {
1916 let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, false);
1917 let id = make_id("tiered-expiry-missing");
1918 let target = new_long_term_revision(&id);
1919 let access_time = Timestamp::now();
1920 let deadline = access_time + Duration::from_hours(1);
1921 if missing_blob {
1922 seed_redirect(&hv.inner, <.inner, &id, &target, deadline).await;
1923 lt.inner.delete_object(&target, access_time).await.unwrap();
1924 }
1925
1926 assert_eq!(
1927 storage
1928 .set_expiry(&id, ExpiryTarget::At(deadline).into(), access_time)
1929 .await
1930 .unwrap(),
1931 SetExpiryResponse::NotFound
1932 );
1933 let expected: &[&str] = if missing_blob { &["lt"] } else { &[] };
1934 assert_eq!(events.lock().unwrap().as_slice(), expected);
1935 }
1936 }
1937
1938 #[test]
1941 fn revision_id_preserves_context() {
1942 let id = make_id("my-key");
1943 let revised = new_long_term_revision(&id);
1944 assert_eq!(revised.context, id.context);
1945 assert!(
1946 revised.key.starts_with("my-key/"),
1947 "revised key should have /<uuid> suffix, got: {}",
1948 revised.key
1949 );
1950 }
1951
1952 #[test]
1953 fn revision_id_roundtrips_storage_path() {
1954 let id = make_id("original");
1955 let revised = new_long_term_revision(&id);
1956 let path = revised.as_storage_path().to_string();
1957 let parsed = ObjectId::from_storage_path(&path)
1958 .unwrap_or_else(|| panic!("failed to parse '{path}'"));
1959 assert_eq!(parsed, revised);
1960 }
1961
1962 #[test]
1963 fn revision_id_is_unique() {
1964 let id = make_id("base-key");
1965 let a = new_long_term_revision(&id);
1966 let b = new_long_term_revision(&id);
1967 assert_ne!(a.key, b.key, "two calls should produce different keys");
1968 }
1969
1970 #[tokio::test]
1973 async fn get_nonexistent_returns_none() {
1974 let (storage, _hv, _lt, _) = make_tiered_storage();
1975 let id = make_id("does-not-exist");
1976
1977 assert!(
1978 storage
1979 .get_object(&id, Timestamp::now(), None)
1980 .await
1981 .unwrap()
1982 .is_none()
1983 );
1984 assert!(
1985 storage
1986 .get_metadata(&id, Timestamp::now())
1987 .await
1988 .unwrap()
1989 .is_none()
1990 );
1991 }
1992
1993 #[tokio::test]
1994 async fn delete_nonexistent_succeeds() {
1995 let (storage, _hv, _lt, _) = make_tiered_storage();
1996 let id = make_id("does-not-exist");
1997
1998 storage.delete_object(&id, Timestamp::now()).await.unwrap();
1999 }
2000
2001 #[tokio::test]
2004 async fn put_small_object_stores_inline() {
2005 let (storage, hv, lt, _) = make_tiered_storage();
2006 let id = make_id("small");
2007 let payload = b"small payload".to_vec();
2008
2009 storage
2010 .put_object(
2011 &id,
2012 &Metadata::default(),
2013 stream::single(payload.clone()),
2014 Timestamp::now(),
2015 )
2016 .await
2017 .unwrap();
2018
2019 assert!(hv.contains(&id), "expected in high-volume");
2020 assert!(!lt.contains(&id), "leaked to long-term");
2021
2022 let (_, _, s) = storage
2023 .get_object(&id, Timestamp::now(), None)
2024 .await
2025 .unwrap()
2026 .unwrap();
2027 let body = stream::read_to_vec(s).await.unwrap();
2028 assert_eq!(body, payload);
2029
2030 assert!(
2031 storage
2032 .get_metadata(&id, Timestamp::now())
2033 .await
2034 .unwrap()
2035 .is_some(),
2036 "get_metadata should return metadata for inline objects"
2037 );
2038 }
2039
2040 #[tokio::test]
2041 async fn put_large_object_creates_tombstone() {
2042 let (storage, hv, lt, _) = make_tiered_storage();
2043 let id = make_id("large");
2044 let payload = vec![0xCDu8; 2 * 1024 * 1024]; let metadata_in = Metadata {
2046 content_type: "image/png".into(),
2047 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2048 time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2049 origin: Some("10.0.0.1".into()),
2050 ..Metadata::default()
2051 };
2052
2053 storage
2054 .put_object(
2055 &id,
2056 &metadata_in,
2057 stream::single(payload.clone()),
2058 Timestamp::now(),
2059 )
2060 .await
2061 .unwrap();
2062
2063 let tombstone = hv.get(&id).expect_tombstone();
2065 assert_eq!(tombstone.time_expires, metadata_in.time_expires);
2066 let lt_id = tombstone.target;
2067 assert!(
2068 lt_id.key().starts_with(id.key()),
2069 "tombstone target key should be a revision of the HV key, got: {}",
2070 lt_id.key()
2071 );
2072
2073 let (lt_meta, _) = lt.get(<_id).expect_object();
2075 assert_eq!(lt_meta.content_type, "image/png");
2076 assert_eq!(lt_meta.expiration_policy, metadata_in.expiration_policy);
2077 assert_eq!(lt_meta.time_expires, tombstone.time_expires);
2078
2079 let (_, _, s) = storage
2081 .get_object(&id, Timestamp::now(), None)
2082 .await
2083 .unwrap()
2084 .unwrap();
2085 let body = stream::read_to_vec(s).await.unwrap();
2086 assert_eq!(body, payload);
2087
2088 let metadata = storage
2090 .get_metadata(&id, Timestamp::now())
2091 .await
2092 .unwrap()
2093 .unwrap();
2094 assert_eq!(metadata.content_type, "image/png");
2095 }
2096
2097 #[tokio::test]
2100 async fn reinsert_small_over_large_swaps_to_inline() {
2101 let (storage, hv, lt, _) = make_tiered_storage();
2102 let id = make_id("reinsert-key");
2103
2104 let large_payload = vec![0xABu8; 2 * 1024 * 1024];
2106 storage
2107 .put_object(
2108 &id,
2109 &Metadata::default(),
2110 stream::single(large_payload),
2111 Timestamp::now(),
2112 )
2113 .await
2114 .unwrap();
2115
2116 let lt_id = hv.get(&id).expect_tombstone().target;
2117
2118 let small_payload = vec![0xCDu8; 100]; storage
2122 .put_object(
2123 &id,
2124 &Metadata::default(),
2125 stream::single(small_payload),
2126 Timestamp::now(),
2127 )
2128 .await
2129 .unwrap();
2130
2131 hv.get(&id).expect_object();
2133
2134 storage.join().await;
2136
2137 lt.get(<_id).expect_not_found();
2139 }
2140
2141 #[tokio::test]
2142 async fn overwrite_large_with_large_replaces_revision() {
2143 let (storage, hv, lt, _) = make_tiered_storage();
2144 let id = make_id("overwrite-large");
2145
2146 let payload1 = vec![0xAAu8; 2 * 1024 * 1024];
2147 storage
2148 .put_object(
2149 &id,
2150 &Metadata::default(),
2151 stream::single(payload1),
2152 Timestamp::now(),
2153 )
2154 .await
2155 .unwrap();
2156 let lt_id_1 = hv.get(&id).expect_tombstone().target;
2157
2158 let payload2 = vec![0xBBu8; 2 * 1024 * 1024];
2159 storage
2160 .put_object(
2161 &id,
2162 &Metadata::default(),
2163 stream::single(payload2.clone()),
2164 Timestamp::now(),
2165 )
2166 .await
2167 .unwrap();
2168 let lt_id_2 = hv.get(&id).expect_tombstone().target;
2169
2170 assert_ne!(
2171 lt_id_1, lt_id_2,
2172 "second write should create a new revision"
2173 );
2174
2175 storage.join().await;
2177
2178 lt.get(<_id_1).expect_not_found();
2179 lt.get(<_id_2).expect_object();
2180
2181 let (_, _, s) = storage
2182 .get_object(&id, Timestamp::now(), None)
2183 .await
2184 .unwrap()
2185 .unwrap();
2186 let body = stream::read_to_vec(s).await.unwrap();
2187 assert_eq!(body, payload2);
2188 }
2189
2190 #[tokio::test]
2193 async fn delete_small_object() {
2194 let (storage, hv, _lt, _) = make_tiered_storage();
2195 let id = make_id("delete-small");
2196
2197 storage
2198 .put_object(
2199 &id,
2200 &Metadata::default(),
2201 stream::single("tiny"),
2202 Timestamp::now(),
2203 )
2204 .await
2205 .unwrap();
2206
2207 storage.delete_object(&id, Timestamp::now()).await.unwrap();
2208
2209 hv.get(&id).expect_not_found();
2210 assert!(
2211 storage
2212 .get_object(&id, Timestamp::now(), None)
2213 .await
2214 .unwrap()
2215 .is_none()
2216 );
2217 }
2218
2219 #[tokio::test]
2220 async fn delete_large_object_cleans_up_both_backends() {
2221 let (storage, hv, lt, _) = make_tiered_storage();
2222 let id = make_id("delete-both");
2223 let payload = vec![0u8; 2 * 1024 * 1024]; storage
2226 .put_object(
2227 &id,
2228 &Metadata::default(),
2229 stream::single(payload),
2230 Timestamp::now(),
2231 )
2232 .await
2233 .unwrap();
2234
2235 let lt_id = hv.get(&id).expect_tombstone().target;
2237
2238 storage.delete_object(&id, Timestamp::now()).await.unwrap();
2239
2240 storage.join().await;
2242
2243 assert!(!hv.contains(&id), "tombstone not cleaned up");
2244 assert!(!lt.contains(<_id), "long-term object not cleaned up");
2245 }
2246
2247 #[derive(Debug)]
2248 struct FailDelete;
2249
2250 #[async_trait::async_trait]
2251 impl Hooks for FailDelete {
2252 async fn delete_object(
2253 &self,
2254 _inner: &InMemoryBackend,
2255 _id: &ObjectId,
2256 _access_time: Timestamp,
2257 ) -> Result<DeleteResponse> {
2258 Err(Error::with_source(
2259 ErrorKind::BackendFailure,
2260 std::io::Error::new(
2261 std::io::ErrorKind::ConnectionRefused,
2262 "simulated long-term delete failure",
2263 ),
2264 ))
2265 }
2266 }
2267
2268 #[tokio::test]
2272 async fn delete_succeeds_when_gcs_cleanup_fails() {
2273 let hv = InMemoryBackend::new("hv");
2274 let lt = TestBackend::new(FailDelete);
2275 let log = NoopChangeLog;
2276 let storage = TieredStorage::new(Box::new(hv.clone()), Box::new(lt), Box::new(log));
2277
2278 let id = make_id("fail-delete");
2279 let payload = vec![0xABu8; 2 * 1024 * 1024]; storage
2281 .put_object(
2282 &id,
2283 &Metadata::default(),
2284 stream::single(payload),
2285 Timestamp::now(),
2286 )
2287 .await
2288 .unwrap();
2289
2290 let result = storage.delete_object(&id, Timestamp::now()).await;
2292 assert!(
2293 result.is_ok(),
2294 "delete should succeed despite GCS cleanup failure"
2295 );
2296
2297 hv.get(&id).expect_not_found();
2299
2300 assert!(
2302 storage
2303 .get_object(&id, Timestamp::now(), None)
2304 .await
2305 .unwrap()
2306 .is_none(),
2307 "object should be unreachable after tombstone is deleted"
2308 );
2309 }
2310
2311 #[derive(Debug, Copy, Clone)]
2314 struct CasConflict;
2315
2316 #[async_trait::async_trait]
2317 impl Hooks for CasConflict {
2318 async fn compare_and_write(
2319 &self,
2320 _inner: &InMemoryBackend,
2321 _id: &ObjectId,
2322 _current: Option<&ObjectId>,
2323 _write: TieredWrite,
2324 _access_time: Timestamp,
2325 ) -> Result<bool> {
2326 Ok(false) }
2328 }
2329
2330 #[tokio::test]
2334 async fn put_large_cas_conflict_cleans_up_new_blob() {
2335 let hv = TestBackend::new(CasConflict);
2336 let lt = InMemoryBackend::new("lt");
2337 let log = NoopChangeLog;
2338 let storage = TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(log));
2339
2340 let id = make_id("cas-conflict-large");
2341 let payload = vec![0xABu8; 2 * 1024 * 1024]; storage
2344 .put_object(
2345 &id,
2346 &Metadata::default(),
2347 stream::single(payload),
2348 Timestamp::now(),
2349 )
2350 .await
2351 .unwrap();
2352
2353 storage.join().await;
2355
2356 assert!(
2357 lt.is_empty(),
2358 "LT blob should be cleaned up after CAS conflict"
2359 );
2360 }
2361
2362 #[tokio::test]
2363 async fn upload_cas_conflict_cleans_up_new_blob() {
2364 let hv = TestBackend::new(CasConflict);
2365 let lt = InMemoryBackend::new("lt");
2366 let storage =
2367 TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(NoopChangeLog));
2368 let id = make_id("upload-cas-conflict");
2369 let payload = vec![0xAB; BACKEND_SIZE_THRESHOLD + 1];
2370 let session =
2371 resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
2372
2373 assert_eq!(
2374 storage
2375 .put_chunk(&session, 0, payload.len() as u64, stream::single(payload))
2376 .await
2377 .unwrap(),
2378 UploadProgress::Complete
2379 );
2380 storage.join().await;
2381 assert!(
2382 lt.is_empty(),
2383 "LT blob should be cleaned up after CAS conflict"
2384 );
2385 }
2386
2387 #[tokio::test]
2391 async fn put_small_over_tombstone_cas_conflict_succeeds() {
2392 let inner = InMemoryBackend::new("hv");
2393 let id = make_id("cas-conflict-small");
2394
2395 let tombstone = Tombstone {
2398 target: make_id("lt-object"),
2399 time_expires: None,
2400 };
2401 inner
2402 .compare_and_write(
2403 &id,
2404 None,
2405 TieredWrite::Tombstone(tombstone),
2406 Timestamp::now(),
2407 )
2408 .await
2409 .unwrap();
2410
2411 let lt = InMemoryBackend::new("lt");
2412 let hv = TestBackend::with_inner(inner, CasConflict);
2413 let log = NoopChangeLog;
2414 let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(log));
2415
2416 storage
2419 .put_object(
2420 &id,
2421 &Metadata::default(),
2422 stream::single("tiny"),
2423 Timestamp::now(),
2424 )
2425 .await
2426 .unwrap();
2427 }
2428
2429 #[derive(Debug)]
2433 struct FailCas(bool);
2434
2435 #[async_trait::async_trait]
2436 impl Hooks for FailCas {
2437 async fn compare_and_write(
2438 &self,
2439 inner: &InMemoryBackend,
2440 id: &ObjectId,
2441 current: Option<&ObjectId>,
2442 write: TieredWrite,
2443 access_time: Timestamp,
2444 ) -> Result<bool> {
2445 if self.0 {
2446 inner
2448 .compare_and_write(id, current, write, access_time)
2449 .await?;
2450 }
2451 Err(Error::with_source(
2452 ErrorKind::BackendFailure,
2453 std::io::Error::new(
2454 std::io::ErrorKind::TimedOut,
2455 "simulated compare_and_write failure",
2456 ),
2457 ))
2458 }
2459 }
2460
2461 #[tokio::test]
2465 async fn no_orphan_when_tombstone_write_fails() {
2466 let lt = InMemoryBackend::new("lt");
2467 let hv = TestBackend::new(FailCas(false));
2468 let log = NoopChangeLog;
2469 let storage = TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(log));
2470
2471 let id = make_id("orphan-test");
2472 let payload = vec![0xABu8; 2 * 1024 * 1024]; let result = storage
2474 .put_object(
2475 &id,
2476 &Metadata::default(),
2477 stream::single(payload),
2478 Timestamp::now(),
2479 )
2480 .await;
2481
2482 assert!(result.is_err());
2483
2484 storage.join().await;
2486
2487 assert!(lt.is_empty(), "long-term object not cleaned up");
2488 }
2489
2490 #[tokio::test]
2491 async fn upload_cas_failure_cleans_up_new_blob() {
2492 let hv = TestBackend::new(FailCas(false));
2493 let lt = InMemoryBackend::new("lt");
2494 let storage =
2495 TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(NoopChangeLog));
2496 let id = make_id("upload-cas-failure");
2497 let payload = vec![0xAB; BACKEND_SIZE_THRESHOLD + 1];
2498 let session =
2499 resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
2500
2501 assert_eq!(
2502 storage
2503 .put_chunk(&session, 0, payload.len() as u64, stream::single(payload))
2504 .await
2505 .unwrap_err()
2506 .kind(),
2507 ErrorKind::BackendFailure
2508 );
2509 storage.join().await;
2510 assert!(
2511 lt.is_empty(),
2512 "LT blob should be cleaned up after CAS failure"
2513 );
2514 }
2515
2516 #[tokio::test]
2520 async fn orphan_tombstone_returns_none() {
2521 let (storage, hv, lt, _) = make_tiered_storage();
2522 let id = make_id("orphan-tombstone");
2523 let payload = vec![0xCDu8; 2 * 1024 * 1024]; storage
2526 .put_object(
2527 &id,
2528 &Metadata::default(),
2529 stream::single(payload),
2530 Timestamp::now(),
2531 )
2532 .await
2533 .unwrap();
2534
2535 let lt_id = hv.get(&id).expect_tombstone().target;
2537
2538 lt.remove(<_id);
2540
2541 assert!(
2542 storage
2543 .get_object(&id, Timestamp::now(), None)
2544 .await
2545 .unwrap()
2546 .is_none(),
2547 "orphan tombstone should resolve to None on get_object"
2548 );
2549 assert!(
2550 storage
2551 .get_metadata(&id, Timestamp::now())
2552 .await
2553 .unwrap()
2554 .is_none(),
2555 "orphan tombstone should resolve to None on get_metadata"
2556 );
2557 }
2558
2559 #[tokio::test]
2564 async fn tombstone_target_is_used_for_reads_and_deletes() {
2565 let hv = InMemoryBackend::new("hv");
2566 let lt = InMemoryBackend::new("lt");
2567 let log = NoopChangeLog;
2568 let storage = TieredStorage::new(Box::new(hv.clone()), Box::new(lt.clone()), Box::new(log));
2569
2570 let hv_id = make_id("hv-key");
2571 let lt_id = make_id("lt-key");
2572 let payload = vec![0xABu8; 100];
2573
2574 lt.put_object(
2576 <_id,
2577 &Metadata::default(),
2578 stream::single(payload.clone()),
2579 Timestamp::now(),
2580 )
2581 .await
2582 .unwrap();
2583 let tombstone = Tombstone {
2584 target: lt_id.clone(),
2585 time_expires: None,
2586 };
2587 hv.compare_and_write(
2588 &hv_id,
2589 None,
2590 TieredWrite::Tombstone(tombstone),
2591 Timestamp::now(),
2592 )
2593 .await
2594 .unwrap();
2595
2596 let (_, _, s) = storage
2598 .get_object(&hv_id, Timestamp::now(), None)
2599 .await
2600 .unwrap()
2601 .unwrap();
2602 let body = stream::read_to_vec(s).await.unwrap();
2603 assert_eq!(body, payload);
2604
2605 storage
2607 .delete_object(&hv_id, Timestamp::now())
2608 .await
2609 .unwrap();
2610 storage.join().await;
2611 assert!(!hv.contains(&hv_id), "tombstone should be removed");
2612 assert!(!lt.contains(<_id), "lt object should be removed");
2613 }
2614
2615 #[tokio::test]
2618 async fn multi_chunk_large_object_chains_buffered_and_remaining() {
2619 let (storage, hv, lt, _) = make_tiered_storage();
2620 let id = make_id("multi-chunk");
2621
2622 let chunk_size = 512 * 1024; let chunk_count = 4; let stream: ClientStream = futures_util::stream::iter(
2627 (0..chunk_count).map(move |i| Ok(Bytes::from(vec![i as u8; chunk_size]))),
2628 )
2629 .boxed();
2630
2631 storage
2632 .put_object(&id, &Metadata::default(), stream, Timestamp::now())
2633 .await
2634 .unwrap();
2635
2636 let lt_id = hv.get(&id).expect_tombstone().target;
2638 let (_, lt_bytes) = lt.get(<_id).expect_object();
2639 assert_eq!(lt_bytes.len(), chunk_size * chunk_count);
2640
2641 for i in 0..chunk_count {
2643 let offset = i * chunk_size;
2644 assert!(
2645 lt_bytes[offset..offset + chunk_size]
2646 .iter()
2647 .all(|&b| b == i as u8),
2648 "data mismatch in chunk {i}"
2649 );
2650 }
2651 }
2652
2653 #[tokio::test]
2659 async fn written_cleanup_after_lost_cas_response() {
2660 let (storage, hv, lt, log) = make_tiered_storage();
2661 let id = make_id("obj");
2662
2663 let payload = vec![0xAAu8; 2 * 1024 * 1024];
2665 storage
2666 .put_object(
2667 &id,
2668 &Metadata::default(),
2669 stream::single(payload.clone()),
2670 Timestamp::now(),
2671 )
2672 .await
2673 .unwrap();
2674 let tombstone1 = hv.get(&id).expect_tombstone().target;
2675
2676 let broken_storage = TieredStorage::new(
2678 Box::new(TestBackend::with_inner(hv.clone(), FailCas(true))),
2679 Box::new(lt.clone()),
2680 Box::new(log.clone()),
2681 );
2682 broken_storage
2683 .put_object(
2684 &id,
2685 &Metadata::default(),
2686 stream::single(payload.clone()),
2687 Timestamp::now(),
2688 )
2689 .await
2690 .unwrap_err(); let tombstone2 = hv.get(&id).expect_tombstone().target;
2692 assert_ne!(tombstone1, tombstone2);
2693
2694 broken_storage.join().await;
2696 lt.get(&tombstone1).expect_not_found();
2697 lt.get(&tombstone2).expect_object();
2698
2699 broken_storage
2701 .delete_object(&id, Timestamp::now())
2702 .await
2703 .unwrap_err();
2704 hv.get(&id).expect_not_found();
2705 broken_storage.join().await;
2706 lt.get(&tombstone2).expect_not_found();
2707
2708 let id = make_id("obj2");
2710 storage
2711 .put_object(
2712 &id,
2713 &Metadata::default(),
2714 stream::single(payload.clone()),
2715 Timestamp::now(),
2716 )
2717 .await
2718 .unwrap();
2719 let tombstone3 = hv.get(&id).expect_tombstone().target;
2720
2721 broken_storage
2723 .put_object(
2724 &id,
2725 &Metadata::default(),
2726 stream::single(&b"small"[..]),
2727 Timestamp::now(),
2728 )
2729 .await
2730 .unwrap_err(); hv.get(&id).expect_object();
2732 broken_storage.join().await;
2733 lt.get(&tombstone3).expect_not_found();
2734 }
2735
2736 #[test]
2740 fn guard_dropped_outside_runtime_does_not_panic() {
2741 let manager = ChangeManager::new(
2742 Box::new(InMemoryBackend::new("hv")),
2743 Box::new(InMemoryBackend::new("lt")),
2744 Box::new(NoopChangeLog),
2745 );
2746
2747 let change = Change {
2748 id: make_id("object-key"),
2749 new: Some(make_id("cleanup-target")),
2750 old: None,
2751 cleanup_after: None,
2752 };
2753
2754 let guard = {
2757 let rt = tokio::runtime::Runtime::new().unwrap();
2758 rt.block_on(manager.record(change)).unwrap()
2759 };
2760
2761 drop(guard); }
2763
2764 #[tokio::test(start_paused = true)]
2769 async fn join_waits_for_cleanup_to_complete() {
2770 let (storage, _hv, _lt, _) = make_tiered_storage();
2771 let change = Change {
2772 id: make_id("object-key"),
2773 new: None,
2774 old: None,
2775 cleanup_after: None,
2776 };
2777 let mut guard = storage.record_change(change).await.unwrap();
2778
2779 tokio::spawn(async move {
2780 tokio::time::sleep(Duration::from_secs(10)).await;
2781 guard.advance(ChangePhase::Completed);
2782 drop(guard);
2783 });
2784
2785 let join_future = tokio::spawn(async move { storage.join().await });
2786
2787 tokio::time::sleep(Duration::from_secs(9)).await;
2788 assert!(!join_future.is_finished(), "finished before guard dropped");
2789
2790 tokio::time::sleep(Duration::from_secs(2)).await;
2791 assert!(join_future.is_finished(), "finish after guard drops");
2792 }
2793
2794 #[derive(Clone, Debug)]
2801 struct PauseAfterPut {
2802 paused: Arc<tokio::sync::Notify>,
2803 resume: Arc<tokio::sync::Notify>,
2804 }
2805
2806 #[async_trait::async_trait]
2807 impl Hooks for PauseAfterPut {
2808 async fn put_object(
2809 &self,
2810 inner: &InMemoryBackend,
2811 id: &ObjectId,
2812 metadata: &Metadata,
2813 stream: ClientStream,
2814 access_time: Timestamp,
2815 ) -> Result<PutResponse> {
2816 inner.put_object(id, metadata, stream, access_time).await?;
2817 self.paused.notify_one();
2818 self.resume.notified().await;
2819 Ok(())
2820 }
2821 }
2822
2823 #[tokio::test]
2826 async fn dropped_future_triggers_cleanup_and_log_entry_removed() {
2827 let paused = Arc::new(tokio::sync::Notify::new());
2828 let hooks = PauseAfterPut {
2829 paused: Arc::clone(&paused),
2830 resume: Arc::new(tokio::sync::Notify::new()),
2831 };
2832
2833 let lt_inner = InMemoryBackend::new("lt");
2834 let log = InMemoryChangeLog::default();
2835 let storage = TieredStorage::new(
2836 Box::new(InMemoryBackend::new("hv")),
2837 Box::new(TestBackend::with_inner(lt_inner.clone(), hooks)),
2838 Box::new(log.clone()),
2839 );
2840
2841 let id = make_id("drop-test");
2842 let metadata = Metadata::default();
2843 let payload = vec![0xABu8; 2 * 1024 * 1024]; tokio::select! {
2847 result = storage.put_object(&id, &metadata, stream::single(payload), Timestamp::now()) => {
2848 panic!("expected put to pause before completing, got: {result:?}");
2849 }
2850 _ = paused.notified() => {
2851 }
2853 }
2854
2855 storage.join().await;
2857
2858 assert!(lt_inner.is_empty(), "orphaned LT blob was not cleaned up");
2860
2861 let entries = log.scan().await.unwrap();
2863 assert!(
2864 entries.is_empty(),
2865 "changelog entry not removed after cleanup"
2866 );
2867 }
2868
2869 #[test]
2872 fn multipart_upload_id_roundtrip() {
2873 let id = TieredUploadId {
2874 revision: "my-key/01924a6f-7e28-7b9a-9c1d-abcdef123456".into(),
2875 upload_id: UploadId::new("upstream-upload-id-abc".into()).unwrap(),
2876 };
2877 let encoded: UploadId = id.clone().try_into().unwrap();
2878 let decoded: TieredUploadId = (&encoded.clone()).try_into().unwrap();
2879 assert_eq!(decoded, id);
2880 }
2881
2882 #[test]
2883 fn malformed_multipart_upload_ids_are_invalid_upload_ids() {
2884 let invalid_base64 = UploadId::new("%%%".into()).unwrap();
2885 let malformed_json = UploadId::new("bm90IGpzb24".into()).unwrap();
2886
2887 for upload_id in [&invalid_base64, &malformed_json] {
2888 let error = TieredUploadId::try_from(upload_id).unwrap_err();
2889 assert_eq!(error.kind(), ErrorKind::InvalidUploadId);
2890 }
2891 }
2892
2893 #[tokio::test]
2894 async fn multipart_single_part_roundtrip() {
2895 let (storage, hv, lt, _) = make_tiered_storage();
2896 let id = make_id("mp-single");
2897 let metadata = Metadata {
2898 content_type: "application/octet-stream".into(),
2899 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2900 time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2901 ..Metadata::default()
2902 };
2903 let payload = vec![0xABu8; 2 * 1024 * 1024]; let upload_id = storage.initiate_multipart(&id, &metadata).await.unwrap();
2906
2907 let etag = storage
2908 .upload_part(
2909 &id,
2910 &upload_id,
2911 NonZeroU32::new(1).unwrap(),
2912 payload.len() as u64,
2913 None,
2914 stream::single(payload.clone()),
2915 )
2916 .await
2917 .unwrap();
2918
2919 let error = storage
2920 .complete_multipart(
2921 &id,
2922 &upload_id,
2923 vec![CompletedPart {
2924 part_number: NonZeroU32::new(1).unwrap(),
2925 etag,
2926 }],
2927 Timestamp::now(),
2928 )
2929 .await
2930 .unwrap();
2931 assert!(
2932 error.is_none(),
2933 "complete_multipart returned error: {error:?}"
2934 );
2935
2936 let (got_meta, _, s) = storage
2938 .get_object(&id, Timestamp::now(), None)
2939 .await
2940 .unwrap()
2941 .unwrap();
2942 let body = stream::read_to_vec(s).await.unwrap();
2943 assert_eq!(body, payload);
2944 assert_eq!(got_meta.content_type, "application/octet-stream");
2945
2946 let tombstone = hv.get(&id).expect_tombstone();
2948 assert!(
2949 tombstone.target.key().starts_with(id.key()),
2950 "tombstone target should be a revision key"
2951 );
2952 assert_eq!(tombstone.time_expires, metadata.time_expires);
2953 assert_eq!(
2954 lt.get(&tombstone.target).expect_object().0.time_expires,
2955 tombstone.time_expires
2956 );
2957 }
2958
2959 #[tokio::test]
2960 async fn multipart_upload() {
2961 let (storage, _hv, _lt, _) = make_tiered_storage();
2962 let id = make_id("multipart");
2963
2964 let upload_id = storage
2965 .initiate_multipart(&id, &Metadata::default())
2966 .await
2967 .unwrap();
2968
2969 let part1 = vec![0xAAu8; 512 * 1024];
2970 let part2 = vec![0xBBu8; 512 * 1024];
2971 let part3 = vec![0xCCu8; 512 * 1024];
2972
2973 let etag3 = storage
2974 .upload_part(
2975 &id,
2976 &upload_id,
2977 NonZeroU32::new(3).unwrap(),
2978 part3.len() as u64,
2979 None,
2980 stream::single(part3.clone()),
2981 )
2982 .await
2983 .unwrap();
2984 let etag2 = storage
2985 .upload_part(
2986 &id,
2987 &upload_id,
2988 NonZeroU32::new(2).unwrap(),
2989 part2.len() as u64,
2990 None,
2991 stream::single(part2.clone()),
2992 )
2993 .await
2994 .unwrap();
2995 let etag1 = storage
2996 .upload_part(
2997 &id,
2998 &upload_id,
2999 NonZeroU32::new(1).unwrap(),
3000 part1.len() as u64,
3001 None,
3002 stream::single(part1.clone()),
3003 )
3004 .await
3005 .unwrap();
3006
3007 let error = storage
3008 .complete_multipart(
3009 &id,
3010 &upload_id,
3011 vec![
3012 CompletedPart {
3013 part_number: NonZeroU32::new(1).unwrap(),
3014 etag: etag1,
3015 },
3016 CompletedPart {
3017 part_number: NonZeroU32::new(2).unwrap(),
3018 etag: etag2,
3019 },
3020 CompletedPart {
3021 part_number: NonZeroU32::new(3).unwrap(),
3022 etag: etag3,
3023 },
3024 ],
3025 Timestamp::now(),
3026 )
3027 .await
3028 .unwrap();
3029 assert!(error.is_none());
3030
3031 let (_, _, s) = storage
3032 .get_object(&id, Timestamp::now(), None)
3033 .await
3034 .unwrap()
3035 .unwrap();
3036 let body = stream::read_to_vec(s).await.unwrap();
3037
3038 let mut expected = Vec::new();
3039 expected.extend_from_slice(&part1);
3040 expected.extend_from_slice(&part2);
3041 expected.extend_from_slice(&part3);
3042 assert_eq!(body, expected);
3043 }
3044
3045 #[tokio::test]
3046 async fn multipart_abort() {
3047 let (storage, hv, _lt, _) = make_tiered_storage();
3048 let id = make_id("mp-abort");
3049
3050 let upload_id = storage
3051 .initiate_multipart(&id, &Metadata::default())
3052 .await
3053 .unwrap();
3054
3055 let payload = vec![0xABu8; 100];
3057 storage
3058 .upload_part(
3059 &id,
3060 &upload_id,
3061 NonZeroU32::new(1).unwrap(),
3062 payload.len() as u64,
3063 None,
3064 stream::single(payload),
3065 )
3066 .await
3067 .unwrap();
3068
3069 storage.abort_multipart(&id, &upload_id).await.unwrap();
3070
3071 hv.get(&id).expect_not_found();
3073
3074 assert!(
3076 storage
3077 .get_object(&id, Timestamp::now(), None)
3078 .await
3079 .unwrap()
3080 .is_none()
3081 );
3082 }
3083
3084 #[tokio::test]
3085 async fn multipart_list_parts() {
3086 let (storage, _hv, _lt, _) = make_tiered_storage();
3087 let id = make_id("mp-list");
3088
3089 let upload_id = storage
3090 .initiate_multipart(&id, &Metadata::default())
3091 .await
3092 .unwrap();
3093
3094 let part1 = vec![0xAAu8; 100];
3095 let part2 = vec![0xBBu8; 200];
3096 storage
3097 .upload_part(
3098 &id,
3099 &upload_id,
3100 NonZeroU32::new(1).unwrap(),
3101 part1.len() as u64,
3102 None,
3103 stream::single(part1),
3104 )
3105 .await
3106 .unwrap();
3107 storage
3108 .upload_part(
3109 &id,
3110 &upload_id,
3111 NonZeroU32::new(2).unwrap(),
3112 part2.len() as u64,
3113 None,
3114 stream::single(part2),
3115 )
3116 .await
3117 .unwrap();
3118
3119 let resp = storage
3120 .list_parts(&id, &upload_id, None, None)
3121 .await
3122 .unwrap();
3123 assert_eq!(resp.parts.len(), 2);
3124 assert_eq!(resp.parts[0].part_number.get(), 1);
3125 assert_eq!(resp.parts[0].size, 100);
3126 assert_eq!(resp.parts[1].part_number.get(), 2);
3127 assert_eq!(resp.parts[1].size, 200);
3128 }
3129
3130 #[tokio::test]
3131 async fn multipart_overwrites_existing_tombstone() {
3132 let (storage, hv, lt, _) = make_tiered_storage();
3133 let id = make_id("mp-overwrite");
3134
3135 let payload1 = vec![0xAAu8; 2 * 1024 * 1024];
3137 storage
3138 .put_object(
3139 &id,
3140 &Metadata::default(),
3141 stream::single(payload1),
3142 Timestamp::now(),
3143 )
3144 .await
3145 .unwrap();
3146 let old_lt_id = hv.get(&id).expect_tombstone().target;
3147
3148 let upload_id = storage
3150 .initiate_multipart(&id, &Metadata::default())
3151 .await
3152 .unwrap();
3153
3154 let payload2 = vec![0xBBu8; 2 * 1024 * 1024];
3155 let etag = storage
3156 .upload_part(
3157 &id,
3158 &upload_id,
3159 NonZeroU32::new(1).unwrap(),
3160 payload2.len() as u64,
3161 None,
3162 stream::single(payload2.clone()),
3163 )
3164 .await
3165 .unwrap();
3166
3167 let lt_id = hv.get(&id).expect_tombstone().target;
3170 assert_eq!(old_lt_id, lt_id);
3171
3172 let error = storage
3173 .complete_multipart(
3174 &id,
3175 &upload_id,
3176 vec![CompletedPart {
3177 part_number: NonZeroU32::new(1).unwrap(),
3178 etag,
3179 }],
3180 Timestamp::now(),
3181 )
3182 .await
3183 .unwrap();
3184 assert!(error.is_none());
3185
3186 let new_lt_id = hv.get(&id).expect_tombstone().target;
3188 assert_ne!(old_lt_id, new_lt_id);
3189
3190 storage.join().await;
3192
3193 lt.get(&old_lt_id).expect_not_found();
3195 lt.get(&new_lt_id).expect_object();
3196
3197 let (_, _, s) = storage
3199 .get_object(&id, Timestamp::now(), None)
3200 .await
3201 .unwrap()
3202 .unwrap();
3203 let body = stream::read_to_vec(s).await.unwrap();
3204 assert_eq!(body, payload2);
3205 }
3206
3207 #[derive(Debug)]
3213 struct CompleteMultipartButReturnError;
3214
3215 #[async_trait::async_trait]
3216 impl Hooks for CompleteMultipartButReturnError {
3217 async fn complete_multipart(
3218 &self,
3219 inner: &InMemoryBackend,
3220 id: &ObjectId,
3221 upload_id: &UploadId,
3222 parts: Vec<CompletedPart>,
3223 access_time: Timestamp,
3224 ) -> Result<CompleteMultipartResponse> {
3225 inner
3226 .complete_multipart(id, upload_id, parts, access_time)
3227 .await
3228 .unwrap();
3229 Err(Error::with_source(
3230 ErrorKind::BackendFailure,
3231 std::io::Error::new(
3232 std::io::ErrorKind::TimedOut,
3233 "simulated network error on complete_multipart",
3234 ),
3235 ))
3236 }
3237
3238 async fn get_metadata(
3239 &self,
3240 _inner: &InMemoryBackend,
3241 _id: &ObjectId,
3242 _access_time: Timestamp,
3243 ) -> Result<MetadataResponse> {
3244 Err(Error::with_source(
3245 ErrorKind::BackendFailure,
3246 std::io::Error::new(
3247 std::io::ErrorKind::TimedOut,
3248 "simulated network error on get_metadata",
3249 ),
3250 ))
3251 }
3252 }
3253
3254 #[tokio::test]
3258 async fn cleans_up_orphan_after_failed_multipart_complete() {
3259 let hv = InMemoryBackend::new("hv");
3260 let lt_inner = InMemoryBackend::new("lt");
3261 let log = InMemoryChangeLog::default();
3262 let storage = TieredStorage::new(
3263 Box::new(hv.clone()),
3264 Box::new(TestBackend::with_inner(
3265 lt_inner.clone(),
3266 CompleteMultipartButReturnError {},
3267 )),
3268 Box::new(log.clone()),
3269 );
3270
3271 let id = make_id("mp-orphan");
3272 let upload_id = storage
3273 .initiate_multipart(&id, &Metadata::default())
3274 .await
3275 .unwrap();
3276
3277 let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3278 let physical = ObjectId {
3279 context: id.context.clone(),
3280 key: tiered_id.revision,
3281 };
3282
3283 let payload = vec![0xABu8; 2 * 1024 * 1024];
3284 let etag = storage
3285 .upload_part(
3286 &id,
3287 &upload_id,
3288 NonZeroU32::new(1).unwrap(),
3289 payload.len() as u64,
3290 None,
3291 stream::single(payload),
3292 )
3293 .await
3294 .unwrap();
3295
3296 let result = storage
3297 .complete_multipart(
3298 &id,
3299 &upload_id,
3300 vec![CompletedPart {
3301 part_number: NonZeroU32::new(1).unwrap(),
3302 etag,
3303 }],
3304 Timestamp::now(),
3305 )
3306 .await;
3307 assert!(result.is_err());
3308 storage.join().await;
3309
3310 lt_inner.get(&physical).expect_object();
3313 hv.get(&id).expect_not_found();
3314
3315 log.expire_all();
3317 let manager = ChangeManager::new(
3318 Box::new(hv.clone()),
3319 Box::new(lt_inner.clone()),
3320 Box::new(log.clone()),
3321 );
3322 manager.recover().await.unwrap();
3323
3324 lt_inner.get(&physical).expect_not_found();
3326 let remaining = log.scan().await.unwrap();
3328 assert!(remaining.is_empty());
3329 }
3330
3331 #[derive(Debug)]
3332 struct FailOnFirstCompleteMultipartAttempt {
3333 attempt: Mutex<u32>,
3334 }
3335
3336 impl FailOnFirstCompleteMultipartAttempt {
3337 fn new() -> Self {
3338 Self {
3339 attempt: Mutex::new(0),
3340 }
3341 }
3342 }
3343
3344 #[async_trait::async_trait]
3345 impl Hooks for FailOnFirstCompleteMultipartAttempt {
3346 async fn complete_multipart(
3347 &self,
3348 inner: &InMemoryBackend,
3349 id: &ObjectId,
3350 upload_id: &UploadId,
3351 parts: Vec<CompletedPart>,
3352 access_time: Timestamp,
3353 ) -> Result<CompleteMultipartResponse> {
3354 let mut attempt = self.attempt.lock().await;
3355 *attempt += 1;
3356 if *attempt == 1 {
3357 Err(Error::with_source(
3358 ErrorKind::BackendFailure,
3359 std::io::Error::new(std::io::ErrorKind::TimedOut, "simulated network error"),
3360 ))
3361 } else {
3362 Ok(inner
3363 .complete_multipart(id, upload_id, parts, access_time)
3364 .await
3365 .unwrap())
3366 }
3367 }
3368 }
3369
3370 #[tokio::test]
3375 async fn multipart_complete_succeeds_on_retry_and_leaves_state_consistent() {
3376 let hv = InMemoryBackend::new("hv");
3377 let lt_inner = InMemoryBackend::new("lt");
3378 let log = InMemoryChangeLog::default();
3379 let storage = TieredStorage::new(
3380 Box::new(hv.clone()),
3381 Box::new(TestBackend::with_inner(
3382 lt_inner.clone(),
3383 FailOnFirstCompleteMultipartAttempt::new(),
3384 )),
3385 Box::new(log.clone()),
3386 );
3387
3388 let id = make_id("mp-retry");
3389 let upload_id = storage
3390 .initiate_multipart(&id, &Metadata::default())
3391 .await
3392 .unwrap();
3393
3394 let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3395 let physical = ObjectId {
3396 context: id.context.clone(),
3397 key: tiered_id.revision,
3398 };
3399
3400 let payload = vec![0xABu8; 2 * 1024 * 1024];
3401 let etag = storage
3402 .upload_part(
3403 &id,
3404 &upload_id,
3405 NonZeroU32::new(1).unwrap(),
3406 payload.len() as u64,
3407 None,
3408 stream::single(payload.clone()),
3409 )
3410 .await
3411 .unwrap();
3412
3413 let result = storage
3415 .complete_multipart(
3416 &id,
3417 &upload_id,
3418 vec![CompletedPart {
3419 part_number: NonZeroU32::new(1).unwrap(),
3420 etag: etag.clone(),
3421 }],
3422 Timestamp::now(),
3423 )
3424 .await;
3425 assert!(result.is_err());
3426 storage.join().await;
3427
3428 let result = storage
3430 .complete_multipart(
3431 &id,
3432 &upload_id,
3433 vec![CompletedPart {
3434 part_number: NonZeroU32::new(1).unwrap(),
3435 etag,
3436 }],
3437 Timestamp::now(),
3438 )
3439 .await;
3440 assert!(result.is_ok());
3441 storage.join().await;
3442
3443 let (_, _, s) = storage
3445 .get_object(&id, Timestamp::now(), None)
3446 .await
3447 .unwrap()
3448 .unwrap();
3449 let body = stream::read_to_vec(s).await.unwrap();
3450 assert_eq!(body, payload);
3451
3452 log.expire_all();
3454 let manager = ChangeManager::new(
3455 Box::new(hv.clone()),
3456 Box::new(lt_inner.clone()),
3457 Box::new(log.clone()),
3458 );
3459 manager.recover().await.unwrap();
3460
3461 lt_inner.get(&physical).expect_object();
3463 let tombstone = hv.get(&id).expect_tombstone();
3465 assert_eq!(tombstone.target, physical);
3466 let remaining = log.scan().await.unwrap();
3468 assert!(remaining.is_empty());
3469
3470 let (_, _, s) = storage
3472 .get_object(&id, Timestamp::now(), None)
3473 .await
3474 .unwrap()
3475 .unwrap();
3476 let body = stream::read_to_vec(s).await.unwrap();
3477 assert_eq!(body, payload);
3478 }
3479
3480 #[derive(Debug)]
3481 struct FailOnFirstGetMetadataAttempt {
3482 attempt: Mutex<u32>,
3483 }
3484
3485 impl FailOnFirstGetMetadataAttempt {
3486 fn new() -> Self {
3487 Self {
3488 attempt: Mutex::new(0),
3489 }
3490 }
3491 }
3492
3493 #[async_trait::async_trait]
3494 impl Hooks for FailOnFirstGetMetadataAttempt {
3495 async fn get_metadata(
3496 &self,
3497 inner: &InMemoryBackend,
3498 id: &ObjectId,
3499 access_time: Timestamp,
3500 ) -> Result<MetadataResponse> {
3501 let mut attempt = self.attempt.lock().await;
3502 *attempt += 1;
3503 if *attempt == 1 {
3504 Err(Error::with_source(
3505 ErrorKind::BackendFailure,
3506 std::io::Error::new(std::io::ErrorKind::TimedOut, "simulated network error"),
3507 ))
3508 } else {
3509 inner.get_metadata(id, access_time).await
3510 }
3511 }
3512 }
3513
3514 #[tokio::test]
3521 async fn multipart_complete_succeeds_on_retry_if_get_metadata_errs_and_leaves_state_consistent()
3522 {
3523 let hv = InMemoryBackend::new("hv");
3524 let lt_inner = InMemoryBackend::new("lt");
3525 let log = InMemoryChangeLog::default();
3526 let storage = TieredStorage::new(
3527 Box::new(hv.clone()),
3528 Box::new(TestBackend::with_inner(
3529 lt_inner.clone(),
3530 FailOnFirstGetMetadataAttempt::new(),
3531 )),
3532 Box::new(log.clone()),
3533 );
3534
3535 let id = make_id("mp-retry-meta");
3536 let upload_id = storage
3537 .initiate_multipart(&id, &Metadata::default())
3538 .await
3539 .unwrap();
3540
3541 let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3542 let physical = ObjectId {
3543 context: id.context.clone(),
3544 key: tiered_id.revision,
3545 };
3546
3547 let payload = vec![0xABu8; 2 * 1024 * 1024];
3548 let etag = storage
3549 .upload_part(
3550 &id,
3551 &upload_id,
3552 NonZeroU32::new(1).unwrap(),
3553 payload.len() as u64,
3554 None,
3555 stream::single(payload.clone()),
3556 )
3557 .await
3558 .unwrap();
3559
3560 let result = storage
3563 .complete_multipart(
3564 &id,
3565 &upload_id,
3566 vec![CompletedPart {
3567 part_number: NonZeroU32::new(1).unwrap(),
3568 etag: etag.clone(),
3569 }],
3570 Timestamp::now(),
3571 )
3572 .await;
3573 assert!(result.is_err());
3574 storage.join().await;
3575
3576 let result = storage
3578 .complete_multipart(
3579 &id,
3580 &upload_id,
3581 vec![CompletedPart {
3582 part_number: NonZeroU32::new(1).unwrap(),
3583 etag,
3584 }],
3585 Timestamp::now(),
3586 )
3587 .await;
3588 assert!(result.is_ok());
3589 storage.join().await;
3590
3591 let (_, _, s) = storage
3593 .get_object(&id, Timestamp::now(), None)
3594 .await
3595 .unwrap()
3596 .unwrap();
3597 let body = stream::read_to_vec(s).await.unwrap();
3598 assert_eq!(body, payload);
3599
3600 log.expire_all();
3602 let manager = ChangeManager::new(
3603 Box::new(hv.clone()),
3604 Box::new(lt_inner.clone()),
3605 Box::new(log.clone()),
3606 );
3607 manager.recover().await.unwrap();
3608
3609 lt_inner.get(&physical).expect_object();
3611 let tombstone = hv.get(&id).expect_tombstone();
3613 assert_eq!(tombstone.target, physical);
3614 let remaining = log.scan().await.unwrap();
3616 assert!(remaining.is_empty());
3617
3618 let (_, _, s) = storage
3620 .get_object(&id, Timestamp::now(), None)
3621 .await
3622 .unwrap()
3623 .unwrap();
3624 let body = stream::read_to_vec(s).await.unwrap();
3625 assert_eq!(body, payload);
3626 }
3627}