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 6202e6f42 feat(studio): paginate system alerts, alert rules, Apache
consumer groups and ACL users (#2581)
6202e6f42 is described below
commit 6202e6f42b1cffce55926cd08723d4a74d0b33cd
Author: xdz997 <[email protected]>
AuthorDate: Thu Aug 27 14:15:50 2026 +0800
feat(studio): paginate system alerts, alert rules, Apache consumer groups
and ACL users (#2581)
* feat: paginate system alerts
* fix: paginate Apache ACL users in the database
* fix: normalize ACL user keyword search
* feat: paginate alert rules
* fix: paginate Apache consumer groups in the database
---
.../studio/instance/acl/AclRepository.java | 2 +
.../rocketmq/studio/instance/acl/AclService.java | 20 +-
.../instance/acl/MybatisPlusAclRepository.java | 17 ++
.../studio/instance/topic/MetadataService.java | 14 +-
.../rocketmq/studio/ops/alert/AlertRepository.java | 13 ++
.../studio/ops/alert/AlertRuleController.java | 11 +
.../rocketmq/studio/ops/alert/AlertService.java | 50 ++++-
.../ops/alert/MybatisPlusAlertRepository.java | 73 +++++-
.../studio/ops/alert/SystemAlertController.java | 9 +
.../rocketmq/studio/provider/InstanceProvider.java | 10 +
.../provider/apache/ApacheInstanceProvider.java | 6 +
.../studio/provider/apache/MetadataProvider.java | 15 ++
.../provider/apache/RocketMQMetadataProvider.java | 53 +++--
.../studio/instance/acl/AclServiceTest.java | 27 +++
.../instance/acl/MybatisPlusAclRepositoryTest.java | 36 +++
.../studio/instance/topic/MetadataServiceTest.java | 20 +-
.../studio/ops/alert/AlertRuleControllerTest.java | 26 +++
.../studio/ops/alert/AlertServiceTest.java | 83 ++++++-
.../ops/alert/MybatisPlusAlertRepositoryTest.java | 70 ++++++
.../apache/ApacheInstanceProviderTest.java | 12 +
.../apache/RocketMQMetadataProviderTest.java | 37 +++
web/src/api/ops.ts | 33 +++
web/src/i18n/translations.ts | 2 +
web/src/pages/ops/__tests__/AlertsPage.test.tsx | 65 +++++-
.../pages/ops/__tests__/SystemAlertsPage.test.tsx | 144 +++++++-----
web/src/pages/ops/alerts.tsx | 97 ++++++--
web/src/pages/ops/systemAlerts.tsx | 248 ++++++++++++---------
web/src/services/opsService.ts | 46 ++++
28 files changed, 979 insertions(+), 260 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclRepository.java
index 944cd4d80..8e247cbf0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclRepository.java
@@ -37,6 +37,8 @@ public interface AclRepository {
List<AclUserVO> findUsers();
+ PageResult<AclUserVO> findUserPage(String keyword, int page, int pageSize);
+
Optional<AclUserVO> findUserById(Long id);
AclUserVO saveUser(AclUserVO user);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
index 716ad88a9..31bf49c19 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/AclService.java
@@ -149,19 +149,25 @@ public class AclService {
if (page < 1 || pageSize < 1 || pageSize > 100) {
throw new BusinessException(400, "page must be >= 1 and pageSize
must be between 1 and 100");
}
- String query = keyword == null ? "" :
keyword.trim().toLowerCase(Locale.ROOT);
- List<AclUserVO> users = (isTencentInstance(instanceId)
- ? tencentAclService.listUsers(instanceId) :
aclRepository.findUsers()).stream()
+ if (isTencentInstance(instanceId)) {
+ String query = keyword == null ? "" :
keyword.trim().toLowerCase(Locale.ROOT);
+ List<AclUserVO> users =
tencentAclService.listUsers(instanceId).stream()
.filter(user -> query.isEmpty()
|| containsIgnoreCase(user.getUsername(), query)
|| containsIgnoreCase(user.getAccessKey(), query))
.sorted(Comparator.comparing(AclUserVO::getGmtCreate,
Comparator.nullsLast(Comparator.reverseOrder()))
.thenComparing(AclUserVO::getId,
Comparator.nullsLast(Comparator.naturalOrder())))
.toList();
- int from = (int) Math.min(Pagination.pageOffset(page, pageSize),
users.size());
- int to = Math.min(from + pageSize, users.size());
- return PageResult.of(users.subList(from,
to).stream().map(this::maskCredentials).toList(),
- users.size(), page, pageSize);
+ int from = (int) Math.min(Pagination.pageOffset(page, pageSize),
users.size());
+ int to = Math.min(from + pageSize, users.size());
+ return PageResult.of(
+ users.subList(from,
to).stream().map(this::maskCredentials).toList(),
+ users.size(), page, pageSize);
+ }
+ PageResult<AclUserVO> result = aclRepository.findUserPage(
+ StringUtils.hasText(keyword) ? keyword.trim() : null, page,
pageSize);
+ return
PageResult.of(result.getItems().stream().map(this::maskCredentials).toList(),
+ result.getTotal(), result.getPage(), result.getSize());
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
index 76d7520f0..5f463c9ed 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
@@ -103,6 +103,23 @@ public class MybatisPlusAclRepository implements
AclRepository {
.collect(Collectors.toList());
}
+ @Override
+ public PageResult<AclUserVO> findUserPage(String keyword, int page, int
pageSize) {
+ String search = StringUtils.hasText(keyword) ?
keyword.trim().toLowerCase() : null;
+ QueryWrapper<RmqAclUser> query = new QueryWrapper<RmqAclUser>()
+ .and(search != null, w -> w
+ .like("username", search)
+ .or().like("access_key", search))
+ .orderByDesc("gmt_create")
+ .orderByDesc("id");
+ IPage<RmqAclUser> mapperPage = userMapper.selectPage(new Page<>(page,
pageSize), query);
+ List<AclUserVO> items = mapperPage.getRecords().stream()
+ .map(MybatisPlusAclRepository::toUserVO)
+ .collect(Collectors.toList());
+ return PageResult.of(items, mapperPage.getTotal(), (int)
mapperPage.getCurrent(),
+ (int) mapperPage.getSize());
+ }
+
@Override
public Optional<AclUserVO> findUserById(Long id) {
if (id == null) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index ae77ed0e9..ded09bee3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -18,7 +18,6 @@ package org.apache.rocketmq.studio.instance.topic;
import org.apache.rocketmq.studio.provider.apache.AdminClient;
import org.apache.rocketmq.studio.provider.apache.MetadataProvider;
-import org.apache.rocketmq.studio.common.util.Pagination;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -186,12 +185,13 @@ public class MetadataService {
public PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String clusterId, String search,
int page, int
pageSize) {
validatePagination(page, pageSize);
-
- List<ConsumerGroupVO> groups = listConsumerGroups(instanceId,
clusterId, search);
- int total = groups.size();
- int from = (int) Math.min(Pagination.pageOffset(page, pageSize),
total);
- int to = Math.min(from + pageSize, total);
- return PageResult.of(groups.subList(from, to), total, page, pageSize);
+ instanceId = normalizeInstanceId(instanceId);
+ if (!StringUtils.hasText(instanceId) &&
StringUtils.hasText(clusterId)) {
+ return
metadataProvider.listConsumerGroupsPage(normalizeFilter(clusterId),
+ normalizeFilter(search), page, pageSize);
+ }
+ return resolve(instanceId).listConsumerGroupsPage(instanceId,
normalizeFilter(search),
+ page, pageSize);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRepository.java
index ae4e4b59e..92ca40b05 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRepository.java
@@ -18,10 +18,19 @@ package org.apache.rocketmq.studio.ops.alert;
import java.util.List;
+import java.util.Optional;
+
+import org.apache.rocketmq.studio.common.domain.PageResult;
public interface AlertRepository {
List<AlertRuleVO> findAllRules();
+ PageResult<AlertRuleVO> findRulePage(String search, Boolean enabled, int
page, int pageSize);
+
+ Optional<AlertRuleVO> findRuleById(Long id);
+
+ List<AlertRuleVO> findRulesByIds(List<Long> ids);
+
AlertRuleVO saveRule(AlertRuleVO rule);
boolean replaceRule(AlertRuleVO rule);
@@ -30,6 +39,10 @@ public interface AlertRepository {
List<SystemAlertVO> findAlerts(String level);
+ PageResult<SystemAlertVO> findAlerts(String level, int page, int pageSize);
+
+ Optional<SystemAlertVO> findAlertById(Long id);
+
boolean acknowledgeAlert(SystemAlertVO alert);
int deleteAcknowledgedAlerts();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleController.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleController.java
index 044b99982..36d75c52c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleController.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.alert;
import org.apache.rocketmq.studio.common.domain.Result;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
@@ -24,6 +25,7 @@ import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
@@ -40,6 +42,15 @@ public class AlertRuleController {
return Result.ok(alertService.listRules());
}
+ @GetMapping("/page")
+ public Result<PageResult<AlertRuleVO>> listRulesPage(
+ @RequestParam(required = false) String search,
+ @RequestParam(required = false) Boolean enabled,
+ @RequestParam(defaultValue = "1") int page,
+ @RequestParam(defaultValue = "20") int pageSize) {
+ return Result.ok(alertService.listRules(search, enabled, page,
pageSize));
+ }
+
@GetMapping("/export")
public Result<AlertRulesYamlVO> exportRules() {
return Result.ok(new
AlertRulesYamlVO(alertService.exportPrometheusRulesYaml()));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
index 9d98ab672..d8cd43fa6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
@@ -17,7 +17,9 @@
package org.apache.rocketmq.studio.ops.alert;
import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.audit.OperationAuditService;
+import org.springframework.util.StringUtils;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
@@ -28,7 +30,6 @@ import java.util.List;
import java.util.HashSet;
import java.util.Locale;
import java.util.Map;
-import java.util.Objects;
import java.util.Set;
import java.util.regex.Pattern;
@@ -51,6 +52,15 @@ public class AlertService {
return alertRepository.findAllRules();
}
+ public PageResult<AlertRuleVO> listRules(String search, Boolean enabled,
int page,
+ int pageSize) {
+ validateRulePagination(page, pageSize);
+ String normalizedSearch = StringUtils.hasText(search) ? search.trim()
: null;
+ log.info("Listing alert rules, search={}, enabled={}, page={},
pageSize={}",
+ normalizedSearch, enabled, page, pageSize);
+ return alertRepository.findRulePage(normalizedSearch, enabled, page,
pageSize);
+ }
+
public String exportPrometheusRulesYaml() {
List<AlertRuleVO> rules = alertRepository.findAllRules().stream()
.filter(AlertRuleVO::isEnabled)
@@ -128,10 +138,7 @@ public class AlertService {
public AlertRuleVO toggleRule(Long id, boolean enabled) {
log.info("Toggling alert rule id={}, enabled={}", id, enabled);
validateRuleId(id);
- List<AlertRuleVO> rules = alertRepository.findAllRules();
- AlertRuleVO rule = rules.stream()
- .filter(r -> Objects.equals(r.getId(), id))
- .findFirst()
+ AlertRuleVO rule = alertRepository.findRuleById(id)
.orElseThrow(() -> new
org.apache.rocketmq.studio.common.exception.BusinessException(404, "Alert rule
not found: " + id));
rule.setEnabled(enabled);
AlertRuleVO saved = alertRepository.saveRule(rule);
@@ -152,7 +159,7 @@ public class AlertService {
public AlertRuleBulkResultVO bulkToggleRules(List<Long> ids, boolean
enabled) {
List<Long> normalizedIds = normalizeBulkIds(ids);
Map<Long, AlertRuleVO> rulesById = new LinkedHashMap<>();
- for (AlertRuleVO rule : alertRepository.findAllRules()) {
+ for (AlertRuleVO rule : alertRepository.findRulesByIds(normalizedIds))
{
rulesById.put(rule.getId(), rule);
}
List<Long> succeeded = new ArrayList<>();
@@ -218,22 +225,36 @@ public class AlertService {
return normalized;
}
+ private static void validateRulePagination(int page, int pageSize) {
+ if (page < 1) {
+ throw new BusinessException(400, "page must be greater than zero");
+ }
+ if (pageSize < 1 || pageSize > 100) {
+ throw new BusinessException(400, "pageSize must be between 1 and
100");
+ }
+ }
+
public List<SystemAlertVO> listAlerts(String level) {
log.info("Listing system alerts, level={}", level);
return alertRepository.findAlerts(level);
}
+ public PageResult<SystemAlertVO> listAlerts(String level, int page, int
pageSize) {
+ validateAlertPagination(page, pageSize);
+ String normalizedLevel = StringUtils.hasText(level) ? level.trim() :
level;
+ log.info("Listing system alerts, level={}, page={}, pageSize={}",
+ normalizedLevel, page, pageSize);
+ return alertRepository.findAlerts(normalizedLevel, page, pageSize);
+ }
+
public SystemAlertVO acknowledgeAlert(Long id) {
log.info("Acknowledging system alert id={}", id);
if (id == null) {
throw new BusinessException(400, "System alert ID is required");
}
- List<SystemAlertVO> alerts = alertRepository.findAlerts(null);
- SystemAlertVO alert = alerts.stream()
- .filter(a -> Objects.equals(a.getId(), id))
- .findFirst()
+ SystemAlertVO alert = alertRepository.findAlertById(id)
.orElseThrow(() -> new
org.apache.rocketmq.studio.common.exception.BusinessException(404, "System
alert not found: " + id));
alert.setAcknowledged(true);
if (!alertRepository.acknowledgeAlert(alert)) {
@@ -257,6 +278,15 @@ public class AlertService {
return alertRuleAssetService.loadDefaultRules();
}
+ private static void validateAlertPagination(int page, int pageSize) {
+ if (page < 1) {
+ throw new BusinessException(400, "page must be greater than 0");
+ }
+ if (pageSize < 1 || pageSize > 100) {
+ throw new BusinessException(400, "pageSize must be between 1 and
100");
+ }
+ }
+
private PrometheusAlertRule toPrometheusRule(AlertRuleVO rule) {
String team = inferTeam(rule.getMetric());
return new PrometheusAlertRule(
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
index c3c6d4538..8755f4ce6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
@@ -17,6 +17,8 @@
package org.apache.rocketmq.studio.ops.alert;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.AlertLevel;
import org.apache.rocketmq.studio.persistence.entity.RmqAlertRule;
import org.apache.rocketmq.studio.persistence.entity.RmqSystemAlert;
@@ -31,6 +33,7 @@ import java.time.LocalDateTime;
import java.util.Arrays;
import java.util.List;
import java.util.Locale;
+import java.util.Optional;
import java.util.stream.Collectors;
/**
@@ -50,6 +53,45 @@ public class MybatisPlusAlertRepository implements
AlertRepository {
.collect(Collectors.toList());
}
+ @Override
+ public PageResult<AlertRuleVO> findRulePage(String search, Boolean enabled,
+ int page, int pageSize) {
+ QueryWrapper<RmqAlertRule> query = ruleQuery(search, enabled);
+ Page<RmqAlertRule> result = ruleMapper.selectPage(new Page<>(page,
pageSize), query);
+ List<AlertRuleVO> items = result.getRecords().stream()
+ .map(MybatisPlusAlertRepository::toRuleVO)
+ .toList();
+ return PageResult.of(items, result.getTotal(), page, pageSize);
+ }
+
+ @Override
+ public Optional<AlertRuleVO> findRuleById(Long id) {
+ if (id == null) {
+ return Optional.empty();
+ }
+ return Optional.ofNullable(ruleMapper.selectById(id))
+ .map(MybatisPlusAlertRepository::toRuleVO);
+ }
+
+ @Override
+ public List<AlertRuleVO> findRulesByIds(List<Long> ids) {
+ if (ids == null || ids.isEmpty()) {
+ return List.of();
+ }
+ return ruleMapper.selectList(new QueryWrapper<RmqAlertRule>()
+ .in("id", ids))
+ .stream()
+ .map(MybatisPlusAlertRepository::toRuleVO)
+ .toList();
+ }
+
+ private QueryWrapper<RmqAlertRule> ruleQuery(String search, Boolean
enabled) {
+ return new QueryWrapper<RmqAlertRule>()
+ .like(StringUtils.hasText(search), "name", search)
+ .eq(enabled != null, "enabled", enabled)
+ .orderByAsc("name", "id");
+ }
+
@Override
@Transactional
public AlertRuleVO saveRule(AlertRuleVO rule) {
@@ -78,14 +120,37 @@ public class MybatisPlusAlertRepository implements
AlertRepository {
@Override
public List<SystemAlertVO> findAlerts(String level) {
- QueryWrapper<RmqSystemAlert> query = new QueryWrapper<RmqSystemAlert>()
- .eq(StringUtils.hasText(level), "level", level == null ? null
: level.toLowerCase(Locale.ROOT))
- .orderByDesc("time");
- return alertMapper.selectList(query).stream()
+ return alertMapper.selectList(alertQuery(level)).stream()
.map(MybatisPlusAlertRepository::toAlertVO)
.collect(Collectors.toList());
}
+ @Override
+ public PageResult<SystemAlertVO> findAlerts(String level, int page, int
pageSize) {
+ Page<RmqSystemAlert> result = alertMapper.selectPage(
+ new Page<>(page, pageSize), alertQuery(level));
+ List<SystemAlertVO> items = result.getRecords().stream()
+ .map(MybatisPlusAlertRepository::toAlertVO)
+ .toList();
+ return PageResult.of(items, result.getTotal(), page, pageSize);
+ }
+
+ @Override
+ public Optional<SystemAlertVO> findAlertById(Long id) {
+ if (id == null) {
+ return Optional.empty();
+ }
+ return Optional.ofNullable(alertMapper.selectById(id))
+ .map(MybatisPlusAlertRepository::toAlertVO);
+ }
+
+ private QueryWrapper<RmqSystemAlert> alertQuery(String level) {
+ return new QueryWrapper<RmqSystemAlert>()
+ .eq(StringUtils.hasText(level), "level",
+ level == null ? null : level.toLowerCase(Locale.ROOT))
+ .orderByDesc("time", "id");
+ }
+
@Override
public boolean acknowledgeAlert(SystemAlertVO alert) {
return alertMapper.updateById(toAlertEntity(alert)) > 0;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/SystemAlertController.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/SystemAlertController.java
index b1ef42a67..d731153df 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/SystemAlertController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/SystemAlertController.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.alert;
import org.apache.rocketmq.studio.common.domain.Result;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
@@ -43,6 +44,14 @@ public class SystemAlertController {
return Result.ok(alertService.listAlerts(level));
}
+ @GetMapping("/page")
+ public Result<PageResult<SystemAlertVO>> listAlertsPage(
+ @RequestParam(required = false) String level,
+ @RequestParam(defaultValue = "1") int page,
+ @RequestParam(defaultValue = "20") int pageSize) {
+ return Result.ok(alertService.listAlerts(level, page, pageSize));
+ }
+
@PostMapping("/acknowledge")
public Result<SystemAlertVO> acknowledgeAlert(
@Valid @RequestBody(required = false) AcknowledgeSystemAlertDTO
request) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
index d83c8e377..e7aceb622 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
@@ -86,6 +86,16 @@ public interface InstanceProvider {
List<ConsumerGroupVO> listConsumerGroups(String instanceId, String search);
+ default PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String search,
+ int page, int pageSize) {
+ List<ConsumerGroupVO> groups = listConsumerGroups(instanceId, search);
+ int total = groups.size();
+ long offset = Pagination.pageOffset(page, pageSize);
+ int from = (int) Math.min(offset, total);
+ int to = from + (int) Math.min(pageSize, total - from);
+ return PageResult.of(groups.subList(from, to), total, page, pageSize);
+ }
+
ConsumerGroupVO createConsumerGroup(String instanceId, ConsumerGroupVO
group);
void deleteConsumerGroup(String instanceId, String groupName);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index a9e94e7dc..1694901fa 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -122,6 +122,12 @@ public class ApacheInstanceProvider implements
InstanceProvider {
return metadataProvider.listConsumerGroups(instanceId, null, search);
}
+ @Override
+ public PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String search,
+ int page, int pageSize) {
+ return metadataProvider.listConsumerGroupsPage(instanceId, null,
search, page, pageSize);
+ }
+
@Override
public ConsumerGroupVO createConsumerGroup(String instanceId,
ConsumerGroupVO group) {
return adminClient.createConsumerGroup(group);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
index 7e16d005e..66454d34e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/MetadataProvider.java
@@ -56,6 +56,21 @@ public interface MetadataProvider {
return listConsumerGroups(clusterId, search);
}
+ default PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
clusterId, String search,
+ int page, int pageSize) {
+ List<ConsumerGroupVO> groups = listConsumerGroups(clusterId, search);
+ int total = groups.size();
+ long offset = Pagination.pageOffset(page, pageSize);
+ int from = (int) Math.min(offset, total);
+ int to = from + (int) Math.min(pageSize, total - from);
+ return PageResult.of(groups.subList(from, to), total, page, pageSize);
+ }
+
+ default PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String clusterId,
+ String search, int page, int pageSize) {
+ return listConsumerGroupsPage(clusterId, search, page, pageSize);
+ }
+
List<BrokerRouteVO> getTopicRoutes(String instanceId, String name);
List<TopicConsumerVO> getTopicConsumers(String instanceId, String name);
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 9ed5a8d01..fedc63108 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
@@ -227,29 +227,28 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
List<ConsumerGroupVO> result = new ArrayList<>();
for (RmqGroup entity : groupMapper.selectList(query)) {
- ConsumerGroupVO vo = new ConsumerGroupVO();
- vo.setId(entity.getId());
- vo.setName(entity.getName());
- vo.setClusterId(entity.getClusterId());
- vo.setInstanceId(entity.getInstanceId());
- // consumeType stores the real ConsumeType
("CLUSTERING"/"BROADCASTING"); messageModel
- // holds the subscription mode ("Push"/"Pop") and is not a
ConsumeType.
- vo.setConsumeType(parseConsumeType(entity.getConsumeType()));
- // messageModel stores the subscription mode ("Push"/"Pop");
surface it so read paths
- // (web detail, AI rmq.group.list) never see a null
subscriptionMode.
-
vo.setSubscriptionMode(parseSubscriptionMode(entity.getMessageModel()));
- vo.setRetryMaxTimes(entity.getMaxRetry() == null ? 0 :
entity.getMaxRetry());
- vo.setGmtCreate(entity.getGmtCreate());
- vo.setGmtModified(entity.getGmtModified());
-
- // Live stats (online clients, lag, delay) are enriched in
parallel below with a
- // bounded executor and per-group timeout instead of blocking the
listing.
- result.add(vo);
+ result.add(toConsumerGroupVO(entity));
}
enrichLiveStats(instanceId, result);
return result;
}
+ @Override
+ public PageResult<ConsumerGroupVO> listConsumerGroupsPage(String
instanceId, String clusterId,
+ String search, int page, int pageSize) {
+ LambdaQueryWrapper<RmqGroup> query = new LambdaQueryWrapper<RmqGroup>()
+ .eq(instanceId != null, RmqGroup::getInstanceId,
normalizeMetadataScope(instanceId))
+ .eq(StringUtils.hasText(clusterId), RmqGroup::getClusterId,
clusterId)
+ .like(StringUtils.hasText(search), RmqGroup::getName, search)
+ .orderByAsc(RmqGroup::getName, RmqGroup::getId);
+ Page<RmqGroup> result = groupMapper.selectPage(new Page<>(page,
pageSize), query);
+ List<ConsumerGroupVO> groups = result.getRecords().stream()
+ .map(this::toConsumerGroupVO)
+ .toList();
+ enrichLiveStats(instanceId, groups);
+ return PageResult.of(groups, result.getTotal(), page, pageSize);
+ }
+
private static final int ONLINE_ENRICHMENT_THREADS = 8;
private static final long ONLINE_ENRICHMENT_TIMEOUT_SECONDS = 3;
@@ -265,6 +264,24 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
onlineEnrichmentExecutor.shutdownNow();
}
+ private ConsumerGroupVO toConsumerGroupVO(RmqGroup entity) {
+ ConsumerGroupVO vo = new ConsumerGroupVO();
+ vo.setId(entity.getId());
+ vo.setName(entity.getName());
+ vo.setClusterId(entity.getClusterId());
+ vo.setInstanceId(entity.getInstanceId());
+ // consumeType stores the real ConsumeType
("CLUSTERING"/"BROADCASTING"); messageModel
+ // holds the subscription mode ("Push"/"Pop") and is not a ConsumeType.
+ vo.setConsumeType(parseConsumeType(entity.getConsumeType()));
+ // messageModel stores the subscription mode ("Push"/"Pop"); surface
it so read paths
+ // (web detail, AI rmq.group.list) never see a null subscriptionMode.
+
vo.setSubscriptionMode(parseSubscriptionMode(entity.getMessageModel()));
+ vo.setRetryMaxTimes(entity.getMaxRetry() == null ? 0 :
entity.getMaxRetry());
+ vo.setGmtCreate(entity.getGmtCreate());
+ vo.setGmtModified(entity.getGmtModified());
+ return vo;
+ }
+
private void enrichLiveStats(String instanceId, List<ConsumerGroupVO>
groups) {
boolean noLiveSource = !StringUtils.hasText(instanceId) && !hasAdmin();
if (groups.isEmpty() || noLiveSource) {
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
index 33f318e59..472d60ee3 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/AclServiceTest.java
@@ -187,6 +187,33 @@ class AclServiceTest {
verifyNoInteractions(tencentAclService);
}
+ @Test
+ void pageUsersShouldUseDatabasePaginationForApacheUsers() {
+ AclUserVO user = AclUserVO.builder()
+ .id(1L)
+ .username("orders")
+ .accessKey("access-key-123456")
+ .secretKey("secret-key-987654")
+ .admin(false)
+ .clusters(List.of("cluster-a"))
+ .gmtCreate(LocalDateTime.now())
+ .build();
+ when(aclRepository.findUserPage("orders", 2, 20))
+ .thenReturn(PageResult.of(List.of(user), 21, 2, 20));
+
+ PageResult<AclUserVO> result = aclService.pageUsers(null, 2, 20, "
orders ");
+
+ assertThat(result.getItems()).singleElement().satisfies(listed -> {
+ assertThat(listed.getUsername()).isEqualTo("orders");
+ assertThat(listed.getAccessKey()).isEqualTo("acce****3456");
+ assertThat(listed.getSecretKey()).isEqualTo("secr****7654");
+ });
+ assertThat(result.getTotal()).isEqualTo(21);
+ verify(aclRepository).findUserPage("orders", 2, 20);
+ verify(aclRepository, never()).findUsers();
+ verifyNoInteractions(instanceRepository, tencentAclService);
+ }
+
@Test
void
listRulesShouldReturnEmptyPageWhenTencentPageOffsetExceedsIntegerRange() {
InstanceVO instance = InstanceVO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
index 21626d723..0e7e075f3 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
@@ -116,6 +116,42 @@ class MybatisPlusAclRepositoryTest {
.contains("principal", "resource", "scope", "decision",
"acl_version");
}
+ @Test
+ void findUserPageShouldNormalizeKeywordFiltersAndApplyDatabasePagination()
{
+ RmqAclUser entity = new RmqAclUser();
+ entity.setId(11L);
+ entity.setUsername("user-orders");
+ entity.setAccessKey("access-key-orders");
+ entity.setSecretKey(CredentialUtils.encodeBase64("secret-key"));
+ entity.setAdmin(false);
+ entity.setGmtCreate(LocalDateTime.of(2026, 8, 20, 9, 30));
+ Page<RmqAclUser> mapperPage = new Page<RmqAclUser>(3, 20)
+ .setRecords(List.of(entity))
+ .setTotal(51);
+ when(userMapper.selectPage(any(IPage.class),
any(Wrapper.class))).thenReturn(mapperPage);
+
+ PageResult<AclUserVO> result = repository.findUserPage(" ORDERS ", 3,
20);
+
+ ArgumentCaptor<IPage<RmqAclUser>> pageCaptor =
ArgumentCaptor.forClass(IPage.class);
+ ArgumentCaptor<Wrapper<RmqAclUser>> queryCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ verify(userMapper).selectPage(pageCaptor.capture(),
queryCaptor.capture());
+ assertThat(pageCaptor.getValue().getCurrent()).isEqualTo(3);
+ assertThat(pageCaptor.getValue().getSize()).isEqualTo(20);
+ assertThat(result.getTotal()).isEqualTo(51);
+ assertThat(result.getPage()).isEqualTo(3);
+ assertThat(result.getSize()).isEqualTo(20);
+ assertThat(result.getItems()).singleElement().satisfies(user -> {
+ assertThat(user.getUsername()).isEqualTo("user-orders");
+ assertThat(user.getAccessKey()).isEqualTo("access-key-orders");
+ assertThat(user.getSecretKey()).isEqualTo("secret-key");
+ });
+ assertThat(queryCaptor.getValue().getSqlSegment())
+ .contains("username", "access_key", "ORDER BY gmt_create
DESC,id DESC");
+ assertThat(((QueryWrapper<RmqAclUser>)
queryCaptor.getValue()).getParamNameValuePairs())
+ .containsValue("%orders%");
+ verify(userMapper, never()).selectList(any(QueryWrapper.class));
+ }
+
@Test
void replaceRuleShouldReturnEmptyWhenConcurrentDeleteWins() {
RmqAclRule existing = new RmqAclRule();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index aa6e713ae..595be3523 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -287,14 +287,10 @@ class MetadataServiceTest {
@Test
void listConsumerGroupsPageShouldPaginateFromOneBasedIndexes() {
- ConsumerGroupVO first = new ConsumerGroupVO();
- first.setName("cg-a");
- ConsumerGroupVO second = new ConsumerGroupVO();
- second.setName("cg-b");
ConsumerGroupVO third = new ConsumerGroupVO();
third.setName("cg-c");
- when(apacheProvider.listConsumerGroups("instance-a", "order"))
- .thenReturn(List.of(first, second, third));
+ when(apacheProvider.listConsumerGroupsPage("instance-a", "order", 2,
2))
+ .thenReturn(PageResult.of(List.of(third), 3, 2, 2));
PageResult<ConsumerGroupVO> result =
metadataService.listConsumerGroupsPage("instance-a", null,
"order", 2, 2);
@@ -303,23 +299,23 @@ class MetadataServiceTest {
assertThat(result.getTotal()).isEqualTo(3);
assertThat(result.getPage()).isEqualTo(2);
assertThat(result.getSize()).isEqualTo(2);
- verify(apacheProvider).listConsumerGroups("instance-a", "order");
+ verify(apacheProvider).listConsumerGroupsPage("instance-a", "order",
2, 2);
+ verify(apacheProvider,
org.mockito.Mockito.never()).listConsumerGroups("instance-a", "order");
}
@Test
void
listConsumerGroupsPageShouldReturnEmptyItemsWhenPageStartsPastFilteredTotal() {
- ConsumerGroupVO first = new ConsumerGroupVO();
- first.setName("cg-a");
- when(metadataProvider.listConsumerGroups("cluster-1",
"order")).thenReturn(List.of(first));
+ when(metadataProvider.listConsumerGroupsPage("cluster-1", "order", 2,
1))
+ .thenReturn(PageResult.empty(2, 1));
PageResult<ConsumerGroupVO> result =
metadataService.listConsumerGroupsPage(null, "cluster-1",
"order", 2, 1);
assertThat(result.getItems()).isEmpty();
- assertThat(result.getTotal()).isEqualTo(1);
+ assertThat(result.getTotal()).isZero();
assertThat(result.getPage()).isEqualTo(2);
assertThat(result.getSize()).isEqualTo(1);
- verify(metadataProvider).listConsumerGroups("cluster-1", "order");
+ verify(metadataProvider).listConsumerGroupsPage("cluster-1", "order",
2, 1);
}
@Test
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleControllerTest.java
index 40762f7db..c7442d4d3 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleControllerTest.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.alert;
import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import
org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc;
@@ -68,6 +69,31 @@ class AlertRuleControllerTest {
.andExpect(jsonPath("$.data[0].enabled").value(true));
}
+ @Test
+ void listRulesPageShouldPassFiltersAndReturnPageContract() throws
Exception {
+ AlertRuleVO rule = AlertRuleVO.builder()
+ .id(1L)
+ .name("High Lag")
+ .enabled(true)
+ .build();
+ when(alertService.listRules("lag", true, 2, 20))
+ .thenReturn(PageResult.of(List.of(rule), 21, 2, 20));
+
+ mockMvc.perform(get("/api/alert-rules/page")
+ .param("search", "lag")
+ .param("enabled", "true")
+ .param("page", "2")
+ .param("pageSize", "20"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200))
+ .andExpect(jsonPath("$.data.total").value(21))
+ .andExpect(jsonPath("$.data.page").value(2))
+ .andExpect(jsonPath("$.data.size").value(20))
+ .andExpect(jsonPath("$.data.items[0].id").value(1));
+
+ verify(alertService).listRules("lag", true, 2, 20);
+ }
+
@Test
void exportRulesShouldReturnGeneratedYaml() throws Exception {
when(alertService.exportPrometheusRulesYaml()).thenReturn("groups:\n");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
index a59ee9927..051105355 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.ops.alert;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.AlertLevel;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.audit.OperationAuditService;
@@ -31,10 +32,12 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Locale;
+import java.util.Optional;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq;
@@ -84,6 +87,32 @@ class AlertServiceTest {
assertThat(result).isEmpty();
}
+ @Test
+ void listRulesShouldNormalizeSearchAndDelegateFiltersToRepository() {
+ PageResult<AlertRuleVO> repositoryPage = PageResult.of(List.of(), 0,
2, 20);
+ when(alertRepository.findRulePage("lag", true, 2,
20)).thenReturn(repositoryPage);
+
+ PageResult<AlertRuleVO> result = alertService.listRules(" lag ",
true, 2, 20);
+
+ assertThat(result).isSameAs(repositoryPage);
+ verify(alertRepository).findRulePage("lag", true, 2, 20);
+ }
+
+ @Test
+ void listRulesShouldRejectInvalidPaginationBeforeRepositoryAccess() {
+ assertThatThrownBy(() -> alertService.listRules("lag", true, 0, 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("page must be greater than zero");
+ assertThatThrownBy(() -> alertService.listRules("lag", true, 1, 0))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("pageSize must be between 1 and 100");
+ assertThatThrownBy(() -> alertService.listRules("lag", true, 1, 101))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("pageSize must be between 1 and 100");
+
+ verify(alertRepository, never()).findRulePage(any(), any(), anyInt(),
anyInt());
+ }
+
@Test
void
exportPrometheusRulesYamlShouldReturnDefaultRulesWhenRepositoryIsEmpty() {
when(alertRepository.findAllRules()).thenReturn(Collections.emptyList());
@@ -641,7 +670,7 @@ class AlertServiceTest {
@Test
void toggleRuleShouldEnableRule() {
AlertRuleVO existing = AlertRuleVO.builder().id(1L).name("CPU
Alert").enabled(false).build();
- when(alertRepository.findAllRules()).thenReturn(List.of(existing));
+
when(alertRepository.findRuleById(1L)).thenReturn(Optional.of(existing));
when(alertRepository.saveRule(any(AlertRuleVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
AlertRuleVO result = alertService.toggleRule(1L, true);
@@ -655,7 +684,7 @@ class AlertServiceTest {
@Test
void toggleRuleShouldDisableRule() {
AlertRuleVO existing = AlertRuleVO.builder().id(1L).name("CPU
Alert").enabled(true).build();
- when(alertRepository.findAllRules()).thenReturn(List.of(existing));
+
when(alertRepository.findRuleById(1L)).thenReturn(Optional.of(existing));
when(alertRepository.saveRule(any(AlertRuleVO.class))).thenAnswer(invocation ->
invocation.getArgument(0));
AlertRuleVO result = alertService.toggleRule(1L, false);
@@ -671,11 +700,12 @@ class AlertServiceTest {
.satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
verify(alertRepository, never()).findAllRules();
+ verify(alertRepository, never()).findRuleById(any());
}
@Test
- void toggleRuleShouldIgnorePersistedRulesWithNullIds() {
-
when(alertRepository.findAllRules()).thenReturn(List.of(AlertRuleVO.builder().name("corrupt").build()));
+ void toggleRuleShouldUseADirectIdLookupWhenRuleIsMissing() {
+ when(alertRepository.findRuleById(999L)).thenReturn(Optional.empty());
assertThatThrownBy(() -> alertService.toggleRule(999L, true))
.isInstanceOf(BusinessException.class)
@@ -684,7 +714,7 @@ class AlertServiceTest {
@Test
void toggleRuleShouldThrowWhenRuleNotFound() {
-
when(alertRepository.findAllRules()).thenReturn(Collections.emptyList());
+ when(alertRepository.findRuleById(999L)).thenReturn(Optional.empty());
assertThatThrownBy(() -> alertService.toggleRule(999L, true))
.isInstanceOf(BusinessException.class)
@@ -724,7 +754,7 @@ class AlertServiceTest {
@Test
void bulkToggleShouldDeduplicateIdsAndReportMissingRules() {
AlertRuleVO rule = AlertRuleVO.builder().id(1L).name("High
CPU").enabled(false).build();
- when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+ when(alertRepository.findRulesByIds(List.of(1L,
999L))).thenReturn(List.of(rule));
when(alertRepository.replaceRule(any(AlertRuleVO.class))).thenReturn(true);
AlertRuleBulkResultVO result = alertService.bulkToggleRules(
@@ -740,7 +770,7 @@ class AlertServiceTest {
@Test
void
bulkToggleShouldReportRulesDeletedConcurrentlyInsteadOfRecreatingThem() {
AlertRuleVO rule = AlertRuleVO.builder().id(1L).name("High
CPU").enabled(false).build();
- when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
when(alertRepository.findRulesByIds(List.of(1L))).thenReturn(List.of(rule));
when(alertRepository.replaceRule(any(AlertRuleVO.class))).thenReturn(false);
AlertRuleBulkResultVO result =
alertService.bulkToggleRules(List.of(1L), true);
@@ -797,11 +827,40 @@ class AlertServiceTest {
verify(alertRepository).findAlerts(null);
}
+ @Test
+ void listAlertsPageShouldReturnTheRequestedPage() {
+ SystemAlertVO alert =
SystemAlertVO.builder().id(7L).level(AlertLevel.warning)
+ .title("Slow Consumer").build();
+ when(alertRepository.findAlerts("warning", 2, 20))
+ .thenReturn(PageResult.of(List.of(alert), 21, 2, 20));
+
+ PageResult<SystemAlertVO> result = alertService.listAlerts(" warning
", 2, 20);
+
+ assertThat(result.getItems()).containsExactly(alert);
+ assertThat(result.getTotal()).isEqualTo(21);
+ assertThat(result.getPage()).isEqualTo(2);
+ assertThat(result.getSize()).isEqualTo(20);
+ verify(alertRepository).findAlerts("warning", 2, 20);
+ }
+
+ @Test
+ void listAlertsPageShouldRejectInvalidPaginationBeforeRepositoryAccess() {
+ assertThatThrownBy(() -> alertService.listAlerts(null, 0, 20))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("page must be greater than 0");
+ assertThatThrownBy(() -> alertService.listAlerts(null, 1, 101))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("pageSize must be between 1 and 100");
+
+ verify(alertRepository, never()).findAlerts(any(),
org.mockito.ArgumentMatchers.anyInt(),
+ org.mockito.ArgumentMatchers.anyInt());
+ }
+
@Test
void acknowledgeAlertShouldSetAcknowledgedTrue() {
SystemAlertVO existing =
SystemAlertVO.builder().id(1L).level(AlertLevel.error)
.title("Broker Down").acknowledged(false).build();
- when(alertRepository.findAlerts(null)).thenReturn(List.of(existing));
+
when(alertRepository.findAlertById(1L)).thenReturn(java.util.Optional.of(existing));
when(alertRepository.acknowledgeAlert(any(SystemAlertVO.class))).thenReturn(true);
SystemAlertVO result = alertService.acknowledgeAlert(1L);
@@ -819,12 +878,12 @@ class AlertServiceTest {
.hasMessage("System alert ID is required")
.satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
- verify(alertRepository, never()).findAlerts(any());
+ verify(alertRepository, never()).findAlertById(any());
}
@Test
void acknowledgeAlertShouldIgnorePersistedAlertsWithNullIds() {
-
when(alertRepository.findAlerts(null)).thenReturn(List.of(SystemAlertVO.builder().title("corrupt").build()));
+
when(alertRepository.findAlertById(999L)).thenReturn(java.util.Optional.empty());
assertThatThrownBy(() -> alertService.acknowledgeAlert(999L))
.isInstanceOf(BusinessException.class)
@@ -833,7 +892,7 @@ class AlertServiceTest {
@Test
void acknowledgeAlertShouldThrowWhenAlertNotFound() {
-
when(alertRepository.findAlerts(null)).thenReturn(Collections.emptyList());
+
when(alertRepository.findAlertById(999L)).thenReturn(java.util.Optional.empty());
assertThatThrownBy(() -> alertService.acknowledgeAlert(999L))
.isInstanceOf(BusinessException.class)
@@ -844,7 +903,7 @@ class AlertServiceTest {
void acknowledgeAlertShouldRejectConcurrentRemoval() {
SystemAlertVO existing =
SystemAlertVO.builder().id(1L).level(AlertLevel.error)
.title("Broker Down").acknowledged(false).build();
- when(alertRepository.findAlerts(null)).thenReturn(List.of(existing));
+
when(alertRepository.findAlertById(1L)).thenReturn(java.util.Optional.of(existing));
when(alertRepository.acknowledgeAlert(any(SystemAlertVO.class))).thenReturn(false);
assertThatThrownBy(() -> alertService.acknowledgeAlert(1L))
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
index 0b81c7e4c..0d00083b7 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
@@ -18,12 +18,16 @@ package org.apache.rocketmq.studio.ops.alert;
import com.baomidou.mybatisplus.core.conditions.Wrapper;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.baomidou.mybatisplus.core.metadata.IPage;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.persistence.entity.RmqAlertRule;
import org.apache.rocketmq.studio.persistence.entity.RmqSystemAlert;
import org.apache.rocketmq.studio.persistence.mapper.RmqAlertRuleMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqSystemAlertMapper;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.ArgumentCaptor;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
@@ -60,6 +64,72 @@ class MybatisPlusAlertRepositoryTest {
assertThat(repository.replaceRule(rule)).isFalse();
}
+ @Test
+ void findRulePageShouldApplyFiltersOrderingAndDatabasePagination() {
+ RmqAlertRule entity = new RmqAlertRule();
+ entity.setId(1L);
+ entity.setName("High Lag");
+ entity.setEnabled(true);
+ entity.setChannels("email");
+ Page<RmqAlertRule> mapperPage = new Page<RmqAlertRule>(3, 20)
+ .setRecords(List.of(entity))
+ .setTotal(51);
+ when(ruleMapper.selectPage(any(IPage.class),
any(Wrapper.class))).thenReturn(mapperPage);
+
+ PageResult<AlertRuleVO> result = repository.findRulePage("lag", true,
3, 20);
+
+ ArgumentCaptor<IPage<RmqAlertRule>> pageCaptor =
ArgumentCaptor.forClass(IPage.class);
+ ArgumentCaptor<Wrapper<RmqAlertRule>> queryCaptor =
ArgumentCaptor.forClass(Wrapper.class);
+ verify(ruleMapper).selectPage(pageCaptor.capture(),
queryCaptor.capture());
+ assertThat(pageCaptor.getValue().getCurrent()).isEqualTo(3);
+ assertThat(pageCaptor.getValue().getSize()).isEqualTo(20);
+ assertThat(result.getTotal()).isEqualTo(51);
+ assertThat(result.getPage()).isEqualTo(3);
+ assertThat(result.getSize()).isEqualTo(20);
+ assertThat(result.getItems()).singleElement()
+ .satisfies(rule -> {
+ assertThat(rule.getId()).isEqualTo(1L);
+ assertThat(rule.getName()).isEqualTo("High Lag");
+ assertThat(rule.isEnabled()).isTrue();
+ });
+ QueryWrapper<RmqAlertRule> query = (QueryWrapper<RmqAlertRule>)
queryCaptor.getValue();
+ query.getCustomSqlSegment();
+ assertThat(query.getSqlSegment())
+ .contains("name", "enabled", "ORDER BY name ASC,id ASC");
+ assertThat(query.getParamNameValuePairs())
+ .containsValue("%lag%")
+ .containsValue(true);
+ verify(ruleMapper, never()).selectList(any());
+ }
+
+ @Test
+ void findRuleByIdShouldUsePrimaryKeyLookupWithoutFullListRead() {
+ RmqAlertRule entity = new RmqAlertRule();
+ entity.setId(9L);
+ entity.setName("High Lag");
+ entity.setEnabled(true);
+ when(ruleMapper.selectById(9L)).thenReturn(entity);
+
+ assertThat(repository.findRuleById(9L)).hasValueSatisfying(rule ->
+ assertThat(rule.getId()).isEqualTo(9L));
+ assertThat(repository.findRuleById(null)).isEmpty();
+ verify(ruleMapper, never()).selectList(any());
+ }
+
+ @Test
+ void findRulesByIdsShouldUseBoundedIdQuery() {
+ RmqAlertRule entity = new RmqAlertRule();
+ entity.setId(1L);
+ entity.setName("High Lag");
+ entity.setEnabled(true);
+
when(ruleMapper.selectList(any(Wrapper.class))).thenReturn(List.of(entity));
+
+ assertThat(repository.findRulesByIds(List.of(1L, 2L))).singleElement()
+ .satisfies(rule -> assertThat(rule.getId()).isEqualTo(1L));
+ assertThat(repository.findRulesByIds(null)).isEmpty();
+ assertThat(repository.findRulesByIds(List.of())).isEmpty();
+ }
+
@Test
void findAlertsShouldNormalizeStoredLevelValues() {
RmqSystemAlert entity = new RmqSystemAlert();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
index ed63f46fb..8c7a913ed 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
@@ -16,9 +16,11 @@
*/
package org.apache.rocketmq.studio.provider.apache;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.message.MessageProvider;
import org.apache.rocketmq.studio.provider.InstanceCapability;
import org.junit.jupiter.api.Test;
@@ -104,4 +106,14 @@ class ApacheInstanceProviderTest {
verify(metadataProvider).listConsumerGroups("inst-1", null, "orders");
}
+
+ @Test
+ void listConsumerGroupsPageShouldRouteThroughDatabasePaginationTest() {
+ PageResult<ConsumerGroupVO> page = PageResult.of(java.util.List.of(),
0, 1, 20);
+ when(metadataProvider.listConsumerGroupsPage("inst-1", null, "orders",
1, 20)).thenReturn(page);
+
+ assertThat(provider.listConsumerGroupsPage("inst-1", "orders", 1,
20)).isSameAs(page);
+
+ verify(metadataProvider).listConsumerGroupsPage("inst-1", null,
"orders", 1, 20);
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
index 543ddea0a..eec2cbe81 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
@@ -19,6 +19,8 @@ package org.apache.rocketmq.studio.provider.apache;
import com.baomidou.mybatisplus.core.MybatisConfiguration;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.metadata.TableInfoHelper;
+import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.ibatis.builder.MapperBuilderAssistant;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.common.message.MessageQueue;
@@ -197,6 +199,41 @@ class RocketMQMetadataProviderTest {
assertThat(captor.getValue().getSqlSegment()).contains("instance_id",
"cluster_id");
}
+ @Test
+ void listConsumerGroupsPageShouldUseDatabasePaginationAndStableOrdering() {
+ TableInfoHelper.initTableInfo(new MapperBuilderAssistant(new
MybatisConfiguration(), ""), RmqGroup.class);
+ RmqGroup entity = new RmqGroup();
+ entity.setId(21L);
+ entity.setName("group-b");
+ entity.setInstanceId("instance-a");
+ entity.setClusterId("cluster-1");
+ entity.setConsumeType("CLUSTERING");
+ entity.setMessageModel("Pop");
+ entity.setMaxRetry(5);
+ Page<RmqGroup> databasePage = new Page<>(2, 1, 21);
+ databasePage.setRecords(List.of(entity));
+ when(groupMapper.selectPage(any(Page.class),
any(LambdaQueryWrapper.class))).thenReturn(databasePage);
+ RocketMQMetadataProvider provider = newProvider();
+
+ PageResult<ConsumerGroupVO> result = provider.listConsumerGroupsPage(
+ "instance-a", "cluster-1", "group", 2, 1);
+
+ assertThat(result.getItems()).hasSize(1);
+ assertThat(result.getItems().get(0).getName()).isEqualTo("group-b");
+
assertThat(result.getItems().get(0).getSubscriptionMode()).isEqualTo(SubscriptionMode.Pop);
+ assertThat(result.getTotal()).isEqualTo(21);
+ assertThat(result.getPage()).isEqualTo(2);
+ assertThat(result.getSize()).isEqualTo(1);
+
+ org.mockito.ArgumentCaptor<LambdaQueryWrapper<RmqGroup>> captor =
+ org.mockito.ArgumentCaptor.forClass(LambdaQueryWrapper.class);
+ verify(groupMapper).selectPage(any(Page.class), captor.capture());
+ assertThat(captor.getValue().getSqlSegment())
+ .contains("instance_id", "cluster_id", "ORDER BY name ASC,id
ASC");
+ verify(groupMapper, never()).selectList(any());
+ verify(runtimeAdminClientResolver, times(2)).execute(eq("instance-a"),
any());
+ }
+
@Test
void getTopicRoutesShouldUseSelectedInstanceRuntimeClient() {
List<BrokerRouteVO> routes =
List.of(BrokerRouteVO.builder().brokerName("broker-a").build());
diff --git a/web/src/api/ops.ts b/web/src/api/ops.ts
index 83a25d043..d7061ac66 100644
--- a/web/src/api/ops.ts
+++ b/web/src/api/ops.ts
@@ -21,6 +21,13 @@ export interface AlertRuleBulkResult {
updatedRules: AlertRule[];
}
+export interface AlertRulePage {
+ items: AlertRule[];
+ total: number;
+ page: number;
+ size: number;
+}
+
// Matches mock/dashboard.ts systemAlerts
export interface SystemAlert {
id: number;
@@ -31,6 +38,13 @@ export interface SystemAlert {
acknowledged: boolean;
}
+export interface SystemAlertPage {
+ items: SystemAlert[];
+ total: number;
+ page: number;
+ size: number;
+}
+
// Matches mock/audit.ts (inferred from data)
export interface AuditRecord {
id: number;
@@ -70,6 +84,16 @@ export async function listAlertRules() {
return res.data.data;
}
+export async function listAlertRulesPage(params: {
+ search?: string;
+ enabled?: boolean;
+ page?: number;
+ pageSize?: number;
+}) {
+ const res = await client.get<{ data: AlertRulePage }>('/alert-rules/page', {
params });
+ return res.data.data;
+}
+
export async function createAlertRule(data: Partial<AlertRule>) {
const res = await client.post<{ data: AlertRule }>('/alert-rules/create',
data);
return res.data.data;
@@ -108,6 +132,15 @@ export async function listSystemAlerts() {
return res.data.data;
}
+export async function listSystemAlertsPage(params: {
+ level?: string;
+ page?: number;
+ pageSize?: number;
+}) {
+ const res = await client.get<{ data: SystemAlertPage
}>('/system-alerts/page', { params });
+ return res.data.data;
+}
+
export async function acknowledgeAlert(id: number) {
await client.post('/system-alerts/acknowledge', { id });
}
diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts
index 41ff6fb9d..455a0e185 100644
--- a/web/src/i18n/translations.ts
+++ b/web/src/i18n/translations.ts
@@ -229,6 +229,7 @@ const translations: Record<string, Record<Lang, string>> = {
},
'alerts.totalRules': { zh: '规则总数', en: 'Total Rules' },
'alerts.enabled': { zh: '已启用', en: 'Enabled' },
+ 'alerts.disabled': { zh: '已禁用', en: 'Disabled' },
'alerts.triggered24h': { zh: '24h 触发', en: 'Triggered (24h)' },
'alerts.newRule': { zh: '新建规则', en: 'New Rule' },
'alerts.ruleName': { zh: '规则名称', en: 'Rule Name' },
@@ -240,6 +241,7 @@ const translations: Record<string, Record<Lang, string>> = {
'alerts.neverTriggered': { zh: '从未触发', en: 'Never' },
'alerts.ruleCreated': { zh: '规则创建成功', en: 'Rule created successfully' },
'alerts.selectedRules': { zh: '已选择 {count} 条告警规则', en: 'Selected alert
rules: {count}' },
+ 'alerts.searchPlaceholder': { zh: '搜索规则名称或指标', en: 'Search rule name or
metric' },
'alerts.bulkEnable': { zh: '批量启用', en: 'Enable Selected' },
'alerts.bulkDisable': { zh: '批量禁用', en: 'Disable Selected' },
'alerts.bulkDelete': { zh: '批量删除', en: 'Delete Selected' },
diff --git a/web/src/pages/ops/__tests__/AlertsPage.test.tsx
b/web/src/pages/ops/__tests__/AlertsPage.test.tsx
index f8fea48fc..cc4cf7807 100644
--- a/web/src/pages/ops/__tests__/AlertsPage.test.tsx
+++ b/web/src/pages/ops/__tests__/AlertsPage.test.tsx
@@ -25,14 +25,14 @@ import AlertsPage from '../alerts';
import {
bulkDeleteAlertRules,
bulkToggleAlertRules,
- listAlertRules,
+ listAlertRulesPage,
toggleAlertRule,
} from '../../../services/opsService';
vi.mock('../../../services/opsService', () => ({
createAlertRule: vi.fn(),
deleteAlertRule: vi.fn(),
- listAlertRules: vi.fn(),
+ listAlertRulesPage: vi.fn(),
toggleAlertRule: vi.fn(),
bulkToggleAlertRules: vi.fn(),
bulkDeleteAlertRules: vi.fn(),
@@ -110,7 +110,12 @@ function getRuleRow(ruleName: string) {
describe('AlertsPage', () => {
beforeEach(() => {
vi.clearAllMocks();
- vi.mocked(listAlertRules).mockResolvedValue(alertRules.map(cloneRule));
+ vi.mocked(listAlertRulesPage).mockResolvedValue({
+ items: alertRules.map(cloneRule),
+ total: alertRules.length,
+ page: 1,
+ size: 20,
+ });
vi.mocked(toggleAlertRule).mockImplementation(async (id, enabled) => {
const rule = alertRules.find((item) => item.id === id);
if (!rule) throw new Error(`Rule not found: ${id}`);
@@ -158,10 +163,58 @@ describe('AlertsPage', () => {
expect(within(getRuleRow('Consumer
lag')).getByRole('checkbox')).not.toBeChecked();
});
- it('keeps only failed alert rules selected after a partial bulk failure',
async () => {
- vi.mocked(listAlertRules).mockResolvedValue(
- alertRules.map((rule) => ({ ...cloneRule(rule), enabled: true })),
+ it('loads a server-side page and filters by status and search', async () => {
+ vi.mocked(listAlertRulesPage).mockClear();
+ vi.mocked(listAlertRulesPage).mockResolvedValue({
+ items: [cloneRule(alertRules[0])],
+ total: 21,
+ page: 2,
+ size: 20,
+ });
+ const user = userEvent.setup();
+ renderPage();
+
+ await screen.findByText('Broker disk usage');
+ expect(listAlertRulesPage).toHaveBeenLastCalledWith({
+ enabled: undefined,
+ page: 1,
+ pageSize: 20,
+ search: undefined,
+ });
+
+ await user.type(screen.getByPlaceholderText('搜索规则名称或指标'), 'disk');
+ await waitFor(() =>
+ expect(listAlertRulesPage).toHaveBeenLastCalledWith({
+ enabled: undefined,
+ page: 1,
+ pageSize: 20,
+ search: 'disk',
+ }),
+ );
+
+ await user.click(screen.getAllByRole('combobox')[0]);
+ await user.click(
+ await screen.findByText('已启用', { selector:
'.ant-select-item-option-content' }),
+ );
+ await waitFor(() =>
+ expect(listAlertRulesPage).toHaveBeenLastCalledWith({
+ enabled: true,
+ page: 1,
+ pageSize: 20,
+ search: 'disk',
+ }),
);
+ expect(screen.getByText('规则总数')).toBeInTheDocument();
+ expect(screen.getByText('21')).toBeInTheDocument();
+ });
+
+ it('keeps only failed alert rules selected after a partial bulk failure',
async () => {
+ vi.mocked(listAlertRulesPage).mockResolvedValue({
+ items: alertRules.map((rule) => ({ ...cloneRule(rule), enabled: true })),
+ total: alertRules.length,
+ page: 1,
+ size: 20,
+ });
vi.mocked(bulkToggleAlertRules).mockResolvedValue({
succeededIds: [1],
failures: { '2': 'network error' },
diff --git a/web/src/pages/ops/__tests__/SystemAlertsPage.test.tsx
b/web/src/pages/ops/__tests__/SystemAlertsPage.test.tsx
index 681003c75..83a73e1b1 100644
--- a/web/src/pages/ops/__tests__/SystemAlertsPage.test.tsx
+++ b/web/src/pages/ops/__tests__/SystemAlertsPage.test.tsx
@@ -10,13 +10,13 @@ import { render, screen, waitFor } from
'@testing-library/react';
import userEvent from '@testing-library/user-event';
import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest';
import { LangProvider } from '../../../i18n/LangContext';
-import { acknowledgeAlert, listSystemAlerts } from
'../../../services/opsService';
+import { acknowledgeAlert, listSystemAlertsPage } from
'../../../services/opsService';
import SystemAlertsPage from '../systemAlerts';
vi.mock('../../../services/opsService', () => ({
acknowledgeAlert: vi.fn(),
clearAcknowledgedAlerts: vi.fn(),
- listSystemAlerts: vi.fn(),
+ listSystemAlertsPage: vi.fn(),
}));
beforeAll(() => {
@@ -44,41 +44,53 @@ const renderPage = () =>
</App>,
);
+const mixedCaseAlerts = {
+ total: 2,
+ page: 1,
+ size: 20,
+ items: [
+ {
+ id: 4,
+ level: 'Error',
+ title: 'Mixed-case error',
+ description: 'error',
+ time: '2026-08-10 01:00',
+ acknowledged: false,
+ },
+ {
+ id: 5,
+ level: 'WARNING',
+ title: 'Mixed-case warning',
+ description: 'warning',
+ time: '2026-08-10 01:01',
+ acknowledged: false,
+ },
+ ],
+};
+
describe('SystemAlertsPage', () => {
beforeEach(() => {
vi.clearAllMocks();
- vi.mocked(listSystemAlerts).mockResolvedValue([
- {
- id: 1,
- level: 'error',
- title: 'Broker unavailable',
- description: 'broker a',
- time: '2026-08-10 01:00',
- acknowledged: false,
- },
- {
- id: 2,
- level: 'warning',
- title: 'Consumer lag',
- description: 'consumer b',
- time: '2026-08-10 01:01',
- acknowledged: false,
- },
- ]);
+ vi.mocked(listSystemAlertsPage).mockReset();
+ vi.mocked(listSystemAlertsPage).mockResolvedValue(mixedCaseAlerts);
});
it('renders an alert with an unknown backend level', async () => {
- vi.mocked(listSystemAlerts).mockResolvedValue([
- {
- id: 3,
- level: 'critical',
- title: 'Critical broker condition',
- description: 'A newer backend emitted this level',
- time: '2026-08-10 01:00',
- acknowledged: false,
- },
- ]);
-
+ vi.mocked(listSystemAlertsPage).mockResolvedValue({
+ total: 1,
+ page: 1,
+ size: 20,
+ items: [
+ {
+ id: 3,
+ level: 'critical',
+ title: 'Critical broker condition',
+ description: 'A newer backend emitted this level',
+ time: '2026-08-10 01:00',
+ acknowledged: false,
+ },
+ ],
+ });
renderPage();
expect(await screen.findByText('Critical broker
condition')).toBeInTheDocument();
@@ -87,36 +99,64 @@ describe('SystemAlertsPage', () => {
});
it('filters backend alert levels case-insensitively', async () => {
- vi.mocked(listSystemAlerts).mockResolvedValue([
- {
- id: 4,
- level: 'Error',
- title: 'Mixed-case error',
- description: 'error',
- time: '2026-08-10 01:00',
- acknowledged: false,
- },
- {
- id: 5,
- level: 'WARNING',
- title: 'Mixed-case warning',
- description: 'warning',
- time: '2026-08-10 01:01',
- acknowledged: false,
- },
- ]);
+ vi.mocked(listSystemAlertsPage)
+ .mockResolvedValueOnce(mixedCaseAlerts)
+ .mockResolvedValueOnce({
+ total: 1,
+ page: 1,
+ size: 20,
+ items: [
+ {
+ id: 4,
+ level: 'Error',
+ title: 'Mixed-case error',
+ description: 'error',
+ time: '2026-08-10 01:00',
+ acknowledged: false,
+ },
+ ],
+ });
const user = userEvent.setup();
renderPage();
await screen.findByText('Mixed-case error');
- await user.click(screen.getByRole('button', { name: /严重/ }));
-
+ await user.click(screen.getByRole('button', { name: /严\s*重/ }));
expect(screen.getByText('Mixed-case error')).toBeInTheDocument();
- expect(screen.queryByText('Mixed-case warning')).not.toBeInTheDocument();
+ await waitFor(() =>
+ expect(listSystemAlertsPage).toHaveBeenLastCalledWith({
+ level: 'error',
+ page: 1,
+ pageSize: 20,
+ }),
+ );
+ await waitFor(() => expect(screen.queryByText('Mixed-case
warning')).not.toBeInTheDocument());
});
it('tracks simultaneous acknowledgements independently', async () => {
vi.mocked(acknowledgeAlert).mockImplementation(() => new Promise(() =>
{}));
+ vi.mocked(listSystemAlertsPage).mockResolvedValue({
+ total: 2,
+ page: 1,
+ size: 20,
+ items: [
+ {
+ id: 1,
+ level: 'error',
+ title: 'Broker unavailable',
+ description: 'broker a',
+ time: '2026-08-10 01:00',
+ acknowledged: false,
+ },
+ {
+ id: 2,
+ level: 'warning',
+ title: 'Consumer lag',
+ description: 'consumer b',
+ time: '2026-08-10 01:01',
+ acknowledged: false,
+ },
+ ],
+ });
const user = userEvent.setup();
renderPage();
diff --git a/web/src/pages/ops/alerts.tsx b/web/src/pages/ops/alerts.tsx
index a68058c3f..2ae9e5567 100644
--- a/web/src/pages/ops/alerts.tsx
+++ b/web/src/pages/ops/alerts.tsx
@@ -15,7 +15,7 @@
* limitations under the License.
*/
-import { useEffect, useState, type Key } from 'react';
+import { useCallback, useEffect, useRef, useState, type Key } from 'react';
import { Plus, Pencil, Trash } from '@phosphor-icons/react';
import {
Button,
@@ -43,7 +43,7 @@ import {
bulkDeleteAlertRules,
bulkToggleAlertRules,
deleteAlertRule,
- listAlertRules,
+ listAlertRulesPage,
toggleAlertRule,
updateAlertRule,
} from '../../services/opsService';
@@ -66,6 +66,11 @@ const AlertsPage = () => {
const { t } = useLang();
const { token } = theme.useToken();
const [rules, setRules] = useState<AlertRule[]>([]);
+ const [totalRules, setTotalRules] = useState(0);
+ const [page, setPage] = useState(1);
+ const [pageSize, setPageSize] = useState(20);
+ const [search, setSearch] = useState('');
+ const [enabledFilter, setEnabledFilter] = useState<'all' | 'enabled' |
'disabled'>('all');
const [loading, setLoading] = useState(true);
const [modalVisible, setModalVisible] = useState(false);
const [editingRule, setEditingRule] = useState<AlertRule | null>(null);
@@ -74,6 +79,7 @@ const AlertsPage = () => {
const [selectedRuleIds, setSelectedRuleIds] = useState<Key[]>([]);
const [bulkAction, setBulkAction] = useState<'enable' | 'disable' | 'delete'
| null>(null);
const [form] = Form.useForm();
+ const requestIdRef = useRef(0);
const channelLabels: Record<string, string> = {
dingtalk: 'DingTalk',
@@ -81,24 +87,35 @@ const AlertsPage = () => {
sms: 'SMS',
};
- useEffect(() => {
- let cancelled = false;
-
- void listAlertRules()
- .then((nextRules) => {
- if (!cancelled) setRules(nextRules);
- })
- .catch(() => {
- if (!cancelled) message.error('告警规则加载失败,请稍后重试');
- })
- .finally(() => {
- if (!cancelled) setLoading(false);
+ const loadRules = useCallback(async () => {
+ const requestId = ++requestIdRef.current;
+ setLoading(true);
+ try {
+ const result = await listAlertRulesPage({
+ search: search.trim() || undefined,
+ enabled: enabledFilter === 'all' ? undefined : enabledFilter ===
'enabled',
+ page,
+ pageSize,
});
+ if (requestId !== requestIdRef.current) return;
+ setRules(result.items);
+ setTotalRules(result.total);
+ } catch {
+ if (requestId === requestIdRef.current) {
+ message.error('告警规则加载失败,请稍后重试');
+ }
+ } finally {
+ if (requestId === requestIdRef.current) setLoading(false);
+ }
+ }, [enabledFilter, page, pageSize, search]);
+ useEffect(() => {
+ const timer = window.setTimeout(() => void loadRules(), 0);
return () => {
- cancelled = true;
+ window.clearTimeout(timer);
+ requestIdRef.current += 1;
};
- }, []);
+ }, [loadRules]);
const enabledCount = rules.filter((r) => r.enabled).length;
const selectedCount = selectedRuleIds.length;
@@ -165,7 +182,6 @@ const AlertsPage = () => {
if (updatedRules.size > 0) {
setRules((previous) => previous.map((rule) =>
updatedRules.get(rule.id) ?? rule));
}
-
setSelectedRuleIds(failedIds.map(Number));
if (failedIds.length === 0) {
@@ -212,6 +228,7 @@ const AlertsPage = () => {
const succeeded = new Set(result.succeededIds);
const failedIds = Object.keys(result.failures);
setRules((previous) => previous.filter((rule) =>
!succeeded.has(rule.id)));
+ setTotalRules((previous) => Math.max(previous -
result.succeededIds.length, 0));
setSelectedRuleIds(failedIds.map(Number));
if (failedIds.length === 0)
message.success(t('alerts.bulkDeleteSuccess'));
else
@@ -364,9 +381,7 @@ const AlertsPage = () => {
<Flex gap={16}>
<Flex align="center" gap={4}>
<span style={{ fontSize: 14, color: '#999'
}}>{t('alerts.totalRules')}</span>
- <span style={{ fontSize: 18, fontWeight: 600, color: '#3b82f6'
}}>
- {rules.length}
- </span>
+ <span style={{ fontSize: 18, fontWeight: 600, color: '#3b82f6'
}}>{totalRules}</span>
</Flex>
<Flex align="center" gap={4}>
<span style={{ fontSize: 14, color: '#999'
}}>{t('alerts.enabled')}</span>
@@ -402,7 +417,30 @@ const AlertsPage = () => {
borderBottom: `1px solid ${token.colorBorderSecondary}`,
}}
>
- <span style={{ color: token.colorTextSecondary }}>
+ <Input.Search
+ allowClear
+ placeholder={t('alerts.searchPlaceholder')}
+ style={{ width: 260 }}
+ value={search}
+ onChange={(event) => {
+ setPage(1);
+ setSearch(event.target.value);
+ }}
+ />
+ <Select
+ value={enabledFilter}
+ onChange={(value) => {
+ setPage(1);
+ setEnabledFilter(value);
+ }}
+ style={{ width: 130 }}
+ options={[
+ { value: 'all', label: t('common.all') },
+ { value: 'enabled', label: t('alerts.enabled') },
+ { value: 'disabled', label: t('alerts.disabled') },
+ ]}
+ />
+ <span style={{ color: token.colorTextSecondary, marginLeft: 8 }}>
{t('alerts.selectedRules', { count: selectedCount })}
</span>
<Flex gap={8}>
@@ -440,7 +478,22 @@ const AlertsPage = () => {
size="small"
loading={loading}
rowSelection={rowSelection}
- pagination={false}
+ pagination={{
+ current: page,
+ pageSize,
+ total: totalRules,
+ showSizeChanger: true,
+ pageSizeOptions: ['20', '50', '100'],
+ showTotal: (count) => `${t('common.total')} ${count}`,
+ onChange: (nextPage, nextPageSize) => {
+ if (nextPageSize !== pageSize) {
+ setPage(1);
+ setPageSize(nextPageSize);
+ } else {
+ setPage(nextPage);
+ }
+ },
+ }}
scroll={{ x: tableScrollX(columns, { selection: true }) }}
/>
</Card>
diff --git a/web/src/pages/ops/systemAlerts.tsx
b/web/src/pages/ops/systemAlerts.tsx
index d5179d710..2482206ce 100644
--- a/web/src/pages/ops/systemAlerts.tsx
+++ b/web/src/pages/ops/systemAlerts.tsx
@@ -15,15 +15,15 @@
* limitations under the License.
*/
-import { useEffect, useState } from 'react';
-import { Card, Tag, Flex, Typography, Badge, Button, message } from 'antd';
+import { useCallback, useEffect, useRef, useState } from 'react';
+import { Card, Tag, Typography, Button, Flex, Table, message } from 'antd';
import { CheckCircle, Trash } from '@phosphor-icons/react';
import PageHeader from '../../components/PageHeader';
import { useLang } from '../../i18n/LangContext';
import {
acknowledgeAlert,
clearAcknowledgedAlerts,
- listSystemAlerts,
+ listSystemAlertsPage,
} from '../../services/opsService';
import type { SystemAlert } from '../../api/ops';
@@ -41,34 +41,43 @@ const SystemAlertsPage = () => {
};
const [alerts, setAlerts] = useState<SystemAlert[]>([]);
- const [levelFilter, setLevelFilter] = useState<string>('all');
+ const [total, setTotal] = useState(0);
+ const [page, setPage] = useState(1);
+ const [pageSize, setPageSize] = useState(20);
+ const [levelFilter, setLevelFilter] = useState<string |
undefined>(undefined);
const [loading, setLoading] = useState(true);
const [acknowledgingIds, setAcknowledgingIds] = useState<Set<number>>(() =>
new Set());
const [clearing, setClearing] = useState(false);
+ const requestIdRef = useRef(0);
- useEffect(() => {
- let cancelled = false;
-
- void listSystemAlerts()
- .then((data) => {
- if (!cancelled) setAlerts(data);
- })
- .catch(() => {
- if (!cancelled) message.error('系统告警加载失败,请稍后重试');
- })
- .finally(() => {
- if (!cancelled) setLoading(false);
+ const loadAlerts = useCallback(async () => {
+ const requestId = ++requestIdRef.current;
+ setLoading(true);
+ try {
+ const result = await listSystemAlertsPage({
+ level: levelFilter,
+ page,
+ pageSize,
});
+ if (requestId !== requestIdRef.current) return;
+ setAlerts(result.items);
+ setTotal(result.total);
+ } catch {
+ if (requestId === requestIdRef.current) message.error('系统告警加载失败,请稍后重试');
+ } finally {
+ if (requestId === requestIdRef.current) setLoading(false);
+ }
+ }, [levelFilter, page, pageSize]);
+ useEffect(() => {
+ const timer = window.setTimeout(() => void loadAlerts(), 0);
return () => {
- cancelled = true;
+ window.clearTimeout(timer);
+ requestIdRef.current += 1;
};
- }, []);
+ }, [loadAlerts]);
- const filtered =
- levelFilter === 'all'
- ? alerts
- : alerts.filter((a) => normalizeAlertLevel(a.level) === levelFilter);
+ const filtered = alerts;
const unackCount = alerts.filter((a) => !a.acknowledged).length;
@@ -93,8 +102,7 @@ const SystemAlertsPage = () => {
setClearing(true);
try {
await clearAcknowledgedAlerts();
- const fresh = await listSystemAlerts();
- setAlerts(fresh);
+ await loadAlerts();
message.success(t('sysAlerts.cleared'));
} catch {
message.error('清理已确认告警失败,请稍后重试');
@@ -121,104 +129,124 @@ const SystemAlertsPage = () => {
/>
<Flex gap={8} style={{ marginBottom: 16 }}>
- {['all', 'error', 'warning', 'info'].map((level) => (
+ {[undefined, 'error', 'warning', 'info'].map((level) => (
<Button
- key={level}
+ key={level ?? 'all'}
type={levelFilter === level ? 'primary' : 'default'}
size="small"
- onClick={() => setLevelFilter(level)}
+ onClick={() => {
+ setPage(1);
+ setLevelFilter(level);
+ }}
>
- {level === 'all' ? t('common.all') :
alertLevelConfig[level]?.label}
- {level !== 'all' && (
- <Badge
- count={alerts.filter((a) => normalizeAlertLevel(a.level) ===
level).length}
- style={{
- marginLeft: 4,
- backgroundColor:
- level === 'error' ? '#ff4d4f' : level === 'warning' ?
'#fa8c16' : '#1677ff',
- }}
- size="small"
- />
- )}
+ {level === undefined ? t('common.all') :
alertLevelConfig[level]?.label}
</Button>
))}
</Flex>
- <Flex vertical gap={12}>
- {loading && <Card loading />}
- {!loading &&
- filtered.map((alert) => {
- const normalizedLevel = normalizeAlertLevel(alert.level);
- const cfg = alertLevelConfig[normalizedLevel] ?? {
- color: '#8c8c8c',
- bg: '#fafafa',
- label: alert.level || t('common.na'),
- };
- return (
- <div
- key={alert.id}
- style={{
- display: 'flex',
- alignItems: 'flex-start',
- gap: 12,
- padding: '12px 16px',
- borderRadius: 8,
- background: cfg.bg,
- borderLeft: `3px solid ${cfg.color}`,
- opacity: alert.acknowledged ? 0.6 : 1,
- }}
- >
- <div style={{ flex: 1 }}>
- <Flex align="center" gap={8}>
- <Text strong style={{ fontSize: 14 }}>
- {alert.title}
- </Text>
- <Tag
- color={
- normalizedLevel === 'error'
- ? 'error'
- : normalizedLevel === 'warning'
- ? 'warning'
- : normalizedLevel === 'info'
- ? 'processing'
- : 'default'
- }
- style={{ fontSize: 14, lineHeight: '18px', padding: '0
6px' }}
- >
- {cfg.label}
- </Tag>
- </Flex>
- <Text type="secondary" style={{ fontSize: 14 }}>
- {alert.description}
+ <Card styles={{ body: { padding: 0 } }}>
+ <Table<SystemAlert>
+ rowKey="id"
+ loading={loading}
+ dataSource={filtered}
+ pagination={{
+ current: page,
+ pageSize,
+ total,
+ showSizeChanger: true,
+ pageSizeOptions: ['20', '50', '100'],
+ showTotal: (count) => `${t('common.total')} ${count}`,
+ onChange: (nextPage, nextPageSize) => {
+ if (nextPageSize !== pageSize) {
+ setPage(1);
+ setPageSize(nextPageSize);
+ } else {
+ setPage(nextPage);
+ }
+ },
+ }}
+ columns={[
+ {
+ title: t('sysAlerts.severe'),
+ dataIndex: 'level',
+ width: 110,
+ render: (level: string) => {
+ const normalizedLevel = normalizeAlertLevel(level);
+ const cfg = alertLevelConfig[normalizedLevel] ?? {
+ color: '#8c8c8c',
+ bg: '#fafafa',
+ label: level || t('common.na'),
+ };
+ return (
+ <Tag
+ color={
+ normalizedLevel === 'error'
+ ? 'error'
+ : normalizedLevel === 'warning'
+ ? 'warning'
+ : normalizedLevel === 'info'
+ ? 'processing'
+ : 'default'
+ }
+ style={{ fontSize: 14, lineHeight: '18px', padding: '0
6px' }}
+ >
+ {cfg.label}
+ </Tag>
+ );
+ },
+ },
+ {
+ title: t('sysAlerts.title'),
+ dataIndex: 'title',
+ render: (_: string, alert) => (
+ <>
+ <Text strong style={{ fontSize: 14 }}>
+ {alert.title}
</Text>
- </div>
- <Flex align="center" gap={8} style={{ flexShrink: 0 }}>
<Text type="secondary" style={{ fontSize: 14 }}>
- {alert.time}
+ {alert.description}
</Text>
- {!alert.acknowledged && (
- <Button
- size="small"
- type="link"
- icon={<CheckCircle size={14} />}
- onClick={() => handleAck(alert.id)}
- loading={acknowledgingIds.has(alert.id)}
- >
- {t('sysAlerts.acknowledge')}
- </Button>
- )}
- </Flex>
- </div>
- );
+ </>
+ ),
+ },
+ {
+ title: t('audit.time'),
+ dataIndex: 'time',
+ width: 180,
+ render: (time: string) => (
+ <Text type="secondary" style={{ fontSize: 14 }}>
+ {time}
+ </Text>
+ ),
+ },
+ {
+ title: t('common.actions'),
+ key: 'actions',
+ width: 130,
+ render: (_: unknown, alert: SystemAlert) =>
+ alert.acknowledged ? (
+ <Tag>{t('sysAlerts.acknowledged')}</Tag>
+ ) : (
+ <Button
+ size="small"
+ type="link"
+ icon={<CheckCircle size={14} />}
+ onClick={() => handleAck(alert.id)}
+ loading={acknowledgingIds.has(alert.id)}
+ >
+ {t('sysAlerts.acknowledge')}
+ </Button>
+ ),
+ },
+ ]}
+ onRow={(alert) => ({
+ style: {
+ background:
alertLevelConfig[normalizeAlertLevel(alert.level)]?.bg ?? '#fafafa',
+ opacity: alert.acknowledged ? 0.6 : 1,
+ },
})}
- {!loading && filtered.length === 0 && (
- <Card>
- <Flex justify="center" style={{ padding: 40 }}>
- <Text type="secondary">{t('sysAlerts.noAlerts')}</Text>
- </Flex>
- </Card>
- )}
- </Flex>
+ />
+ </Card>
</div>
);
};
diff --git a/web/src/services/opsService.ts b/web/src/services/opsService.ts
index b2d3ca069..60ccd3393 100644
--- a/web/src/services/opsService.ts
+++ b/web/src/services/opsService.ts
@@ -5,7 +5,9 @@ import * as opsApi from '../api/ops';
import type {
AlertRule,
AlertRuleBulkResult,
+ AlertRulePage,
SystemAlert,
+ SystemAlertPage,
AuditQuery,
AuditRecord,
PageResult,
@@ -101,6 +103,31 @@ export async function listAlertRules():
Promise<AlertRule[]> {
return opsApi.listAlertRules();
}
+export async function listAlertRulesPage(
+ query: { search?: string; enabled?: boolean; page?: number; pageSize?:
number } = {},
+): Promise<AlertRulePage> {
+ if (!isMockMode()) return opsApi.listAlertRulesPage(query);
+ const search = query.search?.trim().toLowerCase();
+ const rules = alertRulesState
+ .filter((rule) => !query.enabled || rule.enabled === query.enabled)
+ .filter(
+ (rule) =>
+ !search ||
+ rule.name.toLowerCase().includes(search) ||
+ rule.metric.toLowerCase().includes(search),
+ )
+ .map(copyAlertRule);
+ const page = query.page ?? 1;
+ const pageSize = query.pageSize ?? 20;
+ const from = Math.min((page - 1) * pageSize, rules.length);
+ return {
+ items: rules.slice(from, from + pageSize),
+ total: rules.length,
+ page,
+ size: pageSize,
+ };
+}
+
export async function createAlertRule(data: Partial<AlertRule>):
Promise<AlertRule> {
if (isMockMode()) {
const rule: AlertRule = {
@@ -195,6 +222,25 @@ export async function listSystemAlerts():
Promise<SystemAlert[]> {
return opsApi.listSystemAlerts();
}
+export async function listSystemAlertsPage(
+ query: { level?: string; page?: number; pageSize?: number } = {},
+): Promise<SystemAlertPage> {
+ if (!isMockMode()) return opsApi.listSystemAlertsPage(query);
+ const level = query.level?.toLowerCase();
+ const alerts = (mockSystemAlerts as unknown as SystemAlert[])
+ .filter((alert) => !level || alert.level.toLowerCase() === level)
+ .map(copySystemAlert);
+ const page = query.page ?? 1;
+ const pageSize = query.pageSize ?? 20;
+ const from = Math.min((page - 1) * pageSize, alerts.length);
+ return {
+ items: alerts.slice(from, from + pageSize),
+ total: alerts.length,
+ page,
+ size: pageSize,
+ };
+}
+
export async function acknowledgeAlert(id: number): Promise<void> {
if (isMockMode()) {
const a = mockSystemAlerts.find((a: Record<string, unknown>) => a.id ===
id);