This is an automated email from the ASF dual-hosted git repository.
dengliming 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 b3c3169d6b test: add dedicated unit tests for zero-test
sync/register/protocol classes (#6813) (#7051)
b3c3169d6b is described below
commit b3c3169d6b819878d3b18b747b97703fe6fc7ef3
Author: wy471x <[email protected]>
AuthorDate: Mon Sep 14 13:06:00 2026 +0800
test: add dedicated unit tests for zero-test sync/register/protocol classes
(#6813) (#7051)
Co-authored-by: moremind <[email protected]>
---
.../tcp/handler/TcpUpstreamDataHandlerTest.java | 116 +++++++++
.../protocol/mqtt/MqttBootstrapServerTest.java | 101 ++++++++
.../shenyu/protocol/mqtt/MqttContextTest.java | 97 ++++++++
.../shenyu/protocol/mqtt/MqttFactoryTest.java | 186 +++++++++++++++
.../DefaultConnectionConfigProviderTest.java | 76 ++++++
.../tcp/connection/TcpConnectionBridgeTest.java | 132 ++++++++++
.../http/HttpClientRegisterRepositoryTest.java | 265 +++++++++++++++++++++
.../data/http/refresh/AbstractDataRefreshTest.java | 152 ++++++++++++
.../http/refresh/AiProxyApiKeyDataRefreshTest.java | 135 +++++++++++
.../data/http/refresh/DataRefreshFactoryTest.java | 107 +++++++++
.../handler/AiProxyApiKeyDataHandlerTest.java | 124 ++++++++++
.../handler/ProxySelectorDataHandlerTest.java | 116 +++++++++
12 files changed, 1607 insertions(+)
diff --git
a/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
new file mode 100644
index 0000000000..4a7b375c9e
--- /dev/null
+++
b/shenyu-plugin/shenyu-plugin-proxy/shenyu-plugin-tcp/src/test/java/org/apache/shenyu/plugin/tcp/handler/TcpUpstreamDataHandlerTest.java
@@ -0,0 +1,116 @@
+/*
+ * 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.plugin.tcp.handler;
+
+import org.apache.shenyu.common.dto.DiscoverySyncData;
+import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
+import org.apache.shenyu.protocol.tcp.BootstrapServer;
+import org.apache.shenyu.protocol.tcp.UpstreamProvider;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Test cases for {@link TcpUpstreamDataHandler}.
+ */
+public final class TcpUpstreamDataHandlerTest {
+
+ private static final String SELECTOR = "tcp-upstream-selector";
+
+ private final TcpBootstrapFactory factory =
TcpBootstrapFactory.getSingleton();
+
+ private final TcpUpstreamDataHandler dataHandler = new
TcpUpstreamDataHandler();
+
+ @BeforeEach
+ public void setUp() {
+ factory.clearCache();
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
+ }
+
+ @AfterEach
+ public void tearDown() {
+ factory.clearCache();
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
+ }
+
+ @Test
+ public void pluginNameShouldReturnTcp() {
+ assertEquals("tcp", dataHandler.pluginName());
+ }
+
+ @Test
+ public void handlerShouldRemoveUpstreamsMissingFromTheNewSnapshot() {
+ DiscoveryUpstreamData kept = upstream("127.0.0.1:10001");
+ DiscoveryUpstreamData removed = upstream("127.0.0.1:10002");
+ DiscoveryUpstreamData added = upstream("127.0.0.1:10003");
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Arrays.asList(kept, removed));
+ BootstrapServer bootstrapServer = mock(BootstrapServer.class);
+ factory.cache(SELECTOR, bootstrapServer);
+
+ dataHandler.handlerDiscoveryUpstreamData(
+ syncData(kept, added));
+
+ @SuppressWarnings("unchecked")
+ ArgumentCaptor<List<DiscoveryUpstreamData>> captor =
ArgumentCaptor.forClass(List.class);
+ verify(bootstrapServer).removeCommonUpstream(captor.capture());
+ List<DiscoveryUpstreamData> removeList = captor.getValue();
+ assertEquals(1, removeList.size());
+ assertEquals("127.0.0.1:10002", removeList.get(0).getUrl());
+ assertEquals(Arrays.asList("127.0.0.1:10001", "127.0.0.1:10003"),
+ UpstreamProvider.getSingleton().provide(SELECTOR).stream()
+ .map(DiscoveryUpstreamData::getUrl)
+ .collect(Collectors.toList()));
+ }
+
+ @Test
+ public void handlerShouldStillRefreshCacheWhenServerIsAbsent() {
+ DiscoveryUpstreamData kept = upstream("127.0.0.1:10001");
+ DiscoveryUpstreamData added = upstream("127.0.0.1:10003");
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.singletonList(kept));
+
+ assertDoesNotThrow(() ->
dataHandler.handlerDiscoveryUpstreamData(syncData(kept, added)));
+
+ List<String> urls =
UpstreamProvider.getSingleton().provide(SELECTOR).stream()
+ .map(DiscoveryUpstreamData::getUrl)
+ .collect(Collectors.toList());
+ assertTrue(urls.contains("127.0.0.1:10003"));
+ }
+
+ private DiscoverySyncData syncData(final DiscoveryUpstreamData...
upstreams) {
+ DiscoverySyncData discoverySyncData = new DiscoverySyncData();
+ discoverySyncData.setSelectorName(SELECTOR);
+ discoverySyncData.setUpstreamDataList(Arrays.asList(upstreams));
+ return discoverySyncData;
+ }
+
+ private DiscoveryUpstreamData upstream(final String url) {
+ return
DiscoveryUpstreamData.builder().url(url).status(0).weight(1).build();
+ }
+}
diff --git
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttBootstrapServerTest.java
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttBootstrapServerTest.java
new file mode 100644
index 0000000000..66e941c36e
--- /dev/null
+++
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttBootstrapServerTest.java
@@ -0,0 +1,101 @@
+/*
+ * 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.protocol.mqtt;
+
+import io.netty.channel.ChannelFuture;
+import io.netty.channel.EventLoopGroup;
+import org.apache.shenyu.common.utils.Singleton;
+import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
+import org.apache.shenyu.protocol.mqtt.repositories.SubscribeRepository;
+import org.apache.shenyu.protocol.mqtt.repositories.TopicRepository;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.lang.reflect.Field;
+import java.time.Duration;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test cases for {@link MqttBootstrapServer}.
+ */
+public final class MqttBootstrapServerTest {
+
+ @BeforeEach
+ public void setUp() {
+ MqttContext context = new MqttContext();
+ context.setPort(0);
+ context.setBossGroupThreadCount(1);
+ context.setWorkerGroupThreadCount(1);
+ context.setMaxPayloadSize(1024 * 1024);
+ context.setUserName("test-user");
+ context.setPassword("test-password");
+ context.setLeakDetectorLevel("disabled");
+ }
+
+ @AfterEach
+ public void tearDown() {
+ MqttContext context = new MqttContext();
+ context.setPort(0);
+ context.setBossGroupThreadCount(0);
+ context.setWorkerGroupThreadCount(0);
+ context.setMaxPayloadSize(0);
+ context.setUserName(null);
+ context.setPassword(null);
+ context.setLeakDetectorLevel(null);
+ }
+
+ @Test
+ public void initShouldRegisterAllRepositories() {
+ MqttBootstrapServer server = new MqttBootstrapServer();
+
+ server.init();
+
+ assertNotNull(Singleton.INST.get(ChannelRepository.class));
+ assertNotNull(Singleton.INST.get(SubscribeRepository.class));
+ assertNotNull(Singleton.INST.get(TopicRepository.class));
+ }
+
+ @Test
+ public void startAndShutdownShouldReleaseChannelAndEventLoops() throws
Exception {
+ MqttBootstrapServer server = new MqttBootstrapServer();
+
+ server.start();
+
+ ChannelFuture future = getField(server, "future", ChannelFuture.class);
+ assertTrue(future.channel().isActive());
+
+ server.shutdown();
+
+ assertFalse(future.channel().isActive());
+ EventLoopGroup bossGroup = getField(server, "bossGroup",
EventLoopGroup.class);
+ EventLoopGroup workerGroup = getField(server, "workerGroup",
EventLoopGroup.class);
+ await().atMost(Duration.ofSeconds(5)).until(bossGroup::isTerminated);
+ await().atMost(Duration.ofSeconds(5)).until(workerGroup::isTerminated);
+ }
+
+ private <T> T getField(final Object target, final String name, final
Class<T> type) throws Exception {
+ Field field = MqttBootstrapServer.class.getDeclaredField(name);
+ field.setAccessible(true);
+ return type.cast(field.get(target));
+ }
+}
diff --git
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttContextTest.java
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttContextTest.java
new file mode 100644
index 0000000000..6575362d6c
--- /dev/null
+++
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttContextTest.java
@@ -0,0 +1,97 @@
+/*
+ * 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.protocol.mqtt;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.nio.charset.StandardCharsets;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test cases for {@link MqttContext}.
+ */
+public final class MqttContextTest {
+
+ @BeforeEach
+ public void setUp() {
+ MqttContext context = new MqttContext();
+ context.setPort(1883);
+ context.setBossGroupThreadCount(1);
+ context.setWorkerGroupThreadCount(2);
+ context.setMaxPayloadSize(1024);
+ context.setUserName("test-user");
+ context.setPassword("test-password");
+ context.setLeakDetectorLevel("disabled");
+ }
+
+ @AfterEach
+ public void tearDown() {
+ MqttContext context = new MqttContext();
+ context.setPort(0);
+ context.setBossGroupThreadCount(0);
+ context.setWorkerGroupThreadCount(0);
+ context.setMaxPayloadSize(0);
+ context.setUserName(null);
+ context.setPassword(null);
+ context.setLeakDetectorLevel(null);
+ }
+
+ @Test
+ public void settersShouldUpdateStaticState() {
+ MqttContext context = new MqttContext();
+
+ assertEquals(1883, context.getPort());
+ assertEquals(1, context.getBossGroupThreadCount());
+ assertEquals(2, context.getWorkerGroupThreadCount());
+ assertEquals(1024, context.getMaxPayloadSize());
+ assertEquals("test-user", context.getUserName());
+ assertEquals("test-password", context.getPassword());
+ assertEquals("disabled", context.getLeakDetectorLevel());
+ }
+
+ @Test
+ public void emptyUserNameOrPasswordShouldBeRejected() {
+ byte[] password = "test-password".getBytes(StandardCharsets.UTF_8);
+
+ assertFalse(MqttContext.isValid("", password));
+ assertFalse(MqttContext.isValid("test-user", new byte[0]));
+ }
+
+ @Test
+ public void mismatchedCredentialsShouldBeRejected() {
+ byte[] password = "test-password".getBytes(StandardCharsets.UTF_8);
+
+ assertFalse(MqttContext.isValid("another-user", password));
+ assertFalse(MqttContext.isValid("test-user",
"another-password".getBytes(StandardCharsets.UTF_8)));
+ }
+
+ @Test
+ public void validCredentialsShouldBeAccepted() {
+ assertTrue(MqttContext.isValid("test-user",
"test-password".getBytes(StandardCharsets.UTF_8)));
+ }
+
+ @Test
+ public void nullUserNameArgumentShouldBeRejectedWithoutNpe() {
+ assertFalse(MqttContext.isValid(null,
"test-password".getBytes(StandardCharsets.UTF_8)));
+ }
+}
diff --git
a/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttFactoryTest.java
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttFactoryTest.java
new file mode 100644
index 0000000000..deb0550121
--- /dev/null
+++
b/shenyu-protocol/shenyu-protocol-mqtt/src/test/java/org/apache/shenyu/protocol/mqtt/MqttFactoryTest.java
@@ -0,0 +1,186 @@
+/*
+ * 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.protocol.mqtt;
+
+import io.netty.buffer.Unpooled;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.ChannelInboundHandlerAdapter;
+import io.netty.channel.embedded.EmbeddedChannel;
+import io.netty.handler.codec.mqtt.MqttConnectMessage;
+import io.netty.handler.codec.mqtt.MqttConnectPayload;
+import io.netty.handler.codec.mqtt.MqttConnectVariableHeader;
+import io.netty.handler.codec.mqtt.MqttFixedHeader;
+import io.netty.handler.codec.mqtt.MqttMessage;
+import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader;
+import io.netty.handler.codec.mqtt.MqttMessageType;
+import io.netty.handler.codec.mqtt.MqttPubAckMessage;
+import io.netty.handler.codec.mqtt.MqttPublishMessage;
+import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
+import io.netty.handler.codec.mqtt.MqttQoS;
+import io.netty.handler.codec.mqtt.MqttSubscribeMessage;
+import io.netty.handler.codec.mqtt.MqttSubscribePayload;
+import io.netty.handler.codec.mqtt.MqttTopicSubscription;
+import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage;
+import io.netty.handler.codec.mqtt.MqttUnsubscribePayload;
+import io.netty.handler.codec.mqtt.MqttVersion;
+import org.apache.shenyu.common.utils.Singleton;
+import org.apache.shenyu.protocol.mqtt.repositories.ChannelRepository;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test cases for {@link MqttFactory}.
+ */
+public final class MqttFactoryTest {
+
+ private static final String CLIENT_ID = "factory-client";
+
+ private static final String USER_NAME = "factory-user";
+
+ private static final String PASSWORD = "factory-password";
+
+ @BeforeEach
+ public void setUp() {
+ Singleton.INST.single(ChannelRepository.class, new
ChannelRepository());
+ new MqttContext().setUserName(USER_NAME);
+ new MqttContext().setPassword(PASSWORD);
+ }
+
+ @AfterEach
+ public void tearDown() {
+ new MqttContext().setUserName(null);
+ new MqttContext().setPassword(null);
+ }
+
+ @Test
+ public void messageWithoutFixedHeaderShouldBeIgnored() {
+ EmbeddedChannel channel = new EmbeddedChannel(new
ChannelInboundHandlerAdapter());
+ ChannelHandlerContext ctx = channel.pipeline().lastContext();
+
+ new MqttFactory(new MqttMessage(null, null), ctx).connect();
+
+ assertTrue(channel.isActive());
+ channel.finishAndReleaseAll();
+ }
+
+ @Test
+ public void connectShouldBeDispatchedToConnectHandler() {
+ EmbeddedChannel channel = new EmbeddedChannel(new
ChannelInboundHandlerAdapter());
+ ChannelHandlerContext ctx = channel.pipeline().lastContext();
+
+ new MqttFactory(connectMessage(), ctx).connect();
+
+ assertNotNull(channel.readOutbound());
+ assertTrue(channel.isActive());
+ channel.finishAndReleaseAll();
+ }
+
+ @Test
+ public void publishBeforeConnectShouldBeDispatchedAndCloseChannel() {
+ MqttPublishMessage publish = new MqttPublishMessage(
+ fixedHeader(MqttMessageType.PUBLISH),
+ new MqttPublishVariableHeader("topic", 0),
+ Unpooled.EMPTY_BUFFER);
+
+ assertDispatchedMessageClosesChannel(publish);
+ }
+
+ @Test
+ public void subscribeBeforeConnectShouldBeDispatchedAndCloseChannel() {
+ MqttSubscribeMessage subscribe = new MqttSubscribeMessage(
+ fixedHeader(MqttMessageType.SUBSCRIBE),
+ MqttMessageIdVariableHeader.from(1),
+ new MqttSubscribePayload(Collections.singletonList(
+ new MqttTopicSubscription("topic",
MqttQoS.AT_MOST_ONCE))));
+
+ assertDispatchedMessageClosesChannel(subscribe);
+ }
+
+ @Test
+ public void unsubscribeBeforeConnectShouldBeDispatchedAndCloseChannel() {
+ MqttUnsubscribeMessage unsubscribe = new MqttUnsubscribeMessage(
+ fixedHeader(MqttMessageType.UNSUBSCRIBE),
+ MqttMessageIdVariableHeader.from(1),
+ new
MqttUnsubscribePayload(Collections.singletonList("topic")));
+
+ assertDispatchedMessageClosesChannel(unsubscribe);
+ }
+
+ @Test
+ public void pingReqBeforeConnectShouldBeDispatchedAndCloseChannel() {
+ assertDispatchedMessageClosesChannel(new
MqttMessage(fixedHeader(MqttMessageType.PINGREQ)));
+ }
+
+ @Test
+ public void pubAckShouldFallThroughToNoOp() {
+ EmbeddedChannel channel = new EmbeddedChannel(new
ChannelInboundHandlerAdapter());
+ ChannelHandlerContext ctx = channel.pipeline().lastContext();
+
+ MqttPubAckMessage pubAck = new MqttPubAckMessage(
+ new MqttFixedHeader(MqttMessageType.PUBACK, false,
MqttQoS.AT_LEAST_ONCE, false, 0),
+ MqttMessageIdVariableHeader.from(1));
+ new MqttFactory(pubAck, ctx).connect();
+
+ assertTrue(channel.isActive());
+ channel.finishAndReleaseAll();
+ }
+
+ @Test
+ public void disconnectShouldFallThroughToNoOp() {
+ EmbeddedChannel channel = new EmbeddedChannel(new
ChannelInboundHandlerAdapter());
+ ChannelHandlerContext ctx = channel.pipeline().lastContext();
+
+ new MqttFactory(new
MqttMessage(fixedHeader(MqttMessageType.DISCONNECT)), ctx).connect();
+
+ assertTrue(channel.isActive());
+ channel.finishAndReleaseAll();
+ }
+
+ private void assertDispatchedMessageClosesChannel(final MqttMessage
message) {
+ EmbeddedChannel channel = new EmbeddedChannel(new
ChannelInboundHandlerAdapter());
+ ChannelHandlerContext ctx = channel.pipeline().lastContext();
+
+ new MqttFactory(message, ctx).connect();
+ channel.runPendingTasks();
+
+ assertFalse(channel.isActive());
+ channel.finishAndReleaseAll();
+ }
+
+ private MqttFixedHeader fixedHeader(final MqttMessageType messageType) {
+ return new MqttFixedHeader(messageType, false, MqttQoS.AT_MOST_ONCE,
false, 0);
+ }
+
+ private MqttConnectMessage connectMessage() {
+ MqttFixedHeader fixedHeader = new
MqttFixedHeader(MqttMessageType.CONNECT, false, MqttQoS.AT_MOST_ONCE, false, 0);
+ MqttConnectVariableHeader variableHeader = new
MqttConnectVariableHeader(
+ MqttVersion.MQTT_3_1_1.protocolName(),
MqttVersion.MQTT_3_1_1.protocolLevel(),
+ true, true, false, 0, false, false, 60);
+ MqttConnectPayload payload = new MqttConnectPayload(CLIENT_ID, null,
null,
+ USER_NAME, PASSWORD.getBytes(StandardCharsets.UTF_8));
+ return new MqttConnectMessage(fixedHeader, variableHeader, payload);
+ }
+}
diff --git
a/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/DefaultConnectionConfigProviderTest.java
b/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/DefaultConnectionConfigProviderTest.java
new file mode 100644
index 0000000000..82fc251cd2
--- /dev/null
+++
b/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/DefaultConnectionConfigProviderTest.java
@@ -0,0 +1,76 @@
+/*
+ * 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.protocol.tcp.connection;
+
+import org.apache.shenyu.common.dto.DiscoveryUpstreamData;
+import org.apache.shenyu.common.exception.ShenyuException;
+import org.apache.shenyu.protocol.tcp.UpstreamProvider;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import java.net.URI;
+import java.sql.Timestamp;
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test cases for {@link DefaultConnectionConfigProvider}.
+ */
+public final class DefaultConnectionConfigProviderTest {
+
+ private static final String SELECTOR = "tcp-selector";
+
+ @AfterEach
+ public void cleanUpstreams() {
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
+ }
+
+ @Test
+ public void getProxiedServiceShouldBuildUriFromSelectedUpstream() {
+ DiscoveryUpstreamData upstream = DiscoveryUpstreamData.builder()
+ .protocol("tcp")
+ .url("127.0.0.1:20000")
+ .status(0)
+ .weight(1)
+ .props("{\"warmupTime\":0}")
+ .dateCreated(new Timestamp(System.currentTimeMillis()))
+ .build();
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.singletonList(upstream));
+ DefaultConnectionConfigProvider provider =
+ new DefaultConnectionConfigProvider("random", SELECTOR);
+
+ URI proxiedService = provider.getProxiedService("127.0.0.1");
+
+ assertEquals(URI.create("tcp://127.0.0.1:20000"), proxiedService);
+ }
+
+ @Test
+ public void getProxiedServiceShouldThrowWhenNoUpstreamExists() {
+ UpstreamProvider.getSingleton().createUpstreams(SELECTOR,
Collections.emptyList());
+ DefaultConnectionConfigProvider provider =
+ new DefaultConnectionConfigProvider("random", SELECTOR);
+
+ ShenyuException exception =
+ assertThrows(ShenyuException.class, () ->
provider.getProxiedService("127.0.0.1"));
+
+ assertTrue(exception.getMessage().contains("don't have any upstream"));
+ }
+}
diff --git
a/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/TcpConnectionBridgeTest.java
b/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/TcpConnectionBridgeTest.java
new file mode 100644
index 0000000000..c43045f6dd
--- /dev/null
+++
b/shenyu-protocol/shenyu-protocol-tcp/src/test/java/org/apache/shenyu/protocol/tcp/connection/TcpConnectionBridgeTest.java
@@ -0,0 +1,132 @@
+/*
+ * 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.protocol.tcp.connection;
+
+import io.netty.channel.Channel;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.reactivestreams.Publisher;
+import reactor.core.Disposable;
+import reactor.core.publisher.Mono;
+import reactor.netty.ByteBufFlux;
+import reactor.netty.Connection;
+import reactor.netty.NettyInbound;
+import reactor.netty.NettyOutbound;
+
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentCaptor.forClass;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Test cases for {@link TcpConnectionBridge}.
+ */
+public final class TcpConnectionBridgeTest {
+
+ private Connection server;
+
+ private Connection client;
+
+ private NettyInbound serverInbound;
+
+ private NettyInbound clientInbound;
+
+ private ByteBufFlux serverReceiveFlux;
+
+ private ByteBufFlux clientReceiveFlux;
+
+ private NettyOutbound serverOutbound;
+
+ private NettyOutbound clientOutbound;
+
+ private Channel serverChannel;
+
+ private Channel clientChannel;
+
+ @BeforeEach
+ public void setUp() {
+ server = mock(Connection.class);
+ client = mock(Connection.class);
+ serverInbound = mock(NettyInbound.class);
+ clientInbound = mock(NettyInbound.class);
+ serverOutbound = mock(NettyOutbound.class);
+ clientOutbound = mock(NettyOutbound.class);
+ serverReceiveFlux = mock(ByteBufFlux.class);
+ clientReceiveFlux = mock(ByteBufFlux.class);
+ serverChannel = mock(Channel.class);
+ clientChannel = mock(Channel.class);
+
+ when(serverReceiveFlux.retain()).thenReturn(serverReceiveFlux);
+ when(clientReceiveFlux.retain()).thenReturn(clientReceiveFlux);
+ when(server.inbound()).thenReturn(serverInbound);
+ when(server.outbound()).thenReturn(serverOutbound);
+ when(client.inbound()).thenReturn(clientInbound);
+ when(client.outbound()).thenReturn(clientOutbound);
+ when(serverInbound.receive()).thenReturn(serverReceiveFlux);
+ when(clientInbound.receive()).thenReturn(clientReceiveFlux);
+
doReturn(serverOutbound).when(serverOutbound).send(any(Publisher.class));
+
doReturn(clientOutbound).when(clientOutbound).send(any(Publisher.class));
+ doReturn(Mono.empty()).when(serverOutbound).then();
+ doReturn(Mono.empty()).when(clientOutbound).then();
+ when(server.channel()).thenReturn(serverChannel);
+ when(client.channel()).thenReturn(clientChannel);
+ when(server.onDispose(any(Disposable.class))).thenReturn(server);
+ when(client.onDispose(any(Disposable.class))).thenReturn(client);
+ }
+
+ @Test
+ public void bridgeShouldRelayTrafficInBothDirections() {
+ TcpConnectionBridge bridge = new TcpConnectionBridge();
+
+ bridge.bridge(server, client);
+
+ verify(serverInbound).receive();
+ verify(clientInbound).receive();
+ verify(serverOutbound).send(any(Publisher.class));
+ verify(clientOutbound).send(any(Publisher.class));
+ verify(serverOutbound).then();
+ verify(clientOutbound).then();
+ }
+
+ @Test
+ public void disposingServerShouldCloseClientChannel() {
+ ArgumentCaptor<Disposable> captor = forClass(Disposable.class);
+ when(server.onDispose(captor.capture())).thenReturn(server);
+
+ new TcpConnectionBridge().bridge(server, client);
+
+ assertNotNull(captor.getValue());
+ captor.getValue().dispose();
+ verify(clientChannel).close();
+ }
+
+ @Test
+ public void disposingClientShouldCloseServerChannel() {
+ ArgumentCaptor<Disposable> captor = forClass(Disposable.class);
+ when(client.onDispose(captor.capture())).thenReturn(client);
+
+ new TcpConnectionBridge().bridge(server, client);
+
+ captor.getValue().dispose();
+ verify(serverChannel).close();
+ }
+}
diff --git
a/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
new file mode 100644
index 0000000000..e10db3a70b
--- /dev/null
+++
b/shenyu-register-center/shenyu-register-client/shenyu-register-client-http/src/test/java/org/apache/shenyu/register/client/http/HttpClientRegisterRepositoryTest.java
@@ -0,0 +1,265 @@
+/*
+ * 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.register.client.http;
+
+import org.apache.shenyu.common.constant.Constants;
+import org.apache.shenyu.register.client.http.utils.RegisterUtils;
+import org.apache.shenyu.register.client.http.utils.RuntimeUtils;
+import org.apache.shenyu.register.common.config.ShenyuRegisterCenterConfig;
+import org.apache.shenyu.register.common.dto.ApiDocRegisterDTO;
+import org.apache.shenyu.register.common.dto.DiscoveryConfigRegisterDTO;
+import org.apache.shenyu.register.common.dto.McpToolsRegisterDTO;
+import org.apache.shenyu.register.common.dto.MetaDataRegisterDTO;
+import org.apache.shenyu.register.common.dto.URIRegisterDTO;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.lang.reflect.Field;
+import java.util.Optional;
+import java.util.Properties;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+
+/**
+ * Test cases for {@link HttpClientRegisterRepository}.
+ */
+public final class HttpClientRegisterRepositoryTest {
+
+ private static final String FIRST_SERVER = "http://localhost:9095";
+
+ private static final String SECOND_SERVER = "http://localhost:9096";
+
+ private static final String TOKEN = "test-token";
+
+ private HttpClientRegisterRepository repository;
+
+ @BeforeEach
+ public void setUp() throws Exception {
+ resetStatics();
+ repository = new HttpClientRegisterRepository(config(FIRST_SERVER));
+ }
+
+ @AfterEach
+ public void tearDown() throws Exception {
+ resetStatics();
+ }
+
+ @Test
+ public void persistUriShouldRegisterToEveryServer() {
+ HttpClientRegisterRepository multiServerRepository = new
HttpClientRegisterRepository(config(
+ FIRST_SERVER + "," + SECOND_SERVER));
+ URIRegisterDTO uriRegisterDTO = uriRegisterDTO();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+
+ multiServerRepository.persistURI(uriRegisterDTO);
+
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.URI_PATH), eq(Constants.URI),
eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(SECOND_SERVER + Constants.URI_PATH), eq(Constants.URI),
eq(TOKEN)));
+ }
+ }
+
+ @Test
+ public void
persistUriShouldSkipRegistrationWhenPortIsUsedByAnotherProcess() {
+ URIRegisterDTO uriRegisterDTO = uriRegisterDTO();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(true);
+
+ repository.doPersistURI(uriRegisterDTO);
+
+ registerUtils.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ public void
persistInterfaceApiDocMcpToolsAndDiscoveryShouldUseExpectedPaths() {
+ MetaDataRegisterDTO metaData = metaDataRegisterDTO();
+ ApiDocRegisterDTO apiDocRegisterDTO = ApiDocRegisterDTO.builder()
+ .contextPath("/demo")
+ .apiPath("/hello")
+ .httpMethod(1)
+ .rpcType("http")
+ .build();
+ McpToolsRegisterDTO mcpToolsRegisterDTO = new McpToolsRegisterDTO();
+ mcpToolsRegisterDTO.setMetaDataRegisterDTO(metaData);
+ DiscoveryConfigRegisterDTO discoveryConfigRegisterDTO =
DiscoveryConfigRegisterDTO.builder().build();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+
+ repository.persistInterface(metaData);
+ repository.persistApiDoc(apiDocRegisterDTO);
+ repository.persistMcpTools(mcpToolsRegisterDTO);
+ repository.doPersistDiscoveryConfig(discoveryConfigRegisterDTO);
+
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.META_PATH),
eq(Constants.META_TYPE), eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.API_DOC_PATH),
eq(Constants.API_DOC_TYPE), eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.MCP_TOOLS_PATH),
eq(Constants.MCP_TOOLS_TYPE), eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.DISCOVERY_CONFIG_PATH),
eq(Constants.DISCOVERY_CONFIG_TYPE), eq(TOKEN)));
+ }
+ }
+
+ @Test
+ public void heartbeatAndOfflineShouldSendExpectedRequests() {
+ URIRegisterDTO heartbeatDTO = uriRegisterDTO();
+ URIRegisterDTO offlineDTO = uriRegisterDTO();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+
+ repository.sendHeartbeat(heartbeatDTO);
+ repository.offline(offlineDTO);
+
+ registerUtils.verify(() -> RegisterUtils.doHeartBeat(anyString(),
+ eq(FIRST_SERVER + Constants.URI_PATH),
eq(Constants.HEARTBEAT), eq(TOKEN)));
+ registerUtils.verify(() -> RegisterUtils.doUnregister(anyString(),
+ eq(FIRST_SERVER + Constants.OFFLINE_PATH), eq(TOKEN)));
+ }
+ }
+
+ @Test
+ public void heartbeatShouldSkipWhenPortIsUsedByAnotherProcess() {
+ URIRegisterDTO heartbeatDTO = uriRegisterDTO();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(true);
+
+ repository.sendHeartbeat(heartbeatDTO);
+
+ registerUtils.verifyNoInteractions();
+ }
+ }
+
+ @Test
+ public void closeRepositoryShouldUnregisterLastUriAndApiDoc() {
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+
+ repository.persistURI(uriRegisterDTO());
+
repository.persistApiDoc(ApiDocRegisterDTO.builder().apiPath("/hello").build());
+ repository.closeRepository();
+
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.URI_PATH), eq(Constants.URI),
eq(TOKEN)), times(2));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.API_DOC_PATH),
eq(Constants.API_DOC_TYPE), eq(TOKEN)), times(2));
+ }
+ }
+
+ @Test
+ public void doPersistUriShouldThrowWhenEveryServerFails() throws
IOException {
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.of(TOKEN));
+ registerUtils.when(() -> RegisterUtils.doRegister(anyString(),
anyString(), anyString(), anyString()))
+ .thenThrow(new IOException("register failed"));
+
+ RuntimeException exception =
+ assertThrows(RuntimeException.class, () ->
repository.doPersistURI(uriRegisterDTO()));
+
+ assertTrue(exception.getCause() instanceof IOException);
+ assertEquals("register failed", exception.getCause().getMessage());
+ }
+ }
+
+ @Test
+ public void loginFailureShouldSkipRegistration() {
+ URIRegisterDTO uriRegisterDTO = uriRegisterDTO();
+
+ try (MockedStatic<RegisterUtils> registerUtils =
mockStatic(RegisterUtils.class);
+ MockedStatic<RuntimeUtils> runtimeUtils =
mockStatic(RuntimeUtils.class)) {
+ runtimeUtils.when(() ->
RuntimeUtils.listenByOther(anyInt())).thenReturn(false);
+ registerUtils.when(() -> RegisterUtils.doLogin(anyString(),
anyString(), anyString()))
+ .thenReturn(Optional.empty());
+
+ assertThrows(RuntimeException.class, () ->
repository.doPersistURI(uriRegisterDTO));
+ registerUtils.verify(() -> RegisterUtils.doRegister(anyString(),
+ eq(FIRST_SERVER + Constants.URI_PATH), eq(Constants.URI),
eq(TOKEN)), never());
+ }
+ }
+
+ private ShenyuRegisterCenterConfig config(final String serverLists) {
+ Properties props = new Properties();
+ props.setProperty(Constants.USER_NAME, "admin");
+ props.setProperty(Constants.PASS_WORD, "123456");
+ return new ShenyuRegisterCenterConfig("http", serverLists, props);
+ }
+
+ private URIRegisterDTO uriRegisterDTO() {
+ return URIRegisterDTO.builder()
+ .appName("demo")
+ .rpcType("http")
+ .host("127.0.0.1")
+ .port(18080)
+ .build();
+ }
+
+ private MetaDataRegisterDTO metaDataRegisterDTO() {
+ MetaDataRegisterDTO metaDataRegisterDTO = new MetaDataRegisterDTO();
+ metaDataRegisterDTO.setAppName("demo");
+ metaDataRegisterDTO.setRpcType("http");
+ metaDataRegisterDTO.setHost("127.0.0.1");
+ metaDataRegisterDTO.setPort(18080);
+ metaDataRegisterDTO.setPath("/demo");
+ return metaDataRegisterDTO;
+ }
+
+ private void resetStatics() throws Exception {
+ Field uriField =
HttpClientRegisterRepository.class.getDeclaredField("uriRegisterDTO");
+ uriField.setAccessible(true);
+ uriField.set(null, null);
+ Field apiDocField =
HttpClientRegisterRepository.class.getDeclaredField("apiDocRegisterDTO");
+ apiDocField.setAccessible(true);
+ apiDocField.set(null, null);
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AbstractDataRefreshTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AbstractDataRefreshTest.java
new file mode 100644
index 0000000000..c73680de09
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AbstractDataRefreshTest.java
@@ -0,0 +1,152 @@
+/*
+ * 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.sync.data.http.refresh;
+
+import com.google.gson.JsonObject;
+import org.apache.shenyu.common.dto.ConfigData;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Test cases for {@link AbstractDataRefresh}.
+ */
+public final class AbstractDataRefreshTest {
+
+ private static final ConfigGroupEnum GROUP = ConfigGroupEnum.META_DATA;
+
+ @BeforeEach
+ public void clearGroupCache() {
+ AbstractDataRefresh.GROUP_CACHE.remove(GROUP);
+ }
+
+ @AfterEach
+ public void tearDown() {
+ AbstractDataRefresh.GROUP_CACHE.remove(GROUP);
+ }
+
+ @Test
+ public void refreshShouldReturnFalseWhenConvertReturnsNull() {
+ StubDataRefresh dataRefresh = new StubDataRefresh();
+ dataRefresh.skipConvert = true;
+
+ assertFalse(dataRefresh.refresh(new JsonObject()));
+ assertFalse(dataRefresh.refreshed);
+ }
+
+ @Test
+ public void refreshShouldUpdateWhenGroupCacheIsEmpty() {
+ StubDataRefresh dataRefresh = new StubDataRefresh();
+ ConfigData<PluginData> config = config("md5-new", 100L);
+ dataRefresh.parsed = config;
+
+ assertTrue(dataRefresh.refresh(new JsonObject()));
+ assertTrue(dataRefresh.refreshed);
+ assertSame(config, AbstractDataRefresh.GROUP_CACHE.get(GROUP));
+ }
+
+ @Test
+ public void refreshShouldIgnoreSameMd5EvenWithNewerModifyTime() {
+ StubDataRefresh dataRefresh = new StubDataRefresh();
+ ConfigData<PluginData> original = config("md5-same", 100L);
+ dataRefresh.parsed = original;
+ assertTrue(dataRefresh.refresh(new JsonObject()));
+
+ dataRefresh.refreshed = false;
+ dataRefresh.parsed = config("md5-same", 200L);
+ assertFalse(dataRefresh.refresh(new JsonObject()));
+ assertFalse(dataRefresh.refreshed);
+ assertSame(original, AbstractDataRefresh.GROUP_CACHE.get(GROUP));
+ }
+
+ @Test
+ public void refreshShouldIgnoreNewerMd5WhenModifyTimeIsNotNewer() {
+ StubDataRefresh dataRefresh = new StubDataRefresh();
+ ConfigData<PluginData> original = config("md5-old", 200L);
+ dataRefresh.parsed = original;
+ assertTrue(dataRefresh.refresh(new JsonObject()));
+
+ dataRefresh.refreshed = false;
+ dataRefresh.parsed = config("md5-new", 100L);
+ assertFalse(dataRefresh.refresh(new JsonObject()));
+ assertFalse(dataRefresh.refreshed);
+ assertSame(original, AbstractDataRefresh.GROUP_CACHE.get(GROUP));
+ }
+
+ @Test
+ public void refreshShouldUpdateWhenMd5AndModifyTimeAreNewer() {
+ StubDataRefresh dataRefresh = new StubDataRefresh();
+ ConfigData<PluginData> original = config("md5-old", 100L);
+ dataRefresh.parsed = original;
+ assertTrue(dataRefresh.refresh(new JsonObject()));
+
+ dataRefresh.refreshed = false;
+ ConfigData<PluginData> latest = config("md5-new", 200L);
+ dataRefresh.parsed = latest;
+ assertTrue(dataRefresh.refresh(new JsonObject()));
+ assertTrue(dataRefresh.refreshed);
+ assertSame(latest, AbstractDataRefresh.GROUP_CACHE.get(GROUP));
+ }
+
+ private ConfigData<PluginData> config(final String md5, final long
lastModifyTime) {
+ return new ConfigData<>(md5, lastModifyTime,
Collections.<PluginData>emptyList());
+ }
+
+ private static final class StubDataRefresh extends
AbstractDataRefresh<PluginData> {
+
+ private boolean skipConvert;
+
+ private boolean refreshed;
+
+ private ConfigData<PluginData> parsed;
+
+ @Override
+ protected JsonObject convert(final JsonObject data) {
+ return skipConvert ? null : data;
+ }
+
+ @Override
+ protected ConfigData<PluginData> fromJson(final JsonObject data) {
+ return parsed;
+ }
+
+ @Override
+ protected void refresh(final List<PluginData> data) {
+ refreshed = true;
+ }
+
+ @Override
+ protected boolean updateCacheIfNeed(final ConfigData<PluginData>
result) {
+ return updateCacheIfNeed(result, GROUP);
+ }
+
+ @Override
+ public ConfigData<?> cacheConfigData() {
+ return GROUP_CACHE.get(GROUP);
+ }
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AiProxyApiKeyDataRefreshTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AiProxyApiKeyDataRefreshTest.java
new file mode 100644
index 0000000000..b4762b4609
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/AiProxyApiKeyDataRefreshTest.java
@@ -0,0 +1,135 @@
+/*
+ * 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.sync.data.http.refresh;
+
+import com.google.gson.JsonObject;
+import org.apache.shenyu.common.dto.ConfigData;
+import org.apache.shenyu.common.dto.ProxyApiKeyData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.sync.data.api.AiProxyApiKeyDataSubscriber;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Test cases for {@link AiProxyApiKeyDataRefresh}.
+ */
+public final class AiProxyApiKeyDataRefreshTest {
+
+ private final AiProxyApiKeyDataSubscriber subscriber =
mock(AiProxyApiKeyDataSubscriber.class);
+
+ private final AiProxyApiKeyDataRefresh dataRefresh =
+ new
AiProxyApiKeyDataRefresh(Collections.singletonList(subscriber));
+
+ @BeforeEach
+ public void clearGroupCache() {
+
AbstractDataRefresh.GROUP_CACHE.remove(ConfigGroupEnum.AI_PROXY_API_KEY);
+ }
+
+ @AfterEach
+ public void tearDown() {
+
AbstractDataRefresh.GROUP_CACHE.remove(ConfigGroupEnum.AI_PROXY_API_KEY);
+ }
+
+ @Test
+ public void convertShouldReturnGroupJson() {
+ JsonObject data = new JsonObject();
+ JsonObject groupJson = new JsonObject();
+ data.add(ConfigGroupEnum.AI_PROXY_API_KEY.name(), groupJson);
+
+ assertEquals(groupJson, dataRefresh.convert(data));
+ assertNull(dataRefresh.convert(new JsonObject()));
+ }
+
+ @Test
+ public void fromJsonShouldParseConfigData() {
+ ProxyApiKeyData apiKeyData = ProxyApiKeyData.builder()
+ .realApiKey("real-key")
+ .proxyApiKey("proxy-key")
+ .enabled(true)
+ .build();
+ ConfigData<ProxyApiKeyData> config =
+ new ConfigData<>("md5", 100L,
Collections.singletonList(apiKeyData));
+ JsonObject json =
GsonUtils.getGson().fromJson(GsonUtils.getGson().toJson(config),
JsonObject.class);
+
+ assertEquals(config, dataRefresh.fromJson(json));
+ }
+
+ @Test
+ public void refreshShouldClearThenSubscribeAllItems() {
+ ProxyApiKeyData first =
ProxyApiKeyData.builder().proxyApiKey("first").build();
+ ProxyApiKeyData second =
ProxyApiKeyData.builder().proxyApiKey("second").build();
+
+ dataRefresh.refresh(List.of(first, second));
+
+ verify(subscriber).refresh();
+ verify(subscriber).onSubscribe(first);
+ verify(subscriber).onSubscribe(second);
+ }
+
+ @Test
+ public void refreshWithEmptyListShouldOnlyClear() {
+ dataRefresh.refresh(Collections.emptyList());
+
+ verify(subscriber).refresh();
+ verify(subscriber, never()).onSubscribe(any());
+ }
+
+ @Test
+ public void refreshShouldTolerateNoSubscribers() {
+ AiProxyApiKeyDataRefresh emptyRefresh =
+ new AiProxyApiKeyDataRefresh(Collections.emptyList());
+
+ assertDoesNotThrow(() ->
emptyRefresh.refresh(Collections.singletonList(new ProxyApiKeyData())));
+ }
+
+ @Test
+ public void refreshJsonShouldUpdateCacheAndNotifySubscribers() {
+ ProxyApiKeyData apiKeyData =
ProxyApiKeyData.builder().proxyApiKey("proxy-key").build();
+ ConfigData<ProxyApiKeyData> config = new ConfigData<>("md5-new",
System.currentTimeMillis(),
+ Collections.singletonList(apiKeyData));
+ JsonObject groupJson =
GsonUtils.getGson().fromJson(GsonUtils.getGson().toJson(config),
JsonObject.class);
+ JsonObject data = new JsonObject();
+ data.add(ConfigGroupEnum.AI_PROXY_API_KEY.name(), groupJson);
+
+ assertTrue(dataRefresh.refresh(data));
+ verify(subscriber).refresh();
+ verify(subscriber).onSubscribe(apiKeyData);
+ assertEquals(config, dataRefresh.cacheConfigData());
+ }
+
+ @Test
+ public void cacheConfigDataShouldBeEmptyBeforeFirstRefresh() {
+ assertFalse(dataRefresh.refresh(new JsonObject()));
+ assertNull(dataRefresh.cacheConfigData());
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DataRefreshFactoryTest.java
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DataRefreshFactoryTest.java
new file mode 100644
index 0000000000..29a38434d2
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-http/src/test/java/org/apache/shenyu/sync/data/http/refresh/DataRefreshFactoryTest.java
@@ -0,0 +1,107 @@
+/*
+ * 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.sync.data.http.refresh;
+
+import com.google.gson.JsonObject;
+import org.apache.shenyu.common.dto.ConfigData;
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.enums.ConfigGroupEnum;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Test cases for {@link DataRefreshFactory}.
+ */
+public final class DataRefreshFactoryTest {
+
+ private final PluginDataSubscriber pluginDataSubscriber =
mock(PluginDataSubscriber.class);
+
+ private final DataRefreshFactory dataRefreshFactory;
+
+ public DataRefreshFactoryTest() {
+ dataRefreshFactory = new DataRefreshFactory(pluginDataSubscriber,
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyList(),
+ Collections.emptyList());
+ }
+
+ @BeforeEach
+ public void clearPluginCache() {
+ AbstractDataRefresh.GROUP_CACHE.remove(ConfigGroupEnum.PLUGIN);
+ }
+
+ @AfterEach
+ public void tearDown() {
+ AbstractDataRefresh.GROUP_CACHE.remove(ConfigGroupEnum.PLUGIN);
+ }
+
+ @Test
+ public void executorShouldReturnFalseWhenNoGroupDataIsPresent() {
+ assertFalse(dataRefreshFactory.executor(new JsonObject()));
+ assertNull(dataRefreshFactory.cacheConfigData(ConfigGroupEnum.PLUGIN));
+ }
+
+ @Test
+ public void executorShouldRefreshRegisteredGroup() {
+ PluginData pluginData =
PluginData.builder().name("sign-plugin").enabled(true).build();
+ ConfigData<PluginData> config = new ConfigData<>("md5-new",
System.currentTimeMillis(),
+ Collections.singletonList(pluginData));
+ JsonObject groupJson =
GsonUtils.getGson().fromJson(GsonUtils.getGson().toJson(config),
JsonObject.class);
+ JsonObject data = new JsonObject();
+ data.add(ConfigGroupEnum.PLUGIN.name(), groupJson);
+
+ assertTrue(dataRefreshFactory.executor(data));
+
+ verify(pluginDataSubscriber).refreshPluginDataAll();
+ verify(pluginDataSubscriber).onSubscribe(pluginData);
+ assertEquals("md5-new",
dataRefreshFactory.cacheConfigData(ConfigGroupEnum.PLUGIN).getMd5());
+ }
+
+ @Test
+ public void executorShouldReturnFalseWhenTheSameConfigIsRepeated() {
+ PluginData pluginData =
PluginData.builder().name("sign-plugin").build();
+ long lastModifyTime = System.currentTimeMillis();
+ ConfigData<PluginData> config = new ConfigData<>("md5-same",
lastModifyTime,
+ Collections.singletonList(pluginData));
+ JsonObject groupJson =
GsonUtils.getGson().fromJson(GsonUtils.getGson().toJson(config),
JsonObject.class);
+ JsonObject data = new JsonObject();
+ data.add(ConfigGroupEnum.PLUGIN.name(), groupJson);
+
+ assertTrue(dataRefreshFactory.executor(data));
+ assertFalse(dataRefreshFactory.executor(data));
+ verify(pluginDataSubscriber).refreshPluginDataAll();
+ verify(pluginDataSubscriber).onSubscribe(pluginData);
+ verify(pluginDataSubscriber, never()).unSubscribe(any());
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AiProxyApiKeyDataHandlerTest.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AiProxyApiKeyDataHandlerTest.java
new file mode 100644
index 0000000000..89e9964037
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/AiProxyApiKeyDataHandlerTest.java
@@ -0,0 +1,124 @@
+/*
+ * 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.plugin.sync.data.websocket.handler;
+
+import com.google.gson.Gson;
+import org.apache.shenyu.common.dto.ProxyApiKeyData;
+import org.apache.shenyu.sync.data.api.AiProxyApiKeyDataSubscriber;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.LinkedList;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Test cases for {@link AiProxyApiKeyDataHandler}.
+ */
+public final class AiProxyApiKeyDataHandlerTest {
+
+ private final List<AiProxyApiKeyDataSubscriber> subscribers;
+
+ private final AiProxyApiKeyDataHandler dataHandler;
+
+ public AiProxyApiKeyDataHandlerTest() {
+ subscribers = new LinkedList<>();
+ subscribers.add(mock(AiProxyApiKeyDataSubscriber.class));
+ subscribers.add(mock(AiProxyApiKeyDataSubscriber.class));
+ dataHandler = new AiProxyApiKeyDataHandler(subscribers);
+ }
+
+ @Test
+ public void testConvert() {
+ ProxyApiKeyData apiKeyData = ProxyApiKeyData.builder()
+ .realApiKey("real-key")
+ .proxyApiKey("proxy-key")
+ .description("description")
+ .enabled(true)
+ .namespaceId("namespace")
+ .selectorId("selector")
+ .build();
+ List<ProxyApiKeyData> sources = Collections.singletonList(apiKeyData);
+
+ List<ProxyApiKeyData> result = dataHandler.convert(new
Gson().toJson(sources));
+
+ assertEquals(sources, result);
+ }
+
+ @Test
+ public void testDoRefresh() {
+ List<ProxyApiKeyData> apiKeyDataList = createFakeApiKeyData(3);
+
+ dataHandler.doRefresh(apiKeyDataList);
+
+ subscribers.forEach(subscriber -> verify(subscriber).refresh());
+ apiKeyDataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).onSubscribe(data)));
+ }
+
+ @Test
+ public void testDoUpdate() {
+ List<ProxyApiKeyData> apiKeyDataList = createFakeApiKeyData(4);
+
+ dataHandler.doUpdate(apiKeyDataList);
+
+ apiKeyDataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).onSubscribe(data)));
+ }
+
+ @Test
+ public void testDoDelete() {
+ List<ProxyApiKeyData> apiKeyDataList = createFakeApiKeyData(3);
+
+ dataHandler.doDelete(apiKeyDataList);
+
+ apiKeyDataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).unSubscribe(data)));
+ }
+
+ @Test
+ public void testNullDataListShouldBeIgnored() {
+ assertDoesNotThrow(() -> dataHandler.doUpdate(null));
+ assertDoesNotThrow(() -> dataHandler.doDelete(null));
+ }
+
+ @Test
+ public void testNullSubscribersShouldBeIgnored() {
+ AiProxyApiKeyDataHandler handlerWithoutSubscribers = new
AiProxyApiKeyDataHandler(null);
+
+ assertDoesNotThrow(() ->
handlerWithoutSubscribers.doRefresh(createFakeApiKeyData(1)));
+ assertDoesNotThrow(() ->
handlerWithoutSubscribers.doUpdate(createFakeApiKeyData(1)));
+ assertDoesNotThrow(() ->
handlerWithoutSubscribers.doDelete(createFakeApiKeyData(1)));
+ }
+
+ private List<ProxyApiKeyData> createFakeApiKeyData(final int count) {
+ List<ProxyApiKeyData> result = new LinkedList<>();
+ for (int i = 1; i <= count; i++) {
+ result.add(ProxyApiKeyData.builder()
+ .realApiKey("real-key-" + i)
+ .proxyApiKey("proxy-key-" + i)
+ .enabled(true)
+ .build());
+ }
+ return result;
+ }
+}
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/ProxySelectorDataHandlerTest.java
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/ProxySelectorDataHandlerTest.java
new file mode 100644
index 0000000000..8bb30a2d91
--- /dev/null
+++
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/ProxySelectorDataHandlerTest.java
@@ -0,0 +1,116 @@
+/*
+ * 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.plugin.sync.data.websocket.handler;
+
+import com.google.gson.Gson;
+import org.apache.shenyu.common.dto.ProxySelectorData;
+import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedList;
+import java.util.List;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Test cases for {@link ProxySelectorDataHandler}.
+ */
+public final class ProxySelectorDataHandlerTest {
+
+ private final List<ProxySelectorDataSubscriber> subscribers;
+
+ private final ProxySelectorDataHandler dataHandler;
+
+ public ProxySelectorDataHandlerTest() {
+ subscribers = new LinkedList<>();
+ subscribers.add(mock(ProxySelectorDataSubscriber.class));
+ subscribers.add(mock(ProxySelectorDataSubscriber.class));
+ dataHandler = new ProxySelectorDataHandler(subscribers);
+ }
+
+ @Test
+ public void testConvert() {
+ ProxySelectorData proxySelectorData = proxySelectorData("selector-1",
20000);
+ List<ProxySelectorData> sources =
Collections.singletonList(proxySelectorData);
+
+ List<ProxySelectorData> result = dataHandler.convert(new
Gson().toJson(sources));
+
+ assertEquals(1, result.size());
+ assertEquals("selector-1", result.get(0).getName());
+ assertEquals(20000, result.get(0).getForwardPort());
+ }
+
+ @Test
+ public void testDoRefresh() {
+ List<ProxySelectorData> dataList = createFakeProxySelectorData(3);
+
+ dataHandler.doRefresh(dataList);
+
+ subscribers.forEach(subscriber -> verify(subscriber).refresh());
+ dataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).onSubscribe(data)));
+ }
+
+ @Test
+ public void testDoUpdate() {
+ List<ProxySelectorData> dataList = createFakeProxySelectorData(4);
+
+ dataHandler.doUpdate(dataList);
+
+ dataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).onSubscribe(data)));
+ }
+
+ @Test
+ public void testDoDelete() {
+ List<ProxySelectorData> dataList = createFakeProxySelectorData(3);
+
+ dataHandler.doDelete(dataList);
+
+ dataList.forEach(data ->
+ subscribers.forEach(subscriber ->
verify(subscriber).unSubscribe(data)));
+ }
+
+ @Test
+ public void testDoDeleteShouldTolerateNullItems() {
+ ProxySelectorData proxySelectorData = proxySelectorData("selector-1",
20000);
+
+ dataHandler.doDelete(Arrays.asList(proxySelectorData, null));
+
+ subscribers.forEach(subscriber ->
verify(subscriber).unSubscribe(proxySelectorData));
+ }
+
+ private List<ProxySelectorData> createFakeProxySelectorData(final int
count) {
+ List<ProxySelectorData> result = new LinkedList<>();
+ for (int i = 1; i <= count; i++) {
+ result.add(proxySelectorData("selector-" + i, 20000 + i));
+ }
+ return result;
+ }
+
+ private ProxySelectorData proxySelectorData(final String name, final int
forwardPort) {
+ ProxySelectorData proxySelectorData = new ProxySelectorData();
+ proxySelectorData.setName(name);
+ proxySelectorData.setForwardPort(forwardPort);
+ return proxySelectorData;
+ }
+}