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)); }
