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

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 66214e4600 [Fix][Connector-V2] Do not fail FTP listings on logout 
reset (#11695)
66214e4600 is described below

commit 66214e4600cc2baa4b07441eb08c08f3c97a78d1
Author: Jast <[email protected]>
AuthorDate: Thu Aug 20 23:11:48 2026 +0800

    [Fix][Connector-V2] Do not fail FTP listings on logout reset (#11695)
    
    Co-authored-by: zhangshenghang <[email protected]>
    Co-authored-by: davidzollo <[email protected]>
---
 .../file/ftp/system/SeaTunnelFTPFileSystem.java    | 32 ++++++++++++++--------
 .../ftp/system/SeaTunnelFTPFileSystemTest.java     | 18 ++++++++++++
 2 files changed, 38 insertions(+), 12 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
index f5cb9bfaf0..7cf322fd87 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
@@ -249,19 +249,31 @@ public class SeaTunnelFTPFileSystem extends FileSystem 
implements StreamingFileS
      * Logout and disconnect the given FTPClient. *
      *
      * @param client FTPClient
-     * @throws IOException IOException
      */
-    private void disconnect(FTPClient client) throws IOException {
-        if (client != null) {
-            if (!client.isConnected()) {
-                throw new FTPException("Client not connected");
-            }
+    void disconnect(FTPClient client) {
+        if (client == null || !client.isConnected()) {
+            return;
+        }
+
+        try {
             boolean logoutSuccess = client.logout();
-            client.disconnect();
             if (!logoutSuccess) {
                 LOG.warn(
                         "Logout failed while disconnecting, error code - " + 
client.getReplyCode());
             }
+        } catch (IOException e) {
+            // Some FTP servers close the control connection before responding 
to QUIT. The
+            // preceding operation has already completed, so do not turn a 
successful operation
+            // into a failure while releasing the connection.
+            LOG.warn("Failed to logout from FTP server while disconnecting", 
e);
+        } finally {
+            if (client.isConnected()) {
+                try {
+                    client.disconnect();
+                } catch (IOException e) {
+                    LOG.warn("Failed to disconnect from FTP server", e);
+                }
+            }
         }
     }
 
@@ -854,11 +866,7 @@ public class SeaTunnelFTPFileSystem extends FileSystem 
implements StreamingFileS
         } catch (IOException ioe) {
             throw new FTPException("Failed to get home directory", ioe);
         } finally {
-            try {
-                disconnect(client);
-            } catch (IOException ioe) {
-                throw new FTPException("Failed to disconnect", ioe);
-            }
+            disconnect(client);
         }
     }
 
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
index a29351bc8b..6c75ca9bd1 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
@@ -19,6 +19,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.file.ftp.system;
 
 import 
org.apache.seatunnel.connectors.seatunnel.file.hadoop.FileStatusListingSession;
 
+import org.apache.commons.net.ftp.FTPClient;
 import org.apache.commons.net.ftp.FTPFile;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.fs.FSDataInputStream;
@@ -38,17 +39,23 @@ import org.mockftpserver.fake.filesystem.FileSystem;
 import org.mockftpserver.fake.filesystem.UnixFakeFileSystem;
 
 import java.io.IOException;
+import java.net.SocketException;
 import java.net.URI;
 import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
 import java.util.Calendar;
 import java.util.List;
 
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 /** Unit tests for SeaTunnelFTPFileSystem. */
 public class SeaTunnelFTPFileSystemTest {
@@ -169,6 +176,17 @@ public class SeaTunnelFTPFileSystemTest {
         ftpFileSystem.delete(testDir, true);
     }
 
+    @Test
+    public void testDisconnectIgnoresConnectionResetDuringLogout() throws 
IOException {
+        FTPClient client = mock(FTPClient.class);
+        when(client.isConnected()).thenReturn(true);
+        doThrow(new SocketException("Connection reset")).when(client).logout();
+
+        assertDoesNotThrow(() -> ftpFileSystem.disconnect(client));
+
+        verify(client).disconnect();
+    }
+
     @Test
     public void testStreamingListingSkipsUnparseableEntries() throws 
IOException {
         FTPFile file = new FTPFile();

Reply via email to