diff --git a/store/src/main/java/org/apache/rocketmq/store/dledger/DLedgerCommitLog.java b/store/src/main/java/org/apache/rocketmq/store/dledger/DLedgerCommitLog.java index 142aab7af8..2c37331905 100644 --- a/store/src/main/java/org/apache/rocketmq/store/dledger/DLedgerCommitLog.java +++ b/store/src/main/java/org/apache/rocketmq/store/dledger/DLedgerCommitLog.java @@ -162,7 +162,11 @@ public class DLedgerCommitLog extends CommitLog { if (committedPos > 0) { return committedPos; } - if (committedPos == 0 && dLedgerFileList.getMinOffset() > 0) { + // getCommittedPos() yields a non-positive value when no entry has been committed yet, + // which is the normal state right after a DLedger store is layered on top of legacy + // commitlog data. In that case the divided commitlog offset, i.e. the min offset of the + // DLedger file list, is the real max offset. + if (dLedgerFileList.getMinOffset() > 0) { return dLedgerFileList.getMinOffset(); } return 0; diff --git a/store/src/test/java/org/apache/rocketmq/store/StoreTestBase.java b/store/src/test/java/org/apache/rocketmq/store/StoreTestBase.java index ae0841b8b8..2a7d9d8f5e 100644 --- a/store/src/test/java/org/apache/rocketmq/store/StoreTestBase.java +++ b/store/src/test/java/org/apache/rocketmq/store/StoreTestBase.java @@ -24,8 +24,10 @@ import org.apache.rocketmq.common.message.MessageExtBrokerInner; import org.junit.After; import java.io.File; +import java.io.IOException; import java.net.InetAddress; import java.net.InetSocketAddress; +import java.net.ServerSocket; import java.net.SocketAddress; import java.net.UnknownHostException; import java.util.ArrayList; @@ -47,8 +49,27 @@ public class StoreTestBase { private static AtomicInteger port = new AtomicInteger(30000); + private static final int MAX_PORT_PROBE_ATTEMPTS = 200; + public static synchronized int nextPort() { - return port.addAndGet(5); + for (int i = 0; i < MAX_PORT_PROBE_ATTEMPTS; i++) { + int candidate = port.addAndGet(5); + if (isPortAvailable(candidate)) { + return candidate; + } + } + throw new IllegalStateException("Failed to find an available port after " + + MAX_PORT_PROBE_ATTEMPTS + " attempts, last tried " + port.get()); + } + + private static boolean isPortAvailable(int candidate) { + try (ServerSocket serverSocket = new ServerSocket()) { + serverSocket.setReuseAddress(false); + serverSocket.bind(new InetSocketAddress(candidate)); + return true; + } catch (IOException e) { + return false; + } } protected MessageExtBatch buildBatchMessage(int size) {