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 796c02e18 fix(server): harden settings writes, alert acknowledgement
and offline group grading (#4284)
796c02e18 is described below
commit 796c02e18c23dc654a6c599f52d5a41a40e4eba2
Author: 烤化の初雪 <[email protected]>
AuthorDate: Wed Sep 16 15:53:36 2026 +0800
fix(server): harden settings writes, alert acknowledgement and offline
group grading (#4284)
Incorporates five related fixes from the same author:
- #4243 acknowledge the alert state when a reminder event is acknowledged
- #4255 build general settings saves from a fresh read
- #4259 treat client-reported offline and broadcast groups as empty results
- #4261 grade never-connected and broadcast groups as offline states
- #4339 answer disabled-account logins like invalid credentials
---
.../apache/rocketmq/studio/auth/AuthService.java | 13 ++-
.../rocketmq/studio/ops/ai/LlmConfigService.java | 5 +
.../rocketmq/studio/ops/alert/AlertService.java | 8 +-
.../persistence/mapper/RmqAlertStateMapper.java | 2 +-
.../provider/apache/RocketMQAdminClientImpl.java | 10 ++
.../provider/apache/RocketMQMetadataProvider.java | 14 ++-
.../studio/auth/AuthServiceDatabaseTest.java | 129 +++++++++++++++++++++
.../studio/ops/ai/LlmConfigServiceTest.java | 27 +++++
.../studio/ops/alert/AlertServiceTest.java | 14 +++
.../alert/RmqAlertStateMapperIntegrationTest.java | 93 +++++++++++++++
.../apache/RocketMQAdminClientImplTest.java | 54 +++++++++
.../apache/RocketMQMetadataProviderTest.java | 29 +++++
web/src/pages/settings/GeneralSettingsTab.tsx | 21 +++-
.../settings/__tests__/GeneralSettingsTab.test.tsx | 49 ++++++++
14 files changed, 460 insertions(+), 8 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
index ce93012f9..dbc0b7dcc 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/auth/AuthService.java
@@ -89,6 +89,12 @@ public class AuthService {
private final PasswordHasher passwordHasher;
private final LoginRateLimiter loginRateLimiter;
+ // A well-formed PBKDF2 hash of a random preimage that is not a real
credential.
+ // Disabled accounts verify against it so an attempt costs the same as a
+ // wrong-password attempt on an enabled account while their stored hash is
never used.
+ private static final String DUMMY_PASSWORD_HASH =
+
"pbkdf2$210000$kPEa0eO/wGZsulkmR6fTEA==$pvP6IpsLryF+jTLiMRjHX+vwXcwOsT+WB8tJPMhmCpU=";
+
// Retained only for narrow unit tests that construct the legacy service
directly.
private final Map<String, AuthSession> activeTokens = new
ConcurrentHashMap<>();
@@ -352,7 +358,12 @@ public class AuthService {
RmqStudioUser user = findUserByUsername(request.getUsername())
.orElseThrow(() -> new BusinessException(401, "Invalid
username or password"));
if (!Boolean.TRUE.equals(user.getEnabled())) {
- throw new BusinessException(403, "User account is disabled");
+ // Answer exactly like a wrong password on an enabled account:
burn one dummy
+ // derivation so the response timing matches, and never touch this
account's
+ // real hash. A distinguishable answer here would let an
unauthenticated
+ // caller enumerate accounts and their enabled state.
+ passwordHasher.matches(request.getPassword(), DUMMY_PASSWORD_HASH);
+ throw new BusinessException(401, "Invalid username or password");
}
if (!passwordHasher.matches(request.getPassword(),
user.getPasswordHash())) {
throw new BusinessException(401, "Invalid username or password");
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
index 9266e09de..4525f18a0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigService.java
@@ -119,6 +119,11 @@ public class LlmConfigService {
.notifySound(current.isNotifySound())
.sessionTimeout(current.getSessionTimeout())
.requireLogin(current.isRequireLogin())
+ // The LLM form never edits notification channels; keep the
stored values
+ // instead of persisting nulls over them.
+ .dingtalkWebhook(current.getDingtalkWebhook())
+ .smsWebhook(current.getSmsWebhook())
+ .emailRecipients(current.getEmailRecipients())
.llmProvider(normalized.getProvider())
.llmEngine(normalized.getEngine())
.apiKey(persistedApiKey)
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 f456991e7..1f2aa3f4f 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
@@ -571,8 +571,12 @@ public class AlertService {
if (!alertRepository.acknowledgeAlert(alert)) {
throw new BusinessException(404, "System alert not found: " + id);
}
- if ("FIRING".equalsIgnoreCase(alert.getTransition())
- && alert.getRuleId() != null &&
hasText(alert.getFingerprint()) && alert.getTime() != null) {
+ // FIRING and REMINDER events belong to the same firing episode, so
acknowledging
+ // either must ACK the active state; the repository only honors events
whose time
+ // is not older than the current episode, which keeps stale events
inert.
+ String transition = alert.getTransition();
+ if (alert.getRuleId() != null && hasText(alert.getFingerprint()) &&
alert.getTime() != null
+ && ("FIRING".equalsIgnoreCase(transition) ||
"REMINDER".equalsIgnoreCase(transition))) {
alertStateRepository.acknowledge(new
AlertStateKey(alert.getRuleId(), alert.getFingerprint()),
alert.getTime().toInstant(ZoneOffset.UTC));
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqAlertStateMapper.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqAlertStateMapper.java
index 40290e7cf..bf41f13c3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqAlertStateMapper.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/mapper/RmqAlertStateMapper.java
@@ -33,7 +33,7 @@ public interface RmqAlertStateMapper extends
BaseMapper<RmqAlertState> {
@Update("UPDATE rmq_alert_state SET status = 'ACKED', gmt_modified =
#{now}, version = version + 1 "
+ "WHERE rule_id = #{ruleId} AND fingerprint = #{fingerprint} AND
status = 'FIRING' "
- + "AND fired_at = #{firedAt}")
+ + "AND fired_at <= #{firedAt}")
int acknowledgeFiring(@Param("ruleId") Long ruleId, @Param("fingerprint")
String fingerprint,
@Param("firedAt") LocalDateTime firedAt,
@Param("now") LocalDateTime now);
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 f1131b7c4..0b34bd39f 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
@@ -35,6 +35,7 @@ import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfi
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.studio.common.util.MqResponseCodes;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
@@ -239,7 +240,16 @@ public class RocketMQAdminClientImpl implements
AdminClient {
&& brokerException.getResponseCode() ==
ResponseCode.CONSUMER_NOT_ONLINE) {
return true;
}
+ // rocketmq-tools locates a group through the %RETRY%<group> topic
route before any
+ // broker call, so a group that never connected fails with
TOPIC_NOT_EXIST for that
+ // retry topic, and a broadcast group with an empty offset table fails
with
+ // BROADCAST_CONSUMPTION. Both mean "no live data", not a lookup
failure.
String message = exception.getMessage();
+ if (MqResponseCodes.hasResponseCode(exception,
ResponseCode.BROADCAST_CONSUMPTION)
+ || MqResponseCodes.hasResponseCode(exception,
ResponseCode.TOPIC_NOT_EXIST)
+ && message != null && message.contains("%RETRY%")) {
+ return true;
+ }
return message != null && (message.contains("not online") ||
message.contains("CODE: 206"));
}
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 9d7ea8cf8..97fbbd5b1 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
@@ -744,8 +744,20 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return brokerException.getResponseCode()
==
org.apache.rocketmq.remoting.protocol.ResponseCode.CONSUMER_NOT_ONLINE;
}
+ // rocketmq-tools grades these business states as MQClientException,
so the typed code
+ // only survives in the message text: examineConsumerConnectionInfo
throws
+ // CONSUMER_NOT_ONLINE for a group whose clients all disconnected, and
+ // examineConsumeStats throws BROADCAST_CONSUMPTION for a broadcast
group with an
+ // empty offset table. Both mean "no live data", not a connectivity
failure
+ // (same grading as RocketMQClientProvider.isGroupConnectionAbsent).
+ if (MqResponseCodes.hasResponseCode(e,
+
org.apache.rocketmq.remoting.protocol.ResponseCode.CONSUMER_NOT_ONLINE,
+
org.apache.rocketmq.remoting.protocol.ResponseCode.BROADCAST_CONSUMPTION)) {
+ return true;
+ }
String message = e.getMessage();
- return message != null && message.contains("not online");
+ return message != null && (message.contains("not online")
+ || message.contains("Not found the consumer group
connection"));
}
// ── Helper methods ──────────────────────────────────────────────────
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
index d359046fe..965e304ca 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/auth/AuthServiceDatabaseTest.java
@@ -42,7 +42,11 @@ import java.util.Map;
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.anyString;
+import static org.mockito.ArgumentMatchers.argThat;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -422,6 +426,131 @@ class AuthServiceDatabaseTest {
.hasMessageStartingWith("Too many failed login attempts");
}
+ @Test
+ void
disabledAccountsWithWrongPasswordsGetTheUniformInvalidCredentialsResponse() {
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+ when(userMapper.selectOne(any(Wrapper.class)))
+ .thenReturn(user(1L, "disabled-user", false, false,
"password-1"));
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("disabled-user");
+ request.setPassword("totally-wrong");
+
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(401))
+ .hasMessage("Invalid username or password");
+ }
+
+ @Test
+ void
disabledAccountsWithCorrectPasswordsGetTheUniformInvalidCredentialsResponse() {
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+ when(userMapper.selectOne(any(Wrapper.class)))
+ .thenReturn(user(1L, "disabled-user", false, false,
"password-1"));
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("disabled-user");
+ request.setPassword("password-1");
+
+ // Deliberately uniform with unknown users and wrong passwords: the
previous
+ // 403 "User account is disabled" revealed the account state before any
+ // credential check. No session may be created either.
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(401))
+ .hasMessage("Invalid username or password");
+ verify(sessionMapper, never()).insert(any(RmqStudioSession.class));
+ }
+
+ @Test
+ void disabledAccountVerificationNeverTouchesTheAccountPasswordHash() {
+ PasswordHasher hasherSpy = mock(PasswordHasher.class);
+ SettingsRepository repository = mock(SettingsRepository.class);
+ when(repository.loadGeneralSettings())
+
.thenReturn(GeneralSettingsVO.builder().sessionTimeout(30).build());
+ authService = new AuthService(new AuthProperties(), repository,
+ Clock.fixed(Instant.parse("2026-08-13T00:00:00Z"),
ZoneOffset.UTC), userMapper,
+ sessionMapper, hasherSpy);
+ String storedHash =
"pbkdf2$210000$cmVhbFNhbHQxNjJ5dGVzRQ==$cmVhbERpZ2VzdEJhc2U2NDQ0NDQ9";
+ RmqStudioUser disabledUser = user(1L, "disabled-user", false, false,
"password-1");
+ disabledUser.setPasswordHash(storedHash);
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+
when(userMapper.selectOne(any(Wrapper.class))).thenReturn(disabledUser);
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("disabled-user");
+ request.setPassword("password-1");
+
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(401))
+ .hasMessage("Invalid username or password");
+ verify(hasherSpy, never()).matches(anyString(), eq(storedHash));
+ // The dummy verification must be a full-cost PBKDF2 hash so the
response
+ // timing matches a wrong-password attempt on an enabled account.
+ verify(hasherSpy, times(1)).matches(anyString(), argThat(hash ->
+ hash != null && hash.startsWith("pbkdf2$210000$")));
+ }
+
+ @Test
+ void unknownUsersShareTheUniformInvalidCredentialsResponse() {
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+ when(userMapper.selectOne(any(Wrapper.class))).thenReturn(null);
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("no-such-user");
+ request.setPassword("password-1");
+
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(401))
+ .hasMessage("Invalid username or password");
+ }
+
+ @Test
+ void
enabledAccountsWithWrongPasswordsGetTheUniformInvalidCredentialsResponse() {
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+ when(userMapper.selectOne(any(Wrapper.class)))
+ .thenReturn(user(1L, "operator", false, true, "password-1"));
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("operator");
+ request.setPassword("totally-wrong");
+
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(401))
+ .hasMessage("Invalid username or password");
+ }
+
+ @Test
+ void disabledAccountFailuresAreRateLimitedLikeOtherFailedLogins() {
+ when(userMapper.selectCount(isNull())).thenReturn(1L);
+ when(userMapper.selectOne(any(Wrapper.class)))
+ .thenReturn(user(1L, "disabled-user", false, false,
"password-1"));
+
+ LoginDTO request = new LoginDTO();
+ request.setUsername("disabled-user");
+ request.setPassword("totally-wrong");
+
+ for (int attempt = 0; attempt < LoginRateLimiter.MAX_FAILED_ATTEMPTS;
attempt++) {
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Invalid username or password");
+ }
+
+ assertThatThrownBy(() -> authService.login(request))
+ .isInstanceOf(BusinessException.class)
+ .satisfies(exception ->
+ assertThat(((BusinessException)
exception).getCode()).isEqualTo(429))
+ .hasMessageStartingWith("Too many failed login attempts");
+ }
+
private RmqStudioSession activeSession(Long id, Long userId, LocalDateTime
lastSeenAt) {
RmqStudioSession session = new RmqStudioSession();
session.setId(id);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
index a927c1516..04064a476 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigServiceTest.java
@@ -330,6 +330,33 @@ class LlmConfigServiceTest {
assertThat(captor.getValue().getBaseUrl()).isEqualTo("https://api.openai.com/v1");
}
+ @Test
+ void saveConfigShouldPreserveNotificationChannelFields() {
+
when(settingsService.getGeneralSettings()).thenReturn(GeneralSettingsVO.builder()
+ .theme("dark")
+
.dingtalkWebhook("https://oapi.dingtalk.com/robot/send?access_token=abc")
+ .smsWebhook("https://sms.example.com/notify")
+ .emailRecipients("[email protected],[email protected]")
+ .build());
+ LlmConfigVO config = LlmConfigVO.builder()
+ .provider("deepseek")
+ .apiKey("sk-deepseek")
+ .apiBase("https://api.deepseek.com/v1")
+ .model("deepseek-chat")
+ .enabled(true)
+ .build();
+
+ llmConfigService.saveConfig(config);
+
+ ArgumentCaptor<GeneralSettingsVO> captor =
ArgumentCaptor.forClass(GeneralSettingsVO.class);
+ verify(settingsService).saveGeneralSettings(captor.capture());
+ assertThat(captor.getValue().getDingtalkWebhook())
+
.isEqualTo("https://oapi.dingtalk.com/robot/send?access_token=abc");
+
assertThat(captor.getValue().getSmsWebhook()).isEqualTo("https://sms.example.com/notify");
+ assertThat(captor.getValue().getEmailRecipients())
+ .isEqualTo("[email protected],[email protected]");
+ }
+
@Test
void saveConfigShouldRejectInvalidApiBase() {
assertThatThrownBy(() ->
llmConfigService.saveConfig(LlmConfigVO.builder()
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 ebde0b84f..46c95023b 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
@@ -1388,6 +1388,20 @@ class AlertServiceTest {
LocalDateTime.of(2026, 8, 22, 12,
0).toInstant(ZoneOffset.UTC));
}
+ @Test
+ void acknowledgingReminderEventShouldAcknowledgeItsActiveRuleStateTest() {
+ SystemAlertVO alert =
SystemAlertVO.builder().id(1L).ruleId(7L).fingerprint("fingerprint")
+ .time(LocalDateTime.of(2026, 8, 22, 12, 30))
+ .transition("REMINDER").acknowledged(false).build();
+ when(alertRepository.findAlertById(1L)).thenReturn(Optional.of(alert));
+
when(alertRepository.acknowledgeAlert(any(SystemAlertVO.class))).thenReturn(true);
+
+ alertService.acknowledgeAlert(1L);
+
+ verify(alertStateRepository).acknowledge(new AlertStateKey(7L,
"fingerprint"),
+ LocalDateTime.of(2026, 8, 22, 12,
30).toInstant(ZoneOffset.UTC));
+ }
+
@Test
void acknowledgingResolvedEventMustNotAcknowledgeANewerFiringStateTest() {
SystemAlertVO resolved =
SystemAlertVO.builder().id(1L).ruleId(7L).fingerprint("fingerprint")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/RmqAlertStateMapperIntegrationTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/RmqAlertStateMapperIntegrationTest.java
new file mode 100644
index 000000000..881dd388a
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/RmqAlertStateMapperIntegrationTest.java
@@ -0,0 +1,93 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.ops.alert;
+
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import org.apache.rocketmq.studio.persistence.entity.RmqAlertState;
+import org.apache.rocketmq.studio.persistence.mapper.RmqAlertStateMapper;
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+
+import java.time.LocalDateTime;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@SpringBootTest(properties = {"studio.auth.login-required=false"})
+class RmqAlertStateMapperIntegrationTest {
+ private static final long RULE_ID_BASE = 2749000L;
+
+ @Autowired
+ private RmqAlertStateMapper mapper;
+
+ @Test
+ void acknowledgeFiringShouldAcceptReminderTimesOfTheCurrentEpisodeTest() {
+ long ruleId = RULE_ID_BASE + 1;
+ LocalDateTime firedAt = LocalDateTime.of(2026, 9, 1, 12, 0);
+ insertFiringState(ruleId, firedAt);
+ try {
+ LocalDateTime reminderTime = firedAt.plusMinutes(30);
+
+ int updated = mapper.acknowledgeFiring(ruleId, "fingerprint",
reminderTime,
+ LocalDateTime.now());
+
+ assertThat(updated).isEqualTo(1);
+ assertThat(mapper.selectList(new QueryWrapper<RmqAlertState>()
+ .eq("rule_id",
ruleId)).stream().map(RmqAlertState::getStatus))
+ .containsExactly("ACKED");
+ } finally {
+ cleanup(ruleId);
+ }
+ }
+
+ @Test
+ void acknowledgeFiringShouldRejectEventsOlderThanTheCurrentEpisodeTest() {
+ long ruleId = RULE_ID_BASE + 2;
+ LocalDateTime firedAt = LocalDateTime.of(2026, 9, 1, 12, 0);
+ insertFiringState(ruleId, firedAt);
+ try {
+ LocalDateTime staleEpisodeEventTime = firedAt.minusHours(1);
+
+ int updated = mapper.acknowledgeFiring(ruleId, "fingerprint",
staleEpisodeEventTime,
+ LocalDateTime.now());
+
+ assertThat(updated).isEqualTo(0);
+ assertThat(mapper.selectList(new QueryWrapper<RmqAlertState>()
+ .eq("rule_id",
ruleId)).stream().map(RmqAlertState::getStatus))
+ .containsExactly("FIRING");
+ } finally {
+ cleanup(ruleId);
+ }
+ }
+
+ private void insertFiringState(long ruleId, LocalDateTime firedAt) {
+ cleanup(ruleId);
+ RmqAlertState state = new RmqAlertState();
+ state.setRuleId(ruleId);
+ state.setFingerprint("fingerprint");
+ state.setStatus("FIRING");
+ state.setConsecutiveHits(3);
+ state.setCurrentValue(90.0);
+ state.setFiredAt(firedAt);
+ state.setVersion(0);
+ mapper.insert(state);
+ }
+
+ private void cleanup(long ruleId) {
+ mapper.delete(new QueryWrapper<RmqAlertState>().eq("rule_id", ruleId));
+ }
+}
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 487cc16ac..0a0fde32e 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
@@ -184,6 +184,21 @@ class RocketMQAdminClientImplTest {
assertThat(group.getOnlineInstances()).isZero();
}
+ @Test
+ void getConsumerGroupReturnsOfflineDetailForGroupWithoutRetryRouteTest()
throws Exception {
+ // rocketmq-tools examineConsumerConnectionInfo locates the group
through the
+ // %RETRY%<group> route, so a group created but never connected fails
with
+ // TOPIC_NOT_EXIST for that retry topic before any broker is contacted.
+ when(adminExt.examineConsumerConnectionInfo("orders")).thenThrow(
+ new MQClientException(ResponseCode.TOPIC_NOT_EXIST,
+ "No topic route info in name server for the topic:
%RETRY%orders"));
+
+ ConsumerGroupVO group = adminClient.getConsumerGroup(null, "orders");
+
+ assertThat(group.getName()).isEqualTo("orders");
+ assertThat(group.getOnlineInstances()).isZero();
+ }
+
@Test
void
getConsumerGroupFillsProxySideConnectionsWhenBrokerReportsOfflineTest() throws
Exception {
when(adminExt.examineConsumerConnectionInfo("orders"))
@@ -437,6 +452,45 @@ class RocketMQAdminClientImplTest {
anyLong(), anyBoolean());
}
+ @Test
+ void
previewResetOffsetShouldReturnEmptyPreviewForGroupWithoutRetryRouteTest()
throws Exception {
+ when(adminExt.examineConsumeStats("cg-orders")).thenThrow(
+ new MQClientException(ResponseCode.TOPIC_NOT_EXIST,
+ "No topic route info in name server for the topic:
%RETRY%cg-orders"));
+
+ ResetConsumerOffsetPreviewVO preview = adminClient.previewResetOffset(
+ null, "cg-orders", 1784246400000L, "orders");
+
+ assertThat(preview.isComplete()).isFalse();
+ assertThat(preview.isAllowReset()).isFalse();
+ assertThat(preview.getQueueCount()).isZero();
+ assertThat(preview.getWarnings())
+ .containsExactly("Consumer group is not online and no consume
offset data is available");
+ verify(adminExt, never()).resetOffsetByTimestamp(anyString(),
anyString(), anyString(),
+ anyLong(), anyBoolean());
+ }
+
+ @Test
+ void previewResetOffsetShouldReturnEmptyPreviewForBroadcastGroupTest()
throws Exception {
+ // rocketmq-tools examineConsumeStats grades a broadcast group's empty
offset table as
+ // MQClientException(BROADCAST_CONSUMPTION); the code only survives in
the message text.
+ when(adminExt.examineConsumeStats("cg-orders")).thenThrow(
+ new MQClientException(ResponseCode.BROADCAST_CONSUMPTION,
+ "Not found the consumer group consume stats, because
return offset table is empty, "
+ + "the consumer is under the broadcast mode"));
+
+ ResetConsumerOffsetPreviewVO preview = adminClient.previewResetOffset(
+ null, "cg-orders", 1784246400000L, "orders");
+
+ assertThat(preview.isComplete()).isFalse();
+ assertThat(preview.isAllowReset()).isFalse();
+ assertThat(preview.getQueueCount()).isZero();
+ assertThat(preview.getWarnings())
+ .containsExactly("Consumer group is not online and no consume
offset data is available");
+ verify(adminExt, never()).resetOffsetByTimestamp(anyString(),
anyString(), anyString(),
+ anyLong(), anyBoolean());
+ }
+
@Test
void previewResetOffsetShouldRejectBlankTopicBeforeResolvingAdmin() {
assertThatThrownBy(() -> adminClient.previewResetOffset("instance-a",
"cg-orders", 1784246400000L, " "))
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 d71dc5276..d25b2c81d 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
@@ -563,6 +563,21 @@ class RocketMQMetadataProviderTest {
assertThat(newLiveProvider(admin).getGroupProgress(null,
"group-offline")).isEmpty();
}
+ @Test
+ void getGroupProgressShouldReturnEmptyForBroadcastGroupTest() throws
Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ // rocketmq-tools examineConsumeStats throws
MQClientException(BROADCAST_CONSUMPTION)
+ // for a broadcast group with an empty offset table; the code only
survives in the
+ // message text ("CODE: 213 DESC: ... the consumer is under the
broadcast mode").
+ when(admin.examineConsumeStats("group-broadcast")).thenThrow(
+ new org.apache.rocketmq.client.exception.MQClientException(
+
org.apache.rocketmq.remoting.protocol.ResponseCode.BROADCAST_CONSUMPTION,
+ "Not found the consumer group consume stats, because
return offset table is empty, "
+ + "the consumer is under the broadcast mode"));
+
+ assertThat(newLiveProvider(admin).getGroupProgress(null,
"group-broadcast")).isEmpty();
+ }
+
@Test
void getGroupSubscriptionsSurfacesAdminFailure() throws Exception {
DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
@@ -635,6 +650,20 @@ class RocketMQMetadataProviderTest {
assertThat(newLiveProvider(admin).getGroupSubscriptions(null,
"group-proxy")).isEmpty();
}
+ @Test
+ void
getGroupSubscriptionsShouldReturnEmptyWhenClientReportsGroupOfflineTest()
throws Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ // rocketmq-tools examineConsumerConnectionInfo throws
MQClientException(CONSUMER_NOT_ONLINE)
+ // when the broker returns an empty connection set; the typed code
only survives in the
+ // message text ("CODE: 206 DESC: Not found the consumer group
connection ...").
+ when(admin.examineConsumerConnectionInfo("group-offline")).thenThrow(
+ new org.apache.rocketmq.client.exception.MQClientException(
+
org.apache.rocketmq.remoting.protocol.ResponseCode.CONSUMER_NOT_ONLINE,
+ "Not found the consumer group connection"));
+
+ assertThat(newLiveProvider(admin).getGroupSubscriptions(null,
"group-offline")).isEmpty();
+ }
+
@Test
void getGroupSubscriptionsShouldFallBackToProxyConnectionsTest() throws
Exception {
DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
diff --git a/web/src/pages/settings/GeneralSettingsTab.tsx
b/web/src/pages/settings/GeneralSettingsTab.tsx
index 7e347f92f..bfc204cad 100644
--- a/web/src/pages/settings/GeneralSettingsTab.tsx
+++ b/web/src/pages/settings/GeneralSettingsTab.tsx
@@ -99,13 +99,27 @@ export const GeneralSettingsTab = () => {
};
}, [message, notifyForm, securityForm, t]);
+ // Other tabs (AI assistant settings) write the same settings record while
this tab stays
+ // mounted, so every save must be built from a fresh read instead of the
mount-time snapshot.
+ const loadFreshSettings = async (fallback: GeneralSettings):
Promise<GeneralSettings> => {
+ try {
+ const fresh = await getGeneralSettings();
+ setSettings(fresh);
+ return fresh;
+ } catch {
+ return fallback;
+ }
+ };
+
const persistPreference = async (patch: Partial<GeneralSettings>) => {
if (!settings) return;
const next = { ...settings, ...patch };
setSettings(next);
setSavingPreference(true);
try {
- await saveGeneralSettings(buildPayload(next));
+ const base = await loadFreshSettings(settings);
+ await saveGeneralSettings(buildPayload({ ...base, ...patch }));
+ setSettings({ ...base, ...patch });
message.success(t('settings.saveSuccess'));
} catch {
message.error(t('settings.saveFailed'));
@@ -117,10 +131,11 @@ export const GeneralSettingsTab = () => {
const mergeAndSave = async (patch: Partial<GeneralSettings>) => {
if (!settings) return false;
try {
- await saveGeneralSettings(buildPayload({ ...settings, ...patch }));
+ const base = await loadFreshSettings(settings);
+ await saveGeneralSettings(buildPayload({ ...base, ...patch }));
const statePatch = { ...patch };
delete statePatch.clearDingtalkSigningSecret;
- setSettings({ ...settings, ...statePatch });
+ setSettings({ ...base, ...statePatch });
return true;
} catch {
message.error(t('settings.saveFailed'));
diff --git a/web/src/pages/settings/__tests__/GeneralSettingsTab.test.tsx
b/web/src/pages/settings/__tests__/GeneralSettingsTab.test.tsx
index 7ca24f249..6faa0e9ef 100644
--- a/web/src/pages/settings/__tests__/GeneralSettingsTab.test.tsx
+++ b/web/src/pages/settings/__tests__/GeneralSettingsTab.test.tsx
@@ -149,6 +149,55 @@ describe('GeneralSettingsTab', () => {
),
);
expect(localStorage.getItem('rocketmq-studio-theme')).toBe('dark');
+ // The fresh read that guards unmanaged fields must not revert the visible
preference.
+ await waitFor(() =>
+
expect(screen.getByText('深色').closest('.ant-segmented-item')).toHaveClass(
+ 'ant-segmented-item-selected',
+ ),
+ );
+ });
+
+ it('saves llm fields from a fresh read instead of the mount-time snapshot',
async () => {
+ const baseSettings = {
+ theme: 'system',
+ compact: false,
+ desktopNotify: false,
+ notifySound: false,
+ sessionTimeout: 30,
+ requireLogin: true,
+ apiKeyConfigured: false,
+ };
+ vi.mocked(getGeneralSettings)
+ .mockResolvedValueOnce({
+ ...baseSettings,
+ llmProvider: 'openai',
+ model: 'gpt-test',
+ baseUrl: 'https://openai.example/v1',
+ })
+ .mockResolvedValue({
+ ...baseSettings,
+ llmProvider: 'deepseek',
+ model: 'deepseek-chat',
+ baseUrl: 'https://deepseek.example/v1',
+ });
+ vi.mocked(saveGeneralSettings).mockResolvedValue();
+ renderTab();
+
+ const saveButtons = await screen.findAllByRole('button', { name: '保存设置' });
+ await waitFor(() => expect(saveButtons[0]).toBeEnabled());
+ // The AI assistant tab saved a different provider while this tab stayed
mounted;
+ // submitting the security form must not resurrect the stale provider.
+ fireEvent.submit(saveButtons[0].closest('form')!);
+
+ await waitFor(() => expect(saveGeneralSettings).toHaveBeenCalledTimes(1));
+ expect(saveGeneralSettings).toHaveBeenCalledWith(
+ expect.objectContaining({
+ llmProvider: 'deepseek',
+ model: 'deepseek-chat',
+ baseUrl: 'https://deepseek.example/v1',
+ sessionTimeout: 30,
+ }),
+ );
});
it('can explicitly clear a configured DingTalk signing secret', async () => {