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 19b4067d5 fix(instance): release the pooled endpoint client after 
commit when deleting an instance (#3127)
19b4067d5 is described below

commit 19b4067d5305ce4b41c057888e102c34fdef6c60
Author: cyberslack_lee <[email protected]>
AuthorDate: Mon Sep 7 20:35:04 2026 +0800

    fix(instance): release the pooled endpoint client after commit when 
deleting an instance (#3127)
    
    fix instanceservice bug
---
 .../rocketmq/studio/instance/InstanceService.java  | 17 +++++++++++-
 .../studio/instance/InstanceServiceTest.java       | 32 ++++++++++++++++++++++
 2 files changed, 48 insertions(+), 1 deletion(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java 
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
index f8de5b996..c5b78c10f 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
@@ -42,6 +42,8 @@ import lombok.extern.slf4j.Slf4j;
 import org.springframework.dao.DataIntegrityViolationException;
 import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
+import org.springframework.transaction.support.TransactionSynchronization;
+import 
org.springframework.transaction.support.TransactionSynchronizationManager;
 
 import java.time.LocalDateTime;
 import java.util.ArrayList;
@@ -625,9 +627,22 @@ public class InstanceService {
             throw new BusinessException(404, "InstanceVO not found: " + id);
         }
         removeDataSourceBindings(existing.getName());
-        releaseApacheEndpointIfUnused(existing, null);
         recordAudit("DELETE_INSTANCE", "INSTANCE", String.valueOf(id), null,
                 instanceAuditDetail(existing));
+        releaseApacheEndpointAfterCommit(existing);
+    }
+
+    private void releaseApacheEndpointAfterCommit(InstanceVO existing) {
+        if (!TransactionSynchronizationManager.isSynchronizationActive()) {
+            releaseApacheEndpointIfUnused(existing, null);
+            return;
+        }
+        TransactionSynchronizationManager.registerSynchronization(new 
TransactionSynchronization() {
+            @Override
+            public void afterCommit() {
+                releaseApacheEndpointIfUnused(existing, null);
+            }
+        });
     }
 
     /**
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
index 6ef5a5277..54d3d8f37 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
@@ -34,6 +34,8 @@ import 
org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
 import org.apache.rocketmq.studio.provider.InstanceProvider;
 import org.apache.rocketmq.studio.settings.DataSourceVO;
 import org.apache.rocketmq.studio.settings.SettingsRepository;
+import org.springframework.transaction.support.TransactionSynchronization;
+import 
org.springframework.transaction.support.TransactionSynchronizationManager;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.ArgumentCaptor;
@@ -997,6 +999,36 @@ class InstanceServiceTest {
         verify(adminFactory).release("namesrv:9876");
     }
 
+    @Test
+    void deleteInstanceShouldDeferEndpointReleaseUntilAfterCommitTest() {
+        InstanceVO existing = InstanceVO.builder()
+                .name("to-delete")
+                .endpoint("namesrv:9876")
+                .build();
+        existing.setId(1L);
+        
when(instanceRepository.findById(1L)).thenReturn(Optional.of(existing));
+        when(instanceRepository.findAll()).thenReturn(List.of());
+        
when(providerRegistry.forVendor(InstanceVendor.APACHE)).thenReturn(instanceProvider);
+        when(instanceRepository.deleteById(1L)).thenReturn(true);
+
+        TransactionSynchronizationManager.initSynchronization();
+        try {
+            instanceService.deleteInstance(1L);
+
+            verify(adminFactory, never()).release(any());
+            verify(clientPool, never()).release(any());
+
+            for (TransactionSynchronization sync : 
TransactionSynchronizationManager.getSynchronizations()) {
+                sync.afterCommit();
+            }
+        } finally {
+            TransactionSynchronizationManager.clearSynchronization();
+        }
+
+        verify(adminFactory).release("namesrv:9876");
+        verify(clientPool).release("namesrv:9876");
+    }
+
     @Test
     void deleteInstanceShouldRemoveMetricsDataSourceBindingTest() {
         InstanceVO existing = InstanceVO.builder().name("to-delete").build();

Reply via email to