fix: handle DLedger control entries during recovery

This commit is contained in:
通融
2026-08-15 19:37:16 +08:00
parent 802a8b0f48
commit f337446439
5 changed files with 748 additions and 33 deletions
@@ -40,6 +40,7 @@ import org.junit.Test;
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.assertNotNull;
import static org.junit.Assert.assertTrue;
public class DLedgerControllerTest {
private static int port = 30000;
private List<String> baseDirs;
private List<DLedgerController> 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 class DLedgerControllerTest {
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 class DLedgerControllerTest {
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<String> selfIds = controllers.stream()
.map(controller -> controller.getControllerConfig().getControllerDLegerSelfId())
.collect(Collectors.toList());
List<DLedgerController> 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 class DLedgerControllerTest {
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);
+1
View File
@@ -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",
@@ -2741,16 +2741,19 @@ public class DefaultMessageStore implements MessageStore {
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());
@@ -22,6 +22,7 @@ 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;
@@ -103,12 +104,16 @@ public class DLedgerCommitLog extends CommitLog {
dLedgerConfig.setFileReservedHours(defaultMessageStore.getMessageStoreConfig().getFileReservedTime() + 1);
dLedgerConfig.setPreferredLeaderId(defaultMessageStore.getMessageStoreConfig().getPreferredLeaderId());
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);
@@ -406,19 +411,24 @@ public class DLedgerCommitLog extends CommitLog {
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());
@@ -486,16 +496,53 @@ public class DLedgerCommitLog extends CommitLog {
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;
}
int magicCode = byteBuffer.getInt(DLedgerEntry.BODY_OFFSET + MessageDecoder.MESSAGE_MAGIC_CODE_POSITION);
if (magicCode != MESSAGE_MAGIC_CODE) {
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(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 {
@@ -503,11 +550,15 @@ public class DLedgerCommitLog extends CommitLog {
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()) {
@@ -537,30 +588,54 @@ public class DLedgerCommitLog extends CommitLog {
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 */);
@@ -0,0 +1,592 @@
/*
* 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 batchResult = putBatch(messageStore, topic, 3, 1);
Assert.assertTrue(singleResult.getAppendMessageResult().getWroteOffset() > 0);
Assert.assertTrue(batchResult.getAppendMessageResult().getWroteOffset()
> singleResult.getAppendMessageResult().getWroteOffset());
Assert.assertEquals(3, batchResult.getAppendMessageResult().getMsgNum());
Assert.assertNotNull(singleResult.getAppendMessageResult().getMsgId());
Assert.assertNotNull(batchResult.getAppendMessageResult().getMsgId());
Assert.assertEquals(3, batchResult.getAppendMessageResult().getMsgId().split(",").length);
awaitStoreReady(messageStore, topic, 4);
Assert.assertEquals(0, messageStore.getMinOffsetInQueue(topic, QUEUE_ID));
Assert.assertTrue(commitLog(messageStore).getCommittedPos()
> batchResult.getAppendMessageResult().getWroteOffset());
doGetMessages(messageStore, topic, QUEUE_ID, 4, 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<DefaultMessageStore> 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<DefaultMessageStore> 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<byte[]> 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<DefaultMessageStore> stores) {
AtomicReference<DefaultMessageStore> 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<byte[]> readMessageBodies(DefaultMessageStore messageStore, String topic, int queueId,
int messageCount) {
List<byte[]> 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<byte[]> expectedBodies) {
List<byte[]> 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<Status> 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();
}
}
}