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 566b0e4b44b19..431591b72d65a 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 @@ -42,7 +42,6 @@ import org.apache.kafka.common.ConsumerGroupState; import org.apache.kafka.common.ElectionType; import org.apache.kafka.common.KafkaException; -import org.apache.kafka.common.KafkaFuture; import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.Node; @@ -102,6 +101,8 @@ import org.apache.kafka.common.message.DeleteRecordsResponseData.DeleteRecordsTopicResult; import org.apache.kafka.common.message.DeleteTopicsRequestData; import org.apache.kafka.common.message.DeleteTopicsResponseData.DeletableTopicResult; +import org.apache.kafka.common.message.DescribeConfigsRequestData; +import org.apache.kafka.common.message.DescribeConfigsResponseData; import org.apache.kafka.common.message.DescribeGroupsRequestData; import org.apache.kafka.common.message.DescribeGroupsResponseData.DescribedGroup; import org.apache.kafka.common.message.DescribeGroupsResponseData.DescribedGroupMember; @@ -1916,129 +1917,95 @@ void handleFailure(Throwable throwable) { @Override public DescribeConfigsResult describeConfigs(Collection configResources, final DescribeConfigsOptions options) { - final Map> unifiedRequestFutures = new HashMap<>(); - final Map> brokerFutures = new HashMap<>(configResources.size()); - - // The BROKER resources which we want to describe. We must make a separate DescribeConfigs - // request for every BROKER resource we want to describe. - final Collection brokerResources = new ArrayList<>(); - - // The non-BROKER resources which we want to describe. These resources can be described by a - // single, unified DescribeConfigs request. - final Collection unifiedRequestResources = new ArrayList<>(configResources.size()); + // Partition the requested config resources based on which broker they must be sent to with the + // null broker being used for config resources which can be obtained from any broker + final Map>> brokerFutures = new HashMap<>(configResources.size()); for (ConfigResource resource : configResources) { - if (dependsOnSpecificNode(resource)) { - brokerFutures.put(resource, new KafkaFutureImpl<>()); - brokerResources.add(resource); - } else { - unifiedRequestFutures.put(resource, new KafkaFutureImpl<>()); - unifiedRequestResources.add(resource); - } + Integer broker = nodeFor(resource); + brokerFutures.compute(broker, (key, value) -> { + if (value == null) { + value = new HashMap<>(); + } + value.put(resource, new KafkaFutureImpl<>()); + return value; + }); } final long now = time.milliseconds(); - if (!unifiedRequestResources.isEmpty()) { + for (Map.Entry>> entry : brokerFutures.entrySet()) { + Integer broker = entry.getKey(); + Map> unified = entry.getValue(); + runnable.call(new Call("describeConfigs", calcDeadlineMs(now, options.timeoutMs()), - new LeastLoadedNodeProvider()) { + broker != null ? new ConstantNodeIdProvider(broker) : new LeastLoadedNodeProvider()) { @Override DescribeConfigsRequest.Builder createRequest(int timeoutMs) { - return new DescribeConfigsRequest.Builder(unifiedRequestResources) - .includeSynonyms(options.includeSynonyms()) - .includeDocumentation(options.includeDocumentation()); + return new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() + .setResources(unified.keySet().stream() + .map(config -> + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceName(config.name()) + .setResourceType(config.type().id())) + .collect(Collectors.toList())) + .setIncludeSynonyms(options.includeSynonyms()) + .setIncludeDocumentation(options.includeDocumentation())); } @Override void handleResponse(AbstractResponse abstractResponse) { DescribeConfigsResponse response = (DescribeConfigsResponse) abstractResponse; - for (Map.Entry> entry : unifiedRequestFutures.entrySet()) { + for (Map.Entry entry : response.resultMap().entrySet()) { ConfigResource configResource = entry.getKey(); - KafkaFutureImpl future = entry.getValue(); - DescribeConfigsResponse.Config config = response.config(configResource); - if (config == null) { - future.completeExceptionally(new UnknownServerException( - "Malformed broker response: missing config for " + configResource)); - continue; - } - if (config.error().isFailure()) { - future.completeExceptionally(config.error().exception()); - continue; - } - List configEntries = new ArrayList<>(); - for (DescribeConfigsResponse.ConfigEntry configEntry : config.entries()) { - configEntries.add(new ConfigEntry(configEntry.name(), - configEntry.value(), configSource(configEntry.source()), - configEntry.isSensitive(), configEntry.isReadOnly(), - configSynonyms(configEntry), configType(configEntry.type()), - configEntry.documentation())); + DescribeConfigsResponseData.DescribeConfigsResult describeConfigsResult = entry.getValue(); + KafkaFutureImpl future = unified.get(configResource); + if (future == null) { + if (broker != null) { + log.warn("The config {} in the response from broker {} is not in the request", + configResource, broker); + } else { + log.warn("The config {} in the response from the least loaded broker is not in the request", + configResource); + } + } else { + if (describeConfigsResult.errorCode() != Errors.NONE.code()) { + future.completeExceptionally(Errors.forCode(describeConfigsResult.errorCode()) + .exception(describeConfigsResult.errorMessage())); + } else { + future.complete(describeConfigResult(describeConfigsResult)); + } } - future.complete(new Config(configEntries)); } + completeUnrealizedFutures( + unified.entrySet().stream(), + configResource -> "The broker response did not contain a result for config resource " + configResource); } @Override void handleFailure(Throwable throwable) { - completeAllExceptionally(unifiedRequestFutures.values(), throwable); + completeAllExceptionally(unified.values(), throwable); } }, now); } - for (Map.Entry> entry : brokerFutures.entrySet()) { - final KafkaFutureImpl brokerFuture = entry.getValue(); - final ConfigResource resource = entry.getKey(); - final int nodeId = Integer.parseInt(resource.name()); - runnable.call(new Call("describeBrokerConfigs", calcDeadlineMs(now, options.timeoutMs()), - new ConstantNodeIdProvider(nodeId)) { - - @Override - DescribeConfigsRequest.Builder createRequest(int timeoutMs) { - return new DescribeConfigsRequest.Builder(Collections.singleton(resource)) - .includeSynonyms(options.includeSynonyms()) - .includeDocumentation(options.includeDocumentation()); - } - - @Override - void handleResponse(AbstractResponse abstractResponse) { - DescribeConfigsResponse response = (DescribeConfigsResponse) abstractResponse; - DescribeConfigsResponse.Config config = response.configs().get(resource); - - if (config == null) { - brokerFuture.completeExceptionally(new UnknownServerException( - "Malformed broker response: missing config for " + resource)); - return; - } - if (config.error().isFailure()) - brokerFuture.completeExceptionally(config.error().exception()); - else { - List configEntries = new ArrayList<>(); - for (DescribeConfigsResponse.ConfigEntry configEntry : config.entries()) { - configEntries.add(new ConfigEntry(configEntry.name(), configEntry.value(), - configSource(configEntry.source()), configEntry.isSensitive(), configEntry.isReadOnly(), - configSynonyms(configEntry), configType(configEntry.type()), configEntry.documentation())); - } - brokerFuture.complete(new Config(configEntries)); - } - } - - @Override - void handleFailure(Throwable throwable) { - brokerFuture.completeExceptionally(throwable); - } - }, now); - } - final Map> allFutures = new HashMap<>(); - allFutures.putAll(brokerFutures); - allFutures.putAll(unifiedRequestFutures); - return new DescribeConfigsResult(allFutures); + return new DescribeConfigsResult(new HashMap<>(brokerFutures.entrySet().stream() + .flatMap(x -> x.getValue().entrySet().stream()) + .collect(Collectors.toMap(Entry::getKey, Entry::getValue)))); } - private List configSynonyms(DescribeConfigsResponse.ConfigEntry configEntry) { - List synonyms = new ArrayList<>(configEntry.synonyms().size()); - for (DescribeConfigsResponse.ConfigSynonym synonym : configEntry.synonyms()) { - synonyms.add(new ConfigEntry.ConfigSynonym(synonym.name(), synonym.value(), configSource(synonym.source()))); - } - return synonyms; + private Config describeConfigResult(DescribeConfigsResponseData.DescribeConfigsResult describeConfigsResult) { + return new Config(describeConfigsResult.configs().stream().map(config -> new ConfigEntry( + config.name(), + config.value(), + DescribeConfigsResponse.ConfigSource.forId(config.configSource()).source(), + config.isSensitive(), + config.readOnly(), + (config.synonyms().stream().map(synonym -> new ConfigEntry.ConfigSynonym(synonym.name(), synonym.value(), + DescribeConfigsResponse.ConfigSource.forId(synonym.source()).source()))).collect(Collectors.toList()), + DescribeConfigsResponse.ConfigType.forId(config.configType()).type(), + config.documentation() + )).collect(Collectors.toList())); } private ConfigEntry.ConfigSource configSource(DescribeConfigsResponse.ConfigSource source) { @@ -2068,46 +2035,6 @@ private ConfigEntry.ConfigSource configSource(DescribeConfigsResponse.ConfigSour return configSource; } - private ConfigEntry.ConfigType configType(DescribeConfigsResponse.ConfigType type) { - if (type == null) { - return ConfigEntry.ConfigType.UNKNOWN; - } - - ConfigEntry.ConfigType configType; - switch (type) { - case BOOLEAN: - configType = ConfigEntry.ConfigType.BOOLEAN; - break; - case CLASS: - configType = ConfigEntry.ConfigType.CLASS; - break; - case DOUBLE: - configType = ConfigEntry.ConfigType.DOUBLE; - break; - case INT: - configType = ConfigEntry.ConfigType.INT; - break; - case LIST: - configType = ConfigEntry.ConfigType.LIST; - break; - case LONG: - configType = ConfigEntry.ConfigType.LONG; - break; - case PASSWORD: - configType = ConfigEntry.ConfigType.PASSWORD; - break; - case SHORT: - configType = ConfigEntry.ConfigType.SHORT; - break; - case STRING: - configType = ConfigEntry.ConfigType.STRING; - break; - default: - configType = ConfigEntry.ConfigType.UNKNOWN; - } - return configType; - } - @Override @Deprecated public AlterConfigsResult alterConfigs(Map configs, final AlterConfigsOptions options) { @@ -2118,8 +2045,9 @@ public AlterConfigsResult alterConfigs(Map configs, fina final Collection unifiedRequestResources = new ArrayList<>(); for (ConfigResource resource : configs.keySet()) { - if (dependsOnSpecificNode(resource)) { - NodeProvider nodeProvider = new ConstantNodeIdProvider(Integer.parseInt(resource.name())); + Integer node = nodeFor(resource); + if (node != null) { + NodeProvider nodeProvider = new ConstantNodeIdProvider(node); allFutures.putAll(alterConfigs(configs, options, Collections.singleton(resource), nodeProvider)); } else unifiedRequestResources.add(resource); @@ -2183,8 +2111,9 @@ public AlterConfigsResult incrementalAlterConfigs(Map unifiedRequestResources = new ArrayList<>(); for (ConfigResource resource : configs.keySet()) { - if (dependsOnSpecificNode(resource)) { - NodeProvider nodeProvider = new ConstantNodeIdProvider(Integer.parseInt(resource.name())); + Integer node = nodeFor(resource); + if (node != null) { + NodeProvider nodeProvider = new ConstantNodeIdProvider(node); allFutures.putAll(incrementalAlterConfigs(configs, options, Collections.singleton(resource), nodeProvider)); } else unifiedRequestResources.add(resource); @@ -3685,11 +3614,16 @@ private void handleNotControllerError(Errors error) throws ApiException { } /** - * Returns a boolean indicating whether the resource needs to go to a specific node + * Returns the broker id pertaining to the given resource, or null if the resource is not associated + * with a particular broker. */ - private boolean dependsOnSpecificNode(ConfigResource resource) { - return (resource.type() == ConfigResource.Type.BROKER && !resource.isDefault()) - || resource.type() == ConfigResource.Type.BROKER_LOGGER; + private Integer nodeFor(ConfigResource resource) { + if ((resource.type() == ConfigResource.Type.BROKER && !resource.isDefault()) + || resource.type() == ConfigResource.Type.BROKER_LOGGER) { + return Integer.valueOf(resource.name()); + } else { + return null; + } } private List getMembersFromGroup(String groupId) { 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 4cb9119d325a0..670de3755764a 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 @@ -53,6 +53,8 @@ import org.apache.kafka.common.message.DescribeAclsResponseData; import org.apache.kafka.common.message.DescribeClientQuotasRequestData; import org.apache.kafka.common.message.DescribeClientQuotasResponseData; +import org.apache.kafka.common.message.DescribeConfigsRequestData; +import org.apache.kafka.common.message.DescribeConfigsResponseData; import org.apache.kafka.common.message.DescribeDelegationTokenRequestData; import org.apache.kafka.common.message.DescribeDelegationTokenResponseData; import org.apache.kafka.common.message.DescribeGroupsRequestData; @@ -112,8 +114,6 @@ import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.protocol.types.Type; import org.apache.kafka.common.record.RecordBatch; -import org.apache.kafka.common.requests.DescribeConfigsRequest; -import org.apache.kafka.common.requests.DescribeConfigsResponse; import org.apache.kafka.common.requests.FetchRequest; import org.apache.kafka.common.requests.FetchResponse; import org.apache.kafka.common.requests.ListOffsetRequest; @@ -184,8 +184,8 @@ public Struct parseResponse(short version, ByteBuffer buffer) { DESCRIBE_ACLS(29, "DescribeAcls", DescribeAclsRequestData.SCHEMAS, DescribeAclsResponseData.SCHEMAS), CREATE_ACLS(30, "CreateAcls", CreateAclsRequestData.SCHEMAS, CreateAclsResponseData.SCHEMAS), DELETE_ACLS(31, "DeleteAcls", DeleteAclsRequestData.SCHEMAS, DeleteAclsResponseData.SCHEMAS), - DESCRIBE_CONFIGS(32, "DescribeConfigs", DescribeConfigsRequest.schemaVersions(), - DescribeConfigsResponse.schemaVersions()), + DESCRIBE_CONFIGS(32, "DescribeConfigs", DescribeConfigsRequestData.SCHEMAS, + DescribeConfigsResponseData.SCHEMAS), ALTER_CONFIGS(33, "AlterConfigs", AlterConfigsRequestData.SCHEMAS, AlterConfigsResponseData.SCHEMAS), ALTER_REPLICA_LOG_DIRS(34, "AlterReplicaLogDirs", AlterReplicaLogDirsRequestData.SCHEMAS, 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 f17af31178108..32581421bb940 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 @@ -151,7 +151,7 @@ public static AbstractResponse parseResponse(ApiKeys apiKey, Struct struct, shor case DELETE_ACLS: return new DeleteAclsResponse(struct, version); case DESCRIBE_CONFIGS: - return new DescribeConfigsResponse(struct); + return new DescribeConfigsResponse(struct, version); case ALTER_CONFIGS: return new AlterConfigsResponse(struct, version); case ALTER_REPLICA_LOG_DIRS: diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsRequest.java index 8ea76301a29c3..f17896e66e0ca 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsRequest.java @@ -16,194 +16,64 @@ */ package org.apache.kafka.common.requests; -import org.apache.kafka.common.config.ConfigResource; +import org.apache.kafka.common.message.DescribeConfigsRequestData; +import org.apache.kafka.common.message.DescribeConfigsResponseData; import org.apache.kafka.common.protocol.ApiKeys; -import org.apache.kafka.common.protocol.types.ArrayOf; -import org.apache.kafka.common.protocol.types.Field; -import org.apache.kafka.common.protocol.types.Schema; +import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.protocol.types.Struct; import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Objects; - -import static org.apache.kafka.common.protocol.types.Type.BOOLEAN; -import static org.apache.kafka.common.protocol.types.Type.INT8; -import static org.apache.kafka.common.protocol.types.Type.STRING; +import java.util.stream.Collectors; public class DescribeConfigsRequest extends AbstractRequest { - private static final String RESOURCES_KEY_NAME = "resources"; - private static final String INCLUDE_SYNONYMS = "include_synonyms"; - private static final String RESOURCE_TYPE_KEY_NAME = "resource_type"; - private static final String RESOURCE_NAME_KEY_NAME = "resource_name"; - private static final String CONFIG_NAMES_KEY_NAME = "config_names"; - private static final String INCLUDE_DOCUMENTATION = "include_documentation"; - - private static final Schema DESCRIBE_CONFIGS_REQUEST_RESOURCE_V0 = new Schema( - new Field(RESOURCE_TYPE_KEY_NAME, INT8), - new Field(RESOURCE_NAME_KEY_NAME, STRING), - new Field(CONFIG_NAMES_KEY_NAME, ArrayOf.nullable(STRING))); - - private static final Schema DESCRIBE_CONFIGS_REQUEST_V0 = new Schema( - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_REQUEST_RESOURCE_V0), "An array of config resources to be returned.")); - - private static final Schema DESCRIBE_CONFIGS_REQUEST_V1 = new Schema( - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_REQUEST_RESOURCE_V0), "An array of config resources to be returned."), - new Field(INCLUDE_SYNONYMS, BOOLEAN)); - - /** - * The version number is bumped to indicate that on quota violation brokers send out responses before throttling. - */ - private static final Schema DESCRIBE_CONFIGS_REQUEST_V2 = DESCRIBE_CONFIGS_REQUEST_V1; - - private static final Schema DESCRIBE_CONFIGS_REQUEST_V3 = new Schema( - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_REQUEST_RESOURCE_V0), "An array of config resources to be returned."), - new Field(INCLUDE_SYNONYMS, BOOLEAN), - new Field(INCLUDE_DOCUMENTATION, BOOLEAN)); - - public static Schema[] schemaVersions() { - return new Schema[] { - DESCRIBE_CONFIGS_REQUEST_V0, - DESCRIBE_CONFIGS_REQUEST_V1, - DESCRIBE_CONFIGS_REQUEST_V2, - DESCRIBE_CONFIGS_REQUEST_V3 - }; - } - public static class Builder extends AbstractRequest.Builder { - private final Map> resourceToConfigNames; - private boolean includeSynonyms; - private boolean includeDocumentation; + private final DescribeConfigsRequestData data; - public Builder(Map> resourceToConfigNames) { + public Builder(DescribeConfigsRequestData data) { super(ApiKeys.DESCRIBE_CONFIGS); - this.resourceToConfigNames = Objects.requireNonNull(resourceToConfigNames, "resourceToConfigNames"); - } - - public Builder includeSynonyms(boolean includeSynonyms) { - this.includeSynonyms = includeSynonyms; - return this; - } - - public Builder includeDocumentation(boolean includeDocumentation) { - this.includeDocumentation = includeDocumentation; - return this; - } - - public Builder(Collection resources) { - this(toResourceToConfigNames(resources)); - } - - private static Map> toResourceToConfigNames(Collection resources) { - Map> result = new HashMap<>(resources.size()); - for (ConfigResource resource : resources) - result.put(resource, null); - return result; + this.data = data; } @Override public DescribeConfigsRequest build(short version) { - return new DescribeConfigsRequest( - version, resourceToConfigNames, includeSynonyms, includeDocumentation); + return new DescribeConfigsRequest(data, version); } } - private final Map> resourceToConfigNames; - private final boolean includeSynonyms; - private final boolean includeDocumentation; + private final DescribeConfigsRequestData data; - public DescribeConfigsRequest( - short version, Map> resourceToConfigNames, - boolean includeSynonyms) { - this(version, resourceToConfigNames, includeSynonyms, false); - } - public DescribeConfigsRequest( - short version, Map> resourceToConfigNames, - boolean includeSynonyms, boolean includeDocumentation) { + public DescribeConfigsRequest(DescribeConfigsRequestData data, short version) { super(ApiKeys.DESCRIBE_CONFIGS, version); - this.resourceToConfigNames = Objects.requireNonNull(resourceToConfigNames, "resourceToConfigNames"); - this.includeSynonyms = includeSynonyms; - this.includeDocumentation = includeDocumentation; + this.data = data; } public DescribeConfigsRequest(Struct struct, short version) { super(ApiKeys.DESCRIBE_CONFIGS, version); - Object[] resourcesArray = struct.getArray(RESOURCES_KEY_NAME); - resourceToConfigNames = new HashMap<>(resourcesArray.length); - for (Object resourceObj : resourcesArray) { - Struct resourceStruct = (Struct) resourceObj; - ConfigResource.Type resourceType = ConfigResource.Type.forId(resourceStruct.getByte(RESOURCE_TYPE_KEY_NAME)); - String resourceName = resourceStruct.getString(RESOURCE_NAME_KEY_NAME); - - Object[] configNamesArray = resourceStruct.getArray(CONFIG_NAMES_KEY_NAME); - List configNames = null; - if (configNamesArray != null) { - configNames = new ArrayList<>(configNamesArray.length); - for (Object configNameObj : configNamesArray) - configNames.add((String) configNameObj); - } - - resourceToConfigNames.put(new ConfigResource(resourceType, resourceName), configNames); - } - this.includeSynonyms = struct.hasField(INCLUDE_SYNONYMS) ? struct.getBoolean(INCLUDE_SYNONYMS) : false; - this.includeDocumentation = struct.hasField(INCLUDE_DOCUMENTATION) ? struct.getBoolean(INCLUDE_DOCUMENTATION) : false; - } - - public Collection resources() { - return resourceToConfigNames.keySet(); + this.data = new DescribeConfigsRequestData(struct, version); } - /** - * Return null if all config names should be returned. - */ - public Collection configNames(ConfigResource resource) { - return resourceToConfigNames.get(resource); - } - - public boolean includeSynonyms() { - return includeSynonyms; - } - - public boolean includeDocumentation() { - return includeDocumentation; + public DescribeConfigsRequestData data() { + return data; } @Override protected Struct toStruct() { - Struct struct = new Struct(ApiKeys.DESCRIBE_CONFIGS.requestSchema(version())); - List resourceStructs = new ArrayList<>(resources().size()); - for (Map.Entry> entry : resourceToConfigNames.entrySet()) { - ConfigResource resource = entry.getKey(); - Struct resourceStruct = struct.instance(RESOURCES_KEY_NAME); - resourceStruct.set(RESOURCE_TYPE_KEY_NAME, resource.type().id()); - resourceStruct.set(RESOURCE_NAME_KEY_NAME, resource.name()); - - String[] configNames = entry.getValue() == null ? null : entry.getValue().toArray(new String[0]); - resourceStruct.set(CONFIG_NAMES_KEY_NAME, configNames); - - resourceStructs.add(resourceStruct); - } - struct.set(RESOURCES_KEY_NAME, resourceStructs.toArray(new Struct[0])); - struct.setIfExists(INCLUDE_SYNONYMS, includeSynonyms); - struct.setIfExists(INCLUDE_DOCUMENTATION, includeDocumentation); - return struct; + return data.toStruct(version()); } @Override public DescribeConfigsResponse getErrorResponse(int throttleTimeMs, Throwable e) { - ApiError error = ApiError.fromThrowable(e); - Map errors = new HashMap<>(resources().size()); - DescribeConfigsResponse.Config config = new DescribeConfigsResponse.Config(error, - Collections.emptyList()); - for (ConfigResource resource : resources()) - errors.put(resource, config); - return new DescribeConfigsResponse(throttleTimeMs, errors); + Errors error = Errors.forException(e); + return new DescribeConfigsResponse(new DescribeConfigsResponseData() + .setThrottleTimeMs(throttleTimeMs) + .setResults(data.resources().stream().map(result -> { + return new DescribeConfigsResponseData.DescribeConfigsResult().setErrorCode(error.code()) + .setErrorMessage(error.message()) + .setResourceName(result.resourceName()) + .setResourceType(result.resourceType()); + }).collect(Collectors.toList()) + )); } public static DescribeConfigsRequest parse(ByteBuffer buffer, short version) { diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java index bd46bc32d7800..3767226395e81 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeConfigsResponse.java @@ -18,127 +18,22 @@ package org.apache.kafka.common.requests; import org.apache.kafka.common.config.ConfigResource; +import org.apache.kafka.common.message.DescribeConfigsResponseData; +import org.apache.kafka.common.message.DescribeConfigsResponseData.DescribeConfigsResult; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; -import org.apache.kafka.common.protocol.types.ArrayOf; -import org.apache.kafka.common.protocol.types.Field; -import org.apache.kafka.common.protocol.types.Schema; import org.apache.kafka.common.protocol.types.Struct; import java.nio.ByteBuffer; -import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.HashMap; -import java.util.List; import java.util.Map; import java.util.Objects; - -import static org.apache.kafka.common.protocol.CommonFields.ERROR_CODE; -import static org.apache.kafka.common.protocol.CommonFields.ERROR_MESSAGE; -import static org.apache.kafka.common.protocol.CommonFields.THROTTLE_TIME_MS; -import static org.apache.kafka.common.protocol.types.Type.BOOLEAN; -import static org.apache.kafka.common.protocol.types.Type.INT8; -import static org.apache.kafka.common.protocol.types.Type.NULLABLE_STRING; -import static org.apache.kafka.common.protocol.types.Type.STRING; +import java.util.function.Function; +import java.util.stream.Collectors; public class DescribeConfigsResponse extends AbstractResponse { - private static final String RESOURCES_KEY_NAME = "resources"; - - private static final String RESOURCE_TYPE_KEY_NAME = "resource_type"; - private static final String RESOURCE_NAME_KEY_NAME = "resource_name"; - - private static final String CONFIG_ENTRIES_KEY_NAME = "config_entries"; - - private static final String CONFIG_NAME_KEY_NAME = "config_name"; - private static final String CONFIG_VALUE_KEY_NAME = "config_value"; - private static final String IS_SENSITIVE_KEY_NAME = "is_sensitive"; - private static final String IS_DEFAULT_KEY_NAME = "is_default"; - private static final String READ_ONLY_KEY_NAME = "read_only"; - private static final String CONFIG_TYPE_KEY_NAME = "config_type"; - private static final String CONFIG_DOCUMENTATION_KEY_NAME = "config_documentation"; - - private static final String CONFIG_SYNONYMS_KEY_NAME = "config_synonyms"; - private static final String CONFIG_SOURCE_KEY_NAME = "config_source"; - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTRY_V0 = new Schema( - new Field(CONFIG_NAME_KEY_NAME, STRING), - new Field(CONFIG_VALUE_KEY_NAME, NULLABLE_STRING), - new Field(READ_ONLY_KEY_NAME, BOOLEAN), - new Field(IS_DEFAULT_KEY_NAME, BOOLEAN), - new Field(IS_SENSITIVE_KEY_NAME, BOOLEAN)); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_SYNONYM_V1 = new Schema( - new Field(CONFIG_NAME_KEY_NAME, STRING), - new Field(CONFIG_VALUE_KEY_NAME, NULLABLE_STRING), - new Field(CONFIG_SOURCE_KEY_NAME, INT8)); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTRY_V1 = new Schema( - new Field(CONFIG_NAME_KEY_NAME, STRING), - new Field(CONFIG_VALUE_KEY_NAME, NULLABLE_STRING), - new Field(READ_ONLY_KEY_NAME, BOOLEAN), - new Field(CONFIG_SOURCE_KEY_NAME, INT8), - new Field(IS_SENSITIVE_KEY_NAME, BOOLEAN), - new Field(CONFIG_SYNONYMS_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_SYNONYM_V1))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTRY_V3 = new Schema( - new Field(CONFIG_NAME_KEY_NAME, STRING), - new Field(CONFIG_VALUE_KEY_NAME, NULLABLE_STRING), - new Field(READ_ONLY_KEY_NAME, BOOLEAN), - new Field(CONFIG_SOURCE_KEY_NAME, INT8), - new Field(IS_SENSITIVE_KEY_NAME, BOOLEAN), - new Field(CONFIG_SYNONYMS_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_SYNONYM_V1)), - new Field(CONFIG_TYPE_KEY_NAME, INT8), - new Field(CONFIG_DOCUMENTATION_KEY_NAME, NULLABLE_STRING)); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTITY_V0 = new Schema( - ERROR_CODE, - ERROR_MESSAGE, - new Field(RESOURCE_TYPE_KEY_NAME, INT8), - new Field(RESOURCE_NAME_KEY_NAME, STRING), - new Field(CONFIG_ENTRIES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTRY_V0))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTITY_V1 = new Schema( - ERROR_CODE, - ERROR_MESSAGE, - new Field(RESOURCE_TYPE_KEY_NAME, INT8), - new Field(RESOURCE_NAME_KEY_NAME, STRING), - new Field(CONFIG_ENTRIES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTRY_V1))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_ENTITY_V3 = new Schema( - ERROR_CODE, - ERROR_MESSAGE, - new Field(RESOURCE_TYPE_KEY_NAME, INT8), - new Field(RESOURCE_NAME_KEY_NAME, STRING), - new Field(CONFIG_ENTRIES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTRY_V3))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_V0 = new Schema( - THROTTLE_TIME_MS, - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTITY_V0))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_V1 = new Schema( - THROTTLE_TIME_MS, - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTITY_V1))); - - private static final Schema DESCRIBE_CONFIGS_RESPONSE_V3 = new Schema( - THROTTLE_TIME_MS, - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_CONFIGS_RESPONSE_ENTITY_V3))); - - /** - * The version number is bumped to indicate that on quota violation brokers send out responses before throttling. - */ - private static final Schema DESCRIBE_CONFIGS_RESPONSE_V2 = DESCRIBE_CONFIGS_RESPONSE_V1; - - public static Schema[] schemaVersions() { - return new Schema[]{ - DESCRIBE_CONFIGS_RESPONSE_V0, - DESCRIBE_CONFIGS_RESPONSE_V1, - DESCRIBE_CONFIGS_RESPONSE_V2, - DESCRIBE_CONFIGS_RESPONSE_V3 - }; - } - public static class Config { private final ApiError error; private final Collection entries; @@ -219,47 +114,63 @@ public String documentation() { } public enum ConfigSource { - UNKNOWN_CONFIG((byte) 0), - TOPIC_CONFIG((byte) 1), - DYNAMIC_BROKER_CONFIG((byte) 2), - DYNAMIC_DEFAULT_BROKER_CONFIG((byte) 3), - STATIC_BROKER_CONFIG((byte) 4), - DEFAULT_CONFIG((byte) 5), - DYNAMIC_BROKER_LOGGER_CONFIG((byte) 6); + UNKNOWN((byte) 0, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.UNKNOWN), + TOPIC_CONFIG((byte) 1, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.DYNAMIC_TOPIC_CONFIG), + DYNAMIC_BROKER_CONFIG((byte) 2, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.DYNAMIC_BROKER_CONFIG), + DYNAMIC_DEFAULT_BROKER_CONFIG((byte) 3, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.DYNAMIC_DEFAULT_BROKER_CONFIG), + STATIC_BROKER_CONFIG((byte) 4, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.STATIC_BROKER_CONFIG), + DEFAULT_CONFIG((byte) 5, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.DEFAULT_CONFIG), + DYNAMIC_BROKER_LOGGER_CONFIG((byte) 6, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource.DYNAMIC_BROKER_LOGGER_CONFIG); final byte id; + private final org.apache.kafka.clients.admin.ConfigEntry.ConfigSource source; private static final ConfigSource[] VALUES = values(); - ConfigSource(byte id) { + ConfigSource(byte id, org.apache.kafka.clients.admin.ConfigEntry.ConfigSource source) { this.id = id; + this.source = source; + } + + public byte id() { + return id; } public static ConfigSource forId(byte id) { if (id < 0) throw new IllegalArgumentException("id should be positive, id: " + id); if (id >= VALUES.length) - return UNKNOWN_CONFIG; + return UNKNOWN; return VALUES[id]; } + + public org.apache.kafka.clients.admin.ConfigEntry.ConfigSource source() { + return source; + } } public enum ConfigType { - UNKNOWN((byte) 0), - BOOLEAN((byte) 1), - STRING((byte) 2), - INT((byte) 3), - SHORT((byte) 4), - LONG((byte) 5), - DOUBLE((byte) 6), - LIST((byte) 7), - CLASS((byte) 8), - PASSWORD((byte) 9); + UNKNOWN((byte) 0, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.UNKNOWN), + BOOLEAN((byte) 1, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.BOOLEAN), + STRING((byte) 2, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.STRING), + INT((byte) 3, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.INT), + SHORT((byte) 4, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.SHORT), + LONG((byte) 5, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.LONG), + DOUBLE((byte) 6, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.DOUBLE), + LIST((byte) 7, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.LIST), + CLASS((byte) 8, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.CLASS), + PASSWORD((byte) 9, org.apache.kafka.clients.admin.ConfigEntry.ConfigType.PASSWORD); final byte id; + final org.apache.kafka.clients.admin.ConfigEntry.ConfigType type; private static final ConfigType[] VALUES = values(); - ConfigType(byte id) { + ConfigType(byte id, org.apache.kafka.clients.admin.ConfigEntry.ConfigType type) { this.id = id; + this.type = type; + } + + public byte id() { + return id; } public static ConfigType forId(byte id) { @@ -269,6 +180,10 @@ public static ConfigType forId(byte id) { return UNKNOWN; return VALUES[id]; } + + public org.apache.kafka.clients.admin.ConfigEntry.ConfigType type() { + return type; + } } public static class ConfigSynonym { @@ -293,158 +208,66 @@ public ConfigSource source() { } } + public Map resultMap() { + return data().results().stream().collect(Collectors.toMap( + configsResult -> + new ConfigResource(ConfigResource.Type.forId(configsResult.resourceType()), + configsResult.resourceName()), + Function.identity())); + } - private final int throttleTimeMs; - private final Map configs; + private final DescribeConfigsResponseData data; - public DescribeConfigsResponse(int throttleTimeMs, Map configs) { - this.throttleTimeMs = throttleTimeMs; - this.configs = Objects.requireNonNull(configs, "configs"); + public DescribeConfigsResponse(DescribeConfigsResponseData data) { + this.data = data; } - public DescribeConfigsResponse(Struct struct) { - throttleTimeMs = struct.get(THROTTLE_TIME_MS); - Object[] resourcesArray = struct.getArray(RESOURCES_KEY_NAME); - configs = new HashMap<>(resourcesArray.length); - for (Object resourceObj : resourcesArray) { - Struct resourceStruct = (Struct) resourceObj; - - ApiError error = new ApiError(resourceStruct); - ConfigResource.Type resourceType = ConfigResource.Type.forId(resourceStruct.getByte(RESOURCE_TYPE_KEY_NAME)); - String resourceName = resourceStruct.getString(RESOURCE_NAME_KEY_NAME); - ConfigResource resource = new ConfigResource(resourceType, resourceName); - - Object[] configEntriesArray = resourceStruct.getArray(CONFIG_ENTRIES_KEY_NAME); - List configEntries = new ArrayList<>(configEntriesArray.length); - for (Object configEntriesObj: configEntriesArray) { - Struct configEntriesStruct = (Struct) configEntriesObj; - String configName = configEntriesStruct.getString(CONFIG_NAME_KEY_NAME); - String configValue = configEntriesStruct.getString(CONFIG_VALUE_KEY_NAME); - boolean isSensitive = configEntriesStruct.getBoolean(IS_SENSITIVE_KEY_NAME); - ConfigType type = ConfigType.UNKNOWN; - if (configEntriesStruct.hasField(CONFIG_TYPE_KEY_NAME)) { - type = ConfigType.forId(configEntriesStruct.getByte(CONFIG_TYPE_KEY_NAME)); - } - String documentation = null; - if (configEntriesStruct.hasField(CONFIG_DOCUMENTATION_KEY_NAME)) { - documentation = configEntriesStruct.getString(CONFIG_DOCUMENTATION_KEY_NAME); - } - ConfigSource configSource; - if (configEntriesStruct.hasField(CONFIG_SOURCE_KEY_NAME)) - configSource = ConfigSource.forId(configEntriesStruct.getByte(CONFIG_SOURCE_KEY_NAME)); - else if (configEntriesStruct.hasField(IS_DEFAULT_KEY_NAME)) { - if (configEntriesStruct.getBoolean(IS_DEFAULT_KEY_NAME)) - configSource = ConfigSource.DEFAULT_CONFIG; - else { - switch (resourceType) { - case BROKER: - configSource = ConfigSource.STATIC_BROKER_CONFIG; - break; - case TOPIC: - configSource = ConfigSource.TOPIC_CONFIG; - break; - default: - configSource = ConfigSource.UNKNOWN_CONFIG; - break; + public DescribeConfigsResponse(Struct struct, short version) { + this.data = new DescribeConfigsResponseData(struct, version); + if (version == 0) { + for (DescribeConfigsResult result : data.results()) { + for (DescribeConfigsResponseData.DescribeConfigsResourceResult config : result.configs()) { + if (config.isDefault()) { + config.setConfigSource(ConfigSource.DEFAULT_CONFIG.id); + } else { + if (result.resourceType() == ConfigResource.Type.BROKER.id()) { + config.setConfigSource(ConfigSource.STATIC_BROKER_CONFIG.id); + } else if (result.resourceType() == ConfigResource.Type.TOPIC.id()) { + config.setConfigSource(ConfigSource.TOPIC_CONFIG.id); + } else { + config.setConfigSource(ConfigSource.UNKNOWN.id); } } - } else - throw new IllegalStateException("Config entry should contain either is_default or config_source"); - boolean readOnly = configEntriesStruct.getBoolean(READ_ONLY_KEY_NAME); - Collection synonyms; - if (configEntriesStruct.hasField(CONFIG_SYNONYMS_KEY_NAME)) { - Object[] synonymsArray = configEntriesStruct.getArray(CONFIG_SYNONYMS_KEY_NAME); - synonyms = new ArrayList<>(synonymsArray.length); - for (Object synonymObj: synonymsArray) { - Struct synonymStruct = (Struct) synonymObj; - String synonymConfigName = synonymStruct.getString(CONFIG_NAME_KEY_NAME); - String synonymConfigValue = synonymStruct.getString(CONFIG_VALUE_KEY_NAME); - ConfigSource source = ConfigSource.forId(synonymStruct.getByte(CONFIG_SOURCE_KEY_NAME)); - synonyms.add(new ConfigSynonym(synonymConfigName, synonymConfigValue, source)); - } - } else { - synonyms = Collections.emptyList(); } - configEntries.add(new ConfigEntry(configName, configValue, configSource, isSensitive, readOnly, synonyms, type, documentation)); } - Config config = new Config(error, configEntries); - configs.put(resource, config); } } - public Map configs() { - return configs; - } - - public Config config(ConfigResource resource) { - return configs.get(resource); + public DescribeConfigsResponseData data() { + return data; } @Override public int throttleTimeMs() { - return throttleTimeMs; + return data.throttleTimeMs(); } @Override public Map errorCounts() { Map errorCounts = new HashMap<>(); - configs.values().forEach(response -> - updateErrorCounts(errorCounts, response.error.error()) + data.results().forEach(response -> + updateErrorCounts(errorCounts, Errors.forCode(response.errorCode())) ); return errorCounts; } @Override protected Struct toStruct(short version) { - Struct struct = new Struct(ApiKeys.DESCRIBE_CONFIGS.responseSchema(version)); - struct.set(THROTTLE_TIME_MS, throttleTimeMs); - List resourceStructs = new ArrayList<>(configs.size()); - for (Map.Entry entry : configs.entrySet()) { - Struct resourceStruct = struct.instance(RESOURCES_KEY_NAME); - - ConfigResource resource = entry.getKey(); - resourceStruct.set(RESOURCE_TYPE_KEY_NAME, resource.type().id()); - resourceStruct.set(RESOURCE_NAME_KEY_NAME, resource.name()); - - Config config = entry.getValue(); - config.error.write(resourceStruct); - - List configEntryStructs = new ArrayList<>(config.entries.size()); - for (ConfigEntry configEntry : config.entries) { - Struct configEntriesStruct = resourceStruct.instance(CONFIG_ENTRIES_KEY_NAME); - configEntriesStruct.set(CONFIG_NAME_KEY_NAME, configEntry.name); - configEntriesStruct.set(CONFIG_VALUE_KEY_NAME, configEntry.value); - configEntriesStruct.set(IS_SENSITIVE_KEY_NAME, configEntry.isSensitive); - configEntriesStruct.setIfExists(CONFIG_SOURCE_KEY_NAME, configEntry.source.id); - configEntriesStruct.setIfExists(IS_DEFAULT_KEY_NAME, configEntry.source == ConfigSource.DEFAULT_CONFIG); - configEntriesStruct.set(READ_ONLY_KEY_NAME, configEntry.readOnly); - if (configEntriesStruct.hasField(CONFIG_TYPE_KEY_NAME) && configEntry.type != null) { - configEntriesStruct.set(CONFIG_TYPE_KEY_NAME, configEntry.type.id); - } - configEntriesStruct.setIfExists(CONFIG_DOCUMENTATION_KEY_NAME, configEntry.documentation); - configEntryStructs.add(configEntriesStruct); - if (configEntriesStruct.hasField(CONFIG_SYNONYMS_KEY_NAME)) { - List configSynonymStructs = new ArrayList<>(configEntry.synonyms.size()); - for (ConfigSynonym synonym : configEntry.synonyms) { - Struct configSynonymStruct = configEntriesStruct.instance(CONFIG_SYNONYMS_KEY_NAME); - configSynonymStruct.set(CONFIG_NAME_KEY_NAME, synonym.name); - configSynonymStruct.set(CONFIG_VALUE_KEY_NAME, synonym.value); - configSynonymStruct.set(CONFIG_SOURCE_KEY_NAME, synonym.source.id); - configSynonymStructs.add(configSynonymStruct); - } - configEntriesStruct.set(CONFIG_SYNONYMS_KEY_NAME, configSynonymStructs.toArray(new Struct[0])); - } - } - resourceStruct.set(CONFIG_ENTRIES_KEY_NAME, configEntryStructs.toArray(new Struct[0])); - - resourceStructs.add(resourceStruct); - } - struct.set(RESOURCES_KEY_NAME, resourceStructs.toArray(new Struct[0])); - return struct; + return data.toStruct(version); } public static DescribeConfigsResponse parse(ByteBuffer buffer, short version) { - return new DescribeConfigsResponse(ApiKeys.DESCRIBE_CONFIGS.parseResponse(version, buffer)); + return new DescribeConfigsResponse(ApiKeys.DESCRIBE_CONFIGS.parseResponse(version, buffer), version); } @Override 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 32c6c084a32fb..1052bcaaad92a 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 @@ -81,6 +81,7 @@ import org.apache.kafka.common.message.DeleteTopicsResponseData; import org.apache.kafka.common.message.DeleteTopicsResponseData.DeletableTopicResult; import org.apache.kafka.common.message.DescribeAclsResponseData; +import org.apache.kafka.common.message.DescribeConfigsResponseData; import org.apache.kafka.common.message.DescribeGroupsResponseData; import org.apache.kafka.common.message.DescribeGroupsResponseData.DescribedGroupMember; import org.apache.kafka.common.message.ElectLeadersResponseData.PartitionResult; @@ -174,6 +175,7 @@ import java.util.stream.Stream; import static java.util.Arrays.asList; +import static java.util.Collections.emptyList; import static java.util.Collections.singletonList; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignablePartitionResponse; import static org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.ReassignableTopicResponse; @@ -969,15 +971,87 @@ public void testElectLeaders() throws Exception { } @Test - public void testDescribeConfigs() throws Exception { + public void testDescribeBrokerConfigs() throws Exception { + ConfigResource broker0Resource = new ConfigResource(ConfigResource.Type.BROKER, "0"); + ConfigResource broker1Resource = new ConfigResource(ConfigResource.Type.BROKER, "1"); try (AdminClientUnitTestEnv env = mockClientEnv()) { env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); - env.kafkaClient().prepareResponse(new DescribeConfigsResponse(0, - Collections.singletonMap(new ConfigResource(ConfigResource.Type.BROKER, "0"), - new DescribeConfigsResponse.Config(ApiError.NONE, Collections.emptySet())))); - DescribeConfigsResult result2 = env.adminClient().describeConfigs(Collections.singleton( - new ConfigResource(ConfigResource.Type.BROKER, "0"))); - result2.all().get(); + env.kafkaClient().prepareResponseFrom(new DescribeConfigsResponse( + new DescribeConfigsResponseData().setResults(asList(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(broker0Resource.name()).setResourceType(broker0Resource.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList())))), env.cluster().nodeById(0)); + env.kafkaClient().prepareResponseFrom(new DescribeConfigsResponse( + new DescribeConfigsResponseData().setResults(asList(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(broker1Resource.name()).setResourceType(broker1Resource.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList())))), env.cluster().nodeById(1)); + Map> result = env.adminClient().describeConfigs(asList( + broker0Resource, + broker1Resource)).values(); + assertEquals(new HashSet<>(asList(broker0Resource, broker1Resource)), result.keySet()); + result.get(broker0Resource).get(); + result.get(broker1Resource).get(); + } + } + + @Test + public void testDescribeBrokerAndLogConfigs() throws Exception { + ConfigResource brokerResource = new ConfigResource(ConfigResource.Type.BROKER, "0"); + ConfigResource brokerLoggerResource = new ConfigResource(ConfigResource.Type.BROKER_LOGGER, "0"); + try (AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + env.kafkaClient().prepareResponseFrom(new DescribeConfigsResponse( + new DescribeConfigsResponseData().setResults(asList(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(brokerResource.name()).setResourceType(brokerResource.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList()), + new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(brokerLoggerResource.name()).setResourceType(brokerLoggerResource.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList())))), env.cluster().nodeById(0)); + Map> result = env.adminClient().describeConfigs(asList( + brokerResource, + brokerLoggerResource)).values(); + assertEquals(new HashSet<>(asList(brokerResource, brokerLoggerResource)), result.keySet()); + result.get(brokerResource).get(); + result.get(brokerLoggerResource).get(); + } + } + + @Test + public void testDescribeConfigsPartialResponse() throws Exception { + ConfigResource topic = new ConfigResource(ConfigResource.Type.TOPIC, "topic"); + ConfigResource topic2 = new ConfigResource(ConfigResource.Type.TOPIC, "topic2"); + try (AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + env.kafkaClient().prepareResponse(new DescribeConfigsResponse( + new DescribeConfigsResponseData().setResults(asList(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(topic.name()).setResourceType(topic.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList()))))); + Map> result = env.adminClient().describeConfigs(asList( + topic, + topic2)).values(); + assertEquals(new HashSet<>(asList(topic, topic2)), result.keySet()); + result.get(topic); + TestUtils.assertFutureThrows(result.get(topic2), ApiException.class); + } + } + + @Test + public void testDescribeConfigsUnrequested() throws Exception { + ConfigResource topic = new ConfigResource(ConfigResource.Type.TOPIC, "topic"); + ConfigResource unrequested = new ConfigResource(ConfigResource.Type.TOPIC, "unrequested"); + try (AdminClientUnitTestEnv env = mockClientEnv()) { + env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); + env.kafkaClient().prepareResponse(new DescribeConfigsResponse( + new DescribeConfigsResponseData().setResults(asList(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(topic.name()).setResourceType(topic.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList()), + new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(unrequested.name()).setResourceType(unrequested.type().id()).setErrorCode(Errors.NONE.code()) + .setConfigs(emptyList()))))); + Map> result = env.adminClient().describeConfigs(asList( + topic)).values(); + assertEquals(new HashSet<>(asList(topic)), result.keySet()); + assertNotNull(result.get(topic).get()); + assertNull(result.get(unrequested)); } } 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 0c22cadf2fb40..ad7a1502d8825 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 @@ -80,6 +80,10 @@ import org.apache.kafka.common.message.DescribeAclsResponseData; import org.apache.kafka.common.message.DescribeAclsResponseData.AclDescription; import org.apache.kafka.common.message.DescribeAclsResponseData.DescribeAclsResource; +import org.apache.kafka.common.message.DescribeConfigsRequestData; +import org.apache.kafka.common.message.DescribeConfigsResponseData; +import org.apache.kafka.common.message.DescribeConfigsResponseData.DescribeConfigsResourceResult; +import org.apache.kafka.common.message.DescribeConfigsResponseData.DescribeConfigsResult; import org.apache.kafka.common.message.DescribeGroupsRequestData; import org.apache.kafka.common.message.DescribeGroupsResponseData; import org.apache.kafka.common.message.DescribeGroupsResponseData.DescribedGroup; @@ -167,7 +171,6 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; -import java.util.Collection; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; @@ -436,14 +439,14 @@ public void testSerialization() throws Exception { checkRequest(createDescribeConfigsRequest(0), true); checkRequest(createDescribeConfigsRequestWithConfigEntries(0), false); checkErrorResponse(createDescribeConfigsRequest(0), new UnknownServerException(), true); - checkResponse(createDescribeConfigsResponse(), 0, false); + checkResponse(createDescribeConfigsResponse((short) 0), 0, false); checkRequest(createDescribeConfigsRequest(1), true); checkRequest(createDescribeConfigsRequestWithConfigEntries(1), false); checkRequest(createDescribeConfigsRequestWithDocumentation(1), false); checkRequest(createDescribeConfigsRequestWithDocumentation(2), false); checkRequest(createDescribeConfigsRequestWithDocumentation(3), false); checkErrorResponse(createDescribeConfigsRequest(1), new UnknownServerException(), true); - checkResponse(createDescribeConfigsResponse(), 1, false); + checkResponse(createDescribeConfigsResponse((short) 1), 1, false); checkDescribeConfigsResponseVersions(); checkRequest(createCreatePartitionsRequest(), true); checkRequest(createCreatePartitionsRequestWithAssignments(), false); @@ -500,39 +503,48 @@ private void checkOlderFetchVersions() throws Exception { } private void verifyDescribeConfigsResponse(DescribeConfigsResponse expected, DescribeConfigsResponse actual, int version) throws Exception { - for (ConfigResource resource : expected.configs().keySet()) { - Collection deserializedEntries1 = actual.config(resource).entries(); - Iterator expectedEntries = expected.config(resource).entries().iterator(); - for (DescribeConfigsResponse.ConfigEntry entry : deserializedEntries1) { - DescribeConfigsResponse.ConfigEntry expectedEntry = expectedEntries.next(); - assertEquals(expectedEntry.name(), entry.name()); - assertEquals(expectedEntry.value(), entry.value()); - assertEquals(expectedEntry.isReadOnly(), entry.isReadOnly()); - assertEquals(expectedEntry.isSensitive(), entry.isSensitive()); + for (Map.Entry resource : expected.resultMap().entrySet()) { + List actualEntries = actual.resultMap().get(resource.getKey()).configs(); + Iterator expectedEntries = expected.resultMap().get(resource.getKey()).configs().iterator(); + for (DescribeConfigsResourceResult actualEntry : actualEntries) { + DescribeConfigsResourceResult expectedEntry = expectedEntries.next(); + assertEquals(expectedEntry.name(), actualEntry.name()); + assertEquals("Non-matching values for " + actualEntry.name() + " in version " + version, + expectedEntry.value(), actualEntry.value()); + assertEquals("Non-matching readonly for " + actualEntry.name() + " in version " + version, + expectedEntry.readOnly(), actualEntry.readOnly()); + assertEquals("Non-matching isSensitive for " + actualEntry.name() + " in version " + version, + expectedEntry.isSensitive(), actualEntry.isSensitive()); if (version < 3) { - assertEquals(ConfigType.UNKNOWN, entry.type()); + assertEquals("Non-matching configType for " + actualEntry.name() + " in version " + version, + ConfigType.UNKNOWN.id(), actualEntry.configType()); } else { - assertEquals(expectedEntry.type(), entry.type()); + assertEquals("Non-matching configType for " + actualEntry.name() + " in version " + version, + expectedEntry.configType(), actualEntry.configType()); } - if (version == 1 || version == 3 || (expectedEntry.source() != DescribeConfigsResponse.ConfigSource.DYNAMIC_BROKER_CONFIG && - expectedEntry.source() != DescribeConfigsResponse.ConfigSource.DYNAMIC_DEFAULT_BROKER_CONFIG)) - assertEquals(expectedEntry.source(), entry.source()); + if (version == 1 || version == 3 || (expectedEntry.configSource() != DescribeConfigsResponse.ConfigSource.DYNAMIC_BROKER_CONFIG.id() && + expectedEntry.configSource() != DescribeConfigsResponse.ConfigSource.DYNAMIC_DEFAULT_BROKER_CONFIG.id())) + assertEquals("Non-matching configSource for " + actualEntry.name() + " in version " + version, + expectedEntry.configSource(), actualEntry.configSource()); else - assertEquals(DescribeConfigsResponse.ConfigSource.STATIC_BROKER_CONFIG, entry.source()); + assertEquals("Non matching configSource for " + actualEntry.name() + " in version " + version, + DescribeConfigsResponse.ConfigSource.STATIC_BROKER_CONFIG.id(), actualEntry.configSource()); } } } private void checkDescribeConfigsResponseVersions() throws Exception { - DescribeConfigsResponse response = createDescribeConfigsResponse(); + DescribeConfigsResponse response = createDescribeConfigsResponse((short) 0); DescribeConfigsResponse deserialized0 = (DescribeConfigsResponse) deserialize(response, response.toStruct((short) 0), (short) 0); verifyDescribeConfigsResponse(response, deserialized0, 0); + response = createDescribeConfigsResponse((short) 1); DescribeConfigsResponse deserialized1 = (DescribeConfigsResponse) deserialize(response, response.toStruct((short) 1), (short) 1); verifyDescribeConfigsResponse(response, deserialized1, 1); + response = createDescribeConfigsResponse((short) 3); DescribeConfigsResponse deserialized3 = (DescribeConfigsResponse) deserialize(response, response.toStruct((short) 3), (short) 3); verifyDescribeConfigsResponse(response, deserialized3, 3); @@ -1879,43 +1891,85 @@ private DeleteAclsResponse createDeleteAclsResponse() { } private DescribeConfigsRequest createDescribeConfigsRequest(int version) { - return new DescribeConfigsRequest.Builder(asList( - new ConfigResource(ConfigResource.Type.BROKER, "0"), - new ConfigResource(ConfigResource.Type.TOPIC, "topic"))) + return new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() + .setResources(asList( + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.BROKER.id()) + .setResourceName("0"), + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.TOPIC.id()) + .setResourceName("topic")))) .build((short) version); } private DescribeConfigsRequest createDescribeConfigsRequestWithConfigEntries(int version) { - Map> resources = new HashMap<>(); - resources.put(new ConfigResource(ConfigResource.Type.BROKER, "0"), asList("foo", "bar")); - resources.put(new ConfigResource(ConfigResource.Type.TOPIC, "topic"), null); - resources.put(new ConfigResource(ConfigResource.Type.TOPIC, "topic a"), emptyList()); - return new DescribeConfigsRequest.Builder(resources).build((short) version); + return new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() + .setResources(asList( + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.BROKER.id()) + .setResourceName("0") + .setConfigurationKeys(asList("foo", "bar")), + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.TOPIC.id()) + .setResourceName("topic") + .setConfigurationKeys(null), + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.TOPIC.id()) + .setResourceName("topic a") + .setConfigurationKeys(emptyList())))).build((short) version); } private DescribeConfigsRequest createDescribeConfigsRequestWithDocumentation(int version) { - Map> resources = new HashMap<>(); - resources.put(new ConfigResource(ConfigResource.Type.BROKER, "0"), asList("foo", "bar")); - return new DescribeConfigsRequest.Builder(resources).includeDocumentation(true).build((short) version); - } - - private DescribeConfigsResponse createDescribeConfigsResponse() { - Map configs = new HashMap<>(); - List synonyms = emptyList(); - List configEntries = asList( - new DescribeConfigsResponse.ConfigEntry("config_name", "config_value", - DescribeConfigsResponse.ConfigSource.DYNAMIC_BROKER_CONFIG, true, false, synonyms), - new DescribeConfigsResponse.ConfigEntry("another_name", "another value", - DescribeConfigsResponse.ConfigSource.DEFAULT_CONFIG, false, true, synonyms), - new DescribeConfigsResponse.ConfigEntry("yet_another_name", "yet another value", - DescribeConfigsResponse.ConfigSource.DEFAULT_CONFIG, false, true, synonyms, - ConfigType.BOOLEAN, "some description") - ); - configs.put(new ConfigResource(ConfigResource.Type.BROKER, "0"), new DescribeConfigsResponse.Config( - ApiError.NONE, configEntries)); - configs.put(new ConfigResource(ConfigResource.Type.TOPIC, "topic"), new DescribeConfigsResponse.Config( - ApiError.NONE, Collections.emptyList())); - return new DescribeConfigsResponse(200, configs); + DescribeConfigsRequestData data = new DescribeConfigsRequestData() + .setResources(asList( + new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.BROKER.id()) + .setResourceName("0") + .setConfigurationKeys(asList("foo", "bar")))); + if (version == 3) { + data.setIncludeDocumentation(true); + } + return new DescribeConfigsRequest.Builder(data).build((short) version); + } + + private DescribeConfigsResponse createDescribeConfigsResponse(short version) { + return new DescribeConfigsResponse(new DescribeConfigsResponseData().setResults(asList( + new DescribeConfigsResult() + .setErrorCode(Errors.NONE.code()) + .setResourceType(ConfigResource.Type.BROKER.id()) + .setResourceName("0") + .setConfigs(asList( + new DescribeConfigsResourceResult() + .setName("config_name") + .setValue("config_value") + // Note: the v0 default for this field that should be exposed to callers is + // context-dependent. For example, if the resource is a broker, this should default to 4. + // -1 is just a placeholder value. + .setConfigSource(version == 0 ? DescribeConfigsResponse.ConfigSource.STATIC_BROKER_CONFIG.id() : DescribeConfigsResponse.ConfigSource.DYNAMIC_BROKER_CONFIG.id) + .setIsSensitive(true).setReadOnly(false) + .setSynonyms(emptyList()), + new DescribeConfigsResourceResult() + .setName("yet_another_name") + .setValue("yet another value") + .setConfigSource(version == 0 ? DescribeConfigsResponse.ConfigSource.STATIC_BROKER_CONFIG.id() : DescribeConfigsResponse.ConfigSource.DEFAULT_CONFIG.id) + .setIsSensitive(false).setReadOnly(true) + .setSynonyms(emptyList()) + .setConfigType(ConfigType.BOOLEAN.id()) + .setDocumentation("some description"), + new DescribeConfigsResourceResult() + .setName("another_name") + .setValue("another value") + .setConfigSource(version == 0 ? DescribeConfigsResponse.ConfigSource.STATIC_BROKER_CONFIG.id() : DescribeConfigsResponse.ConfigSource.DEFAULT_CONFIG.id) + .setIsSensitive(false).setReadOnly(true) + .setSynonyms(emptyList()) + )), + new DescribeConfigsResult() + .setErrorCode(Errors.NONE.code()) + .setResourceType(ConfigResource.Type.TOPIC.id()) + .setResourceName("topic") + .setConfigs(emptyList()) + ))); + } private AlterConfigsRequest createAlterConfigsRequest() { diff --git a/core/src/main/scala/kafka/server/AdminManager.scala b/core/src/main/scala/kafka/server/AdminManager.scala index 742156a985316..2e16c978b8854 100644 --- a/core/src/main/scala/kafka/server/AdminManager.scala +++ b/core/src/main/scala/kafka/server/AdminManager.scala @@ -34,6 +34,8 @@ import org.apache.kafka.common.internals.Topic import org.apache.kafka.common.message.CreatePartitionsRequestData.CreatePartitionsTopic import org.apache.kafka.common.message.CreateTopicsRequestData.CreatableTopic import org.apache.kafka.common.message.CreateTopicsResponseData.{CreatableTopicConfigs, CreatableTopicResult} +import org.apache.kafka.common.message.DescribeConfigsResponseData +import org.apache.kafka.common.message.DescribeConfigsRequestData.DescribeConfigsResource import org.apache.kafka.common.metrics.Metrics import org.apache.kafka.common.network.ListenerName import org.apache.kafka.server.policy.{AlterConfigPolicy, CreateTopicPolicy} @@ -167,15 +169,12 @@ class AdminManager(val config: KafkaConfig, val createEntry = createTopicConfigEntry(logConfig, configs, includeSynonyms = false, includeDocumentation = false)(_, _) val topicConfigs = logConfig.values.asScala.map { case (k, v) => val entry = createEntry(k, v) - val source = ConfigSource.values.indices.map(_.toByte) - .find(i => ConfigSource.forId(i.toByte) == entry.source) - .getOrElse(0.toByte) new CreatableTopicConfigs() .setName(k) .setValue(entry.value) .setIsSensitive(entry.isSensitive) - .setReadOnly(entry.isReadOnly) - .setConfigSource(source) + .setReadOnly(entry.readOnly) + .setConfigSource(entry.configSource) }.toList.asJava result.setConfigs(topicConfigs) result.setNumPartitions(assignments.size) @@ -347,28 +346,28 @@ class AdminManager(val config: KafkaConfig, } } - def describeConfigs(resourceToConfigNames: Map[ConfigResource, Option[Set[String]]], includeSynonyms: Boolean, includeDocumentation: Boolean): Map[ConfigResource, DescribeConfigsResponse.Config] = { - resourceToConfigNames.map { case (resource, configNames) => + def describeConfigs(resourceToConfigNames: List[DescribeConfigsResource], includeSynonyms: Boolean, includeDocumentation: Boolean): List[DescribeConfigsResponseData.DescribeConfigsResult] = { + resourceToConfigNames.map { case resource => def allConfigs(config: AbstractConfig) = { config.originals.asScala.filter(_._2 != null) ++ config.values.asScala } def createResponseConfig(configs: Map[String, Any], - createConfigEntry: (String, Any) => DescribeConfigsResponse.ConfigEntry): DescribeConfigsResponse.Config = { + createConfigEntry: (String, Any) => DescribeConfigsResponseData.DescribeConfigsResourceResult): DescribeConfigsResponseData.DescribeConfigsResult = { val filteredConfigPairs = configs.filter { case (configName, _) => /* Always returns true if configNames is None */ - configNames.forall(_.contains(configName)) + resource.configurationKeys.asScala.forall(_.contains(configName)) }.toBuffer val configEntries = filteredConfigPairs.map { case (name, value) => createConfigEntry(name, value) } - new DescribeConfigsResponse.Config(ApiError.NONE, configEntries.asJava) + new DescribeConfigsResponseData.DescribeConfigsResult().setErrorCode(Errors.NONE.code) + .setConfigs(configEntries.asJava) } try { - val resourceConfig = resource.`type` match { - + val configResult = ConfigResource.Type.forId(resource.resourceType) match { case ConfigResource.Type.TOPIC => - val topic = resource.name + val topic = resource.resourceName Topic.validate(topic) if (metadataCache.contains(topic)) { // Consider optimizing this by caching the configs or retrieving them from the `Log` when possible @@ -376,30 +375,33 @@ class AdminManager(val config: KafkaConfig, val logConfig = LogConfig.fromProps(KafkaServer.copyKafkaConfigToLog(config), topicProps) createResponseConfig(allConfigs(logConfig), createTopicConfigEntry(logConfig, topicProps, includeSynonyms, includeDocumentation)) } else { - new DescribeConfigsResponse.Config(new ApiError(Errors.UNKNOWN_TOPIC_OR_PARTITION, null), Collections.emptyList[DescribeConfigsResponse.ConfigEntry]) + new DescribeConfigsResponseData.DescribeConfigsResult().setErrorCode(Errors.UNKNOWN_TOPIC_OR_PARTITION.code) + .setConfigs(Collections.emptyList[DescribeConfigsResponseData.DescribeConfigsResourceResult]) } case ConfigResource.Type.BROKER => - if (resource.name == null || resource.name.isEmpty) + if (resource.resourceName == null || resource.resourceName.isEmpty) createResponseConfig(config.dynamicConfig.currentDynamicDefaultConfigs, - createBrokerConfigEntry(perBrokerConfig = false, includeSynonyms, includeDocumentation)) - else if (resourceNameToBrokerId(resource.name) == config.brokerId) + createBrokerConfigEntry(perBrokerConfig = false, includeSynonyms, includeDocumentation)) + else if (resourceNameToBrokerId(resource.resourceName) == config.brokerId) createResponseConfig(allConfigs(config), - createBrokerConfigEntry(perBrokerConfig = true, includeSynonyms, includeDocumentation)) + createBrokerConfigEntry(perBrokerConfig = true, includeSynonyms, includeDocumentation)) else - throw new InvalidRequestException(s"Unexpected broker id, expected ${config.brokerId} or empty string, but received ${resource.name}") + throw new InvalidRequestException(s"Unexpected broker id, expected ${config.brokerId} or empty string, but received ${resource.resourceName}") case ConfigResource.Type.BROKER_LOGGER => - if (resource.name == null || resource.name.isEmpty) + if (resource.resourceName == null || resource.resourceName.isEmpty) throw new InvalidRequestException("Broker id must not be empty") - else if (resourceNameToBrokerId(resource.name) != config.brokerId) - throw new InvalidRequestException(s"Unexpected broker id, expected ${config.brokerId} but received ${resource.name}") + else if (resourceNameToBrokerId(resource.resourceName) != config.brokerId) + throw new InvalidRequestException(s"Unexpected broker id, expected ${config.brokerId} but received ${resource.resourceName}") else createResponseConfig(Log4jController.loggers, - (name, value) => new DescribeConfigsResponse.ConfigEntry(name, value.toString, ConfigSource.DYNAMIC_BROKER_LOGGER_CONFIG, false, false, List.empty.asJava)) + (name, value) => new DescribeConfigsResponseData.DescribeConfigsResourceResult().setName(name) + .setValue(value.toString).setConfigSource(ConfigSource.DYNAMIC_BROKER_LOGGER_CONFIG.id) + .setIsSensitive(false).setReadOnly(false).setSynonyms(List.empty.asJava)) case resourceType => throw new InvalidRequestException(s"Unsupported resource type: $resourceType") } - resource -> resourceConfig + configResult.setResourceName(resource.resourceName).setResourceType(resource.resourceType) } catch { case e: Throwable => // Log client errors at a lower level than unexpected exceptions @@ -408,9 +410,15 @@ class AdminManager(val config: KafkaConfig, info(message, e) else error(message, e) - resource -> new DescribeConfigsResponse.Config(ApiError.fromThrowable(e), Collections.emptyList[DescribeConfigsResponse.ConfigEntry]) + val err = ApiError.fromThrowable(e) + new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(resource.resourceName) + .setResourceType(resource.resourceType) + .setErrorMessage(err.message) + .setErrorCode(err.error.code) + .setConfigs(Collections.emptyList[DescribeConfigsResponseData.DescribeConfigsResourceResult]) } - }.toMap + }.toList } def alterConfigs(configs: Map[ConfigResource, AlterConfigsRequest.Config], validateOnly: Boolean): Map[ConfigResource, ApiError] = { @@ -681,15 +689,15 @@ class AdminManager(val config: KafkaConfig, case _ => DescribeConfigsResponse.ConfigType.UNKNOWN } } - - private def configSynonyms(name: String, synonyms: List[String], isSensitive: Boolean): List[DescribeConfigsResponse.ConfigSynonym] = { + + private def configSynonyms(name: String, synonyms: List[String], isSensitive: Boolean): List[DescribeConfigsResponseData.DescribeConfigsSynonym] = { val dynamicConfig = config.dynamicConfig - val allSynonyms = mutable.Buffer[DescribeConfigsResponse.ConfigSynonym]() + val allSynonyms = mutable.Buffer[DescribeConfigsResponseData.DescribeConfigsSynonym]() def maybeAddSynonym(map: Map[String, String], source: ConfigSource)(name: String): Unit = { map.get(name).map { value => val configValue = if (isSensitive) null else value - allSynonyms += new DescribeConfigsResponse.ConfigSynonym(name, configValue, source) + allSynonyms += new DescribeConfigsResponseData.DescribeConfigsSynonym().setName(name).setValue(configValue).setSource(source.id) } } @@ -701,7 +709,7 @@ class AdminManager(val config: KafkaConfig, } private def createTopicConfigEntry(logConfig: LogConfig, topicProps: Properties, includeSynonyms: Boolean, includeDocumentation: Boolean) - (name: String, value: Any): DescribeConfigsResponse.ConfigEntry = { + (name: String, value: Any): DescribeConfigsResponseData.DescribeConfigsResourceResult = { val configEntryType = LogConfig.configType(name) val isSensitive = KafkaConfig.maybeSensitive(configEntryType) val valueAsString = if (isSensitive) null else ConfigDef.convertToString(value, configEntryType.orNull) @@ -712,17 +720,21 @@ class AdminManager(val config: KafkaConfig, if (!topicProps.containsKey(name)) list else - new DescribeConfigsResponse.ConfigSynonym(name, valueAsString, ConfigSource.TOPIC_CONFIG) +: list + new DescribeConfigsResponseData.DescribeConfigsSynonym().setName(name).setValue(valueAsString) + .setSource(ConfigSource.TOPIC_CONFIG.id) +: list } - val source = if (allSynonyms.isEmpty) ConfigSource.DEFAULT_CONFIG else allSynonyms.head.source + val source = if (allSynonyms.isEmpty) ConfigSource.DEFAULT_CONFIG.id else allSynonyms.head.source val synonyms = if (!includeSynonyms) List.empty else allSynonyms val dataType = configResponseType(configEntryType) val configDocumentation = if (includeDocumentation) brokerDocumentation(name) else null - new DescribeConfigsResponse.ConfigEntry(name, valueAsString, source, isSensitive, false, synonyms.asJava, dataType, configDocumentation) + new DescribeConfigsResponseData.DescribeConfigsResourceResult() + .setName(name).setValue(valueAsString).setConfigSource(source) + .setIsSensitive(isSensitive).setReadOnly(false).setSynonyms(synonyms.asJava) + .setDocumentation(configDocumentation).setConfigType(dataType.id) } private def createBrokerConfigEntry(perBrokerConfig: Boolean, includeSynonyms: Boolean, includeDocumentation: Boolean) - (name: String, value: Any): DescribeConfigsResponse.ConfigEntry = { + (name: String, value: Any): DescribeConfigsResponseData.DescribeConfigsResourceResult = { val allNames = brokerSynonyms(name) val configEntryType = KafkaConfig.configType(name) val isSensitive = KafkaConfig.maybeSensitive(configEntryType) @@ -733,13 +745,16 @@ class AdminManager(val config: KafkaConfig, case _ => ConfigDef.convertToString(value, configEntryType.orNull) } val allSynonyms = configSynonyms(name, allNames, isSensitive) - .filter(perBrokerConfig || _.source == ConfigSource.DYNAMIC_DEFAULT_BROKER_CONFIG) + .filter(perBrokerConfig || _.source == ConfigSource.DYNAMIC_DEFAULT_BROKER_CONFIG.id) val synonyms = if (!includeSynonyms) List.empty else allSynonyms - val source = if (allSynonyms.isEmpty) ConfigSource.DEFAULT_CONFIG else allSynonyms.head.source + val source = if (allSynonyms.isEmpty) ConfigSource.DEFAULT_CONFIG.id else allSynonyms.head.source val readOnly = !DynamicBrokerConfig.AllDynamicConfigs.contains(name) + val dataType = configResponseType(configEntryType) val configDocumentation = if (includeDocumentation) brokerDocumentation(name) else null - new DescribeConfigsResponse.ConfigEntry(name, valueAsString, source, isSensitive, readOnly, synonyms.asJava, dataType, configDocumentation) + new DescribeConfigsResponseData.DescribeConfigsResourceResult().setName(name).setValue(valueAsString).setConfigSource(source) + .setIsSensitive(isSensitive).setReadOnly(readOnly).setSynonyms(synonyms.asJava) + .setDocumentation(configDocumentation).setConfigType(dataType.id) } private def sanitizeEntityName(entityName: String): String = @@ -827,32 +842,32 @@ class AdminManager(val config: KafkaConfig, val excludeClientId = wantExcluded(clientIdComponent) val userEntries = if (exactUser && excludeClientId) - Map(((Some(user.get), None) -> adminZkClient.fetchEntityConfig(ConfigType.User, sanitizedUser))) + Map((Some(user.get), None) -> adminZkClient.fetchEntityConfig(ConfigType.User, sanitizedUser)) else if (!excludeUser && !exactClientId) adminZkClient.fetchAllEntityConfigs(ConfigType.User).map { case (name, props) => - ((Some(desanitizeEntityName(name)), None) -> props) + (Some(desanitizeEntityName(name)), None) -> props } else Map.empty val clientIdEntries = if (excludeUser && exactClientId) - Map(((None, Some(clientId.get)) -> adminZkClient.fetchEntityConfig(ConfigType.Client, sanitizedClientId))) + Map((None, Some(clientId.get)) -> adminZkClient.fetchEntityConfig(ConfigType.Client, sanitizedClientId)) else if (!exactUser && !excludeClientId) adminZkClient.fetchAllEntityConfigs(ConfigType.Client).map { case (name, props) => - ((None, Some(desanitizeEntityName(name))) -> props) + (None, Some(desanitizeEntityName(name))) -> props } else Map.empty val bothEntries = if (exactUser && exactClientId) - Map(((Some(user.get), Some(clientId.get)) -> - adminZkClient.fetchEntityConfig(ConfigType.User, s"${sanitizedUser}/clients/${sanitizedClientId}"))) + Map((Some(user.get), Some(clientId.get)) -> + adminZkClient.fetchEntityConfig(ConfigType.User, s"${sanitizedUser}/clients/${sanitizedClientId}")) else if (!excludeUser && !excludeClientId) adminZkClient.fetchAllChildEntityConfigs(ConfigType.User, ConfigType.Client).map { case (name, props) => val components = name.split("/") if (components.size != 3 || components(1) != "clients") throw new IllegalArgumentException(s"Unexpected config path: ${name}") - ((Some(desanitizeEntityName(components(0))), Some(desanitizeEntityName(components(2)))) -> props) + (Some(desanitizeEntityName(components(0))), Some(desanitizeEntityName(components(2)))) -> props } else Map.empty @@ -873,13 +888,13 @@ class AdminManager(val config: KafkaConfig, case _: NumberFormatException => throw new IllegalStateException(s"Unexpected client quota configuration value: ${key} -> ${value}") } - (key -> doubleValue) + key -> doubleValue } } (userEntries ++ clientIdEntries ++ bothEntries).map { case ((u, c), p) => if (!p.isEmpty && matches(userComponent, u) && matches(clientIdComponent, c)) - Some((userClientIdToEntity(u, c) -> fromProps(p))) + Some(userClientIdToEntity(u, c) -> fromProps(p)) else None }.flatten.toMap @@ -929,7 +944,7 @@ class AdminManager(val config: KafkaConfig, info(s"Error encountered while updating client quotas", e) ApiError.fromThrowable(e) } - (entry.entity -> apiError) + entry.entity -> apiError }.toMap } } diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 36acd8cc45263..0a49d944936fd 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -51,7 +51,7 @@ import org.apache.kafka.common.internals.Topic.{GROUP_METADATA_TOPIC_NAME, TRANS import org.apache.kafka.common.message.AlterConfigsResponseData.AlterConfigsResourceResponse import org.apache.kafka.common.message.CreateTopicsRequestData.CreatableTopic import org.apache.kafka.common.message.CreatePartitionsResponseData.CreatePartitionsTopicResult -import org.apache.kafka.common.message.{AddOffsetsToTxnResponseData, AlterConfigsResponseData, AlterPartitionReassignmentsResponseData, AlterReplicaLogDirsResponseData, CreateAclsResponseData, CreatePartitionsResponseData, CreateTopicsResponseData, DeleteAclsResponseData, DeleteGroupsResponseData, DeleteRecordsResponseData, DeleteTopicsResponseData, DescribeAclsResponseData, DescribeGroupsResponseData, DescribeLogDirsResponseData, EndTxnResponseData, ExpireDelegationTokenResponseData, FindCoordinatorResponseData, HeartbeatResponseData, InitProducerIdResponseData, JoinGroupResponseData, LeaveGroupResponseData, ListGroupsResponseData, ListPartitionReassignmentsResponseData, OffsetCommitRequestData, OffsetCommitResponseData, OffsetDeleteResponseData, RenewDelegationTokenResponseData, SaslAuthenticateResponseData, SaslHandshakeResponseData, StopReplicaResponseData, SyncGroupResponseData, UpdateMetadataResponseData} +import org.apache.kafka.common.message.{AddOffsetsToTxnResponseData, AlterConfigsResponseData, AlterPartitionReassignmentsResponseData, AlterReplicaLogDirsResponseData, CreateAclsResponseData, CreatePartitionsResponseData, CreateTopicsResponseData, DeleteAclsResponseData, DeleteGroupsResponseData, DeleteRecordsResponseData, DeleteTopicsResponseData, DescribeAclsResponseData, DescribeConfigsResponseData, DescribeGroupsResponseData, DescribeLogDirsResponseData, EndTxnResponseData, ExpireDelegationTokenResponseData, FindCoordinatorResponseData, HeartbeatResponseData, InitProducerIdResponseData, JoinGroupResponseData, LeaveGroupResponseData, ListGroupsResponseData, ListPartitionReassignmentsResponseData, OffsetCommitRequestData, OffsetCommitResponseData, OffsetDeleteResponseData, RenewDelegationTokenResponseData, SaslAuthenticateResponseData, SaslHandshakeResponseData, StopReplicaResponseData, SyncGroupResponseData, UpdateMetadataResponseData} import org.apache.kafka.common.message.CreateTopicsResponseData.{CreatableTopicResult, CreatableTopicResultCollection} import org.apache.kafka.common.message.DeleteGroupsResponseData.{DeletableGroupResult, DeletableGroupResultCollection} import org.apache.kafka.common.message.AlterPartitionReassignmentsResponseData.{ReassignablePartitionResponse, ReassignableTopicResponse} @@ -2568,25 +2568,32 @@ class KafkaApis(val requestChannel: RequestChannel, def handleDescribeConfigsRequest(request: RequestChannel.Request): Unit = { val describeConfigsRequest = request.body[DescribeConfigsRequest] - val (authorizedResources, unauthorizedResources) = describeConfigsRequest.resources.asScala.toBuffer.partition { resource => - resource.`type` match { + val (authorizedResources, unauthorizedResources) = describeConfigsRequest.data.resources.asScala.toBuffer.partition { resource => + ConfigResource.Type.forId(resource.resourceType) match { case ConfigResource.Type.BROKER | ConfigResource.Type.BROKER_LOGGER => authorize(request.context, DESCRIBE_CONFIGS, CLUSTER, CLUSTER_NAME) case ConfigResource.Type.TOPIC => - authorize(request.context, DESCRIBE_CONFIGS, TOPIC, resource.name) - case rt => throw new InvalidRequestException(s"Unexpected resource type $rt for resource ${resource.name}") + authorize(request.context, DESCRIBE_CONFIGS, TOPIC, resource.resourceName) + case rt => throw new InvalidRequestException(s"Unexpected resource type $rt for resource ${resource.resourceName}") } } - val authorizedConfigs = adminManager.describeConfigs(authorizedResources.map { resource => - resource -> Option(describeConfigsRequest.configNames(resource)).map(_.asScala.toSet) - }.toMap, describeConfigsRequest.includeSynonyms, describeConfigsRequest.includeDocumentation) + val authorizedConfigs = adminManager.describeConfigs(authorizedResources.toList, describeConfigsRequest.data.includeSynonyms, describeConfigsRequest.data.includeDocumentation) val unauthorizedConfigs = unauthorizedResources.map { resource => - val error = configsAuthorizationApiError(resource) - resource -> new DescribeConfigsResponse.Config(error, util.Collections.emptyList[DescribeConfigsResponse.ConfigEntry]) + val error = ConfigResource.Type.forId(resource.resourceType) match { + case ConfigResource.Type.BROKER | ConfigResource.Type.BROKER_LOGGER => Errors.CLUSTER_AUTHORIZATION_FAILED + case ConfigResource.Type.TOPIC => Errors.TOPIC_AUTHORIZATION_FAILED + case rt => throw new InvalidRequestException(s"Unexpected resource type $rt for resource ${resource.resourceName}") + } + new DescribeConfigsResponseData.DescribeConfigsResult().setErrorCode(error.code) + .setErrorMessage(error.message) + .setConfigs(Collections.emptyList[DescribeConfigsResponseData.DescribeConfigsResourceResult]) + .setResourceName(resource.resourceName) + .setResourceType(resource.resourceType) } sendResponseMaybeThrottle(request, requestThrottleMs => - new DescribeConfigsResponse(requestThrottleMs, (authorizedConfigs ++ unauthorizedConfigs).asJava)) + new DescribeConfigsResponse(new DescribeConfigsResponseData().setThrottleTimeMs(requestThrottleMs) + .setResults((authorizedConfigs ++ unauthorizedConfigs).asJava))) } def handleAlterReplicaLogDirsRequest(request: RequestChannel.Request): Unit = { diff --git a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala index f3b259b634a3f..d97a2c64d05c8 100644 --- a/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala +++ b/core/src/test/scala/integration/kafka/api/AuthorizerIntegrationTest.scala @@ -44,7 +44,7 @@ import org.apache.kafka.common.message.LeaderAndIsrRequestData.LeaderAndIsrParti import org.apache.kafka.common.message.LeaveGroupRequestData.MemberIdentity import org.apache.kafka.common.message.StopReplicaRequestData.{StopReplicaPartitionState, StopReplicaTopicState} import org.apache.kafka.common.message.UpdateMetadataRequestData.{UpdateMetadataBroker, UpdateMetadataEndpoint, UpdateMetadataPartitionState} -import org.apache.kafka.common.message.{AlterPartitionReassignmentsRequestData, AlterReplicaLogDirsRequestData, ControlledShutdownRequestData, CreateAclsRequestData, CreatePartitionsRequestData, CreateTopicsRequestData, DeleteAclsRequestData, DeleteGroupsRequestData, DeleteRecordsRequestData, DeleteTopicsRequestData, DescribeGroupsRequestData, DescribeLogDirsRequestData, FindCoordinatorRequestData, HeartbeatRequestData, IncrementalAlterConfigsRequestData, JoinGroupRequestData, ListPartitionReassignmentsRequestData, OffsetCommitRequestData, SyncGroupRequestData} +import org.apache.kafka.common.message.{AlterPartitionReassignmentsRequestData, AlterReplicaLogDirsRequestData, ControlledShutdownRequestData, CreateAclsRequestData, CreatePartitionsRequestData, CreateTopicsRequestData, DeleteAclsRequestData, DeleteGroupsRequestData, DeleteRecordsRequestData, DeleteTopicsRequestData, DescribeConfigsRequestData, DescribeGroupsRequestData, DescribeLogDirsRequestData, FindCoordinatorRequestData, HeartbeatRequestData, IncrementalAlterConfigsRequestData, JoinGroupRequestData, ListPartitionReassignmentsRequestData, OffsetCommitRequestData, SyncGroupRequestData} import org.apache.kafka.common.network.ListenerName import org.apache.kafka.common.protocol.{ApiKeys, Errors} import org.apache.kafka.common.record.{CompressionType, MemoryRecords, RecordBatch, Records, SimpleRecord} @@ -181,7 +181,7 @@ class AuthorizerIntegrationTest extends BaseRequestTest { resp.data.topics.find(tp.topic).partitions.find(tp.partition).errorCode)), ApiKeys.OFFSET_FOR_LEADER_EPOCH -> ((resp: OffsetsForLeaderEpochResponse) => resp.responses.get(tp).error), ApiKeys.DESCRIBE_CONFIGS -> ((resp: DescribeConfigsResponse) => - resp.configs.get(new ConfigResource(ConfigResource.Type.TOPIC, tp.topic)).error.error), + Errors.forCode(resp.resultMap.get(new ConfigResource(ConfigResource.Type.TOPIC, tp.topic)).errorCode)), ApiKeys.ALTER_CONFIGS -> ((resp: AlterConfigsResponse) => resp.errors.get(new ConfigResource(ConfigResource.Type.TOPIC, tp.topic)).error), ApiKeys.INIT_PRODUCER_ID -> ((resp: InitProducerIdResponse) => resp.error), @@ -491,7 +491,9 @@ class AuthorizerIntegrationTest extends BaseRequestTest { .setOffset(0L)))))).build() private def describeConfigsRequest = - new DescribeConfigsRequest.Builder(Collections.singleton(new ConfigResource(ConfigResource.Type.TOPIC, tp.topic))).build() + new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData().setResources(Collections.singletonList( + new DescribeConfigsRequestData.DescribeConfigsResource().setResourceType(ConfigResource.Type.TOPIC.id) + .setResourceName(tp.topic)))).build() private def alterConfigsRequest = new AlterConfigsRequest.Builder( diff --git a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala index be64f00e66e43..3832a698debdc 100644 --- a/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala +++ b/core/src/test/scala/unit/kafka/server/KafkaApisTest.scala @@ -256,17 +256,22 @@ class KafkaApisTest { expectNoThrottling() val configResource = new ConfigResource(ConfigResource.Type.TOPIC, resourceName) - val config = new DescribeConfigsResponse.Config(ApiError.NONE, Collections.emptyList[DescribeConfigsResponse.ConfigEntry]) EasyMock.expect(adminManager.describeConfigs(anyObject(), EasyMock.eq(true), EasyMock.eq(false))) - .andReturn(Map(configResource -> config)) + .andReturn( + List(new DescribeConfigsResponseData.DescribeConfigsResult() + .setResourceName(configResource.name) + .setResourceType(configResource.`type`.id) + .setErrorCode(Errors.NONE.code) + .setConfigs(Collections.emptyList()))) EasyMock.replay(replicaManager, clientRequestQuotaManager, requestChannel, authorizer, adminManager) - val resourceToConfigNames = Map[ConfigResource, util.Collection[String]]( - configResource -> Collections.emptyList[String]) - val request = buildRequest(new DescribeConfigsRequest(requestHeader.apiVersion, - resourceToConfigNames.asJava, true)) + val request = buildRequest(new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() + .setIncludeSynonyms(true) + .setResources(List(new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceName("topic-1") + .setResourceType(ConfigResource.Type.TOPIC.id)).asJava)).build(requestHeader.apiVersion)) createKafkaApis(authorizer = Some(authorizer)).handleDescribeConfigsRequest(request) verify(authorizer, adminManager) diff --git a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala index 0e4ef315e02ba..b6a6b1e9e4270 100644 --- a/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala +++ b/core/src/test/scala/unit/kafka/server/RequestQuotaTest.scala @@ -458,7 +458,10 @@ class RequestQuotaTest extends BaseRequestTest { .setOperation(AclOperation.ANY.code) .setPermissionType(AclPermissionType.DENY.code)))) case ApiKeys.DESCRIBE_CONFIGS => - new DescribeConfigsRequest.Builder(Collections.singleton(new ConfigResource(ConfigResource.Type.TOPIC, tp.topic))) + new DescribeConfigsRequest.Builder(new DescribeConfigsRequestData() + .setResources(Collections.singletonList(new DescribeConfigsRequestData.DescribeConfigsResource() + .setResourceType(ConfigResource.Type.TOPIC.id) + .setResourceName(tp.topic)))) case ApiKeys.ALTER_CONFIGS => new AlterConfigsRequest.Builder(