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
Original file line number Diff line number Diff line change
Expand Up @@ -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()));
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,117 +18,115 @@

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<DescribeAclsRequest> {
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);
return new DescribeAclsRequest(data, version);
}

@Override
public String toString() {
return "(type=DescribeAclsRequest, filter=" + filter + ")";
return data.toString();
}
}

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;
validate(version);
}

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());
return new DescribeAclsResponse(response);
}

public static DescribeAclsRequest parse(ByteBuffer buffer, short version) {
return new DescribeAclsRequest(ApiKeys.DESCRIBE_ACLS.parseRequest(version, buffer), 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");
private void validate(short version) {
if (version == 0) {
if (data.resourcePatternType() == PatternType.ANY.code()) {
data.setResourcePatternType(PatternType.LITERAL.code());

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I tried changing the JSON request definition but I couldn't get it to generate the right Data class so I've had to put this logic :(
As far as I can tell DeleteAclsRequest has the same requirement.

}
if (data.resourcePatternType() != PatternType.LITERAL.code()) {
throw new UnsupportedVersionException("Version 0 only supports literal resource pattern types");
}
}

if (filter.isUnknown()) {
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");
}
}

}
Loading