fix(storage): mark async writer idle before invoking flush callbacks - #16504
kalragauri wants to merge 3 commits into
Conversation
There was a problem hiding this comment.
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.
Codecov Report❌ Patch coverage is
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. 🚀 New features to boost your workflow:
|
64d9ee4 to
c0f601c
Compare
When a
Flush()completes onAsyncWriterConnectionResumedorAsyncWriterConnectionBuffered,OnQuery()delegates toSetFlushed()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::kWritinginWriterConnectionResumed,writing_ == trueinWriterConnectionBuffered), andOnQuery()only restarted the write loop afterSetFlushed()returned. As a result, any operation (Flush,Write,Finalize, orClose) chained synchronously from aFlush().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_andflush_handlers_were updated:Write()could cancel a queued flush.HandleNewData()recalculatedflush_from the buffer size alone, so aWrite()below the low watermark issued while anotherFlush()was still pending (including one chained from a callback) setflush_ = 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 everyflush_handlers_entry after any flush, even when the buffer was still at or above the low watermark. It also clearedflush_once noFlush()calls were pending, regardless of how much data remained buffered.Changes
SetFlushed()(writer_connection_resumed.cc,writer_connection_buffered.cc):flush_handlers_to release and the satisfiedpending_flush_promises_, updateflush_, and mark the writer idle (state_ = State::kIdle/writing_ = false) while holdingmu_.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.StartWriting()after callbacks finish to drain any remaining buffered data (a no-op if a callback already restarted the write loop).flush_handlers_only once the buffer is below the low watermark or empty (ClearHandlersIfEmpty()instead ofClearHandlers()), and setflush_ = NeedsFlush(lk)so the writer keeps flushing until those handlers can be released.OnQuery()(is_resumebranch): Apply the same sequence when notifyingflush_handlers_on stream resume.NeedsFlush()(new helper): Returns true whileFlush()calls are pending, or while the buffer is non-empty and at or above the low watermark.HandleNewData()andSetFlushed()both use it, so a smallWrite()can no longer clearflush_while a flush is pending.ClearHandlersIfEmpty(): An empty buffer now always releasesflush_handlers_. With a low watermark of 0, the previoussize >= lwmcheck 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)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