1use super::tracing::{TracingObjectDescriptor, TracingResponse};
16use crate::Result;
17use crate::model::{Object, ReadObjectRequest};
18use crate::model_ext::WriteObjectRequest;
19use crate::read_object::ReadObjectResponse;
20use crate::storage::client::StorageInner;
21use crate::storage::info::INSTRUMENTATION;
22use crate::storage::perform_upload::PerformUpload;
23use crate::storage::read_object::Reader;
24use crate::storage::request_options::RequestOptions;
25use crate::storage::streaming_source::{Seek, StreamingSource};
26#[cfg(google_cloud_unstable_storage_bidi)]
27use crate::{
28 appendable_object_writer::AppendableObjectWriter,
29 model_ext::{OpenAppendableObjectRequest, ReopenAppendableObjectRequest},
30 storage::bidi_write::{
31 connector::Connector as BidiWriteConnector, transport::AppendableObjectWriterTransport,
32 },
33};
34use crate::{
35 model_ext::OpenObjectRequest, object_descriptor::ObjectDescriptor,
36 storage::bidi::connector::Connector, storage::bidi::transport::ObjectDescriptorTransport,
37};
38use gaxi::observability::{ClientRequestAttributes, DurationMetric, RequestRecorder};
39use std::sync::Arc;
40
41#[derive(Clone, Debug)]
57pub struct Storage {
58 inner: Arc<StorageInner>,
59 tracing: bool,
60 metric: DurationMetric,
61}
62
63impl Storage {
64 pub(crate) fn new(inner: Arc<StorageInner>, tracing: bool) -> Arc<Self> {
65 let metric = DurationMetric::new(&INSTRUMENTATION);
66 Arc::new(Self {
67 inner,
68 tracing,
69 metric,
70 })
71 }
72
73 async fn read_object_plain(
74 &self,
75 request: ReadObjectRequest,
76 options: RequestOptions,
77 ) -> Result<ReadObjectResponse> {
78 let reader = Reader {
79 inner: self.inner.clone(),
80 request,
81 options,
82 };
83 reader.response().await
84 }
85
86 #[tracing::instrument(name = "read_object", level = tracing::Level::DEBUG, ret, err(Debug))]
87 async fn read_object_tracing(
88 &self,
89 request: ReadObjectRequest,
90 options: RequestOptions,
91 ) -> Result<ReadObjectResponse> {
92 let resource_name = format!("//storage.googleapis.com/{}", request.bucket);
93 let (span, pending) = gaxi::client_request_signals!(
94 metric: self.metric.clone(),
95 info: *INSTRUMENTATION,
96 method: "client::Storage::read_object",
97 async {
98 if let Some(recorder) = RequestRecorder::current() {
99 recorder.on_client_request(
100 ClientRequestAttributes::default()
101 .set_url_template("/storage/v1/b/{bucket}/o/{object}")
102 .set_resource_name(resource_name),
103 );
104 }
105 self.read_object_plain(request, options).await
106 });
107
108 let response = pending.await?;
109 let inner = TracingResponse::new(response.into_parts(), span);
110 Ok(ReadObjectResponse::new(Box::new(inner)))
111 }
112
113 async fn write_object_buffered_plain<P>(
114 &self,
115 payload: P,
116 request: WriteObjectRequest,
117 options: RequestOptions,
118 ) -> Result<Object>
119 where
120 P: StreamingSource + Send + Sync + 'static,
121 {
122 PerformUpload::new(
123 payload,
124 self.inner.clone(),
125 request.spec,
126 request.params,
127 options,
128 )
129 .send()
130 .await
131 }
132
133 #[tracing::instrument(name = "write_object_buffered", level = tracing::Level::DEBUG, ret, err(Debug), skip(payload))]
134 async fn write_object_buffered_tracing<P>(
135 &self,
136 payload: P,
137 request: WriteObjectRequest,
138 options: RequestOptions,
139 ) -> Result<Object>
140 where
141 P: StreamingSource + Send + Sync + 'static,
142 {
143 let resource_name = format!(
144 "//storage.googleapis.com/{}",
145 request
146 .spec
147 .resource
148 .as_ref()
149 .map(|r| r.bucket.as_str())
150 .unwrap_or_default()
151 );
152 let (_span, pending) = gaxi::client_request_signals!(
153 metric: self.metric.clone(),
154 info: *INSTRUMENTATION,
155 method: "client::Storage::write_object",
156 async {
157 if let Some(recorder) = RequestRecorder::current() {
158 recorder.on_client_request(
159 ClientRequestAttributes::default()
160 .set_url_template("/upload/storage/v1/b/{bucket}/o")
161 .set_resource_name(resource_name),
162 );
163 }
164 self.write_object_buffered_plain(payload, request, options).await
165 }
166 );
167 pending.await
168 }
169
170 async fn write_object_unbuffered_plain<P>(
171 &self,
172 payload: P,
173 request: WriteObjectRequest,
174 options: RequestOptions,
175 ) -> Result<Object>
176 where
177 P: StreamingSource + Seek + Send + Sync + 'static,
178 {
179 PerformUpload::new(
180 payload,
181 self.inner.clone(),
182 request.spec,
183 request.params,
184 options,
185 )
186 .send_unbuffered(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 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 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 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 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 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 tx.send(Ok(response.clone())).await?;
773 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 #[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 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 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 #[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 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 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!(
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 #[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 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 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 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_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 #[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 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 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 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}