diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index aae8c5a098b83..9e0d2e8e1dd43 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -625,7 +625,9 @@ public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadE PENDING_READ_OPS_UPDATER.incrementAndGet(this); OpReadEntry op = OpReadEntry.create(this, readPosition, numOfEntriesToRead, callback, ctx, maxPosition); - ledger.asyncReadEntries(op); + if (op != null) { + ledger.asyncReadEntries(op); + } } @Override @@ -763,6 +765,9 @@ public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntrie } else { OpReadEntry op = OpReadEntry.create(this, readPosition, numberOfEntriesToRead, callback, ctx, maxPosition); + if (op == null) { + return; + } if (!WAITING_READ_OP_UPDATER.compareAndSet(this, null, op)) { callback.readEntriesFailed(new ManagedLedgerException("We can only have a single waiting callback"), 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 629d10d48d64a..c1fb3eef639cf 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 @@ -2154,6 +2154,7 @@ PositionImpl startReadOperationOnLedger(PositionImpl position, OpReadEntry opRea if (null == ledgerId) { opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(K key) method is used to return the " + "least key greater than or equal to the given key, or null if there is no such key"), null); + return null; } if (ledgerId != position.getLedgerId()) { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java index 91a6e26f567d0..78cf0d01390f9 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java @@ -48,7 +48,6 @@ class OpReadEntry implements ReadEntriesCallback { public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPositionRef, int count, ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { OpReadEntry op = RECYCLER.get(); - op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); op.cursor = cursor; op.count = count; op.callback = callback; @@ -58,6 +57,11 @@ public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPosi } op.maxPosition = maxPosition; op.ctx = ctx; + PositionImpl position = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); + if (position == null) { + return null; + } + op.readPosition = position; op.nextReadPosition = PositionImpl.get(op.readPosition); return op; } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index d837651f1fac5..74a0cf65b8246 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -3393,4 +3393,31 @@ public void testCancellationOfScheduledTasks() throws Exception { assertTrue(timeoutTask2.isCancelled()); assertTrue(checkLedgerRollTask2.isCancelled()); } + + @Test(timeOut = 20000) + public void testReadNonExistentLedger() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("test-non-existent-ledger", config); + ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("test-non-existent-cursor"); + + CompletableFuture future = new CompletableFuture<>(); + PositionImpl pos = PositionImpl.latest; + cursor.seek(pos); + cursor.asyncReadEntries(1, new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + future.complete(true); + } + }, null, pos); + + cursor.close(); + ledger.close(); + assert(future.get()); + assertEquals(cursor.getPendingReadOpsCount(), 0); + } }