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 e3990f1041 fix(logging): flush buffered logs on close (#7127)
e3990f1041 is described below

commit e3990f1041f39fd23fb913c61ddc8c70f57c4de6
Author: Liming Deng <[email protected]>
AuthorDate: Thu Oct 1 07:05:19 2026 +0800

    fix(logging): flush buffered logs on close (#7127)
---
 .../common/collector/AbstractLogCollector.java     | 35 ++++++++++-
 .../common/collector/AbstractLogCollectorTest.java | 72 ++++++++++++++++++++++
 2 files changed, 104 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 03d2dae9f2..535022155a 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
@@ -86,7 +86,7 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
 
     @Override
     public void collect(final L log) {
-        if (Objects.isNull(log) || 
Objects.isNull(getLogConsumeClient(log.getSelectorId()))) {
+        if (!started.get() || Objects.isNull(log) || 
Objects.isNull(getLogConsumeClient(log.getSelectorId()))) {
             return;
         }
         if (getMultiClient()) {
@@ -181,6 +181,31 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
         }
     }
 
+    private void flushBufferQueues() throws Exception {
+        if (getMultiClient()) {
+            for (Map.Entry<String, BlockingQueue<L>> entry : 
bufferQueueS.entrySet()) {
+                flushBufferQueue(entry.getValue(), 
getLogConsumeClient(entry.getKey()));
+            }
+        } else {
+            flushBufferQueue(bufferQueue, getLogConsumeClient());
+        }
+    }
+
+    private void flushBufferQueue(final BlockingQueue<L> queue, final 
AbstractLogConsumeClient<?, L> logConsumeClient) throws Exception {
+        if (Objects.isNull(queue) || Objects.isNull(logConsumeClient)) {
+            return;
+        }
+        int batchSize = 100;
+        while (!queue.isEmpty()) {
+            List<L> logs = new ArrayList<>(batchSize);
+            queue.drainTo(logs, batchSize);
+            if (logs.isEmpty()) {
+                return;
+            }
+            logConsumeClient.consume(logs);
+        }
+    }
+
     private void desensitizeShenyuRequestLog(final L logInfo, final 
KeyWordMatch keyWordMatch, final String desensitizedAlg) {
         
logInfo.setClientIp(desensitizeForSingleWord(GenericLoggingConstant.CLIENT_IP, 
logInfo.getClientIp(), keyWordMatch, desensitizedAlg));
         
logInfo.setTimeLocal(desensitizeForSingleWord(GenericLoggingConstant.TIME_LOCAL,
 logInfo.getTimeLocal(), keyWordMatch, desensitizedAlg));
@@ -262,8 +287,12 @@ public abstract class AbstractLogCollector<T extends 
AbstractLogConsumeClient<?,
     public void close() throws Exception {
         started.set(false);
         AbstractLogConsumeClient<?, ?> logCollectClient = 
getLogConsumeClient();
-        if (Objects.nonNull(logCollectClient)) {
-            logCollectClient.close();
+        try {
+            flushBufferQueues();
+        } finally {
+            if (Objects.nonNull(logCollectClient)) {
+                logCollectClient.close();
+            }
         }
     }
 }
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 7b5973cfdd..55540e804e 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
@@ -44,6 +44,9 @@ 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.assertTrue;
+import static org.mockito.ArgumentMatchers.argThat;
+import static org.mockito.Mockito.inOrder;
 import static org.mockito.ArgumentMatchers.anyList;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
@@ -247,6 +250,75 @@ public class AbstractLogCollectorTest {
         assertEquals(15L, log.getUpstreamResponseTime());
     }
 
+    @Test
+    public void testCloseFlushesBufferedLogsBeforeClosingClient() throws 
Exception {
+        BlockingQueue<ShenyuRequestLog> bufferQueue = new 
LinkedBlockingDeque<>(2);
+        ShenyuRequestLog first = new ShenyuRequestLog();
+        ShenyuRequestLog second = new ShenyuRequestLog();
+        bufferQueue.add(first);
+        bufferQueue.add(second);
+        setField(collector, "bufferQueue", bufferQueue);
+
+        collector.close();
+
+        org.mockito.InOrder closeOrder = inOrder(logConsumeClient);
+        closeOrder.verify(logConsumeClient).consume(argThat(logs -> 
logs.size() == 2
+                && logs.get(0) == first && logs.get(1) == second));
+        closeOrder.verify(logConsumeClient).close();
+        assertTrue(bufferQueue.isEmpty());
+    }
+
+    @Test
+    public void testCloseFlushesEveryMultiClientBuffer() throws Exception {
+        AbstractLogConsumeClient<?, ShenyuRequestLog> firstClient = 
mock(AbstractLogConsumeClient.class);
+        AbstractLogConsumeClient<?, ShenyuRequestLog> secondClient = 
mock(AbstractLogConsumeClient.class);
+        Map<String, AbstractLogConsumeClient<?, ShenyuRequestLog>> clients = 
new HashMap<>();
+        clients.put("first", firstClient);
+        clients.put("second", secondClient);
+        AbstractLogCollector<AbstractLogConsumeClient<?, ShenyuRequestLog>, 
ShenyuRequestLog, GenericGlobalConfig> multiClientCollector =
+                new AbstractLogCollector<>() {
+                    @Override
+                    protected AbstractLogConsumeClient<?, ShenyuRequestLog> 
getLogConsumeClient() {
+                        return logConsumeClient;
+                    }
+
+                    @Override
+                    protected AbstractLogConsumeClient<?, ShenyuRequestLog> 
getLogConsumeClient(final String selectorId) {
+                        return clients.get(selectorId);
+                    }
+
+                    @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 first = new ShenyuRequestLog();
+        ShenyuRequestLog second = new ShenyuRequestLog();
+        BlockingQueue<ShenyuRequestLog> firstQueue = new 
LinkedBlockingDeque<>(1);
+        BlockingQueue<ShenyuRequestLog> secondQueue = new 
LinkedBlockingDeque<>(1);
+        firstQueue.add(first);
+        secondQueue.add(second);
+        getBufferQueues(multiClientCollector).put("first", firstQueue);
+        getBufferQueues(multiClientCollector).put("second", secondQueue);
+
+        multiClientCollector.close();
+
+        verify(firstClient).consume(argThat(logs -> logs.size() == 1 && 
logs.get(0) == first));
+        verify(secondClient).consume(argThat(logs -> logs.size() == 1 && 
logs.get(0) == second));
+        verify(logConsumeClient).close();
+        assertTrue(firstQueue.isEmpty());
+        assertTrue(secondQueue.isEmpty());
+    }
+
     private static void setField(final Object target, final String fieldName, 
final Object value) throws Exception {
         Field field = AbstractLogCollector.class.getDeclaredField(fieldName);
         field.setAccessible(true);

Reply via email to