Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<DescribeConfigsRequest> {
private final Map<ConfigResource, Collection<String>> resourceToConfigNames;
private boolean includeSynonyms;
private boolean includeDocumentation;
private final DescribeConfigsRequestData data;

public Builder(Map<ConfigResource, Collection<String>> 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<ConfigResource> resources) {
this(toResourceToConfigNames(resources));
}

private static Map<ConfigResource, Collection<String>> toResourceToConfigNames(Collection<ConfigResource> resources) {
Map<ConfigResource, Collection<String>> result = new HashMap<>(resources.size());
for (ConfigResource resource : resources)
result.put(resource, null);
return result;
this.data = data;
Comment thread
tombentley marked this conversation as resolved.
Outdated
}

@Override
public DescribeConfigsRequest build(short version) {
return new DescribeConfigsRequest(
version, resourceToConfigNames, includeSynonyms, includeDocumentation);
return new DescribeConfigsRequest(data, version);
}
}

private final Map<ConfigResource, Collection<String>> resourceToConfigNames;
private final boolean includeSynonyms;
private final boolean includeDocumentation;
private final DescribeConfigsRequestData data;

public DescribeConfigsRequest(
short version, Map<ConfigResource, Collection<String>> resourceToConfigNames,
boolean includeSynonyms) {
this(version, resourceToConfigNames, includeSynonyms, false);
}
public DescribeConfigsRequest(
short version, Map<ConfigResource, Collection<String>> 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<String> 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<ConfigResource> resources() {
return resourceToConfigNames.keySet();
this.data = new DescribeConfigsRequestData(struct, version);
}

/**
* Return null if all config names should be returned.
*/
public Collection<String> 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<Struct> resourceStructs = new ArrayList<>(resources().size());
for (Map.Entry<ConfigResource, Collection<String>> 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<ConfigResource, DescribeConfigsResponse.Config> 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) {
Expand Down
Loading