diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 0455f0efa8bb6..7a2d4953f84a2 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -764,6 +764,29 @@ default ManagedLedgerAttributes getManagedLedgerAttributes() { void asyncReadEntry(Position position, AsyncCallbacks.ReadEntryCallback callback, Object ctx); + /** + * Create a standalone cursorless streaming reader (cache-friendly, default). Equivalent to + * {@link #newRandomReader(boolean) newRandomReader(true)}. + * + * @throws UnsupportedOperationException when the managed-ledger implementation does not support random reads + */ + default RandomReader newRandomReader() { + return newRandomReader(true); + } + + /** + * Create a standalone cursorless reader with the requested cache strategy. + * + *
{@code streaming=true} admits read misses to the cache and seeds the write tail for high hit rates. + * {@code streaming=false} serves cache hits but neither writes back misses nor seeds the tail, so + * low-frequency random reads do not pollute the cache. + * + * @throws UnsupportedOperationException when the managed-ledger implementation does not support random reads + */ + default RandomReader newRandomReader(boolean streaming) { + throw new UnsupportedOperationException("RandomReader is not supported by this ManagedLedger implementation"); + } + /** * Get all the managed ledgers. */ diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java new file mode 100644 index 0000000000000..e2c03983ce549 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger; + +import java.io.Closeable; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import org.apache.bookkeeper.common.annotation.InterfaceAudience; +import org.apache.bookkeeper.common.annotation.InterfaceStability; +import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; + +/** + * A cursorless, stateless reader for entries that are currently available in a managed ledger. + * + *
A random reader does not maintain a read position, acknowledge entries, contribute to backlog, or prevent ledger + * trimming. A read can therefore fail, or start at the next retained ledger, when ledgers are trimmed concurrently. + * Reads do not wait for future entries. + * + *
Callers must release every returned {@link Entry}. Completion can run on a BookKeeper, Netty, or managed-ledger
+ * thread; callers that mutate thread-confined state must explicitly select an appropriate executor.
+ */
+@InterfaceAudience.LimitedPrivate
+@InterfaceStability.Evolving
+public interface RandomReader extends Closeable {
+
+ /**
+ * Read up to {@code numberOfEntries} starting at {@code startPosition}, inclusive.
+ */
+ CompletableFuture {@code maxPosition} is inclusive. A null value is equivalent to {@link PositionFactory#LATEST}.
+ * {@code maxSizeBytes} uses the same estimate-based cap as {@link ManagedCursor}; at least one entry can be
+ * returned even when that entry exceeds the requested size.
+ *
+ * If a storage error occurs after entries have been collected, the future completes successfully with that
+ * partial list and the read stops. An error before the first entry completes the future exceptionally.
+ */
+ CompletableFuture> read(Position startPosition, int numberOfEntries);
+
+ /**
+ * Read up to {@code maxPosition}, inclusive, without a size limit.
+ */
+ default CompletableFuture
> read(Position startPosition, int numberOfEntries, Position maxPosition) {
+ return read(startPosition, numberOfEntries, maxPosition, ManagedLedgerUtils.NO_MAX_SIZE_LIMIT);
+ }
+
+ /**
+ * Read using an estimated-size limit and no position limit.
+ */
+ default CompletableFuture
> read(Position startPosition, int numberOfEntries, long maxSizeBytes) {
+ return read(startPosition, numberOfEntries, PositionFactory.LATEST, maxSizeBytes);
+ }
+
+ /**
+ * Read entries subject to count, position, and estimated-size limits.
+ *
+ *
> read(Position startPosition, int numberOfEntries, Position maxPosition,
+ long maxSizeBytes);
+
+ /**
+ * Unregister this reader. Closing does not cancel reads already in progress.
+ */
+ @Override
+ void close();
+}
diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
index 26fdb458a21b2..d868a9f5fcaff 100644
--- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
+++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java
@@ -122,6 +122,7 @@
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.PositionBound;
import org.apache.bookkeeper.mledger.PositionFactory;
+import org.apache.bookkeeper.mledger.RandomReader;
import org.apache.bookkeeper.mledger.WaitingEntryCallBack;
import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback;
import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback;
@@ -190,6 +191,7 @@ public Logger getLogger() {
// ordered by read position (when cacheEvictionByMarkDeletedPosition=false) or by mark delete position
// (when cacheEvictionByMarkDeletedPosition=true)
private final ActiveManagedCursorContainer activeCursors;
+ private final RandomReaders randomReaders;
// Ever-increasing counter of entries added
@@ -393,6 +395,7 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper
} else {
activeCursors = new ManagedCursorContainerImpl();
}
+ randomReaders = new RandomReaders(this);
this.factory = factory;
this.bookKeeper = bookKeeper;
this.config = config;
@@ -1650,11 +1653,13 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) {
public synchronized void asyncClose(final CloseCallback callback, final Object ctx) {
State state = STATE_UPDATER.get(this);
if (state.isFenced()) {
+ randomReaders.closeAll();
cancelScheduledTasks();
factory.close(this);
callback.closeFailed(new ManagedLedgerFencedException(), ctx);
return;
} else if (state == State.Closed) {
+ randomReaders.closeAll();
log.debug("Ignoring request to close a closed managed ledger");
callback.closeComplete(ctx);
return;
@@ -1664,6 +1669,7 @@ public synchronized void asyncClose(final CloseCallback callback, final Object c
factory.close(this);
STATE_UPDATER.set(this, State.Closed);
+ randomReaders.closeAll();
clearPendingAddEntries(new ManagedLedgerAlreadyClosedException("Managed ledger is closed"));
cancelScheduledTasks();
@@ -2344,6 +2350,11 @@ public void asyncReadEntry(Position position, ReadEntryCallback callback, Object
}
+ @Override
+ public RandomReader newRandomReader(boolean streaming) {
+ return randomReaders.create(streaming);
+ }
+
private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) {
if (opReadEntry.readPosition.compareTo(opReadEntry.maxPosition) > 0) {
@@ -2465,6 +2476,27 @@ protected void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry
}
}
+ void asyncReadEntryForRandomReader(ReadHandle ledger, long firstEntry, long lastEntry,
+ IntSupplier expectedReadCount, ReadEntriesCallback callback) {
+ // expectedReadCount is resolved by the caller: streaming readers weight misses (>0), pure-random bypasses (0).
+ asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, callback, null);
+ }
+
+ private void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry,
+ IntSupplier expectedReadCount, ReadEntriesCallback callback, Object ctx) {
+ if (config.getReadEntryTimeoutSeconds() > 0) {
+ // set readOpCount to uniquely validate if ReadEntryCallbackWrapper is already recycled
+ long readOpCount = READ_OP_COUNT_UPDATER.incrementAndGet(this);
+ long createdTime = System.nanoTime();
+ ReadEntryCallbackWrapper readCallback = ReadEntryCallbackWrapper.create(name, ledger.getId(), firstEntry,
+ callback, readOpCount, createdTime, ctx);
+ lastReadCallback = readCallback;
+ entryCache.asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, readCallback, readOpCount);
+ } else {
+ entryCache.asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, callback, ctx);
+ }
+ }
+
static final class ReadEntryCallbackWrapper implements ReadEntryCallback, ReadEntriesCallback {
volatile ReadEntryCallback readEntryCallback;
@@ -4509,6 +4541,7 @@ public synchronized void setFenced() {
log.info().log("Moving to Fenced state");
State prev = STATE_UPDATER.getAndSet(this, State.Fenced);
if (prev != State.Fenced) {
+ randomReaders.closeAll();
clearPendingAddEntries(new ManagedLedgerFencedException("ManagedLedger "
+ name + " is fenced"));
}
@@ -4518,6 +4551,7 @@ synchronized void setFencedForDeletion() {
log.info().log("Moving to FencedForDeletion state");
State prev = STATE_UPDATER.getAndSet(this, State.FencedForDeletion);
if (prev != State.FencedForDeletion) {
+ randomReaders.closeAll();
clearPendingAddEntries(new ManagedLedgerFencedException("ManagedLedger "
+ name + " is fenced"));
}
@@ -5270,7 +5304,17 @@ public void waitForPendingCacheEvictions() {
}
boolean shouldCacheAddedEntry() {
- // Avoid caching entries if no cursor has been created
- return getActiveCursors().shouldCacheAddedEntry();
+ // Only streaming random readers seed the tail; pure-random readers do not.
+ return getActiveCursors().shouldCacheAddedEntry() || randomReaders.hasStreamingReaders();
+ }
+
+ @VisibleForTesting
+ int getActiveRandomReaderCount() {
+ return randomReaders.size();
+ }
+
+ // Package-private: streaming readers contribute to expectedReadCount on the read and add paths.
+ int getActiveStreamingRandomReaderCount() {
+ return randomReaders.streamingCount();
}
}
diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java
index 2079caf08a36a..d5fa74b708ccf 100644
--- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java
+++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java
@@ -268,12 +268,14 @@ public void run() {
long ledgerId = ledger != null ? ledger.getId() : ((Position) ctx).getLedgerId();
// Handle caching for tailing reads
+ // Streaming readers (like cursors) consume the tail, so they count toward expectedReadCount and eviction
+ // priority. With cacheEvictionByExpectedReadCount disabled the handler is intentionally null.
if (ml.shouldCacheAddedEntry()) {
int expectedReadCount = 0;
// only use expectedReadCount if cache eviction is enabled by expected read count
if (ml.getConfig().isCacheEvictionByExpectedReadCount()) {
- // use the number of active cursors as the expected read count
- expectedReadCount = ml.getActiveCursors().size();
+ // active cursors + streaming random readers all read the tail entry
+ expectedReadCount = ml.getActiveCursors().size() + ml.getActiveStreamingRandomReaderCount();
}
EntryImpl entry = EntryImpl.create(ledgerId, entryId, data, expectedReadCount);
entry.setDecreaseReadCountOnRelease(false);
diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java
new file mode 100644
index 0000000000000..ca14b4bfac9ee
--- /dev/null
+++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java
@@ -0,0 +1,275 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.bookkeeper.mledger.impl;
+
+import static java.lang.Math.min;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.atomic.AtomicBoolean;
+import lombok.CustomLog;
+import org.apache.bookkeeper.client.LedgerHandle;
+import org.apache.bookkeeper.client.api.ReadHandle;
+import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback;
+import org.apache.bookkeeper.mledger.Entry;
+import org.apache.bookkeeper.mledger.ManagedLedgerException;
+import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerFencedException;
+import org.apache.bookkeeper.mledger.Position;
+import org.apache.bookkeeper.mledger.PositionFactory;
+import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.State;
+import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo;
+import org.apache.pulsar.common.util.FutureUtil;
+
+@CustomLog
+class OpReadEntries implements ReadEntriesCallback {
+ private final ManagedLedgerImpl ledger;
+ private final Position maxPosition;
+ private final int count;
+ private final CompletableFuture
> promise = new CompletableFuture<>();
+ private final List
> read(ManagedLedgerImpl ledger, Position readPosition, int count,
+ Position maxPosition) {
+ return read(ledger, readPosition, count, maxPosition, true);
+ }
+
+ static CompletableFuture
> read(ManagedLedgerImpl ledger, Position readPosition, int count,
+ Position maxPosition, boolean streaming) {
+ OpReadEntries op = new OpReadEntries(ledger, readPosition, count, maxPosition, streaming);
+ op.readEntries();
+ return op.promise;
+ }
+
+ void readEntries() {
+ if (terminal.get()) {
+ return;
+ }
+ final State state = ManagedLedgerImpl.STATE_UPDATER.get(ledger);
+ if (state.isFenced() || state == State.Closed) {
+ readEntriesFailed(new ManagedLedgerFencedException(), null);
+ return;
+ }
+
+ if (readPosition.compareTo(maxPosition) > 0) {
+ checkReadCompletion();
+ return;
+ }
+
+ long ledgerId = readPosition.getLedgerId();
+ LedgerHandle currentLedger = ledger.currentLedger;
+
+ if (currentLedger != null && ledgerId == currentLedger.getId()) {
+ // Current writing ledger is not in the cache (since we don't want
+ // it to be automatically evicted), and we cannot use 2 different
+ // ledger handles (read & write)for the same ledger.
+ internalReadFromLedger(currentLedger);
+ } else {
+ LedgerInfo ledgerInfo = ledger.ledgers.get(ledgerId);
+ if (ledgerInfo == null || ledgerInfo.getEntries() == 0) {
+ updateReadPosition(getNextLedgerPosition(ledgerId));
+ checkReadCompletion();
+ return;
+ }
+
+ ledger.getLedgerHandle(ledgerId).thenAccept(this::internalReadFromLedger)
+ .exceptionally(ex -> {
+ ledger.log.error().attr("position", readPosition).exceptionMessage(ex)
+ .log("Error opening ledger for reading");
+ readEntriesFailed(ManagedLedgerException.getManagedLedgerException(
+ FutureUtil.unwrapCompletionException(ex)), null);
+ return null;
+ });
+ }
+ }
+
+ private void internalReadFromLedger(ReadHandle readHandle) {
+ long firstEntry = readPosition.getEntryId();
+ long lastEntryInLedger;
+
+ Position lastPosition = ledger.lastConfirmedEntry;
+
+ if (readHandle.getId() == lastPosition.getLedgerId()) {
+ // For the current ledger, we only give read visibility to the last entry we have received a confirmation in
+ // the managed ledger layer
+ lastEntryInLedger = lastPosition.getEntryId();
+ } else {
+ // For other ledgers, already closed the BK lastAddConfirmed is appropriate
+ lastEntryInLedger = readHandle.getLastAddConfirmed();
+ }
+
+ if (readHandle.getId() == maxPosition.getLedgerId()) {
+ lastEntryInLedger = min(maxPosition.getEntryId(), lastEntryInLedger);
+ }
+
+ if (firstEntry > lastEntryInLedger) {
+ log.debug().attr("ledgerId", readHandle.getId())
+ .attr("lastEntry", lastEntryInLedger)
+ .attr("readEntry", firstEntry)
+ .log("No more messages to read from ledger");
+
+ LedgerHandle currentLedger = ledger.currentLedger;
+ if (currentLedger == null || readHandle.getId() != currentLedger.getId()) {
+ updateReadPosition(getNextLedgerPosition(readHandle.getId()));
+ } else {
+ updateReadPosition(readPosition);
+ }
+
+ checkReadCompletion();
+ return;
+ }
+
+ long lastEntry = min(firstEntry + getNumberOfEntriesToRead() - 1, lastEntryInLedger);
+
+ log.debug().attr("ledgerId", readHandle.getId())
+ .attr("firstEntry", firstEntry)
+ .attr("lastEntry", lastEntry)
+ .log("Reading entries from ledger");
+ if (streaming) {
+ ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry,
+ ledger::getActiveStreamingRandomReaderCount, this);
+ } else {
+ // Pure-random: do not write back misses (expectedReadCount = 0).
+ ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry, () -> 0, this);
+ }
+ }
+
+ private Position getNextLedgerPosition(long ledgerId) {
+ Long nextLedgerId = ledger.ledgers.ceilingKey(ledgerId + 1);
+ return PositionFactory.create(nextLedgerId != null ? nextLedgerId : ledgerId + 1, 0);
+ }
+
+ @Override
+ public void readEntriesComplete(List
> read(Position startPosition, int numberOfEntries) {
+ return read(startPosition, numberOfEntries, PositionFactory.LATEST,
+ ManagedLedgerUtils.NO_MAX_SIZE_LIMIT);
+ }
+
+ @Override
+ public CompletableFuture
> read(Position startPosition, int numberOfEntries, Position maxPosition,
+ long maxSizeBytes) {
+ if (closed.get() || owner.isClosed()) {
+ return CompletableFuture.failedFuture(
+ new ManagedLedgerAlreadyClosedException("Random reader is already closed"));
+ }
+ if (startPosition == null || numberOfEntries <= 0) {
+ return CompletableFuture.failedFuture(new IllegalArgumentException("Invalid parameters"));
+ }
+
+ Position normalizedMaxPosition = maxPosition != null ? maxPosition : PositionFactory.LATEST;
+ Position normalizedStartPosition = normalizeStartPosition(startPosition);
+ if (normalizedStartPosition == null || normalizedMaxPosition.compareTo(normalizedStartPosition) < 0) {
+ return CompletableFuture.completedFuture(Collections.emptyList());
+ }
+
+ int effectiveCount = numberOfEntries;
+ if (maxSizeBytes != ManagedLedgerUtils.NO_MAX_SIZE_LIMIT) {
+ effectiveCount = Math.min(numberOfEntries,
+ estimateEntryCountByBytesSize(numberOfEntries, maxSizeBytes, normalizedStartPosition, ledger));
+ }
+
+ return OpReadEntries.read(ledger, normalizedStartPosition, effectiveCount, normalizedMaxPosition, streaming);
+ }
+
+ private Position normalizeStartPosition(Position startPosition) {
+ if (PositionFactory.EARLIEST.equals(startPosition)) {
+ Map.Entry
> readAfterLast = reader.read(p0.getNext(), 10);
+ CompletableFuture
> readLatest = reader.read(PositionFactory.LATEST, 10);
+
+ assertEquals(readAfterLast.get(5, TimeUnit.SECONDS), Collections.emptyList());
+ assertEquals(readLatest.get(5, TimeUnit.SECONDS), Collections.emptyList());
+
+ ledger.addEntry("entry-1".getBytes(Encoding));
+ assertEquals(readAfterLast.get(5, TimeUnit.SECONDS), Collections.emptyList());
+ assertEquals(readLatest.get(5, TimeUnit.SECONDS), Collections.emptyList());
+
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderExactCountBoundaries() throws Exception {
+ ManagedLedger ledger = factory.open("testRandomReaderExactCountBoundaries",
+ new ManagedLedgerConfig().setMaxEntriesPerLedger(3));
+ @Cleanup RandomReader reader = ledger.newRandomReader();
+
+ Position p0 = ledger.addEntry("entry-0".getBytes(Encoding));
+ Position p1 = ledger.addEntry("entry-1".getBytes(Encoding));
+ Position p2 = ledger.addEntry("entry-2".getBytes(Encoding));
+ Position p3 = ledger.addEntry("entry-3".getBytes(Encoding));
+ Position p4 = ledger.addEntry("entry-4".getBytes(Encoding));
+
+ assertEquals(p0.getLedgerId(), p2.getLedgerId());
+ assertNotEquals(p2.getLedgerId(), p3.getLedgerId());
+
+ assertEntryPositionsAndRelease(reader.read(p1, 1).get(5, TimeUnit.SECONDS), p1);
+ assertEntryPositionsAndRelease(reader.read(p1, 2).get(5, TimeUnit.SECONDS), p1, p2);
+ assertEntryPositionsAndRelease(reader.read(p1, 10).get(5, TimeUnit.SECONDS), p1, p2, p3, p4);
+
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderSkipsEmptyLedgers() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderSkipsEmptyLedgers",
+ new ManagedLedgerConfig().setMaxEntriesPerLedger(10));
+ @Cleanup RandomReader reader = ledger.newRandomReader();
+
+ Position p0 = ledger.addEntry("entry-0".getBytes(Encoding));
+ ledger.ledgerClosed(ledger.currentLedger, 0L);
+ Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger),
+ ManagedLedgerImpl.State.LedgerOpened));
+ LedgerHandle emptyLedger = ledger.currentLedger;
+ ledger.ledgerClosed(emptyLedger, -1L);
+ Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger),
+ ManagedLedgerImpl.State.LedgerOpened));
+ Position p1 = ledger.addEntry("entry-1".getBytes(Encoding));
+
+ assertNotEquals(p0.getLedgerId(), p1.getLedgerId());
+ assertFalse(ledger.getLedgersInfo().containsKey(emptyLedger.getId()));
+ assertEntryPositionsAndRelease(
+ reader.read(PositionFactory.create(emptyLedger.getId(), 0), 1).get(5, TimeUnit.SECONDS),
+ p1);
+
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStopsOnErrorAndReturnsPartialEntries() throws Exception {
+ ManagedLedger ledger = factory.open("testRandomReaderStopsOnErrorAndReturnsPartialEntries",
+ new ManagedLedgerConfig().setMaxEntriesPerLedger(1));
+ @Cleanup RandomReader reader = ledger.newRandomReader();
+
+ Position p0 = ledger.addEntry("entry-0".getBytes(Encoding));
+ Position p1 = ledger.addEntry("entry-1".getBytes(Encoding));
+ Position p2 = ledger.addEntry("entry-2".getBytes(Encoding));
+
+ assertNotEquals(p0.getLedgerId(), p1.getLedgerId());
+ assertNotEquals(p1.getLedgerId(), p2.getLedgerId());
+
+ bkc.deleteLedger(p1.getLedgerId());
+
+ assertEntryPositionsAndRelease(reader.read(p0, 3).get(5, TimeUnit.SECONDS), p0);
+
+ assertTrue(expectFutureFailure(reader.read(p1, 3)) instanceof ManagedLedgerException);
+
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderFailsWhenClosed() throws Exception {
+ ManagedLedger ledger = factory.open("testRandomReaderFailsWhenClosed");
+ @Cleanup RandomReader reader = ledger.newRandomReader();
+ Position position = ledger.addEntry("entry-0".getBytes(Encoding));
+
+ ledger.close();
+
+ assertTrue(expectFutureFailure(reader.read(position, 1))
+ instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException);
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderIgnoresCursorAckState() throws Exception {
+ ManagedLedger ledger = factory.open("testRandomReaderIgnoresCursorAckState");
+ @Cleanup RandomReader reader = ledger.newRandomReader();
+ ManagedCursor cursor = ledger.openCursor("c1");
+
+ Position p0 = ledger.addEntry("entry-0".getBytes(Encoding));
+ Position p1 = ledger.addEntry("entry-1".getBytes(Encoding));
+
+ List
> first = reader.read(p0, 2);
+ CompletableFuture
> second = reader.read(p2, 2);
+ assertEntryPositionsAndRelease(first.get(5, TimeUnit.SECONDS), p0, p1);
+ assertEntryPositionsAndRelease(second.get(5, TimeUnit.SECONDS), p2, p3);
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testOpReadEntriesSkipsEmptyLedger() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testOpReadEntriesSkipsEmptyLedger",
+ new ManagedLedgerConfig().setMaxEntriesPerLedger(10));
+ ledger.addEntry("entry-0".getBytes(Encoding));
+ ledger.ledgerClosed(ledger.currentLedger, 0L);
+ Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger),
+ ManagedLedgerImpl.State.LedgerOpened));
+ LedgerHandle emptyLedger = ledger.currentLedger;
+ ledger.ledgerClosed(emptyLedger, -1L);
+ Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger),
+ ManagedLedgerImpl.State.LedgerOpened));
+ Position nextPosition = ledger.addEntry("entry-1".getBytes(Encoding));
+ ledger.ledgers.put(emptyLedger.getId(), new LedgerInfo()
+ .setLedgerId(emptyLedger.getId()).setEntries(0).setSize(0));
+
+ CompletableFuture
> promise = OpReadEntries.read(ledger,
+ PositionFactory.create(emptyLedger.getId(), 0), 1, PositionFactory.LATEST);
+ assertEntryPositionsAndRelease(promise.get(5, TimeUnit.SECONDS), nextPosition);
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingFalseDoesNotSeedCache() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderStreamingFalseDoesNotSeedCache");
+ assertEquals(ledger.getActiveRandomReaderCount(), 0);
+ assertFalse(ledger.shouldCacheAddedEntry());
+ try (RandomReader reader = ledger.newRandomReader(false)) {
+ assertEquals(ledger.getActiveRandomReaderCount(), 1);
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 0);
+ assertFalse(ledger.shouldCacheAddedEntry());
+ ledger.addEntry("entry-0".getBytes(Encoding));
+ assertEquals(ledger.getCacheSize(), 0);
+ }
+ assertEquals(ledger.getActiveRandomReaderCount(), 0);
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingFalseDoesNotPopulateCacheOnMiss() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderStreamingFalseDoesNotPopulateCacheOnMiss");
+ Position position = ledger.addEntry("entry".getBytes(Encoding));
+ assertEquals(ledger.getCacheSize(), 0);
+ try (RandomReader reader = ledger.newRandomReader(false)) {
+ long missesBefore = factory.getMbean().getCacheMissesTotal();
+ assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore + 1);
+ assertEquals(ledger.getCacheSize(), 0);
+ long missesBefore2 = factory.getMbean().getCacheMissesTotal();
+ assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore2 + 1);
+ }
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingTruePopulatesCacheOnMiss() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderStreamingTruePopulatesCacheOnMiss");
+ Position position = ledger.addEntry("entry".getBytes(Encoding));
+ assertEquals(ledger.getCacheSize(), 0);
+ try (RandomReader reader = ledger.newRandomReader(true)) {
+ long missesBefore = factory.getMbean().getCacheMissesTotal();
+ assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore + 1);
+ assertTrue(ledger.getCacheSize() > 0);
+ long hitsBefore = factory.getMbean().getCacheHitsTotal();
+ assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ assertEquals(factory.getMbean().getCacheHitsTotal(), hitsBefore + 1);
+ }
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingModesCoexistForSeeding() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderStreamingModesCoexistForSeeding");
+ RandomReader bypass = ledger.newRandomReader(false);
+ assertFalse(ledger.shouldCacheAddedEntry());
+ RandomReader streaming = ledger.newRandomReader(true);
+ assertTrue(ledger.shouldCacheAddedEntry());
+ streaming.close();
+ assertFalse(ledger.shouldCacheAddedEntry());
+ bypass.close();
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testNewRandomReaderDefaultsToStreaming() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testNewRandomReaderDefaultsToStreaming");
+ try (RandomReader defaultReader = ledger.newRandomReader()) {
+ try (RandomReader explicitReader = ledger.newRandomReader(true)) {
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 2);
+ }
+ }
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 0);
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingFalseServesExistingCacheHit() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderStreamingFalseServesExistingCacheHit");
+ Position position = ledger.addEntry("entry".getBytes(Encoding));
+ try (RandomReader seeder = ledger.newRandomReader(true)) {
+ assertEntryDataAndRelease(seeder.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ }
+ assertTrue(ledger.getCacheSize() > 0);
+ try (RandomReader bypass = ledger.newRandomReader(false)) {
+ long missesBefore = factory.getMbean().getCacheMissesTotal();
+ assertEntryDataAndRelease(bypass.read(position, 1).get(5, TimeUnit.SECONDS), "entry");
+ assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore);
+ }
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testStreamingReaderCountAccessorWiring() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testStreamingReaderCountAccessorWiring");
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 0);
+ RandomReader s1 = ledger.newRandomReader(true);
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 1);
+ RandomReader b1 = ledger.newRandomReader(false);
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 1);
+ RandomReader s2 = ledger.newRandomReader(true);
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 2);
+ s1.close();
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 1);
+ s2.close();
+ b1.close();
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 0);
+ ledger.close();
+ }
+
+ @Test(timeOut = 20000)
+ public void testFenceZeroesStreamingReaders() throws Exception {
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testFenceZeroesStreamingReaders");
+ RandomReader streaming = ledger.newRandomReader(true);
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 1);
+ ledger.setFenced();
+ assertEquals(ledger.getActiveStreamingRandomReaderCount(), 0);
+ assertTrue(((RandomReaderImpl) streaming).isClosed());
+ // Do not close(): a fenced ledger cannot be closed (throws ManagedLedgerFencedException); the factory
+ // tears it down. Mirrors testRandomReaderLifecycleTracksManagedLedgerCloseAndFence.
+ }
+
+ // These tests assert the expectedReadCount value entries carry, not eviction under pressure.
+
+ @Test(timeOut = 20000)
+ public void testRandomReaderStreamingSeedsTailWithEvictionWeight() throws Exception {
+ ManagedLedgerConfig config = new ManagedLedgerConfig();
+ config.setCacheEvictionByExpectedReadCount(true);
+ ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open(
+ "testRandomReaderStreamingSeedsTailWithEvictionWeight", config);
+ try (RandomReader reader = ledger.newRandomReader(true)) {
+ Position position = ledger.addEntry("entry-0".getBytes(Encoding));
+ List