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