diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java index ed6615ee17eec..ff7f4e661d692 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/Admin.java @@ -1663,6 +1663,26 @@ default FenceProducersResult fenceProducers(Collection transactionalIds) FenceProducersResult fenceProducers(Collection transactionalIds, FenceProducersOptions options); + /** + * List the client metrics configuration resources available in the cluster. + * + * @param options The options to use when listing the client metrics resources. + * @return The ListClientMetricsResourcesResult. + */ + ListClientMetricsResourcesResult listClientMetricsResources(ListClientMetricsResourcesOptions options); + + /** + * List the client metrics configuration resources available in the cluster with the default options. + *

+ * This is a convenience method for {@link #listClientMetricsResources(ListClientMetricsResourcesOptions)} + * with default options. See the overload for more details. + * + * @return The ListClientMetricsResourcesResult. + */ + default ListClientMetricsResourcesResult listClientMetricsResources() { + return listClientMetricsResources(new ListClientMetricsResourcesOptions()); + } + /** * Determines the client's unique client instance ID used for telemetry. This ID is unique to * this specific client instance and will not change after it is initially generated. diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ClientMetricsResourceListing.java b/clients/src/main/java/org/apache/kafka/clients/admin/ClientMetricsResourceListing.java new file mode 100644 index 0000000000000..0d902acd1c433 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ClientMetricsResourceListing.java @@ -0,0 +1,54 @@ +/* + * 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.common.annotation.InterfaceStability; + +import java.util.Objects; + +@InterfaceStability.Evolving +public class ClientMetricsResourceListing { + private final String name; + + public ClientMetricsResourceListing(String name) { + this.name = name; + } + + public String name() { + return name; + } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + ClientMetricsResourceListing that = (ClientMetricsResourceListing) o; + return Objects.equals(name, that.name); + } + + @Override + public int hashCode() { + return Objects.hash(name); + } + + @Override + public String toString() { + return "ClientMetricsResourceListing(" + + "name='" + name + + ')'; + } +} \ No newline at end of file diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java b/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java index 157bd70656406..9fc809dbddd81 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ForwardingAdmin.java @@ -278,6 +278,11 @@ public FenceProducersResult fenceProducers(Collection transactionalIds, return delegate.fenceProducers(transactionalIds, options); } + @Override + public ListClientMetricsResourcesResult listClientMetricsResources(ListClientMetricsResourcesOptions options) { + return delegate.listClientMetricsResources(options); + } + @Override public Uuid clientInstanceId(Duration timeout) { return delegate.clientInstanceId(timeout); diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index 9329b29b9785c..18bd108a7891a 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -141,6 +141,7 @@ import org.apache.kafka.common.message.DescribeUserScramCredentialsResponseData; import org.apache.kafka.common.message.ExpireDelegationTokenRequestData; import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity; +import org.apache.kafka.common.message.ListClientMetricsResourcesRequestData; import org.apache.kafka.common.message.ListGroupsRequestData; import org.apache.kafka.common.message.ListGroupsResponseData; import org.apache.kafka.common.message.ListPartitionReassignmentsRequestData; @@ -210,6 +211,8 @@ import org.apache.kafka.common.requests.IncrementalAlterConfigsRequest; import org.apache.kafka.common.requests.IncrementalAlterConfigsResponse; import org.apache.kafka.common.requests.JoinGroupRequest; +import org.apache.kafka.common.requests.ListClientMetricsResourcesRequest; +import org.apache.kafka.common.requests.ListClientMetricsResourcesResponse; import org.apache.kafka.common.requests.ListGroupsRequest; import org.apache.kafka.common.requests.ListGroupsResponse; import org.apache.kafka.common.requests.ListOffsetsRequest; @@ -4385,6 +4388,36 @@ public FenceProducersResult fenceProducers(Collection transactionalIds, return new FenceProducersResult(future.all()); } + @Override + public ListClientMetricsResourcesResult listClientMetricsResources(ListClientMetricsResourcesOptions options) { + final long now = time.milliseconds(); + final KafkaFutureImpl> future = new KafkaFutureImpl<>(); + runnable.call(new Call("listClientMetricsResources", calcDeadlineMs(now, options.timeoutMs()), + new LeastLoadedNodeProvider()) { + + @Override + ListClientMetricsResourcesRequest.Builder createRequest(int timeoutMs) { + return new ListClientMetricsResourcesRequest.Builder(new ListClientMetricsResourcesRequestData()); + } + + @Override + void handleResponse(AbstractResponse abstractResponse) { + ListClientMetricsResourcesResponse response = (ListClientMetricsResourcesResponse) abstractResponse; + if (response.error().isFailure()) { + future.completeExceptionally(response.error().exception()); + } else { + future.complete(response.clientMetricsResources()); + } + } + + @Override + void handleFailure(Throwable throwable) { + future.completeExceptionally(throwable); + } + }, now); + return new ListClientMetricsResourcesResult(future); + } + @Override public Uuid clientInstanceId(Duration timeout) { throw new UnsupportedOperationException(); diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesOptions.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesOptions.java new file mode 100644 index 0000000000000..333863e243e60 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesOptions.java @@ -0,0 +1,29 @@ +/* + * 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.common.annotation.InterfaceStability; + +/** + * Options for {@link Admin#listClientMetricsResources()}. + * + * The API of this class is evolving, see {@link Admin} for details. + */ +@InterfaceStability.Evolving +public class ListClientMetricsResourcesOptions extends AbstractOptions { +} \ No newline at end of file diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesResult.java b/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesResult.java new file mode 100644 index 0000000000000..c8609564fa5d7 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/admin/ListClientMetricsResourcesResult.java @@ -0,0 +1,57 @@ +/* + * 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.common.KafkaFuture; +import org.apache.kafka.common.annotation.InterfaceStability; +import org.apache.kafka.common.internals.KafkaFutureImpl; + +import java.util.Collection; + +/** + * The result of the {@link Admin#listClientMetricsResources()} call. + *

+ * The API of this class is evolving, see {@link Admin} for details. + */ +@InterfaceStability.Evolving +public class ListClientMetricsResourcesResult { + private final KafkaFuture> future; + + ListClientMetricsResourcesResult(KafkaFuture> future) { + this.future = future; + } + + /** + * Returns a future that yields either an exception, or the full set of client metrics + * listings. + * + * In the event of a failure, the future yields nothing but the first exception which + * occurred. + */ + public KafkaFuture> all() { + final KafkaFutureImpl> result = new KafkaFutureImpl<>(); + future.whenComplete((listings, throwable) -> { + if (throwable != null) { + result.completeExceptionally(throwable); + } else { + result.complete(listings); + } + }); + return result; + } +} \ No newline at end of file diff --git a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java index 1635da2e200da..a5c6ef5ea832a 100644 --- a/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java +++ b/clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java @@ -116,7 +116,8 @@ public enum ApiKeys { CONTROLLER_REGISTRATION(ApiMessageType.CONTROLLER_REGISTRATION), GET_TELEMETRY_SUBSCRIPTIONS(ApiMessageType.GET_TELEMETRY_SUBSCRIPTIONS), PUSH_TELEMETRY(ApiMessageType.PUSH_TELEMETRY), - ASSIGN_REPLICAS_TO_DIRS(ApiMessageType.ASSIGN_REPLICAS_TO_DIRS); + ASSIGN_REPLICAS_TO_DIRS(ApiMessageType.ASSIGN_REPLICAS_TO_DIRS), + LIST_CLIENT_METRICS_RESOURCES(ApiMessageType.LIST_CLIENT_METRICS_RESOURCES); private static final Map> APIS_BY_LISTENER = new EnumMap<>(ApiMessageType.ListenerType.class); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java index e71c5debf62c7..23f67cb5273e5 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AbstractRequest.java @@ -322,6 +322,8 @@ private static AbstractRequest doParseRequest(ApiKeys apiKey, short apiVersion, return PushTelemetryRequest.parse(buffer, apiVersion); case ASSIGN_REPLICAS_TO_DIRS: return AssignReplicasToDirsRequest.parse(buffer, apiVersion); + case LIST_CLIENT_METRICS_RESOURCES: + return ListClientMetricsResourcesRequest.parse(buffer, apiVersion); default: throw new AssertionError(String.format("ApiKey %s is not currently handled in `parseRequest`, the " + "code should be updated to do so.", apiKey)); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java index d849a0740b3a3..f99da4e2119df 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/AbstractResponse.java @@ -259,6 +259,8 @@ public static AbstractResponse parseResponse(ApiKeys apiKey, ByteBuffer response return PushTelemetryResponse.parse(responseBuffer, version); case ASSIGN_REPLICAS_TO_DIRS: return AssignReplicasToDirsResponse.parse(responseBuffer, version); + case LIST_CLIENT_METRICS_RESOURCES: + return ListClientMetricsResourcesResponse.parse(responseBuffer, version); default: throw new AssertionError(String.format("ApiKey %s is not currently handled in `parseResponse`, the " + "code should be updated to do so.", apiKey)); diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesRequest.java new file mode 100644 index 0000000000000..f2e1c7b7341de --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesRequest.java @@ -0,0 +1,77 @@ +/* + * 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.common.requests; + +import org.apache.kafka.common.message.ListClientMetricsResourcesRequestData; +import org.apache.kafka.common.message.ListClientMetricsResourcesResponseData; +import org.apache.kafka.common.protocol.ApiKeys; +import org.apache.kafka.common.protocol.ByteBufferAccessor; +import org.apache.kafka.common.protocol.Errors; + +import java.nio.ByteBuffer; + +public class ListClientMetricsResourcesRequest extends AbstractRequest { + public static class Builder extends AbstractRequest.Builder { + public final ListClientMetricsResourcesRequestData data; + + public Builder(ListClientMetricsResourcesRequestData data) { + super(ApiKeys.LIST_CLIENT_METRICS_RESOURCES); + this.data = data; + } + + @Override + public ListClientMetricsResourcesRequest build(short version) { + return new ListClientMetricsResourcesRequest(data, version); + } + + @Override + public String toString() { + return data.toString(); + } + } + + private final ListClientMetricsResourcesRequestData data; + + private ListClientMetricsResourcesRequest(ListClientMetricsResourcesRequestData data, short version) { + super(ApiKeys.LIST_CLIENT_METRICS_RESOURCES, version); + this.data = data; + } + + public ListClientMetricsResourcesRequestData data() { + return data; + } + + @Override + public ListClientMetricsResourcesResponse getErrorResponse(int throttleTimeMs, Throwable e) { + Errors error = Errors.forException(e); + ListClientMetricsResourcesResponseData response = new ListClientMetricsResourcesResponseData() + .setErrorCode(error.code()) + .setThrottleTimeMs(throttleTimeMs); + return new ListClientMetricsResourcesResponse(response); + } + + public static ListClientMetricsResourcesRequest parse(ByteBuffer buffer, short version) { + return new ListClientMetricsResourcesRequest(new ListClientMetricsResourcesRequestData( + new ByteBufferAccessor(buffer), version), version); + } + + @Override + public String toString(boolean verbose) { + return data.toString(); + } + +} \ No newline at end of file diff --git a/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesResponse.java new file mode 100644 index 0000000000000..1ea196058c9c7 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/common/requests/ListClientMetricsResourcesResponse.java @@ -0,0 +1,77 @@ +/* + * 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.common.requests; + +import org.apache.kafka.clients.admin.ClientMetricsResourceListing; +import org.apache.kafka.common.message.ListClientMetricsResourcesResponseData; +import org.apache.kafka.common.protocol.ApiKeys; +import org.apache.kafka.common.protocol.ByteBufferAccessor; +import org.apache.kafka.common.protocol.Errors; + +import java.nio.ByteBuffer; +import java.util.Collection; +import java.util.stream.Collectors; +import java.util.Map; + +public class ListClientMetricsResourcesResponse extends AbstractResponse { + private final ListClientMetricsResourcesResponseData data; + + public ListClientMetricsResourcesResponse(ListClientMetricsResourcesResponseData data) { + super(ApiKeys.LIST_CLIENT_METRICS_RESOURCES); + this.data = data; + } + + public ListClientMetricsResourcesResponseData data() { + return data; + } + + public ApiError error() { + return new ApiError(Errors.forCode(data.errorCode())); + } + + @Override + public Map errorCounts() { + return errorCounts(Errors.forCode(data.errorCode())); + } + + public static ListClientMetricsResourcesResponse parse(ByteBuffer buffer, short version) { + return new ListClientMetricsResourcesResponse(new ListClientMetricsResourcesResponseData( + new ByteBufferAccessor(buffer), version)); + } + + @Override + public String toString() { + return data.toString(); + } + + @Override + public int throttleTimeMs() { + return data.throttleTimeMs(); + } + + @Override + public void maybeSetThrottleTimeMs(int throttleTimeMs) { + data.setThrottleTimeMs(throttleTimeMs); + } + + public Collection clientMetricsResources() { + return data.clientMetricsResources() + .stream() + .map(entry -> new ClientMetricsResourceListing(entry.name())) + .collect(Collectors.toList()); + } +} \ No newline at end of file diff --git a/clients/src/main/resources/common/message/ListClientMetricsResourcesRequest.json b/clients/src/main/resources/common/message/ListClientMetricsResourcesRequest.json new file mode 100644 index 0000000000000..4b5640956b765 --- /dev/null +++ b/clients/src/main/resources/common/message/ListClientMetricsResourcesRequest.json @@ -0,0 +1,25 @@ +// 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. + +{ + "apiKey": 74, + "type": "request", + "listeners": ["broker"], + "name": "ListClientMetricsResourcesRequest", + "validVersions": "0", + "flexibleVersions": "0+", + "fields": [ + ] +} diff --git a/clients/src/main/resources/common/message/ListClientMetricsResourcesResponse.json b/clients/src/main/resources/common/message/ListClientMetricsResourcesResponse.json new file mode 100644 index 0000000000000..6d3321c17b61e --- /dev/null +++ b/clients/src/main/resources/common/message/ListClientMetricsResourcesResponse.json @@ -0,0 +1,30 @@ +// 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. + +{ + "apiKey": 74, + "type": "response", + "name": "ListClientMetricsResourcesResponse", + "validVersions": "0", + "flexibleVersions": "0+", + "fields": [ + { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+", + "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." }, + { "name": "ErrorCode", "type": "int16", "versions": "0+" }, + { "name": "ClientMetricsResources", "type": "[]ClientMetricsResource", "versions": "0+", "fields": [ + { "name": "Name", "type": "string", "versions": "0+" } + ]} + ] +} 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 19653017a985b..2c8ada5b55655 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 @@ -116,6 +116,7 @@ import org.apache.kafka.common.message.InitProducerIdResponseData; import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity; import org.apache.kafka.common.message.LeaveGroupResponseData; +import org.apache.kafka.common.message.ListClientMetricsResourcesResponseData; import org.apache.kafka.common.message.LeaveGroupResponseData.MemberResponse; import org.apache.kafka.common.message.ListGroupsResponseData; import org.apache.kafka.common.message.ListOffsetsResponseData; @@ -182,6 +183,8 @@ import org.apache.kafka.common.requests.JoinGroupRequest; import org.apache.kafka.common.requests.LeaveGroupRequest; import org.apache.kafka.common.requests.LeaveGroupResponse; +import org.apache.kafka.common.requests.ListClientMetricsResourcesRequest; +import org.apache.kafka.common.requests.ListClientMetricsResourcesResponse; import org.apache.kafka.common.requests.ListGroupsRequest; import org.apache.kafka.common.requests.ListGroupsResponse; import org.apache.kafka.common.requests.ListOffsetsRequest; @@ -7090,6 +7093,68 @@ private static MemberDescription convertToMemberDescriptions(DescribedGroupMembe assignment); } + @Test + public void testListClientMetricsResources() throws Exception { + try (AdminClientUnitTestEnv env = mockClientEnv()) { + List expected = Arrays.asList( + new ClientMetricsResourceListing("one"), + new ClientMetricsResourceListing("two") + ); + + ListClientMetricsResourcesResponseData responseData = + new ListClientMetricsResourcesResponseData().setErrorCode(Errors.NONE.code()); + + responseData.clientMetricsResources() + .add(new ListClientMetricsResourcesResponseData.ClientMetricsResource().setName("one")); + responseData.clientMetricsResources() + .add((new ListClientMetricsResourcesResponseData.ClientMetricsResource()).setName("two")); + + env.kafkaClient().prepareResponse( + request -> request instanceof ListClientMetricsResourcesRequest, + new ListClientMetricsResourcesResponse(responseData)); + + ListClientMetricsResourcesResult result = env.adminClient().listClientMetricsResources(); + assertEquals(new HashSet<>(expected), new HashSet<>(result.all().get())); + } + } + + @Test + public void testListClientMetricsResourcesEmpty() throws Exception { + try (AdminClientUnitTestEnv env = mockClientEnv()) { + List expected = Collections.emptyList(); + + ListClientMetricsResourcesResponseData responseData = + new ListClientMetricsResourcesResponseData().setErrorCode(Errors.NONE.code()); + + env.kafkaClient().prepareResponse( + request -> request instanceof ListClientMetricsResourcesRequest, + new ListClientMetricsResourcesResponse(responseData)); + + ListClientMetricsResourcesResult result = env.adminClient().listClientMetricsResources(); + assertEquals(new HashSet<>(expected), new HashSet<>(result.all().get())); + } + } + + @Test + public void testListClientMetricsResourcesNotSupported() throws Exception { + try (AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().prepareResponse( + request -> request instanceof ListClientMetricsResourcesRequest, + prepareListClientMetricsResourcesResponse(Errors.UNSUPPORTED_VERSION)); + + ListClientMetricsResourcesResult result = env.adminClient().listClientMetricsResources(); + + // Validate response + assertNotNull(result.all()); + TestUtils.assertFutureThrows(result.all(), Errors.UNSUPPORTED_VERSION.exception().getClass()); + } + } + + private static ListClientMetricsResourcesResponse prepareListClientMetricsResourcesResponse(Errors error) { + return new ListClientMetricsResourcesResponse(new ListClientMetricsResourcesResponseData() + .setErrorCode(error.code())); + } + @SafeVarargs private static void assertCollectionIs(Collection collection, T... elements) { for (T element : elements) { diff --git a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java index 52188c68df40d..2a33f7f324c56 100644 --- a/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java +++ b/clients/src/test/java/org/apache/kafka/clients/admin/MockAdminClient.java @@ -1301,6 +1301,11 @@ public FenceProducersResult fenceProducers(Collection transactionalIds, throw new UnsupportedOperationException("Not implemented yet"); } + @Override + public ListClientMetricsResourcesResult listClientMetricsResources(ListClientMetricsResourcesOptions options) { + throw new UnsupportedOperationException("Not implemented yet"); + } + @Override synchronized public void close(Duration timeout) {} diff --git a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java index a802ee3e58035..46be917dfb8eb 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/RequestResponseTest.java @@ -167,6 +167,8 @@ import org.apache.kafka.common.message.LeaderAndIsrResponseData.LeaderAndIsrTopicErrorCollection; import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity; import org.apache.kafka.common.message.LeaveGroupResponseData; +import org.apache.kafka.common.message.ListClientMetricsResourcesRequestData; +import org.apache.kafka.common.message.ListClientMetricsResourcesResponseData; import org.apache.kafka.common.message.ListGroupsRequestData; import org.apache.kafka.common.message.ListGroupsResponseData; import org.apache.kafka.common.message.ListOffsetsRequestData.ListOffsetsPartition; @@ -1074,6 +1076,7 @@ private AbstractRequest getRequest(ApiKeys apikey, short version) { case GET_TELEMETRY_SUBSCRIPTIONS: return createGetTelemetrySubscriptionsRequest(version); case PUSH_TELEMETRY: return createPushTelemetryRequest(version); case ASSIGN_REPLICAS_TO_DIRS: return createAssignReplicasToDirsRequest(version); + case LIST_CLIENT_METRICS_RESOURCES: return createListClientMetricsResourcesRequest(version); default: throw new IllegalArgumentException("Unknown API key " + apikey); } } @@ -1154,6 +1157,7 @@ private AbstractResponse getResponse(ApiKeys apikey, short version) { case GET_TELEMETRY_SUBSCRIPTIONS: return createGetTelemetrySubscriptionsResponse(); case PUSH_TELEMETRY: return createPushTelemetryResponse(); case ASSIGN_REPLICAS_TO_DIRS: return createAssignReplicasToDirsResponse(); + case LIST_CLIENT_METRICS_RESOURCES: return createListClientMetricsResourcesResponse(); default: throw new IllegalArgumentException("Unknown API key " + apikey); } } @@ -3616,6 +3620,17 @@ private PushTelemetryResponse createPushTelemetryResponse() { return new PushTelemetryResponse(response); } + private ListClientMetricsResourcesRequest createListClientMetricsResourcesRequest(short version) { + return new ListClientMetricsResourcesRequest.Builder(new ListClientMetricsResourcesRequestData()).build(version); + } + + private ListClientMetricsResourcesResponse createListClientMetricsResourcesResponse() { + ListClientMetricsResourcesResponseData response = new ListClientMetricsResourcesResponseData(); + response.setErrorCode(Errors.NONE.code()); + response.setThrottleTimeMs(0); + return new ListClientMetricsResourcesResponse(response); + } + @Test public void testInvalidSaslHandShakeRequest() { AbstractRequest request = new SaslHandshakeRequest.Builder( diff --git a/core/src/main/scala/kafka/network/RequestConvertToJson.scala b/core/src/main/scala/kafka/network/RequestConvertToJson.scala index 6844acf66cba7..4d08a6243f113 100644 --- a/core/src/main/scala/kafka/network/RequestConvertToJson.scala +++ b/core/src/main/scala/kafka/network/RequestConvertToJson.scala @@ -32,15 +32,19 @@ object RequestConvertToJson { case req: AllocateProducerIdsRequest => AllocateProducerIdsRequestDataJsonConverter.write(req.data, request.version) case req: AlterClientQuotasRequest => AlterClientQuotasRequestDataJsonConverter.write(req.data, request.version) case req: AlterConfigsRequest => AlterConfigsRequestDataJsonConverter.write(req.data, request.version) - case req: AlterPartitionRequest => AlterPartitionRequestDataJsonConverter.write(req.data, request.version) case req: AlterPartitionReassignmentsRequest => AlterPartitionReassignmentsRequestDataJsonConverter.write(req.data, request.version) + case req: AlterPartitionRequest => AlterPartitionRequestDataJsonConverter.write(req.data, request.version) case req: AlterReplicaLogDirsRequest => AlterReplicaLogDirsRequestDataJsonConverter.write(req.data, request.version) case res: AlterUserScramCredentialsRequest => AlterUserScramCredentialsRequestDataJsonConverter.write(res.data, request.version) case req: ApiVersionsRequest => ApiVersionsRequestDataJsonConverter.write(req.data, request.version) + case req: AssignReplicasToDirsRequest => AssignReplicasToDirsRequestDataJsonConverter.write(req.data, request.version) case req: BeginQuorumEpochRequest => BeginQuorumEpochRequestDataJsonConverter.write(req.data, request.version) case req: BrokerHeartbeatRequest => BrokerHeartbeatRequestDataJsonConverter.write(req.data, request.version) case req: BrokerRegistrationRequest => BrokerRegistrationRequestDataJsonConverter.write(req.data, request.version) + case req: ConsumerGroupDescribeRequest => ConsumerGroupDescribeRequestDataJsonConverter.write(req.data, request.version) + case req: ConsumerGroupHeartbeatRequest => ConsumerGroupHeartbeatRequestDataJsonConverter.write(req.data, request.version) case req: ControlledShutdownRequest => ControlledShutdownRequestDataJsonConverter.write(req.data, request.version) + case req: ControllerRegistrationRequest => ControllerRegistrationRequestDataJsonConverter.write(req.data, request.version) case req: CreateAclsRequest => CreateAclsRequestDataJsonConverter.write(req.data, request.version) case req: CreateDelegationTokenRequest => CreateDelegationTokenRequestDataJsonConverter.write(req.data, request.version) case req: CreatePartitionsRequest => CreatePartitionsRequestDataJsonConverter.write(req.data, request.version) @@ -51,18 +55,22 @@ object RequestConvertToJson { case req: DeleteTopicsRequest => DeleteTopicsRequestDataJsonConverter.write(req.data, request.version) case req: DescribeAclsRequest => DescribeAclsRequestDataJsonConverter.write(req.data, request.version) case req: DescribeClientQuotasRequest => DescribeClientQuotasRequestDataJsonConverter.write(req.data, request.version) + case req: DescribeClusterRequest => DescribeClusterRequestDataJsonConverter.write(req.data, request.version) case req: DescribeConfigsRequest => DescribeConfigsRequestDataJsonConverter.write(req.data, request.version) case req: DescribeDelegationTokenRequest => DescribeDelegationTokenRequestDataJsonConverter.write(req.data, request.version) case req: DescribeGroupsRequest => DescribeGroupsRequestDataJsonConverter.write(req.data, request.version) case req: DescribeLogDirsRequest => DescribeLogDirsRequestDataJsonConverter.write(req.data, request.version) + case req: DescribeProducersRequest => DescribeProducersRequestDataJsonConverter.write(req.data, request.version) case req: DescribeQuorumRequest => DescribeQuorumRequestDataJsonConverter.write(req.data, request.version) - case res: DescribeUserScramCredentialsRequest => DescribeUserScramCredentialsRequestDataJsonConverter.write(res.data, request.version) + case req: DescribeTransactionsRequest => DescribeTransactionsRequestDataJsonConverter.write(req.data, request.version) + case req: DescribeUserScramCredentialsRequest => DescribeUserScramCredentialsRequestDataJsonConverter.write(req.data, request.version) case req: ElectLeadersRequest => ElectLeadersRequestDataJsonConverter.write(req.data, request.version) - case req: EndTxnRequest => EndTxnRequestDataJsonConverter.write(req.data, request.version) case req: EndQuorumEpochRequest => EndQuorumEpochRequestDataJsonConverter.write(req.data, request.version) + case req: EndTxnRequest => EndTxnRequestDataJsonConverter.write(req.data, request.version) case req: EnvelopeRequest => EnvelopeRequestDataJsonConverter.write(req.data, request.version) case req: ExpireDelegationTokenRequest => ExpireDelegationTokenRequestDataJsonConverter.write(req.data, request.version) case req: FetchRequest => FetchRequestDataJsonConverter.write(req.data, request.version) + case req: FetchSnapshotRequest => FetchSnapshotRequestDataJsonConverter.write(req.data, request.version) case req: FindCoordinatorRequest => FindCoordinatorRequestDataJsonConverter.write(req.data, request.version) case req: GetTelemetrySubscriptionsRequest => GetTelemetrySubscriptionsRequestDataJsonConverter.write(req.data, request.version) case req: HeartbeatRequest => HeartbeatRequestDataJsonConverter.write(req.data, request.version) @@ -71,9 +79,11 @@ object RequestConvertToJson { case req: JoinGroupRequest => JoinGroupRequestDataJsonConverter.write(req.data, request.version) case req: LeaderAndIsrRequest => LeaderAndIsrRequestDataJsonConverter.write(req.data, request.version) case req: LeaveGroupRequest => LeaveGroupRequestDataJsonConverter.write(req.data, request.version) + case req: ListClientMetricsResourcesRequest => ListClientMetricsResourcesRequestDataJsonConverter.write(req.data, request.version) case req: ListGroupsRequest => ListGroupsRequestDataJsonConverter.write(req.data, request.version) case req: ListOffsetsRequest => ListOffsetsRequestDataJsonConverter.write(req.data, request.version) case req: ListPartitionReassignmentsRequest => ListPartitionReassignmentsRequestDataJsonConverter.write(req.data, request.version) + case req: ListTransactionsRequest => ListTransactionsRequestDataJsonConverter.write(req.data, request.version) case req: MetadataRequest => MetadataRequestDataJsonConverter.write(req.data, request.version) case req: OffsetCommitRequest => OffsetCommitRequestDataJsonConverter.write(req.data, request.version) case req: OffsetDeleteRequest => OffsetDeleteRequestDataJsonConverter.write(req.data, request.version) @@ -92,15 +102,6 @@ object RequestConvertToJson { case req: UpdateMetadataRequest => UpdateMetadataRequestDataJsonConverter.write(req.data, request.version) case req: VoteRequest => VoteRequestDataJsonConverter.write(req.data, request.version) case req: WriteTxnMarkersRequest => WriteTxnMarkersRequestDataJsonConverter.write(req.data, request.version) - case req: FetchSnapshotRequest => FetchSnapshotRequestDataJsonConverter.write(req.data, request.version) - case req: DescribeClusterRequest => DescribeClusterRequestDataJsonConverter.write(req.data, request.version) - case req: DescribeProducersRequest => DescribeProducersRequestDataJsonConverter.write(req.data, request.version) - case req: DescribeTransactionsRequest => DescribeTransactionsRequestDataJsonConverter.write(req.data, request.version) - case req: ListTransactionsRequest => ListTransactionsRequestDataJsonConverter.write(req.data, request.version) - case req: ConsumerGroupHeartbeatRequest => ConsumerGroupHeartbeatRequestDataJsonConverter.write(req.data, request.version) - case req: ConsumerGroupDescribeRequest => ConsumerGroupDescribeRequestDataJsonConverter.write(req.data, request.version) - case req: ControllerRegistrationRequest => ControllerRegistrationRequestDataJsonConverter.write(req.data, request.version) - case req: AssignReplicasToDirsRequest => AssignReplicasToDirsRequestDataJsonConverter.write(req.data, request.version) case _ => throw new IllegalStateException(s"ApiKey ${request.apiKey} is not currently handled in `request`, the " + "code should be updated to do so."); } @@ -113,15 +114,19 @@ object RequestConvertToJson { case res: AllocateProducerIdsResponse => AllocateProducerIdsResponseDataJsonConverter.write(res.data, version) case res: AlterClientQuotasResponse => AlterClientQuotasResponseDataJsonConverter.write(res.data, version) case res: AlterConfigsResponse => AlterConfigsResponseDataJsonConverter.write(res.data, version) - case res: AlterPartitionResponse => AlterPartitionResponseDataJsonConverter.write(res.data, version) case res: AlterPartitionReassignmentsResponse => AlterPartitionReassignmentsResponseDataJsonConverter.write(res.data, version) + case res: AlterPartitionResponse => AlterPartitionResponseDataJsonConverter.write(res.data, version) case res: AlterReplicaLogDirsResponse => AlterReplicaLogDirsResponseDataJsonConverter.write(res.data, version) case res: AlterUserScramCredentialsResponse => AlterUserScramCredentialsResponseDataJsonConverter.write(res.data, version) case res: ApiVersionsResponse => ApiVersionsResponseDataJsonConverter.write(res.data, version) + case res: AssignReplicasToDirsResponse => AssignReplicasToDirsResponseDataJsonConverter.write(res.data, version) case res: BeginQuorumEpochResponse => BeginQuorumEpochResponseDataJsonConverter.write(res.data, version) case res: BrokerHeartbeatResponse => BrokerHeartbeatResponseDataJsonConverter.write(res.data, version) case res: BrokerRegistrationResponse => BrokerRegistrationResponseDataJsonConverter.write(res.data, version) + case res: ConsumerGroupDescribeResponse => ConsumerGroupDescribeResponseDataJsonConverter.write(res.data, version) + case res: ConsumerGroupHeartbeatResponse => ConsumerGroupHeartbeatResponseDataJsonConverter.write(res.data, version) case res: ControlledShutdownResponse => ControlledShutdownResponseDataJsonConverter.write(res.data, version) + case req: ControllerRegistrationResponse => ControllerRegistrationResponseDataJsonConverter.write(req.data, version) case res: CreateAclsResponse => CreateAclsResponseDataJsonConverter.write(res.data, version) case res: CreateDelegationTokenResponse => CreateDelegationTokenResponseDataJsonConverter.write(res.data, version) case res: CreatePartitionsResponse => CreatePartitionsResponseDataJsonConverter.write(res.data, version) @@ -132,18 +137,22 @@ object RequestConvertToJson { case res: DeleteTopicsResponse => DeleteTopicsResponseDataJsonConverter.write(res.data, version) case res: DescribeAclsResponse => DescribeAclsResponseDataJsonConverter.write(res.data, version) case res: DescribeClientQuotasResponse => DescribeClientQuotasResponseDataJsonConverter.write(res.data, version) + case res: DescribeClusterResponse => DescribeClusterResponseDataJsonConverter.write(res.data, version) case res: DescribeConfigsResponse => DescribeConfigsResponseDataJsonConverter.write(res.data, version) case res: DescribeDelegationTokenResponse => DescribeDelegationTokenResponseDataJsonConverter.write(res.data, version) case res: DescribeGroupsResponse => DescribeGroupsResponseDataJsonConverter.write(res.data, version) case res: DescribeLogDirsResponse => DescribeLogDirsResponseDataJsonConverter.write(res.data, version) + case res: DescribeProducersResponse => DescribeProducersResponseDataJsonConverter.write(res.data, version) case res: DescribeQuorumResponse => DescribeQuorumResponseDataJsonConverter.write(res.data, version) + case res: DescribeTransactionsResponse => DescribeTransactionsResponseDataJsonConverter.write(res.data, version) case res: DescribeUserScramCredentialsResponse => DescribeUserScramCredentialsResponseDataJsonConverter.write(res.data, version) case res: ElectLeadersResponse => ElectLeadersResponseDataJsonConverter.write(res.data, version) - case res: EndTxnResponse => EndTxnResponseDataJsonConverter.write(res.data, version) case res: EndQuorumEpochResponse => EndQuorumEpochResponseDataJsonConverter.write(res.data, version) + case res: EndTxnResponse => EndTxnResponseDataJsonConverter.write(res.data, version) case res: EnvelopeResponse => EnvelopeResponseDataJsonConverter.write(res.data, version) case res: ExpireDelegationTokenResponse => ExpireDelegationTokenResponseDataJsonConverter.write(res.data, version) case res: FetchResponse => FetchResponseDataJsonConverter.write(res.data, version, false) + case res: FetchSnapshotResponse => FetchSnapshotResponseDataJsonConverter.write(res.data, version) case res: FindCoordinatorResponse => FindCoordinatorResponseDataJsonConverter.write(res.data, version) case res: GetTelemetrySubscriptionsResponse => GetTelemetrySubscriptionsResponseDataJsonConverter.write(res.data, version) case res: HeartbeatResponse => HeartbeatResponseDataJsonConverter.write(res.data, version) @@ -152,9 +161,11 @@ object RequestConvertToJson { case res: JoinGroupResponse => JoinGroupResponseDataJsonConverter.write(res.data, version) case res: LeaderAndIsrResponse => LeaderAndIsrResponseDataJsonConverter.write(res.data, version) case res: LeaveGroupResponse => LeaveGroupResponseDataJsonConverter.write(res.data, version) + case res: ListClientMetricsResourcesResponse => ListClientMetricsResourcesResponseDataJsonConverter.write(res.data, version) case res: ListGroupsResponse => ListGroupsResponseDataJsonConverter.write(res.data, version) case res: ListOffsetsResponse => ListOffsetsResponseDataJsonConverter.write(res.data, version) case res: ListPartitionReassignmentsResponse => ListPartitionReassignmentsResponseDataJsonConverter.write(res.data, version) + case res: ListTransactionsResponse => ListTransactionsResponseDataJsonConverter.write(res.data, version) case res: MetadataResponse => MetadataResponseDataJsonConverter.write(res.data, version) case res: OffsetCommitResponse => OffsetCommitResponseDataJsonConverter.write(res.data, version) case res: OffsetDeleteResponse => OffsetDeleteResponseDataJsonConverter.write(res.data, version) @@ -173,15 +184,6 @@ object RequestConvertToJson { case res: UpdateMetadataResponse => UpdateMetadataResponseDataJsonConverter.write(res.data, version) case res: WriteTxnMarkersResponse => WriteTxnMarkersResponseDataJsonConverter.write(res.data, version) case res: VoteResponse => VoteResponseDataJsonConverter.write(res.data, version) - case res: FetchSnapshotResponse => FetchSnapshotResponseDataJsonConverter.write(res.data, version) - case res: DescribeClusterResponse => DescribeClusterResponseDataJsonConverter.write(res.data, version) - case res: DescribeProducersResponse => DescribeProducersResponseDataJsonConverter.write(res.data, version) - case res: DescribeTransactionsResponse => DescribeTransactionsResponseDataJsonConverter.write(res.data, version) - case res: ListTransactionsResponse => ListTransactionsResponseDataJsonConverter.write(res.data, version) - case res: ConsumerGroupHeartbeatResponse => ConsumerGroupHeartbeatResponseDataJsonConverter.write(res.data, version) - case res: ConsumerGroupDescribeResponse => ConsumerGroupDescribeResponseDataJsonConverter.write(res.data, version) - case req: ControllerRegistrationResponse => ControllerRegistrationResponseDataJsonConverter.write(req.data, version) - case res: AssignReplicasToDirsResponse => AssignReplicasToDirsResponseDataJsonConverter.write(res.data, version) case _ => throw new IllegalStateException(s"ApiKey ${response.apiKey} is not currently handled in `response`, the " + "code should be updated to do so."); } diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 1ac38a0895d21..c3168a188c9a1 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -247,6 +247,7 @@ class KafkaApis(val requestChannel: RequestChannel, case ApiKeys.CONSUMER_GROUP_DESCRIBE => handleConsumerGroupDescribe(request).exceptionally(handleError) case ApiKeys.GET_TELEMETRY_SUBSCRIPTIONS => handleGetTelemetrySubscriptionsRequest(request) case ApiKeys.PUSH_TELEMETRY => handlePushTelemetryRequest(request) + case ApiKeys.LIST_CLIENT_METRICS_RESOURCES => handleListClientMetricsResources(request) case _ => throw new IllegalStateException(s"No handler for request api key ${request.header.apiKey}") } } catch { @@ -3787,6 +3788,19 @@ class KafkaApis(val requestChannel: RequestChannel, CompletableFuture.completedFuture[Unit](()) } + // Just a placeholder for now. + def handleListClientMetricsResources(request: RequestChannel.Request): Unit = { + val listClientMetricsResourcesRequest = request.body[ListClientMetricsResourcesRequest] + + if (!authHelper.authorize(request.context, DESCRIBE_CONFIGS, CLUSTER, CLUSTER_NAME)) { + requestHelper.sendMaybeThrottle(request, listClientMetricsResourcesRequest.getErrorResponse(Errors.CLUSTER_AUTHORIZATION_FAILED.exception)) + } else { + // Just return an empty list in the placeholder + val data = new ListClientMetricsResourcesResponseData() + requestHelper.sendMaybeThrottle(request, new ListClientMetricsResourcesResponse(data)) + } + } + private def updateRecordConversionStats(request: RequestChannel.Request, tp: TopicPartition, conversionStats: RecordValidationStats): Unit = { diff --git a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala index a817881b17a7c..255ada6012bbe 100644 --- a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala +++ b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala @@ -722,6 +722,9 @@ class RequestQuotaTest extends BaseRequestTest { case ApiKeys.ASSIGN_REPLICAS_TO_DIRS => new AssignReplicasToDirsRequest.Builder(new AssignReplicasToDirsRequestData()) + case ApiKeys.LIST_CLIENT_METRICS_RESOURCES => + new ListClientMetricsResourcesRequest.Builder(new ListClientMetricsResourcesRequestData()) + case _ => throw new IllegalArgumentException("Unsupported API key " + apiKey) }