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 22c26d258 fix: DLQ contract, topic filtering, catalog paging and
runtime robustness (#2365)
22c26d258 is described below
commit 22c26d258dae68ae517abd0a4878e63c8c8f41f7
Author: aias00 <[email protected]>
AuthorDate: Wed Aug 19 17:32:55 2026 +0800
fix: DLQ contract, topic filtering, catalog paging and runtime robustness
(#2365)
* fix: skip null brokers in Apache runtime topology
* fix: scope query history to authenticated subject
* fix: complete Aliyun cloud instance catalogs
* fix: align paged Topic filtering with canonical system topics
Signed-off-by: liuhy <[email protected]>
* Keep DLQ search and pagination behavior consistent
Constraint: limit the fix to DLQ search semantics, pagination reset, and
documented API contract
Rejected: backend-only search fix | frontend could still strand users on
stale later pages
Confidence: high
Scope-risk: narrow
Tested:
JAVA_HOME=/Users/aias/Library/Java/JavaVirtualMachines/openjdk-21.0.2/Contents/Home
PATH="/Library/Java/JavaVirtualMachines/jdk-17.0.3+7/Contents/Home/bin:/Users/aias/.qoder/entry:/Users/aias/.nvm/versions/node/v22.23.1/lib/node_modules/@openai/codex/node_modules/@openai/codex-darwin-arm64/vendor/aarch64-apple-darwin/codex-path:/Users/aias/.strix/bin:/Users/aias/.grok/bin:/Users/aias/.opencode/bin:/Users/aias/.antigravity/antigravity/bin:/Users/aias/.local/bin:/Users/aias/.bun
[...]
Tested: cd web && npm run test -- src/api/dlq.test.ts
src/pages/instance/__tests__/DLQPage.test.tsx
Tested: cd web && npm run build
Signed-off-by: liuhy <[email protected]>
* fix: commit DLQ rows and cleared loading state in one render
Cherry-picking the DLQ pagination-reset test from #2365 exposed a render
race: rows became visible while the spin overlay (loading) still blocked
pointer events, because setGroups landed in .then and setLoading(false)
in a separate .finally render. Fold the loading flag into the data
callbacks so both commit atomically.
Signed-off-by: lizhimins <[email protected]>
---------
Signed-off-by: liuhy <[email protected]>
Signed-off-by: lizhimins <[email protected]>
Co-authored-by: lizhimins <[email protected]>
---
docs/api-spec.md | 37 ++++++++++---
.../studio/common/util/SystemTopicFilter.java | 17 ++----
.../instance/message/QueryHistoryService.java | 13 ++++-
.../provider/alibaba/AliyunCatalogService.java | 10 +++-
.../provider/apache/RocketMQAdminClientImpl.java | 2 +-
.../provider/apache/RocketMQDLQProvider.java | 6 ++-
.../provider/apache/RocketMQMetadataProvider.java | 5 +-
.../studio/common/util/SystemTopicFilterTest.java | 9 ++--
.../studio/instance/dlq/DLQControllerTest.java | 17 ++++++
.../instance/message/QueryHistoryServiceTest.java | 13 +++++
.../provider/alibaba/AliyunCatalogServiceTest.java | 23 +++++---
.../apache/RocketMQAdminClientImplTest.java | 24 +++++++++
.../provider/apache/RocketMQDLQProviderTest.java | 16 ++++++
web/src/api/dlq.test.ts | 30 +++++++++--
web/src/pages/instance/__tests__/DLQPage.test.tsx | 63 ++++++++++++++++++----
web/src/pages/instance/dlq.tsx | 17 ++++--
16 files changed, 245 insertions(+), 57 deletions(-)
diff --git a/docs/api-spec.md b/docs/api-spec.md
index 7c7051f9a..12a0882a2 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -1377,10 +1377,24 @@ GET /api/messages/:msgId/trace
### 9.1 获取 DLQ 列表
```
-GET /api/dlq?clusterId={clusterId}
+GET
/api/dlq?instanceId={instanceId}&search={keyword}&page={page}&pageSize={pageSize}
```
-**Response `data`:** `DLQGroup[]`
+| 参数 | 类型 | 必填 | 说明 |
+|------|------|------|------|
+| `instanceId` | `string` | 是 | 实例 ID(全局唯一字符串) |
+| `search` | `string` | 否 | 按 Group 名称或 DLQ Topic 搜索,大小写不敏感 |
+| `page` | `number` | 否 | 页码,默认 `1` |
+| `pageSize` | `number` | 否 | 每页条数,默认 `20`,最大 `100` |
+
+**Response `data`:** `PageResult<DLQGroup>`
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `items` | `DLQGroup[]` | 当前页数据 |
+| `total` | `number` | 匹配总数 |
+| `page` | `number` | 当前页码 |
+| `size` | `number` | 当前页大小 |
| 字段 | 类型 | 说明 |
|------|------|------|
@@ -1389,7 +1403,8 @@ GET /api/dlq?clusterId={clusterId}
| `messageCount` | `number` | 死信消息数量 |
| `lastEnqueueTime` | `string` | 最后入队时间 (ISO 8601) |
| `retryCount` | `number` | 已重试次数 |
-| `status` | `string` | 状态: `active` / `empty` |
+| `status` | `string` | 状态: `ACTIVE` / `EMPTY` / `UNAVAILABLE` |
+| `statsAvailable` | `boolean` | 是否成功读取 DLQ Topic 统计信息 |
### 9.2 重发死信消息
@@ -1401,12 +1416,22 @@ POST /api/dlq/resend
| 字段 | 类型 | 必填 | 说明 |
|------|------|------|------|
+| `instanceId` | `string` | 是 | 实例 ID(全局唯一字符串) |
| `groupName` | `string` | 是 | 消费组名称 |
-| `startTime` | `string` | 是 | 重投时间范围起始 (ISO 8601) |
-| `endTime` | `string` | 是 | 重投时间范围结束 (ISO 8601) |
+| `startTime` | `number` | 否 | 重投时间范围起始(Unix 毫秒时间戳) |
+| `endTime` | `number` | 否 | 重投时间范围结束(Unix 毫秒时间戳) |
| `targetTopic` | `string` | 否 | 目标 Topic,不传则重投回原 Topic |
-**Response `data`:** `null`
+**Response `data`:** `DLQResendResult`
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `matched` | `number` | 命中的死信消息数 |
+| `resent` | `number` | 成功重投条数 |
+| `failed` | `number` | 重投失败条数 |
+| `outcome` | `string` | 结果: `SUCCESS` / `PARTIAL` / `FAILED` / `NO_MESSAGES` |
+| `scanIncomplete` | `boolean` | 是否有部分队列扫描失败 |
+| `failedQueueCount` | `number` | 扫描失败的队列数 |
---
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/util/SystemTopicFilter.java
b/server/src/main/java/org/apache/rocketmq/studio/common/util/SystemTopicFilter.java
index ae21e1443..db2458ec1 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/util/SystemTopicFilter.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/util/SystemTopicFilter.java
@@ -16,6 +16,8 @@
*/
package org.apache.rocketmq.studio.common.util;
+import org.apache.rocketmq.common.topic.TopicValidator;
+
import java.util.Set;
/**
@@ -27,16 +29,7 @@ import java.util.Set;
*/
public final class SystemTopicFilter {
- private static final Set<String> SYSTEM_TOPIC_PREFIXES = Set.of(
- "RMQ_SYS_", "rmq_sys_", "SCHEDULE_TOPIC_", "%RETRY%", "%DLQ%",
- "CID_", "broker_", "BenchmarkTest"
- );
-
- private static final Set<String> SYSTEM_TOPICS = Set.of(
- "TBW102", "SELF_TEST_TOPIC", "DefaultCluster",
"OFFSET_MOVED_EVENT",
- "broker", "SCHEDULE_TOPIC_XXXX", "RMQ_SYS_TRANS_HALF_TOPIC",
- "RMQ_SYS_TRACE_TOPIC", "RMQ_SYS_TRANS_OP_HALF_TOPIC"
- );
+ private static final Set<String> RETRY_AND_DLQ_PREFIXES =
Set.of("%RETRY%", "%DLQ%");
private SystemTopicFilter() {
}
@@ -54,10 +47,10 @@ public final class SystemTopicFilter {
if (topicName == null || topicName.isEmpty()) {
return true;
}
- if (SYSTEM_TOPICS.contains(topicName)) {
+ if (TopicValidator.isSystemTopic(topicName)) {
return true;
}
- for (String prefix : SYSTEM_TOPIC_PREFIXES) {
+ for (String prefix : RETRY_AND_DLQ_PREFIXES) {
if (topicName.startsWith(prefix)) {
return true;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryService.java
index b97a8f107..0bb21e87f 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/QueryHistoryService.java
@@ -100,8 +100,10 @@ public class QueryHistoryService {
public PageResult<MessageQueryHistoryVO> listMessageQueries(String
clusterId, String queryType,
String
search, int page, int pageSize) {
String pattern = escapeLike(search);
+ String queriedBy = AuthenticatedUserContext.currentUsernameOrSystem();
QueryWrapper<RmqMessageQuery> query = new
QueryWrapper<RmqMessageQuery>()
.eq(StringUtils.hasText(clusterId), "cluster_id", clusterId)
+ .eq("queried_by", queriedBy)
.eq(StringUtils.hasText(queryType), "query_type", queryType)
.and(StringUtils.hasText(search), nested -> nested
.like("topic", pattern)
@@ -118,8 +120,10 @@ public class QueryHistoryService {
public PageResult<TraceQueryHistoryVO> listTraceQueries(String clusterId,
String search,
int page, int
pageSize) {
String pattern = escapeLike(search);
+ String queriedBy = AuthenticatedUserContext.currentUsernameOrSystem();
QueryWrapper<RmqTraceQuery> query = new QueryWrapper<RmqTraceQuery>()
.eq(StringUtils.hasText(clusterId), "cluster_id", clusterId)
+ .eq("queried_by", queriedBy)
.and(StringUtils.hasText(search), nested -> nested
.like("topic", pattern)
.or().like("msg_id", pattern)
@@ -132,19 +136,24 @@ public class QueryHistoryService {
}
public QueryHistorySummaryVO summarize(String clusterId) {
+ String queriedBy = AuthenticatedUserContext.currentUsernameOrSystem();
QueryWrapper<RmqMessageQuery> messageFilter = new
QueryWrapper<RmqMessageQuery>()
- .eq(StringUtils.hasText(clusterId), "cluster_id", clusterId);
+ .eq(StringUtils.hasText(clusterId), "cluster_id", clusterId)
+ .eq("queried_by", queriedBy);
QueryWrapper<RmqTraceQuery> traceFilter = new
QueryWrapper<RmqTraceQuery>()
- .eq(StringUtils.hasText(clusterId), "cluster_id", clusterId);
+ .eq(StringUtils.hasText(clusterId), "cluster_id", clusterId)
+ .eq("queried_by", queriedBy);
long messageCount = messageQueryMapper.selectCount(messageFilter);
long traceCount = traceQueryMapper.selectCount(traceFilter);
RmqMessageQuery latestMessage = messageQueryMapper.selectOne(
new QueryWrapper<RmqMessageQuery>()
.eq(StringUtils.hasText(clusterId), "cluster_id",
clusterId)
+ .eq("queried_by", queriedBy)
.orderByDesc("gmt_create", "id").last("LIMIT 1"));
RmqTraceQuery latestTrace = traceQueryMapper.selectOne(
new QueryWrapper<RmqTraceQuery>()
.eq(StringUtils.hasText(clusterId), "cluster_id",
clusterId)
+ .eq("queried_by", queriedBy)
.orderByDesc("gmt_create", "id").last("LIMIT 1"));
LocalDateTime latest = latestOf(
latestMessage == null ? null : latestMessage.getGmtCreate(),
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
index f6f957938..cc011df88 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
@@ -115,7 +115,7 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
private List<ListInstancesResponseBody.List> fetchAllInstances(Long
credentialId, String regionId) {
List<ListInstancesResponseBody.List> all = new ArrayList<>();
- for (int page = 1; page <= AliyunConverters.MAX_PAGES; page++) {
+ for (int page = 1; ; page++) {
ListInstancesRequest request = ListInstancesRequest.builder()
.pageNumber(page)
.pageSize(AliyunConverters.PAGE_SIZE)
@@ -129,7 +129,13 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
break;
}
all.addAll(list);
- if (list.size() < AliyunConverters.PAGE_SIZE) {
+ Long totalCount = data.getTotalCount();
+ if (totalCount != null && totalCount < 0) {
+ throw new BusinessException(502, "Aliyun instance catalog
returned a negative totalCount");
+ }
+ if (list.size() < AliyunConverters.PAGE_SIZE
+ || totalCount != null && all.size() >= totalCount
+ || totalCount == null && page >=
AliyunConverters.MAX_PAGES) {
break;
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
index a16e8b849..895bcd0f3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
@@ -630,7 +630,7 @@ public class RocketMQAdminClientImpl implements AdminClient
{
}
for (BrokerData brokerData :
clusterInfo.getBrokerAddrTable().values()) {
- if (brokerData.getBrokerAddrs() == null) {
+ if (brokerData == null || brokerData.getBrokerAddrs() == null) {
continue;
}
// Use master address (brokerId = 0) preferentially
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
index 7415d94dc..7c4aca641 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
@@ -55,6 +55,7 @@ import java.util.ArrayList;
import java.util.Base64;
import java.util.Collections;
import java.util.Comparator;
+import java.util.Locale;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -93,6 +94,7 @@ public class RocketMQDLQProvider implements DLQProvider {
throws Exception {
TopicList topicList = adminExt.fetchAllTopicList();
Set<String> topics = topicList == null ? Collections.emptySet() :
topicList.getTopicList();
+ String normalizedSearch = StringUtils.hasText(search) ?
search.trim().toLowerCase(Locale.ROOT) : null;
List<String> dlqTopics = new ArrayList<>();
for (String topic : topics) {
@@ -103,7 +105,9 @@ public class RocketMQDLQProvider implements DLQProvider {
if (!StringUtils.hasText(groupName)) {
continue;
}
- if (search == null || groupName.contains(search) ||
topic.contains(search)) {
+ if (normalizedSearch == null
+ ||
groupName.toLowerCase(Locale.ROOT).contains(normalizedSearch)
+ ||
topic.toLowerCase(Locale.ROOT).contains(normalizedSearch)) {
dlqTopics.add(topic);
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 48bfdc935..d4e93d7e3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -42,6 +42,7 @@ import
org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.common.util.SystemGroupFilter;
import org.apache.rocketmq.studio.common.util.SystemTopicFilter;
+import org.apache.rocketmq.common.topic.TopicValidator;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
@@ -146,8 +147,10 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
.eq(StringUtils.hasText(clusterId), RmqTopic::getClusterId,
clusterId)
.eq(StringUtils.hasText(type), RmqTopic::getTopicType, type)
.like(StringUtils.hasText(search), RmqTopic::getName, search)
- .notLikeRight(RmqTopic::getName, "RMQ_SYS_")
+ .notIn(RmqTopic::getName, TopicValidator.getSystemTopicSet())
.notLikeRight(RmqTopic::getName, "rmq_sys_")
+ .notLikeRight(RmqTopic::getName, "%RETRY%")
+ .notLikeRight(RmqTopic::getName, "%DLQ%")
.orderByAsc(RmqTopic::getName, RmqTopic::getId);
Page<RmqTopic> result = topicMapper.selectPage(new Page<>(page,
pageSize), query);
return
PageResult.of(result.getRecords().stream().map(this::toTopicVO).toList(),
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/common/util/SystemTopicFilterTest.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/SystemTopicFilterTest.java
index a3e0f5fa1..4a03c2fd9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/common/util/SystemTopicFilterTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/SystemTopicFilterTest.java
@@ -25,7 +25,7 @@ import static org.assertj.core.api.Assertions.assertThat;
class SystemTopicFilterTest {
@Test
- void shouldRecognizeSharedSystemTopicPrefixesAndBrokerNamesTest() {
+ void shouldRecognizeCanonicalSystemTopicsAndBrokerNamesTest() {
Set<String> brokerNames = Set.of("broker-prod-a", "broker-prod-b");
assertThat(SystemTopicFilter.isSystem(null, brokerNames)).isTrue();
@@ -35,14 +35,15 @@ class SystemTopicFilterTest {
assertThat(SystemTopicFilter.isSystem("SCHEDULE_TOPIC_XXXX",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("%RETRY%consumer-a",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("%DLQ%consumer-a",
brokerNames)).isTrue();
- assertThat(SystemTopicFilter.isSystem("CID_RMQ_SYS_TRANS",
brokerNames)).isTrue();
- assertThat(SystemTopicFilter.isSystem("broker_config",
brokerNames)).isTrue();
- assertThat(SystemTopicFilter.isSystem("BenchmarkTestTopic",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("TBW102", brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("SELF_TEST_TOPIC",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("OFFSET_MOVED_EVENT",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("broker-prod-a",
brokerNames)).isTrue();
assertThat(SystemTopicFilter.isSystem("orders",
brokerNames)).isFalse();
+ assertThat(SystemTopicFilter.isSystem("CID_orders",
brokerNames)).isFalse();
+ assertThat(SystemTopicFilter.isSystem("broker_events",
brokerNames)).isFalse();
+ assertThat(SystemTopicFilter.isSystem("BenchmarkTestOrders",
brokerNames)).isFalse();
+ assertThat(SystemTopicFilter.isSystem("SCHEDULE_TOPIC_orders",
brokerNames)).isFalse();
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
index c87cfd56d..caf8c5ab1 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
@@ -94,6 +94,23 @@ class DLQControllerTest {
verify(dlqService).listDLQGroups(eq("instance-1"), isNull(), eq(1),
eq(20));
}
+ @Test
+ void listDLQGroupsShouldPassSearchAndPaging() throws Exception {
+ when(dlqService.listDLQGroups(eq("instance-1"), eq("order"), eq(2),
eq(50)))
+ .thenReturn(PageResult.empty(2, 50));
+
+ mockMvc.perform(get("/api/dlq")
+ .param("instanceId", "instance-1")
+ .param("search", "order")
+ .param("page", "2")
+ .param("pageSize", "50"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.page").value(2))
+ .andExpect(jsonPath("$.data.size").value(50));
+
+ verify(dlqService).listDLQGroups(eq("instance-1"), eq("order"), eq(2),
eq(50));
+ }
+
@Test
void resendMessagesShouldReturnSuccess() throws Exception {
Map<String, Object> body = Map.of(
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryServiceTest.java
index 940e27c64..4260c2702 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/QueryHistoryServiceTest.java
@@ -107,6 +107,7 @@ class QueryHistoryServiceTest {
@Test
void listsMessageHistoryAsNewestFirstPage() {
+ AuthenticatedUserContext.setUsername("alice");
RmqMessageQuery entity = new RmqMessageQuery();
entity.setId(9L);
entity.setQueryType("KEY");
@@ -132,10 +133,15 @@ class QueryHistoryServiceTest {
assertThat(item.getMessageKey()).isEqualTo("order-1");
assertThat(item.getQueriedBy()).isEqualTo("alice");
});
+
+ ArgumentCaptor<Wrapper<RmqMessageQuery>> queryCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ verify(messageQueryMapper).selectPage(any(Page.class),
queryCaptor.capture());
+
assertThat(queryCaptor.getValue().getCustomSqlSegment()).contains("queried_by");
}
@Test
void summarizesBothHistoryStreams() {
+ AuthenticatedUserContext.setUsername("alice");
RmqMessageQuery message = new RmqMessageQuery();
message.setGmtCreate(LocalDateTime.of(2026, 8, 5, 10, 0));
RmqTraceQuery trace = new RmqTraceQuery();
@@ -150,5 +156,12 @@ class QueryHistoryServiceTest {
assertThat(summary.getMessageQueries()).isEqualTo(7);
assertThat(summary.getTraceQueries()).isEqualTo(4);
assertThat(summary.getLatestQueryAt()).isEqualTo(LocalDateTime.of(2026, 8, 5,
12, 0));
+
+ ArgumentCaptor<Wrapper<RmqMessageQuery>> messageCountCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ ArgumentCaptor<Wrapper<RmqTraceQuery>> traceCountCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ verify(messageQueryMapper).selectCount(messageCountCaptor.capture());
+ verify(traceQueryMapper).selectCount(traceCountCaptor.capture());
+
assertThat(messageCountCaptor.getValue().getCustomSqlSegment()).contains("queried_by");
+
assertThat(traceCountCaptor.getValue().getCustomSqlSegment()).contains("queried_by");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
index 74a61f36c..2936e2ebb 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
@@ -122,7 +122,8 @@ class AliyunCatalogServiceTest {
@Test
void listCloudInstancesShouldAggregatePagesTest() {
- ListInstancesResponse firstPage =
instancesResponse(instanceRows(AliyunConverters.PAGE_SIZE, 0));
+ ListInstancesResponse firstPage = instancesResponse(
+ instanceRows(AliyunConverters.PAGE_SIZE, 0),
AliyunConverters.PAGE_SIZE + 3L);
ListInstancesResponse secondPage = instancesResponse(instanceRows(3,
AliyunConverters.PAGE_SIZE));
when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION), any()))
.thenReturn(firstPage, secondPage);
@@ -134,14 +135,18 @@ class AliyunCatalogServiceTest {
}
@Test
- void listCloudInstancesShouldStopAtMaxPagesTest() {
- ListInstancesResponse fullPage =
instancesResponse(instanceRows(AliyunConverters.PAGE_SIZE, 0));
- when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION),
any())).thenReturn(fullPage);
+ void
listCloudInstancesShouldContinuePastFallbackPageCapWhenTotalCountRequiresIt() {
+ ListInstancesResponse fullPage = instancesResponse(
+ instanceRows(AliyunConverters.PAGE_SIZE, 0),
AliyunConverters.PAGE_SIZE * 6L + 1);
+ ListInstancesResponse finalPage = instancesResponse(instanceRows(1,
AliyunConverters.PAGE_SIZE * 6),
+ AliyunConverters.PAGE_SIZE * 6L + 1);
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION), any()))
+ .thenReturn(fullPage, fullPage, fullPage, fullPage, fullPage,
fullPage, finalPage);
List<CloudInstanceOptionVO> options =
service.listCloudInstances(CREDENTIAL_ID, REGION, null);
- assertThat(options).hasSize(AliyunConverters.PAGE_SIZE *
AliyunConverters.MAX_PAGES);
- verify(clientFactory,
times(AliyunConverters.MAX_PAGES)).call(eq(CREDENTIAL_ID), eq(REGION), any());
+ assertThat(options).hasSize(AliyunConverters.PAGE_SIZE * 6 + 1);
+ verify(clientFactory, times(7)).call(eq(CREDENTIAL_ID), eq(REGION),
any());
}
@Test
@@ -228,6 +233,10 @@ class AliyunCatalogServiceTest {
}
private static ListInstancesResponse
instancesResponse(List<ListInstancesResponseBody.List> rows) {
+ return instancesResponse(rows, (long) rows.size());
+ }
+
+ private static ListInstancesResponse
instancesResponse(List<ListInstancesResponseBody.List> rows, long totalCount) {
return ListInstancesResponse.create().toBuilder()
.statusCode(200)
.body(ListInstancesResponseBody.builder()
@@ -235,7 +244,7 @@ class AliyunCatalogServiceTest {
.list(rows)
.pageNumber(1L)
.pageSize((long) AliyunConverters.PAGE_SIZE)
- .totalCount((long) rows.size())
+ .totalCount(totalCount)
.build())
.build())
.build();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
index ccfcf8e2c..26f31deb5 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
@@ -247,6 +247,30 @@ class RocketMQAdminClientImplTest {
org.mockito.ArgumentMatchers.eq("10.0.1.1:10911"),
any(TopicConfig.class));
}
+ @Test
+ void createTopicSkipsNullBrokerDataWhenFallingBackToAllBrokers() throws
Exception {
+ TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new
MybatisConfiguration(), ""), RmqTopic.class);
+ ClusterInfo clusterInfo = new ClusterInfo();
+ Map<String, BrokerData> brokerAddrTable = new HashMap<>();
+ brokerAddrTable.put("missing-broker", null);
+ BrokerData brokerData = new BrokerData();
+ brokerData.setBrokerName("broker-1");
+ brokerData.setBrokerAddrs(new HashMap<>(Map.of(0L, "10.0.0.1:10911")));
+ brokerAddrTable.put("broker-1", brokerData);
+ clusterInfo.setBrokerAddrTable(brokerAddrTable);
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(topicMapper.selectOne(any())).thenReturn(null);
+ doNothing().when(adminExt).createAndUpdateTopicConfig(anyString(),
any(TopicConfig.class));
+
+ TopicVO topic = new TopicVO();
+ topic.setName("orders");
+
+ adminClient.createTopic(topic);
+
+ verify(adminExt).createAndUpdateTopicConfig(
+ org.mockito.ArgumentMatchers.eq("10.0.0.1:10911"),
any(TopicConfig.class));
+ }
+
@Test
void topicWritesUseSelectedInstanceAdmin() throws Exception {
TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new
MybatisConfiguration(), ""), RmqTopic.class);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
index 383b6830f..607d1bddd 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
@@ -162,6 +162,22 @@ class RocketMQDLQProviderTest {
.containsExactly("order-b");
}
+ @Test
+ void listDLQGroupsShouldFilterCaseInsensitivelyTest() throws Exception {
+ TopicList topicList = new TopicList();
+ topicList.setTopicList(Set.of(
+ MixAll.DLQ_GROUP_TOPIC_PREFIX + "Order-Consumer",
+ MixAll.DLQ_GROUP_TOPIC_PREFIX + "payment-consumer"));
+ when(adminExt.fetchAllTopicList()).thenReturn(topicList);
+ when(adminExt.examineTopicStats(anyString())).thenReturn(new
TopicStatsTable());
+
+ PageResult<DLQGroupVO> filtered = provider.listDLQGroups("instance-a",
"order", 1, 20);
+
+ assertThat(filtered.getTotal()).isEqualTo(1);
+ assertThat(filtered.getItems()).extracting(DLQGroupVO::getGroupName)
+ .containsExactly("Order-Consumer");
+ }
+
@Test
void resendMessagesShouldRejectInvertedTimeRangeBeforeCreatingConsumers() {
diff --git a/web/src/api/dlq.test.ts b/web/src/api/dlq.test.ts
index 796063a5a..343df8de9 100644
--- a/web/src/api/dlq.test.ts
+++ b/web/src/api/dlq.test.ts
@@ -44,15 +44,35 @@ describe('DLQ API', () => {
it('loads and unwraps DLQ groups', async () => {
const pageData = { items: [group], total: 1, page: 1, size: 20 };
- mock.onGet('/dlq').reply((config) =>
- config.params?.instanceId === 'instance-1'
- ? [200, { code: 200, data: pageData }]
- : [404, {}],
- );
+ mock
+ .onGet('/dlq')
+ .reply((config) =>
+ config.params?.instanceId === 'instance-1' &&
+ config.params?.page === 1 &&
+ config.params?.pageSize === 20
+ ? [200, { code: 200, data: pageData }]
+ : [404, {}],
+ );
await expect(listDLQGroups('instance-1')).resolves.toEqual(pageData);
});
+ it('passes paged search parameters through to the backend contract', async
() => {
+ const pageData = { items: [group], total: 1, page: 2, size: 50 };
+ mock
+ .onGet('/dlq')
+ .reply((config) =>
+ config.params?.instanceId === 'instance-1' &&
+ config.params?.search === 'order' &&
+ config.params?.page === 2 &&
+ config.params?.pageSize === 50
+ ? [200, { code: 200, data: pageData }]
+ : [404, {}],
+ );
+
+ await expect(listDLQGroups('instance-1', 'order', 2,
50)).resolves.toEqual(pageData);
+ });
+
it('sends epoch milliseconds for the resend time range', async () => {
const payload = {
instanceId: 'instance-1',
diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx
b/web/src/pages/instance/__tests__/DLQPage.test.tsx
index a6a247513..50ef82460 100644
--- a/web/src/pages/instance/__tests__/DLQPage.test.tsx
+++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx
@@ -121,6 +121,9 @@ describe('DLQ page', () => {
});
beforeEach(() => {
+ vi.mocked(messageService.listDLQGroups).mockReset();
+ vi.mocked(messageService.resendDLQ).mockReset();
+ vi.mocked(messageService.exportDLQMessages).mockReset();
createObjectURL = vi.fn().mockReturnValue('blob:dlq');
revokeObjectURL = vi.fn();
Object.defineProperty(URL, 'createObjectURL', {
@@ -148,6 +151,43 @@ describe('DLQ page', () => {
expect(messageService.listDLQGroups).toHaveBeenCalledWith('instance-1',
undefined, 1, 20);
});
+ it('resets pagination to the first page when the search term changes', async
() => {
+ vi.mocked(messageService.listDLQGroups)
+ .mockResolvedValueOnce({
+ items: [dlqGroup],
+ total: 40,
+ page: 1,
+ size: 20,
+ })
+ .mockResolvedValueOnce({
+ items: [dlqGroup],
+ total: 40,
+ page: 2,
+ size: 20,
+ })
+ .mockResolvedValueOnce({
+ items: [dlqGroup],
+ total: 1,
+ page: 1,
+ size: 20,
+ });
+ const user = userEvent.setup();
+ renderWithProviders(<DLQPage />);
+
+ await screen.findByText('cg-order');
+ await user.click(screen.getByTitle('2'));
+ await waitFor(() =>
+
expect(messageService.listDLQGroups).toHaveBeenLastCalledWith('instance-1',
undefined, 2, 20),
+ );
+
+ const searchInput = screen.getByPlaceholderText('搜索 Group 名称或 DLQ Topic');
+ await user.type(searchInput, 'ord');
+
+ await waitFor(() =>
+
expect(messageService.listDLQGroups).toHaveBeenLastCalledWith('instance-1',
'ord', 1, 20),
+ );
+ });
+
it('surfaces unavailable DLQ provider errors when loading groups', async ()
=> {
vi.mocked(messageService.listDLQGroups).mockRejectedValue(
new Error('DLQ provider is not configured'),
@@ -171,10 +211,9 @@ describe('DLQ page', () => {
});
it('sorts DLQ rows with missing enqueue timestamps', async () => {
- vi.mocked(messageService.listDLQGroups).mockResolvedValue(pageOf([
- dlqGroup,
- { ...secondDlqGroup, lastEnqueueTime: null },
- ]));
+ vi.mocked(messageService.listDLQGroups).mockResolvedValue(
+ pageOf([dlqGroup, { ...secondDlqGroup, lastEnqueueTime: null }]),
+ );
const user = userEvent.setup();
renderWithProviders(<DLQPage />);
@@ -250,13 +289,15 @@ describe('DLQ page', () => {
});
it('neutralizes formulas hidden behind a leading line feed in CSV summary
exports', async () => {
- vi.mocked(messageService.listDLQGroups).mockResolvedValue(pageOf([
- {
- ...dlqGroup,
- groupName: '\n=1+1',
- dlqTopic: '%DLQ%formula',
- },
- ]));
+ vi.mocked(messageService.listDLQGroups).mockResolvedValue(
+ pageOf([
+ {
+ ...dlqGroup,
+ groupName: '\n=1+1',
+ dlqTopic: '%DLQ%formula',
+ },
+ ]),
+ );
const user = userEvent.setup();
renderWithProviders(<DLQPage />);
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index b670b36b7..7bb1efbf2 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -167,6 +167,9 @@ const DLQPage = () => {
};
}
+ // Clear `loading` inside the same callback as the data updates so rows and
+ // the cleared spinner commit in one batched render — otherwise rows can be
+ // visible for a render while the spin overlay still blocks pointer events.
void listDLQGroups(selectedInstanceId, search || undefined, page, pageSize)
.then((result) => {
if (!cancelled) {
@@ -179,13 +182,14 @@ const DLQPage = () => {
setSelectedGroupNames((selected) =>
selected.filter((groupName) => availableGroups.has(groupName)),
);
+ setLoading(false);
}
})
.catch((error) => {
- if (!cancelled) setLoadError(getErrorMessage(error,
DEFAULT_LOAD_ERROR));
- })
- .finally(() => {
- if (!cancelled) setLoading(false);
+ if (!cancelled) {
+ setLoadError(getErrorMessage(error, DEFAULT_LOAD_ERROR));
+ setLoading(false);
+ }
});
return () => {
@@ -402,7 +406,10 @@ const DLQPage = () => {
placeholder="搜索 Group 名称或 DLQ Topic"
allowClear
value={search}
- onChange={(e) => setSearch(e.target.value)}
+ onChange={(e) => {
+ setSearch(e.target.value);
+ setPage(1);
+ }}
onSearch={(value) => {
setSearch(value);
setPage(1);