This is an automated email from the ASF dual-hosted git repository.

Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git


The following commit(s) were added to refs/heads/master by this push:
     new cd8e0c2c8c fix(admin): publish batch upstream changes after commit 
(#7242)
cd8e0c2c8c is described below

commit cd8e0c2c8c6d3a4ff71672ae528b84155ce654bc
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:59:14 2026 +0800

    fix(admin): publish batch upstream changes after commit (#7242)
---
 .../service/impl/DiscoveryUpstreamServiceImpl.java | 13 +++++++-
 .../service/DiscoveryUpstreamServiceTest.java      | 39 ++++++++++++++++++++++
 2 files changed, 51 insertions(+), 1 deletion(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/DiscoveryUpstreamServiceImpl.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/DiscoveryUpstreamServiceImpl.java
index 24675f44a4..2c7e154e91 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/DiscoveryUpstreamServiceImpl.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/DiscoveryUpstreamServiceImpl.java
@@ -46,6 +46,8 @@ import org.apache.shenyu.common.dto.DiscoverySyncData;
 import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
 import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
+import org.springframework.transaction.support.TransactionSynchronization;
+import 
org.springframework.transaction.support.TransactionSynchronizationManager;
 import org.springframework.util.StringUtils;
 
 import java.util.Collections;
@@ -136,7 +138,16 @@ public class DiscoveryUpstreamServiceImpl implements 
DiscoveryUpstreamService {
             discoveryUpstreamDO.setDiscoveryHandlerId(discoveryHandlerId);
             discoveryUpstreamMapper.insert(discoveryUpstreamDO);
         }
-        this.fetchAll(discoveryHandlerId);
+        if (TransactionSynchronizationManager.isSynchronizationActive()) {
+            TransactionSynchronizationManager.registerSynchronization(new 
TransactionSynchronization() {
+                @Override
+                public void afterCommit() {
+                    fetchAll(discoveryHandlerId);
+                }
+            });
+        } else {
+            this.fetchAll(discoveryHandlerId);
+        }
         return 0;
     }
 
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/DiscoveryUpstreamServiceTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/DiscoveryUpstreamServiceTest.java
index 7e187e1963..95615e10a0 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/DiscoveryUpstreamServiceTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/DiscoveryUpstreamServiceTest.java
@@ -17,6 +17,9 @@
 
 package org.apache.shenyu.admin.service;
 
+import org.springframework.transaction.support.TransactionSynchronization;
+import 
org.springframework.transaction.support.TransactionSynchronizationManager;
+
 import org.apache.shenyu.admin.discovery.DiscoveryProcessor;
 import org.apache.shenyu.admin.discovery.DiscoveryProcessorHolder;
 import org.apache.shenyu.admin.mapper.DiscoveryHandlerMapper;
@@ -59,6 +62,9 @@ import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.BDDMockito.given;
 import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.never;
 
 /**
  * Test cases for DiscoveryUpstreamService.
@@ -229,6 +235,39 @@ public final class DiscoveryUpstreamServiceTest {
         when(discoveryMapper.selectById(any())).thenReturn(buildDiscoveryDO());
         
when(discoveryUpstreamMapper.deleteByDiscoveryHandlerId(anyString())).thenReturn(0);
         discoveryUpstreamService.updateBatch("123", 
Collections.singletonList(buildDiscoveryUpstreamDTO("")));
+        verify(discoveryProcessor).changeUpstream(any(), any());
+    }
+
+    @Test
+    public void testUpdateBatchPublishesOnlyAfterCommit() {
+        
when(discoveryProcessorHolder.chooseProcessor(anyString())).thenReturn(discoveryProcessor);
+        
when(selectorMapper.selectByDiscoveryHandlerId(any())).thenReturn(buildSelectorDO());
+        
when(discoveryHandlerMapper.selectById(any())).thenReturn(buildDiscoveryHandlerDO());
+        when(pluginMapper.selectById(any())).thenReturn(buildPluginDO());
+        when(discoveryMapper.selectById(any())).thenReturn(buildDiscoveryDO());
+        TransactionSynchronizationManager.initSynchronization();
+        try {
+            discoveryUpstreamService.updateBatch("123", 
Collections.singletonList(buildDiscoveryUpstreamDTO("")));
+            verifyNoInteractions(discoveryProcessor);
+            verify(discoveryUpstreamMapper, 
never()).selectByDiscoveryHandlerId(any());
+            
TransactionSynchronizationManager.getSynchronizations().forEach(TransactionSynchronization::afterCommit);
+            verify(discoveryProcessor).changeUpstream(any(), any());
+        } finally {
+            TransactionSynchronizationManager.clearSynchronization();
+        }
+    }
+
+    @Test
+    public void testRolledBackBatchDoesNotPublish() {
+        TransactionSynchronizationManager.initSynchronization();
+        try {
+            discoveryUpstreamService.updateBatch("123", 
Collections.emptyList());
+            
TransactionSynchronizationManager.getSynchronizations().forEach(sync -> 
sync.afterCompletion(TransactionSynchronization.STATUS_ROLLED_BACK));
+            verifyNoInteractions(discoveryProcessor);
+            verify(discoveryUpstreamMapper, 
never()).selectByDiscoveryHandlerId(any());
+        } finally {
+            TransactionSynchronizationManager.clearSynchronization();
+        }
     }
 
     private void testUpdate() {

Reply via email to