From f01fc315c5af4e32c7ac4be62625b8475bd5a058 Mon Sep 17 00:00:00 2001 From: Oneby Wang <891734032@qq.com> Date: Wed, 14 Jan 2026 08:59:02 +0800 Subject: [PATCH 1/4] [fix][test] Initial commit to reproduce mockito spy and volatile field problem --- .../mledger/impl/cache/RangeEntryCacheManagerImpl.java | 2 ++ .../org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java | 3 ++- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheManagerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheManagerImpl.java index 6d048e11389cd..eac0073c6d9bb 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheManagerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheManagerImpl.java @@ -95,6 +95,8 @@ public EntryCache getEntryCache(ManagedLedger ml) { rangeCacheRemovalQueue, entryLengthFunction); EntryCache currentEntryCache = caches.putIfAbsent(ml.getName(), newEntryCache); if (currentEntryCache != null) { + log.warn("Entry cache for {} already exists, newEntryCache: {}, currentEntryCache: {}", ml.getName(), + newEntryCache, currentEntryCache); return currentEntryCache; } else { return newEntryCache; 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 12bbd66670b50..0278ded7b3871 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 @@ -3957,7 +3957,8 @@ public void testLockReleaseWhenTrimLedger() throws Exception { initManagedLedgerConfig(config); config.setMaxEntriesPerLedger(1); - ManagedLedgerImpl ledger = spy((ManagedLedgerImpl) factory.open("testLockReleaseWhenTrimLedger", config)); + ManagedLedgerImpl originalLedger = (ManagedLedgerImpl) factory.open("testLockReleaseWhenTrimLedger", config); + ManagedLedgerImpl ledger = spy(originalLedger); doThrow(new ManagedLedgerException.LedgerNotExistException("First non deleted Ledger is not found")) .when(ledger).advanceCursorsIfNecessary(any()); final int entries = 10; From 55685421e31a3dc42652045724c0ccc128d405d2 Mon Sep 17 00:00:00 2001 From: Oneby Wang <891734032@qq.com> Date: Wed, 14 Jan 2026 14:32:36 +0800 Subject: [PATCH 2/4] [fix][test] Fix testLockReleaseWhenTrimLedger flaky test --- .../ManagedLedgerSpyUsingOverrideTest.java | 137 ++++++++++++++++++ .../mledger/impl/ManagedLedgerTest.java | 36 ----- .../test/MockedBookKeeperTestCase.java | 12 +- 3 files changed, 147 insertions(+), 38 deletions(-) create mode 100644 managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java new file mode 100644 index 0000000000000..af0bc0f81353e --- /dev/null +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java @@ -0,0 +1,137 @@ +/* + * 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 org.assertj.core.api.Assertions.assertThat; +import static org.testng.Assert.assertEquals; +import java.nio.charset.Charset; +import java.nio.charset.StandardCharsets; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Supplier; +import org.apache.bookkeeper.client.BookKeeper; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; +import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.bookkeeper.mledger.ManagedLedgerFactoryConfig; +import org.apache.bookkeeper.mledger.proto.MLDataFormats; +import org.apache.bookkeeper.mledger.util.Futures; +import org.apache.bookkeeper.test.MockedBookKeeperTestCase; +import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; +import org.awaitility.Awaitility; +import org.testng.annotations.Factory; +import org.testng.annotations.Test; + +public class ManagedLedgerSpyUsingOverrideTest { + + private static final Charset Encoding = StandardCharsets.UTF_8; + + @Factory + public Object[] createNestedTestInstances() { + return new Object[]{new AdvanceCursorsIfNecessaryThrowsExceptionTest()}; + } + + static class AdvanceCursorsIfNecessaryThrowsExceptionTest extends MockedBookKeeperTestCase { + + private AtomicInteger advanceCursorsIfNecessaryCallTimes; + + @Override + protected ManagedLedgerFactoryImpl initManagedLedgerFactory(MetadataStoreExtended metadataStore, + BookKeeper bookKeeper, + ManagedLedgerFactoryConfig managedLedgerFactoryConfig, + ManagedLedgerConfig managedLedgerConfig) + throws Exception { + return new ManagedLedgerFactoryImpl(metadataStore, bookKeeper, managedLedgerFactoryConfig, + managedLedgerConfig) { + @Override + protected ManagedLedgerImpl createManagedLedger(BookKeeper bk, MetaStore store, String name, + ManagedLedgerConfig config, + Supplier> mlOwnershipChecker) { + return new ManagedLedgerImpl(this, bk, store, config, scheduledExecutor, name, mlOwnershipChecker) { + @Override + void advanceCursorsIfNecessary(List ledgersToDelete) + throws ManagedLedgerException.LedgerNotExistException { + advanceCursorsIfNecessaryCallTimes.incrementAndGet(); + throw new ManagedLedgerException.LedgerNotExistException( + "First non deleted Ledger is not found"); + } + }; + } + }; + } + + @Override + protected void setUpTestCase() throws Exception { + advanceCursorsIfNecessaryCallTimes = new AtomicInteger(0); + super.setUpTestCase(); + } + + @Override + protected void cleanUpTestCase() throws Exception { + advanceCursorsIfNecessaryCallTimes = new AtomicInteger(0); + super.cleanUpTestCase(); + } + + @Test + public void testLockReleaseWhenTrimLedger() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + initManagedLedgerConfig(config); + config.setMaxEntriesPerLedger(1); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testLockReleaseWhenTrimLedger", config); + final int entries = 10; + ManagedCursor cursor = ledger.openCursor("test-cursor" + UUID.randomUUID()); + for (int i = 0; i < entries; i++) { + ledger.addEntry(String.valueOf(i).getBytes(Encoding)); + } + + // Wait for new ledger created to avoid flaky test, so currentLedger is the last emptyLedger. + // This make sure that we will read entries by calling ReadHandle instead of using entry cache. + // If we don't wait for new ledger created, we may lose the last valid ReadHandle(emptyLedgerLedgerId-1) + // in ledgerCache. + Awaitility.await().untilAsserted(() -> { + assertEquals(ledger.ledgers.size() - 1, entries); + }); + + List entryList = cursor.readEntries(entries); + assertEquals(entryList.size(), entries); + + // We don't need to wait here since read operation is completed. + assertEquals(ledger.ledgerCache.size() - 1, entries - 1); + + cursor.clearBacklog(); + ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + + // Cleanup fails because ManagedLedgerNotFoundException is thrown + assertEquals(ledger.ledgers.size() - 1, entries); + assertEquals(ledger.ledgerCache.size() - 1, entries - 1); + + // The lock is released even if an ManagedLedgerNotFoundException occurs, so it can be called repeatedly + Awaitility.await().untilAsserted(() -> assertThat(advanceCursorsIfNecessaryCallTimes.get()).isEqualTo(3)); + cursor.close(); + ledger.close(); + } + } + +} 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 0278ded7b3871..3c5ed4731319d 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 @@ -25,11 +25,9 @@ import static org.mockito.ArgumentMatchers.anyMap; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.atLeast; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; @@ -3951,40 +3949,6 @@ public void testInvalidateReadHandleWhenDeleteLedger() throws Exception { ledger.close(); } - @Test - public void testLockReleaseWhenTrimLedger() throws Exception { - ManagedLedgerConfig config = new ManagedLedgerConfig(); - initManagedLedgerConfig(config); - config.setMaxEntriesPerLedger(1); - - ManagedLedgerImpl originalLedger = (ManagedLedgerImpl) factory.open("testLockReleaseWhenTrimLedger", config); - ManagedLedgerImpl ledger = spy(originalLedger); - doThrow(new ManagedLedgerException.LedgerNotExistException("First non deleted Ledger is not found")) - .when(ledger).advanceCursorsIfNecessary(any()); - final int entries = 10; - ManagedCursor cursor = ledger.openCursor("test-cursor" + UUID.randomUUID()); - for (int i = 0; i < entries; i++) { - ledger.addEntry(String.valueOf(i).getBytes(Encoding)); - } - List entryList = cursor.readEntries(entries); - assertEquals(entryList.size(), entries); - assertEquals(ledger.ledgers.size() - 1, entries); - assertEquals(ledger.ledgerCache.size() - 1, entries - 1); - cursor.clearBacklog(); - ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); - ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); - // Cleanup fails because ManagedLedgerNotFoundException is thrown - Awaitility.await().untilAsserted(() -> { - assertEquals(ledger.ledgers.size() - 1, entries); - assertEquals(ledger.ledgerCache.size() - 1, entries - 1); - }); - // The lock is released even if an ManagedLedgerNotFoundException occurs, so it can be called repeatedly - Awaitility.await().untilAsserted(() -> - verify(ledger, atLeast(2)).advanceCursorsIfNecessary(any())); - cursor.close(); - ledger.close(); - } - @Test public void testInvalidateReadHandleWhenConsumed() throws Exception { ManagedLedgerConfig config = new ManagedLedgerConfig(); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java b/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java index f4619955ea18b..9194dc5505c8b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/test/MockedBookKeeperTestCase.java @@ -24,6 +24,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import lombok.SneakyThrows; +import org.apache.bookkeeper.client.BookKeeper; import org.apache.bookkeeper.client.PulsarMockBookKeeper; import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; @@ -87,8 +88,7 @@ public final void setUp(Method method) throws Exception { initManagedLedgerFactoryConfig(managedLedgerFactoryConfig); ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); initManagedLedgerConfig(managedLedgerConfig); - factory = - new ManagedLedgerFactoryImpl(metadataStore, bkc, managedLedgerFactoryConfig, managedLedgerConfig); + factory = initManagedLedgerFactory(metadataStore, bkc, managedLedgerFactoryConfig, managedLedgerConfig); setUpTestCase(); } @@ -103,6 +103,14 @@ protected void initManagedLedgerFactoryConfig(ManagedLedgerFactoryConfig config) config.setCacheEvictionIntervalMs(200); } + protected ManagedLedgerFactoryImpl initManagedLedgerFactory(MetadataStoreExtended metadataStore, + BookKeeper bookKeeper, + ManagedLedgerFactoryConfig managedLedgerFactoryConfig, + ManagedLedgerConfig managedLedgerConfig) + throws Exception { + return new ManagedLedgerFactoryImpl(metadataStore, bookKeeper, managedLedgerFactoryConfig, managedLedgerConfig); + } + protected void setUpTestCase() throws Exception { } From e6ef842d677221f9f97fbbc16803d94b679252d2 Mon Sep 17 00:00:00 2001 From: Oneby Wang <891734032@qq.com> Date: Wed, 14 Jan 2026 14:39:13 +0800 Subject: [PATCH 3/4] [fix][test] Fix checkstyle --- .../impl/ManagedLedgerSpyUsingOverrideTest.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java index af0bc0f81353e..d1327813eec0d 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java @@ -57,16 +57,16 @@ static class AdvanceCursorsIfNecessaryThrowsExceptionTest extends MockedBookKeep @Override protected ManagedLedgerFactoryImpl initManagedLedgerFactory(MetadataStoreExtended metadataStore, - BookKeeper bookKeeper, - ManagedLedgerFactoryConfig managedLedgerFactoryConfig, - ManagedLedgerConfig managedLedgerConfig) + BookKeeper bookKeeper, + ManagedLedgerFactoryConfig managedLedgerFactoryConfig, + ManagedLedgerConfig managedLedgerConfig) throws Exception { return new ManagedLedgerFactoryImpl(metadataStore, bookKeeper, managedLedgerFactoryConfig, managedLedgerConfig) { @Override protected ManagedLedgerImpl createManagedLedger(BookKeeper bk, MetaStore store, String name, - ManagedLedgerConfig config, - Supplier> mlOwnershipChecker) { + ManagedLedgerConfig config, + Supplier> mlOwnershipChecker) { return new ManagedLedgerImpl(this, bk, store, config, scheduledExecutor, name, mlOwnershipChecker) { @Override void advanceCursorsIfNecessary(List ledgersToDelete) From d4e4a92b9c66bdde0d2b4f9a8710d41582c64325 Mon Sep 17 00:00:00 2001 From: Oneby Wang <891734032@qq.com> Date: Wed, 14 Jan 2026 16:27:20 +0800 Subject: [PATCH 4/4] [fix][test] Optimize test --- .../mledger/impl/ManagedLedgerSpyUsingOverrideTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java index d1327813eec0d..6913d41c8df0c 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerSpyUsingOverrideTest.java @@ -88,11 +88,11 @@ protected void setUpTestCase() throws Exception { @Override protected void cleanUpTestCase() throws Exception { - advanceCursorsIfNecessaryCallTimes = new AtomicInteger(0); + advanceCursorsIfNecessaryCallTimes.set(0); super.cleanUpTestCase(); } - @Test + @Test(invocationCount = 1000) public void testLockReleaseWhenTrimLedger() throws Exception { ManagedLedgerConfig config = new ManagedLedgerConfig(); initManagedLedgerConfig(config);