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