Skip to content
Merged
Show file tree
Hide file tree
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
50 changes: 29 additions & 21 deletions google/cloud/storage/internal/async/writer_connection_buffered.cc
Original file line number Diff line number Diff line change
Expand Up @@ -169,9 +169,12 @@ class AsyncWriterConnectionBufferedState
// Create a new promise for this flush operation.
promise<Status> current_flush_promise;
auto f = current_flush_promise.get_future();
pending_flush_promises_.push_back(std::move(current_flush_promise));

resend_buffer_.Append(WritePayloadImpl::GetImpl(p));
auto const target_offset =
buffer_offset_ + static_cast<std::int64_t>(resend_buffer_.size());
pending_flush_promises_.push_back(
PendingFlush{std::move(current_flush_promise), target_offset});

flush_ = true;
HandleNewData(std::move(lk), true);
// Return the future associated with the new promise.
Expand Down Expand Up @@ -396,7 +399,7 @@ class AsyncWriterConnectionBufferedState
return;
}
// SetFlushed will release the lock before returning.
SetFlushed(std::move(lk), Status{});
SetFlushed(std::move(lk), Status{}, persisted_size);
// Re-acquire the lock to re-enter the write loop.
WriteLoop(std::unique_lock<std::mutex>(mu_));
// The notifications are deferred until the lock is released, as they might
Expand Down Expand Up @@ -529,7 +532,7 @@ class AsyncWriterConnectionBufferedState
lk.unlock();
// Notify handlers and pending flushes *after* releasing the lock.
for (auto& h : handlers) h->Execute(Status{});
for (auto& pf : pending_flushes) pf.set_value(Status{}); // Success
for (auto& pf : pending_flushes) pf.p.set_value(Status{}); // Success
p.set_value(std::move(object)); // Set value on the moved promise
}

Expand All @@ -553,34 +556,30 @@ class AsyncWriterConnectionBufferedState
lk.unlock();
// Notify handlers and pending flushes after releasing the lock.
for (auto& h : handlers) h->Execute(status);
for (auto& pf : pending_flushes) pf.set_value(status);
for (auto& pf : pending_flushes) pf.p.set_value(status);
p.set_value(status); // Set value on the moved promise.
}

