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 4b96f8321 fix(message): enumerate read queues when listing queue 
offsets (#3309)
4b96f8321 is described below

commit 4b96f8321dbcc69e847576c2f386184370e38ce3
Author: Zhao Jianing <[email protected]>
AuthorDate: Mon Sep 7 16:36:28 2026 +0800

    fix(message): enumerate read queues when listing queue offsets (#3309)
    
    getQueueOffsets iterated queueData.getWriteQueueNums() while every other
    browse path is read-queue based: the broker's PullMessageProcessor
    rejects queueId >= readQueueNums with SYSTEM_ERROR, and
    fetchSubscribeMessageQueues (used by queryByTopic and the classic
    console) enumerates [0, readQueueNums). When read != write:
    
    - read > write (shrink draining): queues in [write, read) still hold
      browsable messages but were missing from the QueueBrowser
    - write > read (queues not yet readable): queues in [read, write) were
      listed but every pull on them fails with "queueId is illegal"
    
    Enumerate from readQueueNums to match the broker's pull validation.
---
 .../provider/apache/RocketMQMessageProvider.java   |  2 +-
 .../apache/RocketMQMessageProviderTest.java        | 52 ++++++++++++++++++++++
 2 files changed, 53 insertions(+), 1 deletion(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index 1c112bfa6..2191c6ea4 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -204,7 +204,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                     return Collections.emptyList();
                 }
                 for (QueueData queueData : route.getQueueDatas()) {
-                    for (int queueId = 0; queueId < 
queueData.getWriteQueueNums(); queueId++) {
+                    for (int queueId = 0; queueId < 
queueData.getReadQueueNums(); queueId++) {
                         MessageQueue queue = new MessageQueue(topic, 
queueData.getBrokerName(), queueId);
                         result.add(QueueOffsetVO.builder()
                                 .brokerName(queue.getBrokerName())
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 4a4949d8e..e38d1524e 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -31,6 +31,8 @@ import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.body.ConsumeMessageDirectlyResult;
 import org.apache.rocketmq.remoting.protocol.body.CMResult;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.remoting.protocol.route.QueueData;
+import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
 import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
@@ -40,6 +42,7 @@ import 
org.apache.rocketmq.studio.instance.message.MessageRecordVO;
 import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
 import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
 import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.instance.message.QueueOffsetVO;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExtImpl;
 import org.junit.jupiter.api.BeforeEach;
@@ -776,6 +779,44 @@ class RocketMQMessageProviderTest {
         assertThat(pulledOffsets).allMatch(offset -> offset >= 
expectedFirstOffset);
     }
 
+    @Test
+    void getQueueOffsetsListsReadQueuesWhenReadCountExceedsWriteCount() throws 
Exception {
+        // Shrinking first lowers writeQueueNums while reads keep draining the 
tail
+        // queues, so queues in [writeQueueNums, readQueueNums) still hold 
browsable
+        // messages and must appear in the browser 
(fetchSubscribeMessageQueues, used
+        // by queryByTopic, enumerates the same read queues).
+        when(adminExt.examineTopicRouteInfo("TopicA"))
+                .thenReturn(routeWithQueueCounts("broker-a", 2, 4));
+        
when(adminExt.minOffset(any(MessageQueue.class))).thenAnswer(invocation ->
+                invocation.<MessageQueue>getArgument(0).getQueueId() * 10L);
+        
when(adminExt.maxOffset(any(MessageQueue.class))).thenAnswer(invocation ->
+                invocation.<MessageQueue>getArgument(0).getQueueId() * 10L + 
5L);
+
+        List<QueueOffsetVO> offsets = provider.getQueueOffsets("instance-a", 
"TopicA");
+
+        assertThat(offsets).extracting(QueueOffsetVO::getBrokerName)
+                .containsOnly("broker-a");
+        assertThat(offsets).extracting(QueueOffsetVO::getQueueId)
+                .containsExactly(0, 1, 2, 3);
+        assertThat(offsets).extracting(QueueOffsetVO::getMinOffset)
+                .containsExactly(0L, 10L, 20L, 30L);
+        assertThat(offsets).extracting(QueueOffsetVO::getMaxOffset)
+                .containsExactly(5L, 15L, 25L, 35L);
+    }
+
+    @Test
+    void getQueueOffsetsSkipsWriteOnlyQueuesWhenWriteCountExceedsReadCount() 
throws Exception {
+        // The broker's PullMessageProcessor rejects queueId >= readQueueNums 
with
+        // SYSTEM_ERROR, so write-only queues cannot be browsed and must not 
be listed.
+        when(adminExt.examineTopicRouteInfo("TopicA"))
+                .thenReturn(routeWithQueueCounts("broker-a", 4, 2));
+
+        List<QueueOffsetVO> offsets = provider.getQueueOffsets("instance-a", 
"TopicA");
+
+        assertThat(offsets).extracting(QueueOffsetVO::getQueueId)
+                .containsExactly(0, 1);
+    }
+
     private MQClientAPIImpl mockOffsetLookupClient() {
         DefaultMQAdminExtImpl adminExtImpl = mock(DefaultMQAdminExtImpl.class);
         MQClientInstance clientInstance = mock(MQClientInstance.class);
@@ -787,6 +828,17 @@ class RocketMQMessageProviderTest {
     }
 
 
+    private static TopicRouteData routeWithQueueCounts(
+            String brokerName, int writeQueueNums, int readQueueNums) {
+        QueueData queueData = new QueueData();
+        queueData.setBrokerName(brokerName);
+        queueData.setWriteQueueNums(writeQueueNums);
+        queueData.setReadQueueNums(readQueueNums);
+        TopicRouteData route = new TopicRouteData();
+        route.setQueueDatas(List.of(queueData));
+        return route;
+    }
+
     private static ClusterInfo clusterInfoWithBrokerAddresses(String... 
brokerAddresses) {
         ClusterInfo clusterInfo = new ClusterInfo();
         Map<String, BrokerData> brokerAddrTable = new HashMap<>();

Reply via email to