From d7fea1b52097386f67ae2e867b846da69a92bcbf Mon Sep 17 00:00:00 2001 From: HTHou Date: Wed, 2 Sep 2026 19:11:36 +0800 Subject: [PATCH 1/6] RATIS-2681. Add a gRPC peer data transfer listener --- .../dev-support/findbugsExcludeFile.xml | 5 + .../org/apache/ratis/grpc/GrpcConfigKeys.java | 18 ++ .../ratis/grpc/GrpcDataTransferEvent.java | 91 +++++++ .../org/apache/ratis/grpc/GrpcFactory.java | 10 +- .../ratis/grpc/server/GrpcLogAppender.java | 181 +++++++++++-- .../TestGrpcDataTransferEventListener.java | 254 ++++++++++++++++++ 6 files changed, 540 insertions(+), 19 deletions(-) create mode 100644 ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java create mode 100644 ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java diff --git a/ratis-grpc/dev-support/findbugsExcludeFile.xml b/ratis-grpc/dev-support/findbugsExcludeFile.xml index 3f10ecb4e6..ba57c27c2b 100644 --- a/ratis-grpc/dev-support/findbugsExcludeFile.xml +++ b/ratis-grpc/dev-support/findbugsExcludeFile.xml @@ -28,4 +28,9 @@ + + + + + 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..61402332a5 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,24 @@ static void setCredentials(Parameters parameters, ServerCredentials credentials) parameters.put(CREDENTIALS_PARAMETER, credentials, CREDENTIALS_CLASS); } + String DATA_TRANSFER_EVENT_CONSUMER_PARAMETER = PREFIX + ".data.transfer.event.consumer"; + Class DATA_TRANSFER_EVENT_CONSUMER_CLASS = Consumer.class; + @SuppressWarnings("unchecked") + static Consumer dataTransferEventConsumer(Parameters parameters) { + return parameters == null ? null + : (Consumer) parameters.get( + DATA_TRANSFER_EVENT_CONSUMER_PARAMETER, DATA_TRANSFER_EVENT_CONSUMER_CLASS); + } + + /** + * Sets the consumer for peer data transfer events. The consumer is invoked on Ratis internal + * threads and must not block. + */ + static void setDataTransferEventConsumer( + Parameters parameters, Consumer consumer) { + parameters.put(DATA_TRANSFER_EVENT_CONSUMER_PARAMETER, consumer, DATA_TRANSFER_EVENT_CONSUMER_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/GrpcDataTransferEvent.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java new file mode 100644 index 0000000000..6770c54831 --- /dev/null +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java @@ -0,0 +1,91 @@ +/* + * 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.protocol.RaftPeerId; + +import java.time.Instant; +import java.util.Objects; + +/** Information about the outcome of a data transfer between Ratis peers. */ +public final class GrpcDataTransferEvent { + /** The protection method used by the peer connection. */ + public enum ProtectionMethod { + TLS, + NONE + } + + /** The transfer outcome. */ + public enum Result { + SUCCESS, + FAILURE + } + + private final Instant timestamp; + private final RaftPeerId source; + private final RaftPeerId destination; + private final ProtectionMethod protectionMethod; + private final Result result; + private final Throwable error; + + private GrpcDataTransferEvent(Instant timestamp, RaftPeerId source, RaftPeerId destination, + ProtectionMethod protectionMethod, Result result, Throwable error) { + this.timestamp = Objects.requireNonNull(timestamp, "timestamp"); + this.source = Objects.requireNonNull(source, "source"); + this.destination = Objects.requireNonNull(destination, "destination"); + this.protectionMethod = Objects.requireNonNull(protectionMethod, "protectionMethod"); + this.result = Objects.requireNonNull(result, "result"); + this.error = error; + } + + public static GrpcDataTransferEvent success(RaftPeerId source, RaftPeerId destination, + ProtectionMethod protectionMethod) { + return new GrpcDataTransferEvent(Instant.now(), source, destination, + protectionMethod, Result.SUCCESS, null); + } + + public static GrpcDataTransferEvent failure(RaftPeerId source, RaftPeerId destination, + ProtectionMethod protectionMethod, Throwable error) { + return new GrpcDataTransferEvent(Instant.now(), source, destination, + protectionMethod, Result.FAILURE, Objects.requireNonNull(error, "error")); + } + + public Instant getTimestamp() { + return timestamp; + } + + public RaftPeerId getSource() { + return source; + } + + public RaftPeerId getDestination() { + return destination; + } + + public ProtectionMethod getProtectionMethod() { + return protectionMethod; + } + + public Result getResult() { + return result; + } + + public Throwable getError() { + return error; + } +} 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..93bf7f416c 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 Consumer dataTransferEventConsumer; 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.dataTransferEventConsumer(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, + Consumer dataTransferEventConsumer, GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, GrpcTlsConfig clientTlsConfig, GrpcTlsConfig serverTlsConfig) { this.servicesCustomizer = servicesCustomizer; this.serverCredentials = serverCredentials; + this.dataTransferEventConsumer = dataTransferEventConsumer; this.forServerSupplier = MemoizedSupplier.valueOf(() -> new SslContexts( tlsConfig, adminTlsConfig, clientTlsConfig, serverTlsConfig, BUILD_SSL_CONTEXT_FOR_SERVER)); @@ -127,7 +131,11 @@ public SupportedRpcType getRpcType() { @Override public LogAppender newLogAppender(RaftServer.Division server, LeaderState state, FollowerInfo f) { - return new GrpcLogAppender(server, state, f); + final GrpcDataTransferEvent.ProtectionMethod protectionMethod = + dataTransferEventConsumer != null && forClientSupplier.get().serverSslContext != null + ? GrpcDataTransferEvent.ProtectionMethod.TLS + : GrpcDataTransferEvent.ProtectionMethod.NONE; + return new GrpcLogAppender(server, state, f, dataTransferEventConsumer, protectionMethod); } @Override 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..4e0da1958f 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,11 +19,13 @@ import org.apache.ratis.conf.RaftProperties; import org.apache.ratis.grpc.GrpcConfigKeys; +import org.apache.ratis.grpc.GrpcDataTransferEvent; import org.apache.ratis.grpc.GrpcUtil; import org.apache.ratis.grpc.metrics.GrpcServerMetrics; import org.apache.ratis.metrics.Timekeeper; import org.apache.ratis.proto.RaftProtos.InstallSnapshotResult; import org.apache.ratis.protocol.RaftPeerId; +import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.retry.RetryPolicy; import org.apache.ratis.server.RaftServer; import org.apache.ratis.server.RaftServerConfigKeys; @@ -50,8 +52,11 @@ import java.io.IOException; import java.io.InterruptedIOException; +import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; import java.util.LinkedList; +import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; @@ -61,7 +66,9 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Consumer; import static org.apache.ratis.server.raftlog.LogProtoUtils.toLogEntryTermIndexString; @@ -161,6 +168,8 @@ synchronized int process(Event event) { @SuppressWarnings({"squid:S3077"}) // Suppress volatile for generic type private volatile StreamObservers appendLogRequestObserver; private final boolean useSeparateHBChannel; + private final Consumer dataTransferEventConsumer; + private final GrpcDataTransferEvent.ProtectionMethod dataTransferProtectionMethod; private final GrpcServerMetrics grpcServerMetrics; @@ -170,6 +179,12 @@ synchronized int process(Event event) { private final ReplyState replyState = new ReplyState(); public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f) { + this(server, leaderState, f, null, GrpcDataTransferEvent.ProtectionMethod.NONE); + } + + public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f, + Consumer dataTransferEventConsumer, + GrpcDataTransferEvent.ProtectionMethod dataTransferProtectionMethod) { super(server, leaderState, f); Objects.requireNonNull(getServerRpc(), "getServerRpc() == null"); @@ -183,6 +198,9 @@ public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, Foll this.logMessageBatchDuration = GrpcConfigKeys.Server.logMessageBatchDuration(properties); this.installSnapshotEnabled = RaftServerConfigKeys.Log.Appender.installSnapshotEnabled(properties); this.useSeparateHBChannel = GrpcConfigKeys.Server.heartbeatChannel(properties); + this.dataTransferEventConsumer = dataTransferEventConsumer; + this.dataTransferProtectionMethod = Objects.requireNonNull( + dataTransferProtectionMethod, "dataTransferProtectionMethod"); grpcServerMetrics = new GrpcServerMetrics(server.getMemberId().toString()); grpcServerMetrics.addPendingRequestsCount(getFollowerId().toString(), pendingRequests::logRequestsSize); @@ -194,6 +212,41 @@ public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, Foll RaftServerConfigKeys.Log.Appender.RETRY_POLICY_KEY); } + private void notifyDataTransferSuccess() { + notifyDataTransferEvent(GrpcDataTransferEvent.Result.SUCCESS, null); + } + + private void notifyDataTransferFailure(Throwable error) { + notifyDataTransferEvent(GrpcDataTransferEvent.Result.FAILURE, Objects.requireNonNull(error, "error")); + } + + private void notifyDataTransferEvent(GrpcDataTransferEvent.Result result, Throwable error) { + if (dataTransferEventConsumer == null) { + return; + } + final GrpcDataTransferEvent event = result == GrpcDataTransferEvent.Result.SUCCESS + ? GrpcDataTransferEvent.success(getServer().getId(), getFollowerId(), dataTransferProtectionMethod) + : GrpcDataTransferEvent.failure(getServer().getId(), getFollowerId(), dataTransferProtectionMethod, error); + try { + dataTransferEventConsumer.accept(event); + } catch (Throwable listenerFailure) { + LOG.warn("gRPC data transfer event consumer threw an exception", listenerFailure); + } + } + + private void notifyAppendEntriesFailure(AppendEntriesRequest request, Throwable error) { + if (request != null && request.containsStateMachineData()) { + request.stopRequestTimer(); + notifyDataTransferFailure(error); + } + } + + private void notifyAppendEntriesFailures(List requests, Throwable error) { + requests.stream() + .filter(AppendEntriesRequest::isSent) + .forEach(request -> notifyAppendEntriesFailure(request, error)); + } + @Override public GrpcServicesImpl getServerRpc() { return (GrpcServicesImpl)super.getServerRpc(); @@ -203,7 +256,8 @@ 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 pendingRequestFailure) { + List discarded = Collections.emptyList(); try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { getClient().resetConnectBackoff(); if (appendLogRequestObserver != null) { @@ -212,7 +266,7 @@ private void resetClient(AppendEntriesRequest request, Event event) { } final int errorCount = replyState.process(event); // clear the pending requests queue and reset the next index of follower - pendingRequests.clear(); + discarded = pendingRequests.clear(dataTransferEventConsumer != null); final FollowerInfo f = getFollower(); final long nextIndex = 1 + Optional.ofNullable(request) .map(AppendEntriesRequest::getPreviousLog) @@ -223,15 +277,16 @@ private void resetClient(AppendEntriesRequest request, Event event) { BatchLogger.print(BatchLogKey.RESET_CLIENT, f.getId() + "-" + followerNextIndex, suffix -> LOG.warn("{}: Follower failed (request=null, errorCount={}); keep nextIndex ({}) unchanged and retry.{}", this, errorCount, followerNextIndex, suffix), logMessageBatchDuration); - return; - } - if (request != null && request.isHeartbeat()) { - return; + } else if (request == null || !request.isHeartbeat()) { + getFollower().computeNextIndex(getNextIndexForError(nextIndex)); } - getFollower().computeNextIndex(getNextIndexForError(nextIndex)); } catch (IOException ie) { LOG.warn("{}: Failed to resetClient for {}", this, getFollowerId(), ie); } + if (pendingRequestFailure != null) { + notifyAppendEntriesFailure(request, pendingRequestFailure); + notifyAppendEntriesFailures(discarded, pendingRequestFailure); + } } private boolean isFollowerCommitBehindLastCommitIndex() { @@ -399,7 +454,8 @@ private void appendLog(boolean heartbeat) throws IOException { if (pending == null) { return; } - request = new AppendEntriesRequest(pending, getFollowerId(), grpcServerMetrics); + request = new AppendEntriesRequest( + pending, getFollowerId(), grpcServerMetrics, dataTransferEventConsumer != null); pendingRequests.put(request); increaseNextIndex(pending); if (appendLogRequestObserver == null) { @@ -413,7 +469,16 @@ private void appendLog(boolean heartbeat) throws IOException { sleep(remaining, heartbeat); } if (isRunning()) { - sendRequest(request, pending); + try { + sendRequest(request, pending); + } catch (IOException | RuntimeException e) { + final AppendEntriesRequest failed = pendingRequests.remove(request.getCallId(), request.isHeartbeat()); + if (failed != null) { + failed.stopRequestTimer(); + notifyAppendEntriesFailure(failed, e); + } + throw e; + } } } @@ -453,6 +518,8 @@ private void timeoutAppendRequest(long cid, boolean heartbeat) { this, heartbeat ? "HEARTBEAT " : "", errorCount, pending); grpcServerMetrics.onRequestTimeout(getFollowerId().toString(), heartbeat); pending.stopRequestTimer(); + notifyAppendEntriesFailure(pending, + new TimeoutIOException("Timed out appendEntries request " + pending)); } } @@ -489,6 +556,9 @@ public void onNext(AppendEntriesReplyProto reply) { if (request != null) { request.stopRequestTimer(); // Update completion time getFollower().updateLastRespondedAppendEntriesSendTime(request.getSendTime()); + if (request.containsStateMachineData()) { + notifyDataTransferSuccess(); + } } getFollower().updateLastRpcResponseTime(); @@ -554,13 +624,14 @@ 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() { LOG.info("{}: follower responses appendEntries COMPLETED", this); - resetClient(null, Event.COMPLETE); + resetClient(null, Event.COMPLETE, + new IOException("AppendEntries response stream completed with pending requests")); } @Override @@ -570,10 +641,15 @@ public String toString() { } private void updateNextIndex(long replyNextIndex) { + final List discarded; try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - pendingRequests.clear(); + discarded = pendingRequests.clear(dataTransferEventConsumer != null); getFollower().setNextIndex(replyNextIndex); } + if (!discarded.isEmpty()) { + notifyAppendEntriesFailures(discarded, + new IOException("AppendEntries request was invalidated by an inconsistency response")); + } } private class InstallSnapshotResponseHandler implements StreamObserver { @@ -581,6 +657,7 @@ private class InstallSnapshotResponseHandler implements StreamObserver pending = new LinkedList<>(); private final CompletableFuture done = new CompletableFuture<>(); private final boolean isNotificationOnly; + private final AtomicBoolean transferReported = new AtomicBoolean(); InstallSnapshotResponseHandler() { this(false); @@ -645,11 +722,32 @@ void waitForResponse() { done.get(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); + final InterruptedIOException interrupted = + new InterruptedIOException("Interrupted while waiting for InstallSnapshot responses"); + interrupted.initCause(e); + notifyTransferFailure(interrupted); } catch (ExecutionException e) { + notifyTransferFailure(e); throw new IllegalStateException("Failed to complete " + name, e); } } + void notifyTransferSuccess() { + if (!isNotificationOnly && transferReported.compareAndSet(false, true)) { + notifyDataTransferSuccess(); + } + } + + void notifyTransferFailure(Throwable error) { + if (!isNotificationOnly && transferReported.compareAndSet(false, true)) { + notifyDataTransferFailure(error); + } + } + + private void notifyTransferFailure(InstallSnapshotReplyProto reply) { + notifyTransferFailure(new IOException("InstallSnapshot failed with result " + reply.getResult())); + } + void close() { done.complete(null); notifyLogAppender(); @@ -694,9 +792,11 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case NOT_LEADER: + notifyTransferFailure(reply); onFollowerTerm(reply.getTerm()); break; case CONF_MISMATCH: + notifyTransferFailure(reply); LOG.error("{}: CONF_MISMATCH ({}): Leader {} has it set to {} but follower {} has it set to {}", this, RaftServerConfigKeys.Log.Appender.INSTALL_SNAPSHOT_ENABLED_KEY, getServer().getId(), installSnapshotEnabled, getFollowerId(), !installSnapshotEnabled); @@ -712,6 +812,7 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case SNAPSHOT_UNAVAILABLE: + notifyTransferFailure(reply); BatchLogger.print(BatchLogKey.SNAPSHOT_UNAVAILABLE, name, suffix -> LOG.info("{}: Follower failed since the snapshot is unavailable {}", this, suffix)); getFollower().setAttemptedToInstallSnapshot(); @@ -719,10 +820,12 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case UNRECOGNIZED: + notifyTransferFailure(reply); LOG.error("{}: Reply result {}, {}", name, reply.getResult(), ServerStringUtils.toInstallSnapshotReplyString(reply)); break; case SNAPSHOT_EXPIRED: + notifyTransferFailure(reply); LOG.warn("{}: Follower failed since the request expired, {}", name, ServerStringUtils.toInstallSnapshotReplyString(reply)); default: @@ -732,13 +835,14 @@ public void onNext(InstallSnapshotReplyProto reply) { @Override public void onError(Throwable t) { + notifyTransferFailure(t); if (!isRunning()) { LOG.info("{} is stopped", GrpcLogAppender.this); return; } GrpcUtil.warn(LOG, () -> this + ": Failed InstallSnapshot", t); grpcServerMetrics.onRequestRetry(); // Update try counter - resetClient(null, Event.ERROR); + resetClient(null, Event.ERROR, t); close(); } @@ -768,6 +872,7 @@ private void installSnapshot(SnapshotInfo snapshot) { final InstallSnapshotResponseHandler responseHandler = new InstallSnapshotResponseHandler(); StreamObserver snapshotRequestObserver = null; final String requestId = UUID.randomUUID().toString(); + boolean sentAllRequests = true; try { snapshotRequestObserver = getClient().installSnapshot( getFollower().getName() + "-installSnapshot-" + requestId, @@ -778,6 +883,7 @@ private void installSnapshot(SnapshotInfo snapshot) { getFollower().updateLastRpcSendTime(false); responseHandler.addPending(request); } else { + sentAllRequests = false; break; } } @@ -785,6 +891,7 @@ private void installSnapshot(SnapshotInfo snapshot) { grpcServerMetrics.onInstallSnapshot(); } catch (Exception e) { LOG.warn(this + ": failed to installSnapshot " + snapshot, e); + responseHandler.notifyTransferFailure(e); if (snapshotRequestObserver != null) { snapshotRequestObserver.onError(e); } @@ -792,9 +899,16 @@ private void installSnapshot(SnapshotInfo snapshot) { } responseHandler.waitForResponse(); - if (responseHandler.hasAllResponse()) { + if (!sentAllRequests) { + responseHandler.notifyTransferFailure( + new IOException("InstallSnapshot stopped before all requests were sent")); + } else if (responseHandler.hasAllResponse()) { getFollower().setSnapshotIndex(snapshot.getTermIndex().getIndex()); LOG.info("{}: installed snapshot {} successfully", this, snapshot); + responseHandler.notifyTransferSuccess(); + } else { + responseHandler.notifyTransferFailure( + new IOException("InstallSnapshot completed with pending responses")); } } @@ -841,16 +955,25 @@ static class AppendEntriesRequest { private final long callId; private final TermIndex previousLog; private final int entriesCount; + private final boolean containsStateMachineData; private final TermIndex firstEntry; private final TermIndex lastEntry; @SuppressWarnings({"squid:S3077"}) // Suppress volatile for generic type private volatile Timestamp sendTime; - AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, GrpcServerMetrics grpcServerMetrics) { + AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, + GrpcServerMetrics grpcServerMetrics) { + this(proto, followerId, grpcServerMetrics, false); + } + + AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, + GrpcServerMetrics grpcServerMetrics, boolean checkStateMachineData) { this.callId = proto.getServerRequest().getCallId(); this.previousLog = proto.hasPreviousLog()? TermIndex.valueOf(proto.getPreviousLog()): null; this.entriesCount = proto.getEntriesCount(); + this.containsStateMachineData = checkStateMachineData && proto.getEntriesList().stream() + .anyMatch(entry -> entry.hasStateMachineLogEntry()); this.firstEntry = entriesCount > 0? TermIndex.valueOf(proto.getEntries(0)): null; this.lastEntry = entriesCount > 0? TermIndex.valueOf(proto.getEntries(entriesCount - 1)): null; @@ -880,13 +1003,24 @@ void startRequestTimer() { } void stopRequestTimer() { - timerContext.stop(); + if (timerContext != null) { + timerContext.stop(); + timerContext = null; + } } boolean isHeartbeat() { return entriesCount == 0; } + boolean containsStateMachineData() { + return containsStateMachineData; + } + + boolean isSent() { + return sendTime != null; + } + @Override public String toString() { return JavaUtils.getClassSimpleName(getClass()) @@ -903,9 +1037,20 @@ int logRequestsSize() { return logRequests.size(); } - void clear() { - logRequests.clear(); + List clear(boolean collectLogRequests) { + if (!collectLogRequests) { + logRequests.clear(); + heartbeats.clear(); + return Collections.emptyList(); + } + final List removed = new ArrayList<>(); + logRequests.forEach((callId, request) -> { + if (logRequests.remove(callId, request)) { + removed.add(request); + } + }); heartbeats.clear(); + return removed; } void put(AppendEntriesRequest request) { diff --git a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java new file mode 100644 index 0000000000..2665de2863 --- /dev/null +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java @@ -0,0 +1,254 @@ +/* + * 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.AppendEntriesRequestProto; +import org.apache.ratis.proto.RaftProtos.ReplicationLevel; +import org.apache.ratis.protocol.RaftClientReply; +import org.apache.ratis.protocol.RaftPeerId; +import org.apache.ratis.security.SecurityTestUtils; +import org.apache.ratis.server.RaftServer; +import org.apache.ratis.server.RaftServerConfigKeys; +import org.apache.ratis.server.impl.MiniRaftCluster; +import org.apache.ratis.server.impl.PeerChanges; +import org.apache.ratis.server.impl.RaftServerTestUtil; +import org.apache.ratis.statemachine.SnapshotInfo; +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.apache.ratis.util.SizeInBytes; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import javax.net.ssl.KeyManager; +import javax.net.ssl.TrustManager; + +import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; + +public class TestGrpcDataTransferEventListener extends BaseTest { + private static RaftProperties newProperties() { + final RaftProperties properties = new RaftProperties(); + properties.setClass(MiniRaftCluster.STATEMACHINE_CLASS_KEY, + SimpleStateMachine4Testing.class, StateMachine.class); + return properties; + } + + private static Parameters newParameters(boolean tls, + ConcurrentLinkedQueue events) throws Exception { + final Parameters parameters = new Parameters(); + GrpcConfigKeys.Server.setDataTransferEventConsumer(parameters, events::add); + if (tls) { + final KeyManager serverKeyManager = + SecurityTestUtils.getKeyManager(SecurityTestUtils::getServerKeyStore); + final TrustManager serverTrustManager = + SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore); + final KeyManager clientKeyManager = + SecurityTestUtils.getKeyManager(SecurityTestUtils::getClientKeyStore); + final TrustManager clientTrustManager = + SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore); + + GrpcConfigKeys.Server.setTlsConf(parameters, + new GrpcTlsConfig(serverKeyManager, serverTrustManager, true)); + final GrpcTlsConfig clientConfig = + new GrpcTlsConfig(clientKeyManager, clientTrustManager, true); + GrpcConfigKeys.Admin.setTlsConf(parameters, clientConfig); + GrpcConfigKeys.Client.setTlsConf(parameters, clientConfig); + } + return parameters; + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + @Timeout(value = 60, unit = TimeUnit.SECONDS) + public void testSuccessfulAppendEntries(boolean tls) throws Exception { + final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); + final Parameters parameters = newParameters(tls, events); + final String[] ids = MiniRaftCluster.generateIds(3, 10); + + try (MiniRaftClusterWithGrpc cluster = + new MiniRaftClusterWithGrpc(ids, new String[0], newProperties(), parameters)) { + cluster.start(); + final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); + final RaftPeerId leaderId = leader.getId(); + events.clear(); + + try (RaftClient client = cluster.createClient(leaderId)) { + final RaftClientReply write = + client.io().send(new RaftTestUtil.SimpleMessage("data-transfer")); + Assertions.assertTrue(write.isSuccess()); + Assertions.assertTrue( + client.io().watch(write.getLogIndex(), ReplicationLevel.ALL_COMMITTED).isSuccess()); + } + + final Set followers = cluster.getFollowers().stream() + .map(RaftServer.Division::getId) + .collect(Collectors.toSet()); + JavaUtils.attempt(() -> { + final Set destinations = successfulEvents(events, leaderId).stream() + .map(GrpcDataTransferEvent::getDestination) + .collect(Collectors.toSet()); + Assertions.assertEquals(followers, destinations); + }, 10, HUNDRED_MILLIS, "data transfer events", LOG); + + final GrpcDataTransferEvent.ProtectionMethod expected = tls + ? GrpcDataTransferEvent.ProtectionMethod.TLS + : GrpcDataTransferEvent.ProtectionMethod.NONE; + successfulEvents(events, leaderId).forEach(event -> { + Assertions.assertNotNull(event.getTimestamp()); + Assertions.assertEquals(expected, event.getProtectionMethod()); + Assertions.assertNull(event.getError()); + }); + } + } + + private static List successfulEvents( + ConcurrentLinkedQueue events, RaftPeerId source) { + return events.stream() + .filter(event -> event.getSource().equals(source)) + .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.SUCCESS) + .collect(Collectors.toList()); + } + + @Test + @Timeout(value = 120, unit = TimeUnit.SECONDS) + public void testSuccessfulInstallSnapshot() throws Exception { + final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); + final Parameters parameters = newParameters(false, events); + final RaftProperties properties = newProperties(); + RaftServerConfigKeys.Snapshot.setAutoTriggerEnabled(properties, true); + RaftServerConfigKeys.Snapshot.setAutoTriggerThreshold(properties, 64); + RaftServerConfigKeys.Log.setPurgeGap(properties, 8); + RaftServerConfigKeys.Log.Appender.setSnapshotChunkSizeMax(properties, SizeInBytes.ONE_KB); + RaftServerConfigKeys.LeaderElection.setMemberMajorityAdd(properties, true); + + final String[] ids = MiniRaftCluster.generateIds(1, 30); + try (MiniRaftClusterWithGrpc cluster = + new MiniRaftClusterWithGrpc(ids, new String[0], properties, parameters)) { + cluster.start(); + final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); + final RaftPeerId leaderId = leader.getId(); + + try (RaftClient client = cluster.createClient(leaderId)) { + for (int i = 0; i < 127; i++) { + Assertions.assertTrue( + client.io().send(new RaftTestUtil.SimpleMessage("snapshot-" + i)).isSuccess()); + } + Assertions.assertTrue(client.getSnapshotManagementApi(leaderId).create(3000).isSuccess()); + } + final SnapshotInfo leaderSnapshot = leader.getStateMachine().getLatestSnapshot(); + Assertions.assertNotNull(leaderSnapshot); + events.clear(); + + final PeerChanges change = cluster.addNewPeers(1, true); + final RaftPeerId addedPeer = change.getAddedPeers().get(0).getId(); + cluster.setConfiguration(change.getPeersInNewConf()); + RaftServerTestUtil.waitAndCheckNewConf(cluster, change.getPeersInNewConf(), 0, null); + + JavaUtils.attempt(() -> { + final SnapshotInfo installed = + cluster.getDivision(addedPeer).getStateMachine().getLatestSnapshot(); + Assertions.assertNotNull(installed); + Assertions.assertTrue(successfulEvents(events, leaderId).stream() + .anyMatch(event -> event.getDestination().equals(addedPeer))); + }, 20, ONE_SECOND, "snapshot data transfer event", LOG); + } + } + + @Test + @Timeout(value = 60, unit = TimeUnit.SECONDS) + public void testAppendEntriesFailureAndConsumerException() throws Exception { + final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); + final AtomicInteger consumerCalls = new AtomicInteger(); + final Parameters parameters = new Parameters(); + GrpcConfigKeys.Server.setDataTransferEventConsumer(parameters, event -> { + events.add(event); + consumerCalls.incrementAndGet(); + if (event.getResult() == GrpcDataTransferEvent.Result.FAILURE) { + throw new IllegalStateException("Injected consumer failure"); + } + }); + + final String[] ids = MiniRaftCluster.generateIds(2, 20); + try (MiniRaftClusterWithGrpc cluster = + new MiniRaftClusterWithGrpc(ids, new String[0], newProperties(), parameters)) { + cluster.start(); + final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); + final RaftPeerId leaderId = leader.getId(); + events.clear(); + consumerCalls.set(0); + + final AtomicBoolean injectFailure = new AtomicBoolean(true); + try { + CodeInjectionForTesting.put( + GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST, (localId, remoteId, args) -> { + if (leaderId.equals(localId) + && args.length > 0 + && args[0] instanceof AppendEntriesRequestProto) { + final AppendEntriesRequestProto request = (AppendEntriesRequestProto) args[0]; + final boolean containsStateMachineData = request.getEntriesList().stream() + .anyMatch(entry -> entry.hasStateMachineLogEntry()); + if (containsStateMachineData && injectFailure.compareAndSet(true, false)) { + throw new IllegalStateException("Injected AppendEntries failure"); + } + } + return false; + }); + + try (RaftClient client = cluster.createClient(leaderId)) { + Assertions.assertTrue( + client.io().send(new RaftTestUtil.SimpleMessage("retry-after-failure")).isSuccess()); + } + + JavaUtils.attempt(() -> Assertions.assertTrue(events.stream() + .anyMatch(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE)), + 10, HUNDRED_MILLIS, "failed data transfer event", LOG); + } finally { + CodeInjectionForTesting.remove(GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST); + } + + final GrpcDataTransferEvent failure = events.stream() + .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE) + .findFirst() + .orElseThrow(AssertionError::new); + Assertions.assertEquals(leaderId, failure.getSource()); + Assertions.assertEquals(GrpcDataTransferEvent.ProtectionMethod.NONE, + failure.getProtectionMethod()); + Assertions.assertTrue(failure.getError() instanceof IllegalStateException); + Assertions.assertEquals(1, events.stream() + .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE) + .count()); + Assertions.assertTrue(consumerCalls.get() > 0); + } + } +} From c939b32357b122f921f9fb07dead7d9ecaa56d3b Mon Sep 17 00:00:00 2001 From: HTHou Date: Mon, 7 Sep 2026 11:51:35 +0800 Subject: [PATCH 2/6] RATIS-2681. Expose generic gRPC log appender lifecycle callbacks --- .../dev-support/findbugsExcludeFile.xml | 5 - .../org/apache/ratis/grpc/GrpcConfigKeys.java | 21 +- .../ratis/grpc/GrpcDataTransferEvent.java | 91 ------ .../org/apache/ratis/grpc/GrpcFactory.java | 22 +- .../ratis/grpc/GrpcLogAppenderListener.java | 70 +++++ .../ratis/grpc/server/GrpcLogAppender.java | 264 +++++++----------- .../TestGrpcDataTransferEventListener.java | 254 ----------------- .../grpc/TestGrpcLogAppenderListener.java | 208 ++++++++++++++ .../server/TestGrpcLogAppenderCallbacks.java | 248 ++++++++++++++++ 9 files changed, 642 insertions(+), 541 deletions(-) delete mode 100644 ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java create mode 100644 ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java delete mode 100644 ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java create mode 100644 ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java create mode 100644 ratis-test/src/test/java/org/apache/ratis/grpc/server/TestGrpcLogAppenderCallbacks.java diff --git a/ratis-grpc/dev-support/findbugsExcludeFile.xml b/ratis-grpc/dev-support/findbugsExcludeFile.xml index ba57c27c2b..3f10ecb4e6 100644 --- a/ratis-grpc/dev-support/findbugsExcludeFile.xml +++ b/ratis-grpc/dev-support/findbugsExcludeFile.xml @@ -28,9 +28,4 @@ - - - - - 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 61402332a5..6a0056e1e0 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,22 +304,15 @@ static void setCredentials(Parameters parameters, ServerCredentials credentials) parameters.put(CREDENTIALS_PARAMETER, credentials, CREDENTIALS_CLASS); } - String DATA_TRANSFER_EVENT_CONSUMER_PARAMETER = PREFIX + ".data.transfer.event.consumer"; - Class DATA_TRANSFER_EVENT_CONSUMER_CLASS = Consumer.class; - @SuppressWarnings("unchecked") - static Consumer dataTransferEventConsumer(Parameters parameters) { - return parameters == null ? null - : (Consumer) parameters.get( - DATA_TRANSFER_EVENT_CONSUMER_PARAMETER, DATA_TRANSFER_EVENT_CONSUMER_CLASS); + String LOG_APPENDER_LISTENER_FACTORY_PARAMETER = PREFIX + ".log.appender.listener.factory"; + static GrpcLogAppenderListener.Factory logAppenderListenerFactory(Parameters parameters) { + return parameters == null ? null : parameters.get( + LOG_APPENDER_LISTENER_FACTORY_PARAMETER, GrpcLogAppenderListener.Factory.class); } - /** - * Sets the consumer for peer data transfer events. The consumer is invoked on Ratis internal - * threads and must not block. - */ - static void setDataTransferEventConsumer( - Parameters parameters, Consumer consumer) { - parameters.put(DATA_TRANSFER_EVENT_CONSUMER_PARAMETER, consumer, DATA_TRANSFER_EVENT_CONSUMER_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, GrpcLogAppenderListener.Factory.class); } String TLS_CONF_PARAMETER = PREFIX + ".tls.conf"; diff --git a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java deleted file mode 100644 index 6770c54831..0000000000 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcDataTransferEvent.java +++ /dev/null @@ -1,91 +0,0 @@ -/* - * 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.protocol.RaftPeerId; - -import java.time.Instant; -import java.util.Objects; - -/** Information about the outcome of a data transfer between Ratis peers. */ -public final class GrpcDataTransferEvent { - /** The protection method used by the peer connection. */ - public enum ProtectionMethod { - TLS, - NONE - } - - /** The transfer outcome. */ - public enum Result { - SUCCESS, - FAILURE - } - - private final Instant timestamp; - private final RaftPeerId source; - private final RaftPeerId destination; - private final ProtectionMethod protectionMethod; - private final Result result; - private final Throwable error; - - private GrpcDataTransferEvent(Instant timestamp, RaftPeerId source, RaftPeerId destination, - ProtectionMethod protectionMethod, Result result, Throwable error) { - this.timestamp = Objects.requireNonNull(timestamp, "timestamp"); - this.source = Objects.requireNonNull(source, "source"); - this.destination = Objects.requireNonNull(destination, "destination"); - this.protectionMethod = Objects.requireNonNull(protectionMethod, "protectionMethod"); - this.result = Objects.requireNonNull(result, "result"); - this.error = error; - } - - public static GrpcDataTransferEvent success(RaftPeerId source, RaftPeerId destination, - ProtectionMethod protectionMethod) { - return new GrpcDataTransferEvent(Instant.now(), source, destination, - protectionMethod, Result.SUCCESS, null); - } - - public static GrpcDataTransferEvent failure(RaftPeerId source, RaftPeerId destination, - ProtectionMethod protectionMethod, Throwable error) { - return new GrpcDataTransferEvent(Instant.now(), source, destination, - protectionMethod, Result.FAILURE, Objects.requireNonNull(error, "error")); - } - - public Instant getTimestamp() { - return timestamp; - } - - public RaftPeerId getSource() { - return source; - } - - public RaftPeerId getDestination() { - return destination; - } - - public ProtectionMethod getProtectionMethod() { - return protectionMethod; - } - - public Result getResult() { - return result; - } - - public Throwable getError() { - return error; - } -} 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 93bf7f416c..ea9610a254 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,7 +93,7 @@ private SslContexts(GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, private final GrpcServices.Customizer servicesCustomizer; private final ServerCredentials serverCredentials; - private final Consumer dataTransferEventConsumer; + private final GrpcLogAppenderListener.Factory logAppenderListenerFactory; private final Supplier forServerSupplier; private final Supplier forClientSupplier; @@ -101,7 +101,7 @@ private SslContexts(GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, public GrpcFactory(Parameters parameters) { this(GrpcConfigKeys.Server.servicesCustomizer(parameters), GrpcConfigKeys.Server.credentials(parameters), - GrpcConfigKeys.Server.dataTransferEventConsumer(parameters), + GrpcConfigKeys.Server.logAppenderListenerFactory(parameters), GrpcConfigKeys.TLS.conf(parameters), GrpcConfigKeys.Admin.tlsConf(parameters), GrpcConfigKeys.Client.tlsConf(parameters), @@ -111,12 +111,12 @@ public GrpcFactory(Parameters parameters) { private GrpcFactory(GrpcServices.Customizer servicesCustomizer, ServerCredentials serverCredentials, - Consumer dataTransferEventConsumer, + GrpcLogAppenderListener.Factory logAppenderListenerFactory, GrpcTlsConfig tlsConfig, GrpcTlsConfig adminTlsConfig, GrpcTlsConfig clientTlsConfig, GrpcTlsConfig serverTlsConfig) { this.servicesCustomizer = servicesCustomizer; this.serverCredentials = serverCredentials; - this.dataTransferEventConsumer = dataTransferEventConsumer; + this.logAppenderListenerFactory = logAppenderListenerFactory; this.forServerSupplier = MemoizedSupplier.valueOf(() -> new SslContexts( tlsConfig, adminTlsConfig, clientTlsConfig, serverTlsConfig, BUILD_SSL_CONTEXT_FOR_SERVER)); @@ -131,11 +131,15 @@ public SupportedRpcType getRpcType() { @Override public LogAppender newLogAppender(RaftServer.Division server, LeaderState state, FollowerInfo f) { - final GrpcDataTransferEvent.ProtectionMethod protectionMethod = - dataTransferEventConsumer != null && forClientSupplier.get().serverSslContext != null - ? GrpcDataTransferEvent.ProtectionMethod.TLS - : GrpcDataTransferEvent.ProtectionMethod.NONE; - return new GrpcLogAppender(server, state, f, dataTransferEventConsumer, protectionMethod); + 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", 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..5f99753c59 --- /dev/null +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java @@ -0,0 +1,70 @@ +/* + * 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.proto.RaftProtos.InstallSnapshotReplyProto; +import org.apache.ratis.proto.RaftProtos.InstallSnapshotRequestProto; +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, + * filtering messages and aggregating snapshot chunks. + * Append request, matched-reply and reset callbacks are serialized per appender. Failure + * callbacks may race with these callbacks. 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); + } + + /** An append attempt is registered, before establishing or writing its stream. */ + default void onAppendEntriesRequest(AppendEntriesRequestProto request) { } + + /** A response was matched to a pending append request. */ + default void onAppendEntriesReply(AppendEntriesReplyProto reply) { } + + /** A local send error or request timeout occurred; a later stream notification may follow. */ + default void onAppendEntriesFailure(long callId, Throwable error) { } + + /** Pending append attempts are invalidated by a reset or by a stopped response stream. */ + default void onAppendEntriesReset(Throwable error) { } + + /** A snapshot attempt starts before stream creation, including notification-only attempts. */ + default void onInstallSnapshotStart(String requestId, boolean notificationOnly) { } + + /** A snapshot chunk or notification is about to be sent. */ + default void onInstallSnapshotRequest(String requestId, InstallSnapshotRequestProto request) { } + + /** A snapshot response was received; this is not necessarily a terminal response. */ + default void onInstallSnapshotReply(String requestId, InstallSnapshotReplyProto reply) { } + + /** + * The stream completed (null error), failed, or its sender was interrupted. Completion alone + * does not imply that all chunks were sent or acknowledged. Racing notifications may repeat. + */ + default void onInstallSnapshotEnd(String requestId, Throwable error) { } +} 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 4e0da1958f..2413a48db2 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,7 +19,7 @@ import org.apache.ratis.conf.RaftProperties; import org.apache.ratis.grpc.GrpcConfigKeys; -import org.apache.ratis.grpc.GrpcDataTransferEvent; +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; @@ -52,11 +52,8 @@ import java.io.IOException; import java.io.InterruptedIOException; -import java.util.ArrayList; -import java.util.Collections; import java.util.Comparator; import java.util.LinkedList; -import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; @@ -66,7 +63,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; @@ -168,10 +164,9 @@ synchronized int process(Event event) { @SuppressWarnings({"squid:S3077"}) // Suppress volatile for generic type private volatile StreamObservers appendLogRequestObserver; private final boolean useSeparateHBChannel; - private final Consumer dataTransferEventConsumer; - private final GrpcDataTransferEvent.ProtectionMethod dataTransferProtectionMethod; private final GrpcServerMetrics grpcServerMetrics; + private final GrpcLogAppenderListener listener; private final AutoCloseableReadWriteLock lock; private final StackTraceElement caller; @@ -179,13 +174,13 @@ synchronized int process(Event event) { private final ReplyState replyState = new ReplyState(); public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f) { - this(server, leaderState, f, null, GrpcDataTransferEvent.ProtectionMethod.NONE); + this(server, leaderState, f, null); } public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, FollowerInfo f, - Consumer dataTransferEventConsumer, - GrpcDataTransferEvent.ProtectionMethod dataTransferProtectionMethod) { + GrpcLogAppenderListener listener) { super(server, leaderState, f); + this.listener = listener; Objects.requireNonNull(getServerRpc(), "getServerRpc() == null"); @@ -198,9 +193,6 @@ public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, Foll this.logMessageBatchDuration = GrpcConfigKeys.Server.logMessageBatchDuration(properties); this.installSnapshotEnabled = RaftServerConfigKeys.Log.Appender.installSnapshotEnabled(properties); this.useSeparateHBChannel = GrpcConfigKeys.Server.heartbeatChannel(properties); - this.dataTransferEventConsumer = dataTransferEventConsumer; - this.dataTransferProtectionMethod = Objects.requireNonNull( - dataTransferProtectionMethod, "dataTransferProtectionMethod"); grpcServerMetrics = new GrpcServerMetrics(server.getMemberId().toString()); grpcServerMetrics.addPendingRequestsCount(getFollowerId().toString(), pendingRequests::logRequestsSize); @@ -212,41 +204,22 @@ public GrpcLogAppender(RaftServer.Division server, LeaderState leaderState, Foll RaftServerConfigKeys.Log.Appender.RETRY_POLICY_KEY); } - private void notifyDataTransferSuccess() { - notifyDataTransferEvent(GrpcDataTransferEvent.Result.SUCCESS, null); - } - - private void notifyDataTransferFailure(Throwable error) { - notifyDataTransferEvent(GrpcDataTransferEvent.Result.FAILURE, Objects.requireNonNull(error, "error")); - } - - private void notifyDataTransferEvent(GrpcDataTransferEvent.Result result, Throwable error) { - if (dataTransferEventConsumer == null) { - return; - } - final GrpcDataTransferEvent event = result == GrpcDataTransferEvent.Result.SUCCESS - ? GrpcDataTransferEvent.success(getServer().getId(), getFollowerId(), dataTransferProtectionMethod) - : GrpcDataTransferEvent.failure(getServer().getId(), getFollowerId(), dataTransferProtectionMethod, error); - try { - dataTransferEventConsumer.accept(event); - } catch (Throwable listenerFailure) { - LOG.warn("gRPC data transfer event consumer threw an exception", listenerFailure); + private void notifyListener(Consumer notification) { + if (listener != null) { + try { + notification.accept(listener); + } catch (Throwable t) { + LOG.warn("gRPC log appender listener threw an exception", t); + } } } - private void notifyAppendEntriesFailure(AppendEntriesRequest request, Throwable error) { - if (request != null && request.containsStateMachineData()) { - request.stopRequestTimer(); - notifyDataTransferFailure(error); + private void notifyReset(Throwable error) { + try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + notifyListener(l -> l.onAppendEntriesReset(error)); } } - private void notifyAppendEntriesFailures(List requests, Throwable error) { - requests.stream() - .filter(AppendEntriesRequest::isSent) - .forEach(request -> notifyAppendEntriesFailure(request, error)); - } - @Override public GrpcServicesImpl getServerRpc() { return (GrpcServicesImpl)super.getServerRpc(); @@ -256,9 +229,9 @@ private GrpcServerProtocolClient getClient() throws IOException { return getServerRpc().getProxies().getProxy(getFollowerId()); } - private void resetClient(AppendEntriesRequest request, Event event, Throwable pendingRequestFailure) { - List discarded = Collections.emptyList(); + private void resetClient(AppendEntriesRequest request, Event event, Throwable error) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { + notifyListener(l -> l.onAppendEntriesReset(error)); getClient().resetConnectBackoff(); if (appendLogRequestObserver != null) { appendLogRequestObserver.stop(); @@ -266,7 +239,7 @@ private void resetClient(AppendEntriesRequest request, Event event, Throwable pe } final int errorCount = replyState.process(event); // clear the pending requests queue and reset the next index of follower - discarded = pendingRequests.clear(dataTransferEventConsumer != null); + pendingRequests.clear(); final FollowerInfo f = getFollower(); final long nextIndex = 1 + Optional.ofNullable(request) .map(AppendEntriesRequest::getPreviousLog) @@ -277,16 +250,15 @@ private void resetClient(AppendEntriesRequest request, Event event, Throwable pe BatchLogger.print(BatchLogKey.RESET_CLIENT, f.getId() + "-" + followerNextIndex, suffix -> LOG.warn("{}: Follower failed (request=null, errorCount={}); keep nextIndex ({}) unchanged and retry.{}", this, errorCount, followerNextIndex, suffix), logMessageBatchDuration); - } else if (request == null || !request.isHeartbeat()) { - getFollower().computeNextIndex(getNextIndexForError(nextIndex)); + return; + } + if (request != null && request.isHeartbeat()) { + return; } + getFollower().computeNextIndex(getNextIndexForError(nextIndex)); } catch (IOException ie) { LOG.warn("{}: Failed to resetClient for {}", this, getFollowerId(), ie); } - if (pendingRequestFailure != null) { - notifyAppendEntriesFailure(request, pendingRequestFailure); - notifyAppendEntriesFailures(discarded, pendingRequestFailure); - } } private boolean isFollowerCommitBehindLastCommitIndex() { @@ -444,41 +416,48 @@ 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, dataTransferEventConsumer != null); - 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); + notifyListener(l -> l.onAppendEntriesRequest(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()) { - try { + final TimeDuration remaining = getRemainingWaitTime(); + if (remaining.isPositive()) { + sleep(remaining, heartbeat); + } + if (isRunning()) { sendRequest(request, pending); - } catch (IOException | RuntimeException e) { - final AppendEntriesRequest failed = pendingRequests.remove(request.getCallId(), request.isHeartbeat()); + } else { + final long cid = request.getCallId(); + pendingRequests.remove(cid, request.isHeartbeat()); + notifyListener(l -> l.onAppendEntriesFailure(cid, new IOException("Log appender stopped before send"))); + } + } 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(); - notifyAppendEntriesFailure(failed, e); } - throw e; + notifyListener(l -> l.onAppendEntriesFailure(cid, e)); } + throw e; } } @@ -510,7 +489,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); @@ -518,8 +497,8 @@ private void timeoutAppendRequest(long cid, boolean heartbeat) { this, heartbeat ? "HEARTBEAT " : "", errorCount, pending); grpcServerMetrics.onRequestTimeout(getFollowerId().toString(), heartbeat); pending.stopRequestTimer(); - notifyAppendEntriesFailure(pending, - new TimeoutIOException("Timed out appendEntries request " + pending)); + notifyListener(l -> l.onAppendEntriesFailure(cid, + new TimeoutIOException("Timed out appendEntries request " + pending))); } } @@ -539,7 +518,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()); /** @@ -552,13 +531,16 @@ private class AppendLogResponseHandler implements StreamObserver l.onAppendEntriesReply(reply)); + } + } if (request != null) { request.stopRequestTimer(); // Update completion time getFollower().updateLastRespondedAppendEntriesSendTime(request.getSendTime()); - if (request.containsStateMachineData()) { - notifyDataTransferSuccess(); - } } getFollower().updateLastRpcResponseTime(); @@ -616,6 +598,7 @@ private void onNextImpl(AppendEntriesRequest request, AppendEntriesReplyProto re @Override public void onError(Throwable t) { if (!isRunning()) { + notifyReset(t); LOG.info("{} is already stopped", GrpcLogAppender.this); return; } @@ -630,8 +613,12 @@ public void onError(Throwable t) { @Override public void onCompleted() { LOG.info("{}: follower responses appendEntries COMPLETED", this); - resetClient(null, Event.COMPLETE, - new IOException("AppendEntries response stream completed with pending requests")); + final IOException error = new IOException("AppendEntries response stream completed"); + if (!isRunning()) { + notifyReset(error); + return; + } + resetClient(null, Event.COMPLETE, error); } @Override @@ -641,23 +628,19 @@ public String toString() { } private void updateNextIndex(long replyNextIndex) { - final List discarded; try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - discarded = pendingRequests.clear(dataTransferEventConsumer != null); + notifyListener(l -> l.onAppendEntriesReset(new IOException("AppendEntries invalidated by inconsistency"))); + pendingRequests.clear(); getFollower().setNextIndex(replyNextIndex); } - if (!discarded.isEmpty()) { - notifyAppendEntriesFailures(discarded, - new IOException("AppendEntries request was invalidated by an inconsistency response")); - } } - private class InstallSnapshotResponseHandler implements StreamObserver { + class InstallSnapshotResponseHandler implements StreamObserver { private final String name; private final Queue pending = new LinkedList<>(); private final CompletableFuture done = new CompletableFuture<>(); private final boolean isNotificationOnly; - private final AtomicBoolean transferReported = new AtomicBoolean(); + private final String requestId = UUID.randomUUID().toString(); InstallSnapshotResponseHandler() { this(false); @@ -666,6 +649,7 @@ private class InstallSnapshotResponseHandler implements StreamObserver l.onInstallSnapshotStart(requestId, notifyOnly)); } void addPending(InstallSnapshotRequestProto request) { @@ -722,32 +706,13 @@ void waitForResponse() { done.get(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); - final InterruptedIOException interrupted = - new InterruptedIOException("Interrupted while waiting for InstallSnapshot responses"); - interrupted.initCause(e); - notifyTransferFailure(interrupted); + notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); } catch (ExecutionException e) { - notifyTransferFailure(e); + notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); throw new IllegalStateException("Failed to complete " + name, e); } } - void notifyTransferSuccess() { - if (!isNotificationOnly && transferReported.compareAndSet(false, true)) { - notifyDataTransferSuccess(); - } - } - - void notifyTransferFailure(Throwable error) { - if (!isNotificationOnly && transferReported.compareAndSet(false, true)) { - notifyDataTransferFailure(error); - } - } - - private void notifyTransferFailure(InstallSnapshotReplyProto reply) { - notifyTransferFailure(new IOException("InstallSnapshot failed with result " + reply.getResult())); - } - void close() { done.complete(null); notifyLogAppender(); @@ -761,6 +726,7 @@ boolean hasAllResponse() { @Override public void onNext(InstallSnapshotReplyProto reply) { + notifyListener(l -> l.onInstallSnapshotReply(requestId, reply)); BatchLogger.print(BatchLogKey.INSTALL_SNAPSHOT_REPLY, name, suffix -> LOG.info("{}: received {} reply {} {}", this, replyState.isFirstReplyReceived() ? "a" : "the first", @@ -792,11 +758,9 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case NOT_LEADER: - notifyTransferFailure(reply); onFollowerTerm(reply.getTerm()); break; case CONF_MISMATCH: - notifyTransferFailure(reply); LOG.error("{}: CONF_MISMATCH ({}): Leader {} has it set to {} but follower {} has it set to {}", this, RaftServerConfigKeys.Log.Appender.INSTALL_SNAPSHOT_ENABLED_KEY, getServer().getId(), installSnapshotEnabled, getFollowerId(), !installSnapshotEnabled); @@ -812,7 +776,6 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case SNAPSHOT_UNAVAILABLE: - notifyTransferFailure(reply); BatchLogger.print(BatchLogKey.SNAPSHOT_UNAVAILABLE, name, suffix -> LOG.info("{}: Follower failed since the snapshot is unavailable {}", this, suffix)); getFollower().setAttemptedToInstallSnapshot(); @@ -820,12 +783,10 @@ public void onNext(InstallSnapshotReplyProto reply) { removePending(reply); break; case UNRECOGNIZED: - notifyTransferFailure(reply); LOG.error("{}: Reply result {}, {}", name, reply.getResult(), ServerStringUtils.toInstallSnapshotReplyString(reply)); break; case SNAPSHOT_EXPIRED: - notifyTransferFailure(reply); LOG.warn("{}: Follower failed since the request expired, {}", name, ServerStringUtils.toInstallSnapshotReplyString(reply)); default: @@ -835,8 +796,9 @@ public void onNext(InstallSnapshotReplyProto reply) { @Override public void onError(Throwable t) { - notifyTransferFailure(t); + notifyListener(l -> l.onInstallSnapshotEnd(requestId, t)); if (!isRunning()) { + close(); LOG.info("{} is stopped", GrpcLogAppender.this); return; } @@ -848,6 +810,7 @@ public void onError(Throwable t) { @Override public void onCompleted() { + notifyListener(l -> l.onInstallSnapshotEnd(requestId, null)); if (!isNotificationOnly || LOG.isDebugEnabled()) { LOG.info("{}: follower responded installSnapshot COMPLETED", this); } @@ -865,33 +828,32 @@ public String toString() { * Send installSnapshot request to Follower with a snapshot. * @param snapshot the snapshot to be sent to Follower */ - private void installSnapshot(SnapshotInfo snapshot) { + void installSnapshot(SnapshotInfo snapshot) { LOG.info("{}: followerNextIndex = {} but logStartIndex = {}, send snapshot {} to follower", this, getFollower().getNextIndex(), getRaftLog().getStartIndex(), snapshot); final InstallSnapshotResponseHandler responseHandler = new InstallSnapshotResponseHandler(); StreamObserver snapshotRequestObserver = null; - final String requestId = UUID.randomUUID().toString(); - boolean sentAllRequests = true; + final String requestId = responseHandler.requestId; try { snapshotRequestObserver = getClient().installSnapshot( getFollower().getName() + "-installSnapshot-" + requestId, installSnapshotStreamTimeout, maxOutstandingInstallSnapshots, responseHandler); for (InstallSnapshotRequestProto request : newInstallSnapshotRequests(requestId, snapshot)) { if (isRunning()) { + notifyListener(l -> l.onInstallSnapshotRequest(requestId, request)); snapshotRequestObserver.onNext(request); getFollower().updateLastRpcSendTime(false); responseHandler.addPending(request); } else { - sentAllRequests = false; break; } } snapshotRequestObserver.onCompleted(); grpcServerMetrics.onInstallSnapshot(); } catch (Exception e) { + notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); LOG.warn(this + ": failed to installSnapshot " + snapshot, e); - responseHandler.notifyTransferFailure(e); if (snapshotRequestObserver != null) { snapshotRequestObserver.onError(e); } @@ -899,16 +861,9 @@ private void installSnapshot(SnapshotInfo snapshot) { } responseHandler.waitForResponse(); - if (!sentAllRequests) { - responseHandler.notifyTransferFailure( - new IOException("InstallSnapshot stopped before all requests were sent")); - } else if (responseHandler.hasAllResponse()) { + if (responseHandler.hasAllResponse()) { getFollower().setSnapshotIndex(snapshot.getTermIndex().getIndex()); LOG.info("{}: installed snapshot {} successfully", this, snapshot); - responseHandler.notifyTransferSuccess(); - } else { - responseHandler.notifyTransferFailure( - new IOException("InstallSnapshot completed with pending responses")); } } @@ -932,11 +887,13 @@ private void notifyInstallSnapshot(TermIndex firstAvailable) { snapshotRequestObserver = getClient().installSnapshot(getFollower().getName() + "-notifyInstallSnapshot", requestTimeoutDuration, 0, responseHandler); + notifyListener(l -> l.onInstallSnapshotRequest(responseHandler.requestId, request)); snapshotRequestObserver.onNext(request); getFollower().updateLastRpcSendTime(false); responseHandler.addPending(request); snapshotRequestObserver.onCompleted(); } catch (Exception e) { + notifyListener(l -> l.onInstallSnapshotEnd(responseHandler.requestId, e)); GrpcUtil.warn(LOG, () -> this + ": Failed to notify follower to install snapshot.", e); if (snapshotRequestObserver != null) { snapshotRequestObserver.onError(e); @@ -955,25 +912,16 @@ static class AppendEntriesRequest { private final long callId; private final TermIndex previousLog; private final int entriesCount; - private final boolean containsStateMachineData; private final TermIndex firstEntry; private final TermIndex lastEntry; @SuppressWarnings({"squid:S3077"}) // Suppress volatile for generic type private volatile Timestamp sendTime; - AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, - GrpcServerMetrics grpcServerMetrics) { - this(proto, followerId, grpcServerMetrics, false); - } - - AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, - GrpcServerMetrics grpcServerMetrics, boolean checkStateMachineData) { + AppendEntriesRequest(AppendEntriesRequestProto proto, RaftPeerId followerId, GrpcServerMetrics grpcServerMetrics) { this.callId = proto.getServerRequest().getCallId(); this.previousLog = proto.hasPreviousLog()? TermIndex.valueOf(proto.getPreviousLog()): null; this.entriesCount = proto.getEntriesCount(); - this.containsStateMachineData = checkStateMachineData && proto.getEntriesList().stream() - .anyMatch(entry -> entry.hasStateMachineLogEntry()); this.firstEntry = entriesCount > 0? TermIndex.valueOf(proto.getEntries(0)): null; this.lastEntry = entriesCount > 0? TermIndex.valueOf(proto.getEntries(entriesCount - 1)): null; @@ -1005,7 +953,6 @@ void startRequestTimer() { void stopRequestTimer() { if (timerContext != null) { timerContext.stop(); - timerContext = null; } } @@ -1013,14 +960,6 @@ boolean isHeartbeat() { return entriesCount == 0; } - boolean containsStateMachineData() { - return containsStateMachineData; - } - - boolean isSent() { - return sendTime != null; - } - @Override public String toString() { return JavaUtils.getClassSimpleName(getClass()) @@ -1037,20 +976,9 @@ int logRequestsSize() { return logRequests.size(); } - List clear(boolean collectLogRequests) { - if (!collectLogRequests) { - logRequests.clear(); - heartbeats.clear(); - return Collections.emptyList(); - } - final List removed = new ArrayList<>(); - logRequests.forEach((callId, request) -> { - if (logRequests.remove(callId, request)) { - removed.add(request); - } - }); + void clear() { + logRequests.clear(); heartbeats.clear(); - return removed; } void put(AppendEntriesRequest request) { diff --git a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java deleted file mode 100644 index 2665de2863..0000000000 --- a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcDataTransferEventListener.java +++ /dev/null @@ -1,254 +0,0 @@ -/* - * 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.AppendEntriesRequestProto; -import org.apache.ratis.proto.RaftProtos.ReplicationLevel; -import org.apache.ratis.protocol.RaftClientReply; -import org.apache.ratis.protocol.RaftPeerId; -import org.apache.ratis.security.SecurityTestUtils; -import org.apache.ratis.server.RaftServer; -import org.apache.ratis.server.RaftServerConfigKeys; -import org.apache.ratis.server.impl.MiniRaftCluster; -import org.apache.ratis.server.impl.PeerChanges; -import org.apache.ratis.server.impl.RaftServerTestUtil; -import org.apache.ratis.statemachine.SnapshotInfo; -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.apache.ratis.util.SizeInBytes; -import org.junit.jupiter.api.Assertions; -import org.junit.jupiter.api.Test; -import org.junit.jupiter.api.Timeout; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.ValueSource; - -import javax.net.ssl.KeyManager; -import javax.net.ssl.TrustManager; - -import java.util.List; -import java.util.Set; -import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; -import java.util.stream.Collectors; - -public class TestGrpcDataTransferEventListener extends BaseTest { - private static RaftProperties newProperties() { - final RaftProperties properties = new RaftProperties(); - properties.setClass(MiniRaftCluster.STATEMACHINE_CLASS_KEY, - SimpleStateMachine4Testing.class, StateMachine.class); - return properties; - } - - private static Parameters newParameters(boolean tls, - ConcurrentLinkedQueue events) throws Exception { - final Parameters parameters = new Parameters(); - GrpcConfigKeys.Server.setDataTransferEventConsumer(parameters, events::add); - if (tls) { - final KeyManager serverKeyManager = - SecurityTestUtils.getKeyManager(SecurityTestUtils::getServerKeyStore); - final TrustManager serverTrustManager = - SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore); - final KeyManager clientKeyManager = - SecurityTestUtils.getKeyManager(SecurityTestUtils::getClientKeyStore); - final TrustManager clientTrustManager = - SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore); - - GrpcConfigKeys.Server.setTlsConf(parameters, - new GrpcTlsConfig(serverKeyManager, serverTrustManager, true)); - final GrpcTlsConfig clientConfig = - new GrpcTlsConfig(clientKeyManager, clientTrustManager, true); - GrpcConfigKeys.Admin.setTlsConf(parameters, clientConfig); - GrpcConfigKeys.Client.setTlsConf(parameters, clientConfig); - } - return parameters; - } - - @ParameterizedTest - @ValueSource(booleans = {false, true}) - @Timeout(value = 60, unit = TimeUnit.SECONDS) - public void testSuccessfulAppendEntries(boolean tls) throws Exception { - final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); - final Parameters parameters = newParameters(tls, events); - final String[] ids = MiniRaftCluster.generateIds(3, 10); - - try (MiniRaftClusterWithGrpc cluster = - new MiniRaftClusterWithGrpc(ids, new String[0], newProperties(), parameters)) { - cluster.start(); - final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); - final RaftPeerId leaderId = leader.getId(); - events.clear(); - - try (RaftClient client = cluster.createClient(leaderId)) { - final RaftClientReply write = - client.io().send(new RaftTestUtil.SimpleMessage("data-transfer")); - Assertions.assertTrue(write.isSuccess()); - Assertions.assertTrue( - client.io().watch(write.getLogIndex(), ReplicationLevel.ALL_COMMITTED).isSuccess()); - } - - final Set followers = cluster.getFollowers().stream() - .map(RaftServer.Division::getId) - .collect(Collectors.toSet()); - JavaUtils.attempt(() -> { - final Set destinations = successfulEvents(events, leaderId).stream() - .map(GrpcDataTransferEvent::getDestination) - .collect(Collectors.toSet()); - Assertions.assertEquals(followers, destinations); - }, 10, HUNDRED_MILLIS, "data transfer events", LOG); - - final GrpcDataTransferEvent.ProtectionMethod expected = tls - ? GrpcDataTransferEvent.ProtectionMethod.TLS - : GrpcDataTransferEvent.ProtectionMethod.NONE; - successfulEvents(events, leaderId).forEach(event -> { - Assertions.assertNotNull(event.getTimestamp()); - Assertions.assertEquals(expected, event.getProtectionMethod()); - Assertions.assertNull(event.getError()); - }); - } - } - - private static List successfulEvents( - ConcurrentLinkedQueue events, RaftPeerId source) { - return events.stream() - .filter(event -> event.getSource().equals(source)) - .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.SUCCESS) - .collect(Collectors.toList()); - } - - @Test - @Timeout(value = 120, unit = TimeUnit.SECONDS) - public void testSuccessfulInstallSnapshot() throws Exception { - final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); - final Parameters parameters = newParameters(false, events); - final RaftProperties properties = newProperties(); - RaftServerConfigKeys.Snapshot.setAutoTriggerEnabled(properties, true); - RaftServerConfigKeys.Snapshot.setAutoTriggerThreshold(properties, 64); - RaftServerConfigKeys.Log.setPurgeGap(properties, 8); - RaftServerConfigKeys.Log.Appender.setSnapshotChunkSizeMax(properties, SizeInBytes.ONE_KB); - RaftServerConfigKeys.LeaderElection.setMemberMajorityAdd(properties, true); - - final String[] ids = MiniRaftCluster.generateIds(1, 30); - try (MiniRaftClusterWithGrpc cluster = - new MiniRaftClusterWithGrpc(ids, new String[0], properties, parameters)) { - cluster.start(); - final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); - final RaftPeerId leaderId = leader.getId(); - - try (RaftClient client = cluster.createClient(leaderId)) { - for (int i = 0; i < 127; i++) { - Assertions.assertTrue( - client.io().send(new RaftTestUtil.SimpleMessage("snapshot-" + i)).isSuccess()); - } - Assertions.assertTrue(client.getSnapshotManagementApi(leaderId).create(3000).isSuccess()); - } - final SnapshotInfo leaderSnapshot = leader.getStateMachine().getLatestSnapshot(); - Assertions.assertNotNull(leaderSnapshot); - events.clear(); - - final PeerChanges change = cluster.addNewPeers(1, true); - final RaftPeerId addedPeer = change.getAddedPeers().get(0).getId(); - cluster.setConfiguration(change.getPeersInNewConf()); - RaftServerTestUtil.waitAndCheckNewConf(cluster, change.getPeersInNewConf(), 0, null); - - JavaUtils.attempt(() -> { - final SnapshotInfo installed = - cluster.getDivision(addedPeer).getStateMachine().getLatestSnapshot(); - Assertions.assertNotNull(installed); - Assertions.assertTrue(successfulEvents(events, leaderId).stream() - .anyMatch(event -> event.getDestination().equals(addedPeer))); - }, 20, ONE_SECOND, "snapshot data transfer event", LOG); - } - } - - @Test - @Timeout(value = 60, unit = TimeUnit.SECONDS) - public void testAppendEntriesFailureAndConsumerException() throws Exception { - final ConcurrentLinkedQueue events = new ConcurrentLinkedQueue<>(); - final AtomicInteger consumerCalls = new AtomicInteger(); - final Parameters parameters = new Parameters(); - GrpcConfigKeys.Server.setDataTransferEventConsumer(parameters, event -> { - events.add(event); - consumerCalls.incrementAndGet(); - if (event.getResult() == GrpcDataTransferEvent.Result.FAILURE) { - throw new IllegalStateException("Injected consumer failure"); - } - }); - - final String[] ids = MiniRaftCluster.generateIds(2, 20); - try (MiniRaftClusterWithGrpc cluster = - new MiniRaftClusterWithGrpc(ids, new String[0], newProperties(), parameters)) { - cluster.start(); - final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); - final RaftPeerId leaderId = leader.getId(); - events.clear(); - consumerCalls.set(0); - - final AtomicBoolean injectFailure = new AtomicBoolean(true); - try { - CodeInjectionForTesting.put( - GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST, (localId, remoteId, args) -> { - if (leaderId.equals(localId) - && args.length > 0 - && args[0] instanceof AppendEntriesRequestProto) { - final AppendEntriesRequestProto request = (AppendEntriesRequestProto) args[0]; - final boolean containsStateMachineData = request.getEntriesList().stream() - .anyMatch(entry -> entry.hasStateMachineLogEntry()); - if (containsStateMachineData && injectFailure.compareAndSet(true, false)) { - throw new IllegalStateException("Injected AppendEntries failure"); - } - } - return false; - }); - - try (RaftClient client = cluster.createClient(leaderId)) { - Assertions.assertTrue( - client.io().send(new RaftTestUtil.SimpleMessage("retry-after-failure")).isSuccess()); - } - - JavaUtils.attempt(() -> Assertions.assertTrue(events.stream() - .anyMatch(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE)), - 10, HUNDRED_MILLIS, "failed data transfer event", LOG); - } finally { - CodeInjectionForTesting.remove(GrpcServicesImpl.GRPC_SEND_SERVER_REQUEST); - } - - final GrpcDataTransferEvent failure = events.stream() - .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE) - .findFirst() - .orElseThrow(AssertionError::new); - Assertions.assertEquals(leaderId, failure.getSource()); - Assertions.assertEquals(GrpcDataTransferEvent.ProtectionMethod.NONE, - failure.getProtectionMethod()); - Assertions.assertTrue(failure.getError() instanceof IllegalStateException); - Assertions.assertEquals(1, events.stream() - .filter(event -> event.getResult() == GrpcDataTransferEvent.Result.FAILURE) - .count()); - Assertions.assertTrue(consumerCalls.get() > 0); - } - } -} 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..f7437353fc --- /dev/null +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java @@ -0,0 +1,208 @@ +/* + * 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.InstallSnapshotReplyProto; +import org.apache.ratis.proto.RaftProtos.InstallSnapshotRequestProto; +import org.apache.ratis.proto.RaftProtos.ReplicationLevel; +import org.apache.ratis.protocol.RaftClientReply; +import org.apache.ratis.protocol.RaftPeerId; +import org.apache.ratis.security.SecurityTestUtils; +import org.apache.ratis.server.RaftServer; +import org.apache.ratis.server.RaftServerConfigKeys; +import org.apache.ratis.server.impl.MiniRaftCluster; +import org.apache.ratis.server.impl.PeerChanges; +import org.apache.ratis.server.impl.RaftServerTestUtil; +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.apache.ratis.util.SizeInBytes; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +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.concurrent.atomic.AtomicInteger; +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; + } + + private static void setTls(Parameters parameters) throws Exception { + GrpcConfigKeys.Server.setTlsConf(parameters, new GrpcTlsConfig( + SecurityTestUtils.getKeyManager(SecurityTestUtils::getServerKeyStore), + SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore), true)); + final GrpcTlsConfig client = new GrpcTlsConfig( + SecurityTestUtils.getKeyManager(SecurityTestUtils::getClientKeyStore), + SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore), true); + GrpcConfigKeys.Admin.setTlsConf(parameters, client); + GrpcConfigKeys.Client.setTlsConf(parameters, client); + } + + @ParameterizedTest + @ValueSource(booleans = {false, true}) + @Timeout(value = 60, unit = TimeUnit.SECONDS) + public void testAppendEntries(boolean tls) 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() { + private final Set pending = ConcurrentHashMap.newKeySet(); + + @Override + public void onAppendEntriesRequest(AppendEntriesRequestProto request) { + if (request.getEntriesList().stream().anyMatch(entry -> entry.hasStateMachineLogEntry())) { + pending.add(request.getServerRequest().getCallId()); + } + } + + @Override + public void onAppendEntriesReply(AppendEntriesReplyProto reply) { + if (pending.remove(reply.getServerReply().getCallId())) { + destinations.add(destination.getId()); + } + } + + @Override + public void onAppendEntriesFailure(long callId, Throwable error) { + if (pending.remove(callId)) { + failures.add(error); + throw new IllegalStateException("Injected listener failure"); + } + } + }); + if (tls) { + setTls(parameters); + } + 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()); + } + } + + @Test + @Timeout(value = 120, unit = TimeUnit.SECONDS) + public void testSnapshotCallbacks() throws Exception { + final Parameters parameters = new Parameters(); + final AtomicInteger chunks = new AtomicInteger(); + final AtomicInteger replies = new AtomicInteger(); + final Set started = ConcurrentHashMap.newKeySet(); + final Set completed = ConcurrentHashMap.newKeySet(); + GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (source, destination) -> + new GrpcLogAppenderListener() { + @Override + public void onInstallSnapshotStart(String requestId, boolean notificationOnly) { + if (!notificationOnly) { + started.add(requestId); + } + } + + @Override + public void onInstallSnapshotRequest(String requestId, InstallSnapshotRequestProto request) { + if (started.contains(requestId)) { + chunks.incrementAndGet(); + } + } + + @Override + public void onInstallSnapshotReply(String requestId, InstallSnapshotReplyProto reply) { + if (started.contains(requestId)) { + replies.incrementAndGet(); + } + } + + @Override + public void onInstallSnapshotEnd(String requestId, Throwable error) { + if (error == null && started.contains(requestId)) { + completed.add(requestId); + } + } + }); + final RaftProperties properties = newProperties(); + RaftServerConfigKeys.Snapshot.setAutoTriggerEnabled(properties, true); + RaftServerConfigKeys.Snapshot.setAutoTriggerThreshold(properties, 64); + RaftServerConfigKeys.Log.setPurgeGap(properties, 8); + RaftServerConfigKeys.Log.Appender.setSnapshotChunkSizeMax(properties, SizeInBytes.ONE_KB); + RaftServerConfigKeys.LeaderElection.setMemberMajorityAdd(properties, true); + try (MiniRaftClusterWithGrpc cluster = new MiniRaftClusterWithGrpc( + MiniRaftCluster.generateIds(1, 30), new String[0], properties, parameters)) { + cluster.start(); + final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); + try (RaftClient client = cluster.createClient(leader.getId())) { + for (int i = 0; i < 127; i++) { + Assertions.assertTrue(client.io().send(new RaftTestUtil.SimpleMessage("snapshot-" + i)).isSuccess()); + } + Assertions.assertTrue(client.getSnapshotManagementApi(leader.getId()).create(3000).isSuccess()); + } + final PeerChanges change = cluster.addNewPeers(1, true); + cluster.setConfiguration(change.getPeersInNewConf()); + RaftServerTestUtil.waitAndCheckNewConf(cluster, change.getPeersInNewConf(), 0, null); + JavaUtils.attempt(() -> { + Assertions.assertFalse(completed.isEmpty()); + Assertions.assertTrue(chunks.get() > 1); + Assertions.assertEquals(chunks.get(), replies.get()); + }, 20, ONE_SECOND, "snapshot callbacks", LOG); + } + } +} 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..f015287223 --- /dev/null +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/server/TestGrpcLogAppenderCallbacks.java @@ -0,0 +1,248 @@ +/* + * 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.protocol.exceptions.TimeoutIOException; +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.statemachine.SnapshotInfo; +import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +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.ArgumentCaptor; +import org.mockito.InOrder; + +import java.io.IOException; +import java.util.concurrent.CompletableFuture; +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.ArgumentMatchers.isA; +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.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 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); + 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(listener).onAppendEntriesReset(error); + Assertions.assertDoesNotThrow(() -> responses.onError(error)); + verify(listener).onAppendEntriesReset(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"); + responses.onError(error); + verify(listener).onAppendEntriesReset(error); + } + + @Test + public void testCompletionAfterStop() throws Exception { + addPending(1); + appender.stopAsync().get(5, TimeUnit.SECONDS); + clearInvocations(client, leaderState); + responses.onCompleted(); + verify(listener).onAppendEntriesReset(isA(IOException.class)); + 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(listener); + order.verify(listener).onAppendEntriesRequest(request); + order.verify(listener).onAppendEntriesFailure(1, error); + Assertions.assertFalse(appender.hasPendingDataRequests()); + } + + @Test + public void testTimeout() { + addPending(1); + appender.timeoutAppendRequest(1, false); + appender.timeoutAppendRequest(1, false); + verify(listener).onAppendEntriesFailure(eq(1L), isA(TimeoutIOException.class)); + } + + @ParameterizedTest + @EnumSource(value = AppendResult.class, names = {"SUCCESS", "NOT_LEADER", "INCONSISTENCY"}) + public void testMatchedReply(AppendResult result) { + addPending(1); + final AppendEntriesReplyProto reply = AppendEntriesReplyProto.newBuilder() + .setServerReply(RaftRpcReplyProto.newBuilder().setCallId(1)).setResult(result).build(); + responses.onNext(reply); + responses.onNext(reply); + verify(listener).onAppendEntriesReply(reply); + appender.timeoutAppendRequest(1, false); + verify(listener, never()).onAppendEntriesFailure(eq(1L), any()); + } + + @Test + public void testFactoryFailureIsIsolated() throws Exception { + final Parameters parameters = new Parameters(); + GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (member, peer) -> { + Assertions.assertEquals(server.getMemberId(), member); + Assertions.assertEquals(follower.getPeer(), peer); + throw new IllegalStateException("Injected factory failure"); + }); + final LogAppender created = new GrpcFactory(parameters).newLogAppender(server, leaderState, follower); + Assertions.assertNotNull(created); + created.stopAsync().get(5, TimeUnit.SECONDS); + } + + @Test + public void testSnapshotStreamCreationFailure() throws Exception { + final IOException error = new IOException("Cannot create snapshot connection"); + when(serverRpc.getProxies().getProxy(destination)).thenThrow(error); + appender.installSnapshot(mock(SnapshotInfo.class, RETURNS_DEEP_STUBS)); + final ArgumentCaptor requestId = ArgumentCaptor.forClass(String.class); + final InOrder order = inOrder(listener); + order.verify(listener).onInstallSnapshotStart(requestId.capture(), eq(false)); + order.verify(listener).onInstallSnapshotEnd(requestId.getValue(), error); + verify(listener, never()).onInstallSnapshotRequest(anyString(), any()); + } + + @Test + public void testSnapshotErrorAfterStopUnblocksWaiter() throws Exception { + final GrpcLogAppender.InstallSnapshotResponseHandler snapshot = appender.new InstallSnapshotResponseHandler(); + final ArgumentCaptor requestId = ArgumentCaptor.forClass(String.class); + verify(listener).onInstallSnapshotStart(requestId.capture(), eq(false)); + final CompletableFuture done = (CompletableFuture) RaftTestUtil.getDeclaredField(snapshot, "done"); + appender.stopAsync().get(5, TimeUnit.SECONDS); + final IOException error = new IOException("Snapshot stream failed after stop"); + doThrow(new IllegalStateException("Listener failure")).when(listener).onInstallSnapshotEnd(anyString(), any()); + snapshot.onError(error); + Assertions.assertTrue(done.isDone()); + verify(listener).onInstallSnapshotEnd(requestId.getValue(), error); + } +} From 476231a9f9bd3a1bcc9bcffd8be79dd1026ca97b Mon Sep 17 00:00:00 2001 From: HTHou Date: Mon, 7 Sep 2026 12:14:39 +0800 Subject: [PATCH 3/6] RATIS-2681. Simplify log appender listener tests --- .../grpc/TestGrpcLogAppenderListener.java | 22 ++----------------- 1 file changed, 2 insertions(+), 20 deletions(-) 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 index f7437353fc..3a62f32b47 100644 --- a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java @@ -30,7 +30,6 @@ import org.apache.ratis.proto.RaftProtos.ReplicationLevel; import org.apache.ratis.protocol.RaftClientReply; import org.apache.ratis.protocol.RaftPeerId; -import org.apache.ratis.security.SecurityTestUtils; import org.apache.ratis.server.RaftServer; import org.apache.ratis.server.RaftServerConfigKeys; import org.apache.ratis.server.impl.MiniRaftCluster; @@ -44,8 +43,6 @@ import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.ValueSource; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @@ -63,21 +60,9 @@ private static RaftProperties newProperties() { return properties; } - private static void setTls(Parameters parameters) throws Exception { - GrpcConfigKeys.Server.setTlsConf(parameters, new GrpcTlsConfig( - SecurityTestUtils.getKeyManager(SecurityTestUtils::getServerKeyStore), - SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore), true)); - final GrpcTlsConfig client = new GrpcTlsConfig( - SecurityTestUtils.getKeyManager(SecurityTestUtils::getClientKeyStore), - SecurityTestUtils.getTrustManager(SecurityTestUtils::getTrustStore), true); - GrpcConfigKeys.Admin.setTlsConf(parameters, client); - GrpcConfigKeys.Client.setTlsConf(parameters, client); - } - - @ParameterizedTest - @ValueSource(booleans = {false, true}) + @Test @Timeout(value = 60, unit = TimeUnit.SECONDS) - public void testAppendEntries(boolean tls) throws Exception { + public void testAppendEntries() throws Exception { final Parameters parameters = new Parameters(); final Set destinations = ConcurrentHashMap.newKeySet(); final ConcurrentLinkedQueue failures = new ConcurrentLinkedQueue<>(); @@ -108,9 +93,6 @@ public void onAppendEntriesFailure(long callId, Throwable error) { } } }); - if (tls) { - setTls(parameters); - } try (MiniRaftClusterWithGrpc cluster = new MiniRaftClusterWithGrpc( MiniRaftCluster.generateIds(3, 10), new String[0], newProperties(), parameters)) { cluster.start(); From 80b96f6d6419e98d7de71eb7536489606e655ae6 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 8 Sep 2026 09:41:15 +0800 Subject: [PATCH 4/6] RATIS-2681. Scope log appender callbacks to AppendEntries --- .../dev-support/findbugsExcludeFile.xml | 6 + .../org/apache/ratis/grpc/GrpcConfigKeys.java | 5 +- .../org/apache/ratis/grpc/GrpcFactory.java | 2 +- .../ratis/grpc/GrpcLogAppenderListener.java | 48 ++++---- .../ratis/grpc/server/GrpcLogAppender.java | 103 ++++++++-------- .../grpc/TestGrpcLogAppenderListener.java | 114 ++++-------------- .../server/TestGrpcLogAppenderCallbacks.java | 89 +++++++------- 7 files changed, 159 insertions(+), 208 deletions(-) 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 6a0056e1e0..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 @@ -305,14 +305,15 @@ static void setCredentials(Parameters parameters, ServerCredentials credentials) } 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, GrpcLogAppenderListener.Factory.class); + 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, GrpcLogAppenderListener.Factory.class); + parameters.put(LOG_APPENDER_LISTENER_FACTORY_PARAMETER, factory, LOG_APPENDER_LISTENER_FACTORY_CLASS); } String TLS_CONF_PARAMETER = PREFIX + ".tls.conf"; 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 ea9610a254..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 @@ -136,7 +136,7 @@ public LogAppender newLogAppender(RaftServer.Division server, LeaderState state, try { listener = logAppenderListenerFactory.create(server.getMemberId(), f.getPeer()); } catch (Throwable t) { - LOG.warn("Failed to create gRPC log appender listener", t); + LOG.warn("{}->{}: Failed to create gRPC log appender listener", server.getMemberId(), f.getId(), t); } } return new GrpcLogAppender(server, state, f, listener); 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 index 5f99753c59..4310cc6b0f 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java @@ -19,8 +19,6 @@ import org.apache.ratis.proto.RaftProtos.AppendEntriesReplyProto; import org.apache.ratis.proto.RaftProtos.AppendEntriesRequestProto; -import org.apache.ratis.proto.RaftProtos.InstallSnapshotReplyProto; -import org.apache.ratis.proto.RaftProtos.InstallSnapshotRequestProto; import org.apache.ratis.protocol.RaftGroupMemberId; import org.apache.ratis.protocol.RaftPeer; @@ -29,8 +27,8 @@ * 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, - * filtering messages and aggregating snapshot chunks. - * Append request, matched-reply and reset callbacks are serialized per appender. Failure + * and filtering messages. This interface currently observes AppendEntries, not InstallSnapshot. + * Append request, reply and reset callbacks are serialized per appender. Failure * callbacks may race with these callbacks. No exactly-once terminal notification is guaranteed. */ public interface GrpcLogAppenderListener { @@ -41,30 +39,38 @@ interface Factory { GrpcLogAppenderListener create(RaftGroupMemberId source, RaftPeer destination); } - /** An append attempt is registered, before establishing or writing its stream. */ - default void onAppendEntriesRequest(AppendEntriesRequestProto request) { } + /** @return the AppendEntries listener, or null to disable its callbacks. Called once per appender. */ + default AppendEntries appendEntries() { + return null; + } - /** A response was matched to a pending append request. */ - default void onAppendEntriesReply(AppendEntriesReplyProto reply) { } + /** 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 local send error or request timeout occurred; a later stream notification may follow. */ - default void onAppendEntriesFailure(long callId, Throwable error) { } + /** A response was received, possibly after its request timed out or was invalidated. */ + default void onReply(AppendEntriesReplyProto reply) { } - /** Pending append attempts are invalidated by a reset or by a stopped response stream. */ - default void onAppendEntriesReset(Throwable error) { } + /** A local send error occurred; a later stream notification may follow. */ + default void onFailure(long callId, Throwable error) { } - /** A snapshot attempt starts before stream creation, including notification-only attempts. */ - default void onInstallSnapshotStart(String requestId, boolean notificationOnly) { } + /** A pending request timed out. */ + default void onTimeout(long callId) { } - /** A snapshot chunk or notification is about to be sent. */ - default void onInstallSnapshotRequest(String requestId, InstallSnapshotRequestProto request) { } + /** A response stream completed, including after the appender stopped. */ + default void onCompleted() { } - /** A snapshot response was received; this is not necessarily a terminal response. */ - default void onInstallSnapshotReply(String requestId, InstallSnapshotReplyProto reply) { } + /** A response stream failed, including after the appender stopped. */ + default void onError(Throwable error) { } + } /** - * The stream completed (null error), failed, or its sender was interrupted. Completion alone - * does not imply that all chunks were sent or acknowledged. Racing notifications may repeat. + * Pending attempts are invalidated before resetting the client or reconciling an inconsistent log. + * The reason is diagnostic text, not a stable identifier; error may be null. */ - default void onInstallSnapshotEnd(String requestId, Throwable error) { } + default void onReset(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 2413a48db2..c5d027a2b1 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 @@ -25,7 +25,6 @@ import org.apache.ratis.metrics.Timekeeper; import org.apache.ratis.proto.RaftProtos.InstallSnapshotResult; import org.apache.ratis.protocol.RaftPeerId; -import org.apache.ratis.protocol.exceptions.TimeoutIOException; import org.apache.ratis.retry.RetryPolicy; import org.apache.ratis.server.RaftServer; import org.apache.ratis.server.RaftServerConfigKeys; @@ -167,20 +166,18 @@ synchronized int process(Event event) { 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) { - this(server, leaderState, f, null); - } - 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"); @@ -209,17 +206,35 @@ private void notifyListener(Consumer notification) { try { notification.accept(listener); } catch (Throwable t) { - LOG.warn("gRPC log appender listener threw an exception", t); + LOG.warn("{}: gRPC log appender listener threw an exception", this, t); } } } - private void notifyReset(Throwable error) { - try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - notifyListener(l -> l.onAppendEntriesReset(error)); + 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); + } } } + /** Invoked with the appender write lock held. */ + private void notifyReset(String reason, Throwable error) { + notifyListener(l -> l.onReset(reason, error)); + } + @Override public GrpcServicesImpl getServerRpc() { return (GrpcServicesImpl)super.getServerRpc(); @@ -231,7 +246,7 @@ private GrpcServerProtocolClient getClient() throws IOException { private void resetClient(AppendEntriesRequest request, Event event, Throwable error) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - notifyListener(l -> l.onAppendEntriesReset(error)); + notifyReset("resetClient: " + event, error); getClient().resetConnectBackoff(); if (appendLogRequestObserver != null) { appendLogRequestObserver.stop(); @@ -285,16 +300,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() { @@ -429,7 +449,7 @@ void appendLog(boolean heartbeat) throws IOException { } request = new AppendEntriesRequest(pending, getFollowerId(), grpcServerMetrics); pendingRequests.put(request); - notifyListener(l -> l.onAppendEntriesRequest(pending)); + notifyAppendEntriesListener(l -> l.onRequest(pending)); increaseNextIndex(pending); if (appendLogRequestObserver == null) { appendLogRequestObserver = new StreamObservers( @@ -443,10 +463,6 @@ void appendLog(boolean heartbeat) throws IOException { } if (isRunning()) { sendRequest(request, pending); - } else { - final long cid = request.getCallId(); - pendingRequests.remove(cid, request.isHeartbeat()); - notifyListener(l -> l.onAppendEntriesFailure(cid, new IOException("Log appender stopped before send"))); } } catch (IOException | RuntimeException e) { if (request != null) { @@ -455,7 +471,7 @@ void appendLog(boolean heartbeat) throws IOException { if (failed != null) { failed.stopRequestTimer(); } - notifyListener(l -> l.onAppendEntriesFailure(cid, e)); + notifyAppendEntriesListener(l -> l.onFailure(cid, e)); } throw e; } @@ -497,8 +513,7 @@ void timeoutAppendRequest(long cid, boolean heartbeat) { this, heartbeat ? "HEARTBEAT " : "", errorCount, pending); grpcServerMetrics.onRequestTimeout(getFollowerId().toString(), heartbeat); pending.stopRequestTimer(); - notifyListener(l -> l.onAppendEntriesFailure(cid, - new TimeoutIOException("Timed out appendEntries request " + pending))); + notifyAppendEntriesListener(l -> l.onTimeout(cid)); } } @@ -534,9 +549,7 @@ public void onNext(AppendEntriesReplyProto reply) { final AppendEntriesRequest request; try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { request = pendingRequests.remove(reply); - if (request != null) { - notifyListener(l -> l.onAppendEntriesReply(reply)); - } + notifyAppendEntriesListener(l -> l.onReply(reply)); } if (request != null) { request.stopRequestTimer(); // Update completion time @@ -597,8 +610,8 @@ private void onNextImpl(AppendEntriesRequest request, AppendEntriesReplyProto re */ @Override public void onError(Throwable t) { + notifyAppendEntriesListener(l -> l.onError(t)); if (!isRunning()) { - notifyReset(t); LOG.info("{} is already stopped", GrpcLogAppender.this); return; } @@ -612,13 +625,12 @@ public void onError(Throwable t) { @Override public void onCompleted() { + notifyAppendEntriesListener(GrpcLogAppenderListener.AppendEntries::onCompleted); LOG.info("{}: follower responses appendEntries COMPLETED", this); - final IOException error = new IOException("AppendEntries response stream completed"); if (!isRunning()) { - notifyReset(error); return; } - resetClient(null, Event.COMPLETE, error); + resetClient(null, Event.COMPLETE, null); } @Override @@ -629,18 +641,17 @@ public String toString() { private void updateNextIndex(long replyNextIndex) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - notifyListener(l -> l.onAppendEntriesReset(new IOException("AppendEntries invalidated by inconsistency"))); + notifyReset("AppendEntries INCONSISTENCY", null); pendingRequests.clear(); getFollower().setNextIndex(replyNextIndex); } } - class InstallSnapshotResponseHandler implements StreamObserver { + private class InstallSnapshotResponseHandler implements StreamObserver { private final String name; private final Queue pending = new LinkedList<>(); private final CompletableFuture done = new CompletableFuture<>(); private final boolean isNotificationOnly; - private final String requestId = UUID.randomUUID().toString(); InstallSnapshotResponseHandler() { this(false); @@ -649,7 +660,6 @@ class InstallSnapshotResponseHandler implements StreamObserver l.onInstallSnapshotStart(requestId, notifyOnly)); } void addPending(InstallSnapshotRequestProto request) { @@ -706,9 +716,7 @@ void waitForResponse() { done.get(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); - notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); } catch (ExecutionException e) { - notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); throw new IllegalStateException("Failed to complete " + name, e); } } @@ -726,7 +734,6 @@ boolean hasAllResponse() { @Override public void onNext(InstallSnapshotReplyProto reply) { - notifyListener(l -> l.onInstallSnapshotReply(requestId, reply)); BatchLogger.print(BatchLogKey.INSTALL_SNAPSHOT_REPLY, name, suffix -> LOG.info("{}: received {} reply {} {}", this, replyState.isFirstReplyReceived() ? "a" : "the first", @@ -796,9 +803,7 @@ public void onNext(InstallSnapshotReplyProto reply) { @Override public void onError(Throwable t) { - notifyListener(l -> l.onInstallSnapshotEnd(requestId, t)); if (!isRunning()) { - close(); LOG.info("{} is stopped", GrpcLogAppender.this); return; } @@ -810,7 +815,6 @@ public void onError(Throwable t) { @Override public void onCompleted() { - notifyListener(l -> l.onInstallSnapshotEnd(requestId, null)); if (!isNotificationOnly || LOG.isDebugEnabled()) { LOG.info("{}: follower responded installSnapshot COMPLETED", this); } @@ -828,20 +832,19 @@ public String toString() { * Send installSnapshot request to Follower with a snapshot. * @param snapshot the snapshot to be sent to Follower */ - void installSnapshot(SnapshotInfo snapshot) { + private void installSnapshot(SnapshotInfo snapshot) { LOG.info("{}: followerNextIndex = {} but logStartIndex = {}, send snapshot {} to follower", this, getFollower().getNextIndex(), getRaftLog().getStartIndex(), snapshot); final InstallSnapshotResponseHandler responseHandler = new InstallSnapshotResponseHandler(); StreamObserver snapshotRequestObserver = null; - final String requestId = responseHandler.requestId; + final String requestId = UUID.randomUUID().toString(); try { snapshotRequestObserver = getClient().installSnapshot( getFollower().getName() + "-installSnapshot-" + requestId, installSnapshotStreamTimeout, maxOutstandingInstallSnapshots, responseHandler); for (InstallSnapshotRequestProto request : newInstallSnapshotRequests(requestId, snapshot)) { if (isRunning()) { - notifyListener(l -> l.onInstallSnapshotRequest(requestId, request)); snapshotRequestObserver.onNext(request); getFollower().updateLastRpcSendTime(false); responseHandler.addPending(request); @@ -852,7 +855,6 @@ void installSnapshot(SnapshotInfo snapshot) { snapshotRequestObserver.onCompleted(); grpcServerMetrics.onInstallSnapshot(); } catch (Exception e) { - notifyListener(l -> l.onInstallSnapshotEnd(requestId, e)); LOG.warn(this + ": failed to installSnapshot " + snapshot, e); if (snapshotRequestObserver != null) { snapshotRequestObserver.onError(e); @@ -887,13 +889,11 @@ private void notifyInstallSnapshot(TermIndex firstAvailable) { snapshotRequestObserver = getClient().installSnapshot(getFollower().getName() + "-notifyInstallSnapshot", requestTimeoutDuration, 0, responseHandler); - notifyListener(l -> l.onInstallSnapshotRequest(responseHandler.requestId, request)); snapshotRequestObserver.onNext(request); getFollower().updateLastRpcSendTime(false); responseHandler.addPending(request); snapshotRequestObserver.onCompleted(); } catch (Exception e) { - notifyListener(l -> l.onInstallSnapshotEnd(responseHandler.requestId, e)); GrpcUtil.warn(LOG, () -> this + ": Failed to notify follower to install snapshot.", e); if (snapshotRequestObserver != null) { snapshotRequestObserver.onError(e); @@ -951,8 +951,9 @@ void startRequestTimer() { } void stopRequestTimer() { - if (timerContext != null) { - timerContext.stop(); + final Timekeeper.Context context = timerContext; + if (context != null) { + context.stop(); } } 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 index 3a62f32b47..d837882ccd 100644 --- a/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java +++ b/ratis-test/src/test/java/org/apache/ratis/grpc/TestGrpcLogAppenderListener.java @@ -25,21 +25,15 @@ 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.InstallSnapshotReplyProto; -import org.apache.ratis.proto.RaftProtos.InstallSnapshotRequestProto; 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.RaftServerConfigKeys; import org.apache.ratis.server.impl.MiniRaftCluster; -import org.apache.ratis.server.impl.PeerChanges; -import org.apache.ratis.server.impl.RaftServerTestUtil; 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.apache.ratis.util.SizeInBytes; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; @@ -49,7 +43,6 @@ import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; public class TestGrpcLogAppenderListener extends BaseTest { @@ -69,28 +62,33 @@ public void testAppendEntries() throws Exception { final AtomicBoolean injectFailure = new AtomicBoolean(true); GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (source, destination) -> new GrpcLogAppenderListener() { - private final Set pending = ConcurrentHashMap.newKeySet(); - @Override - public void onAppendEntriesRequest(AppendEntriesRequestProto request) { - if (request.getEntriesList().stream().anyMatch(entry -> entry.hasStateMachineLogEntry())) { - pending.add(request.getServerRequest().getCallId()); - } - } + public AppendEntries appendEntries() { + return new AppendEntries() { + private final Set pending = ConcurrentHashMap.newKeySet(); - @Override - public void onAppendEntriesReply(AppendEntriesReplyProto reply) { - if (pending.remove(reply.getServerReply().getCallId())) { - destinations.add(destination.getId()); - } - } + @Override + public void onRequest(AppendEntriesRequestProto request) { + if (request.getEntriesList().stream().anyMatch(entry -> entry.hasStateMachineLogEntry())) { + pending.add(request.getServerRequest().getCallId()); + } + } - @Override - public void onAppendEntriesFailure(long callId, Throwable error) { - if (pending.remove(callId)) { - failures.add(error); - throw new IllegalStateException("Injected listener failure"); - } + @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( @@ -123,68 +121,4 @@ public void onAppendEntriesFailure(long callId, Throwable error) { } } - @Test - @Timeout(value = 120, unit = TimeUnit.SECONDS) - public void testSnapshotCallbacks() throws Exception { - final Parameters parameters = new Parameters(); - final AtomicInteger chunks = new AtomicInteger(); - final AtomicInteger replies = new AtomicInteger(); - final Set started = ConcurrentHashMap.newKeySet(); - final Set completed = ConcurrentHashMap.newKeySet(); - GrpcConfigKeys.Server.setLogAppenderListenerFactory(parameters, (source, destination) -> - new GrpcLogAppenderListener() { - @Override - public void onInstallSnapshotStart(String requestId, boolean notificationOnly) { - if (!notificationOnly) { - started.add(requestId); - } - } - - @Override - public void onInstallSnapshotRequest(String requestId, InstallSnapshotRequestProto request) { - if (started.contains(requestId)) { - chunks.incrementAndGet(); - } - } - - @Override - public void onInstallSnapshotReply(String requestId, InstallSnapshotReplyProto reply) { - if (started.contains(requestId)) { - replies.incrementAndGet(); - } - } - - @Override - public void onInstallSnapshotEnd(String requestId, Throwable error) { - if (error == null && started.contains(requestId)) { - completed.add(requestId); - } - } - }); - final RaftProperties properties = newProperties(); - RaftServerConfigKeys.Snapshot.setAutoTriggerEnabled(properties, true); - RaftServerConfigKeys.Snapshot.setAutoTriggerThreshold(properties, 64); - RaftServerConfigKeys.Log.setPurgeGap(properties, 8); - RaftServerConfigKeys.Log.Appender.setSnapshotChunkSizeMax(properties, SizeInBytes.ONE_KB); - RaftServerConfigKeys.LeaderElection.setMemberMajorityAdd(properties, true); - try (MiniRaftClusterWithGrpc cluster = new MiniRaftClusterWithGrpc( - MiniRaftCluster.generateIds(1, 30), new String[0], properties, parameters)) { - cluster.start(); - final RaftServer.Division leader = RaftTestUtil.waitForLeader(cluster); - try (RaftClient client = cluster.createClient(leader.getId())) { - for (int i = 0; i < 127; i++) { - Assertions.assertTrue(client.io().send(new RaftTestUtil.SimpleMessage("snapshot-" + i)).isSuccess()); - } - Assertions.assertTrue(client.getSnapshotManagementApi(leader.getId()).create(3000).isSuccess()); - } - final PeerChanges change = cluster.addNewPeers(1, true); - cluster.setConfiguration(change.getPeersInNewConf()); - RaftServerTestUtil.waitAndCheckNewConf(cluster, change.getPeersInNewConf(), 0, null); - JavaUtils.attempt(() -> { - Assertions.assertFalse(completed.isEmpty()); - Assertions.assertTrue(chunks.get() > 1); - Assertions.assertEquals(chunks.get(), replies.get()); - }, 20, ONE_SECOND, "snapshot callbacks", LOG); - } - } } 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 index f015287223..d134891ba9 100644 --- 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 @@ -34,12 +34,10 @@ import org.apache.ratis.protocol.RaftGroupId; import org.apache.ratis.protocol.RaftGroupMemberId; import org.apache.ratis.protocol.RaftPeerId; -import org.apache.ratis.protocol.exceptions.TimeoutIOException; 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.statemachine.SnapshotInfo; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assertions; @@ -49,17 +47,15 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.EnumSource; import org.junit.jupiter.params.provider.ValueSource; -import org.mockito.ArgumentCaptor; import org.mockito.InOrder; import java.io.IOException; -import java.util.concurrent.CompletableFuture; +import java.io.InterruptedIOException; 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.ArgumentMatchers.isA; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.clearInvocations; import static org.mockito.Mockito.doReturn; @@ -68,6 +64,7 @@ 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; @@ -77,6 +74,7 @@ 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); @@ -99,6 +97,7 @@ public void setup() throws Exception { 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"); @@ -140,9 +139,9 @@ public void testStreamErrorAfterStopOrLeaderChange(boolean stop) throws Exceptio } clearInvocations(client, leaderState, follower, follower.getErrorState()); final IOException error = new IOException("Stream failed after leadership ended"); - doThrow(new IllegalStateException("Listener failure")).when(listener).onAppendEntriesReset(error); + doThrow(new IllegalStateException("Listener failure")).when(appendEntries).onError(error); Assertions.assertDoesNotThrow(() -> responses.onError(error)); - verify(listener).onAppendEntriesReset(error); + verify(appendEntries).onError(error); verifyNoInteractions(client, leaderState, follower.getErrorState()); verify(follower, never()).computeNextIndex(any()); } @@ -152,8 +151,11 @@ 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).onReset(anyString(), eq(error)); responses.onError(error); - verify(listener).onAppendEntriesReset(error); + final InOrder order = inOrder(appendEntries, listener); + order.verify(appendEntries).onError(error); + order.verify(listener).onReset(anyString(), eq(error)); } @Test @@ -162,7 +164,8 @@ public void testCompletionAfterStop() throws Exception { appender.stopAsync().get(5, TimeUnit.SECONDS); clearInvocations(client, leaderState); responses.onCompleted(); - verify(listener).onAppendEntriesReset(isA(IOException.class)); + verify(appendEntries).onCompleted(); + verify(listener, never()).onReset(anyString(), any()); verifyNoInteractions(client, leaderState); } @@ -180,9 +183,9 @@ public void testStreamCreationFailure(boolean failProxy) throws Exception { when(client.appendEntries(any(), eq(false))).thenThrow(error); } Assertions.assertSame(error, Assertions.assertThrows(Exception.class, () -> appender.appendLog(false))); - final InOrder order = inOrder(listener); - order.verify(listener).onAppendEntriesRequest(request); - order.verify(listener).onAppendEntriesFailure(1, error); + final InOrder order = inOrder(appendEntries); + order.verify(appendEntries).onRequest(request); + order.verify(appendEntries).onFailure(1, error); Assertions.assertFalse(appender.hasPendingDataRequests()); } @@ -191,58 +194,58 @@ public void testTimeout() { addPending(1); appender.timeoutAppendRequest(1, false); appender.timeoutAppendRequest(1, false); - verify(listener).onAppendEntriesFailure(eq(1L), isA(TimeoutIOException.class)); + verify(appendEntries).onTimeout(1); } @ParameterizedTest @EnumSource(value = AppendResult.class, names = {"SUCCESS", "NOT_LEADER", "INCONSISTENCY"}) - public void testMatchedReply(AppendResult result) { + 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(listener).onAppendEntriesReply(reply); + verify(appendEntries, times(2)).onReply(reply); + if (result == AppendResult.INCONSISTENCY) { + verify(listener, times(2)).onReset(anyString(), eq(null)); + Assertions.assertFalse(appender.hasPendingDataRequests()); + } appender.timeoutAppendRequest(1, false); - verify(listener, never()).onAppendEntriesFailure(eq(1L), any()); + verify(appendEntries, never()).onTimeout(1); } - @Test - public void testFactoryFailureIsIsolated() throws Exception { + @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); - throw new IllegalStateException("Injected factory failure"); + 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); } - @Test - public void testSnapshotStreamCreationFailure() throws Exception { - final IOException error = new IOException("Cannot create snapshot connection"); - when(serverRpc.getProxies().getProxy(destination)).thenThrow(error); - appender.installSnapshot(mock(SnapshotInfo.class, RETURNS_DEEP_STUBS)); - final ArgumentCaptor requestId = ArgumentCaptor.forClass(String.class); - final InOrder order = inOrder(listener); - order.verify(listener).onInstallSnapshotStart(requestId.capture(), eq(false)); - order.verify(listener).onInstallSnapshotEnd(requestId.getValue(), error); - verify(listener, never()).onInstallSnapshotRequest(anyString(), any()); - } - - @Test - public void testSnapshotErrorAfterStopUnblocksWaiter() throws Exception { - final GrpcLogAppender.InstallSnapshotResponseHandler snapshot = appender.new InstallSnapshotResponseHandler(); - final ArgumentCaptor requestId = ArgumentCaptor.forClass(String.class); - verify(listener).onInstallSnapshotStart(requestId.capture(), eq(false)); - final CompletableFuture done = (CompletableFuture) RaftTestUtil.getDeclaredField(snapshot, "done"); - appender.stopAsync().get(5, TimeUnit.SECONDS); - final IOException error = new IOException("Snapshot stream failed after stop"); - doThrow(new IllegalStateException("Listener failure")).when(listener).onInstallSnapshotEnd(anyString(), any()); - snapshot.onError(error); - Assertions.assertTrue(done.isDone()); - verify(listener).onInstallSnapshotEnd(requestId.getValue(), error); + @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(); } } From 785a1eaaa9a7bbf05891d941fb3c8da4d4ee59eb Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 8 Sep 2026 14:44:29 +0800 Subject: [PATCH 5/6] RATIS-2681. Retrigger CI From 7c6604cd4bc67c3459f89c7f961eea4b2c13305d Mon Sep 17 00:00:00 2001 From: HTHou Date: Wed, 9 Sep 2026 09:28:07 +0800 Subject: [PATCH 6/6] RATIS-2681. Separate inconsistency callbacks from client resets --- .../ratis/grpc/GrpcLogAppenderListener.java | 15 ++++++-- .../ratis/grpc/server/GrpcLogAppender.java | 16 ++------ .../server/TestGrpcLogAppenderCallbacks.java | 37 +++++++++++++++++-- 3 files changed, 48 insertions(+), 20 deletions(-) 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 index 4310cc6b0f..b497f9efd4 100644 --- a/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java +++ b/ratis-grpc/src/main/java/org/apache/ratis/grpc/GrpcLogAppenderListener.java @@ -28,8 +28,9 @@ * 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, reply and reset callbacks are serialized per appender. Failure - * callbacks may race with these callbacks. No exactly-once terminal notification is guaranteed. + * 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. */ @@ -52,6 +53,12 @@ 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) { } @@ -66,10 +73,10 @@ default void onError(Throwable error) { } } /** - * Pending attempts are invalidated before resetting the client or reconciling an inconsistent log. + * 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 onReset(String reason, Throwable error) { } + 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 c5d027a2b1..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 @@ -230,11 +230,6 @@ private void notifyAppendEntriesListener(Consumer l.onReset(reason, error)); - } - @Override public GrpcServicesImpl getServerRpc() { return (GrpcServicesImpl)super.getServerRpc(); @@ -246,7 +241,7 @@ private GrpcServerProtocolClient getClient() throws IOException { private void resetClient(AppendEntriesRequest request, Event event, Throwable error) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - notifyReset("resetClient: " + event, error); + notifyListener(l -> l.onResetClient("resetClient: " + event, error)); getClient().resetConnectBackoff(); if (appendLogRequestObserver != null) { appendLogRequestObserver.stop(); @@ -546,11 +541,8 @@ class AppendLogResponseHandler implements StreamObserver l.onReply(reply)); - } + final AppendEntriesRequest request = pendingRequests.remove(reply); + notifyAppendEntriesListener(l -> l.onReply(reply)); if (request != null) { request.stopRequestTimer(); // Update completion time getFollower().updateLastRespondedAppendEntriesSendTime(request.getSendTime()); @@ -641,7 +633,7 @@ public String toString() { private void updateNextIndex(long replyNextIndex) { try (AutoCloseableLock writeLock = lock.writeLock(caller, LOG::trace)) { - notifyReset("AppendEntries INCONSISTENCY", null); + notifyAppendEntriesListener(GrpcLogAppenderListener.AppendEntries::onReplyInconsistency); pendingRequests.clear(); getFollower().setNextIndex(replyNextIndex); } 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 index d134891ba9..f5b959e5a5 100644 --- 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 @@ -39,6 +39,8 @@ 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; @@ -51,6 +53,9 @@ 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; @@ -151,11 +156,11 @@ 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).onReset(anyString(), eq(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).onReset(anyString(), eq(error)); + order.verify(listener).onResetClient(anyString(), eq(error)); } @Test @@ -165,7 +170,7 @@ public void testCompletionAfterStop() throws Exception { clearInvocations(client, leaderState); responses.onCompleted(); verify(appendEntries).onCompleted(); - verify(listener, never()).onReset(anyString(), any()); + verify(listener, never()).onResetClient(anyString(), any()); verifyNoInteractions(client, leaderState); } @@ -208,13 +213,37 @@ public void testReplyAndInvalidation(AppendResult result) { responses.onNext(reply); verify(appendEntries, times(2)).onReply(reply); if (result == AppendResult.INCONSISTENCY) { - verify(listener, times(2)).onReset(anyString(), eq(null)); + 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 {