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