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<>();