From 11d48d309b96f21a8e7f79349ae7824494f93ade Mon Sep 17 00:00:00 2001 From: Scott Hart Date: Wed, 30 Sep 2026 17:33:32 -0400 Subject: [PATCH] ci(storage): deflake connection_impl_test --- .../storage/internal/connection_impl_test.cc | 62 +++++++++++++------ 1 file changed, 43 insertions(+), 19 deletions(-) diff --git a/google/cloud/storage/internal/connection_impl_test.cc b/google/cloud/storage/internal/connection_impl_test.cc index d9e96bc683b46..f5d61c959f10a 100644 --- a/google/cloud/storage/internal/connection_impl_test.cc +++ b/google/cloud/storage/internal/connection_impl_test.cc @@ -55,6 +55,7 @@ using ::testing::ByMove; using ::testing::ElementsAre; using ::testing::Eq; using ::testing::HasSubstr; +using ::testing::Ne; using ::testing::Not; using ::testing::Property; using ::testing::Return; @@ -852,7 +853,6 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) { // primary attempt. Any other thread is running a hedge. auto primary_thread = std::make_shared(); auto unblock_primary = std::make_shared>(); - auto primary_closed = std::make_shared>(); auto mock = std::make_unique(); EXPECT_CALL(*mock, options).Times(AtLeast(0)); EXPECT_CALL(*mock, ReadObject) @@ -865,24 +865,36 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) { }); return StatusOr>(std::move(source)); }) - .WillRepeatedly([primary_thread, unblock_primary, primary_closed]( - auto&, auto const&, ReadObjectRangeRequest const&) { + .WillOnce([primary_thread, unblock_primary]( + auto&, auto const&, ReadObjectRangeRequest const&) { + EXPECT_THAT(std::this_thread::get_id(), Eq(*primary_thread)); auto source = std::make_unique(); - if (std::this_thread::get_id() == *primary_thread) { - EXPECT_CALL(*source, Read) - .WillOnce([unblock_primary](char* buf, std::size_t) { - unblock_primary->get_future().wait(); - return MakeReadResult("slow", buf); - }); - EXPECT_CALL(*source, Close).WillOnce([primary_closed] { - primary_closed->set_value(); - return make_status_or(HttpResponse{HttpStatusCode::kOk, {}, {}}); - }); - } else { - EXPECT_CALL(*source, Read).WillOnce([](char* buf, std::size_t) { - return MakeReadResult("hedge", buf); - }); - } + EXPECT_CALL(*source, Read) + .WillOnce([unblock_primary](char* buf, std::size_t) { + unblock_primary->get_future().wait(); + return MakeReadResult("slow", buf); + }); + EXPECT_CALL(*source, Close).WillOnce([] { + return make_status_or(HttpResponse{HttpStatusCode::kOk, {}, {}}); + }); + return StatusOr>(std::move(source)); + }) + .WillOnce([primary_thread](auto&, auto const&, + ReadObjectRangeRequest const&) { + EXPECT_THAT(std::this_thread::get_id(), Ne(*primary_thread)); + auto source = std::make_unique(); + EXPECT_CALL(*source, Read).WillOnce([](char* buf, std::size_t) { + return MakeReadResult("hedge", buf); + }); + return StatusOr>(std::move(source)); + }) + .WillOnce([primary_thread](auto&, auto const&, + ReadObjectRangeRequest const&) { + EXPECT_THAT(std::this_thread::get_id(), Eq(*primary_thread)); + auto source = std::make_unique(); + EXPECT_CALL(*source, Read).WillOnce([](char* buf, std::size_t) { + return MakeReadResult("sync", buf); + }); return StatusOr>(std::move(source)); }); @@ -916,9 +928,21 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) { StatusOr result = (*source)->Read(buffer.data(), buffer.size()); unblock_primary->set_value(); - primary_closed->get_future().wait(); ASSERT_THAT(result, IsOk()); EXPECT_THAT(std::string(buffer.data(), result->bytes_received), Eq("hedge")); + + // Flush the single-threaded read pool to ensure the losing primary attempt + // and its captured references to `client` have completely finished executing + // before tearing down the test. + { + google::cloud::internal::OptionsSpan const span( + client->options().set( + std::chrono::seconds(30))); + StatusOr> sync_source = + client->ReadObject(ReadObjectRangeRequest("test-bucket", "sync")); + ASSERT_THAT(sync_source, IsOk()); + ASSERT_THAT((*sync_source)->Read(buffer.data(), buffer.size()), IsOk()); + } } } // namespace