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

Reply via email to