Skip to content

fix(storage): mark async writer idle before invoking flush callbacks - #16504

Open
kalragauri wants to merge 3 commits into
googleapis:mainfrom
kalragauri:fix/flush-order
Open

kalragauri wants to merge 3 commits into
googleapis:mainfrom
kalragauri:fix/flush-order

Conversation

@kalragauri

@kalragauri kalragauri commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

When a Flush() completes on AsyncWriterConnectionResumed or AsyncWriterConnectionBuffered, OnQuery() delegates to SetFlushed() to satisfy pending flush promises (pending_flush_promises_) and notify handlers waiting for the buffer to shrink (flush_handlers_).

Previously, SetFlushed() satisfied the flush promises while the writer was still marked busy (state_ == State::kWriting in WriterConnectionResumed, writing_ == true in WriterConnectionBuffered), and OnQuery() only restarted the write loop after SetFlushed() returned. As a result, any operation (Flush, Write, Finalize, or Close) chained synchronously from a Flush().then(...) continuation observed a busy writer and was queued rather than dispatched immediately before the callback returned.

Running callbacks while the writer is idle also exposed two related issues in how flush_ and flush_handlers_ were updated:

  • A small Write() could cancel a queued flush. HandleNewData() recalculated flush_ from the buffer size alone, so a Write() below the low watermark issued while another Flush() was still pending (including one chained from a callback) set flush_ = false. The queued flush was then sent as a plain write, and its future was never satisfied.
  • Write() calls blocked at the high watermark were released too early. SetFlushed() released every flush_handlers_ entry after any flush, even when the buffer was still at or above the low watermark. It also cleared flush_ once no Flush() calls were pending, regardless of how much data remained buffered.

Changes

  • SetFlushed() (writer_connection_resumed.cc, writer_connection_buffered.cc):
    1. Collect the flush_handlers_ to release and the satisfied pending_flush_promises_, update flush_, and mark the writer idle (state_ = State::kIdle / writing_ = false) while holding mu_.
    2. Unlock mu_ and invoke the handlers and flush promises before restarting the write loop, so chained operations see an idle writer and are dispatched immediately, and earlier flush futures cannot be overtaken by a queued flush that completes inline.
    3. Call StartWriting() after callbacks finish to drain any remaining buffered data (a no-op if a callback already restarted the write loop).
    4. Release flush_handlers_ only once the buffer is below the low watermark or empty (ClearHandlersIfEmpty() instead of ClearHandlers()), and set flush_ = NeedsFlush(lk) so the writer keeps flushing until those handlers can be released.
  • OnQuery() (is_resume branch): Apply the same sequence when notifying flush_handlers_ on stream resume.
  • NeedsFlush() (new helper): Returns true while Flush() calls are pending, or while the buffer is non-empty and at or above the low watermark. HandleNewData() and SetFlushed() both use it, so a small Write() can no longer clear flush_ while a flush is pending.
  • ClearHandlersIfEmpty(): An empty buffer now always releases flush_handlers_. With a low watermark of 0, the previous size >= lwm check was always true, so the handlers were never released.

Call Flow: Before vs. After

Before this PR

sequenceDiagram
    participant Impl as Underlying Stream (impl_)
    participant Writer as Buffered / Resumed Writer
    participant User as User f1.then() Callback

    Impl->>Writer: Flush(data1) completes -> OnQuery() -> SetFlushed()
    Note over Writer: Writer is still busy (kWriting / writing_ == true)<br/>Unlocks mu_
    Writer->>User: f1.set_value(OK)
    activate User
    User->>Writer: Flush(data2)
    Note over Writer: HandleNewData() queues data2<br/>StartWriting() sees writer busy -> returns early
    Writer-->>User: returns future f2 (data2 NOT dispatched yet)
    User-->>Writer: callback returns
    deactivate User
    Note over Writer: OnQuery() restarts the write loop after SetFlushed() returns
    Writer->>Impl: impl_->Flush(data2) (dispatched after callback returns)
Loading

After this PR

sequenceDiagram
    participant Impl as Underlying Stream (impl_)
    participant Writer as Buffered / Resumed Writer
    participant User as User f1.then() Callback

    Impl->>Writer: Flush(data1) completes -> OnQuery() -> SetFlushed()
    Note over Writer: Marks writer IDLE under mu_ (kIdle / writing_ = false)<br/>Unlocks mu_
    Writer->>User: f1.set_value(OK)
    activate User
    User->>Writer: Flush(data2)
    Note over Writer: HandleNewData() queues data2<br/>StartWriting() sees writer idle -> dispatches immediately
    Writer->>Impl: impl_->Flush(data2) (in flight before callback returns)
    Writer-->>User: returns future f2
    User-->>Writer: callback returns
    deactivate User
    Note over Writer: StartWriting() sees writer busy -> no-op
Loading

@product-auto-label product-auto-label Bot added the api: storage Issues related to the Cloud Storage API. label Sep 29, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Code Review

This pull request refactors the asynchronous buffered and resumed writer connections to mark the writer as idle under the lock before releasing it and notifying handlers, ensuring chained operations are dispatched immediately. It also adds corresponding unit tests to verify these callback behaviors. The review feedback points out a potential data race and crash in the newly added cross-thread flush tests, where worker.join() can be called before the continuation thread finishes initializing the worker thread object, and suggests waiting for the callback to complete before joining.

@kalragauri
kalragauri marked this pull request as ready for review September 29, 2026 11:44
@kalragauri
kalragauri requested review from a team as code owners September 29, 2026 11:44
@kalragauri
kalragauri requested a review from v-pratap September 29, 2026 11:44
@codecov

codecov Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 99.88426% with 1 line in your changes missing coverage. Please review.
✅ Project coverage is 92.38%. Comparing base (54869ea) to head (c0f601c).

Files with missing lines Patch % Lines
...e/internal/async/writer_connection_resumed_test.cc 99.77% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #16504      +/-   ##
==========================================
+ Coverage   92.34%   92.38%   +0.03%     
==========================================
  Files        2262     2262              
  Lines      217173   218014     +841     
==========================================
+ Hits       200549   201404     +855     
+ Misses      16624    16610      -14     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

Comment thread google/cloud/storage/internal/async/writer_connection_resumed.cc Outdated
Comment thread google/cloud/storage/internal/async/writer_connection_resumed.cc

This branch was successfully deployed

1 active deployment
false — c0f601c0 Deployed Oct 5, 2026 by kalragauri via Save PR ref #12262
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

api: storage Issues related to the Cloud Storage API.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants