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