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

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 0def1f305bf Keep pooled SFTP connection on application-level SFTP 
errors (#73288)
0def1f305bf is described below

commit 0def1f305bf2e0490bccaf97ecf9f16c1f5fc1c1
Author: Mathieu Monet <[email protected]>
AuthorDate: Wed Sep 23 22:19:56 2026 +0200

    Keep pooled SFTP connection on application-level SFTP errors (#73288)
    
    An SFTP status reply such as SFTPNoSuchFile or SFTPPermissionDenied means 
the
    server answered over a live channel, so the pooled connection behind it is
    still usable. Treating every exception as a transport fault made one missing
    file cost a full reconnect and re-auth, and because callers keep iterating 
the
    cost compounds: a production task turned 7,598 misses into 9,244 SSH
    connections and 135-190 MB of task logs per attempt.
    
    * Import SFTP exception classes directly and tighten pool tests
    
    The exception classes read better imported by name, which also lets asyncssh
    go back under TYPE_CHECKING since its remaining uses are annotations only.
    
    Two assertions could not fail as written: one ranged over a call list that 
is
    empty if the release never happens, the other compared client identities 
that
    the shared fixture guarantees whether or not the connection is pooled.
    
    The docstring claimed every other SFTPError is a server status reply, but
    asyncssh raises some client-side, so it now describes the property that
    actually matters for reuse.
---
 .../sftp/src/airflow/providers/sftp/pools/sftp.py  | 32 ++++++++-
 providers/sftp/tests/unit/sftp/pools/test_sftp.py  | 76 +++++++++++++++++++++-
 2 files changed, 105 insertions(+), 3 deletions(-)

diff --git a/providers/sftp/src/airflow/providers/sftp/pools/sftp.py 
b/providers/sftp/src/airflow/providers/sftp/pools/sftp.py
index 1115b4e19b4..8f449f54bc6 100644
--- a/providers/sftp/src/airflow/providers/sftp/pools/sftp.py
+++ b/providers/sftp/src/airflow/providers/sftp/pools/sftp.py
@@ -24,12 +24,33 @@ from threading import Lock
 from typing import TYPE_CHECKING
 from weakref import WeakKeyDictionary
 
+from asyncssh import SFTPBadMessage, SFTPConnectionLost, SFTPError, 
SFTPNoConnection
+
 from airflow.providers.sftp.hooks.sftp import SFTPHookAsync
 from airflow.utils.log.logging_mixin import LoggingMixin
 
 if TYPE_CHECKING:
     import asyncssh
 
+#: SFTP status codes that report a broken or desynchronised transport.
+_TRANSPORT_SFTP_ERRORS: tuple[type[SFTPError], ...] = (
+    SFTPNoConnection,
+    SFTPConnectionLost,
+    SFTPBadMessage,
+)
+
+
+def _is_connection_faulty(exc: BaseException) -> bool:
+    """
+    Whether ``exc`` means the pooled SSH connection can no longer be reused.
+
+    Any other ``SFTPError`` is a status-level failure that leaves the channel 
usable,
+    so the connection stays good; a non-``SFTPError`` is assumed to have 
broken it.
+    """
+    if isinstance(exc, SFTPError):
+        return isinstance(exc, _TRANSPORT_SFTP_ERRORS)
+    return True
+
 
 @dataclass
 class _LoopState:
@@ -236,9 +257,16 @@ class SFTPClientPool(LoggingMixin):
                 await self._release_pair(pair, state, faulty=True)
             raise
         except Exception as e:
-            self.log.warning("Dropping faulty connection for '%s': %s", 
self.sftp_conn_id, e)
+            faulty = _is_connection_faulty(e)
+            if faulty:
+                self.log.warning("Dropping faulty connection for '%s': %s", 
self.sftp_conn_id, e)
+            else:
+                # Debug: fires once per missing file when iterating many paths.
+                self.log.debug(
+                    "Keeping pooled connection for '%s' after SFTP error: %s", 
self.sftp_conn_id, e
+                )
             if pair:
-                await self._release_pair(pair, state, faulty=True)
+                await self._release_pair(pair, state, faulty=faulty)
             raise
         else:
             await self._release_pair(pair, state, faulty=False)
diff --git a/providers/sftp/tests/unit/sftp/pools/test_sftp.py 
b/providers/sftp/tests/unit/sftp/pools/test_sftp.py
index 98125f47198..44df80d693a 100644
--- a/providers/sftp/tests/unit/sftp/pools/test_sftp.py
+++ b/providers/sftp/tests/unit/sftp/pools/test_sftp.py
@@ -19,8 +19,17 @@ from __future__ import annotations
 import asyncio
 
 import pytest
+from asyncssh import (
+    SFTPBadMessage,
+    SFTPConnectionLost,
+    SFTPFailure,
+    SFTPNoConnection,
+    SFTPNoSuchFile,
+    SFTPOpUnsupported,
+    SFTPPermissionDenied,
+)
 
-from airflow.providers.sftp.pools.sftp import SFTPClientPool
+from airflow.providers.sftp.pools.sftp import SFTPClientPool, 
_is_connection_faulty
 
 
 class TestSFTPClientPool:
@@ -88,6 +97,71 @@ class TestSFTPClientPool:
 
         assert any(call.kwargs.get("faulty") is True for call in 
release_spy.call_args_list)
 
+    @pytest.mark.parametrize(
+        ("exc", "expected"),
+        [
+            pytest.param(SFTPNoSuchFile("missing"), False, id="no-such-file"),
+            pytest.param(SFTPPermissionDenied("denied"), False, 
id="permission-denied"),
+            pytest.param(SFTPFailure("failure"), False, id="failure"),
+            pytest.param(SFTPOpUnsupported("unsupported"), False, 
id="op-unsupported"),
+            pytest.param(SFTPConnectionLost("lost"), True, 
id="connection-lost"),
+            pytest.param(SFTPNoConnection("no connection"), True, 
id="no-connection"),
+            pytest.param(SFTPBadMessage("bad message"), True, 
id="bad-message"),
+            pytest.param(ValueError("boom"), True, id="not-an-sftp-error"),
+        ],
+    )
+    def test_is_connection_faulty(self, exc, expected):
+        assert _is_connection_faulty(exc) is expected
+
+    @pytest.mark.asyncio
+    async def 
test_get_sftp_client_keeps_connection_on_application_sftp_error(self, 
sftp_hook_mocked, mocker):
+        """An SFTP status reply means the server answered, so the connection 
stays pooled."""
+        pool = SFTPClientPool("app_error_conn", pool_size=1)
+        release_spy = mocker.spy(pool, "_release_pair")
+
+        with pytest.raises(SFTPNoSuchFile):
+            async with pool.get_sftp_client():
+                raise SFTPNoSuchFile("missing")
+
+        assert release_spy.call_args_list
+        assert all(call.kwargs.get("faulty") is False for call in 
release_spy.call_args_list)
+
+    @staticmethod
+    async def _miss(pool, seen: list) -> None:
+        """One failed lookup against the pool, recording the client it was 
served."""
+        async with pool.get_sftp_client() as sftp:
+            seen.append(sftp)
+            raise SFTPNoSuchFile("missing")
+
+    @pytest.mark.asyncio
+    async def 
test_get_sftp_client_reuses_connection_after_application_sftp_error(
+        self, sftp_hook_mocked, mocker
+    ):
+        """Repeated misses must not cost a reconnect each: the pooled pair is 
handed back."""
+        async with SFTPClientPool("reuse_conn", pool_size=1) as pool:
+            close_spy = mocker.spy(pool, "_close_connection_pair")
+            seen: list = []
+
+            for _ in range(5):
+                with pytest.raises(SFTPNoSuchFile):
+                    await self._miss(pool, seen)
+
+            assert close_spy.call_count == 0
+            assert len(seen) == 5
+
+    @pytest.mark.asyncio
+    async def 
test_get_sftp_client_marks_connection_faulty_on_transport_sftp_error(
+        self, sftp_hook_mocked, mocker
+    ):
+        pool = SFTPClientPool("transport_error_conn", pool_size=1)
+        release_spy = mocker.spy(pool, "_release_pair")
+
+        with pytest.raises(SFTPConnectionLost):
+            async with pool.get_sftp_client():
+                raise SFTPConnectionLost("lost")
+
+        assert any(call.kwargs.get("faulty") is True for call in 
release_spy.call_args_list)
+
     @pytest.mark.asyncio
     async def 
test_get_sftp_client_marks_connection_faulty_on_cancellation(self, 
sftp_hook_mocked, mocker):
         pool = SFTPClientPool("cancel_conn", pool_size=1)

Reply via email to