This is an automated email from the ASF dual-hosted git repository.
dengliming 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 66557f1d07 fix: use offer for bounded log collector queues (#7069)
66557f1d07 is described below
commit 66557f1d072a0888f9439bf782442e26150fa687
Author: Southern <[email protected]>
AuthorDate: Thu Sep 17 12:16:13 2026 +0800
fix: use offer for bounded log collector queues (#7069)
Replace the non-atomic size-check and add operation in AbstractLogCollector
with BlockingQueue.offer for both single-
client and multi-client queues. Add regression tests covering successful
enqueue and full-queue handling without
IllegalStateException.
Co-authored-by: Liming Deng <[email protected]>
---
.../common/collector/AbstractLogCollector.java | 8 +-
.../common/collector/AbstractLogCollectorTest.java | 94 +++++++++++++++++++++-
2 files changed, 95 insertions(+), 7 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 9f8024cac0..560ff13c94 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
@@ -92,13 +92,9 @@ public abstract class AbstractLogCollector<T extends
AbstractLogConsumeClient<?,
if (getMultiClient()) {
String selectorId = log.getSelectorId();
BlockingQueue<L> bufferQueue =
bufferQueueS.computeIfAbsent(selectorId, bufferQueueS -> initQueue(selectorId));
- if (bufferQueue.size() < bufferSize) {
- bufferQueue.add(log);
- }
+ bufferQueue.offer(log);
} else {
- if (bufferQueue.size() < bufferSize) {
- bufferQueue.add(log);
- }
+ bufferQueue.offer(log);
}
}
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 fb7404b131..b6db3e0893 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
@@ -25,25 +25,33 @@ import
org.apache.shenyu.plugin.logging.desensitize.api.enums.DataDesensitizeEnu
import org.apache.shenyu.plugin.logging.desensitize.api.matcher.KeyWordMatch;
import org.junit.jupiter.api.Test;
+import java.lang.reflect.Field;
import java.util.Collections;
import java.util.HashSet;
+import java.util.Map;
import java.util.Set;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingDeque;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.mock;
/**
* The Test Case For AbstractLogCollector.
*/
public class AbstractLogCollectorTest {
+ private final AbstractLogConsumeClient<?, ShenyuRequestLog>
logConsumeClient = mock(AbstractLogConsumeClient.class);
+
private final AbstractLogCollector<AbstractLogConsumeClient<?,
ShenyuRequestLog>, ShenyuRequestLog, GenericGlobalConfig> collector =
new AbstractLogCollector<>() {
@Override
protected AbstractLogConsumeClient<?, ShenyuRequestLog>
getLogConsumeClient() {
- return null;
+ return logConsumeClient;
}
@Override
@@ -56,6 +64,65 @@ public class AbstractLogCollectorTest {
}
};
+ @Test
+ public void testCollectAddsLogWhenBufferQueueHasCapacity() throws
Exception {
+ BlockingQueue<ShenyuRequestLog> bufferQueue = new
LinkedBlockingDeque<>(1);
+ setField(collector, "bufferSize", 1);
+ setField(collector, "bufferQueue", bufferQueue);
+ ShenyuRequestLog log = new ShenyuRequestLog();
+
+ collector.collect(log);
+
+ assertSame(log, bufferQueue.peek());
+ }
+
+ @Test
+ public void testCollectDoesNotThrowWhenBufferQueueIsFull() throws
Exception {
+ ShenyuRequestLog bufferedLog = new ShenyuRequestLog();
+ BlockingQueue<ShenyuRequestLog> bufferQueue = new
StaleSizeLinkedBlockingDeque();
+ bufferQueue.add(bufferedLog);
+ setField(collector, "bufferSize", 1);
+ setField(collector, "bufferQueue", bufferQueue);
+
+ assertDoesNotThrow(() -> collector.collect(new ShenyuRequestLog()));
+ assertSame(bufferedLog, bufferQueue.peek());
+ }
+
+ @Test
+ public void testCollectDoesNotThrowWhenMultiClientBufferQueueIsFull()
throws Exception {
+ AbstractLogCollector<AbstractLogConsumeClient<?, ShenyuRequestLog>,
ShenyuRequestLog, GenericGlobalConfig> multiClientCollector =
+ new AbstractLogCollector<>() {
+ @Override
+ protected AbstractLogConsumeClient<?, ShenyuRequestLog>
getLogConsumeClient() {
+ return logConsumeClient;
+ }
+
+ @Override
+ protected boolean getMultiClient() {
+ return true;
+ }
+
+ @Override
+ protected GenericGlobalConfig getLogCollectConfig() {
+ return null;
+ }
+
+ @Override
+ protected void desensitizeLog(final ShenyuRequestLog log,
final KeyWordMatch keyWordMatch, final String desensitizeAlg) {
+ }
+ };
+ ShenyuRequestLog bufferedLog = new ShenyuRequestLog();
+ BlockingQueue<ShenyuRequestLog> bufferQueue = new
StaleSizeLinkedBlockingDeque();
+ bufferQueue.add(bufferedLog);
+ setField(multiClientCollector, "bufferSize", 1);
+ getBufferQueues(multiClientCollector).put("selector", bufferQueue);
+ ShenyuRequestLog log = new ShenyuRequestLog();
+ log.setSelectorId("selector");
+
+ assertDoesNotThrow(() -> multiClientCollector.collect(log));
+ assertSame(bufferedLog, bufferQueue.peek());
+ }
+
@Test
public void testDesensitizeToleratesNullBoxedNumericFields() {
// a chunked byte-type response reaches desensitize with
responseContentLength,
@@ -84,4 +151,29 @@ public class AbstractLogCollectorTest {
assertEquals(200, log.getStatus());
assertEquals(15L, log.getUpstreamResponseTime());
}
+
+ private static void setField(final Object target, final String fieldName,
final Object value) throws Exception {
+ Field field = AbstractLogCollector.class.getDeclaredField(fieldName);
+ field.setAccessible(true);
+ field.set(target, value);
+ }
+
+ @SuppressWarnings("unchecked")
+ private static Map<String, BlockingQueue<ShenyuRequestLog>>
getBufferQueues(final Object target) throws Exception {
+ Field field =
AbstractLogCollector.class.getDeclaredField("bufferQueueS");
+ field.setAccessible(true);
+ return (Map<String, BlockingQueue<ShenyuRequestLog>>)
field.get(target);
+ }
+
+ private static final class StaleSizeLinkedBlockingDeque extends
LinkedBlockingDeque<ShenyuRequestLog> {
+
+ private StaleSizeLinkedBlockingDeque() {
+ super(1);
+ }
+
+ @Override
+ public int size() {
+ return 0;
+ }
+ }
}