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 779c9900..fe079597 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 6fd7f5aa..d18589e9 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 4e79c6b1..a192d6a8 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(() -> {