From 4956976dc2b938022a453008d5b7f73eb1da986e Mon Sep 17 00:00:00 2001 From: horizonzy Date: Sun, 10 Apr 2022 13:35:35 +0800 Subject: [PATCH 1/6] patch #5809: Fix the ledgerID not found cause NPE. --- .../mledger/impl/ManagedLedgerImpl.java | 11 +++-- .../bookkeeper/mledger/impl/OpReadEntry.java | 15 ++++++- .../mledger/impl/ManagedCursorTest.java | 44 +++++++++++++++++++ 3 files changed, 65 insertions(+), 5 deletions(-) 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 28d6b63a99f43..dec709bf1b7db 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 @@ -1750,6 +1750,12 @@ void asyncReadEntries(OpReadEntry opReadEntry) { opReadEntry.readEntriesFailed(new ManagedLedgerFencedException(), opReadEntry.ctx); return; } + if (opReadEntry.isInvalid()) { + opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(opReadEntry.readPosition" + + ") method is used to return the least key greater than or equal to the given key, " + + "or null if there is no such key"), opReadEntry.ctx); + return; + } long ledgerId = opReadEntry.readPosition.getLedgerId(); @@ -2230,9 +2236,8 @@ void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { PositionImpl startReadOperationOnLedger(PositionImpl position, OpReadEntry opReadEntry) { Long ledgerId = ledgers.ceilingKey(position.getLedgerId()); 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); + opReadEntry.makeInvalid(); + 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 178ee2dedabbe..ce7756e0f5ce5 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 @@ -39,6 +39,7 @@ class OpReadEntry implements ReadEntriesCallback { private int count; private ReadEntriesCallback callback; Object ctx; + private boolean invalid; // Results private List entries; @@ -48,7 +49,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,7 +58,8 @@ public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPosi } op.maxPosition = maxPosition; op.ctx = ctx; - op.nextReadPosition = PositionImpl.get(op.readPosition); + op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); + op.nextReadPosition = op.readPosition == null ? null : PositionImpl.get(op.readPosition); return op; } @@ -161,6 +162,14 @@ void checkReadCompletion() { } } + public void makeInvalid() { + this.invalid = true; + } + + public boolean isInvalid() { + return this.invalid; + } + public int getNumberOfEntriesToRead() { return count - entries.size(); } @@ -185,8 +194,10 @@ protected OpReadEntry newObject(Recycler.Handle recyclerHandle) { public void recycle() { cursor = null; readPosition = null; + count = 0; callback = null; ctx = null; + invalid = false; entries = null; nextReadPosition = null; maxPosition = null; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 6202a0f2a4ef1..2156e5fffaa48 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -3750,5 +3750,49 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { Awaitility.await().untilAsserted(() -> assertTrue(flag.get())); } + @Test + public void testReadInvalidPositionEntry() throws Exception { + ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); + managedLedgerConfig.setMaxEntriesPerLedger(1); + managedLedgerConfig.setMetadataMaxEntriesPerLedger(1); + managedLedgerConfig.setMinimumRolloverTime(0, TimeUnit.MILLISECONDS); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory + .open("testReadInvalidPositionEntry", managedLedgerConfig); + ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("test"); + + PositionImpl lastPosition = (PositionImpl) ledger.addEntry("test".getBytes(Encoding)); + ledger.rollCurrentLedgerIfFull(); + + AtomicReference exceptionRef = new AtomicReference<>(); + AtomicReference ctxRef = new AtomicReference<>(); + + ReadEntriesCallback callback = new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + exceptionRef.set(exception); + ctxRef.set(ctx); + } + }; + //Invalid position + PositionImpl invalidPosition = new PositionImpl(ledger.lastConfirmedEntry.getLedgerId() + 1, 0); + + Object ctx = new Object(); + OpReadEntry opReadEntry = OpReadEntry.create(cursor, invalidPosition, 10, callback, + ctx, PositionImpl.get(lastPosition.getLedgerId(), -1)); + ledger.asyncReadEntries(opReadEntry); + + // when readPosition's ledger didn't exist, throw NoMoreEntriesToReadException. + Awaitility.await().untilAsserted(() -> assertNotNull(exceptionRef.get())); + Awaitility.await().untilAsserted(() -> assertEquals(exceptionRef.get().getClass(), ManagedLedgerException.NoMoreEntriesToReadException.class)); + + Awaitility.await().untilAsserted(() -> assertNotNull(ctxRef.get())); + Awaitility.await().untilAsserted(() -> assertEquals(ctxRef.get(), ctx)); + } + private static final Logger log = LoggerFactory.getLogger(ManagedCursorTest.class); } From f7c52fa35c0be39faceb5de958c04a72686c9bba Mon Sep 17 00:00:00 2001 From: horizonzy Date: Sun, 10 Apr 2022 14:47:50 +0800 Subject: [PATCH 2/6] code style fix. --- .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 dec709bf1b7db..5cf274763de2f 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 @@ -1751,9 +1751,9 @@ void asyncReadEntries(OpReadEntry opReadEntry) { return; } if (opReadEntry.isInvalid()) { - opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(opReadEntry.readPosition" - + ") method is used to return the least key greater than or equal to the given key, " - + "or null if there is no such key"), opReadEntry.ctx); + opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(" + + "opReadEntry.readPosition method is used to return the least key greater than or equal to the " + + "given key, or null if there is no such key"), opReadEntry.ctx); return; } From f9f066741cea2e5b05318e7fdd26b47865f494b6 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Sun, 10 Apr 2022 15:07:56 +0800 Subject: [PATCH 3/6] error msg tweak. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 5cf274763de2f..4c2d955611fe2 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 @@ -1752,7 +1752,7 @@ void asyncReadEntries(OpReadEntry opReadEntry) { } if (opReadEntry.isInvalid()) { opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(" - + "opReadEntry.readPosition method is used to return the least key greater than or equal to the " + + "opReadEntry.readPosition) method is used to return the least key greater than or equal to the " + "given key, or null if there is no such key"), opReadEntry.ctx); return; } From 8da51bb6dcbd82444e2ed166fece346e98493a31 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Sun, 10 Apr 2022 19:23:48 +0800 Subject: [PATCH 4/6] for wait opReadEntry, the readPosition will recalculate, invalid should re check. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 3 +++ .../java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java | 4 ++++ 2 files changed, 7 insertions(+) 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 4c2d955611fe2..77498ed3448f2 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 @@ -2238,6 +2238,9 @@ PositionImpl startReadOperationOnLedger(PositionImpl position, OpReadEntry opRea if (null == ledgerId) { opReadEntry.makeInvalid(); return null; + } else { + // for wait opReadEntry, the readPosition will recalculate. + opReadEntry.makeValid(); } 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 ce7756e0f5ce5..22f9ff38d5f4f 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 @@ -166,6 +166,10 @@ public void makeInvalid() { this.invalid = true; } + public void makeValid() { + this.invalid = false; + } + public boolean isInvalid() { return this.invalid; } From 23b3dd76f02bba8ab690b8aa7c0fdd886dc0c7ed Mon Sep 17 00:00:00 2001 From: horizonzy Date: Mon, 11 Apr 2022 12:54:22 +0800 Subject: [PATCH 5/6] add test for wait OpReadEntry. --- .../mledger/impl/ManagedCursorImpl.java | 5 +- .../mledger/impl/ManagedCursorTest.java | 72 ++++++++++++++++++- 2 files changed, 73 insertions(+), 4 deletions(-) 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 199868b543fdd..a0a95c2af5319 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 @@ -821,6 +821,8 @@ private void checkForNewEntries(OpReadEntry op, ReadEntriesCallback callback, Ob name, op.readPosition); } PENDING_READ_OPS_UPDATER.incrementAndGet(this); + //Here maybe in schedule thread, recalculate readPosition, ledger maybe add some new entry. + op.readPosition = ledger.startReadOperationOnLedger((PositionImpl) getReadPosition(), op); ledger.asyncReadEntries(op); } else { if (log.isDebugEnabled()) { @@ -2761,7 +2763,8 @@ void notifyEntriesAvailable() { } PENDING_READ_OPS_UPDATER.incrementAndGet(this); - opReadEntry.readPosition = (PositionImpl) getReadPosition(); + //recalculate readPosition, cause ledger's add some new entry. + opReadEntry.readPosition = ledger.startReadOperationOnLedger((PositionImpl) getReadPosition(), opReadEntry); ledger.asyncReadEntries(opReadEntry); } else { // No one is waiting to be notified. Ignore diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 2156e5fffaa48..1a38edd1b003b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -3779,11 +3779,11 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { } }; //Invalid position - PositionImpl invalidPosition = new PositionImpl(ledger.lastConfirmedEntry.getLedgerId() + 1, 0); + PositionImpl invalidPosition = new PositionImpl(lastPosition.getLedgerId() + 1, 0); Object ctx = new Object(); - OpReadEntry opReadEntry = OpReadEntry.create(cursor, invalidPosition, 10, callback, - ctx, PositionImpl.get(lastPosition.getLedgerId(), -1)); + OpReadEntry opReadEntry = OpReadEntry.create(cursor, invalidPosition, 1, callback, + ctx, PositionImpl.LATEST); ledger.asyncReadEntries(opReadEntry); // when readPosition's ledger didn't exist, throw NoMoreEntriesToReadException. @@ -3794,5 +3794,71 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { Awaitility.await().untilAsserted(() -> assertEquals(ctxRef.get(), ctx)); } + @Test + public void testReadInvalidPositionEntryWhenNotify() throws Exception { + ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); + managedLedgerConfig.setMaxEntriesPerLedger(1); + managedLedgerConfig.setMetadataMaxEntriesPerLedger(1); + managedLedgerConfig.setMinimumRolloverTime(0, TimeUnit.MILLISECONDS); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory + .open("testReadInvalidPositionEntryWhenNotify", managedLedgerConfig); + ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("test"); + + PositionImpl lastPosition = (PositionImpl) ledger.addEntry("test".getBytes(Encoding)); + + //make cursor read position greater than ledger's lastConfirmedEntry. + //It guarantees readPosition is valid and cursor has not more entries. + cursor.setReadPosition(new PositionImpl(lastPosition.getLedgerId() + 1, 0)); + + AtomicReference exceptionRef = new AtomicReference<>(); + AtomicReference exceptionCtxRef = new AtomicReference<>(); + + AtomicReference entryRef = new AtomicReference<>(); + AtomicReference ctxRef = new AtomicReference<>(); + + + ReadEntriesCallback callback = new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + assertEquals(entries.size(), 1); + entryRef.set(entries.get(0)); + ctxRef.set(ctx); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + exceptionRef.set(exception); + exceptionCtxRef.set(ctx); + } + }; + + Object ctx = new Object(); + + cursor.asyncReadEntriesOrWait(1, -1, callback, ctx, PositionImpl.LATEST); + + // here wait checkForNewEntries to add cursor to ledger waitingCursors. + if (cursor.config.getNewEntriesCheckDelayInMillis() > 0) { + Thread.sleep(cursor.config.getNewEntriesCheckDelayInMillis() + 50); + } + + //Add entry to notify waiting cursor. + ledger.addEntry("test1".getBytes(Encoding)); + + //Cause add new entry, the wait opReadEntry can work normally. + Awaitility.await().untilAsserted(() -> assertNull(exceptionRef.get())); + Awaitility.await().untilAsserted(() -> assertNull(exceptionCtxRef.get())); + + + Awaitility.await().untilAsserted(() -> assertNotNull(entryRef.get())); + Awaitility.await().untilAsserted(() -> assertEquals(entryRef.get().getData(), "test1".getBytes(Encoding))); + Awaitility.await().untilAsserted(() -> assertEquals(ctxRef.get(), ctx)); + + entryRef.get().release(); + } + + + + + private static final Logger log = LoggerFactory.getLogger(ManagedCursorTest.class); } From 4ab367e04942b2fd6a2cb1d123c0cb1c5671b137 Mon Sep 17 00:00:00 2001 From: horizonzy Date: Mon, 11 Apr 2022 15:27:17 +0800 Subject: [PATCH 6/6] add unit for notify wait opReadEntry still invalid situation. --- .../mledger/impl/ManagedCursorTest.java | 69 ++++++++++++++++++- 1 file changed, 66 insertions(+), 3 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index 1a38edd1b003b..577a5482cdbb7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -3795,20 +3795,22 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { } @Test - public void testReadInvalidPositionEntryWhenNotify() throws Exception { + public void testReadEntryWhenNotifyNormally() throws Exception { ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); managedLedgerConfig.setMaxEntriesPerLedger(1); managedLedgerConfig.setMetadataMaxEntriesPerLedger(1); managedLedgerConfig.setMinimumRolloverTime(0, TimeUnit.MILLISECONDS); ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory - .open("testReadInvalidPositionEntryWhenNotify", managedLedgerConfig); + .open("testReadEntryWhenNotifyNormally", managedLedgerConfig); ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("test"); PositionImpl lastPosition = (PositionImpl) ledger.addEntry("test".getBytes(Encoding)); //make cursor read position greater than ledger's lastConfirmedEntry. //It guarantees readPosition is valid and cursor has not more entries. - cursor.setReadPosition(new PositionImpl(lastPosition.getLedgerId() + 1, 0)); + // maxEntriesPerLedger == 1, so add two entry will occupy 2 ledgers. and cursor will occupy one ledger, + // so here need + 2. + cursor.setReadPosition(new PositionImpl(lastPosition.getLedgerId() + 2, 0)); AtomicReference exceptionRef = new AtomicReference<>(); AtomicReference exceptionCtxRef = new AtomicReference<>(); @@ -3856,6 +3858,67 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { entryRef.get().release(); } + @Test + public void testReadEntryWhenNotifyInvalid() throws Exception { + ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); + managedLedgerConfig.setMaxEntriesPerLedger(1); + managedLedgerConfig.setMetadataMaxEntriesPerLedger(1); + managedLedgerConfig.setMinimumRolloverTime(0, TimeUnit.MILLISECONDS); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory + .open("testReadEntryWhenNotifyInvalid", managedLedgerConfig); + ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("test"); + + PositionImpl lastPosition = (PositionImpl) ledger.addEntry("test".getBytes(Encoding)); + + //make cursor read position greater than ledger's lastConfirmedEntry. + //It guarantees readPosition is valid and cursor has not more entries. + // maxEntriesPerLedger == 1, so add two entry will occupy 2 ledgers. and cursor will occupy one ledger. + //so here + 3, make opReadEntry still invalid. + cursor.setReadPosition(new PositionImpl(lastPosition.getLedgerId() + 3, 0)); + + AtomicReference exceptionRef = new AtomicReference<>(); + AtomicReference exceptionCtxRef = new AtomicReference<>(); + + AtomicReference entryRef = new AtomicReference<>(); + AtomicReference ctxRef = new AtomicReference<>(); + + + ReadEntriesCallback callback = new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + exceptionRef.set(exception); + exceptionCtxRef.set(ctx); + } + }; + + Object ctx = new Object(); + + cursor.asyncReadEntriesOrWait(1, -1, callback, ctx, PositionImpl.LATEST); + + // here wait checkForNewEntries to add cursor to ledger waitingCursors. + if (cursor.config.getNewEntriesCheckDelayInMillis() > 0) { + Thread.sleep(cursor.config.getNewEntriesCheckDelayInMillis() + 50); + } + + //Add entry to notify waiting cursor. + ledger.addEntry("test1".getBytes(Encoding)); + + + // when readPosition's ledger didn't exist, throw NoMoreEntriesToReadException. + Awaitility.await().untilAsserted(() -> assertNotNull(exceptionRef.get())); + Awaitility.await().untilAsserted(() -> assertEquals(exceptionRef.get().getClass(), ManagedLedgerException.NoMoreEntriesToReadException.class)); + + Awaitility.await().untilAsserted(() -> assertNotNull(exceptionCtxRef.get())); + Awaitility.await().untilAsserted(() -> assertEquals(exceptionCtxRef.get(), ctx)); + + Awaitility.await().untilAsserted(() -> assertNull(entryRef.get())); + Awaitility.await().untilAsserted(() -> assertNull(ctxRef.get())); + } +