1use crate::attempts::count_retry_attempt;
6use crate::immutable_write::{readback, ImmutableReadback};
7use crate::keyspace::{
8 normalize_key_prefix, scope_list_prefix, scope_object_key, unscope_listed_key,
9};
10use crate::object_store::Result;
11use crate::retry::{
12 transport_retry_backoff, transport_retry_pause, OperationDeadline, TransportRetryPolicy,
13};
14use crate::timing::{MonotonicTimer, StdMonotonicTimer};
15use crate::{
16 ByteRange, ByteStream, ObjectBody, ObjectMetadata, ObjectStore, ObjectStoreError, PutMode,
17};
18use async_trait::async_trait;
19use bytes::Bytes;
20use futures::stream::{self, BoxStream, FuturesUnordered, StreamExt};
21use object_store as provider_store;
22use provider_store::multipart::{MultipartStore, PartId};
23use provider_store::path::Path;
24use provider_store::{
25 GetOptions, GetRange, ObjectMeta, PutOptions, PutPayload, PutResult, UpdateVersion,
26};
27use std::fmt;
28use std::ops::Range;
29use std::sync::Arc;
30use std::time::Duration;
31
32#[derive(Debug, Clone, PartialEq, Eq)]
34pub struct ProviderObjectStoreConfig {
35 pub key_prefix: Option<String>,
37}
38
39pub const PROVIDER_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(30);
44
45pub const PROVIDER_CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
47
48pub const PROVIDER_OPERATION_DEADLINE: Duration = Duration::from_secs(120);
64
65pub const PROVIDER_MULTIPART_THRESHOLD_BYTES: u64 = 8 * 1024 * 1024;
74
75pub const PROVIDER_MULTIPART_PART_BYTES: u64 = 8 * 1024 * 1024;
81
82pub const PROVIDER_MULTIPART_PART_WINDOW: usize = 4;
84
85pub const PROVIDER_STREAMED_PART_WINDOW: usize = 1;
92
93pub(crate) const MAX_PROVIDER_MULTIPART_PARTS: usize = 10_000;
97
98pub const PROVIDER_TRANSFER_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(120);
105
106pub(crate) const PROVIDER_TRANSFER_BODY_MIN_BYTES: u64 = 1024 * 1024;
112
113pub(crate) fn request_phase_bound(request_body_bytes: u64) -> Duration {
117 if request_body_bytes >= PROVIDER_TRANSFER_BODY_MIN_BYTES {
118 PROVIDER_TRANSFER_ATTEMPT_TIMEOUT
119 } else {
120 PROVIDER_ATTEMPT_TIMEOUT
121 }
122}
123
124pub(crate) fn provider_client_options() -> provider_store::ClientOptions {
131 provider_store::ClientOptions::new()
132 .with_timeout(PROVIDER_ATTEMPT_TIMEOUT)
133 .with_connect_timeout(PROVIDER_CONNECT_TIMEOUT)
134}
135
136pub(crate) fn provider_retry_config() -> provider_store::RetryConfig {
140 provider_store::RetryConfig {
141 retry_timeout: PROVIDER_OPERATION_DEADLINE,
142 ..Default::default()
143 }
144}
145
146#[derive(Debug, Clone, Copy, PartialEq, Eq)]
151struct MultipartGeometry {
152 threshold_bytes: u64,
153 part_bytes: u64,
154}
155
156impl MultipartGeometry {
157 const DEFAULT: Self = Self {
158 threshold_bytes: PROVIDER_MULTIPART_THRESHOLD_BYTES,
159 part_bytes: PROVIDER_MULTIPART_PART_BYTES,
160 };
161}
162
163#[derive(Clone)]
165pub struct ProviderObjectStore {
166 inner: Arc<dyn provider_store::ObjectStore>,
167 multipart: Option<Arc<dyn MultipartStore>>,
172 multipart_geometry: MultipartGeometry,
173 key_prefix: Option<String>,
174 transport_retry: TransportRetryPolicy,
175 timer: Arc<dyn MonotonicTimer>,
176}
177
178impl fmt::Debug for ProviderObjectStore {
179 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
180 f.debug_struct("ProviderObjectStore")
181 .field("key_prefix", &self.key_prefix)
182 .field("multipart_upload", &self.multipart.is_some())
183 .finish_non_exhaustive()
184 }
185}
186
187impl ProviderObjectStore {
188 pub fn new(
193 inner: Arc<dyn provider_store::ObjectStore>,
194 multipart: Option<Arc<dyn MultipartStore>>,
195 config: ProviderObjectStoreConfig,
196 ) -> Result<Self> {
197 Ok(Self {
198 inner,
199 multipart,
200 multipart_geometry: MultipartGeometry::DEFAULT,
201 key_prefix: normalize_key_prefix(config.key_prefix.as_deref())?,
202 transport_retry: TransportRetryPolicy::DEFAULT,
203 timer: Arc::new(StdMonotonicTimer::default()),
204 })
205 }
206
207 #[cfg(test)]
208 fn transport_retry(mut self, transport_retry: TransportRetryPolicy) -> Self {
209 self.transport_retry = transport_retry;
210 self
211 }
212
213 #[cfg(test)]
214 fn multipart_geometry(mut self, threshold_bytes: u64, part_bytes: u64) -> Self {
215 self.multipart_geometry = MultipartGeometry {
216 threshold_bytes,
217 part_bytes,
218 };
219 self
220 }
221
222 #[cfg(test)]
223 fn monotonic_timer(mut self, timer: Arc<dyn MonotonicTimer>) -> Self {
224 self.timer = timer;
225 self
226 }
227
228 fn to_path(&self, key: &str) -> Result<Path> {
229 let scoped = scope_object_key(self.key_prefix.as_deref(), key)?;
230 Path::parse(scoped).map_err(|err| ObjectStoreError::InvalidKey {
231 object_key: key.to_owned(),
232 message: err.to_string(),
233 })
234 }
235
236 pub(crate) fn validate_key(&self, key: &str) -> Result<()> {
237 self.to_path(key).map(|_| ())
238 }
239
240 fn list_path(&self, prefix: &str) -> Result<Option<Path>> {
241 let scoped = scope_list_prefix(self.key_prefix.as_deref(), prefix)?;
242 if scoped.is_empty() {
243 return Ok(None);
244 }
245 Path::parse(scoped)
246 .map(Some)
247 .map_err(|err| ObjectStoreError::InvalidKey {
248 object_key: prefix.to_owned(),
249 message: err.to_string(),
250 })
251 }
252
253 fn from_meta(meta: ObjectMeta) -> ObjectMetadata {
254 ObjectMetadata {
255 etag: meta.e_tag,
256 version: meta.version,
257 size_bytes: meta.size,
258 last_modified_ms: u64::try_from(meta.last_modified.timestamp_millis()).ok(),
259 }
260 }
261
262 fn from_put_result(result: PutResult, size_bytes: u64) -> ObjectMetadata {
263 ObjectMetadata {
264 etag: result.e_tag,
265 version: result.version,
266 size_bytes,
267 last_modified_ms: None,
268 }
269 }
270
271 async fn ranged_get(&self, path: &Path, start: u64, end: u64) -> RangedGet {
272 let options = GetOptions {
273 range: Some(GetRange::Bounded(Range { start, end })),
274 ..Default::default()
275 };
276 match self.inner.get_opts(path, options).await {
277 Ok(result) => match result.bytes().await {
278 Ok(bytes) => RangedGet::Bytes(bytes),
279 Err(err) => RangedGet::Refused(err),
280 },
281 Err(err) if provider_not_found(&err) => RangedGet::NotFound,
282 Err(err) => RangedGet::Refused(err),
283 }
284 }
285
286 async fn put_large_multipart(
303 &self,
304 multipart: &dyn MultipartStore,
305 key: &str,
306 path: &Path,
307 bytes: Bytes,
308 ) -> Result<ObjectMetadata> {
309 let size_bytes = bytes.len() as u64;
310 let upload = MultipartWrite {
311 store: self,
312 multipart,
313 key,
314 path,
315 };
316
317 let upload_id = upload.create(size_bytes).await?;
318 let result = upload.upload_parts_and_complete(&upload_id, &bytes).await;
319 match result {
320 Ok(metadata) => Ok(metadata),
321 Err(err) => {
322 upload.abort(&upload_id).await;
326 Err(err)
327 }
328 }
329 }
330
331 fn next_write_backoff(
340 &self,
341 key: &str,
342 operation: &'static str,
343 payload_bytes: u64,
344 retries: &mut u32,
345 deadline: Option<&OperationDeadline<'_>>,
346 err: &provider_store::Error,
347 ) -> Option<Duration> {
348 if *retries >= self.transport_retry.max_retries {
349 tracing::warn!(
350 object_key = key,
351 operation,
352 retry = *retries,
353 payload_bytes,
354 error = %err,
355 "object store write retry budget exhausted; not retrying",
356 );
357 return None;
358 }
359 let mut remaining = Duration::MAX;
360 if let Some(deadline) = deadline {
361 let Some(deadline_remaining) = deadline.remaining() else {
362 tracing::warn!(
363 object_key = key,
364 operation,
365 retry = *retries,
366 payload_bytes,
367 error = %err,
368 "object store operation deadline exhausted; not retrying",
369 );
370 return None;
371 };
372 remaining = deadline_remaining;
373 }
374 *retries += 1;
375 count_retry_attempt();
379 let backoff = transport_retry_backoff(&self.transport_retry, *retries).min(remaining);
380 tracing::info!(
381 object_key = key,
382 operation,
383 retry = *retries,
384 max_retries = self.transport_retry.max_retries,
385 backoff_ms = u64::try_from(backoff.as_millis()).unwrap_or(u64::MAX),
386 error = %err,
387 "transient object store write failure, backing off before retry",
388 );
389 Some(backoff)
390 }
391}
392
393struct PartReader {
399 body: ByteStream,
400 carry: Option<Bytes>,
402 part_bytes: usize,
403 exhausted: bool,
404}
405
406impl PartReader {
407 fn new(body: ByteStream, part_bytes: usize) -> Self {
408 Self {
409 body,
410 carry: None,
411 part_bytes,
412 exhausted: false,
413 }
414 }
415
416 async fn next_part(&mut self) -> Result<Option<Bytes>> {
423 let mut buffer = bytes::BytesMut::with_capacity(self.part_bytes);
424 while buffer.len() < self.part_bytes {
425 let mut chunk = match self.carry.take() {
426 Some(chunk) => chunk,
427 None if self.exhausted => break,
428 None => match self.body.next().await {
429 Some(chunk) => chunk?,
430 None => {
431 self.exhausted = true;
432 break;
433 }
434 },
435 };
436 let take = (self.part_bytes - buffer.len()).min(chunk.len());
437 buffer.extend_from_slice(&chunk.split_to(take));
438 if !chunk.is_empty() {
439 self.carry = Some(chunk);
440 }
441 }
442 Ok((!buffer.is_empty()).then(|| buffer.freeze()))
443 }
444
445 fn exhausted(&self) -> bool {
447 self.exhausted && self.carry.is_none()
448 }
449}
450
451struct AbortUploadOnDrop {
461 multipart: Option<Arc<dyn MultipartStore>>,
462 path: Path,
463 upload_id: provider_store::MultipartId,
464}
465
466impl AbortUploadOnDrop {
467 fn disarm(&mut self) {
468 self.multipart = None;
469 }
470}
471
472impl Drop for AbortUploadOnDrop {
473 fn drop(&mut self) {
474 let Some(multipart) = self.multipart.take() else {
475 return;
476 };
477 let Ok(handle) = tokio::runtime::Handle::try_current() else {
478 tracing::warn!(
479 object_key = %self.path,
480 operation = "abort_multipart",
481 "abandoned streamed write has no runtime to abort its multipart upload on; \
482 parts remain until the bucket lifecycle rule collects them",
483 );
484 return;
485 };
486 let path = self.path.clone();
487 let upload_id = std::mem::take(&mut self.upload_id);
488 handle.spawn(async move {
489 if let Err(err) = multipart.abort_multipart(&path, &upload_id).await {
490 tracing::warn!(
491 object_key = %path,
492 operation = "abort_multipart",
493 error = %err,
494 "failed to abort the multipart upload of an abandoned streamed write",
495 );
496 }
497 });
498 }
499}
500
501struct MultipartWrite<'op> {
504 store: &'op ProviderObjectStore,
505 multipart: &'op dyn MultipartStore,
506 key: &'op str,
507 path: &'op Path,
508}
509
510impl MultipartWrite<'_> {
511 async fn create(&self, payload_bytes: u64) -> Result<provider_store::MultipartId> {
512 let mut retries: u32 = 0;
513 loop {
514 let err = match self.multipart.create_multipart(self.path).await {
515 Ok(upload_id) => return Ok(upload_id),
516 Err(err) => err,
517 };
518 if !provider_transport_retryable(&err) {
519 return Err(map_provider_error(self.key, err));
520 }
521 let Some(backoff) = self.store.next_write_backoff(
522 self.key,
523 "create_multipart",
524 payload_bytes,
525 &mut retries,
526 None,
527 &err,
528 ) else {
529 return Err(map_provider_error(self.key, err));
530 };
531 transport_retry_pause(backoff).await;
532 }
533 }
534
535 async fn upload_stream_and_complete(
549 &self,
550 upload_id: &provider_store::MultipartId,
551 head: Bytes,
552 mut parts_reader: PartReader,
553 ) -> Result<u64> {
554 let mut size_bytes = head.len() as u64;
555 let mut parts = vec![self.upload_part(upload_id, 0, head).await?];
556
557 while let Some(payload) = parts_reader.next_part().await? {
558 if parts.len() >= MAX_PROVIDER_MULTIPART_PARTS {
559 return Err(ObjectStoreError::transport(
560 self.key,
561 format!(
562 "streamed payload needs more than the provider's \
563 {MAX_PROVIDER_MULTIPART_PARTS}-part limit at this part size"
564 ),
565 ));
566 }
567 size_bytes += payload.len() as u64;
568 parts.push(self.upload_part(upload_id, parts.len(), payload).await?);
569 }
570
571 match self
572 .multipart
573 .complete_multipart(self.path, upload_id, parts)
574 .await
575 {
576 Ok(_) => Ok(size_bytes),
577 Err(err) => Err(map_provider_error(self.key, err)),
578 }
579 }
580
581 async fn upload_parts_and_complete(
582 &self,
583 upload_id: &provider_store::MultipartId,
584 bytes: &Bytes,
585 ) -> Result<ObjectMetadata> {
586 let part_size = self.store.multipart_geometry.part_bytes as usize;
587 let part_count = bytes.len().div_ceil(part_size);
588 let mut part_ids: Vec<Option<PartId>> = vec![None; part_count];
589 let mut in_flight = FuturesUnordered::new();
590 let mut next_part = 0usize;
591
592 loop {
593 while in_flight.len() < PROVIDER_MULTIPART_PART_WINDOW && next_part < part_count {
594 let part_index = next_part;
595 let start = part_index * part_size;
596 let end = (start + part_size).min(bytes.len());
597 let payload = bytes.slice(start..end);
598 in_flight.push(async move {
599 let uploaded = self.upload_part(upload_id, part_index, payload).await;
600 (part_index, uploaded)
601 });
602 next_part += 1;
603 }
604 match in_flight.next().await {
605 Some((part_index, Ok(part_id))) => part_ids[part_index] = Some(part_id),
606 Some((_, Err(err))) => return Err(err),
609 None => break,
610 }
611 }
612
613 let parts = part_ids
614 .into_iter()
615 .map(|part_id| part_id.expect("every part completed before the window drained"))
616 .collect();
617 self.complete(upload_id, parts, bytes).await
618 }
619
620 async fn upload_part(
621 &self,
622 upload_id: &provider_store::MultipartId,
623 part_index: usize,
624 payload: Bytes,
625 ) -> Result<PartId> {
626 let payload_bytes = payload.len() as u64;
627 let mut retries: u32 = 0;
628 loop {
629 let err = match self
630 .multipart
631 .put_part(
632 self.path,
633 upload_id,
634 part_index,
635 PutPayload::from(payload.clone()),
636 )
637 .await
638 {
639 Ok(part_id) => return Ok(part_id),
640 Err(err) => err,
641 };
642 if !provider_transport_retryable(&err) {
643 return Err(map_provider_error(self.key, err));
644 }
645 let Some(backoff) = self.store.next_write_backoff(
646 self.key,
647 "put_part",
648 payload_bytes,
649 &mut retries,
650 None,
651 &err,
652 ) else {
653 return Err(map_provider_error(self.key, err));
654 };
655 transport_retry_pause(backoff).await;
656 }
657 }
658
659 async fn complete(
660 &self,
661 upload_id: &provider_store::MultipartId,
662 parts: Vec<PartId>,
663 bytes: &Bytes,
664 ) -> Result<ObjectMetadata> {
665 let size_bytes = bytes.len() as u64;
666 match self
667 .multipart
668 .complete_multipart(self.path, upload_id, parts)
669 .await
670 {
671 Ok(result) => Ok(ProviderObjectStore::from_put_result(result, size_bytes)),
672 Err(err) if provider_transport_retryable(&err) => {
673 self.resolve_ambiguous_completion(upload_id, bytes, err)
674 .await
675 }
676 Err(err) => Err(map_provider_error(self.key, err)),
677 }
678 }
679
680 async fn resolve_ambiguous_completion(
693 &self,
694 upload_id: &provider_store::MultipartId,
695 bytes: &Bytes,
696 final_err: provider_store::Error,
697 ) -> Result<ObjectMetadata> {
698 match readback(self.store, self.key, bytes).await {
699 Ok(ImmutableReadback::Identical(metadata)) => {
700 self.abort(upload_id).await;
701 Ok(metadata)
702 }
703 Ok(ImmutableReadback::Different) => self.unproven_completion(
704 final_err,
705 "the object at the key does not hold the payload bytes",
706 ),
707 Ok(ImmutableReadback::Missing) => {
708 self.unproven_completion(final_err, "no object exists at the key")
709 }
710 Err(verify_err) => {
711 let original = map_provider_error(self.key, final_err).message();
712 Err(ObjectStoreError::transport(
713 self.key,
714 format!(
715 "{original}; failed to verify multipart completion outcome: {verify_err}"
716 ),
717 ))
718 }
719 }
720 }
721
722 fn unproven_completion(
723 &self,
724 final_err: provider_store::Error,
725 outcome: &'static str,
726 ) -> Result<ObjectMetadata> {
727 tracing::warn!(
728 object_key = self.key,
729 operation = "complete_multipart",
730 outcome,
731 "ambiguous multipart completion did not land",
732 );
733 Err(map_provider_error(self.key, final_err))
734 }
735
736 async fn abort(&self, upload_id: &provider_store::MultipartId) {
737 match self.multipart.abort_multipart(self.path, upload_id).await {
743 Ok(()) => {}
744 Err(err) if provider_not_found(&err) => {}
745 Err(err) => {
746 tracing::warn!(
747 object_key = self.key,
748 operation = "abort_multipart",
749 error = %err,
750 "failed to abort multipart upload; parts remain until the bucket lifecycle rule collects them",
751 );
752 }
753 }
754 }
755}
756
757#[async_trait]
758impl ObjectStore for ProviderObjectStore {
759 async fn head(&self, key: &str) -> Result<Option<ObjectMetadata>> {
760 let path = self.to_path(key)?;
761 match self.inner.head(&path).await {
762 Ok(meta) => Ok(Some(Self::from_meta(meta))),
763 Err(err) if provider_not_found(&err) => Ok(None),
764 Err(err) => Err(map_provider_error(key, err)),
765 }
766 }
767
768 async fn get_with_metadata(&self, key: &str) -> Result<Option<ObjectBody>> {
769 let path = self.to_path(key)?;
770 match self.inner.get(&path).await {
771 Ok(result) => {
772 let metadata = Self::from_meta(result.meta.clone());
773 let bytes = result
774 .bytes()
775 .await
776 .map_err(|err| map_provider_error(key, err))?;
777 Ok(Some(ObjectBody {
778 metadata,
779 bytes: bytes.to_vec(),
780 }))
781 }
782 Err(err) if provider_not_found(&err) => Ok(None),
783 Err(err) => Err(map_provider_error(key, err)),
784 }
785 }
786
787 async fn get(&self, key: &str, range: Option<ByteRange>) -> Result<Option<Bytes>> {
795 let path = self.to_path(key)?;
796 let Some(range) = range else {
797 return match self.inner.get(&path).await {
798 Ok(result) => result
799 .bytes()
800 .await
801 .map(Some)
802 .map_err(|err| map_provider_error(key, err)),
803 Err(err) if provider_not_found(&err) => Ok(None),
804 Err(err) => Err(map_provider_error(key, err)),
805 };
806 };
807 if range.end_exclusive < range.start_inclusive {
808 return Err(ObjectStoreError::InvalidRange {
809 object_key: key.to_owned(),
810 });
811 }
812 if range.end_exclusive == range.start_inclusive {
813 return match self.head(key).await? {
816 None => Ok(None),
817 Some(metadata) if range.start_inclusive > metadata.size_bytes => {
818 Err(ObjectStoreError::InvalidRange {
819 object_key: key.to_owned(),
820 })
821 }
822 Some(_) => Ok(Some(Bytes::new())),
823 };
824 }
825 match self
826 .ranged_get(&path, range.start_inclusive, range.end_exclusive)
827 .await
828 {
829 RangedGet::Bytes(bytes) => Ok(Some(bytes)),
830 RangedGet::NotFound => Ok(None),
831 RangedGet::Refused(err) => {
832 match self.head(key).await? {
835 None => Ok(None),
836 Some(metadata) if range.start_inclusive > metadata.size_bytes => {
837 Err(ObjectStoreError::InvalidRange {
838 object_key: key.to_owned(),
839 })
840 }
841 Some(metadata) if range.start_inclusive == metadata.size_bytes => {
842 Ok(Some(Bytes::new()))
843 }
844 Some(metadata) if range.end_exclusive > metadata.size_bytes => {
845 match self
848 .ranged_get(&path, range.start_inclusive, metadata.size_bytes)
849 .await
850 {
851 RangedGet::Bytes(bytes) => Ok(Some(bytes)),
852 RangedGet::NotFound => Ok(None),
853 RangedGet::Refused(err) => Err(map_provider_error(key, err)),
854 }
855 }
856 Some(_) => Err(map_provider_error(key, err)),
857 }
858 }
859 }
860 }
861
862 async fn put(&self, key: &str, bytes: Bytes, mode: PutMode) -> Result<ObjectMetadata> {
863 let path = self.to_path(key)?;
864 let size_bytes = bytes.len() as u64;
865 if matches!(mode, PutMode::Overwrite)
866 && size_bytes >= self.multipart_geometry.threshold_bytes
867 {
868 if let Some(multipart) = self.multipart.clone() {
869 return self
870 .put_large_multipart(multipart.as_ref(), key, &path, bytes)
871 .await;
872 }
873 }
874 let compare_and_swap = matches!(mode, PutMode::CompareAndSwap { .. });
879 let options = PutOptions {
880 mode: map_put_mode(mode),
881 ..Default::default()
882 };
883 match self
884 .inner
885 .put_opts(&path, PutPayload::from(bytes), options)
886 .await
887 {
888 Ok(result) => Ok(Self::from_put_result(result, size_bytes)),
889 Err(err) if compare_and_swap && provider_not_found(&err) => {
890 Err(ObjectStoreError::PreconditionFailed {
891 object_key: key.to_owned(),
892 })
893 }
894 Err(err) => Err(map_provider_error(key, err)),
895 }
896 }
897
898 async fn put_streamed(&self, key: &str, body: ByteStream, mode: PutMode) -> Result<u64> {
909 let path = self.to_path(key)?;
910 let mut reader = PartReader::new(body, self.multipart_geometry.part_bytes as usize);
911 let head = reader.next_part().await?.unwrap_or_else(Bytes::new);
912 let Some(multipart) = self.multipart.clone() else {
913 let mut bytes = bytes::BytesMut::from(head.as_ref());
916 while let Some(part) = reader.next_part().await? {
917 bytes.extend_from_slice(&part);
918 }
919 let bytes = bytes.freeze();
920 let size_bytes = bytes.len() as u64;
921 self.put(key, bytes, mode).await?;
922 return Ok(size_bytes);
923 };
924 if reader.exhausted() {
925 let size_bytes = head.len() as u64;
926 self.put(key, head, mode).await?;
927 return Ok(size_bytes);
928 }
929
930 let upload = MultipartWrite {
931 store: self,
932 multipart: multipart.as_ref(),
933 key,
934 path: &path,
935 };
936 let upload_id = upload.create(0).await?;
937 let mut abort_on_drop = AbortUploadOnDrop {
938 multipart: Some(Arc::clone(&multipart)),
939 path: path.clone(),
940 upload_id: upload_id.clone(),
941 };
942 let result = upload
943 .upload_stream_and_complete(&upload_id, head, reader)
944 .await;
945 abort_on_drop.disarm();
946 match result {
947 Ok(size_bytes) => Ok(size_bytes),
948 Err(err) => {
949 upload.abort(&upload_id).await;
952 Err(err)
953 }
954 }
955 }
956
957 async fn delete(&self, key: &str) -> Result<()> {
958 let path = self.to_path(key)?;
959 let deadline =
960 OperationDeadline::start(self.timer.as_ref(), self.transport_retry.operation_deadline);
961 let mut retries: u32 = 0;
962 loop {
963 let err = match self.inner.delete(&path).await {
967 Ok(()) => return Ok(()),
968 Err(err) if provider_not_found(&err) => return Ok(()),
969 Err(err) => err,
970 };
971 if !provider_transport_retryable(&err) {
972 return Err(map_provider_error(key, err));
973 }
974 let Some(backoff) =
975 self.next_write_backoff(key, "delete", 0, &mut retries, Some(&deadline), &err)
976 else {
977 return Err(map_provider_error(key, err));
978 };
979 transport_retry_pause(backoff).await;
980 }
981 }
982
983 fn list_prefix_stream(&self, prefix: &str) -> BoxStream<'static, Result<String>> {
984 let prefix_path = match self.list_path(prefix) {
985 Ok(prefix_path) => prefix_path,
986 Err(err) => return stream::once(async { Err(err) }).boxed(),
987 };
988 let key_prefix = self.key_prefix.clone();
989 let listed_prefix = prefix.to_owned();
990 self.inner
991 .list(prefix_path.as_ref())
992 .filter_map(move |result| {
993 let key_prefix = key_prefix.clone();
994 let listed_prefix = listed_prefix.clone();
995 async move {
996 match result {
997 Ok(meta) => {
998 let key = meta.location.as_ref();
999 match key_prefix.as_deref() {
1000 Some(prefix) => unscope_listed_key(Some(prefix), key).map(Ok),
1001 None => Some(Ok(key.to_owned())),
1002 }
1003 }
1004 Err(err) => Some(Err(map_provider_error(&listed_prefix, err))),
1005 }
1006 }
1007 })
1008 .boxed()
1009 }
1010}
1011
1012fn map_put_mode(mode: PutMode) -> provider_store::PutMode {
1013 match mode {
1014 PutMode::Overwrite => provider_store::PutMode::Overwrite,
1015 PutMode::CreateIfAbsent => provider_store::PutMode::Create,
1016 PutMode::CompareAndSwap { expected_etag } => {
1017 provider_store::PutMode::Update(UpdateVersion {
1022 e_tag: Some(expected_etag.clone()),
1023 version: Some(expected_etag),
1024 })
1025 }
1026 }
1027}
1028
1029enum RangedGet {
1030 Bytes(Bytes),
1031 NotFound,
1032 Refused(provider_store::Error),
1033}
1034
1035fn provider_not_found(err: &provider_store::Error) -> bool {
1036 matches!(err, provider_store::Error::NotFound { .. })
1037}
1038
1039fn provider_transport_retryable(err: &provider_store::Error) -> bool {
1040 matches!(err, provider_store::Error::Generic { .. })
1048}
1049
1050fn map_provider_error(object_key: &str, err: provider_store::Error) -> ObjectStoreError {
1051 match err {
1052 provider_store::Error::NotFound { .. } => ObjectStoreError::NotFound {
1053 object_key: object_key.to_owned(),
1054 },
1055 provider_store::Error::AlreadyExists { .. }
1056 | provider_store::Error::Precondition { .. }
1057 | provider_store::Error::NotModified { .. } => ObjectStoreError::PreconditionFailed {
1058 object_key: object_key.to_owned(),
1059 },
1060 provider_store::Error::InvalidPath { source } => ObjectStoreError::InvalidKey {
1061 object_key: object_key.to_owned(),
1062 message: source.to_string(),
1063 },
1064 provider_store::Error::NotSupported { .. } | provider_store::Error::NotImplemented => {
1065 ObjectStoreError::Unsupported("provider object store operation")
1066 }
1067 provider_store::Error::UnknownConfigurationKey { key, store } => {
1068 ObjectStoreError::transport(
1069 object_key,
1070 format!("unknown {store} configuration key `{key}`"),
1071 )
1072 }
1073 provider_store::Error::Generic { source, .. } => {
1074 ObjectStoreError::transport(object_key, sanitize_provider_message(&source.to_string()))
1075 }
1076 provider_store::Error::JoinError { source } => {
1077 ObjectStoreError::transport(object_key, sanitize_provider_message(&source.to_string()))
1078 }
1079 provider_store::Error::PermissionDenied { source, .. }
1080 | provider_store::Error::Unauthenticated { source, .. } => {
1081 ObjectStoreError::PermissionDenied {
1082 object_key: object_key.to_owned(),
1083 message: sanitize_provider_message(&source.to_string()),
1084 }
1085 }
1086 other => {
1087 ObjectStoreError::transport(object_key, sanitize_provider_message(&other.to_string()))
1088 }
1089 }
1090}
1091
1092const CREDENTIAL_QUERY_PARAMS: &[&str] = &[
1095 "X-Amz-Signature",
1096 "X-Amz-Credential",
1097 "X-Amz-Security-Token",
1098 "AWSAccessKeyId",
1099 "Signature",
1100 "sig",
1101];
1102
1103const CREDENTIAL_XML_ELEMENTS: &[&str] = &[
1106 "StringToSign",
1107 "StringToSignBytes",
1108 "CanonicalRequest",
1109 "SignatureProvided",
1110 "AWSAccessKeyId",
1111];
1112
1113fn sanitize_provider_message(message: &str) -> String {
1118 let mut sanitized = message.to_owned();
1119 for param in CREDENTIAL_QUERY_PARAMS {
1120 sanitized = mask_query_param_values(&sanitized, param);
1121 }
1122 for element in CREDENTIAL_XML_ELEMENTS {
1123 sanitized = mask_xml_element_text(&sanitized, element);
1124 }
1125 sanitized
1126}
1127
1128fn mask_query_param_values(message: &str, param: &str) -> String {
1132 let needle = format!("{param}=");
1133 let mut out = String::with_capacity(message.len());
1134 let mut cursor = 0;
1135 while let Some(found) = message[cursor..].find(&needle) {
1136 let start = cursor + found;
1137 let value_start = start + needle.len();
1138 out.push_str(&message[cursor..value_start]);
1139 cursor = value_start;
1140 let at_boundary = start > 0 && matches!(message.as_bytes()[start - 1], b'?' | b'&');
1141 if at_boundary {
1142 let value_len = message[cursor..]
1143 .find(|c: char| {
1144 matches!(c, '&' | '"' | '\'' | ')' | '<' | '>' | ':' | ',') || c.is_whitespace()
1145 })
1146 .unwrap_or(message.len() - cursor);
1147 out.push_str("<redacted>");
1148 cursor += value_len;
1149 }
1150 }
1151 out.push_str(&message[cursor..]);
1152 out
1153}
1154
1155fn mask_xml_element_text(message: &str, element: &str) -> String {
1159 let open = format!("<{element}>");
1160 let close = format!("</{element}>");
1161 let mut out = String::with_capacity(message.len());
1162 let mut rest = message;
1163 while let Some(position) = rest.find(&open) {
1164 let text_start = position + open.len();
1165 out.push_str(&rest[..text_start]);
1166 out.push_str("<redacted>");
1167 rest = &rest[text_start..];
1168 match rest.find(&close) {
1169 Some(text_end) => rest = &rest[text_end..],
1170 None => rest = "",
1171 }
1172 }
1173 out.push_str(rest);
1174 out
1175}
1176
1177#[cfg(test)]
1178mod tests {
1179 use super::*;
1180 use crate::metrics::{
1181 InstrumentedObjectStore, ObjectStoreOperation, VecObjectStoreMetricsRecorder,
1182 };
1183 use crate::test_support::SteppingTimer;
1184 use futures::StreamExt;
1185 use object_store::memory::InMemory;
1186
1187 fn memory_store() -> ProviderObjectStore {
1188 let inner = Arc::new(InMemory::default());
1189 ProviderObjectStore::new(
1190 Arc::clone(&inner) as Arc<dyn provider_store::ObjectStore>,
1191 Some(inner),
1192 ProviderObjectStoreConfig {
1193 key_prefix: Some("tenant-a".to_owned()),
1194 },
1195 )
1196 .expect("provider store")
1197 }
1198
1199 #[test]
1200 fn provider_messages_drop_credential_material_and_keep_the_diagnosis() {
1201 let presigned = sanitize_provider_message(
1202 "Generic S3 error: error sending request for url \
1203 (https://bucket.s3.amazonaws.com/k?X-Amz-Algorithm=AWS4-HMAC-SHA256\
1204 &X-Amz-Credential=AKIAIOSFODNN7EXAMPLE%2F20260726%2Fus-east-1%2Fs3%2Faws4_request\
1205 &X-Amz-Signature=deadbeefcafe): operation timed out",
1206 );
1207 assert!(!presigned.contains("AKIAIOSFODNN7EXAMPLE"), "{presigned}");
1208 assert!(!presigned.contains("deadbeefcafe"), "{presigned}");
1209 assert!(
1210 presigned.contains("X-Amz-Signature=<redacted>"),
1211 "{presigned}"
1212 );
1213 assert!(presigned.contains("operation timed out"), "{presigned}");
1214 assert!(
1215 presigned.contains("bucket.s3.amazonaws.com/k"),
1216 "{presigned}"
1217 );
1218
1219 let signature_mismatch = sanitize_provider_message(
1220 "Client error with status 403 Forbidden: <Error>\
1221 <Code>SignatureDoesNotMatch</Code>\
1222 <StringToSign>AWS4-HMAC-SHA256 20260726T000000Z scope digest</StringToSign>\
1223 <SignatureProvided>cafe0123</SignatureProvided>\
1224 <AWSAccessKeyId>AKIAIOSFODNN7EXAMPLE</AWSAccessKeyId></Error>",
1225 );
1226 assert!(
1227 !signature_mismatch.contains("AKIAIOSFODNN7EXAMPLE"),
1228 "{signature_mismatch}"
1229 );
1230 assert!(
1231 !signature_mismatch.contains("cafe0123"),
1232 "{signature_mismatch}"
1233 );
1234 assert!(
1235 !signature_mismatch.contains("20260726T000000Z"),
1236 "{signature_mismatch}"
1237 );
1238 assert!(
1239 signature_mismatch.contains("SignatureDoesNotMatch"),
1240 "{signature_mismatch}"
1241 );
1242 assert!(
1243 signature_mismatch.contains("403 Forbidden"),
1244 "{signature_mismatch}"
1245 );
1246
1247 let azure_sas = sanitize_provider_message(
1248 "error for url https://account.blob.core.windows.net/c/k?sv=2021-08-06\
1249 &se=2026-07-26&sig=aGVsbG8: 403",
1250 );
1251 assert!(!azure_sas.contains("aGVsbG8"), "{azure_sas}");
1252 assert!(azure_sas.contains("sig=<redacted>"), "{azure_sas}");
1253
1254 let unrelated = sanitize_provider_message("policy?ResponseSignature=keep&sigil=keep2");
1256 assert!(unrelated.contains("keep2"), "{unrelated}");
1257 assert!(!unrelated.contains("Signature=<redacted>"), "{unrelated}");
1258 }
1259
1260 #[tokio::test]
1261 async fn provider_store_preserves_put_get_head_and_prefix_scoping() {
1262 let store = memory_store();
1263 let key = "namespaces/demo/wal/head.json";
1264
1265 let metadata = store
1266 .put_if_absent(key, Bytes::from_static(b"head"))
1267 .await
1268 .expect("put");
1269 assert_eq!(metadata.size_bytes, 4);
1270 assert!(metadata.etag.is_some());
1271
1272 let head = store.head(key).await.expect("head").expect("head exists");
1273 assert_eq!(head.size_bytes, 4);
1274 assert_eq!(
1275 store.get(key, None).await.expect("get"),
1276 Some(Bytes::from_static(b"head"))
1277 );
1278 assert_eq!(
1279 store.list_prefix("namespaces/demo/").await.expect("list"),
1280 vec![key.to_owned()]
1281 );
1282 }
1283
1284 #[tokio::test]
1285 async fn ranged_reads_match_the_reference_contract_in_one_round_trip() {
1286 let store = memory_store();
1287 let key = "namespaces/demo/metadata/tables/tbl_abc.sst.zst";
1288 store
1289 .put_if_absent(key, Bytes::from_static(b"0123456789"))
1290 .await
1291 .expect("put");
1292
1293 let range = |start, end| {
1294 Some(ByteRange {
1295 start_inclusive: start,
1296 end_exclusive: end,
1297 })
1298 };
1299
1300 assert_eq!(
1302 store.get(key, range(2, 6)).await.expect("bounded"),
1303 Some(Bytes::from_static(b"2345"))
1304 );
1305 assert_eq!(
1307 store.get(key, range(6, 99)).await.expect("clamped"),
1308 Some(Bytes::from_static(b"6789"))
1309 );
1310 assert_eq!(
1312 store.get(key, range(10, 12)).await.expect("at end"),
1313 Some(Bytes::new())
1314 );
1315 assert!(matches!(
1317 store.get(key, range(11, 12)).await,
1318 Err(ObjectStoreError::InvalidRange { .. })
1319 ));
1320 assert!(matches!(
1322 store.get(key, range(6, 2)).await,
1323 Err(ObjectStoreError::InvalidRange { .. })
1324 ));
1325 assert_eq!(
1327 store.get(key, range(4, 4)).await.expect("zero length"),
1328 Some(Bytes::new())
1329 );
1330
1331 let missing = "namespaces/demo/metadata/tables/tbl_missing.sst.zst";
1334 assert_eq!(
1335 store.get(missing, range(0, 4)).await.expect("missing"),
1336 None
1337 );
1338 assert_eq!(
1339 store.get(missing, range(3, 3)).await.expect("missing zero"),
1340 None
1341 );
1342 }
1343
1344 #[tokio::test]
1345 async fn provider_store_enforces_create_and_cas_preconditions() {
1346 let store = memory_store();
1347 let key = "namespaces/demo/wal/head.json";
1348 let first = store
1349 .put_if_absent(key, Bytes::from_static(b"one"))
1350 .await
1351 .expect("first put");
1352
1353 assert!(matches!(
1354 store.put_if_absent(key, Bytes::from_static(b"two")).await,
1355 Err(ObjectStoreError::PreconditionFailed { .. })
1356 ));
1357 assert!(matches!(
1358 store
1359 .compare_and_swap(key, "stale", Bytes::from_static(b"two"))
1360 .await,
1361 Err(ObjectStoreError::PreconditionFailed { .. })
1362 ));
1363 assert!(matches!(
1364 store
1365 .compare_and_swap(
1366 "namespaces/demo/control/missing-head.json",
1367 "missing",
1368 Bytes::from_static(b"two")
1369 )
1370 .await,
1371 Err(ObjectStoreError::PreconditionFailed { .. })
1372 ));
1373 let etag = first.etag.expect("etag");
1374 store
1375 .compare_and_swap(key, &etag, Bytes::from_static(b"two"))
1376 .await
1377 .expect("cas");
1378 assert_eq!(
1379 store.get(key, None).await.expect("get"),
1380 Some(Bytes::from_static(b"two"))
1381 );
1382 }
1383
1384 #[tokio::test]
1385 async fn provider_store_range_semantics_match_blocking_contract() {
1386 let store = memory_store();
1387 let key = "content-stores/cs_0123456789abcdef0123456789abcdef/objects/ab/con_abcdef0123456789abcdef0123456789";
1388 store
1389 .put_overwrite(key, Bytes::from_static(b"abcdef"))
1390 .await
1391 .expect("put");
1392
1393 assert_eq!(
1394 store
1395 .get(
1396 key,
1397 Some(ByteRange {
1398 start_inclusive: 2,
1399 end_exclusive: 4,
1400 }),
1401 )
1402 .await
1403 .expect("range"),
1404 Some(Bytes::from_static(b"cd"))
1405 );
1406 assert_eq!(
1407 store
1408 .get(
1409 key,
1410 Some(ByteRange {
1411 start_inclusive: 6,
1412 end_exclusive: 10,
1413 }),
1414 )
1415 .await
1416 .expect("empty"),
1417 Some(Bytes::new())
1418 );
1419 assert!(matches!(
1420 store
1421 .get(
1422 key,
1423 Some(ByteRange {
1424 start_inclusive: 7,
1425 end_exclusive: 8,
1426 }),
1427 )
1428 .await,
1429 Err(ObjectStoreError::InvalidRange { .. })
1430 ));
1431 }
1432
1433 #[tokio::test]
1434 async fn provider_stream_reports_invalid_prefix() {
1435 let store = memory_store();
1436 let mut stream = store.list_prefix_stream("../");
1437 assert!(matches!(
1438 stream.next().await,
1439 Some(Err(ObjectStoreError::InvalidKey { .. }))
1440 ));
1441 }
1442
1443 use provider_store::{GetResult, ListResult, MultipartUpload, PutMultipartOptions};
1444 use std::collections::{BTreeMap, HashMap, VecDeque};
1445 use std::sync::atomic::{AtomicUsize, Ordering};
1446 use std::sync::Mutex;
1447
1448 #[derive(Debug)]
1449 enum WriteScript {
1450 FailWithoutLanding,
1451 LandThenFail,
1452 FailAuth,
1453 VanishThenFail,
1457 }
1458
1459 #[derive(Debug)]
1460 enum ReadScript {
1461 Transport,
1462 }
1463
1464 #[derive(Default)]
1472 struct FlakyStore {
1473 inner: InMemory,
1474 put_script: Mutex<VecDeque<WriteScript>>,
1475 get_script: Mutex<VecDeque<ReadScript>>,
1476 delete_script: Mutex<VecDeque<WriteScript>>,
1477 part_script: Mutex<HashMap<usize, VecDeque<WriteScript>>>,
1478 complete_script: Mutex<VecDeque<WriteScript>>,
1479 puts: AtomicUsize,
1480 gets: AtomicUsize,
1481 deletes: AtomicUsize,
1482 multipart_creates: AtomicUsize,
1483 part_attempts: Mutex<HashMap<usize, usize>>,
1484 multipart_completes: AtomicUsize,
1485 multipart_aborts: AtomicUsize,
1486 next_upload_id: AtomicUsize,
1487 multipart_uploads: Mutex<HashMap<String, BTreeMap<usize, Bytes>>>,
1488 }
1489
1490 impl fmt::Debug for FlakyStore {
1491 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1492 f.debug_struct("FlakyStore").finish_non_exhaustive()
1493 }
1494 }
1495
1496 impl fmt::Display for FlakyStore {
1497 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1498 write!(f, "FlakyStore")
1499 }
1500 }
1501
1502 fn transport_glitch() -> provider_store::Error {
1503 provider_store::Error::Generic {
1504 store: "flaky",
1505 source: "error sending request".into(),
1506 }
1507 }
1508
1509 fn auth_rejection(location: &Path) -> provider_store::Error {
1510 provider_store::Error::PermissionDenied {
1511 path: location.to_string(),
1512 source: "access denied".into(),
1513 }
1514 }
1515
1516 #[async_trait]
1517 impl provider_store::ObjectStore for FlakyStore {
1518 async fn put_opts(
1519 &self,
1520 location: &Path,
1521 payload: PutPayload,
1522 opts: PutOptions,
1523 ) -> provider_store::Result<PutResult> {
1524 self.puts.fetch_add(1, Ordering::SeqCst);
1525 let script = self.put_script.lock().expect("put script").pop_front();
1526 match script {
1527 Some(WriteScript::FailWithoutLanding) => Err(transport_glitch()),
1528 Some(WriteScript::LandThenFail) => {
1529 self.inner.put_opts(location, payload, opts).await?;
1530 Err(transport_glitch())
1531 }
1532 Some(WriteScript::FailAuth) => Err(auth_rejection(location)),
1533 Some(WriteScript::VanishThenFail) => {
1534 unreachable!("VanishThenFail is a completion script")
1535 }
1536 None => self.inner.put_opts(location, payload, opts).await,
1537 }
1538 }
1539
1540 async fn put_multipart_opts(
1541 &self,
1542 location: &Path,
1543 opts: PutMultipartOptions,
1544 ) -> provider_store::Result<Box<dyn MultipartUpload>> {
1545 self.inner.put_multipart_opts(location, opts).await
1546 }
1547
1548 async fn get_opts(
1549 &self,
1550 location: &Path,
1551 options: GetOptions,
1552 ) -> provider_store::Result<GetResult> {
1553 self.gets.fetch_add(1, Ordering::SeqCst);
1554 let script = self.get_script.lock().expect("get script").pop_front();
1555 match script {
1556 Some(ReadScript::Transport) => Err(transport_glitch()),
1557 None => self.inner.get_opts(location, options).await,
1558 }
1559 }
1560
1561 async fn delete(&self, location: &Path) -> provider_store::Result<()> {
1562 self.deletes.fetch_add(1, Ordering::SeqCst);
1563 let script = self
1564 .delete_script
1565 .lock()
1566 .expect("delete script")
1567 .pop_front();
1568 match script {
1569 Some(WriteScript::FailWithoutLanding) => Err(transport_glitch()),
1570 Some(WriteScript::LandThenFail) => {
1571 self.inner.delete(location).await?;
1572 Err(transport_glitch())
1573 }
1574 Some(WriteScript::FailAuth) => Err(auth_rejection(location)),
1575 Some(WriteScript::VanishThenFail) => {
1576 unreachable!("VanishThenFail is a completion script")
1577 }
1578 None => self.inner.delete(location).await,
1579 }
1580 }
1581
1582 fn list(
1583 &self,
1584 prefix: Option<&Path>,
1585 ) -> BoxStream<'static, provider_store::Result<ObjectMeta>> {
1586 self.inner.list(prefix)
1587 }
1588
1589 async fn list_with_delimiter(
1590 &self,
1591 prefix: Option<&Path>,
1592 ) -> provider_store::Result<ListResult> {
1593 self.inner.list_with_delimiter(prefix).await
1594 }
1595
1596 async fn copy(&self, from: &Path, to: &Path) -> provider_store::Result<()> {
1597 self.inner.copy(from, to).await
1598 }
1599
1600 async fn copy_if_not_exists(&self, from: &Path, to: &Path) -> provider_store::Result<()> {
1601 self.inner.copy_if_not_exists(from, to).await
1602 }
1603 }
1604
1605 impl FlakyStore {
1606 fn store_part(
1607 &self,
1608 id: &provider_store::MultipartId,
1609 part_idx: usize,
1610 data: PutPayload,
1611 ) -> provider_store::Result<()> {
1612 let mut uploads = self.multipart_uploads.lock().expect("uploads");
1613 let upload =
1614 uploads
1615 .get_mut(id.as_str())
1616 .ok_or_else(|| provider_store::Error::NotFound {
1617 path: id.clone(),
1618 source: "no such upload".into(),
1619 })?;
1620 upload.insert(part_idx, Bytes::from(data));
1621 Ok(())
1622 }
1623
1624 async fn land_completion(
1625 &self,
1626 path: &Path,
1627 id: &provider_store::MultipartId,
1628 parts: &[PartId],
1629 ) -> provider_store::Result<PutResult> {
1630 let upload = self
1631 .multipart_uploads
1632 .lock()
1633 .expect("uploads")
1634 .remove(id.as_str())
1635 .ok_or_else(|| provider_store::Error::NotFound {
1636 path: id.clone(),
1637 source: "no such upload".into(),
1638 })?;
1639 assert_eq!(
1640 upload.len(),
1641 parts.len(),
1642 "completion must list exactly the uploaded parts"
1643 );
1644 let mut buf = Vec::new();
1645 for part in upload.values() {
1646 buf.extend_from_slice(part);
1647 }
1648 provider_store::ObjectStore::put_opts(
1649 &self.inner,
1650 path,
1651 buf.into(),
1652 PutOptions::default(),
1653 )
1654 .await
1655 }
1656 }
1657
1658 #[async_trait]
1659 impl MultipartStore for FlakyStore {
1660 async fn create_multipart(
1661 &self,
1662 _path: &Path,
1663 ) -> provider_store::Result<provider_store::MultipartId> {
1664 self.multipart_creates.fetch_add(1, Ordering::SeqCst);
1665 let id = self
1666 .next_upload_id
1667 .fetch_add(1, Ordering::SeqCst)
1668 .to_string();
1669 self.multipart_uploads
1670 .lock()
1671 .expect("uploads")
1672 .insert(id.clone(), BTreeMap::new());
1673 Ok(id)
1674 }
1675
1676 async fn put_part(
1677 &self,
1678 path: &Path,
1679 id: &provider_store::MultipartId,
1680 part_idx: usize,
1681 data: PutPayload,
1682 ) -> provider_store::Result<PartId> {
1683 *self
1684 .part_attempts
1685 .lock()
1686 .expect("part attempts")
1687 .entry(part_idx)
1688 .or_default() += 1;
1689 let script = self
1690 .part_script
1691 .lock()
1692 .expect("part script")
1693 .get_mut(&part_idx)
1694 .and_then(VecDeque::pop_front);
1695 match script {
1696 Some(WriteScript::FailWithoutLanding) => Err(transport_glitch()),
1697 Some(WriteScript::LandThenFail) => {
1698 self.store_part(id, part_idx, data)?;
1699 Err(transport_glitch())
1700 }
1701 Some(WriteScript::FailAuth) => Err(auth_rejection(path)),
1702 Some(WriteScript::VanishThenFail) => {
1703 unreachable!("VanishThenFail is a completion script")
1704 }
1705 None => {
1706 self.store_part(id, part_idx, data)?;
1707 Ok(PartId {
1708 content_id: part_idx.to_string(),
1709 })
1710 }
1711 }
1712 }
1713
1714 async fn complete_multipart(
1715 &self,
1716 path: &Path,
1717 id: &provider_store::MultipartId,
1718 parts: Vec<PartId>,
1719 ) -> provider_store::Result<PutResult> {
1720 self.multipart_completes.fetch_add(1, Ordering::SeqCst);
1721 let script = self
1722 .complete_script
1723 .lock()
1724 .expect("complete script")
1725 .pop_front();
1726 match script {
1727 Some(WriteScript::FailWithoutLanding) => Err(transport_glitch()),
1728 Some(WriteScript::LandThenFail) => {
1729 self.land_completion(path, id, &parts).await?;
1730 Err(transport_glitch())
1731 }
1732 Some(WriteScript::FailAuth) => Err(auth_rejection(path)),
1733 Some(WriteScript::VanishThenFail) => {
1734 self.multipart_uploads
1735 .lock()
1736 .expect("uploads")
1737 .remove(id.as_str());
1738 Err(transport_glitch())
1739 }
1740 None => self.land_completion(path, id, &parts).await,
1741 }
1742 }
1743
1744 async fn abort_multipart(
1745 &self,
1746 _path: &Path,
1747 id: &provider_store::MultipartId,
1748 ) -> provider_store::Result<()> {
1749 self.multipart_aborts.fetch_add(1, Ordering::SeqCst);
1750 self.multipart_uploads
1751 .lock()
1752 .expect("uploads")
1753 .remove(id.as_str());
1754 Ok(())
1755 }
1756 }
1757
1758 fn retrying_store(flaky: Arc<FlakyStore>) -> ProviderObjectStore {
1759 ProviderObjectStore::new(
1760 Arc::clone(&flaky) as Arc<dyn provider_store::ObjectStore>,
1761 Some(flaky),
1762 ProviderObjectStoreConfig {
1763 key_prefix: Some("tenant-a".to_owned()),
1764 },
1765 )
1766 .expect("provider store")
1767 .transport_retry(TransportRetryPolicy {
1768 max_retries: 4,
1769 initial_backoff: Duration::from_millis(1),
1770 max_backoff: Duration::from_millis(1),
1771 operation_deadline: PROVIDER_OPERATION_DEADLINE,
1772 })
1773 }
1774
1775 fn script_puts(flaky: &FlakyStore, script: impl IntoIterator<Item = WriteScript>) {
1776 flaky.put_script.lock().expect("put script").extend(script);
1777 }
1778
1779 #[test]
1780 fn transport_retry_backoff_doubles_and_caps() {
1781 let policy = TransportRetryPolicy::DEFAULT;
1782 assert_eq!(
1783 transport_retry_backoff(&policy, 1),
1784 Duration::from_millis(100)
1785 );
1786 assert_eq!(
1787 transport_retry_backoff(&policy, 2),
1788 Duration::from_millis(200)
1789 );
1790 assert_eq!(
1791 transport_retry_backoff(&policy, 8),
1792 Duration::from_millis(12_800)
1793 );
1794 assert_eq!(transport_retry_backoff(&policy, 9), Duration::from_secs(15));
1795 assert_eq!(
1796 transport_retry_backoff(&policy, 10),
1797 Duration::from_secs(15)
1798 );
1799 }
1800
1801 #[tokio::test]
1802 async fn mutable_overwrite_transport_failure_is_not_retried() {
1803 let flaky = Arc::new(FlakyStore::default());
1804 let store = retrying_store(Arc::clone(&flaky));
1805 script_puts(&flaky, [WriteScript::FailWithoutLanding]);
1806 let key = "namespaces/demo/uploads/upl_1.json";
1807
1808 let error = store
1809 .put_overwrite(key, Bytes::from_static(b"session"))
1810 .await
1811 .expect_err("mutable overwrite must surface an ambiguous outcome");
1812
1813 assert!(matches!(error, ObjectStoreError::Transport { .. }));
1814 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
1815 }
1816
1817 #[tokio::test]
1818 async fn immutable_small_write_uses_create_if_absent() {
1819 let flaky = Arc::new(FlakyStore::default());
1820 let store = retrying_store(Arc::clone(&flaky));
1821 let key = "namespaces/demo/uploads/upl_2.json";
1822
1823 store
1824 .put_immutable_verified(key, Bytes::from_static(b"immutable bytes"))
1825 .await
1826 .expect("small immutable write");
1827
1828 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
1829 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 0);
1830 assert_eq!(
1831 store.get(key, None).await.expect("get"),
1832 Some(Bytes::from_static(b"immutable bytes"))
1833 );
1834 }
1835
1836 #[tokio::test]
1837 async fn immutable_already_present_identical_is_accepted_without_rewrite() {
1838 let flaky = Arc::new(FlakyStore::default());
1839 let store = retrying_store(Arc::clone(&flaky));
1840 let key = "namespaces/demo/wal/00000001.cbor.zst";
1841 let bytes = Bytes::from_static(b"identical immutable bytes");
1842 seed_scoped_object(&flaky, key, bytes.clone()).await;
1843
1844 store
1845 .put_immutable_verified(key, bytes)
1846 .await
1847 .expect("identical object is accepted");
1848
1849 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
1850 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
1851 }
1852
1853 #[tokio::test]
1854 async fn immutable_different_bytes_at_key_are_corruption_class() {
1855 let flaky = Arc::new(FlakyStore::default());
1856 let store = retrying_store(Arc::clone(&flaky));
1857 let key = "namespaces/demo/wal/00000002.cbor.zst";
1858 seed_scoped_object(&flaky, key, Bytes::from_static(b"theirs")).await;
1859
1860 let error = store
1861 .put_immutable_verified(key, Bytes::from_static(b"mine"))
1862 .await
1863 .expect_err("different immutable bytes are rejected");
1864
1865 assert!(matches!(
1866 error,
1867 crate::ImmutableWriteError::DifferentObject { object_key } if object_key == key
1868 ));
1869 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
1870 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
1871 }
1872
1873 #[tokio::test(start_paused = true)]
1874 async fn immutable_ambiguous_landed_write_is_accepted_by_readback() {
1875 let flaky = Arc::new(FlakyStore::default());
1876 let store = retrying_store(Arc::clone(&flaky));
1877 script_puts(&flaky, [WriteScript::LandThenFail]);
1878 let key = "namespaces/demo/uploads/upl_3.json";
1879
1880 store
1881 .put_immutable_verified(key, Bytes::from_static(b"payload"))
1882 .await
1883 .expect("readback proves the first attempt landed");
1884
1885 assert_eq!(flaky.puts.load(Ordering::SeqCst), 2);
1886 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
1887 }
1888
1889 #[tokio::test(start_paused = true)]
1890 async fn immutable_ambiguous_outcome_rejects_different_readback() {
1891 let flaky = Arc::new(FlakyStore::default());
1892 let store = retrying_store(Arc::clone(&flaky));
1893 let key = "namespaces/demo/wal/00000003.cbor.zst";
1894 seed_scoped_object(&flaky, key, Bytes::from_static(b"theirs")).await;
1895 script_puts(&flaky, [WriteScript::FailWithoutLanding]);
1896
1897 let error = store
1898 .put_immutable_verified(key, Bytes::from_static(b"mine"))
1899 .await
1900 .expect_err("ambiguous write cannot adopt different bytes");
1901
1902 assert!(matches!(
1903 error,
1904 crate::ImmutableWriteError::DifferentObject { .. }
1905 ));
1906 assert_eq!(flaky.puts.load(Ordering::SeqCst), 2);
1907 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
1908 }
1909
1910 #[tokio::test(start_paused = true)]
1911 async fn immutable_transport_failures_retry_inside_the_operation() {
1912 let flaky = Arc::new(FlakyStore::default());
1913 let store = retrying_store(Arc::clone(&flaky));
1914 script_puts(
1915 &flaky,
1916 [
1917 WriteScript::FailWithoutLanding,
1918 WriteScript::FailWithoutLanding,
1919 ],
1920 );
1921 let key = "namespaces/demo/uploads/upl_4.json";
1922
1923 store
1924 .put_immutable_verified(key, Bytes::from_static(b"payload"))
1925 .await
1926 .expect("transient failures are retried");
1927
1928 assert_eq!(flaky.puts.load(Ordering::SeqCst), 3);
1929 assert_eq!(flaky.gets.load(Ordering::SeqCst), 0);
1930 }
1931
1932 #[tokio::test(start_paused = true)]
1933 async fn immutable_transport_failure_surfaces_after_the_retry_budget() {
1934 let flaky = Arc::new(FlakyStore::default());
1935 let store = retrying_store(Arc::clone(&flaky));
1936 script_puts(&flaky, (0..11).map(|_| WriteScript::FailWithoutLanding));
1937 let key = "namespaces/demo/uploads/upl_9.json";
1938
1939 let error = store
1940 .put_immutable_verified(key, Bytes::from_static(b"payload"))
1941 .await
1942 .expect_err("persistent failure exhausts the immutable retry budget");
1943
1944 assert!(matches!(
1945 error,
1946 crate::ImmutableWriteError::Transport {
1947 source: ObjectStoreError::Transport { .. },
1948 ..
1949 }
1950 ));
1951 assert_eq!(flaky.puts.load(Ordering::SeqCst), 11);
1952 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
1953 }
1954
1955 #[tokio::test]
1956 async fn compare_and_swap_never_retries_transport_failures() {
1957 let flaky = Arc::new(FlakyStore::default());
1958 let store = retrying_store(Arc::clone(&flaky));
1959 let key = "namespaces/demo/wal/head.json";
1960 let seeded = store
1961 .put_overwrite(key, Bytes::from_static(b"one"))
1962 .await
1963 .expect("seed head");
1964 let etag = seeded.etag.expect("etag");
1965 script_puts(&flaky, [WriteScript::FailWithoutLanding]);
1966
1967 let error = store
1968 .compare_and_swap(key, &etag, Bytes::from_static(b"two"))
1969 .await
1970 .expect_err("compare-and-swap surfaces the transport failure");
1971
1972 assert!(matches!(error, ObjectStoreError::Transport { .. }));
1973 assert_eq!(flaky.puts.load(Ordering::SeqCst), 2);
1974 assert_eq!(
1975 store.get(key, None).await.expect("get"),
1976 Some(Bytes::from_static(b"one"))
1977 );
1978 }
1979
1980 #[tokio::test]
1981 async fn delete_retries_stop_once_the_operation_deadline_is_spent() {
1982 let flaky = Arc::new(FlakyStore::default());
1983 let store = retrying_store(Arc::clone(&flaky))
1984 .monotonic_timer(Arc::new(SteppingTimer::new(45_000)));
1985 let key = "namespaces/demo/uploads/upl_10.json";
1986 store
1987 .put_overwrite(key, Bytes::from_static(b"payload"))
1988 .await
1989 .expect("seed object");
1990 for _ in 0..6 {
1991 flaky
1992 .delete_script
1993 .lock()
1994 .expect("delete script")
1995 .push_back(WriteScript::FailWithoutLanding);
1996 }
1997
1998 let error = store
1999 .delete(key)
2000 .await
2001 .expect_err("deadline exhaustion surfaces the transport failure");
2002
2003 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2004 let attempts = flaky.deletes.load(Ordering::SeqCst);
2005 assert!(
2006 attempts < 5,
2007 "deadline must stop the loop before the count budget ({attempts} attempts)"
2008 );
2009 }
2010
2011 #[tokio::test]
2012 async fn non_transport_provider_errors_are_not_retried() {
2013 let flaky = Arc::new(FlakyStore::default());
2014 let store = retrying_store(Arc::clone(&flaky));
2015 script_puts(&flaky, [WriteScript::FailAuth]);
2016 let key = "namespaces/demo/uploads/upl_5.json";
2017
2018 let error = store
2019 .put_overwrite(key, Bytes::from_static(b"payload"))
2020 .await
2021 .expect_err("auth rejection surfaces immediately");
2022
2023 assert!(matches!(error, ObjectStoreError::PermissionDenied { .. }));
2024 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
2025 }
2026
2027 #[tokio::test]
2028 async fn delete_retries_transient_failures_and_landed_deletes() {
2029 let flaky = Arc::new(FlakyStore::default());
2030 let store = retrying_store(Arc::clone(&flaky));
2031
2032 let landed_key = "namespaces/demo/uploads/upl_6.json";
2033 store
2034 .put_overwrite(landed_key, Bytes::from_static(b"payload"))
2035 .await
2036 .expect("seed object");
2037 flaky
2038 .delete_script
2039 .lock()
2040 .expect("delete script")
2041 .push_back(WriteScript::LandThenFail);
2042 store
2043 .delete(landed_key)
2044 .await
2045 .expect("landed delete converges to success");
2046 assert_eq!(flaky.deletes.load(Ordering::SeqCst), 2);
2047 assert!(store.head(landed_key).await.expect("head").is_none());
2048
2049 let transient_key = "namespaces/demo/uploads/upl_7.json";
2050 store
2051 .put_overwrite(transient_key, Bytes::from_static(b"payload"))
2052 .await
2053 .expect("seed object");
2054 flaky
2055 .delete_script
2056 .lock()
2057 .expect("delete script")
2058 .push_back(WriteScript::FailWithoutLanding);
2059 store
2060 .delete(transient_key)
2061 .await
2062 .expect("transient delete failure is retried");
2063 assert_eq!(flaky.deletes.load(Ordering::SeqCst), 4);
2064 assert!(store.head(transient_key).await.expect("head").is_none());
2065 }
2066
2067 #[tokio::test]
2070 async fn a_retried_call_reports_its_attempts_to_the_metrics_wrapper() {
2071 let flaky = Arc::new(FlakyStore::default());
2072 let recorder = Arc::new(VecObjectStoreMetricsRecorder::default());
2073 let store =
2074 InstrumentedObjectStore::new(retrying_store(Arc::clone(&flaky)), recorder.clone());
2075 let key = "namespaces/demo/uploads/upl_8.json";
2076
2077 store
2078 .put_overwrite(key, Bytes::from_static(b"payload"))
2079 .await
2080 .expect("seed object");
2081 flaky
2082 .delete_script
2083 .lock()
2084 .expect("delete script")
2085 .push_back(WriteScript::FailWithoutLanding);
2086 store.delete(key).await.expect("the delete converges");
2087
2088 let samples = recorder.samples();
2089 let delete = samples
2090 .iter()
2091 .find(|sample| sample.operation == ObjectStoreOperation::Delete)
2092 .expect("the delete is sampled");
2093 assert_eq!(
2094 delete.attempts, 2,
2095 "one failed attempt, then the one that landed"
2096 );
2097 let seed = samples
2098 .iter()
2099 .find(|sample| sample.operation == ObjectStoreOperation::Put)
2100 .expect("the seed write is sampled");
2101 assert_eq!(
2102 seed.attempts, 1,
2103 "a call that never retried made one attempt"
2104 );
2105 }
2106
2107 #[test]
2108 fn request_phase_bound_has_two_flat_tiers() {
2109 assert_eq!(request_phase_bound(0), PROVIDER_ATTEMPT_TIMEOUT);
2110 assert_eq!(
2111 request_phase_bound(PROVIDER_TRANSFER_BODY_MIN_BYTES - 1),
2112 PROVIDER_ATTEMPT_TIMEOUT
2113 );
2114 assert_eq!(
2115 request_phase_bound(PROVIDER_TRANSFER_BODY_MIN_BYTES),
2116 PROVIDER_TRANSFER_ATTEMPT_TIMEOUT
2117 );
2118 assert_eq!(
2119 request_phase_bound(PROVIDER_MULTIPART_PART_BYTES),
2120 PROVIDER_TRANSFER_ATTEMPT_TIMEOUT
2121 );
2122 }
2123
2124 const MULTIPART_TEST_THRESHOLD: u64 = 1024;
2125 const MULTIPART_TEST_PART: u64 = 512;
2126 const MULTIPART_KEY: &str =
2127 "content-stores/cs_0123456789abcdef0123456789abcdef/objects/ab/con_abcdef0123456789abcdef0123456789";
2128
2129 fn multipart_test_store(flaky: Arc<FlakyStore>) -> ProviderObjectStore {
2132 retrying_store(flaky).multipart_geometry(MULTIPART_TEST_THRESHOLD, MULTIPART_TEST_PART)
2133 }
2134
2135 fn multipart_payload(len: usize) -> Vec<u8> {
2136 (0..len).map(|index| (index % 251) as u8).collect()
2137 }
2138
2139 fn script_part(
2140 flaky: &FlakyStore,
2141 part_index: usize,
2142 script: impl IntoIterator<Item = WriteScript>,
2143 ) {
2144 flaky
2145 .part_script
2146 .lock()
2147 .expect("part script")
2148 .entry(part_index)
2149 .or_default()
2150 .extend(script);
2151 }
2152
2153 fn part_attempts(flaky: &FlakyStore, part_index: usize) -> usize {
2154 flaky
2155 .part_attempts
2156 .lock()
2157 .expect("part attempts")
2158 .get(&part_index)
2159 .copied()
2160 .unwrap_or(0)
2161 }
2162
2163 fn script_complete(flaky: &FlakyStore, script: impl IntoIterator<Item = WriteScript>) {
2164 flaky
2165 .complete_script
2166 .lock()
2167 .expect("complete script")
2168 .extend(script);
2169 }
2170
2171 async fn seed_scoped_object(flaky: &FlakyStore, key: &str, bytes: Bytes) {
2174 provider_store::ObjectStore::put_opts(
2175 &flaky.inner,
2176 &Path::from(format!("tenant-a/{key}")),
2177 bytes.into(),
2178 PutOptions::default(),
2179 )
2180 .await
2181 .expect("seed object");
2182 }
2183
2184 #[tokio::test]
2185 async fn immutable_large_write_routes_through_existing_multipart_path() {
2186 let flaky = Arc::new(FlakyStore::default());
2187 let store = retrying_store(Arc::clone(&flaky));
2188 let payload =
2189 multipart_payload(usize::try_from(PROVIDER_MULTIPART_THRESHOLD_BYTES).expect("usize"));
2190
2191 store
2192 .put_immutable_verified(MULTIPART_KEY, Bytes::from(payload.clone()))
2193 .await
2194 .expect("large immutable write");
2195
2196 assert_eq!(flaky.puts.load(Ordering::SeqCst), 0);
2197 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 1);
2198 assert_eq!(part_attempts(&flaky, 0), 1);
2199 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2200 assert_eq!(
2201 store.get(MULTIPART_KEY, None).await.expect("get"),
2202 Some(Bytes::from(payload))
2203 );
2204 }
2205
2206 #[tokio::test(start_paused = true)]
2207 async fn immutable_large_write_owns_completion_retry() {
2208 let flaky = Arc::new(FlakyStore::default());
2209 let store = retrying_store(Arc::clone(&flaky));
2210 let payload =
2211 multipart_payload(usize::try_from(PROVIDER_MULTIPART_THRESHOLD_BYTES).expect("usize"));
2212 script_complete(&flaky, [WriteScript::FailWithoutLanding]);
2213
2214 store
2215 .put_immutable_verified(MULTIPART_KEY, Bytes::from(payload.clone()))
2216 .await
2217 .expect("immutable operation retries the whole multipart write");
2218
2219 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 2);
2220 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 2);
2221 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2222 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
2223 assert_eq!(
2224 store.get(MULTIPART_KEY, None).await.expect("get"),
2225 Some(Bytes::from(payload))
2226 );
2227 }
2228
2229 #[tokio::test]
2230 async fn large_put_routes_through_multipart_and_preserves_bytes() {
2231 let flaky = Arc::new(FlakyStore::default());
2232 let store = multipart_test_store(Arc::clone(&flaky));
2233 let payload = multipart_payload(1300);
2235
2236 let metadata = store
2237 .put_overwrite(MULTIPART_KEY, Bytes::from(payload.clone()))
2238 .await
2239 .expect("multipart put");
2240
2241 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 1);
2242 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2243 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 0);
2244 assert_eq!(
2245 flaky.puts.load(Ordering::SeqCst),
2246 0,
2247 "no whole-object PUT for a payload above the threshold"
2248 );
2249 assert_eq!(part_attempts(&flaky, 0), 1);
2250 assert_eq!(part_attempts(&flaky, 1), 1);
2251 assert_eq!(part_attempts(&flaky, 2), 1);
2252 assert_eq!(metadata.size_bytes, 1300);
2253 assert_eq!(
2254 store.get(MULTIPART_KEY, None).await.expect("get"),
2255 Some(Bytes::from(payload))
2256 );
2257 }
2258
2259 #[tokio::test]
2260 async fn multipart_threshold_boundary_routes_exactly() {
2261 let flaky = Arc::new(FlakyStore::default());
2262 let store = multipart_test_store(Arc::clone(&flaky));
2263
2264 store
2265 .put_overwrite(
2266 "namespaces/demo/uploads/upl_small.bin",
2267 Bytes::from(multipart_payload(MULTIPART_TEST_THRESHOLD as usize - 1)),
2268 )
2269 .await
2270 .expect("below-threshold put");
2271 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 0);
2272 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
2273
2274 store
2275 .put_overwrite(
2276 MULTIPART_KEY,
2277 Bytes::from(multipart_payload(MULTIPART_TEST_THRESHOLD as usize)),
2278 )
2279 .await
2280 .expect("at-threshold put");
2281 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 1);
2282 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
2283 }
2284
2285 #[tokio::test]
2289 async fn large_create_if_absent_stays_single_request() {
2290 let flaky = Arc::new(FlakyStore::default());
2291 let store = multipart_test_store(Arc::clone(&flaky));
2292 let payload = multipart_payload(1300);
2293
2294 store
2295 .put_if_absent(MULTIPART_KEY, Bytes::from(payload.clone()))
2296 .await
2297 .expect("create absent large object");
2298 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 0);
2299 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
2300
2301 let error = store
2302 .put_if_absent(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2303 .await
2304 .expect_err("existing object fails the create precondition");
2305 assert!(matches!(error, ObjectStoreError::PreconditionFailed { .. }));
2306 assert_eq!(
2307 flaky.multipart_creates.load(Ordering::SeqCst),
2308 0,
2309 "the conflict is decided by the provider precondition, not a pre-check"
2310 );
2311 }
2312
2313 fn streamed(payload: &[u8], chunk_bytes: usize) -> ByteStream {
2316 let chunks: Vec<Bytes> = payload
2317 .chunks(chunk_bytes)
2318 .map(Bytes::copy_from_slice)
2319 .collect();
2320 stream::iter(chunks.into_iter().map(Ok)).boxed()
2321 }
2322
2323 #[tokio::test]
2326 async fn a_streamed_put_cuts_the_stream_into_parts_and_preserves_bytes() {
2327 let flaky = Arc::new(FlakyStore::default());
2328 let store = multipart_test_store(Arc::clone(&flaky));
2329 let payload = multipart_payload(1300);
2332
2333 let size_bytes = store
2334 .put_streamed(
2335 MULTIPART_KEY,
2336 streamed(&payload, 100),
2337 PutMode::CreateIfAbsent,
2338 )
2339 .await
2340 .expect("streamed multipart put");
2341
2342 assert_eq!(size_bytes, 1300);
2343 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 1);
2344 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2345 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 0);
2346 assert_eq!(
2347 flaky.puts.load(Ordering::SeqCst),
2348 0,
2349 "a payload past one part never becomes a whole-object PUT"
2350 );
2351 assert_eq!(part_attempts(&flaky, 0), 1);
2352 assert_eq!(part_attempts(&flaky, 1), 1);
2353 assert_eq!(part_attempts(&flaky, 2), 1);
2354 assert_eq!(
2355 store.get(MULTIPART_KEY, None).await.expect("get"),
2356 Some(Bytes::from(payload))
2357 );
2358 }
2359
2360 #[tokio::test]
2364 async fn a_short_streamed_put_is_one_request_that_keeps_its_precondition() {
2365 let flaky = Arc::new(FlakyStore::default());
2366 let store = multipart_test_store(Arc::clone(&flaky));
2367 let payload = multipart_payload(MULTIPART_TEST_PART as usize - 1);
2368
2369 store
2370 .put_streamed(
2371 MULTIPART_KEY,
2372 streamed(&payload, 64),
2373 PutMode::CreateIfAbsent,
2374 )
2375 .await
2376 .expect("short streamed put");
2377 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 0);
2378 assert_eq!(flaky.puts.load(Ordering::SeqCst), 1);
2379
2380 let error = store
2381 .put_streamed(
2382 MULTIPART_KEY,
2383 streamed(&payload, 64),
2384 PutMode::CreateIfAbsent,
2385 )
2386 .await
2387 .expect_err("the key is taken and create-only means it");
2388 assert!(matches!(error, ObjectStoreError::PreconditionFailed { .. }));
2389 assert_eq!(
2390 store.get(MULTIPART_KEY, None).await.expect("get"),
2391 Some(Bytes::from(payload))
2392 );
2393 }
2394
2395 #[tokio::test]
2398 async fn an_empty_streamed_put_writes_an_empty_object() {
2399 let flaky = Arc::new(FlakyStore::default());
2400 let store = multipart_test_store(Arc::clone(&flaky));
2401
2402 let size_bytes = store
2403 .put_streamed(
2404 MULTIPART_KEY,
2405 stream::empty().boxed(),
2406 PutMode::CreateIfAbsent,
2407 )
2408 .await
2409 .expect("empty streamed put");
2410
2411 assert_eq!(size_bytes, 0);
2412 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 0);
2413 assert_eq!(
2414 store.get(MULTIPART_KEY, None).await.expect("get"),
2415 Some(Bytes::new())
2416 );
2417 }
2418
2419 #[tokio::test]
2423 async fn a_streamed_put_that_fails_mid_stream_abandons_its_upload() {
2424 let flaky = Arc::new(FlakyStore::default());
2425 let store = multipart_test_store(Arc::clone(&flaky));
2426 let head = multipart_payload(MULTIPART_TEST_PART as usize);
2427 let body = stream::iter([
2428 Ok(Bytes::from(head)),
2429 Ok(Bytes::from(multipart_payload(64))),
2430 Err(ObjectStoreError::transport(
2431 MULTIPART_KEY,
2432 "the client stopped sending",
2433 )),
2434 ])
2435 .boxed();
2436
2437 let error = store
2438 .put_streamed(MULTIPART_KEY, body, PutMode::CreateIfAbsent)
2439 .await
2440 .expect_err("a payload that stops is not a write");
2441
2442 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2443 assert_eq!(flaky.multipart_creates.load(Ordering::SeqCst), 1);
2444 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 0);
2445 assert_eq!(
2446 flaky.multipart_aborts.load(Ordering::SeqCst),
2447 1,
2448 "the abandoned upload is aborted, not left holding parts"
2449 );
2450 assert_eq!(store.get(MULTIPART_KEY, None).await.expect("get"), None);
2451 }
2452
2453 #[tokio::test]
2454 async fn multipart_part_failures_are_retried_in_place() {
2455 let flaky = Arc::new(FlakyStore::default());
2456 let store = multipart_test_store(Arc::clone(&flaky));
2457 script_part(
2458 &flaky,
2459 1,
2460 [
2461 WriteScript::FailWithoutLanding,
2462 WriteScript::FailWithoutLanding,
2463 ],
2464 );
2465 let payload = multipart_payload(1300);
2466
2467 store
2468 .put_overwrite(MULTIPART_KEY, Bytes::from(payload.clone()))
2469 .await
2470 .expect("multipart put survives transient part failures");
2471
2472 assert_eq!(part_attempts(&flaky, 0), 1);
2473 assert_eq!(
2474 part_attempts(&flaky, 1),
2475 3,
2476 "the failing part retries in place under the same index"
2477 );
2478 assert_eq!(part_attempts(&flaky, 2), 1);
2479 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 0);
2480 assert_eq!(
2481 store.get(MULTIPART_KEY, None).await.expect("get"),
2482 Some(Bytes::from(payload))
2483 );
2484 }
2485
2486 #[tokio::test]
2487 async fn multipart_part_budget_exhaustion_aborts_the_upload() {
2488 let flaky = Arc::new(FlakyStore::default());
2489 let store = multipart_test_store(Arc::clone(&flaky));
2490 script_part(&flaky, 0, (0..6).map(|_| WriteScript::FailWithoutLanding));
2491
2492 let error = store
2493 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2494 .await
2495 .expect_err("persistent part failure surfaces after the retry budget");
2496
2497 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2498 assert_eq!(part_attempts(&flaky, 0), 5, "1 attempt + max_retries");
2499 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 0);
2500 assert_eq!(
2501 flaky.multipart_aborts.load(Ordering::SeqCst),
2502 1,
2503 "a failed upload is aborted so no parts are stranded"
2504 );
2505 assert!(store.head(MULTIPART_KEY).await.expect("head").is_none());
2506 }
2507
2508 #[tokio::test]
2509 async fn multipart_auth_failure_is_not_retried() {
2510 let flaky = Arc::new(FlakyStore::default());
2511 let store = multipart_test_store(Arc::clone(&flaky));
2512 script_part(&flaky, 0, [WriteScript::FailAuth]);
2513
2514 let error = store
2515 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2516 .await
2517 .expect_err("auth rejection surfaces immediately");
2518
2519 assert!(matches!(error, ObjectStoreError::PermissionDenied { .. }));
2522 assert_eq!(part_attempts(&flaky, 0), 1);
2523 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2524 }
2525
2526 #[tokio::test]
2527 async fn multipart_complete_transport_failure_resolves_landed_completion() {
2528 let flaky = Arc::new(FlakyStore::default());
2529 let store = multipart_test_store(Arc::clone(&flaky));
2530 script_complete(&flaky, [WriteScript::LandThenFail]);
2531 let payload = multipart_payload(1300);
2532
2533 let metadata = store
2534 .put_overwrite(MULTIPART_KEY, Bytes::from(payload.clone()))
2535 .await
2536 .expect("landed completion reported as the success it was");
2537
2538 assert_eq!(
2539 flaky.multipart_completes.load(Ordering::SeqCst),
2540 1,
2541 "completion is attempted once before byte-identity resolution"
2542 );
2543 assert_eq!(
2544 flaky.gets.load(Ordering::SeqCst),
2545 1,
2546 "one read-back proves the landed write by byte identity"
2547 );
2548 assert_eq!(
2549 flaky.multipart_aborts.load(Ordering::SeqCst),
2550 1,
2551 "the proven completion still aborts the gone upload id best-effort"
2552 );
2553 assert_eq!(metadata.size_bytes, 1300);
2554 let head = store
2555 .head(MULTIPART_KEY)
2556 .await
2557 .expect("head")
2558 .expect("object exists");
2559 assert_eq!(
2560 metadata.etag, head.etag,
2561 "resolution reports the landed object's own metadata"
2562 );
2563 assert_eq!(
2564 store.get(MULTIPART_KEY, None).await.expect("get"),
2565 Some(Bytes::from(payload))
2566 );
2567 }
2568
2569 #[tokio::test]
2570 async fn raw_multipart_complete_failure_without_landing_is_not_retried() {
2571 let flaky = Arc::new(FlakyStore::default());
2572 let store = multipart_test_store(Arc::clone(&flaky));
2573 script_complete(&flaky, [WriteScript::FailWithoutLanding]);
2574 let payload = multipart_payload(1300);
2575
2576 let error = store
2577 .put_overwrite(MULTIPART_KEY, Bytes::from(payload.clone()))
2578 .await
2579 .expect_err("raw overwrite surfaces the ambiguous completion");
2580
2581 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2582 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2583 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2584 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
2585 assert!(store.head(MULTIPART_KEY).await.expect("head").is_none());
2586 }
2587
2588 #[tokio::test]
2593 async fn multipart_complete_rejects_stale_same_size_object() {
2594 let flaky = Arc::new(FlakyStore::default());
2595 let store = multipart_test_store(Arc::clone(&flaky));
2596 let stale = Bytes::from(vec![0xAA_u8; 1300]);
2597 seed_scoped_object(&flaky, MULTIPART_KEY, stale.clone()).await;
2598 script_complete(&flaky, [WriteScript::FailWithoutLanding]);
2599
2600 let error = store
2601 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2602 .await
2603 .expect_err("an unproven completion fails instead of adopting the stale object");
2604
2605 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2606 assert_eq!(
2607 flaky.multipart_completes.load(Ordering::SeqCst),
2608 1,
2609 "raw overwrite does not replay completion"
2610 );
2611 assert_eq!(
2612 flaky.gets.load(Ordering::SeqCst),
2613 1,
2614 "one read-back tested the outcome"
2615 );
2616 assert_eq!(
2617 flaky.multipart_aborts.load(Ordering::SeqCst),
2618 1,
2619 "the failed upload is aborted so no parts are stranded"
2620 );
2621 assert_eq!(
2622 store.get(MULTIPART_KEY, None).await.expect("get"),
2623 Some(stale),
2624 "the stale object is untouched"
2625 );
2626 }
2627
2628 #[tokio::test]
2633 async fn multipart_complete_accepts_identical_object_and_aborts() {
2634 let flaky = Arc::new(FlakyStore::default());
2635 let store = multipart_test_store(Arc::clone(&flaky));
2636 let payload = multipart_payload(1300);
2637 seed_scoped_object(&flaky, MULTIPART_KEY, Bytes::from(payload.clone())).await;
2638 script_complete(&flaky, [WriteScript::FailWithoutLanding]);
2639
2640 let metadata = store
2641 .put_overwrite(MULTIPART_KEY, Bytes::from(payload))
2642 .await
2643 .expect("byte-identical object proves the outcome");
2644
2645 assert_eq!(metadata.size_bytes, 1300);
2646 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
2647 assert_eq!(
2648 flaky.multipart_aborts.load(Ordering::SeqCst),
2649 1,
2650 "the dangling upload is aborted on proven success"
2651 );
2652 assert!(
2653 flaky.multipart_uploads.lock().expect("uploads").is_empty(),
2654 "no parts remain stranded"
2655 );
2656 }
2657
2658 #[tokio::test]
2662 async fn multipart_complete_gone_upload_with_stale_object_fails_as_transport() {
2663 let flaky = Arc::new(FlakyStore::default());
2664 let store = multipart_test_store(Arc::clone(&flaky));
2665 let stale = Bytes::from(vec![0xAA_u8; 1300]);
2666 seed_scoped_object(&flaky, MULTIPART_KEY, stale.clone()).await;
2667 script_complete(&flaky, [WriteScript::VanishThenFail]);
2668
2669 let error = store
2670 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2671 .await
2672 .expect_err("a vanished upload with a stale object is a failed write");
2673
2674 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2675 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2676 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
2677 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2678 assert_eq!(
2679 store.get(MULTIPART_KEY, None).await.expect("get"),
2680 Some(stale),
2681 "the stale object is untouched"
2682 );
2683 }
2684
2685 #[tokio::test]
2688 async fn multipart_complete_first_attempt_rejection_skips_verification() {
2689 let flaky = Arc::new(FlakyStore::default());
2690 let store = multipart_test_store(Arc::clone(&flaky));
2691 script_complete(&flaky, [WriteScript::FailAuth]);
2692
2693 let error = store
2694 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2695 .await
2696 .expect_err("auth rejection surfaces immediately");
2697
2698 assert!(matches!(error, ObjectStoreError::PermissionDenied { .. }));
2699 assert_eq!(flaky.multipart_completes.load(Ordering::SeqCst), 1);
2700 assert_eq!(flaky.gets.load(Ordering::SeqCst), 0);
2701 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2702 }
2703
2704 #[tokio::test]
2708 async fn multipart_complete_unverifiable_outcome_surfaces_both_failures() {
2709 let flaky = Arc::new(FlakyStore::default());
2710 let store = multipart_test_store(Arc::clone(&flaky));
2711 script_complete(&flaky, [WriteScript::FailWithoutLanding]);
2712 flaky
2713 .get_script
2714 .lock()
2715 .expect("get script")
2716 .push_back(ReadScript::Transport);
2717
2718 let error = store
2719 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2720 .await
2721 .expect_err("an unverifiable outcome is an error, not a success");
2722
2723 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2724 let message = error.message();
2725 assert!(
2726 message.contains("failed to verify multipart completion outcome"),
2727 "message names the verification failure: {message}"
2728 );
2729 assert_eq!(flaky.gets.load(Ordering::SeqCst), 1);
2730 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2731 }
2732
2733 #[tokio::test]
2737 async fn multipart_part_retries_are_count_bounded_not_clock_bounded() {
2738 let flaky = Arc::new(FlakyStore::default());
2739 let store = multipart_test_store(Arc::clone(&flaky))
2740 .monotonic_timer(Arc::new(SteppingTimer::new(45_000)));
2741 script_part(&flaky, 0, (0..6).map(|_| WriteScript::FailWithoutLanding));
2742
2743 let error = store
2744 .put_overwrite(MULTIPART_KEY, Bytes::from(multipart_payload(1300)))
2745 .await
2746 .expect_err("persistent part failure surfaces after the retry budget");
2747
2748 assert!(matches!(error, ObjectStoreError::Transport { .. }));
2749 assert_eq!(
2750 part_attempts(&flaky, 0),
2751 5,
2752 "1 attempt + max_retries, unaffected by elapsed time"
2753 );
2754 assert_eq!(flaky.multipart_aborts.load(Ordering::SeqCst), 1);
2755 }
2756}