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;
+    }
+}

Reply via email to