Skip to main content

objectstore_service/backend/
bigtable.rs

1//! BigTable backend for high-volume, low-latency storage of small objects.
2//!
3//! # Row Format
4//!
5//! Object row keys use the object's storage path. An object row contains either an **object** or a
6//! **tombstone** — never both. The two layouts are mutually exclusive and distinguished by
7//! column presence:
8//!
9//! | Column | Family    | Content                     | Present when       |
10//! |--------|-----------|-----------------------------|--------------------|
11//! | `p`    | `fg`/`fm` | Compressed payload bytes    | Object row only    |
12//! | `m`    | `fg`/`fm` | [`Metadata`] JSON           | Object row only    |
13//! | `r`    | `fg`/`fm` | Redirect path to LT storage | Tombstone row only |
14//!
15//! The `r` column signals a tombstone row: its **value** is the long-term `ObjectId`
16//! serialized via `as_storage_path()`. Callers can resolve the LT object directly from the
17//! `r` value without reconstructing it from the row key. Its column family and timestamp
18//! carry the tombstone's concrete expiration deadline.
19//!
20//! `p`/`m` and `r` are mutually exclusive. Every write begins with a `DeleteFromRow`
21//! mutation that clears all columns before writing the new cells, so mixed rows cannot exist.
22//!
23//! Upload markers use a separate `uploads` namespace and contain only a one-byte
24//! `u` cell in `fg`, timestamped with the upload deadline.
25//!
26//! ## Legacy Tombstone Format
27//!
28//! Tombstones written before the `r` column layout used the object-row format with an
29//! empty `p` column and `"is_redirect_tombstone": true` in the `m` JSON. Both formats are
30//! supported for reading. A `bigtable.legacy_tombstone_read` metric is emitted on each legacy
31//! read. Legacy tombstones expire naturally by TTL/GC; a successful conditional
32//! expiry extension upgrades them to the `r` format. Tombstone metadata in the historical `t`
33//! column is ignored; the corresponding `r` cell contains all information needed by readers.
34
35use 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/// Configuration for [`BigTableBackend`].
65///
66/// Stores objects in [Google Cloud Bigtable], a NoSQL wide-column database optimized for
67/// high-throughput, low-latency workloads with small objects. Authentication uses Application
68/// Default Credentials (ADC).
69///
70/// **Note**: The table must be pre-created with the following column families:
71/// - `fg`: timestamp-based garbage collection (`maxage=1s`)
72/// - `fm`: manual garbage collection (`no GC policy`)
73///
74/// [Google Cloud Bigtable]: https://cloud.google.com/bigtable
75///
76/// # Example
77///
78/// ```yaml
79/// storage:
80///   type: bigtable
81///   project_id: my-project
82///   instance_name: objectstore
83///   table_name: objectstore
84/// ```
85#[derive(Debug, Clone, Deserialize, Serialize)]
86pub struct BigTableConfig {
87    /// Optional custom Bigtable endpoint.
88    ///
89    /// Useful for testing with emulators. If `None`, uses the default Bigtable endpoint.
90    ///
91    /// # Default
92    ///
93    /// `None` (uses default Bigtable endpoint)
94    ///
95    /// # Environment Variables
96    ///
97    /// - `OS__STORAGE__TYPE=bigtable`
98    /// - `OS__STORAGE__ENDPOINT=localhost:8086` (optional)
99    pub endpoint: Option<String>,
100
101    /// GCP project ID.
102    ///
103    /// The Google project ID (not project number) containing the Bigtable instance.
104    ///
105    /// # Environment Variables
106    ///
107    /// - `OS__STORAGE__PROJECT_ID=my-project`
108    pub project_id: String,
109
110    /// Bigtable instance name.
111    ///
112    /// # Environment Variables
113    ///
114    /// - `OS__STORAGE__INSTANCE_NAME=my-instance`
115    pub instance_name: String,
116
117    /// Bigtable table name.
118    ///
119    /// The table must exist before starting the server.
120    ///
121    /// # Environment Variables
122    ///
123    /// - `OS__STORAGE__TABLE_NAME=objectstore`
124    pub table_name: String,
125
126    /// Optional number of connections to maintain to Bigtable.
127    ///
128    /// # Default
129    ///
130    /// `None` (defaults to 1)
131    ///
132    /// # Environment Variables
133    ///
134    /// - `OS__STORAGE__CONNECTIONS=16` (optional)
135    pub connections: Option<usize>,
136
137    /// Timeout for an individual Bigtable RPC attempt.
138    ///
139    /// # Default
140    ///
141    /// `2s`
142    ///
143    /// # Environment Variables
144    ///
145    /// - `OS__STORAGE__RPC_TIMEOUT=2s`
146    /// - `OS__STORAGE__HIGH_VOLUME__RPC_TIMEOUT=2s` (tiered storage)
147    #[serde(default = "default_rpc_timeout", with = "humantime_serde")]
148    pub rpc_timeout: Duration,
149
150    /// Reports what this backend stores, for per-usecase cost attribution.
151    ///
152    /// # Default
153    ///
154    /// `None`, which disables reporting for this backend.
155    ///
156    /// # Environment Variables
157    ///
158    /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=bigtable_objectstore`
159    /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional)
160    #[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
168/// Maximum age for connections (GRPC channels) to Bigtable, after which they will be swapped with
169/// new ones in the background.
170/// This is intended to avoid latency spikes that could occur every hour or so, when the server
171/// closes long standing connections ([source](https://web.archive.org/web/20260211140930/https://docs.cloud.google.com/bigtable/docs/performance#cold-starts:~:text=return%20an%20error.-,Cold%20start,-at%20client%20initialization)).
172/// `tonic` already handles reconnections transparently, but lazily, meaning that the first requests
173/// that attempt to use a certain channel after the server has closed it will pay the cost of the
174/// reconnection, resulting in increased latency for those requests.
175const MAX_CHANNEL_AGE: Option<Duration> = Some(Duration::from_mins(50));
176/// Permission scopes required for accessing the BigTable data API.
177const TOKEN_SCOPES: &[&str] = &["https://www.googleapis.com/auth/bigtable.data"];
178
179/// How often to retry failed requests.
180const REQUEST_RETRY_COUNT: usize = 2;
181/// How many times to retry a CAS mutation before giving up and returning an error.
182const CAS_RETRY_COUNT: usize = 3;
183
184/// Column that stores the raw payload (compressed).
185const COLUMN_PAYLOAD: &[u8] = b"p";
186/// Column that stores metadata in JSON.
187const COLUMN_METADATA: &[u8] = b"m";
188/// Column that stores the redirect path for tombstone rows.
189const COLUMN_REDIRECT: &[u8] = b"r";
190/// Regex to match all non-payload columns (`m`, `r`) for metadata-only reads.
191const FILTER_META: &[u8] = b"^[mr]$";
192
193/// Column which marks the presence of an ongoing resumable upload.
194const COLUMN_UPLOAD: &[u8] = b"u";
195
196/// Column family that uses timestamp-based garbage collection.
197///
198/// We require a GC rule on this family to automatically delete rows.
199/// See: <https://cloud.google.com/bigtable/docs/gc-cell-level>
200const FAMILY_GC: &str = "fg";
201/// Column family that uses manual garbage collection.
202const FAMILY_MANUAL: &str = "fm";
203
204/// BigTable storage backend for high-volume, low-latency object storage.
205pub 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
225/// Creates a row filter that matches a single column by exact qualifier.
226fn 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
234/// Creates a row filter matching the legacy tombstone format: `m` column JSON starts with
235/// `{"is_redirect_tombstone":true`.
236///
237/// After legacy tombstones expire naturally this filter becomes dead code in both callers.
238fn 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
253/// Wraps `inner` so that it only matches live (non-expired) cells.
254///
255/// Uses the rounded access time directly. Legacy fractional cell timestamps can be excluded
256/// up to one second before their rounded metadata deadline; new writes are second-aligned.
257fn 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                    // Manual family: never expires.
263                    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                    // GC family: only match non-expired cells.
276                    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
302/// Builds a raw row filter that matches any live tombstone row, new- or legacy-format.
303///
304/// New format: presence of the `r` column.
305/// Legacy format: `is_redirect_tombstone: true` in the `m` column JSON.
306///
307/// After legacy tombstones expire naturally this simplifies to just
308/// `column_filter(COLUMN_REDIRECT)`.
309fn 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
320/// Returns a [`MutatePredicate`] that matches any live tombstone row.
321///
322/// Mutations will not run on live tombstones (`predicate_matched == false`). They _will_
323/// run on expired tombstones as well as non-tombstones. Used by
324/// [`BigTableBackend::put_non_tombstone`] and [`BigTableBackend::compare_and_write`] as
325/// the `CheckAndMutateRow` predicate.
326///
327/// This predicate cannot distinguish an empty row from a row holding a regular object; a caller
328/// that needs to know whether its mutation hit anything wants [`non_tombstone_predicate`].
329fn tombstone_predicate(access_time: Timestamp) -> MutatePredicate {
330    MutatePredicate::Exclude(tombstone_filter(access_time))
331}
332
333/// Returns a [`MutatePredicate`] that is the logical negation of [`tombstone_predicate`];
334/// it matches everything _except_ live tombstones.
335///
336/// Mutations run only when the predicate matches (`predicate_matched == true`). They will
337/// run on expired tombstones as well as non-tombstones. The match result doubles as a
338/// "was a row removed?" signal. Used by [`BigTableBackend::delete_non_tombstone`] as the
339/// `CheckAndMutateRow` predicate.
340///
341/// Built as a `Condition` filter:
342/// - Predicate: [`tombstone_filter`] -> is a live tombstone present?
343/// - True branch: `BlockAllFilter` -> match nothing; live tombstones should be preserved
344/// - False branch: `PassAllFilter` -> match everything: expired tombstones and non-tombstones
345fn 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
361/// Builds an anchored regex pattern (`^…$`) that matches `value` literally.
362///
363/// Uses [`regex::escape`] so that metacharacters in storage paths (`.`, `/`, etc.)
364/// are treated as literal bytes.
365fn exact_value_regex(value: &str) -> Vec<u8> {
366    format!("^{}$", regex::escape(value)).into_bytes()
367}
368
369/// Matches tombstones whose redirect resolves to `target`.
370///
371/// ## Predicate Matches
372///
373/// Must be used with `true_mutations` and `predicate_matched == true`.
374///
375/// ## Details
376///
377/// Always includes an exact match on the `r` (redirect) column:
378/// - Chain: `r` column present AND value == `target` storage path
379///
380/// When `target == own_id` (the caller expects a legacy identity redirect), the
381/// exact match is wrapped in an Interleave with two additional fallbacks:
382/// - Chain: `r` column present AND value == `b""` (empty-sentinel written before the redirect
383///   column stored the path)
384/// - Chain: `m` column present AND value matches `{"is_redirect_tombstone":true...}` regex
385///   (legacy metadata format predating the dedicated `r` column)
386fn 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    // Also match legacy tombstones that resolve to the HV id:
420    // - empty `r` value (written before the redirect column stored the path)
421    // - legacy `m` column format (`is_redirect_tombstone: true`)
422    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
432/// Returns a [`MutatePredicate`] that matches tombstones whose redirect resolves to either `old` or `new`.
433///
434/// Mutations run only when the predicate matches (`predicate_matched == true`):
435/// equivalent to `t == old || t == new`. Built as an Interleave of two
436/// [`redirect_target_filter`] calls — yields cells iff at least one branch matches.
437/// An absent row or non-tombstone row yields 0 cells, so `predicate_matched = false` (conflict).
438fn 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
456/// Returns a [`MutatePredicate`] that matches rows where no conflicting tombstone exists.
457///
458/// Mutations run only when the row is conflict-free (`predicate_matched == false`):
459/// no tombstone is present, or the tombstone's redirect already points to `target`.
460///
461/// Built as an inverted `Condition` filter:
462/// - Predicate: [`redirect_target_filter`]`(target)` — tombstone already points to `target`?
463/// - True branch: `BlockAllFilter` → 0 cells (already at target, safe state).
464/// - False branch: [`tombstone_filter`] → 0 cells when no tombstone exists.
465///
466/// Both safe states yield 0 cells, so `predicate_matched = false` in both cases.
467fn 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
517/// Matches an inline row whose metadata cell has the observed expiry timestamp.
518fn 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/// The condition under which a [`BigTableBackend::check_and_mutate`] write proceeds.
561///
562/// Each variant pairs a row filter with the state that makes the write safe:
563/// `Include` writes when the row matches; `Exclude` writes when it does not.
564#[derive(Clone, Debug)]
565enum MutatePredicate {
566    /// Write proceeds when the filter matches the row.
567    ///
568    /// Mutations run in `true_mutations`; succeeds when `predicate_matched == true`.
569    Include(v2::RowFilter),
570    /// Write proceeds when the filter does not match the row.
571    ///
572    /// Mutations run in `false_mutations`; succeeds when `predicate_matched == false`.
573    Exclude(v2::RowFilter),
574}
575
576/// Creates a row filter that reads all non-payload columns (`m`, `r`).
577///
578/// Used by metadata-only reads to avoid fetching the (potentially large) payload column
579/// while still being able to detect both new- and legacy-format tombstones.
580fn 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
594/// Creates a `DeleteFromRow` mutation wrapped in the outer [`v2::Mutation`] envelope.
595fn delete_row_mutation() -> v2::Mutation {
596    mutation(mutation::Mutation::DeleteFromRow(
597        mutation::DeleteFromRow {},
598    ))
599}
600
601/// Builds the three mutations that write an object row: clear existing data,
602/// then set the payload and metadata cells. Returns them with the resulting row size.
603///
604/// Used by both [`BigTableBackend::put_row`] (unconditional write) and
605/// [`BigTableBackend::put_non_tombstone`] (conditional write).
606fn 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    // Record the payload size in the metadata before persisting it.
617    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        // NB: We explicitly delete the row to clear metadata on overwrite.
624        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
643/// Approximates the bytes a row occupies, as its key plus every cell value written.
644///
645/// This function does not distinguish between object rows and tombstone rows. It does not
646/// include Bigtable's own overhead.
647fn 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
659/// Builds the two mutations that write a tombstone row: clear existing data,
660/// then set the redirect cell.
661///
662/// Used by both unconditional tombstone writes and the conditional expiry-extension paths.
663fn 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/// Subset of [`Metadata`] that indicates a row is a tombstone instead of a real object.
681///
682/// Used to construct [`RowData`].
683#[derive(Debug, Deserialize)]
684struct LegacyTombstoneMeta {
685    /// Internal redirect tombstone marker.
686    ///
687    /// When `true`, this object is a legacy tombstone. This implies:
688    ///  - the payload is empty
689    ///  - metadata is not meaningful
690    ///  - the `r` column is not present
691    #[serde(default)]
692    is_redirect_tombstone: bool,
693}
694
695/// Parsed data from a BigTable row's cells.
696enum RowData {
697    /// A regular object row with payload and metadata.
698    Object {
699        metadata: Metadata,
700        payload: Vec<u8>,
701        /// Original GC timestamp for exact CAS matching, including legacy fractional seconds.
702        expiry_micros: i64,
703    },
704    /// A tombstone row indicating the real payload lives on the long-term backend.
705    Tombstone {
706        target: Vec<u8>,
707        time_expires: Option<Timestamp>,
708        /// Original GC timestamp for exact CAS matching, including legacy fractional seconds.
709        expiry_micros: i64,
710    },
711}
712
713impl RowData {
714    /// Parses a set of row cells into a [`RowData`].
715    ///
716    /// New-format tombstones are identified by the presence of the `r` column.
717    /// Legacy tombstones (written before the column migration) are identified by
718    /// `is_redirect_tombstone: true` in the `m` column JSON; a
719    /// `bigtable.legacy_tombstone_read` metric is emitted on each such read.
720    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            // NB: All cells are written with the same timestamp; last write is safe.
730
731            // Only derive expiration from GC-family cells — manual-family cells
732            // use server-assigned timestamps that don't represent expiration.
733            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            // Metadata may have been skipped by a payload-only internal read.
775            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    /// Returns the resolved expiration timestamp for this row, regardless of variant.
786    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    /// Returns `true` if this row is expired as of the given `time`.
794    ///
795    /// Only applies to rows with an expiration deadline.
796    fn expires_before(&self, time: Timestamp) -> bool {
797        self.time_expires().is_some_and(|ts| ts < time)
798    }
799}
800
801/// Parses the raw `r` column bytes into a redirect target [`ObjectId`].
802///
803/// For tombstones with an empty `r` value, falls back to the ID of the tombstone
804/// itself and emits a `bigtable.empty_redirect_read` metric so deployments can
805/// track when it is safe to remove the legacy empty-value code path.
806fn 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    /// Creates a new [`BigTableBackend`] from the given `config`.
820    ///
821    /// Pass an `endpoint` in the config to connect to a local emulator; omit it to use real GCP
822    /// credentials. `connections` controls the gRPC connection pool size (defaults to 1).
823    /// A `PingAndWarm` request is sent through the pool every 10 seconds in an effort
824    /// to keep the connections active.
825    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, // is_read_only
846                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, // is_read_only
855                Some(rpc_timeout),
856                Arc::new(token_provider),
857                connections.unwrap_or(1),
858                true, // prime_channels
859                None, // app_profile_id
860                MAX_CHANNEL_AGE,
861                Some(Duration::from_secs(10)), // periodic PingAndWarm
862            )
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    /// Reads a single row by key, returning parsed row data.
878    ///
879    /// Returns `None` if the row is absent or has expired.
880    #[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    /// Writes an object row, returning the size of the row it wrote.
941    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    /// Executes a `CheckAndMutateRow` request.
954    #[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            // A tombstone was present: read its data for the caller.
1159            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                // Race: Tombstone was replaced by an object, retry to overwrite
1180                Some(RowData::Object { .. }) => continue,
1181                // Race: Tombstone was deleted, retry to write.
1182                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                    // If object size wasn't written into the metadata, re-compute it now
1225                    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        // Read metadata and tombstone columns — skip the (potentially large) payload.
1244        // NB: `metadata.size` will only be populated if the size was added to the metadata before
1245        // writing to Bigtable.
1246        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        // Inline extension needs metadata and payload from the same read so a
1283        // successful conditional rewrite can preserve the payload verbatim.
1284        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); // wrong row kind
1299                }
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); // already expired
1306                }
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)); // already satisfied
1313                }
1314
1315                // Observing a live cell here is not atomic with wall-clock
1316                // expiry or Bigtable GC. The conditional write may still lose
1317                // to either and then returns false.
1318                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); // wrong row kind
1337                };
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); // wrong target
1348                }
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)); // already satisfied
1354                }
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            // Nothing was deleted: either a tombstone is in the way, or the row is absent.
1407            // Read the row to find out which, and to hand the tombstone to the caller.
1408            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                // Race: An object appeared since the predicate ran, delete the new object now.
1429                Some(RowData::Object { .. }) => continue,
1430                // The row is absent or expired, nothing left to do.
1431                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        // Get the correct set of mutations to apply as well as the new expiration date.
1460        // If we're deleting something, `expires_at` is `None`. If we're writing something
1461        // without an expiration date, `expires_at` is `Some(None)`.
1462        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            // Don't record anything if the write didn't succeed
1486            (false, _) => {}
1487            // We wrote something (the inner `expires_at` is `None` for manual GC)
1488            (true, Some(expires_at)) => {
1489                self.change_stream
1490                    .write(id, row_size(&path, &mutations), expires_at)
1491            }
1492            // We deleted something
1493            (true, None) => self.change_stream.delete(id),
1494        }
1495
1496        Ok(written)
1497    }
1498}
1499
1500/// Retries a BigTable RPC on transient errors.
1501async 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                // Non-RPC error; the error event carries the details.
1522                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        // Transient errors on auth token refresh
1548        BigTableError::GCPAuthError(_) => true,
1549        // Transient GRPC network failures
1550        BigTableError::TransportError(_) => true,
1551        // These could also indicate transient network failures
1552        BigTableError::IoError(_) => true,
1553        BigTableError::TimeoutError(_) => true,
1554
1555        // See https://docs.cloud.google.com/bigtable/docs/status-codes
1556        BigTableError::RpcError(status) => match status.code() {
1557            // Generic retriable status
1558            Code::Unavailable => true,
1559            // Timeouts
1560            Code::Cancelled => true,
1561            Code::DeadlineExceeded => true,
1562            // Token might have refreshed too late
1563            Code::Unauthenticated => true,
1564            // Unspecified, attempt to retry anyways
1565            Code::Aborted => true,
1566            Code::Internal => true,
1567            Code::FailedPrecondition => true,
1568            Code::Unknown => true,
1569            _ => false,
1570        },
1571        _ => false,
1572    }
1573}
1574
1575/// Resolves an optional byte range against a payload buffer, returning the
1576/// applicable content range and the (potentially narrowed) payload.
1577///
1578/// When `range` is `None`, returns the full payload unchanged. Uses
1579/// `Bytes::slice` to avoid copying data.
1580fn 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    // NB: Most of these tests require a BigTable emulator running. This is done
1610    // automatically in CI.
1611    //
1612    // Refer to the readme for how to set up the emulator.
1613
1614    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        // Resolve `time_expires` from `now` (as `from_insert_headers` does) unless the test set
1660        // it explicitly, so `object_mutations` has an expiration to persist.
1661        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    /// Writes a legacy-format tombstone row directly into Bigtable.
1682    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    /// Writes a historical `r`/`t` tombstone row with an empty `r` value directly.
1717    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(), // empty — legacy format
1728            })),
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    // --- Section 1: Object Operations ---
1743
1744    /// Verifies the full roundtrip: put → get_object (payload + metadata) → get_metadata (metadata).
1745    #[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    /// Verifies that a server-resolved `time_expires` is persisted verbatim, not recomputed.
1783    #[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    /// Verifies that absent rows return None or succeed silently for all read/delete operations.
1804    #[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    /// Backend reads are side-effect-free; explicit extension preserves payload.
1876    #[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            // Backdate `now` so the written expiry (past_now + tti) is stale but not expired.
1891            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            // Simulate a legacy fractional deadline. Renewal must match the raw GC timestamp.
1899            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    // --- Section 2: Expiration ---
2112
2113    #[tokio::test]
2114    async fn test_ttl_immediate() -> Result<()> {
2115        // NB: We create a TTL that immediately expires in this test. This might be optimized away
2116        // in a future implementation, so we will have to update this test accordingly.
2117
2118        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        // NB: We create a TTI that immediately expires in this test. This might be optimized away
2141        // in a future implementation, so we will have to update this test accordingly.
2142
2143        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    // --- Section 3: Tiered Operations ---
2164
2165    /// Covers all three row states for `get_tiered_object` and `get_tiered_metadata`.
2166    ///
2167    /// - **empty**: both return NotFound.
2168    /// - **object**: put_object, both return the Object variant with correct payload/metadata.
2169    /// - **tombstone**: CAS-write with a distinct `lt_id`, both return the Tombstone variant
2170    ///   with `target == lt_id`.
2171    #[tokio::test]
2172    async fn test_tiered_get() -> Result<()> {
2173        let backend = create_test_backend().await?;
2174
2175        // empty
2176        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        // object
2189        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        // tombstone
2217        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    /// Covers all three row states for `put_non_tombstone`.
2244    ///
2245    /// - **empty**: returns None, object is readable.
2246    /// - **object**: overwrites with new payload, returns None.
2247    /// - **tombstone**: returns Some(Tombstone) with the correct target; tombstone still intact.
2248    #[tokio::test]
2249    async fn test_put_non_tombstone() -> Result<()> {
2250        let backend = create_test_backend().await?;
2251
2252        // empty: put_non_tombstone on absent row succeeds and makes object readable.
2253        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        // object: put_non_tombstone on existing object replaces payload, returns None.
2271        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        // tombstone: put_non_tombstone returns Some(Tombstone) and leaves tombstone intact.
2284        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    /// Covers all three row states for `delete_non_tombstone`.
2310    ///
2311    /// - **empty**: returns None.
2312    /// - **object**: returns None, row gone.
2313    /// - **tombstone**: returns Some(Tombstone) with correct target; tombstone still intact.
2314    ///
2315    /// Verifies that the `r` column is correctly detected by both the `ReadRows` column
2316    /// filter and the `CheckAndMutate` `non_tombstone_predicate`.
2317    #[tokio::test]
2318    async fn test_delete_non_tombstone() -> Result<()> {
2319        let backend = create_test_backend().await?;
2320
2321        // empty
2322        let id = make_id();
2323        assert_eq!(
2324            backend.delete_non_tombstone(&id, Timestamp::now()).await?,
2325            None
2326        );
2327
2328        // object
2329        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        // tombstone
2344        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    // --- Section 4: Compare-and-Write ---
2367
2368    /// Creating a tombstone on an empty row succeeds; a retry of the same CAS also succeeds.
2369    ///
2370    /// After creation, both tiered and legacy APIs reflect the tombstone.
2371    #[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        // First create succeeds.
2384        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        // Tiered reads must see the tombstone with the correct target and deadline.
2395        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        // Legacy reads must error rather than leak tombstone data.
2412        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        // Idempotent retry: retry with the same target succeeds
2426        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    /// Swapping a tombstone target: wrong expected → false, correct expected → true.
2440    #[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        // Wrong target: CAS fails, tombstone unchanged.
2456        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        // Correct target: CAS succeeds, target updated.
2473        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        // Idempotent retry: same A→B swap returns true.
2486        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    /// Swapping a tombstone for inline object data: wrong expected → false, correct → true.
2495    #[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        // Wrong target: CAS fails, tombstone intact.
2510        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        // Correct target: CAS succeeds, row becomes an inline object.
2521        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(&lt_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        // Idempotent retry: row is already inline (no tombstone), same CAS returns true.
2536        let retry = backend
2537            .compare_and_write(&id, Some(&lt_id), write, Timestamp::now())
2538            .await?;
2539        assert!(retry, "idempotent retry");
2540
2541        Ok(())
2542    }
2543
2544    /// CAS-write an object onto an empty row (expected=None, write=Object) succeeds.
2545    #[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    /// CAS-delete: wrong expected → false; correct expected → true, row gone.
2569    #[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        // Wrong target: fails, row preserved.
2584        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        // Correct target: succeeds, row gone.
2594        let deleted = backend
2595            .compare_and_write(&id, Some(&lt_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        // Idempotent retry: row is already absent (no tombstone), same delete returns true.
2604        let retry = backend
2605            .compare_and_write(&id, Some(&lt_id), TieredWrite::Delete, Timestamp::now())
2606            .await?;
2607        assert!(retry, "idempotent retry");
2608
2609        // Inline object replaced tombstone: Safe to delete since it is an idempotent operation.
2610        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    // --- Section 5: Legacy Tombstone Compatibility ---
2628
2629    /// Legacy Manual and TTL tombstones are correctly read via the tiered APIs.
2630    ///
2631    /// Uses `Manual` expiration so `timestamp_micros = -1` (server-assigned ≈ write time)
2632    /// does not trigger immediate expiry.
2633    #[tokio::test]
2634    async fn test_legacy_tombstone_reads() -> Result<()> {
2635        let backend = create_test_backend().await?;
2636
2637        // Manual policy: get_tiered_metadata returns a non-expiring tombstone.
2638        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        // TTL policy: get_tiered_metadata reconstructs the concrete deadline.
2655        //
2656        // A future cell timestamp (now + TTL) is required so `expires_before` does not
2657        // immediately filter the row.
2658        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    /// A conditional extension upgrades a legacy TTI tombstone to `r`.
2673    #[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        // Place time_expires near expiry but still in the future.
2682        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        // A read observes the legacy row but leaves it unchanged.
2692        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        // After extension, the row uses the requested timestamp.
2719        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    /// Legacy tombstones are handled correctly by all conditional write operations.
2736    ///
2737    /// Covers: `put_non_tombstone`, `delete_non_tombstone`, CAS-delete for both the
2738    /// legacy-metadata format and the empty-redirect format.
2739    #[tokio::test]
2740    async fn test_legacy_tombstone_conditional_ops() -> Result<()> {
2741        let backend = create_test_backend().await?;
2742
2743        // put_non_tombstone returns Some(target == id) for a legacy tombstone.
2744        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        // delete_non_tombstone returns Some(target == id) for a legacy tombstone.
2752        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        // CAS-delete succeeds on a legacy-metadata tombstone (target resolves to hv_id).
2758        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        // CAS-delete succeeds on an empty-redirect tombstone (target resolves to hv_id).
2773        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    /// An empty `r` value falls back to the HV id when resolving the tombstone target.
2791    #[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    // --- Section 6: Expired Tombstone Handling ---
2806
2807    /// CAS with `current=None` must succeed when the row holds an expired
2808    /// tombstone. The physical row still exists but is logically gone.
2809    #[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    /// `put_non_tombstone` must succeed when the row holds only an expired
2850    /// tombstone — the expired row is logically absent.
2851    #[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    // --- Range Request Tests ---
2886
2887    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        // The key, the 10-byte payload, and the serialized metadata. `object_mutations`
3005        // stamps the size into the metadata before serializing it, so the expected length
3006        // has to account for that too.
3007        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        // Key plus payload plus metadata, so strictly more than the payload alone.
3067        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        // An expired tombstone is past its lifetime: reclaimed, not handed to the caller.
3083        // This test serves as documentation of that potentially surprising behavior.
3084        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        // The same holds for an object row (expired or otherwise).
3100        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}