This is an automated email from the ASF dual-hosted git repository.

RongtongJin pushed a commit to branch codex/dledger-latest-pr336-adapter
in repository https://gitbox.apache.org/repos/asf/rocketmq.git

commit 802a8b0f48efabd5c5b491fd616042a6b166c5be
Author: 通融 <[email protected]>
AuthorDate: Sat Aug 15 17:39:55 2026 +0800

    build: adapt RocketMQ to latest DLedger
---
 WORKSPACE                                          |   6 +-
 .../broker/dledger/DLedgerRoleChangeHandler.java   |   2 +-
 common/pom.xml                                     |   4 -
 .../controller/impl/DLedgerController.java         |   2 +-
 .../impl/DLedgerControllerStateMachine.java        |   8 +-
 pom.xml                                            |  10 +-
 remoting/BUILD.bazel                               |   1 -
 .../protocol/RemotingSerializableCompatTest.java   | 152 ++++++++++++++++++---
 .../rocketmq/store/dledger/DLedgerCommitLog.java   |  94 +++++++++++--
 9 files changed, 228 insertions(+), 51 deletions(-)

diff --git a/WORKSPACE b/WORKSPACE
index 0e95cd42e8..eb8b88c9dd 100644
--- a/WORKSPACE
+++ b/WORKSPACE
@@ -40,8 +40,7 @@ load("@rules_jvm_external//:defs.bzl", "maven_install")
 maven_install(
     artifacts = [
         "junit:junit:4.13.2",
-        "com.alibaba:fastjson:1.2.83",
-        "com.alibaba.fastjson2:fastjson2:2.0.59",
+        "com.alibaba.fastjson2:fastjson2:2.0.64",
         "org.hamcrest:hamcrest-library:1.3",
         "io.netty:netty-all:4.1.130.Final",
         "org.assertj:assertj-core:3.22.0",
@@ -54,7 +53,7 @@ maven_install(
         "commons-validator:commons-validator:1.10.0",
         "org.apache.commons:commons-lang3:3.20.0",
         "org.hamcrest:hamcrest-core:1.3",
-        "io.openmessaging.storage:dledger:0.3.2",
+        "io.openmessaging.storage:dledger:0.3.3-pr336-f2-64-SNAPSHOT",
         "net.java.dev.jna:jna:4.2.2",
         "ch.qos.logback:logback-classic:1.2.10",
         "ch.qos.logback:logback-core:1.2.10",
@@ -117,6 +116,7 @@ maven_install(
         "org.slf4j:slf4j-api:2.0.3",
         "org.javassist:javassist:3.20.0-GA",
     ],
+    excluded_artifacts = ["org.apache.rocketmq:rocketmq-remoting"],
     fetch_sources = False,
     repositories = [
         "https://repo1.maven.org/maven2";,
diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/dledger/DLedgerRoleChangeHandler.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/dledger/DLedgerRoleChangeHandler.java
index e6cb97640b..4daac4eca5 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/dledger/DLedgerRoleChangeHandler.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/dledger/DLedgerRoleChangeHandler.java
@@ -80,7 +80,7 @@ public class DLedgerRoleChangeHandler implements 
DLedgerLeaderElector.RoleChange
                                 if 
(dLegerServer.getDLedgerStore().getLedgerEndIndex() == -1) {
                                     break;
                                 }
-                                if 
(dLegerServer.getDLedgerStore().getLedgerEndIndex() == 
dLegerServer.getDLedgerStore().getCommittedIndex()
+                                if 
(dLegerServer.getDLedgerStore().getLedgerEndIndex() == 
dLegerServer.getMemberState().getCommittedIndex()
                                     && messageStore.dispatchBehindBytes() == 
0) {
                                     break;
                                 }
diff --git a/common/pom.xml b/common/pom.xml
index b931afcb89..caa9d6b537 100644
--- a/common/pom.xml
+++ b/common/pom.xml
@@ -32,10 +32,6 @@
     </properties>
 
     <dependencies>
-        <dependency>
-            <groupId>com.alibaba</groupId>
-            <artifactId>fastjson</artifactId>
-        </dependency>
         <dependency>
             <groupId>com.alibaba.fastjson2</groupId>
             <artifactId>fastjson2</artifactId>
diff --git 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerController.java
 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerController.java
index 3421010340..8fd060445c 100644
--- 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerController.java
+++ 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerController.java
@@ -17,11 +17,11 @@
 package org.apache.rocketmq.controller.impl;
 
 import com.google.common.base.Stopwatch;
-import io.openmessaging.storage.dledger.AppendFuture;
 import io.openmessaging.storage.dledger.DLedgerConfig;
 import io.openmessaging.storage.dledger.DLedgerLeaderElector;
 import io.openmessaging.storage.dledger.DLedgerServer;
 import io.openmessaging.storage.dledger.MemberState;
+import io.openmessaging.storage.dledger.common.AppendFuture;
 import io.openmessaging.storage.dledger.protocol.AppendEntryRequest;
 import io.openmessaging.storage.dledger.protocol.AppendEntryResponse;
 import io.openmessaging.storage.dledger.protocol.BatchAppendEntryRequest;
diff --git 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerControllerStateMachine.java
 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerControllerStateMachine.java
index f67967e960..c30792c5a1 100644
--- 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerControllerStateMachine.java
+++ 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerControllerStateMachine.java
@@ -20,7 +20,8 @@ import io.openmessaging.storage.dledger.entry.DLedgerEntry;
 import io.openmessaging.storage.dledger.exception.DLedgerException;
 import io.openmessaging.storage.dledger.snapshot.SnapshotReader;
 import io.openmessaging.storage.dledger.snapshot.SnapshotWriter;
-import io.openmessaging.storage.dledger.statemachine.CommittedEntryIterator;
+import io.openmessaging.storage.dledger.statemachine.ApplyEntry;
+import io.openmessaging.storage.dledger.statemachine.ApplyEntryIterator;
 import io.openmessaging.storage.dledger.statemachine.StateMachine;
 import org.apache.rocketmq.common.constant.LoggerName;
 import org.apache.rocketmq.controller.impl.event.EventMessage;
@@ -51,12 +52,13 @@ public class DLedgerControllerStateMachine implements 
StateMachine {
     }
 
     @Override
-    public void onApply(CommittedEntryIterator iterator) {
+    public void onApply(ApplyEntryIterator iterator) {
         int applyingSize = 0;
         long firstApplyIndex = -1;
         long lastApplyIndex = -1;
         while (iterator.hasNext()) {
-            final DLedgerEntry entry = iterator.next();
+            final ApplyEntry applyEntry = iterator.next();
+            final DLedgerEntry entry = applyEntry.getEntry();
             final byte[] body = entry.getBody();
             if (body != null && body.length > 0) {
                 final EventMessage event = 
this.eventSerializer.deserialize(body);
diff --git a/pom.xml b/pom.xml
index 645ad51225..c084989117 100644
--- a/pom.xml
+++ b/pom.xml
@@ -104,8 +104,7 @@
         <netty.version>4.1.130.Final</netty.version>
         <netty.tcnative.version>2.0.53.Final</netty.tcnative.version>
         <bcpkix-jdk18on.version>1.83</bcpkix-jdk18on.version>
-        <fastjson.version>1.2.83</fastjson.version>
-        <fastjson2.version>2.0.63</fastjson2.version>
+        <fastjson2.version>2.0.64</fastjson2.version>
         <javassist.version>3.20.0-GA</javassist.version>
         <jna.version>4.2.2</jna.version>
         <commons-lang3.version>3.20.0</commons-lang3.version>
@@ -123,7 +122,7 @@
         <lz4-java.version>1.10.3</lz4-java.version>
         <opentracing.version>0.33.0</opentracing.version>
         <jaeger.version>1.8.1</jaeger.version>
-        <dleger.version>0.3.2</dleger.version>
+        <dleger.version>0.3.3-pr336-f2-64-SNAPSHOT</dleger.version>
         <annotations-api.version>6.0.53</annotations-api.version>
         <extra-enforcer-rules.version>1.0-beta-4</extra-enforcer-rules.version>
         
<concurrentlinkedhashmap-lru.version>1.4.2</concurrentlinkedhashmap-lru.version>
@@ -688,11 +687,6 @@
                 <type>jar</type>
                 <version>${bcpkix-jdk18on.version}</version>
             </dependency>
-            <dependency>
-                <groupId>com.alibaba</groupId>
-                <artifactId>fastjson</artifactId>
-                <version>${fastjson.version}</version>
-            </dependency>
             <dependency>
                 <groupId>com.alibaba.fastjson2</groupId>
                 <artifactId>fastjson2</artifactId>
diff --git a/remoting/BUILD.bazel b/remoting/BUILD.bazel
index 62273e5e9d..2d91b53c9e 100644
--- a/remoting/BUILD.bazel
+++ b/remoting/BUILD.bazel
@@ -52,7 +52,6 @@ java_library(
         "//common",
         "//:test_deps",
         "@maven//:org_objenesis_objenesis",
-        "@maven//:com_alibaba_fastjson",
         "@maven//:com_alibaba_fastjson2_fastjson2",
         "@maven//:com_google_code_gson_gson",
         "@maven//:com_google_guava_guava",
diff --git 
a/remoting/src/test/java/org/apache/rocketmq/remoting/protocol/RemotingSerializableCompatTest.java
 
b/remoting/src/test/java/org/apache/rocketmq/remoting/protocol/RemotingSerializableCompatTest.java
index 35c1c7b891..901e5a0dff 100644
--- 
a/remoting/src/test/java/org/apache/rocketmq/remoting/protocol/RemotingSerializableCompatTest.java
+++ 
b/remoting/src/test/java/org/apache/rocketmq/remoting/protocol/RemotingSerializableCompatTest.java
@@ -17,17 +17,27 @@
 
 package org.apache.rocketmq.remoting.protocol;
 
-import com.alibaba.fastjson.annotation.JSONField;
-import com.alibaba.fastjson2.JSON;
+import com.alibaba.fastjson2.annotation.JSONField;
+import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
 import org.apache.rocketmq.remoting.protocol.body.BatchAck;
+import org.apache.rocketmq.remoting.protocol.body.Connection;
+import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection;
+import org.apache.rocketmq.remoting.protocol.heartbeat.ConsumeType;
+import org.apache.rocketmq.remoting.protocol.heartbeat.MessageModel;
+import org.apache.rocketmq.remoting.protocol.heartbeat.SubscriptionData;
 import org.junit.Test;
 import org.objenesis.ObjenesisStd;
 import org.reflections.Reflections;
 
+import java.io.BufferedReader;
+import java.io.File;
+import java.io.InputStreamReader;
 import java.lang.reflect.Array;
 import java.lang.reflect.Field;
 import java.lang.reflect.Modifier;
+import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.BitSet;
 import java.util.HashMap;
 import java.util.HashSet;
@@ -38,12 +48,31 @@ import java.util.Random;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
 
 import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
 
 public class RemotingSerializableCompatTest {
+
+    private static final String FASTJSON1_BATCH_ACK =
+        "{\"b\":\"Kg==\",\"c\":\"fixture-consumer\",\"it\":60000,"
+            + "\"pt\":1700000000123,\"q\":7,\"r\":\"0\",\"rq\":3,"
+            + "\"so\":1234567890123,\"t\":\"FixtureTopic\"}";
+    private static final String FASTJSON1_SUBSCRIPTION_DATA =
+        
"{\"classFilterMode\":true,\"codeSet\":[101,202],\"expressionType\":\"SQL92\","
+            + "\"subString\":\"TagA || TagB\",\"subVersion\":1700000000456,"
+            + "\"tagsSet\":[\"TagA\",\"TagB\"],\"topic\":\"FixtureTopic\"}";
+    private static final String FASTJSON1_CONSUMER_CONNECTION =
+        "{\"connectionSet\":[{\"clientAddr\":\"127.0.0.1:10911\","
+            + 
"\"clientId\":\"fixture-client@instance-1\",\"language\":\"GO\",\"version\":433}],"
+            + "\"consumeFromWhere\":\"CONSUME_FROM_TIMESTAMP\","
+            + 
"\"consumeType\":\"CONSUME_PASSIVELY\",\"messageModel\":\"CLUSTERING\","
+            + "\"subscriptionTable\":{\"FixtureTopic\":"
+            + FASTJSON1_SUBSCRIPTION_DATA + "}}";
     
     @Test
     public void testCompatibilityCheck() {
@@ -65,28 +94,106 @@ public class RemotingSerializableCompatTest {
                 fillDefaultFields(instance, clazz);
                 assertTrue(checkCompatible(instance, clazz));
             } catch (Exception e) {
-                System.err.printf("Class %s: incompatible, error: %s\n", 
clazz.getName(), e.getMessage());
+                throw new AssertionError("Class " + clazz.getName() + " could 
not be checked", e);
             }
         }
     }
 
     @Test
-    public void testCompatibilityCheckWithBitSet() {
+    public void testFastjson1BatchAckFixture() {
         BitSet bitSet = new BitSet();
         bitSet.set(1);
         bitSet.set(3);
         bitSet.set(5);
-        String fastjson1Str = 
"{\"b\":\"Kg==\",\"c\":\"DEFAULT_CONSUMER\",\"it\":5000,\"pt\":1760694281326,\"q\":1,\"r\":\"0\",\"rq\":2,\"so\":100,\"t\":\"myTopic\"}";
-        BatchAck batchAck = JSON.parseObject(fastjson1Str, BatchAck.class);
+        BatchAck batchAck = RemotingSerializable.fromJson(FASTJSON1_BATCH_ACK, 
BatchAck.class);
         assertEquals(bitSet, batchAck.getBitSet());
-        assertEquals("DEFAULT_CONSUMER", batchAck.getConsumerGroup());
-        assertEquals(5000, batchAck.getInvisibleTime());
-        assertEquals(1760694281326L, batchAck.getPopTime());
-        assertEquals(1, batchAck.getQueueId());
+        assertEquals("fixture-consumer", batchAck.getConsumerGroup());
+        assertEquals(60000, batchAck.getInvisibleTime());
+        assertEquals(1700000000123L, batchAck.getPopTime());
+        assertEquals(7, batchAck.getQueueId());
         assertEquals("0", batchAck.getRetry());
-        assertEquals(2, batchAck.getReviveQueueId());
-        assertEquals(100, batchAck.getStartOffset());
-        assertEquals("myTopic", batchAck.getTopic());
+        assertEquals(3, batchAck.getReviveQueueId());
+        assertEquals(1234567890123L, batchAck.getStartOffset());
+        assertEquals("FixtureTopic", batchAck.getTopic());
+    }
+
+    @Test
+    public void testFastjson1SubscriptionDataFixture() {
+        SubscriptionData subscriptionData = RemotingSerializable.fromJson(
+            FASTJSON1_SUBSCRIPTION_DATA, SubscriptionData.class);
+        assertSubscriptionData(subscriptionData);
+    }
+
+    @Test
+    public void testFastjson1ConsumerConnectionFixture() {
+        ConsumerConnection consumerConnection = RemotingSerializable.fromJson(
+            FASTJSON1_CONSUMER_CONNECTION, ConsumerConnection.class);
+        assertEquals(ConsumeFromWhere.CONSUME_FROM_TIMESTAMP, 
consumerConnection.getConsumeFromWhere());
+        assertEquals(ConsumeType.CONSUME_PASSIVELY, 
consumerConnection.getConsumeType());
+        assertEquals(MessageModel.CLUSTERING, 
consumerConnection.getMessageModel());
+        assertEquals(1, consumerConnection.getConnectionSet().size());
+        Connection connection = 
consumerConnection.getConnectionSet().iterator().next();
+        assertEquals("127.0.0.1:10911", connection.getClientAddr());
+        assertEquals("fixture-client@instance-1", connection.getClientId());
+        assertEquals(LanguageCode.GO, connection.getLanguage());
+        assertEquals(433, connection.getVersion());
+        assertEquals(433, consumerConnection.computeMinVersion());
+        assertEquals(new HashSet<>(Arrays.asList("FixtureTopic")),
+            consumerConnection.getSubscriptionTable().keySet());
+        
assertSubscriptionData(consumerConnection.getSubscriptionTable().get("FixtureTopic"));
+    }
+
+    @Test
+    public void testRemotingCodecColdStart() throws Exception {
+        String javaExecutable = System.getProperty("java.home")
+            + File.separator + "bin" + File.separator + "java";
+        String classPath = System.getProperty(
+            "surefire.test.class.path", System.getProperty("java.class.path"));
+
+        ProcessBuilder processBuilder = new ProcessBuilder(
+            javaExecutable, "-cp", classPath, ColdStartProbe.class.getName());
+        processBuilder.environment().remove("JAVA_TOOL_OPTIONS");
+        processBuilder.environment().remove("_JAVA_OPTIONS");
+        processBuilder.environment().remove("JDK_JAVA_OPTIONS");
+        Process process = processBuilder.redirectErrorStream(true).start();
+        boolean finished = process.waitFor(10, TimeUnit.SECONDS);
+        if (!finished) {
+            process.destroyForcibly();
+            fail("Cold-start probe did not finish");
+        }
+
+        StringBuilder output = new StringBuilder();
+        try (BufferedReader reader = new BufferedReader(
+            new InputStreamReader(process.getInputStream(), 
StandardCharsets.UTF_8))) {
+            String line;
+            while ((line = reader.readLine()) != null) {
+                output.append(line).append(System.lineSeparator());
+            }
+        }
+        assertEquals(output.toString(), 0, process.exitValue());
+    }
+
+    public static final class ColdStartProbe {
+        private ColdStartProbe() {
+        }
+
+        public static void main(String[] args) {
+            try {
+                ConsumerConnection connection = new ConsumerConnection();
+                String json = RemotingSerializable.toJson(connection, false);
+                ConsumerConnection decoded = RemotingSerializable.fromJson(
+                    json, ConsumerConnection.class);
+                if (decoded == null || decoded.getConnectionSet() == null) {
+                    throw new AssertionError(
+                        "Remoting codec returned an incomplete 
ConsumerConnection: " + json);
+                }
+                Runtime.getRuntime().halt(0);
+            } catch (Throwable t) {
+                t.printStackTrace(System.err);
+                System.err.flush();
+                Runtime.getRuntime().halt(1);
+            }
+        }
     }
     
     private void fillDefaultFields(final Object obj, final Class<?> clazz) 
throws Exception {
@@ -94,7 +201,7 @@ public class RemotingSerializableCompatTest {
             return;
         }
         for (Field field : clazz.getDeclaredFields()) {
-            if (Modifier.isStatic(field.getModifiers())) {
+            if (Modifier.isStatic(field.getModifiers()) || 
Modifier.isTransient(field.getModifiers())) {
                 continue;
             }
             field.setAccessible(true);
@@ -273,7 +380,7 @@ public class RemotingSerializableCompatTest {
         Class<?> clazz = original.getClass();
         boolean result = true;
         for (Field field : clazz.getDeclaredFields()) {
-            if (Modifier.isStatic(field.getModifiers())) {
+            if (Modifier.isStatic(field.getModifiers()) || 
Modifier.isTransient(field.getModifiers())) {
                 continue;
             }
             JSONField jsonField = field.getAnnotation(JSONField.class);
@@ -408,16 +515,27 @@ public class RemotingSerializableCompatTest {
     }
     
     private boolean checkCompatible(final Object original, final Class<?> 
clazz) {
-        String json = com.alibaba.fastjson.JSON.toJSONString(original);
+        String json = RemotingSerializable.toJson(original, false);
         Object deserialized;
         try {
-            deserialized = com.alibaba.fastjson2.JSON.parseObject(json, clazz);
+            deserialized = RemotingSerializable.fromJson(json, clazz);
         } catch (Exception e) {
             System.err.printf("Deserialization failed for %s: %s\n", 
clazz.getName(), e.getMessage());
             return false;
         }
         return checkCompatible(original, deserialized, clazz.getSimpleName(), 
new HashMap<>());
     }
+
+    private void assertSubscriptionData(final SubscriptionData 
subscriptionData) {
+        assertTrue(subscriptionData.isClassFilterMode());
+        assertEquals("FixtureTopic", subscriptionData.getTopic());
+        assertEquals("TagA || TagB", subscriptionData.getSubString());
+        assertEquals(new HashSet<>(Arrays.asList("TagA", "TagB")), 
subscriptionData.getTagsSet());
+        assertEquals(new HashSet<>(Arrays.asList(101, 202)), 
subscriptionData.getCodeSet());
+        assertEquals(1700000000456L, subscriptionData.getSubVersion());
+        assertEquals("SQL92", subscriptionData.getExpressionType());
+        assertNull(subscriptionData.getFilterClassSource());
+    }
     
     private <T> T allocateInstance(final Class<T> clazz) {
         return new ObjenesisStd().newInstance(clazz);
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 34fdcf1b6c..06f9e0dc42 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
@@ -16,11 +16,13 @@
  */
 package org.apache.rocketmq.store.dledger;
 
-import io.openmessaging.storage.dledger.AppendFuture;
-import io.openmessaging.storage.dledger.BatchAppendFuture;
 import io.openmessaging.storage.dledger.DLedgerConfig;
 import io.openmessaging.storage.dledger.DLedgerServer;
+import io.openmessaging.storage.dledger.common.AppendFuture;
+import io.openmessaging.storage.dledger.common.BatchAppendFuture;
 import io.openmessaging.storage.dledger.entry.DLedgerEntry;
+import io.openmessaging.storage.dledger.entry.DLedgerEntryCoder;
+import io.openmessaging.storage.dledger.entry.DLedgerIndexEntry;
 import io.openmessaging.storage.dledger.protocol.AppendEntryRequest;
 import io.openmessaging.storage.dledger.protocol.AppendEntryResponse;
 import io.openmessaging.storage.dledger.protocol.BatchAppendEntryRequest;
@@ -69,6 +71,8 @@ public class DLedgerCommitLog extends CommitLog {
     private final DLedgerConfig dLedgerConfig;
     private final DLedgerMmapFileStore dLedgerFileStore;
     private final MmapFileList dLedgerFileList;
+    private volatile long cachedCommittedIndex = Long.MIN_VALUE;
+    private volatile long cachedCommittedPos = -1;
 
     //The id identifies the broker role, 0 means master, others means slave
     private final int id;
@@ -98,7 +102,7 @@ public class DLedgerCommitLog extends CommitLog {
         
dLedgerConfig.setDeleteWhen(defaultMessageStore.getMessageStoreConfig().getDeleteWhen());
         
dLedgerConfig.setFileReservedHours(defaultMessageStore.getMessageStoreConfig().getFileReservedTime()
 + 1);
         
dLedgerConfig.setPreferredLeaderId(defaultMessageStore.getMessageStoreConfig().getPreferredLeaderId());
-        
dLedgerConfig.setEnableBatchPush(defaultMessageStore.getMessageStoreConfig().isEnableBatchPush());
+        
dLedgerConfig.setEnableBatchAppend(defaultMessageStore.getMessageStoreConfig().isEnableBatchPush());
         
dLedgerConfig.setDiskSpaceRatioToCheckExpired(defaultMessageStore.getMessageStoreConfig().getDiskMaxUsedSpaceRatio()
 / 100f);
 
         id = Integer.parseInt(dLedgerConfig.getSelfId().substring(1)) + 1;
@@ -149,15 +153,50 @@ public class DLedgerCommitLog extends CommitLog {
 
     @Override
     public long getMaxOffset() {
-        if (dLedgerFileStore.getCommittedPos() > 0) {
-            return dLedgerFileStore.getCommittedPos();
+        long committedPos = getCommittedPos();
+        if (committedPos > 0) {
+            return committedPos;
         }
-        if (dLedgerFileList.getMinOffset() > 0) {
+        if (committedPos == 0 && dLedgerFileList.getMinOffset() > 0) {
             return dLedgerFileList.getMinOffset();
         }
         return 0;
     }
 
+    long getCommittedPos() {
+        long committedIndex = 
dLedgerServer.getMemberState().getCommittedIndex();
+        if (committedIndex < 0) {
+            return -1;
+        }
+        if (committedIndex == cachedCommittedIndex) {
+            return cachedCommittedPos;
+        }
+
+        SelectMmapBufferResult indexBuffer = null;
+        try {
+            indexBuffer = dLedgerFileStore.getIndexFileList().getData(
+                committedIndex * DLedgerMmapFileStore.INDEX_UNIT_SIZE,
+                DLedgerMmapFileStore.INDEX_UNIT_SIZE);
+            if (indexBuffer == null) {
+                return -1;
+            }
+            DLedgerIndexEntry indexEntry =
+                DLedgerEntryCoder.decodeIndex(indexBuffer.getByteBuffer());
+            if (indexEntry.getIndex() != committedIndex) {
+                return -1;
+            }
+            long committedPos = indexEntry.getPosition() + 
indexEntry.getSize();
+            cachedCommittedPos = committedPos;
+            cachedCommittedIndex = committedIndex;
+            return committedPos;
+        } catch (RuntimeException e) {
+            log.warn("Failed to resolve committed position for index={}", 
committedIndex, e);
+            return -1;
+        } finally {
+            SelectMmapBufferResult.release(indexBuffer);
+        }
+    }
+
     @Override
     public long getMinOffset() {
         if (!mappedFileQueue.getMappedFiles().isEmpty()) {
@@ -232,11 +271,19 @@ public class DLedgerCommitLog extends CommitLog {
     }
 
     public SelectMmapBufferResult truncate(SelectMmapBufferResult sbr) {
-        long committedPos = dLedgerFileStore.getCommittedPos();
-        if (sbr == null || sbr.getStartOffset() == committedPos) {
+        long committedPos = getCommittedPos();
+        return truncate(sbr, committedPos);
+    }
+
+    private SelectMmapBufferResult truncate(SelectMmapBufferResult sbr, long 
committedPos) {
+        if (sbr == null) {
+            return null;
+        }
+        if (committedPos < 0 || sbr.getStartOffset() >= committedPos) {
+            SelectMmapBufferResult.release(sbr);
             return null;
         }
-        if (sbr.getStartOffset() + sbr.getSize() <= committedPos) {
+        if (sbr.getSize() <= committedPos - sbr.getStartOffset()) {
             return sbr;
         } else {
             sbr.setSize((int) (committedPos - sbr.getStartOffset()));
@@ -257,7 +304,8 @@ public class DLedgerCommitLog extends CommitLog {
         if (offset < dividedCommitlogOffset) {
             return super.getData(offset, returnFirstOnNotFound);
         }
-        if (offset >= dLedgerFileStore.getCommittedPos()) {
+        long committedPos = getCommittedPos();
+        if (committedPos < 0 || offset >= committedPos) {
             return null;
         }
         int mappedFileSize = 
this.dLedgerServer.getdLedgerConfig().getMappedFileSizeForEntryData();
@@ -265,7 +313,7 @@ public class DLedgerCommitLog extends CommitLog {
         if (mappedFile != null) {
             int pos = (int) (offset % mappedFileSize);
             SelectMmapBufferResult sbr = mappedFile.selectMappedBuffer(pos);
-            return convertSbr(truncate(sbr));
+            return convertSbr(truncate(sbr, committedPos));
         }
 
         return null;
@@ -276,14 +324,26 @@ public class DLedgerCommitLog extends CommitLog {
         if (offset < dividedCommitlogOffset) {
             return super.getData(offset, size, byteBuffer);
         }
-        if (offset >= dLedgerFileStore.getCommittedPos()) {
+        long committedPos = getCommittedPos();
+        if (committedPos < 0 || offset >= committedPos || size < 0 || 
byteBuffer.remaining() < size
+            || size > committedPos - offset) {
             return false;
         }
         int mappedFileSize = 
this.dLedgerServer.getdLedgerConfig().getMappedFileSizeForEntryData();
         MmapFile mappedFile = 
this.dLedgerFileList.findMappedFileByOffset(offset, offset == 0);
         if (mappedFile != null) {
+            long selectedOffset = mappedFile.getFileFromOffset() + (offset % 
mappedFileSize);
+            if (selectedOffset >= committedPos || size > committedPos - 
selectedOffset) {
+                return false;
+            }
             int pos = (int) (offset % mappedFileSize);
-            return mappedFile.getData(pos, size, byteBuffer);
+            int originalLimit = byteBuffer.limit();
+            byteBuffer.limit(byteBuffer.position() + size);
+            try {
+                return mappedFile.getData(pos, size, byteBuffer);
+            } finally {
+                byteBuffer.limit(originalLimit);
+            }
         }
         return false;
     }
@@ -787,9 +847,17 @@ public class DLedgerCommitLog extends CommitLog {
         if (offset < dividedCommitlogOffset) {
             return super.getMessage(offset, size);
         }
+        long committedPos = getCommittedPos();
+        if (committedPos < 0 || offset >= committedPos || size < 0 || size > 
committedPos - offset) {
+            return null;
+        }
         int mappedFileSize = 
this.dLedgerServer.getdLedgerConfig().getMappedFileSizeForEntryData();
         MmapFile mappedFile = 
this.dLedgerFileList.findMappedFileByOffset(offset, offset == 0);
         if (mappedFile != null) {
+            long selectedOffset = mappedFile.getFileFromOffset() + (offset % 
mappedFileSize);
+            if (selectedOffset >= committedPos || size > committedPos - 
selectedOffset) {
+                return null;
+            }
             int pos = (int) (offset % mappedFileSize);
             return convertSbr(mappedFile.selectMappedBuffer(pos, size));
         }

Reply via email to