Repository navigation
fix(qwp): fix sender close() and drain() timing out after the server acknowledged all data - #106
Closed
bluestreak01 wants to merge 1 commit into
Closed
bluestreak01 wants to merge 1 commit into
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.
Member
Author
|
Superseded by #107: same commit (50d1221), on branch |
Contributor
[PR Coverage check]😍 pass : 18 / 18 (100.00%) file detail
|
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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.