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

Reply via email to