Skip to main content

google_cloud_storage/storage/
transport.rs

1// Copyright 2025 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     https://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// An implementation of [`stub::Storage`][crate::storage::stub::Storage] that
42/// interacts with the Cloud Storage service.
43///
44/// This is the default implementation of a
45/// [`client::Storage<T>`][crate::storage::client::Storage].
46///
47/// ## Example
48///
49/// ```
50/// # async fn sample() -> anyhow::Result<()> {
51/// use google_cloud_storage::client::Storage;
52/// use google_cloud_storage::stub::DefaultStorage;
53/// let client: Storage<DefaultStorage> = Storage::builder().build().await?;
54/// # Ok(()) }
55/// ```
56#[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    /// Implements [crate::client::Storage::read_object].
367    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    /// Implements [crate::client::Storage::write_object].
379    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    /// Implements [crate::client::Storage::write_object].
398    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        // The span counts one more event: the EOF
551        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        // TODO(#4772) - Move these `use` declarations and constants once the tracing APIs are stable.
667        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        // This is the initial response.
703        tx.send(Ok(response.clone())).await?;
704        // These simulate the calls to ObjectDescriptor::read_range(). The data is wrong, but this
705        // test is about the spans.
706        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    /// Models a complete lifecycle ending in close: `open` -> `append` -> `flush` -> `close`.
820    #[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        // The first response is the initial handshake/metadata response expected by
853        // `connector.rs` immediately upon opening the stream, before any data is sent.
854        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                        // The second response is sent ONLY when the client explicitly requests a flush.
863                        // `writer.append()` does not wait for a response, but `writer.flush()` does.
864                        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    /// Models a complete lifecycle ending in finalize: `open` -> `append` -> `flush` -> `finalize`.
901    #[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    /// Models a complete lifecycle: `reopen` -> `append` -> `flush` -> `finalize`.
1027    #[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        // The first response is the initial handshake/metadata response expected by
1049        // `connector.rs` immediately upon opening the stream, before any data is sent.
1050        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                        // The second response is sent ONLY when the client explicitly requests a flush.
1071                        // `writer.append()` does not wait for a response, but `writer.flush()` does.
1072                        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        // This is a subset of the fields, but good enough to catch most
1150        // mistakes. Recall that we use a macro, which is already tested.
1151        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}