kevin-wu24 commented on code in PR #22531:
URL: https://github.com/apache/kafka/pull/22531#discussion_r3483236786
##########
raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientReconfigTest.java:
##########
@@ -364,6 +364,170 @@ public void testAddVoter() throws Exception {
context.assertSentAddVoterResponse(Errors.NONE);
}
+ @Test
+ public void testAddVoterDerivesDirectoryIdAndEndpoints() throws Exception {
+ ReplicaKey local = replicaKey(randomReplicaId(), true);
+ ReplicaKey follower = replicaKey(local.id() + 1, true);
+ VoterSet voters = VoterSetTest.voterSet(Stream.of(local, follower));
+
+ RaftClientTestContext context = new
RaftClientTestContext.Builder(local.id(), local.directoryId().get())
+ .withRaftProtocol(RaftProtocol.KIP_1141_PROTOCOL)
+ .withBootstrapSnapshot(Optional.of(voters))
+ .withUnknownLeader(3)
+ .build();
+
+ context.unattachedToLeader();
+ int epoch = context.currentEpoch();
+
+ ReplicaKey newVoter = replicaKey(local.id() + 2, true);
+ InetSocketAddress newAddress = InetSocketAddress.createUnresolved(
+ "localhost",
+ 9990 + newVoter.id()
+ );
+ Endpoints newListeners = Endpoints.fromInetSocketAddresses(
+ Map.of(context.channel.listenerName(), newAddress)
+ );
+ context.client.setNodeEndpointProvider(nodeId ->
+ nodeId == newVoter.id() ? newListeners : Endpoints.empty()
+ );
+
+ prepareLeaderToReceiveAddVoter(context, epoch, local, follower,
newVoter);
+
+ context.deliverRequest(
+ context.addVoterRequest(
+ Integer.MAX_VALUE,
+ ReplicaKey.of(newVoter.id(), Uuid.ZERO_UUID),
+ Endpoints.empty()
+ )
+ );
+
+ completeApiVersionsForAddVoter(context, newVoter, newAddress);
+
+ // Handle the API_VERSIONS response
+ context.client.poll();
+ // Append new VotersRecord to log
+ context.client.poll();
+
+ commitNewVoterSetForAddVoter(context, local, follower, newVoter,
epoch);
+
+ // Expect reply for AddVoter request
+ context.pollUntilResponse();
+ context.assertSentAddVoterResponse(Errors.NONE);
+ }
+
+ @Test
+ public void testAddVoterFailsToDeriveDefaultListenerEndpoint() throws
Exception {
+ ReplicaKey local = replicaKey(randomReplicaId(), true);
+ ReplicaKey follower = replicaKey(local.id() + 1, true);
+ VoterSet voters = VoterSetTest.voterSet(Stream.of(local, follower));
+
+ RaftClientTestContext context = new
RaftClientTestContext.Builder(local.id(), local.directoryId().get())
+ .withRaftProtocol(RaftProtocol.KIP_1141_PROTOCOL)
+ .withBootstrapSnapshot(Optional.of(voters))
+ .withUnknownLeader(3)
+ .build();
+
+ context.unattachedToLeader();
+ int epoch = context.currentEpoch();
+ ReplicaKey newVoter = replicaKey(local.id() + 2, true);
+
+ prepareLeaderToReceiveAddVoter(context, epoch, local, follower,
newVoter);
+
+ context.deliverRequest(
+ context.addVoterRequest(
+ Integer.MAX_VALUE,
+ ReplicaKey.of(newVoter.id(), Uuid.ZERO_UUID),
+ Endpoints.empty()
+ )
+ );
+
+ context.pollUntilResponse();
+ context.assertSentAddVoterResponse(Errors.INVALID_REQUEST);
+ }
+
+ @Test
+ public void testAddVoterDoesNotReplaceExplicitDirectoryId() throws
Exception {
+ ReplicaKey local = replicaKey(randomReplicaId(), true);
+ ReplicaKey follower = replicaKey(local.id() + 1, true);
+ VoterSet voters = VoterSetTest.voterSet(Stream.of(local, follower));
+
+ RaftClientTestContext context = new
RaftClientTestContext.Builder(local.id(), local.directoryId().get())
+ .withRaftProtocol(RaftProtocol.KIP_1141_PROTOCOL)
+ .withBootstrapSnapshot(Optional.of(voters))
+ .withUnknownLeader(3)
+ .build();
+
+ context.unattachedToLeader();
+ int epoch = context.currentEpoch();
+ ReplicaKey observer = replicaKey(local.id() + 2, true);
+ ReplicaKey requestedVoter = replicaKey(observer.id(), true);
+ InetSocketAddress newAddress = InetSocketAddress.createUnresolved(
+ "localhost",
+ 9990 + observer.id()
+ );
+ Endpoints newListeners = Endpoints.fromInetSocketAddresses(
+ Map.of(context.channel.listenerName(), newAddress)
+ );
+
+ prepareLeaderToReceiveAddVoter(context, epoch, local, follower,
observer);
+ context.deliverRequest(
+ context.addVoterRequest(Integer.MAX_VALUE, requestedVoter,
newListeners)
+ );
+ completeApiVersionsForAddVoter(context, requestedVoter, newAddress);
+
+ context.pollUntilResponse();
+ context.assertSentAddVoterResponse(Errors.REQUEST_TIMED_OUT);
Review Comment:
I think you're hitting this code in
`AddVoterHandler#handleApiVersionsResponse`:
```
if (!leaderState.isReplicaCaughtUp(current.voterKey(), currentTimeMs)) {
logger.info(
"Aborting add voter operation for {} at {} since it is
lagging behind: {}",
current.voterKey(),
current.voterEndpoints(),
leaderState.getReplicaState(current.voterKey())
);
```
We do a fetch via the `observer` replica key to catch it up before sending
the add voter request, but not the `requestedVoter` replica key.
Ideally, putting the wrong directory id in the `ADD_RAFT_VOTER` request
should not even get to this point. We should return `INVALID_REQUEST` before
even entering `AddVoterHandler`. We should be returning a similar error to the
test below.
Because you're running `add-voter` from the node that is being added (in
both manual + auto join case), the directory id would be correct or
`ZERO_UUID`, but never incorrect.
--
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]