kevin-wu24 commented on code in PR #22620:
URL: https://github.com/apache/kafka/pull/22620#discussion_r3468304438
##########
raft/src/main/java/org/apache/kafka/raft/internals/UpdateVoterHandler.java:
##########
@@ -155,50 +140,193 @@ public CompletionStage<UpdateRaftVoterResponseData>
handleUpdateVoterRequest(
);
}
- // Check that endpoints includes the default listener
- if (voterEndpoints.address(defaultListenerName).isEmpty()) {
+ // Send API_VERSIONS request to new voter to test new default endpoint
+ var timeout = requestSender.send(
+ voterEndpoints
+ .address(requestSender.listenerName())
+ .map(address -> new Node(voterKey.id(), address.getHostName(),
address.getPort()))
+ .orElseThrow(
+ () -> new IllegalStateException(
+ String.format(
+ "Provided listeners %s do not contain a listener
for %s",
+ voterEndpoints,
+ requestSender.listenerName()
+ )
+ )
+ ),
+ this::buildApiVersionsRequest,
+ currentTimeMs
+ );
+ if (timeout.isEmpty()) {
return CompletableFuture.completedFuture(
RaftUtil.updateVoterResponse(
- Errors.INVALID_REQUEST,
+ Errors.REQUEST_TIMED_OUT,
requestListenerName,
leaderState.leaderAndEpoch(),
leaderState.leaderEndpoints()
)
);
}
+ var state = new UpdateVoterHandlerState(
+ voterKey,
+ voterEndpoints,
+ requestListenerName,
+ new SupportedVersionRange(
+ supportedKraftVersions.minSupportedVersion(),
+ supportedKraftVersions.maxSupportedVersion()
+ ),
+ time.timer(timeout.getAsLong())
+ );
+ changeVoterState.resetUpdateVoterHandlerState(
+ Errors.UNKNOWN_SERVER_ERROR,
+ leaderState.leaderAndEpoch(),
+ leaderState.leaderEndpoints(),
+ Optional.of(state)
+ );
+
+ return state.future();
+ }
+
+ /**
+ * Handle the API_VERSIONS response for an update voter operation.
+ *
+ * @param leaderState the leader state
+ * @param source the node that sent the response
+ * @param error the error from the response
+ * @param supportedKraftVersions the supported kraft version range from
the response
+ * @param currentTimeMs the current time in milliseconds
+ * @return true if the update voter operation should continue, false if it
was aborted
+ */
+ public boolean handleApiVersionsResponse(
+ LeaderState<?> leaderState,
+ Node source,
+ Errors error,
+ Optional<ApiVersionsResponseData.SupportedFeatureKey>
supportedKraftVersions,
+ long currentTimeMs
+ ) {
+ var changeVoterState = leaderState.changeVoterState();
+ var handlerState = changeVoterState.updateVoterHandlerState();
+ if (handlerState.isEmpty()) {
+ // There are no pending add operation just ignore the api response
+ return true;
+ }
+
+ // Check that the API_VERSIONS response matches the id of the voter
getting added
+ var current = handlerState.get();
+ if (!current.expectingApiResponse(source.id())) {
+ logger.info(
+ "API_VERSIONS response is not expected from {}: voterKey is
{}, lastOffset is {}",
+ source,
+ current.voterKey(),
+ current.lastOffset()
+ );
+
+ return true;
+ } else if (error != Errors.NONE) {
+ // Abort operation if the API_VERSIONS returned an error
+ logger.info(
+ "Aborting update voter operation for {} at {} since
API_VERSIONS returned an error {}",
+ current.voterKey(),
+ current.voterEndpoints(),
+ error
+ );
+
+ changeVoterState.resetUpdateVoterHandlerState(
+ Errors.REQUEST_TIMED_OUT,
+ leaderState.leaderAndEpoch(),
+ leaderState.leaderEndpoints(),
+ Optional.empty()
+ );
+
+ return false;
+ } else if (
+ !Optional.of(current.supportedKraftVersions())
+
.equals(supportedKraftVersions.map(this::convertToVersionRange))
Review Comment:
Ah, I think I read `current` wrong. Apologies.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]