Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand All @@ -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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand All @@ -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;
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
* <p>
* 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<Throwable> 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(() -> {
Expand Down
Loading