diff --git a/ratis-grpc/dev-support/findbugsExcludeFile.xml b/ratis-grpc/dev-support/findbugsExcludeFile.xml index 3f10ecb4e6..d10bf0fa2b 100644 --- a/ratis-grpc/dev-support/findbugsExcludeFile.xml +++ b/ratis-grpc/dev-support/findbugsExcludeFile.xml @@ -15,6 +15,12 @@ limitations under the License. --> + + + + + + diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcConfigKeys.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcConfigKeys.java index 9ef7ee498e..4c8ade45dd 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcConfigKeys.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcConfigKeys.java @@ -304,6 +304,18 @@ static void setCredentials(Parameters parameters, ServerCredentials credentials) parameters.put(CREDENTIALS_PARAMETER, credentials, CREDENTIALS_CLASS); } + String LOG_APPENDER_LISTENER_FACTORY_PARAMETER = PREFIX + ".log.appender.listener.factory"; + Class LOG_APPENDER_LISTENER_FACTORY_CLASS = GrpcLogAppenderListener.Factory.class; + static GrpcLogAppenderListener.Factory logAppenderListenerFactory(Parameters parameters) { + return parameters == null ? null : parameters.get( + LOG_APPENDER_LISTENER_FACTORY_PARAMETER, LOG_APPENDER_LISTENER_FACTORY_CLASS); + } + + /** Sets an optional factory for observing the lifecycle of each peer log appender. */ + static void setLogAppenderListenerFactory(Parameters parameters, GrpcLogAppenderListener.Factory factory) { + parameters.put(LOG_APPENDER_LISTENER_FACTORY_PARAMETER, factory, LOG_APPENDER_LISTENER_FACTORY_CLASS); + } + String TLS_CONF_PARAMETER = PREFIX + ".tls.conf"; Class TLS_CONF_CLASS = TLS.CONF_CLASS; static GrpcTlsConfig tlsConf(Parameters parameters) { diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcFactory.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcFactory.java index 6168d14c95..27fe56b8d7 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcFactory.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcFactory.java @@ -93,6 +93,7 @@ private SslContexts(GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, private final GrpcServices.Customizer servicesCustomizer; private final ServerCredentials serverCredentials; + private final GrpcLogAppenderListener.Factory logAppenderListenerFactory; private final Supplier forServerSupplier; private final Supplier forClientSupplier; @@ -100,6 +101,7 @@ private SslContexts(GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, public GrpcFactory(Parameters parameters) { this(GrpcConfigKeys.Server.servicesCustomizer(parameters), GrpcConfigKeys.Server.credentials(parameters), + GrpcConfigKeys.Server.logAppenderListenerFactory(parameters), GrpcConfigKeys.TLS.conf(parameters), GrpcConfigKeys.Admin.tlsConf(parameters), GrpcConfigKeys.Client.tlsConf(parameters), @@ -109,10 +111,12 @@ public GrpcFactory(Parameters parameters) { private GrpcFactory(GrpcServices.Customizer servicesCustomizer, ServerCredentials serverCredentials, + GrpcLogAppenderListener.Factory logAppenderListenerFactory, GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, GrpcTlsConfig clientTlsConfig, GrpcTlsConfig serverTlsConfig) { this.servicesCustomizer = servicesCustomizer; this.serverCredentials = serverCredentials; + this.logAppenderListenerFactory = logAppenderListenerFactory; this.forServerSupplier = MemoizedSupplier.valueOf(() -> new SslContexts( tlsConfig, adminTlsConfig, clientTlsConfig, serverTlsConfig, BUILD_SSL_CONTEXT_FOR_SERVER)); @@ -127,7 +131,15 @@ public SupportedRpcType getRpcType() { @Override public LogAppender newLogAppender(RaftServer.Division server, LeaderState state, FollowerInfo f) { - return new GrpcLogAppender(server, state, f); + GrpcLogAppenderListener listener = null; + if (logAppenderListenerFactory != null) { + try { + listener = logAppenderListenerFactory.create(server.getMemberId(), f.getPeer()); + } catch (Throwable t) { + LOG.warn("{}->{}: Failed to create gRPC log appender listener", server.getMemberId(), f.getId(), t); + } + } + return new GrpcLogAppender(server, state, f, listener); } @Override diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java new file mode 100644 index 0000000000..b497f9efd4 --- /dev/null +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java @@ -0,0 +1,83 @@ +/* + * 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.grpc; + +import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto; +import org.apache.ratis.proto.RaftProtos.AppendEntriesRequestProto; +import org.apache.ratis.protocol.RaftGroupMemberId; +import org.apache.ratis.protocol.RaftPeer; + +/** + * Observes a single peer log appender. Callbacks run on Ratis threads, may be concurrent and may + * hold an appender lock. They must not block or retain request payloads. Exceptions are isolated + * from replication. Callbacks describe lifecycle activity, not application-level outcomes: + * consumers are responsible for correlating requests, handling racing terminal notifications, + * and filtering messages. This interface currently observes AppendEntries, not InstallSnapshot. + * Append request registration, client reset and reply inconsistency callbacks are serialized per + * appender. Replies and other terminal callbacks may race with these callbacks. Consumers must + * handle these races; no exactly-once terminal notification is guaranteed. + */ +public interface GrpcLogAppenderListener { + /** Creates a separate listener for each appender, including after leadership changes. */ + @FunctionalInterface + interface Factory { + /** @return the listener, or null to disable observation for this appender. */ + GrpcLogAppenderListener create(RaftGroupMemberId source, RaftPeer destination); + } + + /** @return the AppendEntries listener, or null to disable its callbacks. Called once per appender. */ + default AppendEntries appendEntries() { + return null; + } + + /** Observes append attempts and their response streams, including the separate heartbeat stream. */ + interface AppendEntries { + /** An append attempt is registered, before establishing or writing its stream. */ + default void onRequest(AppendEntriesRequestProto request) { } + + /** A response was received, possibly after its request timed out or was invalidated. */ + default void onReply(AppendEntriesReplyProto reply) { } + + /** + * An INCONSISTENCY reply is being handled and pending append requests are about to be cleared. + * Called after onReply for that reply, with the appender write lock held, without resetting the client. + */ + default void onReplyInconsistency() { } + + /** A local send error occurred; a later stream notification may follow. */ + default void onFailure(long callId, Throwable error) { } + + /** A pending request timed out. */ + default void onTimeout(long callId) { } + + /** A response stream completed, including after the appender stopped. */ + default void onCompleted() { } + + /** A response stream failed, including after the appender stopped. */ + default void onError(Throwable error) { } + } + + /** + * The client is about to be reset, which may invalidate pending append attempts. + * The reason is diagnostic text, not a stable identifier; error may be null. + */ + default void onResetClient(String reason, Throwable error) { } + + /** The appender run loop exited, normally or exceptionally. Pending attempts may remain. */ + default void onNotRunning() { } +} 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 29867b544e..e6870b9b89 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 @@ -19,6 +19,7 @@ import org.apache.ratis.conf.RaftProperties; import org.apache.ratis.grpc.GrpcConfigKeys; +import org.apache.ratis.grpc.GrpcLogAppenderListener; import org.apache.ratis.grpc.GrpcUtil; import org.apache.ratis.grpc.metrics.GrpcServerMetrics; import org.apache.ratis.metrics.Timekeeper; @@ -62,6 +63,7 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Consumer; import static org.apache.ratis.server.raftlog.LogProtoUtils.toLogEntryTermIndexString; @@ -163,14 +165,19 @@ synchronized int process(Event event) { private final boolean useSeparateHBChannel; private final GrpcServerMetrics grpcServerMetrics; + private final GrpcLogAppenderListener listener; + private final GrpcLogAppenderListener.AppendEntries appendEntriesListener; private final AutoCloseableReadWriteLock lock; private final StackTraceElement caller; private final RetryPolicy errorRetryWaitPolicy; private final ReplyState replyState = new ReplyState(); - public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f) { + public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f, + GrpcLogAppenderListener listener) { super(server, leaderState, f); + this.listener = listener; + this.appendEntriesListener = getAppendEntriesListener(listener); Objects.requireNonNull(getServerRpc(), "getServerRpc() == null"); @@ -194,6 +201,35 @@ public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, Foll RaftServerConfigKeys.Log.Appender.RETRY_POLICY_KEY); } + private void notifyListener(Consumer notification) { + if (listener != null) { + try { + notification.accept(listener); + } catch (Throwable t) { + LOG.warn("{}: gRPC log appender listener threw an exception", this, t); + } + } + } + + private GrpcLogAppenderListener.AppendEntries getAppendEntriesListener(GrpcLogAppenderListener logListener) { + try { + return logListener == null ? null : logListener.appendEntries(); + } catch (Throwable t) { + LOG.warn("{}: Failed to get AppendEntries listener", this, t); + return null; + } + } + + private void notifyAppendEntriesListener(Consumer notification) { + if (appendEntriesListener != null) { + try { + notification.accept(appendEntriesListener); + } catch (Throwable t) { + LOG.warn("{}: AppendEntries listener threw an exception", this, t); + } + } + } + @Override public GrpcServicesImpl getServerRpc() { return (GrpcServicesImpl)super.getServerRpc(); @@ -203,8 +239,9 @@ private GrpcServerProtocolClient getClient() throws IOException { return getServerRpc().getProxies().getProxy(getFollowerId()); } - private void resetClient(AppendEntriesRequest request, Event event) { + private void resetClient(AppendEntriesRequest request, Event event, Throwable error) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + notifyListener(l -> l.onResetClient("resetClient: " + event, error)); getClient().resetConnectBackoff(); if (appendLogRequestObserver != null) { appendLogRequestObserver.stop(); @@ -258,16 +295,21 @@ private boolean installSnapshot() { @Override public void run() throws IOException { - for(; isRunning(); mayWait()) { - //HB period is expired OR we have messages OR follower is behind with commit index - if (shouldSendAppendEntries() || isFollowerCommitBehindLastCommitIndex()) { - final boolean installingSnapshot = installSnapshot(); - appendLog(installingSnapshot || haveTooManyPendingRequests()); + try { + for(; isRunning(); mayWait()) { + //HB period is expired OR we have messages OR follower is behind with commit index + if (shouldSendAppendEntries() || isFollowerCommitBehindLastCommitIndex()) { + final boolean installingSnapshot = installSnapshot(); + appendLog(installingSnapshot || haveTooManyPendingRequests()); + } + getLeaderState().checkHealth(getFollower()); + } + Optional.ofNullable(appendLogRequestObserver).ifPresent(StreamObservers::onCompleted); + } finally { + try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + notifyListener(GrpcLogAppenderListener::onNotRunning); } - getLeaderState().checkHealth(getFollower()); } - - Optional.ofNullable(appendLogRequestObserver).ifPresent(StreamObservers::onCompleted); } public long getWaitTimeMs() { @@ -389,31 +431,44 @@ public Comparator getCallIdComparator() { return CALL_ID_COMPARATOR; } - private void appendLog(boolean heartbeat) throws IOException { + void appendLog(boolean heartbeat) throws IOException { final AppendEntriesRequestProto pending; - final AppendEntriesRequest request; - try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - // Prepare and send the append request. - // Note changes on follower's nextIndex and ops on pendingRequests should always be done under the write-lock - pending = newAppendEntriesRequest(callId.getAndIncrement(), heartbeat); - if (pending == null) { - return; - } - request = new AppendEntriesRequest(pending, getFollowerId(), grpcServerMetrics); - pendingRequests.put(request); - increaseNextIndex(pending); - if (appendLogRequestObserver == null) { - appendLogRequestObserver = new StreamObservers( - getClient(), new AppendLogResponseHandler(), useSeparateHBChannel, getWaitTimeMin()); + AppendEntriesRequest request = null; + try { + try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + // Prepare and send the append request. + // Note changes on follower's nextIndex and ops on pendingRequests should always be done under the write-lock + pending = newAppendEntriesRequest(callId.getAndIncrement(), heartbeat); + if (pending == null) { + return; + } + request = new AppendEntriesRequest(pending, getFollowerId(), grpcServerMetrics); + pendingRequests.put(request); + notifyAppendEntriesListener(l -> l.onRequest(pending)); + increaseNextIndex(pending); + if (appendLogRequestObserver == null) { + appendLogRequestObserver = new StreamObservers( + getClient(), new AppendLogResponseHandler(), useSeparateHBChannel, getWaitTimeMin()); + } } - } - final TimeDuration remaining = getRemainingWaitTime(); - if (remaining.isPositive()) { - sleep(remaining, heartbeat); - } - if (isRunning()) { - sendRequest(request, pending); + final TimeDuration remaining = getRemainingWaitTime(); + if (remaining.isPositive()) { + sleep(remaining, heartbeat); + } + if (isRunning()) { + sendRequest(request, pending); + } + } catch (IOException | RuntimeException e) { + if (request != null) { + final long cid = request.getCallId(); + final AppendEntriesRequest failed = pendingRequests.remove(cid, request.isHeartbeat()); + if (failed != null) { + failed.stopRequestTimer(); + } + notifyAppendEntriesListener(l -> l.onFailure(cid, e)); + } + throw e; } } @@ -445,7 +500,7 @@ private void sendRequest(AppendEntriesRequest request, } } - private void timeoutAppendRequest(long cid, boolean heartbeat) { + void timeoutAppendRequest(long cid, boolean heartbeat) { final AppendEntriesRequest pending = pendingRequests.remove(cid, heartbeat); if (pending != null) { final int errorCount = replyState.process(Event.TIMEOUT); @@ -453,6 +508,7 @@ private void timeoutAppendRequest(long cid, boolean heartbeat) { this, heartbeat ? "HEARTBEAT " : "", errorCount, pending); grpcServerMetrics.onRequestTimeout(getFollowerId().toString(), heartbeat); pending.stopRequestTimer(); + notifyAppendEntriesListener(l -> l.onTimeout(cid)); } } @@ -472,7 +528,7 @@ private void increaseNextIndex(final long installedSnapshotIndex, Object reason) /** * StreamObserver for handling responses from the follower */ - private class AppendLogResponseHandler implements StreamObserver { + class AppendLogResponseHandler implements StreamObserver { private final String name = getFollower().getName() + "-" + JavaUtils.getClassSimpleName(getClass()); /** @@ -485,7 +541,8 @@ private class AppendLogResponseHandler implements StreamObserver l.onReply(reply)); if (request != null) { request.stopRequestTimer(); // Update completion time getFollower().updateLastRespondedAppendEntriesSendTime(request.getSendTime()); @@ -545,6 +602,7 @@ private void onNextImpl(AppendEntriesRequest request, AppendEntriesReplyProto re */ @Override public void onError(Throwable t) { + notifyAppendEntriesListener(l -> l.onError(t)); if (!isRunning()) { LOG.info("{} is already stopped", GrpcLogAppender.this); return; @@ -554,13 +612,17 @@ public void onError(Throwable t) { logMessageBatchDuration, t instanceof StatusRuntimeException); grpcServerMetrics.onRequestRetry(); // Update try counter AppendEntriesRequest request = pendingRequests.remove(GrpcUtil.getCallId(t), GrpcUtil.isHeartbeat(t)); - resetClient(request, Event.ERROR); + resetClient(request, Event.ERROR, t); } @Override public void onCompleted() { + notifyAppendEntriesListener(GrpcLogAppenderListener.AppendEntries::onCompleted); LOG.info("{}: follower responses appendEntries COMPLETED", this); - resetClient(null, Event.COMPLETE); + if (!isRunning()) { + return; + } + resetClient(null, Event.COMPLETE, null); } @Override @@ -571,6 +633,7 @@ public String toString() { private void updateNextIndex(long replyNextIndex) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + notifyAppendEntriesListener(GrpcLogAppenderListener.AppendEntries::onReplyInconsistency); pendingRequests.clear(); getFollower().setNextIndex(replyNextIndex); } @@ -738,7 +801,7 @@ public void onError(Throwable t) { } GrpcUtil.warn(LOG, () -> this + ": Failed InstallSnapshot", t); grpcServerMetrics.onRequestRetry(); // Update try counter - resetClient(null, Event.ERROR); + resetClient(null, Event.ERROR, t); close(); } @@ -880,7 +943,10 @@ void startRequestTimer() { } void stopRequestTimer() { - timerContext.stop(); + final Timekeeper.Context context = timerContext; + if (context != null) { + context.stop(); + } } boolean isHeartbeat() { diff --git a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java new file mode 100644 index 0000000000..d837882ccd --- /dev/null +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java @@ -0,0 +1,124 @@ +/* + * 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.grpc; + +import org.apache.ratis.BaseTest; +import org.apache.ratis.RaftTestUtil; +import org.apache.ratis.client.RaftClient; +import org.apache.ratis.conf.Parameters; +import org.apache.ratis.conf.RaftProperties; +import org.apache.ratis.grpc.server.GrpcServicesImpl; +import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto; +import org.apache.ratis.proto.RaftProtos.AppendEntriesRequestProto; +import org.apache.ratis.proto.RaftProtos.ReplicationLevel; +import org.apache.ratis.protocol.RaftClientReply; +import org.apache.ratis.protocol.RaftPeerId; +import org.apache.ratis.server.RaftServer; +import org.apache.ratis.server.impl.MiniRaftCluster; +import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.statemachine.impl.SimpleStateMachine4Testing; +import org.apache.ratis.util.CodeInjectionForTesting; +import org.apache.ratis.util.JavaUtils; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; + +public class TestGrpcLogAppenderListener extends BaseTest { + private static RaftProperties newProperties() { + final RaftProperties properties = new RaftProperties(); + properties.setClass(MiniRaftCluster.STATEMACHINE_CLASS_KEY, + SimpleStateMachine4Testing.class, StateMachine.class); + return properties; + } + + @Test + @Timeout(value = 60, unit = TimeUnit.SECONDS) + public void testAppendEntries() throws Exception { + final Parameters parameters = new Parameters(); + final Set destinations = ConcurrentHashMap.newKeySet(); + final ConcurrentLinkedQueue failures = new ConcurrentLinkedQueue<>(); + final AtomicBoolean injectFailure = new AtomicBoolean(true); + GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (source, destination) -> + new GrpcLogAppenderListener() { + @Override + public AppendEntries appendEntries() { + return new AppendEntries() { + private final Set pending = ConcurrentHashMap.newKeySet(); + + @Override + public void onRequest(AppendEntriesRequestProto request) { + if (request.getEntriesList().stream().anyMatch(entry -> entry.hasStateMachineLogEntry())) { + pending.add(request.getServerRequest().getCallId()); + } + } + + @Override + public void onReply(AppendEntriesReplyProto reply) { + if (pending.remove(reply.getServerReply().getCallId())) { + destinations.add(destination.getId()); + } + } + + @Override + public void onFailure(long callId, Throwable error) { + if (pending.remove(callId)) { + failures.add(error); + throw new IllegalStateException("Injected listener failure"); + } + } + }; + } + }); + try (MiniRaftClusterWithGrpc cluster = new MiniRaftClusterWithGrpc( + MiniRaftCluster.generateIds(3, 10), new String[0], newProperties(), parameters)) { + cluster.start(); + final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); + CodeInjectionForTesting.put(GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST, (local, remote, args) -> { + if (leader.getId().equals(local) && args[0] instanceof AppendEntriesRequestProto) { + final AppendEntriesRequestProto request = (AppendEntriesRequestProto) args[0]; + if (request.getEntriesList().stream().anyMatch(entry -> entry.hasStateMachineLogEntry()) + && injectFailure.compareAndSet(true, false)) { + throw new IllegalStateException("Injected send failure"); + } + } + return false; + }); + try (RaftClient client = cluster.createClient(leader.getId())) { + final RaftClientReply write = client.io().send(new RaftTestUtil.SimpleMessage("listener")); + Assertions.assertTrue(write.isSuccess()); + Assertions.assertTrue(client.io().watch(write.getLogIndex(), ReplicationLevel.ALL_COMMITTED).isSuccess()); + } finally { + CodeInjectionForTesting.remove(GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST); + } + final Set followers = cluster.getFollowers().stream() + .map(RaftServer.Division::getId).collect(Collectors.toSet()); + JavaUtils.attempt(() -> Assertions.assertEquals(followers, destinations), + 10, HUNDRED_MILLIS, "append replies", LOG); + Assertions.assertEquals(1, failures.size()); + Assertions.assertInstanceOf(IllegalStateException.class, failures.element()); + } + } + +} diff --git a/ratis-test/src/test/java/org/apache/ratis/grpc/server/TestGrpcLogAppenderCallbacks.java b/ratis-test/src/test/java/org/apache/ratis/grpc/server/TestGrpcLogAppenderCallbacks.java new file mode 100644 index 0000000000..f5b959e5a5 --- /dev/null +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/server/TestGrpcLogAppenderCallbacks.java @@ -0,0 +1,280 @@ +/* + * 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.grpc.server; + +import org.apache.ratis.RaftTestUtil; +import org.apache.ratis.conf.RaftProperties; +import org.apache.ratis.conf.Parameters; +import org.apache.ratis.grpc.GrpcConfigKeys; +import org.apache.ratis.grpc.GrpcFactory; +import org.apache.ratis.grpc.GrpcLogAppenderListener; +import org.apache.ratis.grpc.metrics.GrpcServerMetrics; +import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto; +import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto.AppendResult; +import org.apache.ratis.proto.RaftProtos.AppendEntriesRequestProto; +import org.apache.ratis.proto.RaftProtos.LogEntryProto; +import org.apache.ratis.proto.RaftProtos.RaftRpcReplyProto; +import org.apache.ratis.proto.RaftProtos.RaftRpcRequestProto; +import org.apache.ratis.proto.RaftProtos.StateMachineLogEntryProto; +import org.apache.ratis.protocol.RaftGroupId; +import org.apache.ratis.protocol.RaftGroupMemberId; +import org.apache.ratis.protocol.RaftPeerId; +import org.apache.ratis.server.RaftServer; +import org.apache.ratis.server.leader.FollowerInfo; +import org.apache.ratis.server.leader.LeaderState; +import org.apache.ratis.server.leader.LogAppender; +import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.AutoCloseableLock; +import org.apache.ratis.util.AutoCloseableReadWriteLock; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.ValueSource; +import org.mockito.InOrder; + +import java.io.IOException; +import java.io.InterruptedIOException; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.clearInvocations; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +@Timeout(value = 10, unit = TimeUnit.SECONDS) +public class TestGrpcLogAppenderCallbacks { + private final RaftPeerId source = RaftPeerId.valueOf("source"); + private final RaftPeerId destination = RaftPeerId.valueOf("destination"); + private final GrpcLogAppenderListener listener = mock(GrpcLogAppenderListener.class); + private final GrpcLogAppenderListener.AppendEntries appendEntries = mock(GrpcLogAppenderListener.AppendEntries.class); + private final RaftServer.Division server = mock(RaftServer.Division.class, RETURNS_DEEP_STUBS); + private final FollowerInfo follower = mock(FollowerInfo.class, RETURNS_DEEP_STUBS); + private final LeaderState leaderState = mock(LeaderState.class); + private final GrpcServicesImpl serverRpc = mock(GrpcServicesImpl.class, RETURNS_DEEP_STUBS); + private final GrpcServerProtocolClient client = mock(GrpcServerProtocolClient.class); + private GrpcLogAppender appender; + private GrpcLogAppender.RequestMap pending; + private GrpcServerMetrics metrics; + private StreamObserver responses; + + @BeforeEach + public void setup() throws Exception { + when(server.getRaftServer().getProperties()).thenReturn(new RaftProperties()); + when(server.getRaftServer().getServerRpc()).thenReturn(serverRpc); + when(server.getId()).thenReturn(source); + when(server.getMemberId()).thenReturn(RaftGroupMemberId.valueOf(source, RaftGroupId.randomId())); + when(server.getInfo().isAlive()).thenReturn(true); + when(server.getInfo().isLeader()).thenReturn(true); + when(server.getRaftLog().isOpened()).thenReturn(true); + when(follower.getId()).thenReturn(destination); + when(follower.getName()).thenReturn("test-follower"); + when(serverRpc.getProxies().getProxy(destination)).thenReturn(client); + when(listener.appendEntries()).thenReturn(appendEntries); + appender = spy(new GrpcLogAppender(server, leaderState, follower, listener)); + pending = (GrpcLogAppender.RequestMap) RaftTestUtil.getDeclaredField(appender, "pendingRequests"); + metrics = (GrpcServerMetrics) RaftTestUtil.getDeclaredField(appender, "grpcServerMetrics"); + responses = appender.new AppendLogResponseHandler(); + Assertions.assertTrue(appender.isRunning()); + } + + @AfterEach + public void cleanup() throws Exception { + if (appender != null) { + appender.stopAsync().get(5, TimeUnit.SECONDS); + } + } + + private AppendEntriesRequestProto request(long callId) { + return AppendEntriesRequestProto.newBuilder() + .setServerRequest(RaftRpcRequestProto.newBuilder().setCallId(callId)) + .addEntries(LogEntryProto.newBuilder().setIndex(callId) + .setStateMachineLogEntry(StateMachineLogEntryProto.getDefaultInstance())) + .build(); + } + + private void addPending(long callId) { + final GrpcLogAppender.AppendEntriesRequest request = + new GrpcLogAppender.AppendEntriesRequest(request(callId), destination, metrics); + request.startRequestTimer(); + pending.put(request); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testStreamErrorAfterStopOrLeaderChange(boolean stop) throws Exception { + addPending(1); + addPending(2); + if (stop) { + appender.stopAsync().get(5, TimeUnit.SECONDS); + } else { + when(server.getInfo().isLeader()).thenReturn(false); + } + clearInvocations(client, leaderState, follower, follower.getErrorState()); + final IOException error = new IOException("Stream failed after leadership ended"); + doThrow(new IllegalStateException("Listener failure")).when(appendEntries).onError(error); + Assertions.assertDoesNotThrow(() -> responses.onError(error)); + verify(appendEntries).onError(error); + verifyNoInteractions(client, leaderState, follower.getErrorState()); + verify(follower, never()).computeNextIndex(any()); + } + + @Test + public void testResetNotificationPrecedesClientFailure() throws Exception { + addPending(1); + when(serverRpc.getProxies().getProxy(destination)).thenThrow(new IOException("Closed client")); + final IOException error = new IOException("Original stream error"); + doThrow(new IllegalStateException("Listener failure")).when(listener).onResetClient(anyString(), eq(error)); + responses.onError(error); + final InOrder order = inOrder(appendEntries, listener); + order.verify(appendEntries).onError(error); + order.verify(listener).onResetClient(anyString(), eq(error)); + } + + @Test + public void testCompletionAfterStop() throws Exception { + addPending(1); + appender.stopAsync().get(5, TimeUnit.SECONDS); + clearInvocations(client, leaderState); + responses.onCompleted(); + verify(appendEntries).onCompleted(); + verify(listener, never()).onResetClient(anyString(), any()); + verifyNoInteractions(client, leaderState); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testStreamCreationFailure(boolean failProxy) throws Exception { + final AppendEntriesRequestProto request = request(1); + doReturn(request).when(appender).newAppendEntriesRequest(0, false); + final Exception error; + if (failProxy) { + error = new IOException("Failed to create peer client"); + when(serverRpc.getProxies().getProxy(destination)).thenThrow(error); + } else { + error = new IllegalStateException("Failed to create append stream"); + when(client.appendEntries(any(), eq(false))).thenThrow(error); + } + Assertions.assertSame(error, Assertions.assertThrows(Exception.class, () -> appender.appendLog(false))); + final InOrder order = inOrder(appendEntries); + order.verify(appendEntries).onRequest(request); + order.verify(appendEntries).onFailure(1, error); + Assertions.assertFalse(appender.hasPendingDataRequests()); + } + + @Test + public void testTimeout() { + addPending(1); + appender.timeoutAppendRequest(1, false); + appender.timeoutAppendRequest(1, false); + verify(appendEntries).onTimeout(1); + } + + @ParameterizedTest + @EnumSource(value = AppendResult.class, names = {"SUCCESS", "NOT_LEADER", "INCONSISTENCY"}) + public void testReplyAndInvalidation(AppendResult result) { + addPending(1); + addPending(2); + final AppendEntriesReplyProto reply = AppendEntriesReplyProto.newBuilder() + .setServerReply(RaftRpcReplyProto.newBuilder().setCallId(1)).setResult(result).build(); + responses.onNext(reply); + responses.onNext(reply); + verify(appendEntries, times(2)).onReply(reply); + if (result == AppendResult.INCONSISTENCY) { + final InOrder order = inOrder(appendEntries); + order.verify(appendEntries).onReply(reply); + order.verify(appendEntries).onReplyInconsistency(); + verify(appendEntries, times(2)).onReplyInconsistency(); + Assertions.assertFalse(appender.hasPendingDataRequests()); + } else { + verify(appendEntries, never()).onReplyInconsistency(); + } + verify(listener, never()).onResetClient(anyString(), any()); + appender.timeoutAppendRequest(1, false); + verify(appendEntries, never()).onTimeout(1); + } + + @Test + public void testReplyDoesNotAcquireAppenderWriteLock() throws Exception { + addPending(1); + final AppendEntriesReplyProto reply = AppendEntriesReplyProto.newBuilder() + .setServerReply(RaftRpcReplyProto.newBuilder().setCallId(1)).setResult(AppendResult.SUCCESS).build(); + final AutoCloseableReadWriteLock appenderLock = + (AutoCloseableReadWriteLock) RaftTestUtil.getDeclaredField(appender, "lock"); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + try (AutoCloseableLock ignored = appenderLock.writeLock(null, null)) { + CompletableFuture.runAsync(() -> responses.onNext(reply), executor).get(5, TimeUnit.SECONDS); + verify(appendEntries).onReply(reply); + Assertions.assertFalse(appender.hasPendingDataRequests()); + } finally { + executor.shutdownNow(); + Assertions.assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testListenerInitializationFailureIsIsolated(boolean failFactory) throws Exception { + when(listener.appendEntries()).thenThrow(new IllegalStateException("Injected accessor failure")); + final Parameters parameters = new Parameters(); + GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (member, peer) -> { + Assertions.assertEquals(server.getMemberId(), member); + Assertions.assertEquals(follower.getPeer(), peer); + if (failFactory) { + throw new IllegalStateException("Injected factory failure"); + } + return listener; + }); + final LogAppender created = new GrpcFactory(parameters).newLogAppender(server, leaderState, follower); + Assertions.assertNotNull(created); + created.stopAsync().get(5, TimeUnit.SECONDS); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testRunLoopExit(boolean exceptional) throws Exception { + addPending(1); + doThrow(new IllegalStateException("Listener failure")).when(listener).onNotRunning(); + if (exceptional) { + final InterruptedIOException error = new InterruptedIOException("Interrupted send"); + doThrow(error).when(appender).appendLog(true); + Assertions.assertSame(error, Assertions.assertThrows(IOException.class, appender::run)); + } else { + when(server.getInfo().isLeader()).thenReturn(false); + appender.run(); + } + verify(listener).onNotRunning(); + } +}