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 ce2a2e37b4 fix: offload websocket reconnect from shared timer to
dedicated executor (#6898)
ce2a2e37b4 is described below
commit ce2a2e37b4790cf1f0f1bfe49c66b1523935d8b1
Author: wy471x <[email protected]>
AuthorDate: Mon Aug 17 09:18:50 2026 +0800
fix: offload websocket reconnect from shared timer to dedicated executor
(#6898)
* fix: offload websocket reconnect from shared timer to dedicated executor
healthCheck() called reconnectBlocking() on the shared wheel timer's
single-thread executor, blocking all other timer tasks (other clients'
health checks, masterCheck) during TCP connect timeout. Offload
reconnect to a dedicated CachedThreadPool executor with exponential
backoff + jitter (1s-60s) and an AtomicBoolean guard to prevent
concurrent reconnect attempts.
Co-Authored-By: Claude Opus 4.7 <[email protected]>
* fix: bound reconnect executor and enforce backoff in websocket client
Address review feedback on reconnect changes:
- replace unbounded cached pool with bounded ShenyuThreadPoolExecutor
(core 1, max 8, 60s keep-alive) using repo-standard task queue
- stamp lastReconnectAttemptTime after the attempt completes so a slow
connect no longer consumes the backoff window
- stub reconnectBlocking() in tests instead of attempting real socket
connections
Co-Authored-By: Claude Opus 4.7 <[email protected]>
---------
Co-authored-by: Claude Opus 4.7 <[email protected]>
Co-authored-by: aias00 <[email protected]>
---
.../websocket/client/ShenyuWebsocketClient.java | 90 +++++++-
.../client/ShenyuWebsocketClientTest.java | 239 +++++++++++++++++++--
2 files changed, 307 insertions(+), 22 deletions(-)
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 4913433cd5..ddf5878055 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
@@ -43,11 +43,20 @@ import org.java_websocket.handshake.ServerHandshake;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import org.apache.shenyu.common.concurrent.MemorySafeTaskQueue;
+import org.apache.shenyu.common.concurrent.ShenyuThreadFactory;
+import org.apache.shenyu.common.concurrent.ShenyuThreadPoolExecutor;
+
import java.net.URI;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.ThreadLocalRandom;
+import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
/**
* The type shenyu websocket client.
@@ -58,23 +67,48 @@ public final class ShenyuWebsocketClient extends
WebSocketClient {
* logger.
*/
private static final Logger LOG =
LoggerFactory.getLogger(ShenyuWebsocketClient.class);
-
+
+ private static final int RECONNECT_EXECUTOR_CORE_POOL_SIZE = 1;
+
+ private static final int RECONNECT_EXECUTOR_MAX_POOL_SIZE = 8;
+
+ private static final long RECONNECT_EXECUTOR_KEEP_ALIVE_MS =
TimeUnit.SECONDS.toMillis(60);
+
+ private static final ExecutorService RECONNECT_EXECUTOR = new
ShenyuThreadPoolExecutor(
+ RECONNECT_EXECUTOR_CORE_POOL_SIZE,
+ RECONNECT_EXECUTOR_MAX_POOL_SIZE,
+ RECONNECT_EXECUTOR_KEEP_ALIVE_MS,
+ TimeUnit.MILLISECONDS,
+ new MemorySafeTaskQueue<>(Constants.THE_256_MB),
+ ShenyuThreadFactory.create("websocket-reconnect", true),
+ new ThreadPoolExecutor.AbortPolicy());
+
+ private static final long MIN_RECONNECT_BACKOFF_MS =
TimeUnit.SECONDS.toMillis(1);
+
+ private static final long MAX_RECONNECT_BACKOFF_MS =
TimeUnit.SECONDS.toMillis(60);
+
private volatile boolean alreadySync = Boolean.FALSE;
-
+
private final WebsocketDataHandler websocketDataHandler;
-
+
private final Timer timer;
-
+
private TimerTask timerTask;
-
+
private String runningMode;
-
+
private String masterUrl;
-
+
private volatile boolean isConnectedToMaster;
-
+
private final String namespaceId;
+ private final AtomicBoolean reconnecting = new AtomicBoolean(false);
+
+ private volatile long lastReconnectAttemptTime;
+
+ private final AtomicInteger reconnectBackoff = new AtomicInteger(0);
+
/**
* Instantiates a new shenyu websocket client.
*
@@ -238,11 +272,13 @@ public final class ShenyuWebsocketClient extends
WebSocketClient {
private void healthCheck() {
try {
if (!this.isOpen()) {
- this.reconnectBlocking();
+ if (this.reconnecting.compareAndSet(false, true)) {
+ RECONNECT_EXECUTOR.submit(this::doReconnect);
+ }
} else {
+ this.reconnectBackoff.set(0);
this.sendPing();
send(getInstanceInfo());
-// send(DataEventTypeEnum.RUNNING_MODE.name());
LOG.debug("websocket send to [{}] ping message successful",
this.getURI());
}
} catch (Exception e) {
@@ -250,6 +286,40 @@ public final class ShenyuWebsocketClient extends
WebSocketClient {
}
}
+ private void doReconnect() {
+ try {
+ long backoff = calculateBackoff();
+ long since = System.currentTimeMillis() - lastReconnectAttemptTime;
+ long waitMs = backoff - since;
+ if (waitMs > 0) {
+ Thread.sleep(waitMs);
+ }
+ try {
+ this.reconnectBlocking();
+ } finally {
+ lastReconnectAttemptTime = System.currentTimeMillis();
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ } catch (Exception e) {
+ reconnectBackoff.set(Math.min(reconnectBackoff.get() + 1, 10));
+ LOG.error("websocket reconnect server[{}] error", this.getURI(),
e);
+ } finally {
+ this.reconnecting.set(false);
+ }
+ }
+
+ private long calculateBackoff() {
+ int failures = reconnectBackoff.get();
+ if (failures <= 0) {
+ return 0;
+ }
+ long base = Math.min(
+ MIN_RECONNECT_BACKOFF_MS * (1L << Math.min(failures - 1, 10)),
+ MAX_RECONNECT_BACKOFF_MS);
+ return base + (long) (base * 0.5 *
ThreadLocalRandom.current().nextDouble());
+ }
+
private String getInstanceInfo() {
// Combine instance and host information
Map<String, Object> combinedInfo = Map.of(
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 0936a435dd..dcb42b7db6 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
@@ -35,43 +35,56 @@ import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
+import java.lang.reflect.Field;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
import java.net.URI;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import org.mockito.Answers;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.withSettings;
/**
* add test case for {@link ShenyuWebsocketClient}.
*/
@ExtendWith(MockitoExtension.class)
public class ShenyuWebsocketClientTest {
-
+
@InjectMocks
private ShenyuWebsocketClient shenyuWebsocketClient;
-
+
@Mock
private URI serverUri;
-
+
@Mock
private PluginDataSubscriber pluginDataSubscriber;
-
+
@Mock
private List<MetaDataSubscriber> metaDataSubscribers;
-
+
@Mock
private List<AuthDataSubscriber> authDataSubscribers;
-
+
@Mock
private ScheduledThreadPoolExecutor executor;
-
+
private WebsocketData<PluginData> websocketData;
-
+
@BeforeEach
public void setUp() {
websocketData = new WebsocketData<>();
@@ -83,7 +96,7 @@ public class ShenyuWebsocketClientTest {
list.add(pluginData);
websocketData.setData(list);
}
-
+
@Test
public void testOnOpen() {
shenyuWebsocketClient = spy(shenyuWebsocketClient);
@@ -94,7 +107,7 @@ public class ShenyuWebsocketClientTest {
verify(shenyuWebsocketClient).send(DataEventTypeEnum.RUNNING_MODE.name());
verify(shenyuWebsocketClient).send(DataEventTypeEnum.MYSELF.name());
}
-
+
@Test
public void testOnMessage() {
doNothing().when(pluginDataSubscriber).onSubscribe(any());
@@ -102,7 +115,7 @@ public class ShenyuWebsocketClientTest {
shenyuWebsocketClient.onMessage(json);
verify(pluginDataSubscriber).onSubscribe(any());
}
-
+
@Test
public void testOnClose() {
shenyuWebsocketClient = spy(shenyuWebsocketClient);
@@ -110,10 +123,212 @@ public class ShenyuWebsocketClientTest {
shenyuWebsocketClient.onClose(1, "shenyu-plugin-grpc", true);
verify(shenyuWebsocketClient).close();
}
-
+
@Test
public void testOnError() {
shenyuWebsocketClient = spy(shenyuWebsocketClient);
Assertions.assertDoesNotThrow(() -> shenyuWebsocketClient.onError(new
ShenyuException("test")));
}
+
+ // ========== reconnect/backoff tests ==========
+
+ private ShenyuWebsocketClient createMockClient() {
+ ShenyuWebsocketClient client = mock(ShenyuWebsocketClient.class,
+ withSettings().defaultAnswer(Answers.CALLS_REAL_METHODS));
+ setField(client, "reconnecting", new AtomicBoolean(false));
+ setField(client, "reconnectBackoff", new AtomicInteger(0));
+ setField(client, "lastReconnectAttemptTime", 0L);
+ return client;
+ }
+
+ private void setField(final Object target, final String name, final Object
value) {
+ try {
+ Field field = ShenyuWebsocketClient.class.getDeclaredField(name);
+ field.setAccessible(true);
+ field.set(target, value);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private Object getField(final Object target, final String name) {
+ try {
+ Field field = ShenyuWebsocketClient.class.getDeclaredField(name);
+ field.setAccessible(true);
+ return field.get(target);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private Object invokePrivate(final Object target, final String methodName)
{
+ try {
+ Method method =
ShenyuWebsocketClient.class.getDeclaredMethod(methodName);
+ method.setAccessible(true);
+ return method.invoke(target);
+ } catch (InvocationTargetException e) {
+ Throwable cause = e.getCause();
+ if (cause instanceof RuntimeException) {
+ throw (RuntimeException) cause;
+ }
+ throw new RuntimeException(cause);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ // ---------- calculateBackoff tests ----------
+
+ @Test
+ void testCalculateBackoffReturnsZeroForNoFailures() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(0));
+
+ long backoff = (long) invokePrivate(client, "calculateBackoff");
+
+ assertEquals(0, backoff);
+ }
+
+ @Test
+ void testCalculateBackoffExponentialGrowth() {
+ ShenyuWebsocketClient client = createMockClient();
+
+ setField(client, "reconnectBackoff", new AtomicInteger(1));
+ long backoff1 = (long) invokePrivate(client, "calculateBackoff");
+ assertTrue(backoff1 >= 1000 && backoff1 <= 1500,
+ () -> "Expected [1000, 1500] but got " + backoff1);
+
+ setField(client, "reconnectBackoff", new AtomicInteger(2));
+ long backoff2 = (long) invokePrivate(client, "calculateBackoff");
+ assertTrue(backoff2 >= 2000 && backoff2 <= 3000,
+ () -> "Expected [2000, 3000] but got " + backoff2);
+
+ setField(client, "reconnectBackoff", new AtomicInteger(4));
+ long backoff4 = (long) invokePrivate(client, "calculateBackoff");
+ assertTrue(backoff4 >= 8000 && backoff4 <= 12000,
+ () -> "Expected [8000, 12000] but got " + backoff4);
+ }
+
+ @Test
+ void testCalculateBackoffMaxCap() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(10));
+
+ long backoff = (long) invokePrivate(client, "calculateBackoff");
+
+ assertTrue(backoff >= 60000 && backoff <= 90000,
+ () -> "Expected [60000, 90000] but got " + backoff);
+ }
+
+ @Test
+ void testCalculateBackoffIncludesJitter() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(1));
+ boolean varied = false;
+ long first = (long) invokePrivate(client, "calculateBackoff");
+ for (int i = 0; i < 20; i++) {
+ if ((long) invokePrivate(client, "calculateBackoff") != first) {
+ varied = true;
+ break;
+ }
+ }
+ assertTrue(varied, "Backoff should vary due to jitter");
+ }
+
+ // ---------- healthCheck tests ----------
+
+ @Test
+ void testHealthCheckDoesNotDoubleSubmitWhenAlreadyReconnecting() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnecting", new AtomicBoolean(true));
+ doReturn(false).when(client).isOpen();
+
+ invokePrivate(client, "healthCheck");
+
+ assertTrue(((AtomicBoolean) getField(client, "reconnecting")).get());
+ verify(client).isOpen();
+ }
+
+ @Test
+ void testHealthCheckResetsBackoffAndSendsPingWhenOpen() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(5));
+ doReturn(true).when(client).isOpen();
+ doNothing().when(client).sendPing();
+ doNothing().when(client).send(anyString());
+ doReturn(URI.create("ws://localhost:9090")).when(client).getURI();
+
+ invokePrivate(client, "healthCheck");
+
+ assertEquals(0, ((AtomicInteger) getField(client,
"reconnectBackoff")).get());
+ verify(client).sendPing();
+ }
+
+ // ---------- doReconnect tests ----------
+ // reconnectBlocking() is stubbed to throw, so no real socket connection
is attempted.
+
+ @Test
+ void testDoReconnectIncrementsBackoffOnFailure() throws
InterruptedException {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(0));
+ doThrow(new RuntimeException("test")).when(client).reconnectBlocking();
+ doReturn(URI.create("ws://localhost:9090")).when(client).getURI();
+
+ invokePrivate(client, "doReconnect");
+
+ assertEquals(1, ((AtomicInteger) getField(client,
"reconnectBackoff")).get());
+ assertFalse(((AtomicBoolean) getField(client, "reconnecting")).get());
+ }
+
+ @Test
+ void testDoReconnectBackoffCappedAtTen() throws InterruptedException {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(10));
+ doThrow(new RuntimeException("test")).when(client).reconnectBlocking();
+ doReturn(URI.create("ws://localhost:9090")).when(client).getURI();
+
+ invokePrivate(client, "doReconnect");
+
+ assertEquals(10, ((AtomicInteger) getField(client,
"reconnectBackoff")).get());
+ }
+
+ @Test
+ void testDoReconnectResetsReconnectingOnFailure() throws
InterruptedException {
+ ShenyuWebsocketClient client = createMockClient();
+ doThrow(new RuntimeException("test")).when(client).reconnectBlocking();
+ doReturn(URI.create("ws://localhost:9090")).when(client).getURI();
+
+ invokePrivate(client, "doReconnect");
+
+ assertFalse(((AtomicBoolean) getField(client, "reconnecting")).get());
+ }
+
+ @Test
+ void testDoReconnectAppliesBackoffSleep() throws InterruptedException {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(1));
+ setField(client, "lastReconnectAttemptTime",
System.currentTimeMillis());
+ doThrow(new RuntimeException("test")).when(client).reconnectBlocking();
+ doReturn(URI.create("ws://localhost:9090")).when(client).getURI();
+
+ long start = System.currentTimeMillis();
+ invokePrivate(client, "doReconnect");
+ long elapsed = System.currentTimeMillis() - start;
+
+ assertTrue(elapsed >= 800,
+ () -> "Expected >= 800ms backoff sleep, got " + elapsed +
"ms");
+ }
+
+ @Test
+ void testDoReconnectPreservesInterruptStatus() {
+ ShenyuWebsocketClient client = createMockClient();
+ setField(client, "reconnectBackoff", new AtomicInteger(1));
+ setField(client, "lastReconnectAttemptTime",
System.currentTimeMillis());
+
+ Thread.currentThread().interrupt();
+ invokePrivate(client, "doReconnect");
+
+ assertTrue(Thread.interrupted(), "Interrupt status should be preserved
after reconnect");
+ assertFalse(((AtomicBoolean) getField(client, "reconnecting")).get());
+ }
}