Skip to main content

objectstore_service/backend/
gcs.rs

1//! Google Cloud Storage backend for long-term storage of large objects.
2
3use std::borrow::Cow;
4use std::collections::BTreeMap;
5use std::future::Future;
6use std::num::NonZeroU64;
7use std::sync::Arc;
8use std::time::SystemTime;
9use std::{fmt, io};
10
11use futures_util::{StreamExt, TryStreamExt};
12use gcp_auth::TokenProvider;
13use objectstore_types::headers;
14use objectstore_types::metadata::{ExpirationPolicy, Metadata};
15use objectstore_types::range::{ByteRange, ContentRange};
16use objectstore_types::time::{Rfc3339Timestamp, Timestamp};
17use reqwest::header::{HeaderMap, HeaderName};
18use reqwest::{Body, IntoUrl, Method, RequestBuilder, StatusCode, Url, header, multipart};
19use serde::{Deserialize, Serialize};
20
21use crate::backend::common::{
22    self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, MetadataResponse,
23    MultipartUploadBackend, PutResponse, SetExpiryResponse,
24};
25use crate::backend::extensions::{ReqwestResultExt, ResponseExt, SendTraced};
26use crate::change_stream::{
27    ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream,
28};
29use crate::error::{Error, ErrorKind, Result, ResultExt as _};
30use crate::gcp_auth::PrefetchingTokenProvider;
31use crate::id::ObjectId;
32use crate::multipart::{
33    AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
34    ListPartsResponse, PartNumber, UploadId, UploadPartResponse,
35};
36use crate::resumable::{BackendToken, Session, UploadProgress};
37use crate::stream::ClientStream;
38
39/// Configuration for [`GcsBackend`].
40///
41/// Stores objects in [Google Cloud Storage]. Authentication uses Application Default Credentials
42/// (ADC), which can be provided via the `GOOGLE_APPLICATION_CREDENTIALS` environment variable or
43/// the GCE/GKE metadata service.
44///
45/// **Note**: The bucket must be pre-created with the following lifecycle policy:
46/// - `daysSinceCustomTime`: 1 day
47/// - `action`: delete
48///
49/// [Google Cloud Storage]: https://cloud.google.com/storage
50///
51/// # Example
52///
53/// ```yaml
54/// storage:
55///   type: gcs
56///   bucket: objectstore-bucket
57/// ```
58#[derive(Debug, Clone, Deserialize, Serialize)]
59pub struct GcsConfig {
60    /// Optional custom GCS endpoint URL.
61    ///
62    /// Useful for testing with emulators. If `None`, uses the default GCS endpoint.
63    ///
64    /// # Default
65    ///
66    /// `None` (uses default GCS endpoint)
67    ///
68    /// # Environment Variables
69    ///
70    /// - `OS__STORAGE__TYPE=gcs`
71    /// - `OS__STORAGE__ENDPOINT=http://localhost:9000` (optional)
72    pub endpoint: Option<String>,
73
74    /// GCS bucket name.
75    ///
76    /// The bucket must exist before starting the server.
77    ///
78    /// # Environment Variables
79    ///
80    /// - `OS__STORAGE__BUCKET=my-gcs-bucket`
81    pub bucket: String,
82
83    /// Reports what this backend stores, for per-usecase cost attribution.
84    ///
85    /// # Default
86    ///
87    /// `None`, which disables reporting for this backend.
88    ///
89    /// # Environment Variables
90    ///
91    /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=gcs_objectstore`
92    /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional)
93    #[serde(default, skip_serializing_if = "Option::is_none")]
94    pub cogs: Option<CostTrackerStreamConfig>,
95}
96
97/// Response header carrying the size GCS stored, in bytes.
98///
99/// Differs from `Content-Length`, which describes the transfer, once an object is
100/// content-encoded.
101const STORED_CONTENT_LENGTH: &str = "x-goog-stored-content-length";
102
103/// Reads how many payload bytes GCS stored for an object, consuming the response.
104///
105/// Prefers the [`STORED_CONTENT_LENGTH`] response header. If it is missing or malformed, falls back
106/// to the `size` of the [`GcsObject`] in the response body. Returns `None` if neither source yields
107/// a size.
108///
109/// Either way the body is read to the end, so reqwest can return the connection to its pool.
110async fn read_stored_content_length(response: reqwest::Response) -> Option<u64> {
111    let header = response
112        .headers()
113        .get(STORED_CONTENT_LENGTH)
114        .and_then(|value| value.to_str().ok())
115        .and_then(|value| value.parse().ok());
116
117    if let Some(size) = header {
118        response.drain_body().await;
119        return Some(size);
120    }
121
122    response.json::<GcsObject>().await.ok()?.size?.parse().ok()
123}
124
125/// Default endpoint used to access the GCS JSON API.
126const DEFAULT_ENDPOINT: &str = "https://storage.googleapis.com";
127/// Permission scopes required for accessing GCS.
128const TOKEN_SCOPES: &[&str] = &["https://www.googleapis.com/auth/devstorage.read_write"];
129/// How many times to retry failed operations.
130const REQUEST_RETRY_COUNT: usize = 2;
131
132/// Prefix for our built-in metadata stored in GCS metadata field
133const BUILTIN_META_PREFIX: &str = "x-sn-";
134/// Prefix for user custom metadata stored in GCS metadata field
135const CUSTOM_META_PREFIX: &str = "x-snme-";
136
137/// GCS object resource.
138///
139/// This is the representation of the object resource in GCS JSON API without its payload. Where no
140/// dedicated fields are available, we encode both built-in and custom metadata in the `metadata`
141/// field.
142#[derive(Debug, Serialize, Deserialize)]
143#[serde(rename_all = "camelCase")]
144struct GcsObject {
145    /// Content-Type of the object data. If an object is stored without a Content-Type, it is served
146    /// as application/octet-stream.
147    pub content_type: Cow<'static, str>,
148
149    /// Content encoding, used to store [`Metadata::compression`].
150    #[serde(default, skip_serializing_if = "Option::is_none")]
151    pub content_encoding: Option<String>,
152
153    /// Custom time stamp used for time-based expiration.
154    #[serde(default, skip_serializing_if = "Option::is_none")]
155    pub custom_time: Option<Rfc3339Timestamp>,
156
157    /// The `Content-Length` of the data in bytes. GCS returns this as a string.
158    ///
159    /// GCS sets this in metadata responses. We can use it to know the size of an object
160    /// without having to stream it.
161    pub size: Option<String>,
162
163    /// Timestamp of when this object was created.
164    #[serde(default, skip_serializing_if = "Option::is_none")]
165    pub time_created: Option<Rfc3339Timestamp>,
166
167    /// User-provided metadata, including our built-in metadata.
168    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
169    pub metadata: BTreeMap<GcsMetaKey, String>,
170
171    /// Version of the object's contents.
172    #[serde(skip_serializing)]
173    pub generation: String,
174
175    /// Version of the object's metadata.
176    #[serde(skip_serializing)]
177    pub metageneration: String,
178}
179
180impl GcsObject {
181    /// Bytes this object's custom metadata occupies, keys included.
182    ///
183    /// GCS stores metadata alongside the payload, so it counts toward an object's size.
184    fn metadata_size(&self) -> u64 {
185        self.metadata
186            .iter()
187            .filter(|(key, _)| !matches!(key, GcsMetaKey::EmulatorIgnored))
188            .map(|(key, value)| key.to_string().len() as u64 + value.len() as u64)
189            .sum()
190    }
191
192    /// Returns `true` if the object is expired at the given access time.
193    pub fn is_expired(&self, access_time: Timestamp) -> bool {
194        match self.custom_time {
195            Some(expires_at) => access_time > expires_at.into_inner(),
196            None => false,
197        }
198    }
199
200    /// Returns the generation and metageneration of this object.
201    pub fn generations(&self) -> GcsGenerations<'_> {
202        (&self.generation, &self.metageneration)
203    }
204
205    /// Converts our Metadata type to GCS JSON object metadata.
206    pub fn from_metadata(metadata: &Metadata) -> Self {
207        let mut gcs_object = GcsObject {
208            content_type: metadata.content_type.clone(),
209            size: metadata.size.map(|size| size.to_string()),
210            content_encoding: None,
211            custom_time: None,
212            time_created: metadata.time_created.map(Timestamp::as_rfc3339),
213            metadata: BTreeMap::new(),
214            generation: String::new(),
215            metageneration: String::new(),
216        };
217
218        // For time-based expiration, set the `customTime` field. The bucket must have a
219        // `daysSinceCustomTime` lifecycle rule configured to delete objects with this field set.
220        // This rule automatically skips objects without `customTime` set.
221        gcs_object.custom_time = metadata.time_expires.map(Timestamp::as_rfc3339);
222
223        if let Some(compression) = metadata.compression {
224            gcs_object.content_encoding = Some(compression.to_string());
225        }
226
227        if metadata.expiration_policy != ExpirationPolicy::default() {
228            gcs_object.metadata.insert(
229                GcsMetaKey::Expiration,
230                metadata.expiration_policy.to_string(),
231            );
232        }
233
234        // Free-form strings are stored escaped, even though this JSON representation could carry
235        // them verbatim. See `insert_gcs_meta_header` for why, and why both writers must agree.
236        if let Some(origin) = &metadata.origin {
237            gcs_object.metadata.insert(
238                GcsMetaKey::Origin,
239                headers::encode_header_str(origin).into(),
240            );
241        }
242
243        if let Some(filename) = &metadata.filename {
244            gcs_object.metadata.insert(
245                GcsMetaKey::Filename,
246                headers::encode_header_str(filename).into(),
247            );
248        }
249
250        for (key, value) in &metadata.custom {
251            gcs_object.metadata.insert(
252                GcsMetaKey::Custom(key.clone()),
253                headers::encode_header_str(value).into(),
254            );
255        }
256
257        gcs_object
258    }
259
260    /// Converts GCS JSON object metadata to our Metadata type.
261    pub fn into_metadata(mut self) -> Result<Metadata> {
262        // Remove ignored metadata keys that are set by the GCS emulator.
263        self.metadata.remove(&GcsMetaKey::EmulatorIgnored);
264
265        let expiration_policy = self
266            .metadata
267            .remove(&GcsMetaKey::Expiration)
268            .map(|s| s.parse())
269            .transpose()
270            .context(ErrorKind::CorruptData, "decoding GCS expiration policy")?
271            .unwrap_or_default();
272
273        let origin = self
274            .metadata
275            .remove(&GcsMetaKey::Origin)
276            .map(|value| decode_gcs_meta_value(&value, "decoding GCS origin metadata"))
277            .transpose()?;
278        let filename = self
279            .metadata
280            .remove(&GcsMetaKey::Filename)
281            .map(|value| decode_gcs_meta_value(&value, "decoding GCS filename metadata"))
282            .transpose()?;
283
284        let content_type = self.content_type;
285        let compression = self
286            .content_encoding
287            .map(|s| s.parse())
288            .transpose()
289            .context(ErrorKind::CorruptData, "decoding GCS compression")?;
290        let size = self
291            .size
292            .map(|size| size.parse())
293            .transpose()
294            .context(ErrorKind::CorruptData, "decoding GCS object size")?;
295        let time_created = self.time_created.map(Rfc3339Timestamp::into_inner);
296
297        // At this point, all built-in metadata should have been removed from self.metadata.
298        let mut custom = BTreeMap::new();
299        for (key, value) in self.metadata {
300            if let GcsMetaKey::Custom(custom_key) = key {
301                custom.insert(
302                    custom_key,
303                    decode_gcs_meta_value(&value, "decoding GCS custom metadata")?,
304                );
305            } else {
306                return Err(Error::new(
307                    ErrorKind::CorruptData,
308                    format!("unexpected GCS metadata key: {key}"),
309                ));
310            }
311        }
312
313        Ok(Metadata {
314            content_type,
315            expiration_policy,
316            compression,
317            origin,
318            filename,
319            size,
320            custom,
321            time_created,
322            time_expires: self.custom_time.map(Rfc3339Timestamp::into_inner),
323        })
324    }
325}
326
327/// The object and metadata versions of a GCS object, used for conditional updates.
328type GcsGenerations<'a> = (&'a str, &'a str);
329
330/// Key for [`GcsObject::metadata`].
331#[derive(Clone, Debug, PartialEq, Eq, Ord, PartialOrd)]
332enum GcsMetaKey {
333    /// Built-in metadata key for [`Metadata::expiration_policy`].
334    Expiration,
335    /// Built-in metadata key for [`Metadata::origin`].
336    Origin,
337    /// Built-in metadata key for [`Metadata::filename`].
338    Filename,
339    /// Ignored metadata set by the GCS emulator.
340    EmulatorIgnored,
341    /// User-defined custom metadata key.
342    Custom(String),
343}
344
345impl std::str::FromStr for GcsMetaKey {
346    type Err = anyhow::Error;
347
348    fn from_str(s: &str) -> Result<Self, Self::Err> {
349        if s.starts_with("x_emulator_") || s.starts_with("x_testbench_") {
350            return Ok(GcsMetaKey::EmulatorIgnored);
351        }
352
353        Ok(match s.strip_prefix(BUILTIN_META_PREFIX) {
354            Some("expiration") => GcsMetaKey::Expiration,
355            Some("origin") => GcsMetaKey::Origin,
356            Some("filename") => GcsMetaKey::Filename,
357            Some(unknown) => anyhow::bail!("unknown builtin metadata key: {unknown}"),
358            None => match s.strip_prefix(CUSTOM_META_PREFIX) {
359                Some(key) => GcsMetaKey::Custom(key.to_string()),
360                None => anyhow::bail!("invalid GCS metadata key format: {s}"),
361            },
362        })
363    }
364}
365
366impl fmt::Display for GcsMetaKey {
367    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
368        match self {
369            Self::Expiration => write!(f, "{BUILTIN_META_PREFIX}expiration"),
370            Self::Origin => write!(f, "{BUILTIN_META_PREFIX}origin"),
371            Self::Filename => write!(f, "{BUILTIN_META_PREFIX}filename"),
372            Self::EmulatorIgnored => unreachable!("do not serialize emulator metadata"),
373            Self::Custom(key) => write!(f, "{CUSTOM_META_PREFIX}{key}"),
374        }
375    }
376}
377
378impl<'de> Deserialize<'de> for GcsMetaKey {
379    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
380    where
381        D: serde::Deserializer<'de>,
382    {
383        let s = Cow::<'de, str>::deserialize(deserializer)?;
384        s.parse().map_err(serde::de::Error::custom)
385    }
386}
387
388impl Serialize for GcsMetaKey {
389    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
390    where
391        S: serde::Serializer,
392    {
393        serializer.collect_str(self)
394    }
395}
396
397/// Builds HTTP headers that encode `metadata` for a GCS XML API request.
398fn metadata_to_gcs_headers(metadata: &Metadata) -> Result<header::HeaderMap> {
399    let mut headers = header::HeaderMap::new();
400
401    if let Some(custom_time) = metadata.time_expires {
402        let formatted = custom_time.as_rfc3339();
403        headers.insert(
404            HeaderName::from_static("x-goog-custom-time"),
405            formatted
406                .to_string()
407                .parse()
408                .context(ErrorKind::Internal, "encoding GCS custom-time header")?,
409        );
410    }
411
412    if let Some(compression) = metadata.compression {
413        headers.insert(
414            header::CONTENT_ENCODING,
415            compression
416                .to_string()
417                .parse()
418                .context(ErrorKind::Internal, "encoding GCS content-encoding header")?,
419        );
420    }
421
422    if metadata.expiration_policy != ExpirationPolicy::default() {
423        insert_gcs_meta_header(
424            &mut headers,
425            &GcsMetaKey::Expiration,
426            &metadata.expiration_policy.to_string(),
427        )?;
428    }
429
430    if let Some(origin) = &metadata.origin {
431        insert_gcs_meta_header(&mut headers, &GcsMetaKey::Origin, origin)?;
432    }
433
434    if let Some(filename) = &metadata.filename {
435        insert_gcs_meta_header(&mut headers, &GcsMetaKey::Filename, filename)?;
436    }
437
438    for (key, value) in &metadata.custom {
439        insert_gcs_meta_header(&mut headers, &GcsMetaKey::Custom(key.clone()), value)?;
440    }
441
442    Ok(headers)
443}
444
445/// Decodes a stored GCS metadata value into its logical string.
446fn decode_gcs_meta_value(value: &str, context: &'static str) -> Result<String> {
447    headers::decode_header_str(value).context(ErrorKind::CorruptData, context)
448}
449
450/// Inserts a single `x-goog-meta-*` header, escaping the value for transport.
451///
452/// Google: "you should generally avoid non-ascii characters, because they are not permitted in
453/// HTTP headers, which the XML API uses" ([docs]). Real GCS does preserve raw UTF-8 here, but that
454/// is undocumented, and it still drops leading whitespace and turns invalid UTF-8 into `U+FFFD`.
455///
456/// [`GcsObject::from_metadata`] escapes the same values on the JSON path: reads always come back
457/// through the JSON API and cannot tell which writer produced an object, so both must agree.
458///
459/// [docs]: https://docs.cloud.google.com/storage/docs/metadata
460fn insert_gcs_meta_header(
461    headers: &mut header::HeaderMap,
462    key: &GcsMetaKey,
463    value: &str,
464) -> Result<()> {
465    let header_name = format!("x-goog-meta-{key}");
466    headers.insert(
467        HeaderName::try_from(&header_name).context(
468            ErrorKind::Internal,
469            format!("encoding GCS metadata header {header_name}"),
470        )?,
471        headers::encode_header_value(value),
472    );
473    Ok(())
474}
475
476/// Special status code returned by GCS when a Resumable Upload is canceled successfully or when
477/// making other requests to a session that was recently canceled.
478const CLIENT_CLOSED_REQUEST_STATUS: u16 = 499;
479
480enum GcsUploadProgress {
481    Incomplete(u64),
482    Complete(GcsObject),
483}
484
485impl From<GcsUploadProgress> for UploadProgress {
486    fn from(value: GcsUploadProgress) -> Self {
487        match value {
488            GcsUploadProgress::Incomplete(offset) => UploadProgress::Incomplete { offset },
489            GcsUploadProgress::Complete(_) => UploadProgress::Complete,
490        }
491    }
492}
493
494/// Returns `true` if the error is a transient backend failure worth retrying.
495fn error_is_retryable(error: &Error) -> bool {
496    matches!(
497        error.kind(),
498        ErrorKind::BackendRateLimited | ErrorKind::BackendTimeout | ErrorKind::BackendUnavailable
499    )
500}
501
502/// GCS JSON API backend for long-term storage of large objects.
503pub struct GcsBackend {
504    client: reqwest::Client,
505    endpoint: Url,
506    bucket: String,
507    token_provider: Option<PrefetchingTokenProvider>,
508
509    change_stream: Arc<dyn ChangeStream>,
510}
511
512impl GcsBackend {
513    /// Creates an authenticated GCS JSON API backend bound to the bucket in `config`.
514    pub async fn new(config: GcsConfig, streams: &ChangeStreamFactory) -> anyhow::Result<Self> {
515        let GcsConfig {
516            endpoint,
517            bucket,
518            cogs,
519        } = config;
520        let change_stream = streams.build(cogs.as_ref());
521
522        let token_provider = if endpoint.is_none() {
523            Some(PrefetchingTokenProvider::gcp_auth(TOKEN_SCOPES).await?)
524        } else {
525            None
526        };
527
528        let endpoint_str = endpoint.as_deref().unwrap_or(DEFAULT_ENDPOINT);
529
530        Ok(Self {
531            client: common::reqwest_client(),
532            endpoint: endpoint_str
533                .parse()
534                .map_err(|e| anyhow::Error::new(e).context("invalid GCS endpoint URL"))?,
535            bucket,
536            token_provider,
537            change_stream,
538        })
539    }
540
541    /// Formats the GCS object (metadata) URL for the given key.
542    fn object_url(&self, id: &ObjectId) -> Result<Url> {
543        let mut url = self.endpoint.clone();
544
545        let path = id.as_storage_path().to_string();
546        url.path_segments_mut()
547            .map_err(|()| {
548                Error::new(
549                    ErrorKind::Internal,
550                    format!("building GCS object URL from {}", self.endpoint),
551                )
552            })?
553            .extend(&["storage", "v1", "b", &self.bucket, "o", &path]);
554
555        Ok(url)
556    }
557
558    /// Formats the GCS upload URL for the given upload type.
559    fn upload_url(&self, id: &ObjectId, upload_type: &str) -> Result<Url> {
560        let mut url = self.endpoint.clone();
561
562        url.path_segments_mut()
563            .map_err(|()| {
564                Error::new(
565                    ErrorKind::Internal,
566                    format!("building GCS object URL from {}", self.endpoint),
567                )
568            })?
569            .extend(&["upload", "storage", "v1", "b", &self.bucket, "o"]);
570
571        url.query_pairs_mut()
572            .append_pair("uploadType", upload_type)
573            .append_pair("name", &id.as_storage_path().to_string());
574
575        Ok(url)
576    }
577
578    /// Formats a GCS XML API URL for the given object.
579    ///
580    /// Unlike [`object_url`](Self::object_url) (JSON API at
581    /// `/storage/v1/b/{bucket}/o/{name}`), this produces
582    /// `/{bucket}/{path_segments}` for the S3-compatible XML API used by
583    /// multipart uploads.
584    fn xml_object_url(&self, id: &ObjectId) -> Result<Url> {
585        let mut url = self.endpoint.clone();
586        {
587            let mut segments = url.path_segments_mut().map_err(|()| {
588                Error::new(
589                    ErrorKind::Internal,
590                    format!("building GCS object URL from {}", self.endpoint),
591                )
592            })?;
593            segments.push(&self.bucket);
594            for part in id.as_storage_path().to_string().split('/') {
595                segments.push(part);
596            }
597        }
598        Ok(url)
599    }
600
601    /// Creates a request builder with the appropriate authentication.
602    async fn request(&self, method: Method, url: impl IntoUrl) -> Result<RequestBuilder> {
603        let mut builder = self.client.request(method, url);
604        if let Some(provider) = &self.token_provider {
605            let token = provider.token(TOKEN_SCOPES).await.context(
606                ErrorKind::BackendFailure,
607                "getting GCS authentication token",
608            )?;
609            builder = builder.bearer_auth(token.as_str());
610        }
611        Ok(builder)
612    }
613
614    /// Retries a GCS request on transient errors.
615    async fn with_retry<T, F>(&self, action: &'static str, f: impl Fn() -> F) -> Result<T>
616    where
617        F: Future<Output = Result<T>> + Send,
618    {
619        let mut retry_count = 0usize;
620        loop {
621            match f().await {
622                Ok(res) => return Ok(res),
623                Err(ref e) if retry_count < REQUEST_RETRY_COUNT && error_is_retryable(e) => {
624                    retry_count += 1;
625                    objectstore_metrics::count!("gcs.retries", action = action);
626                    objectstore_log::warn!(!!e, retry_count, action, "Retrying request");
627                }
628                Err(e) => {
629                    objectstore_metrics::count!("gcs.failures", action = action);
630                    return Err(e);
631                }
632            }
633        }
634    }
635
636    /// Fetches GCS object metadata without modifying the object.
637    #[tracing::instrument(level = "debug", fields(%object_url), skip(self))]
638    async fn get_gcs_metadata(
639        &self,
640        object_url: &Url,
641        access_time: Timestamp,
642    ) -> Result<Option<GcsObject>> {
643        let metadata_opt = self
644            .with_retry("get_metadata", || async {
645                let resp = self
646                    .request(Method::GET, object_url.clone())
647                    .await?
648                    .send_traced()
649                    .await
650                    .reqwest_context("getting GCS object metadata")?;
651
652                if resp.status() == StatusCode::NOT_FOUND {
653                    resp.drain_body().await;
654                    return Ok(None);
655                }
656
657                let metadata: GcsObject = resp
658                    .check_error("getting GCS object metadata")
659                    .await?
660                    .json()
661                    .await
662                    .reqwest_context("getting GCS object metadata")?;
663
664                Ok(Some(metadata))
665            })
666            .await?;
667
668        let Some(gcs_metadata) = metadata_opt else {
669            objectstore_log::debug!("Object not found");
670            return Ok(None);
671        };
672
673        // Filter already expired objects but leave them to garbage collection
674        if gcs_metadata.is_expired(access_time) {
675            objectstore_log::debug!("Object found but past expiry");
676            return Ok(None);
677        }
678
679        Ok(Some(gcs_metadata))
680    }
681
682    /// Moves an object's `customTime` and optionally updates its metadata in the same PATCH.
683    ///
684    /// Lifecycle expiry is anchored to `customTime`; TTL updates accompany that deadline.
685    ///
686    /// Returns the expiry update outcome, distinguishing absence from a conflict.
687    #[tracing::instrument(level = "debug", fields(%object_url), skip(self))]
688    async fn update_custom_time(
689        &self,
690        object_url: Url,
691        custom_time: Timestamp,
692        generations: GcsGenerations<'_>,
693        metadata: Option<BTreeMap<GcsMetaKey, String>>,
694    ) -> Result<SetExpiryResponse> {
695        #[derive(Debug, Serialize)]
696        #[serde(rename_all = "camelCase")]
697        struct CustomTimeRequest {
698            custom_time: Rfc3339Timestamp,
699            #[serde(skip_serializing_if = "Option::is_none")]
700            metadata: Option<BTreeMap<GcsMetaKey, String>>,
701        }
702
703        let mut object_url = object_url;
704        object_url
705            .query_pairs_mut()
706            .append_pair("ifGenerationMatch", generations.0)
707            .append_pair("ifMetagenerationMatch", generations.1);
708
709        let request = CustomTimeRequest {
710            custom_time: custom_time.as_rfc3339(),
711            metadata,
712        };
713
714        self.with_retry("update_custom_time", || async {
715            let response = self
716                .request(Method::PATCH, object_url.clone())
717                .await?
718                .json(&request)
719                .send_traced()
720                .await
721                .reqwest_context("updating GCS custom time")?;
722
723            let outcome = match response.status() {
724                StatusCode::NOT_FOUND => Some(SetExpiryResponse::NotFound),
725                StatusCode::PRECONDITION_FAILED => Some(SetExpiryResponse::Rejected),
726                _ => None,
727            };
728            if let Some(outcome) = outcome {
729                response.drain_body().await;
730                return Ok(outcome);
731            }
732
733            response
734                .check_error("updating GCS custom time")
735                .await?
736                .drain_body()
737                .await;
738
739            Ok(SetExpiryResponse::Satisfied(custom_time))
740        })
741        .await
742    }
743
744    /// Reports an object write to the [`ChangeStream`].
745    fn report_object_write(
746        &self,
747        id: &ObjectId,
748        stored_size: Option<u64>,
749        metadata_size: u64,
750        expires_at: Option<Timestamp>,
751    ) {
752        match stored_size {
753            Some(stored_size) => {
754                self.change_stream
755                    .write(id, stored_size + metadata_size, expires_at)
756            }
757            None => {
758                objectstore_metrics::count!("change_stream.unreported", reason = "no_stored_size")
759            }
760        }
761    }
762}
763
764impl fmt::Debug for GcsBackend {
765    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
766        f.debug_struct("GcsJsonApi")
767            .field("endpoint", &self.endpoint)
768            .field("bucket", &self.bucket)
769            .finish_non_exhaustive()
770    }
771}
772
773/// Converts GCS's inclusive `Range: bytes=0-N` acknowledgement into the next offset.
774fn range_header_to_offset(value: &str, upload_length: NonZeroU64) -> Result<u64> {
775    let end = value.strip_prefix("bytes=0-").ok_or_else(|| {
776        Error::new(
777            ErrorKind::Internal,
778            "malformed GCS Range header for resumable progress",
779        )
780    })?;
781    let end = end.parse::<u64>().context(
782        ErrorKind::Internal,
783        "invalid GCS Range header for resumable progress",
784    )?;
785
786    let offset = end
787        .checked_add(1)
788        .ok_or_else(|| Error::new(ErrorKind::Internal, "GCS offset overflows u64"))?;
789    if offset > upload_length.get() {
790        return Err(Error::new(
791            ErrorKind::Internal,
792            "GCS offset exceeds upload length",
793        ));
794    }
795    Ok(offset)
796}
797
798/// Interprets a Resumable Upload response, shared by chunk writes and status queries.
799///
800/// Returns the progress GCS reported, plus the completed object when this is the response that
801/// finished the upload.
802async fn range_response_to_upload_progress(
803    upload_length: NonZeroU64,
804    response: reqwest::Response,
805) -> Result<GcsUploadProgress> {
806    let status = response.status();
807
808    match status {
809        StatusCode::NOT_FOUND => {
810            response.drain_body().await;
811            return Err(ErrorKind::UnknownUploadSession.into());
812        }
813        status if status == StatusCode::GONE || status.as_u16() == CLIENT_CLOSED_REQUEST_STATUS => {
814            response.drain_body().await;
815            return Err(ErrorKind::UploadSessionGone.into());
816        }
817        _ => {}
818    }
819
820    let response = response
821        .check_error("processing a GCS resumable upload response")
822        .await?;
823
824    match status {
825        StatusCode::OK | StatusCode::CREATED => {
826            let body = response
827                .bytes()
828                .await
829                .reqwest_context("reading a completed GCS resumable upload response")?;
830            let object = serde_json::from_slice::<GcsObject>(&body).context(
831                ErrorKind::CorruptData,
832                "parsing a completed GCS resumable upload response",
833            )?;
834            Ok(GcsUploadProgress::Complete(object))
835        }
836        // GCS calls it "308 Resume Incomplete"
837        StatusCode::PERMANENT_REDIRECT => {
838            let offset = match response.headers().get(header::RANGE) {
839                Some(range) => {
840                    let range = range.to_str().map_err(|_| {
841                        Error::new(
842                            ErrorKind::BackendFailure,
843                            "invalid GCS resumable upload Range header",
844                        )
845                    })?;
846                    range_header_to_offset(range, upload_length)?
847                }
848                // GCS omits this header while it holds nothing
849                None => 0,
850            };
851            response.drain_body().await;
852            Ok(GcsUploadProgress::Incomplete(offset))
853        }
854        _ => {
855            response.drain_body().await;
856            Err(Error::new(
857                ErrorKind::BackendFailure,
858                format!("unexpected GCS resumable upload status {status}"),
859            ))
860        }
861    }
862}
863
864#[async_trait::async_trait]
865impl Backend for GcsBackend {
866    fn name(&self) -> &'static str {
867        "gcs"
868    }
869
870    fn upload_granularity(&self) -> u64 {
871        256 * 1024
872    }
873
874    fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
875        Ok(self)
876    }
877
878    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
879    async fn put_object(
880        &self,
881        id: &ObjectId,
882        metadata: &Metadata,
883        stream: ClientStream,
884        _access_time: Timestamp,
885    ) -> Result<PutResponse> {
886        objectstore_log::debug!("Writing to GCS backend");
887        let gcs_metadata = GcsObject::from_metadata(metadata);
888
889        // NB: Ensure the order of these fields and that a content-type is attached to them. Both
890        // are required by the GCS API.
891        let metadata_json = serde_json::to_string(&gcs_metadata)
892            .context(ErrorKind::Internal, "encoding GCS upload metadata")?;
893
894        let content_type = metadata
895            .content_type
896            .parse()
897            .context(ErrorKind::InvalidMetadata, "encoding GCS content type")?;
898
899        let multipart = multipart::Form::new()
900            .part(
901                "metadata",
902                multipart::Part::text(metadata_json)
903                    .mime_str("application/json")
904                    .expect("application/json is a valid mime type"),
905            )
906            .part(
907                "media",
908                multipart::Part::stream(Body::wrap_stream(stream.boxed()))
909                    // Ensure the content type goes into headers unchanged including case
910                    // sensitivity, as required by GCS. This means we cannot use `mime_str`.
911                    .headers(HeaderMap::from_iter([(header::CONTENT_TYPE, content_type)])),
912            );
913
914        // GCS requires a multipart/related request. Its body looks identical to
915        // multipart/form-data, but the Content-Type header is different. Hence, we have to manually
916        // set the header *after* writing the multipart form into the request.
917        let content_type = format!("multipart/related; boundary={}", multipart.boundary());
918
919        let response = self
920            .request(Method::POST, self.upload_url(id, "multipart")?)
921            .await?
922            .multipart(multipart)
923            .header(header::CONTENT_TYPE, content_type)
924            .send_traced()
925            .await
926            .check_error("uploading a GCS object")
927            .await?;
928
929        let stored_size = read_stored_content_length(response).await;
930        self.report_object_write(
931            id,
932            stored_size,
933            gcs_metadata.metadata_size(),
934            metadata.time_expires,
935        );
936
937        Ok(())
938    }
939
940    #[tracing::instrument(level = "debug", skip(self))]
941    async fn get_object(
942        &self,
943        id: &ObjectId,
944        access_time: Timestamp,
945        range: Option<ByteRange>,
946    ) -> Result<GetResponse> {
947        objectstore_log::debug!("Reading from GCS backend");
948        let object_url = self.object_url(id)?;
949        let mut generation_retry_count = 0usize;
950
951        loop {
952            let Some(gcs_metadata) = self.get_gcs_metadata(&object_url, access_time).await? else {
953                return Ok(None);
954            };
955
956            let mut download_url = object_url.clone();
957            download_url
958                .query_pairs_mut()
959                .append_pair("alt", "media")
960                .append_pair("ifGenerationMatch", &gcs_metadata.generation);
961
962            let payload_response = self
963                .with_retry("get_payload", || async {
964                    let mut req = self.request(Method::GET, download_url.clone()).await?;
965                    if let Some(r) = range {
966                        req = req.header(header::RANGE, r.to_header_value());
967                    }
968
969                    let resp = req
970                        .send_traced()
971                        .await
972                        .reqwest_context("getting a GCS object payload")?;
973
974                    if resp.status() == StatusCode::PRECONDITION_FAILED {
975                        resp.drain_body().await;
976                        return Ok(None);
977                    }
978
979                    if resp.status() == StatusCode::RANGE_NOT_SATISFIABLE {
980                        let raw = resp
981                            .headers()
982                            .get(header::CONTENT_RANGE)
983                            .and_then(|v| v.to_str().ok());
984                        let total = raw.and_then(ContentRange::parse_unsatisfiable_total);
985                        let err = match total {
986                            Some(total) => ErrorKind::RangeNotSatisfiable { total }.into(),
987                            None => Error::new(
988                                ErrorKind::BackendFailure,
989                                "invalid GCS 416 Content-Range",
990                            ),
991                        };
992                        resp.drain_body().await;
993                        return Err(err);
994                    }
995
996                    resp.check_error("getting a GCS object payload")
997                        .await
998                        .map(Some)
999                })
1000                .await?;
1001
1002            let Some(payload_response) = payload_response else {
1003                if generation_retry_count >= REQUEST_RETRY_COUNT {
1004                    objectstore_metrics::count!("gcs.failures", action = "get_object");
1005                    return Err(Error::new(
1006                        ErrorKind::BackendFailure,
1007                        "GCS object kept changing while being read",
1008                    ));
1009                }
1010
1011                generation_retry_count += 1;
1012                objectstore_metrics::count!("gcs.retries", action = "get_object");
1013                continue;
1014            };
1015
1016            let content_range = if payload_response.status() == StatusCode::PARTIAL_CONTENT {
1017                Some(
1018                    payload_response
1019                        .headers()
1020                        .get(header::CONTENT_RANGE)
1021                        .and_then(|v| v.to_str().ok())
1022                        .and_then(|s| s.parse::<ContentRange>().ok())
1023                        .ok_or_else(|| {
1024                            Error::new(ErrorKind::BackendFailure, "missing GCS 206 Content-Range")
1025                        })?,
1026                )
1027            } else {
1028                None
1029            };
1030
1031            let stream = payload_response
1032                .bytes_stream()
1033                .map_err(io::Error::other)
1034                .boxed();
1035
1036            let metadata = gcs_metadata.into_metadata()?;
1037            return Ok(Some((metadata, content_range, stream)));
1038        }
1039    }
1040
1041    #[tracing::instrument(level = "debug", skip(self))]
1042    async fn get_metadata(
1043        &self,
1044        id: &ObjectId,
1045        access_time: Timestamp,
1046    ) -> Result<MetadataResponse> {
1047        objectstore_log::debug!("Reading metadata from GCS backend");
1048        let object_url = self.object_url(id)?;
1049        match self.get_gcs_metadata(&object_url, access_time).await? {
1050            Some(gcs_metadata) => Ok(Some(gcs_metadata.into_metadata()?)),
1051            None => Ok(None),
1052        }
1053    }
1054
1055    #[tracing::instrument(level = "debug", skip(self))]
1056    async fn set_expiry(
1057        &self,
1058        id: &ObjectId,
1059        target: ExpiryUpdate,
1060        access_time: Timestamp,
1061    ) -> Result<SetExpiryResponse> {
1062        let object_url = self.object_url(id)?;
1063        let Some(object) = self.get_gcs_metadata(&object_url, access_time).await? else {
1064            return Ok(SetExpiryResponse::NotFound);
1065        };
1066        let Some(current_expiry) = object.custom_time else {
1067            return Ok(SetExpiryResponse::Rejected);
1068        };
1069        let time_created = object.time_created.map(Rfc3339Timestamp::into_inner);
1070        let Some(expire_at) = target.resolve(time_created, access_time)? else {
1071            return Ok(SetExpiryResponse::Rejected);
1072        };
1073        if current_expiry.into_inner() >= expire_at {
1074            return Ok(SetExpiryResponse::Satisfied(expire_at)); // already satisfied
1075        }
1076
1077        let policy = object
1078            .metadata
1079            .get(&GcsMetaKey::Expiration)
1080            .map(|value| value.parse())
1081            .transpose()
1082            .context(ErrorKind::CorruptData, "decoding GCS expiration policy")?
1083            .unwrap_or_default();
1084        let updated_policy = common::extended_expiration_policy(
1085            policy,
1086            time_created,
1087            current_expiry.into_inner(),
1088            expire_at,
1089        )?;
1090        let metadata = if updated_policy != policy {
1091            Some(BTreeMap::from([(
1092                GcsMetaKey::Expiration,
1093                updated_policy.to_string(),
1094            )]))
1095        } else {
1096            None
1097        };
1098
1099        let outcome = self
1100            .update_custom_time(object_url, expire_at, object.generations(), metadata)
1101            .await?;
1102        if matches!(outcome, SetExpiryResponse::Satisfied(_)) {
1103            self.change_stream.update(id, Some(expire_at));
1104        }
1105
1106        Ok(outcome)
1107    }
1108
1109    #[tracing::instrument(level = "debug", skip(self))]
1110    async fn delete_object(
1111        &self,
1112        id: &ObjectId,
1113        _access_time: Timestamp,
1114    ) -> Result<DeleteResponse> {
1115        objectstore_log::debug!("Deleting from GCS backend");
1116        let object_url = self.object_url(id)?;
1117
1118        let deleted = self
1119            .with_retry("delete", || async {
1120                let resp = self
1121                    .request(Method::DELETE, object_url.clone())
1122                    .await?
1123                    .send_traced()
1124                    .await
1125                    .reqwest_context("deleting a GCS object")?;
1126
1127                // Do not error for objects that do not exist
1128                if resp.status() == StatusCode::NOT_FOUND {
1129                    resp.drain_body().await;
1130                    return Ok(false);
1131                }
1132
1133                resp.check_error("deleting a GCS object")
1134                    .await?
1135                    .drain_body()
1136                    .await;
1137
1138                Ok(true)
1139            })
1140            .await?;
1141
1142        if deleted {
1143            self.change_stream.delete(id);
1144        }
1145
1146        Ok(())
1147    }
1148
1149    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1150    async fn create_upload_session(
1151        &self,
1152        id: &ObjectId,
1153        metadata: &Metadata,
1154        upload_length: NonZeroU64,
1155    ) -> Result<Option<BackendToken>> {
1156        objectstore_log::debug!("Creating resumable upload session on GCS backend");
1157        let url = self.upload_url(id, "resumable")?;
1158        let metadata_json = serde_json::to_vec(&GcsObject::from_metadata(metadata)).context(
1159            ErrorKind::Internal,
1160            "serializing GCS resumable upload metadata",
1161        )?;
1162        let content_type = metadata.content_type.clone();
1163
1164        let location = self
1165            .with_retry("create_resumable_upload", || {
1166                let url = url.clone();
1167                let metadata_json = metadata_json.clone();
1168                let content_type = content_type.clone();
1169                async move {
1170                    let response = self
1171                        .request(Method::POST, url)
1172                        .await?
1173                        .header(header::CONTENT_TYPE, "application/json")
1174                        .header("x-upload-content-type", content_type.as_ref())
1175                        .header("x-upload-content-length", upload_length.get())
1176                        .body(metadata_json)
1177                        .send_traced()
1178                        .await
1179                        .check_error("creating a GCS resumable upload")
1180                        .await?;
1181
1182                    if response.status() != StatusCode::OK {
1183                        let status = response.status();
1184                        response.drain_body().await;
1185                        return Err(Error::new(
1186                            ErrorKind::BackendFailure,
1187                            format!("unexpected GCS resumable upload creation status {status}"),
1188                        ));
1189                    }
1190
1191                    let location = response
1192                        .headers()
1193                        .get(header::LOCATION)
1194                        .and_then(|value| value.to_str().ok())
1195                        .map(str::to_owned)
1196                        .ok_or_else(|| {
1197                            Error::new(
1198                                ErrorKind::BackendFailure,
1199                                "missing valid Location header in GCS resumable upload creation response",
1200                            )
1201                        })?;
1202                    response.drain_body().await;
1203                    Ok(location)
1204                }
1205            })
1206            .await?;
1207
1208        let session_uri = Url::parse(&location).map_err(|_| {
1209            Error::new(
1210                ErrorKind::BackendFailure,
1211                "invalid Location URL in GCS resumable upload creation response",
1212            )
1213        })?;
1214        Ok(Some(session_uri.into()))
1215    }
1216
1217    #[tracing::instrument(level = "debug", fields(?session, offset, content_length), skip_all)]
1218    async fn put_chunk(
1219        &self,
1220        session: &Session,
1221        offset: u64,
1222        content_length: u64,
1223        stream: ClientStream,
1224    ) -> Result<UploadProgress> {
1225        objectstore_log::debug!("Uploading resumable chunk to GCS backend");
1226        let session_uri =
1227            Url::parse(&session.backend_token).map_err(|_| ErrorKind::UnknownUploadSession)?;
1228
1229        let end = offset
1230            .checked_add(content_length)
1231            .filter(|end| *end <= session.upload_length.get())
1232            .ok_or(ErrorKind::ChunkExceedsUploadLength {
1233                offset,
1234                content_length,
1235                upload_length: session.upload_length.get(),
1236            })?;
1237        let granularity = self.upload_granularity();
1238        if content_length > 0 && content_length < granularity && end != session.upload_length.get()
1239        {
1240            return Err(ErrorKind::ChunkTooSmall {
1241                chunk_length: content_length,
1242                upload_granularity: granularity,
1243            }
1244            .into());
1245        }
1246
1247        let content_range = match content_length {
1248            // An empty chunk is equivalent to an offset query.
1249            0 => format!("bytes */{}", session.upload_length),
1250            _ => format!("bytes {offset}-{}/{}", end - 1, session.upload_length),
1251        };
1252
1253        let response = self
1254            .request(Method::PUT, session_uri.as_str())
1255            .await?
1256            .header(header::CONTENT_LENGTH, content_length)
1257            .header(header::CONTENT_RANGE, content_range)
1258            .body(Body::wrap_stream(stream))
1259            .send_traced()
1260            .await
1261            .reqwest_context("uploading a GCS resumable chunk")?;
1262
1263        let progress = range_response_to_upload_progress(session.upload_length, response).await?;
1264        if let GcsUploadProgress::Complete(ref object) = progress {
1265            let stored_size = object.size.as_deref().and_then(|size| size.parse().ok());
1266            let expires_at = object.custom_time.map(Rfc3339Timestamp::into_inner);
1267            self.report_object_write(
1268                &session.object_id,
1269                stored_size,
1270                object.metadata_size(),
1271                expires_at,
1272            );
1273        }
1274        Ok(progress.into())
1275    }
1276
1277    #[tracing::instrument(level = "debug", fields(?session), skip_all)]
1278    async fn upload_offset(&self, session: &Session) -> Result<UploadProgress> {
1279        objectstore_log::debug!("Querying resumable upload offset on GCS backend");
1280        let session_uri =
1281            Url::parse(&session.backend_token).map_err(|_| ErrorKind::UnknownUploadSession)?;
1282
1283        self.with_retry("query_resumable_upload", || async {
1284            let response = self
1285                .request(Method::PUT, session_uri.as_str())
1286                .await?
1287                .header(
1288                    header::CONTENT_RANGE,
1289                    format!("bytes */{}", session.upload_length),
1290                )
1291                .send_traced()
1292                .await
1293                .reqwest_context("querying a GCS resumable upload")?;
1294
1295            let progress =
1296                range_response_to_upload_progress(session.upload_length, response).await?;
1297            // The final `put_chunk` may have persisted the object but failed while
1298            // reading its response, so completion observed here must be reported too.
1299            if let GcsUploadProgress::Complete(ref object) = progress {
1300                let stored_size = object.size.as_deref().and_then(|size| size.parse().ok());
1301                let expires_at = object.custom_time.map(Rfc3339Timestamp::into_inner);
1302                self.report_object_write(
1303                    &session.object_id,
1304                    stored_size,
1305                    object.metadata_size(),
1306                    expires_at,
1307                );
1308            }
1309            Ok(progress.into())
1310        })
1311        .await
1312    }
1313
1314    #[tracing::instrument(level = "debug", fields(?session), skip_all)]
1315    async fn cancel_upload(&self, session: &Session) -> Result<()> {
1316        objectstore_log::debug!("Cancelling resumable upload on GCS backend");
1317        let session_uri =
1318            Url::parse(&session.backend_token).map_err(|_| ErrorKind::UnknownUploadSession)?;
1319        self.with_retry("cancel_resumable_upload", || {
1320            let session_uri = session_uri.clone();
1321            async move {
1322                let response = self
1323                    .request(Method::DELETE, session_uri)
1324                    .await?
1325                    .send_traced()
1326                    .await
1327                    .reqwest_context("canceling a GCS resumable upload")?;
1328                match response.status() {
1329                    // Expected status code when canceling a recently created upload.
1330                    status if status.as_u16() == CLIENT_CLOSED_REQUEST_STATUS => {
1331                        response.drain_body().await;
1332                        Ok(())
1333                    }
1334                    // The upload was already canceled or never existed.
1335                    StatusCode::GONE | StatusCode::NOT_FOUND => {
1336                        response.drain_body().await;
1337                        Ok(())
1338                    }
1339                    _ => {
1340                        response
1341                            .check_error("canceling a GCS resumable upload")
1342                            .await?
1343                            .drain_body()
1344                            .await;
1345                        Err(Error::new(
1346                            ErrorKind::BackendFailure,
1347                            "unexpected GCS resumable upload cancellation status",
1348                        ))
1349                    }
1350                }
1351            }
1352        })
1353        .await
1354    }
1355
1356    async fn join(&self) {
1357        flush_change_stream(&self.change_stream).await;
1358    }
1359}
1360
1361#[derive(Debug, Deserialize)]
1362#[serde(rename_all = "PascalCase")]
1363struct XmlInitiateMultipartUploadResponse {
1364    upload_id: String,
1365}
1366
1367impl TryFrom<XmlInitiateMultipartUploadResponse> for InitiateMultipartResponse {
1368    type Error = Error;
1369
1370    fn try_from(r: XmlInitiateMultipartUploadResponse) -> Result<Self> {
1371        Ok(UploadId::new(r.upload_id)?)
1372    }
1373}
1374
1375#[derive(Debug, Deserialize)]
1376#[serde(rename_all = "PascalCase")]
1377struct XmlListPartsResponse {
1378    #[serde(default)]
1379    is_truncated: bool,
1380    next_part_number_marker: Option<PartNumber>,
1381    #[serde(default, rename = "Part")]
1382    parts: Vec<XmlPart>,
1383}
1384
1385impl From<XmlListPartsResponse> for ListPartsResponse {
1386    fn from(xml: XmlListPartsResponse) -> Self {
1387        Self {
1388            parts: xml.parts.into_iter().map(Into::into).collect(),
1389            is_truncated: xml.is_truncated,
1390            next_part_number_marker: xml.next_part_number_marker,
1391        }
1392    }
1393}
1394
1395#[derive(Debug, Deserialize)]
1396#[serde(rename_all = "PascalCase")]
1397struct XmlPart {
1398    part_number: PartNumber,
1399    #[serde(rename = "ETag")]
1400    e_tag: String,
1401    #[serde(with = "humantime_serde")]
1402    last_modified: SystemTime,
1403    size: u64,
1404}
1405
1406impl From<XmlPart> for crate::multipart::Part {
1407    fn from(p: XmlPart) -> Self {
1408        Self {
1409            part_number: p.part_number,
1410            etag: p.e_tag,
1411            last_modified: p.last_modified,
1412            size: p.size,
1413        }
1414    }
1415}
1416
1417#[derive(Debug, Serialize)]
1418#[serde(rename = "CompleteMultipartUpload")]
1419struct XmlCompleteMultipartUpload {
1420    #[serde(rename = "Part")]
1421    parts: Vec<XmlCompletePart>,
1422}
1423
1424impl From<Vec<CompletedPart>> for XmlCompleteMultipartUpload {
1425    fn from(parts: Vec<CompletedPart>) -> Self {
1426        Self {
1427            parts: parts.into_iter().map(Into::into).collect(),
1428        }
1429    }
1430}
1431
1432#[derive(Debug, Serialize)]
1433#[serde(rename_all = "PascalCase")]
1434struct XmlCompletePart {
1435    part_number: PartNumber,
1436    #[serde(rename = "ETag")]
1437    e_tag: String,
1438}
1439
1440impl From<CompletedPart> for XmlCompletePart {
1441    fn from(p: CompletedPart) -> Self {
1442        Self {
1443            part_number: p.part_number,
1444            e_tag: p.etag,
1445        }
1446    }
1447}
1448
1449#[derive(Debug, Deserialize)]
1450#[serde(rename = "Error", rename_all = "PascalCase")]
1451struct XmlError {
1452    code: String,
1453    message: String,
1454}
1455
1456impl From<XmlError> for crate::multipart::CompleteMultipartError {
1457    fn from(e: XmlError) -> Self {
1458        Self {
1459            code: e.code,
1460            message: e.message,
1461        }
1462    }
1463}
1464
1465/// XXX: Any change that affects this implementation should be manually tested against real GCS.
1466/// That's because the fork of [storage-testbench](https://github.com/googleapis/storage-testbench)
1467/// that we test against has an incomplete implementation of the XML multipart API that likely doesn't match GCS's behavior in many cases.
1468#[async_trait::async_trait]
1469impl MultipartUploadBackend for GcsBackend {
1470    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
1471    async fn initiate_multipart(
1472        &self,
1473        id: &ObjectId,
1474        metadata: &Metadata,
1475    ) -> Result<InitiateMultipartResponse> {
1476        objectstore_log::debug!("Initiating multipart upload on GCS backend");
1477        let mut url = self.xml_object_url(id)?;
1478        url.set_query(Some("uploads"));
1479
1480        let mut headers = metadata_to_gcs_headers(metadata)?;
1481        headers.insert(
1482            header::CONTENT_TYPE,
1483            metadata
1484                .content_type
1485                .parse()
1486                .context(ErrorKind::InvalidMetadata, "encoding GCS content type")?,
1487        );
1488        headers.insert(
1489            header::CONTENT_LENGTH,
1490            header::HeaderValue::from_static("0"),
1491        );
1492
1493        self.with_retry("initiate_multipart", || {
1494            let url = url.clone();
1495            let headers = headers.clone();
1496            async move {
1497                let resp = self
1498                    .request(Method::POST, url)
1499                    .await?
1500                    .headers(headers)
1501                    .send_traced()
1502                    .await
1503                    .check_error("initiating a GCS multipart upload")
1504                    .await?;
1505
1506                let body = resp
1507                    .bytes()
1508                    .await
1509                    .reqwest_context("reading GCS initiate-multipart response")?;
1510
1511                let xml: XmlInitiateMultipartUploadResponse =
1512                    quick_xml::de::from_reader(body.as_ref()).context(
1513                        ErrorKind::CorruptData,
1514                        "decoding GCS initiate-multipart response",
1515                    )?;
1516
1517                xml.try_into()
1518            }
1519        })
1520        .await
1521    }
1522
1523    #[tracing::instrument(level = "debug", skip(self, content_md5, body))]
1524    async fn upload_part(
1525        &self,
1526        id: &ObjectId,
1527        upload_id: &UploadId,
1528        part_number: PartNumber,
1529        content_length: u64,
1530        content_md5: Option<&str>,
1531        body: ClientStream,
1532    ) -> Result<UploadPartResponse> {
1533        objectstore_log::debug!("Uploading part to GCS backend");
1534        let mut url = self.xml_object_url(id)?;
1535        url.query_pairs_mut()
1536            .append_pair("partNumber", &part_number.to_string())
1537            .append_pair("uploadId", upload_id);
1538
1539        let mut builder = self
1540            .request(Method::PUT, url)
1541            .await?
1542            .header(header::CONTENT_LENGTH, content_length)
1543            .body(Body::wrap_stream(body));
1544
1545        if let Some(md5) = content_md5 {
1546            builder = builder.header("content-md5", md5);
1547        }
1548
1549        let resp = builder
1550            .send_traced()
1551            .await
1552            .check_error("uploading a GCS multipart part")
1553            .await?;
1554
1555        let etag = resp
1556            .headers()
1557            .get(header::ETAG)
1558            .and_then(|v| v.to_str().ok())
1559            .map(|s| s.to_owned())
1560            .ok_or_else(|| {
1561                Error::new(
1562                    ErrorKind::BackendFailure,
1563                    "GCS upload-part response missing ETag",
1564                )
1565            })?;
1566
1567        resp.drain_body().await;
1568
1569        Ok(etag)
1570    }
1571
1572    #[tracing::instrument(level = "debug", skip(self))]
1573    async fn list_parts(
1574        &self,
1575        id: &ObjectId,
1576        upload_id: &UploadId,
1577        max_parts: Option<u32>,
1578        part_number_marker: Option<PartNumber>,
1579    ) -> Result<ListPartsResponse> {
1580        objectstore_log::debug!("Listing parts on GCS backend");
1581        let mut url = self.xml_object_url(id)?;
1582        {
1583            let mut pairs = url.query_pairs_mut();
1584            pairs.append_pair("uploadId", upload_id);
1585            if let Some(max) = max_parts {
1586                pairs.append_pair("max-parts", &max.to_string());
1587            }
1588            if let Some(marker) = part_number_marker {
1589                pairs.append_pair("part-number-marker", &marker.to_string());
1590            }
1591        }
1592
1593        self.with_retry("list_parts", || {
1594            let url = url.clone();
1595            async move {
1596                let resp = self
1597                    .request(Method::GET, url)
1598                    .await?
1599                    .send_traced()
1600                    .await
1601                    .check_error("listing GCS multipart parts")
1602                    .await?;
1603
1604                let body = resp
1605                    .bytes()
1606                    .await
1607                    .reqwest_context("reading GCS list-parts response")?;
1608
1609                let xml: XmlListPartsResponse = quick_xml::de::from_reader(body.as_ref())
1610                    .context(ErrorKind::CorruptData, "decoding GCS list-parts response")?;
1611
1612                Ok(xml.into())
1613            }
1614        })
1615        .await
1616    }
1617
1618    #[tracing::instrument(level = "debug", skip(self))]
1619    async fn abort_multipart(
1620        &self,
1621        id: &ObjectId,
1622        upload_id: &UploadId,
1623    ) -> Result<AbortMultipartResponse> {
1624        objectstore_log::debug!("Aborting multipart upload on GCS backend");
1625        let mut url = self.xml_object_url(id)?;
1626        url.query_pairs_mut().append_pair("uploadId", upload_id);
1627
1628        self.with_retry("abort_multipart", || {
1629            let url = url.clone();
1630            async move {
1631                let resp = self
1632                    .request(Method::DELETE, url)
1633                    .await?
1634                    .send_traced()
1635                    .await
1636                    .reqwest_context("aborting a GCS multipart upload")?;
1637
1638                // XXX: real S3 would return 404 here if the upload has been recently completed and we
1639                // would have to handle it. It turns out GCS returns 204 instead, so we don't need to
1640                // handle that case.
1641
1642                resp.check_error("aborting a GCS multipart upload")
1643                    .await?
1644                    .drain_body()
1645                    .await;
1646
1647                Ok(())
1648            }
1649        })
1650        .await
1651    }
1652
1653    #[tracing::instrument(level = "debug", skip(self, parts))]
1654    async fn complete_multipart(
1655        &self,
1656        id: &ObjectId,
1657        upload_id: &UploadId,
1658        parts: Vec<CompletedPart>,
1659        _access_time: Timestamp,
1660    ) -> Result<CompleteMultipartResponse> {
1661        objectstore_log::debug!("Completing multipart upload on GCS backend");
1662        let mut url = self.xml_object_url(id)?;
1663        url.query_pairs_mut().append_pair("uploadId", upload_id);
1664
1665        let body = XmlCompleteMultipartUpload::from(parts);
1666        let xml = quick_xml::se::to_string(&body).context(
1667            ErrorKind::Internal,
1668            "encoding GCS complete-multipart request",
1669        )?;
1670
1671        self.with_retry("complete_multipart", || {
1672            let url = url.clone();
1673            let xml = xml.clone();
1674            async move {
1675                let resp = self
1676                    .request(Method::POST, url)
1677                    .await?
1678                    .header(header::CONTENT_TYPE, "application/xml")
1679                    .body(xml)
1680                    .send_traced()
1681                    .await
1682                    .check_error("completing a GCS multipart upload")
1683                    .await?;
1684
1685                // XXX: real S3 would return 404 here if the upload has been recently completed and we
1686                // would have to handle it. It turns out GCS returns 200 instead, so we don't need to
1687                // handle that case.
1688
1689                let body = resp
1690                    .bytes()
1691                    .await
1692                    .reqwest_context("reading GCS complete-multipart response")?;
1693
1694                let error = quick_xml::de::from_reader::<_, XmlError>(body.as_ref())
1695                    .ok()
1696                    .map(Into::into);
1697
1698                Ok(error)
1699            }
1700        })
1701        .await
1702    }
1703}
1704
1705#[cfg(test)]
1706mod tests {
1707    use std::collections::BTreeMap;
1708    use std::io::{Read, Write};
1709    use std::net::{TcpListener, TcpStream};
1710    use std::num::{NonZeroU32, NonZeroU64};
1711    use std::sync::mpsc;
1712    use std::thread;
1713    use std::time::Duration;
1714
1715    use anyhow::Result;
1716    #[cfg(feature = "storage-cogs")]
1717    use objectstore_inventory_tracker::{OpType, test_utils::DummyProducer};
1718    use objectstore_types::scope::{Scope, Scopes};
1719    use reqwest::header::{HeaderMap, HeaderValue};
1720
1721    use super::*;
1722    use crate::backend::common::ExpiryTarget;
1723    use crate::id::ObjectContext;
1724    use crate::multipart::CompletedPart;
1725    use crate::stream;
1726    #[cfg(feature = "storage-cogs")]
1727    use crate::stream::ClientError;
1728
1729    impl GcsBackend {
1730        async fn create_upload_session(
1731            &self,
1732            id: &ObjectId,
1733            metadata: &Metadata,
1734            upload_length: NonZeroU64,
1735        ) -> Result<Session> {
1736            let backend_token =
1737                <Self as Backend>::create_upload_session(self, id, metadata, upload_length)
1738                    .await?
1739                    .ok_or_else(|| Error::from(ErrorKind::Unsupported))?;
1740            Ok(Session {
1741                object_id: id.clone(),
1742                upload_length,
1743                backend_token,
1744            })
1745        }
1746    }
1747
1748    fn read_http_request(connection: &mut TcpStream) -> String {
1749        let mut bytes = Vec::new();
1750        let mut byte = [0];
1751        while !bytes.ends_with(b"\r\n\r\n") {
1752            connection.read_exact(&mut byte).unwrap();
1753            bytes.push(byte[0]);
1754        }
1755        String::from_utf8(bytes).unwrap()
1756    }
1757
1758    fn write_http_response(
1759        connection: &mut TcpStream,
1760        status: &str,
1761        content_type: &str,
1762        body: &str,
1763    ) {
1764        write!(
1765            connection,
1766            "HTTP/1.1 {status}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
1767            body.len()
1768        )
1769        .unwrap();
1770    }
1771
1772    fn start_generation_race_server()
1773    -> (String, mpsc::Receiver<Vec<String>>, thread::JoinHandle<()>) {
1774        let listener = TcpListener::bind(("127.0.0.1", 0)).unwrap();
1775        let endpoint = format!("http://{}", listener.local_addr().unwrap());
1776        let (request_tx, request_rx) = mpsc::channel();
1777        let server = thread::spawn(move || {
1778            let responses = [
1779                (
1780                    "200 OK",
1781                    "application/json",
1782                    r#"{"contentType":"text/old","generation":"1","metageneration":"1"}"#,
1783                ),
1784                ("412 Precondition Failed", "text/plain", "stale generation"),
1785                (
1786                    "200 OK",
1787                    "application/json",
1788                    r#"{"contentType":"text/new","generation":"2","metageneration":"1"}"#,
1789                ),
1790                ("200 OK", "text/plain", "new"),
1791            ];
1792            let mut requests = Vec::new();
1793
1794            for (status, content_type, body) in responses {
1795                let (mut connection, _) = listener.accept().unwrap();
1796                requests.push(read_http_request(&mut connection));
1797                write_http_response(&mut connection, status, content_type, body);
1798            }
1799
1800            request_tx.send(requests).unwrap();
1801        });
1802        (endpoint, request_rx, server)
1803    }
1804
1805    const RESUMABLE_CHUNK_SIZE: usize = 256 * 1024;
1806
1807    // NB: Not run any of these tests, you need to have a GCS emulator running. This is done
1808    // automatically in CI.
1809    //
1810    // Refer to the readme for how to set up the emulator.
1811
1812    fn test_config() -> GcsConfig {
1813        GcsConfig {
1814            endpoint: Some("http://localhost:8087".into()),
1815            bucket: "test-bucket".into(),
1816            cogs: None,
1817        }
1818    }
1819
1820    async fn create_test_backend() -> Result<GcsBackend> {
1821        GcsBackend::new(test_config(), &ChangeStreamFactory::default()).await
1822    }
1823
1824    #[cfg(feature = "storage-cogs")]
1825    async fn create_test_backend_with_change_stream() -> Result<(GcsBackend, DummyProducer)> {
1826        let (streams, producer) = crate::change_stream::dummy_factory();
1827        let config = GcsConfig {
1828            cogs: Some(CostTrackerStreamConfig {
1829                shared_resource_id: "gcs_objectstore".into(),
1830                sample_rate: 1.0,
1831            }),
1832            ..test_config()
1833        };
1834
1835        Ok((GcsBackend::new(config, &streams).await?, producer))
1836    }
1837
1838    #[derive(Deserialize)]
1839    struct RetryTestResource {
1840        id: String,
1841    }
1842
1843    /// Configures one storage-testbench failure and sends its ID on subsequent backend requests.
1844    async fn inject_retry_test(
1845        backend: &mut GcsBackend,
1846        method: &str,
1847        instruction: &str,
1848    ) -> Result<()> {
1849        let retry_test: RetryTestResource = reqwest::Client::new()
1850            .post(backend.endpoint.join("retry_test")?)
1851            .json(&serde_json::json!({
1852                "instructions": { method: [instruction] },
1853                "transport": "HTTP",
1854            }))
1855            .send()
1856            .await?
1857            .error_for_status()?
1858            .json()
1859            .await?;
1860
1861        let mut headers = HeaderMap::new();
1862        headers.insert("x-retry-test-id", HeaderValue::from_str(&retry_test.id)?);
1863        backend.client = reqwest::Client::builder()
1864            .default_headers(headers)
1865            .build()?;
1866        Ok(())
1867    }
1868
1869    fn make_id() -> ObjectId {
1870        ObjectId::random(ObjectContext {
1871            usecase: "testing".into(),
1872            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1873        })
1874    }
1875
1876    fn make_id_with_key(key: &str) -> ObjectId {
1877        ObjectId::new(
1878            ObjectContext {
1879                usecase: "testing".into(),
1880                scopes: Scopes::from_iter([
1881                    Scope::create("organization", "42").unwrap(),
1882                    Scope::create("project", "7").unwrap(),
1883                ]),
1884            },
1885            key.into(),
1886        )
1887    }
1888
1889    fn nonzero(value: u64) -> NonZeroU64 {
1890        NonZeroU64::new(value).unwrap()
1891    }
1892
1893    #[tokio::test]
1894    async fn get_object_retries_from_metadata_after_generation_conflict() -> Result<()> {
1895        let (endpoint, request_rx, server) = start_generation_race_server();
1896        let backend = GcsBackend::new(
1897            GcsConfig {
1898                endpoint: Some(endpoint),
1899                bucket: "bucket".into(),
1900                cogs: None,
1901            },
1902            &ChangeStreamFactory::default(),
1903        )
1904        .await?;
1905
1906        let (metadata, _, stream) = backend
1907            .get_object(&make_id(), Timestamp::now(), None)
1908            .await?
1909            .unwrap();
1910        assert_eq!(metadata.content_type, "text/new");
1911        assert_eq!(stream::read_to_vec(stream).await?, b"new");
1912
1913        let requests = request_rx.recv().unwrap();
1914        assert_eq!(requests.len(), 4);
1915        assert!(!requests[0].contains("alt=media"));
1916        assert!(requests[1].contains("alt=media&ifGenerationMatch=1"));
1917        assert!(!requests[2].contains("alt=media"));
1918        assert!(requests[3].contains("alt=media&ifGenerationMatch=2"));
1919        server.join().unwrap();
1920
1921        Ok(())
1922    }
1923
1924    #[test]
1925    fn resumable_range_reports_next_offset_and_rejects_malformed_values() -> Result<()> {
1926        assert_eq!(range_header_to_offset("bytes=0-0", nonzero(10))?, 1);
1927        assert_eq!(
1928            range_header_to_offset("bytes=0-262143", nonzero(300_000))?,
1929            262_144
1930        );
1931
1932        for malformed in [
1933            "",
1934            "bytes=1-2",
1935            "bytes=0-",
1936            "bytes=0-*",
1937            "bytes=0-9 ",
1938            "bytes=0-9,bytes=20-30",
1939        ] {
1940            assert!(
1941                range_header_to_offset(malformed, nonzero(100)).is_err(),
1942                "accepted {malformed:?}"
1943            );
1944        }
1945        assert_eq!(range_header_to_offset("bytes=0-9", nonzero(10))?, 10);
1946        assert!(range_header_to_offset("bytes=0-18446744073709551615", nonzero(u64::MAX)).is_err());
1947        Ok(())
1948    }
1949
1950    #[tokio::test]
1951    async fn test_resumable_empty_chunk_reports_offset_without_writing() -> Result<()> {
1952        let backend = create_test_backend().await?;
1953        let id = make_id_with_key("resumable-empty-chunk");
1954        let upload_length = 2 * (RESUMABLE_CHUNK_SIZE + 2);
1955        let token = backend
1956            .create_upload_session(&id, &Metadata::default(), nonzero(upload_length as u64))
1957            .await?;
1958
1959        // An empty chunk cannot advance a session that still expects bytes. It reports the
1960        // authoritative offset instead of failing, and leaves what GCS holds untouched.
1961        assert_eq!(
1962            backend
1963                .put_chunk(&token, 0, 0, stream::single(Vec::new()))
1964                .await?,
1965            UploadProgress::Incomplete { offset: 0 }
1966        );
1967        // The backend rejects undersized non-final chunks, so use an unaligned chunk larger
1968        // than one granularity unit to exercise the emulator's alignment behavior.
1969        let chunk_length = RESUMABLE_CHUNK_SIZE + 2;
1970        let after_write = backend
1971            .put_chunk(
1972                &token,
1973                0,
1974                chunk_length as u64,
1975                stream::single(vec![b'a'; chunk_length]),
1976            )
1977            .await?;
1978        assert!(matches!(after_write, UploadProgress::Incomplete { .. }));
1979
1980        // Whichever prefix GCS acknowledged, an empty chunk reports that same position rather
1981        // than moving it. GCS documents that a chunk "should be a multiple of 256 KiB ... unless
1982        // it's the last chunk", and that a client "should not assume that the server received all
1983        // bytes sent in any given request". The emulator acknowledges the unaligned tail, so the
1984        // position itself is not asserted here.
1985        assert_eq!(
1986            backend
1987                .put_chunk(&token, chunk_length as u64, 0, stream::single(Vec::new()))
1988                .await?,
1989            after_write
1990        );
1991        assert_eq!(backend.upload_offset(&token).await?, after_write);
1992        Ok(())
1993    }
1994
1995    #[tokio::test]
1996    async fn test_resumable_single_and_multi_chunk_uploads() -> Result<()> {
1997        let backend = create_test_backend().await?;
1998
1999        let single_id = make_id_with_key("resumable-single");
2000        let single = b"single chunk".to_vec();
2001        let token = backend
2002            .create_upload_session(
2003                &single_id,
2004                &Metadata::default(),
2005                nonzero(single.len() as u64),
2006            )
2007            .await?;
2008        assert_eq!(
2009            backend
2010                .put_chunk(
2011                    &token,
2012                    0,
2013                    single.len() as u64,
2014                    stream::single(single.clone()),
2015                )
2016                .await?,
2017            UploadProgress::Complete
2018        );
2019        let (_, _, payload) = backend
2020            .get_object(&single_id, Timestamp::now(), None)
2021            .await?
2022            .unwrap();
2023        assert_eq!(stream::read_to_vec(payload).await?, single);
2024
2025        let multi_id = make_id_with_key("resumable-multi");
2026        let mut expected = vec![b'a'; RESUMABLE_CHUNK_SIZE];
2027        expected.extend_from_slice(b"final");
2028        let token = backend
2029            .create_upload_session(
2030                &multi_id,
2031                &Metadata::default(),
2032                nonzero(expected.len() as u64),
2033            )
2034            .await?;
2035        assert_eq!(
2036            backend.upload_offset(&token).await?,
2037            UploadProgress::Incomplete { offset: 0 }
2038        );
2039        let error = backend
2040            .put_chunk(&token, 0, 1, stream::single(b"a".to_vec()))
2041            .await
2042            .unwrap_err();
2043        assert_eq!(
2044            error.kind(),
2045            ErrorKind::ChunkTooSmall {
2046                chunk_length: 1,
2047                upload_granularity: RESUMABLE_CHUNK_SIZE as u64,
2048            }
2049        );
2050        assert_eq!(
2051            backend
2052                .put_chunk(
2053                    &token,
2054                    0,
2055                    RESUMABLE_CHUNK_SIZE as u64,
2056                    stream::single(expected[..RESUMABLE_CHUNK_SIZE].to_vec()),
2057                )
2058                .await?,
2059            UploadProgress::Incomplete {
2060                offset: RESUMABLE_CHUNK_SIZE as u64
2061            }
2062        );
2063        assert_eq!(
2064            backend.upload_offset(&token).await?,
2065            UploadProgress::Incomplete {
2066                offset: RESUMABLE_CHUNK_SIZE as u64
2067            }
2068        );
2069        assert_eq!(
2070            backend
2071                .put_chunk(
2072                    &token,
2073                    RESUMABLE_CHUNK_SIZE as u64,
2074                    5,
2075                    stream::single(b"final".to_vec()),
2076                )
2077                .await?,
2078            UploadProgress::Complete
2079        );
2080        let (_, _, payload) = backend
2081            .get_object(&multi_id, Timestamp::now(), None)
2082            .await?
2083            .unwrap();
2084        assert_eq!(stream::read_to_vec(payload).await?, expected);
2085        Ok(())
2086    }
2087
2088    #[tokio::test]
2089    async fn test_resumable_unaligned_chunk_and_rewind_preserve_persisted_bytes() -> Result<()> {
2090        let backend = create_test_backend().await?;
2091        let id = make_id_with_key("resumable-rewind");
2092        let upload_length = RESUMABLE_CHUNK_SIZE + 3;
2093        let token = backend
2094            .create_upload_session(&id, &Metadata::default(), nonzero(upload_length as u64))
2095            .await?;
2096
2097        let prefix = vec![b'a'; RESUMABLE_CHUNK_SIZE];
2098        assert_eq!(
2099            backend
2100                .put_chunk(&token, 0, prefix.len() as u64, stream::single(prefix))
2101                .await?,
2102            UploadProgress::Incomplete {
2103                offset: RESUMABLE_CHUNK_SIZE as u64
2104            }
2105        );
2106
2107        // Resend three already-persisted positions with different bytes. GCS ignores that overlap
2108        // without comparing it, then appends the suffix from its authoritative offset.
2109        assert_eq!(
2110            backend
2111                .put_chunk(
2112                    &token,
2113                    (RESUMABLE_CHUNK_SIZE - 3) as u64,
2114                    6,
2115                    stream::single(b"BADxyz".to_vec()),
2116                )
2117                .await?,
2118            UploadProgress::Complete
2119        );
2120
2121        let (_, _, payload) = backend
2122            .get_object(&id, Timestamp::now(), None)
2123            .await?
2124            .unwrap();
2125        let payload = stream::read_to_vec(payload).await?;
2126        assert_eq!(&payload[RESUMABLE_CHUNK_SIZE - 3..], b"aaaxyz");
2127        Ok(())
2128    }
2129
2130    #[tokio::test]
2131    async fn test_resumable_rejects_oversized_chunks() -> Result<()> {
2132        let backend = create_test_backend().await?;
2133        let id = make_id_with_key("resumable-validation");
2134        let token = backend
2135            .create_upload_session(&id, &Metadata::default(), nonzero(10))
2136            .await?;
2137
2138        let error = backend
2139            .put_chunk(&token, 8, 3, stream::single(b"abc".to_vec()))
2140            .await
2141            .unwrap_err();
2142        assert_eq!(
2143            error.kind(),
2144            ErrorKind::ChunkExceedsUploadLength {
2145                offset: 8,
2146                content_length: 3,
2147                upload_length: 10
2148            }
2149        );
2150
2151        let error = backend
2152            .put_chunk(&token, u64::MAX, 2, stream::single(b"ab".to_vec()))
2153            .await
2154            .unwrap_err();
2155        assert!(matches!(
2156            error.kind(),
2157            ErrorKind::ChunkExceedsUploadLength { .. }
2158        ));
2159        Ok(())
2160    }
2161
2162    #[tokio::test]
2163    async fn test_resumable_cancel_is_idempotent_and_session_is_not_found() -> Result<()> {
2164        let backend = create_test_backend().await?;
2165        let id = make_id_with_key("resumable-cancel");
2166        let token = backend
2167            .create_upload_session(&id, &Metadata::default(), nonzero(10))
2168            .await?;
2169
2170        backend.cancel_upload(&token).await?;
2171        backend.cancel_upload(&token).await?;
2172        assert!(matches!(
2173            backend.upload_offset(&token).await,
2174            Err(error) if error.kind() == ErrorKind::UnknownUploadSession
2175        ));
2176        Ok(())
2177    }
2178
2179    #[tokio::test]
2180    async fn test_resumable_retries_only_replayable_operations() -> Result<()> {
2181        let mut backend = create_test_backend().await?;
2182        inject_retry_test(&mut backend, "storage.objects.insert", "return-503").await?;
2183        let id = make_id_with_key("resumable-retries");
2184
2185        // Creation consumes the injected 503 and succeeds on the backend's retry.
2186        let token = backend
2187            .create_upload_session(&id, &Metadata::default(), nonzero(10))
2188            .await?;
2189
2190        inject_retry_test(&mut backend, "storage.objects.insert", "return-503").await?;
2191        assert_eq!(
2192            backend.upload_offset(&token).await?,
2193            UploadProgress::Incomplete { offset: 0 }
2194        );
2195
2196        inject_retry_test(&mut backend, "storage.objects.delete", "return-503").await?;
2197        backend.cancel_upload(&token).await?;
2198        Ok(())
2199    }
2200
2201    #[tokio::test]
2202    async fn test_resumable_stream_failures_require_explicit_offset_recovery() -> Result<()> {
2203        // A failure before persistence is returned directly: put_chunk must not retry the body.
2204        let mut backend = create_test_backend().await?;
2205        let id = make_id_with_key("resumable-failure-before");
2206        let token = backend
2207            .create_upload_session(&id, &Metadata::default(), nonzero(4))
2208            .await?;
2209        inject_retry_test(&mut backend, "storage.objects.insert", "return-503").await?;
2210        assert!(matches!(
2211            backend
2212                .put_chunk(&token, 0, 4, stream::single(b"data".to_vec()))
2213                .await,
2214            Err(error) if error.kind() == ErrorKind::BackendUnavailable
2215        ));
2216        assert_eq!(
2217            backend.upload_offset(&token).await?,
2218            UploadProgress::Incomplete { offset: 0 }
2219        );
2220
2221        // The emulator persists the first KiB, then fails. Recovery observes that prefix and the
2222        // caller resumes exactly from the returned offset.
2223        let mut backend = create_test_backend().await?;
2224        let id = make_id_with_key("resumable-failure-partial");
2225        let data = vec![b'p'; 2048];
2226        let token = backend
2227            .create_upload_session(&id, &Metadata::default(), nonzero(data.len() as u64))
2228            .await?;
2229        inject_retry_test(
2230            &mut backend,
2231            "storage.objects.insert",
2232            "return-503-after-1K",
2233        )
2234        .await?;
2235        assert!(matches!(
2236            backend
2237                .put_chunk(
2238                    &token,
2239                    0,
2240                    data.len() as u64,
2241                    stream::single(data.clone()),
2242                )
2243                .await,
2244            Err(error) if error.kind() == ErrorKind::BackendUnavailable
2245        ));
2246        assert_eq!(
2247            backend.upload_offset(&token).await?,
2248            UploadProgress::Incomplete { offset: 1024 }
2249        );
2250        assert_eq!(
2251            backend
2252                .put_chunk(&token, 1024, 1024, stream::single(data[1024..].to_vec()))
2253                .await?,
2254            UploadProgress::Complete
2255        );
2256
2257        // GCS persists the final bytes, but storage-testbench truncates the successful JSON
2258        // response. The chunk is not retried; an explicit status query observes completion.
2259        let mut backend = create_test_backend().await?;
2260        let id = make_id_with_key("resumable-failure-final");
2261        let token = backend
2262            .create_upload_session(&id, &Metadata::default(), nonzero(5))
2263            .await?;
2264        inject_retry_test(
2265            &mut backend,
2266            "storage.objects.insert",
2267            "return-broken-stream-final-chunk-after-0B",
2268        )
2269        .await?;
2270        assert!(matches!(
2271            backend
2272                .put_chunk(&token, 0, 5, stream::single(b"final".to_vec()))
2273                .await,
2274            Err(error) if error.kind() == ErrorKind::CorruptData
2275        ));
2276        assert_eq!(
2277            backend.upload_offset(&token).await?,
2278            UploadProgress::Complete
2279        );
2280        Ok(())
2281    }
2282
2283    async fn get_gcs_generations(
2284        backend: &GcsBackend,
2285        object_url: Url,
2286    ) -> Result<(String, String)> {
2287        Ok(backend
2288            .request(Method::GET, object_url)
2289            .await?
2290            .send_traced()
2291            .await
2292            .check_error("getting GCS object metadata")
2293            .await?
2294            .json::<GcsObject>()
2295            .await
2296            .context(
2297                ErrorKind::BackendFailure,
2298                "decoding GCS object metadata response",
2299            )
2300            .map(|object| (object.generation, object.metageneration))?)
2301    }
2302
2303    #[tokio::test]
2304    async fn test_roundtrip() -> Result<()> {
2305        let backend = create_test_backend().await?;
2306
2307        let id = make_id();
2308        let metadata = Metadata {
2309            content_type: "text/plain".into(),
2310            expiration_policy: ExpirationPolicy::Manual,
2311            compression: None,
2312            origin: Some("203.0.113.42".into()),
2313            filename: Some("hello.txt".into()),
2314            custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
2315            time_created: Some(Timestamp::now()),
2316            time_expires: None,
2317            size: None,
2318        };
2319
2320        backend
2321            .put_object(
2322                &id,
2323                &metadata,
2324                stream::single("hello, world"),
2325                Timestamp::now(),
2326            )
2327            .await?;
2328
2329        let (meta, _, stream) = backend
2330            .get_object(&id, Timestamp::now(), None)
2331            .await?
2332            .unwrap();
2333
2334        let payload = stream::read_to_vec(stream).await?;
2335        let str_payload = str::from_utf8(&payload).unwrap();
2336        assert_eq!(str_payload, "hello, world");
2337        assert_eq!(meta.content_type, metadata.content_type);
2338        assert_eq!(meta.origin, metadata.origin);
2339        assert_eq!(meta.filename, metadata.filename);
2340        assert_eq!(meta.custom, metadata.custom);
2341        assert!(metadata.time_created.is_some());
2342
2343        Ok(())
2344    }
2345
2346    /// Metadata with a non-ASCII filename and custom metadata value.
2347    fn unicode_metadata() -> Metadata {
2348        Metadata {
2349            filename: Some("réport-📄.pdf".into()),
2350            custom: BTreeMap::from_iter([("release".into(), "vérsion-1.0-🚀".into())]),
2351            ..Default::default()
2352        }
2353    }
2354
2355    /// Both GCS write paths must agree on how a logical string is stored, because every read goes
2356    /// through the JSON API: `put_object` writes metadata as JSON, `initiate_multipart` writes it
2357    /// as `x-goog-meta-*` request headers.
2358    #[tokio::test]
2359    async fn test_unicode_metadata_roundtrip_json_upload() -> Result<()> {
2360        let backend = create_test_backend().await?;
2361        let id = make_id();
2362        let metadata = unicode_metadata();
2363
2364        backend
2365            .put_object(
2366                &id,
2367                &metadata,
2368                stream::single("hello, world"),
2369                Timestamp::now(),
2370            )
2371            .await?;
2372
2373        let meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
2374        assert_eq!(meta.filename, metadata.filename);
2375        assert_eq!(meta.custom, metadata.custom);
2376
2377        Ok(())
2378    }
2379
2380    #[tokio::test]
2381    async fn test_unicode_metadata_roundtrip_multipart_upload() -> Result<()> {
2382        let backend = create_test_backend().await?;
2383        let id = make_id();
2384        let metadata = unicode_metadata();
2385
2386        multipart_put(&backend, &id, &metadata, "hello, world").await?;
2387
2388        let meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
2389        assert_eq!(meta.filename, metadata.filename);
2390        assert_eq!(meta.custom, metadata.custom);
2391
2392        Ok(())
2393    }
2394
2395    #[test]
2396    fn from_metadata_uses_provided_time_expires() {
2397        let created = Timestamp::now();
2398        let expires = created + Duration::from_hours(1);
2399        let metadata = Metadata {
2400            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2401            time_created: Some(created),
2402            time_expires: Some(expires),
2403            ..Default::default()
2404        };
2405
2406        let gcs_object = GcsObject::from_metadata(&metadata);
2407        let custom_time = gcs_object.custom_time.map(Rfc3339Timestamp::into_inner);
2408        assert_eq!(custom_time, Some(expires));
2409
2410        let roundtripped = gcs_object.into_metadata().unwrap();
2411        assert_eq!(roundtripped.time_created, Some(created));
2412        assert_eq!(roundtripped.time_expires, Some(expires));
2413    }
2414
2415    #[tokio::test]
2416    async fn test_get_nonexistent() -> Result<()> {
2417        let backend = create_test_backend().await?;
2418
2419        let id = make_id();
2420        let result = backend.get_object(&id, Timestamp::now(), None).await?;
2421        assert!(result.is_none());
2422
2423        Ok(())
2424    }
2425
2426    #[tokio::test]
2427    async fn test_delete_nonexistent() -> Result<()> {
2428        let backend = create_test_backend().await?;
2429
2430        let id = make_id();
2431        backend.delete_object(&id, Timestamp::now()).await?;
2432
2433        Ok(())
2434    }
2435
2436    #[tokio::test]
2437    async fn test_overwrite() -> Result<()> {
2438        let backend = create_test_backend().await?;
2439
2440        let id = make_id();
2441        let metadata = Metadata {
2442            custom: BTreeMap::from_iter([("invalid".into(), "invalid".into())]),
2443            ..Default::default()
2444        };
2445
2446        backend
2447            .put_object(&id, &metadata, stream::single("hello"), Timestamp::now())
2448            .await?;
2449
2450        let metadata = Metadata {
2451            custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
2452            ..Default::default()
2453        };
2454
2455        backend
2456            .put_object(&id, &metadata, stream::single("world"), Timestamp::now())
2457            .await?;
2458
2459        let (meta, _, stream) = backend
2460            .get_object(&id, Timestamp::now(), None)
2461            .await?
2462            .unwrap();
2463
2464        let payload = stream::read_to_vec(stream).await?;
2465        let str_payload = str::from_utf8(&payload).unwrap();
2466        assert_eq!(str_payload, "world");
2467        assert_eq!(meta.custom, metadata.custom);
2468
2469        Ok(())
2470    }
2471
2472    #[tokio::test]
2473    async fn test_read_after_delete() -> Result<()> {
2474        let backend = create_test_backend().await?;
2475
2476        let id = make_id();
2477        let metadata = Metadata::default();
2478
2479        backend
2480            .put_object(
2481                &id,
2482                &metadata,
2483                stream::single("hello, world"),
2484                Timestamp::now(),
2485            )
2486            .await?;
2487
2488        backend.delete_object(&id, Timestamp::now()).await?;
2489
2490        let result = backend.get_object(&id, Timestamp::now(), None).await?;
2491        assert!(result.is_none());
2492
2493        Ok(())
2494    }
2495
2496    #[tokio::test]
2497    async fn test_ttl_immediate() -> Result<()> {
2498        // NB: We create a TTL that immediately expires in this tests. This might be optimized away
2499        // in a future implementation, so we will have to update this test accordingly.
2500
2501        let backend = create_test_backend().await?;
2502
2503        let id = make_id();
2504        let metadata = Metadata {
2505            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(0)),
2506            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2507            ..Default::default()
2508        };
2509
2510        backend
2511            .put_object(
2512                &id,
2513                &metadata,
2514                stream::single("hello, world"),
2515                Timestamp::now(),
2516            )
2517            .await?;
2518
2519        let result = backend.get_object(&id, Timestamp::now(), None).await?;
2520        assert!(result.is_none());
2521
2522        Ok(())
2523    }
2524
2525    #[tokio::test]
2526    async fn test_tti_immediate() -> Result<()> {
2527        // NB: We create a TTI that immediately expires in this tests. This might be optimized away
2528        // in a future implementation, so we will have to update this test accordingly.
2529
2530        let backend = create_test_backend().await?;
2531
2532        let id = make_id();
2533        let metadata = Metadata {
2534            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(0)),
2535            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
2536            ..Default::default()
2537        };
2538
2539        backend
2540            .put_object(
2541                &id,
2542                &metadata,
2543                stream::single("hello, world"),
2544                Timestamp::now(),
2545            )
2546            .await?;
2547
2548        let result = backend.get_object(&id, Timestamp::now(), None).await?;
2549        assert!(result.is_none());
2550
2551        Ok(())
2552    }
2553
2554    #[tokio::test]
2555    async fn test_get_metadata_returns_metadata() -> Result<()> {
2556        let backend = create_test_backend().await?;
2557
2558        let id = make_id();
2559        let metadata = Metadata {
2560            content_type: "text/plain".into(),
2561            origin: Some("203.0.113.42".into()),
2562            custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
2563            ..Default::default()
2564        };
2565
2566        backend
2567            .put_object(
2568                &id,
2569                &metadata,
2570                stream::single("hello, world"),
2571                Timestamp::now(),
2572            )
2573            .await?;
2574
2575        let meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
2576        assert_eq!(meta.content_type, metadata.content_type);
2577        assert_eq!(meta.origin, metadata.origin);
2578        assert_eq!(meta.custom, metadata.custom);
2579
2580        Ok(())
2581    }
2582
2583    #[tokio::test]
2584    async fn test_get_metadata_nonexistent() -> Result<()> {
2585        let backend = create_test_backend().await?;
2586
2587        let id = make_id();
2588        let result = backend.get_metadata(&id, Timestamp::now()).await?;
2589        assert!(result.is_none());
2590
2591        Ok(())
2592    }
2593
2594    #[tokio::test]
2595    async fn test_set_expiry() -> Result<()> {
2596        for is_ttl in [false, true] {
2597            let backend = create_test_backend().await?;
2598
2599            let id = make_id();
2600            let tti = Duration::from_hours(2 * 24);
2601            let metadata = Metadata {
2602                content_type: "text/plain".into(),
2603                custom: [("preserved".into(), "yes".into())].into(),
2604                expiration_policy: if is_ttl {
2605                    ExpirationPolicy::TimeToLive(tti)
2606                } else {
2607                    ExpirationPolicy::TimeToIdle(tti)
2608                },
2609                time_expires: Some(Timestamp::now() + tti),
2610                ..Default::default()
2611            };
2612
2613            backend
2614                .put_object(
2615                    &id,
2616                    &metadata,
2617                    stream::single("hello, world"),
2618                    Timestamp::now(),
2619                )
2620                .await?;
2621
2622            // Backdate custom_time while keeping the object live.
2623            let object_url = backend.object_url(&id)?;
2624            let old_deadline = Timestamp::now() + Duration::from_mins(1);
2625            let generations = get_gcs_generations(&backend, object_url.clone()).await?;
2626            backend
2627                .update_custom_time(
2628                    object_url,
2629                    old_deadline,
2630                    (&generations.0, &generations.1),
2631                    None,
2632                )
2633                .await?;
2634
2635            // Backend reads return the stored deadline without modifying it.
2636            let pre_meta = backend.get_metadata(&id, Timestamp::now()).await?.unwrap();
2637            let pre_expiry = pre_meta.time_expires.unwrap();
2638            let created = pre_meta.time_created.unwrap();
2639            assert_eq!(
2640                backend
2641                    .get_metadata(&id, Timestamp::now())
2642                    .await?
2643                    .unwrap()
2644                    .time_expires,
2645                Some(pre_expiry)
2646            );
2647
2648            let requested = created + tti + tti;
2649            for target in [
2650                ExpiryTarget::At(requested),
2651                ExpiryTarget::FromCreation(tti + tti),
2652            ] {
2653                assert_eq!(
2654                    backend
2655                        .set_expiry(&id, target.into(), Timestamp::now())
2656                        .await?,
2657                    SetExpiryResponse::Satisfied(requested)
2658                );
2659            }
2660            assert_eq!(
2661                backend
2662                    .get_metadata(&id, Timestamp::now())
2663                    .await?
2664                    .unwrap()
2665                    .time_expires,
2666                Some(requested)
2667            );
2668
2669            let (updated, _, stream) = backend
2670                .get_object(&id, Timestamp::now(), None)
2671                .await?
2672                .unwrap();
2673            assert_eq!(
2674                updated.expiration_policy,
2675                if is_ttl {
2676                    ExpirationPolicy::TimeToLive(tti + tti)
2677                } else {
2678                    ExpirationPolicy::TimeToIdle(tti)
2679                }
2680            );
2681            assert_eq!(updated.custom, metadata.custom);
2682            let payload = stream::read_to_vec(stream).await?;
2683            assert_eq!(&payload, b"hello, world");
2684        }
2685        Ok(())
2686    }
2687
2688    #[tokio::test]
2689    async fn test_expiry_conflict() -> Result<()> {
2690        let backend = create_test_backend().await?;
2691        let id = make_id();
2692        let metadata = Metadata {
2693            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
2694            time_expires: Some(Timestamp::now() + Duration::from_mins(10)),
2695            ..Default::default()
2696        };
2697        backend
2698            .put_object(&id, &metadata, stream::single("payload"), Timestamp::now())
2699            .await?;
2700        let object_url = backend.object_url(&id)?;
2701        let generations = get_gcs_generations(&backend, object_url.clone()).await?;
2702
2703        let deadline = Timestamp::now() + Duration::from_hours(1);
2704        assert_eq!(
2705            backend
2706                .update_custom_time(
2707                    object_url.clone(),
2708                    deadline,
2709                    (&generations.0, &generations.1),
2710                    None,
2711                )
2712                .await?,
2713            SetExpiryResponse::Satisfied(deadline)
2714        );
2715        assert_eq!(
2716            backend
2717                .update_custom_time(
2718                    object_url.clone(),
2719                    Timestamp::now() + Duration::from_hours(2),
2720                    (&generations.0, &generations.1),
2721                    None,
2722                )
2723                .await?,
2724            SetExpiryResponse::Rejected
2725        );
2726
2727        backend.delete_object(&id, Timestamp::now()).await?;
2728        assert_eq!(
2729            backend
2730                .update_custom_time(
2731                    object_url,
2732                    Timestamp::now() + Duration::from_hours(2),
2733                    (&generations.0, &generations.1),
2734                    None,
2735                )
2736                .await?,
2737            SetExpiryResponse::NotFound
2738        );
2739        Ok(())
2740    }
2741
2742    #[tokio::test]
2743    async fn test_compressed_payload_roundtrip() -> Result<()> {
2744        use objectstore_types::metadata::Compression;
2745
2746        let backend = create_test_backend().await?;
2747
2748        let plaintext = b"hello, world (but compressed with zstd)";
2749        let compressed = zstd::encode_all(&plaintext[..], 3)?;
2750
2751        let id = make_id();
2752        let metadata = Metadata {
2753            content_type: "text/plain".into(),
2754            compression: Some(Compression::Zstd),
2755            ..Default::default()
2756        };
2757
2758        backend
2759            .put_object(
2760                &id,
2761                &metadata,
2762                stream::single(compressed.clone()),
2763                Timestamp::now(),
2764            )
2765            .await?;
2766
2767        let (meta, _, stream) = backend
2768            .get_object(&id, Timestamp::now(), None)
2769            .await?
2770            .unwrap();
2771        let payload = stream::read_to_vec(stream).await?;
2772
2773        assert_eq!(meta.compression, Some(Compression::Zstd));
2774        assert_eq!(
2775            payload, compressed,
2776            "Payload should be returned still compressed, not auto-decompressed"
2777        );
2778
2779        Ok(())
2780    }
2781
2782    #[tokio::test]
2783    async fn test_multipart_single_part() -> Result<()> {
2784        let backend = create_test_backend().await?;
2785        let id = make_id();
2786        let metadata = Metadata {
2787            content_type: "text/plain".into(),
2788            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_mins(33)),
2789            origin: Some("203.0.113.42".into()),
2790            custom: BTreeMap::from_iter([("hello".into(), "world".into())]),
2791            ..Default::default()
2792        };
2793
2794        let upload_id = backend.initiate_multipart(&id, &metadata).await?;
2795
2796        let data = b"hello, multipart world!";
2797        let etag = backend
2798            .upload_part(
2799                &id,
2800                &upload_id,
2801                NonZeroU32::new(1).unwrap(),
2802                data.len() as u64,
2803                None,
2804                stream::single(data.to_vec()),
2805            )
2806            .await?;
2807
2808        let result = backend
2809            .complete_multipart(
2810                &id,
2811                &upload_id,
2812                vec![CompletedPart {
2813                    part_number: NonZeroU32::new(1).unwrap(),
2814                    etag,
2815                }],
2816                Timestamp::now(),
2817            )
2818            .await?;
2819        assert!(result.is_none(), "expected no error on complete");
2820
2821        let (meta, _, stream) = backend
2822            .get_object(&id, Timestamp::now(), None)
2823            .await?
2824            .unwrap();
2825        let payload = stream::read_to_vec(stream).await?;
2826        assert_eq!(payload, data);
2827        assert_eq!(meta.content_type, "text/plain".to_string());
2828        assert_eq!(
2829            meta.expiration_policy,
2830            ExpirationPolicy::TimeToLive(Duration::from_mins(33))
2831        );
2832        assert_eq!(meta.origin, Some("203.0.113.42".into()));
2833        assert_eq!(
2834            meta.custom,
2835            BTreeMap::from_iter([("hello".into(), "world".into())])
2836        );
2837
2838        Ok(())
2839    }
2840
2841    #[tokio::test]
2842    async fn test_multipart_multiple_parts() -> Result<()> {
2843        let backend = create_test_backend().await?;
2844        let id = make_id();
2845        let metadata = Metadata::default();
2846
2847        let upload_id = backend.initiate_multipart(&id, &metadata).await?;
2848
2849        // Non-final parts must be >= 5 MiB.
2850        const MIN_PART: usize = 5 * 1024 * 1024;
2851        let part1 = vec![b'a'; MIN_PART];
2852        let part2 = vec![b'b'; MIN_PART];
2853        let part3 = b"cccc".to_vec();
2854
2855        let etag1 = backend
2856            .upload_part(
2857                &id,
2858                &upload_id,
2859                NonZeroU32::new(1).unwrap(),
2860                part1.len() as u64,
2861                None,
2862                stream::single(part1.clone()),
2863            )
2864            .await?;
2865        let etag2 = backend
2866            .upload_part(
2867                &id,
2868                &upload_id,
2869                NonZeroU32::new(2).unwrap(),
2870                part2.len() as u64,
2871                None,
2872                stream::single(part2.clone()),
2873            )
2874            .await?;
2875        let etag3 = backend
2876            .upload_part(
2877                &id,
2878                &upload_id,
2879                NonZeroU32::new(3).unwrap(),
2880                part3.len() as u64,
2881                None,
2882                stream::single(part3.clone()),
2883            )
2884            .await?;
2885
2886        let result = backend
2887            .complete_multipart(
2888                &id,
2889                &upload_id,
2890                vec![
2891                    CompletedPart {
2892                        part_number: NonZeroU32::new(1).unwrap(),
2893                        etag: etag1,
2894                    },
2895                    CompletedPart {
2896                        part_number: NonZeroU32::new(2).unwrap(),
2897                        etag: etag2,
2898                    },
2899                    CompletedPart {
2900                        part_number: NonZeroU32::new(3).unwrap(),
2901                        etag: etag3,
2902                    },
2903                ],
2904                Timestamp::now(),
2905            )
2906            .await?;
2907        assert!(result.is_none(), "expected no error on complete");
2908
2909        // Object exists after complete
2910        let (_meta, _, stream) = backend
2911            .get_object(&id, Timestamp::now(), None)
2912            .await?
2913            .unwrap();
2914        let payload = stream::read_to_vec(stream).await?;
2915        let mut expected = Vec::new();
2916        expected.extend_from_slice(&part1);
2917        expected.extend_from_slice(&part2);
2918        expected.extend_from_slice(&part3);
2919        assert_eq!(payload, expected);
2920
2921        Ok(())
2922    }
2923
2924    #[tokio::test]
2925    async fn test_multipart_out_of_order_upload() -> Result<()> {
2926        let backend = create_test_backend().await?;
2927        let id = make_id();
2928        let metadata = Metadata::default();
2929
2930        let upload_id = backend.initiate_multipart(&id, &metadata).await?;
2931
2932        // Non-final parts must be >= 5 MiB.
2933        const MIN_PART: usize = 5 * 1024 * 1024;
2934        let part1 = vec![b'a'; MIN_PART];
2935        let part2 = vec![b'b'; MIN_PART];
2936        let part3 = b"cccc".to_vec();
2937
2938        // Upload parts out of order: 2, 3, 1.
2939        let etag2 = backend
2940            .upload_part(
2941                &id,
2942                &upload_id,
2943                NonZeroU32::new(2).unwrap(),
2944                part2.len() as u64,
2945                None,
2946                stream::single(part2.clone()),
2947            )
2948            .await?;
2949        let etag3 = backend
2950            .upload_part(
2951                &id,
2952                &upload_id,
2953                NonZeroU32::new(3).unwrap(),
2954                part3.len() as u64,
2955                None,
2956                stream::single(part3.clone()),
2957            )
2958            .await?;
2959        let etag1 = backend
2960            .upload_part(
2961                &id,
2962                &upload_id,
2963                NonZeroU32::new(1).unwrap(),
2964                part1.len() as u64,
2965                None,
2966                stream::single(part1.clone()),
2967            )
2968            .await?;
2969
2970        // Complete with parts listed in order.
2971        let result = backend
2972            .complete_multipart(
2973                &id,
2974                &upload_id,
2975                vec![
2976                    CompletedPart {
2977                        part_number: NonZeroU32::new(1).unwrap(),
2978                        etag: etag1,
2979                    },
2980                    CompletedPart {
2981                        part_number: NonZeroU32::new(2).unwrap(),
2982                        etag: etag2,
2983                    },
2984                    CompletedPart {
2985                        part_number: NonZeroU32::new(3).unwrap(),
2986                        etag: etag3,
2987                    },
2988                ],
2989                Timestamp::now(),
2990            )
2991            .await?;
2992        assert!(result.is_none(), "expected no error on complete");
2993
2994        // Verify reassembly order matches part numbers, not upload order.
2995        let (_meta, _, stream) = backend
2996            .get_object(&id, Timestamp::now(), None)
2997            .await?
2998            .unwrap();
2999        let payload = stream::read_to_vec(stream).await?;
3000        let mut expected = Vec::new();
3001        expected.extend_from_slice(&part1);
3002        expected.extend_from_slice(&part2);
3003        expected.extend_from_slice(&part3);
3004        assert_eq!(payload, expected);
3005
3006        Ok(())
3007    }
3008
3009    #[tokio::test]
3010    async fn test_multipart_list_parts() -> Result<()> {
3011        let backend = create_test_backend().await?;
3012        let id = make_id();
3013        let metadata = Metadata::default();
3014
3015        let upload_id = backend.initiate_multipart(&id, &metadata).await?;
3016
3017        let etag1 = backend
3018            .upload_part(
3019                &id,
3020                &upload_id,
3021                NonZeroU32::new(1).unwrap(),
3022                3,
3023                None,
3024                stream::single(b"aaa".to_vec()),
3025            )
3026            .await?;
3027        let etag2 = backend
3028            .upload_part(
3029                &id,
3030                &upload_id,
3031                NonZeroU32::new(2).unwrap(),
3032                3,
3033                None,
3034                stream::single(b"bbb".to_vec()),
3035            )
3036            .await?;
3037
3038        // List all parts.
3039        let list = backend.list_parts(&id, &upload_id, None, None).await?;
3040        assert_eq!(list.parts.len(), 2);
3041        assert_eq!(list.parts[0].part_number.get(), 1);
3042        assert_eq!(list.parts[0].etag, etag1);
3043        assert_eq!(list.parts[0].size, 3);
3044        assert_eq!(list.parts[1].part_number.get(), 2);
3045        assert_eq!(list.parts[1].etag, etag2);
3046        assert_eq!(list.parts[1].size, 3);
3047
3048        // List with max_parts=1 to test pagination.
3049        let page1 = backend.list_parts(&id, &upload_id, Some(1), None).await?;
3050        assert_eq!(page1.parts.len(), 1);
3051        assert_eq!(page1.parts[0].part_number.get(), 1);
3052        assert!(page1.is_truncated);
3053        assert!(page1.next_part_number_marker.is_some());
3054
3055        let page2 = backend
3056            .list_parts(&id, &upload_id, Some(1), page1.next_part_number_marker)
3057            .await?;
3058        assert_eq!(page2.parts.len(), 1);
3059        assert_eq!(page2.parts[0].part_number.get(), 2);
3060
3061        // Clean up.
3062        backend.abort_multipart(&id, &upload_id).await?;
3063
3064        Ok(())
3065    }
3066
3067    #[tokio::test]
3068    async fn test_multipart_abort() -> Result<()> {
3069        let backend = create_test_backend().await?;
3070        let id = make_id();
3071        let metadata = Metadata::default();
3072
3073        let upload_id = backend.initiate_multipart(&id, &metadata).await?;
3074
3075        backend
3076            .upload_part(
3077                &id,
3078                &upload_id,
3079                NonZeroU32::new(1).unwrap(),
3080                5,
3081                None,
3082                stream::single(b"hello".to_vec()),
3083            )
3084            .await?;
3085
3086        backend.abort_multipart(&id, &upload_id).await?;
3087
3088        // Object should not exist after abort.
3089        let result = backend.get_object(&id, Timestamp::now(), None).await?;
3090        assert!(result.is_none(), "object should not exist after abort");
3091
3092        // A second abort should still succeed (idempotent 404 handling).
3093        backend.abort_multipart(&id, &upload_id).await?;
3094
3095        Ok(())
3096    }
3097
3098    async fn multipart_put(
3099        backend: &GcsBackend,
3100        id: &ObjectId,
3101        metadata: &Metadata,
3102        payload: impl Into<bytes::Bytes>,
3103    ) -> Result<()> {
3104        let payload: bytes::Bytes = payload.into();
3105        let upload_id = backend.initiate_multipart(id, metadata).await?;
3106        let etag = backend
3107            .upload_part(
3108                id,
3109                &upload_id,
3110                NonZeroU32::new(1).unwrap(),
3111                payload.len() as u64,
3112                None,
3113                stream::single(payload),
3114            )
3115            .await?;
3116        let error = backend
3117            .complete_multipart(
3118                id,
3119                &upload_id,
3120                vec![CompletedPart {
3121                    part_number: NonZeroU32::new(1).unwrap(),
3122                    etag,
3123                }],
3124                Timestamp::now(),
3125            )
3126            .await?;
3127        assert!(
3128            error.is_none(),
3129            "complete_multipart returned error: {error:?}"
3130        );
3131        Ok(())
3132    }
3133
3134    #[tokio::test]
3135    async fn test_multipart_ttl_immediate() -> Result<()> {
3136        let backend = create_test_backend().await?;
3137        let id = make_id();
3138        let metadata = Metadata {
3139            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(0)),
3140            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
3141            ..Default::default()
3142        };
3143
3144        multipart_put(&backend, &id, &metadata, "hello, world").await?;
3145
3146        let result = backend.get_object(&id, Timestamp::now(), None).await?;
3147        assert!(result.is_none());
3148
3149        Ok(())
3150    }
3151
3152    #[tokio::test]
3153    async fn test_multipart_tti_immediate() -> Result<()> {
3154        let backend = create_test_backend().await?;
3155        let id = make_id();
3156        let metadata = Metadata {
3157            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(0)),
3158            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
3159            ..Default::default()
3160        };
3161
3162        multipart_put(&backend, &id, &metadata, "hello, world").await?;
3163
3164        let result = backend.get_object(&id, Timestamp::now(), None).await?;
3165        assert!(result.is_none());
3166
3167        Ok(())
3168    }
3169
3170    #[tokio::test]
3171    async fn test_multipart_compressed_payload_roundtrip() -> Result<()> {
3172        use objectstore_types::metadata::Compression;
3173
3174        let backend = create_test_backend().await?;
3175
3176        let plaintext = b"hello, world (but compressed with zstd)";
3177        let compressed = zstd::encode_all(&plaintext[..], 3)?;
3178
3179        let id = make_id();
3180        let metadata = Metadata {
3181            content_type: "text/plain".into(),
3182            compression: Some(Compression::Zstd),
3183            ..Default::default()
3184        };
3185
3186        multipart_put(&backend, &id, &metadata, compressed.clone()).await?;
3187
3188        let (meta, _, stream) = backend
3189            .get_object(&id, Timestamp::now(), None)
3190            .await?
3191            .unwrap();
3192        let payload = stream::read_to_vec(stream).await?;
3193
3194        assert_eq!(meta.compression, Some(Compression::Zstd));
3195        assert_eq!(
3196            payload, compressed,
3197            "Payload should be returned still compressed, not auto-decompressed"
3198        );
3199
3200        Ok(())
3201    }
3202
3203    #[cfg(feature = "storage-cogs")]
3204    #[tokio::test]
3205    async fn change_stream_reports_the_size_gcs_stored() -> Result<()> {
3206        let (backend, producer) = create_test_backend_with_change_stream().await?;
3207        let id = make_id();
3208        let payload = vec![b'x'; 4096];
3209        let metadata = Metadata {
3210            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(3600)),
3211            time_expires: Some(Timestamp::now() + Duration::from_secs(3600)),
3212            ..Default::default()
3213        };
3214
3215        backend
3216            .put_object(
3217                &id,
3218                &metadata,
3219                stream::single::<ClientError>(payload.clone()),
3220                Timestamp::now(),
3221            )
3222            .await?;
3223
3224        let records = producer.records();
3225        assert_eq!(records.len(), 1);
3226        assert_eq!(records[0].op_type, OpType::Write);
3227        assert_eq!(records[0].shared_resource_id, "gcs_objectstore");
3228        assert_eq!(records[0].app_feature, "testing");
3229        assert_eq!(
3230            records[0].size,
3231            Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size())
3232        );
3233        assert!(records[0].expiration_time.is_some());
3234
3235        Ok(())
3236    }
3237
3238    #[cfg(feature = "storage-cogs")]
3239    #[tokio::test]
3240    async fn resumable_completion_reports_to_change_stream() -> Result<()> {
3241        let (backend, producer) = create_test_backend_with_change_stream().await?;
3242        let id = make_id();
3243        let payload = b"resumable payload".to_vec();
3244        let metadata = Metadata {
3245            time_expires: Some(Timestamp::now() + Duration::from_secs(3600)),
3246            ..Default::default()
3247        };
3248        let token = backend
3249            .create_upload_session(&id, &metadata, nonzero(payload.len() as u64))
3250            .await?;
3251
3252        assert_eq!(
3253            backend
3254                .put_chunk(
3255                    &token,
3256                    0,
3257                    payload.len() as u64,
3258                    stream::single::<ClientError>(payload.clone()),
3259                )
3260                .await?,
3261            UploadProgress::Complete
3262        );
3263
3264        let records = producer.records();
3265        assert_eq!(records.len(), 1);
3266        assert_eq!(records[0].op_type, OpType::Write);
3267        assert_eq!(
3268            records[0].size,
3269            Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size())
3270        );
3271        assert!(records[0].expiration_time.is_some());
3272        Ok(())
3273    }
3274
3275    #[cfg(feature = "storage-cogs")]
3276    #[tokio::test]
3277    async fn resumable_completion_after_corrupt_response_reports_to_change_stream() -> Result<()> {
3278        let (mut backend, producer) = create_test_backend_with_change_stream().await?;
3279        let id = make_id_with_key("resumable-change-stream-after-corrupt-response");
3280        let payload = b"final".to_vec();
3281        let metadata = Metadata::default();
3282        let token = backend
3283            .create_upload_session(&id, &metadata, nonzero(payload.len() as u64))
3284            .await?;
3285        inject_retry_test(
3286            &mut backend,
3287            "storage.objects.insert",
3288            "return-broken-stream-final-chunk-after-0B",
3289        )
3290        .await?;
3291
3292        assert!(matches!(
3293            backend
3294                .put_chunk(
3295                    &token,
3296                    0,
3297                    payload.len() as u64,
3298                    stream::single::<ClientError>(payload.clone()),
3299                )
3300                .await,
3301            Err(error) if error.kind() == ErrorKind::CorruptData
3302        ));
3303        assert_eq!(
3304            backend.upload_offset(&token).await?,
3305            UploadProgress::Complete
3306        );
3307
3308        let records = producer.records();
3309        assert_eq!(records.len(), 1);
3310        assert_eq!(records[0].op_type, OpType::Write);
3311        assert_eq!(
3312            records[0].size,
3313            Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size())
3314        );
3315        Ok(())
3316    }
3317
3318    #[cfg(feature = "storage-cogs")]
3319    #[tokio::test]
3320    async fn change_stream_size_includes_metadata_keys_and_values() -> Result<()> {
3321        let (backend, producer) = create_test_backend_with_change_stream().await?;
3322        let payload = b"tiny".to_vec();
3323
3324        let bare = Metadata::default();
3325        backend
3326            .put_object(
3327                &make_id(),
3328                &bare,
3329                stream::single::<ClientError>(payload.clone()),
3330                Timestamp::now(),
3331            )
3332            .await?;
3333
3334        let annotated = Metadata {
3335            custom: BTreeMap::from_iter([("a-fairly-long-metadata-key".into(), "value".into())]),
3336            ..Default::default()
3337        };
3338        backend
3339            .put_object(
3340                &make_id(),
3341                &annotated,
3342                stream::single::<ClientError>(payload.clone()),
3343                Timestamp::now(),
3344            )
3345            .await?;
3346
3347        let records = producer.records();
3348        assert_eq!(records.len(), 2);
3349
3350        let bare_size = records[0].size.unwrap();
3351        let annotated_size = records[1].size.unwrap();
3352        assert_eq!(
3353            bare_size,
3354            payload.len() as u64,
3355            "default metadata contributes no custom keys"
3356        );
3357        assert!(
3358            annotated_size > bare_size,
3359            "same payload, more metadata: {annotated_size} should exceed {bare_size}"
3360        );
3361
3362        Ok(())
3363    }
3364
3365    #[cfg(feature = "storage-cogs")]
3366    #[tokio::test]
3367    async fn change_stream_reports_nothing_when_the_object_was_already_gone() -> Result<()> {
3368        let (backend, producer) = create_test_backend_with_change_stream().await?;
3369
3370        backend.delete_object(&make_id(), Timestamp::now()).await?;
3371
3372        assert!(producer.records().is_empty());
3373
3374        Ok(())
3375    }
3376
3377    #[cfg(feature = "storage-cogs")]
3378    #[tokio::test]
3379    async fn change_stream_reports_deletes() -> Result<()> {
3380        let (backend, producer) = create_test_backend_with_change_stream().await?;
3381        let id = make_id();
3382
3383        backend
3384            .put_object(
3385                &id,
3386                &Metadata::default(),
3387                stream::single::<ClientError>(b"hi".to_vec()),
3388                Timestamp::now(),
3389            )
3390            .await?;
3391        producer.clear();
3392
3393        backend.delete_object(&id, Timestamp::now()).await?;
3394
3395        let records = producer.records();
3396        assert_eq!(records.len(), 1, "a retried delete must report only once");
3397        assert_eq!(records[0].op_type, OpType::Delete);
3398        assert_eq!(records[0].size, None);
3399
3400        Ok(())
3401    }
3402
3403    #[cfg(feature = "storage-cogs")]
3404    #[tokio::test]
3405    async fn change_stream_reports_expiry_extension_as_an_update() -> Result<()> {
3406        let (backend, producer) = create_test_backend_with_change_stream().await?;
3407        let id = make_id();
3408        let metadata = Metadata {
3409            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_secs(3600)),
3410            time_expires: Some(Timestamp::now() + Duration::from_secs(1)),
3411            ..Default::default()
3412        };
3413
3414        backend
3415            .put_object(
3416                &id,
3417                &metadata,
3418                stream::single::<ClientError>(b"hi".to_vec()),
3419                Timestamp::now(),
3420            )
3421            .await?;
3422        producer.clear();
3423
3424        backend.get_metadata(&id, Timestamp::now()).await?;
3425        assert!(producer.records().is_empty());
3426
3427        backend
3428            .set_expiry(
3429                &id,
3430                ExpiryTarget::At(Timestamp::now() + Duration::from_secs(3600)).into(),
3431                Timestamp::now(),
3432            )
3433            .await?;
3434
3435        let records = producer.records();
3436        assert_eq!(records.len(), 1);
3437        assert_eq!(records[0].op_type, OpType::Update);
3438        assert_eq!(records[0].size, None);
3439        assert!(records[0].expiration_time.is_some());
3440
3441        Ok(())
3442    }
3443}