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 891bc5deb9 fix(logging): make collector startup idempotent (#7227)
891bc5deb9 is described below

commit 891bc5deb969591224e2a6416f30ecbc39518116
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 12:18:14 2026 +0800

    fix(logging): make collector startup idempotent (#7227)
---
 .../common/collector/AbstractLogCollector.java     | 52 ++++++++++---
 .../common/collector/AbstractLogCollectorTest.java | 88 +++++++++++++++++++++-
 2 files changed, 128 insertions(+), 12 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 535022155a..3576996a76 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
@@ -57,9 +57,9 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
 
     private static final Logger LOG = 
LoggerFactory.getLogger(AbstractLogCollector.class);
 
-    private int bufferSize;
+    private volatile int bufferSize;
 
-    private BlockingQueue<L> bufferQueue;
+    private volatile BlockingQueue<L> bufferQueue;
 
     private final Map<String, BlockingQueue<L>> bufferQueueS = 
Maps.newConcurrentMap();
 
@@ -67,12 +67,18 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
 
     private long lastPushTime;
 
-    private final AtomicBoolean started = new AtomicBoolean(true);
+    private final AtomicBoolean started = new AtomicBoolean(false);
+
+    private ShenyuThreadPoolExecutor executor;
 
     @Override
-    public void start() {
-        bufferSize = getLogCollectConfig().getBufferQueueSize();
-        bufferQueue = new LinkedBlockingDeque<>(bufferSize);
+    public synchronized void start() {
+        refreshBufferSize();
+        if (started.get()) {
+            return;
+        }
+        // Admission is bounded in collect(), allowing live capacity changes 
without dropping queued logs.
+        bufferQueue = new LinkedBlockingDeque<>();
         ShenyuConfig config = 
Optional.ofNullable(Singleton.INST.get(ShenyuConfig.class)).orElse(new 
ShenyuConfig());
         final ShenyuConfig.SharedPool sharedPool = config.getSharedPool();
         ShenyuThreadPoolExecutor threadExecutor = new 
ShenyuThreadPoolExecutor(sharedPool.getCorePoolSize(),
@@ -81,7 +87,22 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
                 ShenyuThreadFactory.create(config.getSharedPool().getPrefix(), 
true),
                 new ThreadPoolExecutor.AbortPolicy());
         started.set(true);
-        threadExecutor.execute(this::consume);
+        executor = threadExecutor;
+        try {
+            threadExecutor.execute(() -> consume(threadExecutor));
+        } catch (RuntimeException e) {
+            started.set(false);
+            threadExecutor.shutdownNow();
+            throw e;
+        }
+    }
+
+    private void refreshBufferSize() {
+        int configuredSize = getLogCollectConfig().getBufferQueueSize();
+        if (configuredSize <= 0) {
+            throw new IllegalArgumentException("bufferQueueSize must be 
positive");
+        }
+        bufferSize = configuredSize;
     }
 
     @Override
@@ -94,7 +115,13 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
             BlockingQueue<L> bufferQueue = 
bufferQueueS.computeIfAbsent(selectorId, bufferQueueS -> initQueue(selectorId));
             bufferQueue.offer(log);
         } else {
-            bufferQueue.offer(log);
+            BlockingQueue<L> queue = bufferQueue;
+            synchronized (queue) {
+                // On shrink, retain the backlog but reject new logs until it 
falls below the new limit.
+                if (queue.size() < bufferSize) {
+                    queue.offer(log);
+                }
+            }
         }
     }
 
@@ -107,8 +134,8 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
     /**
      * batch and async consume.
      */
-    private void consume() {
-        while (started.get()) {
+    private void consume(final ThreadPoolExecutor consumerExecutor) {
+        while (!consumerExecutor.isShutdown()) {
             int diffTimeMSForPush = 100;
             try {
                 List<L> logs = new ArrayList<>();
@@ -284,8 +311,11 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
     protected abstract void desensitizeLog(L log, KeyWordMatch keyWordMatch, 
String desensitizeAlg);
 
     @Override
-    public void close() throws Exception {
+    public synchronized void close() throws Exception {
         started.set(false);
+        if (Objects.nonNull(executor)) {
+            executor.shutdownNow();
+        }
         AbstractLogConsumeClient<?, ?> logCollectClient = 
getLogConsumeClient();
         try {
             flushBufferQueues();
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 55540e804e..205682125f 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
@@ -32,18 +32,26 @@ import java.util.HashMap;
 import java.util.HashSet;
 import java.util.Map;
 import java.util.Set;
+import java.util.List;
+import java.util.ArrayList;
+import java.util.stream.IntStream;
+import org.apache.shenyu.common.concurrent.ShenyuThreadPoolExecutor;
+import org.mockito.MockedConstruction;
 import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.LinkedBlockingDeque;
+import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 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.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.argThat;
 import static org.mockito.Mockito.inOrder;
@@ -51,6 +59,9 @@ import static org.mockito.ArgumentMatchers.anyList;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.times;
+import static org.mockito.ArgumentMatchers.any;
 
 /**
  * The Test Case For AbstractLogCollector.
@@ -59,6 +70,8 @@ public class AbstractLogCollectorTest {
 
     private final AbstractLogConsumeClient<?, ShenyuRequestLog> 
logConsumeClient = mock(AbstractLogConsumeClient.class);
 
+    private final GenericGlobalConfig collectorConfig = new 
GenericGlobalConfig();
+
     private final AbstractLogCollector<AbstractLogConsumeClient<?, 
ShenyuRequestLog>, ShenyuRequestLog, GenericGlobalConfig> collector =
             new AbstractLogCollector<>() {
                 @Override
@@ -68,7 +81,7 @@ public class AbstractLogCollectorTest {
 
                 @Override
                 protected GenericGlobalConfig getLogCollectConfig() {
-                    return null;
+                    return collectorConfig;
                 }
 
                 @Override
@@ -76,6 +89,76 @@ public class AbstractLogCollectorTest {
                 }
             };
 
+    @Test
+    public void testRepeatedStartAndRestartDoNotLeakConsumers() throws 
Exception {
+        try {
+            collector.start();
+            ThreadPoolExecutor first = (ThreadPoolExecutor) 
ReflectionTestUtils.getField(collector, "executor");
+            Object queue = ReflectionTestUtils.getField(collector, 
"bufferQueue");
+            collector.start();
+            assertSame(first, ReflectionTestUtils.getField(collector, 
"executor"));
+            assertSame(queue, ReflectionTestUtils.getField(collector, 
"bufferQueue"));
+            collector.close();
+            collector.start();
+            assertNotSame(first, ReflectionTestUtils.getField(collector, 
"executor"));
+            assertTrue(first.awaitTermination(2, TimeUnit.SECONDS));
+        } finally {
+            collector.close();
+        }
+        ThreadPoolExecutor last = (ThreadPoolExecutor) 
ReflectionTestUtils.getField(collector, "executor");
+        assertTrue(last.awaitTermination(2, TimeUnit.SECONDS));
+    }
+
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testLiveCapacityChangesPreserveBacklogAndConsumer() throws 
Exception {
+        try (MockedConstruction<ShenyuThreadPoolExecutor> executors = 
mockConstruction(ShenyuThreadPoolExecutor.class)) {
+            collectorConfig.setBufferQueueSize(1);
+            collector.start();
+            final BlockingQueue<ShenyuRequestLog> queue = 
(BlockingQueue<ShenyuRequestLog>) ReflectionTestUtils.getField(collector, 
"bufferQueue");
+            final ShenyuRequestLog first = new ShenyuRequestLog();
+            final ShenyuRequestLog second = new ShenyuRequestLog();
+            final ShenyuRequestLog third = new ShenyuRequestLog();
+            first.setRequestUri("/first");
+            second.setRequestUri("/second");
+            third.setRequestUri("/third");
+            collector.collect(first);
+            collector.collect(second);
+            assertEquals(List.of(first), new ArrayList<>(queue));
+            collectorConfig.setBufferQueueSize(3);
+            collector.start();
+            collector.collect(second);
+            collector.collect(third);
+            collector.collect(new ShenyuRequestLog());
+            assertEquals(List.of(first, second, third), new 
ArrayList<>(queue));
+            collectorConfig.setBufferQueueSize(1);
+            collector.start();
+            collector.collect(new ShenyuRequestLog());
+            assertEquals(List.of(first, second, third), new 
ArrayList<>(queue));
+            queue.drainTo(new ArrayList<>());
+            collector.collect(first);
+            collector.collect(second);
+            assertEquals(List.of(first), new ArrayList<>(queue));
+            assertSame(queue, ReflectionTestUtils.getField(collector, 
"bufferQueue"));
+            assertEquals(1, executors.constructed().size());
+            verify(executors.constructed().get(0), 
times(1)).execute(any(Runnable.class));
+            collector.close();
+        }
+    }
+
+    @Test
+    public void testConcurrentProducersRespectLiveCapacity() throws Exception {
+        try (MockedConstruction<ShenyuThreadPoolExecutor> executors = 
mockConstruction(ShenyuThreadPoolExecutor.class)) {
+            collectorConfig.setBufferQueueSize(10);
+            collector.start();
+            IntStream.range(0, 1000).parallel().forEach(index -> 
collector.collect(new ShenyuRequestLog()));
+            BlockingQueue<?> queue = (BlockingQueue<?>) 
ReflectionTestUtils.getField(collector, "bufferQueue");
+            assertEquals(10, queue.size());
+            assertEquals(1, executors.constructed().size());
+            collector.close();
+        }
+    }
+
     @Test
     public void testSelectorQueueInitializationDoesNotReplaceGlobalQueue() 
throws Exception {
         final GenericGlobalConfig config = new GenericGlobalConfig();
@@ -115,6 +198,7 @@ public class AbstractLogCollectorTest {
 
     @Test
     public void testCollectAddsLogWhenBufferQueueHasCapacity() throws 
Exception {
+        ((AtomicBoolean) ReflectionTestUtils.getField(collector, 
"started")).set(true);
         BlockingQueue<ShenyuRequestLog> bufferQueue = new 
LinkedBlockingDeque<>(1);
         setField(collector, "bufferSize", 1);
         setField(collector, "bufferQueue", bufferQueue);
@@ -127,6 +211,7 @@ public class AbstractLogCollectorTest {
 
     @Test
     public void testCollectDoesNotThrowWhenBufferQueueIsFull() throws 
Exception {
+        ((AtomicBoolean) ReflectionTestUtils.getField(collector, 
"started")).set(true);
         ShenyuRequestLog bufferedLog = new ShenyuRequestLog();
         BlockingQueue<ShenyuRequestLog> bufferQueue = new 
StaleSizeLinkedBlockingDeque();
         bufferQueue.add(bufferedLog);
@@ -164,6 +249,7 @@ public class AbstractLogCollectorTest {
         BlockingQueue<ShenyuRequestLog> bufferQueue = new 
StaleSizeLinkedBlockingDeque();
         bufferQueue.add(bufferedLog);
         setField(multiClientCollector, "bufferSize", 1);
+        ((AtomicBoolean) ReflectionTestUtils.getField(multiClientCollector, 
"started")).set(true);
         getBufferQueues(multiClientCollector).put("selector", bufferQueue);
         ShenyuRequestLog log = new ShenyuRequestLog();
         log.setSelectorId("selector");

Reply via email to