Repository navigation
fix(qwp): fix sender close() and drain() timing out after the server acknowledged all data - #107
bluestreak01 wants to merge 1 commit into
Conversation
SegmentRing.appendOrFsn() made a frame's bytes visible to the I/O thread (MmapSegment.publishedCursor) before it stored the frame's FSN in publishedFsn. The I/O thread sends a frame as soon as it lies below publishedOffset(), so the server could commit and ACK it while the producer still sat between the two stores, for example when the OS preempted the producer there. acknowledge() clamps every ACK at publishedFsn, so it capped that ACK one frame short and the I/O thread discarded it. QWP server ACKs are cumulative and only follow new frames, so nothing re-delivered it once the producer went quiet: after the final frame, close() and drain() waited out their full timeout although the server had committed every row. QuestDB CI hit this as a 300 s close() drain timeout in QwpSenderE2ETest. MmapSegment.tryAppend() now splits into tryWrite(), which writes the frame and bumps the frame count without publishing it, and publishWritten(), which advances publishedCursor. appendOrFsn() calls tryWrite(), stores publishedFsn, then calls publishWritten(), so publishedFsn covers every frame the I/O thread can send and the clamp never drops a legitimate ACK. publishedFsn may now run one frame ahead of the visible bytes for an instant, never behind them; no concurrent reader depends on the old direction. The clamp stays as the defense-in-depth against bogus server sequence numbers. appendOrFsn() also runs the high-water manager wakeup after publication instead of between the two stores, so the wakeup no longer delays the frame's FSN. tryAppend() keeps its behaviour for every other caller, and the change adds no work to the append path: it reorders the same volatile stores. Two SegmentRingTest tests reproduce the lost ACK, one deterministically through the high-water wakeup and one with a concurrent observer that ACKs each frame as soon as it becomes visible. Both fail on the previous code and pass with this change.
Points the submodule at java-questdb-client 50d12216 (questdb/java-questdb-client#107). The QWP sender's segment ring made a frame's bytes visible to its I/O thread before it published the frame's FSN, so an ACK the server sent inside that window was clamped one frame short and dropped. After the final frame nothing re-delivered it, and close() and drain() waited out their full timeout although the server had committed every row. This branch's linux-other CI leg hit it as a 300 s close() drain timeout in QwpSenderE2ETest.testConcurrentSenders_sameTable_doubleToDecimal. The client now publishes the FSN before the frame becomes visible. Once the client PR merges, this pointer moves to its squashed commit on main.
[PR Coverage check]😍 pass : 18 / 18 (100.00%) file detail
|
|
I found no blocking issues in any of the three PRs, and nothing at Moderate or Minor either. ENT, OSS and client are all approve. Reviewed PR #1265 at level 3 in tandem mode, with questdb/questdb#7770 and the nested #107. Revisions: ENT Submodule scope
PR titles and descriptions: titles follow Conventional Commits with user-facing descriptions, and the labels match. Nothing to fix. FindingsNone at Critical, Moderate or Minor. Test coverageThe test gate passes with no coverage gaps. Runs at the reviewed revisions:
Summary
Behaviour changes to be aware of (deliberate, not defects):
The scratch worktree has been removed, and all three checkouts are clean at the reviewed SHAs. |
The QWP WebSocket sender could lose the server's acknowledgement (ACK) of a frame it had just sent. When that frame was the last one before the sender went quiet,
close()anddrain()waited for their full timeout and then reported unacknowledged data, although the server had committed every row. QuestDB CI hit this as a 300 sclose()drain timeout inQwpSenderE2ETest.testConcurrentSenders_sameTable_doubleToDecimal(build 276712). The server's DEBUG log from that run shows it committed the final frame and sent its ACK 64 µs after receiving it.Cause
SegmentRing.appendOrFsn()made a frame's bytes visible to the I/O thread (MmapSegment.publishedCursor) before it stored the frame's FSN inpublishedFsn. The I/O thread sends a frame as soon as it lies belowpublishedOffset(), so the server can commit and ACK it while the producer is still between the two stores, for example when the OS preempts the producer there.acknowledge()clamps every ACK atpublishedFsn, so it capped that ACK one frame short and the I/O thread discarded it. Server ACKs are cumulative and only follow new frames: mid-stream the next frame's ACK covers the loss, but after the final frame nothing re-delivers it.The window is a few instructions wide, so the hang needs the producer descheduled at exactly that point. Across 305 QuestDB PR builds and 51 macwin builds since 2026-09-07 it appeared once. The defect dates from the original store-and-forward implementation (#17).
Change
MmapSegment.tryAppend()splits into package-privatetryWrite(), which writes the frame and bumps the frame count without publishing it, andpublishWritten(), which advancespublishedCursor.tryAppend()still does both, so its other callers behave as before.appendOrFsn()callstryWrite(), storespublishedFsn, then callspublishWritten(). Every frame the I/O thread can send is already covered bypublishedFsn, so the clamp no longer drops a legitimate ACK. The clamp itself stays, as the defense-in-depth against bogus server sequence numbers thattestAcknowledgeClampsAtPublishedFsnpins.appendOrFsn()runs the high-water manager wakeup after publication rather than between the two stores.Tradeoffs
publishedFsnnow leads the visible bytes by at most one frame for an instant, where it previously lagged them; its javadoc says so. The readers that run concurrently with the producer are the ACK clamp and the error-range reporting. Cursor positioning usesframeCountandpublishedOffset(), whose relative order is unchanged. The producer itself, recovery, engine close and the orphan drainer read it without a concurrent producer and see no difference.CursorEngineAppendLatencyBenchmark(interleaved runs, pinned cores) showed no difference beyond run-to-run noise: 8.25 M vs 8.04 M appends/s with 64-byte payloads (8 runs each, overlapping ranges) and 353.3 k vs 352.9 k appends/s with 4 KiB payloads (10 runs each, about 1% standard deviation), new vs old.Test plan
SegmentRingTest.testAckOfFrameVisibleDuringAppendIsNotClampedAway: the high-water wakeup, which observes the ring once the appended frame is visible, plays the I/O thread and ACKs the newest visible frame. It fails on the previous code every run and passes with the change.SegmentRingTest.testAckOfFrameVisibleToConsumerLandsUnderConcurrentAppends: an observer thread ACKs each frame as soon as it becomes visible while the producer appends 200,000 frames across rotations. It failed 12 of 12 runs on the previous code and passed 50 of 50 with the change, 30 of them restricted to two CPUs.@Ignoreor OS-gated.cutlass/qwptests against this client: 1,555 pass, includingQwpSenderE2ETest137 of 137.Tandem
Pairs with questdb/questdb#7770, which bumps the
java-questdb-clientsubmodule to this branch. That PR is where QuestDB CI hit the hang.