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..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 @@ -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(); } + private 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..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 @@ -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; @@ -40,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; @@ -124,6 +128,38 @@ 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()); + + assertTrue(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(""); + 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(":8081"); + 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"); + assertFalse(this.grpcClientSettingsManager.mergeMetric(Settings.getDefaultInstance()).getMetric().getOn()); + } + @Test public void testOfflineClientLiteSubscription_SettingsNullAndNoCachedSettings() { doReturn(null).when(grpcClientSettingsManager).getRawClientSettings(anyString());