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(); }
