Repository: incubator-ratis
Updated Branches:
  refs/heads/master 05fb38bf1 -> 00f64b4c1


RATIS-217. Support timeout for append entry requests(GrpcLogAppender).  
Contributed by Lokesh Jain


Project: http://git-wip-us.apache.org/repos/asf/incubator-ratis/repo
Commit: http://git-wip-us.apache.org/repos/asf/incubator-ratis/commit/00f64b4c
Tree: http://git-wip-us.apache.org/repos/asf/incubator-ratis/tree/00f64b4c
Diff: http://git-wip-us.apache.org/repos/asf/incubator-ratis/diff/00f64b4c

Branch: refs/heads/master
Commit: 00f64b4c17dfa2f5bd57d25aec5c5eb24c1bb535
Parents: 05fb38b
Author: Tsz-Wo Nicholas Sze <[email protected]>
Authored: Sat Apr 7 21:51:02 2018 +0800
Committer: Tsz-Wo Nicholas Sze <[email protected]>
Committed: Sat Apr 7 21:51:02 2018 +0800

----------------------------------------------------------------------
 .../java/org/apache/ratis/rpc/RpcTimeout.java   | 39 ++++++++++-----
 .../org/apache/ratis/grpc/RaftGrpcUtil.java     | 20 +++++++-
 .../grpc/client/RaftClientProtocolClient.java   | 15 +++---
 .../ratis/grpc/server/GRpcLogAppender.java      | 50 ++++++++++++++++----
 .../grpc/server/RaftServerProtocolService.java  |  2 +-
 .../apache/ratis/server/impl/LeaderState.java   |  7 +--
 .../apache/ratis/server/impl/LogAppender.java   | 11 +++--
 .../ratis/server/impl/RaftServerImpl.java       | 12 ++---
 .../ratis/server/impl/ServerProtoUtils.java     | 14 ++++--
 .../java/org/apache/ratis/RaftAsyncTests.java   | 41 ++++++++++++++++
 .../SimpleStateMachine4Testing.java             | 30 +++++++++++-
 11 files changed, 191 insertions(+), 50 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-common/src/main/java/org/apache/ratis/rpc/RpcTimeout.java
----------------------------------------------------------------------
diff --git a/ratis-common/src/main/java/org/apache/ratis/rpc/RpcTimeout.java 
b/ratis-common/src/main/java/org/apache/ratis/rpc/RpcTimeout.java
index 6c2c7fc..03d7eab 100644
--- a/ratis-common/src/main/java/org/apache/ratis/rpc/RpcTimeout.java
+++ b/ratis-common/src/main/java/org/apache/ratis/rpc/RpcTimeout.java
@@ -17,39 +17,54 @@
  */
 package org.apache.ratis.rpc;
 
