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 817186dc53 fix(admin): make proxy selector imports atomic (#7247)
817186dc53 is described below
commit 817186dc53c51765fa3aa3e75767c9e8b94c339c
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 10:00:52 2026 +0800
fix(admin): make proxy selector imports atomic (#7247)
---
.../shenyu/admin/model/entity/ProxySelectorDO.java | 1 +
.../service/impl/ProxySelectorServiceImpl.java | 21 +++-
.../ProxySelectorImportIntegrationTest.java | 114 +++++++++++++++++++++
3 files changed, 133 insertions(+), 3 deletions(-)
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/entity/ProxySelectorDO.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/entity/ProxySelectorDO.java
index 13f56bc396..ad3ac1d779 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/entity/ProxySelectorDO.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/model/entity/ProxySelectorDO.java
@@ -119,6 +119,7 @@ public class ProxySelectorDO extends BaseDO {
.forwardPort(item.getForwardPort())
.type(item.getType())
.props(JsonUtils.toJson(item.getProps()))
+ .namespaceId(item.getNamespaceId())
.dateUpdated(currentTime).build();
if (StringUtils.hasLength(item.getId())) {
proxySelectorDO.setId(item.getId());
diff --git
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ProxySelectorServiceImpl.java
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ProxySelectorServiceImpl.java
index ef71339450..ed44ff4807 100644
---
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ProxySelectorServiceImpl.java
+++
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/service/impl/ProxySelectorServiceImpl.java
@@ -56,10 +56,13 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
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.CollectionUtils;
import org.springframework.util.StringUtils;
import java.sql.Timestamp;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -468,6 +471,7 @@ public class ProxySelectorServiceImpl implements
ProxySelectorService {
}
@Override
+ @Transactional(rollbackFor = Exception.class)
public ConfigImportResult importData(final List<ProxySelectorData>
proxySelectorList) {
if (CollectionUtils.isEmpty(proxySelectorList)) {
return ConfigImportResult.success();
@@ -507,12 +511,13 @@ public class ProxySelectorServiceImpl implements
ProxySelectorService {
}
@Override
+ @Transactional(rollbackFor = Exception.class)
public ConfigImportResult importData(final String namespace, final
List<ProxySelectorData> proxySelectorList,
final ConfigsImportContext context) {
if (CollectionUtils.isEmpty(proxySelectorList)) {
return ConfigImportResult.success();
}
- Map<String, String> proxySelectorIdMapping =
context.getProxySelectorIdMapping();
+ Map<String, String> proxySelectorIdMapping = new HashMap<>();
Map<String, List<ProxySelectorDO>> pluginProxySelectorMap =
proxySelectorMapper
.selectByNamespaceId(namespace)
.stream()
@@ -536,14 +541,24 @@ public class ProxySelectorServiceImpl implements
ProxySelectorService {
}
String oldProxySelectorId = selectorData.getId();
String newProxySelectorId =
UUIDUtils.getInstance().generateShortUuid();
- selectorData.setId(newProxySelectorId);
- selectorData.setNamespaceId(namespace);
ProxySelectorDO proxySelectorDO =
ProxySelectorDO.buildProxySelectorDO(selectorData);
+ proxySelectorDO.setId(newProxySelectorId);
+ proxySelectorDO.setNamespaceId(namespace);
if (proxySelectorMapper.insert(proxySelectorDO) > 0) {
proxySelectorIdMapping.put(oldProxySelectorId,
newProxySelectorId);
successCount++;
}
}
+ if (TransactionSynchronizationManager.isSynchronizationActive()) {
+ TransactionSynchronizationManager.registerSynchronization(new
TransactionSynchronization() {
+ @Override
+ public void afterCommit() {
+
context.getProxySelectorIdMapping().putAll(proxySelectorIdMapping);
+ }
+ });
+ } else {
+ context.getProxySelectorIdMapping().putAll(proxySelectorIdMapping);
+ }
if (StringUtils.hasLength(errorMsgBuilder)) {
errorMsgBuilder.setLength(errorMsgBuilder.length() - 1);
return ConfigImportResult
diff --git
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ProxySelectorImportIntegrationTest.java
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ProxySelectorImportIntegrationTest.java
new file mode 100644
index 0000000000..b404095096
--- /dev/null
+++
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/service/ProxySelectorImportIntegrationTest.java
@@ -0,0 +1,114 @@
+/*
+ * 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.service.configs.ConfigsImportContext;
+import org.apache.shenyu.common.dto.ProxySelectorData;
+import org.junit.jupiter.api.Test;
+import org.springframework.dao.DataIntegrityViolationException;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.transaction.PlatformTransactionManager;
+import org.springframework.transaction.support.TransactionTemplate;
+
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Verify proxy selector imports are atomic, including their ID mappings.
+ */
+public class ProxySelectorImportIntegrationTest extends
AbstractSpringIntegrationTest {
+
+ @Resource
+ private ProxySelectorService proxySelectorService;
+
+ @Resource
+ private JdbcTemplate jdbcTemplate;
+
+ @Resource
+ private PlatformTransactionManager transactionManager;
+
+ @Test
+ public void testLegacyImportRollsBackEarlierRows() {
+ assertThrows(DataIntegrityViolationException.class,
+ () ->
proxySelectorService.importData(List.of(selector("first", "valid"),
selector("second", null))));
+ assertEquals(0, countImported());
+ }
+
+ @Test
+ public void testNamespaceImportFailurePreservesContextAndInput() {
+ ConfigsImportContext context = new ConfigsImportContext();
+ context.getProxySelectorIdMapping().put("existing", "existing-id");
+ ProxySelectorData first = selector("first", "valid");
+ assertThrows(DataIntegrityViolationException.class,
+ () -> proxySelectorService.importData("import-target",
List.of(first, selector("second", null)), context));
+ assertEquals(0, countImported());
+ assertEquals(Map.of("existing", "existing-id"),
context.getProxySelectorIdMapping());
+ assertEquals("first", first.getId());
+ assertEquals("import-source", first.getNamespaceId());
+ }
+
+ @Test
+ public void testNamespaceImportPublishesCommittedMapping() {
+ ConfigsImportContext context = new ConfigsImportContext();
+ try {
+ assertEquals(1, proxySelectorService.importData("import-target",
List.of(selector("first", "valid")), context).getSuccessCount());
+ String id = context.getProxySelectorIdMapping().get("first");
+ assertNotNull(id);
+ assertEquals("import-target", jdbcTemplate.queryForObject("SELECT
namespace_id FROM proxy_selector WHERE id = ?", String.class, id));
+ } finally {
+ jdbcTemplate.update("DELETE FROM proxy_selector WHERE namespace_id
= 'import-target'");
+ }
+ }
+
+ @Test
+ public void testOuterRollbackDoesNotPublishMapping() {
+ ConfigsImportContext context = new ConfigsImportContext();
+ new
TransactionTemplate(transactionManager).executeWithoutResult(status -> {
+ proxySelectorService.importData("import-target",
List.of(selector("first", "valid")), context);
+ assertEquals(1, countImported());
+ assertTrue(context.getProxySelectorIdMapping().isEmpty());
+ status.setRollbackOnly();
+ });
+ assertEquals(0, countImported());
+ assertTrue(context.getProxySelectorIdMapping().isEmpty());
+ }
+
+ private int countImported() {
+ return jdbcTemplate.queryForObject("SELECT COUNT(*) FROM
proxy_selector WHERE namespace_id IN ('import-source', 'import-target')",
Integer.class);
+ }
+
+ private ProxySelectorData selector(final String id, final String name) {
+ ProxySelectorData data = new ProxySelectorData();
+ data.setId(id);
+ data.setName(name);
+ data.setPluginName("tcp");
+ data.setType("tcp");
+ data.setForwardPort(18080);
+ data.setNamespaceId("import-source");
+ return data;
+ }
+}
+