This is an automated email from the ASF dual-hosted git repository.
szetszwo pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ratis.git
The following commit(s) were added to refs/heads/master by this push:
new 6aafce539 RATIS-2605. Avoid advancing matchIndex from heartbeat
AppendEntries success (#1519)
6aafce539 is described below
commit 6aafce539ca3c280db6e91fe8a6e02b345001952
Author: Ethan Feng <[email protected]>
AuthorDate: Wed Jul 15 02:16:06 2026 +0800
RATIS-2605. Avoid advancing matchIndex from heartbeat AppendEntries success
(#1519)
---
.../ratis/server/leader/LogAppenderDefault.java | 30 ++-
.../server/leader/TestLogAppenderDefault.java | 255 +++++++++++++++++++++
2 files changed, 279 insertions(+), 6 deletions(-)
diff --git
a/ratis-server/src/main/java/org/apache/ratis/server/leader/LogAppenderDefault.java
b/ratis-server/src/main/java/org/apache/ratis/server/leader/LogAppenderDefault.java
index 9d1edd469..66c35873a 100644
---
a/ratis-server/src/main/java/org/apache/ratis/server/leader/LogAppenderDefault.java
+++
b/ratis-server/src/main/java/org/apache/ratis/server/leader/LogAppenderDefault.java
@@ -55,8 +55,24 @@ class LogAppenderDefault extends LogAppenderBase {
return CallId.getComparator();
}
+ private static final class AppendEntriesResult {
+ private final AppendEntriesReplyProto reply;
+ private final long requestFirstIndex;
+ private final long requestPreviousIndex;
+ private final boolean requestHeartbeat;
+
+ private AppendEntriesResult(AppendEntriesReplyProto reply,
AppendEntriesRequestProto request) {
+ this.reply = reply;
+ this.requestFirstIndex = request.getEntriesCount() > 0
+ ? request.getEntries(0).getIndex() : RaftLog.INVALID_LOG_INDEX;
+ this.requestPreviousIndex = request.hasPreviousLog()
+ ? request.getPreviousLog().getIndex() : RaftLog.INVALID_LOG_INDEX;
+ this.requestHeartbeat = request.getEntriesCount() == 0;
+ }
+ }
+
/** Send an appendEntries RPC; retry indefinitely. */
- private AppendEntriesReplyProto sendAppendEntriesWithRetries(AtomicLong
requestFirstIndex)
+ private AppendEntriesResult sendAppendEntriesWithRetries(AtomicLong
requestFirstIndex)
throws InterruptedException, InterruptedIOException, RaftLogIOException {
int retry = 0;
@@ -86,7 +102,7 @@ class LogAppenderDefault extends LogAppenderBase {
getFollower().updateLastRespondedAppendEntriesSendTime(sendTime);
getLeaderState().onFollowerCommitIndex(getFollower(),
reply.getFollowerCommit());
- return reply;
+ return new AppendEntriesResult(reply, proto);
} catch (InterruptedIOException | RaftLogIOException e) {
throw e;
} catch (IOException ioe) {
@@ -160,9 +176,9 @@ class LogAppenderDefault extends LogAppenderBase {
// otherwise if r is null, retry the snapshot installation
} else {
final AtomicLong requestFirstIndex = new
AtomicLong(RaftLog.INVALID_LOG_INDEX);
- final AppendEntriesReplyProto r =
sendAppendEntriesWithRetries(requestFirstIndex);
+ final AppendEntriesResult r =
sendAppendEntriesWithRetries(requestFirstIndex);
if (r != null) {
- handleReply(r, requestFirstIndex.get());
+ handleReply(r.reply, r.requestFirstIndex, r.requestPreviousIndex,
r.requestHeartbeat);
}
}
}
@@ -173,13 +189,15 @@ class LogAppenderDefault extends LogAppenderBase {
}
}
- private void handleReply(AppendEntriesReplyProto reply, long
requestFirstIndex)
+ void handleReply(AppendEntriesReplyProto reply, long requestFirstIndex,
+ long requestPreviousIndex, boolean requestHeartbeat)
throws IllegalArgumentException {
if (reply != null) {
switch (reply.getResult()) {
case SUCCESS:
final long oldNextIndex = getFollower().getNextIndex();
- final long nextIndex = reply.getNextIndex();
+ final long nextIndex = requestHeartbeat
+ ? Math.max(oldNextIndex, requestPreviousIndex + 1) :
reply.getNextIndex();
if (nextIndex < oldNextIndex) {
throw new IllegalStateException("nextIndex=" + nextIndex
+ " < oldNextIndex=" + oldNextIndex
diff --git
a/ratis-server/src/test/java/org/apache/ratis/server/leader/TestLogAppenderDefault.java
b/ratis-server/src/test/java/org/apache/ratis/server/leader/TestLogAppenderDefault.java
new file mode 100644
index 000000000..c17c05241
--- /dev/null
+++
b/ratis-server/src/test/java/org/apache/ratis/server/leader/TestLogAppenderDefault.java
@@ -0,0 +1,255 @@
+/*
+ * 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.ratis.server.leader;
+
+import org.apache.ratis.conf.RaftProperties;
+import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto;
+import org.apache.ratis.protocol.RaftPeer;
+import org.apache.ratis.protocol.RaftPeerId;
+import org.apache.ratis.server.RaftServer;
+import org.apache.ratis.server.raftlog.RaftLog;
+import org.apache.ratis.util.Timestamp;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.LongUnaryOperator;
+
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.when;
+
+class TestLogAppenderDefault {
+ private static final RaftPeerId FOLLOWER_ID = RaftPeerId.valueOf("follower");
+
+ @Test
+ void heartbeatSuccessAdvancesOnlyToRequestPreviousIndex() {
+ final LeaderState leaderState = mock(LeaderState.class);
+ final TestFollowerInfo follower = new TestFollowerInfo(8, 4);
+ final LogAppenderDefault appender = newLogAppender(leaderState, follower);
+ final AppendEntriesReplyProto reply = newSuccessReply(20);
+
+ appender.handleReply(reply, RaftLog.INVALID_LOG_INDEX, 9, true);
+
+ Assertions.assertEquals(9, follower.getMatchIndex());
+ Assertions.assertEquals(10, follower.getNextIndex());
+ Assertions.assertEquals(1, follower.successfulMatchIndexUpdates.get());
+ Assertions.assertEquals(1, follower.nextIndexIncreases.get());
+ verify(leaderState).onFollowerSuccessAppendEntries(follower);
+ verify(leaderState).onAppendEntriesReply(appender, reply);
+ verifyNoMoreInteractions(leaderState);
+ }
+
+ @Test
+ void heartbeatSuccessWithoutPreviousLogDoesNotAdvance() {
+ final LeaderState leaderState = mock(LeaderState.class);
+ final TestFollowerInfo follower = new TestFollowerInfo(8, 4);
+ final LogAppenderDefault appender = newLogAppender(leaderState, follower);
+ final AppendEntriesReplyProto reply = newSuccessReply(20);
+
+ appender.handleReply(reply, RaftLog.INVALID_LOG_INDEX,
RaftLog.INVALID_LOG_INDEX, true);
+
+ Assertions.assertEquals(4, follower.getMatchIndex());
+ Assertions.assertEquals(8, follower.getNextIndex());
+ Assertions.assertEquals(0, follower.successfulMatchIndexUpdates.get());
+ Assertions.assertEquals(0, follower.nextIndexIncreases.get());
+ verify(leaderState).onAppendEntriesReply(appender, reply);
+ verifyNoMoreInteractions(leaderState);
+ }
+
+ @Test
+ void appendSuccessStillUsesReplyNextIndex() {
+ final LeaderState leaderState = mock(LeaderState.class);
+ final TestFollowerInfo follower = new TestFollowerInfo(8, 4);
+ final LogAppenderDefault appender = newLogAppender(leaderState, follower);
+ final AppendEntriesReplyProto reply = newSuccessReply(20);
+
+ appender.handleReply(reply, 8, 7, false);
+
+ Assertions.assertEquals(19, follower.getMatchIndex());
+ Assertions.assertEquals(20, follower.getNextIndex());
+ Assertions.assertEquals(1, follower.successfulMatchIndexUpdates.get());
+ Assertions.assertEquals(1, follower.nextIndexIncreases.get());
+ verify(leaderState).onFollowerSuccessAppendEntries(follower);
+ verify(leaderState).onAppendEntriesReply(appender, reply);
+ verifyNoMoreInteractions(leaderState);
+ }
+
+ private static LogAppenderDefault newLogAppender(LeaderState leaderState,
FollowerInfo follower) {
+ final RaftServer.Division division = mock(RaftServer.Division.class);
+ final RaftServer raftServer = mock(RaftServer.class);
+ when(division.getRaftServer()).thenReturn(raftServer);
+
when(division.getThreadGroup()).thenReturn(Thread.currentThread().getThreadGroup());
+ when(raftServer.getProperties()).thenReturn(new RaftProperties());
+ return new LogAppenderDefault(division, leaderState, follower);
+ }
+
+ private static AppendEntriesReplyProto newSuccessReply(long nextIndex) {
+ return AppendEntriesReplyProto.newBuilder()
+ .setResult(AppendEntriesReplyProto.AppendResult.SUCCESS)
+ .setNextIndex(nextIndex)
+ .build();
+ }
+
+ private static final class TestFollowerInfo implements FollowerInfo {
+ private final RaftPeer peer =
RaftPeer.newBuilder().setId(FOLLOWER_ID).build();
+ private final AtomicInteger successfulMatchIndexUpdates = new
AtomicInteger();
+ private final AtomicInteger nextIndexIncreases = new AtomicInteger();
+ private long matchIndex;
+ private long nextIndex;
+
+ private TestFollowerInfo(long nextIndex, long matchIndex) {
+ this.nextIndex = nextIndex;
+ this.matchIndex = matchIndex;
+ }
+
+ @Override
+ public String getName() {
+ return FOLLOWER_ID.toString();
+ }
+
+ @Override
+ public RaftPeerId getId() {
+ return FOLLOWER_ID;
+ }
+
+ @Override
+ public RaftPeer getPeer() {
+ return peer;
+ }
+
+ @Override
+ public long getMatchIndex() {
+ return matchIndex;
+ }
+
+ @Override
+ public boolean updateMatchIndex(long newMatchIndex) {
+ if (newMatchIndex > matchIndex) {
+ matchIndex = newMatchIndex;
+ successfulMatchIndexUpdates.incrementAndGet();
+ return true;
+ }
+ return false;
+ }
+
+ @Override
+ public long getCommitIndex() {
+ return 0;
+ }
+
+ @Override
+ public boolean updateCommitIndex(long newCommitIndex) {
+ return false;
+ }
+
+ @Override
+ public long getSnapshotIndex() {
+ return 0;
+ }
+
+ @Override
+ public void setSnapshotIndex(long newSnapshotIndex) {
+ }
+
+ @Override
+ public void setAttemptedToInstallSnapshot() {
+ }
+
+ @Override
+ public boolean hasAttemptedToInstallSnapshot() {
+ return false;
+ }
+
+ @Override
+ public long getNextIndex() {
+ return nextIndex;
+ }
+
+ @Override
+ public void increaseNextIndex(long newNextIndex) {
+ if (newNextIndex > nextIndex) {
+ nextIndex = newNextIndex;
+ nextIndexIncreases.incrementAndGet();
+ }
+ }
+
+ @Override
+ public void decreaseNextIndex(long newNextIndex) {
+ nextIndex = Math.min(nextIndex, newNextIndex);
+ }
+
+ @Override
+ public void setNextIndex(long newNextIndex) {
+ nextIndex = newNextIndex;
+ }
+
+ @Override
+ public void updateNextIndex(long newNextIndex) {
+ nextIndex = newNextIndex;
+ }
+
+ @Override
+ public void computeNextIndex(LongUnaryOperator op) {
+ nextIndex = op.applyAsLong(nextIndex);
+ }
+
+ @Override
+ public Timestamp getLastRpcResponseTime() {
+ return Timestamp.currentTime();
+ }
+
+ @Override
+ public Timestamp getLastRpcSendTime() {
+ return Timestamp.currentTime();
+ }
+
+ @Override
+ public void updateLastRpcResponseTime() {
+ }
+
+ @Override
+ public void updateLastRpcSendTime(boolean isHeartbeat) {
+ }
+
+ @Override
+ public Timestamp getLastRpcTime() {
+ return Timestamp.currentTime();
+ }
+
+ @Override
+ public Timestamp getLastHeartbeatSendTime() {
+ return Timestamp.currentTime();
+ }
+
+ @Override
+ public Timestamp getLastRespondedAppendEntriesSendTime() {
+ return Timestamp.currentTime();
+ }
+
+ @Override
+ public void updateLastRespondedAppendEntriesSendTime(Timestamp sendTime) {
+ }
+
+ @Override
+ public ErrorState getErrorState() {
+ return null;
+ }
+ }
+}