From 1469d70e6cf6df55ad7ed58afb84522380b947b0 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Fri, 11 Oct 2019 22:18:18 +0100 Subject: [PATCH 1/3] KAFKA-9026: Use automatic RPC generation in DescribeAcls --- .../apache/kafka/common/protocol/ApiKeys.java | 6 +- .../common/requests/AbstractResponse.java | 2 +- .../common/requests/DescribeAclsRequest.java | 123 ++++++----- .../common/requests/DescribeAclsResponse.java | 198 +++++++----------- .../clients/admin/KafkaAdminClientTest.java | 6 +- .../requests/DescribeAclsRequestTest.java | 18 +- .../requests/DescribeAclsResponseTest.java | 12 +- .../common/requests/RequestResponseTest.java | 2 +- .../main/scala/kafka/server/KafkaApis.scala | 4 +- 9 files changed, 159 insertions(+), 212 deletions(-) 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 e2f7f99a082e7..487adb66f1034 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 @@ -29,6 +29,8 @@ import org.apache.kafka.common.message.DeleteGroupsResponseData; import org.apache.kafka.common.message.DeleteTopicsRequestData; import org.apache.kafka.common.message.DeleteTopicsResponseData; +import org.apache.kafka.common.message.DescribeAclsRequestData; +import org.apache.kafka.common.message.DescribeAclsResponseData; import org.apache.kafka.common.message.DescribeDelegationTokenRequestData; import org.apache.kafka.common.message.DescribeDelegationTokenResponseData; import org.apache.kafka.common.message.DescribeGroupsRequestData; @@ -100,8 +102,6 @@ import org.apache.kafka.common.requests.DeleteAclsResponse; import org.apache.kafka.common.requests.DeleteRecordsRequest; import org.apache.kafka.common.requests.DeleteRecordsResponse; -import org.apache.kafka.common.requests.DescribeAclsRequest; -import org.apache.kafka.common.requests.DescribeAclsResponse; import org.apache.kafka.common.requests.DescribeConfigsRequest; import org.apache.kafka.common.requests.DescribeConfigsResponse; import org.apache.kafka.common.requests.DescribeLogDirsRequest; @@ -178,7 +178,7 @@ public Struct parseResponse(short version, ByteBuffer buffer) { WriteTxnMarkersResponse.schemaVersions()), TXN_OFFSET_COMMIT(28, "TxnOffsetCommit", false, RecordBatch.MAGIC_VALUE_V2, TxnOffsetCommitRequestData.SCHEMAS, TxnOffsetCommitResponseData.SCHEMAS), - DESCRIBE_ACLS(29, "DescribeAcls", DescribeAclsRequest.schemaVersions(), DescribeAclsResponse.schemaVersions()), + DESCRIBE_ACLS(29, "DescribeAcls", DescribeAclsRequestData.SCHEMAS, DescribeAclsResponseData.SCHEMAS), CREATE_ACLS(30, "CreateAcls", CreateAclsRequest.schemaVersions(), CreateAclsResponse.schemaVersions()), DELETE_ACLS(31, "DeleteAcls", DeleteAclsRequest.schemaVersions(), DeleteAclsResponse.schemaVersions()), DESCRIBE_CONFIGS(32, "DescribeConfigs", DescribeConfigsRequest.schemaVersions(), 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 64701529403de..31e0fb370410c 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 @@ -139,7 +139,7 @@ public static AbstractResponse parseResponse(ApiKeys apiKey, Struct struct, shor case TXN_OFFSET_COMMIT: return new TxnOffsetCommitResponse(struct, version); case DESCRIBE_ACLS: - return new DescribeAclsResponse(struct); + return new DescribeAclsResponse(struct, version); case CREATE_ACLS: return new CreateAclsResponse(struct); case DELETE_ACLS: diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java index 6e3f24a1438c6..d6984299029d0 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java @@ -18,98 +18,97 @@ import org.apache.kafka.common.acl.AccessControlEntryFilter; import org.apache.kafka.common.acl.AclBindingFilter; +import org.apache.kafka.common.acl.AclOperation; +import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.errors.UnsupportedVersionException; +import org.apache.kafka.common.message.DescribeAclsRequestData; +import org.apache.kafka.common.message.DescribeAclsResponseData; import org.apache.kafka.common.protocol.ApiKeys; -import org.apache.kafka.common.protocol.types.Schema; import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.resource.PatternType; import org.apache.kafka.common.resource.ResourcePatternFilter; +import org.apache.kafka.common.resource.ResourceType; import java.nio.ByteBuffer; import java.util.Collections; -import static org.apache.kafka.common.protocol.CommonFields.HOST_FILTER; -import static org.apache.kafka.common.protocol.CommonFields.OPERATION; -import static org.apache.kafka.common.protocol.CommonFields.PERMISSION_TYPE; -import static org.apache.kafka.common.protocol.CommonFields.PRINCIPAL_FILTER; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_NAME_FILTER; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_PATTERN_TYPE_FILTER; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_TYPE; public class DescribeAclsRequest extends AbstractRequest { - private static final Schema DESCRIBE_ACLS_REQUEST_V0 = new Schema( - RESOURCE_TYPE, - RESOURCE_NAME_FILTER, - PRINCIPAL_FILTER, - HOST_FILTER, - OPERATION, - PERMISSION_TYPE); - - /** - * V1 sees a new `RESOURCE_PATTERN_TYPE_FILTER` that controls how the filter handles different resource pattern types. - * For more info, see {@link PatternType}. - * - * Also, when the quota is violated, brokers will respond to a version 1 or later request before throttling. - */ - private static final Schema DESCRIBE_ACLS_REQUEST_V1 = new Schema( - RESOURCE_TYPE, - RESOURCE_NAME_FILTER, - RESOURCE_PATTERN_TYPE_FILTER, - PRINCIPAL_FILTER, - HOST_FILTER, - OPERATION, - PERMISSION_TYPE); - - public static Schema[] schemaVersions() { - return new Schema[]{DESCRIBE_ACLS_REQUEST_V0, DESCRIBE_ACLS_REQUEST_V1}; - } public static class Builder extends AbstractRequest.Builder { - private final AclBindingFilter filter; + private final DescribeAclsRequestData data; public Builder(AclBindingFilter filter) { super(ApiKeys.DESCRIBE_ACLS); - this.filter = filter; + ResourcePatternFilter patternFilter = filter.patternFilter(); + AccessControlEntryFilter entryFilter = filter.entryFilter(); + data = new DescribeAclsRequestData() + .setHostFilter(entryFilter.host()) + .setOperation(entryFilter.operation().code()) + .setPermissionType(entryFilter.permissionType().code()) + .setPrincipalFilter(entryFilter.principal()) + .setResourceNameFilter(patternFilter.name()) + .setResourcePatternType(patternFilter.patternType().code()) + .setResourceType(patternFilter.resourceType().code()); } @Override public DescribeAclsRequest build(short version) { - return new DescribeAclsRequest(filter, version); + validate(version); + return new DescribeAclsRequest(data, version); } @Override public String toString() { - return "(type=DescribeAclsRequest, filter=" + filter + ")"; + return data.toString(); + } + + private void validate(short version) { + if (version == 0 + && data.resourcePatternType() != PatternType.LITERAL.code() + && data.resourcePatternType() != PatternType.ANY.code()) { + throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); + } + + if (data.resourcePatternType() == PatternType.UNKNOWN.code() + || data.resourceType() == ResourceType.UNKNOWN.code() + || data.permissionType() == AclPermissionType.UNKNOWN.code() + || data.operation() == AclOperation.UNKNOWN.code()) { + throw new IllegalArgumentException("Filter contain UNKNOWN elements"); + } } } - private final AclBindingFilter filter; + private final DescribeAclsRequestData data; - DescribeAclsRequest(AclBindingFilter filter, short version) { + public DescribeAclsRequest(DescribeAclsRequestData data, short version) { super(ApiKeys.DESCRIBE_ACLS, version); - this.filter = filter; - - validate(filter, version); + this.data = data; } public DescribeAclsRequest(Struct struct, short version) { super(ApiKeys.DESCRIBE_ACLS, version); - ResourcePatternFilter resourceFilter = RequestUtils.resourcePatternFilterFromStructFields(struct); - AccessControlEntryFilter entryFilter = RequestUtils.aceFilterFromStructFields(struct); - this.filter = new AclBindingFilter(resourceFilter, entryFilter); + this.data = new DescribeAclsRequestData(struct, version); + } + + public DescribeAclsRequestData data() { + return data; } @Override protected Struct toStruct() { - Struct struct = new Struct(ApiKeys.DESCRIBE_ACLS.requestSchema(version())); - RequestUtils.resourcePatternFilterSetStructFields(filter.patternFilter(), struct); - RequestUtils.aceFilterSetStructFields(filter.entryFilter(), struct); - return struct; + return data.toStruct(version()); } @Override public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable) { - return new DescribeAclsResponse(throttleTimeMs, ApiError.fromThrowable(throwable), Collections.emptySet()); + ApiError error = ApiError.fromThrowable(throwable); + DescribeAclsResponseData response = new DescribeAclsResponseData() + .setThrottleTimeMs(throttleTimeMs) + .setErrorCode(error.error().code()) + .setErrorMessage(error.message()) + .setResources(Collections.emptyList()); + return new DescribeAclsResponse(response); } public static DescribeAclsRequest parse(ByteBuffer buffer, short version) { @@ -117,18 +116,16 @@ public static DescribeAclsRequest parse(ByteBuffer buffer, short version) { } public AclBindingFilter filter() { - return filter; + ResourcePatternFilter rpf = new ResourcePatternFilter( + ResourceType.fromCode(data.resourceType()), + data.resourceNameFilter(), + PatternType.fromCode(data.resourcePatternType())); + AccessControlEntryFilter acef = new AccessControlEntryFilter( + data.principalFilter(), + data.hostFilter(), + AclOperation.fromCode(data.operation()), + AclPermissionType.fromCode(data.permissionType())); + return new AclBindingFilter(rpf, acef); } - private void validate(AclBindingFilter filter, short version) { - if (version == 0 - && filter.patternFilter().patternType() != PatternType.LITERAL - && filter.patternFilter().patternType() != PatternType.ANY) { - throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); - } - - if (filter.isUnknown()) { - throw new IllegalArgumentException("Filter contain UNKNOWN elements"); - } - } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java index 341845cf2b10a..7c2670810f0bd 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java @@ -17,164 +17,82 @@ package org.apache.kafka.common.requests; -import org.apache.kafka.common.acl.AccessControlEntry; -import org.apache.kafka.common.acl.AclBinding; -import org.apache.kafka.common.resource.PatternType; -import org.apache.kafka.common.resource.ResourcePattern; -import org.apache.kafka.common.errors.UnsupportedVersionException; -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.HashMap; import java.util.List; import java.util.Map; +import java.util.Map.Entry; -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.HOST; -import static org.apache.kafka.common.protocol.CommonFields.OPERATION; -import static org.apache.kafka.common.protocol.CommonFields.PERMISSION_TYPE; -import static org.apache.kafka.common.protocol.CommonFields.PRINCIPAL; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_NAME; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_PATTERN_TYPE; -import static org.apache.kafka.common.protocol.CommonFields.RESOURCE_TYPE; -import static org.apache.kafka.common.protocol.CommonFields.THROTTLE_TIME_MS; +import org.apache.kafka.common.acl.AccessControlEntry; +import org.apache.kafka.common.acl.AclBinding; +import org.apache.kafka.common.acl.AclOperation; +import org.apache.kafka.common.acl.AclPermissionType; +import org.apache.kafka.common.errors.UnsupportedVersionException; +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.protocol.ApiKeys; +import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.protocol.types.Struct; +import org.apache.kafka.common.resource.PatternType; +import org.apache.kafka.common.resource.ResourcePattern; +import org.apache.kafka.common.resource.ResourceType; public class DescribeAclsResponse extends AbstractResponse { - private final static String RESOURCES_KEY_NAME = "resources"; - private final static String ACLS_KEY_NAME = "acls"; - - private static final Schema DESCRIBE_ACLS_RESOURCE_V0 = new Schema( - RESOURCE_TYPE, - RESOURCE_NAME, - new Field(ACLS_KEY_NAME, new ArrayOf(new Schema( - PRINCIPAL, - HOST, - OPERATION, - PERMISSION_TYPE)))); - - /** - * V1 sees a new `RESOURCE_PATTERN_TYPE` that defines the type of the resource pattern. - * - * For more info, see {@link PatternType}. - */ - private static final Schema DESCRIBE_ACLS_RESOURCE_V1 = new Schema( - RESOURCE_TYPE, - RESOURCE_NAME, - RESOURCE_PATTERN_TYPE, - new Field(ACLS_KEY_NAME, new ArrayOf(new Schema( - PRINCIPAL, - HOST, - OPERATION, - PERMISSION_TYPE)))); - - private static final Schema DESCRIBE_ACLS_RESPONSE_V0 = new Schema( - THROTTLE_TIME_MS, - ERROR_CODE, - ERROR_MESSAGE, - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_ACLS_RESOURCE_V0), "The resources and their associated ACLs.")); - - /** - * V1 sees a new `RESOURCE_PATTERN_TYPE` field added to DESCRIBE_ACLS_RESOURCE_V1, that describes how the resource name is interpreted - * and version was bumped to indicate that, on quota violation, brokers send out responses before throttling. - * - * For more info, see {@link PatternType}. - */ - private static final Schema DESCRIBE_ACLS_RESPONSE_V1 = new Schema( - THROTTLE_TIME_MS, - ERROR_CODE, - ERROR_MESSAGE, - new Field(RESOURCES_KEY_NAME, new ArrayOf(DESCRIBE_ACLS_RESOURCE_V1), "The resources and their associated ACLs.")); - - public static Schema[] schemaVersions() { - return new Schema[]{DESCRIBE_ACLS_RESPONSE_V0, DESCRIBE_ACLS_RESPONSE_V1}; - } - private final int throttleTimeMs; - private final ApiError error; - private final Collection acls; + private final DescribeAclsResponseData data; - public DescribeAclsResponse(int throttleTimeMs, ApiError error, Collection acls) { - this.throttleTimeMs = throttleTimeMs; - this.error = error; - this.acls = acls; + public DescribeAclsResponse(DescribeAclsResponseData data) { + this.data = data; } - public DescribeAclsResponse(Struct struct) { - this.throttleTimeMs = struct.get(THROTTLE_TIME_MS); - this.error = new ApiError(struct); - this.acls = new ArrayList<>(); - for (Object resourceStructObj : struct.getArray(RESOURCES_KEY_NAME)) { - Struct resourceStruct = (Struct) resourceStructObj; - ResourcePattern pattern = RequestUtils.resourcePatternromStructFields(resourceStruct); - for (Object aclDataStructObj : resourceStruct.getArray(ACLS_KEY_NAME)) { - Struct aclDataStruct = (Struct) aclDataStructObj; - AccessControlEntry entry = RequestUtils.aceFromStructFields(aclDataStruct); - this.acls.add(new AclBinding(pattern, entry)); - } - } + public DescribeAclsResponse(Struct struct, short version) { + this.data = new DescribeAclsResponseData(struct, version); } @Override protected Struct toStruct(short version) { validate(version); - - Struct struct = new Struct(ApiKeys.DESCRIBE_ACLS.responseSchema(version)); - struct.set(THROTTLE_TIME_MS, throttleTimeMs); - error.write(struct); - - Map> resourceToData = new HashMap<>(); - for (AclBinding acl : acls) { - resourceToData - .computeIfAbsent(acl.pattern(), k -> new ArrayList<>()) - .add(acl.entry()); - } - - List resourceStructs = new ArrayList<>(); - for (Map.Entry> tuple : resourceToData.entrySet()) { - ResourcePattern resource = tuple.getKey(); - Struct resourceStruct = struct.instance(RESOURCES_KEY_NAME); - RequestUtils.resourcePatternSetStructFields(resource, resourceStruct); - List dataStructs = new ArrayList<>(); - for (AccessControlEntry entry : tuple.getValue()) { - Struct dataStruct = resourceStruct.instance(ACLS_KEY_NAME); - RequestUtils.aceSetStructFields(entry, dataStruct); - dataStructs.add(dataStruct); - } - resourceStruct.set(ACLS_KEY_NAME, dataStructs.toArray()); - resourceStructs.add(resourceStruct); - } - struct.set(RESOURCES_KEY_NAME, resourceStructs.toArray()); - return struct; + return data.toStruct(version); } @Override public int throttleTimeMs() { - return throttleTimeMs; + return data.throttleTimeMs(); } public ApiError error() { - return error; + return new ApiError(Errors.forCode(data.errorCode()), data.errorMessage()); } @Override public Map errorCounts() { - return errorCounts(error.error()); + return errorCounts(Errors.forCode(data.errorCode())); } public Collection acls() { + List acls = new ArrayList<>(); + for (DescribeAclsResource resource : data.resources()) { + for (AclDescription acl : resource.acls()) { + ResourcePattern pattern = new ResourcePattern( + ResourceType.fromCode(resource.type()), + resource.name(), + PatternType.fromCode(resource.patternType())); + AccessControlEntry entry = new AccessControlEntry( + acl.principal(), + acl.host(), + AclOperation.fromCode(acl.operation()), + AclPermissionType.fromCode(acl.permissionType())); + acls.add(new AclBinding(pattern, entry)); + } + } return acls; } public static DescribeAclsResponse parse(ByteBuffer buffer, short version) { - return new DescribeAclsResponse(ApiKeys.DESCRIBE_ACLS.responseSchema(version).read(buffer)); + return new DescribeAclsResponse(ApiKeys.DESCRIBE_ACLS.responseSchema(version).read(buffer), version); } @Override @@ -184,7 +102,7 @@ public boolean shouldClientThrottle(short version) { private void validate(short version) { if (version == 0) { - final boolean unsupported = acls.stream() + final boolean unsupported = acls().stream() .map(AclBinding::pattern) .map(ResourcePattern::patternType) .anyMatch(patternType -> patternType != PatternType.LITERAL); @@ -193,9 +111,41 @@ private void validate(short version) { } } - final boolean unknown = acls.stream().anyMatch(AclBinding::isUnknown); + final boolean unknown = acls().stream().anyMatch(AclBinding::isUnknown); if (unknown) { throw new IllegalArgumentException("Contain UNKNOWN elements"); } } + + public static DescribeAclsResponse prepareResponse(int throttleTimeMs, ApiError error, Collection acls) { + Map> map = new HashMap<>(); + for (AclBinding acl : acls) { + map.computeIfAbsent(acl.pattern(), v -> new ArrayList<>()).add(acl.entry()); + } + List resources = new ArrayList<>(); + for (Entry> entry : map.entrySet()) { + ResourcePattern key = entry.getKey(); + List aclDescriptions = new ArrayList<>(); + for (AccessControlEntry ace : entry.getValue()) { + AclDescription ad = new AclDescription() + .setHost(ace.host()) + .setOperation(ace.operation().code()) + .setPermissionType(ace.permissionType().code()) + .setPrincipal(ace.principal()); + aclDescriptions.add(ad); + } + DescribeAclsResource dar = new DescribeAclsResource() + .setName(key.name()) + .setPatternType(key.patternType().code()) + .setType(key.resourceType().code()) + .setAcls(aclDescriptions); + resources.add(dar); + } + DescribeAclsResponseData data = new DescribeAclsResponseData() + .setThrottleTimeMs(throttleTimeMs) + .setErrorCode(error.error().code()) + .setErrorMessage(error.message()) + .setResources(resources); + return new DescribeAclsResponse(data); + } } 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 f23036c1e4bf3..2d691db8b4e63 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 @@ -681,17 +681,17 @@ public void testDescribeAcls() throws Exception { env.kafkaClient().setNodeApiVersions(NodeApiVersions.create()); // Test a call where we get back ACL1 and ACL2. - env.kafkaClient().prepareResponse(new DescribeAclsResponse(0, ApiError.NONE, + env.kafkaClient().prepareResponse(DescribeAclsResponse.prepareResponse(0, ApiError.NONE, asList(ACL1, ACL2))); assertCollectionIs(env.adminClient().describeAcls(FILTER1).values().get(), ACL1, ACL2); // Test a call where we get back no results. - env.kafkaClient().prepareResponse(new DescribeAclsResponse(0, ApiError.NONE, + env.kafkaClient().prepareResponse(DescribeAclsResponse.prepareResponse(0, ApiError.NONE, Collections.emptySet())); assertTrue(env.adminClient().describeAcls(FILTER2).values().get().isEmpty()); // Test a call where we get back an error. - env.kafkaClient().prepareResponse(new DescribeAclsResponse(0, + env.kafkaClient().prepareResponse(DescribeAclsResponse.prepareResponse(0, new ApiError(Errors.SECURITY_DISABLED, "Security is disabled"), Collections.emptySet())); TestUtils.assertFutureError(env.adminClient().describeAcls(FILTER2).values(), SecurityDisabledException.class); diff --git a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsRequestTest.java b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsRequestTest.java index 7d9d1b1416c1e..8be1658f938ce 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsRequestTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsRequestTest.java @@ -48,17 +48,17 @@ public class DescribeAclsRequestTest { @Test(expected = UnsupportedVersionException.class) public void shouldThrowOnV0IfPrefixed() { - new DescribeAclsRequest(PREFIXED_FILTER, V0); + new DescribeAclsRequest.Builder(PREFIXED_FILTER).build(V0); } @Test(expected = IllegalArgumentException.class) public void shouldThrowIfUnknown() { - new DescribeAclsRequest(UNKNOWN_FILTER, V0); + new DescribeAclsRequest.Builder(UNKNOWN_FILTER).build(V0); } @Test public void shouldRoundTripLiteralV0() { - final DescribeAclsRequest original = new DescribeAclsRequest(LITERAL_FILTER, V0); + final DescribeAclsRequest original = new DescribeAclsRequest.Builder(LITERAL_FILTER).build(V0); final Struct struct = original.toStruct(); final DescribeAclsRequest result = new DescribeAclsRequest(struct, V0); @@ -68,13 +68,13 @@ public void shouldRoundTripLiteralV0() { @Test public void shouldRoundTripAnyV0AsLiteral() { - final DescribeAclsRequest original = new DescribeAclsRequest(ANY_FILTER, V0); - final DescribeAclsRequest expected = new DescribeAclsRequest( + final DescribeAclsRequest original = new DescribeAclsRequest.Builder(ANY_FILTER).build(V0); + final DescribeAclsRequest expected = new DescribeAclsRequest.Builder( new AclBindingFilter(new ResourcePatternFilter( ANY_FILTER.patternFilter().resourceType(), ANY_FILTER.patternFilter().name(), PatternType.LITERAL), - ANY_FILTER.entryFilter()), V0); + ANY_FILTER.entryFilter())).build(V0); final Struct struct = original.toStruct(); final DescribeAclsRequest result = new DescribeAclsRequest(struct, V0); @@ -84,7 +84,7 @@ public void shouldRoundTripAnyV0AsLiteral() { @Test public void shouldRoundTripLiteralV1() { - final DescribeAclsRequest original = new DescribeAclsRequest(LITERAL_FILTER, V1); + final DescribeAclsRequest original = new DescribeAclsRequest.Builder(LITERAL_FILTER).build(V1); final Struct struct = original.toStruct(); final DescribeAclsRequest result = new DescribeAclsRequest(struct, V1); @@ -94,7 +94,7 @@ public void shouldRoundTripLiteralV1() { @Test public void shouldRoundTripPrefixedV1() { - final DescribeAclsRequest original = new DescribeAclsRequest(PREFIXED_FILTER, V1); + final DescribeAclsRequest original = new DescribeAclsRequest.Builder(PREFIXED_FILTER).build(V1); final Struct struct = original.toStruct(); final DescribeAclsRequest result = new DescribeAclsRequest(struct, V1); @@ -104,7 +104,7 @@ public void shouldRoundTripPrefixedV1() { @Test public void shouldRoundTripAnyV1() { - final DescribeAclsRequest original = new DescribeAclsRequest(ANY_FILTER, V1); + final DescribeAclsRequest original = new DescribeAclsRequest.Builder(ANY_FILTER).build(V1); final Struct struct = original.toStruct(); final DescribeAclsRequest result = new DescribeAclsRequest(struct, V1); diff --git a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java index 13a3ebb921eba..120c93d091ba8 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java @@ -53,30 +53,30 @@ public class DescribeAclsResponseTest { @Test(expected = UnsupportedVersionException.class) public void shouldThrowOnV0IfNotLiteral() { - new DescribeAclsResponse(10, ApiError.NONE, aclBindings(PREFIXED_ACL1)).toStruct(V0); + DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(PREFIXED_ACL1)).toStruct(V0); } @Test(expected = IllegalArgumentException.class) public void shouldThrowIfUnknown() { - new DescribeAclsResponse(10, ApiError.NONE, aclBindings(UNKNOWN_ACL)).toStruct(V0); + DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(UNKNOWN_ACL)).toStruct(V0); } @Test public void shouldRoundTripV0() { - final DescribeAclsResponse original = new DescribeAclsResponse(10, ApiError.NONE, aclBindings(LITERAL_ACL1, LITERAL_ACL2)); + final DescribeAclsResponse original = DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(LITERAL_ACL1, LITERAL_ACL2)); final Struct struct = original.toStruct(V0); - final DescribeAclsResponse result = new DescribeAclsResponse(struct); + final DescribeAclsResponse result = new DescribeAclsResponse(struct, V0); assertResponseEquals(original, result); } @Test public void shouldRoundTripV1() { - final DescribeAclsResponse original = new DescribeAclsResponse(100, ApiError.NONE, aclBindings(LITERAL_ACL1, PREFIXED_ACL1)); + final DescribeAclsResponse original = DescribeAclsResponse.prepareResponse(100, ApiError.NONE, aclBindings(LITERAL_ACL1, PREFIXED_ACL1)); final Struct struct = original.toStruct(V1); - final DescribeAclsResponse result = new DescribeAclsResponse(struct); + final DescribeAclsResponse result = new DescribeAclsResponse(struct, V1); assertResponseEquals(original, result); } 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 2a36cfe59dd1c..1c1d587270734 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 @@ -1641,7 +1641,7 @@ private DescribeAclsRequest createListAclsRequest() { } private DescribeAclsResponse createDescribeAclsResponse() { - return new DescribeAclsResponse(0, ApiError.NONE, Collections.singleton(new AclBinding( + return DescribeAclsResponse.prepareResponse(0, ApiError.NONE, Collections.singleton(new AclBinding( new ResourcePattern(ResourceType.TOPIC, "mytopic", PatternType.LITERAL), new AccessControlEntry("User:ANONYMOUS", "*", AclOperation.WRITE, AclPermissionType.ALLOW)))); } diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 12621ad6700c2..6783e474405f1 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -2189,14 +2189,14 @@ class KafkaApis(val requestChannel: RequestChannel, authorizer match { case None => sendResponseMaybeThrottle(request, requestThrottleMs => - new DescribeAclsResponse(requestThrottleMs, + DescribeAclsResponse.prepareResponse(requestThrottleMs, new ApiError(Errors.SECURITY_DISABLED, "No Authorizer is configured on the broker"), util.Collections.emptySet())) case Some(auth) => val filter = describeAclsRequest.filter val returnedAcls = new util.HashSet[AclBinding]() auth.acls(filter).asScala.foreach(returnedAcls.add) sendResponseMaybeThrottle(request, requestThrottleMs => - new DescribeAclsResponse(requestThrottleMs, ApiError.NONE, returnedAcls)) + DescribeAclsResponse.prepareResponse(requestThrottleMs, ApiError.NONE, returnedAcls)) } } From a63ab2a22a80c8c1d055ab47ec18c51bbb49d776 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Sun, 24 Nov 2019 21:14:12 +0000 Subject: [PATCH 2/3] Addressed reviews --- .../kafka/clients/admin/KafkaAdminClient.java | 2 +- .../common/requests/DescribeAclsRequest.java | 36 ++++++------- .../common/requests/DescribeAclsResponse.java | 52 +++++++++++-------- .../requests/DescribeAclsResponseTest.java | 5 +- 4 files changed, 52 insertions(+), 43 deletions(-) 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 35f444459a154..70cceb687a851 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 @@ -1705,7 +1705,7 @@ void handleResponse(AbstractResponse abstractResponse) { if (response.error().isFailure()) { future.completeExceptionally(response.error().exception()); } else { - future.complete(response.acls()); + future.complete(DescribeAclsResponse.aclBindings(response.acls())); } } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java index d6984299029d0..59790199351dc 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java @@ -30,7 +30,6 @@ import org.apache.kafka.common.resource.ResourceType; import java.nio.ByteBuffer; -import java.util.Collections; public class DescribeAclsRequest extends AbstractRequest { @@ -54,7 +53,6 @@ public Builder(AclBindingFilter filter) { @Override public DescribeAclsRequest build(short version) { - validate(version); return new DescribeAclsRequest(data, version); } @@ -62,21 +60,6 @@ public DescribeAclsRequest build(short version) { public String toString() { return data.toString(); } - - private void validate(short version) { - if (version == 0 - && data.resourcePatternType() != PatternType.LITERAL.code() - && data.resourcePatternType() != PatternType.ANY.code()) { - throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); - } - - if (data.resourcePatternType() == PatternType.UNKNOWN.code() - || data.resourceType() == ResourceType.UNKNOWN.code() - || data.permissionType() == AclPermissionType.UNKNOWN.code() - || data.operation() == AclOperation.UNKNOWN.code()) { - throw new IllegalArgumentException("Filter contain UNKNOWN elements"); - } - } } private final DescribeAclsRequestData data; @@ -84,6 +67,7 @@ private void validate(short version) { public DescribeAclsRequest(DescribeAclsRequestData data, short version) { super(ApiKeys.DESCRIBE_ACLS, version); this.data = data; + validate(version); } public DescribeAclsRequest(Struct struct, short version) { @@ -106,8 +90,7 @@ public AbstractResponse getErrorResponse(int throttleTimeMs, Throwable throwable DescribeAclsResponseData response = new DescribeAclsResponseData() .setThrottleTimeMs(throttleTimeMs) .setErrorCode(error.error().code()) - .setErrorMessage(error.message()) - .setResources(Collections.emptyList()); + .setErrorMessage(error.message()); return new DescribeAclsResponse(response); } @@ -128,4 +111,19 @@ public AclBindingFilter filter() { return new AclBindingFilter(rpf, acef); } + private void validate(short version) { + if (version == 0 + && data.resourcePatternType() != PatternType.LITERAL.code() + && data.resourcePatternType() != PatternType.ANY.code()) { + throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); + } + + if (data.resourcePatternType() == PatternType.UNKNOWN.code() + || data.resourceType() == ResourceType.UNKNOWN.code() + || data.permissionType() == AclPermissionType.UNKNOWN.code() + || data.operation() == AclOperation.UNKNOWN.code()) { + throw new IllegalArgumentException("Filter contain UNKNOWN elements"); + } + } + } diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java index 7c2670810f0bd..84e5d14d7b92e 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java @@ -72,23 +72,8 @@ public Map errorCounts() { return errorCounts(Errors.forCode(data.errorCode())); } - public Collection acls() { - List acls = new ArrayList<>(); - for (DescribeAclsResource resource : data.resources()) { - for (AclDescription acl : resource.acls()) { - ResourcePattern pattern = new ResourcePattern( - ResourceType.fromCode(resource.type()), - resource.name(), - PatternType.fromCode(resource.patternType())); - AccessControlEntry entry = new AccessControlEntry( - acl.principal(), - acl.host(), - AclOperation.fromCode(acl.operation()), - AclPermissionType.fromCode(acl.permissionType())); - acls.add(new AclBinding(pattern, entry)); - } - } - return acls; + public List acls() { + return data.resources(); } public static DescribeAclsResponse parse(ByteBuffer buffer, short version) { @@ -103,20 +88,45 @@ public boolean shouldClientThrottle(short version) { private void validate(short version) { if (version == 0) { final boolean unsupported = acls().stream() - .map(AclBinding::pattern) - .map(ResourcePattern::patternType) - .anyMatch(patternType -> patternType != PatternType.LITERAL); + .anyMatch(acl -> acl.patternType() != PatternType.LITERAL.code()); if (unsupported) { throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); } } - final boolean unknown = acls().stream().anyMatch(AclBinding::isUnknown); + final boolean unknown = acls().stream() + .map(DescribeAclsResponse::aclBindings) + .anyMatch(bindings -> bindings.stream().anyMatch(AclBinding::isUnknown)); if (unknown) { throw new IllegalArgumentException("Contain UNKNOWN elements"); } } + private static List aclBindings(DescribeAclsResource resource) { + List acls = new ArrayList<>(); + for (AclDescription acl : resource.acls()) { + ResourcePattern pattern = new ResourcePattern( + ResourceType.fromCode(resource.type()), + resource.name(), + PatternType.fromCode(resource.patternType())); + AccessControlEntry entry = new AccessControlEntry( + acl.principal(), + acl.host(), + AclOperation.fromCode(acl.operation()), + AclPermissionType.fromCode(acl.permissionType())); + acls.add(new AclBinding(pattern, entry)); + } + return acls; + } + + public static List aclBindings(List resources) { + List acls = new ArrayList<>(); + for (DescribeAclsResource resource : resources) { + acls.addAll(aclBindings(resource)); + } + return acls; + } + public static DescribeAclsResponse prepareResponse(int throttleTimeMs, ApiError error, Collection acls) { Map> map = new HashMap<>(); for (AclBinding acl : acls) { diff --git a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java index 120c93d091ba8..bf4e7d9f6d4cf 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java @@ -22,6 +22,7 @@ import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.errors.UnsupportedVersionException; +import org.apache.kafka.common.message.DescribeAclsResponseData.DescribeAclsResource; import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.resource.PatternType; import org.apache.kafka.common.resource.ResourcePattern; @@ -82,8 +83,8 @@ public void shouldRoundTripV1() { } private static void assertResponseEquals(final DescribeAclsResponse original, final DescribeAclsResponse actual) { - final Set originalBindings = new HashSet<>(original.acls()); - final Set actualBindings = new HashSet<>(actual.acls()); + final Set originalBindings = new HashSet<>(original.acls()); + final Set actualBindings = new HashSet<>(actual.acls()); assertEquals(originalBindings, actualBindings); } From e951dadbb78112e68eabbaafaa4b744f80aec3d7 Mon Sep 17 00:00:00 2001 From: Mickael Maison Date: Tue, 26 Nov 2019 18:51:11 +0000 Subject: [PATCH 3/3] 2nd round of reviews --- .../common/requests/DescribeAclsRequest.java | 11 +- .../common/requests/DescribeAclsResponse.java | 13 ++- .../common/message/DescribeAclsRequest.json | 5 +- .../common/message/DescribeAclsResponse.json | 5 +- .../requests/DescribeAclsResponseTest.java | 102 ++++++++++++++---- .../common/requests/RequestResponseTest.java | 26 +++-- .../main/scala/kafka/server/KafkaApis.scala | 2 +- 7 files changed, 125 insertions(+), 39 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java index 59790199351dc..f5c58f4e7043d 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsRequest.java @@ -112,10 +112,13 @@ public AclBindingFilter filter() { } private void validate(short version) { - if (version == 0 - && data.resourcePatternType() != PatternType.LITERAL.code() - && data.resourcePatternType() != PatternType.ANY.code()) { - throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); + if (version == 0) { + if (data.resourcePatternType() == PatternType.ANY.code()) { + data.setResourcePatternType(PatternType.LITERAL.code()); + } + if (data.resourcePatternType() != PatternType.LITERAL.code()) { + throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types"); + } } if (data.resourcePatternType() == PatternType.UNKNOWN.code() diff --git a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java index 84e5d14d7b92e..bbf90b2619388 100644 --- a/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java +++ b/clients/src/main/java/org/apache/kafka/common/requests/DescribeAclsResponse.java @@ -94,11 +94,14 @@ private void validate(short version) { } } - final boolean unknown = acls().stream() - .map(DescribeAclsResponse::aclBindings) - .anyMatch(bindings -> bindings.stream().anyMatch(AclBinding::isUnknown)); - if (unknown) { - throw new IllegalArgumentException("Contain UNKNOWN elements"); + for (DescribeAclsResource resource : acls()) { + if (resource.patternType() == PatternType.UNKNOWN.code() || resource.type() == ResourceType.UNKNOWN.code()) + throw new IllegalArgumentException("Contain UNKNOWN elements"); + for (AclDescription acl : resource.acls()) { + if (acl.operation() == AclOperation.UNKNOWN.code() || acl.permissionType() == AclPermissionType.UNKNOWN.code()) { + throw new IllegalArgumentException("Contain UNKNOWN elements"); + } + } } } diff --git a/clients/src/main/resources/common/message/DescribeAclsRequest.json b/clients/src/main/resources/common/message/DescribeAclsRequest.json index 258dddc1f51ef..5cb4f00c997d7 100644 --- a/clients/src/main/resources/common/message/DescribeAclsRequest.json +++ b/clients/src/main/resources/common/message/DescribeAclsRequest.json @@ -18,8 +18,9 @@ "type": "request", "name": "DescribeAclsRequest", // Version 1 adds resource pattern type. - "validVersions": "0-1", - "flexibleVersions": "none", + // Version 2 enables flexible versions. + "validVersions": "0-2", + "flexibleVersions": "2+", "fields": [ { "name": "ResourceType", "type": "int8", "versions": "0+", "about": "The resource type." }, diff --git a/clients/src/main/resources/common/message/DescribeAclsResponse.json b/clients/src/main/resources/common/message/DescribeAclsResponse.json index 9f04de01b8757..d3c626947d1b4 100644 --- a/clients/src/main/resources/common/message/DescribeAclsResponse.json +++ b/clients/src/main/resources/common/message/DescribeAclsResponse.json @@ -19,8 +19,9 @@ "name": "DescribeAclsResponse", // Version 1 adds PatternType. // Starting in version 1, on quota violation, brokers send out responses before throttling. - "validVersions": "0-1", - "flexibleVersions": "none", + // Version 2 enables flexible versions. + "validVersions": "0-2", + "flexibleVersions": "2+", "fields": [ { "name": "ThrottleTimeMs", "type": "int32", "versions": "0+", "about": "The duration in milliseconds for which the request was throttled due to a quota violation, or zero if the request did not violate any quota." }, diff --git a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java index bf4e7d9f6d4cf..ec07bc35d212b 100644 --- a/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java +++ b/clients/src/test/java/org/apache/kafka/common/requests/DescribeAclsResponseTest.java @@ -22,7 +22,10 @@ import org.apache.kafka.common.acl.AclOperation; import org.apache.kafka.common.acl.AclPermissionType; import org.apache.kafka.common.errors.UnsupportedVersionException; +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.protocol.Errors; import org.apache.kafka.common.protocol.types.Struct; import org.apache.kafka.common.resource.PatternType; import org.apache.kafka.common.resource.ResourcePattern; @@ -30,6 +33,7 @@ import org.junit.Test; import java.util.Arrays; +import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Set; @@ -40,46 +44,86 @@ public class DescribeAclsResponseTest { private static final short V0 = 0; private static final short V1 = 1; - private static final AclBinding LITERAL_ACL1 = new AclBinding(new ResourcePattern(ResourceType.TOPIC, "foo", PatternType.LITERAL), - new AccessControlEntry("User:ANONYMOUS", "127.0.0.1", AclOperation.READ, AclPermissionType.DENY)); - - private static final AclBinding LITERAL_ACL2 = new AclBinding(new ResourcePattern(ResourceType.GROUP, "group", PatternType.LITERAL), - new AccessControlEntry("User:*", "127.0.0.1", AclOperation.WRITE, AclPermissionType.ALLOW)); - - private static final AclBinding PREFIXED_ACL1 = new AclBinding(new ResourcePattern(ResourceType.GROUP, "prefix", PatternType.PREFIXED), - new AccessControlEntry("User:*", "127.0.0.1", AclOperation.CREATE, AclPermissionType.ALLOW)); - - private static final AclBinding UNKNOWN_ACL = new AclBinding(new ResourcePattern(ResourceType.UNKNOWN, "foo", PatternType.LITERAL), - new AccessControlEntry("User:ANONYMOUS", "127.0.0.1", AclOperation.READ, AclPermissionType.DENY)); + private static final AclDescription ALLOW_CREATE_ACL = buildAclDescription( + "127.0.0.1", + "User:ANONYMOUS", + AclOperation.CREATE, + AclPermissionType.ALLOW); + + private static final AclDescription DENY_READ_ACL = buildAclDescription( + "127.0.0.1", + "User:ANONYMOUS", + AclOperation.READ, + AclPermissionType.DENY); + + private static final DescribeAclsResource UNKNOWN_ACL = buildResource( + "foo", + ResourceType.UNKNOWN, + PatternType.LITERAL, + Collections.singletonList(DENY_READ_ACL)); + + private static final DescribeAclsResource PREFIXED_ACL1 = buildResource( + "prefix", + ResourceType.GROUP, + PatternType.PREFIXED, + Collections.singletonList(ALLOW_CREATE_ACL)); + + private static final DescribeAclsResource LITERAL_ACL1 = buildResource( + "foo", + ResourceType.TOPIC, + PatternType.LITERAL, + Collections.singletonList(ALLOW_CREATE_ACL)); + + private static final DescribeAclsResource LITERAL_ACL2 = buildResource( + "group", + ResourceType.GROUP, + PatternType.LITERAL, + Collections.singletonList(DENY_READ_ACL)); @Test(expected = UnsupportedVersionException.class) public void shouldThrowOnV0IfNotLiteral() { - DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(PREFIXED_ACL1)).toStruct(V0); + buildResponse(10, Errors.NONE, Collections.singletonList(PREFIXED_ACL1)).toStruct(V0); } @Test(expected = IllegalArgumentException.class) public void shouldThrowIfUnknown() { - DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(UNKNOWN_ACL)).toStruct(V0); + buildResponse(10, Errors.NONE, Collections.singletonList(UNKNOWN_ACL)).toStruct(V0); } @Test public void shouldRoundTripV0() { - final DescribeAclsResponse original = DescribeAclsResponse.prepareResponse(10, ApiError.NONE, aclBindings(LITERAL_ACL1, LITERAL_ACL2)); + List resources = Arrays.asList(LITERAL_ACL1, LITERAL_ACL2); + final DescribeAclsResponse original = buildResponse(10, Errors.NONE, resources); final Struct struct = original.toStruct(V0); final DescribeAclsResponse result = new DescribeAclsResponse(struct, V0); - assertResponseEquals(original, result); + + final DescribeAclsResponse result2 = DescribeAclsResponse.prepareResponse(10, ApiError.NONE, DescribeAclsResponse.aclBindings(resources)); + assertResponseEquals(original, result2); } @Test public void shouldRoundTripV1() { - final DescribeAclsResponse original = DescribeAclsResponse.prepareResponse(100, ApiError.NONE, aclBindings(LITERAL_ACL1, PREFIXED_ACL1)); + List resources = Arrays.asList(LITERAL_ACL1, PREFIXED_ACL1); + final DescribeAclsResponse original = buildResponse(100, Errors.NONE, resources); final Struct struct = original.toStruct(V1); final DescribeAclsResponse result = new DescribeAclsResponse(struct, V1); - assertResponseEquals(original, result); + + final DescribeAclsResponse result2 = DescribeAclsResponse.prepareResponse(100, ApiError.NONE, DescribeAclsResponse.aclBindings(resources)); + assertResponseEquals(original, result2); + } + + @Test + public void testAclBindings() { + final AclBinding original = new AclBinding(new ResourcePattern(ResourceType.TOPIC, "foo", PatternType.LITERAL), + new AccessControlEntry("User:ANONYMOUS", "127.0.0.1", AclOperation.CREATE, AclPermissionType.ALLOW)); + + final List result = DescribeAclsResponse.aclBindings(Collections.singletonList(LITERAL_ACL1)); + assertEquals(1, result.size()); + assertEquals(original, result.get(0)); } private static void assertResponseEquals(final DescribeAclsResponse original, final DescribeAclsResponse actual) { @@ -89,7 +133,27 @@ private static void assertResponseEquals(final DescribeAclsResponse original, fi assertEquals(originalBindings, actualBindings); } - private static List aclBindings(final AclBinding... bindings) { - return Arrays.asList(bindings); + private static DescribeAclsResponse buildResponse(int throttleTimeMs, Errors error, List resources) { + return new DescribeAclsResponse(new DescribeAclsResponseData() + .setThrottleTimeMs(throttleTimeMs) + .setErrorCode(error.code()) + .setErrorMessage(error.message()) + .setResources(resources)); + } + + private static DescribeAclsResource buildResource(String name, ResourceType type, PatternType patternType, List acls) { + return new DescribeAclsResource() + .setName(name) + .setType(type.code()) + .setPatternType(patternType.code()) + .setAcls(acls); + } + + private static AclDescription buildAclDescription(String host, String principal, AclOperation operation, AclPermissionType permission) { + return new AclDescription() + .setHost(host) + .setPrincipal(principal) + .setOperation(operation.code()) + .setPermissionType(permission.code()); } } \ No newline at end of file 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 1c1d587270734..f506de3394053 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 @@ -60,6 +60,9 @@ import org.apache.kafka.common.message.DeleteTopicsRequestData; 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.DescribeAclsResponseData.AclDescription; +import org.apache.kafka.common.message.DescribeAclsResponseData.DescribeAclsResource; import org.apache.kafka.common.message.DescribeGroupsRequestData; import org.apache.kafka.common.message.DescribeGroupsResponseData; import org.apache.kafka.common.message.DescribeGroupsResponseData.DescribedGroup; @@ -362,8 +365,8 @@ public void testSerialization() throws Exception { checkRequest(createTxnOffsetCommitRequest(), true); checkErrorResponse(createTxnOffsetCommitRequest(), new UnknownServerException(), true); checkResponse(createTxnOffsetCommitResponse(), 0, true); - checkRequest(createListAclsRequest(), true); - checkErrorResponse(createListAclsRequest(), new SecurityDisabledException("Security is not enabled."), true); + checkRequest(createDescribeAclsRequest(), true); + checkErrorResponse(createDescribeAclsRequest(), new SecurityDisabledException("Security is not enabled."), true); checkResponse(createDescribeAclsResponse(), ApiKeys.DESCRIBE_ACLS.latestVersion(), true); checkRequest(createCreateAclsRequest(), true); checkErrorResponse(createCreateAclsRequest(), new SecurityDisabledException("Security is not enabled."), true); @@ -1634,16 +1637,27 @@ private TxnOffsetCommitResponse createTxnOffsetCommitResponse() { return new TxnOffsetCommitResponse(0, errorPerPartitions); } - private DescribeAclsRequest createListAclsRequest() { + private DescribeAclsRequest createDescribeAclsRequest() { return new DescribeAclsRequest.Builder(new AclBindingFilter( new ResourcePatternFilter(ResourceType.TOPIC, "mytopic", PatternType.LITERAL), new AccessControlEntryFilter(null, null, AclOperation.ANY, AclPermissionType.ANY))).build(); } private DescribeAclsResponse createDescribeAclsResponse() { - return DescribeAclsResponse.prepareResponse(0, ApiError.NONE, Collections.singleton(new AclBinding( - new ResourcePattern(ResourceType.TOPIC, "mytopic", PatternType.LITERAL), - new AccessControlEntry("User:ANONYMOUS", "*", AclOperation.WRITE, AclPermissionType.ALLOW)))); + DescribeAclsResponseData data = new DescribeAclsResponseData() + .setErrorCode(Errors.NONE.code()) + .setErrorMessage(Errors.NONE.message()) + .setThrottleTimeMs(0) + .setResources(Collections.singletonList(new DescribeAclsResource() + .setType(ResourceType.TOPIC.code()) + .setName("mytopic") + .setPatternType(PatternType.LITERAL.code()) + .setAcls(Collections.singletonList(new AclDescription() + .setHost("*") + .setOperation(AclOperation.WRITE.code()) + .setPermissionType(AclPermissionType.ALLOW.code()) + .setPrincipal("User:ANONYMOUS"))))); + return new DescribeAclsResponse(data); } private CreateAclsRequest createCreateAclsRequest() { diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 6783e474405f1..7669a6f4aa1d1 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -2194,7 +2194,7 @@ class KafkaApis(val requestChannel: RequestChannel, case Some(auth) => val filter = describeAclsRequest.filter val returnedAcls = new util.HashSet[AclBinding]() - auth.acls(filter).asScala.foreach(returnedAcls.add) + auth.acls(filter).forEach(returnedAcls.add) sendResponseMaybeThrottle(request, requestThrottleMs => DescribeAclsResponse.prepareResponse(requestThrottleMs, ApiError.NONE, returnedAcls)) }