Skip to main content

objectstore_service/backend/
tiered.rs

1//! Two-tier storage backend with size-based routing and redirect tombstones.
2//!
3//! [`TieredStorage`] routes objects to a high-volume or long-term backend based
4//! on size and maintains redirect tombstones so that reads never need to probe
5//! both backends. See the [crate-level documentation](crate) for the high-level
6//! motivation, and the [`TieredStorage`] struct docs for routing and tombstone
7//! semantics.
8//!
9//! # Cross-Tier Consistency
10//!
11//! A single logical object may span both backends: a tombstone in HV pointing
12//! to a payload in LT. Mutations keep the two in sync through compare-and-swap
13//! on the high-volume backend (see [`HighVolumeBackend::compare_and_write`]).
14//! Each operation reads the current HV revision, performs its work, then
15//! atomically swaps the HV entry only if the revision is still current —
16//! rolling back on conflict.
17//!
18//! ## Revision Keys
19//!
20//! Every long-term write stores its payload at a **revision key** in the
21//! long-term backend: `{original_key}/{uuid}`. The UUID suffix is random (no
22//! monotonicity is guaranteed), so each write targets a distinct LT path
23//! regardless of whether another write to the same logical key is in progress.
24//! The tombstone in HV then points to this specific revision. Because each
25//! writer owns its own LT blob, the compare-and-swap on the tombstone becomes
26//! an atomic pointer swap: the winner's revision is committed and the loser
27//! can safely delete its own blob without affecting the winner.
28//!
29//! See `new_long_term_revision` for the key construction.
30//!
31//! ## Compare-and-Swap
32//!
33//! All mutating operations follow a common pattern of reading the current
34//! revision, performing the upload, atomically swapping the revision (commit
35//! point), and cleaning up the now-unreferenced LT blob in the background:
36//!
37//! ### Large-Object Write (> 1 MiB)
38//!
39//! 1. **Read HV** to capture the current revision (existing tombstone target,
40//!    or absent).
41//! 2. **Write payload to LT** at a unique revision key.
42//! 3. **Compare-and-swap in HV**: write a tombstone pointing to the new
43//!    revision, only if the current revision still matches step 1.
44//!    - **OK** — schedule background deletion of the old LT blob, if any.
45//!    - **Conflict** — another writer won the race; schedule background deletion
46//!      of our new LT blob.
47//!    - **Error** — reload the tombstone and delete the unreferenced blob or
48//!      blobs.
49//!
50//! ### Small-Object Write (≤ 1 MiB)
51//!
52//! 1. **Write inline to HV**, skipping the write if a tombstone is present.
53//!    - **OK** — done; the object is stored entirely in HV.
54//!    - **Tombstone present** — a large object already occupies this key;
55//!      continue:
56//! 2. **Compare-and-swap in HV**: replace the tombstone with inline data, only
57//!    if the tombstone's revision still matches.
58//!    - **OK** — schedule background deletion of the old LT blob.
59//!    - **Conflict** — another writer won the race; they will clean up the
60//!      LT blob and we have no new LT blob to clean up.
61//!    - **Error** — reload the tombstone and delete the unreferenced blob if
62//!      the write went through.
63//!
64//! ### Delete
65//!
66//! 1. **Delete from HV** if the entry is not a tombstone.
67//!    - **OK** — done; there is no LT data to clean up.
68//!    - **Tombstone present** — a large object is stored here; continue:
69//! 2. **Compare-and-swap in HV**: remove the tombstone, only if its revision
70//!    still matches.
71//!    - **OK** — schedule background deletion of the LT blob.
72//!    - **Conflict** — another writer won the race; they will clean up.
73//!    - **Error** — reload the tombstone and delete the unreferenced blob if
74//!      the write went through.
75//!
76//! Tombstone removal is the commit point for deletes. If the subsequent LT
77//! cleanup fails, an orphan blob remains but the object is already unreachable
78//! through the normal read path.
79//!
80//! ## Last-Writer-Wins
81//!
82//! Concurrent mutations on the same key are inherently a race. Even a write
83//! that returns `Ok` may be immediately overwritten by another caller — there
84//! is no ordering guarantee and objectstore cannot provide a read-your-writes
85//! promise.
86//!
87//! CAS conflicts are therefore **not errors**: the losing writer's data is
88//! cleaned up and `Ok` is returned, because the result is indistinguishable
89//! from having succeeded a moment earlier and then been overwritten.
90//!
91//! ### Idempotency
92//!
93//! `compare_and_write` is idempotent: if the row is already in the target state, it
94//! returns `true` without re-applying the mutation. This is critical for retry
95//! safety. If the server commits a write but the response is lost, a retry sees the
96//! already-mutated state and still returns `true` — so callers do not mistakenly
97//! treat a successful commit as a lost race and clean up data that was actually
98//! persisted.
99//!
100//! # Resumable Uploads
101//!
102//! Resumable uploads are accepted only when their declared size exceeds 1 MiB and are
103//! written to the long-term backend. A resumable upload remains inaccessible until
104//! the long-term backend completes and a high-volume tombstone is committed.
105//!
106//! Each upload has a separate HV marker with a fixed five-day lifetime that serves
107//! as the source of truth for whether the upload logically exists.
108//! This marker is consumed upon cancellation or the first submission of a final chunk, making
109//! finalization one-shot.
110
111use std::num::NonZeroU64;
112use std::sync::Arc;
113use std::sync::atomic::Ordering;
114use std::time::{Duration, SystemTime};
115
116use base64::Engine as _;
117use bytes::Bytes;
118use futures_util::StreamExt;
119use objectstore_types::metadata::Metadata;
120use objectstore_types::range::ByteRange;
121use objectstore_types::resumable::UploadProgress;
122use objectstore_types::time::Timestamp;
123use sentry::{Hub, SentryFutureExt};
124use serde::{Deserialize, Serialize};
125
126use crate::backend::changelog::{Change, ChangeGuard, ChangeLog, ChangeManager, ChangePhase};
127use crate::backend::common::{
128    Backend, DeleteResponse, ExpiryTarget, ExpiryUpdate, GetResponse, HighVolumeBackend,
129    MetadataResponse, MultipartUploadBackend, PutResponse, SetExpiryResponse, TieredGet,
130    TieredMetadata, TieredUpdate, TieredWrite, Tombstone,
131};
132use crate::backend::{HighVolumeStorageConfig, MultipartUploadStorageConfig};
133use crate::error::{Error, ErrorKind, Result, ResultExt as _};
134use crate::id::ObjectId;
135use crate::multipart::{
136    AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
137    ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
138};
139use crate::resumable::{BackendToken, Session};
140use crate::stream::{ClientStream, SizedPeek, counting_stream};
141
142/// The threshold up until which we will go to the "high volume" backend.
143const BACKEND_SIZE_THRESHOLD: usize = 1024 * 1024; // 1 MiB
144
145/// Fixed lifetime of an ongoing upload marker.
146const RESUMABLE_UPLOAD_TTL: Duration = Duration::from_hours(5 * 24);
147
148/// Amount of time for which a `Change` generated by a `complete_multipart` operation is kept in the `Assembling`
149/// state before becoming eligible for cleanup by the `ChangeLog` recovery process.
150/// This allows the client to retry the `complete_multipart` operation upon any failures for at least this long,
151/// avoiding scenarios where the `ChangeLog` recovery would race to delete the assembled LT blob.
152const MULTIPART_COMPLETE_CLEANUP_DELAY: Duration = Duration::from_hours(24);
153
154/// Creates a new [`ObjectId`] with the same context but a unique revision key.
155///
156/// The new key has the format `{original_key}/{uuid_v7}`, producing a distinct
157/// storage path for each large-object write. [`ObjectId::from_storage_path`] parses
158/// the result back correctly because the key portion may contain `/`.
159fn new_long_term_revision(id: &ObjectId) -> ObjectId {
160    ObjectId {
161        context: id.context.clone(),
162        key: format!("{}/{}", id.key, uuid::Uuid::now_v7()),
163    }
164}
165
166/// Configuration for [`TieredStorage`].
167///
168/// Composes two backends into a tiered routing setup: `high_volume` for small
169/// objects and `long_term` for large objects. Nesting [`super::StorageConfig::Tiered`]
170/// inside another tiered config is not supported.
171///
172/// # Example
173///
174/// ```yaml
175/// storage:
176///   type: tiered
177///   high_volume:
178///     type: bigtable
179///     project_id: my-project
180///     instance_name: objectstore
181///     table_name: objectstore
182///   long_term:
183///     type: gcs
184///     bucket: my-objectstore-bucket
185/// ```
186#[derive(Debug, Clone, Deserialize, Serialize)]
187pub struct TieredStorageConfig {
188    /// Backend for high-volume, small objects.
189    ///
190    /// Must be a backend that implements [`HighVolumeBackend`] (currently
191    /// only BigTable).
192    pub high_volume: HighVolumeStorageConfig,
193    /// Backend for large, long-term objects.
194    ///
195    /// Must be a backend that implements [`MultipartUploadBackend`].
196    pub long_term: MultipartUploadStorageConfig,
197}
198
199/// Two-tier storage backend that routes objects by size.
200///
201/// `TieredStorage` implements [`Backend`] and is intended to be used inside a
202/// [`StorageService`](crate::StorageService), which wraps it with task spawning and panic
203/// isolation.
204///
205/// # Size-Based Routing
206///
207/// Objects are routed at write time based on their size relative to a **1 MiB threshold**:
208///
209/// - Objects **≤ 1 MiB** go to the `high_volume` backend — optimized for low-latency reads
210///   and writes of small objects (e.g. BigTable).
211/// - Objects **> 1 MiB** go to the `long_term` backend — optimized for cost-efficient
212///   storage of large objects (e.g. GCS).
213/// - Resumable uploads at or below 1 MiB are declined; larger uploads go to `long_term`.
214///
215/// # Redirect Tombstones
216///
217/// Because the [`ObjectId`] is backend-independent, reads must be able to find an object
218/// without knowing which backend stores it. A naive approach would check the long-term
219/// backend on every read miss in the high-volume backend — but that is slow and expensive.
220///
221/// Instead, when an object is stored in the long-term backend, a **redirect tombstone** is
222/// written in the high-volume backend. It acts as a signpost: "the real data lives in the
223/// other backend at this target." On reads, a single high-volume lookup either returns the
224/// object directly or follows the tombstone to long-term storage, without probing both
225/// backends.
226///
227/// How tombstones are physically stored is determined by the [`HighVolumeBackend`]
228/// implementation — refer to the backend's own documentation for storage format details.
229///
230/// # Consistency
231///
232/// Consistency across the two backends is maintained through compare-and-swap
233/// operations on the high-volume backend (see
234/// [`HighVolumeBackend::compare_and_write`]), not distributed locks. Each
235/// mutating operation reads the current high-volume revision, performs its
236/// work, and then atomically swaps the high-volume entry only if the revision
237/// is still current — rolling back on conflict. Cleanup of unreferenced LT
238/// blobs runs in background tasks so the caller returns as soon as the commit
239/// point is reached. Call [`Backend::join`] during shutdown to wait for
240/// outstanding cleanup.
241///
242/// See the [module-level documentation](self) for per-operation diagrams.
243///
244/// # Usage
245///
246/// `TieredStorage` handles only the routing and consistency logic. Wrap it in a
247/// [`StorageService`](crate::service::StorageService) to add task spawning, panic isolation,
248/// and concurrency limiting.
249#[derive(Debug)]
250pub struct TieredStorage {
251    inner: Arc<ChangeManager>,
252}
253
254impl TieredStorage {
255    /// Creates a new `TieredStorage` with the given backends and change log.
256    pub fn new(
257        high_volume: Box<dyn HighVolumeBackend>,
258        long_term: Box<dyn MultipartUploadBackend>,
259        changelog: Box<dyn ChangeLog>,
260    ) -> Self {
261        let inner = ChangeManager::new(high_volume, long_term, changelog);
262        let hub = Hub::new_from_top(Hub::current());
263        // Note on cancellation: Our `join` method will wait for all tasks tracked by the spawned
264        // recovery job, so we defer shutdown until recovery is complete or times out.
265        tokio::spawn(inner.clone().recover().bind_hub(hub));
266        Self { inner }
267    }
268
269    /// Records the change to the log and returns a guard that cleans up on drop.
270    async fn record_change(&self, change: Change) -> Result<ChangeGuard> {
271        self.inner.clone().record(change).await
272    }
273
274    /// Records the change to the log in the `Assembling` phase, and returns a guard that does
275    /// nothing on drop unless advanced.
276    async fn record_assembling(&self, change: Change) -> Result<ChangeGuard> {
277        self.inner.clone().record_assembling(change).await
278    }
279
280    /// Checks that the upload is still authorized before contacting LT.
281    async fn check_upload_marker(&self, revision: &ObjectId) -> Result<()> {
282        if !self
283            .inner
284            .high_volume
285            .has_upload_marker(revision, Timestamp::now())
286            .await?
287        {
288            return Err(ErrorKind::UploadSessionGone.into());
289        }
290        Ok(())
291    }
292
293    /// Consumes permission to finalize or cancel the upload.
294    async fn delete_upload_marker(&self, revision: &ObjectId) -> Result<()> {
295        if !self
296            .inner
297            .high_volume
298            .delete_upload_marker(revision, Timestamp::now())
299            .await?
300        {
301            return Err(ErrorKind::UploadSessionGone.into());
302        }
303        Ok(())
304    }
305
306    /// Returns the name of the backend corresponding to the given routing choice.
307    fn backend_type(&self, choice: &BackendChoice) -> &'static str {
308        match choice {
309            BackendChoice::HighVolume => self.inner.high_volume.name(),
310            BackendChoice::LongTerm => self.inner.long_term.name(),
311        }
312    }
313
314    /// Puts an object into the high-volume backend.
315    ///
316    /// If a tombstone already exists, attempts to swap it for the new object and delete the old
317    /// long-term object.
318    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
319    async fn put_high_volume(
320        &self,
321        id: &ObjectId,
322        metadata: &Metadata,
323        payload: Bytes,
324        access_time: Timestamp,
325    ) -> Result<()> {
326        let tombstone_opt = self
327            .inner
328            .high_volume
329            .put_non_tombstone(id, metadata, payload.clone(), access_time)
330            .await?;
331
332        let Some(Tombstone { target, .. }) = tombstone_opt else {
333            // No tombstone exists - write succeeded
334            return Ok(());
335        };
336
337        // Tombstone exists — Swap it for inline data
338        let mut guard = self
339            .record_change(Change {
340                id: id.clone(),
341                new: None,
342                old: Some(target.clone()),
343                cleanup_after: None,
344            })
345            .await?;
346
347        let write = TieredWrite::Object(metadata.clone(), payload);
348        guard.advance(ChangePhase::Written);
349
350        let written = self
351            .inner
352            .high_volume
353            .compare_and_write(id, Some(&target), write, access_time)
354            .await?;
355
356        // Update guard and let it schedule cleanup in the background.
357        guard.advance(ChangePhase::compare_and_write(written));
358
359        Ok(())
360    }
361
362    /// Puts an object into the long-term backend with a redirect tombstone in front.
363    ///
364    /// Deletes the previous long-term object if overwriting an existing tombstone. If the tombstone
365    /// write fails, the new long-term object is cleaned up.
366    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
367    async fn put_long_term(
368        &self,
369        id: &ObjectId,
370        metadata: &Metadata,
371        stream: ClientStream,
372        access_time: Timestamp,
373    ) -> Result<()> {
374        // 1. Read current HV revision to establish the write precondition
375        let current = match self
376            .inner
377            .high_volume
378            .get_tiered_metadata(id, access_time)
379            .await?
380        {
381            TieredMetadata::Tombstone(t) => Some(t.target),
382            _ => None,
383        };
384
385        // 2. Write payload to long-term at a unique revision key.
386        let new = new_long_term_revision(id);
387        let mut guard = self
388            .record_change(Change {
389                id: id.clone(),
390                new: Some(new.clone()),
391                old: current.clone(),
392                cleanup_after: None,
393            })
394            .await?;
395
396        self.inner
397            .long_term
398            .put_object(&new, metadata, stream, access_time)
399            .await?;
400        guard.advance(ChangePhase::Written);
401
402        // 3. CAS commit: write tombstone only if HV state matches what we saw.
403        let tombstone = Tombstone {
404            target: new.clone(),
405            time_expires: metadata.time_expires,
406        };
407        let written = self
408            .inner
409            .high_volume
410            .compare_and_write(
411                id,
412                current.as_ref(),
413                TieredWrite::Tombstone(tombstone),
414                access_time,
415            )
416            .await?;
417
418        // Update guard and let it schedule cleanup in the background.
419        guard.advance(ChangePhase::compare_and_write(written));
420
421        Ok(())
422    }
423}
424
425type LongTermBackendToken = BackendToken;
426
427#[derive(Debug, Serialize, Deserialize)]
428struct TieredResumableToken {
429    inner: LongTermBackendToken,
430    revision: String,
431    time_expires: Option<Timestamp>,
432}
433
434impl TieredResumableToken {
435    fn decode(token: &BackendToken) -> Result<Self> {
436        serde_json::from_str(token).map_err(|_| ErrorKind::UnknownUploadSession.into())
437    }
438
439    /// Builds the long-term session with the revision ID and inner token, retaining shared fields.
440    fn into_inner_session(self, session: &Session) -> Session {
441        Session {
442            object_id: ObjectId {
443                context: session.object_id.context.clone(),
444                key: self.revision,
445            },
446            backend_token: self.inner,
447            ..session.clone()
448        }
449    }
450}
451
452#[async_trait::async_trait]
453impl Backend for TieredStorage {
454    fn name(&self) -> &'static str {
455        "tiered"
456    }
457
458    fn upload_granularity(&self) -> u64 {
459        self.inner.long_term.upload_granularity()
460    }
461
462    fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
463        Ok(self)
464    }
465
466    #[tracing::instrument(level = "debug", fields(?id, upload_length), skip_all)]
467    async fn create_upload_session(
468        &self,
469        id: &ObjectId,
470        metadata: &Metadata,
471        upload_length: NonZeroU64,
472    ) -> Result<Option<BackendToken>> {
473        if upload_length.get() <= BACKEND_SIZE_THRESHOLD as u64 {
474            return Ok(None);
475        }
476
477        let revision = new_long_term_revision(id);
478        let Some(inner) = self
479            .inner
480            .long_term
481            .create_upload_session(&revision, metadata, upload_length)
482            .await?
483        else {
484            return Ok(None);
485        };
486
487        if let Err(error) = self
488            .inner
489            .high_volume
490            .create_upload_marker(&revision, Timestamp::now() + RESUMABLE_UPLOAD_TTL)
491            .await
492        {
493            // Best-effort clean-up.
494            let session = Session {
495                object_id: revision,
496                upload_length,
497                backend_token: inner,
498            };
499            if let Err(cleanup_error) = self.inner.long_term.cancel_upload(&session).await {
500                objectstore_log::warn!(
501                    !!&cleanup_error,
502                    "Failed to cancel upload after marker creation failed"
503                );
504            }
505            return Err(error);
506        }
507        let token = TieredResumableToken {
508            revision: revision.key,
509            inner,
510            time_expires: metadata.time_expires,
511        };
512        Ok(Some(serde_json::to_string(&token).context(
513            ErrorKind::Internal,
514            "encoding tiered resumable session",
515        )?))
516    }
517
518    #[tracing::instrument(level = "debug", fields(?session, offset, content_length), skip_all)]
519    async fn put_chunk(
520        &self,
521        session: &Session,
522        offset: u64,
523        content_length: u64,
524        stream: ClientStream,
525    ) -> Result<UploadProgress> {
526        let tiered = TieredResumableToken::decode(&session.backend_token)?;
527        let time_expires = tiered.time_expires;
528        let inner_session = tiered.into_inner_session(session);
529        let id = &session.object_id;
530        let revision = &inner_session.object_id;
531        let end = offset
532            .checked_add(content_length)
533            .filter(|end| *end <= session.upload_length.get())
534            .ok_or(ErrorKind::ChunkExceedsUploadLength {
535                offset,
536                content_length,
537                upload_length: session.upload_length.get(),
538            })?;
539
540        let granularity = self.upload_granularity();
541        if content_length > 0 && content_length < granularity && end != session.upload_length.get()
542        {
543            return Err(ErrorKind::ChunkTooSmall {
544                chunk_length: content_length,
545                upload_granularity: granularity,
546            }
547            .into());
548        }
549
550        // Non-final request; just forward the chunk.
551        if end != session.upload_length.get() {
552            let progress = self
553                .inner
554                .long_term
555                .put_chunk(&inner_session, offset, content_length, stream)
556                .await?;
557            return match progress {
558                UploadProgress::Incomplete { .. } => Ok(progress),
559                UploadProgress::Complete => Err(ErrorKind::UploadSessionGone.into()),
560            };
561        }
562
563        // Read the publication precondition while verifying the LT offset.
564        let (progress, current) = tokio::join!(
565            self.inner.long_term.upload_offset(&inner_session),
566            self.inner
567                .high_volume
568                .get_tiered_metadata(id, Timestamp::now()),
569        );
570        match progress? {
571            UploadProgress::Complete => return Err(ErrorKind::UploadSessionGone.into()),
572            UploadProgress::Incomplete { offset: actual } if actual < offset => {
573                return Err(ErrorKind::UploadOffsetMismatch { offset: actual }.into());
574            }
575            UploadProgress::Incomplete { .. } => {}
576        }
577        let current = match current? {
578            TieredMetadata::Tombstone(t) if t.target == *revision => {
579                return Ok(UploadProgress::Complete);
580            }
581            TieredMetadata::Tombstone(t) => Some(t.target),
582            _ => None,
583        };
584
585        // Delete the marker to claim this upload.
586        self.delete_upload_marker(revision).await?;
587
588        // Follow the regular PUT flow now that this revision has a single owner.
589        let mut guard = self
590            .record_change(Change {
591                id: id.clone(),
592                new: Some(revision.clone()),
593                old: current.clone(),
594                cleanup_after: None,
595            })
596            .await?;
597
598        let progress = self
599            .inner
600            .long_term
601            .put_chunk(&inner_session, offset, content_length, stream)
602            .await?;
603        if progress != UploadProgress::Complete {
604            return Err(ErrorKind::UploadSessionGone.into());
605        }
606        guard.advance(ChangePhase::Written);
607
608        let written = self
609            .inner
610            .high_volume
611            .compare_and_write(
612                id,
613                current.as_ref(),
614                TieredWrite::Tombstone(Tombstone {
615                    target: revision.clone(),
616                    time_expires,
617                }),
618                Timestamp::now(),
619            )
620            .await?;
621        guard.advance(ChangePhase::compare_and_write(written));
622        Ok(UploadProgress::Complete)
623    }
624
625    #[tracing::instrument(level = "debug", fields(?session), skip_all)]
626    async fn upload_offset(&self, session: &Session) -> Result<UploadProgress> {
627        let tiered = TieredResumableToken::decode(&session.backend_token)?;
628        let inner_session = tiered.into_inner_session(session);
629        self.check_upload_marker(&inner_session.object_id).await?;
630        match self.inner.long_term.upload_offset(&inner_session).await? {
631            UploadProgress::Incomplete { offset } => Ok(UploadProgress::Incomplete { offset }),
632            UploadProgress::Complete => Err(ErrorKind::UploadSessionGone.into()),
633        }
634    }
635
636    #[tracing::instrument(level = "debug", fields(?session), skip_all)]
637    async fn cancel_upload(&self, session: &Session) -> Result<()> {
638        let tiered = TieredResumableToken::decode(&session.backend_token)?;
639        let inner_session = tiered.into_inner_session(session);
640        self.delete_upload_marker(&inner_session.object_id).await?;
641        if let Err(error) = self.inner.long_term.cancel_upload(&inner_session).await {
642            objectstore_log::warn!(!!&error, "Failed to cancel upload after deleting marker");
643        }
644        Ok(())
645    }
646
647    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
648    async fn put_object(
649        &self,
650        id: &ObjectId,
651        metadata: &Metadata,
652        stream: ClientStream,
653        access_time: Timestamp,
654    ) -> Result<PutResponse> {
655        let timer = objectstore_metrics::timer!("put.latency", usecase = id.usecase().to_owned());
656        if metadata.origin.is_none() {
657            objectstore_metrics::count!("put.origin_missing", usecase = id.usecase().to_owned());
658        }
659
660        let peeked = SizedPeek::new(stream, BACKEND_SIZE_THRESHOLD).await?;
661        objectstore_metrics::record!(
662            "put.first_chunk.latency" = timer.elapsed(),
663            usecase = id.usecase().to_owned(),
664            complete = if peeked.is_exhausted() { "yes" } else { "no" },
665        );
666
667        let (backend_choice, stored_size) = if peeked.is_exhausted() {
668            let payload = peeked.into_bytes().await?;
669            let payload_len = payload.len() as u64;
670            self.put_high_volume(id, metadata, payload, access_time)
671                .await?;
672            (BackendChoice::HighVolume, payload_len)
673        } else {
674            let (stored_size, stream) = counting_stream(peeked.into_stream());
675            self.put_long_term(id, metadata, stream.boxed(), access_time)
676                .await?;
677            (BackendChoice::LongTerm, stored_size.load(Ordering::Acquire))
678        };
679
680        let backend_ty = self.backend_type(&backend_choice);
681        timer
682            .tag("backend_choice", backend_choice.as_str())
683            .tag("backend_type", backend_ty)
684            .record();
685        objectstore_metrics::record!(
686            "put.size" = stored_size,
687            usecase = id.usecase().to_owned(),
688            backend_choice = backend_choice.as_str(),
689            backend_type = backend_ty,
690            upload_type = "direct",
691        );
692
693        Ok(())
694    }
695
696    #[tracing::instrument(level = "debug", skip(self))]
697    async fn get_object(
698        &self,
699        id: &ObjectId,
700        access_time: Timestamp,
701        range: Option<ByteRange>,
702    ) -> Result<GetResponse> {
703        let timer = objectstore_metrics::timer!(
704            "get.latency.pre-response",
705            usecase = id.usecase().to_owned(),
706        );
707
708        let hv_result = self
709            .inner
710            .high_volume
711            .get_tiered_object(id, access_time, range)
712            .await?;
713        let (result, backend_choice) = match hv_result {
714            TieredGet::NotFound => (None, BackendChoice::HighVolume),
715            TieredGet::Object(metadata, content_range, stream) => (
716                Some((metadata, content_range, stream)),
717                BackendChoice::HighVolume,
718            ),
719            TieredGet::Tombstone(tombstone) => (
720                self.inner
721                    .long_term
722                    .get_object(&tombstone.target, access_time, range)
723                    .await?
724                    .map(|(meta, range, stream)| (align_expiry(meta, &tombstone), range, stream)),
725                BackendChoice::LongTerm,
726            ),
727        };
728
729        let backend_type = self.backend_type(&backend_choice);
730        timer
731            .tag("backend_choice", backend_choice.as_str())
732            .tag("backend_type", backend_type)
733            .record();
734
735        if let Some((ref metadata, ref content_range, _)) = result {
736            let size = content_range.map(|cr| cr.len() as usize).or(metadata.size);
737            if let Some(size) = size {
738                objectstore_metrics::record!(
739                    "get.size" = size,
740                    usecase = id.usecase().to_owned(),
741                    backend_choice = backend_choice.as_str(),
742                    backend_type = backend_type,
743                );
744            }
745        }
746
747        Ok(result)
748    }
749
750    #[tracing::instrument(level = "debug", skip(self))]
751    async fn get_metadata(
752        &self,
753        id: &ObjectId,
754        access_time: Timestamp,
755    ) -> Result<MetadataResponse> {
756        let timer = objectstore_metrics::timer!("head.latency", usecase = id.usecase().to_owned());
757
758        let hv_result = self
759            .inner
760            .high_volume
761            .get_tiered_metadata(id, access_time)
762            .await?;
763        let (result, backend_choice) = match hv_result {
764            TieredMetadata::NotFound => (None, BackendChoice::HighVolume),
765            TieredMetadata::Object(metadata) => (Some(metadata), BackendChoice::HighVolume),
766            TieredMetadata::Tombstone(tombstone) => (
767                self.inner
768                    .long_term
769                    .get_metadata(&tombstone.target, access_time)
770                    .await?
771                    .map(|metadata| align_expiry(metadata, &tombstone)),
772                BackendChoice::LongTerm,
773            ),
774        };
775
776        timer
777            .tag("backend_choice", backend_choice.as_str())
778            .tag("backend_type", self.backend_type(&backend_choice))
779            .record();
780
781        Ok(result)
782    }
783
784    async fn set_expiry(
785        &self,
786        id: &ObjectId,
787        target: ExpiryUpdate,
788        access_time: Timestamp,
789    ) -> Result<SetExpiryResponse> {
790        match self
791            .inner
792            .high_volume
793            .get_tiered_metadata(id, access_time)
794            .await?
795        {
796            TieredMetadata::NotFound => Ok(SetExpiryResponse::NotFound),
797            TieredMetadata::Object(_) => {
798                self.inner
799                    .high_volume
800                    .compare_and_update(id, None, TieredUpdate::SetExpiry(target), access_time)
801                    .await
802            }
803            TieredMetadata::Tombstone(tombstone) => {
804                // Extend LT first. Extending the redirect first could leave it
805                // alive after the blob failed to extend and was reclaimed.
806                let deadline = match self
807                    .inner
808                    .long_term
809                    .set_expiry(&tombstone.target, target, access_time)
810                    .await?
811                {
812                    SetExpiryResponse::Satisfied(deadline) => deadline,
813                    outcome @ (SetExpiryResponse::NotFound | SetExpiryResponse::Rejected) => {
814                        return Ok(outcome);
815                    }
816                };
817
818                // Limits have already been validated at the authoritative LT object.
819                let update = TieredUpdate::SetExpiry(ExpiryTarget::At(deadline).into());
820
821                // NOTE: If this fails, LT may remain extended while the redirect
822                // becomes unreachable earlier. Rolling LT back could interfere
823                // with another renewal that succeeded concurrently. Propagate
824                // the resolved request, not LT's stored deadline, so an
825                // already-later blob cannot over-extend the redirect.
826                self.inner
827                    .high_volume
828                    .compare_and_update(id, Some(&tombstone.target), update, access_time)
829                    .await
830            }
831        }
832    }
833
834    #[tracing::instrument(level = "debug", skip(self))]
835    async fn delete_object(&self, id: &ObjectId, access_time: Timestamp) -> Result<DeleteResponse> {
836        let timer =
837            objectstore_metrics::timer!("delete.latency", usecase = id.usecase().to_owned());
838
839        let mut backend_choice = BackendChoice::HighVolume;
840
841        if let Some(tombstone) = self
842            .inner
843            .high_volume
844            .delete_non_tombstone(id, access_time)
845            .await?
846        {
847            backend_choice = BackendChoice::LongTerm;
848
849            let mut guard = self
850                .record_change(Change {
851                    id: id.clone(),
852                    new: None,
853                    old: Some(tombstone.target.clone()),
854                    cleanup_after: None,
855                })
856                .await?;
857            guard.advance(ChangePhase::Written);
858
859            // Remove the tombstone; the LT blob becomes unreachable at this point.
860            let deleted = self
861                .inner
862                .high_volume
863                .compare_and_write(
864                    id,
865                    Some(&tombstone.target),
866                    TieredWrite::Delete,
867                    access_time,
868                )
869                .await?;
870
871            // Update guard and let it schedule cleanup in the background.
872            guard.advance(ChangePhase::compare_and_write(deleted));
873        }
874
875        timer
876            .tag("backend_choice", backend_choice.as_str())
877            .tag("backend_type", self.backend_type(&backend_choice))
878            .record();
879
880        Ok(())
881    }
882
883    async fn join(&self) {
884        self.inner.tracker.close();
885        tokio::join!(
886            self.inner.high_volume.join(),
887            self.inner.long_term.join(),
888            self.inner.tracker.wait()
889        );
890    }
891}
892
893#[derive(Debug)]
894enum BackendChoice {
895    HighVolume,
896    LongTerm,
897}
898
899impl BackendChoice {
900    fn as_str(&self) -> &'static str {
901        match self {
902            BackendChoice::HighVolume => "high-volume",
903            BackendChoice::LongTerm => "long-term",
904        }
905    }
906}
907
908impl std::fmt::Display for BackendChoice {
909    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
910        f.write_str(self.as_str())
911    }
912}
913
914/// Returns the lower expiry between the redirect and blob, if any.
915///
916/// This is used to ensure correct expiry if they ever drift.
917fn effective_expiry(
918    redirect_expiry: Option<Timestamp>,
919    blob_expiry: Option<Timestamp>,
920) -> Option<Timestamp> {
921    match (redirect_expiry, blob_expiry) {
922        (Some(redirect), Some(blob)) => Some(redirect.min(blob)),
923        (Some(expiry), None) | (None, Some(expiry)) => Some(expiry),
924        (None, None) => None,
925    }
926}
927
928/// Aligns the expiry of the metadata with the expiry of the tombstone.
929///
930/// Keeps the lower expiry between the metadata and the tombstone so that the client sees the most
931/// conservative expiry and automatic expiry bumps still occur.
932fn align_expiry(mut metadata: Metadata, tombstone: &Tombstone) -> Metadata {
933    metadata.time_expires = effective_expiry(tombstone.time_expires, metadata.time_expires);
934    metadata
935}
936
937/// The multipart upload state for TieredStorage.
938#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
939struct TieredUploadId {
940    revision: String,
941    upload_id: UploadId,
942}
943
944impl TryInto<UploadId> for TieredUploadId {
945    type Error = Error;
946
947    fn try_into(self) -> Result<UploadId, Self::Error> {
948        let json =
949            serde_json::to_vec(&self).context(ErrorKind::Internal, "encoding tiered upload ID")?;
950        Ok(UploadId::new(
951            base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(json),
952        )?)
953    }
954}
955
956impl TryFrom<&UploadId> for TieredUploadId {
957    type Error = Error;
958
959    fn try_from(value: &UploadId) -> Result<Self, Self::Error> {
960        let json = base64::engine::general_purpose::URL_SAFE_NO_PAD
961            .decode(value.as_bytes())
962            .kind(ErrorKind::InvalidUploadId)?;
963        serde_json::from_slice(&json).kind(ErrorKind::InvalidUploadId)
964    }
965}
966
967#[async_trait::async_trait]
968impl MultipartUploadBackend for TieredStorage {
969    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
970    async fn initiate_multipart(
971        &self,
972        id: &ObjectId,
973        metadata: &Metadata,
974    ) -> Result<InitiateMultipartResponse> {
975        let timer = objectstore_metrics::timer!(
976            "multipart.initiate.latency",
977            usecase = id.usecase().to_owned(),
978        );
979        let physical = new_long_term_revision(id);
980
981        let upload_id = self
982            .inner
983            .long_term
984            .initiate_multipart(&physical, metadata)
985            .await?;
986
987        let id = TieredUploadId {
988            revision: physical.key,
989            upload_id,
990        };
991        let id = id.try_into()?;
992
993        timer.record();
994        Ok(id)
995    }
996
997    #[tracing::instrument(level = "debug", fields(?id, part_number, content_length), skip_all)]
998    async fn upload_part(
999        &self,
1000        id: &ObjectId,
1001        upload_id: &UploadId,
1002        part_number: PartNumber,
1003        content_length: u64,
1004        content_md5: Option<&str>,
1005        body: ClientStream,
1006    ) -> Result<UploadPartResponse> {
1007        let timer = objectstore_metrics::timer!(
1008            "multipart.upload_part.latency",
1009            usecase = id.usecase().to_owned(),
1010        );
1011        let tiered: TieredUploadId = upload_id.try_into()?;
1012
1013        let physical = ObjectId {
1014            context: id.context.clone(),
1015            key: tiered.revision,
1016        };
1017
1018        let etag = self
1019            .inner
1020            .long_term
1021            .upload_part(
1022                &physical,
1023                &tiered.upload_id,
1024                part_number,
1025                content_length,
1026                content_md5,
1027                body,
1028            )
1029            .await?;
1030
1031        timer.record();
1032        objectstore_metrics::record!(
1033            "multipart.upload_part.size" = content_length,
1034            usecase = id.usecase().to_owned(),
1035        );
1036
1037        Ok(etag)
1038    }
1039
1040    #[tracing::instrument(level = "debug", skip(self, upload_id))]
1041    async fn list_parts(
1042        &self,
1043        id: &ObjectId,
1044        upload_id: &UploadId,
1045        max_parts: Option<u32>,
1046        part_number_marker: Option<PartNumber>,
1047    ) -> Result<ListPartsResponse> {
1048        let timer = objectstore_metrics::timer!(
1049            "multipart.list_parts.latency",
1050            usecase = id.usecase().to_owned(),
1051        );
1052        let tiered: TieredUploadId = upload_id.try_into()?;
1053
1054        let physical = ObjectId {
1055            context: id.context.clone(),
1056            key: tiered.revision,
1057        };
1058
1059        let response = self
1060            .inner
1061            .long_term
1062            .list_parts(&physical, &tiered.upload_id, max_parts, part_number_marker)
1063            .await?;
1064
1065        timer.record();
1066        Ok(response)
1067    }
1068
1069    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1070    async fn abort_multipart(
1071        &self,
1072        id: &ObjectId,
1073        upload_id: &UploadId,
1074    ) -> Result<AbortMultipartResponse> {
1075        let timer = objectstore_metrics::timer!(
1076            "multipart.abort.latency",
1077            usecase = id.usecase().to_owned(),
1078        );
1079        let tiered: TieredUploadId = upload_id.try_into()?;
1080
1081        let physical = ObjectId {
1082            context: id.context.clone(),
1083            key: tiered.revision,
1084        };
1085
1086        let () = self
1087            .inner
1088            .long_term
1089            .abort_multipart(&physical, &tiered.upload_id)
1090            .await?;
1091
1092        timer.record();
1093        Ok(())
1094    }
1095
1096    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1097    async fn complete_multipart(
1098        &self,
1099        id: &ObjectId,
1100        upload_id: &UploadId,
1101        parts: Vec<CompletedPart>,
1102        access_time: Timestamp,
1103    ) -> Result<CompleteMultipartResponse> {
1104        let timer = objectstore_metrics::timer!(
1105            "multipart.complete.latency",
1106            usecase = id.usecase().to_owned(),
1107        );
1108        let part_count = parts.len();
1109        let tiered: TieredUploadId = upload_id.try_into()?;
1110
1111        let physical = ObjectId {
1112            context: id.context.clone(),
1113            key: tiered.revision,
1114        };
1115
1116        // 1. Read current HV revision to establish the write precondition.
1117        let current = match self
1118            .inner
1119            .high_volume
1120            .get_tiered_metadata(id, access_time)
1121            .await?
1122        {
1123            // Optimization: a previous attempt already finalized this revision and tombstone -- report success.
1124            TieredMetadata::Tombstone(t) if t.target == physical => {
1125                timer.record();
1126                return Ok(None);
1127            }
1128            TieredMetadata::Tombstone(t) => Some(t.target),
1129            _ => None,
1130        };
1131
1132        // Register a guard with cleanup deferred to now + `MULTIPART_COMPLETE_CLEANUP_DELAY`,
1133        // so that the user has the chance to retry finalizing the upload in this timeframe.
1134        let mut guard = self
1135            .record_assembling(Change {
1136                id: id.clone(),
1137                new: Some(physical.clone()),
1138                old: current.clone(),
1139                cleanup_after: Some(SystemTime::now() + MULTIPART_COMPLETE_CLEANUP_DELAY),
1140            })
1141            .await?;
1142
1143        // 2. Complete the upload, creating the object at the given revision key.
1144        let maybe_complete_multipart_err = match self
1145            .inner
1146            .long_term
1147            .complete_multipart(&physical, &tiered.upload_id, parts, access_time)
1148            .await
1149        {
1150            // The request went through but we got an error in the response body.
1151            // Transparently proxy the error to the user.
1152            Ok(error) => {
1153                if error.is_some() {
1154                    return Ok(error);
1155                }
1156                None
1157            }
1158            // We got status 4xx/5xx, or a network error.
1159            // Either way, `complete_multipart` might have been completed successfully,
1160            // either now or in a previous attempt (in that case, that's a 404 and we indeed end up
1161            // here).
1162            // We cannot know if that's the case yet, so we continue to the next steps.
1163            Err(err) => Some(err),
1164        };
1165
1166        // 3. Retrieve the metadata of the object, which was determined at initiation time, to
1167        //    get its expiration deadline and size.
1168        //
1169        //    This also serves as an existence check to understand if the LT revision was actually
1170        //    created successfully in this or a previous attempt, in which case we just need to
1171        //    finalize the tombstone.
1172        let metadata = self
1173            .inner
1174            .long_term
1175            .get_metadata(&physical, access_time)
1176            .await;
1177
1178        let metadata = match (metadata, maybe_complete_multipart_err) {
1179            // The LT revision already exists, so we can continue to finalize the tombstone.
1180            (Ok(Some(metadata)), _) => metadata,
1181            // The LT revision doesn't exist, cannot proceed.
1182            (Ok(None), Some(err)) => return Err(err),
1183            // The `complete_multipart` succeeded, creating the object, but the `get_metadata`
1184            // immediately after failed to find the object. This should never happen.
1185            (Ok(None), None) => {
1186                objectstore_log::error!(
1187                    id = ?id,
1188                    upload_id = ?upload_id,
1189                    physical = ?physical,
1190                    "complete_multipart call succeeded on long_term backend, but subsequent get_metadata found no object"
1191                );
1192                return Err(Error::new(
1193                    ErrorKind::BackendFailure,
1194                    "tiered multipart object missing from long-term storage",
1195                ));
1196            }
1197            // Failed to `get_metadata`, cannot proceed.
1198            (Err(get_metadata_err), maybe_complete_multipart_err) => {
1199                // Prefer the `complete_multipart_err`, as it's likely more informative.
1200                // TODO(FS-358): convert this properly. Right now `ApiErrorResponse` will turn this into a 500,
1201                // but we would actually want to transparently surface the original status (and message?) instead.
1202                return Err(maybe_complete_multipart_err.unwrap_or(get_metadata_err));
1203            }
1204        };
1205
1206        // 4. CAS commit: write tombstone only if HV state matches what we saw.
1207        let tombstone = Tombstone {
1208            target: physical.clone(),
1209            time_expires: metadata.time_expires,
1210        };
1211        let written = self
1212            .inner
1213            .high_volume
1214            .compare_and_write(
1215                id,
1216                current.as_ref(),
1217                TieredWrite::Tombstone(tombstone),
1218                access_time,
1219            )
1220            .await?;
1221
1222        // Update guard and let it schedule cleanup in the background.
1223        guard.advance(ChangePhase::compare_and_write(written));
1224
1225        timer.record();
1226        objectstore_metrics::record!(
1227            "multipart.complete.part_count" = part_count as u64,
1228            usecase = id.usecase().to_owned(),
1229        );
1230        if let Some(size) = metadata.size {
1231            objectstore_metrics::record!(
1232                "put.size" = size as u64,
1233                usecase = id.usecase().to_owned(),
1234                backend_choice = BackendChoice::LongTerm.as_str(),
1235                backend_type = self.backend_type(&BackendChoice::LongTerm),
1236                upload_type = "multipart",
1237            );
1238        }
1239
1240        Ok(None)
1241    }
1242}
1243
1244#[cfg(test)]
1245mod tests {
1246    use std::num::{NonZeroU32, NonZeroU64};
1247    use std::sync::Mutex as StdMutex;
1248
1249    use futures::lock::Mutex;
1250    use objectstore_types::metadata::{ExpirationPolicy, Metadata};
1251    use objectstore_types::scope::{Scope, Scopes};
1252
1253    use super::*;
1254    use crate::backend::bigtable::{BigTableBackend, BigTableConfig};
1255    use crate::backend::changelog::{InMemoryChangeLog, NoopChangeLog};
1256    use crate::backend::gcs::{GcsBackend, GcsConfig};
1257    use crate::backend::in_memory::InMemoryBackend;
1258    use crate::backend::testing::{Hooks, TestBackend};
1259    use crate::change_stream::ChangeStreamFactory;
1260    use crate::error::Error;
1261    use crate::id::ObjectContext;
1262    use crate::stream::{self, ClientStream};
1263
1264    fn make_context() -> ObjectContext {
1265        ObjectContext {
1266            usecase: "testing".into(),
1267            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1268        }
1269    }
1270
1271    fn make_id(key: &str) -> ObjectId {
1272        ObjectId::new(make_context(), key.into())
1273    }
1274
1275    fn make_tiered_storage() -> (
1276        TieredStorage,
1277        InMemoryBackend,
1278        InMemoryBackend,
1279        InMemoryChangeLog,
1280    ) {
1281        let hv = InMemoryBackend::new("in-memory-hv");
1282        let lt = InMemoryBackend::new("in-memory-lt");
1283        let changelog = InMemoryChangeLog::default();
1284        let storage = TieredStorage::new(
1285            Box::new(hv.clone()),
1286            Box::new(lt.clone()),
1287            Box::new(changelog.clone()),
1288        );
1289        (storage, hv, lt, changelog)
1290    }
1291
1292    async fn resumable_token(
1293        storage: &TieredStorage,
1294        id: &ObjectId,
1295        metadata: &Metadata,
1296        length: u64,
1297    ) -> Session {
1298        let backend_token = storage
1299            .create_upload_session(id, metadata, NonZeroU64::new(length).unwrap())
1300            .await
1301            .unwrap()
1302            .unwrap();
1303        Session {
1304            object_id: id.clone(),
1305            upload_length: NonZeroU64::new(length).unwrap(),
1306            backend_token,
1307        }
1308    }
1309
1310    fn upload_revision(session: &Session) -> ObjectId {
1311        TieredResumableToken::decode(&session.backend_token)
1312            .unwrap()
1313            .into_inner_session(session)
1314            .object_id
1315    }
1316
1317    #[tokio::test]
1318    async fn resumable_inmemory() -> anyhow::Result<()> {
1319        let (storage, _, _, _) = make_tiered_storage();
1320        let id = make_id("tiered-resumable-inmemory");
1321        let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1];
1322        let token =
1323            resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
1324
1325        let revision = upload_revision(&token);
1326        assert!(
1327            storage
1328                .inner
1329                .high_volume
1330                .has_upload_marker(&revision, Timestamp::now())
1331                .await?
1332        );
1333        assert_eq!(
1334            storage.upload_offset(&token).await?,
1335            UploadProgress::Incomplete { offset: 0 }
1336        );
1337        let split = 256 * 1024;
1338        assert_eq!(
1339            storage
1340                .put_chunk(
1341                    &token,
1342                    0,
1343                    split as u64,
1344                    stream::single(payload[..split].to_vec())
1345                )
1346                .await?,
1347            UploadProgress::Incomplete {
1348                offset: split as u64
1349            }
1350        );
1351        assert!(
1352            storage
1353                .inner
1354                .high_volume
1355                .has_upload_marker(&revision, Timestamp::now())
1356                .await?
1357        );
1358        assert_eq!(
1359            storage.upload_offset(&token).await?,
1360            UploadProgress::Incomplete {
1361                offset: split as u64
1362            }
1363        );
1364        assert!(storage.get_metadata(&id, Timestamp::now()).await?.is_none());
1365
1366        // A completed upload creates a logical object.
1367        assert_eq!(
1368            storage
1369                .put_chunk(
1370                    &token,
1371                    split as u64,
1372                    (payload.len() - split) as u64,
1373                    stream::single(payload[split..].to_vec())
1374                )
1375                .await?,
1376            UploadProgress::Complete
1377        );
1378        assert!(
1379            !storage
1380                .inner
1381                .high_volume
1382                .has_upload_marker(&revision, Timestamp::now())
1383                .await?
1384        );
1385        assert_eq!(
1386            storage.upload_offset(&token).await.unwrap_err().kind(),
1387            ErrorKind::UploadSessionGone
1388        );
1389        let (_, _, body) = storage
1390            .get_object(&id, Timestamp::now(), None)
1391            .await?
1392            .unwrap();
1393        assert_eq!(stream::read_to_vec(body).await?, payload);
1394        Ok(())
1395    }
1396
1397    #[tokio::test]
1398    async fn resumable_bigtable_and_gcs() -> anyhow::Result<()> {
1399        let streams = ChangeStreamFactory::default();
1400        let lt = GcsBackend::new(
1401            GcsConfig {
1402                endpoint: Some("http://localhost:8087".into()),
1403                bucket: "test-bucket".into(),
1404                cogs: None,
1405            },
1406            &streams,
1407        )
1408        .await?;
1409        let hv = BigTableBackend::new(
1410            BigTableConfig {
1411                endpoint: Some("localhost:8086".into()),
1412                project_id: "testing".into(),
1413                instance_name: "objectstore".into(),
1414                table_name: "objectstore".into(),
1415                connections: None,
1416                rpc_timeout: Duration::from_secs(2),
1417                cogs: None,
1418            },
1419            &streams,
1420        )
1421        .await?;
1422        let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(NoopChangeLog));
1423        let id = make_id(&format!("tiered-resumable-{}", uuid::Uuid::now_v7()));
1424        let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1];
1425        let token =
1426            resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
1427
1428        let error = storage
1429            .put_chunk(&token, 0, 1, stream::single("a"))
1430            .await
1431            .unwrap_err();
1432        assert_eq!(
1433            error.kind(),
1434            ErrorKind::ChunkTooSmall {
1435                chunk_length: 1,
1436                upload_granularity: 256 * 1024,
1437            }
1438        );
1439
1440        let revision = upload_revision(&token);
1441        assert!(
1442            storage
1443                .inner
1444                .high_volume
1445                .has_upload_marker(&revision, Timestamp::now())
1446                .await?
1447        );
1448        assert_eq!(
1449            storage.upload_offset(&token).await?,
1450            UploadProgress::Incomplete { offset: 0 }
1451        );
1452        let split = 256 * 1024;
1453        assert_eq!(
1454            storage
1455                .put_chunk(
1456                    &token,
1457                    0,
1458                    split as u64,
1459                    stream::single(payload[..split].to_vec())
1460                )
1461                .await?,
1462            UploadProgress::Incomplete {
1463                offset: split as u64
1464            }
1465        );
1466        assert!(
1467            storage
1468                .inner
1469                .high_volume
1470                .has_upload_marker(&revision, Timestamp::now())
1471                .await?
1472        );
1473        assert_eq!(
1474            storage.upload_offset(&token).await?,
1475            UploadProgress::Incomplete {
1476                offset: split as u64
1477            }
1478        );
1479        assert!(storage.get_metadata(&id, Timestamp::now()).await?.is_none());
1480
1481        // A completed upload creates a logical object.
1482        assert_eq!(
1483            storage
1484                .put_chunk(
1485                    &token,
1486                    split as u64,
1487                    (payload.len() - split) as u64,
1488                    stream::single(payload[split..].to_vec())
1489                )
1490                .await?,
1491            UploadProgress::Complete
1492        );
1493        assert!(
1494            !storage
1495                .inner
1496                .high_volume
1497                .has_upload_marker(&revision, Timestamp::now())
1498                .await?
1499        );
1500        assert_eq!(
1501            storage.upload_offset(&token).await.unwrap_err().kind(),
1502            ErrorKind::UploadSessionGone
1503        );
1504        let (_, _, body) = storage
1505            .get_object(&id, Timestamp::now(), None)
1506            .await?
1507            .unwrap();
1508        assert_eq!(stream::read_to_vec(body).await?, payload);
1509        Ok(())
1510    }
1511
1512    #[tokio::test]
1513    async fn resumable_invalid_chunks() {
1514        let (storage, _, _, _) = make_tiered_storage();
1515        let id = make_id("resumable-invalid");
1516        let length = BACKEND_SIZE_THRESHOLD as u64 + 1;
1517        let token = resumable_token(&storage, &id, &Metadata::default(), length).await;
1518
1519        // Overflow and future offsets are rejected.
1520        assert_eq!(
1521            storage
1522                .put_chunk(&token, u64::MAX, 1, stream::single("x"))
1523                .await
1524                .unwrap_err()
1525                .kind(),
1526            ErrorKind::ChunkExceedsUploadLength {
1527                offset: u64::MAX,
1528                content_length: 1,
1529                upload_length: length
1530            }
1531        );
1532        assert_eq!(
1533            storage
1534                .put_chunk(&token, length - 1, 1, stream::single("x"))
1535                .await
1536                .unwrap_err()
1537                .kind(),
1538            ErrorKind::UploadOffsetMismatch { offset: 0 }
1539        );
1540    }
1541
1542    #[tokio::test]
1543    async fn resumable_cancel() {
1544        let (storage, hv, _, _) = make_tiered_storage();
1545        let id = make_id("resumable-invalid");
1546        let token = resumable_token(
1547            &storage,
1548            &id,
1549            &Metadata::default(),
1550            BACKEND_SIZE_THRESHOLD as u64 + 1,
1551        )
1552        .await;
1553
1554        let revision = upload_revision(&token);
1555        assert!(
1556            hv.has_upload_marker(&revision, Timestamp::now())
1557                .await
1558                .unwrap()
1559        );
1560        storage.cancel_upload(&token).await.unwrap();
1561        assert!(
1562            !hv.has_upload_marker(&revision, Timestamp::now())
1563                .await
1564                .unwrap()
1565        );
1566        assert_eq!(
1567            storage.cancel_upload(&token).await.unwrap_err().kind(),
1568            ErrorKind::UploadSessionGone
1569        );
1570        assert_eq!(
1571            storage.upload_offset(&token).await.unwrap_err().kind(),
1572            ErrorKind::UploadSessionGone
1573        );
1574        assert!(!hv.contains(&id));
1575    }
1576
1577    #[derive(Clone, Debug)]
1578    struct ExpiryHook {
1579        label: &'static str,
1580        events: Arc<StdMutex<Vec<&'static str>>>,
1581        reject: bool,
1582    }
1583
1584    type ExpiryTestStorage = (
1585        TieredStorage,
1586        TestBackend<ExpiryHook>,
1587        TestBackend<ExpiryHook>,
1588        Arc<StdMutex<Vec<&'static str>>>,
1589    );
1590
1591    #[async_trait::async_trait]
1592    impl Hooks for ExpiryHook {
1593        async fn set_expiry(
1594            &self,
1595            inner: &InMemoryBackend,
1596            id: &ObjectId,
1597            target: ExpiryUpdate,
1598            access_time: Timestamp,
1599        ) -> Result<SetExpiryResponse> {
1600            self.events.lock().unwrap().push(self.label);
1601            if self.reject {
1602                Ok(SetExpiryResponse::Rejected)
1603            } else {
1604                inner.set_expiry(id, target, access_time).await
1605            }
1606        }
1607
1608        async fn compare_and_update(
1609            &self,
1610            inner: &InMemoryBackend,
1611            id: &ObjectId,
1612            current: Option<&ObjectId>,
1613            update: TieredUpdate,
1614            access_time: Timestamp,
1615        ) -> Result<SetExpiryResponse> {
1616            self.events.lock().unwrap().push(self.label);
1617            if self.reject {
1618                Ok(SetExpiryResponse::Rejected)
1619            } else {
1620                inner
1621                    .compare_and_update(id, current, update, access_time)
1622                    .await
1623            }
1624        }
1625    }
1626
1627    fn tiered_with_expiry_hooks(hv_reject: bool, lt_reject: bool) -> ExpiryTestStorage {
1628        let events = Arc::new(StdMutex::new(Vec::new()));
1629        let hv = TestBackend::new(ExpiryHook {
1630            label: "hv",
1631            events: Arc::clone(&events),
1632            reject: hv_reject,
1633        });
1634        let lt = TestBackend::new(ExpiryHook {
1635            label: "lt",
1636            events: Arc::clone(&events),
1637            reject: lt_reject,
1638        });
1639        let storage = TieredStorage::new(
1640            Box::new(hv.clone()),
1641            Box::new(lt.clone()),
1642            Box::new(NoopChangeLog),
1643        );
1644        (storage, hv, lt, events)
1645    }
1646
1647    async fn seed_redirect(
1648        hv: &InMemoryBackend,
1649        lt: &InMemoryBackend,
1650        id: &ObjectId,
1651        target: &ObjectId,
1652        expiry: Timestamp,
1653    ) {
1654        let metadata = Metadata {
1655            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_mins(10)),
1656            time_created: expiry.checked_sub(Duration::from_mins(10)),
1657            time_expires: Some(expiry),
1658            ..Default::default()
1659        };
1660        lt.put_object(
1661            target,
1662            &metadata,
1663            stream::single("payload"),
1664            Timestamp::now(),
1665        )
1666        .await
1667        .unwrap();
1668        hv.compare_and_write(
1669            id,
1670            None,
1671            TieredWrite::Tombstone(Tombstone {
1672                target: target.clone(),
1673                time_expires: metadata.time_expires,
1674            }),
1675            Timestamp::now(),
1676        )
1677        .await
1678        .unwrap();
1679    }
1680
1681    #[tokio::test]
1682    async fn set_expiry() {
1683        let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, false);
1684        let id = make_id("tiered-expiry-order");
1685        let target = new_long_term_revision(&id);
1686        let old_expiry = Timestamp::now() + Duration::from_mins(10);
1687        seed_redirect(&hv.inner, &lt.inner, &id, &target, old_expiry).await;
1688
1689        let created = old_expiry - Duration::from_mins(10);
1690        let requested = created + Duration::from_hours(1);
1691        let blob_expiry = requested + Duration::from_hours(1);
1692        lt.inner
1693            .put_object(
1694                &target,
1695                &Metadata {
1696                    expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
1697                    time_created: Some(created),
1698                    time_expires: Some(blob_expiry),
1699                    ..Default::default()
1700                },
1701                stream::single("payload"),
1702                Timestamp::now(),
1703            )
1704            .await
1705            .unwrap();
1706        let update = ExpiryUpdate {
1707            target: ExpiryTarget::FromCreation(Duration::from_hours(1)),
1708            max: Some(Duration::ZERO),
1709        };
1710        assert_eq!(
1711            storage
1712                .set_expiry(&id, update, created)
1713                .await
1714                .unwrap_err()
1715                .kind(),
1716            ErrorKind::InvalidMetadata,
1717        );
1718        assert_eq!(events.lock().unwrap().as_slice(), &["lt"]);
1719        events.lock().unwrap().clear();
1720        assert_eq!(
1721            hv.inner.get(&id).expect_tombstone().time_expires,
1722            Some(old_expiry)
1723        );
1724        assert_eq!(
1725            storage
1726                .set_expiry(
1727                    &id,
1728                    ExpiryUpdate {
1729                        max: Some(Duration::from_hours(1)),
1730                        ..update
1731                    },
1732                    Timestamp::now(),
1733                )
1734                .await
1735                .unwrap(),
1736            SetExpiryResponse::Satisfied(requested)
1737        );
1738        assert_eq!(events.lock().unwrap().as_slice(), &["lt", "hv"]);
1739        assert_eq!(
1740            lt.inner.get(&target).expect_object().0.time_expires,
1741            Some(blob_expiry)
1742        );
1743        assert_eq!(
1744            hv.inner.get(&id).expect_tombstone().time_expires,
1745            Some(requested)
1746        );
1747    }
1748
1749    #[derive(Clone, Debug)]
1750    struct ReplaceInlineOnLookup {
1751        replacement: Metadata,
1752        replaced: Arc<StdMutex<bool>>,
1753    }
1754
1755    #[async_trait::async_trait]
1756    impl Hooks for ReplaceInlineOnLookup {
1757        async fn get_tiered_metadata(
1758            &self,
1759            inner: &InMemoryBackend,
1760            id: &ObjectId,
1761            access_time: Timestamp,
1762        ) -> Result<TieredMetadata> {
1763            let observed = inner.get_tiered_metadata(id, access_time).await?;
1764            let should_replace = {
1765                let mut replaced = self.replaced.lock().unwrap();
1766                !std::mem::replace(&mut *replaced, true)
1767            };
1768            if should_replace {
1769                inner
1770                    .put_object(
1771                        id,
1772                        &self.replacement,
1773                        stream::single("replacement"),
1774                        access_time,
1775                    )
1776                    .await?;
1777            }
1778            Ok(observed)
1779        }
1780    }
1781
1782    #[tokio::test]
1783    async fn inline_creation_target_uses_update_snapshot() {
1784        let access_time = Timestamp::from_unix_secs(1_700_000_000).unwrap();
1785        let old_created = access_time - Duration::from_hours(2);
1786        let new_created = access_time - Duration::from_mins(30);
1787        let old_expiry = access_time + Duration::from_mins(10);
1788        let resolved = new_created + Duration::from_hours(2);
1789        let replacement = Metadata {
1790            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1791            time_created: Some(new_created),
1792            time_expires: Some(old_expiry),
1793            ..Default::default()
1794        };
1795        let hv = TestBackend::new(ReplaceInlineOnLookup {
1796            replacement,
1797            replaced: Arc::new(StdMutex::new(false)),
1798        });
1799        let storage = TieredStorage::new(
1800            Box::new(hv.clone()),
1801            Box::new(InMemoryBackend::new("lt")),
1802            Box::new(NoopChangeLog),
1803        );
1804        let id = make_id("inline-update-snapshot");
1805        hv.inner
1806            .put_object(
1807                &id,
1808                &Metadata {
1809                    expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1810                    time_created: Some(old_created),
1811                    time_expires: Some(old_expiry),
1812                    ..Default::default()
1813                },
1814                stream::single("original"),
1815                access_time,
1816            )
1817            .await
1818            .unwrap();
1819
1820        assert_eq!(
1821            storage
1822                .set_expiry(
1823                    &id,
1824                    ExpiryTarget::FromCreation(Duration::from_hours(2)).into(),
1825                    access_time,
1826                )
1827                .await
1828                .unwrap(),
1829            SetExpiryResponse::Satisfied(resolved)
1830        );
1831        let (metadata, payload) = hv.inner.get(&id).expect_object();
1832        assert_eq!(
1833            metadata.expiration_policy,
1834            ExpirationPolicy::TimeToLive(Duration::from_hours(2))
1835        );
1836        assert_eq!(metadata.time_created, Some(new_created));
1837        assert_eq!(metadata.time_expires, Some(resolved));
1838        assert_eq!(payload, Bytes::from_static(b"replacement"));
1839    }
1840
1841    #[tokio::test]
1842    async fn expiry_hv_conflict() {
1843        let (storage, hv, lt, events) = tiered_with_expiry_hooks(true, false);
1844        let id = make_id("tiered-hv-failure");
1845        let target = new_long_term_revision(&id);
1846        let old_expiry = Timestamp::now() + Duration::from_mins(10);
1847        seed_redirect(&hv.inner, &lt.inner, &id, &target, old_expiry).await;
1848
1849        let requested = old_expiry + Duration::from_mins(50);
1850        assert_eq!(
1851            storage
1852                .set_expiry(
1853                    &id,
1854                    ExpiryTarget::FromCreation(Duration::from_hours(1)).into(),
1855                    Timestamp::now(),
1856                )
1857                .await
1858                .unwrap(),
1859            SetExpiryResponse::Rejected
1860        );
1861        assert_eq!(events.lock().unwrap().as_slice(), &["lt", "hv"]);
1862        assert_eq!(
1863            hv.inner.get(&id).expect_tombstone().time_expires,
1864            Some(old_expiry)
1865        );
1866        assert_eq!(
1867            lt.inner.get(&target).expect_object().0.time_expires,
1868            Some(requested)
1869        );
1870        assert_eq!(
1871            storage
1872                .get_metadata(&id, Timestamp::now())
1873                .await
1874                .unwrap()
1875                .unwrap()
1876                .time_expires,
1877            Some(old_expiry)
1878        );
1879        assert_eq!(
1880            lt.inner.get(&target).expect_object().0.expiration_policy,
1881            ExpirationPolicy::TimeToLive(Duration::from_hours(1))
1882        );
1883    }
1884
1885    #[tokio::test]
1886    async fn expiry_lt_conflict() {
1887        let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, true);
1888        let id = make_id("tiered-lt-conflict");
1889        let target = new_long_term_revision(&id);
1890        seed_redirect(
1891            &hv.inner,
1892            &lt.inner,
1893            &id,
1894            &target,
1895            Timestamp::now() + Duration::from_mins(10),
1896        )
1897        .await;
1898
1899        assert_eq!(
1900            storage
1901                .set_expiry(
1902                    &id,
1903                    ExpiryTarget::At(Timestamp::now() + Duration::from_hours(1)).into(),
1904                    Timestamp::now()
1905                )
1906                .await
1907                .unwrap(),
1908            SetExpiryResponse::Rejected
1909        );
1910        assert_eq!(events.lock().unwrap().as_slice(), &["lt"]);
1911    }
1912
1913    #[tokio::test]
1914    async fn expiry_not_found() {
1915        for missing_blob in [false, true] {
1916            let (storage, hv, lt, events) = tiered_with_expiry_hooks(false, false);
1917            let id = make_id("tiered-expiry-missing");
1918            let target = new_long_term_revision(&id);
1919            let access_time = Timestamp::now();
1920            let deadline = access_time + Duration::from_hours(1);
1921            if missing_blob {
1922                seed_redirect(&hv.inner, &lt.inner, &id, &target, deadline).await;
1923                lt.inner.delete_object(&target, access_time).await.unwrap();
1924            }
1925
1926            assert_eq!(
1927                storage
1928                    .set_expiry(&id, ExpiryTarget::At(deadline).into(), access_time)
1929                    .await
1930                    .unwrap(),
1931                SetExpiryResponse::NotFound
1932            );
1933            let expected: &[&str] = if missing_blob { &["lt"] } else { &[] };
1934            assert_eq!(events.lock().unwrap().as_slice(), expected);
1935        }
1936    }
1937
1938    // --- new_long_term_revision tests ---
1939
1940    #[test]
1941    fn revision_id_preserves_context() {
1942        let id = make_id("my-key");
1943        let revised = new_long_term_revision(&id);
1944        assert_eq!(revised.context, id.context);
1945        assert!(
1946            revised.key.starts_with("my-key/"),
1947            "revised key should have /<uuid> suffix, got: {}",
1948            revised.key
1949        );
1950    }
1951
1952    #[test]
1953    fn revision_id_roundtrips_storage_path() {
1954        let id = make_id("original");
1955        let revised = new_long_term_revision(&id);
1956        let path = revised.as_storage_path().to_string();
1957        let parsed = ObjectId::from_storage_path(&path)
1958            .unwrap_or_else(|| panic!("failed to parse '{path}'"));
1959        assert_eq!(parsed, revised);
1960    }
1961
1962    #[test]
1963    fn revision_id_is_unique() {
1964        let id = make_id("base-key");
1965        let a = new_long_term_revision(&id);
1966        let b = new_long_term_revision(&id);
1967        assert_ne!(a.key, b.key, "two calls should produce different keys");
1968    }
1969
1970    // --- Basic behavior ---
1971
1972    #[tokio::test]
1973    async fn get_nonexistent_returns_none() {
1974        let (storage, _hv, _lt, _) = make_tiered_storage();
1975        let id = make_id("does-not-exist");
1976
1977        assert!(
1978            storage
1979                .get_object(&id, Timestamp::now(), None)
1980                .await
1981                .unwrap()
1982                .is_none()
1983        );
1984        assert!(
1985            storage
1986                .get_metadata(&id, Timestamp::now())
1987                .await
1988                .unwrap()
1989                .is_none()
1990        );
1991    }
1992
1993    #[tokio::test]
1994    async fn delete_nonexistent_succeeds() {
1995        let (storage, _hv, _lt, _) = make_tiered_storage();
1996        let id = make_id("does-not-exist");
1997
1998        storage.delete_object(&id, Timestamp::now()).await.unwrap();
1999    }
2000
2001    // --- Put routing ---
2002
2003    #[tokio::test]
2004    async fn put_small_object_stores_inline() {
2005        let (storage, hv, lt, _) = make_tiered_storage();
2006        let id = make_id("small");
2007        let payload = b"small payload".to_vec();
2008
2009        storage
2010            .put_object(
2011                &id,
2012                &Metadata::default(),
2013                stream::single(payload.clone()),
2014                Timestamp::now(),
2015            )
2016            .await
2017            .unwrap();
2018
2019        assert!(hv.contains(&id), "expected in high-volume");
2020        assert!(!lt.contains(&id), "leaked to long-term");
2021
2022        let (_, _, s) = storage
2023            .get_object(&id, Timestamp::now(), None)
2024            .await
2025            .unwrap()
2026            .unwrap();
2027        let body = stream::read_to_vec(s).await.unwrap();
2028        assert_eq!(body, payload);
2029
2030        assert!(
2031            storage
2032                .get_metadata(&id, Timestamp::now())
2033                .await
2034                .unwrap()
2035                .is_some(),
2036            "get_metadata should return metadata for inline objects"
2037        );
2038    }
2039
2040    #[tokio::test]
2041    async fn put_large_object_creates_tombstone() {
2042        let (storage, hv, lt, _) = make_tiered_storage();
2043        let id = make_id("large");
2044        let payload = vec![0xCDu8; 2 * 1024 * 1024]; // 2 MiB, over threshold
2045        let metadata_in = Metadata {
2046            content_type: "image/png".into(),
2047            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2048            time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2049            origin: Some("10.0.0.1".into()),
2050            ..Metadata::default()
2051        };
2052
2053        storage
2054            .put_object(
2055                &id,
2056                &metadata_in,
2057                stream::single(payload.clone()),
2058                Timestamp::now(),
2059            )
2060            .await
2061            .unwrap();
2062
2063        // Tombstone in HV: correct deadline, target is a revision key.
2064        let tombstone = hv.get(&id).expect_tombstone();
2065        assert_eq!(tombstone.time_expires, metadata_in.time_expires);
2066        let lt_id = tombstone.target;
2067        assert!(
2068            lt_id.key().starts_with(id.key()),
2069            "tombstone target key should be a revision of the HV key, got: {}",
2070            lt_id.key()
2071        );
2072
2073        // LT object at revision key with correct metadata.
2074        let (lt_meta, _) = lt.get(&lt_id).expect_object();
2075        assert_eq!(lt_meta.content_type, "image/png");
2076        assert_eq!(lt_meta.expiration_policy, metadata_in.expiration_policy);
2077        assert_eq!(lt_meta.time_expires, tombstone.time_expires);
2078
2079        // get_object follows the tombstone and returns the correct payload.
2080        let (_, _, s) = storage
2081            .get_object(&id, Timestamp::now(), None)
2082            .await
2083            .unwrap()
2084            .unwrap();
2085        let body = stream::read_to_vec(s).await.unwrap();
2086        assert_eq!(body, payload);
2087
2088        // get_metadata follows the tombstone and returns the correct content_type.
2089        let metadata = storage
2090            .get_metadata(&id, Timestamp::now())
2091            .await
2092            .unwrap()
2093            .unwrap();
2094        assert_eq!(metadata.content_type, "image/png");
2095    }
2096
2097    // --- Put overwrites ---
2098
2099    #[tokio::test]
2100    async fn reinsert_small_over_large_swaps_to_inline() {
2101        let (storage, hv, lt, _) = make_tiered_storage();
2102        let id = make_id("reinsert-key");
2103
2104        // First: insert a large object → creates tombstone in hv, payload in lt at lt_id
2105        let large_payload = vec![0xABu8; 2 * 1024 * 1024];
2106        storage
2107            .put_object(
2108                &id,
2109                &Metadata::default(),
2110                stream::single(large_payload),
2111                Timestamp::now(),
2112            )
2113            .await
2114            .unwrap();
2115
2116        let lt_id = hv.get(&id).expect_tombstone().target;
2117
2118        // Re-insert a SMALL payload with the same key.
2119        // The CAS-swap puts the small object inline in HV and schedules background cleanup.
2120        let small_payload = vec![0xCDu8; 100]; // well under 1 MiB threshold
2121        storage
2122            .put_object(
2123                &id,
2124                &Metadata::default(),
2125                stream::single(small_payload),
2126                Timestamp::now(),
2127            )
2128            .await
2129            .unwrap();
2130
2131        // The small object is now inline in high-volume.
2132        hv.get(&id).expect_object();
2133
2134        // Drain background cleanup tasks before asserting LT state.
2135        storage.join().await;
2136
2137        // The old long-term blob was cleaned up.
2138        lt.get(&lt_id).expect_not_found();
2139    }
2140
2141    #[tokio::test]
2142    async fn overwrite_large_with_large_replaces_revision() {
2143        let (storage, hv, lt, _) = make_tiered_storage();
2144        let id = make_id("overwrite-large");
2145
2146        let payload1 = vec![0xAAu8; 2 * 1024 * 1024];
2147        storage
2148            .put_object(
2149                &id,
2150                &Metadata::default(),
2151                stream::single(payload1),
2152                Timestamp::now(),
2153            )
2154            .await
2155            .unwrap();
2156        let lt_id_1 = hv.get(&id).expect_tombstone().target;
2157
2158        let payload2 = vec![0xBBu8; 2 * 1024 * 1024];
2159        storage
2160            .put_object(
2161                &id,
2162                &Metadata::default(),
2163                stream::single(payload2.clone()),
2164                Timestamp::now(),
2165            )
2166            .await
2167            .unwrap();
2168        let lt_id_2 = hv.get(&id).expect_tombstone().target;
2169
2170        assert_ne!(
2171            lt_id_1, lt_id_2,
2172            "second write should create a new revision"
2173        );
2174
2175        // Drain background cleanup tasks before asserting LT state.
2176        storage.join().await;
2177
2178        lt.get(&lt_id_1).expect_not_found();
2179        lt.get(&lt_id_2).expect_object();
2180
2181        let (_, _, s) = storage
2182            .get_object(&id, Timestamp::now(), None)
2183            .await
2184            .unwrap()
2185            .unwrap();
2186        let body = stream::read_to_vec(s).await.unwrap();
2187        assert_eq!(body, payload2);
2188    }
2189
2190    // --- Delete ---
2191
2192    #[tokio::test]
2193    async fn delete_small_object() {
2194        let (storage, hv, _lt, _) = make_tiered_storage();
2195        let id = make_id("delete-small");
2196
2197        storage
2198            .put_object(
2199                &id,
2200                &Metadata::default(),
2201                stream::single("tiny"),
2202                Timestamp::now(),
2203            )
2204            .await
2205            .unwrap();
2206
2207        storage.delete_object(&id, Timestamp::now()).await.unwrap();
2208
2209        hv.get(&id).expect_not_found();
2210        assert!(
2211            storage
2212                .get_object(&id, Timestamp::now(), None)
2213                .await
2214                .unwrap()
2215                .is_none()
2216        );
2217    }
2218
2219    #[tokio::test]
2220    async fn delete_large_object_cleans_up_both_backends() {
2221        let (storage, hv, lt, _) = make_tiered_storage();
2222        let id = make_id("delete-both");
2223        let payload = vec![0u8; 2 * 1024 * 1024]; // 2 MiB
2224
2225        storage
2226            .put_object(
2227                &id,
2228                &Metadata::default(),
2229                stream::single(payload),
2230                Timestamp::now(),
2231            )
2232            .await
2233            .unwrap();
2234
2235        // Capture lt_id before deleting (it lives at the revision key, not at id).
2236        let lt_id = hv.get(&id).expect_tombstone().target;
2237
2238        storage.delete_object(&id, Timestamp::now()).await.unwrap();
2239
2240        // Drain background cleanup tasks before asserting LT state.
2241        storage.join().await;
2242
2243        assert!(!hv.contains(&id), "tombstone not cleaned up");
2244        assert!(!lt.contains(&lt_id), "long-term object not cleaned up");
2245    }
2246
2247    #[derive(Debug)]
2248    struct FailDelete;
2249
2250    #[async_trait::async_trait]
2251    impl Hooks for FailDelete {
2252        async fn delete_object(
2253            &self,
2254            _inner: &InMemoryBackend,
2255            _id: &ObjectId,
2256            _access_time: Timestamp,
2257        ) -> Result<DeleteResponse> {
2258            Err(Error::with_source(
2259                ErrorKind::BackendFailure,
2260                std::io::Error::new(
2261                    std::io::ErrorKind::ConnectionRefused,
2262                    "simulated long-term delete failure",
2263                ),
2264            ))
2265        }
2266    }
2267
2268    /// When the long-term GCS cleanup fails after the tombstone is deleted, the
2269    /// delete still succeeds (GCS cleanup is best-effort). An orphan blob may
2270    /// remain in LT storage, which is accepted.
2271    #[tokio::test]
2272    async fn delete_succeeds_when_gcs_cleanup_fails() {
2273        let hv = InMemoryBackend::new("hv");
2274        let lt = TestBackend::new(FailDelete);
2275        let log = NoopChangeLog;
2276        let storage = TieredStorage::new(Box::new(hv.clone()), Box::new(lt), Box::new(log));
2277
2278        let id = make_id("fail-delete");
2279        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB -> goes to long-term
2280        storage
2281            .put_object(
2282                &id,
2283                &Metadata::default(),
2284                stream::single(payload),
2285                Timestamp::now(),
2286            )
2287            .await
2288            .unwrap();
2289
2290        // Delete succeeds even though GCS cleanup fails (it is best-effort).
2291        let result = storage.delete_object(&id, Timestamp::now()).await;
2292        assert!(
2293            result.is_ok(),
2294            "delete should succeed despite GCS cleanup failure"
2295        );
2296
2297        // The tombstone in HV is gone (CAS-deleted first, before GCS cleanup).
2298        hv.get(&id).expect_not_found();
2299
2300        // The orphaned GCS blob remains but the object is unreachable through the service.
2301        assert!(
2302            storage
2303                .get_object(&id, Timestamp::now(), None)
2304                .await
2305                .unwrap()
2306                .is_none(),
2307            "object should be unreachable after tombstone is deleted"
2308        );
2309    }
2310
2311    // --- CAS conflicts ---
2312
2313    #[derive(Debug, Copy, Clone)]
2314    struct CasConflict;
2315
2316    #[async_trait::async_trait]
2317    impl Hooks for CasConflict {
2318        async fn compare_and_write(
2319            &self,
2320            _inner: &InMemoryBackend,
2321            _id: &ObjectId,
2322            _current: Option<&ObjectId>,
2323            _write: TieredWrite,
2324            _access_time: Timestamp,
2325        ) -> Result<bool> {
2326            Ok(false) // always conflict
2327        }
2328    }
2329
2330    /// After a large-object write loses the CAS race, the new LT blob must be
2331    /// cleaned up. The put still returns `Ok(())` — from the caller's view, a
2332    /// concurrent write won.
2333    #[tokio::test]
2334    async fn put_large_cas_conflict_cleans_up_new_blob() {
2335        let hv = TestBackend::new(CasConflict);
2336        let lt = InMemoryBackend::new("lt");
2337        let log = NoopChangeLog;
2338        let storage = TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(log));
2339
2340        let id = make_id("cas-conflict-large");
2341        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB -> long-term path
2342
2343        storage
2344            .put_object(
2345                &id,
2346                &Metadata::default(),
2347                stream::single(payload),
2348                Timestamp::now(),
2349            )
2350            .await
2351            .unwrap();
2352
2353        // Drain background cleanup tasks before asserting LT state.
2354        storage.join().await;
2355
2356        assert!(
2357            lt.is_empty(),
2358            "LT blob should be cleaned up after CAS conflict"
2359        );
2360    }
2361
2362    #[tokio::test]
2363    async fn upload_cas_conflict_cleans_up_new_blob() {
2364        let hv = TestBackend::new(CasConflict);
2365        let lt = InMemoryBackend::new("lt");
2366        let storage =
2367            TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(NoopChangeLog));
2368        let id = make_id("upload-cas-conflict");
2369        let payload = vec![0xAB; BACKEND_SIZE_THRESHOLD + 1];
2370        let session =
2371            resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
2372
2373        assert_eq!(
2374            storage
2375                .put_chunk(&session, 0, payload.len() as u64, stream::single(payload))
2376                .await
2377                .unwrap(),
2378            UploadProgress::Complete
2379        );
2380        storage.join().await;
2381        assert!(
2382            lt.is_empty(),
2383            "LT blob should be cleaned up after CAS conflict"
2384        );
2385    }
2386
2387    /// When swapping a tombstone for inline data, a CAS conflict means another
2388    /// writer won. The put still returns `Ok(())` — no LT blob was written, so
2389    /// there is nothing to clean up.
2390    #[tokio::test]
2391    async fn put_small_over_tombstone_cas_conflict_succeeds() {
2392        let inner = InMemoryBackend::new("hv");
2393        let id = make_id("cas-conflict-small");
2394
2395        // Pre-seed a tombstone directly in the inner backend so put_non_tombstone
2396        // returns it instead of writing inline.
2397        let tombstone = Tombstone {
2398            target: make_id("lt-object"),
2399            time_expires: None,
2400        };
2401        inner
2402            .compare_and_write(
2403                &id,
2404                None,
2405                TieredWrite::Tombstone(tombstone),
2406                Timestamp::now(),
2407            )
2408            .await
2409            .unwrap();
2410
2411        let lt = InMemoryBackend::new("lt");
2412        let hv = TestBackend::with_inner(inner, CasConflict);
2413        let log = NoopChangeLog;
2414        let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(log));
2415
2416        // Writing a small object over a tombstone should succeed even when CAS
2417        // conflicts — the other writer's write is accepted.
2418        storage
2419            .put_object(
2420                &id,
2421                &Metadata::default(),
2422                stream::single("tiny"),
2423                Timestamp::now(),
2424            )
2425            .await
2426            .unwrap();
2427    }
2428
2429    // --- Failure / inconsistency ---
2430
2431    /// Simulates compare_and_write failure. If `true`, it fails after commit.
2432    #[derive(Debug)]
2433    struct FailCas(bool);
2434
2435    #[async_trait::async_trait]
2436    impl Hooks for FailCas {
2437        async fn compare_and_write(
2438            &self,
2439            inner: &InMemoryBackend,
2440            id: &ObjectId,
2441            current: Option<&ObjectId>,
2442            write: TieredWrite,
2443            access_time: Timestamp,
2444        ) -> Result<bool> {
2445            if self.0 {
2446                // simulate a network error _after_ commit went through
2447                inner
2448                    .compare_and_write(id, current, write, access_time)
2449                    .await?;
2450            }
2451            Err(Error::with_source(
2452                ErrorKind::BackendFailure,
2453                std::io::Error::new(
2454                    std::io::ErrorKind::TimedOut,
2455                    "simulated compare_and_write failure",
2456                ),
2457            ))
2458        }
2459    }
2460
2461    /// If the tombstone write to the high-volume backend fails after the long-term
2462    /// write succeeds, the long-term object must be cleaned up so we never leave
2463    /// an unreachable orphan in long-term storage.
2464    #[tokio::test]
2465    async fn no_orphan_when_tombstone_write_fails() {
2466        let lt = InMemoryBackend::new("lt");
2467        let hv = TestBackend::new(FailCas(false));
2468        let log = NoopChangeLog;
2469        let storage = TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(log));
2470
2471        let id = make_id("orphan-test");
2472        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB -> long-term path
2473        let result = storage
2474            .put_object(
2475                &id,
2476                &Metadata::default(),
2477                stream::single(payload),
2478                Timestamp::now(),
2479            )
2480            .await;
2481
2482        assert!(result.is_err());
2483
2484        // Drain background cleanup tasks before asserting LT state.
2485        storage.join().await;
2486
2487        assert!(lt.is_empty(), "long-term object not cleaned up");
2488    }
2489
2490    #[tokio::test]
2491    async fn upload_cas_failure_cleans_up_new_blob() {
2492        let hv = TestBackend::new(FailCas(false));
2493        let lt = InMemoryBackend::new("lt");
2494        let storage =
2495            TieredStorage::new(Box::new(hv), Box::new(lt.clone()), Box::new(NoopChangeLog));
2496        let id = make_id("upload-cas-failure");
2497        let payload = vec![0xAB; BACKEND_SIZE_THRESHOLD + 1];
2498        let session =
2499            resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await;
2500
2501        assert_eq!(
2502            storage
2503                .put_chunk(&session, 0, payload.len() as u64, stream::single(payload))
2504                .await
2505                .unwrap_err()
2506                .kind(),
2507            ErrorKind::BackendFailure
2508        );
2509        storage.join().await;
2510        assert!(
2511            lt.is_empty(),
2512            "LT blob should be cleaned up after CAS failure"
2513        );
2514    }
2515
2516    /// If a tombstone exists in high-volume but the corresponding object is
2517    /// missing from long-term storage (e.g. due to a race condition or partial
2518    /// cleanup), reads should gracefully return None rather than error.
2519    #[tokio::test]
2520    async fn orphan_tombstone_returns_none() {
2521        let (storage, hv, lt, _) = make_tiered_storage();
2522        let id = make_id("orphan-tombstone");
2523        let payload = vec![0xCDu8; 2 * 1024 * 1024]; // 2 MiB
2524
2525        storage
2526            .put_object(
2527                &id,
2528                &Metadata::default(),
2529                stream::single(payload),
2530                Timestamp::now(),
2531            )
2532            .await
2533            .unwrap();
2534
2535        // The object is at the revision key in LT, not at id.
2536        let lt_id = hv.get(&id).expect_tombstone().target;
2537
2538        // Remove the long-term object, leaving an orphan tombstone in hv
2539        lt.remove(&lt_id);
2540
2541        assert!(
2542            storage
2543                .get_object(&id, Timestamp::now(), None)
2544                .await
2545                .unwrap()
2546                .is_none(),
2547            "orphan tombstone should resolve to None on get_object"
2548        );
2549        assert!(
2550            storage
2551                .get_metadata(&id, Timestamp::now())
2552                .await
2553                .unwrap()
2554                .is_none(),
2555            "orphan tombstone should resolve to None on get_metadata"
2556        );
2557    }
2558
2559    // --- Redirect target ---
2560
2561    /// A tombstone carrying an explicit `target` is followed correctly on reads and deletes,
2562    /// including when the target ObjectId differs from the HV ObjectId.
2563    #[tokio::test]
2564    async fn tombstone_target_is_used_for_reads_and_deletes() {
2565        let hv = InMemoryBackend::new("hv");
2566        let lt = InMemoryBackend::new("lt");
2567        let log = NoopChangeLog;
2568        let storage = TieredStorage::new(Box::new(hv.clone()), Box::new(lt.clone()), Box::new(log));
2569
2570        let hv_id = make_id("hv-key");
2571        let lt_id = make_id("lt-key");
2572        let payload = vec![0xABu8; 100];
2573
2574        // Write the object under the LT id and a tombstone pointing to it from HV.
2575        lt.put_object(
2576            &lt_id,
2577            &Metadata::default(),
2578            stream::single(payload.clone()),
2579            Timestamp::now(),
2580        )
2581        .await
2582        .unwrap();
2583        let tombstone = Tombstone {
2584            target: lt_id.clone(),
2585            time_expires: None,
2586        };
2587        hv.compare_and_write(
2588            &hv_id,
2589            None,
2590            TieredWrite::Tombstone(tombstone),
2591            Timestamp::now(),
2592        )
2593        .await
2594        .unwrap();
2595
2596        // get_object must follow the tombstone and find the object via the lt_id target.
2597        let (_, _, s) = storage
2598            .get_object(&hv_id, Timestamp::now(), None)
2599            .await
2600            .unwrap()
2601            .unwrap();
2602        let body = stream::read_to_vec(s).await.unwrap();
2603        assert_eq!(body, payload);
2604
2605        // delete_object must clean up both backends using the target.
2606        storage
2607            .delete_object(&hv_id, Timestamp::now())
2608            .await
2609            .unwrap();
2610        storage.join().await;
2611        assert!(!hv.contains(&hv_id), "tombstone should be removed");
2612        assert!(!lt.contains(&lt_id), "lt object should be removed");
2613    }
2614
2615    // --- Multi-chunk ---
2616
2617    #[tokio::test]
2618    async fn multi_chunk_large_object_chains_buffered_and_remaining() {
2619        let (storage, hv, lt, _) = make_tiered_storage();
2620        let id = make_id("multi-chunk");
2621
2622        // Deliver a 2 MiB payload across multiple chunks that individually
2623        // fit under the threshold but collectively exceed it.
2624        let chunk_size = 512 * 1024; // 512 KiB per chunk
2625        let chunk_count = 4; // 4 × 512 KiB = 2 MiB total
2626        let stream: ClientStream = futures_util::stream::iter(
2627            (0..chunk_count).map(move |i| Ok(Bytes::from(vec![i as u8; chunk_size]))),
2628        )
2629        .boxed();
2630
2631        storage
2632            .put_object(&id, &Metadata::default(), stream, Timestamp::now())
2633            .await
2634            .unwrap();
2635
2636        // Should have been routed to long-term (over 1 MiB) at the revision key.
2637        let lt_id = hv.get(&id).expect_tombstone().target;
2638        let (_, lt_bytes) = lt.get(&lt_id).expect_object();
2639        assert_eq!(lt_bytes.len(), chunk_size * chunk_count);
2640
2641        // Verify data integrity — each chunk's fill byte should appear in order.
2642        for i in 0..chunk_count {
2643            let offset = i * chunk_size;
2644            assert!(
2645                lt_bytes[offset..offset + chunk_size]
2646                    .iter()
2647                    .all(|&b| b == i as u8),
2648                "data mismatch in chunk {i}"
2649            );
2650        }
2651    }
2652
2653    // --- Written-phase cleanup ---
2654
2655    /// When a large-object overwrite commits in HV but its response is lost, the guard drops in
2656    /// `Written` phase. Cleanup must read HV to determine the CAS outcome, then delete whichever
2657    /// LT blob is no longer referenced — here the old one, since the new tombstone committed.
2658    #[tokio::test]
2659    async fn written_cleanup_after_lost_cas_response() {
2660        let (storage, hv, lt, log) = make_tiered_storage();
2661        let id = make_id("obj");
2662
2663        // First put: establishes tombstone
2664        let payload = vec![0xAAu8; 2 * 1024 * 1024];
2665        storage
2666            .put_object(
2667                &id,
2668                &Metadata::default(),
2669                stream::single(payload.clone()),
2670                Timestamp::now(),
2671            )
2672            .await
2673            .unwrap();
2674        let tombstone1 = hv.get(&id).expect_tombstone().target;
2675
2676        // Second put: Updates tombstone but fails immediately after committing
2677        let broken_storage = TieredStorage::new(
2678            Box::new(TestBackend::with_inner(hv.clone(), FailCas(true))),
2679            Box::new(lt.clone()),
2680            Box::new(log.clone()),
2681        );
2682        broken_storage
2683            .put_object(
2684                &id,
2685                &Metadata::default(),
2686                stream::single(payload.clone()),
2687                Timestamp::now(),
2688            )
2689            .await
2690            .unwrap_err(); // must fail
2691        let tombstone2 = hv.get(&id).expect_tombstone().target;
2692        assert_ne!(tombstone1, tombstone2);
2693
2694        // The first tombstone's target should be cleaned up, but the second should remain.
2695        broken_storage.join().await;
2696        lt.get(&tombstone1).expect_not_found();
2697        lt.get(&tombstone2).expect_object();
2698
2699        // Now delete the new object with the same tombstone failure
2700        broken_storage
2701            .delete_object(&id, Timestamp::now())
2702            .await
2703            .unwrap_err();
2704        hv.get(&id).expect_not_found();
2705        broken_storage.join().await;
2706        lt.get(&tombstone2).expect_not_found();
2707
2708        // Create a fresh large object
2709        let id = make_id("obj2");
2710        storage
2711            .put_object(
2712                &id,
2713                &Metadata::default(),
2714                stream::single(payload.clone()),
2715                Timestamp::now(),
2716            )
2717            .await
2718            .unwrap();
2719        let tombstone3 = hv.get(&id).expect_tombstone().target;
2720
2721        // Overwrite it with a small object and check again for cleanup
2722        broken_storage
2723            .put_object(
2724                &id,
2725                &Metadata::default(),
2726                stream::single(&b"small"[..]),
2727                Timestamp::now(),
2728            )
2729            .await
2730            .unwrap_err(); // must fail
2731        hv.get(&id).expect_object();
2732        broken_storage.join().await;
2733        lt.get(&tombstone3).expect_not_found();
2734    }
2735
2736    // --- ChangeGuard drop safety tests ---
2737
2738    /// Dropping a guard outside any tokio runtime must not panic.
2739    #[test]
2740    fn guard_dropped_outside_runtime_does_not_panic() {
2741        let manager = ChangeManager::new(
2742            Box::new(InMemoryBackend::new("hv")),
2743            Box::new(InMemoryBackend::new("lt")),
2744            Box::new(NoopChangeLog),
2745        );
2746
2747        let change = Change {
2748            id: make_id("object-key"),
2749            new: Some(make_id("cleanup-target")),
2750            old: None,
2751            cleanup_after: None,
2752        };
2753
2754        // Build the guard inside a temporary runtime, then let the runtime drop
2755        // so that no tokio context is active when the guard drops.
2756        let guard = {
2757            let rt = tokio::runtime::Runtime::new().unwrap();
2758            rt.block_on(manager.record(change)).unwrap()
2759        };
2760
2761        drop(guard); // Must not panic.
2762    }
2763
2764    /// `join` blocks until all in-flight guards have completed cleanup.
2765    ///
2766    /// Time is advanced manually so the test runs at virtual speed. The guard
2767    /// completes after 10 s; `join` must still be waiting at 9 s and done by 11 s.
2768    #[tokio::test(start_paused = true)]
2769    async fn join_waits_for_cleanup_to_complete() {
2770        let (storage, _hv, _lt, _) = make_tiered_storage();
2771        let change = Change {
2772            id: make_id("object-key"),
2773            new: None,
2774            old: None,
2775            cleanup_after: None,
2776        };
2777        let mut guard = storage.record_change(change).await.unwrap();
2778
2779        tokio::spawn(async move {
2780            tokio::time::sleep(Duration::from_secs(10)).await;
2781            guard.advance(ChangePhase::Completed);
2782            drop(guard);
2783        });
2784
2785        let join_future = tokio::spawn(async move { storage.join().await });
2786
2787        tokio::time::sleep(Duration::from_secs(9)).await;
2788        assert!(!join_future.is_finished(), "finished before guard dropped");
2789
2790        tokio::time::sleep(Duration::from_secs(2)).await;
2791        assert!(join_future.is_finished(), "finish after guard drops");
2792    }
2793
2794    // --- Changelog integration tests ---
2795
2796    /// LT backend hook that completes the write, then pauses until resumed.
2797    ///
2798    /// Lets tests cancel the owning future after the blob is committed but
2799    /// before the HV tombstone is set.
2800    #[derive(Clone, Debug)]
2801    struct PauseAfterPut {
2802        paused: Arc<tokio::sync::Notify>,
2803        resume: Arc<tokio::sync::Notify>,
2804    }
2805
2806    #[async_trait::async_trait]
2807    impl Hooks for PauseAfterPut {
2808        async fn put_object(
2809            &self,
2810            inner: &InMemoryBackend,
2811            id: &ObjectId,
2812            metadata: &Metadata,
2813            stream: ClientStream,
2814            access_time: Timestamp,
2815        ) -> Result<PutResponse> {
2816            inner.put_object(id, metadata, stream, access_time).await?;
2817            self.paused.notify_one();
2818            self.resume.notified().await;
2819            Ok(())
2820        }
2821    }
2822
2823    /// When a future is cancelled after the LT write but before the HV tombstone is set,
2824    /// the `ChangeGuard` cleans up the orphaned LT blob and removes the log entry.
2825    #[tokio::test]
2826    async fn dropped_future_triggers_cleanup_and_log_entry_removed() {
2827        let paused = Arc::new(tokio::sync::Notify::new());
2828        let hooks = PauseAfterPut {
2829            paused: Arc::clone(&paused),
2830            resume: Arc::new(tokio::sync::Notify::new()),
2831        };
2832
2833        let lt_inner = InMemoryBackend::new("lt");
2834        let log = InMemoryChangeLog::default();
2835        let storage = TieredStorage::new(
2836            Box::new(InMemoryBackend::new("hv")),
2837            Box::new(TestBackend::with_inner(lt_inner.clone(), hooks)),
2838            Box::new(log.clone()),
2839        );
2840
2841        let id = make_id("drop-test");
2842        let metadata = Metadata::default();
2843        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB → long-term path
2844
2845        // Drive the put until the LT write commits, then cancel before the HV tombstone is set.
2846        tokio::select! {
2847            result = storage.put_object(&id, &metadata, stream::single(payload), Timestamp::now()) => {
2848                panic!("expected put to pause before completing, got: {result:?}");
2849            }
2850            _ = paused.notified() => {
2851                // LT blob stored; cancelling drops the guard in Recorded phase.
2852            }
2853        }
2854
2855        // ChangeGuard dropped → background cleanup task spawned; wait for it.
2856        storage.join().await;
2857
2858        // The orphaned LT blob must have been deleted.
2859        assert!(lt_inner.is_empty(), "orphaned LT blob was not cleaned up");
2860
2861        // The log entry must be gone once cleanup completes.
2862        let entries = log.scan().await.unwrap();
2863        assert!(
2864            entries.is_empty(),
2865            "changelog entry not removed after cleanup"
2866        );
2867    }
2868
2869    // --- Multipart upload ---
2870
2871    #[test]
2872    fn multipart_upload_id_roundtrip() {
2873        let id = TieredUploadId {
2874            revision: "my-key/01924a6f-7e28-7b9a-9c1d-abcdef123456".into(),
2875            upload_id: UploadId::new("upstream-upload-id-abc".into()).unwrap(),
2876        };
2877        let encoded: UploadId = id.clone().try_into().unwrap();
2878        let decoded: TieredUploadId = (&encoded.clone()).try_into().unwrap();
2879        assert_eq!(decoded, id);
2880    }
2881
2882    #[test]
2883    fn malformed_multipart_upload_ids_are_invalid_upload_ids() {
2884        let invalid_base64 = UploadId::new("%%%".into()).unwrap();
2885        let malformed_json = UploadId::new("bm90IGpzb24".into()).unwrap();
2886
2887        for upload_id in [&invalid_base64, &malformed_json] {
2888            let error = TieredUploadId::try_from(upload_id).unwrap_err();
2889            assert_eq!(error.kind(), ErrorKind::InvalidUploadId);
2890        }
2891    }
2892
2893    #[tokio::test]
2894    async fn multipart_single_part_roundtrip() {
2895        let (storage, hv, lt, _) = make_tiered_storage();
2896        let id = make_id("mp-single");
2897        let metadata = Metadata {
2898            content_type: "application/octet-stream".into(),
2899            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2900            time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2901            ..Metadata::default()
2902        };
2903        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB
2904
2905        let upload_id = storage.initiate_multipart(&id, &metadata).await.unwrap();
2906
2907        let etag = storage
2908            .upload_part(
2909                &id,
2910                &upload_id,
2911                NonZeroU32::new(1).unwrap(),
2912                payload.len() as u64,
2913                None,
2914                stream::single(payload.clone()),
2915            )
2916            .await
2917            .unwrap();
2918
2919        let error = storage
2920            .complete_multipart(
2921                &id,
2922                &upload_id,
2923                vec![CompletedPart {
2924                    part_number: NonZeroU32::new(1).unwrap(),
2925                    etag,
2926                }],
2927                Timestamp::now(),
2928            )
2929            .await
2930            .unwrap();
2931        assert!(
2932            error.is_none(),
2933            "complete_multipart returned error: {error:?}"
2934        );
2935
2936        // get_object should follow the tombstone and return the payload.
2937        let (got_meta, _, s) = storage
2938            .get_object(&id, Timestamp::now(), None)
2939            .await
2940            .unwrap()
2941            .unwrap();
2942        let body = stream::read_to_vec(s).await.unwrap();
2943        assert_eq!(body, payload);
2944        assert_eq!(got_meta.content_type, "application/octet-stream");
2945
2946        // HV should have a tombstone, LT should have the object at the physical key.
2947        let tombstone = hv.get(&id).expect_tombstone();
2948        assert!(
2949            tombstone.target.key().starts_with(id.key()),
2950            "tombstone target should be a revision key"
2951        );
2952        assert_eq!(tombstone.time_expires, metadata.time_expires);
2953        assert_eq!(
2954            lt.get(&tombstone.target).expect_object().0.time_expires,
2955            tombstone.time_expires
2956        );
2957    }
2958
2959    #[tokio::test]
2960    async fn multipart_upload() {
2961        let (storage, _hv, _lt, _) = make_tiered_storage();
2962        let id = make_id("multipart");
2963
2964        let upload_id = storage
2965            .initiate_multipart(&id, &Metadata::default())
2966            .await
2967            .unwrap();
2968
2969        let part1 = vec![0xAAu8; 512 * 1024];
2970        let part2 = vec![0xBBu8; 512 * 1024];
2971        let part3 = vec![0xCCu8; 512 * 1024];
2972
2973        let etag3 = storage
2974            .upload_part(
2975                &id,
2976                &upload_id,
2977                NonZeroU32::new(3).unwrap(),
2978                part3.len() as u64,
2979                None,
2980                stream::single(part3.clone()),
2981            )
2982            .await
2983            .unwrap();
2984        let etag2 = storage
2985            .upload_part(
2986                &id,
2987                &upload_id,
2988                NonZeroU32::new(2).unwrap(),
2989                part2.len() as u64,
2990                None,
2991                stream::single(part2.clone()),
2992            )
2993            .await
2994            .unwrap();
2995        let etag1 = storage
2996            .upload_part(
2997                &id,
2998                &upload_id,
2999                NonZeroU32::new(1).unwrap(),
3000                part1.len() as u64,
3001                None,
3002                stream::single(part1.clone()),
3003            )
3004            .await
3005            .unwrap();
3006
3007        let error = storage
3008            .complete_multipart(
3009                &id,
3010                &upload_id,
3011                vec![
3012                    CompletedPart {
3013                        part_number: NonZeroU32::new(1).unwrap(),
3014                        etag: etag1,
3015                    },
3016                    CompletedPart {
3017                        part_number: NonZeroU32::new(2).unwrap(),
3018                        etag: etag2,
3019                    },
3020                    CompletedPart {
3021                        part_number: NonZeroU32::new(3).unwrap(),
3022                        etag: etag3,
3023                    },
3024                ],
3025                Timestamp::now(),
3026            )
3027            .await
3028            .unwrap();
3029        assert!(error.is_none());
3030
3031        let (_, _, s) = storage
3032            .get_object(&id, Timestamp::now(), None)
3033            .await
3034            .unwrap()
3035            .unwrap();
3036        let body = stream::read_to_vec(s).await.unwrap();
3037
3038        let mut expected = Vec::new();
3039        expected.extend_from_slice(&part1);
3040        expected.extend_from_slice(&part2);
3041        expected.extend_from_slice(&part3);
3042        assert_eq!(body, expected);
3043    }
3044
3045    #[tokio::test]
3046    async fn multipart_abort() {
3047        let (storage, hv, _lt, _) = make_tiered_storage();
3048        let id = make_id("mp-abort");
3049
3050        let upload_id = storage
3051            .initiate_multipart(&id, &Metadata::default())
3052            .await
3053            .unwrap();
3054
3055        // Upload a part then abort.
3056        let payload = vec![0xABu8; 100];
3057        storage
3058            .upload_part(
3059                &id,
3060                &upload_id,
3061                NonZeroU32::new(1).unwrap(),
3062                payload.len() as u64,
3063                None,
3064                stream::single(payload),
3065            )
3066            .await
3067            .unwrap();
3068
3069        storage.abort_multipart(&id, &upload_id).await.unwrap();
3070
3071        // No tombstone should have been written.
3072        hv.get(&id).expect_not_found();
3073
3074        // The object should not be reachable.
3075        assert!(
3076            storage
3077                .get_object(&id, Timestamp::now(), None)
3078                .await
3079                .unwrap()
3080                .is_none()
3081        );
3082    }
3083
3084    #[tokio::test]
3085    async fn multipart_list_parts() {
3086        let (storage, _hv, _lt, _) = make_tiered_storage();
3087        let id = make_id("mp-list");
3088
3089        let upload_id = storage
3090            .initiate_multipart(&id, &Metadata::default())
3091            .await
3092            .unwrap();
3093
3094        let part1 = vec![0xAAu8; 100];
3095        let part2 = vec![0xBBu8; 200];
3096        storage
3097            .upload_part(
3098                &id,
3099                &upload_id,
3100                NonZeroU32::new(1).unwrap(),
3101                part1.len() as u64,
3102                None,
3103                stream::single(part1),
3104            )
3105            .await
3106            .unwrap();
3107        storage
3108            .upload_part(
3109                &id,
3110                &upload_id,
3111                NonZeroU32::new(2).unwrap(),
3112                part2.len() as u64,
3113                None,
3114                stream::single(part2),
3115            )
3116            .await
3117            .unwrap();
3118
3119        let resp = storage
3120            .list_parts(&id, &upload_id, None, None)
3121            .await
3122            .unwrap();
3123        assert_eq!(resp.parts.len(), 2);
3124        assert_eq!(resp.parts[0].part_number.get(), 1);
3125        assert_eq!(resp.parts[0].size, 100);
3126        assert_eq!(resp.parts[1].part_number.get(), 2);
3127        assert_eq!(resp.parts[1].size, 200);
3128    }
3129
3130    #[tokio::test]
3131    async fn multipart_overwrites_existing_tombstone() {
3132        let (storage, hv, lt, _) = make_tiered_storage();
3133        let id = make_id("mp-overwrite");
3134
3135        // Put a large object via the normal path.
3136        let payload1 = vec![0xAAu8; 2 * 1024 * 1024];
3137        storage
3138            .put_object(
3139                &id,
3140                &Metadata::default(),
3141                stream::single(payload1),
3142                Timestamp::now(),
3143            )
3144            .await
3145            .unwrap();
3146        let old_lt_id = hv.get(&id).expect_tombstone().target;
3147
3148        // Overwrite via multipart.
3149        let upload_id = storage
3150            .initiate_multipart(&id, &Metadata::default())
3151            .await
3152            .unwrap();
3153
3154        let payload2 = vec![0xBBu8; 2 * 1024 * 1024];
3155        let etag = storage
3156            .upload_part(
3157                &id,
3158                &upload_id,
3159                NonZeroU32::new(1).unwrap(),
3160                payload2.len() as u64,
3161                None,
3162                stream::single(payload2.clone()),
3163            )
3164            .await
3165            .unwrap();
3166
3167        // The multipart upload is not finalized, so the tombstone still points to the old
3168        // revision.
3169        let lt_id = hv.get(&id).expect_tombstone().target;
3170        assert_eq!(old_lt_id, lt_id);
3171
3172        let error = storage
3173            .complete_multipart(
3174                &id,
3175                &upload_id,
3176                vec![CompletedPart {
3177                    part_number: NonZeroU32::new(1).unwrap(),
3178                    etag,
3179                }],
3180                Timestamp::now(),
3181            )
3182            .await
3183            .unwrap();
3184        assert!(error.is_none());
3185
3186        // Now the upload has been finalized, so the new tombstone points to the new revision.
3187        let new_lt_id = hv.get(&id).expect_tombstone().target;
3188        assert_ne!(old_lt_id, new_lt_id);
3189
3190        // Wait for background cleanup.
3191        storage.join().await;
3192
3193        // Old revision should be cleaned up.
3194        lt.get(&old_lt_id).expect_not_found();
3195        lt.get(&new_lt_id).expect_object();
3196
3197        // Assert the contents of the new revision.
3198        let (_, _, s) = storage
3199            .get_object(&id, Timestamp::now(), None)
3200            .await
3201            .unwrap()
3202            .unwrap();
3203        let body = stream::read_to_vec(s).await.unwrap();
3204        assert_eq!(body, payload2);
3205    }
3206
3207    // --- Multipart completion failure handling (consistency, retries, delayed cleanup) ---
3208
3209    /// Assembles the blob via `complete_multipart`, but returns an error to simulate a network
3210    /// failure on the response path. Also fails `get_metadata` so the tiered layer cannot
3211    /// recover by detecting the already-assembled blob.
3212    #[derive(Debug)]
3213    struct CompleteMultipartButReturnError;
3214
3215    #[async_trait::async_trait]
3216    impl Hooks for CompleteMultipartButReturnError {
3217        async fn complete_multipart(
3218            &self,
3219            inner: &InMemoryBackend,
3220            id: &ObjectId,
3221            upload_id: &UploadId,
3222            parts: Vec<CompletedPart>,
3223            access_time: Timestamp,
3224        ) -> Result<CompleteMultipartResponse> {
3225            inner
3226                .complete_multipart(id, upload_id, parts, access_time)
3227                .await
3228                .unwrap();
3229            Err(Error::with_source(
3230                ErrorKind::BackendFailure,
3231                std::io::Error::new(
3232                    std::io::ErrorKind::TimedOut,
3233                    "simulated network error on complete_multipart",
3234                ),
3235            ))
3236        }
3237
3238        async fn get_metadata(
3239            &self,
3240            _inner: &InMemoryBackend,
3241            _id: &ObjectId,
3242            _access_time: Timestamp,
3243        ) -> Result<MetadataResponse> {
3244            Err(Error::with_source(
3245                ErrorKind::BackendFailure,
3246                std::io::Error::new(
3247                    std::io::ErrorKind::TimedOut,
3248                    "simulated network error on get_metadata",
3249                ),
3250            ))
3251        }
3252    }
3253
3254    /// `complete_multipart` on the inner LT backend assembles the blob successfully, but both
3255    /// `complete_multipart` and `get_metadata` return errors, so the tiered layer cannot finalize.
3256    /// After `MULTIPART_COMPLETE_CLEANUP_DELAY`, `ChangeLog` recovery deletes the orphaned blob.
3257    #[tokio::test]
3258    async fn cleans_up_orphan_after_failed_multipart_complete() {
3259        let hv = InMemoryBackend::new("hv");
3260        let lt_inner = InMemoryBackend::new("lt");
3261        let log = InMemoryChangeLog::default();
3262        let storage = TieredStorage::new(
3263            Box::new(hv.clone()),
3264            Box::new(TestBackend::with_inner(
3265                lt_inner.clone(),
3266                CompleteMultipartButReturnError {},
3267            )),
3268            Box::new(log.clone()),
3269        );
3270
3271        let id = make_id("mp-orphan");
3272        let upload_id = storage
3273            .initiate_multipart(&id, &Metadata::default())
3274            .await
3275            .unwrap();
3276
3277        let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3278        let physical = ObjectId {
3279            context: id.context.clone(),
3280            key: tiered_id.revision,
3281        };
3282
3283        let payload = vec![0xABu8; 2 * 1024 * 1024];
3284        let etag = storage
3285            .upload_part(
3286                &id,
3287                &upload_id,
3288                NonZeroU32::new(1).unwrap(),
3289                payload.len() as u64,
3290                None,
3291                stream::single(payload),
3292            )
3293            .await
3294            .unwrap();
3295
3296        let result = storage
3297            .complete_multipart(
3298                &id,
3299                &upload_id,
3300                vec![CompletedPart {
3301                    part_number: NonZeroU32::new(1).unwrap(),
3302                    etag,
3303                }],
3304                Timestamp::now(),
3305            )
3306            .await;
3307        assert!(result.is_err());
3308        storage.join().await;
3309
3310        // The LT blob is orphaned, and no cleanup has been performed (yet), due to the guard being
3311        // dropped while in the `Assembling` state.
3312        lt_inner.get(&physical).expect_object();
3313        hv.get(&id).expect_not_found();
3314
3315        // Simulate the passage of time and run recovery.
3316        log.expire_all();
3317        let manager = ChangeManager::new(
3318            Box::new(hv.clone()),
3319            Box::new(lt_inner.clone()),
3320            Box::new(log.clone()),
3321        );
3322        manager.recover().await.unwrap();
3323
3324        // The orphaned LT blob has been cleaned up.
3325        lt_inner.get(&physical).expect_not_found();
3326        // The change has been removed from the log.
3327        let remaining = log.scan().await.unwrap();
3328        assert!(remaining.is_empty());
3329    }
3330
3331    #[derive(Debug)]
3332    struct FailOnFirstCompleteMultipartAttempt {
3333        attempt: Mutex<u32>,
3334    }
3335
3336    impl FailOnFirstCompleteMultipartAttempt {
3337        fn new() -> Self {
3338            Self {
3339                attempt: Mutex::new(0),
3340            }
3341        }
3342    }
3343
3344    #[async_trait::async_trait]
3345    impl Hooks for FailOnFirstCompleteMultipartAttempt {
3346        async fn complete_multipart(
3347            &self,
3348            inner: &InMemoryBackend,
3349            id: &ObjectId,
3350            upload_id: &UploadId,
3351            parts: Vec<CompletedPart>,
3352            access_time: Timestamp,
3353        ) -> Result<CompleteMultipartResponse> {
3354            let mut attempt = self.attempt.lock().await;
3355            *attempt += 1;
3356            if *attempt == 1 {
3357                Err(Error::with_source(
3358                    ErrorKind::BackendFailure,
3359                    std::io::Error::new(std::io::ErrorKind::TimedOut, "simulated network error"),
3360                ))
3361            } else {
3362                Ok(inner
3363                    .complete_multipart(id, upload_id, parts, access_time)
3364                    .await
3365                    .unwrap())
3366            }
3367        }
3368    }
3369
3370    /// The first attempt to `complete_multipart` fails, which generates a `Change` entry.
3371    /// The second call succeeds.
3372    /// When it's time to clean up, nothing is deleted, as the `complete_multipart` eventually went
3373    /// through before the cleanup deadline.
3374    #[tokio::test]
3375    async fn multipart_complete_succeeds_on_retry_and_leaves_state_consistent() {
3376        let hv = InMemoryBackend::new("hv");
3377        let lt_inner = InMemoryBackend::new("lt");
3378        let log = InMemoryChangeLog::default();
3379        let storage = TieredStorage::new(
3380            Box::new(hv.clone()),
3381            Box::new(TestBackend::with_inner(
3382                lt_inner.clone(),
3383                FailOnFirstCompleteMultipartAttempt::new(),
3384            )),
3385            Box::new(log.clone()),
3386        );
3387
3388        let id = make_id("mp-retry");
3389        let upload_id = storage
3390            .initiate_multipart(&id, &Metadata::default())
3391            .await
3392            .unwrap();
3393
3394        let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3395        let physical = ObjectId {
3396            context: id.context.clone(),
3397            key: tiered_id.revision,
3398        };
3399
3400        let payload = vec![0xABu8; 2 * 1024 * 1024];
3401        let etag = storage
3402            .upload_part(
3403                &id,
3404                &upload_id,
3405                NonZeroU32::new(1).unwrap(),
3406                payload.len() as u64,
3407                None,
3408                stream::single(payload.clone()),
3409            )
3410            .await
3411            .unwrap();
3412
3413        // The first `complete_multipart` call fails.
3414        let result = storage
3415            .complete_multipart(
3416                &id,
3417                &upload_id,
3418                vec![CompletedPart {
3419                    part_number: NonZeroU32::new(1).unwrap(),
3420                    etag: etag.clone(),
3421                }],
3422                Timestamp::now(),
3423            )
3424            .await;
3425        assert!(result.is_err());
3426        storage.join().await;
3427
3428        // The second `complete_multipart` call succeeds.
3429        let result = storage
3430            .complete_multipart(
3431                &id,
3432                &upload_id,
3433                vec![CompletedPart {
3434                    part_number: NonZeroU32::new(1).unwrap(),
3435                    etag,
3436                }],
3437                Timestamp::now(),
3438            )
3439            .await;
3440        assert!(result.is_ok());
3441        storage.join().await;
3442
3443        // The object is there.
3444        let (_, _, s) = storage
3445            .get_object(&id, Timestamp::now(), None)
3446            .await
3447            .unwrap()
3448            .unwrap();
3449        let body = stream::read_to_vec(s).await.unwrap();
3450        assert_eq!(body, payload);
3451
3452        // Simulate the passage of time and run recovery.
3453        log.expire_all();
3454        let manager = ChangeManager::new(
3455            Box::new(hv.clone()),
3456            Box::new(lt_inner.clone()),
3457            Box::new(log.clone()),
3458        );
3459        manager.recover().await.unwrap();
3460
3461        // The LT blob has not been cleaned up, as the write eventually went through.
3462        lt_inner.get(&physical).expect_object();
3463        // The tombstone still points to the blob.
3464        let tombstone = hv.get(&id).expect_tombstone();
3465        assert_eq!(tombstone.target, physical);
3466        // The change has been removed from the log.
3467        let remaining = log.scan().await.unwrap();
3468        assert!(remaining.is_empty());
3469
3470        // The object is still there after recovery.
3471        let (_, _, s) = storage
3472            .get_object(&id, Timestamp::now(), None)
3473            .await
3474            .unwrap()
3475            .unwrap();
3476        let body = stream::read_to_vec(s).await.unwrap();
3477        assert_eq!(body, payload);
3478    }
3479
3480    #[derive(Debug)]
3481    struct FailOnFirstGetMetadataAttempt {
3482        attempt: Mutex<u32>,
3483    }
3484
3485    impl FailOnFirstGetMetadataAttempt {
3486        fn new() -> Self {
3487            Self {
3488                attempt: Mutex::new(0),
3489            }
3490        }
3491    }
3492
3493    #[async_trait::async_trait]
3494    impl Hooks for FailOnFirstGetMetadataAttempt {
3495        async fn get_metadata(
3496            &self,
3497            inner: &InMemoryBackend,
3498            id: &ObjectId,
3499            access_time: Timestamp,
3500        ) -> Result<MetadataResponse> {
3501            let mut attempt = self.attempt.lock().await;
3502            *attempt += 1;
3503            if *attempt == 1 {
3504                Err(Error::with_source(
3505                    ErrorKind::BackendFailure,
3506                    std::io::Error::new(std::io::ErrorKind::TimedOut, "simulated network error"),
3507                ))
3508            } else {
3509                inner.get_metadata(id, access_time).await
3510            }
3511        }
3512    }
3513
3514    /// The first attempt to `complete_multipart` succeeds on the LT backend, but the subsequent
3515    /// `get_metadata` call fails with a network error, causing the overall `complete_multipart`
3516    /// to fail. The second call retries and succeeds (the LT object already exists from the first
3517    /// attempt).
3518    /// When it's time to clean up, nothing is deleted, as the `complete_multipart` eventually went
3519    /// through before the cleanup deadline.
3520    #[tokio::test]
3521    async fn multipart_complete_succeeds_on_retry_if_get_metadata_errs_and_leaves_state_consistent()
3522    {
3523        let hv = InMemoryBackend::new("hv");
3524        let lt_inner = InMemoryBackend::new("lt");
3525        let log = InMemoryChangeLog::default();
3526        let storage = TieredStorage::new(
3527            Box::new(hv.clone()),
3528            Box::new(TestBackend::with_inner(
3529                lt_inner.clone(),
3530                FailOnFirstGetMetadataAttempt::new(),
3531            )),
3532            Box::new(log.clone()),
3533        );
3534
3535        let id = make_id("mp-retry-meta");
3536        let upload_id = storage
3537            .initiate_multipart(&id, &Metadata::default())
3538            .await
3539            .unwrap();
3540
3541        let tiered_id: TieredUploadId = (&upload_id).try_into().unwrap();
3542        let physical = ObjectId {
3543            context: id.context.clone(),
3544            key: tiered_id.revision,
3545        };
3546
3547        let payload = vec![0xABu8; 2 * 1024 * 1024];
3548        let etag = storage
3549            .upload_part(
3550                &id,
3551                &upload_id,
3552                NonZeroU32::new(1).unwrap(),
3553                payload.len() as u64,
3554                None,
3555                stream::single(payload.clone()),
3556            )
3557            .await
3558            .unwrap();
3559
3560        // The first `complete_multipart` call fails (get_metadata network error), even though it
3561        // internally creates the LT blob.
3562        let result = storage
3563            .complete_multipart(
3564                &id,
3565                &upload_id,
3566                vec![CompletedPart {
3567                    part_number: NonZeroU32::new(1).unwrap(),
3568                    etag: etag.clone(),
3569                }],
3570                Timestamp::now(),
3571            )
3572            .await;
3573        assert!(result.is_err());
3574        storage.join().await;
3575
3576        // The second `complete_multipart` call succeeds.
3577        let result = storage
3578            .complete_multipart(
3579                &id,
3580                &upload_id,
3581                vec![CompletedPart {
3582                    part_number: NonZeroU32::new(1).unwrap(),
3583                    etag,
3584                }],
3585                Timestamp::now(),
3586            )
3587            .await;
3588        assert!(result.is_ok());
3589        storage.join().await;
3590
3591        // The object is there.
3592        let (_, _, s) = storage
3593            .get_object(&id, Timestamp::now(), None)
3594            .await
3595            .unwrap()
3596            .unwrap();
3597        let body = stream::read_to_vec(s).await.unwrap();
3598        assert_eq!(body, payload);
3599
3600        // Simulate the passage of time and run recovery.
3601        log.expire_all();
3602        let manager = ChangeManager::new(
3603            Box::new(hv.clone()),
3604            Box::new(lt_inner.clone()),
3605            Box::new(log.clone()),
3606        );
3607        manager.recover().await.unwrap();
3608
3609        // The LT blob has not been cleaned up, as the write eventually went through.
3610        lt_inner.get(&physical).expect_object();
3611        // The tombstone still points to the blob.
3612        let tombstone = hv.get(&id).expect_tombstone();
3613        assert_eq!(tombstone.target, physical);
3614        // The change has been removed from the log.
3615        let remaining = log.scan().await.unwrap();
3616        assert!(remaining.is_empty());
3617
3618        // The object is there after recovery.
3619        let (_, _, s) = storage
3620            .get_object(&id, Timestamp::now(), None)
3621            .await
3622            .unwrap()
3623            .unwrap();
3624        let body = stream::read_to_vec(s).await.unwrap();
3625        assert_eq!(body, payload);
3626    }
3627}