From 2cd19aab75a033ae216db4b4fa619ca2d3e5cbc7 Mon Sep 17 00:00:00 2001 From: liuhy Date: Wed, 29 Jul 2026 00:09:40 -0700 Subject: [PATCH 1/2] [ISSUE #10683] Validate proxy metric collector address --- .../v2/common/GrpcClientSettingsManager.java | 28 ++++++++++++++--- .../common/GrpcClientSettingsManagerTest.java | 31 +++++++++++++++++++ 2 files changed, 55 insertions(+), 4 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java index ac87da8c244..4031602b438 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java @@ -34,6 +34,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.broker.client.ConsumerGroupInfo; import org.apache.rocketmq.common.ServiceThread; import org.apache.rocketmq.common.constant.LoggerName; @@ -117,10 +118,12 @@ protected Settings mergeMetric(Settings settings) { final Metric.Builder metricBuilder = Metric.newBuilder(); switch (metricCollectorMode) { case ON: - final String[] split = metricCollectorAddress.split(":"); - final String host = split[0]; - final int port = Integer.parseInt(split[1]); - Address address = Address.newBuilder().setHost(host).setPort(port).build(); + Address address = parseMetricCollectorAddress(metricCollectorAddress); + if (address == null) { + log.warn("disable client metric collector because metricCollectorAddress is invalid: {}", metricCollectorAddress); + metricBuilder.setOn(false); + break; + } final Endpoints endpoints = Endpoints.newBuilder().setScheme(AddressScheme.IPv4) .addAddresses(address).build(); metricBuilder.setOn(true).setEndpoints(endpoints); @@ -137,6 +140,23 @@ protected Settings mergeMetric(Settings settings) { return settings.toBuilder().setMetric(metric).build(); } + protected static Address parseMetricCollectorAddress(String metricCollectorAddress) { + if (StringUtils.isBlank(metricCollectorAddress)) { + return null; + } + + String[] split = metricCollectorAddress.split(":"); + if (split.length != 2 || StringUtils.isBlank(split[0]) || StringUtils.isBlank(split[1])) { + return null; + } + + try { + return Address.newBuilder().setHost(split[0]).setPort(Integer.parseInt(split[1])).build(); + } catch (NumberFormatException e) { + return null; + } + } + protected static Settings mergeSubscriptionData(Settings settings, SubscriptionGroupConfig groupConfig) { Settings.Builder resultSettingsBuilder = settings.toBuilder(); ProxyConfig proxyConfig = ConfigurationManager.getProxyConfig(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java index 4d0037a272a..96d4ca66ab9 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java @@ -31,6 +31,8 @@ import org.apache.rocketmq.common.lite.LiteSubscriptionDTO; import org.apache.rocketmq.proxy.common.ContextVariable; import org.apache.rocketmq.proxy.common.ProxyContext; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.config.MetricCollectorMode; import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; import org.apache.rocketmq.remoting.protocol.subscription.CustomizedRetryPolicy; import org.apache.rocketmq.remoting.protocol.subscription.ExponentialRetryPolicy; @@ -124,6 +126,35 @@ public void testGetSubscriptionData() { assertNull(this.grpcClientSettingsManager.removeAndGetClientSettings(context)); } + @Test + public void testMergeMetricWithValidCollectorAddress() { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1:8081"); + + Settings settings = this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()); + + assertEquals(true, settings.getMetric().getOn()); + assertEquals("127.0.0.1", settings.getMetric().getEndpoints().getAddresses(0).getHost()); + assertEquals(8081, settings.getMetric().getEndpoints().getAddresses(0).getPort()); + } + + @Test + public void testMergeMetricWithInvalidCollectorAddress() { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); + + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); + assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1"); + assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(":8081"); + assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1:not-a-port"); + assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + } + @Test public void testOfflineClientLiteSubscription_SettingsNullAndNoCachedSettings() { doReturn(null).when(grpcClientSettingsManager).getRawClientSettings(anyString()); From 621a34556224acbd826fb87526dce56b53e8b229 Mon Sep 17 00:00:00 2001 From: liuhy Date: Wed, 29 Jul 2026 03:58:28 -0700 Subject: [PATCH 2/2] [ISSUE #10683] Cover blank metric collector port --- .../grpc/v2/common/GrpcClientSettingsManager.java | 2 +- .../v2/common/GrpcClientSettingsManagerTest.java | 15 ++++++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java index 4031602b438..e72135f2d6f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManager.java @@ -140,7 +140,7 @@ protected Settings mergeMetric(Settings settings) { return settings.toBuilder().setMetric(metric).build(); } - protected static Address parseMetricCollectorAddress(String metricCollectorAddress) { + private static Address parseMetricCollectorAddress(String metricCollectorAddress) { if (StringUtils.isBlank(metricCollectorAddress)) { return null; } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java index 96d4ca66ab9..c5716299785 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java @@ -42,8 +42,10 @@ import org.junit.Test; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotEquals; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; @@ -133,7 +135,7 @@ public void testMergeMetricWithValidCollectorAddress() { Settings settings = this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()); - assertEquals(true, settings.getMetric().getOn()); + assertTrue(settings.getMetric().getOn()); assertEquals("127.0.0.1", settings.getMetric().getEndpoints().getAddresses(0).getHost()); assertEquals(8081, settings.getMetric().getEndpoints().getAddresses(0).getPort()); } @@ -143,16 +145,19 @@ public void testMergeMetricWithInvalidCollectorAddress() { ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); - assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1"); - assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress(":8081"); - assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1: "); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1:not-a-port"); - assertEquals(false, this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); } @Test