Repository: kafka
Updated Branches:
  refs/heads/0.10.2 1f4f9c74c -> 4f6135931


KAFKA-4700: Don't drop security configs in `StreamsKafkaClient`

Author: Ismael Juma <[email protected]>

Reviewers: Eno Thereska, Damian Guy, Guozhang Wang

Closes #2441 from ijuma/streams-kafka-client-drops-security-configs

(cherry picked from commit 31716590454e1bdd53a65836741f316bf4715c9d)
Signed-off-by: Guozhang Wang <[email protected]>


Project: http://git-wip-us.apache.org/repos/asf/kafka/repo
Commit: http://git-wip-us.apache.org/repos/asf/kafka/commit/4f613593
Tree: http://git-wip-us.apache.org/repos/asf/kafka/tree/4f613593
Diff: http://git-wip-us.apache.org/repos/asf/kafka/diff/4f613593

Branch: refs/heads/0.10.2
Commit: 4f61359318f2e7c25ec46d99fd134cf61de28a93
Parents: 1f4f9c7
Author: Ismael Juma <[email protected]>
Authored: Thu Jan 26 14:18:54 2017 -0800
Committer: Guozhang Wang <[email protected]>
Committed: Thu Jan 26 14:19:02 2017 -0800

----------------------------------------------------------------------
 .../org/apache/kafka/streams/StreamsConfig.java |  7 ++++
 .../processor/internals/StreamsKafkaClient.java | 23 ++++++++++-
 .../internals/StreamsKafkaClientTest.java       | 42 ++++++++++++++++++++
 3 files changed, 71 insertions(+), 1 deletion(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
----------------------------------------------------------------------
diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java 
b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
index 956ae0b..d7d6566 100644
--- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
+++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
@@ -388,6 +388,13 @@ public class StreamsConfig extends AbstractConfig {
         return PRODUCER_PREFIX + producerProp;
     }
 
+    /**
+     * Returns a copy of the config definition.
+     */
+    public static ConfigDef configDef() {
+        return new ConfigDef(CONFIG);
+    }
+
     public StreamsConfig(Map<?, ?> props) {
         super(CONFIG, props);
     }

http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java
----------------------------------------------------------------------
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java
 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java
index df590fb..afde63f 100644
--- 
a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java
+++ 
b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClient.java
@@ -24,6 +24,8 @@ import org.apache.kafka.clients.consumer.ConsumerConfig;
 import org.apache.kafka.clients.producer.ProducerConfig;
 import org.apache.kafka.common.Cluster;
 import org.apache.kafka.common.Node;
+import org.apache.kafka.common.config.AbstractConfig;
+import org.apache.kafka.common.config.ConfigDef;
 import org.apache.kafka.common.metrics.JmxReporter;
 import org.apache.kafka.common.metrics.MetricConfig;
 import org.apache.kafka.common.metrics.Metrics;
@@ -52,15 +54,34 @@ import java.util.concurrent.TimeUnit;
 
 public class StreamsKafkaClient {
 
+    private static final ConfigDef CONFIG = StreamsConfig.configDef()
+            .withClientSslSupport()
+            .withClientSaslSupport();
+
+    public static class Config extends AbstractConfig {
+
+        public static Config fromStreamsConfig(StreamsConfig streamsConfig) {
+            return new Config(streamsConfig.originals());
+        }
+
+        public Config(Map<?, ?> originals) {
+            super(CONFIG, originals, false);
+        }
+    }
+
     private final KafkaClient kafkaClient;
     private final List<MetricsReporter> reporters;
-    private final StreamsConfig streamsConfig;
+    private final Config streamsConfig;
 
     private static final int MAX_INFLIGHT_REQUESTS = 100;
 
     public StreamsKafkaClient(final StreamsConfig streamsConfig) {
+        this(Config.fromStreamsConfig(streamsConfig));
+    }
 
+    public StreamsKafkaClient(final Config streamsConfig) {
         this.streamsConfig = streamsConfig;
+
         final Time time = new SystemTime();
 
         final Map<String, String> metricTags = new LinkedHashMap<>();

http://git-wip-us.apache.org/repos/asf/kafka/blob/4f613593/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java
----------------------------------------------------------------------
diff --git 
a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java
 
b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java
new file mode 100644
index 0000000..2fb5724
--- /dev/null
+++ 
b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsKafkaClientTest.java
@@ -0,0 +1,42 @@
+/**
+  * Licensed to the Apache Software Foundation (ASF) under one or more 
contributor license agreements. See the NOTICE
+  * file distributed with this work for additional information regarding 
copyright ownership. The ASF licenses this file
+  * to You under the Apache License, Version 2.0 (the "License"); you may not 
use this file except in compliance with the
+  * License. You may obtain a copy of the License at
+  *
+  * http://www.apache.org/licenses/LICENSE-2.0
+  *
+  * Unless required by applicable law or agreed to in writing, software 
distributed under the License is distributed on
+  * an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either 
express or implied. See the License for the
+  * specific language governing permissions and limitations under the License.
+  */
+
+package org.apache.kafka.streams.processor.internals;
+
+import org.apache.kafka.common.config.AbstractConfig;
+import org.apache.kafka.common.config.SaslConfigs;
+import org.apache.kafka.streams.StreamsConfig;
+import org.junit.Test;
+
+import java.util.Properties;
+
+import static java.util.Arrays.asList;
+import static org.junit.Assert.assertEquals;
+
+public class StreamsKafkaClientTest {
+
+    @Test
+    public void testConfigFromStreamsConfig() {
+        for (final String expectedMechanism : asList("PLAIN", 
"SCRAM-SHA-512")) {
+            final Properties props = new Properties();
+            props.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, 
"some_app_id");
+            props.setProperty(SaslConfigs.SASL_MECHANISM, expectedMechanism);
+            props.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, 
"localhost:9000");
+            final StreamsConfig streamsConfig = new StreamsConfig(props);
+            final AbstractConfig config = 
StreamsKafkaClient.Config.fromStreamsConfig(streamsConfig);
+            assertEquals(expectedMechanism, 
config.values().get(SaslConfigs.SASL_MECHANISM));
+            assertEquals(expectedMechanism, 
config.getString(SaslConfigs.SASL_MECHANISM));
+        }
+    }
+
+}

Reply via email to