This is an automated email from the ASF dual-hosted git repository.
vincbeck 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 38168313b24 Fix prefix handling in Amazon transfer operators (#73269)
38168313b24 is described below
commit 38168313b24efd30e8826854058f33c5c839f6f2
Author: Yuseok Jo <[email protected]>
AuthorDate: Wed Sep 23 22:21:53 2026 +0900
Fix prefix handling in Amazon transfer operators (#73269)
---
providers/amazon/docs/changelog.rst | 16 ++++
.../providers/amazon/aws/transfers/ftp_to_s3.py | 34 +++++--
.../providers/amazon/aws/transfers/s3_to_ftp.py | 35 ++++++-
.../providers/amazon/aws/transfers/s3_to_sftp.py | 35 ++++++-
.../providers/amazon/aws/transfers/sftp_to_s3.py | 22 ++++-
.../unit/amazon/aws/transfers/test_ftp_to_s3.py | 105 ++++++++++++++++++++-
.../unit/amazon/aws/transfers/test_s3_to_ftp.py | 51 ++++++++++
.../unit/amazon/aws/transfers/test_s3_to_sftp.py | 55 +++++++++++
.../unit/amazon/aws/transfers/test_sftp_to_s3.py | 34 +++++++
9 files changed, 364 insertions(+), 23 deletions(-)
diff --git a/providers/amazon/docs/changelog.rst
b/providers/amazon/docs/changelog.rst
index fa3f4e9ee49..378c44d4540 100644
--- a/providers/amazon/docs/changelog.rst
+++ b/providers/amazon/docs/changelog.rst
@@ -35,6 +35,22 @@ Changelog
botocore configuration, which previously never reached it. Set these
explicitly on the operator
if the deferred half needs to differ from the synchronous half.
+.. warning::
+ ``FTPToS3Operator``, ``S3ToFTPOperator``, ``S3ToSFTPOperator``, and
``SFTPToS3Operator`` now match
+ the source ``*_filenames``, when it is a string other than ``"*"``, as a
leading prefix of the
+ file name, as documented, instead of as a substring anywhere in the listed
entry. Entries that
+ contained it only elsewhere are no longer selected, and a warning reports
how many, naming up to
+ ten. Renaming now replaces only that leading prefix, instead of every
occurrence.
+
+ For ``S3ToFTPOperator`` and ``S3ToSFTPOperator`` the prefix is not an S3 key
prefix. Keys under
+ ``s3_key`` are matched on their last path segment and keep their directory
at the destination,
+ and a prefix that contains ``/`` no longer matches. To select a
subdirectory, narrow ``s3_key``
+ to it and append it to ``ftp_path`` or ``sftp_path`` to keep the same
destination.
+
+ With a string ``ftp_filenames``, ``FTPToS3Operator`` now builds the
destination key from the file
+ name alone, so on servers that qualify ``nlst`` entries with the listed
directory the key no
+ longer embeds that directory.
+
9.36.0
......
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/transfers/ftp_to_s3.py
b/providers/amazon/src/airflow/providers/amazon/aws/transfers/ftp_to_s3.py
index 80796534772..e7067d388af 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/transfers/ftp_to_s3.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/transfers/ftp_to_s3.py
@@ -18,6 +18,7 @@
from __future__ import annotations
import ftplib
+import posixpath
from collections.abc import Sequence
from tempfile import NamedTemporaryFile
from typing import TYPE_CHECKING
@@ -30,6 +31,9 @@ if TYPE_CHECKING:
from airflow.sdk import Context
+SKIPPED_SAMPLE_SIZE = 10
+
+
class FTPToS3Operator(BaseOperator):
"""
Transfer of one or more files from an FTP server to S3.
@@ -143,19 +147,37 @@ class FTPToS3Operator(BaseOperator):
path=self.ftp_path,
)
- if self.ftp_filenames == "*":
+ ftp_prefix: str = self.ftp_filenames
+ if ftp_prefix == "*":
files = list_dir
else:
- ftp_filename: str = self.ftp_filenames
- files = [f for f in list_dir if ftp_filename in f]
+ # ``nlst`` may qualify entries with the listed directory,
so the prefix applies
+ # to the file name, while the substring test mirrors the
old rule over the whole entry.
+ files, dropped = [], []
+ for entry in list_dir:
+ if posixpath.basename(entry).startswith(ftp_prefix):
+ files.append(entry)
+ elif ftp_prefix in entry:
+ dropped.append(entry)
+ if dropped:
+ omitted = len(dropped) - SKIPPED_SAMPLE_SIZE
+ self.log.warning(
+ "%d file(s) contain %r but are not selected,
because a string prefix "
+ "matches only at the start of the filename: %s%s",
+ len(dropped),
+ ftp_prefix,
+ dropped[:SKIPPED_SAMPLE_SIZE],
+ f" and {omitted} more" if omitted > 0 else "",
+ )
for file in files:
self.log.info("Moving file %s", file)
+ # The entry is kept as listed for retrieval, but the
destination key is
+ # built from the file name so it never embeds the source
directory.
+ filename = posixpath.basename(file)
if self.s3_filenames and isinstance(self.s3_filenames,
str):
- filename = file.replace(self.ftp_filenames,
self.s3_filenames)
- else:
- filename = file
+ filename = filename.replace(ftp_prefix,
self.s3_filenames, 1)
s3_file_key = f"{self.s3_key}{filename}"
self.__upload_to_s3_from_ftp(file, s3_file_key)
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_ftp.py
b/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_ftp.py
index e89a34a40f7..4666d082b22 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_ftp.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_ftp.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import posixpath
from collections.abc import Sequence
from tempfile import NamedTemporaryFile
from typing import TYPE_CHECKING
@@ -30,6 +31,9 @@ if TYPE_CHECKING:
from airflow.sdk import Context
+SKIPPED_SAMPLE_SIZE = 10
+
+
class S3ToFTPOperator(BaseOperator):
"""
This operator enables the transferring of files from S3 to a FTP server.
@@ -45,7 +49,8 @@ class S3ToFTPOperator(BaseOperator):
``"/"``.
:param s3_filenames: Only used if you want to move multiple files. You can
pass
a list with exact key suffixes present under the s3_key prefix, or a
string
- prefix that all filenames must match. Use ``"*"`` to move all objects
under
+ prefix that all filenames must match. The prefix applies to the file
name, the
+ last segment of each key under s3_key. Use ``"*"`` to move all objects
under
the s3_key prefix.
:param ftp_path: The ftp remote path. For a single file it must include
the file
path. For multiple files it is the destination directory path and must
end
@@ -117,16 +122,36 @@ class S3ToFTPOperator(BaseOperator):
self.log.info("Getting files in s3://%s/%s", self.s3_bucket,
self.s3_key)
all_keys = s3_hook.list_keys(bucket_name=self.s3_bucket,
prefix=self.s3_key) or []
filenames = [k[len(self.s3_key) :] for k in all_keys]
- if self.s3_filenames == "*":
+ s3_prefix: str = self.s3_filenames
+ if s3_prefix == "*":
files = filenames
else:
- s3_prefix: str = self.s3_filenames
- files = [f for f in filenames if s3_prefix in f]
+ # ``list_keys`` recurses, so the prefix applies to the
file name, while the
+ # substring test mirrors the old rule over the whole
relative key.
+ files = [f for f in filenames if
posixpath.basename(f).startswith(s3_prefix)]
+ dropped = [
+ f
+ for f in filenames
+ if s3_prefix in f and not
posixpath.basename(f).startswith(s3_prefix)
+ ]
+ if dropped:
+ omitted = len(dropped) - SKIPPED_SAMPLE_SIZE
+ self.log.warning(
+ "%d file(s) contain %r but are not selected,
because a string prefix "
+ "matches only at the start of the filename: %s%s",
+ len(dropped),
+ s3_prefix,
+ dropped[:SKIPPED_SAMPLE_SIZE],
+ f" and {omitted} more" if omitted > 0 else "",
+ )
for file in files:
self.log.info("Moving file %s", file)
if self.ftp_filenames and isinstance(self.ftp_filenames,
str):
- ftp_filename = file.replace(self.s3_filenames,
self.ftp_filenames)
+ name = posixpath.basename(file)
+ ftp_filename = posixpath.join(
+ posixpath.dirname(file), name.replace(s3_prefix,
self.ftp_filenames, 1)
+ )
else:
ftp_filename = file
self._download_from_s3(
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_sftp.py
b/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_sftp.py
index 004232c2296..7b305362228 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_sftp.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_sftp.py
@@ -17,6 +17,7 @@
# under the License.
from __future__ import annotations
+import posixpath
from collections.abc import Sequence
from tempfile import NamedTemporaryFile
from typing import TYPE_CHECKING
@@ -33,6 +34,9 @@ if TYPE_CHECKING:
from airflow.sdk import Context
+SKIPPED_SAMPLE_SIZE = 10
+
+
class S3ToSFTPOperator(BaseOperator):
"""
This operator enables the transferring of files from S3 to a SFTP server.
@@ -60,7 +64,8 @@ class S3ToSFTPOperator(BaseOperator):
``"/"``.
:param s3_filenames: Only used if you want to move multiple files. You can
pass
a list with exact key suffixes present under the s3_key prefix, or a
string
- prefix that all filenames must match. Use ``"*"`` to move all objects
under
+ prefix that all filenames must match. The prefix applies to the file
name, the
+ last segment of each key under s3_key. Use ``"*"`` to move all objects
under
the s3_key prefix.
:param sftp_filenames: Only used if you want to move multiple files and
name them
differently at the destination. It can be a list of filenames or a
string
@@ -145,16 +150,36 @@ class S3ToSFTPOperator(BaseOperator):
self.log.info("Getting files in s3://%s/%s", self.s3_bucket,
self.s3_key)
all_keys = s3_hook.list_keys(bucket_name=self.s3_bucket,
prefix=self.s3_key) or []
filenames = [k[len(self.s3_key) :] for k in all_keys]
- if self.s3_filenames == "*":
+ s3_prefix: str = self.s3_filenames
+ if s3_prefix == "*":
files = filenames
else:
- s3_prefix: str = self.s3_filenames
- files = [f for f in filenames if s3_prefix in f]
+ # ``list_keys`` recurses, so the prefix applies to the
file name, while the
+ # substring test mirrors the old rule over the whole
relative key.
+ files = [f for f in filenames if
posixpath.basename(f).startswith(s3_prefix)]
+ dropped = [
+ f
+ for f in filenames
+ if s3_prefix in f and not
posixpath.basename(f).startswith(s3_prefix)
+ ]
+ if dropped:
+ omitted = len(dropped) - SKIPPED_SAMPLE_SIZE
+ self.log.warning(
+ "%d file(s) contain %r but are not selected,
because a string prefix "
+ "matches only at the start of the filename: %s%s",
+ len(dropped),
+ s3_prefix,
+ dropped[:SKIPPED_SAMPLE_SIZE],
+ f" and {omitted} more" if omitted > 0 else "",
+ )
for file in files:
self.log.info("Moving file %s", file)
if self.sftp_filenames and isinstance(self.sftp_filenames,
str):
- sftp_filename = file.replace(self.s3_filenames,
self.sftp_filenames)
+ name = posixpath.basename(file)
+ sftp_filename = posixpath.join(
+ posixpath.dirname(file), name.replace(s3_prefix,
self.sftp_filenames, 1)
+ )
else:
sftp_filename = file
self._download_from_s3(
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/transfers/sftp_to_s3.py
b/providers/amazon/src/airflow/providers/amazon/aws/transfers/sftp_to_s3.py
index fbbac10114a..aef21c91068 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/transfers/sftp_to_s3.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/transfers/sftp_to_s3.py
@@ -34,6 +34,9 @@ if TYPE_CHECKING:
from airflow.sdk import Context
+SKIPPED_SAMPLE_SIZE = 10
+
+
class SFTPToS3Operator(BaseOperator):
"""
Transfer files from an SFTP server to Amazon S3.
@@ -183,16 +186,27 @@ class SFTPToS3Operator(BaseOperator):
if isinstance(self.sftp_filenames, str):
self.log.info("Getting files in %s", self.sftp_path)
list_dir = sftp_client.listdir(self.sftp_path)
- if self.sftp_filenames == "*":
+ sftp_prefix: str = self.sftp_filenames
+ if sftp_prefix == "*":
files = list_dir
else:
- sftp_prefix: str = self.sftp_filenames
- files = [f for f in list_dir if sftp_prefix in f]
+ files = [f for f in list_dir if f.startswith(sftp_prefix)]
+ dropped = [f for f in list_dir if sftp_prefix in f and not
f.startswith(sftp_prefix)]
+ if dropped:
+ omitted = len(dropped) - SKIPPED_SAMPLE_SIZE
+ self.log.warning(
+ "%d file(s) contain %r but are not selected,
because a string prefix "
+ "matches only at the start of the filename: %s%s",
+ len(dropped),
+ sftp_prefix,
+ dropped[:SKIPPED_SAMPLE_SIZE],
+ f" and {omitted} more" if omitted > 0 else "",
+ )
for file in files:
self.log.info("Moving file %s", file)
if self.s3_filenames and isinstance(self.s3_filenames,
str):
- s3_filename = file.replace(self.sftp_filenames,
self.s3_filenames)
+ s3_filename = file.replace(sftp_prefix,
self.s3_filenames, 1)
else:
s3_filename = file
self._upload_to_s3(
diff --git a/providers/amazon/tests/unit/amazon/aws/transfers/test_ftp_to_s3.py
b/providers/amazon/tests/unit/amazon/aws/transfers/test_ftp_to_s3.py
index 102969b33c8..83897d52ef1 100644
--- a/providers/amazon/tests/unit/amazon/aws/transfers/test_ftp_to_s3.py
+++ b/providers/amazon/tests/unit/amazon/aws/transfers/test_ftp_to_s3.py
@@ -116,22 +116,121 @@ class TestFTPToS3Operator:
s3_file=operator.s3_key + operator.ftp_filenames[0],
)
- @mock.patch("airflow.providers.ftp.hooks.ftp.FTPHook.list_directory")
+ @mock.patch.object(FTPToS3Operator,
"_FTPToS3Operator__upload_to_s3_from_ftp")
+ @mock.patch(
+ "airflow.providers.ftp.hooks.ftp.FTPHook.list_directory",
+ return_value=["pre_one.txt", "xpre_two.txt", "pre_again_pre_.txt"],
+ )
def test_execute_multiple_files_prefix(
self,
mock_ftp_hook_list_directory,
+ mock_upload_to_s3,
):
operator = FTPToS3Operator(
task_id=TASK_ID,
s3_bucket=BUCKET,
s3_key=S3_KEY_MULTIPLE,
ftp_path=FTP_PATH_MULTIPLE,
- ftp_filenames="test_prefix",
- s3_filenames="s3_prefix",
+ ftp_filenames="pre_",
+ s3_filenames="new_",
+ )
+ with mock.patch.object(operator.log, "warning") as mock_log_warning:
+ operator.execute(None)
+
+
mock_ftp_hook_list_directory.assert_called_once_with(path=FTP_PATH_MULTIPLE)
+ assert mock_upload_to_s3.call_args_list == [
+ mock.call("pre_one.txt", "test/new_one.txt"),
+ mock.call("pre_again_pre_.txt", "test/new_again_pre_.txt"),
+ ]
+ mock_log_warning.assert_called_once_with(mock.ANY, 1, "pre_",
["xpre_two.txt"], "")
+
+ @pytest.mark.parametrize(
+ ("ftp_filenames", "s3_filenames", "expected"),
+ [
+ pytest.param(
+ "pre_",
+ "new_",
+ [
+ ("/tmp/pre_one.txt", "test/new_one.txt"),
+ ("/tmp/pre_again_pre_.txt", "test/new_again_pre_.txt"),
+ ],
+ id="prefix-with-rename",
+ ),
+ pytest.param(
+ "pre_",
+ None,
+ [
+ ("/tmp/pre_one.txt", "test/pre_one.txt"),
+ ("/tmp/pre_again_pre_.txt", "test/pre_again_pre_.txt"),
+ ],
+ id="prefix-without-rename",
+ ),
+ pytest.param(
+ "*",
+ None,
+ [
+ ("/tmp/pre_one.txt", "test/pre_one.txt"),
+ ("/tmp/xpre_two.txt", "test/xpre_two.txt"),
+ ("/tmp/pre_again_pre_.txt", "test/pre_again_pre_.txt"),
+ ],
+ id="wildcard-without-rename",
+ ),
+ ],
+ )
+ @mock.patch.object(FTPToS3Operator,
"_FTPToS3Operator__upload_to_s3_from_ftp")
+ @mock.patch(
+ "airflow.providers.ftp.hooks.ftp.FTPHook.list_directory",
+ return_value=["/tmp/pre_one.txt", "/tmp/xpre_two.txt",
"/tmp/pre_again_pre_.txt"],
+ )
+ def test_execute_builds_s3_key_from_basename_when_nlst_returns_full_paths(
+ self,
+ mock_ftp_hook_list_directory,
+ mock_upload_to_s3,
+ ftp_filenames,
+ s3_filenames,
+ expected,
+ ):
+ """Some FTP servers qualify nlst entries with the listed directory.
The entry is kept for
+ retrieval, but the S3 key comes from the basename so it never embeds
the source directory."""
+ operator = FTPToS3Operator(
+ task_id=TASK_ID,
+ s3_bucket=BUCKET,
+ s3_key=S3_KEY_MULTIPLE,
+ ftp_path=FTP_PATH_MULTIPLE,
+ ftp_filenames=ftp_filenames,
+ s3_filenames=s3_filenames,
)
operator.execute(None)
mock_ftp_hook_list_directory.assert_called_once_with(path=FTP_PATH_MULTIPLE)
+ assert mock_upload_to_s3.call_args_list == [mock.call(*call) for call
in expected]
+
+ @mock.patch.object(FTPToS3Operator,
"_FTPToS3Operator__upload_to_s3_from_ftp")
+ @mock.patch(
+ "airflow.providers.ftp.hooks.ftp.FTPHook.list_directory",
+ return_value=["/srv/data/report.csv", "/srv/data/summary.csv"],
+ )
+ def
test_execute_warns_when_only_the_directory_matched_the_old_substring_rule(
+ self,
+ mock_ftp_hook_list_directory,
+ mock_upload_to_s3,
+ ):
+ """The old rule matched the whole nlst entry, so a prefix hitting only
a directory
+ component used to select every file. Those must be reported, not
dropped in silence."""
+ operator = FTPToS3Operator(
+ task_id=TASK_ID,
+ s3_bucket=BUCKET,
+ s3_key=S3_KEY_MULTIPLE,
+ ftp_path=FTP_PATH_MULTIPLE,
+ ftp_filenames="data",
+ )
+ with mock.patch.object(operator.log, "warning") as mock_log_warning:
+ operator.execute(None)
+
+ mock_upload_to_s3.assert_not_called()
+ mock_log_warning.assert_called_once_with(
+ mock.ANY, 2, "data", ["/srv/data/report.csv",
"/srv/data/summary.csv"], ""
+ )
class TestFTPToS3OperatorInit:
diff --git a/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_ftp.py
b/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_ftp.py
index ca7e084b931..10b91915f5e 100644
--- a/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_ftp.py
+++ b/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_ftp.py
@@ -49,6 +49,57 @@ class TestS3ToFTPOperator:
mock_s3_hook_get_key.return_value.download_fileobj.assert_called_once_with(mock_local_tmp_file_value)
mock_ftp_hook_store_file.assert_called_once_with(operator.ftp_path,
mock_local_tmp_file_value.name)
+ @pytest.mark.parametrize(
+ ("keys", "expected", "skipped"),
+ [
+ pytest.param(
+ ["source/pre_one.txt", "source/xpre_two.txt",
"source/pre_again_pre_.txt"],
+ [
+ ("source/pre_one.txt", "/destination/new_one.txt"),
+ ("source/pre_again_pre_.txt",
"/destination/new_again_pre_.txt"),
+ ],
+ ["xpre_two.txt"],
+ id="prefix-only-at-start",
+ ),
+ pytest.param(
+ ["source/pre_one.txt", "source/pre_dir/pre_two.txt",
"source/pre_dir/other.txt"],
+ [
+ ("source/pre_one.txt", "/destination/new_one.txt"),
+ ("source/pre_dir/pre_two.txt",
"/destination/pre_dir/new_two.txt"),
+ ],
+ ["pre_dir/other.txt"],
+ id="nested-keys-keep-their-directory",
+ ),
+ ],
+ )
+ @mock.patch.object(S3ToFTPOperator, "_download_from_s3")
+ @mock.patch("airflow.providers.amazon.aws.transfers.s3_to_ftp.FTPHook")
+ @mock.patch("airflow.providers.amazon.aws.transfers.s3_to_ftp.S3Hook")
+ def test_execute_matches_the_prefix_on_the_file_name(
+ self, mock_s3_hook_class, mock_ftp_hook_class, mock_download_from_s3,
keys, expected, skipped
+ ):
+ """The prefix matches only the start of the file name, and only that
leading occurrence is
+ renamed. ``list_keys`` recurses, so a nested key is matched on its
file name and keeps its
+ directory, while the old substring rule still decides which skipped
keys are reported."""
+ mock_s3_hook = mock_s3_hook_class.return_value
+ mock_s3_hook.list_keys.return_value = keys
+ operator = S3ToFTPOperator(
+ task_id=TASK_ID,
+ s3_bucket=BUCKET,
+ s3_key="source/",
+ ftp_path="/destination/",
+ s3_filenames="pre_",
+ ftp_filenames="new_",
+ )
+
+ with mock.patch.object(operator.log, "warning") as mock_log_warning:
+ operator.execute(None)
+
+ assert mock_download_from_s3.call_args_list == [
+ mock.call(mock_s3_hook, mock_ftp_hook_class.return_value, *call)
for call in expected
+ ]
+ mock_log_warning.assert_called_once_with(mock.ANY, len(skipped),
"pre_", skipped, "")
+
class TestS3ToFTPOperatorInit:
"""Unit tests for S3ToFTPOperator.__init__ that do not require an FTP
server."""
diff --git
a/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_sftp.py
b/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_sftp.py
index 9d1ef5461d8..3ed273f7d96 100644
--- a/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_sftp.py
+++ b/providers/amazon/tests/unit/amazon/aws/transfers/test_s3_to_sftp.py
@@ -17,6 +17,8 @@
# under the License.
from __future__ import annotations
+from unittest import mock
+
import boto3
import pytest
from moto import mock_aws
@@ -319,6 +321,59 @@ class TestS3ToSFTPOperator:
class TestS3ToSFTPOperatorInit:
"""Unit tests for S3ToSFTPOperator.__init__ that do not require an SSH
server."""
+ @pytest.mark.parametrize(
+ ("keys", "expected", "skipped"),
+ [
+ pytest.param(
+ ["source/pre_one.txt", "source/xpre_two.txt",
"source/pre_again_pre_.txt"],
+ [
+ ("source/pre_one.txt", "/destination/new_one.txt"),
+ ("source/pre_again_pre_.txt",
"/destination/new_again_pre_.txt"),
+ ],
+ ["xpre_two.txt"],
+ id="prefix-only-at-start",
+ ),
+ pytest.param(
+ ["source/pre_one.txt", "source/pre_dir/pre_two.txt",
"source/pre_dir/other.txt"],
+ [
+ ("source/pre_one.txt", "/destination/new_one.txt"),
+ ("source/pre_dir/pre_two.txt",
"/destination/pre_dir/new_two.txt"),
+ ],
+ ["pre_dir/other.txt"],
+ id="nested-keys-keep-their-directory",
+ ),
+ ],
+ )
+ @mock.patch.object(S3ToSFTPOperator, "_download_from_s3")
+ @mock.patch("airflow.providers.amazon.aws.transfers.s3_to_sftp.SSHHook")
+ @mock.patch("airflow.providers.amazon.aws.transfers.s3_to_sftp.S3Hook")
+ def test_execute_matches_the_prefix_on_the_file_name(
+ self, mock_s3_hook_class, mock_ssh_hook_class, mock_download_from_s3,
keys, expected, skipped
+ ):
+ """The prefix matches only the start of the file name, and only that
leading occurrence is
+ renamed. ``list_keys`` recurses, so a nested key is matched on its
file name and keeps its
+ directory, while the old substring rule still decides which skipped
keys are reported."""
+ mock_s3_hook = mock_s3_hook_class.return_value
+ mock_s3_hook.list_keys.return_value = keys
+ sftp_client =
mock_ssh_hook_class.return_value.get_conn.return_value.open_sftp.return_value
+ operator = S3ToSFTPOperator(
+ task_id=TASK_ID,
+ s3_bucket=BUCKET,
+ s3_key="source/",
+ sftp_path="/destination/",
+ sftp_conn_id=SFTP_CONN_ID,
+ s3_filenames="pre_",
+ sftp_filenames="new_",
+ )
+
+ with mock.patch.object(operator.log, "warning") as mock_log_warning:
+ operator.execute(None)
+
+ assert mock_download_from_s3.call_args_list == [
+ mock.call(sftp_client, mock_s3_hook, *call) for call in expected
+ ]
+ mock_log_warning.assert_called_once_with(mock.ANY, len(skipped),
"pre_", skipped, "")
+
@pytest.mark.parametrize(
("s3_filenames", "sftp_filenames"),
[
diff --git
a/providers/amazon/tests/unit/amazon/aws/transfers/test_sftp_to_s3.py
b/providers/amazon/tests/unit/amazon/aws/transfers/test_sftp_to_s3.py
index 9df59839f6b..c715c67fb99 100644
--- a/providers/amazon/tests/unit/amazon/aws/transfers/test_sftp_to_s3.py
+++ b/providers/amazon/tests/unit/amazon/aws/transfers/test_sftp_to_s3.py
@@ -18,6 +18,7 @@
from __future__ import annotations
import warnings
+from unittest import mock
import boto3
import pytest
@@ -216,6 +217,39 @@ class TestSFTPToS3Operator:
class TestSFTPToS3OperatorInit:
"""Unit tests for SFTPToS3Operator.__init__ that do not require an SSH
server."""
+ @mock.patch.object(SFTPToS3Operator, "_upload_to_s3")
+ @mock.patch("airflow.providers.amazon.aws.transfers.sftp_to_s3.SSHHook")
+ @mock.patch("airflow.providers.amazon.aws.transfers.sftp_to_s3.S3Hook")
+ def test_execute_prefix_matches_and_replaces_only_leading_prefix(
+ self, mock_s3_hook_class, mock_ssh_hook_class, mock_upload_to_s3
+ ):
+ mock_s3_hook = mock_s3_hook_class.return_value
+ sftp_client =
mock_ssh_hook_class.return_value.get_conn.return_value.open_sftp.return_value
+ sftp_client.listdir.return_value = ["pre_one.txt", "xpre_two.txt",
"pre_again_pre_.txt"]
+ operator = SFTPToS3Operator(
+ task_id="test_prefix",
+ s3_bucket=BUCKET,
+ s3_key="destination/",
+ sftp_path="/source",
+ sftp_conn_id=SFTP_CONN_ID,
+ sftp_filenames="pre_",
+ s3_filenames="new_",
+ )
+
+ with mock.patch.object(operator.log, "warning") as mock_log_warning:
+ operator.execute(None)
+
+ assert mock_upload_to_s3.call_args_list == [
+ mock.call(sftp_client, mock_s3_hook, "/source/pre_one.txt",
"destination/new_one.txt"),
+ mock.call(
+ sftp_client,
+ mock_s3_hook,
+ "/source/pre_again_pre_.txt",
+ "destination/new_again_pre_.txt",
+ ),
+ ]
+ mock_log_warning.assert_called_once_with(mock.ANY, 1, "pre_",
["xpre_two.txt"], "")
+
def test_s3_conn_id_deprecated(self):
"""s3_conn_id is a deprecated alias for aws_conn_id and must raise
DeprecationWarning."""
with warnings.catch_warnings(record=True) as caught: