From c9a44e62740179b6c490e208baf00fab8b33d332 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Sergio=20Pe=C3=B1a?=
Date: Tue, 28 May 2019 15:25:25 -0500
Subject: [PATCH 1/5] If authorizedOperations is NULL, then allow access to the
Topic
NULL means the Kafka broker does not support authorizedOperations
---
.../engine/AuthorizationTopicAccessValidator.java | 6 +++++-
.../ksql/engine/TopicAccessValidatorFactory.java | 5 +++++
.../AuthorizationTopicAccessValidatorTest.java | 15 +++++++++++++++
3 files changed, 25 insertions(+), 1 deletion(-)
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..a2057aabcef0 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
@@ -38,6 +38,11 @@ public static TopicAccessValidator create(
final MetaStore metaStore
) {
if (isKafkaAuthorizerEnabled(serviceContext.getAdminClient())) {
+ // This service works only on Kakfa 2.3 and newer versions. There was no way to detect
+ // the version of Kafka during this point, so I left the version check during the
+ // AuthorizationTopicAccessValidator validation
+ // (see AuthorizationTopicAccessValidator#hasAcccess)
+
LOG.info("KSQL topic authorization checks enabled.");
return new AuthorizationTopicAccessValidator();
}
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..ccd2f7b9a284 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,21 @@ private Statement givenStatement(final String sql) {
return ksqlEngine.prepare(ksqlEngine.parse(sql).get(0)).getStatement();
}
+ @Test
+ public void shouldAllowAnyOperationIfPermissionsAreNull() {
+ // This test case verifies permissions will not work if the Kafka broker returns NULL
+
+ // 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:
From 219533231751d7654c39fa45f7b185003fea6898 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Sergio=20Pe=C3=B1a?=
Date: Tue, 28 May 2019 16:10:12 -0500
Subject: [PATCH 2/5] Return Dummy validator if Cluster authorizedOperations is
Null
---
.../engine/TopicAccessValidatorFactory.java | 23 +++++++-----
.../ksql/services/KafkaClusterUtil.java | 15 ++++++++
.../TopicAccessValidatorFactoryTest.java | 35 ++++++++++++++++---
3 files changed, 59 insertions(+), 14 deletions(-)
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 a2057aabcef0..6f0103588df2 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,14 +37,17 @@ public static TopicAccessValidator create(
final ServiceContext serviceContext,
final MetaStore metaStore
) {
- if (isKafkaAuthorizerEnabled(serviceContext.getAdminClient())) {
- // This service works only on Kakfa 2.3 and newer versions. There was no way to detect
- // the version of Kafka during this point, so I left the version check during the
- // AuthorizationTopicAccessValidator validation
- // (see AuthorizationTopicAccessValidator#hasAcccess)
+ final AdminClient adminClient = serviceContext.getAdminClient();
- LOG.info("KSQL topic authorization checks enabled.");
- return new AuthorizationTopicAccessValidator();
+ if (isKafkaAuthorizerEnabled(adminClient)) {
+ if (KafkaClusterUtil.isAuthorizedOperationsSupported(adminClient)) {
+ LOG.info("KSQL topic authorization checks enabled.");
+ return new AuthorizationTopicAccessValidator();
+ }
+
+ LOG.info("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
@@ -67,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/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;
}
From 2a60a5bea111391e77daee27ce94735467ea24ef Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Sergio=20Pe=C3=B1a?=
Date: Wed, 29 May 2019 13:39:26 -0500
Subject: [PATCH 3/5] Switch message from info to warn
---
.../io/confluent/ksql/engine/TopicAccessValidatorFactory.java | 2 +-
1 file changed, 1 insertion(+), 1 deletion(-)
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 6f0103588df2..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
@@ -45,7 +45,7 @@ public static TopicAccessValidator create(
return new AuthorizationTopicAccessValidator();
}
- LOG.info("The Kafka broker has an authorization service enabled, but the Kafka "
+ 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.");
}
From 1ef65b1108562701e4c420683acd564669cc1e95 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Sergio=20Pe=C3=B1a?=
Date: Wed, 29 May 2019 14:06:11 -0500
Subject: [PATCH 4/5] Use dummy validator for tests
---
.../confluent/ksql/rest/server/computation/RecoveryTest.java | 4 +++-
.../ksql/rest/server/resources/StreamedQueryResourceTest.java | 4 +++-
2 files changed, 6 insertions(+), 2 deletions(-)
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
From 1ac30330565b6c54e6b2ff605e7a3708bf0a8ee9 Mon Sep 17 00:00:00 2001
From: =?UTF-8?q?Sergio=20Pe=C3=B1a?=
Date: Wed, 29 May 2019 16:33:31 -0500
Subject: [PATCH 5/5] Remove unecessary comment from tests
---
.../ksql/engine/AuthorizationTopicAccessValidatorTest.java | 2 --
1 file changed, 2 deletions(-)
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 ccd2f7b9a284..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
@@ -99,8 +99,6 @@ private Statement givenStatement(final String sql) {
@Test
public void shouldAllowAnyOperationIfPermissionsAreNull() {
- // This test case verifies permissions will not work if the Kafka broker returns NULL
-
// Given:
givenTopicPermissions(TOPIC_1, null);
final Statement statement = givenStatement("SELECT * FROM " + STREAM_TOPIC_1 + ";");