From 004cb2b8cb07ea171e683138f8a9f32ffe058c23 Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Mon, 17 Aug 2026 12:12:53 +0000 Subject: [PATCH 1/3] feat(storage): add range read latency metrics and trace annotations --- .../storage/google_cloud_cpp_storage_grpc.bzl | 2 + .../google_cloud_cpp_storage_grpc.cmake | 2 + .../internal/async/connection_tracing.cc | 28 +++-- .../object_descriptor_connection_tracing.cc | 15 ++- .../object_descriptor_connection_tracing.h | 3 +- ...ject_descriptor_connection_tracing_test.cc | 5 +- .../internal/async/object_descriptor_impl.cc | 7 +- .../async/object_descriptor_reader_tracing.cc | 34 ++++-- .../async/object_descriptor_reader_tracing.h | 3 +- .../object_descriptor_reader_tracing_test.cc | 4 +- .../cloud/storage/internal/async/read_range.h | 4 +- .../async/reader_connection_telemetry.cc | 109 ++++++++++++++++++ .../async/reader_connection_telemetry.h | 52 +++++++++ .../async/reader_connection_tracing.cc | 29 +++-- .../async/reader_connection_tracing.h | 3 +- .../async/reader_connection_tracing_test.cc | 52 ++++++++- 16 files changed, 300 insertions(+), 52 deletions(-) create mode 100644 google/cloud/storage/internal/async/reader_connection_telemetry.cc create mode 100644 google/cloud/storage/internal/async/reader_connection_telemetry.h diff --git a/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl b/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl index 13841ceacfb4d..19d0b585787c6 100644 --- a/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl +++ b/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl @@ -58,6 +58,7 @@ google_cloud_cpp_storage_grpc_hdrs = [ "internal/async/read_range.h", "internal/async/reader_connection_factory.h", "internal/async/reader_connection_impl.h", + "internal/async/reader_connection_telemetry.h", "internal/async/reader_connection_resume.h", "internal/async/reader_connection_tracing.h", "internal/async/rewriter_connection_impl.h", @@ -135,6 +136,7 @@ google_cloud_cpp_storage_grpc_srcs = [ "internal/async/read_range.cc", "internal/async/reader_connection_factory.cc", "internal/async/reader_connection_impl.cc", + "internal/async/reader_connection_telemetry.cc", "internal/async/reader_connection_resume.cc", "internal/async/reader_connection_tracing.cc", "internal/async/rewriter_connection_impl.cc", diff --git a/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake b/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake index 4e194e80a4345..b296c7e86118e 100644 --- a/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake +++ b/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake @@ -137,6 +137,8 @@ add_library( internal/async/reader_connection_factory.h internal/async/reader_connection_impl.cc internal/async/reader_connection_impl.h + internal/async/reader_connection_telemetry.cc + internal/async/reader_connection_telemetry.h internal/async/reader_connection_resume.cc internal/async/reader_connection_resume.h internal/async/reader_connection_tracing.cc diff --git a/google/cloud/storage/internal/async/connection_tracing.cc b/google/cloud/storage/internal/async/connection_tracing.cc index 2e29660276078..4cb8fcc51f1a8 100644 --- a/google/cloud/storage/internal/async/connection_tracing.cc +++ b/google/cloud/storage/internal/async/connection_tracing.cc @@ -56,19 +56,16 @@ class AsyncConnectionTracing : public storage::AsyncConnection { OpenParams p) override { auto span = internal::MakeSpan("storage::AsyncConnection::Open"); internal::OTelScope scope(span); - return impl_->Open(std::move(p)) - .then([oc = opentelemetry::context::RuntimeContext::GetCurrent(), - span = std::move(span)](auto f) - -> StatusOr< - std::shared_ptr> { - auto result = f.get(); - internal::DetachOTelContext(oc); - if (!result) { - return internal::EndSpan(*span, std::move(result).status()); - } - return MakeTracingObjectDescriptorConnection(std::move(span), - *std::move(result)); - }); + auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(), + bucket = p.read_spec.bucket(), span = std::move(span)](auto f) + -> StatusOr> { + auto result = f.get(); + internal::DetachOTelContext(oc); + if (!result) return internal::EndSpan(*span, std::move(result).status()); + return MakeTracingObjectDescriptorConnection( + std::move(span), *std::move(result), std::move(bucket)); + }; + return impl_->Open(std::move(p)).then(std::move(wrap)); } future>> ReadObject( @@ -76,12 +73,13 @@ class AsyncConnectionTracing : public storage::AsyncConnection { auto span = internal::MakeSpan("storage::AsyncConnection::ReadObject"); internal::OTelScope scope(span); auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(), - span = std::move(span)](auto f) + bucket = p.request.bucket(), span = std::move(span)](auto f) -> StatusOr> { auto reader = f.get(); internal::DetachOTelContext(oc); if (!reader) return internal::EndSpan(*span, std::move(reader).status()); - return MakeTracingReaderConnection(std::move(span), *std::move(reader)); + return MakeTracingReaderConnection(std::move(span), *std::move(reader), + std::move(bucket)); }; return impl_->ReadObject(std::move(p)).then(std::move(wrap)); } diff --git a/google/cloud/storage/internal/async/object_descriptor_connection_tracing.cc b/google/cloud/storage/internal/async/object_descriptor_connection_tracing.cc index 4c0d582628575..770ffdb3f62f2 100644 --- a/google/cloud/storage/internal/async/object_descriptor_connection_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_connection_tracing.cc @@ -34,8 +34,11 @@ class AsyncObjectDescriptorConnectionTracing public: explicit AsyncObjectDescriptorConnectionTracing( opentelemetry::nostd::shared_ptr span, - std::shared_ptr impl) - : span_(std::move(span)), impl_(std::move(impl)) {} + std::shared_ptr impl, + std::string bucket_name) + : span_(std::move(span)), + impl_(std::move(impl)), + bucket_name_(std::move(bucket_name)) {} ~AsyncObjectDescriptorConnectionTracing() override { internal::EndSpan(*span_); @@ -55,7 +58,7 @@ class AsyncObjectDescriptorConnectionTracing {{sc::thread::kThreadId, internal::CurrentThreadId()}, {"read-start", p.start}, {"read-length", p.length}}); - return MakeTracingReaderConnection(span_, std::move(result)); + return MakeTracingReaderConnection(span_, std::move(result), bucket_name_); } void MakeSubsequentStream() override { @@ -65,6 +68,7 @@ class AsyncObjectDescriptorConnectionTracing private: opentelemetry::nostd::shared_ptr span_; std::shared_ptr impl_; + std::string bucket_name_; }; } // namespace @@ -72,9 +76,10 @@ class AsyncObjectDescriptorConnectionTracing std::shared_ptr MakeTracingObjectDescriptorConnection( opentelemetry::nostd::shared_ptr span, - std::shared_ptr impl) { + std::shared_ptr impl, + std::string bucket_name) { return std::make_unique( - std::move(span), std::move(impl)); + std::move(span), std::move(impl), std::move(bucket_name)); } GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END diff --git a/google/cloud/storage/internal/async/object_descriptor_connection_tracing.h b/google/cloud/storage/internal/async/object_descriptor_connection_tracing.h index ffe78427a365a..feaadf4ca3f52 100644 --- a/google/cloud/storage/internal/async/object_descriptor_connection_tracing.h +++ b/google/cloud/storage/internal/async/object_descriptor_connection_tracing.h @@ -26,7 +26,8 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN std::shared_ptr MakeTracingObjectDescriptorConnection( opentelemetry::nostd::shared_ptr span, - std::shared_ptr impl); + std::shared_ptr impl, + std::string bucket_name); GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/object_descriptor_connection_tracing_test.cc b/google/cloud/storage/internal/async/object_descriptor_connection_tracing_test.cc index 5719e76c954a7..330843eeb63c0 100644 --- a/google/cloud/storage/internal/async/object_descriptor_connection_tracing_test.cc +++ b/google/cloud/storage/internal/async/object_descriptor_connection_tracing_test.cc @@ -79,7 +79,7 @@ TEST(ObjectDescriptorConnectionTracing, Read) { return std::make_unique(); }); auto actual = MakeTracingObjectDescriptorConnection( - internal::MakeSpan("test-span-name"), std::move(mock)); + internal::MakeSpan("test-span-name"), std::move(mock), "test-bucket"); auto f1 = actual->Read(ObjectDescriptorConnection::ReadParams{100, 200}); actual.reset(); @@ -114,7 +114,8 @@ TEST(ObjectDescriptorConnectionTracing, ReadThenRead) { }); auto connection = MakeTracingObjectDescriptorConnection( - internal::MakeSpan("test-span"), std::move(mock_connection)); + internal::MakeSpan("test-span"), std::move(mock_connection), + "test-bucket"); auto reader = connection->Read({}); auto f = reader->Read().then(expect_no_context); diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index 1358df96e49ba..751bcefb922ad 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -201,7 +201,8 @@ std::unique_ptr ObjectDescriptorImpl::Read( return std::unique_ptr( std::make_unique(std::move(range))); } - return MakeTracingObjectDescriptorReader(std::move(range)); + return MakeTracingObjectDescriptorReader(std::move(range), + read_object_spec_.bucket()); } auto it = stream_manager_->GetLeastBusyStream(); @@ -217,8 +218,8 @@ std::unique_ptr ObjectDescriptorImpl::Read( return std::unique_ptr( std::make_unique(std::move(range))); } - - return MakeTracingObjectDescriptorReader(std::move(range)); + return MakeTracingObjectDescriptorReader(std::move(range), + read_object_spec_.bucket()); } std::shared_ptr diff --git a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc index 8b95c508c3360..b39b8ba01584a 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc @@ -15,9 +15,11 @@ #include "google/cloud/storage/internal/async/object_descriptor_reader_tracing.h" #include "google/cloud/storage/async/reader_connection.h" #include "google/cloud/storage/internal/async/object_descriptor_reader.h" +#include "google/cloud/storage/internal/async/reader_connection_telemetry.h" #include "google/cloud/internal/opentelemetry.h" #include "google/cloud/version.h" #include +#include #include namespace google { @@ -31,8 +33,11 @@ namespace sc = ::opentelemetry::semconv; class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { public: - explicit ObjectDescriptorReaderTracing(std::shared_ptr impl) - : ObjectDescriptorReader(std::move(impl)) {} + explicit ObjectDescriptorReaderTracing(std::shared_ptr impl, + std::string bucket_name) + : ObjectDescriptorReader(std::move(impl)), + bucket_name_( + std::make_shared(std::move(bucket_name))) {} ~ObjectDescriptorReaderTracing() override = default; @@ -41,18 +46,21 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { internal::OTelScope scope(span); return ObjectDescriptorReader::Read().then( [span = std::move(span), - oc = opentelemetry::context::RuntimeContext::GetCurrent()]( - auto f) -> ReadResponse { + oc = opentelemetry::context::RuntimeContext::GetCurrent(), + bucket_name = bucket_name_, + metrics = metrics_](auto f) -> ReadResponse { auto result = f.get(); internal::DetachOTelContext(oc); - if (!absl::holds_alternative(result)) { - auto const& payload = absl::get(result); - + if (auto const* payload = + absl::get_if(&result)) { span->AddEvent( "gl-cpp.read-range", {{/*sc::kRpcMessageType=*/"rpc.message.type", "RECEIVED"}, {sc::thread::kThreadId, internal::CurrentThreadId()}, - {"message.size", static_cast(payload.size())}}); + {"message.size", + static_cast(payload->size())}}); + metrics.RecordRead(*payload, std::chrono::steady_clock::now(), + *bucket_name, span, "gl-cpp.latency.read-range"); } else { span->AddEvent( "gl-cpp.read-range", @@ -64,13 +72,19 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { return result; }); } + + private: + std::shared_ptr bucket_name_; + ReaderConnectionTelemetry metrics_; }; } // namespace std::unique_ptr -MakeTracingObjectDescriptorReader(std::shared_ptr impl) { - return std::make_unique(std::move(impl)); +MakeTracingObjectDescriptorReader(std::shared_ptr impl, + std::string bucket_name) { + return std::make_unique( + std::move(impl), std::move(bucket_name)); } GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END diff --git a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h index 1a010cdbb366b..3282a42c453e3 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h @@ -25,7 +25,8 @@ namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN std::unique_ptr -MakeTracingObjectDescriptorReader(std::shared_ptr impl); +MakeTracingObjectDescriptorReader(std::shared_ptr impl, + std::string bucket_name); GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/object_descriptor_reader_tracing_test.cc b/google/cloud/storage/internal/async/object_descriptor_reader_tracing_test.cc index 98dc46eb9045a..04837f3212ffd 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing_test.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing_test.cc @@ -45,7 +45,7 @@ TEST(ObjectDescriptorReaderTracing, Read) { auto span_catcher = InstallSpanCatcher(); auto impl = std::make_shared(10000, 30); - auto reader = MakeTracingObjectDescriptorReader(impl); + auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket"); auto data = google::storage::v2::ObjectRangeData{}; auto constexpr kData0 = R"pb( @@ -73,7 +73,7 @@ TEST(ObjectDescriptorReaderTracing, Read) { TEST(ObjectDescriptorReaderTracing, ReadError) { auto span_catcher = InstallSpanCatcher(); auto impl = std::make_shared(10000, 30); - auto reader = MakeTracingObjectDescriptorReader(impl); + auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket"); impl->OnFinish(PermanentError()); diff --git a/google/cloud/storage/internal/async/read_range.h b/google/cloud/storage/internal/async/read_range.h index 373b920612997..355b488a84f3a 100644 --- a/google/cloud/storage/internal/async/read_range.h +++ b/google/cloud/storage/internal/async/read_range.h @@ -87,11 +87,11 @@ class ReadRange { #if defined(GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS) || \ defined(GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY) void SetT4(std::chrono::steady_clock::time_point t) { - std::unique_lock lk(self->mu_); + std::unique_lock lk(mu_); t4_ = t; } void SetT5(std::chrono::steady_clock::time_point t) { - std::unique_lock lk(self->mu_); + std::unique_lock lk(mu_); t5_ = t; } #else diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.cc b/google/cloud/storage/internal/async/reader_connection_telemetry.cc new file mode 100644 index 0000000000000..3e40aa162beef --- /dev/null +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.cc @@ -0,0 +1,109 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include "google/cloud/storage/internal/async/reader_connection_telemetry.h" +#include "google/cloud/storage/internal/async/read_payload_impl.h" +#include "google/cloud/internal/opentelemetry.h" +#include "google/cloud/version.h" +#ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS +#include +#include +#endif +#include + +namespace google { +namespace cloud { +namespace storage_internal { +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN + +#ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS +namespace { +struct ReadLatencyMetrics { + opentelemetry::nostd::shared_ptr> + queue_hist; + opentelemetry::nostd::shared_ptr> + network_hist; + opentelemetry::nostd::shared_ptr> + output_hist; + + static ReadLatencyMetrics const& Instance() { + static auto const metrics = [] { + auto meter = + opentelemetry::metrics::Provider::GetMeterProvider()->GetMeter( + "google-cloud-cpp", version::version_string()); + return ReadLatencyMetrics{ + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.queue", + "Read Range Queue Latency", "us"), + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.network", + "Read Range Network Latency", "us"), + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.internal", + "Read Range Internal Overhead", "us"), + }; + }(); + return metrics; + } +}; +} // namespace + +void ReaderConnectionTelemetry::RecordMetrics(std::string const& bucket_name, + double p1, double p2, + double p3) const { + auto const& metrics = ReadLatencyMetrics::Instance(); + if (metrics.queue_hist) + metrics.queue_hist->Record(p1, {{"gcp.storage.bucket", bucket_name}}, + opentelemetry::context::Context{}); + if (metrics.network_hist) + metrics.network_hist->Record(p2, {{"gcp.storage.bucket", bucket_name}}, + opentelemetry::context::Context{}); + if (metrics.output_hist) + metrics.output_hist->Record(p3, {{"gcp.storage.bucket", bucket_name}}, + opentelemetry::context::Context{}); +} +#endif + +void ReaderConnectionTelemetry::RecordRead( + storage::ReadPayload const& payload, + std::chrono::steady_clock::time_point t7, std::string const& bucket_name, + opentelemetry::nostd::shared_ptr const& span, + absl::string_view event_name) const { + auto t4 = ReadPayloadImpl::GetT4(payload); + auto t5 = ReadPayloadImpl::GetT5(payload); + auto t6 = ReadPayloadImpl::GetT6(payload); + + if (t4.time_since_epoch().count() > 0 && t5.time_since_epoch().count() > 0 && + t6.time_since_epoch().count() > 0) { + auto p1 = std::chrono::duration(t5 - t4).count(); + auto p2 = std::chrono::duration(t6 - t5).count(); + auto p3 = std::chrono::duration(t7 - t6).count(); + +#ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS + RecordMetrics(bucket_name, p1, p2, p3); +#else + (void)bucket_name; +#endif + + if (span && span->GetContext().IsValid()) { + span->AddEvent(opentelemetry::nostd::string_view{event_name.data(), + event_name.size()}, + {{"gl-cpp.latency.queue", p1}, + {"gl-cpp.latency.network", p2}, + {"gl-cpp.latency.internal", p3}}); + } + } +} + +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END +} // namespace storage_internal +} // namespace cloud +} // namespace google diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.h b/google/cloud/storage/internal/async/reader_connection_telemetry.h new file mode 100644 index 0000000000000..b73c6013d373f --- /dev/null +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.h @@ -0,0 +1,52 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#ifndef GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READER_CONNECTION_TELEMETRY_H +#define GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READER_CONNECTION_TELEMETRY_H + +#include "google/cloud/storage/internal/async/read_payload_impl.h" +#include "google/cloud/internal/opentelemetry.h" +#include "google/cloud/version.h" +#include "absl/strings/string_view.h" +#include +#include + +namespace google { +namespace cloud { +namespace storage_internal { +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN + +class ReaderConnectionTelemetry { + public: + ReaderConnectionTelemetry() = default; + + void RecordRead( + storage::ReadPayload const& payload, + std::chrono::steady_clock::time_point t7, std::string const& bucket_name, + opentelemetry::nostd::shared_ptr const& span, + absl::string_view event_name = "gl-cpp.latency.read") const; + + private: +#ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS + void RecordMetrics(std::string const& bucket_name, double p1, double p2, + double p3) const; +#endif +}; + +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END +} // namespace storage_internal +} // namespace cloud +} // namespace google + +#endif // GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READER_CONNECTION_TELEMETRY_H diff --git a/google/cloud/storage/internal/async/reader_connection_tracing.cc b/google/cloud/storage/internal/async/reader_connection_tracing.cc index eef65f7bb7208..cc6e1d9fbec07 100644 --- a/google/cloud/storage/internal/async/reader_connection_tracing.cc +++ b/google/cloud/storage/internal/async/reader_connection_tracing.cc @@ -13,8 +13,11 @@ // limitations under the License. #include "google/cloud/storage/internal/async/reader_connection_tracing.h" +#include "google/cloud/storage/internal/async/read_payload_impl.h" +#include "google/cloud/storage/internal/async/reader_connection_telemetry.h" #include "google/cloud/internal/opentelemetry.h" #include +#include #include #include #include @@ -31,8 +34,12 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { public: explicit AsyncReaderConnectionTracing( opentelemetry::nostd::shared_ptr span, - std::unique_ptr impl) - : span_(std::move(span)), impl_(std::move(impl)) {} + std::unique_ptr impl, + std::string bucket_name) + : span_(std::move(span)), + impl_(std::move(impl)), + bucket_name_( + std::make_shared(std::move(bucket_name))) {} void Cancel() override { auto scope = opentelemetry::trace::Scope(span_); @@ -46,9 +53,10 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { future Read() override { internal::OTelScope scope(span_); return impl_->Read() - .then([count = ++count_, span = span_](auto f) -> ReadResponse { + .then([count = ++count_, span = span_, bucket_name = bucket_name_, + metrics = metrics_](auto f) -> ReadResponse { auto r = f.get(); - if (absl::holds_alternative(r)) { + if (auto const* status = absl::get_if(&r)) { span->AddEvent( "gl-cpp.read", { @@ -56,7 +64,7 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { {/*sc::kRpcMessageId=*/"rpc.message.id", count}, {sc::thread::kThreadId, internal::CurrentThreadId()}, }); - return internal::EndSpan(*span, absl::get(std::move(r))); + return internal::EndSpan(*span, *status); } auto const& payload = absl::get(r); span->AddEvent( @@ -67,6 +75,8 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { {sc::thread::kThreadId, internal::CurrentThreadId()}, {"message.starting_offset", payload.offset()}, }); + metrics.RecordRead(payload, std::chrono::steady_clock::now(), + *bucket_name, span); return r; }) .then([oc = opentelemetry::context::RuntimeContext::GetCurrent()]( @@ -85,15 +95,18 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { opentelemetry::nostd::shared_ptr span_; std::unique_ptr impl_; std::int64_t count_ = 0; + std::shared_ptr bucket_name_; + ReaderConnectionTelemetry metrics_; }; } // namespace std::unique_ptr MakeTracingReaderConnection( opentelemetry::nostd::shared_ptr span, - std::unique_ptr impl) { - return std::make_unique(std::move(span), - std::move(impl)); + std::unique_ptr impl, + std::string bucket_name) { + return std::make_unique( + std::move(span), std::move(impl), std::move(bucket_name)); } GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END diff --git a/google/cloud/storage/internal/async/reader_connection_tracing.h b/google/cloud/storage/internal/async/reader_connection_tracing.h index d0f8799dfd83e..44c8608dad8c6 100644 --- a/google/cloud/storage/internal/async/reader_connection_tracing.h +++ b/google/cloud/storage/internal/async/reader_connection_tracing.h @@ -28,7 +28,8 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN std::unique_ptr MakeTracingReaderConnection( opentelemetry::nostd::shared_ptr span, - std::unique_ptr impl); + std::unique_ptr impl, + std::string bucket_name); GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/reader_connection_tracing_test.cc b/google/cloud/storage/internal/async/reader_connection_tracing_test.cc index 9c5abdc4bf703..eda47c67fc301 100644 --- a/google/cloud/storage/internal/async/reader_connection_tracing_test.cc +++ b/google/cloud/storage/internal/async/reader_connection_tracing_test.cc @@ -12,7 +12,10 @@ // See the License for the specific language governing permissions and // limitations under the License. +#ifdef GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY + #include "google/cloud/storage/internal/async/reader_connection_tracing.h" +#include "google/cloud/storage/internal/async/read_payload_impl.h" #include "google/cloud/storage/mocks/mock_async_reader_connection.h" #include "google/cloud/storage/testing/canonical_errors.h" #include "google/cloud/internal/opentelemetry.h" @@ -79,7 +82,7 @@ TEST(ReaderConnectionTracing, WithError) { .WillOnce(expect_context(p1)) .WillOnce(expect_context(p2)); auto actual = MakeTracingReaderConnection( - internal::MakeSpan("test-span-name"), std::move(mock)); + internal::MakeSpan("test-span-name"), std::move(mock), "test-bucket"); auto f1 = actual->Read().then(expect_no_context); p1.set_value(ReadResponse(ReadPayload("m1"))); @@ -138,7 +141,7 @@ TEST(ReaderConnectionTracing, WithSuccess) { .WillOnce(Return(RpcMetadata{{{"hk0", "v0"}, {"hk1", "v1"}}, {{"tk0", "v0"}, {"tk1", "v1"}}})); auto actual = MakeTracingReaderConnection( - internal::MakeSpan("test-span-name"), std::move(mock)); + internal::MakeSpan("test-span-name"), std::move(mock), "test-bucket"); auto f1 = actual->Read().then(expect_no_context); p1.set_value(ReadResponse(ReadPayload("m1"))); @@ -197,8 +200,53 @@ TEST(ReaderConnectionTracing, WithSuccess) { UnorderedElementsAre(Pair("tk0", "v0"), Pair("tk1", "v1"))); } +TEST(ReaderConnectionTracing, ReadLatencyEvent) { + auto span_catcher = InstallSpanCatcher(); + PromiseWithOTelContext p1; + PromiseWithOTelContext p2; + + auto mock = std::make_unique(); + EXPECT_CALL(*mock, Read) + .WillOnce(expect_context(p1)) + .WillOnce(expect_context(p2)); + auto actual = MakeTracingReaderConnection( + internal::MakeSpan("test-span-name"), std::move(mock), "test-bucket"); + + auto f1 = actual->Read().then(expect_no_context); + auto now = std::chrono::steady_clock::now(); + auto payload = ReadPayload("m1"); + ReadPayloadImpl::SetTimestamps(payload, now, + now + std::chrono::microseconds(100), + now + std::chrono::microseconds(300)); + p1.set_value(ReadResponse(std::move(payload))); + (void)f1.get(); + + auto f2 = actual->Read().then(expect_no_context); + p2.set_value(ReadResponse(Status{})); + (void)f2.get(); + + auto spans = span_catcher->GetSpans(); + using EventMatcher = + testing::Matcher; + EXPECT_THAT( + spans, + ElementsAre(AllOf( + SpanNamed("test-span-name"), + SpanHasEvents( + EventMatcher(EventNamed("gl-cpp.read")), + EventMatcher(AllOf( + EventNamed("gl-cpp.latency.read"), + SpanEventAttributesAre( + OTelAttribute("gl-cpp.latency.queue", 100.0), + OTelAttribute("gl-cpp.latency.network", 200.0), + OTelAttribute("gl-cpp.latency.internal", + _)))))))); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal } // namespace cloud } // namespace google + +#endif // GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY From 2ce610a1385820a86a1829adb5d76c481475324a Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Tue, 18 Aug 2026 06:08:19 +0000 Subject: [PATCH 2/3] Address feedback from Gemini code assistant --- .../internal/async/connection_tracing.cc | 6 +++-- .../async/object_descriptor_reader_tracing.cc | 2 +- .../async/reader_connection_telemetry.cc | 26 +++++++++++-------- .../async/reader_connection_telemetry.h | 6 ++--- .../async/reader_connection_tracing.cc | 2 +- 5 files changed, 24 insertions(+), 18 deletions(-) diff --git a/google/cloud/storage/internal/async/connection_tracing.cc b/google/cloud/storage/internal/async/connection_tracing.cc index 4cb8fcc51f1a8..8f8ac51436c67 100644 --- a/google/cloud/storage/internal/async/connection_tracing.cc +++ b/google/cloud/storage/internal/async/connection_tracing.cc @@ -59,7 +59,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection { auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(), bucket = p.read_spec.bucket(), span = std::move(span)](auto f) -> StatusOr> { - auto result = f.get(); + StatusOr> result = + f.get(); internal::DetachOTelContext(oc); if (!result) return internal::EndSpan(*span, std::move(result).status()); return MakeTracingObjectDescriptorConnection( @@ -75,7 +76,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection { auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(), bucket = p.request.bucket(), span = std::move(span)](auto f) -> StatusOr> { - auto reader = f.get(); + StatusOr> reader = + f.get(); internal::DetachOTelContext(oc); if (!reader) return internal::EndSpan(*span, std::move(reader).status()); return MakeTracingReaderConnection(std::move(span), *std::move(reader), diff --git a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc index b39b8ba01584a..a7f07ebeb20de 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc @@ -49,7 +49,7 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { oc = opentelemetry::context::RuntimeContext::GetCurrent(), bucket_name = bucket_name_, metrics = metrics_](auto f) -> ReadResponse { - auto result = f.get(); + ReadResponse result = f.get(); internal::DetachOTelContext(oc); if (auto const* payload = absl::get_if(&result)) { diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.cc b/google/cloud/storage/internal/async/reader_connection_telemetry.cc index 3e40aa162beef..05ff40884398f 100644 --- a/google/cloud/storage/internal/async/reader_connection_telemetry.cc +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.cc @@ -38,8 +38,8 @@ struct ReadLatencyMetrics { output_hist; static ReadLatencyMetrics const& Instance() { - static auto const metrics = [] { - auto meter = + static ReadLatencyMetrics const metrics = [] { + opentelemetry::nostd::shared_ptr meter = opentelemetry::metrics::Provider::GetMeterProvider()->GetMeter( "google-cloud-cpp", version::version_string()); return ReadLatencyMetrics{ @@ -76,21 +76,25 @@ void ReaderConnectionTelemetry::RecordRead( storage::ReadPayload const& payload, std::chrono::steady_clock::time_point t7, std::string const& bucket_name, opentelemetry::nostd::shared_ptr const& span, - absl::string_view event_name) const { - auto t4 = ReadPayloadImpl::GetT4(payload); - auto t5 = ReadPayloadImpl::GetT5(payload); - auto t6 = ReadPayloadImpl::GetT6(payload); + std::string_view event_name) const { + std::chrono::steady_clock::time_point t4 = ReadPayloadImpl::GetT4(payload); + std::chrono::steady_clock::time_point t5 = ReadPayloadImpl::GetT5(payload); + std::chrono::steady_clock::time_point t6 = ReadPayloadImpl::GetT6(payload); - if (t4.time_since_epoch().count() > 0 && t5.time_since_epoch().count() > 0 && - t6.time_since_epoch().count() > 0) { - auto p1 = std::chrono::duration(t5 - t4).count(); - auto p2 = std::chrono::duration(t6 - t5).count(); - auto p3 = std::chrono::duration(t7 - t6).count(); + if (t4 != std::chrono::steady_clock::time_point{} && + t5 != std::chrono::steady_clock::time_point{} && + t6 != std::chrono::steady_clock::time_point{}) { + double p1 = std::chrono::duration(t5 - t4).count(); + double p2 = std::chrono::duration(t6 - t5).count(); + double p3 = std::chrono::duration(t7 - t6).count(); #ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS RecordMetrics(bucket_name, p1, p2, p3); #else (void)bucket_name; + (void)p1; + (void)p2; + (void)p3; #endif if (span && span->GetContext().IsValid()) { diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.h b/google/cloud/storage/internal/async/reader_connection_telemetry.h index b73c6013d373f..9b313217c0642 100644 --- a/google/cloud/storage/internal/async/reader_connection_telemetry.h +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.h @@ -15,12 +15,12 @@ #ifndef GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READER_CONNECTION_TELEMETRY_H #define GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READER_CONNECTION_TELEMETRY_H -#include "google/cloud/storage/internal/async/read_payload_impl.h" +#include "google/cloud/storage/async/reader_connection.h" #include "google/cloud/internal/opentelemetry.h" #include "google/cloud/version.h" -#include "absl/strings/string_view.h" #include #include +#include namespace google { namespace cloud { @@ -35,7 +35,7 @@ class ReaderConnectionTelemetry { storage::ReadPayload const& payload, std::chrono::steady_clock::time_point t7, std::string const& bucket_name, opentelemetry::nostd::shared_ptr const& span, - absl::string_view event_name = "gl-cpp.latency.read") const; + std::string_view event_name = "gl-cpp.latency.read") const; private: #ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS diff --git a/google/cloud/storage/internal/async/reader_connection_tracing.cc b/google/cloud/storage/internal/async/reader_connection_tracing.cc index cc6e1d9fbec07..e93b345e2e73d 100644 --- a/google/cloud/storage/internal/async/reader_connection_tracing.cc +++ b/google/cloud/storage/internal/async/reader_connection_tracing.cc @@ -55,7 +55,7 @@ class AsyncReaderConnectionTracing : public storage::AsyncReaderConnection { return impl_->Read() .then([count = ++count_, span = span_, bucket_name = bucket_name_, metrics = metrics_](auto f) -> ReadResponse { - auto r = f.get(); + ReadResponse r = f.get(); if (auto const* status = absl::get_if(&r)) { span->AddEvent( "gl-cpp.read", From 3c6a2741ecdb225d19e44566c5a9d769f1af379b Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Tue, 18 Aug 2026 06:31:12 +0000 Subject: [PATCH 3/3] Address feedback from code assistant --- .../async/reader_connection_telemetry.cc | 64 ++++++++----------- .../async/reader_connection_telemetry.h | 15 ++++- .../async/reader_connection_tracing_test.cc | 2 +- 3 files changed, 40 insertions(+), 41 deletions(-) diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.cc b/google/cloud/storage/internal/async/reader_connection_telemetry.cc index 05ff40884398f..38b208bdd3bc6 100644 --- a/google/cloud/storage/internal/async/reader_connection_telemetry.cc +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.cc @@ -28,48 +28,35 @@ namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN #ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS -namespace { -struct ReadLatencyMetrics { - opentelemetry::nostd::shared_ptr> - queue_hist; - opentelemetry::nostd::shared_ptr> - network_hist; - opentelemetry::nostd::shared_ptr> - output_hist; - - static ReadLatencyMetrics const& Instance() { - static ReadLatencyMetrics const metrics = [] { - opentelemetry::nostd::shared_ptr meter = - opentelemetry::metrics::Provider::GetMeterProvider()->GetMeter( - "google-cloud-cpp", version::version_string()); - return ReadLatencyMetrics{ - meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.queue", - "Read Range Queue Latency", "us"), - meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.network", - "Read Range Network Latency", "us"), - meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.internal", - "Read Range Internal Overhead", "us"), - }; - }(); - return metrics; - } -}; -} // namespace +ReaderConnectionTelemetry::ReaderConnectionTelemetry() { + opentelemetry::nostd::shared_ptr meter = + opentelemetry::metrics::Provider::GetMeterProvider()->GetMeter( + "google-cloud-cpp", version::version_string()); + metrics_ = { + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.queue", + "Read Range Queue Latency", "us"), + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.network", + "Read Range Network Latency", "us"), + meter->CreateDoubleHistogram("gl-cpp.latency.bidi_read.internal", + "Read Range Internal Overhead", "us"), + }; +} void ReaderConnectionTelemetry::RecordMetrics(std::string const& bucket_name, double p1, double p2, double p3) const { - auto const& metrics = ReadLatencyMetrics::Instance(); - if (metrics.queue_hist) - metrics.queue_hist->Record(p1, {{"gcp.storage.bucket", bucket_name}}, - opentelemetry::context::Context{}); - if (metrics.network_hist) - metrics.network_hist->Record(p2, {{"gcp.storage.bucket", bucket_name}}, - opentelemetry::context::Context{}); - if (metrics.output_hist) - metrics.output_hist->Record(p3, {{"gcp.storage.bucket", bucket_name}}, + if (metrics_.queue_hist) + metrics_.queue_hist->Record(p1, {{"gcp.storage.bucket", bucket_name}}, opentelemetry::context::Context{}); + if (metrics_.network_hist) + metrics_.network_hist->Record(p2, {{"gcp.storage.bucket", bucket_name}}, + opentelemetry::context::Context{}); + if (metrics_.output_hist) + metrics_.output_hist->Record(p3, {{"gcp.storage.bucket", bucket_name}}, + opentelemetry::context::Context{}); } +#else +ReaderConnectionTelemetry::ReaderConnectionTelemetry() = default; #endif void ReaderConnectionTelemetry::RecordRead( @@ -81,9 +68,8 @@ void ReaderConnectionTelemetry::RecordRead( std::chrono::steady_clock::time_point t5 = ReadPayloadImpl::GetT5(payload); std::chrono::steady_clock::time_point t6 = ReadPayloadImpl::GetT6(payload); - if (t4 != std::chrono::steady_clock::time_point{} && - t5 != std::chrono::steady_clock::time_point{} && - t6 != std::chrono::steady_clock::time_point{}) { + if (t4 != std::chrono::steady_clock::time_point{} && t4 <= t5 && t5 <= t6 && + t6 <= t7) { double p1 = std::chrono::duration(t5 - t4).count(); double p2 = std::chrono::duration(t6 - t5).count(); double p3 = std::chrono::duration(t7 - t6).count(); diff --git a/google/cloud/storage/internal/async/reader_connection_telemetry.h b/google/cloud/storage/internal/async/reader_connection_telemetry.h index 9b313217c0642..9aff1d862bd1c 100644 --- a/google/cloud/storage/internal/async/reader_connection_telemetry.h +++ b/google/cloud/storage/internal/async/reader_connection_telemetry.h @@ -18,6 +18,9 @@ #include "google/cloud/storage/async/reader_connection.h" #include "google/cloud/internal/opentelemetry.h" #include "google/cloud/version.h" +#ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS +#include +#endif #include #include #include @@ -29,7 +32,7 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN class ReaderConnectionTelemetry { public: - ReaderConnectionTelemetry() = default; + ReaderConnectionTelemetry(); void RecordRead( storage::ReadPayload const& payload, @@ -41,6 +44,16 @@ class ReaderConnectionTelemetry { #ifdef GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS void RecordMetrics(std::string const& bucket_name, double p1, double p2, double p3) const; + + struct ReadLatencyMetrics { + opentelemetry::nostd::shared_ptr> + queue_hist; + opentelemetry::nostd::shared_ptr> + network_hist; + opentelemetry::nostd::shared_ptr> + output_hist; + }; + ReadLatencyMetrics metrics_; #endif }; diff --git a/google/cloud/storage/internal/async/reader_connection_tracing_test.cc b/google/cloud/storage/internal/async/reader_connection_tracing_test.cc index eda47c67fc301..ffa57997659cf 100644 --- a/google/cloud/storage/internal/async/reader_connection_tracing_test.cc +++ b/google/cloud/storage/internal/async/reader_connection_tracing_test.cc @@ -213,7 +213,7 @@ TEST(ReaderConnectionTracing, ReadLatencyEvent) { internal::MakeSpan("test-span-name"), std::move(mock), "test-bucket"); auto f1 = actual->Read().then(expect_no_context); - auto now = std::chrono::steady_clock::now(); + auto now = std::chrono::steady_clock::now() - std::chrono::milliseconds(2); auto payload = ReadPayload("m1"); ReadPayloadImpl::SetTimestamps(payload, now, now + std::chrono::microseconds(100),