From aa728e391cbaf696163fbc5ed0c65e4fc59b1e3a Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Tue, 8 Sep 2026 08:51:34 +0000 Subject: [PATCH 1/2] fix(storage): support download stall timeout in async BiDi reads --- .../storage/internal/async/open_object.cc | 15 +++- .../internal/async/open_object_test.cc | 82 +++++++++++++++++++ 2 files changed, 96 insertions(+), 1 deletion(-) diff --git a/google/cloud/storage/internal/async/open_object.cc b/google/cloud/storage/internal/async/open_object.cc index 1f3b87d805e3f..0687f2a3abf6a 100644 --- a/google/cloud/storage/internal/async/open_object.cc +++ b/google/cloud/storage/internal/async/open_object.cc @@ -13,6 +13,9 @@ // limitations under the License. #include "google/cloud/storage/internal/async/open_object.h" +#include "google/cloud/storage/internal/grpc/scale_stall_timeout.h" +#include "google/cloud/storage/options.h" +#include "google/cloud/internal/async_read_write_stream_timeout.h" #include "google/cloud/internal/make_status.h" #include "absl/strings/str_cat.h" #include @@ -59,7 +62,17 @@ std::unique_ptr OpenObject::CreateRpc( google::storage::v2::BidiReadObjectRequest const& request) { auto p = RequestParams(request); if (!p.empty()) context->AddMetadata("x-goog-request-params", std::move(p)); - return stub.AsyncBidiReadObject(cq, std::move(context), std::move(options)); + auto timeout = ScaleStallTimeout( + options->get(), + options->get(), + google::storage::v2::ServiceConstants::MAX_READ_CHUNK_BYTES); + auto rpc = + stub.AsyncBidiReadObject(cq, std::move(context), std::move(options)); + return std::make_unique< + google::cloud::internal::AsyncStreamingReadWriteRpcTimeout< + google::storage::v2::BidiReadObjectRequest, + google::storage::v2::BidiReadObjectResponse>>( + cq, timeout, timeout, timeout, std::move(rpc)); } void OpenObject::OnStart(bool ok) { diff --git a/google/cloud/storage/internal/async/open_object_test.cc b/google/cloud/storage/internal/async/open_object_test.cc index aabeafef634be..cade0f86e418c 100644 --- a/google/cloud/storage/internal/async/open_object_test.cc +++ b/google/cloud/storage/internal/async/open_object_test.cc @@ -14,14 +14,17 @@ #include "google/cloud/storage/internal/async/open_object.h" #include "google/cloud/mocks/mock_async_streaming_read_write_rpc.h" +#include "google/cloud/storage/options.h" #include "google/cloud/storage/testing/canonical_errors.h" #include "google/cloud/storage/testing/mock_storage_stub.h" #include "google/cloud/testing_util/async_sequencer.h" #include "google/cloud/testing_util/is_proto_equal.h" +#include "google/cloud/testing_util/mock_completion_queue_impl.h" #include "google/cloud/testing_util/status_matchers.h" #include "google/cloud/testing_util/validate_metadata.h" #include #include +#include #include namespace google { @@ -35,6 +38,7 @@ using ::google::cloud::storage::testing::canonical_errors::PermanentError; using ::google::cloud::testing_util::AsyncSequencer; using ::google::cloud::testing_util::IsOkAndHolds; using ::google::cloud::testing_util::IsProtoEqual; +using ::google::cloud::testing_util::MockCompletionQueueImpl; using ::google::cloud::testing_util::StatusIs; using ::google::protobuf::TextFormat; using ::testing::AllOf; @@ -47,6 +51,16 @@ using MockStream = google::cloud::mocks::MockAsyncStreamingReadWriteRpc< google::storage::v2::BidiReadObjectRequest, google::storage::v2::BidiReadObjectResponse>; +StatusOr CancelledTimer() { + return internal::CancelledError("test-only", GCP_ERROR_INFO()); +} + +StatusOr MakeTimerStatus( + future f) { + if (!f.get()) return CancelledTimer(); + return make_status_or(std::chrono::system_clock::now()); +} + TEST(OpenImpl, RequestParams) { auto constexpr kPlain = R"pb( read_object_spec { @@ -390,6 +404,74 @@ TEST(OpenImpl, UnexpectedFinish) { EXPECT_THAT(response, StatusIs(StatusCode::kInternal)); } +TEST(OpenImpl, TimeoutCancellation) { + AsyncSequencer sequencer; + MockStorageStub mock; + EXPECT_CALL(mock, AsyncBidiReadObject).WillOnce([&sequencer]() { + 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]() { + return sequencer.PushBack("Write").then([](auto f) { return f.get(); }); + }); + EXPECT_CALL(*stream, Read).WillOnce([&sequencer]() { + return sequencer.PushBack("Read").then([](auto) { + return std::optional(); + }); + }); + EXPECT_CALL(*stream, Cancel).Times(testing::AtLeast(1)); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then([](auto) { + return Status(StatusCode::kCancelled, "Stream timeout"); + }); + }); + return std::unique_ptr(std::move(stream)); + }); + + auto mock_cq = std::make_shared(); + EXPECT_CALL(*mock_cq, MakeRelativeTimer).WillRepeatedly([&sequencer](auto) { + return sequencer.PushBack("MakeRelativeTimer").then(MakeTimerStatus); + }); + + CompletionQueue cq(mock_cq); + Options options; + options.set(std::chrono::seconds(1)); + auto coro = std::make_shared( + mock, cq, std::make_shared(), + internal::MakeImmutableOptions(std::move(options)), + google::storage::v2::BidiReadObjectRequest{}); + auto pending = coro->Call(); + + auto timer1 = sequencer.PopFrontWithName(); + EXPECT_EQ(timer1.second, "MakeRelativeTimer"); + auto start = sequencer.PopFrontWithName(); + EXPECT_EQ(start.second, "Start"); + start.first.set_value(true); + timer1.first.set_value(false); + + auto timer2 = sequencer.PopFrontWithName(); + EXPECT_EQ(timer2.second, "MakeRelativeTimer"); + auto write = sequencer.PopFrontWithName(); + EXPECT_EQ(write.second, "Write"); + write.first.set_value(true); + timer2.first.set_value(false); + + auto timer3 = sequencer.PopFrontWithName(); + EXPECT_EQ(timer3.second, "MakeRelativeTimer"); + auto read = sequencer.PopFrontWithName(); + EXPECT_EQ(read.second, "Read"); + timer3.first.set_value(true); + read.first.set_value(true); + + auto finish = sequencer.PopFrontWithName(); + EXPECT_EQ(finish.second, "Finish"); + finish.first.set_value(true); + + auto response = pending.get(); + EXPECT_THAT(response, StatusIs(StatusCode::kCancelled)); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage_internal From 5316c9c8d678ae851b3ce78841f2775dde9bbc29 Mon Sep 17 00:00:00 2001 From: Gauri Kalra Date: Tue, 8 Sep 2026 09:22:44 +0000 Subject: [PATCH 2/2] Address feedback from Gemini code assistant --- google/cloud/storage/internal/async/open_object.cc | 4 ++-- google/cloud/storage/internal/async/open_object_test.cc | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/google/cloud/storage/internal/async/open_object.cc b/google/cloud/storage/internal/async/open_object.cc index 0687f2a3abf6a..148bdfb1cd593 100644 --- a/google/cloud/storage/internal/async/open_object.cc +++ b/google/cloud/storage/internal/async/open_object.cc @@ -62,11 +62,11 @@ std::unique_ptr OpenObject::CreateRpc( google::storage::v2::BidiReadObjectRequest const& request) { auto p = RequestParams(request); if (!p.empty()) context->AddMetadata("x-goog-request-params", std::move(p)); - auto timeout = ScaleStallTimeout( + std::chrono::milliseconds const timeout = ScaleStallTimeout( options->get(), options->get(), google::storage::v2::ServiceConstants::MAX_READ_CHUNK_BYTES); - auto rpc = + std::unique_ptr rpc = stub.AsyncBidiReadObject(cq, std::move(context), std::move(options)); return std::make_unique< google::cloud::internal::AsyncStreamingReadWriteRpcTimeout< diff --git a/google/cloud/storage/internal/async/open_object_test.cc b/google/cloud/storage/internal/async/open_object_test.cc index cade0f86e418c..c1fa3a4dcde02 100644 --- a/google/cloud/storage/internal/async/open_object_test.cc +++ b/google/cloud/storage/internal/async/open_object_test.cc @@ -441,7 +441,7 @@ TEST(OpenImpl, TimeoutCancellation) { mock, cq, std::make_shared(), internal::MakeImmutableOptions(std::move(options)), google::storage::v2::BidiReadObjectRequest{}); - auto pending = coro->Call(); + future> pending = coro->Call(); auto timer1 = sequencer.PopFrontWithName(); EXPECT_EQ(timer1.second, "MakeRelativeTimer"); @@ -468,7 +468,7 @@ TEST(OpenImpl, TimeoutCancellation) { EXPECT_EQ(finish.second, "Finish"); finish.first.set_value(true); - auto response = pending.get(); + StatusOr response = pending.get(); EXPECT_THAT(response, StatusIs(StatusCode::kCancelled)); }