Skip to content
Open
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
74 changes: 53 additions & 21 deletions google/cloud/storage/internal/async/writer_connection_buffered.cc
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,11 @@ class AsyncWriterConnectionBufferedState
bool flush = false) {
if (!resume_status_.ok()) return make_ready_future(resume_status_);
auto const buffer_size = resend_buffer_.size();
flush_ = (buffer_size >= buffer_size_lwm_) || flush;
// `NeedsFlush()` stays true while `Flush()` calls are pending, so a small
// `Write()`, possibly chained from a callback, cannot turn a queued flush
// into a plain write. The explicit low-water mark check preserves the
// existing behavior for an empty `Write()` when the low-water mark is 0.
flush_ = flush || NeedsFlush(lk) || buffer_size >= buffer_size_lwm_;
auto result = make_ready_future(Status{});
if (buffer_size >= buffer_size_hwm_) {
auto p = promise<Status>();
Expand Down Expand Up @@ -345,11 +349,29 @@ class AsyncWriterConnectionBufferedState

auto ClearHandlersIfEmpty(std::unique_lock<std::mutex> const& /* lk */) {
decltype(flush_handlers_) tmp;
if (resend_buffer_.size() >= buffer_size_lwm_) return tmp;
// Release the waiters once the buffer is empty or below the low-water
// mark. The emptiness check matters when the low-water mark is 0.
if (!resend_buffer_.empty() && resend_buffer_.size() >= buffer_size_lwm_) {
return tmp;
}
flush_handlers_.swap(tmp);
return tmp;
}

/**
* Returns true if the write loop must send data with `Flush()`.
*
* That is the case while `Flush()` calls are pending, or while the buffer is
* non-empty and at or above the low-water mark. An empty buffer only needs a
* flush if a `Flush()` is pending, which avoids an endless loop of empty
* flushes when the low-water mark is 0.
*/
bool NeedsFlush(std::unique_lock<std::mutex> const& /* lk */) const {
return !pending_flush_promises_.empty() ||
(!resend_buffer_.empty() &&
resend_buffer_.size() >= buffer_size_lwm_);
}

void OnQuery(std::unique_lock<std::mutex> lk, std::int64_t persisted_size,
bool is_resume = false) {
if (persisted_size < buffer_offset_) {
Expand Down Expand Up @@ -383,28 +405,28 @@ class AsyncWriterConnectionBufferedState
write_offset_ -= static_cast<std::size_t>(n);
}
}
// If the buffer is small enough, collect all the handlers to notify them.
auto const handlers = ClearHandlersIfEmpty(lk);
if (is_resume) {
// We are resuming. The pending flush promises (if any) should not be
// satisfied yet, because we haven't actually flushed the data on the new
// connection. The `WriteLoop` will trigger a flush (potentially empty)
// if `flush_` is still true, which will satisfy the promises when it
// completes. However, we still need to notify any handlers waiting for
// the buffer to shrink, and we need to restart the write loop.
// completes.
auto const handlers = ClearHandlersIfEmpty(lk);
// Mark the writer idle under the lock so any operation chained from a
// handler below sees an idle writer and is dispatched immediately, and
// `writing_` is never modified without holding `mu_`.
resuming_ = false;
lk.unlock();
writing_ = false;
lk.unlock(); // Release lock before notifying
// 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{});
WriteLoop(std::unique_lock<std::mutex>(mu_));
// Re-acquire the lock to restart the write loop. This is a no-op if a
// handler above already restarted it.
StartWriting(std::unique_lock<std::mutex>(mu_));
return;
}
// SetFlushed will release the lock before returning.
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
// call back and try to acquire the lock.
for (auto const& h : handlers) h->Execute(Status{});
}

void WriteStep(std::unique_lock<std::mutex> lk, absl::Cord payload) {
Expand Down Expand Up @@ -564,22 +586,32 @@ class AsyncWriterConnectionBufferedState
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);
// Only release the high-water mark waiters once the buffer drops below the
// low-water mark (or is empty).
auto handlers = ClearHandlersIfEmpty(lk);
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();
}
if (pending_flush_promises_.empty()) {
flush_ = false;
}
lk.unlock(); // Unlock only once before notifying
// Notify handlers and the specific flush promises *after* releasing the
// lock.
// Keep flushing while `NeedsFlush()` holds, so any high-water mark waiters
// not released above are released by a later flush.
flush_ = NeedsFlush(lk);
// Mark the writer idle under the lock so any operation chained from a
// callback below sees an idle writer and is dispatched immediately, and
// `writing_` is never modified without holding `mu_`.
writing_ = false;
lk.unlock(); // Release lock before notifying
// Notify handlers and satisfied flush promises before restarting the
// write loop so callbacks cannot be overtaken by a queued flush that
// completes inline.
for (auto& h : handlers) h->Execute(Status{});
for (auto& f : flushes_to_complete) f.set_value(result);
// Re-acquire the lock to resume writing any remaining buffered data.
// This is a no-op if a callback above already restarted the write loop.
StartWriting(std::unique_lock<std::mutex>(mu_));
}

void SetError(std::unique_lock<std::mutex> lk, Status const& status) {
Expand Down
Loading
Loading