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(request.checksum_precomputation)
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("/storage/v1/b/{bucket}/o/{object}")
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    #[cfg(google_cloud_unstable_storage_bidi)]
365    async fn open_appendable_object_and_append_plain(
366        &self,
367        request: OpenAppendableObjectRequest,
368        chunk: bytes::Bytes,
369        options: RequestOptions,
370    ) -> Result<AppendableObjectWriter> {
371        let connector = BidiWriteConnector::new(options, self.inner.grpc.clone());
372        let transport =
373            AppendableObjectWriterTransport::new_open_and_append(connector, request, chunk).await?;
374        Ok(AppendableObjectWriter::new(transport))
375    }
376
377    #[cfg(google_cloud_unstable_storage_bidi)]
378    #[tracing::instrument(name = "open_appendable_object", level = tracing::Level::DEBUG, ret, err(Debug), skip(chunk))]
379    async fn open_appendable_object_and_append_tracing(
380        &self,
381        request: OpenAppendableObjectRequest,
382        chunk: bytes::Bytes,
383        options: RequestOptions,
384    ) -> Result<AppendableObjectWriter> {
385        let resource_name = format!(
386            "//storage.googleapis.com/{}",
387            request
388                .spec
389                .resource
390                .as_ref()
391                .map(|r| r.bucket.as_str())
392                .unwrap_or_default()
393        );
394        let (_span, pending) = gaxi::client_request_signals!(
395            metric: self.metric.clone(),
396            info: *INSTRUMENTATION,
397            method: "client::Storage::open_appendable_object",
398            async {
399                if let Some(recorder) = RequestRecorder::current() {
400                    recorder.on_client_request(
401                        ClientRequestAttributes::default()
402                            .set_rpc_method("google.storage.v2.Storage/BidiWriteObject")
403                            .set_url_template("/upload/storage/v1/b/{bucket}/o")
404                            .set_resource_name(resource_name),
405                    );
406                }
407                self.open_appendable_object_and_append_plain(request, chunk, options)
408                    .await
409            }
410        );
411        let writer = pending.await?;
412        Ok(AppendableObjectWriter::new(
413            super::tracing::TracingAppendableObjectWriter::new(writer.into_parts()),
414        ))
415    }
416}
417
418impl super::stub::Storage for Storage {
419    /// Implements [crate::client::Storage::read_object].
420    async fn read_object(
421        &self,
422        req: ReadObjectRequest,
423        options: RequestOptions,
424    ) -> Result<ReadObjectResponse> {
425        if self.tracing {
426            return self.read_object_tracing(req, options).await;
427        }
428        self.read_object_plain(req, options).await
429    }
430
431    /// Implements [crate::client::Storage::write_object].
432    async fn write_object_buffered<P>(
433        &self,
434        payload: P,
435        req: WriteObjectRequest,
436        options: RequestOptions,
437    ) -> Result<Object>
438    where
439        P: StreamingSource + Send + Sync + 'static,
440    {
441        if self.tracing {
442            return self
443                .write_object_buffered_tracing(payload, req, options)
444                .await;
445        }
446        self.write_object_buffered_plain(payload, req, options)
447            .await
448    }
449
450    /// Implements [crate::client::Storage::write_object].
451    async fn write_object_unbuffered<P>(
452        &self,
453        payload: P,
454        req: WriteObjectRequest,
455        options: RequestOptions,
456    ) -> Result<Object>
457    where
458        P: StreamingSource + Seek + Send + Sync + 'static,
459    {
460        if self.tracing {
461            return self
462                .write_object_unbuffered_tracing(payload, req, options)
463                .await;
464        }
465        self.write_object_unbuffered_plain(payload, req, options)
466            .await
467    }
468
469    async fn open_object(
470        &self,
471        request: OpenObjectRequest,
472        options: RequestOptions,
473    ) -> Result<(ObjectDescriptor, Vec<ReadObjectResponse>)> {
474        if self.tracing {
475            return self.open_object_tracing(request, options).await;
476        }
477        self.open_object_plain(request, options).await
478    }
479
480    #[cfg(google_cloud_unstable_storage_bidi)]
481    async fn open_appendable_object(
482        &self,
483        request: OpenAppendableObjectRequest,
484        options: RequestOptions,
485    ) -> Result<AppendableObjectWriter> {
486        if self.tracing {
487            return self.open_appendable_object_tracing(request, options).await;
488        }
489        self.open_appendable_object_plain(request, options).await
490    }
491
492    #[cfg(google_cloud_unstable_storage_bidi)]
493    async fn open_appendable_object_and_append(
494        &self,
495        request: OpenAppendableObjectRequest,
496        chunk: bytes::Bytes,
497        options: RequestOptions,
498    ) -> Result<AppendableObjectWriter> {
499        if self.tracing {
500            return self
501                .open_appendable_object_and_append_tracing(request, chunk, options)
502                .await;
503        }
504        self.open_appendable_object_and_append_plain(request, chunk, options)
505            .await
506    }
507
508    #[cfg(google_cloud_unstable_storage_bidi)]
509    async fn reopen_appendable_object(
510        &self,
511        request: ReopenAppendableObjectRequest,
512        options: RequestOptions,
513    ) -> Result<AppendableObjectWriter> {
514        if self.tracing {
515            return self
516                .reopen_appendable_object_tracing(request, options)
517                .await;
518        }
519        self.reopen_appendable_object_plain(request, options).await
520    }
521}
522
523#[cfg(test)]
524mod tests {
525    use super::{Storage, StorageInner};
526    use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
527    use google_cloud_test_utils::test_layer::AttributeValue;
528    use google_cloud_test_utils::test_layer::{CapturedSpan, TestLayer};
529    use httptest::{Expectation, Server, matchers::*, responders::status_code};
530    use pretty_assertions::assert_eq;
531    use std::collections::BTreeMap;
532    use std::sync::Arc;
533
534    impl Storage {
535        pub(crate) fn new_test(inner: Arc<StorageInner>) -> Arc<Self> {
536            Self::new(inner, false)
537        }
538    }
539
540    #[tokio::test]
541    async fn read_object() -> anyhow::Result<()> {
542        let guard = TestLayer::initialize();
543
544        let server = Server::run();
545        server.expect(
546            Expectation::matching(all_of![
547                request::method_path("GET", "/storage/v1/b/test-bucket/o/test-object"),
548                request::query(url_decoded(contains(("alt", "media")))),
549            ])
550            .respond_with(status_code(404)),
551        );
552
553        let client = crate::client::Storage::builder()
554            .with_endpoint(format!("http://{}", server.addr()))
555            .with_credentials(Anonymous::new().build())
556            .with_tracing()
557            .build()
558            .await?;
559        let response = client
560            .read_object("projects/_/buckets/test-bucket", "test-object")
561            .send()
562            .await;
563        assert!(
564            matches!(response, Err(ref e) if e.is_transport()),
565            "{response:?}"
566        );
567
568        let captured = TestLayer::capture(&guard);
569        check_debug_log(&captured, "read_object");
570
571        client_request_span(&captured, "read_object", "404", "http");
572
573        Ok(())
574    }
575
576    #[tokio::test]
577    async fn read_object_success() -> anyhow::Result<()> {
578        let guard = TestLayer::initialize();
579
580        let body = (0..100_000)
581            .map(|i| format!("{i:08} {:1000}", ""))
582            .collect::<Vec<_>>()
583            .join("\n");
584        let server = Server::run();
585        server.expect(
586            Expectation::matching(all_of![
587                request::method_path("GET", "/storage/v1/b/test-bucket/o/test-object"),
588                request::query(url_decoded(contains(("alt", "media")))),
589            ])
590            .respond_with(
591                status_code(200)
592                    .body(body.clone())
593                    .append_header("x-goog-generation", 123456),
594            ),
595        );
596
597        let client = crate::client::Storage::builder()
598            .with_endpoint(format!("http://{}", server.addr()))
599            .with_credentials(Anonymous::new().build())
600            .with_tracing()
601            .build()
602            .await?;
603        let mut got = Vec::new();
604        let mut response = client
605            .read_object("projects/_/buckets/test-bucket", "test-object")
606            .send()
607            .await?;
608        let object = response.object();
609        assert_eq!(object.generation, 123456, "{object:?}");
610        while let Some(b) = response.next().await.transpose()? {
611            got.push(b);
612        }
613
614        let captured = TestLayer::capture(&guard);
615        let span = captured
616            .iter()
617            .find(|s| s.name == "client_request")
618            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
619        // The span counts one more event: the EOF
620        assert_eq!(span.events, got.len() + 1, "{span:?}");
621
622        Ok(())
623    }
624
625    #[tokio::test]
626    async fn write_object_buffered() -> anyhow::Result<()> {
627        let guard = TestLayer::initialize();
628
629        let server = Server::run();
630        server.expect(
631            Expectation::matching(all_of![
632                request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
633                request::query(url_decoded(contains(("uploadType", "multipart")))),
634            ])
635            .respond_with(status_code(404)),
636        );
637
638        let client = crate::client::Storage::builder()
639            .with_endpoint(format!("http://{}", server.addr()))
640            .with_credentials(Anonymous::new().build())
641            .with_tracing()
642            .build()
643            .await?;
644        let response = client
645            .write_object("projects/_/buckets/test-bucket", "test-object", "payload")
646            .send_buffered()
647            .await;
648        assert!(
649            matches!(response, Err(ref e) if e.is_transport()),
650            "{response:?}"
651        );
652
653        let captured = TestLayer::capture(&guard);
654        check_debug_log(&captured, "write_object_buffered");
655
656        client_request_span(&captured, "write_object", "404", "http");
657
658        Ok(())
659    }
660
661    #[tokio::test]
662    async fn write_object_unbuffered() -> anyhow::Result<()> {
663        let guard = TestLayer::initialize();
664
665        let server = Server::run();
666        server.expect(
667            Expectation::matching(all_of![
668                request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
669                request::query(url_decoded(contains(("uploadType", "multipart")))),
670            ])
671            .respond_with(status_code(404)),
672        );
673
674        let client = crate::client::Storage::builder()
675            .with_endpoint(format!("http://{}", server.addr()))
676            .with_credentials(Anonymous::new().build())
677            .with_tracing()
678            .build()
679            .await?;
680        let response = client
681            .write_object("projects/_/buckets/test-bucket", "test-object", "payload")
682            .send_unbuffered()
683            .await;
684        assert!(
685            matches!(response, Err(ref e) if e.is_transport()),
686            "{response:?}"
687        );
688
689        let captured = TestLayer::capture(&guard);
690        check_debug_log(&captured, "write_object_unbuffered");
691
692        client_request_span(&captured, "write_object", "404", "http");
693
694        Ok(())
695    }
696
697    #[tokio::test]
698    async fn open_object() -> anyhow::Result<()> {
699        use gaxi::grpc::tonic::Status as TonicStatus;
700        use google_cloud_gax::error::rpc::Code;
701        use storage_grpc_mock::{MockStorage, start};
702
703        let guard = TestLayer::initialize();
704
705        let mut mock = MockStorage::new();
706        mock.expect_bidi_read_object()
707            .return_once(|_| Err(TonicStatus::not_found("not here")));
708        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
709
710        let client = crate::client::Storage::builder()
711            .with_credentials(Anonymous::new().build())
712            .with_endpoint(endpoint.clone())
713            .with_tracing()
714            .build()
715            .await?;
716        let response = client
717            .open_object("projects/_/buckets/test-bucket", "test-object")
718            .send()
719            .await;
720        assert!(
721            matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
722            "{response:?}"
723        );
724
725        let captured = TestLayer::capture(&guard);
726        check_debug_log(&captured, "open_object");
727
728        client_request_span(&captured, "open_object", "NOT_FOUND", "grpc");
729        Ok(())
730    }
731
732    #[tokio::test]
733    #[ignore = "flaky test, see #5290"]
734    async fn open_object_success() -> anyhow::Result<()> {
735        // TODO(#4772) - Move these `use` declarations and constants once the tracing APIs are stable.
736        use crate::model_ext::ReadRange;
737        use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
738        use storage_grpc_mock::google::storage::v2::{
739            BidiReadObjectResponse, ChecksummedData, Object as ProtoObject, ObjectRangeData,
740            ReadRange as ProtoRange,
741        };
742        use storage_grpc_mock::{MockStorage, start};
743        const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
744        const OBJECT_NAME: &str = "test-object";
745        const BIND_ADDRESS: &str = "0.0.0.0:0";
746        const PAYLOAD: &str = "the quick brown fox jumps over the lazy dog";
747
748        let guard = TestLayer::initialize();
749
750        let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiReadObjectResponse>>(10);
751        let response = BidiReadObjectResponse {
752            metadata: Some(ProtoObject {
753                bucket: BUCKET_NAME.to_string(),
754                name: OBJECT_NAME.to_string(),
755                generation: 123456,
756                ..ProtoObject::default()
757            }),
758            object_data_ranges: vec![ObjectRangeData {
759                read_range: Some(ProtoRange {
760                    read_id: 0_i64,
761                    ..ProtoRange::default()
762                }),
763                range_end: true,
764                checksummed_data: Some(ChecksummedData {
765                    content: PAYLOAD.as_bytes().to_vec(),
766                    crc32c: None,
767                }),
768            }],
769            ..BidiReadObjectResponse::default()
770        };
771        // This is the initial response.
772        tx.send(Ok(response.clone())).await?;
773        // These simulate the calls to ObjectDescriptor::read_range(). The data is wrong, but this
774        // test is about the spans.
775        tx.send(Ok(response.clone())).await?;
776        tx.send(Ok(response.clone())).await?;
777
778        let mut mock = MockStorage::new();
779        mock.expect_bidi_read_object()
780            .return_once(|_| Ok(TonicResponse::from(rx)));
781        let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
782
783        let client = crate::client::Storage::builder()
784            .with_credentials(Anonymous::new().build())
785            .with_endpoint(endpoint.clone())
786            .with_tracing()
787            .build()
788            .await?;
789        let (descriptor, _reader0) = client
790            .open_object(BUCKET_NAME, OBJECT_NAME)
791            .send_and_read(ReadRange::all())
792            .await?;
793        let _reader1 = descriptor.read_range(ReadRange::offset(5)).await;
794        let _reader2 = descriptor.read_range(ReadRange::segment(10, 10)).await;
795        let _reader3 = descriptor.read_range(ReadRange::tail(15)).await;
796
797        let captured = TestLayer::capture(&guard);
798        let _span = captured
799            .iter()
800            .find(|s| s.name == "client_request")
801            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
802
803        let range_spans = captured
804            .iter()
805            .filter(|s| s.name == "read_range")
806            .collect::<Vec<_>>();
807
808        let _span_reader1 = range_spans
809            .clone()
810            .into_iter()
811            .find(|s| {
812                s.attributes
813                    .get("read_range.start")
814                    .and_then(|v| v.as_i64())
815                    == Some(5)
816            })
817            .unwrap_or_else(|| {
818                panic!("missing `read_range` span for ReadRange::offset(5): {range_spans:#?}")
819            });
820
821        let _span_reader2 = range_spans
822            .clone()
823            .into_iter()
824            .find(|s| {
825                s.attributes
826                    .get("read_range.start")
827                    .and_then(|v| v.as_i64())
828                    == Some(10)
829                    && s.attributes
830                        .get("read_range.limit")
831                        .and_then(|v| v.as_i64())
832                        == Some(10)
833            })
834            .unwrap_or_else(|| {
835                panic!("missing `read_range` span for ReadRange::segment(10, 10): {range_spans:#?}")
836            });
837
838        let _span_reader3 = range_spans
839            .clone()
840            .into_iter()
841            .find(|s| {
842                s.attributes
843                    .get("read_range.start")
844                    .and_then(|v| v.as_i64())
845                    == Some(-15)
846            })
847            .unwrap_or_else(|| {
848                panic!("missing `read_range` span for ReadRange::tail(15): {range_spans:#?}")
849            });
850        Ok(())
851    }
852
853    #[cfg(google_cloud_unstable_storage_bidi)]
854    #[tokio::test]
855    async fn open_appendable_object_not_found() -> anyhow::Result<()> {
856        use gaxi::grpc::tonic::Status as TonicStatus;
857        use google_cloud_gax::error::rpc::Code;
858        use storage_grpc_mock::{MockStorage, start};
859
860        let guard = TestLayer::initialize();
861        let mut mock = MockStorage::new();
862        mock.expect_bidi_write_object()
863            .return_once(|_| Err(TonicStatus::not_found("not here")));
864        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
865
866        let client = crate::client::Storage::builder()
867            .with_credentials(Anonymous::new().build())
868            .with_endpoint(endpoint.clone())
869            .with_tracing()
870            .build()
871            .await?;
872        let response = client
873            .open_appendable_object("projects/_/buckets/test-bucket", "test-object")
874            .send()
875            .await;
876        assert!(
877            matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
878            "{response:?}"
879        );
880        let captured = TestLayer::capture(&guard);
881        check_debug_log(&captured, "open_appendable_object");
882
883        client_request_span(&captured, "open_appendable_object", "NOT_FOUND", "grpc");
884
885        Ok(())
886    }
887
888    /// Models a complete lifecycle ending in close: `open` -> `append` -> `flush` -> `close`.
889    #[ignore = "TODO(#6324) - disabled because it was flaky"]
890    #[cfg(google_cloud_unstable_storage_bidi)]
891    #[tokio::test]
892    async fn open_appendable_object_success() -> anyhow::Result<()> {
893        use bytes::Bytes;
894        use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
895        use storage_grpc_mock::google::storage::v2::{
896            BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
897        };
898        use storage_grpc_mock::{MockStorage, start};
899        const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
900        const OBJECT_NAME: &str = "test-object";
901        const BIND_ADDRESS: &str = "0.0.0.0:0";
902
903        let guard = TestLayer::initialize();
904
905        let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
906
907        let initial_response = BidiWriteObjectResponse {
908            write_status: Some(WriteStatus::Resource(Object {
909                bucket: BUCKET_NAME.to_string(),
910                name: OBJECT_NAME.to_string(),
911                generation: 98765,
912                ..Default::default()
913            })),
914            ..BidiWriteObjectResponse::default()
915        };
916
917        let flush_response = BidiWriteObjectResponse {
918            write_status: Some(WriteStatus::PersistedSize(5)),
919            ..BidiWriteObjectResponse::default()
920        };
921
922        // The first response is the initial handshake/metadata response expected by
923        // `connector.rs` immediately upon opening the stream, before any data is sent.
924        tx.send(Ok(initial_response)).await?;
925
926        let mut mock = MockStorage::new();
927        mock.expect_bidi_write_object().return_once(move |req| {
928            let mut stream = req.into_inner();
929            tokio::spawn(async move {
930                while let Some(Ok(msg)) = stream.recv().await {
931                    if msg.flush {
932                        // The second response is sent ONLY when the client explicitly requests a flush.
933                        // `writer.append()` does not wait for a response, but `writer.flush()` does.
934                        let _ = tx.send(Ok(flush_response.clone())).await;
935                    }
936                }
937            });
938            Ok(TonicResponse::new(rx))
939        });
940        let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
941
942        let client = crate::client::Storage::builder()
943            .with_credentials(Anonymous::new().build())
944            .with_endpoint(endpoint.clone())
945            .with_tracing()
946            .build()
947            .await?;
948        let mut writer = client
949            .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
950            .send()
951            .await?;
952
953        writer.append(Bytes::from_static(b"hello")).await?;
954        writer.flush().await?;
955        writer.close().await?;
956
957        let captured = TestLayer::capture(&guard);
958        let _span = captured
959            .iter()
960            .find(|s| s.name == "client_request")
961            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
962
963        check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
964        check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
965        check_bidi_write_span_attributes(&captured, "close", 98765, Some(5));
966
967        Ok(())
968    }
969
970    /// Models a complete lifecycle ending in finalize: `open` -> `append` -> `flush` -> `finalize`.
971    #[ignore = "TODO(#6324) - disabled because it was flaky"]
972    #[cfg(google_cloud_unstable_storage_bidi)]
973    #[tokio::test]
974    async fn open_appendable_object_finalize_success() -> anyhow::Result<()> {
975        use bytes::Bytes;
976        use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
977        use storage_grpc_mock::google::storage::v2::{
978            BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
979        };
980        use storage_grpc_mock::{MockStorage, start};
981        const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
982        const OBJECT_NAME: &str = "test-object";
983        const BIND_ADDRESS: &str = "0.0.0.0:0";
984
985        let guard = TestLayer::initialize();
986
987        let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
988
989        let initial_response = BidiWriteObjectResponse {
990            write_status: Some(WriteStatus::Resource(Object {
991                bucket: BUCKET_NAME.to_string(),
992                name: OBJECT_NAME.to_string(),
993                generation: 98765,
994                ..Default::default()
995            })),
996            ..BidiWriteObjectResponse::default()
997        };
998
999        let flush_response = BidiWriteObjectResponse {
1000            write_status: Some(WriteStatus::PersistedSize(5)),
1001            ..BidiWriteObjectResponse::default()
1002        };
1003
1004        let finalize_response = BidiWriteObjectResponse {
1005            write_status: Some(WriteStatus::Resource(Object {
1006                bucket: BUCKET_NAME.to_string(),
1007                name: OBJECT_NAME.to_string(),
1008                generation: 98765,
1009                size: 5,
1010                ..Default::default()
1011            })),
1012            ..BidiWriteObjectResponse::default()
1013        };
1014
1015        tx.send(Ok(initial_response)).await?;
1016
1017        let mut mock = MockStorage::new();
1018        mock.expect_bidi_write_object().return_once(move |req| {
1019            let mut stream = req.into_inner();
1020            tokio::spawn(async move {
1021                while let Some(Ok(msg)) = stream.recv().await {
1022                    if msg.finish_write {
1023                        let _ = tx.send(Ok(finalize_response.clone())).await;
1024                    } else if msg.flush {
1025                        let _ = tx.send(Ok(flush_response.clone())).await;
1026                    }
1027                }
1028            });
1029            Ok(TonicResponse::new(rx))
1030        });
1031        let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
1032
1033        let client = crate::client::Storage::builder()
1034            .with_credentials(Anonymous::new().build())
1035            .with_endpoint(endpoint.clone())
1036            .with_tracing()
1037            .build()
1038            .await?;
1039        let mut writer = client
1040            .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
1041            .send()
1042            .await?;
1043
1044        writer.append(Bytes::from_static(b"hello")).await?;
1045        writer.flush().await?;
1046        let obj = writer.finalize().await?;
1047        assert_eq!(obj.size, 5);
1048
1049        let captured = TestLayer::capture(&guard);
1050        let _span = captured
1051            .iter()
1052            .find(|s| s.name == "client_request")
1053            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1054
1055        check_bidi_write_span_attributes(&captured, "append", 98765, Some(5));
1056        check_bidi_write_span_attributes(&captured, "flush", 98765, Some(5));
1057        check_bidi_write_span_attributes(&captured, "finalize", 98765, Some(5));
1058
1059        Ok(())
1060    }
1061
1062    #[cfg(google_cloud_unstable_storage_bidi)]
1063    #[tokio::test]
1064    async fn open_appendable_object_and_append_bucket_not_found() -> anyhow::Result<()> {
1065        use gaxi::grpc::tonic::Status as TonicStatus;
1066        use google_cloud_gax::error::rpc::Code;
1067        use storage_grpc_mock::{MockStorage, start};
1068
1069        // Arrange.
1070        let guard = TestLayer::initialize();
1071        let mut mock = MockStorage::new();
1072        mock.expect_bidi_write_object()
1073            .return_once(|_| Err(TonicStatus::not_found("not here")));
1074        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
1075
1076        let client = crate::client::Storage::builder()
1077            .with_credentials(Anonymous::new().build())
1078            .with_endpoint(endpoint.clone())
1079            .with_tracing()
1080            .build()
1081            .await?;
1082
1083        // Act.
1084        let response = client
1085            .open_appendable_object("projects/_/buckets/test-bucket", "test-object")
1086            .send_and_append(bytes::Bytes::from("hello"))
1087            .await;
1088
1089        // Assert.
1090        assert!(
1091            matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
1092            "{response:?}"
1093        );
1094        let captured = TestLayer::capture(&guard);
1095        check_debug_log(&captured, "open_appendable_object");
1096
1097        client_request_span(&captured, "open_appendable_object", "NOT_FOUND", "grpc");
1098
1099        Ok(())
1100    }
1101
1102    /// Models an opening stream with initial payload: `send_and_append` -> `finalize`.
1103    #[ignore = "TODO(#6324) - disabled because it was flaky"]
1104    #[cfg(google_cloud_unstable_storage_bidi)]
1105    #[tokio::test]
1106    async fn open_appendable_object_and_append_success() -> anyhow::Result<()> {
1107        use bytes::Bytes;
1108        use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
1109        use storage_grpc_mock::google::storage::v2::{
1110            BidiWriteObjectResponse, Object, bidi_write_object_response::WriteStatus,
1111        };
1112        use storage_grpc_mock::{MockStorage, start};
1113        const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
1114        const OBJECT_NAME: &str = "test-object";
1115        const BIND_ADDRESS: &str = "0.0.0.0:0";
1116
1117        // Arrange.
1118        let guard = TestLayer::initialize();
1119
1120        let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
1121
1122        let initial_response = BidiWriteObjectResponse {
1123            write_status: Some(WriteStatus::Resource(Object {
1124                bucket: BUCKET_NAME.to_string(),
1125                name: OBJECT_NAME.to_string(),
1126                generation: 98765,
1127                ..Default::default()
1128            })),
1129            ..BidiWriteObjectResponse::default()
1130        };
1131
1132        let finalize_response = BidiWriteObjectResponse {
1133            write_status: Some(WriteStatus::Resource(Object {
1134                bucket: BUCKET_NAME.to_string(),
1135                name: OBJECT_NAME.to_string(),
1136                generation: 98765,
1137                size: 5,
1138                ..Default::default()
1139            })),
1140            ..BidiWriteObjectResponse::default()
1141        };
1142
1143        // Initial handshake response sent immediately upon stream opening.
1144        tx.send(Ok(initial_response)).await?;
1145
1146        let mut mock = MockStorage::new();
1147        mock.expect_bidi_write_object().return_once(move |req| {
1148            let mut stream = req.into_inner();
1149            tokio::spawn(async move {
1150                while let Some(Ok(msg)) = stream.recv().await {
1151                    if msg.finish_write {
1152                        let _ = tx.send(Ok(finalize_response.clone())).await;
1153                    }
1154                }
1155            });
1156            Ok(TonicResponse::new(rx))
1157        });
1158        let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
1159
1160        let client = crate::client::Storage::builder()
1161            .with_credentials(Anonymous::new().build())
1162            .with_endpoint(endpoint.clone())
1163            .with_tracing()
1164            .build()
1165            .await?;
1166
1167        // Act.
1168        let writer = client
1169            .open_appendable_object(BUCKET_NAME, OBJECT_NAME)
1170            .send_and_append(Bytes::from_static(b"hello"))
1171            .await?;
1172
1173        let obj = writer.finalize().await?;
1174
1175        // Assert.
1176        assert_eq!(obj.size, 5);
1177
1178        let captured = TestLayer::capture(&guard);
1179        let _span = captured
1180            .iter()
1181            .find(|s| s.name == "client_request")
1182            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1183
1184        check_bidi_write_span_attributes(&captured, "finalize", 98765, Some(5));
1185
1186        Ok(())
1187    }
1188
1189    #[cfg(google_cloud_unstable_storage_bidi)]
1190    #[tokio::test]
1191    async fn reopen_appendable_object_not_found() -> anyhow::Result<()> {
1192        use gaxi::grpc::tonic::Status as TonicStatus;
1193        use google_cloud_gax::error::rpc::Code;
1194        use storage_grpc_mock::{MockStorage, start};
1195
1196        let guard = TestLayer::initialize();
1197        let mut mock = MockStorage::new();
1198        mock.expect_bidi_write_object()
1199            .return_once(|_| Err(TonicStatus::not_found("not here")));
1200        let (endpoint, _server) = start("0.0.0.0:0", mock).await?;
1201
1202        let client = crate::client::Storage::builder()
1203            .with_credentials(Anonymous::new().build())
1204            .with_endpoint(endpoint.clone())
1205            .with_tracing()
1206            .build()
1207            .await?;
1208        let response = client
1209            .reopen_appendable_object("projects/_/buckets/test-bucket", "test-object", 12345)
1210            .send()
1211            .await;
1212        assert!(
1213            matches!(response, Err(ref e) if e.status().is_some_and(|s| s.code == Code::NotFound)),
1214            "{response:?}"
1215        );
1216        let captured = TestLayer::capture(&guard);
1217        check_debug_log(&captured, "reopen_appendable_object");
1218
1219        client_request_span(&captured, "reopen_appendable_object", "NOT_FOUND", "grpc");
1220
1221        Ok(())
1222    }
1223
1224    /// Models a complete lifecycle: `reopen` -> `append` -> `flush` -> `finalize`.
1225    #[ignore = "TODO(#6324) - disabled because it was flaky"]
1226    #[cfg(google_cloud_unstable_storage_bidi)]
1227    #[tokio::test]
1228    async fn reopen_appendable_object_success() -> anyhow::Result<()> {
1229        use bytes::Bytes;
1230        use gaxi::grpc::tonic::{Response as TonicResponse, Result as TonicResult};
1231        use storage_grpc_mock::google::storage::v2::{
1232            BidiWriteObjectResponse, bidi_write_object_response::WriteStatus,
1233        };
1234        use storage_grpc_mock::{MockStorage, start};
1235        const BUCKET_NAME: &str = "projects/_/buckets/test-bucket";
1236        const OBJECT_NAME: &str = "test-object";
1237        const BIND_ADDRESS: &str = "0.0.0.0:0";
1238
1239        let guard = TestLayer::initialize();
1240
1241        let (tx, rx) = tokio::sync::mpsc::channel::<TonicResult<BidiWriteObjectResponse>>(10);
1242        let response = BidiWriteObjectResponse {
1243            write_status: Some(WriteStatus::PersistedSize(100)),
1244            ..BidiWriteObjectResponse::default()
1245        };
1246
1247        // The first response is the initial handshake/metadata response expected by
1248        // `connector.rs` immediately upon opening the stream, before any data is sent.
1249        tx.send(Ok(response.clone())).await?;
1250
1251        let finalize_response = BidiWriteObjectResponse {
1252            write_status: Some(WriteStatus::Resource(
1253                storage_grpc_mock::google::storage::v2::Object {
1254                    size: 105,
1255                    ..Default::default()
1256                },
1257            )),
1258            ..BidiWriteObjectResponse::default()
1259        };
1260
1261        let mut mock = MockStorage::new();
1262        mock.expect_bidi_write_object().return_once(move |req| {
1263            let mut stream = req.into_inner();
1264            tokio::spawn(async move {
1265                while let Some(Ok(msg)) = stream.recv().await {
1266                    if msg.finish_write {
1267                        let _ = tx.send(Ok(finalize_response.clone())).await;
1268                    } else if msg.flush {
1269                        // The second response is sent ONLY when the client explicitly requests a flush.
1270                        // `writer.append()` does not wait for a response, but `writer.flush()` does.
1271                        let _ = tx.send(Ok(response.clone())).await;
1272                    }
1273                }
1274            });
1275            Ok(TonicResponse::from(rx))
1276        });
1277        let (endpoint, _server) = start(BIND_ADDRESS, mock).await?;
1278
1279        let client = crate::client::Storage::builder()
1280            .with_credentials(Anonymous::new().build())
1281            .with_endpoint(endpoint.clone())
1282            .with_tracing()
1283            .build()
1284            .await?;
1285        let mut writer = client
1286            .reopen_appendable_object(BUCKET_NAME, OBJECT_NAME, 12345)
1287            .send()
1288            .await?;
1289
1290        writer.append(Bytes::from_static(b"hello")).await?;
1291        writer.flush().await?;
1292        let obj = writer.finalize().await?;
1293        assert_eq!(obj.size, 105);
1294
1295        let captured = TestLayer::capture(&guard);
1296        let _span = captured
1297            .iter()
1298            .find(|s| s.name == "client_request")
1299            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1300
1301        check_bidi_write_span_attributes(&captured, "append", 12345, Some(5));
1302        check_bidi_write_span_attributes(&captured, "flush", 12345, Some(100));
1303        check_bidi_write_span_attributes(&captured, "finalize", 12345, Some(105));
1304
1305        Ok(())
1306    }
1307
1308    #[track_caller]
1309    fn check_debug_log(captured: &Vec<CapturedSpan>, method: &'static str) {
1310        let span = captured
1311            .iter()
1312            .find(|s| s.name == method)
1313            .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1314
1315        let got = BTreeMap::from_iter(span.attributes.clone());
1316        let want = ["self", "options", "request"];
1317        let missing = want
1318            .iter()
1319            .filter(|k| !got.contains_key(**k))
1320            .collect::<Vec<_>>();
1321        assert!(
1322            missing.is_empty(),
1323            "missing = {missing:?}\ngot  = {:?}\nwant = {want:?}\nfull = {got:#?}",
1324            got.keys().collect::<Vec<_>>(),
1325        );
1326    }
1327
1328    #[track_caller]
1329    fn client_request_span(
1330        captured: &Vec<CapturedSpan>,
1331        method: &'static str,
1332        error_type: &'static str,
1333        rpc_system: &'static str,
1334    ) {
1335        let expected_attributes: [(&str, &str); 6] = [
1336            ("otel.kind", "Internal"),
1337            ("rpc.system.name", rpc_system),
1338            ("otel.status_code", "ERROR"),
1339            ("gcp.client.service", "storage"),
1340            ("gcp.client.repo", "googleapis/google-cloud-rust"),
1341            ("gcp.client.artifact", "google-cloud-storage"),
1342        ];
1343        let span = captured
1344            .iter()
1345            .find(|s| s.name == "client_request")
1346            .unwrap_or_else(|| panic!("missing `client_request` span in capture: {captured:#?}"));
1347        let got = BTreeMap::from_iter(span.attributes.clone());
1348        // This is a subset of the fields, but good enough to catch most
1349        // mistakes. Recall that we use a macro, which is already tested.
1350        let want = BTreeMap::<String, AttributeValue>::from_iter(
1351            expected_attributes
1352                .iter()
1353                .map(|(k, v)| (k.to_string(), AttributeValue::from(*v)))
1354                .chain(
1355                    [
1356                        (
1357                            "otel.name",
1358                            format!("google_cloud_storage::client::Storage::{method}").into(),
1359                        ),
1360                        ("error.type", error_type.into()),
1361                    ]
1362                    .map(|(k, v)| (k.to_string(), v)),
1363                ),
1364        );
1365        let mismatch = want
1366            .iter()
1367            .filter(|(k, v)| !got.get(k.as_str()).is_some_and(|g| g == *v))
1368            .collect::<Vec<_>>();
1369        assert!(
1370            mismatch.is_empty(),
1371            "mismatch = {mismatch:?}\ngot      = {got:?}\nwant     = {want:?}"
1372        );
1373    }
1374
1375    #[allow(dead_code)]
1376    #[track_caller]
1377    fn check_bidi_write_span_attributes(
1378        captured: &[CapturedSpan],
1379        method: &str,
1380        expected_generation: i64,
1381        expected_size: Option<i64>,
1382    ) {
1383        let span = captured
1384            .iter()
1385            .find(|s| {
1386                s.name == method && s.attributes.contains_key(&format!("{method}.generation"))
1387            })
1388            .unwrap_or_else(|| panic!("missing `{method}` span in capture: {captured:#?}"));
1389
1390        let expected_strs = [
1391            ("otel.kind", "Internal"),
1392            ("rpc.system.name", "grpc"),
1393            ("gcp.client.service", "storage"),
1394            ("gcp.client.version", env!("CARGO_PKG_VERSION")),
1395            ("gcp.client.repo", "googleapis/google-cloud-rust"),
1396            ("gcp.client.artifact", "google-cloud-storage"),
1397            ("gcp.schema.url", "https://opentelemetry.io/schemas/1.39.0"),
1398        ];
1399
1400        for (k, v) in expected_strs {
1401            assert_eq!(
1402                span.attributes.get(k),
1403                Some(&AttributeValue::String(v.to_string().into())),
1404                "mismatched attribute `{k}` in `{method}` span"
1405            );
1406        }
1407
1408        let gen_k = format!("{method}.generation");
1409        assert_eq!(
1410            span.attributes.get(gen_k.as_str()),
1411            Some(&AttributeValue::Int64(expected_generation)),
1412            "mismatched `{gen_k}`"
1413        );
1414
1415        if let Some(s) = expected_size {
1416            if method == "append" {
1417                assert_eq!(
1418                    span.attributes.get("append.chunk_size"),
1419                    Some(&AttributeValue::UInt64(s as u64)),
1420                    "mismatched `append.chunk_size`"
1421                );
1422            } else {
1423                let size_k = format!("{method}.persisted_size");
1424                assert_eq!(
1425                    span.attributes.get(size_k.as_str()),
1426                    Some(&AttributeValue::Int64(s)),
1427                    "mismatched `{size_k}`"
1428                );
1429            }
1430        }
1431    }
1432}