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;
+        }
+    }
 }

Reply via email to