From e49ecc736051a36d03f9666c31bbfbf1ccb2c0c3 Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Wed, 5 Aug 2026 08:29:57 +0000 Subject: [PATCH 1/2] feat(storage): add fast_open_cache_status telemetry for pre-warmed ranges in ObjectDescriptorImpl --- .../internal/async/object_descriptor_impl.cc | 17 +++++++++++++---- .../internal/async/object_descriptor_impl.h | 2 ++ .../async/object_descriptor_reader_tracing.cc | 18 ++++++++++++++---- .../async/object_descriptor_reader_tracing.h | 3 ++- .../object_descriptor_reader_tracing_test.cc | 6 ++++-- 5 files changed, 35 insertions(+), 11 deletions(-) diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index ecd18be5322ed..1914249f78760 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -231,7 +231,10 @@ std::unique_ptr ObjectDescriptorImpl::Read( // Check if this range matches a pre-warmed range. auto cache_key = std::make_pair(p.start, p.length); auto cache_it = prewarmed_ranges_.find(cache_key); + absl::string_view cache_status = "MISS"; + if (cache_it != prewarmed_ranges_.end()) { + cache_status = "HIT"; // Cache hit. Claim the pre-warmed range and return it to the user. auto prewarmed = std::move(cache_it->second); prewarmed_ranges_.erase(cache_it); @@ -247,7 +250,13 @@ std::unique_ptr ObjectDescriptorImpl::Read( return std::unique_ptr( std::make_unique(std::move(prewarmed.range))); } - return MakeTracingObjectDescriptorReader(std::move(prewarmed.range)); + return MakeTracingObjectDescriptorReader(std::move(prewarmed.range), + cache_status); + } + + // If not hit, check if it was evicted earlier due to pacing. + if (evicted_ranges_.erase(cache_key) != 0) { + cache_status = "EVICTED"; } if (stream_manager_->Empty()) { @@ -258,7 +267,7 @@ 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), cache_status); } auto it = stream_manager_->GetLeastBusyStream(); @@ -274,8 +283,7 @@ 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), cache_status); } std::shared_ptr @@ -429,6 +437,7 @@ void ObjectDescriptorImpl::OnRead( max_prewarmed_buffer_size_) { // Evict the range if it exceeds the pacing limit. total_prewarmed_bytes_buffered_ -= unclaimed_it->second.bytes_buffered; + evicted_ranges_.insert(unclaimed_it->second.cache_it->first); prewarmed_ranges_.erase(unclaimed_it->second.cache_it); unclaimed_ranges_.erase(unclaimed_it); diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.h b/google/cloud/storage/internal/async/object_descriptor_impl.h index 0b982b9ca215e..d4ad62a2bcb80 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.h +++ b/google/cloud/storage/internal/async/object_descriptor_impl.h @@ -143,6 +143,8 @@ class ObjectDescriptorImpl // Map of read_id to unclaimed range state (bytes buffered and original key). std::unordered_map unclaimed_ranges_; + // Set capturing tombstones for evicted pre-warmed ranges. + std::set> evicted_ranges_; // Total bytes currently buffered across all unclaimed pre-warmed ranges. std::size_t total_prewarmed_bytes_buffered_ = 0; // Maximum bytes allowed to be buffered across all unclaimed pre-warmed ranges 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..fe9f23cb7154b 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc @@ -31,13 +31,18 @@ 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 cache_status) + : ObjectDescriptorReader(std::move(impl)), + cache_status_(std::move(cache_status)) {} ~ObjectDescriptorReaderTracing() override = default; future Read() override { auto span = internal::MakeSpan("storage::AsyncConnection::ReadRange"); + if (!cache_status_.empty()) { + span->SetAttribute("fast_open_cache_status", cache_status_); + } internal::OTelScope scope(span); return ObjectDescriptorReader::Read().then( [span = std::move(span), @@ -64,13 +69,18 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { return result; }); } + + private: + std::string cache_status_; }; } // namespace std::unique_ptr -MakeTracingObjectDescriptorReader(std::shared_ptr impl) { - return std::make_unique(std::move(impl)); +MakeTracingObjectDescriptorReader(std::shared_ptr impl, + absl::string_view cache_status) { + return std::make_unique( + std::move(impl), std::string(cache_status)); } 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..41fe54547582a 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, + absl::string_view cache_status); 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..9b5c13e5c862b 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,8 @@ TEST(ObjectDescriptorReaderTracing, Read) { auto span_catcher = InstallSpanCatcher(); auto impl = std::make_shared(10000, 30); - auto reader = MakeTracingObjectDescriptorReader(impl); + auto reader = + MakeTracingObjectDescriptorReader(impl, /*cache_status=*/"TEST"); auto data = google::storage::v2::ObjectRangeData{}; auto constexpr kData0 = R"pb( @@ -73,7 +74,8 @@ 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, /*cache_status=*/"TEST"); impl->OnFinish(PermanentError()); From b6c50a820cb2cbc59a643552a17f22811a850cd4 Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Wed, 5 Aug 2026 08:59:09 +0000 Subject: [PATCH 2/2] Address feedback from Gemini code assistant --- .../storage/internal/async/connection_tracing.cc | 5 +++++ .../internal/async/object_descriptor_impl.cc | 4 +++- .../async/object_descriptor_reader_tracing.cc | 13 ++++++------- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/google/cloud/storage/internal/async/connection_tracing.cc b/google/cloud/storage/internal/async/connection_tracing.cc index 469c46468b6e0..6247b5742a300 100644 --- a/google/cloud/storage/internal/async/connection_tracing.cc +++ b/google/cloud/storage/internal/async/connection_tracing.cc @@ -15,6 +15,7 @@ #include "google/cloud/storage/internal/async/connection_tracing.h" #include "google/cloud/storage/async/writer_connection.h" #include "google/cloud/storage/internal/async/object_descriptor_connection_tracing.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/internal/async/reader_connection_tracing.h" #include "google/cloud/storage/internal/async/rewriter_connection_tracing.h" #include "google/cloud/storage/internal/async/writer_connection_tracing.h" @@ -63,6 +64,10 @@ class AsyncConnectionTracing : public storage::AsyncConnection { OpenParams p) override { auto span = internal::MakeSpan("storage::AsyncConnection::Open"); EnrichSpan(*span, p.options, p.read_spec.bucket()); + if (p.options.has()) { + auto const& ranges = p.options.get(); + span->SetAttribute("fast_open_ranges", ranges.size()); + } internal::OTelScope scope(span); return impl_->Open(std::move(p)) .then([oc = opentelemetry::context::RuntimeContext::GetCurrent(), diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index 1914249f78760..9fd847da61137 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -437,7 +437,9 @@ void ObjectDescriptorImpl::OnRead( max_prewarmed_buffer_size_) { // Evict the range if it exceeds the pacing limit. total_prewarmed_bytes_buffered_ -= unclaimed_it->second.bytes_buffered; - evicted_ranges_.insert(unclaimed_it->second.cache_it->first); + if (evicted_ranges_.size() < 1000) { + evicted_ranges_.insert(unclaimed_it->second.cache_it->first); + } prewarmed_ranges_.erase(unclaimed_it->second.cache_it); unclaimed_ranges_.erase(unclaimed_it); 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 fe9f23cb7154b..c173ae0908895 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc @@ -32,16 +32,15 @@ namespace sc = ::opentelemetry::semconv; class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { public: explicit ObjectDescriptorReaderTracing(std::shared_ptr impl, - std::string cache_status) - : ObjectDescriptorReader(std::move(impl)), - cache_status_(std::move(cache_status)) {} + absl::string_view cache_status) + : ObjectDescriptorReader(std::move(impl)), cache_status_(cache_status) {} ~ObjectDescriptorReaderTracing() override = default; future Read() override { auto span = internal::MakeSpan("storage::AsyncConnection::ReadRange"); if (!cache_status_.empty()) { - span->SetAttribute("fast_open_cache_status", cache_status_); + span->SetAttribute("fast_open_cache_status", std::string(cache_status_)); } internal::OTelScope scope(span); return ObjectDescriptorReader::Read().then( @@ -71,7 +70,7 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { } private: - std::string cache_status_; + absl::string_view cache_status_; }; } // namespace @@ -79,8 +78,8 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { std::unique_ptr MakeTracingObjectDescriptorReader(std::shared_ptr impl, absl::string_view cache_status) { - return std::make_unique( - std::move(impl), std::string(cache_status)); + return std::make_unique(std::move(impl), + cache_status); } GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END