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 d36181bfcf748..eecbcba6d6d85 100644 --- a/google/cloud/storage/internal/async/object_descriptor_impl_test.cc +++ b/google/cloud/storage/internal/async/object_descriptor_impl_test.cc @@ -2327,12 +2327,12 @@ TEST(ObjectDescriptorImpl, PrewarmedCacheHit) { return sequencer.PushBack("Read[1]").then([&](auto) { auto response = Response{}; EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); - return absl::make_optional(response); + return std::make_optional(response); }); }) .WillOnce([&sequencer]() { return sequencer.PushBack("Read[2]").then( - [](auto) { return absl::optional{}; }); + [](auto) { return std::optional{}; }); }); EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { return sequencer.PushBack("Finish").then( @@ -2374,7 +2374,7 @@ TEST(ObjectDescriptorImpl, PrewarmedCacheHit) { VariantWith(ResultOf( "contents are", [](storage::ReadPayload const& p) { return p.contents(); }, - ElementsAre(absl::string_view{"Pre-warmed data"})))); + ElementsAre(std::string_view{"Pre-warmed data"})))); EXPECT_THAT(s1->Read().get(), VariantWith(IsOk())); @@ -2426,12 +2426,12 @@ TEST(ObjectDescriptorImpl, PrewarmedCacheMiss) { return sequencer.PushBack("Read[1]").then([&](auto) { auto response = Response{}; EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); - return absl::make_optional(response); + return std::make_optional(response); }); }) .WillOnce([&sequencer]() { return sequencer.PushBack("Read[2]").then( - [](auto) { return absl::optional{}; }); + [](auto) { return std::optional{}; }); }); EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { return sequencer.PushBack("Finish").then( @@ -2477,7 +2477,7 @@ TEST(ObjectDescriptorImpl, PrewarmedCacheMiss) { VariantWith(ResultOf( "contents are", [](storage::ReadPayload const& p) { return p.contents(); }, - ElementsAre(absl::string_view{"Missed range data"})))); + ElementsAre(std::string_view{"Missed range data"})))); EXPECT_THAT(s1->Read().get(), VariantWith(IsOk())); @@ -2530,12 +2530,12 @@ TEST(ObjectDescriptorImpl, PrewarmedPacingEviction) { return sequencer.PushBack("Read[1]").then([&](auto) { auto response = Response{}; EXPECT_TRUE(TextFormat::ParseFromString(kResponse1, &response)); - return absl::make_optional(response); + return std::make_optional(response); }); }) .WillOnce([&sequencer]() { return sequencer.PushBack("Read[2]").then( - [](auto) { return absl::optional{}; }); + [](auto) { return std::optional{}; }); }); EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { return sequencer.PushBack("Finish").then( @@ -2583,6 +2583,340 @@ TEST(ObjectDescriptorImpl, PrewarmedPacingEviction) { next.first.set_value(true); } +TEST(ObjectDescriptorImpl, InitialReadRangesExactRangeCacheHit) { + 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: "Exact hit 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 std::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return std::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 1024}}); + + 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, 1024}); + 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(std::string_view{"Exact hit 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, InitialReadRangesSubRangeCacheMiss) { + 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: 0 read_length: 512 } + )pb"; + + auto constexpr kResponse1 = R"pb( + object_data_ranges { + range_end: true + read_range { read_id: 2 read_offset: 0 } + checksummed_data { content: "Sub 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 std::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return std::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 1024}}); + + 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, 512}); + 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(std::string_view{"Sub 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); +} + +TEST(ObjectDescriptorImpl, PreWarmBufferLimitEvictionTombstone) { + 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 std::make_optional(response); + }); + }) + .WillOnce([&sequencer]() { + return sequencer.PushBack("Read[2]").then( + [](auto) { return std::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + + Options options; + options.set(true); + options.set({{0, 10}}); + // Evict if buffer exceeds 5 bytes. + 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]"); + + // Pre-warmed range should be evicted. Fallback RPC will be requested. + 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); +} + +TEST(ObjectDescriptorImpl, DuplicateInitialRangesDeduplication) { + auto constexpr kResponse0 = R"pb( + metadata { + bucket: "projects/_/buckets/test-bucket" + name: "test-object" + generation: 42 + } + read_handle { handle: "handle-12345" } + )pb"; + + AsyncSequencer sequencer; + auto stream = std::make_unique(); + + // We expect no extra Write calls since duplicates are deduplicated. + EXPECT_CALL(*stream, Write).Times(0); + + EXPECT_CALL(*stream, Read).WillOnce([&sequencer]() { + return sequencer.PushBack("Read[1]").then( + [](auto) { return std::optional{}; }); + }); + EXPECT_CALL(*stream, Finish).WillOnce([&sequencer]() { + return sequencer.PushBack("Finish").then( + [](auto) { return PermanentError(); }); + }); + + MockFactory factory; + EXPECT_CALL(factory, Call).WillOnce([](Request const&) { + return make_ready_future(StatusOr(PermanentError())); + }); + Options options; + options.set(true); + options.set({{0, 1024}, {0, 1024}}); + + 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 next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Read[1]"); + 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/tests/async_client_integration_test.cc b/google/cloud/storage/tests/async_client_integration_test.cc index 27bc9c9f07940..53ad2f3974056 100644 --- a/google/cloud/storage/tests/async_client_integration_test.cc +++ b/google/cloud/storage/tests/async_client_integration_test.cc @@ -1127,6 +1127,157 @@ TEST_F(AsyncClientIntegrationTest, OpenWithChecksumValidation) { storage::Generation(metadata->generation())); } +// Verifies that opening an object with multiple initial disjoint read ranges +// properly pre-warms all ranges and that subsequent Read calls for those exact +// ranges return the expected byte contents. +TEST_F(AsyncClientIntegrationTest, MultiRangeOpenMultipleDisjointRanges) { + if (!UsingEmulator()) GTEST_SKIP(); + auto async = AsyncClient(TestOptions()); + auto client = MakeIntegrationTestClient(true, TestOptions()); + auto object_name = MakeRandomObjectName(); + + StatusOr create = client.CreateBucket( + bucket_name(), storage::BucketMetadata{}.set_location("us-west4")); + if (!create && create.status().code() != StatusCode::kAlreadyExists) { + GTEST_FAIL() << "cannot create bucket: " << create.status(); + } + + // Populate object with random test data (32 KB total payload size). + auto constexpr kSize = 32 * 1024; + auto const block = MakeRandomData(kSize); + + StatusOr> w = + async.StartAppendableObjectUpload(BucketName(bucket_name()), object_name) + .get(); + ASSERT_STATUS_OK(w); + AsyncWriter writer; + AsyncToken token; + std::tie(writer, token) = *std::move(w); + StatusOr p = + writer.Write(std::move(token), WritePayload(block)).get(); + ASSERT_STATUS_OK(p); + token = *std::move(p); + + auto metadata = writer.Finalize(std::move(token)).get(); + ASSERT_STATUS_OK(metadata); + + // Specify 3 non-overlapping, disjoint initial byte ranges with Open: + // Range 1: [0, 1024) + // Range 2: [4096, 5120) + // Range 3: [16384, 17408) + AsyncClient::InitialReadRanges ranges; + ranges.initial_ranges = {{0, 1024}, {4096, 1024}, {16384, 1024}}; + auto descriptor = + async.Open(BucketName(bucket_name()), object_name, std::move(ranges)) + .get(); + ASSERT_STATUS_OK(descriptor); + + std::string actual1, actual2, actual3; + + // Read Range 1 ([0, 1024)) and assert it hits pre-warmed cache with + // correct content. + { + AsyncReader r; + AsyncToken t; + std::tie(r, t) = descriptor->Read(0, 1024); + while (t.valid()) { + auto read = r.Read(std::move(t)).get(); + ASSERT_STATUS_OK(read); + ReadPayload p; + std::tie(p, t) = *std::move(read); + for (std::string_view sv : p.contents()) absl::StrAppend(&actual1, sv); + } + } + EXPECT_EQ(actual1.size(), 1024); + EXPECT_EQ(actual1, block.substr(0, 1024)); + + // Read Range 2 ([4096, 5120)) and assert exact match. + { + AsyncReader r; + AsyncToken t; + std::tie(r, t) = descriptor->Read(4096, 1024); + while (t.valid()) { + auto read = r.Read(std::move(t)).get(); + ASSERT_STATUS_OK(read); + ReadPayload p; + std::tie(p, t) = *std::move(read); + for (std::string_view sv : p.contents()) absl::StrAppend(&actual2, sv); + } + } + EXPECT_EQ(actual2.size(), 1024); + EXPECT_EQ(actual2, block.substr(4096, 1024)); + + // Read Range 3 ([16384, 17408)) and assert exact match. + { + AsyncReader r; + AsyncToken t; + std::tie(r, t) = descriptor->Read(16384, 1024); + while (t.valid()) { + auto read = r.Read(std::move(t)).get(); + ASSERT_STATUS_OK(read); + ReadPayload p; + std::tie(p, t) = *std::move(read); + for (std::string_view sv : p.contents()) absl::StrAppend(&actual3, sv); + } + } + EXPECT_EQ(actual3.size(), 1024); + EXPECT_EQ(actual3, block.substr(16384, 1024)); + client.DeleteObject(bucket_name(), object_name, + storage::Generation(metadata->generation())); +} + +// Verifies teardown safety: when an ObjectDescriptor with active pre-warming +// streams goes out of scope before data is consumed, the underlying gRPC +// streams are safely cancelled and torn down without leaks, or unhandled +// exceptions. +TEST_F(AsyncClientIntegrationTest, MultiRangeOpenDestructorTeardown) { + if (!UsingEmulator()) GTEST_SKIP(); + auto async = AsyncClient(TestOptions()); + auto client = MakeIntegrationTestClient(true, TestOptions()); + auto object_name = MakeRandomObjectName(); + + auto create = client.CreateBucket( + bucket_name(), storage::BucketMetadata{}.set_location("us-west4")); + if (!create && create.status().code() != StatusCode::kAlreadyExists) { + GTEST_FAIL() << "cannot create bucket: " << create.status(); + } + + // Create 1 MB object payload. + auto constexpr kSize = 1024 * 1024; + auto const block = MakeRandomData(kSize); + + auto w = + async.StartAppendableObjectUpload(BucketName(bucket_name()), object_name) + .get(); + ASSERT_STATUS_OK(w); + AsyncWriter writer; + AsyncToken token; + std::tie(writer, token) = *std::move(w); + auto p = writer.Write(std::move(token), WritePayload(block)).get(); + ASSERT_STATUS_OK(p); + token = *std::move(p); + + auto metadata = writer.Finalize(std::move(token)).get(); + ASSERT_STATUS_OK(metadata); + + { + // Open descriptor with a 1 MB initial range. + AsyncClient::InitialReadRanges ranges; + ranges.initial_ranges = {{0, 1024 * 1024}}; + auto descriptor = + async.Open(BucketName(bucket_name()), object_name, std::move(ranges)) + .get(); + ASSERT_STATUS_OK(descriptor); + // Force immediate destructor invocation on descriptor while streams are + // actively pre-warming. This tests proper stream cancellation/finish on + // teardown. + } + + // Ensure bucket cleanup proceeds without issues. + client.DeleteObject(bucket_name(), object_name, + storage::Generation(metadata->generation())); +} + } // namespace GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_END } // namespace storage