From fc63e2ef765c09895f517f95e4740fcaabdd2be7 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Mon, 15 Aug 2022 14:54:20 -0700 Subject: [PATCH 1/2] KAFKA-13940; Return NOT_LEADER_OR_FOLLOWER if DescribeQuorum sent to non-leader --- .../requests/DescribeQuorumResponse.java | 12 ++++++++ .../clients/admin/KafkaAdminClientTest.java | 23 +++++++++++++++ .../apache/kafka/raft/KafkaRaftClient.java | 5 +++- .../kafka/raft/KafkaRaftClientTest.java | 29 +++++++++++++++++++ .../kafka/raft/RaftClientTestContext.java | 27 ++++++++--------- 5 files changed, 82 insertions(+), 14 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java index cbf945b70409a..bc13eaed9d4b9 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java @@ -72,6 +72,18 @@ public int throttleTimeMs() { return DEFAULT_THROTTLE_TIME; } + public static DescribeQuorumResponseData singletonErrorResponse( + TopicPartition topicPartition, + Errors error + ) { + return new DescribeQuorumResponseData() + .setTopics(Collections.singletonList(new DescribeQuorumResponseData.TopicData() + .setTopicName(topicPartition.topic()) + .setPartitions(Collections.singletonList(new DescribeQuorumResponseData.PartitionData() + .setErrorCode(error.code()))))); + } + + public static DescribeQuorumResponseData singletonResponse(TopicPartition topicPartition, int leaderId, int leaderEpoch, diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java index de57813679b99..5faf53f0756fa 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientTest.java @@ -5192,6 +5192,29 @@ public void testDescribeMetadataQuorumSuccess() throws Exception { } } + @Test + public void testDescribeMetadataQuorumRetriableError() throws Exception { + try (final AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create(ApiKeys.DESCRIBE_QUORUM.id, + ApiKeys.DESCRIBE_QUORUM.oldestVersion(), + ApiKeys.DESCRIBE_QUORUM.latestVersion())); + + // First request fails with a NOT_LEADER_OR_FOLLOWER error (which is retriable) + env.kafkaClient().prepareResponse( + body -> body instanceof DescribeQuorumRequest, + prepareDescribeQuorumResponse(Errors.NONE, Errors.NOT_LEADER_OR_FOLLOWER, false, false, false, false, false)); + + // The second request succeeds + env.kafkaClient().prepareResponse( + body -> body instanceof DescribeQuorumRequest, + prepareDescribeQuorumResponse(Errors.NONE, Errors.NONE, false, false, false, false, false)); + + KafkaFuture future = env.adminClient().describeMetadataQuorum().quorumInfo(); + QuorumInfo quorumInfo = future.get(); + assertEquals(defaultQuorumInfo(false), quorumInfo); + } + } + @Test public void testDescribeMetadataQuorumFailure() { try (final AdminClientUnitTestEnv env = mockClientEnv()) { diff --git a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java index cac7a8a3cb998..042a141a760ae 100644 --- a/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java +++ b/raft/src/main/java/org/apache/kafka/raft/KafkaRaftClient.java @@ -1171,7 +1171,10 @@ private DescribeQuorumResponseData handleDescribeQuorumRequest( } if (!quorum.isLeader()) { - return DescribeQuorumRequest.getTopLevelErrorResponse(Errors.INVALID_REQUEST); + return DescribeQuorumResponse.singletonErrorResponse( + log.topicPartition(), + Errors.NOT_LEADER_OR_FOLLOWER + ); } LeaderState leaderState = quorum.leaderStateOrThrow(); diff --git a/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientTest.java b/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientTest.java index 9b2771d2b34e4..a8a346e6dbf5b 100644 --- a/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientTest.java +++ b/raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientTest.java @@ -20,6 +20,7 @@ import org.apache.kafka.common.errors.RecordBatchTooLargeException; import org.apache.kafka.common.memory.MemoryPool; import org.apache.kafka.common.message.BeginQuorumEpochResponseData; +import org.apache.kafka.common.message.DescribeQuorumResponseData; import org.apache.kafka.common.message.DescribeQuorumResponseData.ReplicaState; import org.apache.kafka.common.message.EndQuorumEpochResponseData; import org.apache.kafka.common.message.FetchResponseData; @@ -1946,6 +1947,34 @@ public void testEndQuorumEpochSentBasedOnFetchOffset() throws Exception { ); } + @Test + public void testDescribeQuorumNonLeader() throws Exception { + int localId = 0; + int voter2 = localId + 1; + int voter3 = localId + 2; + int epoch = 2; + Set voters = Utils.mkSet(localId, voter2, voter3); + + RaftClientTestContext context = new RaftClientTestContext.Builder(localId, voters) + .withUnknownLeader(epoch) + .build(); + + context.deliverRequest(DescribeQuorumRequest.singletonRequest(context.metadataPartition)); + context.pollUntilResponse(); + + DescribeQuorumResponseData responseData = context.collectDescribeQuorumResponse(); + assertEquals(Errors.NONE, Errors.forCode(responseData.errorCode())); + + assertEquals(1, responseData.topics().size()); + DescribeQuorumResponseData.TopicData topicData = responseData.topics().get(0); + assertEquals(context.metadataPartition.topic(), topicData.topicName()); + + assertEquals(1, topicData.partitions().size()); + DescribeQuorumResponseData.PartitionData partitionData = topicData.partitions().get(0); + assertEquals(context.metadataPartition.partition(), partitionData.partitionIndex()); + assertEquals(Errors.NOT_LEADER_OR_FOLLOWER, Errors.forCode(partitionData.errorCode())); + } + @Test public void testDescribeQuorum() throws Exception { int localId = 0; diff --git a/raft/src/test/java/org/apache/kafka/raft/RaftClientTestContext.java b/raft/src/test/java/org/apache/kafka/raft/RaftClientTestContext.java index d48e41fb31d0c..b825fc8867a15 100644 --- a/raft/src/test/java/org/apache/kafka/raft/RaftClientTestContext.java +++ b/raft/src/test/java/org/apache/kafka/raft/RaftClientTestContext.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.raft; -import java.util.function.Consumer; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.memory.MemoryPool; @@ -57,8 +56,8 @@ import org.apache.kafka.raft.internals.BatchBuilder; import org.apache.kafka.raft.internals.StringSerde; import org.apache.kafka.server.common.serialization.RecordSerde; -import org.apache.kafka.snapshot.SnapshotReader; import org.apache.kafka.snapshot.RawSnapshotWriter; +import org.apache.kafka.snapshot.SnapshotReader; import org.apache.kafka.test.TestCondition; import org.apache.kafka.test.TestUtils; @@ -76,6 +75,7 @@ import java.util.OptionalInt; import java.util.OptionalLong; import java.util.Set; +import java.util.function.Consumer; import java.util.stream.Collectors; import static org.apache.kafka.raft.RaftUtil.hasValidTopicPartition; @@ -128,7 +128,7 @@ public static final class Builder { private final QuorumStateStore quorumStateStore = new MockQuorumStateStore(); private final MockableRandom random = new MockableRandom(1L); private final LogContext logContext = new LogContext(); - private final MockLog log = new MockLog(METADATA_PARTITION, Uuid.METADATA_TOPIC_ID, logContext); + private final MockLog log = new MockLog(METADATA_PARTITION, Uuid.METADATA_TOPIC_ID, logContext); private final Set voters; private final OptionalInt localId; @@ -440,21 +440,24 @@ void assertResignedLeader(int epoch, int leaderId) throws IOException { assertEquals(ElectionState.withElectedLeader(epoch, leaderId, voters), quorumStateStore.readElectionState()); } - int assertSentDescribeQuorumResponse( - int leaderId, - int leaderEpoch, - long highWatermark, - List voterStates, - List observerStates - ) { + DescribeQuorumResponseData collectDescribeQuorumResponse() { List sentMessages = drainSentResponses(ApiKeys.DESCRIBE_QUORUM); assertEquals(1, sentMessages.size()); RaftResponse.Outbound raftMessage = sentMessages.get(0); assertTrue( raftMessage.data() instanceof DescribeQuorumResponseData, "Unexpected request type " + raftMessage.data()); - DescribeQuorumResponseData response = (DescribeQuorumResponseData) raftMessage.data(); + return (DescribeQuorumResponseData) raftMessage.data(); + } + void assertSentDescribeQuorumResponse( + int leaderId, + int leaderEpoch, + long highWatermark, + List voterStates, + List observerStates + ) { + DescribeQuorumResponseData response = collectDescribeQuorumResponse(); DescribeQuorumResponseData expectedResponse = DescribeQuorumResponse.singletonResponse( metadataPartition, leaderId, @@ -462,9 +465,7 @@ int assertSentDescribeQuorumResponse( highWatermark, voterStates, observerStates); - assertEquals(expectedResponse, response); - return raftMessage.correlationId(); } int assertSentVoteRequest(int epoch, int lastEpoch, long lastEpochOffset, int numVoteReceivers) { From 4172e552bda4fcfbab9457a8af2ab3dcd0a2d444 Mon Sep 17 00:00:00 2001 From: Jason Gustafson Date: Tue, 16 Aug 2022 08:31:43 -0700 Subject: [PATCH 2/2] Set partition index explicitly in DescribeQuorumResponse --- .../apache/kafka/common/requests/DescribeQuorumResponse.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java index bc13eaed9d4b9..06ae681bc5c13 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeQuorumResponse.java @@ -37,7 +37,7 @@ * - {@link Errors#BROKER_NOT_AVAILABLE} * * Partition level errors: - * - {@link Errors#INVALID_REQUEST} + * - {@link Errors#NOT_LEADER_OR_FOLLOWER} * - {@link Errors#UNKNOWN_TOPIC_OR_PARTITION} */ public class DescribeQuorumResponse extends AbstractResponse { @@ -80,6 +80,7 @@ public static DescribeQuorumResponseData singletonErrorResponse( .setTopics(Collections.singletonList(new DescribeQuorumResponseData.TopicData() .setTopicName(topicPartition.topic()) .setPartitions(Collections.singletonList(new DescribeQuorumResponseData.PartitionData() + .setPartitionIndex(topicPartition.partition()) .setErrorCode(error.code()))))); } @@ -94,6 +95,7 @@ public static DescribeQuorumResponseData singletonResponse(TopicPartition topicP .setTopics(Collections.singletonList(new DescribeQuorumResponseData.TopicData() .setTopicName(topicPartition.topic()) .setPartitions(Collections.singletonList(new DescribeQuorumResponseData.PartitionData() + .setPartitionIndex(topicPartition.partition()) .setErrorCode(Errors.NONE.code()) .setLeaderId(leaderId) .setLeaderEpoch(leaderEpoch)