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