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 17bf2f4cf9 fix(logging): backpressure saturated cloud callback 
executors (#7275)
17bf2f4cf9 is described below

commit 17bf2f4cf9f9730bf8b8059b08330bc16d7221b7
Author: Liming Deng <[email protected]>
AuthorDate: Sun Sep 27 11:41:13 2026 +0800

    fix(logging): backpressure saturated cloud callback executors (#7275)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../sls/client/AliyunSlsLogCollectClient.java      |  2 +-
 .../client/CloudLogCallbackBackpressureTest.java   | 69 ++++++++++++++++++++++
 .../lts/client/HuaweiLtsLogCollectClient.java      |  2 +-
 .../client/CloudLogCallbackBackpressureTest.java   | 66 +++++++++++++++++++++
 .../cls/client/TencentClsLogCollectClient.java     |  2 +-
 .../client/CloudLogCallbackBackpressureTest.java   | 66 +++++++++++++++++++++
 6 files changed, 204 insertions(+), 3 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
index 0ca5b6da85..e66418e0ba 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/main/java/org/apache/shenyu/plugin/aliyun/sls/client/AliyunSlsLogCollectClient.java
@@ -191,7 +191,7 @@ public class AliyunSlsLogCollectClient extends 
AbstractLogConsumeClient<AliyunLo
         }
         return new ThreadPoolExecutor(sendThreadCount, 
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
                 new 
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER), 
ShenyuThreadFactory.create("shenyu-aliyun-sls", true),
-                new ThreadPoolExecutor.AbortPolicy());
+                new ThreadPoolExecutor.CallerRunsPolicy());
     }
 
     /**
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..cdaec6d880
--- /dev/null
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-aliyun-sls/src/test/java/org/apache/shenyu/plugin/aliyun/sls/client/CloudLogCallbackBackpressureTest.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.plugin.aliyun.sls.client;
+
+import org.apache.shenyu.plugin.aliyun.sls.config.AliyunLogCollectConfig;
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+    @Test
+    void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+        AliyunLogCollectConfig.AliyunSlsLogConfig config = new 
AliyunLogCollectConfig.AliyunSlsLogConfig();
+        config.setSendThreadCount(1);
+        ThreadPoolExecutor executor = 
ReflectionTestUtils.invokeMethod(AliyunSlsLogCollectClient.class, 
"createThreadPoolExecutor", config);
+        assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+        executor.setMaximumPoolSize(1);
+        CountDownLatch entered = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        try {
+            executor.execute(() -> {
+                entered.countDown();
+                try {
+                    release.await();
+                } catch (InterruptedException ex) {
+                    Thread.currentThread().interrupt();
+                }
+            });
+            assertTrue(entered.await(5, TimeUnit.SECONDS));
+            int capacity = executor.getQueue().remainingCapacity();
+            for (int i = 0; i < capacity; i++) {
+                executor.execute(() -> { });
+            }
+            AtomicReference<Thread> callbackThread = new AtomicReference<>();
+            executor.execute(() -> callbackThread.set(Thread.currentThread()));
+            assertSame(Thread.currentThread(), callbackThread.get());
+            assertEquals(capacity, executor.getQueue().size());
+        } finally {
+            executor.shutdownNow();
+            release.countDown();
+            assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+        }
+    }
+}
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
index b0df92cda4..c6e2d0bfc4 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/main/java/org/apache/shenyu/plugin/huawei/lts/client/HuaweiLtsLogCollectClient.java
@@ -177,7 +177,7 @@ public class HuaweiLtsLogCollectClient extends 
AbstractLogConsumeClient<HuaweiLo
         }
         return new ThreadPoolExecutor(threadCount, 
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
                 new 
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER), 
ShenyuThreadFactory.create("shenyu-huawei-lts", true),
-                new ThreadPoolExecutor.AbortPolicy());
+                new ThreadPoolExecutor.CallerRunsPolicy());
     }
 
     /**
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..3b97cef9d8
--- /dev/null
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-huawei-lts/src/test/java/org/apache/shenyu/plugin/huawei/lts/client/CloudLogCallbackBackpressureTest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.huawei.lts.client;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+    @Test
+    void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+        ThreadPoolExecutor executor = 
ReflectionTestUtils.invokeMethod(HuaweiLtsLogCollectClient.class, 
"createThreadPoolExecutor", 1);
+        assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+        executor.setMaximumPoolSize(1);
+        CountDownLatch entered = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        try {
+            executor.execute(() -> {
+                entered.countDown();
+                try {
+                    release.await();
+                } catch (InterruptedException ex) {
+                    Thread.currentThread().interrupt();
+                }
+            });
+            assertTrue(entered.await(5, TimeUnit.SECONDS));
+            int capacity = executor.getQueue().remainingCapacity();
+            for (int i = 0; i < capacity; i++) {
+                executor.execute(() -> { });
+            }
+            AtomicReference<Thread> callbackThread = new AtomicReference<>();
+            executor.execute(() -> callbackThread.set(Thread.currentThread()));
+            assertSame(Thread.currentThread(), callbackThread.get());
+            assertEquals(capacity, executor.getQueue().size());
+        } finally {
+            executor.shutdownNow();
+            release.countDown();
+            assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+        }
+    }
+}
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
index 4b90205033..ae82a573e3 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/main/java/org/apache/shenyu/plugin/tencent/cls/client/TencentClsLogCollectClient.java
@@ -170,7 +170,7 @@ public class TencentClsLogCollectClient extends 
AbstractLogConsumeClient<Tencent
         }
         return new ThreadPoolExecutor(threadCount, 
GenericLoggingConstant.MAX_ALLOW_THREADS, 60000L, TimeUnit.MILLISECONDS,
                 new 
LinkedBlockingQueue<>(GenericLoggingConstant.MAX_QUEUE_NUMBER), 
ShenyuThreadFactory.create("shenyu-tencent-cls", true),
-                new ThreadPoolExecutor.AbortPolicy());
+                new ThreadPoolExecutor.CallerRunsPolicy());
     }
 
     /**
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
new file mode 100644
index 0000000000..50eeec6870
--- /dev/null
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-tencent-cls/src/test/java/org/apache/shenyu/plugin/tencent/cls/client/CloudLogCallbackBackpressureTest.java
@@ -0,0 +1,66 @@
+/*
+ * 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.tencent.cls.client;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CloudLogCallbackBackpressureTest {
+
+    @Test
+    void saturatedPoolRunsCallbackOnSubmittingThread() throws Exception {
+        ThreadPoolExecutor executor = 
ReflectionTestUtils.invokeMethod(TencentClsLogCollectClient.class, 
"createThreadPoolExecutor", 1);
+        assertEquals(60, executor.getKeepAliveTime(TimeUnit.SECONDS));
+        executor.setMaximumPoolSize(1);
+        CountDownLatch entered = new CountDownLatch(1);
+        CountDownLatch release = new CountDownLatch(1);
+        try {
+            executor.execute(() -> {
+                entered.countDown();
+                try {
+                    release.await();
+                } catch (InterruptedException ex) {
+                    Thread.currentThread().interrupt();
+                }
+            });
+            assertTrue(entered.await(5, TimeUnit.SECONDS));
+            int capacity = executor.getQueue().remainingCapacity();
+            for (int i = 0; i < capacity; i++) {
+                executor.execute(() -> { });
+            }
+            AtomicReference<Thread> callbackThread = new AtomicReference<>();
+            executor.execute(() -> callbackThread.set(Thread.currentThread()));
+            assertSame(Thread.currentThread(), callbackThread.get());
+            assertEquals(capacity, executor.getQueue().size());
+        } finally {
+            executor.shutdownNow();
+            release.countDown();
+            assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+        }
+    }
+}

Reply via email to