1use std::fmt;
36use std::future::Future;
37use std::sync::Arc;
38use std::time::Duration;
39
40use bigtable_rs::bigtable::{BigTableConnection, Error as BigTableError, RowCell};
41use bigtable_rs::google::bigtable::v2::{self, mutation};
42use bytes::Bytes;
43use futures_util::TryStreamExt;
44use objectstore_types::metadata::Metadata;
45use objectstore_types::range::{ByteRange, ContentRange};
46use objectstore_types::time::Timestamp;
47use serde::{Deserialize, Serialize};
48use tonic::Code;
49use tracing::Instrument;
50
51use crate::backend::common::{
52 self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, HighVolumeBackend, MetadataResponse,
53 PutResponse, SetExpiryResponse, TieredGet, TieredMetadata, TieredUpdate, TieredWrite,
54 Tombstone,
55};
56use crate::change_stream::{
57 ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream,
58};
59use crate::error::{Error, ErrorKind, Result, ResultExt as _};
60use crate::gcp_auth::PrefetchingTokenProvider;
61use crate::id::ObjectId;
62use crate::stream::{ChunkedBytes, ClientStream};
63
64#[derive(Debug, Clone, Deserialize, Serialize)]
86pub struct BigTableConfig {
87 pub endpoint: Option<String>,
100
101 pub project_id: String,
109
110 pub instance_name: String,
116
117 pub table_name: String,
125
126 pub connections: Option<usize>,
136
137 #[serde(default = "default_rpc_timeout", with = "humantime_serde")]
148 pub rpc_timeout: Duration,
149
150 #[serde(default, skip_serializing_if = "Option::is_none")]
161 pub cogs: Option<CostTrackerStreamConfig>,
162}
163
164fn default_rpc_timeout() -> Duration {
165 Duration::from_secs(2)
166}
167
168const MAX_CHANNEL_AGE: Option<Duration> = Some(Duration::from_mins(50));
176const TOKEN_SCOPES: &[&str] = &["https://www.googleapis.com/auth/bigtable.data"];
178
179const REQUEST_RETRY_COUNT: usize = 2;
181const CAS_RETRY_COUNT: usize = 3;
183
184const COLUMN_PAYLOAD: &[u8] = b"p";
186const COLUMN_METADATA: &[u8] = b"m";
188const COLUMN_REDIRECT: &[u8] = b"r";
190const FILTER_META: &[u8] = b"^[mr]$";
192
193const COLUMN_UPLOAD: &[u8] = b"u";
195
196const FAMILY_GC: &str = "fg";
201const FAMILY_MANUAL: &str = "fm";
203
204pub struct BigTableBackend {
206 bigtable: BigTableConnection,
207
208 instance_path: String,
209 table_path: String,
210 table_name: String,
211
212 change_stream: Arc<dyn ChangeStream>,
213}
214
215impl fmt::Debug for BigTableBackend {
216 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
217 f.debug_struct("BigTableBackend")
218 .field("instance_path", &self.instance_path)
219 .field("table_path", &self.table_path)
220 .field("table_name", &self.table_name)
221 .finish_non_exhaustive()
222 }
223}
224
225fn column_filter(column: &[u8]) -> v2::RowFilter {
227 v2::RowFilter {
228 filter: Some(v2::row_filter::Filter::ColumnQualifierRegexFilter(
229 [b"^", column, b"$"].concat(),
230 )),
231 }
232}
233
234fn legacy_tombstone_filter() -> v2::RowFilter {
239 v2::RowFilter {
240 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
241 filters: vec![
242 column_filter(COLUMN_METADATA),
243 v2::RowFilter {
244 filter: Some(v2::row_filter::Filter::ValueRegexFilter(
245 b"^\\{\"is_redirect_tombstone\":true[,}].*".to_vec(),
246 )),
247 },
248 ],
249 })),
250 }
251}
252
253fn live_row_filter(inner: v2::RowFilter, now: Timestamp) -> v2::RowFilter {
258 v2::RowFilter {
259 filter: Some(v2::row_filter::Filter::Interleave(
260 v2::row_filter::Interleave {
261 filters: vec![
262 v2::RowFilter {
264 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
265 filters: vec![
266 v2::RowFilter {
267 filter: Some(v2::row_filter::Filter::FamilyNameRegexFilter(
268 format!("^{FAMILY_MANUAL}$"),
269 )),
270 },
271 inner.clone(),
272 ],
273 })),
274 },
275 v2::RowFilter {
277 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
278 filters: vec![
279 v2::RowFilter {
280 filter: Some(v2::row_filter::Filter::FamilyNameRegexFilter(
281 format!("^{FAMILY_GC}$"),
282 )),
283 },
284 v2::RowFilter {
285 filter: Some(v2::row_filter::Filter::TimestampRangeFilter(
286 v2::TimestampRange {
287 start_timestamp_micros: now.as_micros() as i64,
288 end_timestamp_micros: 0,
289 },
290 )),
291 },
292 inner,
293 ],
294 })),
295 },
296 ],
297 },
298 )),
299 }
300}
301
302fn tombstone_filter(access_time: Timestamp) -> v2::RowFilter {
310 let filter = v2::RowFilter {
311 filter: Some(v2::row_filter::Filter::Interleave(
312 v2::row_filter::Interleave {
313 filters: vec![column_filter(COLUMN_REDIRECT), legacy_tombstone_filter()],
314 },
315 )),
316 };
317 live_row_filter(filter, access_time)
318}
319
320fn tombstone_predicate(access_time: Timestamp) -> MutatePredicate {
330 MutatePredicate::Exclude(tombstone_filter(access_time))
331}
332
333fn non_tombstone_predicate(access_time: Timestamp) -> MutatePredicate {
346 MutatePredicate::Include(v2::RowFilter {
347 filter: Some(v2::row_filter::Filter::Condition(Box::new(
348 v2::row_filter::Condition {
349 predicate_filter: Some(Box::new(tombstone_filter(access_time))),
350 true_filter: Some(Box::new(v2::RowFilter {
351 filter: Some(v2::row_filter::Filter::BlockAllFilter(true)),
352 })),
353 false_filter: Some(Box::new(v2::RowFilter {
354 filter: Some(v2::row_filter::Filter::PassAllFilter(true)),
355 })),
356 },
357 ))),
358 })
359}
360
361fn exact_value_regex(value: &str) -> Vec<u8> {
366 format!("^{}$", regex::escape(value)).into_bytes()
367}
368
369fn redirect_target_filter(
387 target: &ObjectId,
388 own_id: &ObjectId,
389 access_time: Timestamp,
390) -> v2::RowFilter {
391 let target_path = exact_value_regex(&target.as_storage_path().to_string());
392
393 let exact_match = v2::RowFilter {
394 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
395 filters: vec![
396 column_filter(COLUMN_REDIRECT),
397 v2::RowFilter {
398 filter: Some(v2::row_filter::Filter::ValueRegexFilter(target_path)),
399 },
400 ],
401 })),
402 };
403
404 if target != own_id {
405 return live_row_filter(exact_match, access_time);
406 }
407
408 let empty_redirect_match = v2::RowFilter {
409 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
410 filters: vec![
411 column_filter(COLUMN_REDIRECT),
412 v2::RowFilter {
413 filter: Some(v2::row_filter::Filter::ValueRegexFilter(b"^$".to_vec())),
414 },
415 ],
416 })),
417 };
418
419 let filter = v2::RowFilter {
423 filter: Some(v2::row_filter::Filter::Interleave(
424 v2::row_filter::Interleave {
425 filters: vec![exact_match, empty_redirect_match, legacy_tombstone_filter()],
426 },
427 )),
428 };
429 live_row_filter(filter, access_time)
430}
431
432fn update_predicate(
439 old: &ObjectId,
440 new: &ObjectId,
441 own_id: &ObjectId,
442 access_time: Timestamp,
443) -> MutatePredicate {
444 MutatePredicate::Include(v2::RowFilter {
445 filter: Some(v2::row_filter::Filter::Interleave(
446 v2::row_filter::Interleave {
447 filters: vec![
448 redirect_target_filter(old, own_id, access_time),
449 redirect_target_filter(new, own_id, access_time),
450 ],
451 },
452 )),
453 })
454}
455
456fn optional_target_predicate(
468 target: &ObjectId,
469 own_id: &ObjectId,
470 access_time: Timestamp,
471) -> MutatePredicate {
472 MutatePredicate::Exclude(v2::RowFilter {
473 filter: Some(v2::row_filter::Filter::Condition(Box::new(
474 v2::row_filter::Condition {
475 predicate_filter: Some(Box::new(redirect_target_filter(
476 target,
477 own_id,
478 access_time,
479 ))),
480 true_filter: Some(Box::new(v2::RowFilter {
481 filter: Some(v2::row_filter::Filter::BlockAllFilter(true)),
482 })),
483 false_filter: Some(Box::new(tombstone_filter(access_time))),
484 },
485 ))),
486 })
487}
488
489fn exact_expiry_filter(start: i64) -> Result<v2::RowFilter> {
490 let end = start.checked_add(1).ok_or_else(|| {
491 Error::new(
492 ErrorKind::Internal,
493 "building Bigtable expiration predicate",
494 )
495 })?;
496 Ok(v2::RowFilter {
497 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
498 filters: vec![
499 v2::RowFilter {
500 filter: Some(v2::row_filter::Filter::FamilyNameRegexFilter(format!(
501 "^{FAMILY_GC}$"
502 ))),
503 },
504 v2::RowFilter {
505 filter: Some(v2::row_filter::Filter::TimestampRangeFilter(
506 v2::TimestampRange {
507 start_timestamp_micros: start,
508 end_timestamp_micros: end,
509 },
510 )),
511 },
512 ],
513 })),
514 })
515}
516
517fn inline_expiry_predicate(
519 observed_expiry: i64,
520 access_time: Timestamp,
521) -> Result<MutatePredicate> {
522 let inline_at_expiry = v2::RowFilter {
523 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
524 filters: vec![
525 column_filter(COLUMN_METADATA),
526 exact_expiry_filter(observed_expiry)?,
527 ],
528 })),
529 };
530
531 Ok(MutatePredicate::Include(v2::RowFilter {
532 filter: Some(v2::row_filter::Filter::Condition(Box::new(
533 v2::row_filter::Condition {
534 predicate_filter: Some(Box::new(tombstone_filter(access_time))),
535 true_filter: Some(Box::new(v2::RowFilter {
536 filter: Some(v2::row_filter::Filter::BlockAllFilter(true)),
537 })),
538 false_filter: Some(Box::new(inline_at_expiry)),
539 },
540 ))),
541 }))
542}
543
544fn redirect_expiry_predicate(
545 target: &ObjectId,
546 own_id: &ObjectId,
547 observed_expiry: i64,
548 access_time: Timestamp,
549) -> Result<MutatePredicate> {
550 Ok(MutatePredicate::Include(v2::RowFilter {
551 filter: Some(v2::row_filter::Filter::Chain(v2::row_filter::Chain {
552 filters: vec![
553 redirect_target_filter(target, own_id, access_time),
554 exact_expiry_filter(observed_expiry)?,
555 ],
556 })),
557 }))
558}
559
560#[derive(Clone, Debug)]
565enum MutatePredicate {
566 Include(v2::RowFilter),
570 Exclude(v2::RowFilter),
574}
575
576fn metadata_filter() -> v2::RowFilter {
581 v2::RowFilter {
582 filter: Some(v2::row_filter::Filter::ColumnQualifierRegexFilter(
583 FILTER_META.to_owned(),
584 )),
585 }
586}
587
588fn mutation(mutation: mutation::Mutation) -> v2::Mutation {
589 v2::Mutation {
590 mutation: Some(mutation),
591 }
592}
593
594fn delete_row_mutation() -> v2::Mutation {
596 mutation(mutation::Mutation::DeleteFromRow(
597 mutation::DeleteFromRow {},
598 ))
599}
600
601fn object_mutations(
607 path: &[u8],
608 mut metadata: Metadata,
609 payload: Vec<u8>,
610) -> Result<([v2::Mutation; 3], u64)> {
611 let (family, timestamp_micros) = match metadata.time_expires {
612 None => (FAMILY_MANUAL, -1),
613 Some(deadline) => (FAMILY_GC, deadline.as_micros() as i64),
614 };
615
616 metadata.size = Some(payload.len());
618
619 let metadata_bytes = serde_json::to_vec(&metadata)
620 .context(ErrorKind::Internal, "encoding Bigtable object metadata")?;
621
622 let mutations = [
623 delete_row_mutation(),
625 mutation(mutation::Mutation::SetCell(mutation::SetCell {
626 family_name: family.to_owned(),
627 column_qualifier: COLUMN_PAYLOAD.to_owned(),
628 timestamp_micros,
629 value: payload,
630 })),
631 mutation(mutation::Mutation::SetCell(mutation::SetCell {
632 family_name: family.to_owned(),
633 column_qualifier: COLUMN_METADATA.to_owned(),
634 timestamp_micros,
635 value: metadata_bytes,
636 })),
637 ];
638
639 let size = row_size(path, &mutations);
640 Ok((mutations, size))
641}
642
643fn row_size(path: &[u8], mutations: &[v2::Mutation]) -> u64 {
648 let cells: usize = mutations
649 .iter()
650 .filter_map(|m| match &m.mutation {
651 Some(mutation::Mutation::SetCell(cell)) => Some(cell.value.len()),
652 _ => None,
653 })
654 .sum();
655
656 (path.len() + cells) as u64
657}
658
659fn tombstone_mutations(tombstone: &Tombstone) -> [v2::Mutation; 2] {
664 let (family, timestamp_micros) = match tombstone.time_expires {
665 None => (FAMILY_MANUAL, -1),
666 Some(deadline) => (FAMILY_GC, deadline.as_micros() as i64),
667 };
668
669 [
670 delete_row_mutation(),
671 mutation(mutation::Mutation::SetCell(mutation::SetCell {
672 family_name: family.to_owned(),
673 column_qualifier: COLUMN_REDIRECT.to_owned(),
674 timestamp_micros,
675 value: tombstone.target.as_storage_path().to_string().into_bytes(),
676 })),
677 ]
678}
679
680#[derive(Debug, Deserialize)]
684struct LegacyTombstoneMeta {
685 #[serde(default)]
692 is_redirect_tombstone: bool,
693}
694
695enum RowData {
697 Object {
699 metadata: Metadata,
700 payload: Vec<u8>,
701 expiry_micros: i64,
703 },
704 Tombstone {
706 target: Vec<u8>,
707 time_expires: Option<Timestamp>,
708 expiry_micros: i64,
710 },
711}
712
713impl RowData {
714 fn from_cells(cells: Vec<RowCell>) -> Result<Self> {
721 let mut metadata_opt: Option<Metadata> = None;
722 let mut redirect_detected = false;
723 let mut redirect_target = Vec::new();
724 let mut expire_at = None;
725 let mut expiry_micros = 0;
726 let mut payload = Vec::new();
727
728 for cell in cells {
729 if cell.family_name == FAMILY_GC {
734 expiry_micros = cell.timestamp_micros;
735 expire_at = Some(Timestamp::from_unix_micros(expiry_micros).context(
736 ErrorKind::CorruptData,
737 "decoding Bigtable expiration timestamp",
738 )?);
739 }
740
741 match cell.qualifier.as_slice() {
742 COLUMN_REDIRECT => {
743 redirect_detected = true;
744 redirect_target = cell.value;
745 }
746 COLUMN_PAYLOAD => {
747 payload = cell.value;
748 }
749 COLUMN_METADATA => {
750 if let Ok(legacy_meta) =
751 serde_json::from_slice::<LegacyTombstoneMeta>(&cell.value)
752 && legacy_meta.is_redirect_tombstone
753 {
754 redirect_detected = true;
755 objectstore_metrics::count!("bigtable.legacy_tombstone_read");
756 } else {
757 metadata_opt = Some(serde_json::from_slice(&cell.value).context(
758 ErrorKind::CorruptData,
759 "decoding Bigtable object metadata",
760 )?);
761 }
762 }
763 _ => {}
764 }
765 }
766
767 Ok(if redirect_detected {
768 RowData::Tombstone {
769 target: redirect_target,
770 time_expires: expire_at,
771 expiry_micros,
772 }
773 } else {
774 let mut metadata = metadata_opt.unwrap_or_default();
776 metadata.time_expires = expire_at;
777 RowData::Object {
778 metadata,
779 payload,
780 expiry_micros,
781 }
782 })
783 }
784
785 fn time_expires(&self) -> Option<Timestamp> {
787 match self {
788 RowData::Object { metadata, .. } => metadata.time_expires,
789 RowData::Tombstone { time_expires, .. } => *time_expires,
790 }
791 }
792
793 fn expires_before(&self, time: Timestamp) -> bool {
797 self.time_expires().is_some_and(|ts| ts < time)
798 }
799}
800
801fn parse_redirect_target(redirect_path: &[u8], tombstone_id: &ObjectId) -> Result<ObjectId> {
807 if redirect_path.is_empty() {
808 objectstore_metrics::count!("bigtable.empty_redirect_read");
809 Ok(tombstone_id.clone())
810 } else {
811 let redirect_str = std::str::from_utf8(redirect_path)
812 .context(ErrorKind::CorruptData, "decoding Bigtable redirect target")?;
813 ObjectId::from_storage_path(redirect_str)
814 .ok_or_else(|| Error::new(ErrorKind::CorruptData, "parsing Bigtable redirect target"))
815 }
816}
817
818impl BigTableBackend {
819 pub async fn new(
826 config: BigTableConfig,
827 streams: &ChangeStreamFactory,
828 ) -> anyhow::Result<Self> {
829 let BigTableConfig {
830 endpoint,
831 project_id,
832 instance_name,
833 table_name,
834 connections,
835 rpc_timeout,
836 cogs,
837 } = config;
838 let change_stream = streams.build(cogs.as_ref());
839
840 let bigtable = if let Some(ref endpoint) = endpoint {
841 BigTableConnection::new_with_emulator(
842 endpoint,
843 &project_id,
844 &instance_name,
845 false, connections.unwrap_or(1),
847 Some(rpc_timeout),
848 )?
849 } else {
850 let token_provider = PrefetchingTokenProvider::gcp_auth(TOKEN_SCOPES).await?;
851 BigTableConnection::new_with_managed_transport(
852 &project_id,
853 &instance_name,
854 false, Some(rpc_timeout),
856 Arc::new(token_provider),
857 connections.unwrap_or(1),
858 true, None, MAX_CHANNEL_AGE,
861 Some(Duration::from_secs(10)), )
863 .await?
864 };
865
866 let client = bigtable.client();
867
868 Ok(Self {
869 bigtable,
870 instance_path: format!("projects/{project_id}/instances/{instance_name}"),
871 table_path: client.get_full_table_name(&table_name),
872 table_name,
873 change_stream,
874 })
875 }
876
877 #[tracing::instrument(level = "debug", fields(action), skip_all)]
881 async fn read_row(
882 &self,
883 path: &[u8],
884 action: &'static str,
885 access_time: Timestamp,
886 filter: Option<v2::RowFilter>,
887 ) -> Result<Option<RowData>> {
888 let request = v2::ReadRowsRequest {
889 table_name: self.table_path.clone(),
890 rows: Some(v2::RowSet {
891 row_keys: vec![path.to_owned()],
892 row_ranges: vec![],
893 }),
894 filter,
895 rows_limit: 1,
896 ..Default::default()
897 };
898
899 let response = retry(action, || async {
900 self.bigtable.client().read_rows(request.clone()).await
901 })
902 .await?;
903 debug_assert!(response.len() <= 1, "Expected at most one row");
904
905 let Some((_, cells)) = response.into_iter().next() else {
906 objectstore_log::debug!("Object not found");
907 return Ok(None);
908 };
909
910 let row = RowData::from_cells(cells)?;
911 Ok(if row.expires_before(access_time) {
912 None
913 } else {
914 Some(row)
915 })
916 }
917
918 #[tracing::instrument(level = "debug", fields(action), skip_all)]
919 async fn mutate(
920 &self,
921 path: Vec<u8>,
922 mutations: impl Into<Vec<v2::Mutation>>,
923 action: &'static str,
924 ) -> Result<v2::MutateRowResponse> {
925 let request = v2::MutateRowRequest {
926 table_name: self.table_path.clone(),
927 row_key: path,
928 mutations: mutations.into(),
929 ..Default::default()
930 };
931
932 let response = retry(action, || async {
933 self.bigtable.client().mutate_row(request.clone()).await
934 })
935 .await?;
936
937 Ok(response.into_inner())
938 }
939
940 async fn put_row(
942 &self,
943 path: Vec<u8>,
944 metadata: Metadata,
945 payload: Vec<u8>,
946 action: &'static str,
947 ) -> Result<(v2::MutateRowResponse, u64)> {
948 let (mutations, size) = object_mutations(&path, metadata, payload)?;
949 let response = self.mutate(path, mutations, action).await?;
950 Ok((response, size))
951 }
952
953 #[tracing::instrument(level = "debug", fields(action = context), skip_all)]
955 async fn check_and_mutate(
956 &self,
957 row_key: Vec<u8>,
958 predicate: MutatePredicate,
959 mutations: impl Into<Vec<v2::Mutation>>,
960 context: &'static str,
961 ) -> Result<bool> {
962 let (filter, true_mutations, false_mutations, success_on_match) = match predicate {
963 MutatePredicate::Include(f) => (f, mutations.into(), vec![], true),
964 MutatePredicate::Exclude(f) => (f, vec![], mutations.into(), false),
965 };
966
967 let request = v2::CheckAndMutateRowRequest {
968 table_name: self.table_path.clone(),
969 row_key,
970 predicate_filter: Some(filter),
971 true_mutations,
972 false_mutations,
973 ..Default::default()
974 };
975
976 let future = retry(context, || async {
977 self.bigtable
978 .client()
979 .check_and_mutate_row(request.clone())
980 .await
981 });
982
983 Ok(future.await?.predicate_matched == success_on_match)
984 }
985}
986
987#[async_trait::async_trait]
988impl Backend for BigTableBackend {
989 fn name(&self) -> &'static str {
990 "bigtable"
991 }
992
993 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
994 async fn put_object(
995 &self,
996 id: &ObjectId,
997 metadata: &Metadata,
998 mut stream: ClientStream,
999 _access_time: Timestamp,
1000 ) -> Result<PutResponse> {
1001 objectstore_log::debug!("Writing to Bigtable backend");
1002 let path = id.as_storage_path().to_string().into_bytes();
1003
1004 let mut payload = ChunkedBytes::new(0);
1005 while let Some(chunk) = stream.try_next().await? {
1006 payload.push(chunk);
1007 }
1008
1009 let (_, size) = self
1010 .put_row(path, metadata.clone(), payload.into_bytes().into(), "put")
1011 .await?;
1012 self.change_stream.write(id, size, metadata.time_expires);
1013
1014 Ok(())
1015 }
1016
1017 #[tracing::instrument(level = "debug", skip(self))]
1018 async fn get_object(
1019 &self,
1020 id: &ObjectId,
1021 access_time: Timestamp,
1022 range: Option<ByteRange>,
1023 ) -> Result<GetResponse> {
1024 match self.get_tiered_object(id, access_time, range).await? {
1025 TieredGet::Object(metadata, content_range, payload) => {
1026 Ok(Some((metadata, content_range, payload)))
1027 }
1028 TieredGet::Tombstone(_) => Err(ErrorKind::UnexpectedTombstone.into()),
1029 TieredGet::NotFound => Ok(None),
1030 }
1031 }
1032
1033 #[tracing::instrument(level = "debug", skip(self))]
1034 async fn get_metadata(
1035 &self,
1036 id: &ObjectId,
1037 access_time: Timestamp,
1038 ) -> Result<MetadataResponse> {
1039 match self.get_tiered_metadata(id, access_time).await? {
1040 TieredMetadata::Object(metadata) => Ok(Some(metadata)),
1041 TieredMetadata::Tombstone(_) => Err(ErrorKind::UnexpectedTombstone.into()),
1042 TieredMetadata::NotFound => Ok(None),
1043 }
1044 }
1045
1046 async fn set_expiry(
1047 &self,
1048 id: &ObjectId,
1049 target: ExpiryUpdate,
1050 access_time: Timestamp,
1051 ) -> Result<SetExpiryResponse> {
1052 self.compare_and_update(id, None, TieredUpdate::SetExpiry(target), access_time)
1053 .await
1054 }
1055
1056 #[tracing::instrument(level = "debug", skip(self))]
1057 async fn delete_object(
1058 &self,
1059 id: &ObjectId,
1060 _access_time: Timestamp,
1061 ) -> Result<DeleteResponse> {
1062 objectstore_log::debug!("Deleting from Bigtable backend");
1063
1064 let path = id.as_storage_path().to_string().into_bytes();
1065 self.mutate(path, [delete_row_mutation()], "delete").await?;
1066 self.change_stream.delete(id);
1067
1068 Ok(())
1069 }
1070
1071 async fn join(&self) {
1072 flush_change_stream(&self.change_stream).await;
1073 }
1074}
1075
1076#[async_trait::async_trait]
1077impl HighVolumeBackend for BigTableBackend {
1078 async fn create_upload_marker(
1079 &self,
1080 revision: &ObjectId,
1081 time_expires: Timestamp,
1082 ) -> Result<()> {
1083 let mutations = vec![mutation(mutation::Mutation::SetCell(mutation::SetCell {
1084 family_name: FAMILY_GC.to_owned(),
1085 column_qualifier: COLUMN_UPLOAD.to_vec(),
1086 timestamp_micros: time_expires.as_micros() as i64,
1087 value: vec![1],
1088 }))];
1089 self.mutate(
1090 revision.as_upload_path().to_string().into_bytes(),
1091 mutations,
1092 "create_upload_marker",
1093 )
1094 .await?;
1095 Ok(())
1096 }
1097
1098 async fn has_upload_marker(&self, revision: &ObjectId, access_time: Timestamp) -> Result<bool> {
1099 let request = v2::ReadRowsRequest {
1100 table_name: self.table_path.clone(),
1101 rows: Some(v2::RowSet {
1102 row_keys: vec![revision.as_upload_path().to_string().into_bytes()],
1103 row_ranges: vec![],
1104 }),
1105 filter: Some(live_row_filter(column_filter(COLUMN_UPLOAD), access_time)),
1106 rows_limit: 1,
1107 ..Default::default()
1108 };
1109 let rows = retry("has_upload_marker", || async {
1110 self.bigtable.client().read_rows(request.clone()).await
1111 })
1112 .await?;
1113 Ok(!rows.is_empty())
1114 }
1115
1116 async fn delete_upload_marker(
1117 &self,
1118 revision: &ObjectId,
1119 access_time: Timestamp,
1120 ) -> Result<bool> {
1121 self.check_and_mutate(
1122 revision.as_upload_path().to_string().into_bytes(),
1123 MutatePredicate::Include(live_row_filter(column_filter(COLUMN_UPLOAD), access_time)),
1124 vec![delete_row_mutation()],
1125 "delete_upload_marker",
1126 )
1127 .await
1128 }
1129
1130 #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1131 async fn put_non_tombstone(
1132 &self,
1133 id: &ObjectId,
1134 metadata: &Metadata,
1135 payload: Bytes,
1136 access_time: Timestamp,
1137 ) -> Result<Option<Tombstone>> {
1138 objectstore_log::debug!("Conditional put to Bigtable backend");
1139
1140 let path = id.as_storage_path().to_string().into_bytes();
1141 let (mutations, size) = object_mutations(&path, metadata.clone(), payload.to_vec())?;
1142
1143 for _ in 0..CAS_RETRY_COUNT {
1144 let write_succeeded = self
1145 .check_and_mutate(
1146 path.clone(),
1147 tombstone_predicate(access_time),
1148 mutations.clone(),
1149 "put_non_tombstone",
1150 )
1151 .await?;
1152
1153 if write_succeeded {
1154 self.change_stream.write(id, size, metadata.time_expires);
1155 return Ok(None);
1156 }
1157
1158 let row = self
1160 .read_row(
1161 &path,
1162 "put_non_tombstone",
1163 access_time,
1164 Some(metadata_filter()),
1165 )
1166 .await?;
1167
1168 match row {
1169 Some(RowData::Tombstone {
1170 target,
1171 time_expires,
1172 ..
1173 }) => {
1174 return Ok(Some(Tombstone {
1175 target: parse_redirect_target(&target, id)?,
1176 time_expires,
1177 }));
1178 }
1179 Some(RowData::Object { .. }) => continue,
1181 None => continue,
1183 }
1184 }
1185
1186 Err(Error::new(
1187 ErrorKind::Internal,
1188 "Bigtable put race exhausted",
1189 ))
1190 }
1191
1192 #[tracing::instrument(level = "debug", skip(self))]
1193 async fn get_tiered_object(
1194 &self,
1195 id: &ObjectId,
1196 access_time: Timestamp,
1197 range: Option<ByteRange>,
1198 ) -> Result<TieredGet> {
1199 objectstore_log::debug!("Reading from Bigtable backend");
1200 let path = id.as_storage_path().to_string().into_bytes();
1201
1202 let Some(row) = self
1203 .read_row(&path, "get_tiered_object", access_time, None)
1204 .await?
1205 else {
1206 return Ok(TieredGet::NotFound);
1207 };
1208
1209 Ok(match row {
1210 RowData::Tombstone {
1211 target,
1212 time_expires,
1213 ..
1214 } => TieredGet::Tombstone(Tombstone {
1215 target: parse_redirect_target(&target, id)?,
1216 time_expires,
1217 }),
1218 RowData::Object {
1219 metadata, payload, ..
1220 } => {
1221 let mut metadata = metadata;
1222 let payload = Bytes::from(payload);
1223 if metadata.size.is_none() {
1224 metadata.size = Some(payload.len());
1226 }
1227
1228 let (content_range, payload) = apply_range(payload, range)?;
1229 TieredGet::Object(metadata, content_range, crate::stream::single(payload))
1230 }
1231 })
1232 }
1233
1234 #[tracing::instrument(level = "debug", skip(self))]
1235 async fn get_tiered_metadata(
1236 &self,
1237 id: &ObjectId,
1238 access_time: Timestamp,
1239 ) -> Result<TieredMetadata> {
1240 objectstore_log::debug!("Reading metadata from Bigtable backend");
1241 let path = id.as_storage_path().to_string().into_bytes();
1242
1243 let row_opt = self
1247 .read_row(
1248 &path,
1249 "get_tiered_metadata",
1250 access_time,
1251 Some(metadata_filter()),
1252 )
1253 .await?;
1254 let Some(row) = row_opt else {
1255 return Ok(TieredMetadata::NotFound);
1256 };
1257
1258 Ok(match row {
1259 RowData::Tombstone {
1260 target,
1261 time_expires,
1262 ..
1263 } => TieredMetadata::Tombstone(Tombstone {
1264 target: parse_redirect_target(&target, id)?,
1265 time_expires,
1266 }),
1267 RowData::Object { metadata, .. } => TieredMetadata::Object(metadata),
1268 })
1269 }
1270
1271 #[tracing::instrument(level = "debug", skip(self))]
1272 async fn compare_and_update(
1273 &self,
1274 id: &ObjectId,
1275 current: Option<&ObjectId>,
1276 update: TieredUpdate,
1277 access_time: Timestamp,
1278 ) -> Result<SetExpiryResponse> {
1279 let TieredUpdate::SetExpiry(expiry_target) = update;
1280 let path = id.as_storage_path().to_string().into_bytes();
1281
1282 let Some(row) = self
1285 .read_row(&path, "set_expiry", access_time, None)
1286 .await?
1287 else {
1288 return Ok(SetExpiryResponse::NotFound);
1289 };
1290
1291 let (expire_at, predicate, mutations): (_, _, Vec<_>) = match row {
1292 RowData::Object {
1293 metadata,
1294 payload,
1295 expiry_micros,
1296 } => {
1297 if current.is_some() {
1298 return Ok(SetExpiryResponse::Rejected); }
1300 let Some(old_expiry) = metadata.time_expires else {
1301 return Ok(SetExpiryResponse::Rejected);
1302 };
1303
1304 if old_expiry < access_time {
1305 return Ok(SetExpiryResponse::NotFound); }
1307 let Some(expire_at) = expiry_target.resolve(metadata.time_created, access_time)?
1308 else {
1309 return Ok(SetExpiryResponse::Rejected);
1310 };
1311 if old_expiry >= expire_at {
1312 return Ok(SetExpiryResponse::Satisfied(expire_at)); }
1314
1315 let predicate = inline_expiry_predicate(expiry_micros, access_time)?;
1319 let mut metadata = metadata;
1320 metadata.expiration_policy = common::extended_expiration_policy(
1321 metadata.expiration_policy,
1322 metadata.time_created,
1323 old_expiry,
1324 expire_at,
1325 )?;
1326 metadata.time_expires = Some(expire_at);
1327 let (mutations, _) = object_mutations(&path, metadata, payload)?;
1328 (expire_at, predicate, mutations.into())
1329 }
1330 RowData::Tombstone {
1331 target,
1332 time_expires,
1333 expiry_micros,
1334 } => {
1335 let Some(expected) = current else {
1336 return Ok(SetExpiryResponse::Rejected); };
1338 let Some(old_expiry) = time_expires else {
1339 return Ok(SetExpiryResponse::Rejected);
1340 };
1341
1342 let redirect_target = parse_redirect_target(&target, id)?;
1343 if old_expiry < access_time {
1344 return Ok(SetExpiryResponse::NotFound);
1345 }
1346 if redirect_target != *expected {
1347 return Ok(SetExpiryResponse::Rejected); }
1349 let Some(expire_at) = expiry_target.resolve(None, access_time)? else {
1350 return Ok(SetExpiryResponse::Rejected);
1351 };
1352 if old_expiry >= expire_at {
1353 return Ok(SetExpiryResponse::Satisfied(expire_at)); }
1355
1356 let predicate =
1357 redirect_expiry_predicate(expected, id, expiry_micros, access_time)?;
1358 let tombstone = Tombstone {
1359 target: redirect_target,
1360 time_expires: Some(expire_at),
1361 };
1362 (expire_at, predicate, tombstone_mutations(&tombstone).into())
1363 }
1364 };
1365
1366 let applied = self
1367 .check_and_mutate(path, predicate, mutations, "set_expiry")
1368 .await?;
1369
1370 if applied {
1371 self.change_stream.update(id, Some(expire_at));
1372 }
1373
1374 Ok(if applied {
1375 SetExpiryResponse::Satisfied(expire_at)
1376 } else {
1377 SetExpiryResponse::Rejected
1378 })
1379 }
1380
1381 #[tracing::instrument(level = "debug", skip(self))]
1382 async fn delete_non_tombstone(
1383 &self,
1384 id: &ObjectId,
1385 access_time: Timestamp,
1386 ) -> Result<Option<Tombstone>> {
1387 objectstore_log::debug!("Conditional delete from Bigtable backend");
1388
1389 let path = id.as_storage_path().to_string().into_bytes();
1390
1391 for _ in 0..CAS_RETRY_COUNT {
1392 let deleted = self
1393 .check_and_mutate(
1394 path.clone(),
1395 non_tombstone_predicate(access_time),
1396 [delete_row_mutation()],
1397 "delete_non_tombstone",
1398 )
1399 .await?;
1400
1401 if deleted {
1402 self.change_stream.delete(id);
1403 return Ok(None);
1404 }
1405
1406 let row = self
1409 .read_row(
1410 &path,
1411 "delete_non_tombstone",
1412 access_time,
1413 Some(metadata_filter()),
1414 )
1415 .await?;
1416
1417 match row {
1418 Some(RowData::Tombstone {
1419 target,
1420 time_expires,
1421 ..
1422 }) => {
1423 return Ok(Some(Tombstone {
1424 target: parse_redirect_target(&target, id)?,
1425 time_expires,
1426 }));
1427 }
1428 Some(RowData::Object { .. }) => continue,
1430 None => return Ok(None),
1432 }
1433 }
1434
1435 Err(Error::new(
1436 ErrorKind::Internal,
1437 "Bigtable delete race exhausted",
1438 ))
1439 }
1440
1441 #[tracing::instrument(level = "debug", skip(self, write))]
1442 async fn compare_and_write(
1443 &self,
1444 id: &ObjectId,
1445 current: Option<&ObjectId>,
1446 write: TieredWrite,
1447 access_time: Timestamp,
1448 ) -> Result<bool> {
1449 objectstore_log::debug!("CAS put to Bigtable backend");
1450
1451 let path = id.as_storage_path().to_string().into_bytes();
1452 let predicate = match (current, write.target()) {
1453 (Some(old), Some(new)) => update_predicate(old, new, id, access_time),
1454 (Some(target), None) => optional_target_predicate(target, id, access_time),
1455 (None, Some(target)) => optional_target_predicate(target, id, access_time),
1456 (None, None) => tombstone_predicate(access_time),
1457 };
1458
1459 let (mutations, expires_at): (Vec<v2::Mutation>, Option<Option<Timestamp>>) = match write {
1463 TieredWrite::Tombstone(tombstone) => (
1464 tombstone_mutations(&tombstone).into(),
1465 Some(tombstone.time_expires),
1466 ),
1467 TieredWrite::Object(m, p) => {
1468 let expires_at = m.time_expires;
1469 let (mutations, _) = object_mutations(&path, m, p.to_vec())?;
1470 (mutations.into(), Some(expires_at))
1471 }
1472 TieredWrite::Delete => (vec![delete_row_mutation()], None),
1473 };
1474
1475 let written = self
1476 .check_and_mutate(
1477 path.clone(),
1478 predicate,
1479 mutations.clone(),
1480 "compare_and_write",
1481 )
1482 .await?;
1483
1484 match (written, expires_at) {
1485 (false, _) => {}
1487 (true, Some(expires_at)) => {
1489 self.change_stream
1490 .write(id, row_size(&path, &mutations), expires_at)
1491 }
1492 (true, None) => self.change_stream.delete(id),
1494 }
1495
1496 Ok(written)
1497 }
1498}
1499
1500async fn retry<T, F>(context: &'static str, f: impl Fn() -> F) -> Result<T>
1502where
1503 F: Future<Output = Result<T, BigTableError>> + Send,
1504{
1505 let mut retry_count = 0usize;
1506
1507 loop {
1508 let attempt_span = tracing::debug_span!(
1509 "bigtable.request",
1510 action = context,
1511 grpc.status = tracing::field::Empty,
1512 );
1513 let attempt = async {
1514 let result = f().await;
1515 let span = tracing::Span::current();
1516 match &result {
1517 Ok(_) => span.record("grpc.status", "ok"),
1518 Err(BigTableError::RpcError(status)) => {
1519 span.record("grpc.status", tracing::field::debug(status.code()))
1520 }
1521 Err(_) => &span,
1523 };
1524 result
1525 };
1526
1527 match attempt.instrument(attempt_span).await {
1528 Ok(res) => return Ok(res),
1529 Err(e) if retry_count >= REQUEST_RETRY_COUNT || !is_retryable(&e) => {
1530 objectstore_metrics::count!("bigtable.failures", action = context);
1531 return Err(e).context(
1532 ErrorKind::BackendFailure,
1533 format!("running Bigtable {context}"),
1534 );
1535 }
1536 Err(e) => {
1537 retry_count += 1;
1538 objectstore_metrics::count!("bigtable.retries", action = context);
1539 objectstore_log::warn!(!!&e, retry_count, context, "Retrying request");
1540 }
1541 }
1542 }
1543}
1544
1545fn is_retryable(error: &BigTableError) -> bool {
1546 match error {
1547 BigTableError::GCPAuthError(_) => true,
1549 BigTableError::TransportError(_) => true,
1551 BigTableError::IoError(_) => true,
1553 BigTableError::TimeoutError(_) => true,
1554
1555 BigTableError::RpcError(status) => match status.code() {
1557 Code::Unavailable => true,
1559 Code::Cancelled => true,
1561 Code::DeadlineExceeded => true,
1562 Code::Unauthenticated => true,
1564 Code::Aborted => true,
1566 Code::Internal => true,
1567 Code::FailedPrecondition => true,
1568 Code::Unknown => true,
1569 _ => false,
1570 },
1571 _ => false,
1572 }
1573}
1574
1575fn apply_range(payload: Bytes, range: Option<ByteRange>) -> Result<(Option<ContentRange>, Bytes)> {
1581 let Some(byte_range) = range else {
1582 return Ok((None, payload));
1583 };
1584
1585 let total = payload.len() as u64;
1586 let content_range = byte_range
1587 .resolve(total)
1588 .ok_or(ErrorKind::RangeNotSatisfiable { total })?;
1589
1590 let sliced = payload.slice(content_range.start as usize..content_range.end as usize + 1);
1591 Ok((Some(content_range), sliced))
1592}
1593
1594#[cfg(test)]
1595mod tests {
1596 use std::collections::BTreeMap;
1597
1598 use anyhow::Result;
1599 #[cfg(feature = "storage-cogs")]
1600 use objectstore_inventory_tracker::{OpType, test_utils::DummyProducer};
1601 use objectstore_types::metadata::ExpirationPolicy;
1602 use objectstore_types::scope::{Scope, Scopes};
1603
1604 use super::*;
1605 use crate::backend::common::ExpiryTarget;
1606 use crate::id::ObjectContext;
1607 use crate::stream;
1608
1609 fn test_config() -> BigTableConfig {
1615 BigTableConfig {
1616 endpoint: Some("localhost:8086".into()),
1617 project_id: "testing".into(),
1618 instance_name: "objectstore".into(),
1619 table_name: "objectstore".into(),
1620 connections: None,
1621 rpc_timeout: default_rpc_timeout(),
1622 cogs: None,
1623 }
1624 }
1625
1626 async fn create_test_backend() -> Result<BigTableBackend> {
1627 BigTableBackend::new(test_config(), &ChangeStreamFactory::default()).await
1628 }
1629
1630 #[cfg(feature = "storage-cogs")]
1631 async fn create_test_backend_with_change_stream() -> Result<(BigTableBackend, DummyProducer)> {
1632 let (streams, producer) = crate::change_stream::dummy_factory();
1633 let config = BigTableConfig {
1634 cogs: Some(CostTrackerStreamConfig {
1635 shared_resource_id: "bigtable_objectstore".into(),
1636 sample_rate: 1.0,
1637 }),
1638 ..test_config()
1639 };
1640
1641 Ok((BigTableBackend::new(config, &streams).await?, producer))
1642 }
1643
1644 fn make_id() -> ObjectId {
1645 ObjectId::random(ObjectContext {
1646 usecase: "testing".into(),
1647 scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1648 })
1649 }
1650
1651 async fn create_object(
1652 backend: &BigTableBackend,
1653 id: &ObjectId,
1654 metadata: &Metadata,
1655 payload: &[u8],
1656 now: Timestamp,
1657 ) -> Result<()> {
1658 let path = id.as_storage_path().to_string().into_bytes();
1659 let mut metadata = metadata.clone();
1662 if metadata.time_expires.is_none() {
1663 metadata.time_expires = metadata.expiration_policy.expires_in().map(|ttl| now + ttl);
1664 }
1665 let (mutations, _) = object_mutations(&path, metadata, payload.to_vec())?;
1666 backend.mutate(path, mutations, "test-setup").await?;
1667 Ok(())
1668 }
1669
1670 async fn create_tombstone(
1671 backend: &BigTableBackend,
1672 id: &ObjectId,
1673 tombstone: &Tombstone,
1674 ) -> Result<()> {
1675 let path = id.as_storage_path().to_string().into_bytes();
1676 let mutations = tombstone_mutations(tombstone);
1677 backend.mutate(path, mutations, "test-setup").await?;
1678 Ok(())
1679 }
1680
1681 async fn write_legacy_tombstone(
1683 backend: &BigTableBackend,
1684 id: &ObjectId,
1685 expiration_policy: ExpirationPolicy,
1686 time_expires: Option<Timestamp>,
1687 ) -> Result<()> {
1688 let meta = if expiration_policy.is_manual() {
1689 r#"{"is_redirect_tombstone":true}"#.to_owned()
1690 } else {
1691 let policy_json = serde_json::to_string(&expiration_policy).unwrap();
1692 format!(r#"{{"is_redirect_tombstone":true,"expiration_policy":{policy_json}}}"#)
1693 };
1694
1695 let (family, timestamp_micros) = if expiration_policy.is_manual() {
1696 (FAMILY_MANUAL, -1)
1697 } else {
1698 let t =
1699 time_expires.unwrap_or(Timestamp::now() + expiration_policy.expires_in().unwrap());
1700 (FAMILY_GC, t.as_micros() as i64)
1701 };
1702
1703 let path = id.as_storage_path().to_string().into_bytes();
1704 let mutations = [mutation(mutation::Mutation::SetCell(mutation::SetCell {
1705 family_name: family.to_owned(),
1706 column_qualifier: COLUMN_METADATA.to_owned(),
1707 timestamp_micros,
1708 value: meta.into_bytes(),
1709 }))];
1710
1711 backend.mutate(path, mutations, "test-setup").await?;
1712
1713 Ok(())
1714 }
1715
1716 async fn write_empty_redirect_tombstone(
1718 backend: &BigTableBackend,
1719 id: &ObjectId,
1720 ) -> Result<()> {
1721 let path = id.as_storage_path().to_string().into_bytes();
1722 let mutations = [
1723 mutation(mutation::Mutation::SetCell(mutation::SetCell {
1724 family_name: FAMILY_MANUAL.to_owned(),
1725 column_qualifier: COLUMN_REDIRECT.to_owned(),
1726 timestamp_micros: -1,
1727 value: b"".to_vec(), })),
1729 mutation(mutation::Mutation::SetCell(mutation::SetCell {
1730 family_name: FAMILY_MANUAL.to_owned(),
1731 column_qualifier: b"t".to_vec(),
1732 timestamp_micros: -1,
1733 value: b"{}".to_vec(),
1734 })),
1735 ];
1736
1737 backend.mutate(path, mutations, "test-setup").await?;
1738
1739 Ok(())
1740 }
1741
1742 #[tokio::test]
1746 async fn test_roundtrip() -> Result<()> {
1747 let backend = create_test_backend().await?;
1748
1749 let id = make_id();
1750 let metadata = Metadata {
1751 content_type: "text/plain".into(),
1752 time_created: Some(Timestamp::now()),
1753 custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
1754 ..Default::default()
1755 };
1756
1757 backend
1758 .put_object(
1759 &id,
1760 &metadata,
1761 stream::single("hello, world"),
1762 Timestamp::now(),
1763 )
1764 .await?;
1765
1766 let (obj_meta, _, stream) = backend
1767 .get_object(&id, Timestamp::now(), None)
1768 .await?
1769 .unwrap();
1770 let payload = stream::read_to_vec(stream).await?;
1771 assert_eq!(payload, b"hello, world");
1772 assert_eq!(obj_meta.content_type, metadata.content_type);
1773 assert_eq!(obj_meta.custom, metadata.custom);
1774
1775 let head_meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
1776 assert_eq!(head_meta.content_type, metadata.content_type);
1777 assert_eq!(head_meta.custom, metadata.custom);
1778
1779 Ok(())
1780 }
1781
1782 #[tokio::test]
1784 async fn test_time_expires_roundtrip() -> Result<()> {
1785 let backend = create_test_backend().await?;
1786
1787 let id = make_id();
1788 let ttl = Duration::from_hours(2 * 24);
1789 let expires = Timestamp::now() + ttl;
1790 let metadata = Metadata {
1791 expiration_policy: ExpirationPolicy::TimeToLive(ttl),
1792 time_expires: Some(expires),
1793 ..Default::default()
1794 };
1795 create_object(&backend, &id, &metadata, b"data", Timestamp::now()).await?;
1796
1797 let meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
1798 assert_eq!(meta.time_expires, Some(expires));
1799
1800 Ok(())
1801 }
1802
1803 #[tokio::test]
1805 async fn test_nonexistent() -> Result<()> {
1806 let backend = create_test_backend().await?;
1807
1808 let id = make_id();
1809 assert!(
1810 backend
1811 .get_object(&id, Timestamp::now(), None)
1812 .await?
1813 .is_none()
1814 );
1815 assert!(backend.get_metadata(&id, Timestamp::now()).await?.is_none());
1816 backend.delete_object(&id, Timestamp::now()).await?;
1817
1818 Ok(())
1819 }
1820
1821 #[tokio::test]
1822 async fn test_overwrite() -> Result<()> {
1823 let backend = create_test_backend().await?;
1824
1825 let id = make_id();
1826 let first_metadata = Metadata {
1827 custom: BTreeMap::from_iter([("invalid".into(), "invalid".into())]),
1828 ..Default::default()
1829 };
1830 create_object(&backend, &id, &first_metadata, b"hello", Timestamp::now()).await?;
1831
1832 let second_metadata = Metadata {
1833 custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
1834 ..Default::default()
1835 };
1836 backend
1837 .put_object(
1838 &id,
1839 &second_metadata,
1840 stream::single("world"),
1841 Timestamp::now(),
1842 )
1843 .await?;
1844
1845 let (meta, _, stream) = backend
1846 .get_object(&id, Timestamp::now(), None)
1847 .await?
1848 .unwrap();
1849 let payload = stream::read_to_vec(stream).await?;
1850 assert_eq!(payload, b"world");
1851 assert_eq!(meta.custom, second_metadata.custom);
1852
1853 Ok(())
1854 }
1855
1856 #[tokio::test]
1857 async fn test_read_after_delete() -> Result<()> {
1858 let backend = create_test_backend().await?;
1859
1860 let id = make_id();
1861 let metadata = Metadata::default();
1862 create_object(&backend, &id, &metadata, b"hello", Timestamp::now()).await?;
1863 backend.delete_object(&id, Timestamp::now()).await?;
1864
1865 assert!(
1866 backend
1867 .get_object(&id, Timestamp::now(), None)
1868 .await?
1869 .is_none()
1870 );
1871
1872 Ok(())
1873 }
1874
1875 #[tokio::test]
1877 async fn test_set_expiry() -> Result<()> {
1878 for is_ttl in [false, true] {
1879 let backend = create_test_backend().await?;
1880 let tti = Duration::from_hours(2 * 24);
1881 let mut metadata = Metadata {
1882 expiration_policy: if is_ttl {
1883 ExpirationPolicy::TimeToLive(tti)
1884 } else {
1885 ExpirationPolicy::TimeToIdle(tti)
1886 },
1887 ..Default::default()
1888 };
1889
1890 let past_now = Timestamp::now() - tti + Duration::from_mins(1);
1892
1893 let id = make_id();
1894 metadata.time_created = Some(past_now);
1895 metadata.time_expires = Some(past_now + tti);
1896 let path = id.as_storage_path().to_string().into_bytes();
1897 let (mutations, _) = object_mutations(&path, metadata, b"hello, world".to_vec())?;
1898 let mutations = mutations.map(|mut mutation| {
1900 if let Some(mutation::Mutation::SetCell(cell)) = &mut mutation.mutation {
1901 cell.timestamp_micros -= 500_000;
1902 }
1903 mutation
1904 });
1905 backend.mutate(path, mutations, "test-setup").await?;
1906
1907 let (observed, _, _) = backend
1908 .get_object(&id, Timestamp::now(), None)
1909 .await?
1910 .unwrap();
1911 let observed_expiry = observed.time_expires.unwrap();
1912 assert_eq!(
1913 backend
1914 .get_metadata(&id, Timestamp::now())
1915 .await?
1916 .unwrap()
1917 .time_expires,
1918 Some(observed_expiry),
1919 "backend reads must not renew TTI"
1920 );
1921
1922 let requested = past_now + tti + tti;
1923 for target in [
1924 ExpiryTarget::At(requested),
1925 ExpiryTarget::FromCreation(tti + tti),
1926 ] {
1927 assert_eq!(
1928 backend
1929 .set_expiry(&id, target.into(), Timestamp::now())
1930 .await?,
1931 SetExpiryResponse::Satisfied(requested)
1932 );
1933 }
1934 assert_eq!(
1935 backend
1936 .get_metadata(&id, Timestamp::now())
1937 .await?
1938 .unwrap()
1939 .time_expires,
1940 Some(requested)
1941 );
1942 let (updated, _, stream) = backend
1943 .get_object(&id, Timestamp::now(), None)
1944 .await?
1945 .unwrap();
1946 assert_eq!(
1947 updated.expiration_policy,
1948 if is_ttl {
1949 ExpirationPolicy::TimeToLive(tti + tti)
1950 } else {
1951 ExpirationPolicy::TimeToIdle(tti)
1952 }
1953 );
1954 let payload = stream::read_to_vec(stream).await?;
1955 assert_eq!(payload, b"hello, world");
1956 }
1957 Ok(())
1958 }
1959
1960 #[tokio::test]
1961 async fn test_expiry_outcomes() -> Result<()> {
1962 let backend = create_test_backend().await?;
1963 let access_time = Timestamp::now();
1964 let deadline = access_time + Duration::from_hours(1);
1965 for (expiry, expected) in [
1966 (None, SetExpiryResponse::Rejected),
1967 (
1968 Some(access_time - Duration::from_secs(1)),
1969 SetExpiryResponse::NotFound,
1970 ),
1971 (Some(deadline), SetExpiryResponse::Rejected),
1972 ] {
1973 let id = make_id();
1974 let metadata = Metadata {
1975 time_expires: expiry,
1976 ..Default::default()
1977 };
1978 create_object(&backend, &id, &metadata, b"payload", access_time).await?;
1979 assert_eq!(
1980 backend
1981 .set_expiry(
1982 &id,
1983 ExpiryTarget::FromCreation(Duration::from_hours(2)).into(),
1984 access_time
1985 )
1986 .await?,
1987 expected
1988 );
1989 }
1990 Ok(())
1991 }
1992
1993 #[tokio::test]
1994 async fn test_expiry_conflict() -> Result<()> {
1995 let backend = create_test_backend().await?;
1996 let missing = make_id();
1997 assert_eq!(
1998 backend
1999 .set_expiry(
2000 &missing,
2001 ExpiryTarget::At(Timestamp::now() + Duration::from_hours(2)).into(),
2002 Timestamp::now()
2003 )
2004 .await?,
2005 SetExpiryResponse::NotFound
2006 );
2007
2008 let id = make_id();
2009 let observed_expiry = Timestamp::now() + Duration::from_hours(1);
2010 let original = Metadata {
2011 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2012 time_expires: Some(observed_expiry),
2013 ..Default::default()
2014 };
2015 create_object(&backend, &id, &original, b"original", Timestamp::now()).await?;
2016
2017 let path = id.as_storage_path().to_string().into_bytes();
2018 let mut extended = original.clone();
2019 extended.time_expires = Some(observed_expiry + Duration::from_hours(1));
2020 let (extension, _) = object_mutations(&path, extended, b"original".to_vec())?;
2021 let predicate =
2022 inline_expiry_predicate(observed_expiry.as_micros() as i64, Timestamp::now())?;
2023
2024 let mut replacement = original.clone();
2025 replacement.time_expires = Some(observed_expiry + Duration::from_mins(1));
2026 create_object(
2027 &backend,
2028 &id,
2029 &replacement,
2030 b"replacement",
2031 Timestamp::now(),
2032 )
2033 .await?;
2034 assert!(
2035 !backend
2036 .check_and_mutate(path, predicate, extension, "test-expiry-conflict")
2037 .await?
2038 );
2039 let (_, _, payload) = backend
2040 .get_object(&id, Timestamp::now(), None)
2041 .await?
2042 .unwrap();
2043 assert_eq!(stream::read_to_vec(payload).await?, b"replacement");
2044 Ok(())
2045 }
2046
2047 #[tokio::test]
2048 async fn test_redirect_expiry() -> Result<()> {
2049 let backend = create_test_backend().await?;
2050 let id = make_id();
2051 let target = ObjectId::random(id.context().clone());
2052 let wrong_target = ObjectId::random(id.context().clone());
2053 let old_expiry = Timestamp::now() + Duration::from_hours(1);
2054 let path = id.as_storage_path().to_string().into_bytes();
2055 let mutations = tombstone_mutations(&Tombstone {
2056 target: target.clone(),
2057 time_expires: Some(old_expiry),
2058 })
2059 .map(|mut mutation| {
2060 if let Some(mutation::Mutation::SetCell(cell)) = &mut mutation.mutation {
2061 cell.timestamp_micros -= 500_000;
2062 }
2063 mutation
2064 });
2065 backend.mutate(path, mutations, "test-setup").await?;
2066
2067 let later = old_expiry + Duration::from_hours(2);
2068 assert_eq!(
2069 backend
2070 .compare_and_update(
2071 &id,
2072 Some(&wrong_target),
2073 TieredUpdate::SetExpiry(ExpiryTarget::At(later).into()),
2074 Timestamp::now(),
2075 )
2076 .await?,
2077 SetExpiryResponse::Rejected
2078 );
2079 assert_eq!(
2080 backend
2081 .compare_and_update(
2082 &id,
2083 Some(&target),
2084 TieredUpdate::SetExpiry(ExpiryTarget::At(later).into()),
2085 Timestamp::now()
2086 )
2087 .await?,
2088 SetExpiryResponse::Satisfied(later)
2089 );
2090 let requested = old_expiry + Duration::from_mins(30);
2091 assert_eq!(
2092 backend
2093 .compare_and_update(
2094 &id,
2095 Some(&target),
2096 TieredUpdate::SetExpiry(ExpiryTarget::At(requested).into()),
2097 Timestamp::now(),
2098 )
2099 .await?,
2100 SetExpiryResponse::Satisfied(requested)
2101 );
2102 let TieredMetadata::Tombstone(tombstone) =
2103 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2104 else {
2105 panic!("expected tombstone");
2106 };
2107 assert_eq!(tombstone.time_expires, Some(later));
2108 Ok(())
2109 }
2110
2111 #[tokio::test]
2114 async fn test_ttl_immediate() -> Result<()> {
2115 let backend = create_test_backend().await?;
2119
2120 let id = make_id();
2121 let metadata = Metadata {
2122 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(0)),
2123 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2124 ..Default::default()
2125 };
2126 create_object(&backend, &id, &metadata, b"hello, world", Timestamp::now()).await?;
2127
2128 assert!(
2129 backend
2130 .get_object(&id, Timestamp::now(), None)
2131 .await?
2132 .is_none()
2133 );
2134
2135 Ok(())
2136 }
2137
2138 #[tokio::test]
2139 async fn test_tti_immediate() -> Result<()> {
2140 let backend = create_test_backend().await?;
2144
2145 let id = make_id();
2146 let metadata = Metadata {
2147 expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(0)),
2148 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2149 ..Default::default()
2150 };
2151 create_object(&backend, &id, &metadata, b"hello, world", Timestamp::now()).await?;
2152
2153 assert!(
2154 backend
2155 .get_object(&id, Timestamp::now(), None)
2156 .await?
2157 .is_none()
2158 );
2159
2160 Ok(())
2161 }
2162
2163 #[tokio::test]
2172 async fn test_tiered_get() -> Result<()> {
2173 let backend = create_test_backend().await?;
2174
2175 let id = make_id();
2177 assert!(matches!(
2178 backend
2179 .get_tiered_object(&id, Timestamp::now(), None)
2180 .await?,
2181 TieredGet::NotFound
2182 ));
2183 assert!(matches!(
2184 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2185 TieredMetadata::NotFound
2186 ));
2187
2188 let id = make_id();
2190 let put_meta = Metadata {
2191 content_type: "text/plain".into(),
2192 custom: BTreeMap::from_iter([("k".into(), "v".into())]),
2193 ..Default::default()
2194 };
2195 create_object(&backend, &id, &put_meta, b"payload", Timestamp::now()).await?;
2196
2197 let TieredGet::Object(obj_meta, _, obj_stream) = backend
2198 .get_tiered_object(&id, Timestamp::now(), None)
2199 .await?
2200 else {
2201 panic!("expected TieredGet::Object");
2202 };
2203 let obj_payload = stream::read_to_vec(obj_stream).await?;
2204 assert_eq!(obj_payload, b"payload");
2205 assert_eq!(obj_meta.content_type, put_meta.content_type);
2206 assert_eq!(obj_meta.custom, put_meta.custom);
2207
2208 let TieredMetadata::Object(head_meta) =
2209 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2210 else {
2211 panic!("expected TieredMetadata::Object");
2212 };
2213 assert_eq!(head_meta.content_type, put_meta.content_type);
2214 assert_eq!(head_meta.custom, put_meta.custom);
2215
2216 let hv_id = make_id();
2218 let lt_id = ObjectId::random(hv_id.context().clone());
2219 let tombstone = Tombstone {
2220 target: lt_id.clone(),
2221 time_expires: None,
2222 };
2223 create_tombstone(&backend, &hv_id, &tombstone).await?;
2224
2225 match backend
2226 .get_tiered_object(&hv_id, Timestamp::now(), None)
2227 .await?
2228 {
2229 TieredGet::Tombstone(get_t) => assert_eq!(get_t.target, lt_id),
2230 other => panic!("expected TieredGet::Tombstone, got {other:?}"),
2231 }
2232 match backend
2233 .get_tiered_metadata(&hv_id, Timestamp::now())
2234 .await?
2235 {
2236 TieredMetadata::Tombstone(meta_t) => assert_eq!(meta_t.target, lt_id,),
2237 other => panic!("expected TieredMetadata::Tombstone, got {other:?}"),
2238 }
2239
2240 Ok(())
2241 }
2242
2243 #[tokio::test]
2249 async fn test_put_non_tombstone() -> Result<()> {
2250 let backend = create_test_backend().await?;
2251
2252 let id = make_id();
2254 let metadata = Metadata::default();
2255 let result = backend
2256 .put_non_tombstone(
2257 &id,
2258 &metadata,
2259 Bytes::from_static(b"first"),
2260 Timestamp::now(),
2261 )
2262 .await?;
2263 assert_eq!(result, None, "expected None on empty row");
2264 let (_, _, stream) = backend
2265 .get_object(&id, Timestamp::now(), None)
2266 .await?
2267 .unwrap();
2268 assert_eq!(&stream::read_to_vec(stream).await?, b"first");
2269
2270 let id = make_id();
2272 create_object(&backend, &id, &metadata, b"old", Timestamp::now()).await?;
2273 let result = backend
2274 .put_non_tombstone(&id, &metadata, Bytes::from_static(b"new"), Timestamp::now())
2275 .await?;
2276 assert_eq!(result, None, "expected None when overwriting object");
2277 let (_, _, stream) = backend
2278 .get_object(&id, Timestamp::now(), None)
2279 .await?
2280 .unwrap();
2281 assert_eq!(&stream::read_to_vec(stream).await?, b"new");
2282
2283 let hv_id = make_id();
2285 let lt_id = ObjectId::random(hv_id.context().clone());
2286 let tombstone = Tombstone {
2287 target: lt_id.clone(),
2288 time_expires: None,
2289 };
2290 create_tombstone(&backend, &hv_id, &tombstone).await?;
2291 let result = backend
2292 .put_non_tombstone(&hv_id, &metadata, Bytes::new(), Timestamp::now())
2293 .await?;
2294 let returned = result.expect("expected Some(Tombstone) when row is a tombstone");
2295 assert_eq!(returned.target, lt_id);
2296 assert!(
2297 matches!(
2298 backend
2299 .get_tiered_metadata(&hv_id, Timestamp::now())
2300 .await?,
2301 TieredMetadata::Tombstone(_)
2302 ),
2303 "tombstone must still exist after put_non_tombstone"
2304 );
2305
2306 Ok(())
2307 }
2308
2309 #[tokio::test]
2318 async fn test_delete_non_tombstone() -> Result<()> {
2319 let backend = create_test_backend().await?;
2320
2321 let id = make_id();
2323 assert_eq!(
2324 backend.delete_non_tombstone(&id, Timestamp::now()).await?,
2325 None
2326 );
2327
2328 let id = make_id();
2330 let metadata = Metadata::default();
2331 create_object(&backend, &id, &metadata, b"hello, world", Timestamp::now()).await?;
2332 assert_eq!(
2333 backend.delete_non_tombstone(&id, Timestamp::now()).await?,
2334 None
2335 );
2336 assert!(
2337 backend
2338 .get_object(&id, Timestamp::now(), None)
2339 .await?
2340 .is_none()
2341 );
2342
2343 let id = make_id();
2345 let tombstone = Tombstone {
2346 target: id.clone(),
2347 time_expires: None,
2348 };
2349 create_tombstone(&backend, &id, &tombstone).await?;
2350 let tombstone = backend
2351 .delete_non_tombstone(&id, Timestamp::now())
2352 .await?
2353 .expect("expected Some(tombstone)");
2354 assert_eq!(tombstone.target, id, "tombstone target must be returned");
2355 assert!(
2356 matches!(
2357 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2358 TieredMetadata::Tombstone(_)
2359 ),
2360 "tombstone must still exist after delete_non_tombstone"
2361 );
2362
2363 Ok(())
2364 }
2365
2366 #[tokio::test]
2372 async fn test_cas_create_tombstone() -> Result<()> {
2373 let backend = create_test_backend().await?;
2374
2375 let hv_id = make_id();
2376 let lt_id = ObjectId::random(hv_id.context().clone());
2377 let time_expires = Some(Timestamp::now() + Duration::from_hours(1));
2378 let tombstone = Tombstone {
2379 target: lt_id.clone(),
2380 time_expires,
2381 };
2382
2383 let committed = backend
2385 .compare_and_write(
2386 &hv_id,
2387 None,
2388 TieredWrite::Tombstone(tombstone.clone()),
2389 Timestamp::now(),
2390 )
2391 .await?;
2392 assert!(committed, "expected CAS success on empty row");
2393
2394 let TieredMetadata::Tombstone(t) = backend
2396 .get_tiered_metadata(&hv_id, Timestamp::now())
2397 .await?
2398 else {
2399 panic!("expected TieredMetadata::Tombstone");
2400 };
2401 assert_eq!(t.target, lt_id, "target must round-trip via r column");
2402 assert_eq!(t.time_expires, time_expires);
2403 match backend
2404 .get_tiered_object(&hv_id, Timestamp::now(), None)
2405 .await?
2406 {
2407 TieredGet::Tombstone(t) => assert_eq!(t.target, lt_id, "round-trip via r column"),
2408 other => panic!("expected TieredGet::Tombstone, got {other:?}"),
2409 }
2410
2411 assert!(
2413 backend
2414 .get_object(&hv_id, Timestamp::now(), None)
2415 .await
2416 .is_err_and(|error| error.kind() == ErrorKind::UnexpectedTombstone)
2417 );
2418 assert!(
2419 backend
2420 .get_metadata(&hv_id, Timestamp::now())
2421 .await
2422 .is_err_and(|error| error.kind() == ErrorKind::UnexpectedTombstone)
2423 );
2424
2425 let second = backend
2427 .compare_and_write(
2428 &hv_id,
2429 None,
2430 TieredWrite::Tombstone(tombstone),
2431 Timestamp::now(),
2432 )
2433 .await?;
2434 assert!(second, "idempotent retry");
2435
2436 Ok(())
2437 }
2438
2439 #[tokio::test]
2441 async fn test_cas_swap_tombstone() -> Result<()> {
2442 let backend = create_test_backend().await?;
2443
2444 let hv_id = make_id();
2445 let old_lt_id = ObjectId::random(hv_id.context().clone());
2446 let wrong_lt_id = ObjectId::random(hv_id.context().clone());
2447 let new_lt_id = ObjectId::random(hv_id.context().clone());
2448
2449 let tombstone = Tombstone {
2450 target: old_lt_id.clone(),
2451 time_expires: None,
2452 };
2453 create_tombstone(&backend, &hv_id, &tombstone).await?;
2454
2455 let write = TieredWrite::Tombstone(Tombstone {
2457 target: new_lt_id.clone(),
2458 time_expires: None,
2459 });
2460 let swapped = backend
2461 .compare_and_write(&hv_id, Some(&wrong_lt_id), write.clone(), Timestamp::now())
2462 .await?;
2463 assert!(!swapped, "expected CAS failure due to wrong target");
2464 match backend
2465 .get_tiered_metadata(&hv_id, Timestamp::now())
2466 .await?
2467 {
2468 TieredMetadata::Tombstone(t) => assert_eq!(t.target, old_lt_id),
2469 other => panic!("expected tombstone, got {other:?}"),
2470 }
2471
2472 let swapped = backend
2474 .compare_and_write(&hv_id, Some(&old_lt_id), write.clone(), Timestamp::now())
2475 .await?;
2476 assert!(swapped, "expected CAS success with correct target");
2477 match backend
2478 .get_tiered_metadata(&hv_id, Timestamp::now())
2479 .await?
2480 {
2481 TieredMetadata::Tombstone(t) => assert_eq!(t.target, new_lt_id),
2482 other => panic!("expected tombstone, got {other:?}"),
2483 }
2484
2485 let retry = backend
2487 .compare_and_write(&hv_id, Some(&old_lt_id), write, Timestamp::now())
2488 .await?;
2489 assert!(retry, "idempotent retry");
2490
2491 Ok(())
2492 }
2493
2494 #[tokio::test]
2496 async fn test_cas_swap_inline() -> Result<()> {
2497 let backend = create_test_backend().await?;
2498
2499 let id = make_id();
2500 let lt_id = ObjectId::random(id.context().clone());
2501 let wrong_id = ObjectId::random(id.context().clone());
2502
2503 let tombstone = Tombstone {
2504 target: lt_id.clone(),
2505 time_expires: None,
2506 };
2507 create_tombstone(&backend, &id, &tombstone).await?;
2508
2509 let write = TieredWrite::Object(Metadata::default(), Bytes::new());
2511 let swapped = backend
2512 .compare_and_write(&id, Some(&wrong_id), write, Timestamp::now())
2513 .await?;
2514 assert!(!swapped, "expected CAS failure with wrong target");
2515 assert!(matches!(
2516 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2517 TieredMetadata::Tombstone(_)
2518 ));
2519
2520 let payload = Bytes::from_static(b"hello inline");
2522 let write = TieredWrite::Object(Metadata::default(), payload.clone());
2523 let swapped = backend
2524 .compare_and_write(&id, Some(<_id), write.clone(), Timestamp::now())
2525 .await?;
2526 assert!(swapped, "expected CAS success with correct target");
2527 let TieredGet::Object(_, _, stream) = backend
2528 .get_tiered_object(&id, Timestamp::now(), None)
2529 .await?
2530 else {
2531 panic!("expected inline object after swap");
2532 };
2533 assert_eq!(&stream::read_to_vec(stream).await?, payload.as_ref());
2534
2535 let retry = backend
2537 .compare_and_write(&id, Some(<_id), write, Timestamp::now())
2538 .await?;
2539 assert!(retry, "idempotent retry");
2540
2541 Ok(())
2542 }
2543
2544 #[tokio::test]
2546 async fn test_cas_create_object_on_empty_row() -> Result<()> {
2547 let backend = create_test_backend().await?;
2548
2549 let id = make_id();
2550 let payload = Bytes::from_static(b"cas object");
2551 let write = TieredWrite::Object(Metadata::default(), payload.clone());
2552 let committed = backend
2553 .compare_and_write(&id, None, write, Timestamp::now())
2554 .await?;
2555 assert!(committed, "expected CAS success on empty row");
2556
2557 let TieredGet::Object(_, _, stream) = backend
2558 .get_tiered_object(&id, Timestamp::now(), None)
2559 .await?
2560 else {
2561 panic!("expected Object after CAS-create");
2562 };
2563 assert_eq!(&stream::read_to_vec(stream).await?, payload.as_ref());
2564
2565 Ok(())
2566 }
2567
2568 #[tokio::test]
2570 async fn test_cas_delete() -> Result<()> {
2571 let backend = create_test_backend().await?;
2572
2573 let id = make_id();
2574 let lt_id = ObjectId::random(id.context().clone());
2575 let wrong_id = ObjectId::random(id.context().clone());
2576
2577 let tombstone = Tombstone {
2578 target: lt_id.clone(),
2579 time_expires: None,
2580 };
2581 create_tombstone(&backend, &id, &tombstone).await?;
2582
2583 let deleted = backend
2585 .compare_and_write(&id, Some(&wrong_id), TieredWrite::Delete, Timestamp::now())
2586 .await?;
2587 assert!(!deleted, "expected CAS failure with wrong target");
2588 assert!(matches!(
2589 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2590 TieredMetadata::Tombstone(_)
2591 ));
2592
2593 let deleted = backend
2595 .compare_and_write(&id, Some(<_id), TieredWrite::Delete, Timestamp::now())
2596 .await?;
2597 assert!(deleted, "expected CAS delete success");
2598 assert!(matches!(
2599 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2600 TieredMetadata::NotFound
2601 ));
2602
2603 let retry = backend
2605 .compare_and_write(&id, Some(<_id), TieredWrite::Delete, Timestamp::now())
2606 .await?;
2607 assert!(retry, "idempotent retry");
2608
2609 let id2 = make_id();
2611 let fake_lt_id = ObjectId::random(id2.context().clone());
2612 let metadata = Metadata::default();
2613 create_object(&backend, &id2, &metadata, b"data", Timestamp::now()).await?;
2614 let deleted = backend
2615 .compare_and_write(
2616 &id2,
2617 Some(&fake_lt_id),
2618 TieredWrite::Delete,
2619 Timestamp::now(),
2620 )
2621 .await?;
2622 assert!(deleted, "expected idempotent deletion");
2623
2624 Ok(())
2625 }
2626
2627 #[tokio::test]
2634 async fn test_legacy_tombstone_reads() -> Result<()> {
2635 let backend = create_test_backend().await?;
2636
2637 let id = make_id();
2639 write_legacy_tombstone(&backend, &id, ExpirationPolicy::Manual, None).await?;
2640
2641 let TieredMetadata::Tombstone(t) =
2642 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2643 else {
2644 panic!("expected tombstone");
2645 };
2646 assert_eq!(t.time_expires, None);
2647 assert!(matches!(
2648 backend
2649 .get_tiered_object(&id, Timestamp::now(), None)
2650 .await?,
2651 TieredGet::Tombstone(_)
2652 ));
2653
2654 let id = make_id();
2659 let ttl = Duration::from_hours(2 * 24);
2660 write_legacy_tombstone(&backend, &id, ExpirationPolicy::TimeToLive(ttl), None).await?;
2661
2662 let TieredMetadata::Tombstone(t) =
2663 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2664 else {
2665 panic!("expected TieredMetadata::Tombstone");
2666 };
2667 assert!(t.time_expires.is_some());
2668
2669 Ok(())
2670 }
2671
2672 #[tokio::test]
2674 async fn test_legacy_tombstone_tti_upgrade() -> Result<()> {
2675 let backend = create_test_backend().await?;
2676 let id = make_id();
2677 let path = id.as_storage_path().to_string().into_bytes();
2678
2679 let tti = Duration::from_hours(2 * 24);
2680
2681 let old_deadline = Timestamp::now() + Duration::from_mins(1);
2683 write_legacy_tombstone(
2684 &backend,
2685 &id,
2686 ExpirationPolicy::TimeToIdle(tti),
2687 Some(old_deadline),
2688 )
2689 .await?;
2690
2691 let TieredMetadata::Tombstone(_) =
2693 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2694 else {
2695 panic!("expected tombstone");
2696 };
2697 assert_eq!(
2698 backend
2699 .read_row(&path, "test-verify", Timestamp::now(), None)
2700 .await?
2701 .and_then(|row| row.time_expires()),
2702 Some(old_deadline)
2703 );
2704
2705 let requested = Timestamp::now() + tti;
2706 assert_eq!(
2707 backend
2708 .compare_and_update(
2709 &id,
2710 Some(&id),
2711 TieredUpdate::SetExpiry(ExpiryTarget::At(requested).into()),
2712 Timestamp::now()
2713 )
2714 .await?,
2715 SetExpiryResponse::Satisfied(requested)
2716 );
2717
2718 let new_deadline = match backend
2720 .read_row(&path, "test-verify", Timestamp::now(), None)
2721 .await?
2722 {
2723 Some(RowData::Tombstone { time_expires, .. }) => time_expires.unwrap(),
2724 _ => panic!("expected tombstone row after extension"),
2725 };
2726
2727 assert!(
2728 new_deadline > old_deadline,
2729 "explicit extension should extend tombstone expiry: {old_deadline:?} -> {new_deadline:?}"
2730 );
2731
2732 Ok(())
2733 }
2734
2735 #[tokio::test]
2740 async fn test_legacy_tombstone_conditional_ops() -> Result<()> {
2741 let backend = create_test_backend().await?;
2742
2743 let id = make_id();
2745 write_legacy_tombstone(&backend, &id, ExpirationPolicy::Manual, None).await?;
2746 let t_opt = backend
2747 .put_non_tombstone(&id, &Metadata::default(), Bytes::new(), Timestamp::now())
2748 .await?;
2749 assert_eq!(t_opt.map(|t| t.target).as_ref(), Some(&id));
2750
2751 let id = make_id();
2753 write_legacy_tombstone(&backend, &id, ExpirationPolicy::Manual, None).await?;
2754 let t_opt = backend.delete_non_tombstone(&id, Timestamp::now()).await?;
2755 assert_eq!(t_opt.map(|t| t.target).as_ref(), Some(&id));
2756
2757 let id = make_id();
2759 write_legacy_tombstone(&backend, &id, ExpirationPolicy::Manual, None).await?;
2760 let deleted = backend
2761 .compare_and_write(&id, Some(&id), TieredWrite::Delete, Timestamp::now())
2762 .await?;
2763 assert!(
2764 deleted,
2765 "CAS-delete must succeed on legacy-metadata tombstone"
2766 );
2767 assert!(matches!(
2768 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2769 TieredMetadata::NotFound
2770 ));
2771
2772 let id = make_id();
2774 write_empty_redirect_tombstone(&backend, &id).await?;
2775 let deleted = backend
2776 .compare_and_write(&id, Some(&id), TieredWrite::Delete, Timestamp::now())
2777 .await?;
2778 assert!(
2779 deleted,
2780 "CAS-delete must succeed on empty-redirect tombstone"
2781 );
2782 assert!(matches!(
2783 backend.get_tiered_metadata(&id, Timestamp::now()).await?,
2784 TieredMetadata::NotFound
2785 ));
2786
2787 Ok(())
2788 }
2789
2790 #[tokio::test]
2792 async fn test_empty_redirect_falls_back_to_hv_id() -> Result<()> {
2793 let backend = create_test_backend().await?;
2794 let id = make_id();
2795
2796 write_empty_redirect_tombstone(&backend, &id).await?;
2797 match backend.get_tiered_metadata(&id, Timestamp::now()).await? {
2798 TieredMetadata::Tombstone(t) => assert_eq!(t.target, id, "must fall back to hv_id"),
2799 other => panic!("expected tombstone, got {other:?}"),
2800 }
2801
2802 Ok(())
2803 }
2804
2805 #[tokio::test]
2810 async fn test_cas_create_tombstone_over_expired() -> Result<()> {
2811 let backend = create_test_backend().await?;
2812
2813 let id = make_id();
2814 let old_lt_id = ObjectId::random(id.context().clone());
2815 let old_tombstone = Tombstone {
2816 target: old_lt_id,
2817 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2818 };
2819 create_tombstone(&backend, &id, &old_tombstone).await?;
2820
2821 let new_lt_id = ObjectId::random(id.context().clone());
2822 let new_tombstone = Tombstone {
2823 target: new_lt_id.clone(),
2824 time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2825 };
2826 let committed = backend
2827 .compare_and_write(
2828 &id,
2829 None,
2830 TieredWrite::Tombstone(new_tombstone),
2831 Timestamp::now(),
2832 )
2833 .await?;
2834 assert!(
2835 committed,
2836 "CAS with current=None must succeed over an expired tombstone"
2837 );
2838
2839 let TieredMetadata::Tombstone(t) =
2840 backend.get_tiered_metadata(&id, Timestamp::now()).await?
2841 else {
2842 panic!("expected new tombstone to be readable");
2843 };
2844 assert_eq!(t.target, new_lt_id);
2845
2846 Ok(())
2847 }
2848
2849 #[tokio::test]
2852 async fn test_put_non_tombstone_over_expired() -> Result<()> {
2853 let backend = create_test_backend().await?;
2854
2855 let id = make_id();
2856 let lt_id = ObjectId::random(id.context().clone());
2857 let tombstone = Tombstone {
2858 target: lt_id,
2859 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2860 };
2861 create_tombstone(&backend, &id, &tombstone).await?;
2862
2863 let result = backend
2864 .put_non_tombstone(
2865 &id,
2866 &Metadata::default(),
2867 Bytes::from_static(b"data"),
2868 Timestamp::now(),
2869 )
2870 .await?;
2871 assert_eq!(
2872 result, None,
2873 "put_non_tombstone must succeed (return None) over an expired tombstone"
2874 );
2875
2876 let (_, _, stream) = backend
2877 .get_object(&id, Timestamp::now(), None)
2878 .await?
2879 .unwrap();
2880 assert_eq!(&stream::read_to_vec(stream).await?, b"data");
2881
2882 Ok(())
2883 }
2884
2885 async fn put_range_test_object(backend: &BigTableBackend) -> Result<ObjectId> {
2888 let id = make_id();
2889 let metadata = Metadata {
2890 content_type: "text/plain".into(),
2891 ..Default::default()
2892 };
2893 let payload = b"Hello, range requests!";
2894 backend
2895 .put_object(
2896 &id,
2897 &metadata,
2898 stream::single(payload.as_slice()),
2899 Timestamp::now(),
2900 )
2901 .await?;
2902 Ok(id)
2903 }
2904
2905 #[tokio::test]
2906 async fn get_object_range_bounded() -> Result<()> {
2907 let backend = create_test_backend().await?;
2908 let id = put_range_test_object(&backend).await?;
2909
2910 let (_, content_range, stream) = backend
2911 .get_object(&id, Timestamp::now(), Some(ByteRange::Bounded(7, 11)))
2912 .await?
2913 .unwrap();
2914 let data = stream::read_to_vec(stream).await?;
2915 assert_eq!(&data, b"range");
2916
2917 let content_range = content_range.unwrap();
2918 assert_eq!(content_range.start, 7);
2919 assert_eq!(content_range.end, 11);
2920 assert_eq!(content_range.total, 22);
2921
2922 Ok(())
2923 }
2924
2925 #[tokio::test]
2926 async fn get_object_range_from() -> Result<()> {
2927 let backend = create_test_backend().await?;
2928 let id = put_range_test_object(&backend).await?;
2929
2930 let (_, content_range, stream) = backend
2931 .get_object(&id, Timestamp::now(), Some(ByteRange::From(7)))
2932 .await?
2933 .unwrap();
2934 let data = stream::read_to_vec(stream).await?;
2935 assert_eq!(&data, b"range requests!");
2936
2937 let content_range = content_range.unwrap();
2938 assert_eq!(content_range.start, 7);
2939 assert_eq!(content_range.end, 21);
2940 assert_eq!(content_range.total, 22);
2941
2942 Ok(())
2943 }
2944
2945 #[tokio::test]
2946 async fn get_object_range_last() -> Result<()> {
2947 let backend = create_test_backend().await?;
2948 let id = put_range_test_object(&backend).await?;
2949
2950 let (_, content_range, stream) = backend
2951 .get_object(&id, Timestamp::now(), Some(ByteRange::Last(9)))
2952 .await?
2953 .unwrap();
2954 let data = stream::read_to_vec(stream).await?;
2955 assert_eq!(&data, b"requests!");
2956
2957 let content_range = content_range.unwrap();
2958 assert_eq!(content_range.start, 13);
2959 assert_eq!(content_range.end, 21);
2960 assert_eq!(content_range.total, 22);
2961
2962 Ok(())
2963 }
2964
2965 #[tokio::test]
2966 async fn get_object_range_unsatisfiable() -> Result<()> {
2967 let backend = create_test_backend().await?;
2968 let id = put_range_test_object(&backend).await?;
2969
2970 match backend
2971 .get_object(&id, Timestamp::now(), Some(ByteRange::From(100)))
2972 .await
2973 {
2974 Err(error) if matches!(error.kind(), ErrorKind::RangeNotSatisfiable { total: 22 }) => {}
2975 Ok(_) => panic!("expected RangeNotSatisfiable, got Ok"),
2976 Err(e) => panic!("expected RangeNotSatisfiable, got {e:?}"),
2977 }
2978
2979 Ok(())
2980 }
2981
2982 #[tokio::test]
2983 async fn get_object_no_range_returns_full_payload() -> Result<()> {
2984 let backend = create_test_backend().await?;
2985 let id = put_range_test_object(&backend).await?;
2986
2987 let (_, content_range, stream) = backend
2988 .get_object(&id, Timestamp::now(), None)
2989 .await?
2990 .unwrap();
2991 let data = stream::read_to_vec(stream).await?;
2992 assert_eq!(&data, b"Hello, range requests!");
2993 assert!(content_range.is_none());
2994
2995 Ok(())
2996 }
2997
2998 #[test]
2999 fn row_size_counts_the_key_and_every_cell() {
3000 let path = b"attachments/org.1/objects/abc";
3001 let (mutations, size) =
3002 object_mutations(path, Metadata::default(), b"0123456789".to_vec()).unwrap();
3003
3004 let stamped = Metadata {
3008 size: Some(10),
3009 ..Default::default()
3010 };
3011 let metadata_len = serde_json::to_vec(&stamped).unwrap().len();
3012 let expected = (path.len() + 10 + metadata_len) as u64;
3013
3014 assert_eq!(row_size(path, &mutations), expected);
3015 assert_eq!(size, expected, "the size handed back matches the mutations");
3016 }
3017
3018 #[test]
3019 fn row_size_is_nonzero_for_tombstones() {
3020 let path = b"attachments/org.1/objects/abc";
3021 let time_expires = Timestamp::now() + Duration::from_secs(60);
3022 let tombstone = Tombstone {
3023 target: ObjectId::from_storage_path("attachments/org.1/objects/abc/0199").unwrap(),
3024 time_expires: Some(time_expires),
3025 };
3026 let mutations = tombstone_mutations(&tombstone);
3027
3028 assert_eq!(mutations.len(), 2);
3029 let set_cell = mutations[1].mutation.as_ref().unwrap();
3030 let mutation::Mutation::SetCell(set_cell) = set_cell else {
3031 panic!("expected redirect SetCell mutation");
3032 };
3033 assert_eq!(set_cell.family_name, FAMILY_GC);
3034 assert_eq!(set_cell.column_qualifier, COLUMN_REDIRECT);
3035 assert_eq!(set_cell.timestamp_micros, time_expires.as_micros() as i64);
3036 assert!(row_size(path, &mutations) > path.len() as u64);
3037 }
3038
3039 #[cfg(feature = "storage-cogs")]
3040 #[tokio::test]
3041 async fn change_stream_reports_writes_and_deletes() -> Result<()> {
3042 let (backend, producer) = create_test_backend_with_change_stream().await?;
3043 let id = make_id();
3044 let metadata = Metadata {
3045 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(3600)),
3046 time_expires: Some(Timestamp::now() + Duration::from_secs(3600)),
3047 ..Default::default()
3048 };
3049
3050 backend
3051 .put_object(
3052 &id,
3053 &metadata,
3054 stream::single::<crate::stream::ClientError>(b"hello".to_vec()),
3055 Timestamp::now(),
3056 )
3057 .await?;
3058 backend.delete_object(&id, Timestamp::now()).await?;
3059
3060 let records = producer.records();
3061 assert_eq!(records.len(), 2);
3062
3063 assert_eq!(records[0].op_type, OpType::Write);
3064 assert_eq!(records[0].app_feature, "testing");
3065 assert_eq!(records[0].shared_resource_id, "bigtable_objectstore");
3066 assert!(records[0].size.unwrap() > b"hello".len() as u64);
3068 assert!(records[0].expiration_time.is_some());
3069
3070 assert_eq!(records[1].op_type, OpType::Delete);
3071 assert_eq!(records[1].size, None);
3072 assert_eq!(records[1].record_id, records[0].record_id);
3073
3074 Ok(())
3075 }
3076
3077 #[cfg(feature = "storage-cogs")]
3078 #[tokio::test]
3079 async fn delete_non_tombstone_reclaims_expired_rows() -> Result<()> {
3080 let (backend, producer) = create_test_backend_with_change_stream().await?;
3081
3082 let id = make_id();
3085 let tombstone = Tombstone {
3086 target: ObjectId::random(id.context().clone()),
3087 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
3088 };
3089 create_tombstone(&backend, &id, &tombstone).await?;
3090 assert_eq!(
3091 backend.delete_non_tombstone(&id, Timestamp::now()).await?,
3092 None,
3093 "an expired tombstone must not be returned to the caller"
3094 );
3095 let records = producer.records();
3096 assert_eq!(records.len(), 1, "the expired tombstone must be reclaimed");
3097 assert_eq!(records[0].op_type, OpType::Delete);
3098
3099 let id = make_id();
3101 let metadata = Metadata {
3102 expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(0)),
3103 time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
3104 ..Default::default()
3105 };
3106 create_object(&backend, &id, &metadata, b"gone", Timestamp::now()).await?;
3107 producer.clear();
3108 assert_eq!(
3109 backend.delete_non_tombstone(&id, Timestamp::now()).await?,
3110 None
3111 );
3112 let records = producer.records();
3113 assert_eq!(records.len(), 1, "the object row must be reclaimed");
3114 assert_eq!(records[0].op_type, OpType::Delete);
3115
3116 Ok(())
3117 }
3118
3119 #[cfg(feature = "storage-cogs")]
3120 #[tokio::test]
3121 async fn change_stream_reports_tombstone_rows() -> Result<()> {
3122 let (backend, producer) = create_test_backend_with_change_stream().await?;
3123 let id = make_id();
3124 let target = new_test_revision(&id);
3125
3126 let tombstone = Tombstone {
3127 target: target.clone(),
3128 time_expires: Some(Timestamp::now() + Duration::from_secs(3600)),
3129 };
3130 let written = backend
3131 .compare_and_write(
3132 &id,
3133 None,
3134 TieredWrite::Tombstone(tombstone),
3135 Timestamp::now(),
3136 )
3137 .await?;
3138 assert!(written);
3139
3140 let records = producer.records();
3141 assert_eq!(records.len(), 1);
3142 assert_eq!(records[0].op_type, OpType::Write);
3143 assert!(
3144 records[0].size.unwrap() > 0,
3145 "tombstone rows occupy storage and must not report zero"
3146 );
3147 assert!(records[0].expiration_time.is_some());
3148
3149 Ok(())
3150 }
3151
3152 #[cfg(feature = "storage-cogs")]
3153 #[tokio::test]
3154 async fn change_stream_reports_expiry_extension_as_an_update() -> Result<()> {
3155 let (backend, producer) = create_test_backend_with_change_stream().await?;
3156 let id = make_id();
3157 let metadata = Metadata {
3158 expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(3600)),
3159 time_expires: Some(Timestamp::now() + Duration::from_secs(1)),
3160 ..Default::default()
3161 };
3162
3163 backend
3164 .put_object(
3165 &id,
3166 &metadata,
3167 stream::single::<crate::stream::ClientError>(b"hello".to_vec()),
3168 Timestamp::now(),
3169 )
3170 .await?;
3171 producer.clear();
3172
3173 backend
3174 .set_expiry(
3175 &id,
3176 ExpiryTarget::At(Timestamp::now() + Duration::from_secs(3600)).into(),
3177 Timestamp::now(),
3178 )
3179 .await?;
3180
3181 let records = producer.records();
3182 assert_eq!(records.len(), 1, "expected exactly one extension report");
3183 assert_eq!(records[0].op_type, OpType::Update);
3184 assert_eq!(
3185 records[0].size, None,
3186 "an extension does not change the size"
3187 );
3188 assert!(records[0].expiration_time.is_some());
3189
3190 Ok(())
3191 }
3192
3193 #[cfg(feature = "storage-cogs")]
3194 fn new_test_revision(id: &ObjectId) -> ObjectId {
3195 ObjectId {
3196 context: id.context.clone(),
3197 key: format!("{}/{}", id.key, uuid::Uuid::now_v7()),
3198 }
3199 }
3200}