From 83ddcfae7ef1b1920ac6388456225359e6e02f82 Mon Sep 17 00:00:00 2001 From: yaozichen2025 Date: Tue, 11 Aug 2026 11:40:49 +0800 Subject: [PATCH] fix(proxy): harden metrics key-value parsing --- .../proxy/metrics/ProxyMetricsManager.java | 39 +++++++------- .../metrics/ProxyMetricsManagerTest.java | 54 +++++++++++++++++++ 2 files changed, 74 insertions(+), 19 deletions(-) create mode 100644 proxy/src/test/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManagerTest.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java index 81db576e3d2..aa09f9457f4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java @@ -34,6 +34,7 @@ import io.opentelemetry.sdk.metrics.export.PeriodicMetricReader; import io.opentelemetry.sdk.resources.Resource; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; @@ -146,15 +147,7 @@ public void start() throws Exception { String labels = proxyConfig.getMetricsLabel(); if (StringUtils.isNotBlank(labels)) { - List kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(labels); - for (String item : kvPairs) { - String[] split = item.split(":"); - if (split.length != 2) { - log.warn("metricsLabel is not valid: {}", labels); - continue; - } - LABEL_MAP.put(split[0], split[1]); - } + LABEL_MAP.putAll(parseKeyValuePairs("metricsLabel", labels)); } if (proxyConfig.isMetricsInDelta()) { LABEL_MAP.put(LABEL_AGGREGATION, AGGREGATION_DELTA); @@ -185,16 +178,7 @@ public void start() throws Exception { String headers = proxyConfig.getMetricsGrpcExporterHeader(); if (StringUtils.isNotBlank(headers)) { - Map headerMap = new HashMap<>(); - List kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(headers); - for (String item : kvPairs) { - String[] split = item.split(":"); - if (split.length != 2) { - log.warn("metricsGrpcExporterHeader is not valid: {}", headers); - continue; - } - headerMap.put(split[0], split[1]); - } + Map headerMap = parseKeyValuePairs("metricsGrpcExporterHeader", headers); headerMap.forEach(metricExporterBuilder::addHeader); } @@ -238,6 +222,23 @@ public void start() throws Exception { initMetrics(proxyMeter, null); } + static Map parseKeyValuePairs(String configName, String config) { + Map result = new LinkedHashMap<>(); + List kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(config); + for (int i = 0; i < kvPairs.size(); i++) { + String item = kvPairs.get(i); + int separatorIndex = item.indexOf(':'); + if (separatorIndex < 0 || StringUtils.isBlank(item.substring(0, separatorIndex))) { + log.warn("{} contains an invalid entry at position {}", configName, i); + continue; + } + String key = item.substring(0, separatorIndex).trim(); + String value = item.substring(separatorIndex + 1).trim(); + result.put(key, value); + } + return result; + } + @Override public void shutdown() throws Exception { if (proxyConfig.getMetricsExporterType() == MetricsExporterType.OTLP_GRPC) { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManagerTest.java new file mode 100644 index 00000000000..8863f155cd1 --- /dev/null +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManagerTest.java @@ -0,0 +1,54 @@ +/* + * 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.rocketmq.proxy.metrics; + +import java.util.Map; +import org.junit.Test; + +import static org.assertj.core.api.Assertions.entry; +import static org.assertj.core.api.Assertions.assertThat; + +public class ProxyMetricsManagerTest { + + @Test + public void parseKeyValuePairsShouldPreserveColonsInValues() { + Map result = ProxyMetricsManager.parseKeyValuePairs( + "metricsGrpcExporterHeader", "Authorization:Bearer token:part,endpoint:https://collector:4317"); + + assertThat(result) + .containsEntry("Authorization", "Bearer token:part") + .containsEntry("endpoint", "https://collector:4317"); + } + + @Test + public void parseKeyValuePairsShouldTrimEntriesAndAllowEmptyValues() { + Map result = ProxyMetricsManager.parseKeyValuePairs( + "metricsLabel", " cluster : production ,optional:"); + + assertThat(result) + .containsEntry("cluster", "production") + .containsEntry("optional", ""); + } + + @Test + public void parseKeyValuePairsShouldSkipMalformedEntriesAndEmptyKeys() { + Map result = ProxyMetricsManager.parseKeyValuePairs( + "metricsGrpcExporterHeader", "missing-separator, :secret,valid:value"); + + assertThat(result).containsExactly(entry("valid", "value")); + } +}