1use 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
37pub type GetResponse = Option<(Metadata, Option<ContentRange>, PayloadStream)>;
39pub type MetadataResponse = Option<Metadata>;
41pub type InsertResponse = ObjectId;
43pub type DeleteResponse = ();
45
46#[derive(Debug)]
48pub struct CreateUploadSessionResponse {
49 pub session: EncryptedSessionToken,
51 pub granularity: u64,
53}
54
55pub const DEFAULT_CONCURRENCY_LIMIT: u32 = 500;
60
61pub const DEFAULT_BACKGROUND_QUEUE_LIMIT: usize = 1_000;
63
64#[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 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 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 pub fn with_background_queue_limit(mut self, limit: usize) -> Self {
134 self.renewals.set_capacity(limit);
135 self
136 }
137
138 pub fn concurrency_limiter(&self) -> &ConcurrencyLimiter {
140 &self.concurrency
141 }
142
143 pub fn tasks_running(&self) -> u32 {
145 self.concurrency.used_permits()
146 }
147
148 pub fn tasks_limit(&self) -> u32 {
150 self.concurrency.total_permits()
151 }
152
153 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 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 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 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 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 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 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 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 pub async fn join(&self) {
361 self.renewals.join().await;
362 self.inner.join().await;
363 }
364
365 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()?; 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 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()?; 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 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()?; 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 pub async fn abort_multipart(
441 &self,
442 id: ObjectId,
443 upload_id: UploadId,
444 ) -> Result<AbortMultipartResponse> {
445 self.inner.as_multipart_upload_backend()?; 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 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()?; 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 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 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 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 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(), 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 let gcs_backend = GcsBackend::new(gcs_config.clone(), &ChangeStreamFactory::default())
765 .await
766 .unwrap();
767
768 let payload_len = 2 * 1024 * 1024;
771 let payload = vec![0xAB; payload_len]; 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 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 service
794 .delete_object(id.clone(), Timestamp::now())
795 .await
796 .unwrap();
797
798 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 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 #[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 #[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]; 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 let paused = Arc::clone(<.hooks.paused);
1270 tokio::select! {
1271 _ = request => panic!("insert should not complete while backend is paused"),
1272 _ = paused.notified() => {}
1273 }
1274
1275 lt.hooks.resume.notify_one();
1279
1280 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 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(<_id), "long-term object missing");
1293 }
1294
1295 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 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 hv.hooks.paused.notified().await;
1323
1324 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 hv.hooks.resume.notify_one();
1344 first.await.unwrap().unwrap();
1345
1346 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 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 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 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 #[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 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}