Skip to main content

objectstore_service/
service.rs

1//! Core storage service and configuration.
2//!
3//! [`StorageService`] is the main entry point for storing and retrieving
4//! objects. Each operation runs in a separate tokio task for panic isolation.
5//!
6//! Callers supply an access timestamp, normally the HTTP request start time.
7//! It stays fixed across admission, backend calls, and TTI calculations. Initial
8//! creation and expiry are already resolved in the supplied metadata. Background
9//! renewals retain the read-derived deadline but check liveness at worker start.
10//!
11//! See the [crate-level documentation](crate) for full architecture details.
12
13use std::future::Future;
14use std::num::NonZeroU64;
15use std::sync::Arc;
16
17use objectstore_types::metadata::Metadata;
18use objectstore_types::range::{ByteRange, ContentRange};
19use objectstore_types::resumable::{SessionToken as EncryptedSessionToken, UploadProgress};
20use objectstore_types::time::Timestamp;
21
22use crate::backend::common::{Backend, ExpiryUpdate, SetExpiryResponse};
23use crate::backend::counting::CountingBackend;
24use crate::background::RenewalScheduler;
25use crate::concurrency::ConcurrencyLimiter;
26use crate::encryption::Cipher;
27use crate::error::{ErrorKind, Result, ResultExt as _};
28use crate::id::{ObjectContext, ObjectId};
29use crate::multipart::{
30    AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
31    ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
32};
33use crate::resumable::Session;
34use crate::stream::{ClientStream, PayloadStream};
35use crate::streaming::StreamExecutor;
36
37/// Service response for [`StorageService::get_object`].
38pub type GetResponse = Option<(Metadata, Option<ContentRange>, PayloadStream)>;
39/// Service response for [`StorageService::get_metadata`].
40pub type MetadataResponse = Option<Metadata>;
41/// Service response for [`StorageService::insert_object`].
42pub type InsertResponse = ObjectId;
43/// Service response for [`StorageService::delete_object`].
44pub type DeleteResponse = ();
45
46/// A newly opened resumable upload session and its upload granularity.
47#[derive(Debug)]
48pub struct CreateUploadSessionResponse {
49    /// The encrypted token used to continue the upload.
50    pub session: EncryptedSessionToken,
51    /// The upload granularity in bytes, or zero when there is none.
52    pub granularity: u64,
53}
54
55/// Default concurrency limit for [`StorageService`].
56///
57/// This value is used when no explicit limiter is set via
58/// [`StorageService::with_concurrency`].
59pub const DEFAULT_CONCURRENCY_LIMIT: u32 = 500;
60
61/// Default number of TTI renewals that may wait for background processing.
62pub const DEFAULT_BACKGROUND_QUEUE_LIMIT: usize = 1_000;
63
64/// Asynchronous storage service wrapping a single [`Backend`].
65///
66/// `StorageService` is the main entry point for storing and retrieving objects.
67/// It delegates all storage operations to the backend supplied at construction,
68/// adding task spawning, panic isolation, and a concurrency limit on top.
69///
70/// The typical backend is [`TieredStorage`](crate::backend::tiered::TieredStorage),
71/// which provides size-based routing to high-volume and long-term backends along
72/// with redirect tombstone management. Any type implementing [`Backend`] can be used.
73///
74/// # Lifecycle
75///
76/// After construction, call [`start`](StorageService::start) to start the
77/// service's background processes.
78///
79/// # Run-to-Completion and Panic Isolation
80///
81/// Each operation runs to completion even if the caller is cancelled (e.g., on
82/// client disconnect). This ensures that multi-step operations in the backend
83/// are never left partially applied. Post-commit cleanup (e.g. deleting
84/// unreferenced long-term blobs) runs in background tasks so callers are not
85/// blocked. Call [`join`](StorageService::join) during shutdown to wait for
86/// outstanding cleanup. Operations are also isolated from panics in backend
87/// code — a failure in one operation does not bring down other in-flight work.
88///
89/// # Concurrency Limit
90///
91/// A [`ConcurrencyLimiter`] caps the number of in-flight backend operations.
92/// Pass a custom limiter via
93/// [`with_concurrency`](StorageService::with_concurrency); without one the
94/// default is [`DEFAULT_CONCURRENCY_LIMIT`] permits with no queue.
95#[derive(Clone, Debug)]
96pub struct StorageService {
97    inner: Arc<dyn Backend>,
98    concurrency: ConcurrencyLimiter,
99    renewals: RenewalScheduler,
100    cipher: Arc<Cipher>,
101}
102
103impl StorageService {
104    /// Creates a new `StorageService` wrapping the given backend.
105    ///
106    /// The backend is wrapped in a [`CountingBackend`] which increments a COGS usage counter for
107    /// each operation run. Single-object operations served directly by `StorageService` are covered
108    /// as we batched operations served by [`StreamExecutor`]. See
109    /// [`backend::counting`](crate::backend::counting) for details.
110    pub fn new(backend: Box<dyn Backend>, cipher: Cipher) -> Self {
111        let inner: Arc<dyn Backend> = Arc::new(CountingBackend::new(backend));
112        let concurrency = ConcurrencyLimiter::new(DEFAULT_CONCURRENCY_LIMIT);
113        Self {
114            inner: Arc::clone(&inner),
115            concurrency: concurrency.clone(),
116            renewals: RenewalScheduler::new(inner, concurrency, DEFAULT_BACKGROUND_QUEUE_LIMIT),
117            cipher: Arc::new(cipher),
118        }
119    }
120
121    /// Replaces the default concurrency limiter.
122    ///
123    /// Must be called before [`start`](Self::start). Without this, the
124    /// service uses a limiter with [`DEFAULT_CONCURRENCY_LIMIT`] permits
125    /// and no queue.
126    pub fn with_concurrency(mut self, limiter: ConcurrencyLimiter) -> Self {
127        self.renewals.set_concurrency(limiter.clone());
128        self.concurrency = limiter;
129        self
130    }
131
132    /// Replaces the default background expiry-renewal queue capacity.
133    pub fn with_background_queue_limit(mut self, limit: usize) -> Self {
134        self.renewals.set_capacity(limit);
135        self
136    }
137
138    /// Returns a reference to the concurrency limiter.
139    pub fn concurrency_limiter(&self) -> &ConcurrencyLimiter {
140        &self.concurrency
141    }
142
143    /// Returns the number of backend tasks currently running.
144    pub fn tasks_running(&self) -> u32 {
145        self.concurrency.used_permits()
146    }
147
148    /// Returns the configured limit for concurrent backend tasks.
149    pub fn tasks_limit(&self) -> u32 {
150        self.concurrency.total_permits()
151    }
152
153    /// Prepares to stream multiple operations concurrently against this service.
154    ///
155    /// Each operation acquires a bulk permit individually via
156    /// [`ConcurrencyLimiter::acquire_bulk`], which caps bulk traffic at a
157    /// configurable percentage of execution slots while allowing operations
158    /// to queue for permits instead of requiring upfront reservation.
159    pub fn stream(&self) -> StreamExecutor {
160        StreamExecutor::new(
161            Arc::clone(&self.inner),
162            self.concurrency.clone(),
163            self.renewals.clone(),
164        )
165    }
166
167    /// Starts background processes for the storage service.
168    ///
169    /// At startup, this tracks the following gauges:
170    ///
171    ///  - `service.concurrency.limit`: concurrent task execution slots
172    ///  - `service.concurrency.queue_limit`: queue size for waiting tasks
173    ///  - `service.concurrency.bulk_limit`: concurrent task execution slots for bulk operations
174    ///
175    /// Also spawns tasks that emit runtime gauges once per second:
176    ///  - `service.concurrency.in_use`: currently running tasks
177    ///  - `service.concurrency.queued`: currently queued tasks
178    ///  - `service.concurrency.bulk_in_use`: currently running bulk tasks
179    ///  - `service.expiry_renewal.queued`: expiry renewals waiting for the background worker
180    pub fn start(&mut self) {
181        self.renewals.start();
182
183        let concurrency = self.concurrency.clone();
184        objectstore_metrics::gauge!("service.concurrency.limit" = concurrency.total_permits());
185        objectstore_metrics::gauge!("service.concurrency.queue_limit" = concurrency.total_queue());
186        objectstore_metrics::gauge!("service.concurrency.bulk_limit" = concurrency.total_bulk());
187
188        tokio::spawn(async move {
189            concurrency
190                .run_emitter(|stats| async move {
191                    objectstore_metrics::gauge!("service.concurrency.in_use" = stats.in_use);
192                    objectstore_metrics::gauge!("service.concurrency.queued" = stats.queued);
193                    objectstore_metrics::gauge!(
194                        "service.concurrency.bulk_in_use" = stats.bulk_in_use
195                    );
196                })
197                .await;
198        });
199
200        let renewals = self.renewals.clone();
201        tokio::spawn(async move {
202            renewals
203                .run_emitter(|queued| async move {
204                    objectstore_metrics::gauge!("service.expiry_renewal.queued" = queued);
205                })
206                .await;
207        });
208    }
209
210    /// Spawns a future in a separate task and awaits its result.
211    ///
212    /// # Observability
213    ///
214    /// This tracks two metrics:
215    ///
216    /// - `service.task.start` (counter) after acquiring a permit
217    /// - `service.task.duration` (distribution) when the task completes
218    ///
219    /// Both are tagged with the given `operation` name and an `outcome`
220    /// of `"success"` or `"error"`.
221    ///
222    /// # Errors
223    ///
224    /// - `AtCapacity` if the concurrency limit is reached
225    /// - `Panic` if the spawned task panics (the panic message is captured for diagnostics)
226    /// - `Dropped` if the task is dropped before sending its result.
227    async fn spawn<T, F>(&self, operation: &'static str, f: F) -> Result<T>
228    where
229        T: Send + 'static,
230        F: Future<Output = Result<T>> + Send + 'static,
231    {
232        let timer = objectstore_metrics::timer!("service.concurrency.wait");
233        let permit = self.concurrency.acquire().await.inspect_err(|_| {
234            objectstore_metrics::count!("service.concurrency.rejected", class = "normal");
235            objectstore_log::warn!("Request rejected: service at capacity");
236        })?;
237
238        timer.record();
239        crate::concurrency::run_metered(operation, permit, f).await
240    }
241
242    /// Creates or overwrites an object.
243    ///
244    /// The object is identified by the components of an [`ObjectId`]. The
245    /// `context` is required, while the `key` can be assigned automatically if
246    /// set to `None`.
247    ///
248    /// # Run-to-completion
249    ///
250    /// Once called, the operation runs to completion even if the returned future
251    /// is dropped (e.g., on client disconnect). This guarantees that partially
252    /// written objects in the backend are never left in an inconsistent state.
253    pub async fn insert_object(
254        &self,
255        context: ObjectContext,
256        key: Option<String>,
257        metadata: Metadata,
258        stream: ClientStream,
259        access_time: Timestamp,
260    ) -> Result<InsertResponse> {
261        metadata.validate().kind(ErrorKind::InvalidMetadata)?;
262        let id = ObjectId::optional(context, key);
263        let inner = Arc::clone(&self.inner);
264        self.spawn("insert", async move {
265            inner
266                .put_object(&id, &metadata, stream, access_time)
267                .await?;
268            Ok(id)
269        })
270        .await
271    }
272
273    /// Retrieves only the metadata for an object, without the payload.
274    pub async fn get_metadata(
275        &self,
276        id: ObjectId,
277        access_time: Timestamp,
278    ) -> Result<MetadataResponse> {
279        let inner = Arc::clone(&self.inner);
280        let renewals = self.renewals.clone();
281        self.spawn("get_metadata", async move {
282            let response = inner.get_metadata(&id, access_time).await?;
283            if let Some(ref metadata) = response
284                && let Some(expire_at) = metadata.check_tti_bump(access_time)
285            {
286                renewals.schedule(id, expire_at);
287            }
288            Ok(response)
289        })
290        .await
291    }
292
293    /// Streams (part of) the contents of an object.
294    pub async fn get_object(
295        &self,
296        id: ObjectId,
297        access_time: Timestamp,
298        range: Option<ByteRange>,
299    ) -> Result<GetResponse> {
300        let inner = Arc::clone(&self.inner);
301        let renewals = self.renewals.clone();
302        self.spawn("get", async move {
303            let response = inner.get_object(&id, access_time, range).await?;
304            if let Some((ref metadata, _, _)) = response
305                && let Some(expire_at) = metadata.check_tti_bump(access_time)
306            {
307                renewals.schedule(id, expire_at);
308            }
309            Ok(response)
310        })
311        .await
312    }
313
314    /// Extends an existing TTL or TTI object's deadline.
315    ///
316    /// Actual extensions adjust TTL duration to approximately match the lifetime since
317    /// creation. TTI duration, payload, and other metadata remain unchanged.
318    ///
319    /// Returns whether the request was satisfied, the object was absent or expired,
320    /// or the update was rejected. See [`SetExpiryResponse`] for details.
321    /// Remaining-lifetime limits in [`ExpiryUpdate`] are enforced against `access_time`;
322    /// exceeding a limit returns [`ErrorKind::InvalidMetadata`].
323    pub async fn set_expiry(
324        &self,
325        id: ObjectId,
326        target: ExpiryUpdate,
327        access_time: Timestamp,
328    ) -> Result<SetExpiryResponse> {
329        let inner = Arc::clone(&self.inner);
330        self.spawn("set_expiry", async move {
331            inner.set_expiry(&id, target, access_time).await
332        })
333        .await
334    }
335
336    /// Deletes an object, if it exists.
337    ///
338    /// # Run-to-completion
339    ///
340    /// Once called, the operation runs to completion even if the returned future
341    /// is dropped. This guarantees that multi-step delete sequences in the backend
342    /// are never left partially applied.
343    pub async fn delete_object(
344        &self,
345        id: ObjectId,
346        access_time: Timestamp,
347    ) -> Result<DeleteResponse> {
348        let inner = Arc::clone(&self.inner);
349        self.spawn("delete", async move {
350            inner.delete_object(&id, access_time).await
351        })
352        .await
353    }
354
355    /// Waits for all outstanding background operations to complete.
356    ///
357    /// Blocks until any pending background cleanup tasks finish, up to the
358    /// backend's configured timeout. Should be called during graceful shutdown
359    /// after the HTTP server has stopped accepting new requests.
360    pub async fn join(&self) {
361        self.renewals.join().await;
362        self.inner.join().await;
363    }
364
365    // --- Multipart upload operations ---
366
367    /// Initiates a new multipart upload.
368    pub async fn initiate_multipart(
369        &self,
370        id: ObjectId,
371        metadata: Metadata,
372    ) -> Result<InitiateMultipartResponse> {
373        metadata.validate().kind(ErrorKind::InvalidMetadata)?;
374        self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported
375        let inner = self.inner.clone();
376        self.spawn("initiate_multipart", async move {
377            inner
378                .as_multipart_upload_backend()?
379                .initiate_multipart(&id, &metadata)
380                .await
381        })
382        .await
383    }
384
385    /// Uploads a single part.
386    ///
387    /// Note that this requires a `content_length`.
388    /// This grants us the broadest and most seamless compatibility when it comes to backends.
389    /// For example, MinIO rejects `UploadPart` requests without a `Content-Length` on plain PUT
390    /// requests.
391    /// This can be worked around by using AWS SigV4 chunked streaming requests, which we could use
392    /// if one day we'll have a usecase where the client doesn't know the part length upfront.
393    pub async fn upload_part(
394        &self,
395        id: ObjectId,
396        upload_id: UploadId,
397        part_number: PartNumber,
398        content_length: u64,
399        content_md5: Option<String>,
400        body: ClientStream,
401    ) -> Result<UploadPartResponse> {
402        self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported
403        let inner = self.inner.clone();
404        self.spawn("upload_part", async move {
405            inner
406                .as_multipart_upload_backend()?
407                .upload_part(
408                    &id,
409                    &upload_id,
410                    part_number,
411                    content_length,
412                    content_md5.as_deref(),
413                    body,
414                )
415                .await
416        })
417        .await
418    }
419
420    /// Lists the parts uploaded so far.
421    pub async fn list_parts(
422        &self,
423        id: ObjectId,
424        upload_id: UploadId,
425        max_parts: Option<u32>,
426        part_number_marker: Option<PartNumber>,
427    ) -> Result<ListPartsResponse> {
428        self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported
429        let inner = self.inner.clone();
430        self.spawn("list_parts", async move {
431            inner
432                .as_multipart_upload_backend()?
433                .list_parts(&id, &upload_id, max_parts, part_number_marker)
434                .await
435        })
436        .await
437    }
438
439    /// Aborts a multipart upload.
440    pub async fn abort_multipart(
441        &self,
442        id: ObjectId,
443        upload_id: UploadId,
444    ) -> Result<AbortMultipartResponse> {
445        self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported
446        let inner = self.inner.clone();
447        self.spawn("abort_multipart", async move {
448            inner
449                .as_multipart_upload_backend()?
450                .abort_multipart(&id, &upload_id)
451                .await
452        })
453        .await
454    }
455
456    /// Finalizes a multipart upload.
457    pub async fn complete_multipart(
458        &self,
459        id: ObjectId,
460        upload_id: UploadId,
461        parts: Vec<CompletedPart>,
462        access_time: Timestamp,
463    ) -> Result<CompleteMultipartResponse> {
464        self.inner.as_multipart_upload_backend()?; // Fail before clone/spawn if unsupported
465        let inner = self.inner.clone();
466        self.spawn("complete_multipart", async move {
467            inner
468                .as_multipart_upload_backend()?
469                .complete_multipart(&id, &upload_id, parts, access_time)
470                .await
471        })
472        .await
473    }
474
475    // --- Resumable upload operations ---
476
477    /// Opens a resumable upload session for an object of `upload_length` bytes.
478    ///
479    /// Returns `Ok(None)` for zero-length objects or when the backend declines resumable uploads
480    /// for this object, in which case the caller should fall back to [`Self::insert_object`].
481    pub async fn create_upload_session(
482        &self,
483        id: ObjectId,
484        metadata: Metadata,
485        upload_length: u64,
486    ) -> Result<Option<CreateUploadSessionResponse>> {
487        let Some(upload_length) = NonZeroU64::new(upload_length) else {
488            return Ok(None);
489        };
490        metadata.validate().kind(ErrorKind::InvalidMetadata)?;
491        let inner = Arc::clone(&self.inner);
492        let cipher = Arc::clone(&self.cipher);
493        self.spawn("create_upload_session", async move {
494            let session = inner
495                .create_upload_session(&id, &metadata, upload_length)
496                .await?;
497            session
498                .map(|backend_token| {
499                    let session = cipher
500                        .encrypt(&Session {
501                            object_id: id,
502                            upload_length,
503                            backend_token,
504                        })
505                        .map(EncryptedSessionToken::new)?;
506                    Ok(CreateUploadSessionResponse {
507                        session,
508                        granularity: inner.upload_granularity(),
509                    })
510                })
511                .transpose()
512        })
513        .await
514    }
515
516    fn session_for(&self, expected_id: &ObjectId, token: EncryptedSessionToken) -> Result<Session> {
517        let session: Session = self
518            .cipher
519            .decrypt(token.as_bytes())
520            .map_err(|_| ErrorKind::UnknownUploadSession)?;
521        if session.object_id != *expected_id {
522            return Err(ErrorKind::UnknownUploadSession.into());
523        }
524        Ok(session)
525    }
526
527    /// Writes a chunk of `content_length` bytes at `offset` into an open session.
528    ///
529    /// Completes the upload once the chunk carrying the last byte is persisted.
530    ///
531    /// # Run-to-completion
532    ///
533    /// Once called, the operation runs to completion even if the returned future is dropped.
534    /// This matters most for the final chunk, which completes the upload.
535    pub async fn put_chunk(
536        &self,
537        id: ObjectId,
538        token: EncryptedSessionToken,
539        offset: u64,
540        content_length: u64,
541        body: ClientStream,
542    ) -> Result<UploadProgress> {
543        let session = self.session_for(&id, token)?;
544        let inner = Arc::clone(&self.inner);
545        self.spawn("put_chunk", async move {
546            inner
547                .put_chunk(&session, offset, content_length, body)
548                .await
549        })
550        .await
551    }
552
553    /// Reports how far a session has progressed.
554    ///
555    /// This can observe completion after the final chunk's response was lost. A composed backend
556    /// may also finish pending publication work, so this requires write permission at the API
557    /// layer.
558    pub async fn upload_offset(
559        &self,
560        id: ObjectId,
561        token: EncryptedSessionToken,
562    ) -> Result<UploadProgress> {
563        let session = self.session_for(&id, token)?;
564        let inner = Arc::clone(&self.inner);
565        self.spawn("upload_offset", async move {
566            inner.upload_offset(&session).await
567        })
568        .await
569    }
570
571    /// Cancels an upload session, discarding whatever was uploaded.
572    pub async fn cancel_upload(&self, id: ObjectId, token: EncryptedSessionToken) -> Result<()> {
573        let session = self.session_for(&id, token)?;
574        let inner = Arc::clone(&self.inner);
575        self.spawn("cancel_upload", async move {
576            inner.cancel_upload(&session).await
577        })
578        .await
579    }
580}
581
582#[cfg(test)]
583mod tests {
584    use std::error::Error as _;
585    use std::sync::atomic::{AtomicUsize, Ordering};
586    use std::sync::{Arc, Mutex};
587    use std::time::Duration;
588
589    use bytes::BytesMut;
590    use futures_util::TryStreamExt;
591    use objectstore_types::metadata::{ExpirationPolicy, Metadata};
592    use objectstore_types::range::ByteRange;
593    use objectstore_types::scope::{Scope, Scopes};
594
595    use super::*;
596    use crate::backend::bigtable::{BigTableBackend, BigTableConfig};
597    use crate::backend::changelog::NoopChangeLog;
598    use crate::backend::common::{ExpiryTarget, HighVolumeBackend, PutResponse, TieredWrite};
599    use crate::backend::gcs::{GcsBackend, GcsConfig};
600    use crate::backend::in_memory::InMemoryBackend;
601    use crate::backend::testing::{Hooks, TestBackend};
602    use crate::backend::tiered::TieredStorage;
603    use crate::change_stream::ChangeStreamFactory;
604    use crate::resumable::BackendToken;
605    use crate::stream::{self, ClientStream};
606
607    #[derive(Clone, Debug, Default)]
608    struct ResumableTokenHooks {
609        seen_tokens: Arc<Mutex<Vec<String>>>,
610    }
611
612    #[async_trait::async_trait]
613    impl Hooks for ResumableTokenHooks {
614        fn upload_granularity(&self, _inner: &InMemoryBackend) -> u64 {
615            256 * 1024
616        }
617
618        async fn create_upload_session(
619            &self,
620            _inner: &InMemoryBackend,
621            _id: &ObjectId,
622            _metadata: &Metadata,
623            _upload_length: NonZeroU64,
624        ) -> Result<Option<BackendToken>> {
625            Ok(Some("backend token".to_owned()))
626        }
627
628        async fn upload_offset(
629            &self,
630            _inner: &InMemoryBackend,
631            session: &Session,
632        ) -> Result<UploadProgress> {
633            assert_eq!(session.upload_length.get(), 4);
634            self.seen_tokens
635                .lock()
636                .unwrap()
637                .push(session.backend_token.clone());
638            Ok(UploadProgress::Incomplete { offset: 0 })
639        }
640    }
641
642    fn make_context() -> ObjectContext {
643        ObjectContext {
644            usecase: "testing".into(),
645            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
646        }
647    }
648
649    fn make_service() -> StorageService {
650        StorageService::new(
651            Box::new(InMemoryBackend::new("in-memory")),
652            Cipher::ephemeral().unwrap(),
653        )
654    }
655
656    #[tokio::test]
657    async fn insert_without_key_generates_unique_id() {
658        let service = make_service();
659
660        let id = service
661            .insert_object(
662                make_context(),
663                None,
664                Metadata::default(),
665                stream::single("auto-keyed"),
666                Timestamp::now(),
667            )
668            .await
669            .unwrap();
670
671        assert!(uuid::Uuid::parse_str(id.key()).is_ok());
672    }
673
674    #[tokio::test]
675    async fn stores_files() {
676        let service = make_service();
677
678        let key = service
679            .insert_object(
680                make_context(),
681                Some("testing".into()),
682                Metadata::default(),
683                stream::single("oh hai!"),
684                Timestamp::now(),
685            )
686            .await
687            .unwrap();
688
689        let (_metadata, _, stream) = service
690            .get_object(key, Timestamp::now(), None)
691            .await
692            .unwrap()
693            .unwrap();
694        let file_contents: BytesMut = stream.try_collect().await.unwrap();
695
696        assert_eq!(file_contents.as_ref(), b"oh hai!");
697    }
698
699    #[tokio::test]
700    async fn works_with_gcs() {
701        let config = GcsConfig {
702            endpoint: Some("http://localhost:8087".into()),
703            bucket: "test-bucket".into(), // aligned with the env var in devservices and CI
704            cogs: None,
705        };
706
707        let backend = GcsBackend::new(config, &ChangeStreamFactory::default())
708            .await
709            .unwrap();
710        let service = StorageService::new(Box::new(backend), Cipher::ephemeral().unwrap());
711
712        let key = service
713            .insert_object(
714                make_context(),
715                Some("testing".into()),
716                Metadata::default(),
717                stream::single("oh hai!"),
718                Timestamp::now(),
719            )
720            .await
721            .unwrap();
722
723        let (_metadata, _, stream) = service
724            .get_object(key, Timestamp::now(), None)
725            .await
726            .unwrap()
727            .unwrap();
728        let file_contents: BytesMut = stream.try_collect().await.unwrap();
729
730        assert_eq!(file_contents.as_ref(), b"oh hai!");
731    }
732
733    #[tokio::test]
734    async fn tombstone_redirect_and_delete() {
735        let bigtable_config = BigTableConfig {
736            endpoint: Some("localhost:8086".into()),
737            project_id: "testing".into(),
738            instance_name: "objectstore".into(),
739            table_name: "objectstore".into(),
740            connections: None,
741            rpc_timeout: Duration::from_secs(2),
742            cogs: None,
743        };
744        let gcs_config = GcsConfig {
745            endpoint: Some("http://localhost:8087".into()),
746            bucket: "test-bucket".into(),
747            cogs: None,
748        };
749
750        let high_volume = Box::new(
751            BigTableBackend::new(bigtable_config, &ChangeStreamFactory::default())
752                .await
753                .unwrap(),
754        );
755        let long_term = Box::new(
756            GcsBackend::new(gcs_config.clone(), &ChangeStreamFactory::default())
757                .await
758                .unwrap(),
759        );
760        let backend = TieredStorage::new(high_volume, long_term, Box::new(NoopChangeLog));
761        let service = StorageService::new(Box::new(backend), Cipher::ephemeral().unwrap());
762
763        // A separate GCS backend to directly inspect the long-term storage.
764        let gcs_backend = GcsBackend::new(gcs_config.clone(), &ChangeStreamFactory::default())
765            .await
766            .unwrap();
767
768        // Insert a >1 MiB object with a key.  This forces the long-term path:
769        // the real payload goes to GCS, and a redirect tombstone is written to BigTable.
770        let payload_len = 2 * 1024 * 1024;
771        let payload = vec![0xAB; payload_len]; // 2 MiB
772        let id = service
773            .insert_object(
774                make_context(),
775                Some("delete-cleanup-test".into()),
776                Metadata::default(),
777                stream::single(payload),
778                Timestamp::now(),
779            )
780            .await
781            .unwrap();
782
783        // Sanity: the object is readable through the service (follows the tombstone).
784        let (_, _, stream) = service
785            .get_object(id.clone(), Timestamp::now(), None)
786            .await
787            .unwrap()
788            .unwrap();
789        let body: BytesMut = stream.try_collect().await.unwrap();
790        assert_eq!(body.len(), payload_len);
791
792        // Delete through the service layer.
793        service
794            .delete_object(id.clone(), Timestamp::now())
795            .await
796            .unwrap();
797
798        // The tombstone in BigTable should be gone, so the service returns None.
799        let after_delete = service
800            .get_object(id.clone(), Timestamp::now(), None)
801            .await
802            .unwrap();
803        assert!(after_delete.is_none(), "tombstone not deleted");
804
805        // The real object in GCS must also be gone — no orphan.
806        let orphan = gcs_backend
807            .get_object(&id, Timestamp::now(), None)
808            .await
809            .unwrap();
810        assert!(orphan.is_none(), "object leaked");
811    }
812
813    // --- Task spawning tests (public API) ---
814
815    #[tokio::test]
816    async fn basic_spawn_insert_and_get() {
817        let service = make_service();
818
819        let id = service
820            .insert_object(
821                make_context(),
822                Some("test-key".into()),
823                Metadata::default(),
824                stream::single("hello world"),
825                Timestamp::now(),
826            )
827            .await
828            .unwrap();
829
830        let (_, _, stream) = service
831            .get_object(id, Timestamp::now(), None)
832            .await
833            .unwrap()
834            .unwrap();
835        let body: BytesMut = stream.try_collect().await.unwrap();
836        assert_eq!(body.as_ref(), b"hello world");
837    }
838
839    #[tokio::test]
840    async fn basic_spawn_metadata_and_delete() {
841        let service = make_service();
842
843        let id = service
844            .insert_object(
845                make_context(),
846                Some("meta-key".into()),
847                Metadata::default(),
848                stream::single("data"),
849                Timestamp::now(),
850            )
851            .await
852            .unwrap();
853
854        let metadata = service
855            .get_metadata(id.clone(), Timestamp::now())
856            .await
857            .unwrap();
858        assert!(metadata.is_some());
859
860        service
861            .delete_object(id.clone(), Timestamp::now())
862            .await
863            .unwrap();
864
865        let after = service
866            .get_object(id, Timestamp::now(), None)
867            .await
868            .unwrap();
869        assert!(after.is_none());
870    }
871
872    #[tokio::test]
873    async fn set_expiry() {
874        let service = make_service();
875        let old_expiry = Timestamp::now() + Duration::from_hours(1);
876        let metadata = Metadata {
877            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
878            time_expires: Some(old_expiry),
879            ..Default::default()
880        };
881        let id = service
882            .insert_object(
883                make_context(),
884                Some("explicit-expiry".into()),
885                metadata,
886                stream::single("payload"),
887                Timestamp::now(),
888            )
889            .await
890            .unwrap();
891        let requested = old_expiry + Duration::from_hours(1);
892
893        assert_eq!(
894            service
895                .set_expiry(
896                    id.clone(),
897                    ExpiryTarget::At(requested).into(),
898                    Timestamp::now()
899                )
900                .await
901                .unwrap(),
902            SetExpiryResponse::Satisfied(requested)
903        );
904        assert_eq!(
905            service
906                .get_metadata(id, Timestamp::now())
907                .await
908                .unwrap()
909                .unwrap()
910                .time_expires,
911            Some(requested)
912        );
913    }
914
915    #[derive(Debug)]
916    struct PanicOnGet;
917
918    #[async_trait::async_trait]
919    impl Hooks for PanicOnGet {
920        async fn get_object(
921            &self,
922            _inner: &InMemoryBackend,
923            _id: &ObjectId,
924            _access_time: Timestamp,
925            _range: Option<ByteRange>,
926        ) -> Result<GetResponse> {
927            panic!("intentional panic in get_object");
928        }
929    }
930
931    #[tokio::test]
932    async fn panic_in_backend_returns_task_failed() {
933        let service = StorageService::new(
934            Box::new(TestBackend::new(PanicOnGet)),
935            Cipher::ephemeral().unwrap(),
936        );
937
938        let id = ObjectId::new(make_context(), "panic-test".into());
939        let result = service.get_object(id, Timestamp::now(), None).await;
940
941        let Err(error) = result else {
942            panic!("expected Panic error");
943        };
944        assert_eq!(error.kind(), ErrorKind::Panic);
945        assert_eq!(error.to_string(), "service task panicked");
946        assert_eq!(
947            std::error::Error::source(&error).unwrap().to_string(),
948            "intentional panic in get_object"
949        );
950    }
951
952    #[derive(Clone, Debug, Default)]
953    struct GateOnExpiry {
954        calls: Arc<AtomicUsize>,
955        access_time: Arc<Mutex<Option<Timestamp>>>,
956        started: Arc<tokio::sync::Notify>,
957        resume: Arc<tokio::sync::Notify>,
958    }
959
960    #[async_trait::async_trait]
961    impl Hooks for GateOnExpiry {
962        async fn set_expiry(
963            &self,
964            inner: &InMemoryBackend,
965            id: &ObjectId,
966            target: ExpiryUpdate,
967            access_time: Timestamp,
968        ) -> Result<SetExpiryResponse> {
969            *self.access_time.lock().unwrap() = Some(access_time);
970            self.calls.fetch_add(1, Ordering::SeqCst);
971            self.started.notify_one();
972            self.resume.notified().await;
973            inner.set_expiry(id, target, access_time).await
974        }
975    }
976
977    fn stale_tti_metadata() -> Metadata {
978        Metadata {
979            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
980            time_expires: Some(Timestamp::now() + Duration::from_mins(1)),
981            ..Default::default()
982        }
983    }
984
985    #[tokio::test]
986    async fn background_renewal() {
987        let now = Timestamp::now();
988        let access_time = now - Duration::from_secs(10);
989        let backend = TestBackend::new(GateOnExpiry::default());
990        let id = ObjectId::new(make_context(), "background-renewal".into());
991        let metadata = stale_tti_metadata();
992        backend
993            .inner
994            .put_object(&id, &metadata, stream::single("payload"), Timestamp::now())
995            .await
996            .unwrap();
997        let mut service =
998            StorageService::new(Box::new(backend.clone()), Cipher::ephemeral().unwrap());
999        service.start();
1000
1001        let response = tokio::time::timeout(
1002            Duration::from_secs(1),
1003            service.get_object(id.clone(), access_time, Some(ByteRange::Bounded(0, 2))),
1004        )
1005        .await
1006        .expect("GET waited for its background renewal")
1007        .unwrap()
1008        .unwrap();
1009        assert_eq!(response.0.time_expires, metadata.time_expires);
1010        backend.hooks.started.notified().await;
1011        assert!(backend.hooks.access_time.lock().unwrap().unwrap() >= now);
1012
1013        let join = tokio::spawn({
1014            let service = service.clone();
1015            async move { service.join().await }
1016        });
1017        tokio::pin!(join);
1018        assert!(
1019            tokio::time::timeout(Duration::from_millis(25), &mut join)
1020                .await
1021                .is_err()
1022        );
1023
1024        backend.hooks.resume.notify_waiters();
1025        tokio::time::timeout(Duration::from_secs(1), &mut join)
1026            .await
1027            .expect("shutdown did not drain renewal")
1028            .unwrap();
1029        assert_eq!(
1030            backend.inner.get(&id).expect_object().0.time_expires,
1031            Some(access_time + Duration::from_hours(1))
1032        );
1033    }
1034
1035    #[tokio::test]
1036    async fn renewal_deduplication() {
1037        let backend = TestBackend::new(GateOnExpiry::default());
1038        let id = ObjectId::new(make_context(), "deduplicated-renewal".into());
1039        backend
1040            .inner
1041            .put_object(
1042                &id,
1043                &stale_tti_metadata(),
1044                stream::single("payload"),
1045                Timestamp::now(),
1046            )
1047            .await
1048            .unwrap();
1049        let mut service =
1050            StorageService::new(Box::new(backend.clone()), Cipher::ephemeral().unwrap());
1051        service.start();
1052
1053        service
1054            .get_metadata(id.clone(), Timestamp::now())
1055            .await
1056            .unwrap();
1057        backend.hooks.started.notified().await;
1058        service.get_metadata(id, Timestamp::now()).await.unwrap();
1059        tokio::task::yield_now().await;
1060        assert_eq!(backend.hooks.calls.load(Ordering::SeqCst), 1);
1061
1062        backend.hooks.resume.notify_waiters();
1063        service.join().await;
1064    }
1065
1066    #[tokio::test]
1067    async fn ttl_read() {
1068        let backend = TestBackend::new(GateOnExpiry::default());
1069        let id = ObjectId::new(make_context(), "ttl-no-renewal".into());
1070        let metadata = Metadata {
1071            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1072            time_expires: Some(Timestamp::now() + Duration::from_mins(1)),
1073            ..Default::default()
1074        };
1075        backend
1076            .inner
1077            .put_object(&id, &metadata, stream::single("payload"), Timestamp::now())
1078            .await
1079            .unwrap();
1080        let mut service =
1081            StorageService::new(Box::new(backend.clone()), Cipher::ephemeral().unwrap());
1082        service.start();
1083
1084        service.get_metadata(id, Timestamp::now()).await.unwrap();
1085        tokio::task::yield_now().await;
1086        assert_eq!(backend.hooks.calls.load(Ordering::SeqCst), 0);
1087        service.join().await;
1088    }
1089
1090    #[tokio::test]
1091    async fn renewal_queueing() {
1092        let backend = TestBackend::new(GateOnExpiry::default());
1093        let first = ObjectId::new(make_context(), "first-renewal".into());
1094        let second = ObjectId::new(make_context(), "queued-renewal".into());
1095        let metadata = stale_tti_metadata();
1096        for id in [&first, &second] {
1097            backend
1098                .inner
1099                .put_object(id, &metadata, stream::single("payload"), Timestamp::now())
1100                .await
1101                .unwrap();
1102        }
1103
1104        let concurrency = ConcurrencyLimiter::new(1);
1105        let mut scheduler = RenewalScheduler::new(Arc::new(backend.clone()), concurrency, 1);
1106        let expire_at = Timestamp::now() + Duration::from_hours(1);
1107        scheduler.schedule(first, expire_at);
1108        assert_eq!(scheduler.queued(), 1);
1109        scheduler.start();
1110        backend.hooks.started.notified().await;
1111        assert_eq!(scheduler.queued(), 0);
1112
1113        scheduler.schedule(second.clone(), expire_at);
1114        tokio::task::yield_now().await;
1115        assert_eq!(backend.hooks.calls.load(Ordering::SeqCst), 1);
1116        assert_eq!(
1117            backend.inner.get(&second).expect_object().0.time_expires,
1118            metadata.time_expires
1119        );
1120
1121        backend.hooks.resume.notify_waiters();
1122        tokio::time::timeout(Duration::from_secs(1), async {
1123            while backend.hooks.calls.load(Ordering::SeqCst) < 2 {
1124                tokio::task::yield_now().await;
1125            }
1126        })
1127        .await
1128        .expect("queued renewal did not start");
1129        backend.hooks.resume.notify_waiters();
1130        scheduler.join().await;
1131        assert_eq!(scheduler.queued(), 0);
1132        assert_eq!(backend.hooks.calls.load(Ordering::SeqCst), 2);
1133        assert!(backend.inner.get(&second).expect_object().0.time_expires > metadata.time_expires);
1134    }
1135
1136    #[derive(Clone, Debug, Default)]
1137    struct FailFirstExpiry {
1138        calls: Arc<AtomicUsize>,
1139        panic: bool,
1140    }
1141
1142    #[async_trait::async_trait]
1143    impl Hooks for FailFirstExpiry {
1144        async fn set_expiry(
1145            &self,
1146            inner: &InMemoryBackend,
1147            id: &ObjectId,
1148            target: ExpiryUpdate,
1149            access_time: Timestamp,
1150        ) -> Result<SetExpiryResponse> {
1151            if self.calls.fetch_add(1, Ordering::SeqCst) == 0 {
1152                assert!(!self.panic, "intentional renewal panic");
1153                return Err(ErrorKind::BackendFailure.into());
1154            }
1155            inner.set_expiry(id, target, access_time).await
1156        }
1157    }
1158
1159    #[tokio::test]
1160    async fn renewal_failure_cleanup() {
1161        for panic in [false, true] {
1162            let backend = TestBackend::new(FailFirstExpiry {
1163                panic,
1164                ..Default::default()
1165            });
1166            let id = ObjectId::new(make_context(), "failed-renewal".into());
1167            let metadata = stale_tti_metadata();
1168            backend
1169                .inner
1170                .put_object(&id, &metadata, stream::single("payload"), Timestamp::now())
1171                .await
1172                .unwrap();
1173            let concurrency = ConcurrencyLimiter::new(1);
1174            let mut scheduler = RenewalScheduler::new(Arc::new(backend.clone()), concurrency, 1);
1175            scheduler.start();
1176            let expire_at = metadata.check_tti_bump(Timestamp::now()).unwrap();
1177
1178            scheduler.schedule(id.clone(), expire_at);
1179            tokio::time::timeout(Duration::from_secs(1), async {
1180                while scheduler.pending() != 0 {
1181                    tokio::task::yield_now().await;
1182                }
1183            })
1184            .await
1185            .expect("failure did not release renewal guards");
1186
1187            scheduler.schedule(id, expire_at);
1188            scheduler.join().await;
1189            assert_eq!(backend.hooks.calls.load(Ordering::SeqCst), 2);
1190        }
1191    }
1192
1193    /// In-memory backend with optional synchronization for `put_object`.
1194    ///
1195    /// When `pause` is enabled, each `put_object` call notifies `paused` and
1196    #[derive(Clone, Debug, Default)]
1197    struct GateOnPut {
1198        pause: bool,
1199        paused: Arc<tokio::sync::Notify>,
1200        resume: Arc<tokio::sync::Notify>,
1201        on_put: Arc<tokio::sync::Notify>,
1202    }
1203
1204    impl GateOnPut {
1205        fn with_pause() -> Self {
1206            Self {
1207                pause: true,
1208                ..Default::default()
1209            }
1210        }
1211    }
1212
1213    #[async_trait::async_trait]
1214    impl Hooks for GateOnPut {
1215        async fn put_object(
1216            &self,
1217            inner: &InMemoryBackend,
1218            id: &ObjectId,
1219            metadata: &Metadata,
1220            stream: ClientStream,
1221            access_time: Timestamp,
1222        ) -> Result<PutResponse> {
1223            if self.pause {
1224                self.paused.notify_one();
1225                self.resume.notified().await;
1226            }
1227            inner.put_object(id, metadata, stream, access_time).await?;
1228            self.on_put.notify_one();
1229            Ok(())
1230        }
1231
1232        async fn compare_and_write(
1233            &self,
1234            inner: &InMemoryBackend,
1235            id: &ObjectId,
1236            current: Option<&ObjectId>,
1237            write: TieredWrite,
1238            access_time: Timestamp,
1239        ) -> Result<bool> {
1240            let notify = matches!(write, TieredWrite::Tombstone(_) | TieredWrite::Object(_, _));
1241            let result = inner
1242                .compare_and_write(id, current, write, access_time)
1243                .await?;
1244            if notify {
1245                self.on_put.notify_one();
1246            }
1247            Ok(result)
1248        }
1249    }
1250
1251    #[tokio::test]
1252    async fn receiver_drop_does_not_prevent_completion() {
1253        let hv = Box::new(TestBackend::new(GateOnPut::default()));
1254        let lt = Box::new(TestBackend::new(GateOnPut::with_pause()));
1255        let backend = TieredStorage::new(hv.clone(), lt.clone(), Box::new(NoopChangeLog));
1256        let service = StorageService::new(Box::new(backend), Cipher::ephemeral().unwrap());
1257
1258        let payload = vec![0xABu8; 2 * 1024 * 1024]; // 2 MiB → long-term path
1259        let request = service.insert_object(
1260            make_context(),
1261            Some("completion-test".into()),
1262            Metadata::default(),
1263            stream::single(payload),
1264            Timestamp::now(),
1265        );
1266
1267        // Start insert through the public API. select! drops the future once the
1268        // backend signals it has paused, simulating a client disconnect mid-write.
1269        let paused = Arc::clone(&lt.hooks.paused);
1270        tokio::select! {
1271            _ = request => panic!("insert should not complete while backend is paused"),
1272            _ = paused.notified() => {}
1273        }
1274
1275        // The spawned task is now blocked inside put_object, and the caller
1276        // request (including the oneshot receiver) has been dropped. Unpause so
1277        // the task can finish writing.
1278        lt.hooks.resume.notify_one();
1279
1280        // Wait for the tombstone write to the high-volume backend, which is the
1281        // last step of the long-term insert path.
1282        let on_put = Arc::clone(&hv.hooks.on_put);
1283        tokio::time::timeout(Duration::from_secs(5), on_put.notified())
1284            .await
1285            .expect("timed out waiting for tombstone write");
1286
1287        // Verify the object was fully written despite the caller being dropped.
1288        // The tombstone in HV points to the revision key in LT.
1289        let id = ObjectId::new(make_context(), "completion-test".into());
1290        let tombstone = hv.inner.get(&id).expect_tombstone();
1291        let lt_id = tombstone.target;
1292        assert!(lt.inner.contains(&lt_id), "long-term object missing");
1293    }
1294
1295    // --- Concurrency limit tests ---
1296
1297    fn make_limited_service(limit: u32) -> (StorageService, TestBackend<GateOnPut>) {
1298        let backend = TestBackend::new(GateOnPut::with_pause());
1299        let service = StorageService::new(Box::new(backend.clone()), Cipher::ephemeral().unwrap())
1300            .with_concurrency(ConcurrencyLimiter::new(limit));
1301        (service, backend)
1302    }
1303
1304    #[tokio::test]
1305    async fn at_capacity_rejects() {
1306        let (service, hv) = make_limited_service(1);
1307
1308        // First insert blocks on the gated backend, holding the single permit.
1309        let svc = service.clone();
1310        let first = tokio::spawn(async move {
1311            svc.insert_object(
1312                make_context(),
1313                Some("first".into()),
1314                Metadata::default(),
1315                stream::single("data"),
1316                Timestamp::now(),
1317            )
1318            .await
1319        });
1320
1321        // Wait for the backend to signal it has paused (permit is held).
1322        hv.hooks.paused.notified().await;
1323
1324        // Second insert should be rejected immediately.
1325        let result = service
1326            .insert_object(
1327                make_context(),
1328                Some("second".into()),
1329                Metadata::default(),
1330                stream::single("data"),
1331                Timestamp::now(),
1332            )
1333            .await;
1334
1335        assert!(
1336            result
1337                .as_ref()
1338                .is_err_and(|error| error.kind() == ErrorKind::AtCapacity),
1339            "expected AtCapacity, got {result:?}"
1340        );
1341
1342        // Unblock the first operation.
1343        hv.hooks.resume.notify_one();
1344        first.await.unwrap().unwrap();
1345
1346        // Now that the permit is released, a new operation should succeed.
1347        service
1348            .get_metadata(
1349                ObjectId::new(make_context(), "first".into()),
1350                Timestamp::now(),
1351            )
1352            .await
1353            .unwrap();
1354    }
1355
1356    #[tokio::test]
1357    async fn tasks_limit_returns_configured_limit() {
1358        let backend = Box::new(InMemoryBackend::new("cap"));
1359        let service = StorageService::new(backend, Cipher::ephemeral().unwrap())
1360            .with_concurrency(ConcurrencyLimiter::new(7));
1361        assert_eq!(service.tasks_limit(), 7);
1362    }
1363
1364    #[tokio::test]
1365    async fn tasks_running_tracks_in_flight() {
1366        let (service, hv) = make_limited_service(5);
1367
1368        assert_eq!(service.tasks_running(), 0);
1369
1370        // Kick off a request that blocks in the backend, holding a permit.
1371        let svc = service.clone();
1372        let _blocked = tokio::spawn(async move {
1373            svc.insert_object(
1374                make_context(),
1375                Some("in-use-test".into()),
1376                Metadata::default(),
1377                stream::single("data"),
1378                Timestamp::now(),
1379            )
1380            .await
1381        });
1382
1383        hv.hooks.paused.notified().await;
1384        assert_eq!(service.tasks_running(), 1);
1385
1386        hv.hooks.resume.notify_one();
1387    }
1388
1389    #[tokio::test]
1390    async fn permits_released_after_panic() {
1391        let service = StorageService::new(
1392            Box::new(TestBackend::new(PanicOnGet)),
1393            Cipher::ephemeral().unwrap(),
1394        )
1395        .with_concurrency(ConcurrencyLimiter::new(1));
1396
1397        // First operation panics — the permit must still be released.
1398        let id = ObjectId::new(make_context(), "panic-permit".into());
1399        let result = service.get_object(id.clone(), Timestamp::now(), None).await;
1400        assert!(result.is_err_and(|error| error.kind() == ErrorKind::Panic));
1401
1402        // Second operation should succeed in acquiring the permit (not AtCapacity).
1403        let result = service.get_object(id, Timestamp::now(), None).await;
1404        assert!(
1405            !result.is_err_and(|error| error.kind() == ErrorKind::AtCapacity),
1406            "permit was not released after panic"
1407        );
1408    }
1409
1410    // --- Resumable uploads ---
1411
1412    #[tokio::test]
1413    async fn resumable_round_trip() {
1414        let service = make_service();
1415        let id = ObjectId::new(make_context(), "resumable".into());
1416        let created = service
1417            .create_upload_session(id.clone(), Metadata::default(), 3)
1418            .await
1419            .unwrap()
1420            .unwrap();
1421        assert_eq!(created.granularity, 0);
1422        let token = created.session;
1423        assert_eq!(
1424            service
1425                .put_chunk(id.clone(), token.clone(), 0, 1, stream::single("a"))
1426                .await
1427                .unwrap(),
1428            UploadProgress::Incomplete { offset: 1 }
1429        );
1430        assert_eq!(
1431            service
1432                .upload_offset(id.clone(), token.clone())
1433                .await
1434                .unwrap(),
1435            UploadProgress::Incomplete { offset: 1 }
1436        );
1437        assert_eq!(
1438            service
1439                .put_chunk(id.clone(), token, 1, 2, stream::single("bc"))
1440                .await
1441                .unwrap(),
1442            UploadProgress::Complete
1443        );
1444        let (metadata, _, body) = service
1445            .get_object(id, Timestamp::now(), None)
1446            .await
1447            .unwrap()
1448            .unwrap();
1449        assert_eq!(metadata.size, Some(3));
1450        assert_eq!(
1451            body.try_collect::<BytesMut>().await.unwrap().as_ref(),
1452            b"abc"
1453        );
1454    }
1455
1456    #[tokio::test]
1457    async fn resumable_create_declines_zero_length() {
1458        let service = StorageService::new(
1459            Box::new(TestBackend::new(ResumableTokenHooks::default())),
1460            Cipher::ephemeral().unwrap(),
1461        );
1462        let id = ObjectId::new(make_context(), "resumable".into());
1463
1464        let result = service
1465            .create_upload_session(id, Metadata::default(), 0)
1466            .await;
1467
1468        assert!(matches!(result, Ok(None)), "{result:?}");
1469    }
1470
1471    #[tokio::test]
1472    async fn resumable_create_validates_metadata() {
1473        let service = make_service();
1474        let id = ObjectId::new(make_context(), "resumable".into());
1475
1476        // A timeout policy with no resolved `time_expires` is rejected before the backend
1477        // is consulted, exactly as it is for a regular insert.
1478        let metadata = Metadata {
1479            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(60)),
1480            ..Default::default()
1481        };
1482
1483        let result = service.create_upload_session(id, metadata, 1024).await;
1484        assert!(result.is_err_and(|error| error.kind() == ErrorKind::InvalidMetadata));
1485    }
1486
1487    #[tokio::test]
1488    async fn resumable_tokens_are_encrypted_by_default() -> Result<()> {
1489        let hooks = ResumableTokenHooks::default();
1490        let service = StorageService::new(
1491            Box::new(TestBackend::new(hooks.clone())),
1492            Cipher::ephemeral().unwrap(),
1493        );
1494        let id = ObjectId::new(make_context(), "resumable".into());
1495
1496        let created = service
1497            .create_upload_session(id.clone(), Metadata::default(), 4)
1498            .await?
1499            .expect("test backend supports resumable uploads");
1500        assert_eq!(created.granularity, 256 * 1024);
1501        let token = created.session;
1502        assert_ne!(token.as_bytes(), b"backend token");
1503        assert!(matches!(
1504            service
1505                .upload_offset(id.clone(), EncryptedSessionToken::new(b"backend token"))
1506                .await,
1507            Err(error) if error.kind() == ErrorKind::UnknownUploadSession
1508        ));
1509        let other_id = ObjectId::new(make_context(), "other".into());
1510        assert!(matches!(
1511            service.upload_offset(other_id, token.clone()).await,
1512            Err(error) if error.kind() == ErrorKind::UnknownUploadSession
1513        ));
1514        assert!(hooks.seen_tokens.lock().unwrap().is_empty());
1515        service.upload_offset(id, token).await?;
1516        assert_eq!(
1517            hooks.seen_tokens.lock().unwrap().as_slice(),
1518            &["backend token"]
1519        );
1520        Ok(())
1521    }
1522
1523    #[tokio::test]
1524    async fn configured_encryption_only_crosses_the_service_boundary() -> Result<()> {
1525        let hooks = ResumableTokenHooks::default();
1526        let encryption = Cipher::new(
1527            "v1",
1528            std::collections::BTreeMap::from([("v1".into(), vec![7; 32])]),
1529        )
1530        .unwrap();
1531        let service = StorageService::new(Box::new(TestBackend::new(hooks.clone())), encryption);
1532        let id = ObjectId::new(make_context(), "resumable".into());
1533
1534        let encrypted = service
1535            .create_upload_session(id.clone(), Metadata::default(), 4)
1536            .await?
1537            .expect("test backend supports resumable uploads")
1538            .session;
1539        assert_ne!(encrypted.as_bytes(), b"backend token");
1540        service.upload_offset(id, encrypted).await?;
1541        assert_eq!(
1542            hooks.seen_tokens.lock().unwrap().as_slice(),
1543            &["backend token"]
1544        );
1545        Ok(())
1546    }
1547
1548    #[tokio::test]
1549    async fn configured_encryption_rejects_plaintext_tokens() {
1550        let hooks = ResumableTokenHooks::default();
1551        let encryption = Cipher::new(
1552            "v1",
1553            std::collections::BTreeMap::from([("v1".into(), vec![7; 32])]),
1554        )
1555        .unwrap();
1556        let service = StorageService::new(Box::new(TestBackend::new(hooks.clone())), encryption);
1557        let id = ObjectId::new(make_context(), "resumable".into());
1558
1559        let result = service
1560            .upload_offset(id, EncryptedSessionToken::new(b"backend token"))
1561            .await;
1562        let error = result.unwrap_err();
1563        assert_eq!(error.kind(), ErrorKind::UnknownUploadSession);
1564        assert_eq!(error.to_string(), "unknown upload session");
1565        assert!(error.source().is_none());
1566        assert!(hooks.seen_tokens.lock().unwrap().is_empty());
1567    }
1568}