Skip to content
Merged
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 @@ -1685,7 +1685,6 @@ public void closeComplete(int rc, LedgerHandle lh, Object o) {
}

ledgerClosed(lh);
createLedgerAfterClosed();
}
}, System.nanoTime());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2237,15 +2237,13 @@ void testFindNewestMatchingAfterLedgerRollover() throws Exception {
ledger.addEntry("fourth".getBytes(Encoding));
Position last = ledger.addEntry("last-expired".getBytes(Encoding));

// roll a new ledger
int numLedgersBefore = ledger.getLedgersInfo().size();
ledger.getConfig().setMaxEntriesPerLedger(1);
Field stateUpdater = ManagedLedgerImpl.class.getDeclaredField("state");
stateUpdater.setAccessible(true);
stateUpdater.set(ledger, ManagedLedgerImpl.State.LedgerOpened);
// roll a new ledger
ledger.rollCurrentLedgerIfFull();
Awaitility.await().atMost(20, TimeUnit.SECONDS)
.until(() -> ledger.getLedgersInfo().size() > numLedgersBefore);
Awaitility.await().untilAsserted(() -> {
Assert.assertEquals(ledger.getLedgersInfo().size(), 1);
Assert.assertEquals(ledger.getState(), ManagedLedgerImpl.State.ClosedLedger);
});

// the algorithm looks for "expired" messages
// starting from the first, then it moves to the last message
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
import static org.mockito.Mockito.when;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertFalse;
import static org.testng.Assert.assertNotEquals;
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.assertNull;
import static org.testng.Assert.assertSame;
Expand Down Expand Up @@ -67,6 +68,7 @@
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Predicate;
import lombok.Cleanup;
Expand Down Expand Up @@ -1811,7 +1813,8 @@ public void testMaximumRolloverTime() throws Exception {
ledger.addEntry("data".getBytes());

Awaitility.await().untilAsserted(() -> {
assertEquals(ledger.getLedgersInfoAsList().size(), 2);
assertEquals(ledger.getLedgersInfoAsList().size(), 1);
assertEquals(ledger.getState(), ManagedLedgerImpl.State.ClosedLedger);
});
}

Expand Down Expand Up @@ -1958,7 +1961,7 @@ public void testDeletionAfterLedgerClosedAndRetention() throws Exception {
ManagedLedgerFactory factory = new ManagedLedgerFactoryImpl(metadataStore, bkc);
ManagedLedgerConfig config = new ManagedLedgerConfig();
config.setRetentionSizeInMB(0);
config.setMaxEntriesPerLedger(1);
config.setMaxEntriesPerLedger(2);
config.setRetentionTime(1, TimeUnit.SECONDS);
config.setMaximumRolloverTime(1, TimeUnit.SECONDS);

Expand All @@ -1968,19 +1971,25 @@ public void testDeletionAfterLedgerClosedAndRetention() throws Exception {
ml.addEntry("iamaverylongmessagethatshouldnotberetained".getBytes());
c1.skipEntries(1, IndividualDeletedEntries.Exclude);
c2.skipEntries(1, IndividualDeletedEntries.Exclude);
// let current ledger close
Field stateUpdater = ManagedLedgerImpl.class.getDeclaredField("state");
stateUpdater.setAccessible(true);
stateUpdater.set(ml, ManagedLedgerImpl.State.LedgerOpened);
long preLedgerId = ml.getLedgersInfoAsList().get(ml.ledgers.size() -1).getLedgerId();
ml.pendingAddEntries.add(OpAddEntry.
createNoRetainBuffer(ml, ByteBufAllocator.DEFAULT.buffer(128).retain(), null, null));
ml.rollCurrentLedgerIfFull();
AtomicLong currentLedgerId = new AtomicLong(-1);
// create a new ledger
Awaitility.await().untilAsserted(() -> {
currentLedgerId.set(ml.getLedgersInfoAsList().get(ml.ledgers.size() -1).getLedgerId());
assertNotEquals(preLedgerId, currentLedgerId.get());
});
// let retention expire
Thread.sleep(1500);
// delete the expired ledger
ml.internalTrimConsumedLedgers(CompletableFuture.completedFuture(null));

// the closed and expired ledger should be deleted
assertTrue(ml.getLedgersInfoAsList().size() <= 1);
assertEquals(ml.getTotalSize(), 0);
assertEquals(ml.getLedgersInfoAsList().size(), 1);
assertEquals(currentLedgerId.get(),
ml.getLedgersInfoAsList().get(ml.getLedgersInfoAsList().size() - 1).getLedgerId());
ml.close();
}

Expand Down Expand Up @@ -2245,10 +2254,12 @@ public void testGetPositionAfterN() throws Exception {
stateUpdater.setAccessible(true);
stateUpdater.set(managedLedger, ManagedLedgerImpl.State.LedgerOpened);
managedLedger.rollCurrentLedgerIfFull();
Awaitility.await().untilAsserted(() -> assertEquals(managedLedger.getLedgersInfo().size(), 3));
Awaitility.await().untilAsserted(() -> {
assertEquals(managedLedger.getLedgersInfo().size(), 2);
assertEquals(managedLedger.getState(), ManagedLedgerImpl.State.ClosedLedger);
});
assertEquals(5, managedLedger.getLedgersInfoAsList().get(0).getEntries());
assertEquals(5, managedLedger.getLedgersInfoAsList().get(1).getEntries());
assertEquals(0, managedLedger.getLedgersInfoAsList().get(2).getEntries());
log.info("### ledgers {}", managedLedger.getLedgersInfo());

long firstLedger = managedLedger.getLedgersInfo().firstKey();
Expand Down Expand Up @@ -3113,12 +3124,10 @@ public void testManagedLedgerRollOverIfFull() throws Exception {

// all the messages have benn acknowledged
// and all the ledgers have been removed except the last ledger
Field stateUpdater = ManagedLedgerImpl.class.getDeclaredField("state");
stateUpdater.setAccessible(true);
stateUpdater.set(ledger, ManagedLedgerImpl.State.LedgerOpened);
ledger.rollCurrentLedgerIfFull();
Awaitility.await().untilAsserted(() -> Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1));
Awaitility.await().untilAsserted(() -> Assert.assertEquals(ledger.getTotalSize(), 0));
Awaitility.await().untilAsserted(() ->
Assert.assertEquals(ledger.getState(), ManagedLedgerImpl.State.ClosedLedger));
}

@Test
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -391,10 +391,13 @@ public void testManagedLedgerTotalSize() throws Exception {
managedLedger.rollCurrentLedgerIfFull();

Awaitility.await().atMost(Duration.ofSeconds(3))
.until(() -> managedLedger.getLedgersInfo().size() > 1);
.untilAsserted(() -> {
Assert.assertEquals(managedLedger.getLedgersInfo().size(), 1);
Assert.assertEquals(managedLedger.getState(), ManagedLedgerImpl.State.ClosedLedger);
});

final List<LedgerInfo> ledgerInfoList = managedLedger.getLedgersInfoAsList();
Assert.assertEquals(ledgerInfoList.size(), 2);
Assert.assertEquals(ledgerInfoList.size(), 1);
Assert.assertEquals(ledgerInfoList.get(0).getSize(), managedLedger.getTotalSize());

cursor.close();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,10 +20,14 @@

import java.lang.reflect.Field;
import java.time.Duration;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.TimeUnit;
import io.netty.buffer.ByteBufAllocator;
import lombok.Cleanup;
import org.apache.bookkeeper.mledger.ManagedLedgerConfig;
import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl;
import org.apache.bookkeeper.mledger.impl.OpAddEntry;
import org.apache.bookkeeper.mledger.proto.MLDataFormats;
import org.apache.pulsar.broker.service.persistent.PersistentTopic;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.Message;
Expand Down Expand Up @@ -89,8 +93,10 @@ public void testCurrentLedgerRolloverIfFull() throws Exception {
consumer.acknowledge(msg);
}

MLDataFormats.ManagedLedgerInfo.LedgerInfo lastLh =
managedLedger.getLedgersInfoAsList().get(managedLedger.getLedgersInfoAsList().size() - 1);
// all the messages have been acknowledged
// and all the ledgers have been removed except the the last ledger
// and all the ledgers have been removed except the last ledger
Awaitility.await()
.pollInterval(Duration.ofMillis(500L))
.untilAsserted(() -> {
Expand All @@ -104,12 +110,41 @@ public void testCurrentLedgerRolloverIfFull() throws Exception {
stateUpdater.set(managedLedger, ManagedLedgerImpl.State.LedgerOpened);
managedLedger.rollCurrentLedgerIfFull();

// the last ledger will be closed and removed and we have one ledger for empty
// If there are no pending write messages, the last ledger will be closed and still held.
Awaitility.await()
.pollInterval(Duration.ofMillis(1000L))
.untilAsserted(() -> {
Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1);
Assert.assertEquals(managedLedger.getTotalSize(), 0);
Assert.assertEquals(lastLh.getLedgerId(),
managedLedger.getLedgersInfoAsList().get(0).getLedgerId());
});
producer.send(new byte[1024 * 1024]);
Message<byte[]> msg = consumer.receive(2, TimeUnit.SECONDS);
Assert.assertNotNull(msg);
consumer.acknowledge(msg);
// Assert that we got a new ledger and all but the current ledger are deleted
Awaitility.await()
.untilAsserted(()-> {
Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1);
Assert.assertNotEquals(lastLh.getLedgerId(),
managedLedger.getLedgersInfoAsList().get(0).getLedgerId());
});
long lastLhIdAfterRolloverAndSendAgain = managedLedger.getLedgersInfoAsList().get(0).getLedgerId();

// Mock pendingAddEntries
OpAddEntry op = OpAddEntry.
createNoRetainBuffer(managedLedger, ByteBufAllocator.DEFAULT.buffer(128).retain(), null, null);
Field pendingAddEntries = managedLedger.getClass().getDeclaredField("pendingAddEntries");
pendingAddEntries.setAccessible(true);
ConcurrentLinkedQueue<OpAddEntry> queue = (ConcurrentLinkedQueue<OpAddEntry>) pendingAddEntries.get(managedLedger);
queue.add(op);
// When ml has pending write messages, ml will create a new ledger and close and delete the previous ledger
Awaitility.await()
.untilAsserted(()-> {
managedLedger.rollCurrentLedgerIfFull();
Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1);
Assert.assertNotEquals(managedLedger.getLedgersInfoAsList().get(0).getLedgerId(),
lastLhIdAfterRolloverAndSendAgain);
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicReferenceFieldUpdater;

import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.State.ClosedLedger;
import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.State.WriteFailed;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertTrue;
Expand Down Expand Up @@ -168,8 +169,9 @@ public void testRecoverSequenceId(boolean isUseManagedLedgerProperties) throws E
stateUpdater.setAccessible(true);
stateUpdater.set(managedLedger, ManagedLedgerImpl.State.LedgerOpened);
managedLedger.rollCurrentLedgerIfFull();
Awaitility.await().until(() -> {
return !managedLedger.ledgerExists(position.getLedgerId());
Awaitility.await().untilAsserted(() -> {
Assert.assertTrue(managedLedger.ledgerExists(position.getLedgerId()));
Assert.assertEquals(managedLedger.getState(), ClosedLedger);
});
}
mlTransactionLog.closeAsync().get();
Expand Down