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);