This is an automated email from the ASF dual-hosted git repository.

mjsax pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 076ea535fa3 KAFKA-17164: Enforce 'application.server' <server>:<port> 
format at config level. (#22202)
076ea535fa3 is described below

commit 076ea535fa39050d27c5ab7d8ea75361d1c4b0a2
Author: ChickenchickenLove <[email protected]>
AuthorDate: Wed Jun 24 02:21:54 2026 +0900

    KAFKA-17164: Enforce 'application.server' <server>:<port> format at config 
level. (#22202)
    
    Implements KIP-1245.
    
    Enforces the `application.server` host:port format during StreamsConfig
    construction.
    
    Reviewers: Matthias J. Sax <[email protected]>, Nikita Shupletsov
    <[email protected]>, Nilesh Kumar <[email protected]>
---
 docs/streams/upgrade-guide.md                      |  2 +
 .../org/apache/kafka/streams/StreamsConfig.java    |  2 +
 .../ApplicationServerConfigValidator.java          | 47 ++++++++++++++++++++++
 .../org/apache/kafka/streams/state/HostInfo.java   | 15 ++++++-
 .../apache/kafka/streams/StreamsConfigTest.java    | 20 ++++++++-
 .../apache/kafka/streams/state/HostInfoTest.java   |  8 ++++
 6 files changed, 91 insertions(+), 3 deletions(-)

diff --git a/docs/streams/upgrade-guide.md b/docs/streams/upgrade-guide.md
index 387bef80920..0951515f865 100644
--- a/docs/streams/upgrade-guide.md
+++ b/docs/streams/upgrade-guide.md
@@ -69,6 +69,8 @@ Since 2.6.0 release, Kafka Streams depends on a RocksDB 
version that requires Ma
 
 Kafka Streams no longer emits a WARN from `KafkaStreams#cleanUp()` when the 
application state directory cannot be deleted only because expected metadata 
files remain, such as `kafka-streams-process-metadata` and/or `.lock`. In this 
case, the local state cleanup is considered successful and the application 
state directory may be retained. Users who require a full local reset including 
persisted process metadata should manually delete the application state 
directory after the Kafka Streams [...]
 
+Kafka Streams now validates the `application.server` configuration when 
`StreamsConfig` is created. The value must be empty or a valid endpoint from 
which Kafka Streams can parse both host and port, such as `host:port` or 
`protocol://host:port`. Invalid values that may previously have failed later 
during startup or assignment now fail earlier with a `ConfigException`. More 
details can be found in 
[KIP-1245](https://cwiki.apache.org/confluence/display/KAFKA/KIP-1245%3A+Enforce+%27applicat
 [...]
+
 ## Streams API changes in 4.3.0
 
 **Note:** Kafka Streams 4.3.0 contains a critical native memory leak in the 
RocksDB state store layer 
([KAFKA-20616](https://issues.apache.org/jira/browse/KAFKA-20616)). The 
`ColumnFamilyOptions` for the offsets column family is not closed, and column 
family handles can leak on close-path exceptions, which under cascading task 
closes (e.g., rebalances or error-triggered recoveries) leads to unbounded 
off-heap memory growth and eventual OOM. Users running Kafka Streams should 
consider upg [...]
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 d9e1569fe36..d6099f4cec8 100644
--- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
+++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java
@@ -44,6 +44,7 @@ import 
org.apache.kafka.streams.errors.LogAndFailProcessingExceptionHandler;
 import org.apache.kafka.streams.errors.ProcessingExceptionHandler;
 import org.apache.kafka.streams.errors.ProductionExceptionHandler;
 import org.apache.kafka.streams.errors.StreamsException;
+import org.apache.kafka.streams.internals.ApplicationServerConfigValidator;
 import org.apache.kafka.streams.internals.StreamsConfigUtils;
 import org.apache.kafka.streams.internals.UpgradeFromValues;
 import org.apache.kafka.streams.kstream.SessionWindowedDeserializer;
@@ -1114,6 +1115,7 @@ public class StreamsConfig extends AbstractConfig {
             .define(APPLICATION_SERVER_CONFIG,
                     Type.STRING,
                     "",
+                    new ApplicationServerConfigValidator(),
                     Importance.LOW,
                     APPLICATION_SERVER_DOC)
             .define(BUFFERED_RECORDS_PER_PARTITION_CONFIG,
diff --git 
a/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
 
b/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
new file mode 100644
index 00000000000..55e4af8fdbe
--- /dev/null
+++ 
b/streams/src/main/java/org/apache/kafka/streams/internals/ApplicationServerConfigValidator.java
@@ -0,0 +1,47 @@
+/*
+ * 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.internals;
+
+import org.apache.kafka.common.config.ConfigDef;
+import org.apache.kafka.common.config.ConfigException;
+import org.apache.kafka.common.utils.Utils;
+import org.apache.kafka.streams.state.HostInfo;
+
+public class ApplicationServerConfigValidator implements ConfigDef.Validator {
+
+    @Override
+    public void ensureValid(final String name, final Object value) {
+        if (!(value instanceof String)) {
+            throw new ConfigException(name + " must be a string");
+        }
+
+        final String endPoint = (String) value;
+        if (Utils.isBlank(endPoint)) {
+            return;
+        }
+        try {
+            HostInfo.buildFromEndpoint(endPoint);
+        } catch (final ConfigException e) {
+            throw new ConfigException(name, value, e.getMessage());
+        }
+    }
+
+    @Override
+    public String toString() {
+        return "A host:port pair, protocol://host:port, or an empty string";
+    }
+}
diff --git a/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java 
b/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
index e2e17df0e58..1393dc8667e 100644
--- a/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
+++ b/streams/src/main/java/org/apache/kafka/streams/state/HostInfo.java
@@ -57,7 +57,20 @@ public class HostInfo {
         }
 
         final String host = getHost(endPoint);
-        final Integer port = getPort(endPoint);
+        if (Utils.isBlank(host)) {
+            throw new ConfigException(
+                String.format("Error parsing host address %s. Expected format 
host:port.", endPoint)
+            );
+        }
+
+        final Integer port;
+        try {
+            port = getPort(endPoint);
+        } catch (final NumberFormatException e) {
+            throw new ConfigException(
+                String.format("Error parsing host address %s. Expected format 
host:port.", endPoint)
+            );
+        }
 
         if (host == null || port == null) {
             throw new ConfigException(
diff --git 
a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java 
b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
index c8f358d9fe6..481c1ad6bae 100644
--- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
+++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java
@@ -251,7 +251,7 @@ public class StreamsConfigTest {
         props.put(StreamsConfig.MAX_WARMUP_REPLICAS_CONFIG, 9);
         props.put(StreamsConfig.PROBING_REBALANCE_INTERVAL_MS_CONFIG, 99_999L);
         
props.put(StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG, 
7L);
-        props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "dummy:host");
+        props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "dummy:8080");
         props.put(StreamsConfig.topicPrefix(TopicConfig.SEGMENT_BYTES_CONFIG), 
1024 * 1024);
         final StreamsConfig streamsConfig = new StreamsConfig(props);
         final Map<String, Object> returnedProps = 
streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx);
@@ -266,7 +266,7 @@ public class StreamsConfigTest {
             
returnedProps.get(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG)
         );
         assertEquals(7L, 
returnedProps.get(StreamsConfig.WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG));
-        assertEquals("dummy:host", 
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
+        assertEquals("dummy:8080", 
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
         assertEquals(1024 * 1024, 
returnedProps.get(StreamsConfig.topicPrefix(TopicConfig.SEGMENT_BYTES_CONFIG)));
     }
 
@@ -1902,6 +1902,22 @@ public class StreamsConfigTest {
         }
     }
 
+    @ParameterizedTest
+    @ValueSource(strings = {"dummy:host", "dummy:9999999999999999999999999", 
"dummy", "dummy:", ":port", ":", ":8080"})
+    public void 
shouldThrowConfigExceptionWithInvalidApplicationServerConfigValue(final String 
applicationServerConfig) {
+        props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, 
applicationServerConfig);
+        assertThrows(ConfigException.class, () -> new StreamsConfig(props));
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"", "127.0.0.1:8080", "localhost:8080", 
"[::1]:8080", "http://localhost:8080"})
+    public void shouldAcceptWithValidApplicationServerConfigValue(final String 
applicationServerConfigValue) {
+        props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, 
applicationServerConfigValue);
+        final StreamsConfig streamsConfig = new StreamsConfig(props);
+        final Map<String, Object> returnedProps = 
streamsConfig.getMainConsumerConfigs(groupId, clientId, threadIdx);
+        assertEquals(applicationServerConfigValue, 
returnedProps.get(StreamsConfig.APPLICATION_SERVER_CONFIG));
+    }
+
     @SuppressWarnings("deprecation")
     @Test
     public void 
shouldNotLogWarningWhenProcessingExceptionHandlerIsEnabledOnGlobalThread() {
diff --git 
a/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java 
b/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
index 10f7749f6a6..26dfeac7c24 100644
--- a/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
+++ b/streams/src/test/java/org/apache/kafka/streams/state/HostInfoTest.java
@@ -19,6 +19,8 @@ package org.apache.kafka.streams.state;
 import org.apache.kafka.common.config.ConfigException;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.is;
@@ -50,4 +52,10 @@ public class HostInfoTest {
     public void shouldThrowConfigExceptionForNonsenseEndPoint() {
         assertThrows(ConfigException.class, () -> 
HostInfo.buildFromEndpoint("nonsense"));
     }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"dummy:host", "dummy:9999999999999999999999999", 
"dummy", "dummy:", ":port", ":", ":8080"})
+    public void shouldThrowConfigExceptionWithInvalidEndpoint(final String 
invalidEndpoint) {
+        assertThrows(ConfigException.class, () -> 
HostInfo.buildFromEndpoint(invalidEndpoint));
+    }
 }

Reply via email to