mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Optimize enableACL configuration
This commit is contained in:
@@ -142,11 +142,8 @@ public class GrpcServerBuilder {
|
||||
|
||||
public GrpcServerBuilder configInterceptor() {
|
||||
// grpc interceptors, including acl, logging etc.
|
||||
if (ConfigurationManager.getProxyConfig().isEnableACL()) {
|
||||
List<AccessValidator> accessValidators = ServiceProvider.load(ServiceProvider.ACL_VALIDATOR_ID, AccessValidator.class);
|
||||
if (accessValidators.isEmpty()) {
|
||||
throw new IllegalArgumentException("Load AccessValidator failed");
|
||||
}
|
||||
List<AccessValidator> accessValidators = ServiceProvider.load(ServiceProvider.ACL_VALIDATOR_ID, AccessValidator.class);
|
||||
if (!accessValidators.isEmpty()) {
|
||||
this.serverBuilder.intercept(new AuthenticationInterceptor(accessValidators));
|
||||
}
|
||||
|
||||
|
||||
+25
-20
@@ -32,6 +32,7 @@ import org.apache.rocketmq.acl.AccessValidator;
|
||||
import org.apache.rocketmq.acl.common.AclException;
|
||||
import org.apache.rocketmq.acl.common.MetadataHeader;
|
||||
import org.apache.rocketmq.acl.plain.PlainAccessResource;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
|
||||
public class AuthenticationInterceptor implements ServerInterceptor {
|
||||
private final List<AccessValidator> accessValidatorList;
|
||||
@@ -46,28 +47,32 @@ public class AuthenticationInterceptor implements ServerInterceptor {
|
||||
return new ForwardingServerCallListener.SimpleForwardingServerCallListener<R>(next.startCall(call, headers)) {
|
||||
@Override
|
||||
public void onMessage(R message) {
|
||||
try {
|
||||
GeneratedMessageV3 messageV3 = (GeneratedMessageV3) message;
|
||||
MetadataHeader metadataHeader = MetadataHeader.builder()
|
||||
.remoteAddress(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.REMOTE_ADDRESS))
|
||||
.namespace(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.NAMESPACE_ID))
|
||||
.authorization(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.AUTHORIZATION))
|
||||
.datetime(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.DATE_TIME))
|
||||
.sessionToken(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.SESSION_TOKEN))
|
||||
.requestId(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.REQUEST_ID))
|
||||
.language(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE))
|
||||
.clientVersion(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.CLIENT_VERSION))
|
||||
.protocol(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.PROTOCOL_VERSION))
|
||||
.requestCode(RequestMapping.map(messageV3.getDescriptorForType().getFullName()))
|
||||
.build();
|
||||
for (AccessValidator accessValidator : accessValidatorList) {
|
||||
AccessResource accessResource = accessValidator.parse(messageV3, metadataHeader);
|
||||
accessValidator.validate(accessResource);
|
||||
addHeader(headers, messageV3, accessResource);
|
||||
if (ConfigurationManager.getProxyConfig().isEnableACL()) {
|
||||
try {
|
||||
GeneratedMessageV3 messageV3 = (GeneratedMessageV3) message;
|
||||
MetadataHeader metadataHeader = MetadataHeader.builder()
|
||||
.remoteAddress(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.REMOTE_ADDRESS))
|
||||
.namespace(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.NAMESPACE_ID))
|
||||
.authorization(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.AUTHORIZATION))
|
||||
.datetime(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.DATE_TIME))
|
||||
.sessionToken(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.SESSION_TOKEN))
|
||||
.requestId(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.REQUEST_ID))
|
||||
.language(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.LANGUAGE))
|
||||
.clientVersion(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.CLIENT_VERSION))
|
||||
.protocol(InterceptorConstants.METADATA.get(Context.current()).get(InterceptorConstants.PROTOCOL_VERSION))
|
||||
.requestCode(RequestMapping.map(messageV3.getDescriptorForType().getFullName()))
|
||||
.build();
|
||||
for (AccessValidator accessValidator : accessValidatorList) {
|
||||
AccessResource accessResource = accessValidator.parse(messageV3, metadataHeader);
|
||||
accessValidator.validate(accessResource);
|
||||
addHeader(headers, messageV3, accessResource);
|
||||
}
|
||||
super.onMessage(message);
|
||||
} catch (AclException aclException) {
|
||||
throw new StatusRuntimeException(Status.PERMISSION_DENIED, headers);
|
||||
}
|
||||
} else {
|
||||
super.onMessage(message);
|
||||
} catch (AclException aclException) {
|
||||
throw new StatusRuntimeException(Status.PERMISSION_DENIED, headers);
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user