This is an automated email from the ASF dual-hosted git repository.

ppkarwasz pushed a commit to branch feat/fix-file-channel-windows
in repository https://gitbox.apache.org/repos/asf/logging-flume.git

commit fe0eeaa3a517166c295b794be73880eb393dbf4e
Author: Piotr P. Karwasz <[email protected]>
AuthorDate: Thu Jul 30 16:52:21 2026 +0200

    Close file channel queues and stores in tests
    
    Several tests leaked FlumeEventQueue instances or reused one temporary
    directory across test methods, which fails on Windows where the leaked
    MapDB files cannot be deleted. Also stop channels whose start() failed,
    now that FileChannel releases the log in that case.
    
    Assisted-By: Claude Fable 5 <[email protected]>
---
 .../apache/flume/channel/file/TestCheckpoint.java  |  9 ++++++++
 .../channel/file/TestCheckpointRebuilder.java      |  9 ++++++--
 .../file/TestEventQueueBackingStoreFactory.java    | 24 +++++++++++++---------
 .../flume/channel/file/TestFileChannelBase.java    |  3 ++-
 .../flume/channel/file/TestFlumeEventQueue.java    | 14 +++++++++++++
 5 files changed, 46 insertions(+), 13 deletions(-)

diff --git 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpoint.java
 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpoint.java
index 54b51222..f37cbd95 100644
--- 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpoint.java
+++ 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpoint.java
@@ -19,6 +19,7 @@ package org.apache.flume.channel.file;
 import java.io.File;
 import java.io.IOException;
 import junit.framework.Assert;
+import org.apache.commons.io.FileUtils;
 import org.apache.flume.channel.file.instrumentation.FileChannelCounter;
 import org.junit.After;
 import org.junit.Before;
@@ -43,6 +44,9 @@ public class TestCheckpoint {
     @After
     public void cleanup() {
         file.delete();
+        inflightPuts.delete();
+        inflightTakes.delete();
+        FileUtils.deleteQuietly(queueSet);
     }
 
     @Test
@@ -52,12 +56,17 @@ public class TestCheckpoint {
         FlumeEventPointer ptrIn = new FlumeEventPointer(10, 20);
         FlumeEventQueue queueIn = new FlumeEventQueue(backingStore, 
inflightTakes, inflightPuts, queueSet);
         queueIn.addHead(ptrIn);
+        // The three queues share one backing store and one queue set 
directory, so each
+        // queue must release the queue set database before the next queue is 
created.
+        queueIn.replayComplete();
         FlumeEventQueue queueOut = new FlumeEventQueue(backingStore, 
inflightTakes, inflightPuts, queueSet);
         Assert.assertEquals(0, queueOut.getLogWriteOrderID());
+        queueOut.replayComplete();
         queueIn.checkpoint(false);
         FlumeEventQueue queueOut2 = new FlumeEventQueue(backingStore, 
inflightTakes, inflightPuts, queueSet);
         FlumeEventPointer ptrOut = queueOut2.removeHead(0L);
         Assert.assertEquals(ptrIn, ptrOut);
         Assert.assertTrue(queueOut2.getLogWriteOrderID() > 0);
+        queueOut2.close();
     }
 }
diff --git 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpointRebuilder.java
 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpointRebuilder.java
