1use super::tracing::{TracingObjectDescriptor, TracingResponse};
16use crate::Result;
17use crate::model::{Object, ReadObjectRequest};
18use crate::model_ext::WriteObjectRequest;
19use crate::read_object::ReadObjectResponse;
20use crate::storage::client::StorageInner;
21use crate::storage::info::INSTRUMENTATION;
22use crate::storage::perform_upload::PerformUpload;
23use crate::storage::read_object::Reader;
24use crate::storage::request_options::RequestOptions;
25use crate::storage::streaming_source::{Seek, StreamingSource};
26#[cfg(google_cloud_unstable_storage_bidi)]
27use crate::{
28 appendable_object_writer::AppendableObjectWriter,
29 model_ext::{OpenAppendableObjectRequest, ReopenAppendableObjectRequest},
30 storage::bidi_write::{
31 connector::Connector as BidiWriteConnector, transport::AppendableObjectWriterTransport,
32 },
33};
34use crate::{
35 model_ext::OpenObjectRequest, object_descriptor::ObjectDescriptor,
36 storage::bidi::connector::Connector, storage::bidi::transport::ObjectDescriptorTransport,
37};
38use gaxi::observability::{ClientRequestAttributes, DurationMetric, RequestRecorder};
39use std::sync::Arc;
40
41#[derive(Clone, Debug)]
57pub struct Storage {
58 inner: Arc<StorageInner>,
59 tracing: bool,
60 metric: DurationMetric,
61}
62
63impl Storage {
64 pub(crate) fn new(inner: Arc<StorageInner>, tracing: bool) -> Arc<Self> {
65 let metric = DurationMetric::new(&INSTRUMENTATION);
66 Arc::new(Self {
67 inner,
68 tracing,
69 metric,
70 })
71 }
72
73 async fn read_object_plain(
74 &self,
75 request: ReadObjectRequest,
76 options: RequestOptions,
77 ) -> Result<ReadObjectResponse> {
78 let reader = Reader {
79 inner: self.inner.clone(),
80 request,
81 options,
82 };
83 reader.response().await
84 }
85
86 #[tracing::instrument(name = "read_object", level = tracing::Level::DEBUG, ret, err(Debug))]
87 async fn read_object_tracing(
88 &self,
89 request: ReadObjectRequest,
90 options: RequestOptions,
91 ) -> Result<ReadObjectResponse> {
92 let resource_name = format!("//storage.googleapis.com/{}", request.bucket);
93 let (span, pending) = gaxi::client_request_signals!(
94 metric: self.metric.clone(),
95 info: *INSTRUMENTATION,
96 method: "client::Storage::read_object",
97 async {
98 if let Some(recorder) = RequestRecorder::current() {
99 recorder.on_client_request(
100 ClientRequestAttributes::default()
101 .set_url_template("/storage/v1/b/{bucket}/o/{object}")
102 .set_resource_name(resource_name),
103 );
104 }
105 self.read_object_plain(request, options).await
106 });
107
108 let response = pending.await?;
109 let inner = TracingResponse::new(response.into_parts(), span);
110 Ok(ReadObjectResponse::new(Box::new(inner)))
111 }
112
113 async fn write_object_buffered_plain<P>(
114 &self,
115 payload: P,
116 request: WriteObjectRequest,
117 options: RequestOptions,
118 ) -> Result<Object>
119 where
120 P: StreamingSource + Send + Sync + 'static,
121 {
122 PerformUpload::new(
123 payload,
124 self.inner.clone(),
125 request.spec,
126 request.params,
127 options,
128 )
129 .send()
130 .await
131 }
132
133 #[tracing::instrument(name = "write_object_buffered", level = tracing::Level::DEBUG, ret, err(Debug), skip(payload))]
134 async fn write_object_buffered_tracing<P>(
135 &self,
136 payload: P,
137 request: WriteObjectRequest,
138 options: RequestOptions,
139 ) -> Result<Object>
140 where
141 P: StreamingSource + Send + Sync + 'static,
142 {
143 let resource_name = format!(
144 "//storage.googleapis.com/{}",
145 request
146 .spec
147 .resource
148 .as_ref()
149 .map(|r| r.bucket.as_str())
150 .unwrap_or_default()
151 );
152 let (_span, pending) = gaxi::client_request_signals!(
153 metric: self.metric.clone(),
154 info: *INSTRUMENTATION,
155 method: "client::Storage::write_object",
156 async {
157 if let Some(recorder) = RequestRecorder::current() {
158 recorder.on_client_request(
159 ClientRequestAttributes::default()
160 .set_url_template("/upload/storage/v1/b/{bucket}/o")
161 .set_resource_name(resource_name),
162 );
163 }
164 self.write_object_buffered_plain(payload, request, options).await
165 }
166 );
167 pending.await
168 }
169
170 async fn write_object_unbuffered_plain<P>(
171 &self,
172 payload: P,
173 request: WriteObjectRequest,
174 options: RequestOptions,
175 ) -> Result<Object>
176 where
177 P: StreamingSource + Seek + Send + Sync + 'static,
178 {
179 PerformUpload::new(
180 payload,
181 self.inner.clone(),
182 request.spec,
183 request.params,
184 options,
185 )
186 .send_unbuffered()
187 .await
188 }
189
190 #[tracing::instrument(name = "write_object_unbuffered", level = tracing::Level::DEBUG, ret, err(Debug), skip(payload))]
191 async fn write_object_unbuffered_tracing<P>(
192 &self,
193 payload: P,
194 request: WriteObjectRequest,
195 options: RequestOptions,
196 ) -> Result<Object>
197 where
198 P: StreamingSource + Seek + Send + Sync + 'static,
199 {
200 let resource_name = format!(
201 "//storage.googleapis.com/{}",
202 request
203 .spec
204 .resource
205 .as_ref()
206 .map(|r| r.bucket.as_str())
207 .unwrap_or_default()
208 );
209 let (_span, pending) = gaxi::client_request_signals!(
210 metric: self.metric.clone(),
211 info: *INSTRUMENTATION,
212 method: "client::Storage::write_object",
213 async {
214 if let Some(recorder) = RequestRecorder::current() {
215 recorder.on_client_request(
216 ClientRequestAttributes::default()
217 .set_url_template("/upload/storage/v1/b/{bucket}/o")
218 .set_resource_name(resource_name),
219 );
220 }
221 self.write_object_unbuffered_plain(payload, request, options).await
222 }
223 );
224 pending.await
225 }
226
227 async fn open_object_plain(
228 &self,
229 request: OpenObjectRequest,
230 options: RequestOptions,
231 ) -> Result<(ObjectDescriptor, Vec<ReadObjectResponse>)> {
232 let (spec, ranges) = request.into_parts();
233 let connector = Connector::new(spec, options, self.inner.grpc.clone());
234 let (transport, readers) = ObjectDescriptorTransport::new(connector, ranges).await?;
235 Ok((ObjectDescriptor::new(transport), readers))
236 }
237
238 #[tracing::instrument(name = "open_object", level = tracing::Level::DEBUG, ret, err(Debug))]
239 async fn open_object_tracing(
240 &self,
241 request: OpenObjectRequest,
242 options: RequestOptions,
243 ) -> Result<(ObjectDescriptor, Vec<ReadObjectResponse>)> {
244 let resource_name = format!("//storage.googleapis.com/{}", request.bucket);
245 let (span, pending) = gaxi::client_request_signals!(
246 metric: self.metric.clone(),
247 info: *INSTRUMENTATION,
248 method: "client::Storage::open_object",
249 async {
250 if let Some(recorder) = RequestRecorder::current() {
251 recorder.on_client_request(
252 ClientRequestAttributes::default()
253 .set_rpc_method("google.storage.v2.Storage/BidiStreamingRead")
254 .set_url_template("/upload/storage/v1/b/{bucket}/o")
255 .set_resource_name(resource_name),
256 );
257 }
258 self.open_object_plain(request, options).await
259 }
260 );
261 let (descriptor, readers) = pending.await?;
262 let descriptor =
263 ObjectDescriptor::new(TracingObjectDescriptor::new(descriptor.into_parts()));
264 let readers = readers
265 .into_iter()
266 .map(|r| {
267 let inner = r.into_parts();
268 ReadObjectResponse::new(Box::new(TracingResponse::new(inner, span.clone())))
269 })
270 .collect::<Vec<_>>();
271 Ok((descriptor, readers))
272 }
273
274 #[cfg(google_cloud_unstable_storage_bidi)]
275 async fn open_appendable_object_plain(
276 &self,
277 request: OpenAppendableObjectRequest,
278 options: RequestOptions,
279 ) -> Result<AppendableObjectWriter> {
280 let connector = BidiWriteConnector::new(options, self.inner.grpc.clone());
281 let transport = AppendableObjectWriterTransport::new_open(connector, request).await?;
282 Ok(AppendableObjectWriter::new(transport))
283 }
284
285 #[cfg(google_cloud_unstable_storage_bidi)]
286 #[tracing::instrument(name = "open_appendable_object", level = tracing::Level::DEBUG, ret, err(Debug))]
287 async fn open_appendable_object_tracing(
288 &self,
289 request: OpenAppendableObjectRequest,
290 options: RequestOptions,
291 ) -> Result<AppendableObjectWriter> {
292 let resource_name = format!(
293 "//storage.googleapis.com/{}",
294 request
295 .spec
296 .resource
297 .as_ref()
298 .map(|r| r.bucket.as_str())
299 .unwrap_or_default()
300 );
301 let (_span, pending) = gaxi::client_request_signals!(
302 metric: self.metric.clone(),
303 info: *INSTRUMENTATION,
304 method: "client::Storage::open_appendable_object",
305 async {
306 if let Some(recorder) = RequestRecorder::current() {
307 recorder.on_client_request(
308 ClientRequestAttributes::default()
309 .set_rpc_method("google.storage.v2.Storage/BidiWriteObject")
310 .set_url_template("/upload/storage/v1/b/{bucket}/o")
311 .set_resource_name(resource_name),
312 );
313 }
314 self.open_appendable_object_plain(request, options).await
315 }
316 );
317 let writer = pending.await?;
318 Ok(AppendableObjectWriter::new(
319 super::tracing::TracingAppendableObjectWriter::new(writer.into_parts()),
320 ))
321 }
322
323 #[cfg(google_cloud_unstable_storage_bidi)]
324 async fn reopen_appendable_object_plain(
325 &self,
326 request: ReopenAppendableObjectRequest,
327 options: RequestOptions,
328 ) -> Result<AppendableObjectWriter> {
329 let connector = BidiWriteConnector::new(options, self.inner.grpc.clone());
330 let transport = AppendableObjectWriterTransport::new_reopen(connector, request).await?;
331 Ok(AppendableObjectWriter::new(transport))
332 }
333
334 #[cfg(google_cloud_unstable_storage_bidi)]
335 #[tracing::instrument(name = "reopen_appendable_object", level = tracing::Level::DEBUG, ret, err(Debug))]
336 async fn reopen_appendable_object_tracing(
337 &self,
338 request: ReopenAppendableObjectRequest,
339 options: RequestOptions,
340 ) -> Result<AppendableObjectWriter> {
341 let resource_name = format!("//storage.googleapis.com/{}", request.bucket);
342 let (_span, pending) = gaxi::client_request_signals!(
343 metric: self.metric.clone(),
344 info: *INSTRUMENTATION,
345 method: "client::Storage::reopen_appendable_object",
346 async {
347 if let Some(recorder) = RequestRecorder::current() {
348 recorder.on_client_request(
349 ClientRequestAttributes::default()
350 .set_rpc_method("google.storage.v2.Storage/BidiWriteObject")
351 .set_url_template("/upload/storage/v1/b/{bucket}/o")
352 .set_resource_name(resource_name),
353 );
354 }
355 self.reopen_appendable_object_plain(request, options).await
356 }
357 );
358 let writer = pending.await?;
359 Ok(AppendableObjectWriter::new(
360 super::tracing::TracingAppendableObjectWriter::new(writer.into_parts()),
361 ))
362 }
363}
364
365impl super::stub::Storage for Storage {
366 async fn read_object(
368 &self,
369 req: ReadObjectRequest,
370 options: RequestOptions,
371 ) -> Result<ReadObjectResponse> {
372 if self.tracing {
373 return self.read_object_tracing(req, options).await;
374 }
375 self.read_object_plain(req, options).await
376 }
377
378 async fn write_object_buffered<P>(
380 &self,
381 payload: P,
382 req: WriteObjectRequest,
383 options: RequestOptions,
384 ) -> Result<Object>
385 where
386 P: StreamingSource + Send + Sync + 'static,
387 {
388 if self.tracing {
389 return self
390 .write_object_buffered_tracing(payload, req, options)
391 .await;
392 }
393 self.write_object_buffered_plain(payload, req, options)
394 .await
395 }
396
397 async fn write_object_unbuffered<P>(
399 &self,
400 payload: P,
401 req: WriteObjectRequest,
402 options: RequestOptions,
403 ) -> Result<Object>
404 where
405 P: StreamingSource + Seek + Send + Sync + 'static,
406 {
407 if self.tracing {
408 return self
409 .write_object_unbuffered_tracing(payload, req, options)
410 .await;
411 }
412 self.write_object_unbuffered_plain(payload, req, options)
413 .await
414 }
415
416 async fn open_object(
417 &self,
418 request: OpenObjectRequest,
419 options: RequestOptions,
420 ) -> Result<(ObjectDescriptor, Vec<ReadObjectResponse>)> {
421 if self.tracing {
422 return self.open_object_tracing(request, options).await;
423 }
424 self.open_object_plain(request, options).await
425 }
426
427 #[cfg(google_cloud_unstable_storage_bidi)]
428 async fn open_appendable_object(
429 &self,
430 request: OpenAppendableObjectRequest,
431 options: RequestOptions,
432 ) -> Result<AppendableObjectWriter> {
433 if self.tracing {
434 return self.open_appendable_object_tracing(request, options).await;
435 }
436 self.open_appendable_object_plain(request, options).await
437 }
438
439 #[cfg(google_cloud_unstable_storage_bidi)]
440 async fn reopen_appendable_object(
441 &self,
442 request: ReopenAppendableObjectRequest,
443 options: RequestOptions,
444 ) -> Result<AppendableObjectWriter> {
445 if self.tracing {
446 return self
447 .reopen_appendable_object_tracing(request, options)
448 .await;
449 }
450 self.reopen_appendable_object_plain(request, options).await
451 }
452}
453
454#[cfg(test)]
455mod tests {
456 use super::{Storage, StorageInner};
457 use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
458 use google_cloud_test_utils::test_layer::AttributeValue;
459 use google_cloud_test_utils::test_layer::{CapturedSpan, TestLayer};
460 use httptest::{Expectation, Server, matchers::*, responders::status_code};
461 use pretty_assertions::assert_eq;
462 use std::collections::BTreeMap;
463 use std::sync::Arc;
464
465 impl Storage {
466 pub(crate) fn new_test(inner: Arc<StorageInner>) -> Arc<Self> {
467 Self::new(inner, false)
468 }
469 }
470
471 #[tokio::test]
472 async fn read_object() -> anyhow::Result<()> {
473 let guard = TestLayer::initialize();
474
475 let server = Server::run();
476 server.expect(
477 Expectation::matching(all_of![
478 request::method_path("GET", "/storage/v1/b/test-bucket/o/test-object"),
479 request::query(url_decoded(contains(("alt", "media")))),
480 ])
481 .respond_with(status_code(404)),
482 );
483
484 let client = crate::client::Storage::builder()
485 .with_endpoint(format!("http://{}", server.addr()))
486 .with_credentials(Anonymous::new().build())
487 .with_tracing()
488 .build()
489 .await?;
490 let response = client
491 .read_object("projects/_/buckets/test-bucket", "test-object")
492 .send()
493 .await;
494 assert!(
495 matches!(response, Err(ref e) if e.is_transport()),
496 "{response:?}"
497 );
498
499 let captured = TestLayer::capture(&guard);
500 check_debug_log(&captured, "read_object");
501
502 client_request_span(&captured, "read_object", "404", "http");
503
504 Ok(())
505 }
506
507 #[tokio::test]
508 async fn read_object_success() -> anyhow::Result<()> {
509 let guard = TestLayer::initialize();
510
511 let body = (0..100_000)
512 .map(|i| format!("{i:08} {:1000}", ""))
513 .collect::<Vec<_>>()
514 .join("\n");
515 let server = Server::run();
516 server.expect(
517 Expectation::matching(all_of![
518 request::method_path("GET", "/storage/v1/b/test-bucket/o/test-object"),
519 request::query(url_decoded(contains(("alt", "media")))),
520 ])
521 .respond_with(
522 status_code(200)
523 .body(body.clone())
524 .append_header("x-goog-generation", 123456),
525 ),
526 );
527
528 let client = crate::client::Storage::builder()
529 .with_endpoint(format!("http://{}", server.addr()))
530 .with_credentials(Anonymous::new().build())
531 .with_tracing()
532 .build()
533 .await?;
534 let mut got = Vec::new();
535 let mut response = client
536 .read_object("projects/_/buckets/test-bucket", "test-object")
537 .send()
538 .await?;
539 let object = response.object();
540 assert_eq!(object.generation, 123456, "{object:?}");
541 while let Some(b) = response.next().await.transpose()? {
542 got.push(b);
543 }
544
545 let captured = TestLayer::capture(&guard);
546 let span = captured
547 .iter()
548 .find(|s| s.name == "client_request")
549 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
550 assert_eq!(span.events, got.len() + 1, "{span:?}");
552
553 Ok(())
554 }
555
556 #[tokio::test]
557 async fn write_object_buffered() -> anyhow::Result<()> {
558 let guard = TestLayer::initialize();
559
560 let server = Server::run();
561 server.expect(
562 Expectation::matching(all_of![
563 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
564 request::query(url_decoded(contains(("uploadType", "multipart")))),
565 ])
566 .respond_with(status_code(404)),
567 );
568
569 let client = crate::client::Storage::builder()
570 .with_endpoint(format!("http://{}", server.addr()))
571 .with_credentials(Anonymous::new().build())
572 .with_tracing()
573 .build()
574 .await?;
575 let response = client
576 .write_object("projects/_/buckets/test-bucket", "test-object", "payload")
577 .send_buffered()
578 .await;
579 assert!(
580 matches!(response, Err(ref e) if e.is_transport()),
581 "{response:?}"
582 );
583
584 let captured = TestLayer::capture(&guard);
585 check_debug_log(&captured, "write_object_buffered");
586
587 client_request_span(&captured, "write_object", "404", "http");
588
589 Ok(())
590 }
591
592 #[tokio::test]
593 async fn write_object_unbuffered() -> anyhow::Result<()> {
594 let guard = TestLayer::initialize();
595
596 let server = Server::run();
597 server.expect(
598 Expectation::matching(all_of![
599 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
600 request::query(url_decoded(contains(("uploadType", "multipart")))),
601 ])
602 .respond_with(status_code(404)),
603 );
604
605 let client = crate::client::Storage::builder()
606 .with_endpoint(format!("http://{}", server.addr()))
607 .with_credentials(Anonymous::new().build())
608 .with_tracing()
609 .build()
610 .await?;
611 let response = client
612 .write_object("projects/_/buckets/test-bucket", "test-object", "payload")
613 .send_unbuffered()
614 .await;
615 assert!(
616 matches!(response, Err(ref e) if e.is_transport()),
617 "{response:?}"
618 );
619
620 let captured = TestLayer::capture(&guard);
621 check_debug_log(&captured, "write_object_unbuffered");
622
623 client_request_span(&captured, "write_object", "404", "http");
624
625 Ok(())
626 }
627
628 #[tokio::test]
629 async fn open_object() -> anyhow::Result<()> {
630 use gaxi::grpc::tonic::Status as TonicStatus;
631 use google_cloud_gax::error::rpc::Code;
632 use storage_grpc_mock::{MockStorage, start};
633
634 let guard = TestLayer::initialize();
635
636 let mut mock = MockStorage::new();
637 mock.expect_bidi_read_object()
638 .return_once(|_| Err(TonicStatus::not_found("not here")));
639 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
640
641 let client = crate::client::Storage::builder()
642 .with_credentials(Anonymous::new().build())
643 .with_endpoint(endpoint.clone())
644 .with_tracing()
645 .build()
646 .await?;
647 let response = client
648 .open_object("projects/_/buckets/test-bucket", "test-object")
649 .send()
650 .await;
651 assert!(
652 matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
653 "{response:?}"
654 );
655
656 let captured = TestLayer::capture(&guard);
657 check_debug_log(&captured, "open_object");
658
659 client_request_span(&captured, "open_object", "NOT_FOUND", "grpc");
660 Ok(())
661 }
662
663 #[tokio::test]
664 #[ignore = "flaky test, see #5290"]
665 async fn open_object_success() -> anyhow::Result<()> {
666 use crate::model_ext::ReadRange;
668 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
669 use storage_grpc_mock::google::storage::v2::{
670 BidiReadObjectResponse, ChecksummedData, Object as ProtoObject, ObjectRangeData,
671 ReadRange as ProtoRange,
672 };
673 use storage_grpc_mock::{MockStorage, start};
674 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
675 const OBJECT_NAME: &str = "test-object";
676 const BIND_ADDRESS: &str = "0.0.0.0:0";
677 const PAYLOAD: &str = "the quick brown fox jumps over the lazy dog";
678
679 let guard = TestLayer::initialize();
680
681 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiReadObjectResponse>>(10);
682 let response = BidiReadObjectResponse {
683 metadata: Some(ProtoObject {
684 bucket: BUCKET_NAME.to_string(),
685 name: OBJECT_NAME.to_string(),
686 generation: 123456,
687 ..ProtoObject::default()
688 }),
689 object_data_ranges: vec![ObjectRangeData {
690 read_range: Some(ProtoRange {
691 read_id: 0_i64,
692 ..ProtoRange::default()
693 }),
694 range_end: true,
695 checksummed_data: Some(ChecksummedData {
696 content: PAYLOAD.as_bytes().to_vec(),
697 crc32c: None,
698 }),
699 }],
700 ..BidiReadObjectResponse::default()
701 };
702 tx.send(Ok(response.clone())).await?;
704 tx.send(Ok(response.clone())).await?;
707 tx.send(Ok(response.clone())).await?;
708
709 let mut mock = MockStorage::new();
710 mock.expect_bidi_read_object()
711 .return_once(|_| Ok(TonicResponse::from(rx)));
712 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
713
714 let client = crate::client::Storage::builder()
715 .with_credentials(Anonymous::new().build())
716 .with_endpoint(endpoint.clone())
717 .with_tracing()
718 .build()
719 .await?;
720 let (descriptor, _reader0) = client
721 .open_object(BUCKET_NAME, OBJECT_NAME)
722 .send_and_read(ReadRange::all())
723 .await?;
724 let _reader1 = descriptor.read_range(ReadRange::offset(5)).await;
725 let _reader2 = descriptor.read_range(ReadRange::segment(10, 10)).await;
726 let _reader3 = descriptor.read_range(ReadRange::tail(15)).await;
727
728 let captured = TestLayer::capture(&guard);
729 let _span = captured
730 .iter()
731 .find(|s| s.name == "client_request")
732 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
733
734 let range_spans = captured
735 .iter()
736 .filter(|s| s.name == "read_range")
737 .collect::<Vec<_>>();
738
739 let _span_reader1 = range_spans
740 .clone()
741 .into_iter()
742 .find(|s| {
743 s.attributes
744 .get("read_range.start")
745 .and_then(|v| v.as_i64())
746 == Some(5)
747 })
748 .unwrap_or_else(|| {
749 panic!("missing `read_range` span for ReadRange::offset(5): {range_spans:#?}")
750 });
751
752 let _span_reader2 = range_spans
753 .clone()
754 .into_iter()
755 .find(|s| {
756 s.attributes
757 .get("read_range.start")
758 .and_then(|v| v.as_i64())
759 == Some(10)
760 && s.attributes
761 .get("read_range.limit")
762 .and_then(|v| v.as_i64())
763 == Some(10)
764 })
765 .unwrap_or_else(|| {
766 panic!("missing `read_range` span for ReadRange::segment(10, 10): {range_spans:#?}")
767 });
768
769 let _span_reader3 = range_spans
770 .clone()
771 .into_iter()
772 .find(|s| {
773 s.attributes
774 .get("read_range.start")
775 .and_then(|v| v.as_i64())
776 == Some(-15)
777 })
778 .unwrap_or_else(|| {
779 panic!("missing `read_range` span for ReadRange::tail(15): {range_spans:#?}")
780 });
781 Ok(())
782 }
783
784 #[cfg(google_cloud_unstable_storage_bidi)]
785 #[tokio::test]
786 async fn open_appendable_object_not_found() -> anyhow::Result<()> {
787 use gaxi::grpc::tonic::Status as TonicStatus;
788 use google_cloud_gax::error::rpc::Code;
789 use storage_grpc_mock::{MockStorage, start};
790
791 let guard = TestLayer::initialize();
792 let mut mock = MockStorage::new();
793 mock.expect_bidi_write_object()
794 .return_once(|_| Err(TonicStatus::not_found("not here")));
795 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
796
797 let client = crate::client::Storage::builder()
798 .with_credentials(Anonymous::new().build())
799 .with_endpoint(endpoint.clone())
800 .with_tracing()
801 .build()
802 .await?;
803 let response = client
804 .open_appendable_object("projects/_/buckets/test-bucket", "test-object")
805 .send()
806 .await;
807 assert!(
808 matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
809 "{response:?}"
810 );
811 let captured = TestLayer::capture(&guard);
812 check_debug_log(&captured, "open_appendable_object");
813
814 client_request_span(&captured, "open_appendable_object", "NOT_FOUND", "grpc");
815
816 Ok(())
817 }
818
819 #[ignore = "TODO(#6324) - disabled because it was flaky"]
821 #[cfg(google_cloud_unstable_storage_bidi)]
822 #[tokio::test]
823 async fn open_appendable_object_success() -> anyhow::Result<()> {
824 use bytes::Bytes;
825 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
826 use storage_grpc_mock::google::storage::v2::{
827 BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
828 };
829 use storage_grpc_mock::{MockStorage, start};
830 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
831 const OBJECT_NAME: &str = "test-object";
832 const BIND_ADDRESS: &str = "0.0.0.0:0";
833
834 let guard = TestLayer::initialize();
835
836 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
837
838 let initial_response = BidiWriteObjectResponse {
839 write_status: Some(WriteStatus::Resource(Object {
840 bucket: BUCKET_NAME.to_string(),
841 name: OBJECT_NAME.to_string(),
842 generation: 98765,
843 ..Default::default()
844 })),
845 ..BidiWriteObjectResponse::default()
846 };
847
848 let flush_response = BidiWriteObjectResponse {
849 write_status: Some(WriteStatus::PersistedSize(5)),
850 ..BidiWriteObjectResponse::default()
851 };
852
853 tx.send(Ok(initial_response)).await?;
856
857 let mut mock = MockStorage::new();
858 mock.expect_bidi_write_object().return_once(move |req| {
859 let mut stream = req.into_inner();
860 tokio::spawn(async move {
861 while let Some(Ok(msg)) = stream.recv().await {
862 if msg.flush {
863 let _ = tx.send(Ok(flush_response.clone())).await;
866 }
867 }
868 });
869 Ok(TonicResponse::new(rx))
870 });
871 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
872
873 let client = crate::client::Storage::builder()
874 .with_credentials(Anonymous::new().build())
875 .with_endpoint(endpoint.clone())
876 .with_tracing()
877 .build()
878 .await?;
879 let mut writer = client
880 .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
881 .send()
882 .await?;
883
884 writer.append(Bytes::from_static(b"hello")).await?;
885 writer.flush().await?;
886 writer.close().await?;
887
888 let captured = TestLayer::capture(&guard);
889 let _span = captured
890 .iter()
891 .find(|s| s.name == "client_request")
892 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
893
894 check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
895 check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
896 check_bidi_write_span_attributes(&captured, "close", 98765, Some(5));
897
898 Ok(())
899 }
900
901 #[ignore = "TODO(#6324) - disabled because it was flaky"]
903 #[cfg(google_cloud_unstable_storage_bidi)]
904 #[tokio::test]
905 async fn open_appendable_object_finalize_success() -> anyhow::Result<()> {
906 use bytes::Bytes;
907 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
908 use storage_grpc_mock::google::storage::v2::{
909 BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
910 };
911 use storage_grpc_mock::{MockStorage, start};
912 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
913 const OBJECT_NAME: &str = "test-object";
914 const BIND_ADDRESS: &str = "0.0.0.0:0";
915
916 let guard = TestLayer::initialize();
917
918 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
919
920 let initial_response = BidiWriteObjectResponse {
921 write_status: Some(WriteStatus::Resource(Object {
922 bucket: BUCKET_NAME.to_string(),
923 name: OBJECT_NAME.to_string(),
924 generation: 98765,
925 ..Default::default()
926 })),
927 ..BidiWriteObjectResponse::default()
928 };
929
930 let flush_response = BidiWriteObjectResponse {
931 write_status: Some(WriteStatus::PersistedSize(5)),
932 ..BidiWriteObjectResponse::default()
933 };
934
935 let finalize_response = BidiWriteObjectResponse {
936 write_status: Some(WriteStatus::Resource(Object {
937 bucket: BUCKET_NAME.to_string(),
938 name: OBJECT_NAME.to_string(),
939 generation: 98765,
940 size: 5,
941 ..Default::default()
942 })),
943 ..BidiWriteObjectResponse::default()
944 };
945
946 tx.send(Ok(initial_response)).await?;
947
948 let mut mock = MockStorage::new();
949 mock.expect_bidi_write_object().return_once(move |req| {
950 let mut stream = req.into_inner();
951 tokio::spawn(async move {
952 while let Some(Ok(msg)) = stream.recv().await {
953 if msg.finish_write {
954 let _ = tx.send(Ok(finalize_response.clone())).await;
955 } else if msg.flush {
956 let _ = tx.send(Ok(flush_response.clone())).await;
957 }
958 }
959 });
960 Ok(TonicResponse::new(rx))
961 });
962 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
963
964 let client = crate::client::Storage::builder()
965 .with_credentials(Anonymous::new().build())
966 .with_endpoint(endpoint.clone())
967 .with_tracing()
968 .build()
969 .await?;
970 let mut writer = client
971 .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
972 .send()
973 .await?;
974
975 writer.append(Bytes::from_static(b"hello")).await?;
976 writer.flush().await?;
977 let obj = writer.finalize().await?;
978 assert_eq!(obj.size, 5);
979
980 let captured = TestLayer::capture(&guard);
981 let _span = captured
982 .iter()
983 .find(|s| s.name == "client_request")
984 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
985
986 check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
987 check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
988 check_bidi_write_span_attributes(&captured, "finalize", 98765, Some(5));
989
990 Ok(())
991 }
992
993 #[cfg(google_cloud_unstable_storage_bidi)]
994 #[tokio::test]
995 async fn reopen_appendable_object_not_found() -> anyhow::Result<()> {
996 use gaxi::grpc::tonic::Status as TonicStatus;
997 use google_cloud_gax::error::rpc::Code;
998 use storage_grpc_mock::{MockStorage, start};
999
1000 let guard = TestLayer::initialize();
1001 let mut mock = MockStorage::new();
1002 mock.expect_bidi_write_object()
1003 .return_once(|_| Err(TonicStatus::not_found("not here")));
1004 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
1005
1006 let client = crate::client::Storage::builder()
1007 .with_credentials(Anonymous::new().build())
1008 .with_endpoint(endpoint.clone())
1009 .with_tracing()
1010 .build()
1011 .await?;
1012 let response = client
1013 .reopen_appendable_object("projects/_/buckets/test-bucket", "test-object", 12345)
1014 .send()
1015 .await;
1016 assert!(
1017 matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
1018 "{response:?}"
1019 );
1020 let captured = TestLayer::capture(&guard);
1021 check_debug_log(&captured, "reopen_appendable_object");
1022
1023 client_request_span(&captured, "reopen_appendable_object", "NOT_FOUND", "grpc");
1024
1025 Ok(())
1026 }
1027
1028 #[cfg(google_cloud_unstable_storage_bidi)]
1030 #[tokio::test]
1031 async fn reopen_appendable_object_success() -> anyhow::Result<()> {
1032 use bytes::Bytes;
1033 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
1034 use storage_grpc_mock::google::storage::v2::{
1035 BidiWriteObjectResponse, bidi_write_object_response::WriteStatus,
1036 };
1037 use storage_grpc_mock::{MockStorage, start};
1038 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
1039 const OBJECT_NAME: &str = "test-object";
1040 const BIND_ADDRESS: &str = "0.0.0.0:0";
1041
1042 let guard = TestLayer::initialize();
1043
1044 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
1045 let response = BidiWriteObjectResponse {
1046 write_status: Some(WriteStatus::PersistedSize(100)),
1047 ..BidiWriteObjectResponse::default()
1048 };
1049
1050 tx.send(Ok(response.clone())).await?;
1053
1054 let finalize_response = BidiWriteObjectResponse {
1055 write_status: Some(WriteStatus::Resource(
1056 storage_grpc_mock::google::storage::v2::Object {
1057 size: 105,
1058 ..Default::default()
1059 },
1060 )),
1061 ..BidiWriteObjectResponse::default()
1062 };
1063
1064 let mut mock = MockStorage::new();
1065 mock.expect_bidi_write_object().return_once(move |req| {
1066 let mut stream = req.into_inner();
1067 tokio::spawn(async move {
1068 while let Some(Ok(msg)) = stream.recv().await {
1069 if msg.finish_write {
1070 let _ = tx.send(Ok(finalize_response.clone())).await;
1071 } else if msg.flush {
1072 let _ = tx.send(Ok(response.clone())).await;
1075 }
1076 }
1077 });
1078 Ok(TonicResponse::from(rx))
1079 });
1080 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
1081
1082 let client = crate::client::Storage::builder()
1083 .with_credentials(Anonymous::new().build())
1084 .with_endpoint(endpoint.clone())
1085 .with_tracing()
1086 .build()
1087 .await?;
1088 let mut writer = client
1089 .reopen_appendable_object(BUCKET_NAME, OBJECT_NAME, 12345)
1090 .send()
1091 .await?;
1092
1093 writer.append(Bytes::from_static(b"hello")).await?;
1094 writer.flush().await?;
1095 let obj = writer.finalize().await?;
1096 assert_eq!(obj.size, 105);
1097
1098 let captured = TestLayer::capture(&guard);
1099 let _span = captured
1100 .iter()
1101 .find(|s| s.name == "client_request")
1102 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1103
1104 check_bidi_write_span_attributes(&captured, "append", 12345, Some(5));
1105 check_bidi_write_span_attributes(&captured, "flush", 12345, Some(100));
1106 check_bidi_write_span_attributes(&captured, "finalize", 12345, Some(105));
1107
1108 Ok(())
1109 }
1110
1111 #[track_caller]
1112 fn check_debug_log(captured: &Vec<CapturedSpan>, method: &'static str) {
1113 let span = captured
1114 .iter()
1115 .find(|s| s.name == method)
1116 .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1117
1118 let got = BTreeMap::from_iter(span.attributes.clone());
1119 let want = ["self", "options", "request"];
1120 let missing = want
1121 .iter()
1122 .filter(|k| !got.contains_key(**k))
1123 .collect::<Vec<_>>();
1124 assert!(
1125 missing.is_empty(),
1126 "missing = {missing:?}\ngot = {:?}\nwant = {want:?}\nfull = {got:#?}",
1127 got.keys().collect::<Vec<_>>(),
1128 );
1129 }
1130
1131 #[track_caller]
1132 fn client_request_span(
1133 captured: &Vec<CapturedSpan>,
1134 method: &'static str,
1135 error_type: &'static str,
1136 rpc_system: &'static str,
1137 ) {
1138 let expected_attributes: [(&str, &str); 6] = [
1139 ("otel.kind", "Internal"),
1140 ("rpc.system.name", rpc_system),
1141 ("otel.status_code", "ERROR"),
1142 ("gcp.client.service", "storage"),
1143 ("gcp.client.repo", "googleapis/google-cloud-rust"),
1144 ("gcp.client.artifact", "google-cloud-storage"),
1145 ];
1146 let span = captured
1147 .iter()
1148 .find(|s| s.name == "client_request")
1149 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1150 let got = BTreeMap::from_iter(span.attributes.clone());
1151 let want = BTreeMap::<String, AttributeValue>::from_iter(
1154 expected_attributes
1155 .iter()
1156 .map(|(k, v)| (k.to_string(), AttributeValue::from(*v)))
1157 .chain(
1158 [
1159 (
1160 "otel.name",
1161 format!("google_cloud_storage::client::Storage::{method}").into(),
1162 ),
1163 ("error.type", error_type.into()),
1164 ]
1165 .map(|(k, v)| (k.to_string(), v)),
1166 ),
1167 );
1168 let mismatch = want
1169 .iter()
1170 .filter(|(k, v)| !got.get(k.as_str()).is_some_and(|g| g == *v))
1171 .collect::<Vec<_>>();
1172 assert!(
1173 mismatch.is_empty(),
1174 "mismatch = {mismatch:?}\ngot = {got:?}\nwant = {want:?}"
1175 );
1176 }
1177
1178 #[allow(dead_code)]
1179 #[track_caller]
1180 fn check_bidi_write_span_attributes(
1181 captured: &[CapturedSpan],
1182 method: &str,
1183 expected_generation: i64,
1184 expected_size: Option<i64>,
1185 ) {
1186 let span = captured
1187 .iter()
1188 .find(|s| {
1189 s.name == method && s.attributes.contains_key(&format!("{method}.generation"))
1190 })
1191 .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1192
1193 let expected_strs = [
1194 ("otel.kind", "Internal"),
1195 ("rpc.system.name", "grpc"),
1196 ("gcp.client.service", "storage"),
1197 ("gcp.client.version", env!("CARGO_PKG_VERSION")),
1198 ("gcp.client.repo", "googleapis/google-cloud-rust"),
1199 ("gcp.client.artifact", "google-cloud-storage"),
1200 ("gcp.schema.url", "https://opentelemetry.io/schemas/1.39.0"),
1201 ];
1202
1203 for (k, v) in expected_strs {
1204 assert_eq!(
1205 span.attributes.get(k),
1206 Some(&AttributeValue::String(v.to_string().into())),
1207 "mismatched attribute `{k}` in `{method}` span"
1208 );
1209 }
1210
1211 let gen_k = format!("{method}.generation");
1212 assert_eq!(
1213 span.attributes.get(gen_k.as_str()),
1214 Some(&AttributeValue::Int64(expected_generation)),
1215 "mismatched `{gen_k}`"
1216 );
1217
1218 if let Some(s) = expected_size {
1219 if method == "append" {
1220 assert_eq!(
1221 span.attributes.get("append.chunk_size"),
1222 Some(&AttributeValue::UInt64(s as u64)),
1223 "mismatched `append.chunk_size`"
1224 );
1225 } else {
1226 let size_k = format!("{method}.persisted_size");
1227 assert_eq!(
1228 span.attributes.get(size_k.as_str()),
1229 Some(&AttributeValue::Int64(s)),
1230 "mismatched `{size_k}`"
1231 );
1232 }
1233 }
1234 }
1235}