This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new 6df86daf fix(metrics): enforce instance data source bindings (#1373)
6df86daf is described below
commit 6df86dafb6615ecb9d3b50c6db600db45b99e687
Author: btlqql <[email protected]>
AuthorDate: Tue Aug 11 15:23:36 2026 +0800
fix(metrics): enforce instance data source bindings (#1373)
---
.../studio/cluster/metrics/MetricsService.java | 14 +++++
.../request/MetricsDataSourceQueryRequest.java | 2 +
.../cluster/metrics/MetricsControllerTest.java | 24 ++++++++
.../studio/cluster/metrics/MetricsServiceTest.java | 72 ++++++++++++++++++++++
web/src/api/metrics.test.ts | 2 +
web/src/api/metrics.ts | 5 +-
web/src/components/MetricsExplorer.tsx | 4 +-
.../components/__tests__/MetricsExplorer.test.tsx | 4 +-
8 files changed, 122 insertions(+), 5 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
index 9750c920..5a3c94a6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
@@ -27,6 +27,7 @@ import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import java.math.BigDecimal;
+import java.util.List;
import java.util.Map;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -82,6 +83,7 @@ public class MetricsService {
MetricQueryDTO resolvedQuery = resolveMetricQuery(request.getQuery());
validateQueryWindow(resolvedQuery);
DataSourceVO dataSource = settingsService.getDataSource(dataSourceKey);
+ validateInstanceBinding(dataSourceKey, dataSource,
request.getInstanceId());
MetricsSource source =
metricsSourceFactory.create(toConfig(dataSource, request));
log.debug("Querying data source {} (type={}): start={}, end={},
step={}",
dataSourceKey, dataSource.getType(),
@@ -89,6 +91,18 @@ public class MetricsService {
return source.query(resolvedQuery);
}
+ private void validateInstanceBinding(String dataSourceKey, DataSourceVO
dataSource, String instanceId) {
+ List<String> bindings = dataSource.getInstanceIds();
+ if (bindings == null || bindings.isEmpty()) {
+ return;
+ }
+ String normalizedInstanceId = StringUtils.hasText(instanceId) ?
instanceId.strip() : null;
+ if (normalizedInstanceId == null ||
!bindings.contains(normalizedInstanceId)) {
+ throw badRequest("Data source " + dataSourceKey + " is not
available for instance "
+ + (normalizedInstanceId == null ? "<missing>" :
normalizedInstanceId));
+ }
+ }
+
private MetricsDataSourceConfig toConfig(DataSourceVO dataSource,
MetricsDataSourceQueryRequest request) {
MetricsDataSourceConfig config = new MetricsDataSourceConfig();
config.setName(dataSource.getName());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/model/request/MetricsDataSourceQueryRequest.java
b/server/src/main/java/org/apache/rocketmq/studio/model/request/MetricsDataSourceQueryRequest.java
index d910af80..61411ce6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/model/request/MetricsDataSourceQueryRequest.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/model/request/MetricsDataSourceQueryRequest.java
@@ -36,6 +36,8 @@ public class MetricsDataSourceQueryRequest {
@NotNull(message = "query is required")
private MetricQueryDTO query;
+ private String instanceId;
+
private String username;
private String password;
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
index 33bead3b..2118aec7 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
@@ -153,6 +153,30 @@ class MetricsControllerTest {
verifyNoInteractions(metricsService);
}
+ @Test
+ void dataSourceQueryShouldForwardSelectedInstance() throws Exception {
+ when(metricsService.queryByDataSource(any(),
any())).thenReturn(MetricDataVO.builder()
+ .resultType("matrix")
+ .series(List.of())
+ .warnings(List.of())
+ .build());
+
+ mockMvc.perform(post("/api/metrics/query/datasource")
+ .param("key", "ds-prom-1")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+ {"instanceId":"instance-a",
+ "query":{"metric":"up","start":1784107658,
+ "end":1784108558,"step":"30s"}}
+ """))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200));
+
+ verify(metricsService).queryByDataSource(
+ org.mockito.ArgumentMatchers.eq("ds-prom-1"),
+ argThat(request ->
"instance-a".equals(request.getInstanceId())));
+ }
+
@Test
void queryShouldReturnBadRequestWhenFieldTypeIsInvalid() throws Exception {
mockMvc.perform(post("/api/metrics/query")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
index 99f9f563..96efdd3e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
@@ -309,6 +309,18 @@ class MetricsServiceTest {
.build();
}
+ private MetricsDataSourceQueryRequest dataSourceRequest(String instanceId)
{
+ MetricsDataSourceQueryRequest request = new
MetricsDataSourceQueryRequest();
+ request.setInstanceId(instanceId);
+ request.setQuery(MetricQueryDTO.builder()
+ .metric("cpu")
+ .start(1700000000L)
+ .end(1700003600L)
+ .step("1m")
+ .build());
+ return request;
+ }
+
private void assertBadRequest(MetricQueryDTO query, String message) {
assertThatExceptionOfType(PrometheusException.class)
.isThrownBy(() -> metricsService.query(query))
@@ -349,6 +361,66 @@ class MetricsServiceTest {
verify(metricsSource).query(any(MetricQueryDTO.class));
}
+ @Test
+ void queryByDataSourceShouldAllowMatchingInstanceBinding() {
+ MetricsDataSourceQueryRequest request =
dataSourceRequest("instance-a");
+ DataSourceVO dataSource = DataSourceVO.builder()
+ .key("ds-1")
+ .name("prometheus-a")
+ .type("prometheus")
+ .url("http://prometheus:9090")
+ .auth("none")
+ .instanceIds(List.of("instance-a", "instance-b"))
+ .build();
+ when(settingsService.getDataSource("ds-1")).thenReturn(dataSource);
+
when(metricsSourceFactory.create(any(MetricsDataSourceConfig.class))).thenReturn(metricsSource);
+
when(metricsSource.query(any(MetricQueryDTO.class))).thenReturn(emptyMetricData());
+
+ metricsService.queryByDataSource("ds-1", request);
+
+ verify(metricsSource).query(any(MetricQueryDTO.class));
+ }
+
+ @Test
+ void queryByDataSourceShouldRejectDifferentInstanceBinding() {
+ MetricsDataSourceQueryRequest request =
dataSourceRequest("instance-b");
+ DataSourceVO dataSource = DataSourceVO.builder()
+ .key("ds-1")
+ .name("prometheus-a")
+ .type("prometheus")
+ .url("http://prometheus:9090")
+ .instanceIds(List.of("instance-a"))
+ .build();
+ when(settingsService.getDataSource("ds-1")).thenReturn(dataSource);
+
+ assertThatExceptionOfType(PrometheusException.class)
+ .isThrownBy(() -> metricsService.queryByDataSource("ds-1",
request))
+ .satisfies(exception -> {
+ assertThat(exception.getStatusCode()).isEqualTo(400);
+ assertThat(exception.getMessage())
+ .isEqualTo("Data source ds-1 is not available for
instance instance-b");
+ });
+ verifyNoInteractions(metricsSourceFactory, metricsSource);
+ }
+
+ @Test
+ void queryByDataSourceShouldRejectMissingInstanceForBoundSource() {
+ MetricsDataSourceQueryRequest request = dataSourceRequest(null);
+ DataSourceVO dataSource = DataSourceVO.builder()
+ .key("ds-1")
+ .name("prometheus-a")
+ .type("prometheus")
+ .url("http://prometheus:9090")
+ .instanceIds(List.of("instance-a"))
+ .build();
+ when(settingsService.getDataSource("ds-1")).thenReturn(dataSource);
+
+ assertThatExceptionOfType(PrometheusException.class)
+ .isThrownBy(() -> metricsService.queryByDataSource("ds-1",
request))
+ .satisfies(exception ->
assertThat(exception.getStatusCode()).isEqualTo(400));
+ verifyNoInteractions(metricsSourceFactory, metricsSource);
+ }
+
@Test
void queryByDataSourceShouldRejectMissingKey() {
MetricsDataSourceQueryRequest request = new
MetricsDataSourceQueryRequest();
diff --git a/web/src/api/metrics.test.ts b/web/src/api/metrics.test.ts
index adb098bc..cd08c512 100644
--- a/web/src/api/metrics.test.ts
+++ b/web/src/api/metrics.test.ts
@@ -101,6 +101,7 @@ describe('metrics API', () => {
it('posts a data-source query by key and returns its result', async () => {
const dsQuery = {
key: 'ds-prom-1',
+ instanceId: 'instance-a',
query: { metric: 'up', start: 1, end: 2, step: '1m' },
};
const result = {
@@ -119,6 +120,7 @@ describe('metrics API', () => {
expect(config.params).toEqual({ key: 'ds-prom-1' });
expect(JSON.parse(config.data)).toEqual({
query: dsQuery.query,
+ instanceId: 'instance-a',
username: undefined,
password: undefined,
bearerToken: undefined,
diff --git a/web/src/api/metrics.ts b/web/src/api/metrics.ts
index 0e48c65c..0b6bf8ed 100644
--- a/web/src/api/metrics.ts
+++ b/web/src/api/metrics.ts
@@ -108,6 +108,7 @@ export async function queryMetrics(query: MetricQuery) {
export interface DataSourceQuery {
key: string;
query: MetricQuery;
+ instanceId?: string;
username?: string;
password?: string;
bearerToken?: string;
@@ -116,10 +117,10 @@ export interface DataSourceQuery {
// Runs a PromQL range query against a configured data source (key identifies
the
// persisted source; credentials are optional and fall back to the stored
config).
export async function queryByDataSource(params: DataSourceQuery) {
- const { key, query, username, password, bearerToken } = params;
+ const { key, query, instanceId, username, password, bearerToken } = params;
const res = await client.post<{ data: MetricData }>(
'/metrics/query/datasource',
- { query, username, password, bearerToken },
+ { query, instanceId, username, password, bearerToken },
{ params: { key } },
);
return res.data.data;
diff --git a/web/src/components/MetricsExplorer.tsx
b/web/src/components/MetricsExplorer.tsx
index eab10121..2e127020 100644
--- a/web/src/components/MetricsExplorer.tsx
+++ b/web/src/components/MetricsExplorer.tsx
@@ -294,7 +294,7 @@ const MetricsExplorer = ({ instanceId }:
MetricsExplorerProps) => {
setQueryError(false);
try {
const result = dataSourceKeyRef.current
- ? await queryByDataSource({ key: dataSourceKeyRef.current, query })
+ ? await queryByDataSource({ key: dataSourceKeyRef.current, query,
instanceId })
: await queryMetrics(query);
if (currentRequest === requestId.current) setData(result);
} catch {
@@ -306,7 +306,7 @@ const MetricsExplorer = ({ instanceId }:
MetricsExplorerProps) => {
if (currentRequest === requestId.current) setQueryLoading(false);
}
},
- [],
+ [instanceId],
);
useEffect(() => {
diff --git a/web/src/components/__tests__/MetricsExplorer.test.tsx
b/web/src/components/__tests__/MetricsExplorer.test.tsx
index 4d75374a..b2b466f0 100644
--- a/web/src/components/__tests__/MetricsExplorer.test.tsx
+++ b/web/src/components/__tests__/MetricsExplorer.test.tsx
@@ -242,11 +242,12 @@ describe('MetricsExplorer', () => {
url: '',
auth: 'None',
status: 'healthy',
+ instanceIds: ['instance-a'],
},
]);
vi.mocked(queryByDataSource).mockResolvedValue(metricData);
- renderWithProviders(<MetricsExplorer />);
+ renderWithProviders(<MetricsExplorer instanceId="instance-a" />);
await screen.findByRole('combobox', { name: '数据源' });
await user.click(screen.getByRole('combobox', { name: '数据源' }));
@@ -259,6 +260,7 @@ describe('MetricsExplorer', () => {
await waitFor(() =>
expect(queryByDataSource).toHaveBeenCalledWith({
key: 'ds-prom-1',
+ instanceId: 'instance-a',
query: {
metric: 'sum(rate(rocketmq_messages_in_total[1m])) by (cluster,
node_id)',
start: 1_799_996_400,