sjhajharia commented on code in PR #22339: URL: https://github.com/apache/kafka/pull/22339#discussion_r3725835663
########## clients/src/test/java/org/apache/kafka/clients/admin/KafkaAdminClientShareGroupTest.java: ########## @@ -0,0 +1,987 @@ +/* + * 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.kafka.clients.admin; + +import org.apache.kafka.clients.ClientRequest; +import org.apache.kafka.clients.MockClient; +import org.apache.kafka.clients.NodeApiVersions; +import org.apache.kafka.common.Cluster; +import org.apache.kafka.common.GroupState; +import org.apache.kafka.common.GroupType; +import org.apache.kafka.common.KafkaException; +import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.Node; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.errors.GroupAuthorizationException; +import org.apache.kafka.common.errors.TopicAuthorizationException; +import org.apache.kafka.common.errors.UnknownServerException; +import org.apache.kafka.common.errors.UnsupportedVersionException; +import org.apache.kafka.common.message.AlterShareGroupOffsetsResponseData; +import org.apache.kafka.common.message.ApiVersionsResponseData.ApiVersion; +import org.apache.kafka.common.message.DeleteShareGroupOffsetsRequestData; +import org.apache.kafka.common.message.DeleteShareGroupOffsetsResponseData; +import org.apache.kafka.common.message.DescribeShareGroupOffsetsRequestData; +import org.apache.kafka.common.message.DescribeShareGroupOffsetsResponseData; +import org.apache.kafka.common.message.FindCoordinatorResponseData; +import org.apache.kafka.common.message.ListGroupsResponseData; +import org.apache.kafka.common.message.ShareGroupDescribeResponseData; +import org.apache.kafka.common.protocol.ApiKeys; +import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.requests.AlterShareGroupOffsetsResponse; +import org.apache.kafka.common.requests.DeleteShareGroupOffsetsRequest; +import org.apache.kafka.common.requests.DeleteShareGroupOffsetsResponse; +import org.apache.kafka.common.requests.DescribeShareGroupOffsetsRequest; +import org.apache.kafka.common.requests.DescribeShareGroupOffsetsResponse; +import org.apache.kafka.common.requests.FindCoordinatorResponse; +import org.apache.kafka.common.requests.ListGroupsResponse; +import org.apache.kafka.common.requests.MetadataResponse; +import org.apache.kafka.common.requests.RequestTestUtils; +import org.apache.kafka.common.requests.ShareGroupDescribeResponse; +import org.apache.kafka.common.utils.MockTime; +import org.apache.kafka.common.utils.Time; +import org.apache.kafka.test.TestUtils; + +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ExecutionException; +import java.util.stream.Collectors; + +import static java.util.Arrays.asList; +import static java.util.Collections.singletonList; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class KafkaAdminClientShareGroupTest extends KafkaAdminClientTestBase { + + @Test + public void testDescribeShareGroups() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + // Retriable FindCoordinatorResponse errors should be retried + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.COORDINATOR_NOT_AVAILABLE, Node.noNode())); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.COORDINATOR_LOAD_IN_PROGRESS, Node.noNode())); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + + ShareGroupDescribeResponseData data = new ShareGroupDescribeResponseData(); + + // Retriable errors should be retried + data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setErrorCode(Errors.COORDINATOR_LOAD_IN_PROGRESS.code())); + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(data)); + + /* + * We need to return two responses here, one with NOT_COORDINATOR error when calling describe share group + * api using coordinator that has moved. This will retry whole operation. So we need to again respond with a + * FindCoordinatorResponse. + * + * And the same reason for COORDINATOR_NOT_AVAILABLE error response + */ + data = new ShareGroupDescribeResponseData(); + data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setErrorCode(Errors.NOT_COORDINATOR.code())); + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(data)); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + + data = new ShareGroupDescribeResponseData(); + data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setErrorCode(Errors.COORDINATOR_NOT_AVAILABLE.code())); + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(data)); + env.kafkaClient().prepareResponse(prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + + data = new ShareGroupDescribeResponseData(); + ShareGroupDescribeResponseData.TopicPartitions topicPartitions = new ShareGroupDescribeResponseData.TopicPartitions() + .setTopicName("my_topic") + .setPartitions(asList(0, 1, 2)); + ShareGroupDescribeResponseData.Assignment memberAssignment = new ShareGroupDescribeResponseData.Assignment() + .setTopicPartitions(asList(topicPartitions)); + ShareGroupDescribeResponseData.Member memberOne = new ShareGroupDescribeResponseData.Member() + .setMemberId("0") + .setClientId("clientId0") + .setClientHost("clientHost") + .setAssignment(memberAssignment); + ShareGroupDescribeResponseData.Member memberTwo = new ShareGroupDescribeResponseData.Member() + .setMemberId("1") + .setClientId("clientId1") + .setClientHost("clientHost") + .setAssignment(memberAssignment); + + ShareGroupDescribeResponseData group0Data = new ShareGroupDescribeResponseData(); + group0Data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setGroupState(GroupState.STABLE.toString()) + .setMembers(asList(memberOne, memberTwo))); + + final List<TopicPartition> expectedTopicPartitions = new ArrayList<>(); + expectedTopicPartitions.add(0, new TopicPartition("my_topic", 0)); + expectedTopicPartitions.add(1, new TopicPartition("my_topic", 1)); + expectedTopicPartitions.add(2, new TopicPartition("my_topic", 2)); + + List<ShareMemberDescription> expectedMemberDescriptions = new ArrayList<>(); + expectedMemberDescriptions.add(convertToShareMemberDescriptions(memberOne, + new ShareMemberAssignment(new HashSet<>(expectedTopicPartitions)))); + expectedMemberDescriptions.add(convertToShareMemberDescriptions(memberTwo, + new ShareMemberAssignment(new HashSet<>(expectedTopicPartitions)))); + data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setGroupState(GroupState.STABLE.toString()) + .setMembers(asList(memberOne, memberTwo))); + + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(data)); + + final DescribeShareGroupsResult result = env.adminClient().describeShareGroups(singletonList(GROUP_ID)); + final ShareGroupDescription groupDescription = result.describedGroups().get(GROUP_ID).get(); + + assertEquals(1, result.describedGroups().size()); + assertEquals(GROUP_ID, groupDescription.groupId()); + assertEquals(2, groupDescription.members().size()); + assertEquals(expectedMemberDescriptions, groupDescription.members()); + } + } + + @Test + public void testDescribeShareGroupsGroupIdNotFound() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + env.kafkaClient().prepareResponse(new FindCoordinatorResponse( + new FindCoordinatorResponseData() + .setCoordinators(asList( + FindCoordinatorResponse.prepareCoordinatorResponse(Errors.NONE, GROUP_ID, env.cluster().controller()), + FindCoordinatorResponse.prepareCoordinatorResponse(Errors.NONE, "group-1", env.cluster().controller()) + )) + )); + + ShareGroupDescribeResponseData.TopicPartitions topicPartitions = new ShareGroupDescribeResponseData.TopicPartitions() + .setTopicName("my_topic") + .setPartitions(asList(0, 1, 2)); + final ShareGroupDescribeResponseData.Assignment memberAssignment = new ShareGroupDescribeResponseData.Assignment() + .setTopicPartitions(asList(topicPartitions)); + ShareGroupDescribeResponseData groupData = new ShareGroupDescribeResponseData(); + groupData.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setGroupState(GroupState.STABLE.toString()) + .setMembers(asList( + new ShareGroupDescribeResponseData.Member() + .setMemberId("0") + .setClientId("clientId0") + .setClientHost("clientHost") + .setAssignment(memberAssignment), + new ShareGroupDescribeResponseData.Member() + .setMemberId("1") + .setClientId("clientId1") + .setClientHost("clientHost") + .setAssignment(memberAssignment)))); + groupData.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId("group-1") + .setGroupState(GroupState.DEAD.toString()) + .setErrorCode(Errors.GROUP_ID_NOT_FOUND.code()) + .setErrorMessage("Group group-1 not found.")); + + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(groupData)); + + Collection<String> groups = new HashSet<>(); + groups.add(GROUP_ID); + groups.add("group-1"); + final DescribeShareGroupsResult result = env.adminClient().describeShareGroups(groups); + assertEquals(2, result.describedGroups().size()); + assertEquals(groups, result.describedGroups().keySet()); + KafkaFuture<Map<String, ShareGroupDescription>> allFuture = result.all(); + assertThrows(ExecutionException.class, allFuture::get); + assertTrue(result.all().isCompletedExceptionally()); + } + } + + @Test + public void testDescribeShareGroupsWithAuthorizedOperationsOmitted() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + env.kafkaClient().prepareResponse( + prepareFindCoordinatorResponse(Errors.NONE, env.cluster().controller())); + + ShareGroupDescribeResponseData data = new ShareGroupDescribeResponseData(); + + data.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setAuthorizedOperations(MetadataResponse.AUTHORIZED_OPERATIONS_OMITTED)); + + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(data)); + + final DescribeShareGroupsResult result = env.adminClient().describeShareGroups(singletonList(GROUP_ID)); + final ShareGroupDescription groupDescription = result.describedGroups().get(GROUP_ID).get(); + + assertNull(groupDescription.authorizedOperations()); + } + } + + @Test + public void testDescribeMultipleShareGroups() { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + env.kafkaClient().prepareResponse(new FindCoordinatorResponse( + new FindCoordinatorResponseData() + .setCoordinators(asList( + FindCoordinatorResponse.prepareCoordinatorResponse(Errors.NONE, GROUP_ID, env.cluster().controller()), + FindCoordinatorResponse.prepareCoordinatorResponse(Errors.NONE, "group-1", env.cluster().controller()) + )) + )); + + ShareGroupDescribeResponseData.TopicPartitions topicPartitions = new ShareGroupDescribeResponseData.TopicPartitions() + .setTopicName("my_topic") + .setPartitions(asList(0, 1, 2)); + final ShareGroupDescribeResponseData.Assignment memberAssignment = new ShareGroupDescribeResponseData.Assignment() + .setTopicPartitions(asList(topicPartitions)); + ShareGroupDescribeResponseData groupData = new ShareGroupDescribeResponseData(); + groupData.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId(GROUP_ID) + .setGroupState(GroupState.STABLE.toString()) + .setMembers(asList( + new ShareGroupDescribeResponseData.Member() + .setMemberId("0") + .setClientId("clientId0") + .setClientHost("clientHost") + .setAssignment(memberAssignment), + new ShareGroupDescribeResponseData.Member() + .setMemberId("1") + .setClientId("clientId1") + .setClientHost("clientHost") + .setAssignment(memberAssignment)))); + groupData.groups().add(new ShareGroupDescribeResponseData.DescribedGroup() + .setGroupId("group-1") + .setGroupState(GroupState.STABLE.toString()) + .setMembers(asList( + new ShareGroupDescribeResponseData.Member() + .setMemberId("0") + .setClientId("clientId0") + .setClientHost("clientHost") + .setAssignment(memberAssignment), + new ShareGroupDescribeResponseData.Member() + .setMemberId("1") + .setClientId("clientId1") + .setClientHost("clientHost") + .setAssignment(memberAssignment)))); + + env.kafkaClient().prepareResponse(new ShareGroupDescribeResponse(groupData)); + + Collection<String> groups = new HashSet<>(); + groups.add(GROUP_ID); + groups.add("group-1"); + final DescribeShareGroupsResult result = env.adminClient().describeShareGroups(groups); + assertEquals(2, result.describedGroups().size()); + assertEquals(groups, result.describedGroups().keySet()); + KafkaFuture<Map<String, ShareGroupDescription>> allFuture = result.all(); + assertDoesNotThrow(() -> allFuture.get()); + assertFalse(allFuture.isCompletedExceptionally()); + } + } + + @Test + public void testListShareGroups() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(4, 0), + AdminClientConfig.RETRIES_CONFIG, "2")) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + // Empty metadata response should be retried + env.kafkaClient().prepareResponse( + RequestTestUtils.metadataResponse( + Collections.emptyList(), + env.cluster().clusterResource().clusterId(), + -1, + Collections.emptyList())); + + env.kafkaClient().prepareResponse( + RequestTestUtils.metadataResponse( + env.cluster().nodes(), + env.cluster().clusterResource().clusterId(), + env.cluster().controller().id(), + Collections.emptyList())); + + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-1") + .setGroupType(GroupType.SHARE.toString()) + .setGroupState("Stable") + ))), + env.cluster().nodeById(0)); + + // handle retriable errors + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.COORDINATOR_NOT_AVAILABLE.code()) + .setGroups(Collections.emptyList()) + ), + env.cluster().nodeById(1)); + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.COORDINATOR_LOAD_IN_PROGRESS.code()) + .setGroups(Collections.emptyList()) + ), + env.cluster().nodeById(1)); + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-2") + .setGroupType(GroupType.SHARE.toString()) + .setGroupState("Stable"), + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-3") + .setGroupType(GroupType.SHARE.toString()) + .setGroupState("Stable") + ))), + env.cluster().nodeById(1)); + + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-4") + .setGroupType(GroupType.SHARE.toString()) + .setGroupState("Stable") + ))), + env.cluster().nodeById(2)); + + // fatal error + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse( + new ListGroupsResponseData() + .setErrorCode(Errors.UNKNOWN_SERVER_ERROR.code()) + .setGroups(Collections.emptyList())), + env.cluster().nodeById(3)); + + final ListGroupsResult result = env.adminClient().listGroups(ListGroupsOptions.forShareGroups()); + TestUtils.assertFutureThrows(UnknownServerException.class, result.all()); + + Collection<GroupListing> listings = result.valid().get(); + assertEquals(4, listings.size()); + + Set<String> groupIds = new HashSet<>(); + for (GroupListing listing : listings) { + groupIds.add(listing.groupId()); + assertTrue(listing.groupState().isPresent()); + } + + assertEquals(Set.of("share-group-1", "share-group-2", "share-group-3", "share-group-4"), groupIds); + assertEquals(1, result.errors().get().size()); + } + } + + @Test + public void testListShareGroupsMetadataFailure() throws Exception { + final Cluster cluster = mockCluster(3, 0); + final Time time = new MockTime(); + + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(time, cluster, + AdminClientConfig.RETRIES_CONFIG, "0")) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + // Empty metadata causes the request to fail since we have no list of brokers + // to send the ListGroups requests to + env.kafkaClient().prepareResponse( + RequestTestUtils.metadataResponse( + Collections.emptyList(), + env.cluster().clusterResource().clusterId(), + -1, + Collections.emptyList())); + + final ListGroupsResult result = env.adminClient().listGroups(ListGroupsOptions.forShareGroups()); + TestUtils.assertFutureThrows(KafkaException.class, result.all()); + } + } + + @Test + public void testListShareGroupsWithStates() throws Exception { + try (AdminClientUnitTestEnv env = new AdminClientUnitTestEnv(mockCluster(1, 0))) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + + env.kafkaClient().prepareResponse(prepareMetadataResponse(env.cluster(), Errors.NONE)); + + env.kafkaClient().prepareResponseFrom( + new ListGroupsResponse(new ListGroupsResponseData() + .setErrorCode(Errors.NONE.code()) + .setGroups(Arrays.asList( + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-1") + .setGroupType(GroupType.SHARE.toString()) + .setProtocolType("share") + .setGroupState("Stable"), + new ListGroupsResponseData.ListedGroup() + .setGroupId("share-group-2") + .setGroupType(GroupType.SHARE.toString()) + .setProtocolType("share") + .setGroupState("Empty")))), + env.cluster().nodeById(0)); + + final ListGroupsResult result = env.adminClient().listGroups(ListGroupsOptions.forShareGroups()); + Collection<GroupListing> listings = result.valid().get(); + + assertEquals(2, listings.size()); + List<GroupListing> expected = new ArrayList<>(); + expected.add(new GroupListing("share-group-1", Optional.of(GroupType.SHARE), "share", Optional.of(GroupState.STABLE))); + expected.add(new GroupListing("share-group-2", Optional.of(GroupType.SHARE), "share", Optional.of(GroupState.EMPTY))); + assertEquals(expected, listings); + assertEquals(0, result.errors().get().size()); + } + } + + @Test + public void testListShareGroupsWithStatesOlderBrokerVersion() { Review Comment: I agree this is a valid typo. However, my intention was to keep this PR purely mechanical - focused only on splitting the tests across files. We can address the additional improvements you suggested in a separate PR, so this one remains limited to the mechanical changes. Does that make sense? -- 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]
