diff --git a/WORKSPACE b/WORKSPACE index 0e95cd42e85..eb8b88c9dd3 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 e6cb97640bd..4daac4eca5e 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 void run() { 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 b931afcb89c..caa9d6b537f 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -32,10 +32,6 @@ - - com.alibaba - fastjson - com.alibaba.fastjson2 fastjson2 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 3421010340a..8fd060445c8 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 f67967e9600..c30792c5a1d 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.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 String generateDLedgerId(String dLedgerGroupId, String dLedgerSelfId) { } @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/controller/src/test/java/org/apache/rocketmq/controller/impl/DLedgerControllerTest.java b/controller/src/test/java/org/apache/rocketmq/controller/impl/DLedgerControllerTest.java index 32e7859a58a..6f448a70752 100644 --- a/controller/src/test/java/org/apache/rocketmq/controller/impl/DLedgerControllerTest.java +++ b/controller/src/test/java/org/apache/rocketmq/controller/impl/DLedgerControllerTest.java @@ -40,6 +40,7 @@ import java.io.File; import java.time.Duration; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashSet; import java.util.List; import java.util.Set; @@ -60,9 +61,15 @@ import static org.junit.Assert.assertTrue; public class DLedgerControllerTest { + private static int port = 30000; private List baseDirs; private List controllers; + private static synchronized int nextPort() { + port += 10; + return port; + } + public DLedgerController launchController(final String group, final String peers, final String selfId, final boolean isEnableElectUncleanMaster) { String tmpdir = System.getProperty("java.io.tmpdir"); @@ -172,7 +179,9 @@ public DLedgerController waitLeader(final List controllers) t public DLedgerController mockMetaData(boolean enableElectUncleanMaster) throws Exception { String group = UUID.randomUUID().toString(); - String peers = String.format("n0-localhost:%d;n1-localhost:%d;n2-localhost:%d", 30000, 30001, 30002); + int basePort = nextPort(); + String peers = String.format("n0-localhost:%d;n1-localhost:%d;n2-localhost:%d", + basePort, basePort + 1, basePort + 2); DLedgerController c0 = launchController(group, peers, "n0", enableElectUncleanMaster); DLedgerController c1 = launchController(group, peers, "n1", enableElectUncleanMaster); DLedgerController c2 = launchController(group, peers, "n2", enableElectUncleanMaster); @@ -237,6 +246,36 @@ public void testElectMaster() throws Exception { assertNotEquals(DEFAULT_IP[0], response.getMasterAddress()); } + @Test + public void testRestartAllControllersRecoversStateBeforeNewEvent() throws Exception { + DLedgerController originalLeader = mockMetaData(false); + String group = originalLeader.getControllerConfig().getControllerDLegerGroup(); + String peers = originalLeader.getControllerConfig().getControllerDLegerPeers(); + List selfIds = controllers.stream() + .map(controller -> controller.getControllerConfig().getControllerDLegerSelfId()) + .collect(Collectors.toList()); + + List originalControllers = new ArrayList<>(controllers); + for (DLedgerController controller : originalControllers) { + controller.shutdown(); + } + controllers.clear(); + for (String selfId : selfIds) { + controllers.add(launchController(group, peers, selfId, false)); + } + DLedgerController restartedLeader = waitLeader(controllers); + + RemotingCommand response = restartedLeader + .getReplicaInfo(new GetReplicaInfoRequestHeader(DEFAULT_BROKER_NAME)).get(10, TimeUnit.SECONDS); + assertEquals(ResponseCode.SUCCESS, response.getCode()); + GetReplicaInfoResponseHeader replicaInfo = + (GetReplicaInfoResponseHeader) response.readCustomHeader(); + SyncStateSet syncStateSet = RemotingSerializable.decode(response.getBody(), SyncStateSet.class); + assertEquals(1L, replicaInfo.getMasterBrokerId().longValue()); + assertEquals(DEFAULT_IP[0], replicaInfo.getMasterAddress()); + assertEquals(new HashSet<>(Arrays.asList(1L, 2L, 3L)), syncStateSet.getSyncStateSet()); + } + @Test public void testBrokerLifecycleListener() throws Exception { final DLedgerController leader = mockMetaData(false); @@ -250,6 +289,11 @@ public void testBrokerLifecycleListener() throws Exception { dLedgerController.shutdown(); controllers.remove(dLedgerController); } + await().atMost(Duration.ofSeconds(10)).until(() -> + leader.getMemberState().getPeersLiveTable().size() == leader.getMemberState().peerSize() - 1 + && leader.getMemberState().getPeersLiveTable().values().stream() + .noneMatch(Boolean.TRUE::equals)); + await().atMost(Duration.ofSeconds(10)).until(() -> !leader.isLeaderState()); final ElectMasterRequestHeader request = ElectMasterRequestHeader.ofControllerTrigger(DEFAULT_BROKER_NAME); setBrokerElectPolicy(leader, 1L); diff --git a/pom.xml b/pom.xml index 645ad51225a..c0849891177 100644 --- a/pom.xml +++ b/pom.xml @@ -104,8 +104,7 @@ 4.1.130.Final 2.0.53.Final 1.83 - 1.2.83 - 2.0.63 + 2.0.64 3.20.0-GA 4.2.2 3.20.0 @@ -123,7 +122,7 @@ 1.10.3 0.33.0 1.8.1 - 0.3.2 + 0.3.3-pr336-f2-64-SNAPSHOT 6.0.53 1.0-beta-4 1.4.2 @@ -688,11 +687,6 @@ jar ${bcpkix-jdk18on.version} - - com.alibaba - fastjson - ${fastjson.version} - com.alibaba.fastjson2 fastjson2 diff --git a/remoting/BUILD.bazel b/remoting/BUILD.bazel index 62273e5e9d0..2d91b53c9e3 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 35c1c7b8912..901e5a0dff9 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.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 void testCompatibilityCheck() { 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 @@ private void fillDefaultFields(final Object obj, final Class clazz) throws Ex 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 @@ private boolean checkCompatible(final Object original, final Object deserialized 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 @@ private boolean isPrimitiveOrWrapper(final Class clazz) { } 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 allocateInstance(final Class clazz) { return new ObjenesisStd().newInstance(clazz); diff --git a/store/BUILD.bazel b/store/BUILD.bazel index 66af7d6b45c..510da0e044d 100644 --- a/store/BUILD.bazel +++ b/store/BUILD.bazel @@ -81,6 +81,7 @@ GenTestRules( "src/test/java/org/apache/rocketmq/store/DefaultMessageStoreTest", "src/test/java/org/apache/rocketmq/store/HATest", "src/test/java/org/apache/rocketmq/store/dledger/DLedgerCommitlogTest", + "src/test/java/org/apache/rocketmq/store/dledger/DLedgerLatestCommitLogTest", "src/test/java/org/apache/rocketmq/store/MappedFileQueueTest", "src/test/java/org/apache/rocketmq/store/queue/BatchConsumeMessageTest", "src/test/java/org/apache/rocketmq/store/dledger/MixCommitlogTest", diff --git a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java index 64ce41e47d5..00576f79f4e 100644 --- a/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java @@ -2741,16 +2741,19 @@ public void doReput() { if (dispatchRequest.isSuccess()) { if (size > 0) { - currentReputTimestamp = dispatchRequest.getStoreTimestamp(); - DefaultMessageStore.this.doDispatch(dispatchRequest); + if (dispatchRequest.getMsgSize() > 0) { + currentReputTimestamp = dispatchRequest.getStoreTimestamp(); + DefaultMessageStore.this.doDispatch(dispatchRequest); - if (isNotifyMessageArriveWhenReput()) { - notifyMessageArriveIfNecessary(dispatchRequest); + if (isNotifyMessageArriveWhenReput()) { + notifyMessageArriveIfNecessary(dispatchRequest); + } } this.reputFromOffset += size; readSize += size; - if (!DefaultMessageStore.this.getMessageStoreConfig().isDuplicationEnable() && + if (dispatchRequest.getMsgSize() > 0 + && !DefaultMessageStore.this.getMessageStoreConfig().isDuplicationEnable() && DefaultMessageStore.this.getMessageStoreConfig().getBrokerRole() == BrokerRole.SLAVE) { DefaultMessageStore.this.storeStatsService .getSinglePutMessageTopicTimesTotal(dispatchRequest.getTopic()).add(dispatchRequest.getBatchSize()); 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 34fdcf1b6c2..142aab7af89 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,14 @@ */ 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.DLedgerEntryType; +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 +72,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,13 +103,17 @@ public DLedgerCommitLog(final DefaultMessageStore defaultMessageStore) { 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.setEnableFastAdvanceCommitIndex(true); dLedgerConfig.setDiskSpaceRatioToCheckExpired(defaultMessageStore.getMessageStoreConfig().getDiskMaxUsedSpaceRatio() / 100f); id = Integer.parseInt(dLedgerConfig.getSelfId().substring(1)) + 1; dLedgerServer = new DLedgerServer(dLedgerConfig); dLedgerFileStore = (DLedgerMmapFileStore) dLedgerServer.getdLedgerStore(); DLedgerMmapFileStore.AppendHook appendHook = (entry, buffer, bodyOffset) -> { + if (entry.getMagic() != DLedgerEntryType.NORMAL.getMagic()) { + return; + } assert bodyOffset == DLedgerEntry.BODY_OFFSET; buffer.position(buffer.position() + bodyOffset + MessageDecoder.PHY_POS_POSITION); buffer.putLong(entry.getPos() + bodyOffset); @@ -149,15 +158,50 @@ public long flush() { @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 +276,19 @@ public SelectMappedBufferResult convertSbr(SelectMmapBufferResult sbr) { } 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 +309,8 @@ public SelectMappedBufferResult getData(final long offset, final boolean returnF 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 +318,7 @@ public SelectMappedBufferResult getData(final long offset, final boolean returnF 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 +329,26 @@ public boolean getData(final long offset, final int size, final ByteBuffer byteB 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; } @@ -346,19 +411,24 @@ private void dledgerRecoverAbnormally(long maxPhyOffsetOfConsumeQueue) throws Ro long mmapFileOffset = 0; while (true) { DispatchRequest dispatchRequest = this.checkMessageAndReturnSize(byteBuffer, checkCRCOnRecover, checkDupInfo); - int size = dispatchRequest.getMsgSize(); + int messageSize = dispatchRequest.getMsgSize(); + int entrySize = dispatchRequest.getBufferSize() == -1 + ? messageSize : dispatchRequest.getBufferSize(); if (dispatchRequest.isSuccess()) { - if (size > 0) { - mmapFileOffset += size; - if (this.defaultMessageStore.getMessageStoreConfig().isDuplicationEnable()) { - if (dispatchRequest.getCommitLogOffset() < this.defaultMessageStore.getConfirmOffset()) { + if (entrySize > 0) { + mmapFileOffset += entrySize; + if (messageSize > 0) { + if (this.defaultMessageStore.getMessageStoreConfig().isDuplicationEnable()) { + if (dispatchRequest.getCommitLogOffset() + < this.defaultMessageStore.getConfirmOffset()) { + this.defaultMessageStore.doDispatch(dispatchRequest); + } + } else { this.defaultMessageStore.doDispatch(dispatchRequest); } - } else { - this.defaultMessageStore.doDispatch(dispatchRequest); } - } else if (size == 0) { + } else if (entrySize == 0) { index++; if (index >= mmapFiles.size()) { log.info("dledger recover physics file over, last mapped file " + mmapFile.getFileName()); @@ -426,16 +496,53 @@ private void setRecoverPosition() { log.info("Will set the initial commitlog offset={} for dledger", dividedCommitlogOffset); } - private boolean isMmapFileMatchedRecover(final MmapFile mmapFile, boolean recoverNormally) throws RocksDBException { - ByteBuffer byteBuffer = mmapFile.sliceByteBuffer(); + private ByteBuffer firstNormalEntryBody(ByteBuffer byteBuffer) { + int limit = byteBuffer.limit(); + int position = byteBuffer.position(); + while (limit - position >= Integer.BYTES * 2) { + int magic = byteBuffer.getInt(position); + int entrySize = byteBuffer.getInt(position + Integer.BYTES); + if (magic == MmapFileList.BLANK_MAGIC_CODE || entrySize < DLedgerEntry.BODY_OFFSET + || entrySize > limit - position) { + return null; + } + if (magic == DLedgerEntryType.NOOP.getMagic()) { + if (entrySize != DLedgerEntry.BODY_OFFSET) { + return null; + } + position += entrySize; + continue; + } + if (magic != DLedgerEntryType.NORMAL.getMagic()) { + return null; + } + ByteBuffer body = byteBuffer.duplicate(); + body.position(position + DLedgerEntry.BODY_OFFSET); + body.limit(position + entrySize); + return body.slice(); + } + return null; + } + + private boolean isMmapFileMatchedRecover(final MmapFile mmapFile, boolean recoverNormally) + throws RocksDBException { + ByteBuffer byteBuffer = firstNormalEntryBody(mmapFile.sliceByteBuffer()); + if (byteBuffer == null + || byteBuffer.limit() < MessageDecoder.MESSAGE_MAGIC_CODE_POSITION + Integer.BYTES) { + return false; + } - int magicCode = byteBuffer.getInt(DLedgerEntry.BODY_OFFSET + MessageDecoder.MESSAGE_MAGIC_CODE_POSITION); - if (magicCode != MESSAGE_MAGIC_CODE) { + int magicCode = byteBuffer.getInt(MessageDecoder.MESSAGE_MAGIC_CODE_POSITION); + if (magicCode != MESSAGE_MAGIC_CODE + && magicCode != MessageDecoder.MESSAGE_MAGIC_CODE_V2) { return false; } + if (byteBuffer.limit() < MessageDecoder.SYSFLAG_POSITION + Integer.BYTES) { + return false; + } int storeTimestampPosition; - int sysFlag = byteBuffer.getInt(DLedgerEntry.BODY_OFFSET + MessageDecoder.SYSFLAG_POSITION); + int sysFlag = byteBuffer.getInt(MessageDecoder.SYSFLAG_POSITION); if ((sysFlag & MessageSysFlag.BORNHOST_V6_FLAG) == 0) { storeTimestampPosition = MessageDecoder.MESSAGE_STORE_TIMESTAMP_POSITION; } else { @@ -443,11 +550,15 @@ private boolean isMmapFileMatchedRecover(final MmapFile mmapFile, boolean recove storeTimestampPosition = MessageDecoder.MESSAGE_STORE_TIMESTAMP_POSITION + 12; } - long storeTimestamp = byteBuffer.getLong(DLedgerEntry.BODY_OFFSET + storeTimestampPosition); + if (byteBuffer.limit() < storeTimestampPosition + Long.BYTES + || byteBuffer.limit() < MessageDecoder.MESSAGE_PHYSIC_OFFSET_POSITION + Long.BYTES) { + return false; + } + long storeTimestamp = byteBuffer.getLong(storeTimestampPosition); if (storeTimestamp == 0) { return false; } - long phyOffset = byteBuffer.getLong(DLedgerEntry.BODY_OFFSET + MessageDecoder.MESSAGE_PHYSIC_OFFSET_POSITION); + long phyOffset = byteBuffer.getLong(MessageDecoder.MESSAGE_PHYSIC_OFFSET_POSITION); if (this.defaultMessageStore.getMessageStoreConfig().isMessageIndexEnable() && this.defaultMessageStore.getMessageStoreConfig().isMessageIndexSafe()) { @@ -477,30 +588,54 @@ public DispatchRequest checkMessageAndReturnSize(ByteBuffer byteBuffer, final bo if (isInrecoveringOldCommitlog) { return super.checkMessageAndReturnSize(byteBuffer, checkCRC, checkDupInfo, readBody); } + int position = byteBuffer.position(); try { - int bodyOffset = DLedgerEntry.BODY_OFFSET; - int pos = byteBuffer.position(); - int magic = byteBuffer.getInt(); + if (byteBuffer.remaining() < Integer.BYTES * 2) { + return new DispatchRequest(-1, false); + } + int magic = byteBuffer.getInt(position); //In dledger, this field is size, it must be gt 0, so it could prevent collision - int magicOld = byteBuffer.getInt(); - if (magicOld == CommitLog.BLANK_MAGIC_CODE - || magicOld == MessageDecoder.MESSAGE_MAGIC_CODE - || magicOld == MessageDecoder.MESSAGE_MAGIC_CODE_V2) { - byteBuffer.position(pos); + int entrySize = byteBuffer.getInt(position + Integer.BYTES); + if (entrySize == CommitLog.BLANK_MAGIC_CODE + || entrySize == MessageDecoder.MESSAGE_MAGIC_CODE + || entrySize == MessageDecoder.MESSAGE_MAGIC_CODE_V2) { return super.checkMessageAndReturnSize(byteBuffer, checkCRC, checkDupInfo, readBody); } if (magic == MmapFileList.BLANK_MAGIC_CODE) { return new DispatchRequest(0, true); } - byteBuffer.position(pos + bodyOffset); - DispatchRequest dispatchRequest = super.checkMessageAndReturnSize(byteBuffer, checkCRC, checkDupInfo, readBody); + if (magic == DLedgerEntryType.NOOP.getMagic()) { + if (entrySize != DLedgerEntry.BODY_OFFSET || entrySize > byteBuffer.remaining()) { + return new DispatchRequest(-1, false); + } + byteBuffer.position(position + entrySize); + DispatchRequest dispatchRequest = new DispatchRequest(0, true); + dispatchRequest.setBufferSize(entrySize); + return dispatchRequest; + } + if (magic != DLedgerEntryType.NORMAL.getMagic() || entrySize < DLedgerEntry.BODY_OFFSET + || entrySize > byteBuffer.remaining()) { + return new DispatchRequest(-1, false); + } + int entryEnd = position + entrySize; + ByteBuffer messageBuffer = byteBuffer.duplicate(); + messageBuffer.position(position + DLedgerEntry.BODY_OFFSET); + messageBuffer.limit(entryEnd); + messageBuffer = messageBuffer.slice(); + DispatchRequest dispatchRequest = super.checkMessageAndReturnSize( + messageBuffer, checkCRC, checkDupInfo, readBody); if (dispatchRequest.isSuccess()) { - dispatchRequest.setBufferSize(dispatchRequest.getMsgSize() + bodyOffset); + if (dispatchRequest.getMsgSize() + DLedgerEntry.BODY_OFFSET != entrySize) { + return new DispatchRequest(-1, false); + } + byteBuffer.position(entryEnd); + dispatchRequest.setBufferSize(entrySize); } else if (dispatchRequest.getMsgSize() > 0) { - dispatchRequest.setBufferSize(dispatchRequest.getMsgSize() + bodyOffset); + dispatchRequest.setBufferSize(dispatchRequest.getMsgSize() + DLedgerEntry.BODY_OFFSET); } return dispatchRequest; } catch (Throwable ignored) { + byteBuffer.position(position); } return new DispatchRequest(-1, false /* success */); @@ -673,7 +808,7 @@ public CompletableFuture asyncPutMessages(MessageExtBatch mess // Back to Results AppendMessageResult appendResult; - BatchAppendFuture dledgerFuture; + AppendFuture dledgerFuture; EncodeResult encodeResult; encodeResult = this.messageSerializer.serialize(messageExtBatch); @@ -705,7 +840,23 @@ public CompletableFuture asyncPutMessages(MessageExtBatch mess log.warn("HandleAppend return false due to error code {}", appendFuture.get().getCode()); return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.OS_PAGE_CACHE_BUSY, new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR))); } - dledgerFuture = (BatchAppendFuture) appendFuture; + dledgerFuture = appendFuture; + + long[] positions; + if (batchNum == 1) { + positions = new long[] {appendFuture.getPos()}; + } else { + if (!(appendFuture instanceof BatchAppendFuture)) { + throw new IllegalStateException("Unexpected append future type for " + batchNum + + "-message batch: " + appendFuture.getClass().getName()); + } + positions = ((BatchAppendFuture) appendFuture).getPositions(); + if (positions == null || positions.length != batchNum + || appendFuture.getPos() != positions[batchNum - 1]) { + throw new IllegalStateException("Inconsistent DLedger batch positions: expected " + batchNum + + " entries ending at " + appendFuture.getPos()); + } + } long wroteOffset = 0; @@ -714,7 +865,7 @@ public CompletableFuture asyncPutMessages(MessageExtBatch mess boolean isFirstOffset = true; long firstWroteOffset = 0; - for (long pos : dledgerFuture.getPositions()) { + for (long pos : positions) { wroteOffset = pos + DLedgerEntry.BODY_OFFSET; if (isFirstOffset) { firstWroteOffset = wroteOffset; @@ -787,9 +938,17 @@ public SelectMappedBufferResult getMessage(final long offset, final int size) { 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)); } diff --git a/store/src/test/java/org/apache/rocketmq/store/dledger/DLedgerLatestCommitLogTest.java b/store/src/test/java/org/apache/rocketmq/store/dledger/DLedgerLatestCommitLogTest.java new file mode 100644 index 00000000000..6f6cbe33031 --- /dev/null +++ b/store/src/test/java/org/apache/rocketmq/store/dledger/DLedgerLatestCommitLogTest.java @@ -0,0 +1,599 @@ +/* + * 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.rocketmq.store.dledger; + +import io.openmessaging.storage.dledger.DLedgerServer; +import io.openmessaging.storage.dledger.common.ReadClosure; +import io.openmessaging.storage.dledger.common.ReadMode; +import io.openmessaging.storage.dledger.common.Status; +import io.openmessaging.storage.dledger.entry.DLedgerEntry; +import io.openmessaging.storage.dledger.entry.DLedgerEntryCoder; +import io.openmessaging.storage.dledger.entry.DLedgerEntryType; +import io.openmessaging.storage.dledger.store.file.DLedgerMmapFileStore; +import io.openmessaging.storage.dledger.store.file.MmapFileList; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.rocketmq.common.message.MessageDecoder; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.message.MessageExtBatch; +import org.apache.rocketmq.common.message.MessageExtBrokerInner; +import org.apache.rocketmq.store.DefaultMessageStore; +import org.apache.rocketmq.store.DispatchRequest; +import org.apache.rocketmq.store.GetMessageResult; +import org.apache.rocketmq.store.PutMessageResult; +import org.apache.rocketmq.store.PutMessageStatus; +import org.junit.Assert; +import org.junit.Test; + +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.awaitility.Awaitility.await; + +public class DLedgerLatestCommitLogTest extends MessageStoreTestBase { + + private static final int QUEUE_ID = 0; + + @Test + public void testUncommittedTailIsNotReadable() throws Exception { + String peers = String.format("n0-localhost:%d;n1-localhost:%d", nextPort(), nextPort()); + DefaultMessageStore leaderStore = null; + try { + leaderStore = createDledgerMessageStore( + createBaseDir(), UUID.randomUUID().toString(), "n0", peers, "n0", false, 0); + String topic = UUID.randomUUID().toString(); + MessageExtBrokerInner message = buildMessage(); + message.setTopic(topic); + message.setQueueId(QUEUE_ID); + + PutMessageResult result = leaderStore.asyncPutMessage(message).get(5, SECONDS); + Assert.assertEquals(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, result.getPutMessageStatus()); + Assert.assertNotNull(result.getAppendMessageResult()); + Assert.assertTrue(result.getAppendMessageResult().getWroteOffset() > 0); + + DLedgerCommitLog commitLog = commitLog(leaderStore); + Assert.assertEquals(-1, commitLog.getdLedgerServer().getMemberState().getCommittedIndex()); + Assert.assertEquals(-1, commitLog.getCommittedPos()); + Assert.assertEquals(0, commitLog.getMaxOffset()); + Assert.assertEquals(0, leaderStore.getMaxOffsetInQueue(topic, QUEUE_ID)); + Assert.assertNull(commitLog.getData(0)); + Assert.assertFalse(commitLog.getData(0, 1, ByteBuffer.allocate(1))); + Assert.assertNull(commitLog.getMessage(result.getAppendMessageResult().getWroteOffset(), 1)); + } finally { + shutdownAndDestroy(leaderStore); + } + } + + @Test + public void testSingleAndBatchAppendPositions() throws Exception { + String peers = String.format("n0-localhost:%d", nextPort()); + DefaultMessageStore messageStore = null; + try { + messageStore = createDledgerMessageStore( + createBaseDir(), UUID.randomUUID().toString(), "n0", peers, null, false, 0); + awaitLeader(Arrays.asList(messageStore)); + String topic = UUID.randomUUID().toString(); + + PutMessageResult singleResult = putSingle(messageStore, topic, 0); + PutMessageResult singleMessageBatchResult = putBatch(messageStore, topic, 1, 1); + PutMessageResult batchResult = putBatch(messageStore, topic, 3, 2); + + Assert.assertTrue(singleResult.getAppendMessageResult().getWroteOffset() > 0); + Assert.assertTrue(singleMessageBatchResult.getAppendMessageResult().getWroteOffset() + > singleResult.getAppendMessageResult().getWroteOffset()); + Assert.assertTrue(batchResult.getAppendMessageResult().getWroteOffset() + > singleMessageBatchResult.getAppendMessageResult().getWroteOffset()); + Assert.assertEquals(1, singleMessageBatchResult.getAppendMessageResult().getMsgNum()); + Assert.assertEquals(3, batchResult.getAppendMessageResult().getMsgNum()); + Assert.assertNotNull(singleResult.getAppendMessageResult().getMsgId()); + Assert.assertNotNull(singleMessageBatchResult.getAppendMessageResult().getMsgId()); + Assert.assertNotNull(batchResult.getAppendMessageResult().getMsgId()); + Assert.assertEquals(1, + singleMessageBatchResult.getAppendMessageResult().getMsgId().split(",").length); + Assert.assertEquals(3, batchResult.getAppendMessageResult().getMsgId().split(",").length); + awaitStoreReady(messageStore, topic, 5); + Assert.assertEquals(0, messageStore.getMinOffsetInQueue(topic, QUEUE_ID)); + Assert.assertTrue(commitLog(messageStore).getCommittedPos() + > batchResult.getAppendMessageResult().getWroteOffset()); + doGetMessages(messageStore, topic, QUEUE_ID, 5, 0); + } finally { + shutdownAndDestroy(messageStore); + } + } + + @Test + public void testThreeNodeElectionAndFailover() throws Exception { + String peers = String.format("n0-localhost:%d;n1-localhost:%d;n2-localhost:%d", + nextPort(), nextPort(), nextPort()); + String group = UUID.randomUUID().toString(); + List allStores = new ArrayList<>(); + try { + allStores.add(createDledgerMessageStore(createBaseDir(), group, "n0", peers, null, false, 0)); + allStores.add(createDledgerMessageStore(createBaseDir(), group, "n1", peers, null, false, 0)); + allStores.add(createDledgerMessageStore(createBaseDir(), group, "n2", peers, null, false, 0)); + List activeStores = new ArrayList<>(allStores); + DefaultMessageStore firstLeader = awaitLeader(activeStores); + String topic = UUID.randomUUID().toString(); + + putSingle(firstLeader, topic, 0); + putSingle(firstLeader, topic, 1); + putSingle(firstLeader, topic, 2); + for (DefaultMessageStore store : activeStores) { + awaitStoreReady(store, topic, 3); + } + long committedBeforeFailover = commitLog(firstLeader).getCommittedPos(); + + firstLeader.shutdown(); + activeStores.remove(firstLeader); + DefaultMessageStore secondLeader = awaitLeader(activeStores); + Assert.assertNotSame(firstLeader, secondLeader); + awaitStoreReady(secondLeader, topic, 3); + Assert.assertTrue(commitLog(secondLeader).getCommittedPos() >= committedBeforeFailover); + doGetMessages(secondLeader, topic, QUEUE_ID, 3, 0); + + // Broker-side DLedgerRoleChangeHandler does this before accepting writes on a new leader. + secondLeader.recoverTopicQueueTable(); + putBatch(secondLeader, topic, 3, 3); + for (DefaultMessageStore store : activeStores) { + awaitStoreReady(store, topic, 6); + } + doGetMessages(secondLeader, topic, QUEUE_ID, 6, 0); + } finally { + for (DefaultMessageStore store : allStores) { + shutdownAndDestroy(store); + } + } + } + + @Test + public void testRestartRecoversCommittedBoundaryBeforeNewWrite() throws Exception { + String base = createBaseDir(); + String peers = String.format("n0-localhost:%d", nextPort()); + String group = UUID.randomUUID().toString(); + String topic = UUID.randomUUID().toString(); + DefaultMessageStore currentStore = null; + try { + currentStore = createDledgerMessageStore(base, group, "n0", peers, null, false, 0); + awaitLeader(Arrays.asList(currentStore)); + doPutMessages(currentStore, topic, QUEUE_ID, 10, 0); + awaitStoreReady(currentStore, topic, 10); + doGetMessages(currentStore, topic, QUEUE_ID, 10, 0); + long physicalBeforeRestart = currentStore.getMaxPhyOffset(); + long maxCqOffsetBeforeRestart = currentStore.getMaxOffsetInQueue(topic, QUEUE_ID); + List bodiesBeforeRestart = readMessageBodies(currentStore, topic, QUEUE_ID, 10); + long committedIndexBeforeRestart = committedIndex(currentStore); + Assert.assertTrue(physicalBeforeRestart > 0); + Assert.assertEquals(10, maxCqOffsetBeforeRestart); + Assert.assertTrue(committedIndexBeforeRestart >= 9); + + currentStore.shutdown(); + currentStore = createDledgerMessageStore(base, group, "n0", peers, null, false, 0); + awaitLeader(Arrays.asList(currentStore)); + awaitCommittedPast(currentStore, committedIndexBeforeRestart); + assertNoopEntry(currentStore, committedIndexBeforeRestart + 1); + awaitStoreReady(currentStore, topic, maxCqOffsetBeforeRestart); + Assert.assertTrue(currentStore.getMaxPhyOffset() >= physicalBeforeRestart); + Assert.assertEquals(0, currentStore.getMinOffsetInQueue(topic, QUEUE_ID)); + Assert.assertEquals(maxCqOffsetBeforeRestart, + currentStore.getMaxOffsetInQueue(topic, QUEUE_ID)); + assertMessageBodies(currentStore, topic, QUEUE_ID, bodiesBeforeRestart); + Assert.assertEquals(commitLog(currentStore).getCommittedPos(), currentStore.getCommitLog().getMaxOffset()); + doGetMessages(currentStore, topic, QUEUE_ID, 10, 0); + + putSingle(currentStore, topic, 10); + awaitStoreReady(currentStore, topic, 11); + doGetMessages(currentStore, topic, QUEUE_ID, 11, 0); + long committedPosBeforeSecondRestart = commitLog(currentStore).getCommittedPos(); + long committedIndexBeforeSecondRestart = committedIndex(currentStore); + + currentStore.shutdown(); + currentStore = createDledgerMessageStore(base, group, "n0", peers, null, true, 0); + awaitLeader(Arrays.asList(currentStore)); + awaitCommittedPast(currentStore, committedIndexBeforeSecondRestart); + assertNoopEntry(currentStore, committedIndexBeforeSecondRestart + 1); + awaitStoreReady(currentStore, topic, 11); + Assert.assertEquals(0, currentStore.getMinOffsetInQueue(topic, QUEUE_ID)); + Assert.assertTrue(commitLog(currentStore).getCommittedPos() >= committedPosBeforeSecondRestart); + Assert.assertEquals(commitLog(currentStore).getCommittedPos(), currentStore.getCommitLog().getMaxOffset()); + doGetMessages(currentStore, topic, QUEUE_ID, 11, 0); + } finally { + shutdownAndDestroy(currentStore); + } + } + + @Test + public void testNoopDispatchContractAndBounds() throws Exception { + String peers = String.format("n0-localhost:%d", nextPort()); + DefaultMessageStore messageStore = null; + try { + messageStore = createDledgerMessageStore( + createBaseDir(), UUID.randomUUID().toString(), "n0", peers, null, false, 0); + DLedgerCommitLog commitLog = commitLog(messageStore); + + ByteBuffer noopBuffer = ByteBuffer.allocate(DLedgerEntry.BODY_OFFSET); + DLedgerEntryCoder.encode(new DLedgerEntry(DLedgerEntryType.NOOP), noopBuffer); + DispatchRequest noop = commitLog.checkMessageAndReturnSize(noopBuffer, true, false, false); + Assert.assertTrue(noop.isSuccess()); + Assert.assertEquals(0, noop.getMsgSize()); + Assert.assertEquals(DLedgerEntry.BODY_OFFSET, noop.getBufferSize()); + Assert.assertEquals(DLedgerEntry.BODY_OFFSET, noopBuffer.position()); + + ByteBuffer undersized = noopHeader(DLedgerEntry.BODY_OFFSET - 1); + DispatchRequest invalidSize = commitLog.checkMessageAndReturnSize(undersized, true, false, false); + Assert.assertFalse(invalidSize.isSuccess()); + Assert.assertEquals(-1, invalidSize.getMsgSize()); + Assert.assertEquals(0, undersized.position()); + + ByteBuffer truncated = noopHeader(DLedgerEntry.BODY_OFFSET + 1); + DispatchRequest invalidBounds = commitLog.checkMessageAndReturnSize(truncated, true, false, false); + Assert.assertFalse(invalidBounds.isSuccess()); + Assert.assertEquals(-1, invalidBounds.getMsgSize()); + Assert.assertEquals(0, truncated.position()); + + awaitLeader(Arrays.asList(messageStore)); + String topic = UUID.randomUUID().toString(); + putSingle(messageStore, topic, 0); + awaitStoreReady(messageStore, topic, 1); + DLedgerServer server = commitLog.getdLedgerServer(); + DLedgerEntry normalEntry = server.getDLedgerStore().get( + server.getDLedgerStore().getLedgerEndIndex()); + Assert.assertEquals(DLedgerEntryType.NORMAL.getMagic(), normalEntry.getMagic()); + byte[] innerMessage = normalEntry.getBody(); + + ByteBuffer legalFollowingEntry = normalEntryBuffer(innerMessage, 0, null); + byte[] legalFollowingBytes = new byte[legalFollowingEntry.remaining()]; + legalFollowingEntry.get(legalFollowingBytes); + byte[] oversizedInnerMessage = Arrays.copyOf(innerMessage, innerMessage.length); + ByteBuffer.wrap(oversizedInnerMessage).putInt(innerMessage.length + Integer.BYTES); + ByteBuffer crossingEntry = normalEntryBuffer( + oversizedInnerMessage, 0, legalFollowingBytes); + DispatchRequest crossingRequest = commitLog.checkMessageAndReturnSize( + crossingEntry, true, false, false); + + ByteBuffer mismatchedEntry = normalEntryBuffer(innerMessage, Integer.BYTES, null); + DispatchRequest mismatchedRequest = commitLog.checkMessageAndReturnSize( + mismatchedEntry, true, false, false); + + Assert.assertFalse(crossingRequest.isSuccess()); + Assert.assertFalse(mismatchedRequest.isSuccess()); + Assert.assertArrayEquals(new int[] {0, 0}, + new int[] {crossingEntry.position(), mismatchedEntry.position()}); + } finally { + shutdownAndDestroy(messageStore); + } + } + + @Test + public void testRaftLogReadNoopDoesNotBuildConsumeQueue() throws Exception { + String peers = String.format("n0-localhost:%d", nextPort()); + DefaultMessageStore messageStore = null; + try { + messageStore = createDledgerMessageStore( + createBaseDir(), UUID.randomUUID().toString(), "n0", peers, null, false, 0); + awaitLeader(Arrays.asList(messageStore)); + Assert.assertTrue(messageStore.getConsumeQueueTable().isEmpty()); + DLedgerServer server = commitLog(messageStore).getdLedgerServer(); + long previousCommittedIndex = committedIndex(messageStore); + long previousLedgerEndIndex = server.getDLedgerStore().getLedgerEndIndex(); + + Status status = appendRaftLogNoop(messageStore); + Assert.assertTrue(status.isOk()); + awaitCommittedPast(messageStore, previousCommittedIndex); + Assert.assertEquals(previousLedgerEndIndex + 1, + server.getDLedgerStore().getLedgerEndIndex()); + assertNoopEntry(messageStore, previousLedgerEndIndex + 1); + awaitNoopConsumed(messageStore); + + Assert.assertTrue(messageStore.getConsumeQueueTable().isEmpty()); + Assert.assertTrue(commitLog(messageStore).getCommittedPos() > 0); + Assert.assertEquals(commitLog(messageStore).getCommittedPos(), messageStore.getCommitLog().getMaxOffset()); + } finally { + shutdownAndDestroy(messageStore); + } + } + + @Test + public void testAbnormalRecoveryAcrossLeadingNoop() throws Exception { + String base = createBaseDir(); + String peers = String.format("n0-localhost:%d", nextPort()); + String group = UUID.randomUUID().toString(); + String topic = String.format("%s%s%s%s", UUID.randomUUID(), UUID.randomUUID(), + UUID.randomUUID(), UUID.randomUUID()); + DefaultMessageStore currentStore = null; + try { + currentStore = createDledgerMessageStore(base, group, "n0", peers, null, false, 0); + awaitLeader(Arrays.asList(currentStore)); + Assert.assertTrue(appendRaftLogNoop(currentStore).isOk()); + awaitCommittedPast(currentStore, -1); + assertNoopEntry(currentStore, 0); + awaitNoopConsumed(currentStore); + Assert.assertTrue(currentStore.getConsumeQueueTable().isEmpty()); + + putSingle(currentStore, topic, 0); + awaitStoreReady(currentStore, topic, 1); + assertMessageMagic(currentStore, topic, QUEUE_ID, + MessageDecoder.MESSAGE_MAGIC_CODE_V2); + doGetMessages(currentStore, topic, QUEUE_ID, 1, 0); + long committedIndexBeforeRestart = committedIndex(currentStore); + long committedPosBeforeRestart = commitLog(currentStore).getCommittedPos(); + + currentStore.shutdown(); + currentStore = createDledgerMessageStore(base, group, "n0", peers, null, true, 0); + awaitLeader(Arrays.asList(currentStore)); + awaitCommittedPast(currentStore, committedIndexBeforeRestart); + assertNoopEntry(currentStore, 0); + assertNoopEntry(currentStore, committedIndexBeforeRestart + 1); + awaitStoreReady(currentStore, topic, 1); + Assert.assertEquals(1, currentStore.getConsumeQueueTable().size()); + Assert.assertTrue(currentStore.getConsumeQueueTable().containsKey(topic)); + Assert.assertTrue(commitLog(currentStore).getCommittedPos() >= committedPosBeforeRestart); + doGetMessages(currentStore, topic, QUEUE_ID, 1, 0); + + putSingle(currentStore, topic, 1); + awaitStoreReady(currentStore, topic, 2); + doGetMessages(currentStore, topic, QUEUE_ID, 2, 0); + } finally { + shutdownAndDestroy(currentStore); + } + } + + @Test + public void testFixedSizeReadsRespectCommittedBoundary() throws Exception { + String peers = String.format("n0-localhost:%d;n1-localhost:%d", nextPort(), nextPort()); + String group = UUID.randomUUID().toString(); + DefaultMessageStore leaderStore = null; + DefaultMessageStore followerStore = null; + try { + leaderStore = createDledgerMessageStore( + createBaseDir(), group, "n0", peers, "n0", false, 0); + followerStore = createDledgerMessageStore( + createBaseDir(), group, "n1", peers, "n0", false, 0); + String topic = UUID.randomUUID().toString(); + DLedgerCommitLog leaderCommitLog = commitLog(leaderStore); + DLedgerMmapFileStore dLedgerStore = + (DLedgerMmapFileStore) leaderCommitLog.getdLedgerServer().getdLedgerStore(); + MmapFileList dataFileList = dLedgerStore.getDataFileList(); + + int messageCount = 0; + while (dataFileList.getMappedFiles().size() < 2 && messageCount < 16) { + MessageExtBrokerInner message = buildMessage(); + message.setTopic(topic); + message.setQueueId(QUEUE_ID); + message.setBody(new byte[16 * 1024]); + PutMessageResult result = leaderStore.asyncPutMessage(message).get(5, SECONDS); + Assert.assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + messageCount++; + } + Assert.assertEquals(2, dataFileList.getMappedFiles().size()); + awaitStoreReady(leaderStore, topic, messageCount); + awaitStoreReady(followerStore, topic, messageCount); + + long selectedBase = dataFileList.getMappedFiles().get(1).getFileFromOffset(); + long committedPos = leaderCommitLog.getCommittedPos(); + Assert.assertTrue(committedPos > selectedBase); + + followerStore.shutdown(); + MessageExtBrokerInner uncommittedMessage = buildMessage(); + uncommittedMessage.setTopic(topic); + uncommittedMessage.setQueueId(QUEUE_ID); + PutMessageResult uncommittedResult = leaderStore.asyncPutMessage(uncommittedMessage).get(5, SECONDS); + Assert.assertEquals(PutMessageStatus.IN_SYNC_REPLICAS_NOT_ENOUGH, + uncommittedResult.getPutMessageStatus()); + Assert.assertEquals(committedPos, leaderCommitLog.getCommittedPos()); + Assert.assertTrue(dataFileList.getMaxWrotePosition() > committedPos); + Assert.assertEquals(2, dataFileList.getMappedFiles().size()); + + ByteBuffer oversizedDestination = ByteBuffer.allocate(64); + Assert.assertTrue(leaderCommitLog.getData(committedPos - 1, 1, oversizedDestination)); + Assert.assertEquals(1, oversizedDestination.position()); + Assert.assertEquals(64, oversizedDestination.limit()); + + Assert.assertEquals(1, dataFileList.deleteExpiredFileByTime(0, 0, 0, true)); + Assert.assertEquals(1, dataFileList.getMappedFiles().size()); + long firstSurvivingBase = dataFileList.getFirstMappedFile().getFileFromOffset(); + Assert.assertEquals(selectedBase, firstSurvivingBase); + Assert.assertTrue(firstSurvivingBase > 0); + Assert.assertTrue(firstSurvivingBase < committedPos); + Assert.assertSame(dataFileList.getFirstMappedFile(), + dataFileList.findMappedFileByOffset(0, true)); + + int crossSize = (int) (committedPos - firstSurvivingBase + 1); + Assert.assertTrue(firstSurvivingBase + crossSize <= dataFileList.getMaxWrotePosition()); + ByteBuffer crossBoundaryDestination = ByteBuffer.allocate(crossSize); + Assert.assertFalse(leaderCommitLog.getData(0, crossSize, crossBoundaryDestination)); + Assert.assertEquals(0, crossBoundaryDestination.position()); + Assert.assertNull(leaderCommitLog.getMessage(0, crossSize)); + } finally { + shutdownAndDestroy(followerStore); + shutdownAndDestroy(leaderStore); + } + } + + private PutMessageResult putSingle(DefaultMessageStore messageStore, String topic, long expectedLogicOffset) + throws Exception { + MessageExtBrokerInner message = buildMessage(); + message.setTopic(topic); + message.setQueueId(QUEUE_ID); + PutMessageResult result = messageStore.asyncPutMessage(message).get(5, SECONDS); + Assert.assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + Assert.assertNotNull(result.getAppendMessageResult()); + Assert.assertEquals(expectedLogicOffset, result.getAppendMessageResult().getLogicsOffset()); + return result; + } + + private PutMessageResult putBatch(DefaultMessageStore messageStore, String topic, int batchSize, + long expectedLogicOffset) throws Exception { + MessageExtBatch batch = buildBatchMessage(batchSize); + batch.setTopic(topic); + batch.setQueueId(QUEUE_ID); + PutMessageResult result = messageStore.asyncPutMessages(batch).get(5, SECONDS); + Assert.assertEquals(PutMessageStatus.PUT_OK, result.getPutMessageStatus()); + Assert.assertNotNull(result.getAppendMessageResult()); + Assert.assertEquals(expectedLogicOffset, result.getAppendMessageResult().getLogicsOffset()); + return result; + } + + private DefaultMessageStore awaitLeader(List stores) { + AtomicReference leaderRef = new AtomicReference<>(); + await().atMost(15, SECONDS).pollInterval(100, MILLISECONDS).until(() -> { + DefaultMessageStore leader = null; + for (DefaultMessageStore store : stores) { + if (commitLog(store).getdLedgerServer().getMemberState().isLeader()) { + if (leader != null) { + return false; + } + leader = store; + } + } + leaderRef.set(leader); + return leader != null; + }); + return leaderRef.get(); + } + + private void awaitStoreReady(DefaultMessageStore messageStore, String topic, long expectedMaxOffset) { + await().atMost(15, SECONDS).pollInterval(100, MILLISECONDS).untilAsserted(() -> { + Assert.assertEquals(expectedMaxOffset, messageStore.getMaxOffsetInQueue(topic, QUEUE_ID)); + Assert.assertEquals(0, messageStore.dispatchBehindBytes()); + }); + } + + private void awaitCommittedPast(DefaultMessageStore messageStore, long previousCommittedIndex) { + DLedgerServer server = commitLog(messageStore).getdLedgerServer(); + await().atMost(15, SECONDS).pollInterval(100, MILLISECONDS).until(() -> { + long committedIndex = server.getMemberState().getCommittedIndex(); + return committedIndex > previousCommittedIndex + && committedIndex == server.getDLedgerStore().getLedgerEndIndex(); + }); + } + + private void awaitNoopConsumed(DefaultMessageStore messageStore) { + await().atMost(15, SECONDS).pollInterval(100, MILLISECONDS).untilAsserted(() -> { + Assert.assertEquals(commitLog(messageStore).getCommittedPos(), messageStore.getCommitLog().getMaxOffset()); + Assert.assertEquals(0, messageStore.dispatchBehindBytes()); + }); + } + + private List readMessageBodies(DefaultMessageStore messageStore, String topic, int queueId, + int messageCount) { + List bodies = new ArrayList<>(messageCount); + for (int i = 0; i < messageCount; i++) { + GetMessageResult result = messageStore.getMessage("group", topic, queueId, i, 1, null); + Assert.assertNotNull(result); + try { + Assert.assertFalse(result.getMessageBufferList().isEmpty()); + MessageExt message = MessageDecoder.decode(result.getMessageBufferList().get(0)); + Assert.assertNotNull(message); + Assert.assertEquals(i, message.getQueueOffset()); + bodies.add(Arrays.copyOf(message.getBody(), message.getBody().length)); + } finally { + result.release(); + } + } + return bodies; + } + + private void assertMessageBodies(DefaultMessageStore messageStore, String topic, int queueId, + List expectedBodies) { + List actualBodies = readMessageBodies(messageStore, topic, queueId, expectedBodies.size()); + for (int i = 0; i < expectedBodies.size(); i++) { + Assert.assertArrayEquals(expectedBodies.get(i), actualBodies.get(i)); + } + } + + private void assertMessageMagic(DefaultMessageStore messageStore, String topic, int queueId, + int expectedMagic) { + GetMessageResult result = messageStore.getMessage("group", topic, queueId, 0, 1, null); + Assert.assertNotNull(result); + try { + Assert.assertFalse(result.getMessageBufferList().isEmpty()); + ByteBuffer messageBuffer = result.getMessageBufferList().get(0).duplicate(); + Assert.assertEquals(expectedMagic, + messageBuffer.getInt(messageBuffer.position() + MessageDecoder.MESSAGE_MAGIC_CODE_POSITION)); + } finally { + result.release(); + } + } + + private Status appendRaftLogNoop(DefaultMessageStore messageStore) throws Exception { + CompletableFuture result = new CompletableFuture<>(); + commitLog(messageStore).getdLedgerServer().handleRead(ReadMode.RAFT_LOG_READ, new ReadClosure() { + @Override + public void done(Status status) { + result.complete(status); + } + }); + return result.get(5, SECONDS); + } + + private ByteBuffer noopHeader(int entrySize) { + ByteBuffer buffer = ByteBuffer.allocate(DLedgerEntry.BODY_OFFSET); + buffer.putInt(DLedgerEntryType.NOOP.getMagic()); + buffer.putInt(entrySize); + buffer.position(0); + buffer.limit(DLedgerEntry.BODY_OFFSET); + return buffer; + } + + private ByteBuffer normalEntryBuffer(byte[] innerMessage, int bodyPadding, byte[] trailingBytes) { + int entrySize = DLedgerEntry.BODY_OFFSET + innerMessage.length + bodyPadding; + int trailingSize = trailingBytes == null ? 0 : trailingBytes.length; + ByteBuffer buffer = ByteBuffer.allocate(entrySize + trailingSize); + buffer.putInt(DLedgerEntryType.NORMAL.getMagic()); + buffer.putInt(entrySize); + buffer.position(DLedgerEntry.BODY_OFFSET); + buffer.put(innerMessage); + buffer.position(entrySize); + if (trailingBytes != null) { + buffer.put(trailingBytes); + } + buffer.flip(); + return buffer; + } + + private void assertNoopEntry(DefaultMessageStore messageStore, long index) { + DLedgerServer server = commitLog(messageStore).getdLedgerServer(); + DLedgerEntry entry = server.getDLedgerStore().get(index); + Assert.assertNotNull(entry); + Assert.assertEquals(DLedgerEntryType.NOOP.getMagic(), entry.getMagic()); + } + + private long committedIndex(DefaultMessageStore messageStore) { + return commitLog(messageStore).getdLedgerServer().getMemberState().getCommittedIndex(); + } + + private DLedgerCommitLog commitLog(DefaultMessageStore messageStore) { + return (DLedgerCommitLog) messageStore.getCommitLog(); + } + + private void shutdownAndDestroy(DefaultMessageStore messageStore) { + if (messageStore == null) { + return; + } + try { + if (!messageStore.isShutdown()) { + messageStore.shutdown(); + } + } finally { + messageStore.destroy(); + } + } +} diff --git a/test/src/test/java/org/apache/rocketmq/test/dledger/DLedgerThreeNodeIT.java b/test/src/test/java/org/apache/rocketmq/test/dledger/DLedgerThreeNodeIT.java new file mode 100644 index 00000000000..2e50e564d25 --- /dev/null +++ b/test/src/test/java/org/apache/rocketmq/test/dledger/DLedgerThreeNodeIT.java @@ -0,0 +1,601 @@ +/* + * 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.rocketmq.test.dledger; + +import java.io.File; +import java.io.IOException; +import java.net.InetAddress; +import java.net.ServerSocket; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.IdentityHashMap; +import java.util.List; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer; +import org.apache.rocketmq.client.consumer.PullResult; +import org.apache.rocketmq.client.consumer.PullStatus; +import org.apache.rocketmq.client.producer.DefaultMQProducer; +import org.apache.rocketmq.client.producer.SendResult; +import org.apache.rocketmq.client.producer.SendStatus; +import org.apache.rocketmq.common.BrokerConfig; +import org.apache.rocketmq.common.MixAll; +import org.apache.rocketmq.common.UtilAll; +import org.apache.rocketmq.common.attribute.CQType; +import org.apache.rocketmq.common.message.Message; +import org.apache.rocketmq.common.message.MessageExt; +import org.apache.rocketmq.common.message.MessageQueue; +import org.apache.rocketmq.namesrv.NamesrvController; +import org.apache.rocketmq.remoting.netty.NettyClientConfig; +import org.apache.rocketmq.remoting.netty.NettyServerConfig; +import org.apache.rocketmq.remoting.protocol.body.ClusterInfo; +import org.apache.rocketmq.remoting.protocol.route.BrokerData; +import org.apache.rocketmq.remoting.protocol.route.TopicRouteData; +import org.apache.rocketmq.store.config.BrokerRole; +import org.apache.rocketmq.store.config.MessageStoreConfig; +import org.apache.rocketmq.store.dledger.DLedgerCommitLog; +import org.apache.rocketmq.test.base.IntegrationTestBase; +import org.apache.rocketmq.tools.admin.DefaultMQAdminExt; +import org.junit.Assert; +import org.junit.Test; + +import static org.awaitility.Awaitility.await; + +public class DLedgerThreeNodeIT { + private static final long AWAIT_SECONDS = 60; + private static final List NODE_IDS = Arrays.asList("n0", "n1", "n2"); + + @Test + public void testProduceFailoverAndRestart() throws Exception { + NamesrvController namesrvController = null; + DefaultMQAdminExt admin = null; + List nodeSpecs = new ArrayList<>(); + List allControllers = new ArrayList<>(); + Set stopped = + Collections.newSetFromMap(new IdentityHashMap()); + try { + namesrvController = IntegrationTestBase.createAndStartNamesrv(); + String namesrvAddr = "127.0.0.1:" + + namesrvController.getNettyServerConfig().getListenPort(); + admin = new DefaultMQAdminExt(); + admin.setInstanceName(UUID.randomUUID().toString()); + admin.setNamesrvAddr(namesrvAddr); + admin.start(); + + String clusterName = "DLedgerCluster-" + UUID.randomUUID(); + String brokerName = "DLedgerBroker-" + UUID.randomUUID(); + String topic = "DLedgerTopic-" + UUID.randomUUID(); + ClusterPorts clusterPorts = allocateClusterPorts(NODE_IDS.size()); + String peers = buildPeers(clusterPorts.dLedgerPorts); + for (int i = 0; i < NODE_IDS.size(); i++) { + nodeSpecs.add(new NodeSpec( + NODE_IDS.get(i), clusterPorts.dLedgerPorts.get(i), + clusterPorts.brokerPorts.get(i), + IntegrationTestBase.createBaseDir())); + } + + List active = startCluster( + nodeSpecs, clusterName, brokerName, namesrvAddr, peers, allControllers); + BrokerController initialLeader = awaitLeader(active); + awaitClusterMaster(admin, brokerName, initialLeader); + Assert.assertTrue(IntegrationTestBase.initTopic( + topic, namesrvAddr, clusterName, 1, CQType.SimpleCQ)); + awaitTopicRouteMaster(admin, topic, brokerName, initialLeader); + awaitTopicOnEveryNode(active, topic); + + List expectedBodies = new ArrayList<>(); + expectedBodies.add("before-single"); + expectedBodies.add("before-batch-0"); + expectedBodies.add("before-batch-1"); + expectedBodies.add("before-batch-2"); + sendInitialSingleAndBatch(namesrvAddr, topic, brokerName); + awaitQueueOffset(active, topic, expectedBodies.size()); + assertBodies(pullExactly( + namesrvAddr, topic, brokerName, expectedBodies.size()), expectedBodies); + + stopController(initialLeader, stopped); + active.remove(initialLeader); + BrokerController failoverLeader = awaitLeader(active); + Assert.assertNotSame(initialLeader, failoverLeader); + awaitClusterMaster(admin, brokerName, failoverLeader); + awaitTopicRouteMaster(admin, topic, brokerName, failoverLeader); + sendOne(namesrvAddr, topic, brokerName, + "after-failover", expectedBodies.size()); + expectedBodies.add("after-failover"); + awaitQueueOffset(active, topic, expectedBodies.size()); + assertBodies(pullExactly( + namesrvAddr, topic, brokerName, expectedBodies.size()), expectedBodies); + + stopControllers(active, stopped); + awaitBrokerRegistrationRemoved(admin, brokerName); + awaitAllNodePortsAvailable(nodeSpecs); + active = startCluster( + nodeSpecs, clusterName, brokerName, namesrvAddr, peers, allControllers); + BrokerController restartedLeader = awaitLeader(active); + awaitClusterMaster(admin, brokerName, restartedLeader); + awaitTopicRouteMaster(admin, topic, brokerName, restartedLeader); + awaitTopicOnEveryNode(active, topic); + awaitQueueOffset(active, topic, expectedBodies.size()); + + // This pull is deliberately before the first post-restart user append. + assertBodies(pullExactly( + namesrvAddr, topic, brokerName, expectedBodies.size()), expectedBodies); + sendOne(namesrvAddr, topic, brokerName, + "after-restart", expectedBodies.size()); + expectedBodies.add("after-restart"); + awaitQueueOffset(active, topic, expectedBodies.size()); + assertBodies(pullExactly( + namesrvAddr, topic, brokerName, expectedBodies.size()), expectedBodies); + } finally { + stopControllers(allControllers, stopped); + if (admin != null) { + admin.shutdown(); + } + if (namesrvController != null) { + namesrvController.shutdown(); + } + for (NodeSpec nodeSpec : nodeSpecs) { + UtilAll.deleteFile(new File(nodeSpec.storeRoot)); + } + } + } + + private static List startCluster(List nodeSpecs, + String clusterName, String brokerName, String namesrvAddr, String peers, + List allControllers) throws Exception { + List controllers = new ArrayList<>(); + for (NodeSpec nodeSpec : nodeSpecs) { + BrokerController controller = startNode( + nodeSpec, clusterName, brokerName, namesrvAddr, peers); + controllers.add(controller); + allControllers.add(controller); + } + return controllers; + } + + private static BrokerController startNode(NodeSpec nodeSpec, String clusterName, + String brokerName, String namesrvAddr, String peers) throws Exception { + BrokerConfig brokerConfig = new BrokerConfig(); + brokerConfig.setBrokerClusterName(clusterName); + brokerConfig.setBrokerName(brokerName); + brokerConfig.setBrokerIP1("127.0.0.1"); + brokerConfig.setBrokerIP2("127.0.0.1"); + brokerConfig.setNamesrvAddr(namesrvAddr); + brokerConfig.setRegisterNameServerPeriod(1000); + brokerConfig.setLoadBalancePollNameServerInterval(500); + + MessageStoreConfig storeConfig = new MessageStoreConfig(); + storeConfig.setStorePathRootDir(nodeSpec.storeRoot); + storeConfig.setStorePathCommitLog( + nodeSpec.storeRoot + File.separator + "commitlog"); + storeConfig.setStorePathDLedgerCommitLog( + nodeSpec.storeRoot + File.separator + "dledger"); + storeConfig.setMappedFileSizeCommitLog(1024 * 1024); + storeConfig.setMaxHashSlotNum(10_000); + storeConfig.setMaxIndexNum(10_000); + storeConfig.setHaListenPort(0); + storeConfig.setEnableDLegerCommitLog(true); + storeConfig.setdLegerGroup(brokerName); + storeConfig.setdLegerSelfId(nodeSpec.selfId); + storeConfig.setdLegerPeers(peers); + storeConfig.setEnableBatchPush(true); + + NettyServerConfig serverConfig = new NettyServerConfig(); + serverConfig.setListenPort(nodeSpec.brokerPort); + BrokerController controller = new BrokerController( + brokerConfig, serverConfig, new NettyClientConfig(), storeConfig); + try { + Assert.assertTrue(controller.initialize()); + controller.start(); + return controller; + } catch (Throwable t) { + try { + controller.shutdown(); + } catch (Throwable ignored) { + } + if (t instanceof Error) { + throw (Error) t; + } + if (t instanceof Exception) { + throw (Exception) t; + } + throw new RuntimeException(t); + } + } + + private static BrokerController awaitLeader(List controllers) { + AtomicReference result = new AtomicReference<>(); + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + BrokerController leader = findLeader(controllers); + if (leader == null) { + return false; + } + result.set(leader); + return true; + }); + return result.get(); + } + + private static BrokerController findLeader(List controllers) { + BrokerController result = null; + for (BrokerController controller : controllers) { + DLedgerCommitLog commitLog = + (DLedgerCommitLog) controller.getMessageStore().getCommitLog(); + boolean dLedgerLeader = + commitLog.getdLedgerServer().getMemberState().isLeader(); + boolean brokerMaster = controller.getMessageStoreConfig().getBrokerRole() + == BrokerRole.SYNC_MASTER; + boolean brokerIdIsMaster = + controller.getBrokerConfig().getBrokerId() == MixAll.MASTER_ID; + if (dLedgerLeader && brokerMaster && brokerIdIsMaster) { + if (result != null) { + return null; + } + result = controller; + } + } + return result; + } + + private static void awaitClusterMaster(DefaultMQAdminExt admin, String brokerName, + BrokerController expectedLeader) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + try { + ClusterInfo clusterInfo = admin.examineBrokerClusterInfo(); + BrokerData brokerData = clusterInfo.getBrokerAddrTable().get(brokerName); + if (brokerData == null) { + return false; + } + return expectedLeader.getBrokerAddr().equals( + brokerData.getBrokerAddrs().get(MixAll.MASTER_ID)); + } catch (Exception ignored) { + return false; + } + }); + } + + private static void awaitBrokerRegistrationRemoved( + DefaultMQAdminExt admin, String brokerName) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS) + .until(() -> !admin.examineBrokerClusterInfo() + .getBrokerAddrTable().containsKey(brokerName)); + } + + private static void awaitTopicRouteMaster(DefaultMQAdminExt admin, String topic, + String brokerName, BrokerController expectedLeader) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + try { + TopicRouteData route = admin.examineTopicRouteInfo(topic); + for (BrokerData brokerData : route.getBrokerDatas()) { + if (brokerName.equals(brokerData.getBrokerName())) { + return expectedLeader.getBrokerAddr().equals( + brokerData.getBrokerAddrs().get(MixAll.MASTER_ID)); + } + } + } catch (Exception ignored) { + } + return false; + }); + } + + private static void awaitTopicOnEveryNode( + List controllers, String topic) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + for (BrokerController controller : controllers) { + if (controller.getTopicConfigManager().selectTopicConfig(topic) == null) { + return false; + } + } + return true; + }); + } + + private static void awaitQueueOffset(List controllers, + String topic, long expectedOffset) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + for (BrokerController controller : controllers) { + if (controller.getMessageStore().getMaxOffsetInQueue(topic, 0) + != expectedOffset + || controller.getMessageStore().dispatchBehindBytes() != 0) { + return false; + } + } + return true; + }); + } + + private static void sendInitialSingleAndBatch( + String namesrvAddr, String topic, String brokerName) throws Exception { + DefaultMQProducer producer = startProducer(namesrvAddr); + try { + MessageQueue queue = awaitPublishQueue(producer, topic, brokerName); + SendResult singleResult = producer.send(new Message( + topic, "before-single".getBytes(StandardCharsets.UTF_8)), queue); + assertSendResult(singleResult, brokerName, 0); + Assert.assertNotNull(singleResult.getOffsetMsgId()); + + List batch = new ArrayList<>(); + for (int i = 0; i < 3; i++) { + batch.add(new Message(topic, + ("before-batch-" + i).getBytes(StandardCharsets.UTF_8))); + } + SendResult batchResult = producer.send(batch, queue); + assertSendResult(batchResult, brokerName, 1); + Assert.assertEquals(3, batchResult.getMsgId().split(",").length); + } finally { + producer.shutdown(); + } + } + + private static void sendOne(String namesrvAddr, String topic, String brokerName, + String body, long expectedQueueOffset) throws Exception { + DefaultMQProducer producer = startProducer(namesrvAddr); + try { + MessageQueue queue = awaitPublishQueue(producer, topic, brokerName); + SendResult result = producer.send( + new Message(topic, body.getBytes(StandardCharsets.UTF_8)), queue); + assertSendResult(result, brokerName, expectedQueueOffset); + Assert.assertNotNull(result.getOffsetMsgId()); + } finally { + producer.shutdown(); + } + } + + private static DefaultMQProducer startProducer(String namesrvAddr) + throws Exception { + DefaultMQProducer producer = + new DefaultMQProducer("dledger-it-" + UUID.randomUUID()); + producer.setInstanceName(UUID.randomUUID().toString()); + producer.setNamesrvAddr(namesrvAddr); + producer.setPollNameServerInterval(500); + producer.setSendMsgTimeout(10_000); + producer.setRetryTimesWhenSendFailed(3); + producer.setVipChannelEnabled(false); + producer.start(); + return producer; + } + + private static MessageQueue awaitPublishQueue(DefaultMQProducer producer, + String topic, String brokerName) { + AtomicReference result = new AtomicReference<>(); + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + try { + MessageQueue queue = selectQueue( + producer.fetchPublishMessageQueues(topic), brokerName); + if (queue == null) { + return false; + } + result.set(queue); + return true; + } catch (Exception ignored) { + return false; + } + }); + return result.get(); + } + + private static List pullExactly(String namesrvAddr, String topic, + String brokerName, int expectedCount) throws Exception { + DefaultMQPullConsumer consumer = + new DefaultMQPullConsumer("dledger-it-" + UUID.randomUUID()); + consumer.setInstanceName(UUID.randomUUID().toString()); + consumer.setNamesrvAddr(namesrvAddr); + consumer.setPollNameServerInterval(500); + consumer.setConsumerPullTimeoutMillis(3_000); + consumer.setVipChannelEnabled(false); + consumer.start(); + AtomicReference resultRef = new AtomicReference<>(); + try { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS).until(() -> { + try { + MessageQueue queue = selectQueue( + consumer.fetchSubscribeMessageQueues(topic), brokerName); + if (queue == null) { + return false; + } + PullResult result = consumer.pull( + queue, "*", 0, Math.max(32, expectedCount)); + if (result.getPullStatus() != PullStatus.FOUND + || result.getMinOffset() != 0 + || result.getMaxOffset() != expectedCount + || result.getMsgFoundList().size() != expectedCount) { + return false; + } + resultRef.set(result); + return true; + } catch (Exception ignored) { + return false; + } + }); + return new ArrayList<>(resultRef.get().getMsgFoundList()); + } finally { + consumer.shutdown(); + } + } + + private static MessageQueue selectQueue( + Iterable queues, String brokerName) { + for (MessageQueue queue : queues) { + if (brokerName.equals(queue.getBrokerName()) && queue.getQueueId() == 0) { + return queue; + } + } + return null; + } + + private static void assertSendResult( + SendResult result, String brokerName, long expectedQueueOffset) { + Assert.assertEquals(SendStatus.SEND_OK, result.getSendStatus()); + Assert.assertEquals(brokerName, result.getMessageQueue().getBrokerName()); + Assert.assertEquals(0, result.getMessageQueue().getQueueId()); + Assert.assertEquals(expectedQueueOffset, result.getQueueOffset()); + Assert.assertNotNull(result.getMsgId()); + } + + private static void assertBodies( + List messages, List expectedBodies) { + Assert.assertEquals(expectedBodies.size(), messages.size()); + for (int i = 0; i < expectedBodies.size(); i++) { + MessageExt message = messages.get(i); + Assert.assertEquals(i, message.getQueueOffset()); + Assert.assertArrayEquals( + expectedBodies.get(i).getBytes(StandardCharsets.UTF_8), message.getBody()); + } + } + + private static void stopControllers(List controllers, + Set stopped) { + for (BrokerController controller : new ArrayList<>(controllers)) { + stopController(controller, stopped); + } + } + + private static void stopController(BrokerController controller, + Set stopped) { + if (controller == null || !stopped.add(controller)) { + return; + } + try { + controller.shutdown(); + } catch (Throwable ignored) { + } + } + + private static ClusterPorts allocateClusterPorts(int count) throws IOException { + List reservations = new ArrayList<>(); + List dLedgerPorts = new ArrayList<>(); + List brokerPorts = new ArrayList<>(); + try { + InetAddress loopback = InetAddress.getByName("127.0.0.1"); + while (dLedgerPorts.size() < count) { + ServerSocket socket = new ServerSocket(0, 50, loopback); + if (socket.getLocalPort() <= 1024) { + socket.close(); + continue; + } + reservations.add(socket); + dLedgerPorts.add(socket.getLocalPort()); + } + int attempts = 0; + while (brokerPorts.size() < count) { + if (++attempts > 1000) { + throw new IOException("Unable to reserve broker and VIP port pairs"); + } + ServerSocket brokerSocket = new ServerSocket(0, 50, loopback); + int brokerPort = brokerSocket.getLocalPort(); + int fastPort = brokerPort - 2; + if (fastPort <= 1024) { + brokerSocket.close(); + continue; + } + try { + ServerSocket fastSocket = new ServerSocket(fastPort, 50, loopback); + reservations.add(brokerSocket); + reservations.add(fastSocket); + brokerPorts.add(brokerPort); + } catch (IOException ignored) { + brokerSocket.close(); + } + } + return new ClusterPorts(dLedgerPorts, brokerPorts); + } finally { + closeSockets(reservations); + } + } + + private static String buildPeers(List ports) { + StringBuilder peers = new StringBuilder(); + for (int i = 0; i < NODE_IDS.size(); i++) { + if (i > 0) { + peers.append(';'); + } + peers.append(NODE_IDS.get(i)).append("-127.0.0.1:").append(ports.get(i)); + } + return peers.toString(); + } + + private static void awaitAllNodePortsAvailable(List nodeSpecs) { + await().atMost(AWAIT_SECONDS, TimeUnit.SECONDS) + .pollInterval(200, TimeUnit.MILLISECONDS) + .until(() -> areAllNodePortsAvailable(nodeSpecs)); + } + + private static boolean areAllNodePortsAvailable(List nodeSpecs) { + List probes = new ArrayList<>(); + try { + InetAddress loopback = InetAddress.getByName("127.0.0.1"); + for (NodeSpec nodeSpec : nodeSpecs) { + probes.add(new ServerSocket(nodeSpec.dLedgerPort, 50, loopback)); + probes.add(new ServerSocket(nodeSpec.brokerPort, 50, loopback)); + probes.add(new ServerSocket(nodeSpec.brokerPort - 2, 50, loopback)); + } + return true; + } catch (IOException ignored) { + return false; + } finally { + closeSockets(probes); + } + } + + private static void closeSockets(List sockets) { + for (ServerSocket socket : sockets) { + try { + socket.close(); + } catch (IOException ignored) { + } + } + } + + private static final class NodeSpec { + private final String selfId; + private final int dLedgerPort; + private final int brokerPort; + private final String storeRoot; + + private NodeSpec(String selfId, int dLedgerPort, int brokerPort, + String storeRoot) { + this.selfId = selfId; + this.dLedgerPort = dLedgerPort; + this.brokerPort = brokerPort; + this.storeRoot = storeRoot; + } + } + + private static final class ClusterPorts { + private final List dLedgerPorts; + private final List brokerPorts; + + private ClusterPorts( + List dLedgerPorts, List brokerPorts) { + this.dLedgerPorts = dLedgerPorts; + this.brokerPorts = brokerPorts; + } + } +}