Skip to content

fix(qwp): fix sender close() and drain() timing out after the server acknowledged all data - #107

Open
bluestreak01 wants to merge 1 commit into
mainfrom
fix/qwp-exec-done-negative-rows-affected
Open

bluestreak01 wants to merge 1 commit into
mainfrom
fix/qwp-exec-done-negative-rows-affected

Conversation

@bluestreak01

Copy link
Copy Markdown
Member

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() and drain() waited for their full timeout and then reported unacknowledged data, although the server had committed every row. QuestDB CI hit this as a 300 s close() drain timeout in QwpSenderE2ETest.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 in publishedFsn. The I/O thread sends a frame as soon as it lies below publishedOffset(), 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 at publishedFsn, 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-private tryWrite(), which writes the frame and bumps the frame count without publishing it, and publishWritten(), which advances publishedCursor. tryAppend() still does both, so its other callers behave as before.
  • appendOrFsn() calls tryWrite(), stores publishedFsn, then calls publishWritten(). Every frame the I/O thread can send is already covered by publishedFsn, so the clamp no longer drops a legitimate ACK. The clamp itself stays, as the defense-in-depth against bogus server sequence numbers that testAcknowledgeClampsAtPublishedFsn pins.
  • appendOrFsn() runs the high-water manager wakeup after publication rather than between the two stores.

Tradeoffs

  • publishedFsn now 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 uses frameCount and publishedOffset(), 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.
  • A bogus server ACK arriving in that instant could cover a frame that is fully written but not yet sent. The clamp never guarded against acknowledging written-but-unsent frames, and the I/O thread's own clamp to frames it actually sent still applies, so this adds no exposure.
  • The change reorders the same three volatile stores and adds no work to the append path. 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

  • New 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.
  • New 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.
  • Full client test suite: 3,477 tests pass; the 7 skips are @Ignore or OS-gated.
  • QuestDB cutlass/qwp tests against this client: 1,555 pass, including QwpSenderE2ETest 137 of 137.
  • With a 20 ms producer pause injected at the old race point, 16 of 16 QuestDB senders hung on the previous code. With the same pause injected between the two stores of the new code, or right after them, 0 of 48 hung and every row arrived.
  • Compiles at Java 8 language level (class file version 52). Not built on a real JDK 8 locally.

Tandem

Pairs with questdb/questdb#7770, which bumps the java-questdb-client submodule to this branch. That PR is where QuestDB CI hit the hang.

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.
@bluestreak01 bluestreak01 added bug Something isn't working QWP tandem labels Oct 7, 2026
bluestreak01 added a commit to questdb/questdb that referenced this pull request Oct 7, 2026
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.
@mtopolnik

Copy link
Copy Markdown
Contributor

[PR Coverage check]

😍 pass : 18 / 18 (100.00%)

file detail

path covered line new line coverage
🔵 io/questdb/client/cutlass/qwp/client/sf/cursor/MmapSegment.java 6 6 100.00%
🔵 io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentRing.java 12 12 100.00%

@bluestreak01

Copy link
Copy Markdown
Member Author

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 6abd2381 (base d37022c7), OSS a5c0d3a4 (merge-base 6c635564), client 50d12216 (base a7e7db3d). CI on all three PRs was still pending when I reviewed them.

Submodule scope

  • questdb (ENT): points at a commit on the OSS PR branch, so the OSS fix commits 5fa772e234..a5c0d3a475 were reviewed. The same pointer move also picks up 6c63556431 (#7706) and a05e90c19e (#7736, which moves the client from 0b9b5766 to a7e7db3d). Both are already on OSS master, so I treated them as out of scope. No code in these diffs calls into them.
  • java-questdb-client (OSS): points at a commit on the client PR branch, so fix(qwp): fix sender close() and drain() timing out after the server acknowledged all data #107 was reviewed.

PR titles and descriptions: titles follow Conventional Commits with user-facing descriptions, and the labels match. Nothing to fix.

Findings

None at Critical, Moderate or Minor.

Test coverage

The test gate passes with no coverage gaps. Runs at the reviewed revisions:

  • Client SegmentRingTest:
    • The two new tests and testAcknowledgeClampsAtPublishedFsn pass at head.
    • With SegmentRing.java and MmapSegment.java reverted to a7e7db3d in a scratch worktree, both new tests fail (expected:<1228> but was:<1227>; ACK of visible FSN 49 was clamped to 48). The concurrent test failed 10 of 10 runs.
    • With only the two stores swapped back, keeping the wakeup after both, the single-threaded test passes. The concurrent test still failed 10 of 10 runs on a multi-core host, so it does guard the store order.
  • OSS and ENT tests: QwpEgressDdlExecTest (10 of 10), QwpEgressAclExecDoneTest (1 of 1) and ColdStorageTestUtilsTest (3 of 3) pass at head, built with -P local-client.
  • Not run against base: the OSS and ENT egress tests. The failure without the fix is clear from source: base CompiledQueryImpl only ever assigns affectedRowsCount = -1, and executeDdl and assertExecDone both assert 0L.

Summary

  • Verdict: ENT approve, OSS approve, client approve.
  • Severity count: 0 Critical, 0 Moderate, 0 Minor; 0 in the diff, 0 breaking unchanged callers.
  • Callers checked: every caller of the code each PR changes.
    • OSS and ENT: the removed CompiledQuery.getAffectedRowsCount() has no remaining callers or implementors.
    • Client: every reader of publishedFsn, publishedOffset() and frameCount, and every acknowledge caller.
    • ENT cold storage: all users of listBucketRelative() and bucketHasDataParquet().

Behaviour changes to be aware of (deliberate, not defects):

  1. Java and .NET clients now see 0 instead of -1 for statements the compiler runs at parse time (TRUNCATE, SET, the ACL statements and similar). The OSS PR documents this, and it matches the client javadoc ("0 for pure DDL"). Nothing in OSS or ENT main code, or in the pinned client, branches on -1.
  2. publishedFsn can briefly run one frame ahead of the bytes the I/O thread can see. Nothing relies on the old order:
    • The send path and reconnect positioning use publishedOffset() and frameCount, whose relative order is unchanged.
    • The I/O thread reads publishedFsn only to bound error-report ranges. The other readers run on the producer thread or with no concurrent producer.
    • ACKs are capped at the last sent frame (nextWireSeq - 1) before reaching the ring, so a real ACK can no longer be clamped away.
  3. Store-and-forward rules: the client change adds no error surfacing, reconnect budget or hard failure. It removes a close()/drain() timeout caused by a dropped ACK.

The scratch worktree has been removed, and all three checkouts are clean at the reviewed SHAs.

@bluestreak01 bluestreak01 added the QUEUED FOR MERGE Approved PR in the merge queue. Do not merge master into this PR. label Oct 7, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working QUEUED FOR MERGE Approved PR in the merge queue. Do not merge master into this PR. QWP tandem

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants