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 248e3869 fix(cluster): keep write/read queue counts in sync on partial 
updates (#1075)
248e3869 is described below

commit 248e3869565ea8a2611ad9d92a6df8e884155ef6
Author: yyqdbngt <[email protected]>
AuthorDate: Thu Aug 6 16:18:03 2026 +0800

    fix(cluster): keep write/read queue counts in sync on partial updates 
(#1075)
    
    Co-authored-by: yyqdbngt <[email protected]>
---
 .../rocketmq/studio/cluster/broker/ClusterService.java   | 13 ++++++++-----
 .../studio/cluster/broker/ClusterServiceTest.java        | 16 ++++++++++++++++
 2 files changed, 24 insertions(+), 5 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 16a8e155..a8099614 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
@@ -154,11 +154,14 @@ public class ClusterService {
         if (command.getFileReservedTime() != null) {
             config.setFileReservedTime(command.getFileReservedTime());
         }
-        if (command.getWriteQueueNums() != null) {
-            config.setWriteQueueNums(command.getWriteQueueNums());
-        }
-        if (command.getReadQueueNums() != null) {
-            config.setReadQueueNums(command.getReadQueueNums());
+        // The broker exposes a single defaultTopicQueueNums property, so a 
partial update of only
+        // one of the write/read queues must mirror the value onto the other 
to keep the stored
+        // config consistent with what the broker actually applies.
+        if (command.getWriteQueueNums() != null || command.getReadQueueNums() 
!= null) {
+            int queueNums = command.getWriteQueueNums() != null
+                    ? command.getWriteQueueNums() : command.getReadQueueNums();
+            config.setWriteQueueNums(queueNums);
+            config.setReadQueueNums(queueNums);
         }
         if (command.getBrokerPermission() != null) {
             config.setBrokerPermission(command.getBrokerPermission());
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 f9ecfe66..6100d4a9 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
@@ -558,6 +558,22 @@ class ClusterServiceTest {
                 .satisfies(ex -> assertThat(((BusinessException) 
ex).getCode()).isEqualTo(404));
     }
 
+    @Test
+    void updateClusterConfigShouldMirrorSingleQueueNumsUpdate() {
+        
when(clusterProvider.refreshClusterDetail("cluster-1")).thenReturn(sampleCluster);
+        UpdateConfigDTO command = UpdateConfigDTO.builder()
+                .id("cluster-1")
+                .writeQueueNums(12)
+                .build();
+
+        ClusterConfigUpdateResultVO result = 
clusterService.updateClusterConfig(command);
+
+        // The broker exposes a single defaultTopicQueueNums property, so the 
stored write and read
+        // queue counts must stay in sync after a partial update.
+        
assertThat(result.getCluster().getConfig().getWriteQueueNums()).isEqualTo(12);
+        
assertThat(result.getCluster().getConfig().getReadQueueNums()).isEqualTo(12);
+    }
+
     @Test
     void updateClusterConfigShouldUseLiveClusterWhenItIsNotPersisted() {
         
when(clusterProvider.refreshClusterDetail("cluster-1")).thenReturn(sampleCluster);

Reply via email to