+import org.apache.ratis.shaded.com.google.common.base.Supplier;
 import org.apache.ratis.util.Preconditions;
 import org.apache.ratis.util.TimeDuration;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 import java.util.concurrent.Executors;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 
 public class RpcTimeout {
+  private static final Logger LOG = LoggerFactory.getLogger(RpcTimeout.class);
   private ScheduledExecutorService timeoutScheduler = null;
-  private TimeDuration callTimeout;
+  private final TimeDuration callTimeout;
+  private int numUsers = 0;
 
-  public RpcTimeout(TimeDuration callTimeout, boolean initialize) {
+  public RpcTimeout(TimeDuration callTimeout) {
     this.callTimeout = callTimeout;
-    if (initialize) {
-      initialize();
+  }
+
+  public synchronized void addUser() {
+    if (timeoutScheduler == null) {
+      timeoutScheduler = Executors.newScheduledThreadPool(1);
     }
+    numUsers++;
   }
 
-  public synchronized void initialize() {
-    timeoutScheduler = Executors.newScheduledThreadPool(1);
+  public synchronized void removeUser() {
+    numUsers--;
+    if (timeoutScheduler != null && numUsers == 0) {
+      timeoutScheduler.shutdown();
+      timeoutScheduler = null;
+    }
   }
 
-  public void onTimeout(Runnable task) {
+  public synchronized void onTimeout(Runnable task, Supplier<String> errorMsg) 
{
     Preconditions.assertTrue(timeoutScheduler != null);
     TimeUnit unit = callTimeout.getUnit();
-    timeoutScheduler.schedule(task, callTimeout.toInt(unit), unit);
+    timeoutScheduler.schedule(() -> {
+      try {
+        task.run();
+      } catch (Throwable t) {
+        LOG.error(errorMsg.get(), t);
+      }
+    }, callTimeout.toInt(unit), unit);
   }
 
   public TimeDuration getCallTimeout() {
     return callTimeout;
   }
-
-  public void close() {
-    timeoutScheduler.shutdown();
-  }
 }

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-grpc/src/main/java/org/apache/ratis/grpc/RaftGrpcUtil.java
----------------------------------------------------------------------
diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/RaftGrpcUtil.java 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/RaftGrpcUtil.java
index 185abbf..bad7a07 100644
--- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/RaftGrpcUtil.java
+++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/RaftGrpcUtil.java
@@ -25,20 +25,27 @@ import org.apache.ratis.shaded.io.grpc.stub.StreamObserver;
 import org.apache.ratis.util.*;
 
 import java.io.IOException;
-import java.util.Objects;
 import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.CompletionException;
 import java.util.function.Function;
 
 public interface RaftGrpcUtil {
   Metadata.Key<String> EXCEPTION_TYPE_KEY =
       Metadata.Key.of("exception-type", Metadata.ASCII_STRING_MARSHALLER);
+  Metadata.Key<String> CALL_ID =
+      Metadata.Key.of("call-id", Metadata.ASCII_STRING_MARSHALLER);
 
   static StatusRuntimeException wrapException(Throwable t) {
+    return wrapException(t, -1);
+  }
+
+  static StatusRuntimeException wrapException(Throwable t, long callId) {
     t = JavaUtils.unwrapCompletionException(t);
 
     Metadata trailers = new Metadata();
     trailers.put(EXCEPTION_TYPE_KEY, t.getClass().getCanonicalName());
+    if (callId > 0) {
+      trailers.put(CALL_ID, String.valueOf(callId));
+    }
     return new StatusRuntimeException(
         Status.INTERNAL.withCause(t).withDescription(t.getMessage()), 
trailers);
   }
@@ -69,6 +76,15 @@ public interface RaftGrpcUtil {
     return new IOException(se);
   }
 
+  static long getCallId(Throwable t) {
+    if (t instanceof StatusRuntimeException) {
+      final Metadata trailers = ((StatusRuntimeException)t).getTrailers();
+      String callId = trailers.get(CALL_ID);
+      return callId != null ? Integer.parseInt(callId) : -1;
+    }
+    return -1;
+  }
+
   static IOException unwrapIOException(Throwable t) {
     final IOException e;
     if (t instanceof StatusRuntimeException) {

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java
----------------------------------------------------------------------
diff --git 
a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java
 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java
index a39e19c..e54cbbd 100644
--- 
a/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java
+++ 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/client/RaftClientProtocolClient.java
@@ -139,8 +139,7 @@ public class RaftClientProtocolClient implements Closeable {
   class AsyncStreamObservers implements Closeable {
     /** Request map: callId -> future */
     private final AtomicReference<Map<Long, 
CompletableFuture<RaftClientReply>>> replies = new AtomicReference<>(new 
ConcurrentHashMap<>());
-    private final RpcTimeout
-        rpcTimeout = new RpcTimeout(timeout, true);
+    private final RpcTimeout rpcTimeout = new RpcTimeout(timeout);
     private final StreamObserver<RaftClientReplyProto> replyStreamObserver = 
new StreamObserver<RaftClientReplyProto>() {
       @Override
       public void onNext(RaftClientReplyProto proto) {
@@ -171,6 +170,10 @@ public class RaftClientProtocolClient implements Closeable 
{
     };
     private final StreamObserver<RaftClientRequestProto> requestStreamObserver 
= append(replyStreamObserver);
 
+    private AsyncStreamObservers() {
+      rpcTimeout.addUser();
+    }
+
     CompletableFuture<RaftClientReply> onNext(RaftClientRequest request) {
       final Map<Long, CompletableFuture<RaftClientReply>> map = replies.get();
       if (map == null) {
@@ -181,7 +184,7 @@ public class RaftClientProtocolClient implements Closeable {
           () -> getName() + ":" + getClass().getSimpleName());
       try {
         
requestStreamObserver.onNext(ClientProtoUtils.toRaftClientRequestProto(request));
-        rpcTimeout.onTimeout(() -> timeoutCheck(request));
+        rpcTimeout.onTimeout(() -> timeoutCheck(request), () -> "Timeout check 
failed for client request: " + request);
       } catch(Throwable t) {
         handleReplyFuture(request.getCallId(), future -> 
future.completeExceptionally(t));
       }
@@ -189,8 +192,8 @@ public class RaftClientProtocolClient implements Closeable {
     }
 
     private void timeoutCheck(RaftClientRequest request) {
-      handleReplyFuture(request.getCallId(),
-          f -> f.completeExceptionally(new IOException("Request timeout " + 
rpcTimeout.getCallTimeout() + ": " + request)));
+      handleReplyFuture(request.getCallId(), f -> f.completeExceptionally(
+          new IOException("Request timeout " + rpcTimeout.getCallTimeout() + 
": " + request)));
     }
 
     private void handleReplyFuture(long callId, 
Consumer<CompletableFuture<RaftClientReply>> handler) {
@@ -203,7 +206,7 @@ public class RaftClientProtocolClient implements Closeable {
     public void close() {
       requestStreamObserver.onCompleted();
       completeReplyExceptionally(null, "close");
-      rpcTimeout.close();
+      rpcTimeout.removeUser();
     }
 
     private void completeReplyExceptionally(Throwable t, String event) {

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/GRpcLogAppender.java
----------------------------------------------------------------------
diff --git 
a/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/GRpcLogAppender.java 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/GRpcLogAppender.java
index 19389bd..eb2db91 100644
--- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/GRpcLogAppender.java
+++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/GRpcLogAppender.java
@@ -20,6 +20,7 @@ package org.apache.ratis.grpc.server;
 import org.apache.ratis.grpc.GrpcConfigKeys;
 import org.apache.ratis.grpc.RaftGRpcService;
 import org.apache.ratis.grpc.RaftGrpcUtil;
+import org.apache.ratis.rpc.RpcTimeout;
 import org.apache.ratis.server.impl.FollowerInfo;
 import org.apache.ratis.server.impl.LeaderState;
 import org.apache.ratis.server.impl.LogAppender;
@@ -34,12 +35,15 @@ import 
org.apache.ratis.shaded.proto.RaftProtos.InstallSnapshotRequestProto;
 import org.apache.ratis.statemachine.SnapshotInfo;
 import org.apache.ratis.util.CodeInjectionForTesting;
 import org.apache.ratis.util.Preconditions;
+import org.apache.ratis.util.TimeDuration;
 
 import java.util.LinkedList;
+import java.util.Map;
 import java.util.Objects;
 import java.util.Queue;
 import java.util.UUID;
-import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 
 /**
@@ -47,12 +51,15 @@ import java.util.concurrent.atomic.AtomicBoolean;
  */
 public class GRpcLogAppender extends LogAppender {
   private final RaftServerProtocolClient client;
-  private final Queue<AppendEntriesRequestProto> pendingRequests;
+  private final Map<Long, AppendEntriesRequestProto> pendingRequests;
   private final int maxPendingRequestsNum;
+  private long callId = 0;
   private volatile boolean firstResponseReceived = false;
 
   private final AppendLogResponseHandler appendResponseHandler;
   private final InstallSnapshotResponseHandler snapshotResponseHandler;
+  private static RpcTimeout rpcTimeout = new RpcTimeout(
+      TimeDuration.valueOf(3, TimeUnit.SECONDS));
 
   private volatile StreamObserver<AppendEntriesRequestProto> 
appendLogRequestObserver;
   private StreamObserver<InstallSnapshotRequestProto> snapshotRequestObserver;
@@ -65,10 +72,11 @@ public class GRpcLogAppender extends LogAppender {
     client = rpcService.getRpcClient(f.getPeer());
     maxPendingRequestsNum = GrpcConfigKeys.Server.leaderOutstandingAppendsMax(
         server.getProxy().getProperties());
-    pendingRequests = new ConcurrentLinkedQueue<>();
+    pendingRequests = new ConcurrentHashMap<>();
 
     appendResponseHandler = new AppendLogResponseHandler();
     snapshotResponseHandler = new InstallSnapshotResponseHandler();
+    rpcTimeout.addUser();
   }
 
   @Override
@@ -134,9 +142,9 @@ public class GRpcLogAppender extends LogAppender {
         // prepare and enqueue the append request. note changes on follower's
         // nextIndex and ops on pendingRequests should always be associated
         // together and protected by the lock
-        pending = createRequest();
+        pending = createRequest(callId++);
         if (pending != null) {
-          Preconditions.assertTrue(pendingRequests.offer(pending));
+          pendingRequests.put(pending.getServerRequest().getCallId(), pending);
           updateNextIndex(pending);
         }
       }
@@ -154,9 +162,18 @@ public class GRpcLogAppender extends LogAppender {
         server.getId(), null, request);
 
     s.onNext(request);
+    rpcTimeout.onTimeout(() -> timeoutAppendRequest(request),
+        () -> "Timeout check failed for append entry request: " + request);
     follower.updateLastRpcSendTime();
   }
 
+  private void timeoutAppendRequest(AppendEntriesRequestProto request) {
+    AppendEntriesRequestProto pendingRequest = 
pendingRequests.remove(request.getServerRequest().getCallId());
+    if (pendingRequest != null) {
+      LOG.info("Timeout executed for append entry request: " + pendingRequest);
+    }
+  }
+
   private void updateNextIndex(AppendEntriesRequestProto request) {
     final int count = request.getEntriesCount();
     if (count > 0) {
@@ -227,6 +244,7 @@ public class GRpcLogAppender extends LogAppender {
             server.getId(), follower.getPeer().getId(), 
RaftGrpcUtil.unwrapThrowable(t));
       }
 
+      long callId = RaftGrpcUtil.getCallId(t);
       synchronized (this) {
         final Status cause = Status.fromThrowable(t);
         if (cause != null && cause.getCode() == Status.Code.INTERNAL) {
@@ -240,7 +258,7 @@ public class GRpcLogAppender extends LogAppender {
         }
 
         // clear the pending requests queue and reset the next index of 
follower
-        AppendEntriesRequestProto request = pendingRequests.peek();
+        AppendEntriesRequestProto request = pendingRequests.get(callId);
         if (request != null) {
           final long nextIndex = request.hasPreviousLog() ?
               request.getPreviousLog().getIndex() + 1 : 
raftLog.getStartIndex();
@@ -262,7 +280,12 @@ public class GRpcLogAppender extends LogAppender {
   }
 
   private void onSuccess(AppendEntriesReplyProto reply) {
-    AppendEntriesRequestProto request = pendingRequests.poll();
+    AppendEntriesRequestProto request = 
pendingRequests.remove(reply.getServerReply().getCallId());
+    if (request == null) {
+      // If reply comes after timeout, the reply is ignored.
+      LOG.warn("Ignoring reply: " + reply);
+      return;
+    }
     updateCommitIndex(request.getLeaderCommit());
 
     final long replyNextIndex = reply.getNextIndex();
@@ -293,13 +316,24 @@ public class GRpcLogAppender extends LogAppender {
   }
 
   private synchronized void onInconsistency(AppendEntriesReplyProto reply) {
-    AppendEntriesRequestProto request = pendingRequests.peek();
+    AppendEntriesRequestProto request = 
pendingRequests.remove(reply.getServerReply().getCallId());
+    if (request == null) {
+      // If reply comes after timeout, the reply is ignored.
+      LOG.warn("Ignoring reply: " + reply);
+      return;
+    }
     Preconditions.assertTrue(request.hasPreviousLog());
     if (request.getPreviousLog().getIndex() >= reply.getNextIndex()) {
       clearPendingRequests(reply.getNextIndex());
     }
   }
 
+  @Override
+  public LogAppender stopSender() {
+    rpcTimeout.removeUser();
+    return super.stopSender();
+  }
+
   private class InstallSnapshotResponseHandler
       implements StreamObserver<InstallSnapshotReplyProto> {
     private final Queue<Integer> pending;

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/RaftServerProtocolService.java
----------------------------------------------------------------------
diff --git 
a/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/RaftServerProtocolService.java
 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/RaftServerProtocolService.java
index 50e2343..e61b64b 100644
--- 
a/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/RaftServerProtocolService.java
+++ 
b/ratis-grpc/src/main/java/org/apache/ratis/grpc/server/RaftServerProtocolService.java
@@ -87,7 +87,7 @@ public class RaftServerProtocolService extends 
RaftServerProtocolServiceImplBase
             LOG.debug("{} got exception when appendEntries {}: {}",
                 getId(), ProtoUtils.toString(request.getServerRequest()), e);
           }
-          responseObserver.onError(RaftGrpcUtil.wrapException(e));
+          responseObserver.onError(RaftGrpcUtil.wrapException(e, 
request.getServerRequest().getCallId()));
           current.completeExceptionally(e);
         }
       }

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/main/java/org/apache/ratis/server/impl/LeaderState.java
----------------------------------------------------------------------
diff --git 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/LeaderState.java 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/LeaderState.java
index 840df08..ca27b7e 100644
--- a/ratis-server/src/main/java/org/apache/ratis/server/impl/LeaderState.java
+++ b/ratis-server/src/main/java/org/apache/ratis/server/impl/LeaderState.java
@@ -266,11 +266,12 @@ public class LeaderState {
         .forEach(protos::add);
   }
 
-  AppendEntriesRequestProto newAppendEntriesRequestProto(
-      RaftPeerId targetId, TermIndex previous, List<LogEntryProto> entries, 
boolean initializing) {
+  AppendEntriesRequestProto newAppendEntriesRequestProto(RaftPeerId targetId,
+      TermIndex previous, List<LogEntryProto> entries, boolean initializing,
+      long callId) {
     return ServerProtoUtils.toAppendEntriesRequestProto(server.getId(), 
targetId,
         server.getGroupId(), currentTerm, entries, 
raftLog.getLastCommittedIndex(),
-        initializing, previous, server.getCommitInfos());
+        initializing, previous, server.getCommitInfos(), callId);
   }
 
   /**

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/main/java/org/apache/ratis/server/impl/LogAppender.java
----------------------------------------------------------------------
diff --git 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/LogAppender.java 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/LogAppender.java
index ed4e843..a6adf8c 100644
--- a/ratis-server/src/main/java/org/apache/ratis/server/impl/LogAppender.java
+++ b/ratis-server/src/main/java/org/apache/ratis/server/impl/LogAppender.java
@@ -43,6 +43,7 @@ import java.io.InterruptedIOException;
 import java.nio.file.Path;
 import java.util.*;
 
+import static org.apache.ratis.server.impl.RaftServerConstants.DEFAULT_CALLID;
 import static 
org.apache.ratis.server.impl.RaftServerConstants.INVALID_LOG_INDEX;
 
 /**
@@ -137,9 +138,9 @@ public class LogAppender extends Daemon {
       return buf.isEmpty();
     }
 
-    AppendEntriesRequestProto getAppendRequest(TermIndex previous) {
+    AppendEntriesRequestProto getAppendRequest(TermIndex previous, long 
callId) {
       final AppendEntriesRequestProto request = 
leaderState.newAppendEntriesRequestProto(
-          getFollowerId(), previous, buf, !follower.isAttendingVote());
+          getFollowerId(), previous, buf, !follower.isAttendingVote(), callId);
       buf.clear();
       totalSize = 0;
       return request;
@@ -164,7 +165,7 @@ public class LogAppender extends Daemon {
     return previous;
   }
 
-  protected AppendEntriesRequestProto createRequest() throws 
RaftLogIOException {
+  protected AppendEntriesRequestProto createRequest(long callId) throws 
RaftLogIOException {
     final TermIndex previous = getPrevious();
     final long leaderNext = raftLog.getNextIndex();
     long next = follower.getNextIndex() + buffer.getPendingEntryNum();
@@ -185,7 +186,7 @@ public class LogAppender extends Daemon {
     }
 
     if (toSend || shouldHeartbeat()) {
-      return buffer.getAppendRequest(previous);
+      return buffer.getAppendRequest(previous, callId);
     }
     return null;
   }
@@ -198,7 +199,7 @@ public class LogAppender extends Daemon {
     while (isAppenderRunning()) { // keep retrying for IOException
       try {
         if (request == null || request.getEntriesCount() == 0) {
-          request = createRequest();
+          request = createRequest(DEFAULT_CALLID);
         }
 
         if (request == null) {

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java
----------------------------------------------------------------------
diff --git 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java
index 595a80a..4e4a811 100644
--- 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java
+++ 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/RaftServerImpl.java
@@ -761,8 +761,8 @@ public class RaftServerImpl implements RaftServerProtocol, 
RaftServerAsynchronou
     final TermIndex previous = r.hasPreviousLog() ?
         ServerProtoUtils.toTermIndex(r.getPreviousLog()) : null;
     return appendEntriesAsync(RaftPeerId.valueOf(request.getRequestorId()),
-        ProtoUtils.toRaftGroupId(request.getRaftGroupId()),
-        r.getLeaderTerm(), previous, r.getLeaderCommit(), r.getInitializing(),
+        ProtoUtils.toRaftGroupId(request.getRaftGroupId()), r.getLeaderTerm(),
+        previous, r.getLeaderCommit(), request.getCallId(), 
r.getInitializing(),
         r.getCommitInfosList(), entries);
   }
 
@@ -780,7 +780,7 @@ public class RaftServerImpl implements RaftServerProtocol, 
RaftServerAsynchronou
 
   private CompletableFuture<AppendEntriesReplyProto> appendEntriesAsync(
       RaftPeerId leaderId, RaftGroupId leaderGroupId, long leaderTerm,
-      TermIndex previous, long leaderCommit, boolean initializing,
+      TermIndex previous, long leaderCommit, long callId, boolean initializing,
       List<CommitInfoProto> commitInfos, LogEntryProto... entries) throws 
IOException {
     CodeInjectionForTesting.execute(APPEND_ENTRIES, getId(),
         leaderId, leaderTerm, previous, leaderCommit, initializing, entries);
@@ -809,7 +809,7 @@ public class RaftServerImpl implements RaftServerProtocol, 
RaftServerAsynchronou
       currentTerm = state.getCurrentTerm();
       if (!recognized) {
         final AppendEntriesReplyProto reply = 
ServerProtoUtils.toAppendEntriesReplyProto(
-            leaderId, getId(), groupId, currentTerm, nextIndex, NOT_LEADER);
+            leaderId, getId(), groupId, currentTerm, nextIndex, NOT_LEADER, 
callId);
         if (LOG.isDebugEnabled()) {
           LOG.debug("{}: Not recognize {} (term={}) as leader, state: {} 
reply: {}",
               getId(), leaderId, leaderTerm, state, 
ProtoUtils.toString(reply));
@@ -835,7 +835,7 @@ public class RaftServerImpl implements RaftServerProtocol, 
RaftServerAsynchronou
       if (previous != null && !containPrevious(previous)) {
         final AppendEntriesReplyProto reply =
             ServerProtoUtils.toAppendEntriesReplyProto(leaderId, getId(), 
groupId,
-                currentTerm, Math.min(nextIndex, previous.getIndex()), 
INCONSISTENCY);
+                currentTerm, Math.min(nextIndex, previous.getIndex()), 
INCONSISTENCY, callId);
         if (LOG.isDebugEnabled()) {
           LOG.debug("{}: inconsistency entries. Leader previous:{}, Reply:{}",
               getId(), previous, ServerProtoUtils.toString(reply));
@@ -855,7 +855,7 @@ public class RaftServerImpl implements RaftServerProtocol, 
RaftServerAsynchronou
       nextIndex = entries[entries.length - 1].getIndex() + 1;
     }
     final AppendEntriesReplyProto reply = 
ServerProtoUtils.toAppendEntriesReplyProto(
-        leaderId, getId(), groupId, currentTerm, nextIndex, SUCCESS);
+        leaderId, getId(), groupId, currentTerm, nextIndex, SUCCESS, callId);
     logAppendEntries(isHeartbeat,
         () -> getId() + ": succeeded to handle AppendEntries. Reply: "
             + ServerProtoUtils.toString(reply));

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/main/java/org/apache/ratis/server/impl/ServerProtoUtils.java
----------------------------------------------------------------------
diff --git 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/ServerProtoUtils.java 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/ServerProtoUtils.java
index cfdf4cf..6855473 100644
--- 
a/ratis-server/src/main/java/org/apache/ratis/server/impl/ServerProtoUtils.java
+++ 
b/ratis-server/src/main/java/org/apache/ratis/server/impl/ServerProtoUtils.java
@@ -173,10 +173,12 @@ public class ServerProtoUtils {
 
   public static AppendEntriesReplyProto toAppendEntriesReplyProto(
       RaftPeerId requestorId, RaftPeerId replyId, RaftGroupId groupId, long 
term,
-      long nextIndex, AppendResult result) {
+      long nextIndex, AppendResult result, long callId) {
+    RaftRpcReplyProto.Builder rpcReply = toRaftRpcReplyProtoBuilder(
+        requestorId, replyId, groupId, result == AppendResult.SUCCESS)
+        .setCallId(callId);
     return AppendEntriesReplyProto.newBuilder()
-        .setServerReply(toRaftRpcReplyProtoBuilder(
-            requestorId, replyId, groupId, result == AppendResult.SUCCESS))
+        .setServerReply(rpcReply)
         .setTerm(term)
         .setNextIndex(nextIndex)
         .setResult(result).build();
@@ -185,10 +187,12 @@ public class ServerProtoUtils {
   public static AppendEntriesRequestProto toAppendEntriesRequestProto(
       RaftPeerId requestorId, RaftPeerId replyId, RaftGroupId groupId, long 
leaderTerm,
       List<LogEntryProto> entries, long leaderCommit, boolean initializing,
-      TermIndex previous, Collection<CommitInfoProto> commitInfos) {
+      TermIndex previous, Collection<CommitInfoProto> commitInfos, long 
callId) {
+    RaftRpcRequestProto.Builder rpcRequest = 
toRaftRpcRequestProtoBuilder(requestorId, replyId, groupId)
+        .setCallId(callId);
     final AppendEntriesRequestProto.Builder b = AppendEntriesRequestProto
         .newBuilder()
-        .setServerRequest(toRaftRpcRequestProtoBuilder(requestorId, replyId, 
groupId))
+        .setServerRequest(rpcRequest)
         .setLeaderTerm(leaderTerm)
         .setLeaderCommit(leaderCommit)
         .setInitializing(initializing);

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/test/java/org/apache/ratis/RaftAsyncTests.java
----------------------------------------------------------------------
diff --git a/ratis-server/src/test/java/org/apache/ratis/RaftAsyncTests.java 
b/ratis-server/src/test/java/org/apache/ratis/RaftAsyncTests.java
index af61aca..3c68469 100644
--- a/ratis-server/src/test/java/org/apache/ratis/RaftAsyncTests.java
+++ b/ratis-server/src/test/java/org/apache/ratis/RaftAsyncTests.java
@@ -256,4 +256,45 @@ public abstract class RaftAsyncTests<CLUSTER extends 
MiniRaftCluster> extends Ba
     RaftBasicTests.testRequestTimeout(true, cluster, LOG, properties);
     cluster.shutdown();
   }
+
+  @Test
+  public void testAppendEntriesTimeout()
+      throws IOException, InterruptedException, ExecutionException {
+    LOG.info("Running testAppendEntriesTimeout");
+    TimeDuration retryCacheExpiryDuration = TimeDuration.valueOf(20, 
TimeUnit.SECONDS);
+    RaftServerConfigKeys.RetryCache.setExpiryTime(properties, 
retryCacheExpiryDuration);
+    final CLUSTER cluster = getFactory().newCluster(NUM_SERVERS, properties);
+    cluster.start();
+    waitForLeader(cluster);
+    long time = System.currentTimeMillis();
+    long waitTime = 5000;
+    try (final RaftClient client = cluster.createClient()) {
+      // block append requests
+      cluster.getServerAliveStream().forEach(raftServer -> {
+        try {
+          if (!raftServer.isLeader()) {
+            ((SimpleStateMachine4Testing) 
raftServer.getStateMachine()).setBlockAppend(true);
+          }
+        } catch (InterruptedException e) {
+          LOG.error("Interrupted while blocking append", e);
+        }
+      });
+      CompletableFuture<RaftClientReply> replyFuture = client.sendAsync(new 
RaftTestUtil.SimpleMessage("abc"));
+      Thread.sleep(waitTime);
+      // replyFuture should not be completed until append request is unblocked.
+      Assert.assertTrue(!replyFuture.isDone());
+      // unblock append request.
+      cluster.getServerAliveStream().forEach(raftServer -> {
+        try {
+          ((SimpleStateMachine4Testing) 
raftServer.getStateMachine()).setBlockAppend(false);
+        } catch (InterruptedException e) {
+          LOG.error("Interrupted while unblocking append", e);
+        }
+      });
+      client.send(new RaftTestUtil.SimpleMessage("abc"));
+      replyFuture.get();
+      Assert.assertTrue(System.currentTimeMillis() - time > waitTime);
+    }
+    cluster.shutdown();
+  }
 }

http://git-wip-us.apache.org/repos/asf/incubator-ratis/blob/00f64b4c/ratis-server/src/test/java/org/apache/ratis/statemachine/SimpleStateMachine4Testing.java
----------------------------------------------------------------------
diff --git 
a/ratis-server/src/test/java/org/apache/ratis/statemachine/SimpleStateMachine4Testing.java
 
b/ratis-server/src/test/java/org/apache/ratis/statemachine/SimpleStateMachine4Testing.java
index 4bf75e1..91643db 100644
--- 
a/ratis-server/src/test/java/org/apache/ratis/statemachine/SimpleStateMachine4Testing.java
+++ 
b/ratis-server/src/test/java/org/apache/ratis/statemachine/SimpleStateMachine4Testing.java
@@ -78,7 +78,8 @@ public class SimpleStateMachine4Testing extends 
BaseStateMachine {
       RaftServerConfigKeys.Log.writeBufferSize(properties).getSizeInt();
 
   private volatile boolean running = true;
-  private boolean blockTransaction = false;
+  private volatile boolean blockTransaction = false;
+  private volatile boolean blockAppend = false;
   private final Semaphore blockingSemaphore = new Semaphore(1);
   private long endIndexLastCkpt = RaftServerConstants.INVALID_LOG_INDEX;
 
@@ -252,7 +253,23 @@ public class SimpleStateMachine4Testing extends 
BaseStateMachine {
     }
     return new TransactionContextImpl(this, request, 
SMLogEntryProto.newBuilder()
         .setData(request.getMessage().getContent())
-        .build());
+        .setStateMachineData(ByteString.copyFromUtf8("StateMachine 
Data")).build());
+  }
+
+  @Override
+  public CompletableFuture<?> writeStateMachineData(LogEntryProto entry) {
+    CompletableFuture f = new CompletableFuture();
+    if (blockAppend) {
+      try {
+        blockingSemaphore.acquire();
+        blockingSemaphore.release();
+      } catch (InterruptedException e) {
+        LOG.error("Could not block writeStateMachineData", e);
+        Thread.currentThread().interrupt();
+      }
+    }
+    f.complete(null);
+    return f;
   }
 
   @Override
@@ -275,4 +292,13 @@ public class SimpleStateMachine4Testing extends 
BaseStateMachine {
       blockingSemaphore.release();
     }
   }
+
+  public void setBlockAppend(boolean blockAppendVal) throws 
InterruptedException {
+    this.blockAppend = blockAppendVal;
+    if (blockAppendVal) {
+      blockingSemaphore.acquire();
+    } else {
+      blockingSemaphore.release();
+    }
+  }
 }

Reply via email to