Skip to main content

loonfs_objectstore/
provider_object_store.rs

1//! The shared provider transport: timeouts, bounded retries for replay-safe
2//! delete and multipart stages, and multipart upload for large immutable
3//! payloads.
4
5use 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/// Configures logical key scoping for a generic provider client.
33#[derive(Debug, Clone, PartialEq, Eq)]
34pub struct ProviderObjectStoreConfig {
35    /// Prefix prepended to every provider key, or `None` to expose the bucket root.
36    pub key_prefix: Option<String>,
37}
38
39/// Bound for one control-plane HTTP attempt's request phase, and the
40/// response-body idle bound for every request. An attempt that makes no
41/// progress for this long fails and counts against the operation deadline
42/// instead of consuming it invisibly.
43pub const PROVIDER_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(30);
44
45/// One HTTP attempt's connect timeout.
46pub const PROVIDER_CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
47
48/// Hard deadline for one logical object-store operation, consumed across
49/// every retry of that operation rather than restarting per attempt. Reads
50/// get it as the provider client's retry timeout; verified immutable writes
51/// and deletes share it through `TransportRetryPolicy`. The deadline gates
52/// starting another attempt, and one outer attempt may itself contain the
53/// inner client's full retry budget for status-code retries — so one
54/// operation's total wall time is bounded by the deadline plus the inner
55/// retry budget plus one attempt bound (worst case roughly six minutes),
56/// still a hard bound. The GC grace window is derived
57/// above this bound (format spec, "Garbage collection", rule 1); multipart
58/// uploads deliberately carry no whole-operation clock (their parts are
59/// individually bounded), which leaves the floor inequality untouched
60/// because everything it times — WAL segments inside the publish budget,
61/// the root compare-and-swap — is a small control object on the
62/// single-request path.
63pub const PROVIDER_OPERATION_DEADLINE: Duration = Duration::from_secs(120);
64
65/// Payload size at and above which overwrite puts use the provider's native
66/// multipart upload instead of one whole-object PUT, matching the multipart
67/// thresholds mainstream storage clients ship. The format spec allows this
68/// for large immutable file data and forbids relying on it for small
69/// mutable control objects: create-if-absent and compare-and-swap puts
70/// never take this path, because providers complete multipart uploads as
71/// unconditional overwrites and those modes exist to carry real provider
72/// preconditions.
73pub const PROVIDER_MULTIPART_THRESHOLD_BYTES: u64 = 8 * 1024 * 1024;
74
75/// Fixed size of every multipart part except the last. Cloudflare R2
76/// requires all non-final parts to share one size, and every supported
77/// provider requires at least 5 MiB per non-final part; 8 MiB matches the
78/// part size mainstream storage clients default to, and keeps every part a
79/// cheap retry that fits comfortably inside one flat attempt bound.
80pub const PROVIDER_MULTIPART_PART_BYTES: u64 = 8 * 1024 * 1024;
81
82/// Concurrent in-flight parts per multipart upload.
83pub const PROVIDER_MULTIPART_PART_WINDOW: usize = 4;
84
85/// Buffered parts a streamed write holds at once.
86///
87/// One. The in-memory path can afford a window because it already holds the
88/// whole payload; a streamed write exists precisely so that it does not, so
89/// it fills a buffer, uploads it, drops it, and only then fills the next.
90/// Peak memory is therefore one part whatever the object's size.
91pub const PROVIDER_STREAMED_PART_WINDOW: usize = 1;
92
93/// Parts one provider multipart upload accepts. Every supported provider
94/// stops at 10,000, which with the part size sets the largest object a
95/// multipart write can produce.
96pub(crate) const MAX_PROVIDER_MULTIPART_PARTS: usize = 10_000;
97
98/// Bound for one payload-bearing HTTP attempt's request phase. A request
99/// body is opaque to progress observation while it uploads, so a flat
100/// generous bound stands in for stall detection: parts are at most
101/// [`PROVIDER_MULTIPART_PART_BYTES`], and an 8 MiB body that cannot finish
102/// inside this bound is moving slower than roughly 70 KiB/s — treated as
103/// stalled and retried on a fresh connection.
104pub const PROVIDER_TRANSFER_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(120);
105
106/// Request bodies at least this large are payload transfers and get
107/// [`PROVIDER_TRANSFER_ATTEMPT_TIMEOUT`] as their request-phase bound;
108/// smaller bodies are control-plane traffic bounded by
109/// [`PROVIDER_ATTEMPT_TIMEOUT`]. Sits well below the part size so multipart
110/// tail parts classify with their siblings.
111pub(crate) const PROVIDER_TRANSFER_BODY_MIN_BYTES: u64 = 1024 * 1024;
112
113/// Bound for one HTTP attempt's request phase (connect, request-body
114/// upload, response headers), by request body size: flat and small for
115/// control-plane requests, flat and generous for payload transfers.
116pub(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
124/// Client options every provider builder applies: an explicit per-attempt
125/// total-request timeout and connect timeout, so a client built from these
126/// options alone is bounded by named constants instead of upstream defaults.
127/// [`crate::transfer_timeouts::TransferTimeoutConnector`] strips the
128/// total-request timeout and replaces it with payload-aware request bounds
129/// and response-body idle bounds.
130pub(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
136/// Retry configuration every provider builder applies: the client's internal
137/// read retries consume [`PROVIDER_OPERATION_DEADLINE`] as one per-operation budget,
138/// matching the write loops above.
139pub(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/// The size routing for multipart writes: payloads at or above the
147/// threshold are uploaded as fixed-size parts. One production value
148/// ([`PROVIDER_MULTIPART_THRESHOLD_BYTES`], [`PROVIDER_MULTIPART_PART_BYTES`]);
149/// tests shrink it to exercise the machinery without allocating gigabytes.
150#[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/// Adapts the upstream `object_store` provider surface to the narrower LoonFS contract.
164#[derive(Clone)]
165pub struct ProviderObjectStore {
166    inner: Arc<dyn provider_store::ObjectStore>,
167    /// The provider's native multipart surface, used for payloads at or
168    /// above [`PROVIDER_MULTIPART_THRESHOLD_BYTES`]. `None` only for
169    /// providers without one; their large puts stay whole-object PUTs under
170    /// the payload-scaled bounds.
171    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    /// Wraps a provider client and optional native multipart surface.
189    ///
190    /// Construction fails when `config.key_prefix` is not a normalized,
191    /// non-escaping logical prefix.
192    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    /// Writes one large payload through the provider's native multipart
287    /// upload: fixed-size parts uploaded through a bounded window, each part
288    /// retried in place on transient failures (part indices are stable, so a
289    /// retry re-sends the same part), and a best-effort abort so a failed
290    /// upload does not strand parts. An ambiguous completion — a transport
291    /// failure whose attempt may have landed — is resolved by the immutable
292    /// operation's shared exact-byte read-back
293    /// ([`MultipartWrite::resolve_ambiguous_completion`]).
294    ///
295    /// There is deliberately no whole-operation clock: every part attempt is
296    /// individually bounded and every retry loop is count-bounded, so a
297    /// healthy transfer takes as long as the link needs while a stuck one
298    /// still fails within one part's retry budget. Completion itself is one
299    /// attempt: only [`ObjectStore::put_immutable_verified`] may replay an
300    /// object-publishing write. Conditional modes stay on the single-request
301    /// path where real provider preconditions exist.
302    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                // Best effort, and harmless when the failure raced a landed
323                // completion: the upload id no longer exists then, and the
324                // abort cannot touch the completed object.
325                upload.abort(&upload_id).await;
326                Err(err)
327            }
328        }
329    }
330
331    /// Applies the shared retry gate for one failed write attempt: `None`
332    /// means the budget is spent and the caller must surface the error;
333    /// `Some` carries the backoff to sleep before the next attempt. The
334    /// budget is the retry count, plus the operation deadline when the
335    /// operation carries one (multipart transfers deliberately do not — see
336    /// [`Self::put_large_multipart`]). Exhaustion logs name the payload
337    /// size so a too-slow-link failure is attributable instead of reading
338    /// as weather.
339    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        // Every write retry this store grants passes through here, so this
376        // is the one place a measuring wrapper above needs counted for its
377        // sample's attempt count.
378        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
393/// Cuts a byte stream into fixed-size parts, holding one at a time.
394///
395/// Chunk boundaries in the source stream carry no meaning, so a chunk that
396/// straddles a part boundary is split and its tail carried into the next
397/// part. A stream that ends exactly on a boundary produces no final part.
398struct PartReader {
399    body: ByteStream,
400    /// The tail of a chunk that overran the part being cut.
401    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    /// Cuts the next part: exactly `part_bytes`, or whatever is left when
417    /// the stream ends. `None` once nothing is left.
418    ///
419    /// A full part is returned without polling the stream again, so a
420    /// caller cannot conclude from a full part that more is coming — only
421    /// a short part proves the stream ended.
422    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    /// Whether the stream has already reported its end.
446    fn exhausted(&self) -> bool {
447        self.exhausted && self.carry.is_none()
448    }
449}
450
451/// Abandons a provider multipart upload if the write that opened it is
452/// dropped before it finishes.
453///
454/// Every path that returns from a streamed write aborts explicitly and
455/// disarms this. What is left is cancellation — a client that disconnects
456/// mid-upload takes the handler's future with it — where there is no
457/// `await` point to run cleanup from, so the abort is spawned. Without a
458/// runtime to spawn on, the upload falls back to the bucket's own
459/// incomplete-upload lifecycle rule, which is what collects it today.
460struct 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
501/// One in-progress multipart write: the store, the provider multipart
502/// surface, and the object being written.
503struct 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    /// Uploads a stream as parts, one buffered part at a time, and
536    /// assembles them.
537    ///
538    /// `head` is the first part, already cut by the caller so it could
539    /// decide that the payload does not fit in one request. Every later part
540    /// is cut into a fresh buffer *after* its predecessor has been uploaded
541    /// and dropped, so peak memory is one part regardless of how large the
542    /// object is. The final part may be short.
543    ///
544    /// Completion is not resolved by read-back the way the in-memory path's
545    /// is: there is no payload left in memory to prove a landed write
546    /// against. An ambiguous completion is a failure, and the caller abandons
547    /// the upload.
548    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                // Dropping the window cancels the sibling part uploads; the
607                // caller aborts the upload so no parts are stranded.
608                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    /// Decides an ambiguous completion by reading the object back: byte
681    /// equality with the payload is the only accepted proof that the write
682    /// landed. Size or etag agreement is never identity — this store serves
683    /// generic overwrite keys, so a stale object of the same length must
684    /// not pass as the new write. The payload is still in memory, so the
685    /// read-back costs one GET on this failure path and proves the put's
686    /// postcondition itself.
687    ///
688    /// A proven completion still aborts the upload id: when this upload's
689    /// completion landed the id is already gone and the abort is a no-op,
690    /// and when the proof came from an identical object some earlier writer
691    /// committed, the abort reclaims this upload's stranded parts.
692    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        // Best effort: an unaborted upload only strands parts until the
738        // bucket's lifecycle rule for incomplete multipart uploads collects
739        // them, so an abort failure is logged rather than surfaced. An
740        // already-gone upload is the no-op success it reads as — typically
741        // its completion landed.
742        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    /// Bounded reads issue the ranged GET directly — one round trip, not a
788    /// sizing HEAD plus a GET — and pay a single HEAD only on the failure
789    /// path, to decide whether the range or the transport was the problem.
790    /// The contract matches the local reference provider exactly: a missing
791    /// object is `Ok(None)` however the request was shaped, an end past the
792    /// object clamps, `start == size` reads empty, and `start > size` is
793    /// `InvalidRange`.
794    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            // A zero-length request needs no bytes; existence and size
814            // alone answer it.
815            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                // The provider refused; one HEAD decides whether the range
833                // was the problem, matching the reference semantics.
834                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                        // A strict provider rejected the over-long end
846                        // instead of clamping; clamp and retry once.
847                        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        // Raw flat writes are deliberately one attempt. In particular, a
875        // mutable overwrite cannot be replayed after an ambiguous transport
876        // outcome without changing write ordering. Immutable callers use
877        // `put_immutable_verified`, whose name supplies the retry invariant.
878        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    /// Cuts the payload into parts as it arrives and uploads them one at a
899    /// time, so a large object costs one part of memory instead of its own
900    /// size.
901    ///
902    /// The first part is cut before anything is decided. A payload that
903    /// ends inside it is an ordinary [`ObjectStore::put`] with the caller's
904    /// mode intact — which is the same size line `put` itself draws, so a
905    /// small streamed write behaves exactly like a small buffered one.
906    /// Anything longer goes through the provider's multipart upload, whose
907    /// completion is an unconditional overwrite.
908    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            // No provider multipart surface: fall back to the buffered
914            // contract rather than pretend, exactly as the default does.
915            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                // Best effort, and harmless when the failure raced a landed
950                // completion: the upload id no longer exists then.
951                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            // Delete is idempotent under this contract: not-found already
964            // reports success, so a retry after an attempt that landed
965            // converges to the same outcome.
966            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            // The compare token is opaque and provider-issued: S3-family
1018            // backends condition on `e_tag`, GCS conditions on `version`
1019            // (its generation). Populate both so each backend reads the
1020            // field it understands.
1021            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    // `Generic` is where the provider client surfaces request failures after
1041    // its own retry policy gives up: for writes that includes mid-flight
1042    // transport errors it refuses to re-send because a non-idempotent HTTP
1043    // request may already have reached the store. The remaining variants are
1044    // definite outcomes (not-found, already-exists, precondition), hard
1045    // rejections (auth, invalid path, unsupported), or the store IO runtime
1046    // shutting down (join error); re-sending those is wrong or futile.
1047    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
1092/// Query parameters whose values are credential material when they appear in
1093/// a URL a provider error echoes back (signed and presigned requests).
1094const 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
1103/// Response-body XML elements that echo signing inputs back to the caller in
1104/// provider auth failures (`SignatureDoesNotMatch` and friends).
1105const CREDENTIAL_XML_ELEMENTS: &[&str] = &[
1106    "StringToSign",
1107    "StringToSignBytes",
1108    "CanonicalRequest",
1109    "SignatureProvided",
1110    "AWSAccessKeyId",
1111];
1112
1113/// Strips credential material from free-text provider errors before they
1114/// enter the error chain: signature/credential query parameters in echoed
1115/// URLs, and the signing-input elements auth-failure bodies quote back.
1116/// The diagnosable parts (status, provider error code, cause) stay.
1117fn 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
1128/// Replaces every `?param=value` / `&param=value` occurrence's value with
1129/// `<redacted>`. Matches only at a query-parameter boundary so `Signature=`
1130/// does not fire inside `X-Amz-Signature=`.
1131fn 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
1155/// Replaces the text inside `<element>...</element>` with `<redacted>`; if
1156/// the closing tag never arrives (truncated body), everything after the
1157/// opening tag goes.
1158fn 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        // A bare name inside a longer parameter is not a boundary match.
1255        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        // In-bounds slice.
1301        assert_eq!(
1302            store.get(key, range(2, 6)).await.expect("bounded"),
1303            Some(Bytes::from_static(b"2345"))
1304        );
1305        // An end past the object clamps.
1306        assert_eq!(
1307            store.get(key, range(6, 99)).await.expect("clamped"),
1308            Some(Bytes::from_static(b"6789"))
1309        );
1310        // Reading at the exact end is empty, not an error.
1311        assert_eq!(
1312            store.get(key, range(10, 12)).await.expect("at end"),
1313            Some(Bytes::new())
1314        );
1315        // A start past the end is an invalid range.
1316        assert!(matches!(
1317            store.get(key, range(11, 12)).await,
1318            Err(ObjectStoreError::InvalidRange { .. })
1319        ));
1320        // An inverted range is rejected without any store call.
1321        assert!(matches!(
1322            store.get(key, range(6, 2)).await,
1323            Err(ObjectStoreError::InvalidRange { .. })
1324        ));
1325        // Zero-length reads answer from existence and size alone.
1326        assert_eq!(
1327            store.get(key, range(4, 4)).await.expect("zero length"),
1328            Some(Bytes::new())
1329        );
1330
1331        // A missing object is `Ok(None)` however the request is shaped —
1332        // the same answer the unranged read and the local provider give.
1333        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        /// The upload vanishes (as a lifecycle rule reaping it would make
1454        /// it) and the attempt reports a transport failure. Only meaningful
1455        /// as a completion script.
1456        VanishThenFail,
1457    }
1458
1459    #[derive(Debug)]
1460    enum ReadScript {
1461        Transport,
1462    }
1463
1464    /// Provider double that fails scripted attempts before delegating to an
1465    /// in-memory store, so retry behavior is observable per attempt.
1466    ///
1467    /// Multipart state is held here rather than delegated: the in-memory
1468    /// provider requires parts to arrive in index order, while real
1469    /// providers (and this transport's retry-driven interleavings) allow
1470    /// any order.
1471    #[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    /// The metrics wrapper sits above this store's retry loops, so without
2068    /// the attempt tally a retried call and a slow one make the same sample.
2069    #[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    /// Retrying store with a test-sized multipart geometry: payloads of
2130    /// 1024+ bytes go multipart in 512-byte parts.
2131    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    /// Places an object at `key` directly on the inner store, bypassing the
2172    /// scripted transport and its counters.
2173    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        // Unaligned tail: parts of 512, 512, and 276 bytes.
2234        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    /// Conditional modes never ride multipart: providers complete multipart
2286    /// uploads as unconditional overwrites, so create-if-absent keeps its
2287    /// real provider precondition on the single-request path at any size.
2288    #[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    /// Cuts a payload into stream chunks that deliberately do not line up
2314    /// with the part size, since a caller's chunk boundaries never do.
2315    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    /// The streamed write's whole job: regroup whatever arrives into fixed
2324    /// parts, upload them one at a time, and assemble exactly the payload.
2325    #[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        // Parts of 512, 512, and 276 bytes, delivered in 100-byte chunks so
2330        // every part boundary falls inside a chunk.
2331        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    /// A payload that ends inside the first buffered part is an ordinary
2361    /// put, with the caller's mode intact — the same size line `put` itself
2362    /// draws, so a small streamed write behaves like a small buffered one.
2363    #[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    /// An empty stream still writes an object, and still writes it the way
2396    /// its mode asks for.
2397    #[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    /// A payload that stops mid-transfer leaves nothing behind: the
2420    /// provider upload is abandoned rather than left holding parts, and no
2421    /// object appears at the key.
2422    #[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        // Auth failures carry their own classification — never mistaken
2520        // for network weather, and never retried.
2521        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    /// The regression this fix exists for: an unproven completion must
2589    /// never adopt a pre-existing object of the same size. The pre-fix code
2590    /// reconciled by size identity and reported success for bytes that were
2591    /// never written.
2592    #[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    /// The content-addressed re-put shape: when the object already holds
2629    /// exactly the payload bytes, byte identity proves the put's
2630    /// postcondition even though this upload's completion never landed —
2631    /// and the dangling upload is reclaimed rather than stranded.
2632    #[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    /// The lifecycle-abort race: the upload vanishes while the completion's
2659    /// outcome is ambiguous and a stale object sits at the key. The stale
2660    /// object is never accepted as this write.
2661    #[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    /// A first-attempt refusal is definite: nothing can have landed, so no
2686    /// read-back runs and the refusal surfaces directly.
2687    #[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    /// When the read-back itself fails, the outcome stays unknown and
2705    /// surfaces as a transport error carrying both failures — never as a
2706    /// success the store cannot prove.
2707    #[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    /// Multipart transfers carry no whole-operation clock: with a stepping
2734    /// timer that would spend the single-request deadline almost instantly,
2735    /// part retries still run to their full count budget.
2736    #[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}