jsancio commented on code in PR #22620:
URL: https://github.com/apache/kafka/pull/22620#discussion_r3813658558


##########
raft/src/main/java/org/apache/kafka/raft/internals/UpdateVoterHandler.java:
##########
@@ -101,112 +171,273 @@ public CompletionStage<UpdateRaftVoterResponseData> 
handleUpdateVoterRequest(
             );
         }
 
-        // Read the voter set from the log or leader state
-        KRaftVersion kraftVersion = partitionState.lastKraftVersion();
-        final Optional<KRaftVersionUpgrade.Voters> inMemoryVoters;
-        final Optional<VoterSet> voters;
-        if (kraftVersion.isReconfigSupported()) {
-            inMemoryVoters = Optional.empty();
-
-            // Check that there are no uncommitted VotersRecord
-            Optional<LogHistory.Entry<VoterSet>> votersEntry = 
partitionState.lastVoterSetEntry();
-            if (votersEntry.isEmpty() || votersEntry.get().offset() >= 
highWatermark.get()) {
-                voters = Optional.empty();
-            } else {
-                voters = votersEntry.map(LogHistory.Entry::value);
-            }
-        } else {
-            inMemoryVoters = leaderState.volatileVoters();
-            if (inMemoryVoters.isEmpty()) {
-                /* This can happen if the remote voter sends an update voter 
request before the
-                 * updated kraft version has been written to the log
-                 */
-                return CompletableFuture.completedFuture(
-                    RaftUtil.updateVoterResponse(
-                        Errors.REQUEST_TIMED_OUT,
-                        requestListenerName,
-                        leaderState.leaderAndEpoch(),
-                        leaderState.leaderEndpoints()
-                    )
-                );
-            }
-            voters = inMemoryVoters.map(KRaftVersionUpgrade.Voters::voters);
-        }
-        if (voters.isEmpty()) {
-            log.info("Unable to read the current voter set with kraft version 
{}", kraftVersion);
+        // Check that endpoints includes the default listener
+        if (voterEndpoints.address(requestSender.listenerName()).isEmpty()) {
             return CompletableFuture.completedFuture(
                 RaftUtil.updateVoterResponse(
-                    Errors.REQUEST_TIMED_OUT,
+                    Errors.INVALID_REQUEST,
                     requestListenerName,
                     leaderState.leaderAndEpoch(),
                     leaderState.leaderEndpoints()
                 )
             );
         }
-        // Check that the supported version range is valid
-        if (!validVersionRange(kraftVersion, supportedKraftVersions)) {
+
+        // 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()
                 )
             );
         }
 
-        // Check that endpoints includes the default listener
-        if (voterEndpoints.address(defaultListenerName).isEmpty()) {
-            return CompletableFuture.completedFuture(
-                RaftUtil.updateVoterResponse(
-                    Errors.INVALID_REQUEST,
-                    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 a pending update voter operation.
+     * <p>
+     * This may abort the pending operation, completing its future with an 
error, if the response
+     * doesn't come from the expected voter, if the API_VERSIONS request 
failed, if the supported
+     * kraft.version range doesn't match the one from the UpdateVoter request, 
if the current voter
+     * set can't be read yet, or if the voter is no longer part of the voter 
set. Otherwise, it
+     * stores the updated voter set, see {@link #storeUpdatedVoters}.
+     *
+     * @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 false only when the API_VERSIONS request itself failed, which 
is the only case where
+     *         the caller (see {@code KafkaRaftClient#handleResponse}) should 
treat this as an
+     *         unsuccessful response for request-tracking purposes; true 
otherwise, including when
+     *         this method aborts the pending update voter operation for 
another reason
+     */
+    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 is no pending update operation; just ignore the 
API_VERSIONS response
+            return true;
+        }
+
+        // Check that the API_VERSIONS response matches the id of the voter 
getting updated
+        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))
+        ) {
+            // Check that the supported version from the ApiVersions response 
matches the supported
+            // version from the UpdateVoter request
+            logger.error(
+                "The supported kraft version from UpdateVoters {} doesn't 
match the supported " +
+                "kraft version from ApiVersions {}",
+                current.supportedKraftVersions(),
+                supportedKraftVersions
+            );
+            changeVoterState.resetUpdateVoterHandlerState(
+                Errors.INVALID_REQUEST,
+                leaderState.leaderAndEpoch(),
+                leaderState.leaderEndpoints(),
+                Optional.empty()
+            );
+            return true;
+        }
+
+        return completeUpdateVoter(leaderState, changeVoterState, current, 
currentTimeMs);
+    }
+
+    /**
+     * Validates the kraft.version range and current voter set, then applies 
the update, once the
+     * API_VERSIONS response has already been matched to a pending update 
voter operation.
+     *
+     * @return true always; matches the caller's return convention where only 
a failed API_VERSIONS

Review Comment:
   That's fair. The method returns void now.



-- 
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]

Reply via email to