Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 43 additions & 19 deletions google/cloud/storage/internal/connection_impl_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -852,7 +853,6 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) {
// primary attempt. Any other thread is running a hedge.
auto primary_thread = std::make_shared<std::thread::id>();
auto unblock_primary = std::make_shared<std::promise<void>>();
auto primary_closed = std::make_shared<std::promise<void>>();
auto mock = std::make_unique<MockGenericStub>();
EXPECT_CALL(*mock, options).Times(AtLeast(0));
EXPECT_CALL(*mock, ReadObject)
Expand All @@ -865,24 +865,36 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) {
});
return StatusOr<std::unique_ptr<ObjectReadSource>>(std::move(source));
})
.WillRepeatedly([primary_thread, unblock_primary, primary_closed](
auto&, auto const&, ReadObjectRangeRequest const&) {
.WillOnce([primary_thread, unblock_primary](
Comment thread
kalragauri marked this conversation as resolved.
auto&, auto const&, ReadObjectRangeRequest const&) {
EXPECT_THAT(std::this_thread::get_id(), Eq(*primary_thread));
auto source = std::make_unique<testing::MockObjectReadSource>();
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::unique_ptr<ObjectReadSource>>(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<testing::MockObjectReadSource>();
EXPECT_CALL(*source, Read).WillOnce([](char* buf, std::size_t) {
return MakeReadResult("hedge", buf);
});
return StatusOr<std::unique_ptr<ObjectReadSource>>(std::move(source));
})
.WillOnce([primary_thread](auto&, auto const&,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The CI is failing because on slower runners, a second hedge timer can occasionally fire. If that happens, this extra hedge request consumes your 3rd .WillOnce() (which was meant for the final sync read) and trips the thread ID assertion since it runs on the background thread.

ReadObjectRangeRequest const&) {
EXPECT_THAT(std::this_thread::get_id(), Eq(*primary_thread));
auto source = std::make_unique<testing::MockObjectReadSource>();
EXPECT_CALL(*source, Read).WillOnce([](char* buf, std::size_t) {
return MakeReadResult("sync", buf);
});
return StatusOr<std::unique_ptr<ObjectReadSource>>(std::move(source));
});

Expand Down Expand Up @@ -916,9 +928,21 @@ TEST(RetryClientTest, HedgedReadRecordsMetricsOnGlobalMeterProvider) {
StatusOr<ReadSourceResult> 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<storage_experimental::ReadHedgeDelayOption>(
std::chrono::seconds(30)));
StatusOr<std::unique_ptr<ObjectReadSource>> sync_source =
client->ReadObject(ReadObjectRangeRequest("test-bucket", "sync"));
ASSERT_THAT(sync_source, IsOk());
ASSERT_THAT((*sync_source)->Read(buffer.data(), buffer.size()), IsOk());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

IIUC, the "sync" read does not eliminate the teardown race because it still runs with hedging enabled, which means it will enqueue a new task on read_pool_ that captures self = shared_from_this(), and it does not flush hedge_pool_.

Since winning tasks on both pools unblock the caller before the worker thread finishes and destroys its captured client reference, a preempted worker on either pool can still outlive the test scope and self-detach.

}
}

} // namespace
Expand Down
Loading