void SetFlushed(std::unique_lock<std::mutex> lk, Status const& result) {
void SetFlushed(std::unique_lock<std::mutex> lk, Status const& result,
std::int64_t persisted_size) {
if (!result.ok()) return SetError(std::move(lk), std::move(result));
// Do NOT reset finalize_ or finalizing_ here.
auto handlers = ClearHandlers(lk);
// Dequeue the promise corresponding to an explicit Flush() call, if any.
if (pending_flush_promises_.empty()) {
// This can happen if SetError cleared the queue first, or if this
// flush was triggered internally by buffer size (not by an explicit
// Flush() call) and thus has no promise in the queue.
flush_ = false;
lk.unlock();
for (auto& h : handlers) h->Execute(Status{});
return;
std::vector<promise<Status>> flushes_to_complete;
while (!pending_flush_promises_.empty() &&
pending_flush_promises_.front().target_offset <= persisted_size) {
flushes_to_complete.push_back(
std::move(pending_flush_promises_.front().p));
pending_flush_promises_.pop_front();
}
auto flushed = std::move(pending_flush_promises_.front());
pending_flush_promises_.pop_front();
if (pending_flush_promises_.empty()) {
flush_ = false;
}
lk.unlock(); // Unlock only once before notifying
// Notify handlers and the specific flush promise *after* releasing the
// Notify handlers and the specific flush promises *after* releasing the
// lock.
for (auto& h : handlers) h->Execute(Status{});
flushed.set_value(result);
for (auto& f : flushes_to_complete) f.set_value(result);
}

void SetError(std::unique_lock<std::mutex> lk, Status const& status) {
Expand Down Expand Up @@ -621,7 +620,7 @@ class AsyncWriterConnectionBufferedState
for (auto& h : handlers) h->Execute(status);
// Set error on all pending flush promises.
for (auto& pf : pending_flushes) {
pf.set_value(status);
pf.p.set_value(status);
}
// Set error on the moved promises *once*.
if (complete_finalized) {
Expand Down Expand Up @@ -684,8 +683,17 @@ class AsyncWriterConnectionBufferedState
// closed_.
future<Status> closed_future_;

// Tracks an outstanding `Flush()` promise alongside the stream offset at the
// time `Flush()` was called. The target offset is used to satisfy promises
// once all data buffered at the time of the `Flush()` call has been
// persisted.
struct PendingFlush {
promise<Status> p;
std::int64_t target_offset;
};

// Queue of promises for outstanding Flush() calls.
std::deque<promise<Status>> pending_flush_promises_;
std::deque<PendingFlush> pending_flush_promises_;

// The resend buffer. If there is an error, this will have all the data since
// the last persisted byte and will be resent.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -824,6 +824,74 @@ TEST(WriteConnectionBuffered, FlushResumesAndDoesNotCompletePrematurely) {
EXPECT_STATUS_OK(f.get());
}

// Verifies that multiple queued Flush() operations track their
// respective cumulative byte offsets and that each flush future is satisfied
// only when the server's persisted_size reaches or exceeds its target offset.
TEST(WriteConnectionBuffered, InterleavedMultiFlush) {
AsyncSequencer<bool> sequencer;

auto expected_write_size = [](std::size_t n) {
return ResultOf(
"payload size", [](auto payload) { return payload.size(); }, Eq(n));
};

auto mock_persisted_size = std::make_shared<std::int64_t>(0);
auto mock = std::make_unique<MockAsyncWriterConnection>();
EXPECT_CALL(*mock, UploadId).WillRepeatedly(Return("test-upload-id"));
EXPECT_CALL(*mock, PersistedState).WillRepeatedly([mock_persisted_size] {
return MakePersistedState(*mock_persisted_size);
});

// Expect two 32 KiB flushes in sequence.
EXPECT_CALL(*mock, Flush(expected_write_size(32 * 1024)))
.WillOnce([&, mock_persisted_size](auto) {
return sequencer.PushBack("Flush1").then([mock_persisted_size](auto) {
*mock_persisted_size = 32 * 1024;
return Status{};
});
})
.WillOnce([&, mock_persisted_size](auto) {
return sequencer.PushBack("Flush2").then([mock_persisted_size](auto) {
*mock_persisted_size = 64 * 1024;
return Status{};
});
});

MockFactory mock_factory;
EXPECT_CALL(mock_factory, Call).Times(0);

auto connection = MakeWriterConnectionBuffered(
mock_factory.AsStdFunction(), std::move(mock), TestOptions());

// Issue first Flush() of 32 KiB.
auto f1 = connection->Flush(TestPayload(32 * 1024));
ASSERT_FALSE(f1.is_ready());

// Issue second Flush() of 32 KiB while the first is still in-flight.
auto f2 = connection->Flush(TestPayload(32 * 1024));
ASSERT_FALSE(f2.is_ready());

// Complete the first flush on the mock (persisted_size reaches 32 KiB).
auto next = sequencer.PopFrontWithName();
EXPECT_EQ(next.second, "Flush1");
next.first.set_value(true);

// f1 must be satisfied with OK, while f2 must remain pending because its
// target offset (64 KiB) has not yet been reached.
ASSERT_TRUE(f1.is_ready());
EXPECT_STATUS_OK(f1.get());
ASSERT_FALSE(f2.is_ready());

// Complete the second flush on the mock (persisted_size reaches 64 KiB).
next = sequencer.PopFrontWithName();
EXPECT_EQ(next.second, "Flush2");
next.first.set_value(true);

// Now f2 must also be satisfied with OK.
ASSERT_TRUE(f2.is_ready());
EXPECT_STATUS_OK(f2.get());
}

TEST(WriteConnectionBuffered, FinalizeWhileFlushing) {
AsyncSequencer<bool> sequencer;

Expand Down
74 changes: 48 additions & 26 deletions google/cloud/storage/internal/async/writer_connection_resumed.cc
Original file line number Diff line number Diff line change
Expand Up @@ -158,9 +158,12 @@ class AsyncWriterConnectionResumedState
// Create a new promise for this flush operation.
promise<Status> current_flush_promise;
auto f = current_flush_promise.get_future();
pending_flush_promises_.push_back(std::move(current_flush_promise));

resend_buffer_.Append(WritePayloadImpl::GetImpl(p));
auto const target_offset =
buffer_offset_ + static_cast<std::int64_t>(resend_buffer_.size());
pending_flush_promises_.push_back(
PendingFlush{std::move(current_flush_promise), target_offset});

flush_ = true;
HandleNewData(std::move(lk), true);
// Return the future associated with the new promise.
Expand Down Expand Up @@ -344,12 +347,12 @@ class AsyncWriterConnectionResumedState
}
lk.unlock();
OnQuery(persisted_size);
SetFlushed(std::unique_lock<std::mutex>(mu_), std::move(result));
}

void OnQuery(StatusOr<std::int64_t> persisted_size) {
if (!persisted_size) return Resume(std::move(persisted_size).status());
return OnQuery(std::unique_lock<std::mutex>(mu_), *persisted_size);
return OnQuery(std::unique_lock<std::mutex>(mu_), *persisted_size,
/*is_resume=*/false);
}

auto ClearHandlers(std::unique_lock<std::mutex> const& /* lk */) {
Expand All @@ -365,7 +368,8 @@ class AsyncWriterConnectionResumedState
return tmp;
}

void OnQuery(std::unique_lock<std::mutex> lk, std::int64_t persisted_size) {
void OnQuery(std::unique_lock<std::mutex> lk, std::int64_t persisted_size,
bool is_resume = false) {
auto handle = impl_->WriteHandle();
if (handle) {
latest_write_handle_ = *std::move(handle);
Expand All @@ -384,7 +388,7 @@ class AsyncWriterConnectionResumedState
}
resend_buffer_.RemovePrefix(static_cast<std::size_t>(n));
buffer_offset_ = persisted_size;
if (state_ == State::kResuming) {
if (state_ == State::kResuming || is_resume) {
// Since the buffer has been modified to start exactly at the point of the
// resume, the next write on this new stream should start from the
// beginning of this truncated buffer.
Expand All @@ -402,8 +406,21 @@ class AsyncWriterConnectionResumedState
}
// If the buffer is small enough, collect all the handlers to notify them.
auto const handlers = ClearHandlersIfEmpty(lk);
if (is_resume) {
state_ = State::kIdle;
StartWriting(std::move(lk));
// The notifications are deferred until the lock is released, as they
// might call back and try to acquire the lock.
for (auto const& h : handlers) {
h->Execute(Status{});
}
return;
}
// SetFlushed will release the lock before returning.
SetFlushed(std::move(lk), Status{}, persisted_size);
// Re-acquire the lock to resume writing now that flush_ has been updated.
state_ = State::kIdle;
StartWriting(std::move(lk));
StartWriting(std::unique_lock<std::mutex>(mu_));
// The notifications are deferred until the lock is released, as they might
// call back and try to acquire the lock.
for (auto const& h : handlers) {
Expand Down Expand Up @@ -535,7 +552,7 @@ class AsyncWriterConnectionResumedState
options_, initial_request_, std::move(res->stream), hash_function_,
persisted_offset, false);
// OnQuery will restart the WriteLoop if necessary.
OnQuery(std::move(lk), persisted_offset);
OnQuery(std::move(lk), persisted_offset, /*is_resume=*/true);
}

void SetFinalized(std::unique_lock<std::mutex> lk,
Expand Down Expand Up @@ -565,7 +582,7 @@ class AsyncWriterConnectionResumedState
lk.unlock();
// Notify handlers and pending flushes *after* releasing the lock.
for (auto& h : handlers) h->Execute(Status{});
for (auto& pf : pending_flushes) pf.set_value(Status{}); // Success
for (auto& pf : pending_flushes) pf.p.set_value(Status{}); // Success
p.set_value(std::move(object)); // Set value on the moved promise
}

Expand All @@ -589,34 +606,30 @@ class AsyncWriterConnectionResumedState
lk.unlock();
// Notify handlers and pending flushes after releasing the lock.
for (auto& h : handlers) h->Execute(status);
for (auto& pf : pending_flushes) pf.set_value(status);
for (auto& pf : pending_flushes) pf.p.set_value(status);
p.set_value(std::move(status)); // Set value on the moved promise.
}

void SetFlushed(std::unique_lock<std::mutex> lk, Status const& result) {
void SetFlushed(std::unique_lock<std::mutex> lk, Status const& result,
std::int64_t persisted_size) {
if (!result.ok()) return SetError(std::move(lk), std::move(result));
// Do NOT reset finalize_ or finalizing_ here.
auto handlers = ClearHandlers(lk);
// Dequeue the promise corresponding to an explicit Flush() call, if any.
if (pending_flush_promises_.empty()) {
// This can happen if SetError cleared the queue first, or if this
// flush was triggered internally by buffer size (not by an explicit
// Flush() call) and thus has no promise in the queue.
flush_ = false;
lk.unlock();
for (auto& h : handlers) h->Execute(Status{});
return;
std::vector<promise<Status>> flushes_to_complete;
while (!pending_flush_promises_.empty() &&
pending_flush_promises_.front().target_offset <= persisted_size) {
flushes_to_complete.push_back(
std::move(pending_flush_promises_.front().p));
pending_flush_promises_.pop_front();
}
auto flushed = std::move(pending_flush_promises_.front());
pending_flush_promises_.pop_front();
if (pending_flush_promises_.empty()) {
flush_ = false;
}
lk.unlock(); // Unlock only once before notifying
// Notify handlers and the specific flush promise *after* releasing the
// Notify handlers and the specific flush promises *after* releasing the
// lock.
for (auto& h : handlers) h->Execute(Status{});
flushed.set_value(result);
for (auto& f : flushes_to_complete) f.set_value(result);
}

void SetError(std::unique_lock<std::mutex> lk, Status const& status) {
Expand Down Expand Up @@ -656,7 +669,7 @@ class AsyncWriterConnectionResumedState
for (auto& h : handlers) h->Execute(status);
// Set error on all pending flush promises.
for (auto& pf : pending_flushes) {
pf.set_value(status);
pf.p.set_value(status);
}
// Set error on the moved promises *once*.
if (complete_finalized) {
Expand Down Expand Up @@ -726,8 +739,17 @@ class AsyncWriterConnectionResumedState
// closed_.
future<Status> closed_future_;

// Tracks an outstanding `Flush()` promise alongside the stream offset at the
// time `Flush()` was called. The target offset is used to satisfy promises
// once all data buffered at the time of the `Flush()` call has been
// persisted.
struct PendingFlush {
promise<Status> p;
std::int64_t target_offset;
};

// Queue of promises for outstanding Flush() calls.
std::deque<promise<Status>> pending_flush_promises_;
std::deque<PendingFlush> pending_flush_promises_;

// The resend buffer. If there is an error, this will have all the data since
// the last persisted byte and will be resent.
Expand Down
Loading
Loading