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

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


The following commit(s) were added to refs/heads/master by this push:
     new b896d7241b [ISSUE #7283] Gate gateway readiness on initial websocket 
synchronization. (#7319)
b896d7241b is described below

commit b896d7241b93f67d92db2475d5bc64f61ea05564
Author: Jerry聊AI <[email protected]>
AuthorDate: Wed Sep 30 16:19:59 2026 +0800

    [ISSUE #7283] Gate gateway readiness on initial websocket synchronization. 
(#7319)
    
    * feature (sync-data-websocket) : gate readiness on initial configuration 
synchronization.
    
    * test(sync-data-websocket): cover readiness with batch refresh callbacks
    
    * fix (sync-data-websocket) : validate sync requests and bound incremental 
tracking.
    
    ---------
    
    Co-authored-by: wy471x <[email protected]>
---
 .../listener/websocket/WebsocketCollector.java     |  59 ++++-
 .../listener/websocket/WebsocketCollectorTest.java |  78 +++++++
 .../shenyu/common/dto/WebsocketSyncFrame.java      |  72 ++++++
 .../common/utils/InitialSyncApplication.java       |  83 +++++++
 .../common/utils/InitialSyncApplicationTest.java   |  56 +++++
 .../k8s/sync/shenyu-bootstrap-websocket.yml        |   5 +-
 shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-cm.yml  |  12 +
 .../base/cache/CommonPluginDataSubscriber.java     |   4 +
 .../base/cache/CommonPluginDataSubscriberTest.java |  40 ++++
 .../pom.xml                                        |   5 +
 .../websocket/WebsocketSyncDataConfiguration.java  |   5 +-
 .../WebsocketSyncHealthConfiguration.java          |  68 ++++++
 .../data/websocket/WebsocketSyncHealthGroup.java   |  69 ++++++
 .../WebsocketSyncDataConfigurationTest.java        |  45 +++-
 .../data/websocket/WebsocketSyncDataService.java   |  15 +-
 .../data/websocket/client/InitialSyncState.java    | 153 ++++++++++++
 .../websocket/client/ShenyuWebsocketClient.java    |  65 +++++-
 .../data/websocket/config/WebsocketConfig.java     |  23 +-
 .../websocket/handler/WebsocketDataHandler.java    |  20 +-
 .../client/InitialSyncConnectionTest.java          | 143 ++++++++++++
 .../websocket/client/InitialSyncStateTest.java     | 256 +++++++++++++++++++++
 .../client/ShenyuWebsocketClientTest.java          |  33 +++
 .../data/websocket/config/WebsocketConfigTest.java |   5 +-
 23 files changed, 1291 insertions(+), 23 deletions(-)

diff --git 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
index c4d7b578c5..968f17bc7c 100644
--- 
a/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
+++ 
b/shenyu-admin/src/main/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollector.java
@@ -31,6 +31,7 @@ import org.apache.shenyu.admin.utils.ThreadLocalUtils;
 import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.constant.InstanceTypeConstants;
 import org.apache.shenyu.common.constant.RunningModeConstants;
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
 import org.apache.shenyu.common.enums.DataEventTypeEnum;
 import org.apache.shenyu.common.enums.RunningModeEnum;
 import org.apache.shenyu.common.exception.ShenyuException;
@@ -53,6 +54,7 @@ import java.util.Objects;
 import java.util.Optional;
 import java.util.Queue;
 import java.util.Set;
+import java.util.UUID;
 import java.util.concurrent.CopyOnWriteArraySet;
 
 /**
@@ -78,6 +80,8 @@ public class WebsocketCollector {
     private static final Map<Session, String> SESSION_NAMESPACE_IDS = 
Maps.newConcurrentMap();
     
     private static final String SESSION_KEY = "sessionKey";
+
+    private static final ThreadLocal<InitialSync> INITIAL_SYNC = new 
ThreadLocal<>();
     
     /**
      * On open.
@@ -155,6 +159,10 @@ public class WebsocketCollector {
      */
     @OnMessage
     public void onMessage(final String message, final Session session) {
+        if (message.startsWith(WebsocketSyncFrame.REQUEST_PREFIX)) {
+            
initialSync(message.substring(WebsocketSyncFrame.REQUEST_PREFIX.length()), 
session);
+            return;
+        }
         if (!Objects.equals(message, DataEventTypeEnum.MYSELF.name())
                 && !Objects.equals(message, 
DataEventTypeEnum.RUNNING_MODE.name())
                 && !message.contains("bootstrapInstanceInfo")) {
@@ -317,8 +325,43 @@ public class WebsocketCollector {
         
     }
     
+    private void initialSync(final String requestId, final Session session) {
+        try {
+            UUID.fromString(requestId);
+        } catch (IllegalArgumentException ex) {
+            LOG.warn("Ignoring initial synchronization request with an invalid 
UUID");
+            return;
+        }
+        ClusterProperties properties = 
SpringBeanUtils.getInstance().getBean(ClusterProperties.class);
+        if (properties.isEnabled()
+                && 
!SpringBeanUtils.getInstance().getBean(ClusterSelectMasterService.class).isMaster())
 {
+            return;
+        }
+        InitialSync sync = new InitialSync(session, requestId);
+        try {
+            INITIAL_SYNC.set(sync);
+            ThreadLocalUtils.put(SESSION_KEY, session);
+            boolean success = 
SpringBeanUtils.getInstance().getBean(SyncDataService.class)
+                    .syncAllByNamespaceId(DataEventTypeEnum.MYSELF, 
getNamespaceId(session));
+            if (success && (!properties.isEnabled()
+                    || 
SpringBeanUtils.getInstance().getBean(ClusterSelectMasterService.class).isMaster()))
 {
+                SESSION_SEND_QUEUES.computeIfAbsent(session, 
SessionSendQueue::new)
+                        .send(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(requestId, sync.sequence, null)));
+            }
+        } finally {
+            INITIAL_SYNC.remove();
+            ThreadLocalUtils.clear();
+        }
+    }
+
     private static void sendMessageBySession(final Session session, final 
String message) {
-        SESSION_SEND_QUEUES.computeIfAbsent(session, 
SessionSendQueue::new).send(message);
+        InitialSync sync = INITIAL_SYNC.get();
+        if (Objects.nonNull(sync) && sync.session == session) {
+            SESSION_SEND_QUEUES.computeIfAbsent(session, SessionSendQueue::new)
+                    .send(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(sync.requestId, sync.sequence++, message)));
+        } else {
+            SESSION_SEND_QUEUES.computeIfAbsent(session, 
SessionSendQueue::new).send(message);
+        }
     }
 
     private static void removeSessionSendQueue(final Session session) {
@@ -369,6 +412,20 @@ public class WebsocketCollector {
         }
     }
 
+    private static final class InitialSync {
+
+        private final Session session;
+
+        private final String requestId;
+
+        private int sequence;
+
+        private InitialSync(final Session session, final String requestId) {
+            this.session = session;
+            this.requestId = requestId;
+        }
+    }
+
     private static final class SessionSendQueue {
 
         private final Session session;
diff --git 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
index 0c3cba451c..bf038a24f7 100644
--- 
a/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
+++ 
b/shenyu-admin/src/test/java/org/apache/shenyu/admin/listener/websocket/WebsocketCollectorTest.java
@@ -29,8 +29,10 @@ import org.apache.shenyu.admin.spring.SpringBeanUtils;
 import org.apache.shenyu.admin.utils.ThreadLocalUtils;
 import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.constant.InstanceTypeConstants;
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
 import org.apache.shenyu.common.enums.DataEventTypeEnum;
 import org.apache.shenyu.common.exception.ShenyuException;
+import org.apache.shenyu.common.utils.GsonUtils;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.BeforeAll;
 import org.junit.jupiter.api.BeforeEach;
@@ -51,7 +53,9 @@ import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Set;
+import java.util.UUID;
 
+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;
@@ -149,6 +153,80 @@ public final class WebsocketCollectorTest {
         websocketCollector.onClose(session);
     }
 
+    @Test
+    void testInvalidInitialSyncRequestIsIgnored() {
+        
when(SpringBeanUtils.getInstance().getBean(ClusterProperties.class)).thenReturn(new
 ClusterProperties());
+        
when(SpringBeanUtils.getInstance().getBean(SyncDataService.class)).thenReturn(syncDataService);
+        RemoteEndpoint.Async async = mockSuccessfulAsyncRemote(session);
+        for (String id : new String[]{"", "not-a-uuid", 
"00000000-0000-0000-0000-00000000000z"}) {
+            assertDoesNotThrow(() -> 
websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + id, session));
+        }
+        verify(syncDataService, never()).syncAllByNamespaceId(any(), 
anyString());
+        verify(async, never()).sendText(anyString(), any(SendHandler.class));
+        assertNull(ThreadLocalUtils.get("sessionKey"));
+        when(syncDataService.syncAllByNamespaceId(DataEventTypeEnum.MYSELF, 
Constants.SYS_DEFAULT_NAMESPACE_ID)).thenReturn(true);
+        websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + 
UUID.randomUUID(), session);
+        verify(async).sendText(anyString(), any(SendHandler.class));
+    }
+
+    @Test
+    void testInitialSyncFramesAndEmptyCompletion() {
+        
when(SpringBeanUtils.getInstance().getBean(ClusterProperties.class)).thenReturn(new
 ClusterProperties());
+        
when(SpringBeanUtils.getInstance().getBean(SyncDataService.class)).thenReturn(syncDataService);
+        final RemoteEndpoint.Async async = mockSuccessfulAsyncRemote(session);
+        websocketCollector.onOpen(session);
+        String id = UUID.randomUUID().toString();
+        when(syncDataService.syncAllByNamespaceId(DataEventTypeEnum.MYSELF, 
Constants.SYS_DEFAULT_NAMESPACE_ID))
+                .thenAnswer(invocation -> {
+                    
WebsocketCollector.send(Constants.SYS_DEFAULT_NAMESPACE_ID, "configuration", 
DataEventTypeEnum.MYSELF);
+                    return true;
+                });
+        websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + id, 
session);
+        ArgumentCaptor<String> messages = 
ArgumentCaptor.forClass(String.class);
+        verify(async, times(2)).sendText(messages.capture(), 
any(SendHandler.class));
+        WebsocketSyncFrame data = 
GsonUtils.getInstance().fromJson(messages.getAllValues().get(0), 
WebsocketSyncFrame.class);
+        final WebsocketSyncFrame end = 
GsonUtils.getInstance().fromJson(messages.getAllValues().get(1), 
WebsocketSyncFrame.class);
+        assertEquals(id, data.getRequestId());
+        assertEquals(0, data.getSequence());
+        assertEquals("configuration", data.getPayload());
+        assertEquals(id, end.getRequestId());
+        assertEquals(1, end.getSequence());
+        assertNull(end.getPayload());
+        when(syncDataService.syncAllByNamespaceId(DataEventTypeEnum.MYSELF, 
Constants.SYS_DEFAULT_NAMESPACE_ID)).thenReturn(true);
+        websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + 
UUID.randomUUID(), session);
+        verify(async, times(3)).sendText(messages.capture(), 
any(SendHandler.class));
+        WebsocketSyncFrame empty = 
GsonUtils.getInstance().fromJson(messages.getValue(), WebsocketSyncFrame.class);
+        assertEquals(0, empty.getSequence());
+        assertNull(empty.getPayload());
+        websocketCollector.onClose(session);
+    }
+
+    @Test
+    void testFailedInitialSyncDoesNotSendCompletion() {
+        
when(SpringBeanUtils.getInstance().getBean(ClusterProperties.class)).thenReturn(new
 ClusterProperties());
+        
when(SpringBeanUtils.getInstance().getBean(SyncDataService.class)).thenReturn(syncDataService);
+        RemoteEndpoint.Async async = mockSuccessfulAsyncRemote(session);
+        when(syncDataService.syncAllByNamespaceId(DataEventTypeEnum.MYSELF, 
Constants.SYS_DEFAULT_NAMESPACE_ID))
+                .thenThrow(new IllegalStateException("snapshot unavailable"));
+        assertThrows(IllegalStateException.class,
+                () -> 
websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + 
UUID.randomUUID(), session));
+        verify(async, never()).sendText(anyString(), any(SendHandler.class));
+        assertNull(ThreadLocalUtils.get("sessionKey"));
+    }
+
+    @Test
+    void testFollowerCannotCompleteInitialSync() {
+        ClusterProperties properties = new ClusterProperties();
+        properties.setEnabled(true);
+        
when(SpringBeanUtils.getInstance().getBean(ClusterProperties.class)).thenReturn(properties);
+        ClusterSelectMasterService master = 
mock(ClusterSelectMasterService.class);
+        
when(SpringBeanUtils.getInstance().getBean(ClusterSelectMasterService.class)).thenReturn(master);
+        RemoteEndpoint.Async async = mockSuccessfulAsyncRemote(session);
+        websocketCollector.onMessage(WebsocketSyncFrame.REQUEST_PREFIX + 
UUID.randomUUID(), session);
+        verify(syncDataService, never()).syncAllByNamespaceId(any(), 
anyString());
+        verify(async, never()).sendText(anyString(), any(SendHandler.class));
+    }
+
     @Test
     void testOnOpenWithBlankNamespaceIdThrows() {
         Map<String, Object> userProperties = new HashMap<>();
diff --git 
a/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketSyncFrame.java
 
b/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketSyncFrame.java
new file mode 100644
index 0000000000..683750fe76
--- /dev/null
+++ 
b/shenyu-common/src/main/java/org/apache/shenyu/common/dto/WebsocketSyncFrame.java
@@ -0,0 +1,72 @@
+/*
+ * 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.common.dto;
+
+/**
+ * Opt-in initial synchronization frame. A null payload marks the end of an 
attempt.
+ */
+public final class WebsocketSyncFrame {
+
+    public static final String REQUEST_PREFIX = "MYSELF_V1:";
+
+    public static final String EVENT_TYPE = "INITIAL_SYNC_V1";
+
+    private final String eventType = EVENT_TYPE;
+
+    private final String requestId;
+
+    private final int sequence;
+
+    private final String payload;
+
+    /**
+     * Create a frame.
+     * @param requestId connection-scoped request identifier
+     * @param sequence number of preceding configuration frames
+     * @param payload original configuration message, or null for completion
+     */
+    public WebsocketSyncFrame(final String requestId, final int sequence, 
final String payload) {
+        this.requestId = requestId;
+        this.sequence = sequence;
+        this.payload = payload;
+    }
+
+    /**
+     * Get the request identifier.
+     * @return request identifier
+     */
+    public String getRequestId() {
+        return requestId;
+    }
+
+    /**
+     * Get the sequence.
+     * @return sequence
+     */
+    public int getSequence() {
+        return sequence;
+    }
+
+    /**
+     * Get the payload.
+     * @return payload
+     */
+    public String getPayload() {
+        return payload;
+    }
+}
diff --git 
a/shenyu-common/src/main/java/org/apache/shenyu/common/utils/InitialSyncApplication.java
 
b/shenyu-common/src/main/java/org/apache/shenyu/common/utils/InitialSyncApplication.java
new file mode 100644
index 0000000000..e4244fede0
--- /dev/null
+++ 
b/shenyu-common/src/main/java/org/apache/shenyu/common/utils/InitialSyncApplication.java
@@ -0,0 +1,83 @@
+/*
+ * 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.common.utils;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionStage;
+
+/**
+ * Tracks configuration application, including asynchronous work registered by 
subscribers.
+ * Subscribers must register their completion stage before returning from the 
callback.
+ * The stage must cover all deferred configuration work, including nested 
tasks. Ongoing
+ * upstream health checks and request-time connections are outside this 
boundary.
+ */
+public final class InitialSyncApplication {
+
+    private static final ThreadLocal<List<CompletableFuture<?>>> ACTIVE = new 
ThreadLocal<>();
+
+    private InitialSyncApplication() {
+    }
+
+    /**
+     * Run callbacks and collect their asynchronous completion stages without 
blocking the socket.
+     * @param action callback
+     * @return completion of the callback and all registered application work
+     */
+    public static CompletableFuture<Void> run(final Runnable action) {
+        List<CompletableFuture<?>> previous = ACTIVE.get();
+        List<CompletableFuture<?>> pending = new ArrayList<>();
+        ACTIVE.set(pending);
+        try {
+            action.run();
+            CompletableFuture<Void> completion = 
CompletableFuture.allOf(pending.toArray(new CompletableFuture<?>[0]));
+            if (Objects.nonNull(previous)) {
+                previous.add(completion);
+            }
+            return completion;
+        } finally {
+            if (Objects.nonNull(previous)) {
+                ACTIVE.set(previous);
+            } else {
+                ACTIVE.remove();
+            }
+        }
+    }
+
+    /**
+     * Register deferred configuration application during the subscriber 
callback.
+     * Outside initial synchronization this leaves the existing asynchronous 
behavior unchanged.
+     * @param completion completion stage, exceptional completion prevents 
readiness
+     */
+    public static void register(final CompletionStage<?> completion) {
+        List<CompletableFuture<?>> pending = ACTIVE.get();
+        if (Objects.nonNull(pending)) {
+            
pending.add(Objects.requireNonNull(completion).toCompletableFuture());
+        }
+    }
+
+    /**
+     * Whether failures must propagate to the synchronization caller.
+     * @return active
+     */
+    public static boolean isActive() {
+        return Objects.nonNull(ACTIVE.get());
+    }
+}
diff --git 
a/shenyu-common/src/test/java/org/apache/shenyu/common/utils/InitialSyncApplicationTest.java
 
b/shenyu-common/src/test/java/org/apache/shenyu/common/utils/InitialSyncApplicationTest.java
new file mode 100644
index 0000000000..ae5121bbd7
--- /dev/null
+++ 
b/shenyu-common/src/test/java/org/apache/shenyu/common/utils/InitialSyncApplicationTest.java
@@ -0,0 +1,56 @@
+/*
+ * 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.common.utils;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.concurrent.CompletableFuture;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class InitialSyncApplicationTest {
+
+    @Test
+    void testNestedApplicationAndContextCleanup() {
+        CompletableFuture<Void> deferred = new CompletableFuture<>();
+        CompletableFuture<Void> result = InitialSyncApplication.run(() -> {
+            assertTrue(InitialSyncApplication.isActive());
+            InitialSyncApplication.run(() -> 
InitialSyncApplication.register(deferred));
+            assertTrue(InitialSyncApplication.isActive());
+        });
+        assertFalse(InitialSyncApplication.isActive());
+        assertFalse(result.isDone());
+        deferred.complete(null);
+        assertTrue(result.isDone());
+        assertFalse(result.isCompletedExceptionally());
+    }
+
+    @Test
+    void testFailureDoesNotLeakContext() {
+        assertThrows(IllegalStateException.class, () -> 
InitialSyncApplication.run(() -> {
+            throw new IllegalStateException("failed callback");
+        }));
+        assertFalse(InitialSyncApplication.isActive());
+        CompletableFuture<Void> deferred = new CompletableFuture<>();
+        CompletableFuture<Void> result = InitialSyncApplication.run(() -> 
InitialSyncApplication.register(deferred));
+        deferred.completeExceptionally(new IllegalStateException("failed 
application"));
+        assertTrue(result.isCompletedExceptionally());
+    }
+}
diff --git a/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-bootstrap-websocket.yml 
b/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-bootstrap-websocket.yml
index b88f5e44dc..c5732f151d 100644
--- a/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-bootstrap-websocket.yml
+++ b/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-bootstrap-websocket.yml
@@ -43,7 +43,7 @@ spec:
             failureThreshold: 3
             httpGet:
               port: 9195
-              path: /actuator/health
+              path: /actuator/health/liveness
           readinessProbe:
             initialDelaySeconds: 30
             periodSeconds: 10
@@ -52,7 +52,7 @@ spec:
             failureThreshold: 3
             httpGet:
               port: 9195
-              path: /actuator/health
+              path: /actuator/health/readiness
           env:
             - name: TZ
               value: Asia/Beijing
@@ -93,4 +93,3 @@ spec:
       port: 9195
       targetPort: 9195
       nodePort: 31195
-
diff --git a/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-cm.yml 
b/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-cm.yml
index 1dc257d5a4..176b4ea1c8 100644
--- a/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-cm.yml
+++ b/shenyu-e2e/shenyu-e2e-case/k8s/sync/shenyu-cm.yml
@@ -519,9 +519,21 @@ data:
             secretKey:
 
   application-bootstrap-sync-websocket.yml: |
+    management:
+      endpoint:
+        health:
+          probes:
+            enabled: true
+          group:
+            readiness:
+              include: readinessState,websocketSync
+            liveness:
+              include: livenessState
     shenyu:
       sync:
         websocket:
+          # Requires an Admin supporting the initial synchronization protocol.
+          initial-sync-readiness: true
           urls: ws://shenyu-admin:9095/websocket
           token: shenyu-sync-token
           allowOrigin: ws://localhost:9195
diff --git 
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
 
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
index 9f2d6f84f9..0ff78bc865 100644
--- 
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
+++ 
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
@@ -25,6 +25,7 @@ import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.common.dto.SelectorData;
 import org.apache.shenyu.common.enums.DataEventTypeEnum;
 import org.apache.shenyu.common.enums.PluginHandlerEventEnum;
+import org.apache.shenyu.common.utils.InitialSyncApplication;
 import org.apache.shenyu.common.utils.JsonUtils;
 import org.apache.shenyu.common.utils.MapUtils;
 import org.apache.shenyu.plugin.base.handler.PluginDataHandler;
@@ -196,6 +197,9 @@ public class CommonPluginDataSubscriber implements 
PluginDataSubscriber {
                         .ifPresent(data -> removeCacheData(classData));
             }
         } catch (Exception e) {
+            if (InitialSyncApplication.isActive()) {
+                throw new IllegalStateException("Initial configuration 
application failed", e);
+            }
             LOG.error("subscribe data handler error, classData: {}, dataType: 
{}", JsonUtils.toJson(classData), dataType, e);
         }
     }
diff --git 
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
 
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
index aa191d558c..a2182f319c 100644
--- 
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
@@ -24,6 +24,7 @@ import org.apache.shenyu.common.dto.PluginData;
 import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.common.dto.SelectorData;
 import org.apache.shenyu.common.enums.PluginHandlerEventEnum;
+import org.apache.shenyu.common.utils.InitialSyncApplication;
 import org.apache.shenyu.plugin.api.utils.SpringBeanUtils;
 import org.apache.shenyu.plugin.base.handler.PluginDataHandler;
 import org.junit.jupiter.api.AfterEach;
@@ -39,13 +40,17 @@ import 
org.springframework.context.ConfigurableApplicationContext;
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.CompletableFuture;
 
+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.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doThrow;
@@ -104,6 +109,41 @@ public final class CommonPluginDataSubscriberTest {
         assertEquals(pluginData, 
baseDataCache.obtainPluginData(pluginData.getName()));
     }
 
+    @Test
+    void 
testInitialSyncPropagatesHandlerFailureWithoutChangingLegacyBehavior() {
+        PluginDataHandler handler = mock(PluginDataHandler.class);
+        org.mockito.Mockito.when(handler.pluginNamed()).thenReturn(mockName1);
+        PluginData data = PluginData.builder().name(mockName1).build();
+        doThrow(new IllegalStateException("handler 
failed")).when(handler).handlerPlugin(data);
+        
commonPluginDataSubscriber.putExtendPluginDataHandler(List.of(handler));
+        assertThrows(IllegalStateException.class, () -> 
InitialSyncApplication.run(() -> commonPluginDataSubscriber.onSubscribe(data)));
+        assertFalse(InitialSyncApplication.isActive());
+        assertDoesNotThrow(() -> commonPluginDataSubscriber.onSubscribe(data));
+    }
+
+    @Test
+    void testInitialSyncTracksDeferredRefreshApplication() {
+        PluginData data = PluginData.builder().name("divide").build();
+        CompletableFuture<Void> deferred = new CompletableFuture<>();
+        doAnswer(invocation -> {
+            InitialSyncApplication.register(deferred);
+            return null;
+        }).when(handler).handlerPlugin(data);
+        CompletableFuture<Void> application = InitialSyncApplication.run(() -> 
commonPluginDataSubscriber.onPluginRefresh(List.of(data)));
+        assertFalse(application.isDone());
+        assertFalse(InitialSyncApplication.isActive());
+        deferred.completeExceptionally(new IllegalStateException("deferred 
refresh failed"));
+        assertTrue(application.isCompletedExceptionally());
+    }
+
+    @Test
+    void testInitialSyncPropagatesRefreshFailure() {
+        PluginData data = PluginData.builder().name("divide").build();
+        doThrow(new IllegalStateException("refresh 
failed")).when(handler).handlerPlugin(data);
+        assertThrows(IllegalStateException.class, () -> 
InitialSyncApplication.run(() -> 
commonPluginDataSubscriber.onPluginRefresh(List.of(data))));
+        assertFalse(InitialSyncApplication.isActive());
+    }
+
     @Test
     public void testUnSubscribe() {
         baseDataCache.cleanPluginData();
diff --git 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/pom.xml
 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/pom.xml
index 3c13a18add..6482e4c817 100644
--- 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/pom.xml
+++ 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/pom.xml
@@ -26,6 +26,11 @@
     <artifactId>shenyu-spring-boot-starter-sync-data-websocket</artifactId>
 
     <dependencies>
+        <dependency>
+            <groupId>org.springframework.boot</groupId>
+            <artifactId>spring-boot-actuator-autoconfigure</artifactId>
+            <optional>true</optional>
+        </dependency>
         <dependency>
             <groupId>org.apache.shenyu</groupId>
             <artifactId>shenyu-sync-data-websocket</artifactId>
diff --git 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfiguration.java
 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfiguration.java
index 1af3f47dc4..f6047ff711 100644
--- 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfiguration.java
+++ 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfiguration.java
@@ -26,7 +26,6 @@ import 
org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
 import org.apache.shenyu.sync.data.api.MetaDataSubscriber;
 import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
 import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
-import org.apache.shenyu.sync.data.api.SyncDataService;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.ObjectProvider;
@@ -36,6 +35,7 @@ import 
org.springframework.boot.autoconfigure.web.ServerProperties;
 import org.springframework.boot.context.properties.ConfigurationProperties;
 import org.springframework.context.annotation.Bean;
 import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.Import;
 
 import java.util.Collections;
 import java.util.List;
@@ -44,6 +44,7 @@ import java.util.List;
  * Websocket sync data configuration for spring boot.
  */
 @Configuration
+@Import(WebsocketSyncHealthConfiguration.class)
 @ConditionalOnClass(WebsocketSyncDataService.class)
 @ConditionalOnProperty(prefix = "shenyu.sync.websocket", name = "urls")
 public class WebsocketSyncDataConfiguration {
@@ -65,7 +66,7 @@ public class WebsocketSyncDataConfiguration {
      * @return the sync data service
      */
     @Bean
-    public SyncDataService websocketSyncDataService(
+    public WebsocketSyncDataService websocketSyncDataService(
             final ObjectProvider<WebsocketConfig> websocketConfig,
             final ShenyuConfig shenyuConfig,
             final ObjectProvider<PluginDataSubscriber> pluginSubscriber,
diff --git 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthConfiguration.java
 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthConfiguration.java
new file mode 100644
index 0000000000..642fb3268f
--- /dev/null
+++ 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthConfiguration.java
@@ -0,0 +1,68 @@
+/*
+ * 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.springboot.starter.sync.data.websocket;
+
+import org.apache.shenyu.plugin.sync.data.websocket.WebsocketSyncDataService;
+import org.springframework.boot.actuate.health.Health;
+import org.springframework.boot.actuate.health.HealthEndpointGroup;
+import org.springframework.boot.actuate.health.HealthEndpointGroups;
+import 
org.springframework.boot.actuate.health.HealthEndpointGroupsPostProcessor;
+import org.springframework.boot.actuate.health.HealthIndicator;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * Optional synchronization health contributor. Include websocketSync in the 
readiness group only.
+ */
+@Configuration(proxyBeanMethods = false)
+@ConditionalOnClass(HealthIndicator.class)
+@ConditionalOnProperty(prefix = "shenyu.sync.websocket", name = 
"initial-sync-readiness", havingValue = "true")
+public class WebsocketSyncHealthConfiguration {
+
+    /**
+     * Restrict synchronization health to readiness, preserving all other 
group members.
+     * This also protects deployments still probing the aggregate endpoint for 
liveness.
+     * @return health group customizer
+     */
+    @Bean
+    public HealthEndpointGroupsPostProcessor websocketSyncHealthGroups() {
+        return groups -> {
+            Map<String, HealthEndpointGroup> customized = new HashMap<>();
+            groups.getNames().forEach(name -> customized.put(name,
+                    new WebsocketSyncHealthGroup(groups.get(name), 
"readiness".equals(name))));
+            customized.putIfAbsent("readiness", new 
WebsocketSyncHealthGroup(groups.getPrimary(), true));
+            return HealthEndpointGroups.of(new 
WebsocketSyncHealthGroup(groups.getPrimary(), false), customized);
+        };
+    }
+
+    /**
+     * Create the startup synchronization health indicator.
+     * @param service synchronization service
+     * @return health indicator
+     */
+    @Bean
+    public HealthIndicator websocketSyncHealthIndicator(final 
WebsocketSyncDataService service) {
+        return () -> service.isInitialSyncReady() ? Health.up().build()
+                : Health.outOfService().withDetail("reason", "Initial 
WebSocket synchronization has not completed; a compatible Admin is 
required").build();
+    }
+}
diff --git 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthGroup.java
 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthGroup.java
new file mode 100644
index 0000000000..cdad64c997
--- /dev/null
+++ 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/main/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncHealthGroup.java
@@ -0,0 +1,69 @@
+/*
+ * 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.springboot.starter.sync.data.websocket;
+
+import org.springframework.boot.actuate.endpoint.SecurityContext;
+import org.springframework.boot.actuate.health.AdditionalHealthEndpointPath;
+import org.springframework.boot.actuate.health.HealthEndpointGroup;
+import org.springframework.boot.actuate.health.HttpCodeStatusMapper;
+import org.springframework.boot.actuate.health.StatusAggregator;
+
+/**
+ * Adds synchronization to readiness while excluding it from aggregate and 
liveness health.
+ */
+final class WebsocketSyncHealthGroup implements HealthEndpointGroup {
+
+    private final HealthEndpointGroup delegate;
+
+    private final boolean readiness;
+
+    WebsocketSyncHealthGroup(final HealthEndpointGroup delegate, final boolean 
readiness) {
+        this.delegate = delegate;
+        this.readiness = readiness;
+    }
+
+    @Override
+    public boolean isMember(final String name) {
+        return "websocketSync".equals(name) ? readiness : 
delegate.isMember(name);
+    }
+
+    @Override
+    public boolean showComponents(final SecurityContext securityContext) {
+        return delegate.showComponents(securityContext);
+    }
+
+    @Override
+    public boolean showDetails(final SecurityContext securityContext) {
+        return delegate.showDetails(securityContext);
+    }
+
+    @Override
+    public StatusAggregator getStatusAggregator() {
+        return delegate.getStatusAggregator();
+    }
+
+    @Override
+    public HttpCodeStatusMapper getHttpCodeStatusMapper() {
+        return delegate.getHttpCodeStatusMapper();
+    }
+
+    @Override
+    public AdditionalHealthEndpointPath getAdditionalPath() {
+        return delegate.getAdditionalPath();
+    }
+}
diff --git 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/test/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfigurationTest.java
 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/test/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfigurationTest.java
index 2f2a766315..3f1bf71841 100644
--- 
a/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/test/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfigurationTest.java
+++ 
b/shenyu-spring-boot-starter/shenyu-spring-boot-starter-sync-data-center/shenyu-spring-boot-starter-sync-data-websocket/src/test/java/org/apache/shenyu/springboot/starter/sync/data/websocket/WebsocketSyncDataConfigurationTest.java
@@ -26,13 +26,23 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.actuate.health.HealthEndpoint;
+import org.springframework.boot.actuate.health.Status;
+import org.springframework.boot.availability.AvailabilityChangeEvent;
+import org.springframework.boot.availability.ReadinessState;
+import org.springframework.context.ApplicationContext;
 import org.springframework.boot.test.context.SpringBootTest;
 import org.springframework.boot.test.mock.mockito.MockBean;
 import org.springframework.test.context.junit.jupiter.SpringExtension;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.atomic.AtomicBoolean;
 
 import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.Matchers.is;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 
 /**
  * Test case for {@link WebsocketSyncDataConfiguration}.
@@ -45,7 +55,11 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
         },
         webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT,
         properties = {
-                "shenyu.sync.websocket.urls=ws://localhost:9095/websocket"
+                "shenyu.sync.websocket.urls=ws://localhost:9095/websocket",
+                "shenyu.sync.websocket.initial-sync-readiness=true",
+                "management.endpoint.health.probes.enabled=true",
+                
"management.endpoint.health.group.readiness.include=readinessState",
+                
"management.endpoint.health.group.liveness.include=livenessState"
         })
 @EnableAutoConfiguration
 @MockBean(PluginDataSubscriber.class)
@@ -56,6 +70,35 @@ public final class WebsocketSyncDataConfigurationTest {
 
     @Autowired
     private WebsocketSyncDataService websocketSyncDataService;
+
+    @Autowired
+    private HealthEndpoint healthEndpoint;
+
+    @Autowired
+    private ApplicationContext applicationContext;
+
+    @Test
+    void testSyncCompletionDoesNotOverrideApplicationReadiness() {
+        AtomicBoolean ready = (AtomicBoolean) 
ReflectionTestUtils.getField(websocketSyncDataService, "initialSyncReady");
+        try {
+            ready.set(true);
+            assertEquals(Status.UP, 
healthEndpoint.healthForPath("readiness").getStatus());
+            AvailabilityChangeEvent.publish(applicationContext, 
ReadinessState.REFUSING_TRAFFIC);
+            assertEquals(Status.OUT_OF_SERVICE, 
healthEndpoint.healthForPath("readiness").getStatus());
+            assertEquals(Status.UP, 
healthEndpoint.healthForPath("liveness").getStatus());
+        } finally {
+            ready.set(false);
+            AvailabilityChangeEvent.publish(applicationContext, 
ReadinessState.ACCEPTING_TRAFFIC);
+        }
+    }
+
+    @Test
+    void testUnavailableAdminKeepsReadinessClosedButLivenessUp() {
+        assertEquals(Status.OUT_OF_SERVICE, 
healthEndpoint.healthForPath("readiness").getStatus());
+        assertEquals(Status.UP, 
healthEndpoint.healthForPath("liveness").getStatus());
+        assertEquals(Status.UP, healthEndpoint.health().getStatus());
+        assertFalse(new WebsocketConfig().isInitialSyncReadiness());
+    }
     
     @Test
     public void testWebsocketSyncDataService() {
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/WebsocketSyncDataService.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/WebsocketSyncDataService.java
index 481ad9854b..4bab04bf75 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/WebsocketSyncDataService.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/WebsocketSyncDataService.java
@@ -47,6 +47,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 /** Websocket sync data service. */
 public class WebsocketSyncDataService implements SyncDataService {
@@ -61,6 +62,9 @@ public class WebsocketSyncDataService implements 
SyncDataService {
     private static final String ORIGIN_HEADER_NAME = "Origin";
     
     private final WebsocketConfig websocketConfig;
+
+    private final AtomicBoolean initialSyncReady = new AtomicBoolean();
+
     
     private final PluginDataSubscriber pluginDataSubscriber;
     
@@ -209,7 +213,8 @@ public class WebsocketSyncDataService implements 
SyncDataService {
                 discoveryUpstreamDataSubscribers,
                 this.aiProxyApiKeyDataSubscribers,
                 namespaceId,
-                serverProperties.getPort());
+                serverProperties.getPort(),
+                websocketConfig.isInitialSyncReadiness() ? initialSyncReady : 
null);
     }
     
     /**
@@ -220,6 +225,14 @@ public class WebsocketSyncDataService implements 
SyncDataService {
     public WebsocketConfig getWebsocketConfig() {
         return websocketConfig;
     }
+
+    /**
+     * Whether initial configuration callbacks have completed successfully.
+     * @return startup readiness
+     */
+    public boolean isInitialSyncReady() {
+        return initialSyncReady.get();
+    }
     
     /**
      * get plugin data subscriber.
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncState.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncState.java
new file mode 100644
index 0000000000..a504214786
--- /dev/null
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncState.java
@@ -0,0 +1,153 @@
+/*
+ * 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.client;
+
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
+import org.apache.shenyu.common.utils.InitialSyncApplication;
+
+import java.util.Objects;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.Consumer;
+
+/**
+ * Tracks one connection's initial synchronization without revoking a 
successful startup.
+ */
+public final class InitialSyncState {
+
+    private final AtomicBoolean ready;
+
+    private String requestId;
+
+    private int sequence;
+
+    private boolean failed;
+
+    private long startedAt;
+
+    private int pending;
+
+    private boolean ended;
+
+    /**
+     * Create connection state.
+     * @param ready shared startup latch
+     */
+    public InitialSyncState(final AtomicBoolean ready) {
+        this.ready = ready;
+    }
+
+    /**
+     * Start an independent attempt.
+     * @return request identifier
+     */
+    public synchronized String begin() {
+        requestId = UUID.randomUUID().toString();
+        sequence = 0;
+        failed = false;
+        pending = 0;
+        ended = false;
+        startedAt = System.nanoTime();
+        return requestId;
+    }
+
+    /**
+     * Whether an incomplete attempt should be retried on a fresh connection.
+     * @return timeout or failed application before startup succeeded
+     */
+    public synchronized boolean needsReconnect() {
+        return !ready.get() && (failed || System.nanoTime() - startedAt >= 
TimeUnit.SECONDS.toNanos(60));
+    }
+
+    /**
+     * Apply a frame before acknowledging it.
+     * @param frame frame
+     * @param apply synchronous configuration application
+     */
+    public synchronized void accept(final WebsocketSyncFrame frame, final 
Consumer<String> apply) {
+        if (failed || ended || Objects.isNull(requestId) || 
!Objects.equals(requestId, frame.getRequestId())) {
+            return;
+        }
+        if (System.nanoTime() - startedAt >= TimeUnit.SECONDS.toNanos(60) || 
frame.getSequence() != sequence) {
+            failed = true;
+            return;
+        }
+        if (Objects.isNull(frame.getPayload())) {
+            ended = true;
+            completeIfApplied();
+            return;
+        }
+        try {
+            pending++;
+            sequence++;
+            String attempt = requestId;
+            InitialSyncApplication.run(() -> 
apply.accept(frame.getPayload())).whenComplete((ignored, error) -> 
applied(attempt, error));
+        } catch (RuntimeException ex) {
+            failed = true;
+            throw ex;
+        }
+    }
+
+    /**
+     * Include incremental application received before the end frame in the 
current attempt.
+     * Failures in this window invalidate the attempt because increments can 
modify the same caches.
+     * Later increments retain legacy behavior and cannot extend the initial 
completion boundary.
+     * @param action incremental callback
+     */
+    public synchronized void applyIncremental(final Runnable action) {
+        if (Objects.isNull(requestId) || failed || ended) {
+            action.run();
+            return;
+        }
+        pending++;
+        String attempt = requestId;
+        try {
+            InitialSyncApplication.run(action).whenComplete((ignored, error) 
-> applied(attempt, error));
+        } catch (RuntimeException ex) {
+            failed = true;
+            throw ex;
+        }
+    }
+
+    private synchronized void applied(final String attempt, final Throwable 
error) {
+        if (!Objects.equals(requestId, attempt)) {
+            return;
+        }
+        pending--;
+        if (Objects.nonNull(error)) {
+            failed = true;
+        }
+        completeIfApplied();
+    }
+
+    private void completeIfApplied() {
+        if (ended && pending == 0 && !failed && System.nanoTime() - startedAt 
< TimeUnit.SECONDS.toNanos(60)) {
+            ready.set(true);
+            requestId = null;
+        }
+    }
+
+    /**
+     * Invalidate the current attempt without revoking an earlier success.
+     */
+    public synchronized void invalidate() {
+        requestId = null;
+        failed = true;
+    }
+}
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
index 20732a78fe..35f9b23253 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClient.java
@@ -21,6 +21,7 @@ import org.apache.shenyu.common.constant.Constants;
 import org.apache.shenyu.common.constant.InstanceTypeConstants;
 import org.apache.shenyu.common.constant.RunningModeConstants;
 import org.apache.shenyu.common.dto.WebsocketData;
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
 import org.apache.shenyu.common.enums.ConfigGroupEnum;
 import org.apache.shenyu.common.enums.DataEventTypeEnum;
 import org.apache.shenyu.common.enums.RunningModeEnum;
@@ -89,6 +90,8 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
 
     private volatile boolean alreadySync = Boolean.FALSE;
 
+    private InitialSyncState initialSyncState;
+
     private final WebsocketDataHandler websocketDataHandler;
 
     private final Timer timer;
@@ -170,7 +173,40 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
                                  final List<AiProxyApiKeyDataSubscriber> 
aiProxyApiKeyDataSubscribers,
                                  final String namespaceId,
                                  final Integer port) {
+        this(serverUri, headers, pluginDataSubscriber, metaDataSubscribers, 
authDataSubscribers,
+                proxySelectorDataSubscribers, 
discoveryUpstreamDataSubscribers, aiProxyApiKeyDataSubscribers,
+                namespaceId, port, null);
+    }
+
+    /**
+     * Create a client with an optional startup readiness latch.
+     * @param serverUri server URI
+     * @param headers headers
+     * @param pluginDataSubscriber plugin subscriber
+     * @param metaDataSubscribers metadata subscribers
+     * @param authDataSubscribers authorization subscribers
+     * @param proxySelectorDataSubscribers proxy selector subscribers
+     * @param discoveryUpstreamDataSubscribers discovery subscribers
+     * @param aiProxyApiKeyDataSubscribers API key subscribers
+     * @param namespaceId namespace
+     * @param port gateway port
+     * @param initialSyncReady startup latch, null for the legacy protocol
+     */
+    public ShenyuWebsocketClient(final URI serverUri,
+                                 final Map<String, String> headers,
+                                 final PluginDataSubscriber 
pluginDataSubscriber,
+                                 final List<MetaDataSubscriber> 
metaDataSubscribers,
+                                 final List<AuthDataSubscriber> 
authDataSubscribers,
+                                 final List<ProxySelectorDataSubscriber> 
proxySelectorDataSubscribers,
+                                 final List<DiscoveryUpstreamDataSubscriber> 
discoveryUpstreamDataSubscribers,
+                                 final List<AiProxyApiKeyDataSubscriber> 
aiProxyApiKeyDataSubscribers,
+                                 final String namespaceId,
+                                 final Integer port,
+                                 final AtomicBoolean initialSyncReady) {
         super(serverUri, headers);
+        if (Objects.nonNull(initialSyncReady)) {
+            this.initialSyncState = new InitialSyncState(initialSyncReady);
+        }
         this.namespaceId = namespaceId;
         LOG.info("shenyu bootstrap websocket namespaceId: {}", namespaceId);
         this.addHeader(Constants.SHENYU_NAMESPACE_ID, namespaceId);
@@ -188,7 +224,12 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
     }
     
     private void connection() {
-        this.connectBlocking();
+        if (Objects.nonNull(initialSyncState)) {
+            // Management endpoints must start even when Admin is unavailable.
+            this.connect();
+        } else {
+            this.connectBlocking();
+        }
         this.timer.add(timerTask = new AbstractRoundTask(null, 
TimeUnit.SECONDS.toMillis(10)) {
             @Override
             public void doRun(final String key, final TimerTask timerTask) {
@@ -218,7 +259,11 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
         LOG.info("websocket connection server[{}] is opened, sending sync 
msg", this.getURI().toString());
         send(DataEventTypeEnum.RUNNING_MODE.name());
         if (!alreadySync) {
-            send(DataEventTypeEnum.MYSELF.name());
+            if (Objects.nonNull(initialSyncState)) {
+                send(WebsocketSyncFrame.REQUEST_PREFIX + 
initialSyncState.begin());
+            } else {
+                send(DataEventTypeEnum.MYSELF.name());
+            }
             alreadySync = true;
         }
     }
@@ -232,6 +277,10 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
         try {
             Map<String, Object> jsonToMap = JsonUtils.jsonToMap(result);
             Object eventType = jsonToMap.get(RunningModeConstants.EVENT_TYPE);
+            if (Objects.equals(WebsocketSyncFrame.EVENT_TYPE, eventType) && 
Objects.nonNull(initialSyncState)) {
+                
initialSyncState.accept(GsonUtils.getInstance().fromJson(result, 
WebsocketSyncFrame.class), this::handleResult);
+                return;
+            }
             if (Objects.equals(DataEventTypeEnum.RUNNING_MODE.name(), 
eventType)) {
                 LOG.info("server[{}] handle running mode result({})", 
this.getURI().toString(), result);
                 this.runningMode = 
String.valueOf(jsonToMap.get(RunningModeConstants.RUNNING_MODE));
@@ -240,10 +289,15 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
                 }
                 this.masterUrl = 
String.valueOf(jsonToMap.get(RunningModeConstants.MASTER_URL));
                 this.isConnectedToMaster = 
Boolean.TRUE.equals(jsonToMap.get(RunningModeConstants.IS_MASTER));
+            } else if (Objects.nonNull(initialSyncState)) {
+                initialSyncState.applyIncremental(() -> handleResult(result));
             } else {
                 handleResult(result);
             }
         } catch (RuntimeException ex) {
+            if (Objects.nonNull(initialSyncState)) {
+                initialSyncState.invalidate();
+            }
             LOG.warn("Failed to handle websocket message from server[{}], the 
message will be ignored", this.getURI(), ex);
         }
     }
@@ -260,6 +314,9 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
     
     @Override
     public void close() {
+        if (Objects.nonNull(initialSyncState)) {
+            initialSyncState.invalidate();
+        }
         alreadySync = false;
         if (this.isOpen()) {
             super.close();
@@ -287,6 +344,10 @@ public final class ShenyuWebsocketClient extends 
WebSocketClient {
             if (this.manuallyClosed.get()) {
                 return;
             }
+            if (Objects.nonNull(initialSyncState) && this.isOpen() && 
initialSyncState.needsReconnect()) {
+                close();
+                return;
+            }
             if (!this.isOpen()) {
                 if (this.reconnecting.compareAndSet(false, true)) {
                     RECONNECT_EXECUTOR.submit(this::doReconnect);
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfig.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfig.java
index 602dc5f23b..e7ce4023d0 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfig.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfig.java
@@ -38,6 +38,24 @@ public class WebsocketConfig {
      */
     private String token;
 
+    private boolean initialSyncReadiness;
+
+    /**
+     * Whether the opt-in initial synchronization protocol is required.
+     * @return enabled
+     */
+    public boolean isInitialSyncReadiness() {
+        return initialSyncReadiness;
+    }
+
+    /**
+     * Enable initial synchronization readiness.
+     * @param initialSyncReadiness enabled
+     */
+    public void setInitialSyncReadiness(final boolean initialSyncReadiness) {
+        this.initialSyncReadiness = initialSyncReadiness;
+    }
+
     /**
      * get urls.
      *
@@ -99,12 +117,13 @@ public class WebsocketConfig {
         WebsocketConfig that = (WebsocketConfig) o;
         return Objects.equals(urls, that.urls)
                 && Objects.equals(allowOrigin, that.allowOrigin)
-                && Objects.equals(token, that.token);
+                && Objects.equals(token, that.token)
+                && initialSyncReadiness == that.initialSyncReadiness;
     }
 
     @Override
     public int hashCode() {
-        return Objects.hash(urls, allowOrigin, token);
+        return Objects.hash(urls, allowOrigin, token, initialSyncReadiness);
     }
 
     @Override
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
index 6ec2b3fc21..f3c5fe91bb 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandler.java
@@ -32,7 +32,7 @@ import 
org.apache.shenyu.sync.data.api.AiProxyApiKeyDataSubscriber;
  */
 public class WebsocketDataHandler {
 
-    private static final EnumMap<ConfigGroupEnum, DataHandler> ENUM_MAP = new 
EnumMap<>(ConfigGroupEnum.class);
+    private final EnumMap<ConfigGroupEnum, DataHandler> handlers = new 
EnumMap<>(ConfigGroupEnum.class);
 
     /**
      * Instantiates a new Websocket data handler.
@@ -47,14 +47,14 @@ public class WebsocketDataHandler {
                                 final List<ProxySelectorDataSubscriber> 
proxySelectorDataSubscribers,
                                 final List<DiscoveryUpstreamDataSubscriber> 
discoveryUpstreamDataSubscribers,
                                 final List<AiProxyApiKeyDataSubscriber> 
aiProxyApiKeyDataSubscribers) {
-        ENUM_MAP.put(ConfigGroupEnum.PLUGIN, new 
PluginDataHandler(pluginDataSubscriber));
-        ENUM_MAP.put(ConfigGroupEnum.SELECTOR, new 
SelectorDataHandler(pluginDataSubscriber));
-        ENUM_MAP.put(ConfigGroupEnum.RULE, new 
RuleDataHandler(pluginDataSubscriber));
-        ENUM_MAP.put(ConfigGroupEnum.APP_AUTH, new 
AuthDataHandler(authDataSubscribers));
-        ENUM_MAP.put(ConfigGroupEnum.META_DATA, new 
MetaDataHandler(metaDataSubscribers));
-        ENUM_MAP.put(ConfigGroupEnum.PROXY_SELECTOR, new 
ProxySelectorDataHandler(proxySelectorDataSubscribers));
-        ENUM_MAP.put(ConfigGroupEnum.DISCOVER_UPSTREAM, new 
DiscoveryUpstreamDataHandler(discoveryUpstreamDataSubscribers));
-        ENUM_MAP.put(ConfigGroupEnum.AI_PROXY_API_KEY, new 
AiProxyApiKeyDataHandler(aiProxyApiKeyDataSubscribers));
+        handlers.put(ConfigGroupEnum.PLUGIN, new 
PluginDataHandler(pluginDataSubscriber));
+        handlers.put(ConfigGroupEnum.SELECTOR, new 
SelectorDataHandler(pluginDataSubscriber));
+        handlers.put(ConfigGroupEnum.RULE, new 
RuleDataHandler(pluginDataSubscriber));
+        handlers.put(ConfigGroupEnum.APP_AUTH, new 
AuthDataHandler(authDataSubscribers));
+        handlers.put(ConfigGroupEnum.META_DATA, new 
MetaDataHandler(metaDataSubscribers));
+        handlers.put(ConfigGroupEnum.PROXY_SELECTOR, new 
ProxySelectorDataHandler(proxySelectorDataSubscribers));
+        handlers.put(ConfigGroupEnum.DISCOVER_UPSTREAM, new 
DiscoveryUpstreamDataHandler(discoveryUpstreamDataSubscribers));
+        handlers.put(ConfigGroupEnum.AI_PROXY_API_KEY, new 
AiProxyApiKeyDataHandler(aiProxyApiKeyDataSubscribers));
     }
 
     /**
@@ -65,7 +65,7 @@ public class WebsocketDataHandler {
      * @param eventType the event type
      */
     public void executor(final ConfigGroupEnum type, final String json, final 
String eventType) {
-        ENUM_MAP.get(type).handle(json, eventType);
+        handlers.get(type).handle(json, eventType);
     }
 
 }
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncConnectionTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncConnectionTest.java
new file mode 100644
index 0000000000..be1a80cb59
--- /dev/null
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncConnectionTest.java
@@ -0,0 +1,143 @@
+/*
+ * 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.client;
+
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
+import org.apache.shenyu.common.utils.GsonUtils;
+import org.apache.shenyu.common.utils.InitialSyncApplication;
+import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
+import org.java_websocket.WebSocket;
+import org.java_websocket.handshake.ClientHandshake;
+import org.java_websocket.server.WebSocketServer;
+import org.junit.jupiter.api.Test;
+
+import java.lang.reflect.Field;
+import java.net.InetSocketAddress;
+import java.net.URI;
+import java.time.Duration;
+import java.util.Collections;
+import java.util.Objects;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+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;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+
+/**
+ * Exercises initial synchronization over a real socket, with deferred 
subscriber application.
+ */
+class InitialSyncConnectionTest {
+
+    @Test
+    void testDeferredApplicationAndDisconnectOverSocket() throws Exception {
+        TestServer server = new TestServer();
+        server.start();
+        assertTrue(server.started.await(5, TimeUnit.SECONDS));
+        AtomicBoolean ready = new AtomicBoolean();
+        CompletableFuture<Void> application = new CompletableFuture<>();
+        PluginDataSubscriber subscriber = mock(PluginDataSubscriber.class);
+        doAnswer(invocation -> {
+            InitialSyncApplication.register(application);
+            return null;
+        }).when(subscriber).onPluginRefresh(any());
+        ShenyuWebsocketClient client = null;
+        try {
+            client = new ShenyuWebsocketClient(URI.create("ws://127.0.0.1:" + 
server.getPort()), Collections.emptyMap(), subscriber,
+                    Collections.emptyList(), Collections.emptyList(), 
Collections.emptyList(), Collections.emptyList(), Collections.emptyList(),
+                    "default", 9195, ready);
+            String request = server.requests.poll(5, TimeUnit.SECONDS);
+            assertNotNull(request);
+            assertTrue(request.startsWith(WebsocketSyncFrame.REQUEST_PREFIX));
+            assertFalse(ready.get());
+            String id = 
request.substring(WebsocketSyncFrame.REQUEST_PREFIX.length());
+            String payload = 
"{\"groupType\":\"PLUGIN\",\"eventType\":\"MYSELF\",\"data\":[{\"id\":\"divide\",\"name\":\"divide\",\"enabled\":true}]}";
+            server.connection.send(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(id, 0, payload)));
+            server.connection.send(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(id, 1, null)));
+            Field stateField = 
ShenyuWebsocketClient.class.getDeclaredField("initialSyncState");
+            stateField.setAccessible(true);
+            InitialSyncState state = (InitialSyncState) stateField.get(client);
+            Field ended = InitialSyncState.class.getDeclaredField("ended");
+            ended.setAccessible(true);
+            await().atMost(Duration.ofSeconds(5)).until(() -> {
+                synchronized (state) {
+                    return ended.getBoolean(state);
+                }
+            });
+            assertFalse(ready.get());
+            application.complete(null);
+            await().atMost(Duration.ofSeconds(5)).untilTrue(ready);
+            server.connection.close();
+            final ShenyuWebsocketClient connectedClient = client;
+            await().atMost(Duration.ofSeconds(5)).until(() -> 
!connectedClient.isOpen());
+            assertTrue(ready.get());
+        } finally {
+            if (Objects.nonNull(client)) {
+                client.nowClose();
+            }
+            server.stop(1000);
+        }
+    }
+
+    private static final class TestServer extends WebSocketServer {
+
+        private final CountDownLatch started = new CountDownLatch(1);
+
+        private final BlockingQueue<String> requests = new 
LinkedBlockingQueue<>();
+
+        private volatile WebSocket connection;
+
+        private TestServer() {
+            super(new InetSocketAddress("127.0.0.1", 0));
+        }
+
+        @Override
+        public void onOpen(final WebSocket socket, final ClientHandshake 
handshake) {
+            connection = socket;
+        }
+
+        @Override
+        public void onClose(final WebSocket socket, final int code, final 
String reason, final boolean remote) {
+        }
+
+        @Override
+        public void onMessage(final WebSocket socket, final String message) {
+            if (message.startsWith(WebsocketSyncFrame.REQUEST_PREFIX)) {
+                requests.add(message);
+            }
+        }
+
+        @Override
+        public void onError(final WebSocket socket, final Exception exception) 
{
+            throw new IllegalStateException(exception);
+        }
+
+        @Override
+        public void onStart() {
+            started.countDown();
+        }
+    }
+}
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncStateTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncStateTest.java
new file mode 100644
index 0000000000..7d98aab7bc
--- /dev/null
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/InitialSyncStateTest.java
@@ -0,0 +1,256 @@
+/*
+ * 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.client;
+
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
+import org.apache.shenyu.common.utils.InitialSyncApplication;
+import org.junit.jupiter.api.Test;
+
+import java.lang.reflect.Field;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests initial synchronization readiness independently of the socket 
handshake.
+ */
+class InitialSyncStateTest {
+
+    @Test
+    void testInterleavedIncrementalMustFinishBeforeReadiness() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CompletableFuture<Void> incremental = new CompletableFuture<>();
+        state.applyIncremental(() -> 
InitialSyncApplication.register(incremental));
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertFalse(ready.get());
+        incremental.complete(null);
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testIncrementalFailureBeforeEndRequiresFreshAttempt() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CompletableFuture<Void> incremental = new CompletableFuture<>();
+        state.applyIncremental(() -> 
InitialSyncApplication.register(incremental));
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        incremental.completeExceptionally(new 
IllegalStateException("incremental application failed"));
+        assertFalse(ready.get());
+        assertTrue(state.needsReconnect());
+        String retry = state.begin();
+        state.accept(new WebsocketSyncFrame(retry, 0, null), value -> { });
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testSynchronousIncrementalFailureBeforeEndBlocksReadiness() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        assertThrows(IllegalStateException.class, () -> 
state.applyIncremental(() -> {
+            throw new IllegalStateException("incremental application failed");
+        }));
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertFalse(ready.get());
+        assertTrue(state.needsReconnect());
+    }
+
+    @Test
+    void testIncrementsAfterEndCannotExtendPendingApplication() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CompletableFuture<Void> initial = new CompletableFuture<>();
+        state.accept(new WebsocketSyncFrame(id, 0, "data"), value -> 
InitialSyncApplication.register(initial));
+        state.accept(new WebsocketSyncFrame(id, 1, null), value -> { });
+        CompletableFuture<Void> later = new CompletableFuture<>();
+        for (int i = 0; i < 100; i++) {
+            state.applyIncremental(() -> {
+                assertFalse(InitialSyncApplication.isActive());
+                InitialSyncApplication.register(later);
+            });
+        }
+        assertFalse(ready.get());
+        initial.complete(null);
+        assertTrue(ready.get());
+        later.completeExceptionally(new IllegalStateException("later update 
failed"));
+        assertTrue(ready.get());
+        assertFalse(state.needsReconnect());
+    }
+
+    @Test
+    void testMultipleAdminsCannotCombinePartialAttempts() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState first = new InitialSyncState(ready);
+        InitialSyncState second = new InitialSyncState(ready);
+        String firstId = first.begin();
+        String secondId = second.begin();
+        first.accept(new WebsocketSyncFrame(firstId, 0, "data"), value -> { });
+        second.accept(new WebsocketSyncFrame(firstId, 1, null), value -> { });
+        second.accept(new WebsocketSyncFrame(secondId, 1, null), value -> { });
+        assertFalse(ready.get());
+        first.accept(new WebsocketSyncFrame(firstId, 1, null), value -> { });
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testEndFrameWaitsForDeferredApplication() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CompletableFuture<Void> application = new CompletableFuture<>();
+        state.accept(new WebsocketSyncFrame(id, 0, "data"), value -> 
InitialSyncApplication.register(application));
+        state.accept(new WebsocketSyncFrame(id, 1, null), value -> { });
+        assertFalse(ready.get());
+        application.complete(null);
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testDeferredFailureAndAbandonedCompletion() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CompletableFuture<Void> application = new CompletableFuture<>();
+        state.accept(new WebsocketSyncFrame(id, 0, "data"), value -> 
InitialSyncApplication.register(application));
+        state.accept(new WebsocketSyncFrame(id, 1, null), value -> { });
+        application.completeExceptionally(new 
IllegalStateException("application failed"));
+        assertFalse(ready.get());
+        assertTrue(state.needsReconnect());
+        String retry = state.begin();
+        CompletableFuture<Void> abandoned = new CompletableFuture<>();
+        state.accept(new WebsocketSyncFrame(retry, 0, "data"), value -> 
InitialSyncApplication.register(abandoned));
+        state.accept(new WebsocketSyncFrame(retry, 1, null), value -> { });
+        state.invalidate();
+        String current = state.begin();
+        abandoned.complete(null);
+        assertFalse(ready.get());
+        state.accept(new WebsocketSyncFrame(current, 0, null), value -> { });
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testUnsupportedPeerTimesOutWithoutOpeningReadiness() throws 
ReflectiveOperationException {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        state.begin();
+        assertFalse(state.needsReconnect());
+        Field startedAt = InitialSyncState.class.getDeclaredField("startedAt");
+        startedAt.setAccessible(true);
+        startedAt.setLong(state, System.nanoTime() - 
TimeUnit.SECONDS.toNanos(61));
+        assertTrue(state.needsReconnect());
+        assertFalse(ready.get());
+    }
+
+    @Test
+    void testEmptySnapshotAndLaterDisconnect() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        assertFalse(ready.get());
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertTrue(ready.get());
+        state.invalidate();
+        assertTrue(ready.get());
+        state.begin();
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testMissingFrameDoesNotOpenReadiness() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        state.accept(new WebsocketSyncFrame(id, 1, null), value -> { });
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertFalse(ready.get());
+    }
+
+    @Test
+    void testFailedApplicationAndFreshAttempt() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        assertThrows(IllegalStateException.class, () -> state.accept(new 
WebsocketSyncFrame(id, 0, "data"), value -> {
+            throw new IllegalStateException("subscriber failed");
+        }));
+        state.accept(new WebsocketSyncFrame(id, 1, null), value -> { });
+        assertFalse(ready.get());
+        String retry = state.begin();
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertFalse(ready.get());
+        state.accept(new WebsocketSyncFrame(retry, 0, "data"), value -> { });
+        assertFalse(ready.get());
+        state.accept(new WebsocketSyncFrame(retry, 1, null), value -> { });
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testDisconnectInvalidatesIncompleteAttempt() {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        state.invalidate();
+        state.accept(new WebsocketSyncFrame(id, 0, null), value -> { });
+        assertFalse(ready.get());
+    }
+
+    @Test
+    void testCompletionWaitsForApplication() throws Exception {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        String id = state.begin();
+        CountDownLatch applying = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        ExecutorService executor = Executors.newFixedThreadPool(2);
+        try {
+            final Future<?> data = executor.submit(() -> state.accept(new 
WebsocketSyncFrame(id, 0, "data"), value -> {
+                applying.countDown();
+                try {
+                    if (!release.await(5, TimeUnit.SECONDS)) {
+                        throw new IllegalStateException("test application 
timed out");
+                    }
+                } catch (InterruptedException ex) {
+                    Thread.currentThread().interrupt();
+                    throw new IllegalStateException(ex);
+                }
+            }));
+            assertTrue(applying.await(5, TimeUnit.SECONDS));
+            final Future<?> completion = executor.submit(() -> 
state.accept(new WebsocketSyncFrame(id, 1, null), value -> { }));
+            assertFalse(ready.get());
+            release.countDown();
+            data.get(5, TimeUnit.SECONDS);
+            completion.get(5, TimeUnit.SECONDS);
+            assertTrue(ready.get());
+        } finally {
+            release.countDown();
+            executor.shutdownNow();
+        }
+    }
+}
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
index ad909e7a4a..e2e345687f 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.plugin.sync.data.websocket.client;
 
 import org.apache.shenyu.common.dto.PluginData;
 import org.apache.shenyu.common.dto.WebsocketData;
+import org.apache.shenyu.common.dto.WebsocketSyncFrame;
 import org.apache.shenyu.common.enums.ConfigGroupEnum;
 import org.apache.shenyu.common.enums.DataEventTypeEnum;
 import org.apache.shenyu.common.exception.ShenyuException;
@@ -120,6 +121,38 @@ public class ShenyuWebsocketClientTest {
         verify(pluginDataSubscriber).onPluginRefresh(any());
     }
 
+    @Test
+    void testInitialSyncRequiresCompletionAfterCallback() throws 
ReflectiveOperationException {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        Field field = 
ShenyuWebsocketClient.class.getDeclaredField("initialSyncState");
+        field.setAccessible(true);
+        field.set(shenyuWebsocketClient, state);
+        String id = state.begin();
+        shenyuWebsocketClient.onMessage(GsonUtils.getInstance().toJson(
+                new WebsocketSyncFrame(id, 0, 
GsonUtils.getInstance().toJson(websocketData))));
+        verify(pluginDataSubscriber).onPluginRefresh(any());
+        assertFalse(ready.get());
+        shenyuWebsocketClient.onMessage(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(id, 1, null)));
+        assertTrue(ready.get());
+    }
+
+    @Test
+    void testInitialSyncCallbackFailureRejectsCompletion() throws 
ReflectiveOperationException {
+        AtomicBoolean ready = new AtomicBoolean();
+        InitialSyncState state = new InitialSyncState(ready);
+        Field field = 
ShenyuWebsocketClient.class.getDeclaredField("initialSyncState");
+        field.setAccessible(true);
+        field.set(shenyuWebsocketClient, state);
+        String id = state.begin();
+        doThrow(new IllegalStateException("apply 
failed")).when(pluginDataSubscriber).onPluginRefresh(any());
+        shenyuWebsocketClient.onMessage(GsonUtils.getInstance().toJson(
+                new WebsocketSyncFrame(id, 0, 
GsonUtils.getInstance().toJson(websocketData))));
+        shenyuWebsocketClient.onMessage(GsonUtils.getInstance().toJson(new 
WebsocketSyncFrame(id, 1, null)));
+        assertFalse(ready.get());
+        assertTrue(state.needsReconnect());
+    }
+
     @Test
     public void testOnMessageShouldIgnoreMalformedJsonAndHandleNextMessage() {
         Assertions.assertDoesNotThrow(() -> 
shenyuWebsocketClient.onMessage("{invalid json"));
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfigTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfigTest.java
index c5f9bf2782..ee14c4d552 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfigTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/config/WebsocketConfigTest.java
@@ -64,13 +64,16 @@ public class WebsocketConfigTest {
         that.setToken(TOKEN);
         assertEquals(websocketConfig, websocketConfig);
         assertEquals(websocketConfig, that);
+        that.setInitialSyncReadiness(true);
+        assertNotEquals(websocketConfig, that);
         assertNotEquals(websocketConfig, null);
         assertNotEquals(websocketConfig, new Object());
     }
 
     @Test
     public void testHashCode() {
-        assertEquals(Objects.hash(websocketConfig.getUrls(), 
websocketConfig.getAllowOrigin(), websocketConfig.getToken()), 
websocketConfig.hashCode());
+        assertEquals(Objects.hash(websocketConfig.getUrls(), 
websocketConfig.getAllowOrigin(),
+                websocketConfig.getToken(), 
websocketConfig.isInitialSyncReadiness()), websocketConfig.hashCode());
     }
 
     @Test

Reply via email to