index 7ebebc7e..649c72f4 100644
--- 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpointRebuilder.java
+++ 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestCheckpointRebuilder.java
@@ -65,8 +65,13 @@ public class TestCheckpointRebuilder extends 
TestFileChannelBase {
         EventQueueBackingStore backingStore =
                 EventQueueBackingStoreFactory.get(checkpointFile, 50, "test", 
new FileChannelCounter("test"));
         FlumeEventQueue queue = new FlumeEventQueue(backingStore, 
inflightTakesFile, inflightPutsFile, queueSetDir);
-        CheckpointRebuilder checkpointRebuilder = new 
CheckpointRebuilder(getAllLogs(dataDirs), queue, true);
-        Assert.assertTrue(checkpointRebuilder.rebuild());
+        try {
+            CheckpointRebuilder checkpointRebuilder = new 
CheckpointRebuilder(getAllLogs(dataDirs), queue, true);
+            Assert.assertTrue(checkpointRebuilder.rebuild());
+        } finally {
+            // Release the checkpoint files before the channel below replays 
them.
+            queue.close();
+        }
         channel = createFileChannel(overrides);
         channel.start();
         Assert.assertTrue(channel.isOpen());
diff --git 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestEventQueueBackingStoreFactory.java
 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestEventQueueBackingStoreFactory.java
index fb26f382..1c8c8f12 100644
--- 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestEventQueueBackingStoreFactory.java
+++ 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestEventQueueBackingStoreFactory.java
@@ -272,16 +272,20 @@ public class TestEventQueueBackingStoreFactory {
     private void verify(EventQueueBackingStore backingStore, long 
expectedVersion, List<Long> expectedPointers)
             throws Exception {
         FlumeEventQueue queue = new FlumeEventQueue(backingStore, 
inflightTakes, inflightPuts, queueSetDir);
-        List<Long> actualPointers = Lists.newArrayList();
-        FlumeEventPointer ptr;
-        while ((ptr = queue.removeHead(0L)) != null) {
-            actualPointers.add(ptr.toLong());
+        try {
+            List<Long> actualPointers = Lists.newArrayList();
+            FlumeEventPointer ptr;
+            while ((ptr = queue.removeHead(0L)) != null) {
+                actualPointers.add(ptr.toLong());
+            }
+            Assert.assertEquals(expectedPointers, actualPointers);
+            Assert.assertEquals(10, backingStore.getCapacity());
+            DataInputStream in = new DataInputStream(new 
FileInputStream(checkpoint));
+            long actualVersion = in.readLong();
+            Assert.assertEquals(expectedVersion, actualVersion);
+            in.close();
+        } finally {
+            queue.close();
         }
-        Assert.assertEquals(expectedPointers, actualPointers);
-        Assert.assertEquals(10, backingStore.getCapacity());
-        DataInputStream in = new DataInputStream(new 
FileInputStream(checkpoint));
-        long actualVersion = in.readLong();
-        Assert.assertEquals(expectedVersion, actualVersion);
-        in.close();
     }
 }
diff --git 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFileChannelBase.java
 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFileChannelBase.java
index 46257d14..bc821af2 100644
--- 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFileChannelBase.java
+++ 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFileChannelBase.java
@@ -70,7 +70,8 @@ public class TestFileChannelBase {
 
     @After
     public void teardown() {
-        if (channel != null && channel.isOpen()) {
+        // Stop the channel even when it failed to start: it still holds open 
files.
+        if (channel != null) {
             channel.stop();
         }
         FileUtils.deleteQuietly(baseDir);
diff --git 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFlumeEventQueue.java
 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFlumeEventQueue.java
index c8c3c2a9..c3c4bb47 100644
--- 
a/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFlumeEventQueue.java
+++ 
b/flume-ng-channels/flume-file-channel/src/test/java/org/apache/flume/channel/file/TestFlumeEventQueue.java
@@ -54,6 +54,15 @@ public class TestFlumeEventQueue {
         File queueSetDir;
 
         EventQueueBackingStoreSupplier() {
+            reset();
+        }
+
+        /**
+         * The supplier instances are shared by all test methods of a 
parameter, so use a
+         * fresh directory for each test: a file leaked by one test (e.g. 
still open on
+         * Windows) must not break the cleanup of the following tests.
+         */
+        void reset() {
             baseDir = Files.createTempDir();
             checkpoint = new File(baseDir, "checkpoint");
             inflightTakes = new File(baseDir, "inflightputs");
@@ -117,11 +126,16 @@ public class TestFlumeEventQueue {
 
     @Before
     public void setup() throws Exception {
+        backingStoreSupplier.reset();
         backingStore = backingStoreSupplier.get();
     }
 
     @After
     public void cleanup() throws IOException {
+        if (queue != null) {
+            queue.close();
+            queue = null;
+        }
         if (backingStore != null) {
             backingStore.close();
         }

Reply via email to