This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12472-e1a7e64b6733d2b2cd4682190c753a81fac93bc2 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 87ed0afe0f783eda1b1f3ac6f0ac9fadedcb4f98 Author: Goutam Adwant <[email protected]> AuthorDate: Sat Oct 3 15:49:48 2026 +0000 [Test][Zeta] Drive FileCollectReader idle close and rediscovery with a manual clock (#12472) --- .../agent/connector/file/FileCollectReader.java | 24 +++++++++--- .../connector/file/cursor/FileTailCursor.java | 17 +++++++-- .../file/FileCollectReaderBehaviorTest.java | 43 +++++++++------------- 3 files changed, 51 insertions(+), 33 deletions(-) diff --git a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReader.java b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReader.java index 31f84d87a0..1db18c601c 100644 --- a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReader.java +++ b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReader.java @@ -46,6 +46,7 @@ import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.function.LongSupplier; public class FileCollectReader implements EdgeInputReader { @@ -54,6 +55,7 @@ public class FileCollectReader implements EdgeInputReader { private final FileCollectConfig config; private final Charset charset; private final EdgeSourcePositionStore sourcePositionStore; + private final LongSupplier clock; private GlobPathResolver globResolver; private final Map<Path, MultilineAssembler> multilineAssemblers = new HashMap<>(); @@ -66,9 +68,21 @@ public class FileCollectReader implements EdgeInputReader { public FileCollectReader( FileCollectConfig config, EdgeSourcePositionStore sourcePositionStore) { + this(config, sourcePositionStore, System::currentTimeMillis); + } + + /** + * @param clock millisecond clock used for glob scan intervals, idle cursor timeouts and line + * timestamps; tests pass a manual clock to step through these transitions deterministically + */ + FileCollectReader( + FileCollectConfig config, + EdgeSourcePositionStore sourcePositionStore, + LongSupplier clock) { this.config = config; this.charset = config.getCharset(); this.sourcePositionStore = sourcePositionStore; + this.clock = clock; } public String id() { @@ -105,13 +119,13 @@ public class FileCollectReader implements EdgeInputReader { lineCounters.put(file, restoredLineNumber(pos)); } - this.lastDiscoveryMs = System.currentTimeMillis(); + this.lastDiscoveryMs = clock.getAsLong(); } @Override public List<EdgeEvent> poll(int maxRecords) throws Exception { List<EdgeEvent> records = new ArrayList<>(); - long now = System.currentTimeMillis(); + long now = clock.getAsLong(); // Periodically scan glob patterns for newly appeared files discoverNewFiles(now); @@ -189,7 +203,7 @@ public class FileCollectReader implements EdgeInputReader { String line, List<EdgeEvent> records) { long lineNum = lineCounters.merge(filePath, 1L, Long::sum); - long ts = System.currentTimeMillis(); + long ts = clock.getAsLong(); MultilineAssembler.LineElement element = new MultilineAssembler.LineElement(line, filePathStr, lineNum, cursor.offset(), ts); @@ -320,7 +334,7 @@ public class FileCollectReader implements EdgeInputReader { } private FileTailCursor openCursor(Path file, EdgeSourcePosition pos) throws IOException { - FileTailCursor cursor = new FileTailCursor(file, charset); + FileTailCursor cursor = new FileTailCursor(file, charset, clock); long seekOffset = 0; if (pos != null && pos.getOffset() > 0) { seekOffset = pos.getOffset(); @@ -331,7 +345,7 @@ public class FileCollectReader implements EdgeInputReader { if (pos != null && inode(pos) != 0 && cursor.inode() != 0 && inode(pos) != cursor.inode()) { cursor.close(); - cursor = new FileTailCursor(file, charset); + cursor = new FileTailCursor(file, charset, clock); cursor.open(0); } return cursor; diff --git a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/cursor/FileTailCursor.java b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/cursor/FileTailCursor.java index b21f84617d..041c52bee5 100644 --- a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/cursor/FileTailCursor.java +++ b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/main/java/org/apache/seatunnel/edge/agent/connector/file/cursor/FileTailCursor.java @@ -24,6 +24,7 @@ import java.io.RandomAccessFile; import java.nio.charset.Charset; import java.nio.file.Files; import java.nio.file.Path; +import java.util.function.LongSupplier; public class FileTailCursor implements Closeable { @@ -31,6 +32,7 @@ public class FileTailCursor implements Closeable { private final Path path; private final Charset charset; + private final LongSupplier clock; private RandomAccessFile raf; private long currentOffset; private long inode; @@ -41,8 +43,17 @@ public class FileTailCursor implements Closeable { private int bufLen; public FileTailCursor(Path path, Charset charset) { + this(path, charset, System::currentTimeMillis); + } + + /** + * @param clock source of the millisecond timestamps recorded as the last activity; tests can + * pass a manual clock to drive idle timeouts deterministically + */ + public FileTailCursor(Path path, Charset charset, LongSupplier clock) { this.path = path; this.charset = charset; + this.clock = clock; this.currentOffset = 0L; this.inode = 0L; this.lastActivityMs = 0L; @@ -60,7 +71,7 @@ public class FileTailCursor implements Closeable { } this.bufPos = 0; this.bufLen = 0; - this.lastActivityMs = System.currentTimeMillis(); + this.lastActivityMs = clock.getAsLong(); } /** @@ -98,7 +109,7 @@ public class FileTailCursor implements Closeable { len--; } this.currentOffset = raf.getFilePointer() - (bufLen - bufPos); - this.lastActivityMs = System.currentTimeMillis(); + this.lastActivityMs = clock.getAsLong(); return new String(arr, 0, len, charset); } baos.write(b); @@ -119,7 +130,7 @@ public class FileTailCursor implements Closeable { this.currentOffset = 0; this.bufPos = 0; this.bufLen = 0; - this.lastActivityMs = System.currentTimeMillis(); + this.lastActivityMs = clock.getAsLong(); } public long offset() { diff --git a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/test/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReaderBehaviorTest.java b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/test/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReaderBehaviorTest.java index 0f15f67181..2eb1a828bb 100644 --- a/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/test/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReaderBehaviorTest.java +++ b/seatunnel-edge-agent/seatunnel-edge-agent-connector/src/test/java/org/apache/seatunnel/edge/agent/connector/file/FileCollectReaderBehaviorTest.java @@ -22,7 +22,6 @@ import org.apache.seatunnel.edge.agent.connector.EdgeEvent; import org.apache.seatunnel.edge.agent.connector.config.FileCollectConfig; import org.apache.seatunnel.edge.agent.connector.config.FileCollectOptions; -import org.awaitility.Awaitility; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -37,6 +36,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import static org.awaitility.Awaitility.await; @@ -45,10 +45,9 @@ public class FileCollectReaderBehaviorTest { private static final long CLOSE_INACTIVE_MS = 500L; private static final long GLOB_SCAN_INTERVAL_MS = 20L; - // Ceiling for the wall-clock choreography below. The common case finishes in milliseconds, - // but the conditions depend on real filesystem scans plus the idle windows above, so keep - // the ceiling well above both: stalled CI runners (Windows runners have shown multi-second - // scheduling stalls) must fail on reader behavior, not on scheduling delay. + // Ceiling for the remaining wall-clock wait, which depends on real filesystem reads. The + // common case finishes in milliseconds, but stalled CI runners (Windows runners have shown + // multi-second scheduling stalls) must fail on reader behavior, not on scheduling delay. private static final long AWAIT_BUDGET_SECONDS = 10L; @TempDir Path tempDir; @@ -63,10 +62,13 @@ public class FileCollectReaderBehaviorTest { map.put(FileCollectOptions.GLOB_SCAN_INTERVAL_MS.key(), GLOB_SCAN_INTERVAL_MS); map.put(FileCollectOptions.READ_FROM_BEGINNING.key(), false); + // Idle close and glob rediscovery are driven by this clock, not by wall-clock time. + AtomicLong clock = new AtomicLong(1_000_000L); FileCollectReader reader = new FileCollectReader( FileCollectConfig.from(ReadonlyConfig.fromMap(map)), - new NoOpPositionStore()); + new NoOpPositionStore(), + clock::get); reader.open(); try { Assertions.assertTrue(reader.poll(10).isEmpty()); @@ -75,26 +77,13 @@ public class FileCollectReaderBehaviorTest { logFile, "first\n".getBytes(StandardCharsets.UTF_8), StandardOpenOption.APPEND); Assertions.assertEquals(1, reader.poll(10).size()); - long idleDeadlineMs = System.currentTimeMillis() + CLOSE_INACTIVE_MS + 30L; - Awaitility.await() - .atMost(AWAIT_BUDGET_SECONDS, TimeUnit.SECONDS) - .pollInterval(GLOB_SCAN_INTERVAL_MS, TimeUnit.MILLISECONDS) - .until( - () -> { - reader.poll(10); - return System.currentTimeMillis() >= idleDeadlineMs; - }); + // Past the idle window: this poll closes the inactive cursor. Discovery runs before + // the close within a poll, so the file is not picked up again by the same poll. + clock.addAndGet(CLOSE_INACTIVE_MS + 1); + Assertions.assertTrue(reader.poll(10).isEmpty()); - // Ensure glob scan rediscovers the file before new data is appended - Awaitility.await() - .atMost(AWAIT_BUDGET_SECONDS, TimeUnit.SECONDS) - .pollInterval(5, TimeUnit.MILLISECONDS) - .until( - () -> { - reader.poll(10); - return System.currentTimeMillis() - >= idleDeadlineMs + GLOB_SCAN_INTERVAL_MS; - }); + // One glob scan interval later the file is rediscovered and reopened at its end. + clock.addAndGet(GLOB_SCAN_INTERVAL_MS); Assertions.assertTrue(reader.poll(10).isEmpty()); Files.write( @@ -113,6 +102,10 @@ public class FileCollectReaderBehaviorTest { events.get(0).getPayload(), StandardCharsets.UTF_8); Assertions.assertTrue(payload.contains("second")); Assertions.assertFalse(payload.contains("first")); + // Line numbering restarts only if the cursor was closed and the + // file was picked up again by the glob scan. + Assertions.assertEquals( + "1", events.get(0).getMetadata().get("line")); }); } finally { reader.close();
