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::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#[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
480#[derive(Debug)]
482struct ResumableUpload {
483 session_uri: Url,
486 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
529fn error_is_retryable(error: &Error) -> bool {
531 matches!(
532 error.kind(),
533 ErrorKind::BackendRateLimited | ErrorKind::BackendTimeout | ErrorKind::BackendUnavailable
534 )
535}
536
537pub 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 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 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 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 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 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 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 #[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 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 #[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 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 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
799fn 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
824async 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 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 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 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 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); }
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 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 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 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 status if status.as_u16() == CLIENT_CLOSED_REQUEST_STATUS => {
1301 response.drain_body().await;
1302 Ok(())
1303 }
1304 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#[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 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 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 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 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 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 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 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 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 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 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let result = backend.get_object(&id, Timestamp::now(), None).await?;
3020 assert!(result.is_none(), "object should not exist after abort");
3021
3022 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}