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 15bcc18757 fix(logging): isolate selector queue initialization (#7228)
15bcc18757 is described below

commit 15bcc187571471c2a4d3ab5fe9931af08e4cd3e9
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 12:01:35 2026 +0800

    fix(logging): isolate selector queue initialization (#7228)
---
 .../common/collector/AbstractLogCollector.java     |  5 ++-
 .../common/collector/AbstractLogCollectorTest.java | 42 ++++++++++++++++++++++
 2 files changed, 44 insertions(+), 3 deletions(-)

diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
index 32116b0b60..03d2dae9f2 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/main/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollector.java
@@ -138,10 +138,9 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
     }
 
     private BlockingQueue<L> initQueue(final String selectorId) {
-        bufferSize = getLogCollectConfig().getBufferQueueSize();
-        bufferQueue = new LinkedBlockingDeque<>(bufferSize);
+        BlockingQueue<L> queue = new 
LinkedBlockingDeque<>(getLogCollectConfig().getBufferQueueSize());
         lastPushTimeS.put(selectorId, System.currentTimeMillis());
-        return bufferQueue;
+        return queue;
     }
 
     private void processBufferQueue(final BlockingQueue<L> bufferQueue, final 
int batchSize,
diff --git 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
index 5883f8acaa..7b5973cfdd 100644
--- 
a/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-logging/shenyu-plugin-logging-common/src/test/java/org/apache/shenyu/plugin/logging/common/collector/AbstractLogCollectorTest.java
@@ -24,6 +24,7 @@ import 
org.apache.shenyu.plugin.logging.common.entity.ShenyuRequestLog;
 import 
org.apache.shenyu.plugin.logging.desensitize.api.enums.DataDesensitizeEnum;
 import org.apache.shenyu.plugin.logging.desensitize.api.matcher.KeyWordMatch;
 import org.junit.jupiter.api.Test;
+import org.springframework.test.util.ReflectionTestUtils;
 
 import java.lang.reflect.Field;
 import java.util.Collections;
@@ -33,6 +34,10 @@ import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -68,6 +73,43 @@ public class AbstractLogCollectorTest {
                 }
             };
 
+    @Test
+    public void testSelectorQueueInitializationDoesNotReplaceGlobalQueue() 
throws Exception {
+        final GenericGlobalConfig config = new GenericGlobalConfig();
+        config.setBufferQueueSize(2);
+        final AbstractLogCollector<?, ShenyuRequestLog, GenericGlobalConfig> 
testCollector = new AbstractLogCollector<>() {
+            @Override
+            protected AbstractLogConsumeClient<?, ShenyuRequestLog> 
getLogConsumeClient() {
+                return logConsumeClient;
+            }
+
+            @Override
+            protected GenericGlobalConfig getLogCollectConfig() {
+                return config;
+            }
+
+            @Override
+            protected void desensitizeLog(final ShenyuRequestLog log, final 
KeyWordMatch keyWordMatch, final String desensitizeAlg) {
+            }
+        };
+        final BlockingQueue<ShenyuRequestLog> global = new 
LinkedBlockingDeque<>(3);
+        setField(testCollector, "bufferQueue", global);
+        setField(testCollector, "bufferSize", 3);
+        final ExecutorService executor = Executors.newFixedThreadPool(2);
+        try {
+            Future<BlockingQueue<ShenyuRequestLog>> first = executor.submit(() 
-> ReflectionTestUtils.invokeMethod(testCollector, "initQueue", "first"));
+            Future<BlockingQueue<ShenyuRequestLog>> second = 
executor.submit(() -> ReflectionTestUtils.invokeMethod(testCollector, 
"initQueue", "second"));
+            BlockingQueue<ShenyuRequestLog> firstQueue = first.get(2, 
TimeUnit.SECONDS);
+            BlockingQueue<ShenyuRequestLog> secondQueue = second.get(2, 
TimeUnit.SECONDS);
+            firstQueue.add(new ShenyuRequestLog());
+            assertEquals(2, secondQueue.remainingCapacity());
+            assertSame(global, ReflectionTestUtils.getField(testCollector, 
"bufferQueue"));
+            assertEquals(3, ReflectionTestUtils.getField(testCollector, 
"bufferSize"));
+        } finally {
+            executor.shutdownNow();
+        }
+    }
+
     @Test
     public void testCollectAddsLogWhenBufferQueueHasCapacity() throws 
Exception {
         BlockingQueue<ShenyuRequestLog> bufferQueue = new 
LinkedBlockingDeque<>(1);

Reply via email to