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

Reply via email to