1use std::future::Future;
4use std::pin::Pin;
5use std::sync::Arc;
6use std::sync::Once;
7use std::time::Duration;
8
9use google_cloud_auth::credentials::anonymous::Builder as AnonymousCredentials;
10use google_cloud_storage::client::{Storage, StorageControl};
11use google_cloud_storage::model::Object;
12use rskit_errors::{AppError, AppResult, ErrorCode};
13use rskit_storage::FileSource;
14use rskit_storage::store::{
15 FileStore, StorageConfig, StorageFactory, StorageRegistry, StoredFile, UploadOptions,
16 content_type_or_default, prefixed_key,
17};
18use serde::{Deserialize, Serialize};
19use tokio::sync::OnceCell;
20
21#[derive(Debug, Clone, Deserialize, Serialize)]
23pub struct Config {
24 pub bucket: String,
26 pub prefix: Option<String>,
28 #[serde(default)]
33 pub anonymous: bool,
34}
35
36struct GcsStore<S = google_cloud_storage::stub::DefaultStorage>
38where
39 S: google_cloud_storage::stub::Storage + 'static,
40{
41 clients: OnceCell<GcsClients<S>>,
42 builder: GcsClientBuilder<S>,
43 config: Config,
44}
45
46struct GcsClients<S>
47where
48 S: google_cloud_storage::stub::Storage + 'static,
49{
50 storage: Storage<S>,
51 control: StorageControl,
52}
53
54type GcsClientFuture<S> = Pin<Box<dyn Future<Output = AppResult<GcsClients<S>>> + Send>>;
55type GcsClientBuilder<S> = Box<dyn Fn() -> GcsClientFuture<S> + Send + Sync>;
56
57impl GcsStore {
58 fn new(config: Config) -> Self {
60 let builder_config = config.clone();
61 let builder = Box::new(move || {
62 let config = builder_config.clone();
63 Box::pin(async move {
64 install_rustls_crypto_provider();
65
66 let (storage, control) = if config.anonymous {
67 let storage = Storage::builder()
68 .with_credentials(AnonymousCredentials::new().build())
69 .build()
70 .await
71 .map_err(|e| config_err("anonymous storage client", e))?;
72 let control = StorageControl::builder()
73 .with_credentials(AnonymousCredentials::new().build())
74 .build()
75 .await
76 .map_err(|e| config_err("anonymous control client", e))?;
77 (storage, control)
78 } else {
79 let storage = Storage::builder()
80 .build()
81 .await
82 .map_err(|e| config_err("storage client", e))?;
83 let control = StorageControl::builder()
84 .build()
85 .await
86 .map_err(|e| config_err("control client", e))?;
87 (storage, control)
88 };
89
90 Ok(GcsClients { storage, control })
91 }) as GcsClientFuture<google_cloud_storage::stub::DefaultStorage>
92 });
93
94 Self {
95 clients: OnceCell::new(),
96 builder,
97 config,
98 }
99 }
100}
101
102fn install_rustls_crypto_provider() {
103 static INSTALL: Once = Once::new();
104 INSTALL.call_once(|| {
105 let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
106 });
107}
108
109fn config_err(stage: &str, e: impl std::fmt::Display) -> AppError {
110 AppError::new(
111 ErrorCode::Internal,
112 format!("GCS {stage} configuration failed: {e}"),
113 )
114}
115
116impl<S> GcsStore<S>
117where
118 S: google_cloud_storage::stub::Storage + 'static,
119{
120 fn full_key(&self, key: &str) -> String {
121 full_key(self.config.prefix.as_deref(), key)
122 }
123
124 fn bucket_resource(&self) -> String {
125 bucket_resource(&self.config.bucket)
126 }
127
128 async fn clients(&self) -> AppResult<&GcsClients<S>> {
129 self.clients.get_or_try_init(|| (self.builder)()).await
130 }
131}
132
133fn full_key(prefix: Option<&str>, key: &str) -> String {
134 prefixed_key(prefix, key)
135}
136
137fn bucket_resource(bucket: &str) -> String {
138 if bucket.starts_with("projects/") {
139 bucket.to_owned()
140 } else {
141 format!("projects/_/buckets/{bucket}")
142 }
143}
144
145fn object_size(size: i64) -> AppResult<u64> {
146 u64::try_from(size).map_err(|_| {
147 AppError::new(
148 ErrorCode::Internal,
149 format!("GCS object size must not be negative: {size}"),
150 )
151 })
152}
153
154fn stored_file_from_object(key: String, obj: Object) -> AppResult<StoredFile> {
155 Ok(
156 StoredFile::new(key, object_size(obj.size)?, Some(&obj.content_type))
157 .with_metadata(obj.metadata),
158 )
159}
160
161#[async_trait::async_trait]
162impl<S> FileStore for GcsStore<S>
163where
164 S: google_cloud_storage::stub::Storage + 'static,
165{
166 async fn upload(
167 &self,
168 source: &FileSource,
169 key: &str,
170 options: UploadOptions,
171 ) -> AppResult<StoredFile> {
172 let data = source.read_all().await?;
173 let size = data.len() as u64;
174 let full_key = self.full_key(key);
175 let bucket = self.bucket_resource();
176 let clients = self.clients().await?;
177
178 let mut request = clients
179 .storage
180 .write_object(bucket, full_key, data)
181 .set_content_type(content_type_or_default(options.content_type()));
182 if !options.metadata.is_empty() {
183 request = request.set_metadata(options.metadata.clone());
184 }
185 Box::pin(request.send_buffered())
186 .await
187 .map_err(|e| AppError::new(ErrorCode::Internal, format!("GCS upload failed: {e}")))?;
188
189 Ok(
190 StoredFile::new(prefixed_key(None, key), size, options.content_type())
191 .with_metadata(options.metadata),
192 )
193 }
194
195 async fn download(&self, key: &str) -> AppResult<FileSource> {
196 let full_key = self.full_key(key);
197 let bucket = self.bucket_resource();
198 let clients = self.clients().await?;
199 let mut response = clients
200 .storage
201 .read_object(bucket, full_key)
202 .send()
203 .await
204 .map_err(|e| AppError::new(ErrorCode::NotFound, format!("GCS download failed: {e}")))?;
205
206 let mut data = Vec::new();
207 while let Some(chunk) = response.next().await {
208 let chunk = chunk.map_err(|e| {
209 AppError::new(
210 ErrorCode::Internal,
211 format!("GCS download stream failed: {e}"),
212 )
213 })?;
214 data.extend_from_slice(&chunk);
215 }
216
217 Ok(FileSource::Bytes(bytes::Bytes::from(data)))
218 }
219
220 async fn delete(&self, key: &str) -> AppResult<()> {
221 let full_key = self.full_key(key);
222 let bucket = self.bucket_resource();
223 let clients = self.clients().await?;
224 clients
225 .control
226 .delete_object()
227 .set_bucket(bucket)
228 .set_object(full_key)
229 .send()
230 .await
231 .map_err(|e| AppError::new(ErrorCode::Internal, format!("GCS delete failed: {e}")))?;
232
233 Ok(())
234 }
235
236 async fn exists(&self, key: &str) -> AppResult<bool> {
237 match self.head(key).await {
238 Ok(_) => Ok(true),
239 Err(_) => Ok(false),
240 }
241 }
242
243 async fn head(&self, key: &str) -> AppResult<StoredFile> {
244 let full_key = self.full_key(key);
245 let clients = self.clients().await?;
246 let obj = clients
247 .control
248 .get_object()
249 .set_bucket(self.bucket_resource())
250 .set_object(full_key)
251 .send()
252 .await
253 .map_err(|e| AppError::new(ErrorCode::NotFound, format!("GCS head failed: {e}")))?;
254
255 stored_file_from_object(prefixed_key(None, key), obj)
256 }
257
258 async fn list(&self, prefix: &str, limit: Option<usize>) -> AppResult<Vec<StoredFile>> {
259 let full_prefix = self.full_key(prefix);
260
261 let max_results = limit.map(i32::try_from).transpose().map_err(|_| {
262 AppError::new(ErrorCode::InvalidInput, "GCS list limit exceeds i32::MAX")
263 })?;
264
265 let clients = self.clients().await?;
266 let mut request = clients
267 .control
268 .list_objects()
269 .set_parent(self.bucket_resource())
270 .set_prefix(full_prefix);
271 if let Some(max_results) = max_results {
272 request = request.set_page_size(max_results);
273 }
274 let resp = request
275 .send()
276 .await
277 .map_err(|e| AppError::new(ErrorCode::Internal, format!("GCS list failed: {e}")))?;
278
279 let items = resp
280 .objects
281 .into_iter()
282 .map(|obj| stored_file_from_object(obj.name.clone(), obj))
283 .collect::<AppResult<Vec<_>>>()?;
284
285 Ok(items)
286 }
287
288 async fn presigned_url(&self, _key: &str, _expires_in: Duration) -> AppResult<String> {
289 Err(AppError::new(
290 ErrorCode::InvalidInput,
291 "GCS presigned URLs are not supported by this backend",
292 ))
293 }
294
295 async fn copy(&self, from_key: &str, to_key: &str) -> AppResult<StoredFile> {
296 let source = self.download(from_key).await?;
297 self.upload(&source, to_key, UploadOptions::new()).await
298 }
299
300 async fn rename(&self, from_key: &str, to_key: &str) -> AppResult<StoredFile> {
301 let result = self.copy(from_key, to_key).await?;
302 self.delete(from_key).await?;
303 Ok(result)
304 }
305}
306
307struct GcsFactory {
308 config: Config,
309}
310
311#[async_trait::async_trait]
312impl StorageFactory for GcsFactory {
313 async fn create(&self, _config: &StorageConfig) -> AppResult<Arc<dyn FileStore>> {
314 Ok(Arc::new(GcsStore::new(self.config.clone())))
315 }
316}
317
318pub fn register(registry: &mut StorageRegistry, config: Config) -> AppResult<()> {
320 registry.register("gcs", Arc::new(GcsFactory { config }))
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326 use google_cloud_gax::options::RequestOptions as GaxRequestOptions;
327 use google_cloud_gax::response::Response;
328 use google_cloud_storage::model::{
329 DeleteObjectRequest, GetObjectRequest, ListObjectsRequest, ListObjectsResponse,
330 ReadObjectRequest,
331 };
332 use google_cloud_storage::model_ext::{ObjectHighlights, WriteObjectRequest};
333 use google_cloud_storage::read_object::ReadObjectResponse;
334 use google_cloud_storage::request_options::RequestOptions as StorageRequestOptions;
335 use google_cloud_storage::streaming_source::StreamingSource;
336 use std::collections::HashMap;
337
338 #[derive(Debug, Default)]
339 struct MockStorage;
340
341 impl google_cloud_storage::stub::Storage for MockStorage {
342 async fn read_object(
343 &self,
344 req: ReadObjectRequest,
345 _options: StorageRequestOptions,
346 ) -> google_cloud_storage::Result<ReadObjectResponse> {
347 assert_eq!(req.bucket, "projects/_/buckets/test-bucket");
348 assert_eq!(req.object, "assets/input.txt");
349 Ok(ReadObjectResponse::from_source(
350 ObjectHighlights::default(),
351 "downloaded",
352 ))
353 }
354
355 async fn write_object_buffered<P>(
356 &self,
357 _payload: P,
358 req: WriteObjectRequest,
359 _options: StorageRequestOptions,
360 ) -> google_cloud_storage::Result<Object>
361 where
362 P: StreamingSource + Send + Sync + 'static,
363 {
364 let resource = req.spec.resource.expect("write resource");
365 assert_eq!(resource.bucket, "projects/_/buckets/test-bucket");
366 assert_eq!(resource.name, "assets/output.txt");
367 assert!(
368 resource.content_type == "text/plain"
369 || resource.content_type == "application/octet-stream"
370 );
371 if resource.content_type == "text/plain" && !resource.metadata.is_empty() {
372 assert_eq!(
373 resource.metadata.get("trace").map(String::as_str),
374 Some("yes")
375 );
376 }
377 Ok(resource)
378 }
379 }
380
381 #[derive(Debug, Default)]
382 struct MockControl;
383
384 impl google_cloud_storage::stub::StorageControl for MockControl {
385 async fn delete_object(
386 &self,
387 req: DeleteObjectRequest,
388 _options: GaxRequestOptions,
389 ) -> google_cloud_storage::Result<Response<()>> {
390 assert_eq!(req.bucket, "projects/_/buckets/test-bucket");
391 assert_eq!(req.object, "assets/input.txt");
392 Ok(Response::from(()))
393 }
394
395 async fn get_object(
396 &self,
397 req: GetObjectRequest,
398 _options: GaxRequestOptions,
399 ) -> google_cloud_storage::Result<Response<Object>> {
400 assert_eq!(req.bucket, "projects/_/buckets/test-bucket");
401 assert_eq!(req.object, "assets/input.txt");
402 Ok(Response::from(object(
403 "assets/input.txt",
404 10,
405 "text/plain",
406 [("origin", "head")],
407 )))
408 }
409
410 async fn list_objects(
411 &self,
412 req: ListObjectsRequest,
413 _options: GaxRequestOptions,
414 ) -> google_cloud_storage::Result<Response<ListObjectsResponse>> {
415 assert_eq!(req.parent, "projects/_/buckets/test-bucket");
416 assert_eq!(req.prefix, "assets/logs");
417 assert_eq!(req.page_size, 2);
418 let response = ListObjectsResponse::new().set_objects([
419 object("assets/logs/a.txt", 1, "text/plain", []),
420 object("assets/logs/b.txt", 2, "text/plain", [("kind", "log")]),
421 ]);
422 Ok(Response::from(response))
423 }
424 }
425
426 fn object<const N: usize>(
427 name: &str,
428 size: i64,
429 content_type: &str,
430 metadata: [(&str, &str); N],
431 ) -> Object {
432 Object::default()
433 .set_bucket("projects/_/buckets/test-bucket")
434 .set_name(name)
435 .set_size(size)
436 .set_content_type(content_type)
437 .set_metadata(
438 metadata
439 .into_iter()
440 .map(|(k, v)| (k.to_owned(), v.to_owned())),
441 )
442 }
443
444 fn unused_builder<S>() -> GcsClientBuilder<S>
445 where
446 S: google_cloud_storage::stub::Storage + 'static,
447 {
448 Box::new(|| {
449 Box::pin(async {
450 Err(AppError::new(
451 ErrorCode::Internal,
452 "test GCS client builder must not be called",
453 ))
454 })
455 })
456 }
457
458 fn store_with<St, Ct>(storage: St, control: Ct) -> GcsStore<St>
459 where
460 St: google_cloud_storage::stub::Storage + 'static,
461 Ct: google_cloud_storage::stub::StorageControl + 'static,
462 {
463 GcsStore {
464 clients: OnceCell::new_with(Some(GcsClients {
465 storage: Storage::from_stub(storage),
466 control: StorageControl::from_stub(control),
467 })),
468 builder: unused_builder(),
469 config: Config {
470 bucket: "test-bucket".into(),
471 prefix: Some("assets".into()),
472 anonymous: true,
473 },
474 }
475 }
476
477 fn test_store() -> GcsStore<MockStorage> {
478 store_with(MockStorage, MockControl)
479 }
480
481 #[test]
482 fn config_deserializes_with_authenticated_default() {
483 let json = r#"{"bucket": "test"}"#;
484 let cfg: Config = serde_json::from_str(json).unwrap();
485 assert_eq!(cfg.bucket, "test");
486 assert!(cfg.prefix.is_none());
487 assert!(!cfg.anonymous);
488 }
489
490 #[test]
491 fn config_deserializes_anonymous_opt_in() {
492 let json = r#"{"bucket": "public-assets", "prefix": "uploads", "anonymous": true}"#;
493 let cfg: Config = serde_json::from_str(json).unwrap();
494 assert_eq!(cfg.bucket, "public-assets");
495 assert_eq!(cfg.prefix.as_deref(), Some("uploads"));
496 assert!(cfg.anonymous);
497 }
498
499 #[test]
500 fn key_and_bucket_helpers_normalize_inputs() {
501 assert_eq!(full_key(Some("/assets/"), "/image.png"), "assets/image.png");
502 assert_eq!(full_key(Some("///"), "image.png"), "image.png");
503 assert_eq!(full_key(None, "/image.png"), "image.png");
504 assert_eq!(
505 bucket_resource("plain-bucket"),
506 "projects/_/buckets/plain-bucket"
507 );
508 assert_eq!(
509 bucket_resource("projects/p/buckets/named"),
510 "projects/p/buckets/named"
511 );
512 }
513
514 #[test]
515 fn object_size_rejects_negative_values() {
516 let err = object_size(-1).unwrap_err();
517 assert_eq!(err.code(), ErrorCode::Internal);
518 assert!(err.message().contains("must not be negative"));
519 }
520
521 #[test]
522 fn stored_file_from_object_preserves_metadata() {
523 let stored = stored_file_from_object(
524 "logical.txt".into(),
525 object("assets/logical.txt", 42, "text/plain", [("owner", "test")]),
526 )
527 .unwrap();
528 assert_eq!(stored.key, "logical.txt");
529 assert_eq!(stored.size, 42);
530 assert_eq!(stored.content_type, "text/plain");
531 assert_eq!(
532 stored.metadata.get("owner").map(String::as_str),
533 Some("test")
534 );
535 }
536
537 #[tokio::test]
538 async fn upload_writes_prefixed_object_and_returns_logical_key() {
539 let store = test_store();
540 let mut metadata = HashMap::new();
541 metadata.insert("trace".into(), "yes".into());
542 let stored = store
543 .upload(
544 &FileSource::Bytes(bytes::Bytes::from_static(b"payload")),
545 "output.txt",
546 UploadOptions::new()
547 .with_content_type("text/plain")
548 .with_metadata(metadata),
549 )
550 .await
551 .unwrap();
552
553 assert_eq!(stored.key, "output.txt");
554 assert_eq!(stored.size, 7);
555 assert_eq!(stored.content_type, "text/plain");
556 assert_eq!(
557 stored.metadata.get("trace").map(String::as_str),
558 Some("yes")
559 );
560 }
561
562 #[tokio::test]
563 async fn upload_defaults_content_type_and_empty_metadata() {
564 let store = test_store();
565 let stored = store
566 .upload(
567 &FileSource::Bytes(bytes::Bytes::from_static(b"payload")),
568 "output.txt",
569 UploadOptions::new(),
570 )
571 .await
572 .unwrap();
573
574 assert_eq!(stored.key, "output.txt");
575 assert_eq!(stored.size, 7);
576 assert_eq!(stored.content_type, "application/octet-stream");
577 assert!(stored.metadata.is_empty());
578 }
579
580 #[tokio::test]
581 async fn upload_accepts_progress_callback() {
582 let store = test_store();
583 let stored = store
584 .upload(
585 &FileSource::Bytes(bytes::Bytes::from_static(b"payload")),
586 "output.txt",
587 UploadOptions::new()
588 .with_content_type("text/plain")
589 .with_progress(Arc::new(|_| {})),
590 )
591 .await
592 .unwrap();
593
594 assert_eq!(stored.key, "output.txt");
595 assert_eq!(stored.size, 7);
596 }
597
598 #[tokio::test]
599 async fn download_reads_all_chunks_into_bytes() {
600 let store = test_store();
601 let source = store.download("input.txt").await.unwrap();
602 let data = source.read_all().await.unwrap();
603 assert_eq!(data.as_ref(), b"downloaded");
604 }
605
606 #[tokio::test]
607 async fn delete_head_exists_and_list_use_control_client() {
608 let store = test_store();
609 store.delete("input.txt").await.unwrap();
610
611 let head = store.head("input.txt").await.unwrap();
612 assert_eq!(head.key, "input.txt");
613 assert_eq!(head.size, 10);
614 assert_eq!(
615 head.metadata.get("origin").map(String::as_str),
616 Some("head")
617 );
618 assert!(store.exists("input.txt").await.unwrap());
619
620 let listed = store.list("logs", Some(2)).await.unwrap();
621 assert_eq!(listed.len(), 2);
622 assert_eq!(listed[0].key, "assets/logs/a.txt");
623 assert_eq!(listed[1].size, 2);
624 }
625
626 #[tokio::test]
627 async fn list_rejects_limits_larger_than_gcs_accepts() {
628 let store = test_store();
629 let err = store
630 .list("logs", Some(i32::MAX as usize + 1))
631 .await
632 .unwrap_err();
633 assert_eq!(err.code(), ErrorCode::InvalidInput);
634 }
635
636 #[tokio::test]
637 async fn presigned_url_reports_unsupported_operation() {
638 let store = test_store();
639 let err = store
640 .presigned_url("input.txt", Duration::from_mins(1))
641 .await
642 .unwrap_err();
643 assert_eq!(err.code(), ErrorCode::InvalidInput);
644 }
645
646 #[tokio::test]
647 async fn copy_and_rename_compose_existing_operations() {
648 let store = test_store();
649 let copied = store.copy("input.txt", "output.txt").await.unwrap();
650 assert_eq!(copied.key, "output.txt");
651 assert_eq!(copied.size, 10);
652
653 let renamed = store.rename("input.txt", "output.txt").await.unwrap();
654 assert_eq!(renamed.key, "output.txt");
655 }
656
657 #[test]
658 fn constructs_offline_without_building_transport() {
659 let store = GcsStore::new(Config {
660 bucket: "public-assets".into(),
661 prefix: Some("uploads".into()),
662 anonymous: true,
663 });
664
665 assert_eq!(store.full_key("image.png"), "uploads/image.png");
666 }
667
668 #[test]
669 fn slash_only_prefix_is_ignored() {
670 let store = GcsStore::new(Config {
671 bucket: "public-assets".into(),
672 prefix: Some("///".into()),
673 anonymous: true,
674 });
675
676 assert_eq!(store.full_key("image.png"), "image.png");
677 }
678
679 #[test]
680 fn register_adds_backend_without_constructing_client() {
681 let mut registry = StorageRegistry::new();
682 register(
683 &mut registry,
684 Config {
685 bucket: "test".into(),
686 prefix: None,
687 anonymous: true,
688 },
689 )
690 .unwrap();
691 assert!(registry.contains("gcs"));
692 }
693
694 #[tokio::test]
695 async fn factory_constructs_anonymous_store_from_config() {
696 let factory = GcsFactory {
697 config: Config {
698 bucket: "public-assets".into(),
699 prefix: Some("uploads".into()),
700 anonymous: true,
701 },
702 };
703
704 factory.create(&StorageConfig::default()).await.unwrap();
705 }
706
707 #[tokio::test]
708 async fn new_builds_anonymous_clients_lazily() {
709 let store = GcsStore::new(Config {
710 bucket: "public-assets".into(),
711 prefix: None,
712 anonymous: true,
713 });
714
715 assert!(store.clients().await.is_ok());
716 }
717
718 #[tokio::test]
719 async fn new_builds_authenticated_clients_lazily() {
720 let store = GcsStore::new(Config {
721 bucket: "private-assets".into(),
722 prefix: Some("uploads".into()),
723 anonymous: false,
724 });
725
726 assert!(store.clients().await.is_ok());
727 }
728
729 fn transport_error() -> google_cloud_storage::Error {
730 google_cloud_gax::error::Error::io(std::io::Error::other("boom"))
731 }
732
733 struct FailingSource;
734
735 impl StreamingSource for FailingSource {
736 type Error = std::io::Error;
737
738 async fn next(&mut self) -> Option<Result<bytes::Bytes, Self::Error>> {
739 Some(Err(std::io::Error::other("chunk boom")))
740 }
741 }
742
743 #[derive(Debug, Default)]
744 struct StreamFailingStorage;
745
746 impl google_cloud_storage::stub::Storage for StreamFailingStorage {
747 async fn read_object(
748 &self,
749 _req: ReadObjectRequest,
750 _options: StorageRequestOptions,
751 ) -> google_cloud_storage::Result<ReadObjectResponse> {
752 Ok(ReadObjectResponse::from_source(
753 ObjectHighlights::default(),
754 FailingSource,
755 ))
756 }
757
758 async fn write_object_buffered<P>(
759 &self,
760 _payload: P,
761 _req: WriteObjectRequest,
762 _options: StorageRequestOptions,
763 ) -> google_cloud_storage::Result<Object>
764 where
765 P: StreamingSource + Send + Sync + 'static,
766 {
767 Err(transport_error())
768 }
769 }
770
771 #[derive(Debug, Default)]
772 struct FailingStorage;
773
774 impl google_cloud_storage::stub::Storage for FailingStorage {
775 async fn read_object(
776 &self,
777 _req: ReadObjectRequest,
778 _options: StorageRequestOptions,
779 ) -> google_cloud_storage::Result<ReadObjectResponse> {
780 Err(transport_error())
781 }
782
783 async fn write_object_buffered<P>(
784 &self,
785 _payload: P,
786 _req: WriteObjectRequest,
787 _options: StorageRequestOptions,
788 ) -> google_cloud_storage::Result<Object>
789 where
790 P: StreamingSource + Send + Sync + 'static,
791 {
792 Err(transport_error())
793 }
794 }
795
796 #[derive(Debug, Default)]
797 struct FailingControl;
798
799 impl google_cloud_storage::stub::StorageControl for FailingControl {
800 async fn delete_object(
801 &self,
802 _req: DeleteObjectRequest,
803 _options: GaxRequestOptions,
804 ) -> google_cloud_storage::Result<Response<()>> {
805 Err(transport_error())
806 }
807
808 async fn get_object(
809 &self,
810 _req: GetObjectRequest,
811 _options: GaxRequestOptions,
812 ) -> google_cloud_storage::Result<Response<Object>> {
813 Err(transport_error())
814 }
815
816 async fn list_objects(
817 &self,
818 _req: ListObjectsRequest,
819 _options: GaxRequestOptions,
820 ) -> google_cloud_storage::Result<Response<ListObjectsResponse>> {
821 Err(transport_error())
822 }
823 }
824
825 fn failing_store() -> GcsStore<FailingStorage> {
826 store_with(FailingStorage, FailingControl)
827 }
828
829 #[tokio::test]
830 async fn upload_maps_transport_failure_to_internal_error() {
831 let err = failing_store()
832 .upload(
833 &FileSource::Bytes(bytes::Bytes::from_static(b"payload")),
834 "output.txt",
835 UploadOptions::new().with_content_type("text/plain"),
836 )
837 .await
838 .unwrap_err();
839 assert_eq!(err.code(), ErrorCode::Internal);
840 assert!(err.message().contains("GCS upload failed"));
841 }
842
843 #[tokio::test]
844 async fn download_maps_transport_failure_to_not_found_error() {
845 let err = failing_store().download("input.txt").await.unwrap_err();
846 assert_eq!(err.code(), ErrorCode::NotFound);
847 assert!(err.message().contains("GCS download failed"));
848 }
849
850 #[tokio::test]
851 async fn download_maps_stream_failure_to_internal_error() {
852 let store = store_with(StreamFailingStorage, MockControl);
853
854 let err = store.download("input.txt").await.unwrap_err();
855 assert_eq!(err.code(), ErrorCode::Internal);
856 assert!(err.message().contains("GCS download stream failed"));
857 }
858
859 #[tokio::test]
860 async fn delete_maps_transport_failure_to_internal_error() {
861 let err = failing_store().delete("input.txt").await.unwrap_err();
862 assert_eq!(err.code(), ErrorCode::Internal);
863 assert!(err.message().contains("GCS delete failed"));
864 }
865
866 #[tokio::test]
867 async fn head_maps_transport_failure_to_not_found_error() {
868 let err = failing_store().head("input.txt").await.unwrap_err();
869 assert_eq!(err.code(), ErrorCode::NotFound);
870 assert!(err.message().contains("GCS head failed"));
871 assert!(!failing_store().exists("input.txt").await.unwrap());
873 }
874
875 #[tokio::test]
876 async fn list_maps_transport_failure_to_internal_error() {
877 let err = failing_store().list("logs", Some(2)).await.unwrap_err();
878 assert_eq!(err.code(), ErrorCode::Internal);
879 assert!(err.message().contains("GCS list failed"));
880 }
881
882 #[tokio::test]
883 async fn operations_propagate_client_builder_failure() {
884 let store: GcsStore<MockStorage> = GcsStore {
885 clients: OnceCell::new(),
886 builder: Box::new(|| {
887 Box::pin(async { Err(AppError::new(ErrorCode::Internal, "builder unavailable")) })
888 }),
889 config: Config {
890 bucket: "test-bucket".into(),
891 prefix: Some("assets".into()),
892 anonymous: true,
893 },
894 };
895
896 let err = store.head("input.txt").await.unwrap_err();
897 assert_eq!(err.code(), ErrorCode::Internal);
898 assert!(err.message().contains("builder unavailable"));
899 }
900
901 #[test]
902 fn config_err_formats_stage_and_source() {
903 let err = config_err("storage client", "no credentials");
904 assert_eq!(err.code(), ErrorCode::Internal);
905 assert_eq!(
906 err.message(),
907 "GCS storage client configuration failed: no credentials"
908 );
909 }
910}