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 c92cf596 fix: retain query history with request context (#1050)
c92cf596 is described below
commit c92cf5969739ed716b1a268afcd1b9f516632292
Author: aias00 <[email protected]>
AuthorDate: Thu Aug 6 01:16:47 2026 -0700
fix: retain query history with request context (#1050)
---
deploy/mysql/upgrade-query-history-context.sql | 36 +++++++
.../studio/persistence/entity/RmqMessageQuery.java | 2 +
.../studio/persistence/entity/RmqTraceQuery.java | 2 +
.../QueryHistoryProperties.java} | 44 ++++-----
.../studio/queryhistory/QueryHistoryService.java | 68 +++++++++++--
.../studio/rocketmq/RocketMQMessageProvider.java | 10 +-
server/src/main/resources/application.yml | 3 +
server/src/main/resources/db/schema.sql | 2 +
.../queryhistory/QueryHistoryServiceTest.java | 105 +++++++++++++++++++++
.../rocketmq/RocketMQMessageProviderTest.java | 3 +-
10 files changed, 237 insertions(+), 38 deletions(-)
diff --git a/deploy/mysql/upgrade-query-history-context.sql
b/deploy/mysql/upgrade-query-history-context.sql
new file mode 100644
index 00000000..8520db44
--- /dev/null
+++ b/deploy/mysql/upgrade-query-history-context.sql
@@ -0,0 +1,36 @@
+-- deploy/mysql/upgrade-query-history-context.sql
+-- Existing MySQL volumes: add query-history context columns.
+-- Fresh volumes already receive these columns from
server/src/main/resources/db/schema.sql.
+--
+-- Run once on existing deployments:
+-- docker exec -i rocketmq-studio-mysql mysql -uroot -pstudio123
rocketmq_studio < upgrade-query-history-context.sql
+
+SET @schema_name := DATABASE();
+
+SET @message_query_column_exists := (
+ SELECT COUNT(*)
+ FROM information_schema.columns
+ WHERE table_schema = @schema_name
+ AND table_name = 'rmq_message_query'
+ AND column_name = 'cluster_id'
+);
+SET @message_query_sql := IF(@message_query_column_exists = 0,
+ 'ALTER TABLE rmq_message_query ADD COLUMN cluster_id VARCHAR(255) AFTER
result_count',
+ 'SELECT 1');
+PREPARE message_query_statement FROM @message_query_sql;
+EXECUTE message_query_statement;
+DEALLOCATE PREPARE message_query_statement;
+
+SET @trace_query_column_exists := (
+ SELECT COUNT(*)
+ FROM information_schema.columns
+ WHERE table_schema = @schema_name
+ AND table_name = 'rmq_trace_query'
+ AND column_name = 'cluster_id'
+);
+SET @trace_query_sql := IF(@trace_query_column_exists = 0,
+ 'ALTER TABLE rmq_trace_query ADD COLUMN cluster_id VARCHAR(255) AFTER
consumer_count',
+ 'SELECT 1');
+PREPARE trace_query_statement FROM @trace_query_sql;
+EXECUTE trace_query_statement;
+DEALLOCATE PREPARE trace_query_statement;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqMessageQuery.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqMessageQuery.java
index 449a8178..d3ac8e7b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqMessageQuery.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqMessageQuery.java
@@ -46,6 +46,8 @@ public class RmqMessageQuery {
private Integer resultCount;
+ private String clusterId;
+
private String queriedBy;
private LocalDateTime queriedAt;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
index 2dc124be..6c645a4c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
@@ -38,6 +38,8 @@ public class RmqTraceQuery {
private Integer consumerCount;
+ private String clusterId;
+
private String queriedBy;
private LocalDateTime queriedAt;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryProperties.java
similarity index 58%
copy from
server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
copy to
server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryProperties.java
index 2dc124be..63387497 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/entity/RmqTraceQuery.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryProperties.java
@@ -14,31 +14,21 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.persistence.entity;
-
-import com.baomidou.mybatisplus.annotation.IdType;
-import com.baomidou.mybatisplus.annotation.TableId;
-import com.baomidou.mybatisplus.annotation.TableName;
-import lombok.Data;
-
-import java.time.LocalDateTime;
-
-@Data
-@TableName("rmq_trace_query")
-public class RmqTraceQuery {
-
- @TableId(type = IdType.AUTO)
- private Long id;
-
- private String msgId;
-
- private String topic;
-
- private Integer nodeCount;
-
- private Integer consumerCount;
-
- private String queriedBy;
-
- private LocalDateTime queriedAt;
+package org.apache.rocketmq.studio.queryhistory;
+
+import lombok.Getter;
+import lombok.Setter;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.stereotype.Component;
+
+@Getter
+@Setter
+@Component
+@ConfigurationProperties(prefix = "studio.query-history")
+public class QueryHistoryProperties {
+
+ /**
+ * Number of days query records are retained. A non-positive value
disables cleanup.
+ */
+ private int retentionDays = 90;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
index 18e512aa..90bad479 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
@@ -16,13 +16,18 @@
*/
package org.apache.rocketmq.studio.queryhistory;
+import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.apache.rocketmq.studio.auth.AuthenticatedUserContext;
import org.apache.rocketmq.studio.persistence.entity.RmqMessageQuery;
import org.apache.rocketmq.studio.persistence.entity.RmqTraceQuery;
import org.apache.rocketmq.studio.persistence.mapper.RmqMessageQueryMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqTraceQueryMapper;
+import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
+import java.time.Clock;
import java.time.LocalDateTime;
@Slf4j
@@ -31,14 +36,27 @@ public class QueryHistoryService {
private final RmqMessageQueryMapper messageQueryMapper;
private final RmqTraceQueryMapper traceQueryMapper;
+ private final QueryHistoryProperties properties;
+ private final Clock clock;
+ @Autowired
public QueryHistoryService(RmqMessageQueryMapper messageQueryMapper,
- RmqTraceQueryMapper traceQueryMapper) {
+ RmqTraceQueryMapper traceQueryMapper,
+ QueryHistoryProperties properties) {
+ this(messageQueryMapper, traceQueryMapper, properties,
Clock.systemUTC());
+ }
+
+ QueryHistoryService(RmqMessageQueryMapper messageQueryMapper,
+ RmqTraceQueryMapper traceQueryMapper,
+ QueryHistoryProperties properties,
+ Clock clock) {
this.messageQueryMapper = messageQueryMapper;
this.traceQueryMapper = traceQueryMapper;
+ this.properties = properties;
+ this.clock = clock;
}
- public void recordMessageQuery(String queryType, String topic, String
msgId,
+ public void recordMessageQuery(String clusterId, String queryType, String
topic, String msgId,
String tag, String key, Long startTime,
Long endTime, int resultCount) {
RmqMessageQuery query = new RmqMessageQuery();
@@ -50,19 +68,55 @@ public class QueryHistoryService {
query.setStartTime(startTime);
query.setEndTime(endTime);
query.setResultCount(resultCount);
- query.setQueriedAt(LocalDateTime.now());
+ query.setClusterId(clusterId);
+ query.setQueriedBy(AuthenticatedUserContext.currentUsernameOrSystem());
+ query.setQueriedAt(LocalDateTime.now(clock));
messageQueryMapper.insert(query);
- log.debug("Message query recorded: type={} topic={}", queryType,
topic);
+ log.debug("Message query recorded: clusterId={} type={} topic={}",
clusterId, queryType, topic);
}
- public void recordTraceQuery(String msgId, String topic, int nodeCount,
int consumerCount) {
+ public void recordTraceQuery(String clusterId, String msgId, String topic,
int nodeCount, int consumerCount) {
RmqTraceQuery query = new RmqTraceQuery();
query.setMsgId(msgId);
query.setTopic(topic);
query.setNodeCount(nodeCount);
query.setConsumerCount(consumerCount);
- query.setQueriedAt(LocalDateTime.now());
+ query.setClusterId(clusterId);
+ query.setQueriedBy(AuthenticatedUserContext.currentUsernameOrSystem());
+ query.setQueriedAt(LocalDateTime.now(clock));
traceQueryMapper.insert(query);
- log.debug("Trace query recorded: msgId={} topic={}", msgId, topic);
+ log.debug("Trace query recorded: clusterId={} msgId={} topic={}",
clusterId, msgId, topic);
+ }
+
+ @Scheduled(fixedDelayString =
"${studio.query-history.cleanup-interval:PT24H}")
+ public void purgeExpiredQueries() {
+ int retentionDays = properties.getRetentionDays();
+ if (retentionDays <= 0) {
+ return;
+ }
+
+ LocalDateTime cutoff =
LocalDateTime.now(clock).minusDays(retentionDays);
+ deleteExpiredMessageQueries(cutoff);
+ deleteExpiredTraceQueries(cutoff);
+ }
+
+ private void deleteExpiredMessageQueries(LocalDateTime cutoff) {
+ try {
+ int deleted =
messageQueryMapper.delete(Wrappers.<RmqMessageQuery>query()
+ .lt("queried_at", cutoff));
+ log.debug("Purged {} expired message query records", deleted);
+ } catch (RuntimeException e) {
+ log.warn("Failed to purge expired message query records: {}",
e.getMessage());
+ }
+ }
+
+ private void deleteExpiredTraceQueries(LocalDateTime cutoff) {
+ try {
+ int deleted =
traceQueryMapper.delete(Wrappers.<RmqTraceQuery>query()
+ .lt("queried_at", cutoff));
+ log.debug("Purged {} expired trace query records", deleted);
+ } catch (RuntimeException e) {
+ log.warn("Failed to purge expired trace query records: {}",
e.getMessage());
+ }
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
index 1f2bac1e..14d6899a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
@@ -489,8 +489,8 @@ public class RocketMQMessageProvider implements
MessageProvider {
private void recordMessageQuery(String queryType, String topic, String
msgId, String tag, String key,
Long startTime, Long endTime, int
resultCount) {
try {
- queryHistoryService.recordMessageQuery(queryType, topic, msgId,
tag, key, startTime, endTime,
- resultCount);
+ queryHistoryService.recordMessageQuery(clusterContext(),
queryType, topic, msgId, tag, key,
+ startTime, endTime, resultCount);
} catch (Exception e) {
log.warn("Failed to record message query history: {}",
e.getMessage());
}
@@ -498,12 +498,16 @@ public class RocketMQMessageProvider implements
MessageProvider {
private void recordTraceQuery(String msgId, String topic, int nodeCount,
int consumerCount) {
try {
- queryHistoryService.recordTraceQuery(msgId, topic, nodeCount,
consumerCount);
+ queryHistoryService.recordTraceQuery(clusterContext(), msgId,
topic, nodeCount, consumerCount);
} catch (Exception e) {
log.warn("Failed to record trace query history: {}",
e.getMessage());
}
}
+ private String clusterContext() {
+ return StringUtils.hasText(properties.getNamesrvAddr()) ?
properties.getNamesrvAddr() : null;
+ }
+
private static TraceRecordVO emptyTrace() {
return TraceRecordVO.builder()
.nodes(Collections.emptyList())
diff --git a/server/src/main/resources/application.yml
b/server/src/main/resources/application.yml
index 7295d3db..ddd36def 100644
--- a/server/src/main/resources/application.yml
+++ b/server/src/main/resources/application.yml
@@ -46,6 +46,9 @@ studio:
bearer-token: ${STUDIO_METRICS_PROMETHEUS_BEARER_TOKEN:}
rocketmq:
namesrv-addr: ${STUDIO_ROCKETMQ_NAMESRV_ADDR:}
+ query-history:
+ retention-days: ${STUDIO_QUERY_HISTORY_RETENTION_DAYS:90}
+ cleanup-interval: ${STUDIO_QUERY_HISTORY_CLEANUP_INTERVAL:PT24H}
llm:
token: ${RMQ_LLM_TOKEN:}
anthropic-base-url: ${RMQ_ANTHROPIC_BASE_URL:}
diff --git a/server/src/main/resources/db/schema.sql
b/server/src/main/resources/db/schema.sql
index a8952cf0..8475f8cd 100644
--- a/server/src/main/resources/db/schema.sql
+++ b/server/src/main/resources/db/schema.sql
@@ -93,6 +93,7 @@ CREATE TABLE IF NOT EXISTS rmq_message_query (
start_time BIGINT,
end_time BIGINT,
result_count INT DEFAULT 0,
+ cluster_id VARCHAR(255),
queried_by VARCHAR(64),
queried_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_queried_at (queried_at),
@@ -106,6 +107,7 @@ CREATE TABLE IF NOT EXISTS rmq_trace_query (
topic VARCHAR(255),
node_count INT DEFAULT 0,
consumer_count INT DEFAULT 0,
+ cluster_id VARCHAR(255),
queried_by VARCHAR(64),
queried_at DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX idx_msg_id (msg_id),
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryServiceTest.java
new file mode 100644
index 00000000..609d4aef
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryServiceTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.studio.queryhistory;
+
+import com.baomidou.mybatisplus.core.conditions.Wrapper;
+import org.apache.rocketmq.studio.auth.AuthenticatedUserContext;
+import org.apache.rocketmq.studio.persistence.entity.RmqMessageQuery;
+import org.apache.rocketmq.studio.persistence.entity.RmqTraceQuery;
+import org.apache.rocketmq.studio.persistence.mapper.RmqMessageQueryMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqTraceQueryMapper;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import java.time.Clock;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneOffset;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class QueryHistoryServiceTest {
+
+ private final RmqMessageQueryMapper messageQueryMapper =
mock(RmqMessageQueryMapper.class);
+ private final RmqTraceQueryMapper traceQueryMapper =
mock(RmqTraceQueryMapper.class);
+ private final QueryHistoryProperties properties = new
QueryHistoryProperties();
+ private final Clock clock =
Clock.fixed(Instant.parse("2026-08-05T12:00:00Z"), ZoneOffset.UTC);
+ private final QueryHistoryService service = new QueryHistoryService(
+ messageQueryMapper, traceQueryMapper, properties, clock);
+
+ @AfterEach
+ void clearUserContext() {
+ AuthenticatedUserContext.clear();
+ }
+
+ @Test
+ void recordsMessageQueryWithClusterAndAuthenticatedOperator() {
+ AuthenticatedUserContext.setUsername("alice");
+
+ service.recordMessageQuery("cluster-a", "TOPIC", "orders", null,
"tag-a", "key-a",
+ 1L, 2L, 3);
+
+ ArgumentCaptor<RmqMessageQuery> captor =
ArgumentCaptor.forClass(RmqMessageQuery.class);
+ verify(messageQueryMapper).insert(captor.capture());
+ RmqMessageQuery query = captor.getValue();
+ assertThat(query.getClusterId()).isEqualTo("cluster-a");
+ assertThat(query.getQueriedBy()).isEqualTo("alice");
+ assertThat(query.getQueriedAt()).isEqualTo(LocalDateTime.of(2026, 8,
5, 12, 0));
+ }
+
+ @Test
+ void recordsTraceQueryWithSystemOperatorWhenNoUserIsAuthenticated() {
+ service.recordTraceQuery("cluster-a", "msg-1", "orders", 2, 1);
+
+ ArgumentCaptor<RmqTraceQuery> captor =
ArgumentCaptor.forClass(RmqTraceQuery.class);
+ verify(traceQueryMapper).insert(captor.capture());
+ assertThat(captor.getValue().getClusterId()).isEqualTo("cluster-a");
+
assertThat(captor.getValue().getQueriedBy()).isEqualTo(AuthenticatedUserContext.SYSTEM_ACTOR);
+ }
+
+ @Test
+ void purgesBothQueryHistoriesUsingConfiguredRetention() {
+ properties.setRetentionDays(7);
+ when(messageQueryMapper.delete(any())).thenReturn(2);
+ when(traceQueryMapper.delete(any())).thenReturn(3);
+
+ service.purgeExpiredQueries();
+
+ ArgumentCaptor<Wrapper<RmqMessageQuery>> messageCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ ArgumentCaptor<Wrapper<RmqTraceQuery>> traceCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ verify(messageQueryMapper).delete(messageCaptor.capture());
+ verify(traceQueryMapper).delete(traceCaptor.capture());
+
assertThat(messageCaptor.getValue().getCustomSqlSegment()).contains("queried_at");
+
assertThat(traceCaptor.getValue().getCustomSqlSegment()).contains("queried_at");
+ }
+
+ @Test
+ void disablesCleanupWhenRetentionIsNonPositive() {
+ properties.setRetentionDays(0);
+
+ service.purgeExpiredQueries();
+
+ verify(messageQueryMapper, never()).delete(any());
+ verify(traceQueryMapper, never()).delete(any());
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
index e8b783b1..ac84429a 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
@@ -87,7 +87,8 @@ class RocketMQMessageProviderTest {
verify(consumer, never()).pull(any(MessageQueue.class),
anyString(), anyLong(), anyInt());
verify(consumer).shutdown();
}
- verify(queryHistoryService).recordMessageQuery("TOPIC", "TopicA",
null, null, null, 100L, 200L, 0);
+ verify(queryHistoryService).recordMessageQuery(null, "TOPIC",
"TopicA", null, null, null,
+ 100L, 200L, 0);
}
@Test