From 50d122166ff0f0c6c7bfe42f4a17d986c0b76ced Mon Sep 17 00:00:00 2001 From: Vlad Ilyushchenko Date: Wed, 7 Oct 2026 19:49:29 +0100 Subject: [PATCH] Publish frame FSN before the frame becomes visible 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. --- .../qwp/client/sf/cursor/MmapSegment.java | 34 ++++- .../qwp/client/sf/cursor/SegmentRing.java | 49 +++++-- .../qwp/client/sf/cursor/SegmentRingTest.java | 128 ++++++++++++++++++ 3 files changed, 192 insertions(+), 19 deletions(-) diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/MmapSegment.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/MmapSegment.java index 779c99000..fe0795977 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/MmapSegment.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/MmapSegment.java @@ -695,6 +695,33 @@ MmapSegment successor() { * just a CRC pass and a memcpy into the mapped region. */ public long tryAppend(long payloadAddr, int payloadLen) { + long offset = tryWrite(payloadAddr, payloadLen); + if (offset != -1L) { + publishWritten(); + } + return offset; + } + + /** + * Makes every frame written by {@link #tryWrite} visible to the consumer: the + * I/O thread may read, and send, any frame below {@link #publishedOffset()}. + * Producer thread only. + */ + void publishWritten() { + // Until this volatile write retires, the consumer cannot see any of the + // bytes tryWrite wrote, and once it does, every one of them is complete. + publishedCursor = appendCursor; + } + + /** + * First half of {@link #tryAppend}: writes the frame and advances the append + * cursor and the frame count, but leaves {@link #publishedOffset()} where it + * was, so the consumer cannot see the frame until {@link #publishWritten}. + * {@link SegmentRing} publishes the frame's FSN in between, so the FSN is + * never behind a frame the I/O thread can send. Returns the frame's offset, + * or -1 with nothing written if it does not fit. Producer thread only. + */ + long tryWrite(long payloadAddr, int payloadLen) { if (payloadLen < 0) { throw new IllegalArgumentException("negative payloadLen: " + payloadLen); } @@ -720,11 +747,10 @@ public long tryAppend(long payloadAddr, int payloadLen) { // Plain read + write of the volatile field. `frameCount++` would // trip the "non-atomic increment of volatile" inspection, but // single-writer invariant (only the producer thread mutates) makes - // the RMW race-free by design. + // the RMW race-free by design. It precedes the publishedCursor store + // in publishWritten, so a consumer that sees the frame's bytes also + // sees the frame count that makes its FSN discoverable. frameCount = frameCount + 1; - // Publish last. Until this volatile write retires, the consumer - // cannot see any of the bytes we just wrote. - publishedCursor = appendCursor; return offset; } diff --git a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentRing.java b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentRing.java index 6fd7f5aaf..d18589e91 100644 --- a/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentRing.java +++ b/core/src/main/java/io/questdb/client/cutlass/qwp/client/sf/cursor/SegmentRing.java @@ -976,7 +976,9 @@ public long ackedFsn() { * server NACK with a bogus wireSeq cannot move {@code ackedFsn} past what * the producer has actually written. If we didn't clamp, the segment * manager could trim segments the I/O thread is still iterating and SEGV - * the JVM on the next {@code Unsafe.getInt} of an unmapped region. + * the JVM on the next {@code Unsafe.getInt} of an unmapped region. The + * clamp never drops a legitimate ACK: {@link #appendOrFsn} publishes a + * frame's FSN before the I/O thread can see, and send, the frame. * * @return {@code true} if the watermark advanced, {@code false} on * no-op (idempotent re-ack or clamped). Callers wishing to fire @@ -1009,8 +1011,9 @@ public boolean acknowledge(long seq) { */ public long appendOrFsn(long payloadAddr, int payloadLen) { checkDurability(); - long offset = active.tryAppend(payloadAddr, payloadLen); - if (offset == -1L) { + MmapSegment target = active; + boolean rotated = false; + if (target.tryWrite(payloadAddr, payloadLen) == -1L) { // Active is full. Try to rotate. MmapSegment spare = hotSpare; if (spare == null) { @@ -1020,7 +1023,7 @@ public long appendOrFsn(long payloadAddr, int payloadLen) { // range durable before the manifest can name its successor. The // manager performs the barrier; the producer uses the existing // backpressure path while it waits. - MmapSegment previous = active; + MmapSegment previous = target; if (requestSyncBeforeRotation(previous)) { wakeManager(); return BACKPRESSURE_NO_SPARE; @@ -1083,25 +1086,38 @@ public long appendOrFsn(long payloadAddr, int payloadLen) { if (wakeup != null) { wakeup.run(); } - offset = active.tryAppend(payloadAddr, payloadLen); - if (offset == -1L) { + target = spare; + if (target.tryWrite(payloadAddr, payloadLen) == -1L) { // Doesn't fit even in a fresh segment -- payload is genuinely too big. return PAYLOAD_TOO_LARGE; } - } else if (!wakeupRequestedForActive + rotated = true; + } + long fsn = nextSeq++; + // The frame is fully written but not yet visible to the I/O thread. + // Publish its FSN first and its bytes second. The I/O thread sends a + // frame the moment it lies below publishedOffset(), and the server can + // ACK it before this method returns; acknowledge() clamps at + // publishedFsn, so an FSN published after the bytes could clamp that + // ACK away. Server ACKs are cumulative and nothing re-delivers a lost + // one, so losing the last frame's ACK would stall close() and drain() + // until their timeouts. publishedFsn may run one frame ahead of + // publishedOffset() for an instant, but never behind it. + publishedFsn = fsn; + target.publishWritten(); + if (!rotated + && !wakeupRequestedForActive && hotSpare == null && managerWakeup != null - && active.publishedOffset() >= signalAtBytes) { + && target.publishedOffset() >= signalAtBytes) { // Backup signal: we're past the high-water mark and still don't // have a spare (manager hasn't caught up yet, or this is the very // first active and rotation hasn't fired the on-rotation wakeup). - // Fire once per active segment. + // Fire once per active segment, after the frame is fully published + // so the wakeup never delays its send. wakeupRequestedForActive = true; managerWakeup.run(); } - long fsn = nextSeq++; - // publishedFsn last so the I/O thread never observes a half-written frame. - publishedFsn = fsn; return fsn; } @@ -1470,9 +1486,12 @@ public long nextSeqHint() { } /** - * Highest FSN whose frame is fully written and visible to consumers (the - * I/O thread). Returns -1 when nothing has been appended yet. Volatile - * read; safe to call from any thread. + * Highest FSN whose frame is fully written. {@link #appendOrFsn} publishes it + * before the frame's bytes become visible below + * {@link MmapSegment#publishedOffset()}, so it covers every frame the I/O + * thread can send: it may run one frame ahead of the visible bytes for an + * instant, never behind them. Returns -1 when nothing has been appended yet. + * Volatile read; safe to call from any thread. */ public long publishedFsn() { return publishedFsn; diff --git a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentRingTest.java b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentRingTest.java index 4e79c6b1a..a192d6a8f 100644 --- a/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentRingTest.java +++ b/core/src/test/java/io/questdb/client/test/cutlass/qwp/client/sf/cursor/SegmentRingTest.java @@ -42,6 +42,7 @@ import java.nio.file.Paths; import java.util.Arrays; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import static org.junit.Assert.assertArrayEquals; @@ -1166,6 +1167,133 @@ public void testAcknowledgeClampsAtPublishedFsn() throws Exception { }); } + /** + * The I/O thread sends a frame as soon as its bytes lie below + * {@link MmapSegment#publishedOffset()}, and the server can ACK it before + * {@code appendOrFsn} returns. Because {@code acknowledge} clamps at + * {@code publishedFsn}, the ring must publish a frame's FSN no later than its bytes: + * an ACK clamped below a frame the I/O thread could already send is consumed and + * lost, and the server's cumulative ACKs never re-deliver it once the producer goes + * quiet. Losing the final frame's ACK this way would stall close() and drain() + * until their timeout. + *

+ * The high-water manager wakeup runs inside {@code appendOrFsn} once the appended + * frame's bytes are visible, so it observes the ring exactly where an I/O-thread + * ACK can race the publication: the wakeup plays the I/O thread and ACKs the + * newest frame visible to it. + */ + @Test + public void testAckOfFrameVisibleDuringAppendIsNotClampedAway() throws Exception { + TestUtils.assertMemoryLeak(() -> { + final int payloadLen = 32; + final long frameSize = MmapSegment.FRAME_HEADER_SIZE + payloadLen; + final long segSize = 64 * 1024; + long buf = Unsafe.malloc(payloadLen, MemoryTag.NATIVE_DEFAULT); + try { + fillPattern(buf, payloadLen, 0); + try (SegmentRing ring = new SegmentRing(MmapSegment.createInMemory(0L, segSize), segSize)) { + final long[] newestVisibleAtWakeup = {-1L}; + final long[] ackedAtWakeup = {-1L}; + final int[] wakeups = {0}; + ring.setManagerWakeup(() -> { + MmapSegment active = ring.getActive(); + long newestVisible = active.baseSeq() + + (active.publishedOffset() - MmapSegment.HEADER_SIZE) / frameSize - 1L; + ring.acknowledge(newestVisible); + newestVisibleAtWakeup[0] = newestVisible; + ackedAtWakeup[0] = ring.ackedFsn(); + wakeups[0]++; + }); + long fsn; + do { + fsn = ring.appendOrFsn(buf, payloadLen); + assertTrue("setup: append must succeed, got " + fsn, fsn >= 0); + } while (wakeups[0] == 0); + + assertEquals("setup: the high-water wakeup fires once per active", 1, wakeups[0]); + assertEquals("setup: the frame appended by the waking call is visible to the I/O thread", + fsn, newestVisibleAtWakeup[0]); + assertEquals("an ACK for a frame the I/O thread could already send must not be clamped away", + fsn, ackedAtWakeup[0]); + assertEquals(fsn, ring.ackedFsn()); + } + } finally { + Unsafe.free(buf, payloadLen, MemoryTag.NATIVE_DEFAULT); + } + }); + } + + /** + * Concurrent form of {@link #testAckOfFrameVisibleDuringAppendIsNotClampedAway}: an + * observer thread plays the I/O thread against a producer that appends across many + * rotations, ACKing each frame the moment its bytes become visible below + * {@link MmapSegment#publishedOffset()}. Every such ACK must land; none may be + * clamped below the frame the observer could already have sent. + */ + @Test(timeout = 60_000L) + public void testAckOfFrameVisibleToConsumerLandsUnderConcurrentAppends() throws Exception { + TestUtils.assertMemoryLeak(() -> { + final int payloadLen = 16; + final long frameSize = MmapSegment.FRAME_HEADER_SIZE + payloadLen; + final long segSize = MmapSegment.HEADER_SIZE + frameSize * 64; + final long frameCount = 200_000L; + long buf = Unsafe.malloc(payloadLen, MemoryTag.NATIVE_DEFAULT); + try { + fillPattern(buf, payloadLen, 0); + try (SegmentRing ring = new SegmentRing(MmapSegment.createInMemory(0L, segSize), segSize)) { + final AtomicBoolean producerDone = new AtomicBoolean(); + final AtomicReference failure = new AtomicReference<>(); + Thread observer = new Thread(() -> { + try { + long lastAcked = -1L; + while (lastAcked < frameCount - 1L && failure.get() == null) { + // Read before the snapshot: once the producer is done, its + // every append happened-before this read, so the snapshot + // below sees all frames and a caught-up observer can stop. + boolean isProducerDone = producerDone.get(); + MmapSegment active = ring.getActive(); + long newestVisible = active.baseSeq() + + (active.publishedOffset() - MmapSegment.HEADER_SIZE) / frameSize - 1L; + if (newestVisible > lastAcked) { + ring.acknowledge(newestVisible); + long acked = ring.ackedFsn(); + if (acked < newestVisible) { + throw new AssertionError("ACK of visible FSN " + newestVisible + + " was clamped to " + acked); + } + lastAcked = newestVisible; + } else if (isProducerDone) { + break; + } + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } + }, "ack-observer"); + observer.start(); + try { + for (long i = 0; i < frameCount && failure.get() == null; i++) { + long fsn; + while ((fsn = ring.appendOrFsn(buf, payloadLen)) == SegmentRing.BACKPRESSURE_NO_SPARE) { + ring.installHotSpare(MmapSegment.createInMemory(ring.nextSeqHint(), segSize)); + } + assertEquals(i, fsn); + } + } catch (Throwable t) { + failure.compareAndSet(null, t); + } finally { + producerDone.set(true); + observer.join(); + } + assertNull("visible frame ACK lost: " + failure.get(), failure.get()); + assertEquals(frameCount - 1L, ring.ackedFsn()); + } + } finally { + Unsafe.free(buf, payloadLen, MemoryTag.NATIVE_DEFAULT); + } + }); + } + @Test public void testNextSealedAfterWalksThousandsOfSegmentsInLinearOperations() throws Exception { TestUtils.assertMemoryLeak(() -> {