diff --git a/ksql-engine/src/main/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidator.java b/ksql-engine/src/main/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidator.java
index 417fdd745bf8..bff4a17abb5e 100644
--- a/ksql-engine/src/main/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidator.java
+++ b/ksql-engine/src/main/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidator.java
@@ -33,6 +33,8 @@
/**
* Checks if a {@link ServiceContext} has access to the source and target topics of transient
* and persistent query statements.
+ *
+ * This validator only works on Kakfa 2.3 or later.
*/
public class AuthorizationTopicAccessValidator implements TopicAccessValidator {
@Override
@@ -121,7 +123,9 @@ private void checkAccess(
final Set authorizedOperations = serviceContext.getTopicClient()
.describeTopic(topicName).authorizedOperations();
- if (!authorizedOperations.contains(operation)) {
+ // Kakfa 2.2 or lower do not support authorizedOperations(). In case of running on a
+ // unsupported broker version, then the authorizeOperation will be null.
+ if (authorizedOperations != null && !authorizedOperations.contains(operation)) {
// This error message is similar to what Kafka throws when it cannot access the topic
// due to an authorization error. I used this message to keep a consistent message.
throw new KsqlException(String.format(
diff --git a/ksql-engine/src/main/java/io/confluent/ksql/engine/TopicAccessValidatorFactory.java b/ksql-engine/src/main/java/io/confluent/ksql/engine/TopicAccessValidatorFactory.java
index 703129741c27..f2403886945c 100644
--- a/ksql-engine/src/main/java/io/confluent/ksql/engine/TopicAccessValidatorFactory.java
+++ b/ksql-engine/src/main/java/io/confluent/ksql/engine/TopicAccessValidatorFactory.java
@@ -37,9 +37,17 @@ public static TopicAccessValidator create(
final ServiceContext serviceContext,
final MetaStore metaStore
) {
- if (isKafkaAuthorizerEnabled(serviceContext.getAdminClient())) {
- LOG.info("KSQL topic authorization checks enabled.");
- return new AuthorizationTopicAccessValidator();
+ final AdminClient adminClient = serviceContext.getAdminClient();
+
+ if (isKafkaAuthorizerEnabled(adminClient)) {
+ if (KafkaClusterUtil.isAuthorizedOperationsSupported(adminClient)) {
+ LOG.info("KSQL topic authorization checks enabled.");
+ return new AuthorizationTopicAccessValidator();
+ }
+
+ LOG.warn("The Kafka broker has an authorization service enabled, but the Kafka "
+ + "version does not support authorizedOperations(). "
+ + "KSQL topic authorization checks will not be enabled.");
}
// Dummy validator if a Kafka authorizer is not enabled
@@ -62,8 +70,10 @@ private static boolean isKafkaAuthorizerEnabled(final AdminClient adminClient) {
if (e.getCause() instanceof ClusterAuthorizationException) {
return true;
}
- }
- return false;
+ // Throw the unknown exception to avoid leaving the Server unsecured if a different
+ // error was thrown
+ throw e;
+ }
}
}
diff --git a/ksql-engine/src/main/java/io/confluent/ksql/services/KafkaClusterUtil.java b/ksql-engine/src/main/java/io/confluent/ksql/services/KafkaClusterUtil.java
index 61e51ff1ce69..4fff2fa48224 100644
--- a/ksql-engine/src/main/java/io/confluent/ksql/services/KafkaClusterUtil.java
+++ b/ksql-engine/src/main/java/io/confluent/ksql/services/KafkaClusterUtil.java
@@ -21,8 +21,11 @@
import java.util.Collection;
import java.util.Collections;
import java.util.Map;
+
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.Config;
+import org.apache.kafka.clients.admin.DescribeClusterOptions;
+import org.apache.kafka.clients.admin.DescribeClusterResult;
import org.apache.kafka.common.Node;
import org.apache.kafka.common.config.ConfigResource;
import org.slf4j.Logger;
@@ -35,6 +38,18 @@ private KafkaClusterUtil() {
}
+ public static boolean isAuthorizedOperationsSupported(final AdminClient adminClient) {
+ try {
+ final DescribeClusterResult authorizedOperations = adminClient.describeCluster(
+ new DescribeClusterOptions().includeAuthorizedOperations(true)
+ );
+
+ return authorizedOperations.authorizedOperations().get() != null;
+ } catch (Exception e) {
+ throw new KsqlServerException("Could not get Kafka authorized operations!", e);
+ }
+ }
+
public static Config getConfig(final AdminClient adminClient) {
try {
final Collection brokers = adminClient.describeCluster().nodes().get();
diff --git a/ksql-engine/src/test/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidatorTest.java b/ksql-engine/src/test/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidatorTest.java
index 1f3782f97a94..a5e63b198991 100644
--- a/ksql-engine/src/test/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidatorTest.java
+++ b/ksql-engine/src/test/java/io/confluent/ksql/engine/AuthorizationTopicAccessValidatorTest.java
@@ -97,6 +97,19 @@ private Statement givenStatement(final String sql) {
return ksqlEngine.prepare(ksqlEngine.parse(sql).get(0)).getStatement();
}
+ @Test
+ public void shouldAllowAnyOperationIfPermissionsAreNull() {
+ // Given:
+ givenTopicPermissions(TOPIC_1, null);
+ final Statement statement = givenStatement("SELECT * FROM " + STREAM_TOPIC_1 + ";");
+
+ // When:
+ accessValidator.validate(serviceContext, metaStore, statement);
+
+ // Then:
+ // Above command should not throw any exception
+ }
+
@Test
public void shouldSingleSelectWithReadPermissionsAllowed() {
// Given:
diff --git a/ksql-engine/src/test/java/io/confluent/ksql/engine/TopicAccessValidatorFactoryTest.java b/ksql-engine/src/test/java/io/confluent/ksql/engine/TopicAccessValidatorFactoryTest.java
index 8e03ba932208..1a3b5f5081f1 100644
--- a/ksql-engine/src/test/java/io/confluent/ksql/engine/TopicAccessValidatorFactoryTest.java
+++ b/ksql-engine/src/test/java/io/confluent/ksql/engine/TopicAccessValidatorFactoryTest.java
@@ -15,10 +15,12 @@
package io.confluent.ksql.engine;
+import static org.easymock.EasyMock.anyObject;
import static org.easymock.EasyMock.expect;
import static org.easymock.EasyMock.mock;
import static org.easymock.EasyMock.replay;
import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.any;
import static org.hamcrest.Matchers.instanceOf;
import static org.hamcrest.Matchers.is;
import static org.hamcrest.Matchers.not;
@@ -29,13 +31,17 @@
import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.Set;
+
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.Config;
import org.apache.kafka.clients.admin.ConfigEntry;
+import org.apache.kafka.clients.admin.DescribeClusterOptions;
import org.apache.kafka.clients.admin.DescribeClusterResult;
import org.apache.kafka.clients.admin.DescribeConfigsResult;
import org.apache.kafka.common.KafkaFuture;
import org.apache.kafka.common.Node;
+import org.apache.kafka.common.acl.AclOperation;
import org.apache.kafka.common.config.ConfigResource;
import org.easymock.EasyMock;
import org.easymock.EasyMockRunner;
@@ -66,7 +72,7 @@ public void setUp() {
@Test
public void shouldReturnAuthorizationValidator() {
// Given:
- givenKafkaAuthorizer("an-authorizer-class");
+ givenKafkaAuthorizer("an-authorizer-class", Collections.emptySet());
// When:
final TopicAccessValidator validator = TopicAccessValidatorFactory.create(serviceContext, null);
@@ -78,7 +84,19 @@ public void shouldReturnAuthorizationValidator() {
@Test
public void shouldReturnDummyValidator() {
// Given:
- givenKafkaAuthorizer("");
+ givenKafkaAuthorizer("", Collections.emptySet());
+
+ // When:
+ final TopicAccessValidator validator = TopicAccessValidatorFactory.create(serviceContext, null);
+
+ // Then
+ assertThat(validator, not(instanceOf(AuthorizationTopicAccessValidator.class)));
+ }
+
+ @Test
+ public void shouldReturnDummyValidatorIfAuthorizedOperationsReturnNull() {
+ // Given:
+ givenKafkaAuthorizer("an-authorizer-class", null);
// When:
final TopicAccessValidator validator = TopicAccessValidatorFactory.create(serviceContext, null);
@@ -87,8 +105,13 @@ public void shouldReturnDummyValidator() {
assertThat(validator, not(instanceOf(AuthorizationTopicAccessValidator.class)));
}
- private void givenKafkaAuthorizer(final String className) {
- expect(adminClient.describeCluster()).andReturn(describeClusterResult());
+ private void givenKafkaAuthorizer(
+ final String className,
+ final Set authOperations
+ ) {
+ expect(adminClient.describeCluster()).andReturn(describeClusterResult(authOperations));
+ expect(adminClient.describeCluster(anyObject()))
+ .andReturn(describeClusterResult(authOperations));
expect(adminClient.describeConfigs(describeBrokerRequest()))
.andReturn(describeBrokerResult(Collections.singletonList(
new ConfigEntry(KAFKA_AUTHORIZER_CLASS_NAME, className)
@@ -97,10 +120,12 @@ private void givenKafkaAuthorizer(final String className) {
replay(adminClient);
}
- private DescribeClusterResult describeClusterResult() {
+ private DescribeClusterResult describeClusterResult(final Set authOperations) {
final Collection nodes = Collections.singletonList(node);
final DescribeClusterResult describeClusterResult = EasyMock.mock(DescribeClusterResult.class);
expect(describeClusterResult.nodes()).andReturn(KafkaFuture.completedFuture(nodes));
+ expect(describeClusterResult.authorizedOperations())
+ .andReturn(KafkaFuture.completedFuture(authOperations));
replay(describeClusterResult);
return describeClusterResult;
}
diff --git a/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/computation/RecoveryTest.java b/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/computation/RecoveryTest.java
index 9ca8c9d38b00..4dea6ab77c26 100644
--- a/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/computation/RecoveryTest.java
+++ b/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/computation/RecoveryTest.java
@@ -171,7 +171,9 @@ private class KsqlServer {
Duration.ofMillis(0),
()->{},
Injectors.DEFAULT,
- TopicAccessValidatorFactory.create(serviceContext, ksqlEngine.getMetaStore()));
+ (sc, metastore, statement) -> {
+ return;
+ });
this.statementExecutor = new StatementExecutor(
ksqlConfig,
ksqlEngine,
diff --git a/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/resources/StreamedQueryResourceTest.java b/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/resources/StreamedQueryResourceTest.java
index d7389f3d1a33..bdc42c1ee589 100644
--- a/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/resources/StreamedQueryResourceTest.java
+++ b/ksql-rest-app/src/test/java/io/confluent/ksql/rest/server/resources/StreamedQueryResourceTest.java
@@ -144,7 +144,9 @@ public void setup() {
DISCONNECT_CHECK_INTERVAL,
COMMAND_QUEUE_CATCHUP_TIMOEUT,
activenessRegistrar,
- TopicAccessValidatorFactory.create(serviceContext, mockKsqlEngine.getMetaStore()));
+ (sc, metastore, statement) -> {
+ return;
+ });
}
@Test