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 #[cfg(google_cloud_unstable_storage_bidi)]
821 #[tokio::test]
822 async fn open_appendable_object_success() -> anyhow::Result<()> {
823 use bytes::Bytes;
824 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
825 use storage_grpc_mock::google::storage::v2::{
826 BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
827 };
828 use storage_grpc_mock::{MockStorage, start};
829 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
830 const OBJECT_NAME: &str = "test-object";
831 const BIND_ADDRESS: &str = "0.0.0.0:0";
832
833 let guard = TestLayer::initialize();
834
835 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
836
837 let initial_response = BidiWriteObjectResponse {
838 write_status: Some(WriteStatus::Resource(Object {
839 bucket: BUCKET_NAME.to_string(),
840 name: OBJECT_NAME.to_string(),
841 generation: 98765,
842 ..Default::default()
843 })),
844 ..BidiWriteObjectResponse::default()
845 };
846
847 let flush_response = BidiWriteObjectResponse {
848 write_status: Some(WriteStatus::PersistedSize(5)),
849 ..BidiWriteObjectResponse::default()
850 };
851
852 tx.send(Ok(initial_response)).await?;
855
856 let mut mock = MockStorage::new();
857 mock.expect_bidi_write_object().return_once(move |req| {
858 let mut stream = req.into_inner();
859 tokio::spawn(async move {
860 while let Some(Ok(msg)) = stream.recv().await {
861 if msg.flush {
862 let _ = tx.send(Ok(flush_response.clone())).await;
865 }
866 }
867 });
868 Ok(TonicResponse::new(rx))
869 });
870 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
871
872 let client = crate::client::Storage::builder()
873 .with_credentials(Anonymous::new().build())
874 .with_endpoint(endpoint.clone())
875 .with_tracing()
876 .build()
877 .await?;
878 let mut writer = client
879 .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
880 .send()
881 .await?;
882
883 writer.append(Bytes::from_static(b"hello")).await?;
884 writer.flush().await?;
885 writer.close().await?;
886
887 let captured = TestLayer::capture(&guard);
888 let _span = captured
889 .iter()
890 .find(|s| s.name == "client_request")
891 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
892
893 check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
894 check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
895 check_bidi_write_span_attributes(&captured, "close", 98765, Some(5));
896
897 Ok(())
898 }
899
900 #[cfg(google_cloud_unstable_storage_bidi)]
902 #[tokio::test]
903 async fn open_appendable_object_finalize_success() -> anyhow::Result<()> {
904 use bytes::Bytes;
905 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
906 use storage_grpc_mock::google::storage::v2::{
907 BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
908 };
909 use storage_grpc_mock::{MockStorage, start};
910 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
911 const OBJECT_NAME: &str = "test-object";
912 const BIND_ADDRESS: &str = "0.0.0.0:0";
913
914 let guard = TestLayer::initialize();
915
916 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
917
918 let initial_response = BidiWriteObjectResponse {
919 write_status: Some(WriteStatus::Resource(Object {
920 bucket: BUCKET_NAME.to_string(),
921 name: OBJECT_NAME.to_string(),
922 generation: 98765,
923 ..Default::default()
924 })),
925 ..BidiWriteObjectResponse::default()
926 };
927
928 let flush_response = BidiWriteObjectResponse {
929 write_status: Some(WriteStatus::PersistedSize(5)),
930 ..BidiWriteObjectResponse::default()
931 };
932
933 let finalize_response = BidiWriteObjectResponse {
934 write_status: Some(WriteStatus::Resource(Object {
935 bucket: BUCKET_NAME.to_string(),
936 name: OBJECT_NAME.to_string(),
937 generation: 98765,
938 size: 5,
939 ..Default::default()
940 })),
941 ..BidiWriteObjectResponse::default()
942 };
943
944 tx.send(Ok(initial_response)).await?;
945
946 let mut mock = MockStorage::new();
947 mock.expect_bidi_write_object().return_once(move |req| {
948 let mut stream = req.into_inner();
949 tokio::spawn(async move {
950 while let Some(Ok(msg)) = stream.recv().await {
951 if msg.finish_write {
952 let _ = tx.send(Ok(finalize_response.clone())).await;
953 } else if msg.flush {
954 let _ = tx.send(Ok(flush_response.clone())).await;
955 }
956 }
957 });
958 Ok(TonicResponse::new(rx))
959 });
960 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
961
962 let client = crate::client::Storage::builder()
963 .with_credentials(Anonymous::new().build())
964 .with_endpoint(endpoint.clone())
965 .with_tracing()
966 .build()
967 .await?;
968 let mut writer = client
969 .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
970 .send()
971 .await?;
972
973 writer.append(Bytes::from_static(b"hello")).await?;
974 writer.flush().await?;
975 let obj = writer.finalize().await?;
976 assert_eq!(obj.size, 5);
977
978 let captured = TestLayer::capture(&guard);
979 let _span = captured
980 .iter()
981 .find(|s| s.name == "client_request")
982 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
983
984 check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
985 check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
986 check_bidi_write_span_attributes(&captured, "finalize", 98765, Some(5));
987
988 Ok(())
989 }
990
991 #[cfg(google_cloud_unstable_storage_bidi)]
992 #[tokio::test]
993 async fn reopen_appendable_object_not_found() -> anyhow::Result<()> {
994 use gaxi::grpc::tonic::Status as TonicStatus;
995 use google_cloud_gax::error::rpc::Code;
996 use storage_grpc_mock::{MockStorage, start};
997
998 let guard = TestLayer::initialize();
999 let mut mock = MockStorage::new();
1000 mock.expect_bidi_write_object()
1001 .return_once(|_| Err(TonicStatus::not_found("not here")));
1002 let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
1003
1004 let client = crate::client::Storage::builder()
1005 .with_credentials(Anonymous::new().build())
1006 .with_endpoint(endpoint.clone())
1007 .with_tracing()
1008 .build()
1009 .await?;
1010 let response = client
1011 .reopen_appendable_object("projects/_/buckets/test-bucket", "test-object", 12345)
1012 .send()
1013 .await;
1014 assert!(
1015 matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
1016 "{response:?}"
1017 );
1018 let captured = TestLayer::capture(&guard);
1019 check_debug_log(&captured, "reopen_appendable_object");
1020
1021 client_request_span(&captured, "reopen_appendable_object", "NOT_FOUND", "grpc");
1022
1023 Ok(())
1024 }
1025
1026 #[cfg(google_cloud_unstable_storage_bidi)]
1028 #[tokio::test]
1029 async fn reopen_appendable_object_success() -> anyhow::Result<()> {
1030 use bytes::Bytes;
1031 use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
1032 use storage_grpc_mock::google::storage::v2::{
1033 BidiWriteObjectResponse, bidi_write_object_response::WriteStatus,
1034 };
1035 use storage_grpc_mock::{MockStorage, start};
1036 const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
1037 const OBJECT_NAME: &str = "test-object";
1038 const BIND_ADDRESS: &str = "0.0.0.0:0";
1039
1040 let guard = TestLayer::initialize();
1041
1042 let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
1043 let response = BidiWriteObjectResponse {
1044 write_status: Some(WriteStatus::PersistedSize(100)),
1045 ..BidiWriteObjectResponse::default()
1046 };
1047
1048 tx.send(Ok(response.clone())).await?;
1051
1052 let finalize_response = BidiWriteObjectResponse {
1053 write_status: Some(WriteStatus::Resource(
1054 storage_grpc_mock::google::storage::v2::Object {
1055 size: 105,
1056 ..Default::default()
1057 },
1058 )),
1059 ..BidiWriteObjectResponse::default()
1060 };
1061
1062 let mut mock = MockStorage::new();
1063 mock.expect_bidi_write_object().return_once(move |req| {
1064 let mut stream = req.into_inner();
1065 tokio::spawn(async move {
1066 while let Some(Ok(msg)) = stream.recv().await {
1067 if msg.finish_write {
1068 let _ = tx.send(Ok(finalize_response.clone())).await;
1069 } else if msg.flush {
1070 let _ = tx.send(Ok(response.clone())).await;
1073 }
1074 }
1075 });
1076 Ok(TonicResponse::from(rx))
1077 });
1078 let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
1079
1080 let client = crate::client::Storage::builder()
1081 .with_credentials(Anonymous::new().build())
1082 .with_endpoint(endpoint.clone())
1083 .with_tracing()
1084 .build()
1085 .await?;
1086 let mut writer = client
1087 .reopen_appendable_object(BUCKET_NAME, OBJECT_NAME, 12345)
1088 .send()
1089 .await?;
1090
1091 writer.append(Bytes::from_static(b"hello")).await?;
1092 writer.flush().await?;
1093 let obj = writer.finalize().await?;
1094 assert_eq!(obj.size, 105);
1095
1096 let captured = TestLayer::capture(&guard);
1097 let _span = captured
1098 .iter()
1099 .find(|s| s.name == "client_request")
1100 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1101
1102 check_bidi_write_span_attributes(&captured, "append", 12345, Some(5));
1103 check_bidi_write_span_attributes(&captured, "flush", 12345, Some(100));
1104 check_bidi_write_span_attributes(&captured, "finalize", 12345, Some(105));
1105
1106 Ok(())
1107 }
1108
1109 #[track_caller]
1110 fn check_debug_log(captured: &Vec<CapturedSpan>, method: &'static str) {
1111 let span = captured
1112 .iter()
1113 .find(|s| s.name == method)
1114 .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1115
1116 let got = BTreeMap::from_iter(span.attributes.clone());
1117 let want = ["self", "options", "request"];
1118 let missing = want
1119 .iter()
1120 .filter(|k| !got.contains_key(**k))
1121 .collect::<Vec<_>>();
1122 assert!(
1123 missing.is_empty(),
1124 "missing = {missing:?}\ngot = {:?}\nwant = {want:?}\nfull = {got:#?}",
1125 got.keys().collect::<Vec<_>>(),
1126 );
1127 }
1128
1129 #[track_caller]
1130 fn client_request_span(
1131 captured: &Vec<CapturedSpan>,
1132 method: &'static str,
1133 error_type: &'static str,
1134 rpc_system: &'static str,
1135 ) {
1136 let expected_attributes: [(&str, &str); 6] = [
1137 ("otel.kind", "Internal"),
1138 ("rpc.system.name", rpc_system),
1139 ("otel.status_code", "ERROR"),
1140 ("gcp.client.service", "storage"),
1141 ("gcp.client.repo", "googleapis/google-cloud-rust"),
1142 ("gcp.client.artifact", "google-cloud-storage"),
1143 ];
1144 let span = captured
1145 .iter()
1146 .find(|s| s.name == "client_request")
1147 .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1148 let got = BTreeMap::from_iter(span.attributes.clone());
1149 let want = BTreeMap::<String, AttributeValue>::from_iter(
1152 expected_attributes
1153 .iter()
1154 .map(|(k, v)| (k.to_string(), AttributeValue::from(*v)))
1155 .chain(
1156 [
1157 (
1158 "otel.name",
1159 format!("google_cloud_storage::client::Storage::{method}").into(),
1160 ),
1161 ("error.type", error_type.into()),
1162 ]
1163 .map(|(k, v)| (k.to_string(), v)),
1164 ),
1165 );
1166 let mismatch = want
1167 .iter()
1168 .filter(|(k, v)| !got.get(k.as_str()).is_some_and(|g| g == *v))
1169 .collect::<Vec<_>>();
1170 assert!(
1171 mismatch.is_empty(),
1172 "mismatch = {mismatch:?}\ngot = {got:?}\nwant = {want:?}"
1173 );
1174 }
1175
1176 #[allow(dead_code)]
1177 #[track_caller]
1178 fn check_bidi_write_span_attributes(
1179 captured: &[CapturedSpan],
1180 method: &str,
1181 expected_generation: i64,
1182 expected_size: Option<i64>,
1183 ) {
1184 let span = captured
1185 .iter()
1186 .find(|s| {
1187 s.name == method && s.attributes.contains_key(&format!("{method}.generation"))
1188 })
1189 .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1190
1191 let expected_strs = [
1192 ("otel.kind", "Internal"),
1193 ("rpc.system.name", "grpc"),
1194 ("gcp.client.service", "storage"),
1195 ("gcp.client.version", env!("CARGO_PKG_VERSION")),
1196 ("gcp.client.repo", "googleapis/google-cloud-rust"),
1197 ("gcp.client.artifact", "google-cloud-storage"),
1198 ("gcp.schema.url", "https://opentelemetry.io/schemas/1.39.0"),
1199 ];
1200
1201 for (k, v) in expected_strs {
1202 assert_eq!(
1203 span.attributes.get(k),
1204 Some(&AttributeValue::String(v.to_string().into())),
1205 "mismatched attribute `{k}` in `{method}` span"
1206 );
1207 }
1208
1209 let gen_k = format!("{method}.generation");
1210 assert_eq!(
1211 span.attributes.get(gen_k.as_str()),
1212 Some(&AttributeValue::Int64(expected_generation)),
1213 "mismatched `{gen_k}`"
1214 );
1215
1216 if let Some(s) = expected_size {
1217 if method == "append" {
1218 assert_eq!(
1219 span.attributes.get("append.chunk_size"),
1220 Some(&AttributeValue::UInt64(s as u64)),
1221 "mismatched `append.chunk_size`"
1222 );
1223 } else {
1224 let size_k = format!("{method}.persisted_size");
1225 assert_eq!(
1226 span.attributes.get(size_k.as_str()),
1227 Some(&AttributeValue::Int64(s)),
1228 "mismatched `{size_k}`"
1229 );
1230 }
1231 }
1232 }
1233}