From d8b83840f2d6b40df1956ca48a8415df4102e983 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 01:00:10 -0700 Subject: [PATCH 1/2] [ISSUE #10788] Tolerate invalid metric collector address --- .../v2/common/GrpcClientSettingsManager.java | 48 ++++++++++++++++--- .../common/GrpcClientSettingsManagerTest.java | 32 +++++++++++++ 2 files changed, 73 insertions(+), 7 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..4eb15a364d0 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,13 +118,13 @@ 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(); - final Endpoints endpoints = Endpoints.newBuilder().setScheme(AddressScheme.IPv4) - .addAddresses(address).build(); - metricBuilder.setOn(true).setEndpoints(endpoints); + Address address = parseMetricCollectorAddress(metricCollectorAddress); + if (address == null) { + metricBuilder.setOn(false); + break; + } + metricBuilder.setOn(true).setEndpoints(Endpoints.newBuilder().setScheme(AddressScheme.IPv4) + .addAddresses(address).build()); break; case PROXY: metricBuilder.setOn(true).setEndpoints(settings.getAccessPoint()); @@ -137,6 +138,39 @@ protected Settings mergeMetric(Settings settings) { return settings.toBuilder().setMetric(metric).build(); } + private Address parseMetricCollectorAddress(String metricCollectorAddress) { + if (StringUtils.isBlank(metricCollectorAddress)) { + log.warn("Disable gRPC client metric collection because metricCollectorAddress is blank"); + return null; + } + String[] addressSegments = metricCollectorAddress.trim().split(":", -1); + if (addressSegments.length != 2) { + log.warn("Disable gRPC client metric collection because metricCollectorAddress is invalid: {}", + metricCollectorAddress); + return null; + } + String host = addressSegments[0].trim(); + String portSegment = addressSegments[1].trim(); + if (StringUtils.isBlank(host) || StringUtils.isBlank(portSegment)) { + log.warn("Disable gRPC client metric collection because metricCollectorAddress is invalid: {}", + metricCollectorAddress); + return null; + } + try { + int port = Integer.parseInt(portSegment); + if (port <= 0 || port > 65535) { + log.warn("Disable gRPC client metric collection because metricCollectorAddress port is out of range: {}", + metricCollectorAddress); + return null; + } + return Address.newBuilder().setHost(host).setPort(port).build(); + } catch (NumberFormatException e) { + log.warn("Disable gRPC client metric collection because metricCollectorAddress port is invalid: {}", + metricCollectorAddress); + 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..9cbf7de6f2f 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; @@ -76,6 +78,36 @@ public void testGetProducerData() { assertNotEquals(settings.getPublishing(), settings.getPublishing().getDefaultInstanceForType()); } + @Test + public void testMergeMetricWithValidCollectorAddress() { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1:9090"); + try { + Settings settings = this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()); + + assertEquals(true, settings.getMetric().getOn()); + assertEquals("127.0.0.1", settings.getMetric().getEndpoints().getAddresses(0).getHost()); + assertEquals(9090, settings.getMetric().getEndpoints().getAddresses(0).getPort()); + } finally { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.OFF.getModeString()); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); + } + } + + @Test + public void testMergeMetricShouldDisableMetricForInvalidCollectorAddress() { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress("localhost"); + try { + Settings settings = this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()); + + assertEquals(false, settings.getMetric().getOn()); + } finally { + ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.OFF.getModeString()); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); + } + } + @Test public void testGetSubscriptionData() { ProxyContext context = ProxyContext.create().withVal(ContextVariable.CLIENT_ID, CLIENT_ID); From f7dddf8c7ddbcee14f81f75df03e36bc7378fc71 Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:54:51 -0700 Subject: [PATCH 2/2] test(proxy): preserve metric collector config --- .../v2/common/GrpcClientSettingsManagerTest.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) 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 9cbf7de6f2f..efdd8719f87 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 @@ -80,6 +80,8 @@ public void testGetProducerData() { @Test public void testMergeMetricWithValidCollectorAddress() { + String originalMetricCollectorMode = ConfigurationManager.getProxyConfig().getMetricCollectorMode(); + String originalMetricCollectorAddress = ConfigurationManager.getProxyConfig().getMetricCollectorAddress(); ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress("127.0.0.1:9090"); try { @@ -89,13 +91,15 @@ public void testMergeMetricWithValidCollectorAddress() { assertEquals("127.0.0.1", settings.getMetric().getEndpoints().getAddresses(0).getHost()); assertEquals(9090, settings.getMetric().getEndpoints().getAddresses(0).getPort()); } finally { - ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.OFF.getModeString()); - ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); + ConfigurationManager.getProxyConfig().setMetricCollectorMode(originalMetricCollectorMode); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(originalMetricCollectorAddress); } } @Test public void testMergeMetricShouldDisableMetricForInvalidCollectorAddress() { + String originalMetricCollectorMode = ConfigurationManager.getProxyConfig().getMetricCollectorMode(); + String originalMetricCollectorAddress = ConfigurationManager.getProxyConfig().getMetricCollectorAddress(); ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.ON.getModeString()); ConfigurationManager.getProxyConfig().setMetricCollectorAddress("localhost"); try { @@ -103,8 +107,8 @@ public void testMergeMetricShouldDisableMetricForInvalidCollectorAddress() { assertEquals(false, settings.getMetric().getOn()); } finally { - ConfigurationManager.getProxyConfig().setMetricCollectorMode(MetricCollectorMode.OFF.getModeString()); - ConfigurationManager.getProxyConfig().setMetricCollectorAddress(""); + ConfigurationManager.getProxyConfig().setMetricCollectorMode(originalMetricCollectorMode); + ConfigurationManager.getProxyConfig().setMetricCollectorAddress(originalMetricCollectorAddress); } }