diff --git a/google/cloud/storage/internal/async/writer_connection_buffered.cc b/google/cloud/storage/internal/async/writer_connection_buffered.cc index c047cece556ea..c65ffed53b350 100644 --- a/google/cloud/storage/internal/async/writer_connection_buffered.cc +++ b/google/cloud/storage/internal/async/writer_connection_buffered.cc @@ -307,10 +307,15 @@ class AsyncWriterConnectionBufferedState std::unique_lock lk(mu_); write_offset_ += write_size; auto impl = Impl(lk); + auto const& state = impl->PersistedState(); + std::int64_t persisted_size = 0; + if (absl::holds_alternative(state)) { + persisted_size = absl::get(state).size(); + } else { + persisted_size = absl::get(state); + } lk.unlock(); - impl->Query().then([w = WeakFromThis()](auto f) { - if (auto self = w.lock()) return self->OnQuery(f.get()); - }); + OnQuery(persisted_size); } void OnQuery(StatusOr persisted_size) { diff --git a/google/cloud/storage/internal/async/writer_connection_buffered_test.cc b/google/cloud/storage/internal/async/writer_connection_buffered_test.cc index a179e4b1e6c6d..b86b8595fb39b 100644 --- a/google/cloud/storage/internal/async/writer_connection_buffered_test.cc +++ b/google/cloud/storage/internal/async/writer_connection_buffered_test.cc @@ -252,21 +252,22 @@ TEST(WriteConnectionBuffered, WriteBuffers) { "payload size", [](auto payload) { return payload.size(); }, Eq(n)); }; + auto mock_persisted_size = std::make_shared(0); auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); EXPECT_CALL(*mock, Write(expected_write_size(8 * 1024))).WillOnce([&](auto) { return sequencer.PushBack("Write").then([](auto) { return Status{}; }); }); - EXPECT_CALL(*mock, Flush(expected_write_size(24 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(32 * 1024)); - }); - }); + EXPECT_CALL(*mock, Flush(expected_write_size(24 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 32 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -303,9 +304,6 @@ TEST(WriteConnectionBuffered, WriteBuffers) { next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); ASSERT_TRUE(w3.is_ready()); EXPECT_STATUS_OK(w3.get()); @@ -319,33 +317,30 @@ TEST(WriteConnectionBuffered, WritePartialFlushAndFinalize) { "payload size", [](auto payload) { return payload.size(); }, Eq(n)); }; + auto mock_persisted_size = std::make_shared(0); auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); EXPECT_CALL(*mock, Write(expected_write_size(8 * 1024))).WillOnce([&](auto) { return sequencer.PushBack("Write").then([](auto) { return Status{}; }); }); - EXPECT_CALL(*mock, Flush(expected_write_size(24 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - { - InSequence seq; - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(24 * 1024)); - }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query2").then([](auto) { - return make_status_or(static_cast(24 * 1024 + 32 * 1024)); + EXPECT_CALL(*mock, Flush(expected_write_size(24 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 24 * 1024; + return Status{}; + }); }); - }); - } // The Finalize() call will flush the remaining data. - EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush2").then([](auto) { return Status{}; }); - }); + EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush2").then([mock_persisted_size](auto) { + *mock_persisted_size = 24 * 1024 + 32 * 1024; + return Status{}; + }); + }); EXPECT_CALL(*mock, Finalize(expected_write_size(0))).WillOnce([&](auto) { return sequencer.PushBack("Finalize").then([](auto) { return make_status_or(TestObject()); @@ -379,9 +374,6 @@ TEST(WriteConnectionBuffered, WritePartialFlushAndFinalize) { next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); // That should free enough data in the buffer to continue writing. ASSERT_TRUE(w1.is_ready()); @@ -395,9 +387,6 @@ TEST(WriteConnectionBuffered, WritePartialFlushAndFinalize) { EXPECT_EQ(next.second, "Flush2"); next.first.set_value(true); next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query2"); - next.first.set_value(true); - next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Finalize"); next.first.set_value(true); @@ -517,20 +506,19 @@ TEST(WriteConnectionBuffered, RewindError) { "payload size", [](auto payload) { return payload.size(); }, Eq(n)); }; + auto mock_persisted_size = std::make_shared(32 * 1024); auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(32 * 1024))); - EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - // This should not happen: the service is rewinding the value of - // `persisted_size`. The client should report this error. - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(16 * 1024)); - }); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 16 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -549,9 +537,6 @@ TEST(WriteConnectionBuffered, RewindError) { auto next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); ASSERT_TRUE(write.is_ready()); auto status = write.get(); @@ -571,20 +556,19 @@ TEST(WriteConnectionBuffered, FastForwardError) { "payload size", [](auto payload) { return payload.size(); }, Eq(n)); }; + auto mock_persisted_size = std::make_shared(32 * 1024); auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(32 * 1024))); - EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - // This should not happen: the service is reporting more data persisted than - // the size of the data sent. The client should report this error. - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(128 * 1024)); - }); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 128 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -603,9 +587,6 @@ TEST(WriteConnectionBuffered, FastForwardError) { auto next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); ASSERT_TRUE(write.is_ready()); auto status = write.get(); @@ -659,16 +640,17 @@ TEST(WriteConnectionBuffered, Flush) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(8 * 1024)); - }); - }); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); + EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 8 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -682,9 +664,6 @@ TEST(WriteConnectionBuffered, Flush) { auto next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); EXPECT_STATUS_OK(w0.get()); } @@ -698,16 +677,19 @@ TEST(WriteConnectionBuffered, FlushWithEmptyPayload) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - // This flush is for the empty payload. - EXPECT_CALL(*mock, Flush(expected_write_size(0))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then( - [](auto) { return make_status_or(static_cast(0)); }); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + // This flush is for the empty payload. + EXPECT_CALL(*mock, Flush(expected_write_size(0))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + // *mock_persisted_size = 0; // Empty payload, so it's 0, which is the + // current value + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -720,9 +702,6 @@ TEST(WriteConnectionBuffered, FlushWithEmptyPayload) { auto next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); EXPECT_STATUS_OK(f.get()); } @@ -793,18 +772,18 @@ TEST(WriteConnectionBuffered, FlushResumesAndDoesNotCompletePrematurely) { EXPECT_CALL(*mock2, UploadId).WillRepeatedly(Return("test-upload-id")); // OnResume will query persisted state. We return 0, meaning the data was // lost. - EXPECT_CALL(*mock2, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - // The resumed connection should receive the Flush call again. - EXPECT_CALL(*mock2, Flush(expected_write_size(8 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush2").then([](auto) { return Status{}; }); - }); - // After Flush2 succeeds, it will query the status. - EXPECT_CALL(*mock2, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(8 * 1024)); - }); + auto mock2_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock2, PersistedState).WillRepeatedly([mock2_persisted_size] { + return MakePersistedState(*mock2_persisted_size); }); + // The resumed connection should receive the Flush call again. + EXPECT_CALL(*mock2, Flush(expected_write_size(8 * 1024))) + .WillOnce([&, mock2_persisted_size](auto) { + return sequencer.PushBack("Flush2").then([mock2_persisted_size](auto) { + *mock2_persisted_size = 8 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).WillOnce([&]() { @@ -839,11 +818,6 @@ TEST(WriteConnectionBuffered, FlushResumesAndDoesNotCompletePrematurely) { EXPECT_EQ(next.second, "Flush2"); next.first.set_value(true); - // Let the query complete. - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - // Now the flush promise should be completed. ASSERT_TRUE(f.is_ready()); EXPECT_STATUS_OK(f.get()); @@ -906,17 +880,18 @@ TEST(WriteConnectionBuffered, CloseWithPayload) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - // The payload is flushed first. - EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(8 * 1024)); - }); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + // The payload is flushed first. + EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 8 * 1024; + return Status{}; + }); + }); // Then the stream is closed with empty payload. EXPECT_CALL(*mock, Close(expected_write_size(0))).WillOnce([&](auto) { return sequencer.PushBack("Close").then([](auto) { return Status{}; }); @@ -935,10 +910,6 @@ TEST(WriteConnectionBuffered, CloseWithPayload) { EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Close"); next.first.set_value(true); @@ -1173,23 +1144,23 @@ TEST(WriteConnectionBuffered, FlushSendsAllBufferedData) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); // The first Write() (8KiB) is immediately sent, not buffered. EXPECT_CALL(*mock, Write(expected_write_size(8 * 1024))).WillOnce([&](auto) { return sequencer.PushBack("Write-8K").then([](auto) { return Status{}; }); }); // The Flush() call has 4KiB of data. Since the buffer is empty, this // 4KiB is flushed. - EXPECT_CALL(*mock, Flush(expected_write_size(4 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush-4K").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or( - static_cast(8 * 1024 + 4 * 1024)); // Total persisted - }); - }); + EXPECT_CALL(*mock, Flush(expected_write_size(4 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush-4K").then([mock_persisted_size](auto) { + *mock_persisted_size = 8 * 1024 + 4 * 1024; + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -1212,9 +1183,6 @@ TEST(WriteConnectionBuffered, FlushSendsAllBufferedData) { next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush-4K"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); EXPECT_STATUS_OK(f1.get()); } @@ -1228,23 +1196,23 @@ TEST(WriteConnectionBuffered, FinalizeSendsAllBufferedData) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); // The first Write() (8KiB) is immediately sent, not buffered. EXPECT_CALL(*mock, Write(expected_write_size(8 * 1024))).WillOnce([&](auto) { return sequencer.PushBack("Write-8K").then([](auto) { return Status{}; }); }); // The Finalize() call has 4KiB of data. Since the buffer is empty, this // 4KiB is flushed. - EXPECT_CALL(*mock, Flush(expected_write_size(4 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush-4K").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or( - static_cast(8 * 1024 + 4 * 1024)); // Total persisted - }); - }); + EXPECT_CALL(*mock, Flush(expected_write_size(4 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush-4K").then([mock_persisted_size](auto) { + *mock_persisted_size = 8 * 1024 + 4 * 1024; + return Status{}; + }); + }); // After the flush, the buffer is empty, and we can finalize. EXPECT_CALL(*mock, Finalize(expected_write_size(0))).WillOnce([&](auto) { return sequencer.PushBack("Finalize").then([](auto) { @@ -1273,9 +1241,6 @@ TEST(WriteConnectionBuffered, FinalizeSendsAllBufferedData) { EXPECT_EQ(next.second, "Flush-4K"); next.first.set_value(true); next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Finalize"); next.first.set_value(true); @@ -1291,33 +1256,29 @@ TEST(WriteConnectionBuffered, WriteTriggersMultipleFlushes) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); { InSequence seq; // Expect two separate flushes as the buffer fills up twice. EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))) - .WillOnce([&](auto) { - return sequencer.PushBack("Flush1").then( - [](auto) { return Status{}; }); + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush1").then([mock_persisted_size](auto) { + *mock_persisted_size = 32 * 1024; + return Status{}; + }); }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query1").then([](auto) { - return make_status_or(static_cast(32 * 1024)); - }); - }); EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024))) - .WillOnce([&](auto) { - return sequencer.PushBack("Flush2").then( - [](auto) { return Status{}; }); + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush2").then([mock_persisted_size](auto) { + *mock_persisted_size = 64 * 1024; + return Status{}; + }); }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query2").then([](auto) { - return make_status_or(static_cast(64 * 1024)); - }); - }); } MockFactory mock_factory; @@ -1334,9 +1295,6 @@ TEST(WriteConnectionBuffered, WriteTriggersMultipleFlushes) { auto next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush1"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query1"); - next.first.set_value(true); EXPECT_STATUS_OK(w1.get()); // Write enough data again to go over the HWM and block. @@ -1347,9 +1305,6 @@ TEST(WriteConnectionBuffered, WriteTriggersMultipleFlushes) { next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Flush2"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query2"); - next.first.set_value(true); EXPECT_STATUS_OK(w2.get()); } @@ -1362,31 +1317,30 @@ TEST(WriteConnectionBuffered, MultipleConcurrentFlushesAreQueued) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); { InSequence seq; // The implementation sends the first Flush() payload immediately. - EXPECT_CALL(*mock, Flush(expected_write_size(4096))).WillOnce([&](auto) { - return sequencer.PushBack("Flush1").then([](auto) { return Status{}; }); - }); - // The Flush is followed by a Query. - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query1").then( - [](auto) { return make_status_or(static_cast(4096)); }); - }); + EXPECT_CALL(*mock, Flush(expected_write_size(4096))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush1").then([mock_persisted_size](auto) { + *mock_persisted_size = 4096; + return Status{}; + }); + }); + EXPECT_CALL(*mock, Flush(expected_write_size(8192))) .Times(1) - .WillOnce([&](auto) { - return sequencer.PushBack("Flush2").then( - [](auto) { return Status{}; }); + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush2").then([mock_persisted_size](auto) { + *mock_persisted_size = 4096 + 8192; + return Status{}; + }); }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query2").then([](auto) { - return make_status_or(static_cast(4096 + 8192)); - }); - }); } MockFactory mock_factory; @@ -1406,11 +1360,6 @@ TEST(WriteConnectionBuffered, MultipleConcurrentFlushesAreQueued) { EXPECT_EQ(next.second, "Flush1"); next.first.set_value(true); - // Satisfy the Query call. - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query1"); - next.first.set_value(true); - // After the Query, the first Flush future should be completed. EXPECT_STATUS_OK(f1.get()); @@ -1419,11 +1368,6 @@ TEST(WriteConnectionBuffered, MultipleConcurrentFlushesAreQueued) { EXPECT_EQ(next.second, "Flush2"); next.first.set_value(true); - // The second Flush is also followed by a Query. - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query2"); - next.first.set_value(true); - EXPECT_STATUS_OK(f2.get()); } diff --git a/google/cloud/storage/internal/async/writer_connection_resumed.cc b/google/cloud/storage/internal/async/writer_connection_resumed.cc index 37faf9983817e..23ffea9395bbc 100644 --- a/google/cloud/storage/internal/async/writer_connection_resumed.cc +++ b/google/cloud/storage/internal/async/writer_connection_resumed.cc @@ -324,14 +324,16 @@ class AsyncWriterConnectionResumedState std::unique_lock lk(mu_); write_offset_ += write_size; auto impl = Impl(lk); + auto const& state = impl->PersistedState(); + std::int64_t persisted_size = 0; + if (absl::holds_alternative(state)) { + persisted_size = absl::get(state).size(); + } else { + persisted_size = absl::get(state); + } lk.unlock(); - impl->Query().then([result, w = WeakFromThis()](auto f) { - auto self = w.lock(); - if (!self) return; - self->OnQuery(f.get()); - self->SetFlushed(std::unique_lock(self->mu_), - std::move(result)); - }); + OnQuery(persisted_size); + SetFlushed(std::unique_lock(mu_), std::move(result)); } void OnQuery(StatusOr persisted_size) { diff --git a/google/cloud/storage/internal/async/writer_connection_resumed_test.cc b/google/cloud/storage/internal/async/writer_connection_resumed_test.cc index 06375cfc9fa59..3968ba8a8fdb0 100644 --- a/google/cloud/storage/internal/async/writer_connection_resumed_test.cc +++ b/google/cloud/storage/internal/async/writer_connection_resumed_test.cc @@ -174,17 +174,19 @@ TEST(WriterConnectionResumed, FlushEmpty) { auto first_response = google::storage::v2::BidiWriteObjectResponse{}; auto mock = std::make_unique(); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - EXPECT_CALL(*mock, WriteHandle).WillRepeatedly(Return(std::nullopt)); - EXPECT_CALL(*mock, Flush).WillRepeatedly([&](auto const& p) { - EXPECT_TRUE(p.payload().empty()); - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then( - [](auto) -> StatusOr { return 0; }); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + EXPECT_CALL(*mock, WriteHandle).WillRepeatedly(Return(std::nullopt)); + EXPECT_CALL(*mock, Flush) + .WillRepeatedly([&, mock_persisted_size](auto const& p) { + EXPECT_TRUE(p.payload().empty()); + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + // Persisted size remains 0 for empty payload + return Status{}; + }); + }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -199,9 +201,6 @@ TEST(WriterConnectionResumed, FlushEmpty) { EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); EXPECT_THAT(flush.get(), StatusIs(StatusCode::kOk)); } @@ -212,14 +211,17 @@ TEST(WriteConnectionResumed, FlushNonEmpty) { auto first_response = google::storage::v2::BidiWriteObjectResponse{}; auto const payload = TestPayload(1024); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); EXPECT_CALL(*mock, WriteHandle).WillRepeatedly(Return(std::nullopt)); EXPECT_CALL(*mock, Flush) - .WillOnce([&](auto const& p) { + .WillOnce([&, mock_persisted_size, payload](auto const& p) { EXPECT_EQ(p.payload(), payload.payload()); - return sequencer.PushBack("Flush").then([](auto f) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto f) { if (!f.get()) return TransientError(); + *mock_persisted_size = 1024; return Status{}; }); }) @@ -227,13 +229,6 @@ TEST(WriteConnectionResumed, FlushNonEmpty) { EXPECT_TRUE(p.payload().empty()); return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then( - [](auto f) -> StatusOr { - if (!f.get()) return TransientError(); - return 1024; - }); - }); MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -253,10 +248,6 @@ TEST(WriteConnectionResumed, FlushNonEmpty) { EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - EXPECT_TRUE(flush.is_ready()); EXPECT_THAT(flush.get(), StatusIs(StatusCode::kOk)); @@ -406,16 +397,14 @@ TEST(WriteConnectionResumed, NoConcurrentWritesWhenFlushAndWriteRace) { auto initial_request = google::storage::v2::BidiWriteObjectRequest{}; auto first_response = google::storage::v2::BidiWriteObjectResponse{}; - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); EXPECT_CALL(*mock, WriteHandle).WillRepeatedly(Return(std::nullopt)); EXPECT_CALL(*mock, Flush(_)).WillRepeatedly([&](auto) { return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then( - [](auto) -> StatusOr { return 0; }); - }); // Make Write detect concurrent invocations. If two writes run concurrently // the compare_exchange will fail and the test will fail. @@ -438,19 +427,13 @@ TEST(WriteConnectionResumed, NoConcurrentWritesWhenFlushAndWriteRace) { // Start a flush which will call impl->Flush() and block. auto flush_future = connection->Flush({}); - // Allow the Flush to complete, this will schedule a Query (but Query will - // remain blocked until we pop it). - auto next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Flush"); - next.first.set_value(true); - // Immediately perform a user Write after the flush completed but before - // Query completes. This can race with the OnQuery-driven write. + // Immediately perform a user Write after the flush started. This can race + // with the OnFlush-driven write continuation. auto write_future = connection->Write(TestPayload(1024)); - // Now allow the Query to complete; OnQuery may schedule a write. - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); + auto next = sequencer.PopFrontWithName(); + EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); // Wait for both futures to complete with a timeout to avoid indefinite hang. @@ -545,8 +528,10 @@ TEST(WriterConnectionResumed, OnQueryUpdatesWriteHandle) { google::storage::v2::BidiWriteObjectResponse first_response; first_response.mutable_write_handle()->set_handle("initial-handle"); - EXPECT_CALL(*mock_ptr, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock_ptr, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); + }); google::storage::v2::BidiWriteHandle new_handle; new_handle.set_handle("updated-handle"); @@ -555,10 +540,13 @@ TEST(WriterConnectionResumed, OnQueryUpdatesWriteHandle) { auto const expected_payload = std::string(1024, 'A'); EXPECT_CALL(*mock_ptr, Flush(_)) - .WillOnce([&](auto const& p) { + .WillOnce([&, mock_persisted_size](auto const& p) { EXPECT_EQ(p.size(), expected_payload.size()); - return sequencer.PushBack("Flush").then([](auto f) { - if (f.get()) return Status{}; + return sequencer.PushBack("Flush").then([mock_persisted_size](auto f) { + if (f.get()) { + *mock_persisted_size = 1024; + return Status{}; + } return TransientError(); }); }) @@ -570,19 +558,6 @@ TEST(WriterConnectionResumed, OnQueryUpdatesWriteHandle) { }); }); - EXPECT_CALL(*mock_ptr, Query) - .WillOnce([&]() { - return sequencer.PushBack("Query").then( - [](auto f) -> StatusOr { - if (!f.get()) return TransientError(); - return 1024; - }); - }) - .WillOnce([&]() { - return sequencer.PushBack("GhostQuery") - .then([](auto) -> StatusOr { return 1024; }); - }); - MockFactory mock_factory; EXPECT_CALL(mock_factory, Call).Times(0); @@ -600,18 +575,10 @@ TEST(WriterConnectionResumed, OnQueryUpdatesWriteHandle) { EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "GhostFlush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "GhostQuery"); - next.first.set_value(true); - EXPECT_THAT(flush.get(), StatusIs(StatusCode::kOk)); current_handle = connection->WriteHandle(); @@ -1032,17 +999,18 @@ TEST(WriterConnectionResumed, CloseWithPayload) { auto mock = std::make_unique(); EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id")); - EXPECT_CALL(*mock, PersistedState) - .WillRepeatedly(Return(MakePersistedState(0))); - // The payload is flushed first. - EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))).WillOnce([&](auto) { - return sequencer.PushBack("Flush").then([](auto) { return Status{}; }); - }); - EXPECT_CALL(*mock, Query).WillOnce([&]() { - return sequencer.PushBack("Query").then([](auto) { - return make_status_or(static_cast(8 * 1024)); - }); + auto mock_persisted_size = std::make_shared(0); + EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] { + return MakePersistedState(*mock_persisted_size); }); + // The payload is flushed first. + EXPECT_CALL(*mock, Flush(expected_write_size(8 * 1024))) + .WillOnce([&, mock_persisted_size](auto) { + return sequencer.PushBack("Flush").then([mock_persisted_size](auto) { + *mock_persisted_size = 8 * 1024; + return Status{}; + }); + }); // Then the stream is closed with empty payload. EXPECT_CALL(*mock, Close(expected_write_size(0))).WillOnce([&](auto) { return sequencer.PushBack("Close").then([](auto) { return Status{}; }); @@ -1062,10 +1030,6 @@ TEST(WriterConnectionResumed, CloseWithPayload) { EXPECT_EQ(next.second, "Flush"); next.first.set_value(true); - next = sequencer.PopFrontWithName(); - EXPECT_EQ(next.second, "Query"); - next.first.set_value(true); - next = sequencer.PopFrontWithName(); EXPECT_EQ(next.second, "Close"); next.first.set_value(true);