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 b0edc606eb fix(admin): activate scale policies after commit (#7253)
b0edc606eb is described below

commit b0edc606ebd7ec029e28256ac45ecff91fe3e57e
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:44:19 2026 +0800

    fix(admin): activate scale policies after commit (#7253)
---
 .../admin/service/impl/ScalePolicyServiceImpl.java |  21 +++-
 .../admin/service/ScalePolicyTransactionTest.java  | 134 +++++++++++++++++++++
 2 files changed, 153 insertions(+), 2 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ScalePolicyServiceImpl.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ScalePolicyServiceImpl.java
index 3996545c23..c20d3aebe4 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ScalePolicyServiceImpl.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ScalePolicyServiceImpl.java
@@ -26,6 +26,9 @@ import 
org.apache.shenyu.admin.scale.scaler.cache.ScalePolicyCache;
 import org.apache.shenyu.admin.service.ScalePolicyService;
 import org.apache.shenyu.common.utils.ListUtil;
 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.util.List;
 
@@ -77,13 +80,27 @@ public class ScalePolicyServiceImpl implements 
ScalePolicyService {
      * @return rows int
      */
     @Override
+    @Transactional(rollbackFor = Exception.class)
     public int update(final ScalePolicyDTO scalePolicyDTO) {
         final ScalePolicyDO scalePolicy = 
ScalePolicyDO.buildScalePolicyDO(scalePolicyDTO);
         int rows = scalePolicyMapper.updateByPrimaryKeySelective(scalePolicy);
         if (rows > 0) {
-            scalePolicyCache.updatePolicy(scalePolicy);
-            scaleService.executeScaling();
+            if (TransactionSynchronizationManager.isSynchronizationActive()) {
+                TransactionSynchronizationManager.registerSynchronization(new 
TransactionSynchronization() {
+                    @Override
+                    public void afterCommit() {
+                        applyPolicy(scalePolicy);
+                    }
+                });
+            } else {
+                applyPolicy(scalePolicy);
+            }
         }
         return rows;
     }
+
+    private void applyPolicy(final ScalePolicyDO scalePolicy) {
+        scalePolicyCache.updatePolicy(scalePolicy);
+        scaleService.executeScaling();
+    }
 }
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ScalePolicyTransactionTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ScalePolicyTransactionTest.java
new file mode 100644
index 0000000000..989c293c27
--- /dev/null
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ScalePolicyTransactionTest.java
@@ -0,0 +1,134 @@
+/*
+ * 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.shenyu.admin.service;
+
+import jakarta.annotation.Resource;
+import org.apache.shenyu.admin.AbstractSpringIntegrationTest;
+import org.apache.shenyu.admin.model.dto.ScalePolicyDTO;
+import org.apache.shenyu.admin.scale.scaler.ScaleService;
+import org.apache.shenyu.admin.scale.scaler.cache.ScalePolicyCache;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.boot.test.mock.mockito.MockBean;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.support.TransactionTemplate;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.Statement;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+/**
+ * Verify scale-policy activation only follows a committed update.
+ */
+public class ScalePolicyTransactionTest extends AbstractSpringIntegrationTest {
+
+    @Resource
+    private ScalePolicyService policyService;
+
+    @Resource
+    private JdbcTemplate jdbcTemplate;
+
+    @Resource
+    private PlatformTransactionManager transactionManager;
+
+    @MockBean
+    private ScalePolicyCache policyCache;
+
+    @MockBean
+    private ScaleService scaleService;
+
+    @BeforeEach
+    public void setup() {
+        jdbcTemplate.update("INSERT INTO scale_policy (id, sort, status, num) 
VALUES ('transaction-policy', 1, 1, 10)");
+        clearInvocations(policyCache, scaleService);
+    }
+
+    @AfterEach
+    public void cleanup() {
+        jdbcTemplate.update("DELETE FROM scale_policy WHERE id = 
'transaction-policy'");
+    }
+
+    @Test
+    public void testRollbackDoesNotChangeCacheOrScale() {
+        new 
TransactionTemplate(transactionManager).executeWithoutResult(status -> {
+            assertEquals(1, policyService.update(policy()));
+            verifyNoInteractions(policyCache, scaleService);
+            status.setRollbackOnly();
+        });
+        assertEquals(10, persistedNum());
+        verifyNoInteractions(policyCache, scaleService);
+    }
+
+    @Test
+    public void testActivationSeesCommittedPolicy() {
+        doAnswer(invocation -> {
+            try (Connection connection = 
jdbcTemplate.getDataSource().getConnection();
+                 Statement statement = connection.createStatement();
+                 ResultSet result = statement.executeQuery("SELECT num FROM 
scale_policy WHERE id = 'transaction-policy'")) {
+                result.next();
+                assertEquals(20, result.getInt(1));
+            }
+            return null;
+        }).when(scaleService).executeScaling();
+        assertEquals(1, policyService.update(policy()));
+        verify(policyCache).updatePolicy(any());
+        verify(scaleService).executeScaling();
+    }
+
+    @Test
+    public void testMissingPolicyDoesNotActivate() {
+        ScalePolicyDTO dto = policy();
+        dto.setId("missing-policy");
+        assertEquals(0, policyService.update(dto));
+        verifyNoInteractions(policyCache, scaleService);
+    }
+
+    @Test
+    public void testExternalFailureOccursAfterDatabaseCommit() {
+        doThrow(new IllegalStateException("scaling 
unavailable")).when(scaleService).executeScaling();
+        assertThrows(IllegalStateException.class, () -> 
policyService.update(policy()));
+        assertEquals(20, persistedNum());
+        verify(policyCache).updatePolicy(any());
+    }
+
+    private int persistedNum() {
+        return jdbcTemplate.queryForObject("SELECT num FROM scale_policy WHERE 
id = 'transaction-policy'", Integer.class);
+    }
+
+    private ScalePolicyDTO policy() {
+        ScalePolicyDTO dto = new ScalePolicyDTO();
+        dto.setId("transaction-policy");
+        dto.setSort(1);
+        dto.setStatus(1);
+        dto.setNum(20);
+        return dto;
+    }
+}
+

Reply via email to