From f47af177dbcda7a0b4330a87e2f81d738c8e620f Mon Sep 17 00:00:00 2001 From: kalragauri Date: Mon, 3 Aug 2026 13:02:29 +0530 Subject: [PATCH 1/4] feat(storage): support pre-warming read ranges in AsyncConnection (#16275) --- .../storage/google_cloud_cpp_storage_grpc.bzl | 1 + .../google_cloud_cpp_storage_grpc.cmake | 1 + .../storage/internal/async/connection_impl.cc | 15 ++ .../async/connection_impl_open_test.cc | 145 +++++++++++++ .../internal/async/object_descriptor_impl.cc | 62 +++++- .../internal/async/object_descriptor_impl.h | 11 + .../async/object_descriptor_impl_test.cc | 192 ++++++++++++++++++ google/cloud/storage/internal/async/options.h | 48 +++++ .../storage/internal/async/read_range.cc | 19 ++ .../cloud/storage/internal/async/read_range.h | 13 ++ .../storage/internal/async/read_range_test.cc | 33 +++ 11 files changed, 535 insertions(+), 5 deletions(-) create mode 100644 google/cloud/storage/internal/async/options.h diff --git a/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl b/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl index 71f2e56b50e2c..a74bcfa54080b 100644 --- a/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl +++ b/google/cloud/storage/google_cloud_cpp_storage_grpc.bzl @@ -52,6 +52,7 @@ google_cloud_cpp_storage_grpc_hdrs = [ "internal/async/open_object.h", "internal/async/open_object_telemetry.h", "internal/async/open_stream.h", + "internal/async/options.h", "internal/async/partial_upload.h", "internal/async/read_payload_fwd.h", "internal/async/read_payload_impl.h", diff --git a/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake b/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake index 966f4c2b015f4..8146ee5147989 100644 --- a/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake +++ b/google/cloud/storage/google_cloud_cpp_storage_grpc.cmake @@ -127,6 +127,7 @@ add_library( internal/async/open_object_telemetry.h internal/async/open_stream.cc internal/async/open_stream.h + internal/async/options.h internal/async/partial_upload.cc internal/async/partial_upload.h internal/async/read_payload_fwd.h diff --git a/google/cloud/storage/internal/async/connection_impl.cc b/google/cloud/storage/internal/async/connection_impl.cc index 5516c23cfcd41..344bd26e86cbf 100644 --- a/google/cloud/storage/internal/async/connection_impl.cc +++ b/google/cloud/storage/internal/async/connection_impl.cc @@ -26,7 +26,9 @@ #include "google/cloud/storage/internal/async/object_descriptor_impl.h" #include "google/cloud/storage/internal/async/open_object.h" #include "google/cloud/storage/internal/async/open_stream.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/internal/async/read_payload_impl.h" +#include "google/cloud/storage/internal/async/read_range.h" #include "google/cloud/storage/internal/async/reader_connection_impl.h" #include "google/cloud/storage/internal/async/reader_connection_resume.h" #include "google/cloud/storage/internal/async/rewriter_connection_impl.h" @@ -58,6 +60,7 @@ #include "google/cloud/internal/make_status.h" #include #include +#include #include namespace google { @@ -220,6 +223,18 @@ AsyncConnectionImpl::Open(OpenParams p) { auto initial_request = google::storage::v2::BidiReadObjectRequest{}; *initial_request.mutable_read_object_spec() = p.read_spec; auto current = internal::MakeImmutableOptions(p.options); + // If pre-warmed ranges are configured, populate the initial request + // with these ranges to start downloading them as soon as the stream opens. + if (current->has()) { + auto const& ranges = current->get(); + for (auto const& r : DeduplicateRanges(ranges)) { + auto* proto_range = initial_request.add_read_ranges(); + proto_range->set_read_offset(r.config.offset); + proto_range->set_read_length(r.config.length); + // Generate sequential IDs starting at 1. The receiver must match these. + proto_range->set_read_id(r.read_id); + } + } // Get the policy factory and immediately create a policy. auto resume_policy = current->get()(); diff --git a/google/cloud/storage/internal/async/connection_impl_open_test.cc b/google/cloud/storage/internal/async/connection_impl_open_test.cc index e619490e6c7c0..32609f98aec4e 100644 --- a/google/cloud/storage/internal/async/connection_impl_open_test.cc +++ b/google/cloud/storage/internal/async/connection_impl_open_test.cc @@ -16,6 +16,7 @@ #include "google/cloud/storage/async/retry_policy.h" #include "google/cloud/storage/internal/async/connection_impl.h" #include "google/cloud/storage/internal/async/default_options.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/internal/grpc/ctype_cord_workaround.h" #include "google/cloud/storage/testing/canonical_errors.h" #include "google/cloud/storage/testing/mock_resume_policy.h" @@ -405,6 +406,150 @@ TEST(AsyncConnectionImplTest, TooManyTransienErrors) { ASSERT_THAT(pending.get(), StatusIs(TransientError().code())); } +TEST(AsyncConnectionImplTest, OpenWithReadRanges) { + auto constexpr kExpectedRequestSpec = R"pb( + bucket: "test-only-invalid" + object: "test-object" + generation: 42 + if_metageneration_match: 7 + )pb"; + auto constexpr kExpectedRequest = R"pb( + read_object_spec { + bucket: "test-only-invalid" + object: "test-object" + generation: 42 + if_metageneration_match: 7 + } + read_ranges { read_offset: 0 read_length: 1024 read_id: 1 } + read_ranges { read_offset: 2048 read_length: 4096 read_id: 2 } + )pb"; + auto constexpr kMetadataText = R"pb( + bucket: "projects/_/buckets/test-bucket" + name: "test-object" + generation: 42 + )pb"; + + AsyncSequencer sequencer; + auto mock = std::make_shared(); + EXPECT_CALL(*mock, AsyncBidiReadObject) + .WillOnce([&](CompletionQueue const&, + std::shared_ptr const&, + google::cloud::internal::ImmutableOptions const& options) { + EXPECT_EQ(options->get(), kAuthority); + + auto stream = std::make_unique(); + EXPECT_CALL(*stream, Start).WillOnce([&sequencer]() { + return sequencer.PushBack("Start").then( + [](auto f) { return f.get(); }); + }); + EXPECT_CALL(*stream, Write) + .WillOnce( + [=, &sequencer]( + google::storage::v2::BidiReadObjectRequest const& request, + grpc::WriteOptions) { + auto expected = google::storage::v2::BidiReadObjectRequest{}; + EXPECT_TRUE( + TextFormat::ParseFromString(kExpectedRequest, &expected)); + EXPECT_THAT(request, IsProtoEqual(expected)); + return sequencer.PushBack("Write").then( + [](auto f) { return f.get(); }); + }); + EXPECT_CALL(*stream, Read) + .WillOnce([&]() { + return sequencer.PushBack("Read").then( + [=](auto f) -> absl::optional< + google::storage::v2::BidiReadObjectResponse> { + if (!f.get()) return absl::nullopt; + auto constexpr kHandleText = R"pb( + handle: "handle-12345" + )pb"; + auto response = + google::storage::v2::BidiReadObjectResponse{}; + EXPECT_TRUE(TextFormat::ParseFromString( + kMetadataText, response.mutable_metadata())); + EXPECT_TRUE(TextFormat::ParseFromString( + kHandleText, response.mutable_read_handle())); + return response; + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[N]").then( + [](auto f) -> absl::optional< + google::storage::v2::BidiReadObjectResponse> { + if (!f.get()) return absl::nullopt; + return google::storage::v2::BidiReadObjectResponse{}; + }); + }); + EXPECT_CALL(*stream, Cancel).WillOnce([&sequencer]() { + (void)sequencer.PushBack("Cancel"); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return Status{}; }); + }); + + return std::unique_ptr(std::move(stream)); + }) + .WillRepeatedly([](CompletionQueue const&, + std::shared_ptr const&, + google::cloud::internal::ImmutableOptions const&) { + auto stream = std::make_unique>(); + ON_CALL(*stream, Start).WillByDefault(InvokeWithoutArgs([] { + return make_ready_future(false); + })); + ON_CALL(*stream, Finish).WillByDefault(InvokeWithoutArgs([] { + return make_ready_future(Status{}); + })); + ON_CALL(*stream, Cancel).WillByDefault([] {}); + return std::unique_ptr(std::move(stream)); + }); + + auto mock_cq = std::make_shared(); + EXPECT_CALL(*mock_cq, MakeRelativeTimer) + .WillRepeatedly([](std::chrono::nanoseconds) { + return make_ready_future( + StatusOr( + std::chrono::system_clock::now())); + }); + auto connection = std::make_shared( + CompletionQueue(mock_cq), std::shared_ptr(), mock, + TestOptions()); + + auto request = google::storage::v2::BidiReadObjectSpec{}; + ASSERT_TRUE(TextFormat::ParseFromString(kExpectedRequestSpec, &request)); + auto options = + connection->options().set({{0, 1024}, {2048, 4096}}); + auto pending = connection->Open({std::move(request), std::move(options)}); + + auto next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Start"); + next.first.set_value(true); + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Write"); + next.first.set_value(true); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Read"); + next.first.set_value(true); + + auto p = pending.get(); + ASSERT_THAT(p, IsOkAndHolds(NotNull())); + auto descriptor = *std::move(p); + + descriptor.reset(); + + auto last_read = sequencer.PopFrontWithName(); + EXPECT_EQ(last_read.second, "Read[N]"); + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Cancel"); + next.first.set_value(true); + last_read.first.set_value(false); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Finish"); + next.first.set_value(true); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index 751bcefb922ad..a6c70ffef205d 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -18,6 +18,8 @@ #include "google/cloud/storage/internal/async/handle_redirect_error.h" #include "google/cloud/storage/internal/async/multi_stream_manager.h" #include "google/cloud/storage/internal/async/object_descriptor_reader_tracing.h" +#include "google/cloud/storage/internal/async/options.h" +#include "google/cloud/storage/internal/async/read_range.h" #include "google/cloud/storage/internal/grpc/object_metadata_parser.h" #include "google/cloud/storage/internal/hash_function.h" #include "google/cloud/storage/internal/hash_function_impl.h" @@ -28,6 +30,7 @@ #include "google/cloud/grpc_error_delegate.h" #include "google/cloud/internal/opentelemetry.h" #include "google/rpc/status.pb.h" +#include #include #include #include @@ -52,6 +55,31 @@ ObjectDescriptorImpl::ObjectDescriptorImpl( []() -> std::shared_ptr { return nullptr; }, // NOLINT std::make_shared(std::move(stream), resume_policy_prototype_->clone())); + // If pre-warmed ranges are specified, initialize their `ReadRange` objects, + // register them as active on the initial stream, and cache them. + if (options_.has()) { + auto const& ranges = options_.get(); + auto it = stream_manager_->GetFirstStream(); + if (it != stream_manager_->End()) { + auto deduped_ranges = DeduplicateRanges(ranges); + for (auto const& dr : deduped_ranges) { + auto range_key = std::make_pair(dr.config.offset, dr.config.length); + auto range = std::make_shared( + dr.config.offset, dr.config.length, read_object_spec_.bucket(), + read_object_spec_.object()); + // Registering on the stream allows `OnRead` to route incoming data to + // these ranges. + it->active_ranges.emplace(dr.read_id, range); + // Cache them so subsequent `Read()` calls can claim them. + prewarmed_ranges_.emplace(range_key, PrewarmedRange{range, dr.read_id}); + } + // Ensure new dynamically requested ranges use IDs that don't conflict + // with pre-warmed ones. + if (!deduped_ranges.empty()) { + read_id_generator_ = deduped_ranges.back().read_id; + } + } + } } ObjectDescriptorImpl::~ObjectDescriptorImpl() { Cancel(); } @@ -193,6 +221,22 @@ std::unique_ptr ObjectDescriptorImpl::Read( read_object_spec_.bucket(), read_object_spec_.object()); std::unique_lock lk(mu_); + // 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); + if (cache_it != prewarmed_ranges_.end()) { + // 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); + lk.unlock(); + if (!internal::TracingEnabled(options_)) { + return std::unique_ptr( + std::make_unique(std::move(prewarmed.range))); + } + return MakeTracingObjectDescriptorReader(std::move(prewarmed.range), + read_object_spec_.bucket()); + } + if (stream_manager_->Empty()) { lk.unlock(); range->OnFinish(Status(StatusCode::kFailedPrecondition, @@ -385,14 +429,22 @@ void ObjectDescriptorImpl::OnRead( auto const l = copy.find(id); if (l == copy.end()) continue; + auto range = l->second; + bool active = false; + lk.lock(); + // Verify the range is still active on this stream. It might have been + // evicted or cancelled during the processing of this batch. + active = it->active_ranges.count(id) != 0; + lk.unlock(); + if (active) { #if defined(GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS) || \ defined(GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY) - l->second->SetT5(t5_stamp); + range->SetT5(t5_stamp); #endif - - // TODO(#15104) - Consider returning if the range is done, and then - // skipping CleanupDoneRanges(). - l->second->OnRead(std::move(range_data), is_transcoded, object_size); + // TODO(#15104) - Consider returning if the range is done, and then + // skipping CleanupDoneRanges(). + range->OnRead(std::move(range_data), is_transcoded, object_size); + } } lk.lock(); stream_manager_->CleanupDoneRanges(it); diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.h b/google/cloud/storage/internal/async/object_descriptor_impl.h index c7b2e2abffc0b..023a0084889b2 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.h +++ b/google/cloud/storage/internal/async/object_descriptor_impl.h @@ -27,9 +27,11 @@ #include "google/storage/v2/storage.pb.h" #include #include +#include #include #include #include +#include #include namespace google { @@ -124,6 +126,15 @@ class ObjectDescriptorImpl std::optional metadata_; std::int64_t read_id_generator_ = 0; + // Information about a pre-warmed range that has not been claimed yet. + struct PrewarmedRange { + std::shared_ptr range; + std::int64_t read_id; + }; + // Cache of pre-warmed ranges, keyed by (offset, length). + std::map, PrewarmedRange> + prewarmed_ranges_; + Options options_; std::unique_ptr stream_manager_; // The future for the proactive background stream. diff --git a/google/cloud/storage/internal/async/object_descriptor_impl_test.cc b/google/cloud/storage/internal/async/object_descriptor_impl_test.cc index 07830a3e63fdb..e29f15e53bedc 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl_test.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl_test.cc @@ -13,6 +13,7 @@ // limitations under the License. #include "google/cloud/storage/internal/async/object_descriptor_impl.h" +#include "google/cloud/storage/internal/async/options.h" // TODO(v-pratap): Remove this when EnableMD5ValidationOption and // EnableCrc32cValidationOption are removed. @@ -2298,6 +2299,197 @@ TEST(ObjectDescriptorImpl, PartialReadChecksumValidationBypassed) { next.first.set_value(true); } +TEST(ObjectDescriptorImpl, PrewarmedCacheHit) { + auto constexpr kResponse0 = R"pb( + metadata { + bucket: "projects/_/buckets/test-bucket" + name: "test-object" + generation: 42 + } + read_handle { handle: "handle-12345" } + )pb"; + + auto constexpr kResponse1 = R"pb( + read_handle { handle: "handle-23456" } + object_data_ranges { + range_end: true + read_range { read_id: 1 read_offset: 0 } + checksummed_data { content: "Pre-warmed data" } + } + )pb"; + + AsyncSequencer sequencer; + auto stream = std::make_unique(); + EXPECT_CALL(*stream, Write).Times(0); + + EXPECT_CALL(*stream, Read) + .WillOnce([=, &sequencer]() { + return sequencer.PushBack("Read[1]").then([&](auto) { + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); + return absl::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return absl::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + EXPECT_CALL(*stream, Cancel).Times(AtMost(1)); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 15}}); + + auto tested = std::make_shared( + NoResume(), factory.AsStdFunction(), + google::storage::v2::BidiReadObjectSpec{}, + std::make_shared(std::move(stream)), options); + + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse0, &response)); + tested->Start(std::move(response)); + + auto read1 = sequencer.PopFrontWithName(); + EXPECT_EQ(read1.second, "Read[1]"); + + auto s1 = tested->Read({0, 15}); + ASSERT_THAT(s1, NotNull()); + + auto s1r1 = s1->Read(); + EXPECT_FALSE(s1r1.is_ready()); + + read1.first.set_value(true); + + EXPECT_TRUE(s1r1.is_ready()); + EXPECT_THAT(s1r1.get(), + VariantWith(ResultOf( + "contents are", + [](storage::ReadPayload const& p) { return p.contents(); }, + ElementsAre(absl::string_view{"Pre-warmed data"})))); + + EXPECT_THAT(s1->Read().get(), VariantWith(IsOk())); + + auto next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Read[2]"); + next.first.set_value(true); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Finish"); + next.first.set_value(true); +} + +TEST(ObjectDescriptorImpl, PrewarmedCacheMiss) { + auto constexpr kResponse0 = R"pb( + metadata { + bucket: "projects/_/buckets/test-bucket" + name: "test-object" + generation: 42 + } + read_handle { handle: "handle-12345" } + )pb"; + + auto constexpr kExpectedRequest = R"pb( + read_ranges { read_id: 2 read_offset: 100 read_length: 15 } + )pb"; + + auto constexpr kResponse1 = R"pb( + object_data_ranges { + range_end: true + read_range { read_id: 2 read_offset: 100 } + checksummed_data { content: "Missed range data" } + } + )pb"; + + AsyncSequencer sequencer; + auto stream = std::make_unique(); + EXPECT_CALL(*stream, Write) + .WillOnce([&](Request const& request, grpc::WriteOptions) { + auto expected = Request{}; + EXPECT_TRUE(TextFormat::ParseFromString(kExpectedRequest, &expected)); + EXPECT_THAT(request, IsProtoEqual(expected)); + return sequencer.PushBack("Write[1]").then([](auto f) { + return f.get(); + }); + }); + + EXPECT_CALL(*stream, Read) + .WillOnce([=, &sequencer]() { + return sequencer.PushBack("Read[1]").then([&](auto) { + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); + return absl::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return absl::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + EXPECT_CALL(*stream, Cancel).Times(AtMost(1)); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 15}}); + + auto tested = std::make_shared( + NoResume(), factory.AsStdFunction(), + google::storage::v2::BidiReadObjectSpec{}, + std::make_shared(std::move(stream)), options); + + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse0, &response)); + tested->Start(std::move(response)); + + auto read1 = sequencer.PopFrontWithName(); + EXPECT_EQ(read1.second, "Read[1]"); + + auto s1 = tested->Read({100, 15}); + ASSERT_THAT(s1, NotNull()); + + auto next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Write[1]"); + next.first.set_value(true); + + auto s1r1 = s1->Read(); + EXPECT_FALSE(s1r1.is_ready()); + + read1.first.set_value(true); + + EXPECT_TRUE(s1r1.is_ready()); + EXPECT_THAT(s1r1.get(), + VariantWith(ResultOf( + "contents are", + [](storage::ReadPayload const& p) { return p.contents(); }, + ElementsAre(absl::string_view{"Missed range data"})))); + + EXPECT_THAT(s1->Read().get(), VariantWith(IsOk())); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Read[2]"); + next.first.set_value(true); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Finish"); + next.first.set_value(true); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/options.h b/google/cloud/storage/internal/async/options.h new file mode 100644 index 0000000000000..f6fe7fa70593c --- /dev/null +++ b/google/cloud/storage/internal/async/options.h @@ -0,0 +1,48 @@ +// 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_OPTIONS_H +#define GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_OPTIONS_H + +#include "google/cloud/options.h" +#include "google/cloud/version.h" +#include +#include + +namespace google { +namespace cloud { +namespace storage_internal { +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN + +// Configuration for a single read range to be pre-warmed. +struct ReadRangeConfig { + std::int64_t offset; + std::int64_t length; +}; + +// Internal option to pass the list of ranges to pre-warm from `AsyncClient` to +// `Connection`. +struct ReadRangesOption { + using Type = std::vector; + static char const* name() { + return "google::cloud::storage::ReadRangesOption"; + } +}; + +GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END +} // namespace storage_internal +} // namespace cloud +} // namespace google + +#endif // GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_OPTIONS_H diff --git a/google/cloud/storage/internal/async/read_range.cc b/google/cloud/storage/internal/async/read_range.cc index dabfe635db73d..f1dd4af2192fd 100644 --- a/google/cloud/storage/internal/async/read_range.cc +++ b/google/cloud/storage/internal/async/read_range.cc @@ -19,12 +19,31 @@ #include "google/cloud/log.h" #include "absl/strings/str_cat.h" #include +#include +#include namespace google { namespace cloud { namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN +std::vector DeduplicateRanges( + std::vector const& ranges, std::int64_t initial_id) { + std::vector deduped; + deduped.reserve(ranges.size()); + std::int64_t id = initial_id; + std::set> seen_ranges; + for (auto const& r : ranges) { + auto range_key = std::make_pair(r.offset, r.length); + if (!seen_ranges.insert(range_key).second) { + continue; + } + ++id; + deduped.push_back({r, id}); + } + return deduped; +} + bool ReadRange::IsDone() const { std::lock_guard lk(mu_); return status_.has_value(); diff --git a/google/cloud/storage/internal/async/read_range.h b/google/cloud/storage/internal/async/read_range.h index 355b488a84f3a..609652ae78f43 100644 --- a/google/cloud/storage/internal/async/read_range.h +++ b/google/cloud/storage/internal/async/read_range.h @@ -16,6 +16,7 @@ #define GOOGLE_CLOUD_CPP_GOOGLE_CLOUD_STORAGE_INTERNAL_ASYNC_READ_RANGE_H #include "google/cloud/storage/async/reader_connection.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/internal/hash_function.h" #include "google/cloud/storage/internal/hash_validator.h" #include "google/cloud/future.h" @@ -27,12 +28,24 @@ #include #include #include +#include namespace google { namespace cloud { namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN +struct DedupedReadRange { + ReadRangeConfig config; + std::int64_t read_id; +}; + +// Deduplicates a list of ReadRangeConfigs and assigns a sequential +// read_id to each unique range, starting with initial_id + 1. +// Returns a vector mapping each assigned read_id to its corresponding range. +std::vector DeduplicateRanges( + std::vector const& ranges, std::int64_t initial_id = 0); + /** * A read range represents a partially completed range download via a * `ObjectDescriptor`. diff --git a/google/cloud/storage/internal/async/read_range_test.cc b/google/cloud/storage/internal/async/read_range_test.cc index 7c2ce1486cd12..1f327dc56db34 100644 --- a/google/cloud/storage/internal/async/read_range_test.cc +++ b/google/cloud/storage/internal/async/read_range_test.cc @@ -591,6 +591,39 @@ TEST(ReadRange, NoResumeIfRequestExceeded) { EXPECT_FALSE(resume.has_value()); } +TEST(ReadRangeDeduplicationTest, Basic) { + std::vector ranges = { + {0, 10}, {10, 10}, {0, 10}, // Duplicate + {20, 10}, {10, 10}, // Duplicate + }; + + auto deduped = DeduplicateRanges(ranges); + ASSERT_EQ(deduped.size(), 3); + EXPECT_EQ(deduped[0].config.offset, 0); + EXPECT_EQ(deduped[0].config.length, 10); + EXPECT_EQ(deduped[0].read_id, 1); + + EXPECT_EQ(deduped[1].config.offset, 10); + EXPECT_EQ(deduped[1].config.length, 10); + EXPECT_EQ(deduped[1].read_id, 2); + + EXPECT_EQ(deduped[2].config.offset, 20); + EXPECT_EQ(deduped[2].config.length, 10); + EXPECT_EQ(deduped[2].read_id, 3); +} + +TEST(ReadRangeDeduplicationTest, InitialId) { + std::vector ranges = { + {0, 10}, + }; + + auto deduped = DeduplicateRanges(ranges, /*initial_id=*/5); + ASSERT_EQ(deduped.size(), 1); + EXPECT_EQ(deduped[0].config.offset, 0); + EXPECT_EQ(deduped[0].config.length, 10); + EXPECT_EQ(deduped[0].read_id, 6); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal From d9c09bec556c1d990e66e88cfba3d3346a2c1e78 Mon Sep 17 00:00:00 2001 From: kalragauri Date: Wed, 5 Aug 2026 10:42:13 +0530 Subject: [PATCH 2/4] feat(storage): implement pacing & eviction for pre-warmed ranges in ObjectDescriptorImpl (#16309) --- .../internal/async/object_descriptor_impl.cc | 71 ++++++++++++-- .../internal/async/object_descriptor_impl.h | 19 +++- .../async/object_descriptor_impl_test.cc | 93 +++++++++++++++++++ google/cloud/storage/internal/async/options.h | 16 +++- .../cloud/storage/internal/async/read_range.h | 8 +- .../storage/internal/async/read_range_test.cc | 4 +- 6 files changed, 191 insertions(+), 20 deletions(-) diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index a6c70ffef205d..ecba87a538c16 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -55,6 +55,10 @@ ObjectDescriptorImpl::ObjectDescriptorImpl( []() -> std::shared_ptr { return nullptr; }, // NOLINT std::make_shared(std::move(stream), resume_policy_prototype_->clone())); + // Initialize the pacing limit from options if configured. + if (options_.has()) { + max_prewarmed_buffer_size_ = options_.get(); + } // If pre-warmed ranges are specified, initialize their `ReadRange` objects, // register them as active on the initial stream, and cache them. if (options_.has()) { @@ -71,7 +75,10 @@ ObjectDescriptorImpl::ObjectDescriptorImpl( // these ranges. it->active_ranges.emplace(dr.read_id, range); // Cache them so subsequent `Read()` calls can claim them. - prewarmed_ranges_.emplace(range_key, PrewarmedRange{range, dr.read_id}); + auto [cache_it, inserted] = prewarmed_ranges_.emplace( + range_key, PrewarmedRange{range, dr.read_id}); + // Mark them as unclaimed for pacing checks and store the iterator. + unclaimed_ranges_.emplace(dr.read_id, UnclaimedRangeState{0, cache_it}); } // Ensure new dynamically requested ranges use IDs that don't conflict // with pre-warmed ones. @@ -228,6 +235,13 @@ std::unique_ptr ObjectDescriptorImpl::Read( // 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); + // Mark as claimed so we stop applying pacing constraints to it, and + // reclaim the buffer budget used by this range atomically under the lock. + auto unclaimed_it = unclaimed_ranges_.find(prewarmed.read_id); + if (unclaimed_it != unclaimed_ranges_.end()) { + total_prewarmed_bytes_buffered_ -= unclaimed_it->second.bytes_buffered; + unclaimed_ranges_.erase(unclaimed_it); + } lk.unlock(); if (!internal::TracingEnabled(options_)) { return std::unique_ptr( @@ -418,32 +432,69 @@ void ObjectDescriptorImpl::OnRead( // Release the lock while notifying the ranges. The notifications may trigger // application code, and that code may callback on this class. lk.unlock(); - #if defined(GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS) || \ defined(GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY) auto t5_stamp = std::chrono::steady_clock::now(); #endif + auto apply_pacing_and_check_eviction = [this](std::int64_t id, + std::size_t chunk_size, + StreamIterator it) { + auto unclaimed_it = unclaimed_ranges_.find(id); + if (unclaimed_it == unclaimed_ranges_.end()) return false; + + if (total_prewarmed_bytes_buffered_ + chunk_size > + max_prewarmed_buffer_size_) { + // Evict the range if it exceeds the pacing limit. + total_prewarmed_bytes_buffered_ -= unclaimed_it->second.bytes_buffered; + prewarmed_ranges_.erase(unclaimed_it->second.cache_it); + unclaimed_ranges_.erase(unclaimed_it); + + // Erasing from active_ranges ensures we ignore any subsequent GCS chunks + // for this range. + it->active_ranges.erase(id); + return true; + } + + // Track buffered data size for pacing. + unclaimed_it->second.bytes_buffered += chunk_size; + total_prewarmed_bytes_buffered_ += chunk_size; + return false; + }; + for (auto& range_data : *response->mutable_object_data_ranges()) { auto id = range_data.read_range().read_id(); auto const l = copy.find(id); if (l == copy.end()) continue; auto range = l->second; - bool active = false; + auto chunk_size = range_data.checksummed_data().content().size(); + + bool evict = false; lk.lock(); - // Verify the range is still active on this stream. It might have been - // evicted or cancelled during the processing of this batch. - active = it->active_ranges.count(id) != 0; + // Verify the range is still active under the lock. Because `OnRead` + // processes chunks in batches, an earlier chunk in the same batch could + // breach the pacing limit and evict a subsequent chunk's range. + bool active = it->active_ranges.count(id) != 0; + if (active) { + evict = apply_pacing_and_check_eviction(id, chunk_size, it); + } lk.unlock(); if (active) { + if (evict) { + // Complete the evicted range with an error. + range->OnFinish(Status(StatusCode::kResourceExhausted, + "Evicted pre-warmed range due to pacing limit")); + } else { #if defined(GOOGLE_CLOUD_CPP_STORAGE_WITH_OTEL_METRICS) || \ defined(GOOGLE_CLOUD_CPP_HAVE_OPENTELEMETRY) - range->SetT5(t5_stamp); + range->SetT5(t5_stamp); #endif - // TODO(#15104) - Consider returning if the range is done, and then - // skipping CleanupDoneRanges(). - range->OnRead(std::move(range_data), is_transcoded, object_size); + // Deliver data to the range. + // TODO(#15104) - Consider returning if the range is done, and then + // skipping CleanupDoneRanges(). + range->OnRead(std::move(range_data), is_transcoded, object_size); + } } } lk.lock(); diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.h b/google/cloud/storage/internal/async/object_descriptor_impl.h index 023a0084889b2..0b982b9ca215e 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.h +++ b/google/cloud/storage/internal/async/object_descriptor_impl.h @@ -132,8 +132,23 @@ class ObjectDescriptorImpl std::int64_t read_id; }; // Cache of pre-warmed ranges, keyed by (offset, length). - std::map, PrewarmedRange> - prewarmed_ranges_; + using PrewarmedRangesMap = + std::map, PrewarmedRange>; + PrewarmedRangesMap prewarmed_ranges_; + + struct UnclaimedRangeState { + std::size_t bytes_buffered; + PrewarmedRangesMap::iterator cache_it; + }; + + // Map of read_id to unclaimed range state (bytes buffered and original key). + std::unordered_map unclaimed_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 + // before we start evicting them. + std::size_t max_prewarmed_buffer_size_ = + 5 * 1024 * 1024; // Default 5 MiB pacing limit Options options_; std::unique_ptr stream_manager_; diff --git a/google/cloud/storage/internal/async/object_descriptor_impl_test.cc b/google/cloud/storage/internal/async/object_descriptor_impl_test.cc index e29f15e53bedc..d36181bfcf748 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl_test.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl_test.cc @@ -2490,6 +2490,99 @@ TEST(ObjectDescriptorImpl, PrewarmedCacheMiss) { next.first.set_value(true); } +TEST(ObjectDescriptorImpl, PrewarmedPacingEviction) { + auto constexpr kResponse0 = R"pb( + metadata { + bucket: "projects/_/buckets/test-bucket" + name: "test-object" + generation: 42 + } + read_handle { handle: "handle-12345" } + )pb"; + + auto constexpr kResponse1 = R"pb( + read_handle { handle: "handle-23456" } + object_data_ranges { + range_end: false + read_range { read_id: 1 read_offset: 0 } + checksummed_data { content: "123456" } + } + )pb"; + + auto constexpr kExpectedRequest = R"pb( + read_ranges { read_id: 2 read_offset: 0 read_length: 10 } + )pb"; + + AsyncSequencer sequencer; + auto stream = std::make_unique(); + EXPECT_CALL(*stream, Write) + .WillOnce([&](Request const& request, grpc::WriteOptions) { + auto expected = Request{}; + EXPECT_TRUE(TextFormat::ParseFromString(kExpectedRequest, &expected)); + EXPECT_THAT(request, IsProtoEqual(expected)); + return sequencer.PushBack("Write[1]").then([](auto f) { + return f.get(); + }); + }); + + EXPECT_CALL(*stream, Read) + .WillOnce([=, &sequencer]() { + return sequencer.PushBack("Read[1]").then([&](auto) { + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); + return absl::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return absl::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + EXPECT_CALL(*stream, Cancel).Times(AtMost(1)); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 10}}); + options.set(5); + + auto tested = std::make_shared( + NoResume(), factory.AsStdFunction(), + google::storage::v2::BidiReadObjectSpec{}, + std::make_shared(std::move(stream)), options); + + auto response = Response{}; + EXPECT_TRUE(TextFormat::ParseFromString(kResponse0, &response)); + tested->Start(std::move(response)); + + auto read1 = sequencer.PopFrontWithName(); + EXPECT_EQ(read1.second, "Read[1]"); + read1.first.set_value(true); + + auto next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Read[2]"); + + auto s1 = tested->Read({0, 10}); + ASSERT_THAT(s1, NotNull()); + + auto write_next = sequencer.PopFrontWithName(); + EXPECT_EQ(write_next.second, "Write[1]"); + write_next.first.set_value(true); + + next.first.set_value(true); + + next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Finish"); + next.first.set_value(true); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal diff --git a/google/cloud/storage/internal/async/options.h b/google/cloud/storage/internal/async/options.h index f6fe7fa70593c..400337118443c 100644 --- a/google/cloud/storage/internal/async/options.h +++ b/google/cloud/storage/internal/async/options.h @@ -25,14 +25,14 @@ namespace cloud { namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN -// Configuration for a single read range to be pre-warmed. +/// Configuration for a single read range to be pre-warmed. struct ReadRangeConfig { std::int64_t offset; std::int64_t length; }; -// Internal option to pass the list of ranges to pre-warm from `AsyncClient` to -// `Connection`. +/// Internal option to pass the list of ranges to pre-warm from `AsyncClient` to +/// `Connection`. struct ReadRangesOption { using Type = std::vector; static char const* name() { @@ -40,6 +40,16 @@ struct ReadRangesOption { } }; +/// Internal option to override the pacing buffer limit (in bytes) for +/// pre-warmed ranges. Primarily used in unit tests to test pacing limits +/// without generating large payloads. +struct PreWarmBufferLimitOption { + using Type = std::size_t; + static char const* name() { + return "google::cloud::storage::PreWarmBufferLimitOption"; + } +}; + GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal } // namespace cloud diff --git a/google/cloud/storage/internal/async/read_range.h b/google/cloud/storage/internal/async/read_range.h index 609652ae78f43..16a747fb2aaa1 100644 --- a/google/cloud/storage/internal/async/read_range.h +++ b/google/cloud/storage/internal/async/read_range.h @@ -40,9 +40,11 @@ struct DedupedReadRange { std::int64_t read_id; }; -// Deduplicates a list of ReadRangeConfigs and assigns a sequential -// read_id to each unique range, starting with initial_id + 1. -// Returns a vector mapping each assigned read_id to its corresponding range. +/** + * Deduplicates a list of ReadRangeConfigs and assigns a sequential + * read_id to each unique range, starting with initial_id + 1. + * Returns a vector mapping each assigned read_id to its corresponding range. + */ std::vector DeduplicateRanges( std::vector const& ranges, std::int64_t initial_id = 0); diff --git a/google/cloud/storage/internal/async/read_range_test.cc b/google/cloud/storage/internal/async/read_range_test.cc index 1f327dc56db34..b0ebbc8a56ae2 100644 --- a/google/cloud/storage/internal/async/read_range_test.cc +++ b/google/cloud/storage/internal/async/read_range_test.cc @@ -591,7 +591,7 @@ TEST(ReadRange, NoResumeIfRequestExceeded) { EXPECT_FALSE(resume.has_value()); } -TEST(ReadRangeDeduplicationTest, Basic) { +TEST(ReadRange, DeduplicateRangesBasic) { std::vector ranges = { {0, 10}, {10, 10}, {0, 10}, // Duplicate {20, 10}, {10, 10}, // Duplicate @@ -612,7 +612,7 @@ TEST(ReadRangeDeduplicationTest, Basic) { EXPECT_EQ(deduped[2].read_id, 3); } -TEST(ReadRangeDeduplicationTest, InitialId) { +TEST(ReadRange, DeduplicateRangesInitialId) { std::vector ranges = { {0, 10}, }; From 75869a754c927a8d39bd1b52ffaba7747796cbb1 Mon Sep 17 00:00:00 2001 From: kalragauri Date: Mon, 10 Aug 2026 16:00:04 +0530 Subject: [PATCH 3/4] feat(storage): add telemetry for pre-warmed ranges in ObjectDescriptorImpl (#16323) --- .../internal/async/connection_tracing.cc | 6 ++ .../internal/async/connection_tracing_test.cc | 35 ++++++++++ .../internal/async/object_descriptor_impl.cc | 50 +++++++++++++- .../internal/async/object_descriptor_impl.h | 3 + .../async/object_descriptor_reader_tracing.cc | 16 +++-- .../async/object_descriptor_reader_tracing.h | 3 +- .../object_descriptor_reader_tracing_test.cc | 67 +++++++++++++++---- 7 files changed, 159 insertions(+), 21 deletions(-) diff --git a/google/cloud/storage/internal/async/connection_tracing.cc b/google/cloud/storage/internal/async/connection_tracing.cc index 8f8ac51436c67..9b8669bfb20e0 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" @@ -55,6 +56,11 @@ class AsyncConnectionTracing : public storage::AsyncConnection { future>> Open( OpenParams p) override { auto span = internal::MakeSpan("storage::AsyncConnection::Open"); + if (p.options.has()) { + auto const& ranges = p.options.get(); + span->SetAttribute("gl-cpp.initial-read-ranges.ranges-count", + ranges.size()); + } internal::OTelScope scope(span); auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(), bucket = p.read_spec.bucket(), span = std::move(span)](auto f) diff --git a/google/cloud/storage/internal/async/connection_tracing_test.cc b/google/cloud/storage/internal/async/connection_tracing_test.cc index bdc562baef97c..992207c698155 100644 --- a/google/cloud/storage/internal/async/connection_tracing_test.cc +++ b/google/cloud/storage/internal/async/connection_tracing_test.cc @@ -15,6 +15,7 @@ #include "google/cloud/storage/internal/async/connection_tracing.h" #include "google/cloud/storage/async/object_descriptor_connection.h" #include "google/cloud/storage/async/reader_connection.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/mocks/mock_async_connection.h" #include "google/cloud/storage/mocks/mock_async_object_descriptor_connection.h" #include "google/cloud/storage/mocks/mock_async_reader_connection.h" @@ -44,8 +45,10 @@ using ::google::cloud::testing_util::EventNamed; using ::google::cloud::testing_util::InstallSpanCatcher; using ::google::cloud::testing_util::IsOk; using ::google::cloud::testing_util::IsOkAndHolds; +using ::google::cloud::testing_util::OTelAttribute; using ::google::cloud::testing_util::OTelContextCaptured; using ::google::cloud::testing_util::PromiseWithOTelContext; +using ::google::cloud::testing_util::SpanHasAttributes; using ::google::cloud::testing_util::SpanHasEvents; using ::google::cloud::testing_util::SpanHasInstrumentationScope; using ::google::cloud::testing_util::SpanKindIsClient; @@ -589,6 +592,38 @@ TEST(ConnectionTracing, OpenSuccess) { SpanHasInstrumentationScope(), SpanKindIsClient()))); } +TEST(ConnectionTracing, OpenSuccessWithInitialReadRanges) { + auto span_catcher = InstallSpanCatcher(); + PromiseWithOTelContext< + StatusOr>> + p; + auto mock = std::make_unique(); + EXPECT_CALL(*mock, options).WillOnce(Return(TracingEnabled())); + EXPECT_CALL(*mock, Open).WillOnce(expect_context(p)); + + auto actual = MakeTracingAsyncConnection(std::move(mock)); + auto open_params = AsyncConnection::OpenParams{}; + open_params.options.set({{0, 100}, {1000, 200}}); + auto f = actual->Open(std::move(open_params)).then(expect_no_context); + + auto mock_descriptor = + std::make_shared(); + p.set_value(StatusOr>( + std::move(mock_descriptor))); + auto result = f.get(); + ASSERT_STATUS_OK(result); + auto descriptor = *std::move(result); + descriptor.reset(); + + auto spans = span_catcher->GetSpans(); + EXPECT_THAT( + spans, ElementsAre( + AllOf(SpanNamed("storage::AsyncConnection::Open"), + SpanHasAttributes(OTelAttribute( + "gl-cpp.initial-read-ranges.ranges-count", 2)), + SpanWithStatus(opentelemetry::trace::StatusCode::kOk)))); +} + TEST(ConnectionTracing, StartAppendableObjectUploadSuccess) { auto span_catcher = InstallSpanCatcher(); PromiseWithOTelContext< diff --git a/google/cloud/storage/internal/async/object_descriptor_impl.cc b/google/cloud/storage/internal/async/object_descriptor_impl.cc index ecba87a538c16..66919442737fe 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl.cc @@ -40,6 +40,31 @@ namespace cloud { namespace storage_internal { GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN +namespace { + +enum class InitialReadRangesCacheStatus { + kNone, + kMiss, + kHit, + kEvicted, +}; + +absl::string_view CacheStatusToString(InitialReadRangesCacheStatus status) { + switch (status) { + case InitialReadRangesCacheStatus::kHit: + return "HIT"; + case InitialReadRangesCacheStatus::kEvicted: + return "EVICTED"; + case InitialReadRangesCacheStatus::kMiss: + return "MISS"; + case InitialReadRangesCacheStatus::kNone: + return ""; + } + return "INVALID_STATUS"; +} + +} // namespace + ObjectDescriptorImpl::ObjectDescriptorImpl( std::unique_ptr resume_policy, OpenStreamFactory make_stream, @@ -50,6 +75,7 @@ ObjectDescriptorImpl::ObjectDescriptorImpl( make_stream_(std::move(make_stream)), read_object_spec_(std::move(read_object_spec)), options_(std::move(options)), + has_initial_read_ranges_(options_.has()), transport_ok_(std::move(transport_ok)) { stream_manager_ = std::make_unique( []() -> std::shared_ptr { return nullptr; }, // NOLINT @@ -231,7 +257,12 @@ 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); + auto cache_status = has_initial_read_ranges_ + ? InitialReadRangesCacheStatus::kMiss + : InitialReadRangesCacheStatus::kNone; + if (cache_it != prewarmed_ranges_.end()) { + cache_status = InitialReadRangesCacheStatus::kHit; // 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); @@ -248,7 +279,13 @@ std::unique_ptr ObjectDescriptorImpl::Read( std::make_unique(std::move(prewarmed.range))); } return MakeTracingObjectDescriptorReader(std::move(prewarmed.range), - read_object_spec_.bucket()); + read_object_spec_.bucket(), + CacheStatusToString(cache_status)); + } + + // If not hit, check if it was evicted earlier due to pacing. + if (evicted_ranges_.erase(cache_key) != 0) { + cache_status = InitialReadRangesCacheStatus::kEvicted; } if (stream_manager_->Empty()) { @@ -260,7 +297,8 @@ std::unique_ptr ObjectDescriptorImpl::Read( std::make_unique(std::move(range))); } return MakeTracingObjectDescriptorReader(std::move(range), - read_object_spec_.bucket()); + read_object_spec_.bucket(), + CacheStatusToString(cache_status)); } auto it = stream_manager_->GetLeastBusyStream(); @@ -277,7 +315,8 @@ std::unique_ptr ObjectDescriptorImpl::Read( std::make_unique(std::move(range))); } return MakeTracingObjectDescriptorReader(std::move(range), - read_object_spec_.bucket()); + read_object_spec_.bucket(), + CacheStatusToString(cache_status)); } std::shared_ptr @@ -447,6 +486,11 @@ 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; + // Cap tombstone set size to prevent unbounded memory growth in long-lived + // descriptors where pre-warmed ranges are evicted but never requested. + 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_impl.h b/google/cloud/storage/internal/async/object_descriptor_impl.h index 0b982b9ca215e..d0cfa98487c05 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 @@ -157,6 +159,7 @@ class ObjectDescriptorImpl google::cloud::StatusOr> pending_stream_; bool cancelled_ = false; + bool has_initial_read_ranges_ = false; std::function transport_ok_; }; 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 9719914f21ce7..34ca76d9e3eba 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.cc @@ -34,14 +34,20 @@ namespace sc = ::opentelemetry::semconv; class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { public: explicit ObjectDescriptorReaderTracing(std::shared_ptr impl, - std::string bucket_name) + std::string bucket_name, + absl::string_view cache_status) : ObjectDescriptorReader(std::move(impl)), - bucket_name_(std::move(bucket_name)) {} + bucket_name_(std::move(bucket_name)), + cache_status_(cache_status) {} ~ObjectDescriptorReaderTracing() override = default; future Read() override { auto span = internal::MakeSpan("storage::AsyncConnection::ReadRange"); + if (!cache_status_.empty()) { + span->SetAttribute("gl-cpp.initial-read-ranges.cache-status", + std::string(cache_status_)); + } internal::OTelScope scope(span); return ObjectDescriptorReader::Read() .then([span = std::move(span), bucket_name = bucket_name_, @@ -78,6 +84,7 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { private: std::string bucket_name_; + absl::string_view cache_status_; ReaderConnectionTelemetry metrics_; }; @@ -85,9 +92,10 @@ class ObjectDescriptorReaderTracing : public ObjectDescriptorReader { std::unique_ptr MakeTracingObjectDescriptorReader(std::shared_ptr impl, - std::string bucket_name) { + std::string bucket_name, + absl::string_view cache_status) { return std::make_unique( - std::move(impl), std::move(bucket_name)); + std::move(impl), std::move(bucket_name), 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 3282a42c453e3..2879c66343d90 100644 --- a/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h +++ b/google/cloud/storage/internal/async/object_descriptor_reader_tracing.h @@ -26,7 +26,8 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN std::unique_ptr MakeTracingObjectDescriptorReader(std::shared_ptr impl, - std::string bucket_name); + std::string bucket_name, + 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 04837f3212ffd..acec73a4f636a 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, "test-bucket"); + auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket", "TEST"); auto data = google::storage::v2::ObjectRangeData{}; auto constexpr kData0 = R"pb( @@ -61,6 +61,8 @@ TEST(ObjectDescriptorReaderTracing, Read) { EXPECT_THAT(spans, ElementsAre(AllOf( SpanNamed("storage::AsyncConnection::ReadRange"), + SpanHasAttributes(OTelAttribute( + "gl-cpp.initial-read-ranges.cache-status", "TEST")), SpanHasEvents(AllOf( EventNamed("gl-cpp.read-range"), SpanEventAttributesAre( @@ -73,23 +75,62 @@ TEST(ObjectDescriptorReaderTracing, Read) { TEST(ObjectDescriptorReaderTracing, ReadError) { auto span_catcher = InstallSpanCatcher(); auto impl = std::make_shared(10000, 30); - auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket"); + auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket", "TEST"); impl->OnFinish(PermanentError()); auto actual = reader->Read().get(); auto spans = span_catcher->GetSpans(); - EXPECT_THAT(spans, - ElementsAre(AllOf( - SpanNamed("storage::AsyncConnection::ReadRange"), - SpanHasAttributes(OTelAttribute( - "gl-cpp.status_code", "NOT_FOUND")), - SpanHasEvents(AllOf( - EventNamed("gl-cpp.read-range"), - SpanEventAttributesAre( - OTelAttribute(sc::thread::kThreadId, _), - OTelAttribute("rpc.message.type", - "RECEIVED"))))))); + EXPECT_THAT( + spans, + ElementsAre(AllOf( + SpanNamed("storage::AsyncConnection::ReadRange"), + SpanHasAttributes( + OTelAttribute("gl-cpp.status_code", "NOT_FOUND"), + OTelAttribute( + "gl-cpp.initial-read-ranges.cache-status", "TEST")), + SpanHasEvents( + AllOf(EventNamed("gl-cpp.read-range"), + SpanEventAttributesAre( + OTelAttribute(sc::thread::kThreadId, _), + OTelAttribute("rpc.message.type", + "RECEIVED"))))))); +} + +TEST(ObjectDescriptorReaderTracing, ReadWithoutInitialReadRanges) { + auto span_catcher = InstallSpanCatcher(); + auto impl = std::make_shared(10000, 30); + // Pass empty string for cache_status when initial read ranges were not + // configured. + auto reader = MakeTracingObjectDescriptorReader(impl, "test-bucket", ""); + + impl->OnFinish(PermanentError()); + + auto actual = reader->Read().get(); + auto spans = span_catcher->GetSpans(); + ASSERT_EQ(spans.size(), 1); + auto const& attributes = spans[0]->GetAttributes(); + EXPECT_EQ(attributes.find("gl-cpp.initial-read-ranges.cache-status"), + attributes.end()); +} + +TEST(ObjectDescriptorReaderTracing, ReadWithCacheStatuses) { + for (auto const* status : {"HIT", "MISS", "EVICTED"}) { + auto span_catcher = InstallSpanCatcher(); + auto impl = std::make_shared(10000, 30); + auto reader = + MakeTracingObjectDescriptorReader(impl, "test-bucket", status); + + impl->OnFinish(PermanentError()); + + auto actual = reader->Read().get(); + auto spans = span_catcher->GetSpans(); + EXPECT_THAT(spans, + ElementsAre(AllOf( + SpanNamed("storage::AsyncConnection::ReadRange"), + SpanHasAttributes(OTelAttribute( + "gl-cpp.initial-read-ranges.cache-status", status))))); + } } } // namespace From 6b3b4306aa6bdb4478e11160e85d1c0f6f2720c3 Mon Sep 17 00:00:00 2001 From: kalragauri Date: Wed, 12 Aug 2026 16:33:13 +0530 Subject: [PATCH 4/4] feat(storage): expose initial read ranges in AsyncClient (#16341) --- google/cloud/storage/async/client.cc | 22 +++++++ google/cloud/storage/async/client.h | 44 ++++++++++++++ google/cloud/storage/async/client_test.cc | 38 ++++++++++++ .../storage/examples/storage_async_samples.cc | 60 +++++++++++++++++++ 4 files changed, 164 insertions(+) diff --git a/google/cloud/storage/async/client.cc b/google/cloud/storage/async/client.cc index 6b01986317b2b..0ed676913569d 100644 --- a/google/cloud/storage/async/client.cc +++ b/google/cloud/storage/async/client.cc @@ -16,6 +16,7 @@ #include "google/cloud/storage/internal/async/connection_impl.h" #include "google/cloud/storage/internal/async/connection_tracing.h" #include "google/cloud/storage/internal/async/default_options.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/internal/grpc/stub.h" #include "google/cloud/grpc_options.h" #include @@ -84,6 +85,27 @@ future> AsyncClient::Open( }); } +future> AsyncClient::Open( + BucketName const& bucket_name, std::string object_name, + InitialReadRanges const& config, Options opts) { + auto spec = google::storage::v2::BidiReadObjectSpec{}; + spec.set_bucket(bucket_name.FullName()); + spec.set_object(std::move(object_name)); + + // Convert the user-facing `InitialReadRanges` to the internal + // `ReadRangesOption` so it can be propagated down to the connection + // implementation. + if (!config.initial_ranges.empty()) { + std::vector internal_ranges; + internal_ranges.reserve(config.initial_ranges.size()); + for (auto const& r : config.initial_ranges) { + internal_ranges.push_back({r.offset, r.length}); + } + opts.set(std::move(internal_ranges)); + } + return Open(std::move(spec), std::move(opts)); +} + future>> AsyncClient::ReadObject( BucketName const& bucket_name, std::string object_name, Options opts) { auto request = google::storage::v2::ReadObjectRequest{}; diff --git a/google/cloud/storage/async/client.h b/google/cloud/storage/async/client.h index cfdc78d044b32..edbfc0f25a28e 100644 --- a/google/cloud/storage/async/client.h +++ b/google/cloud/storage/async/client.h @@ -83,6 +83,27 @@ GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN */ class AsyncClient { public: + /** + * Specifies a byte range for a read request. + */ + struct ByteRange { + std::int64_t offset = 0; + std::int64_t length = 0; + }; + + /** + * Specifies initial byte ranges to request concurrently when opening an + * object. + * + * Passing initial read ranges allows the client to begin fetching expected + * byte ranges during connection setup, which may improve first-byte retrieval + * times for known access patterns. + */ + struct InitialReadRanges { + // The initial read ranges. + std::vector initial_ranges; + }; + /// Create a new client configured with @p options. explicit AsyncClient(Options options = {}); /// Create a new client using @p connection. This is often used for mocking. @@ -273,6 +294,29 @@ class AsyncClient { std::string object_name, Options opts = {}); + /** + * Open an object descriptor, requesting specified initial read ranges + * concurrently. + * + * @par Example + * @snippet storage_async_samples.cc open-object-initial-read-ranges + * + * @par Idempotency + * This is a read-only operation and is always idempotent. The operation will + * retry until the descriptor is successfully created. The descriptor itself + * will resume any incomplete ranged reads if the connection(s) are + * interrupted. Use `ResumePolicyOption` and `ResumePolicy` to control this. + * + * @param bucket_name the name of the bucket that contains the object. + * @param object_name the name of the object to be read. + * @param config initial byte ranges to request during connection setup. + * @param opts options controlling the behavior of this RPC. + */ + future> Open(BucketName const& bucket_name, + std::string object_name, + InitialReadRanges const& config, + Options opts = {}); + /** * Open an object descriptor to perform one or more ranged reads. * diff --git a/google/cloud/storage/async/client_test.cc b/google/cloud/storage/async/client_test.cc index 0c369deb0dd21..3367edc54dda5 100644 --- a/google/cloud/storage/async/client_test.cc +++ b/google/cloud/storage/async/client_test.cc @@ -13,6 +13,7 @@ // limitations under the License. #include "google/cloud/storage/async/client.h" +#include "google/cloud/storage/internal/async/options.h" #include "google/cloud/storage/mocks/mock_async_connection.h" #include "google/cloud/storage/mocks/mock_async_object_descriptor_connection.h" #include "google/cloud/storage/mocks/mock_async_reader_connection.h" @@ -290,6 +291,43 @@ TEST(AsyncClient, Open) { "empty response", [](auto const& p) { return p.size(); }, 0))); } +TEST(AsyncClient, OpenWithInitialReadRanges) { + auto constexpr kExpectedRequest = R"pb( + bucket: "projects/_/buckets/test-bucket" + object: "test-object" + )pb"; + auto mock = std::make_shared(); + EXPECT_CALL(*mock, options).WillRepeatedly(Return(Options{})); + + EXPECT_CALL(*mock, Open).WillOnce([&](AsyncConnection::OpenParams const& p) { + EXPECT_TRUE(p.options.has()); + auto const& ranges = p.options.get(); + EXPECT_EQ(ranges.size(), 2); + if (ranges.size() >= 2) { + EXPECT_EQ(ranges[0].offset, 0); + EXPECT_EQ(ranges[0].length, 100); + EXPECT_EQ(ranges[1].offset, 1000); + EXPECT_EQ(ranges[1].length, 200); + } + + auto expected = google::storage::v2::BidiReadObjectSpec{}; + EXPECT_TRUE(TextFormat::ParseFromString(kExpectedRequest, &expected)); + EXPECT_THAT(p.read_spec, IsProtoEqual(expected)); + + auto descriptor = std::make_shared(); + return make_ready_future(make_status_or( + std::shared_ptr(std::move(descriptor)))); + }); + + auto client = AsyncClient(mock); + AsyncClient::InitialReadRanges config; + config.initial_ranges = {{0, 100}, {1000, 200}}; + auto descriptor = + client.Open(BucketName("test-bucket"), "test-object", std::move(config)) + .get(); + ASSERT_STATUS_OK(descriptor); +} + TEST(AsyncClient, OpenWithInvalidBucket) { auto constexpr kExpectedRequest = R"pb( bucket: "test-only-invalid" diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index 902855efcdfe4..6fa19dd3db50d 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -246,6 +246,55 @@ void OpenObjectMultipleRangedRead(google::cloud::storage::AsyncClient& client, std::cout << "The ranges contain " << count << " newlines\n"; } +void OpenObjectWithInitialReadRanges( + google::cloud::storage::AsyncClient& client, + std::vector const& argv) { + //! [open-object-initial-read-ranges] + // [START storage_open_object_initial_read_ranges] + namespace gcs = google::cloud::storage; + + // Helper coroutine to count newlines returned by an AsyncReader. + auto count_newlines = + [](gcs::AsyncReader reader, + gcs::AsyncToken token) -> google::cloud::future { + std::uint64_t count = 0; + while (token.valid()) { + auto [payload, t] = (co_await reader.Read(std::move(token))).value(); + token = std::move(t); + for (auto const& buffer : payload.contents()) { + count += std::count(buffer.begin(), buffer.end(), '\n'); + } + } + co_return count; + }; + + auto coro = + [&count_newlines]( + gcs::AsyncClient& client, std::string bucket_name, + std::string object_name) -> google::cloud::future { + gcs::AsyncClient::InitialReadRanges config; + config.initial_ranges = {{0, 1024}, {1024, 1024}}; + + auto descriptor = + (co_await client.Open(gcs::BucketName(std::move(bucket_name)), + std::move(object_name), std::move(config))) + .value(); + + auto [r1, t1] = descriptor.Read(0, 1024); + auto [r2, t2] = descriptor.Read(1024, 1024); + + auto c1 = count_newlines(std::move(r1), std::move(t1)); + auto c2 = count_newlines(std::move(r2), std::move(t2)); + co_return (co_await std::move(c1)) + (co_await std::move(c2)); + }; + // [END storage_open_object_initial_read_ranges] + //! [open-object-initial-read-ranges] + // The example is easier to test and run if we call the coroutine and block + // until it completes. + auto const count = coro(client, argv.at(0), argv.at(1)).get(); + std::cout << "The pre-warmed ranges contain " << count << " newlines\n"; +} + void OpenObjectReadFullObject(google::cloud::storage::AsyncClient& client, std::vector const& argv) { //! [open-object-read-full-object] @@ -997,6 +1046,11 @@ void OpenObjectMultipleRangedRead(google::cloud::storage::AsyncClient&, std::cerr << "AsyncClient::Open() example requires coroutines\n"; } +void OpenObjectWithInitialReadRanges(google::cloud::storage::AsyncClient&, + std::vector const&) { + std::cerr << "AsyncClient::Open() example requires coroutines\n"; +} + void OpenMultipleObjectsRangedRead(google::cloud::storage::AsyncClient&, std::vector const&) { std::cerr << "AsyncClient::Open() example requires coroutines\n"; @@ -1306,6 +1360,10 @@ void AutoRun(std::vector const& argv) { << std::endl; OpenObjectMultipleRangedRead(client, {bucket_name, composed_name}); + std::cout << "Running the OpenObjectWithInitialReadRanges() example" + << std::endl; + OpenObjectWithInitialReadRanges(client, {bucket_name, composed_name}); + std::cout << "Running the OpenMultipleObjectsRangedRead() example" << std::endl; auto const multi_read_o1 = @@ -1563,6 +1621,8 @@ int main(int argc, char* argv[]) try { OpenObjectSingleRangedRead), make_entry("open-object-multiple-ranged-read", {}, OpenObjectMultipleRangedRead), + make_entry("open-object-initial-read-ranges", {}, + OpenObjectWithInitialReadRanges), make_entry("open-object-read-full-object", {}, OpenObjectReadFullObject), make_entry("open-multiple-objects-ranged-read", {"", "", ""},