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 b9428465b fix(audit): preserve resource types for operational audits
(#2232)
b9428465b is described below
commit b9428465b24adef16a2e8f98830c3c1b649648f7
Author: shown <[email protected]>
AuthorDate: Wed Aug 19 11:20:14 2026 +0800
fix(audit): preserve resource types for operational audits (#2232)
Signed-off-by: yuluo-yx <[email protected]>
---
.../studio/cluster/broker/ClusterService.java | 3 ++-
.../rocketmq/studio/ops/audit/AuditService.java | 6 ++++++
.../provider/apache/RocketMQAdminClientImpl.java | 12 +++++++++++-
.../apache/RocketMQBrokerConfigService.java | 3 ++-
.../provider/apache/RocketMQDLQProvider.java | 2 +-
.../studio/cluster/broker/ClusterServiceTest.java | 4 +++-
.../studio/ops/audit/AuditServiceTest.java | 12 ++++++++++++
.../apache/RocketMQAdminClientImplTest.java | 22 ++++++++++++----------
.../apache/RocketMQBrokerConfigServiceTest.java | 6 +++---
.../provider/apache/RocketMQDLQProviderTest.java | 13 ++++++++++++-
10 files changed, 64 insertions(+), 19 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
index 458be5c2f..004da7f46 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
@@ -283,7 +283,8 @@ public class ClusterService {
.map(failure -> failure.getAddress() + ": " +
failure.getMessage())
.toList();
try {
- auditService.record("UPDATE_CLUSTER_CONFIG", "CLUSTER:" +
clusterId, clusterId, detail, status.name());
+ auditService.record("UPDATE_CLUSTER_CONFIG", "CLUSTER", "CLUSTER:"
+ clusterId,
+ clusterId, detail, status.name());
} catch (Exception e) {
log.warn("Failed to record cluster config update audit for {}:
{}", clusterId, e.getMessage());
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
index 956206cff..83ebf6635 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
@@ -89,10 +89,16 @@ public class AuditService {
}
public void record(String operationType, String target, String clusterId,
String detail, String result) {
+ record(operationType, null, target, clusterId, detail, result);
+ }
+
+ public void record(String operationType, String resourceType, String
target, String clusterId,
+ String detail, String result) {
AuditRecordVO record = AuditRecordVO.builder()
.timestamp(LocalDateTime.now())
.operator(AuthenticatedUserContext.currentUsernameOrSystem())
.operationType(operationType)
+ .resourceType(resourceType)
.target(target)
.clusterId(clusterId)
.detail(detail)
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 5d9698684..2eefa9c7a 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
@@ -565,13 +565,23 @@ public class RocketMQAdminClientImpl implements
AdminClient {
private void recordAudit(String action, String resource, String detail,
String result) {
try {
- auditService.record(action, resource, detail, result);
+ auditService.record(action, auditResourceType(action), resource,
null, detail, result);
} catch (Exception auditFailure) {
log.warn("Failed to record audit action={} resource={}: {}",
action, resource,
auditFailure.getMessage());
}
}
+ private String auditResourceType(String action) {
+ if (action.endsWith("_TOPIC")) {
+ return "TOPIC";
+ }
+ if ("SEND_MESSAGE".equals(action)) {
+ return "MESSAGE";
+ }
+ return "GROUP";
+ }
+
/**
* Returns the configured default NameServer address, failing fast when
the studio has no
* RocketMQ endpoint configured (equivalent to the former absent admin
bean).
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
index 2452120da..668266c98 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigService.java
@@ -91,7 +91,8 @@ public class RocketMQBrokerConfigService {
private void recordAudit(String clusterId, String detail, String result) {
try {
- auditService.record("UPDATE_BROKER_CONFIG", "CLUSTER:" +
clusterId, clusterId, detail, result);
+ auditService.record("UPDATE_BROKER_CONFIG", "BROKER", "CLUSTER:" +
clusterId,
+ clusterId, detail, result);
} catch (Exception auditFailure) {
log.warn("Failed to record broker config audit for cluster {}:
{}", clusterId,
auditFailure.getMessage());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
index df86a1e58..7415d94dc 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
@@ -456,7 +456,7 @@ public class RocketMQDLQProvider implements DLQProvider {
private void recordAudit(String groupName, String detail, String result) {
try {
- auditService.record("RESEND_DLQ", groupName, detail, result);
+ auditService.record("RESEND_DLQ", "DLQ", groupName, null, detail,
result);
} catch (Exception e) {
log.warn("Failed to record DLQ resend audit: {}", e.getMessage());
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
index 97639b7f0..976d5e747 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
@@ -221,7 +221,7 @@ class ClusterServiceTest {
void updateConfigShouldSucceedWhenAuditRecordingFails() {
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
doThrow(new IllegalStateException("audit storage
unavailable")).when(auditService)
- .record(any(), any(), any(), any(), any());
+ .record(any(), any(), any(), any(), any(), any());
ClusterConfigUpdateResultVO result =
clusterService.updateClusterConfig(UpdateConfigDTO.builder()
.id("cluster-1")
@@ -252,6 +252,7 @@ class ClusterServiceTest {
verifyNoInteractions(brokerConfigService);
verify(auditService).record(
eq("UPDATE_CLUSTER_CONFIG"),
+ eq("CLUSTER"),
eq("CLUSTER:cluster-1"),
eq("cluster-1"),
org.mockito.ArgumentMatchers.contains("No broker address"),
@@ -342,6 +343,7 @@ class ClusterServiceTest {
verify(clusterRepository, never()).updateConfig(eq("cluster-1"),
any());
verify(auditService).record(
eq("UPDATE_CLUSTER_CONFIG"),
+ eq("CLUSTER"),
eq("CLUSTER:cluster-1"),
eq("cluster-1"),
org.mockito.ArgumentMatchers.contains("10.0.0.2:10911"),
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
index 08b8bb51d..5a777bce0 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
@@ -73,6 +73,18 @@ class AuditServiceTest {
assertThat(captor.getValue().getClusterId()).isEqualTo("prod-cn");
}
+ @Test
+ void recordShouldPreserveExplicitResourceClassification() {
+ auditService.record("RESEND_DLQ", "DLQ", "consumer-a", "instance-a",
+ "resent=3", "SUCCESS");
+
+ ArgumentCaptor<AuditRecordVO> captor =
ArgumentCaptor.forClass(AuditRecordVO.class);
+ verify(auditRepository).save(captor.capture());
+ assertThat(captor.getValue().getResourceType()).isEqualTo("DLQ");
+ assertThat(captor.getValue().getTarget()).isEqualTo("consumer-a");
+ assertThat(captor.getValue().getClusterId()).isEqualTo("instance-a");
+ }
+
@Test
void queryLogsDelegatesPaginationAndFiltersToRepository() {
AuditRecordVO record =
AuditRecordVO.builder().operationType("CREATE").build();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
index 5ffb181ad..6786c071c 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
@@ -135,7 +135,7 @@ class RocketMQAdminClientImplTest {
adminClient.resetOffset("instance-a", "cg-orders", 1784246400000L,
"orders");
verify(runtimeAdminClientResolver).execute(org.mockito.ArgumentMatchers.eq("instance-a"),
any());
- verify(auditService).record("RESET_OFFSET", "cg-orders",
+ verify(auditService).record("RESET_OFFSET", "GROUP", "cg-orders", null,
"instanceId=instance-a, topic=orders,
timestamp=1784246400000", "SUCCESS");
}
@@ -161,7 +161,7 @@ class RocketMQAdminClientImplTest {
.isInstanceOf(BusinessException.class)
.hasMessage("Instance not found: missing-instance")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(404));
- verify(auditService).record("RESET_OFFSET", "cg-orders",
+ verify(auditService).record("RESET_OFFSET", "GROUP", "cg-orders", null,
"Instance not found: missing-instance", "FAILED");
}
@@ -175,7 +175,8 @@ class RocketMQAdminClientImplTest {
.isInstanceOf(BusinessException.class)
.hasMessage("Failed to reset offset: broker unavailable")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(500));
- verify(auditService).record("RESET_OFFSET", "cg-orders", "broker
unavailable", "FAILED");
+ verify(auditService).record("RESET_OFFSET", "GROUP", "cg-orders", null,
+ "broker unavailable", "FAILED");
}
@Test
@@ -477,13 +478,14 @@ class RocketMQAdminClientImplTest {
});
doNothing().when(adminExt).createAndUpdateTopicConfig(anyString(),
any(TopicConfig.class));
doThrow(new RuntimeException("audit db down")).when(auditService)
- .record(anyString(), anyString(), anyString(), anyString());
+ .record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
TopicVO topic = new TopicVO();
topic.setName("topicA");
assertThat(adminClient.createTopic(topic).getId()).isEqualTo(1L);
- verify(auditService).record("CREATE_TOPIC", "topicA", "queues=8/8",
"SUCCESS");
+ verify(auditService).record("CREATE_TOPIC", "TOPIC", "topicA", null,
+ "queues=8/8", "SUCCESS");
}
@Test
@@ -491,7 +493,7 @@ class RocketMQAdminClientImplTest {
when(adminExt.examineBrokerClusterInfo())
.thenThrow(new IllegalStateException("broker unavailable"));
doThrow(new RuntimeException("audit db down")).when(auditService)
- .record(anyString(), anyString(), anyString(), anyString());
+ .record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
TopicVO topic = new TopicVO();
topic.setName("topicA");
@@ -538,7 +540,7 @@ class RocketMQAdminClientImplTest {
void sendMessageShouldNotFailWhenAuditRecordingFails() throws Exception {
when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
doThrow(new RuntimeException("audit db down")).when(auditService)
- .record(anyString(), anyString(), anyString(), anyString());
+ .record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
try (MockedConstruction<DefaultMQProducer> mockedProducers =
mockConstruction(DefaultMQProducer.class, (producer,
context) -> {
doNothing().when(producer).start();
@@ -576,7 +578,7 @@ class RocketMQAdminClientImplTest {
}
verifyNoInteractions(runtimeAdminClientResolver);
verify(properties, never()).getNamesrvAddr();
- verify(auditService).record("SEND_MESSAGE", "TopicA",
+ verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
"Message body size 4194306 exceeds the maximum of 4194304
bytes", "FAILED");
}
@@ -664,7 +666,7 @@ class RocketMQAdminClientImplTest {
.isInstanceOf(BusinessException.class)
.hasMessageContaining("FLUSH_DISK_TIMEOUT");
- verify(auditService).record("SEND_MESSAGE", "TopicA",
+ verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA",
null,
"Message send did not succeed: FLUSH_DISK_TIMEOUT",
"FAILED");
}
}
@@ -686,7 +688,7 @@ class RocketMQAdminClientImplTest {
.isInstanceOf(BusinessException.class)
.hasMessageContaining("null");
- verify(auditService).record("SEND_MESSAGE", "TopicA",
+ verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA",
null,
"Message send did not succeed: null", "FAILED");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
index 0c85b1bae..71a6356d7 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQBrokerConfigServiceTest.java
@@ -58,7 +58,7 @@ class RocketMQBrokerConfigServiceTest {
config.setProperty("flushDiskType", "ASYNC_FLUSH");
doNothing().when(adminExt).updateBrokerConfig("broker-a:10911",
config);
doThrow(new IllegalStateException("audit db down")).when(auditService)
- .record(anyString(), anyString(), anyString(), anyString(),
anyString());
+ .record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
brokerConfigService.updateBrokerConfig("broker-a:10911", "cluster-a",
config);
}
@@ -69,7 +69,7 @@ class RocketMQBrokerConfigServiceTest {
doThrow(new IllegalStateException("broker unavailable")).when(adminExt)
.updateBrokerConfig("broker-a:10911", config);
doThrow(new IllegalStateException("audit db down")).when(auditService)
- .record(anyString(), anyString(), anyString(), anyString(),
anyString());
+ .record(anyString(), anyString(), anyString(), any(),
anyString(), anyString());
assertThatThrownBy(() ->
brokerConfigService.updateBrokerConfig("broker-a:10911", "cluster-a", config))
.isInstanceOf(BusinessException.class)
@@ -84,7 +84,7 @@ class RocketMQBrokerConfigServiceTest {
brokerConfigService.updateBrokerConfig("broker-a:10911", "cluster-a",
config);
verify(auditService).record(
- "UPDATE_BROKER_CONFIG", "CLUSTER:cluster-a", "cluster-a",
+ "UPDATE_BROKER_CONFIG", "BROKER", "CLUSTER:cluster-a",
"cluster-a",
"brokerAddr=broker-a:10911, config={}", "SUCCESS");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
index 7fd6c8a80..c4206909d 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
@@ -63,6 +63,7 @@ import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.contains;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.lenient;
@@ -223,7 +224,7 @@ class RocketMQDLQProviderTest {
verify(mockedConsumers.constructed().get(0)).fetchSubscribeMessageQueues(dlqTopic);
assertThat(mockedProducers.constructed()).isEmpty();
}
- verify(auditService).record(eq("RESEND_DLQ"), eq("group-a"),
+ verify(auditService).record(eq("RESEND_DLQ"), eq("DLQ"),
eq("group-a"), eq(null),
contains("group=group-a, dlqTopic=%DLQ%group-a"),
eq("NO_MESSAGES"));
}
@@ -256,7 +257,9 @@ class RocketMQDLQProviderTest {
}
verify(auditService).record(
eq("RESEND_DLQ"),
+ eq("DLQ"),
eq("group-a"),
+ isNull(),
contains("matched=0, resent=0, failed=0"),
eq("NO_MESSAGES"));
verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
@@ -327,7 +330,9 @@ class RocketMQDLQProviderTest {
}
verify(auditService).record(
eq("RESEND_DLQ"),
+ eq("DLQ"),
eq("group-a"),
+ isNull(),
contains("scanFailedQueues=all"),
eq("FAILED"));
}
@@ -359,7 +364,9 @@ class RocketMQDLQProviderTest {
}
verify(auditService).record(
eq("RESEND_DLQ"),
+ eq("DLQ"),
eq("group-a"),
+ isNull(),
contains("scanFailedQueues=1"),
eq("PARTIAL"));
}
@@ -431,7 +438,9 @@ class RocketMQDLQProviderTest {
}
verify(auditService).record(
eq("RESEND_DLQ"),
+ eq("DLQ"),
eq("group-a"),
+ isNull(),
contains("matched=1, resent=1, failed=0"),
eq("SUCCESS"));
}
@@ -474,7 +483,9 @@ class RocketMQDLQProviderTest {
}
verify(auditService).record(
eq("RESEND_DLQ"),
+ eq("DLQ"),
eq("group-a"),
+ isNull(),
contains("matched=1, resent=0, failed=1"),
eq("FAILED"));
}