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

potiuk 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 e56fcdffac6 Skip GCS folder-marker keys in GCSToS3Operator (#72497)
e56fcdffac6 is described below

commit e56fcdffac603dc805c0f44d563632273dc242c9
Author: Yuseok Jo <[email protected]>
AuthorDate: Wed Sep 9 08:24:49 2026 +0900

    Skip GCS folder-marker keys in GCSToS3Operator (#72497)
---
 .../providers/amazon/aws/transfers/gcs_to_s3.py    | 35 +++++++++++
 .../unit/amazon/aws/transfers/test_gcs_to_s3.py    | 67 ++++++++++++++++++++++
 2 files changed, 102 insertions(+)

diff --git 
a/providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py 
b/providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py
index 867f2d3d8cb..0eeb0f35b5d 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py
@@ -164,6 +164,29 @@ class GCSToS3Operator(BaseOperator):
             return os.path.basename(file_path)
         return file_path
 
+    @staticmethod
+    def _strip_overlapping_folder_markers(keys: list[str]) -> tuple[list[str], 
list[str]]:
+        """
+        Drop trailing-slash keys that are strict prefixes of other listed keys.
+
+        Treated as directory markers. A lone trailing-slash key with no overlap
+        (e.g. ``lonely/``) is preserved, and a non-slash key that happens to 
be a
+        strict prefix of another (e.g. ``abc`` of ``abcdef``) is also 
preserved.
+        Returns ``(kept, dropped)``.
+        """
+        if not keys:
+            return [], []
+        ordered = sorted(set(keys))
+        kept: list[str] = []
+        dropped: list[str] = []
+        for current, nxt in zip(ordered, ordered[1:]):
+            if current.endswith("/") and nxt.startswith(current):
+                dropped.append(current)
+            else:
+                kept.append(current)
+        kept.append(ordered[-1])
+        return kept, dropped
+
     def execute(self, context: Context) -> list[str]:
         # list all files in an Google Cloud Storage bucket
         gcs_hook = GCSHook(
@@ -187,6 +210,18 @@ class GCSToS3Operator(BaseOperator):
 
         gcs_files = gcs_hook.list(**list_kwargs)  # type: ignore
 
+        gcs_files, dropped_keys = 
self._strip_overlapping_folder_markers(gcs_files)
+        if self.flatten_structure:
+            # A kept marker like lonely/ has no basename, so flattening would 
hit the destination prefix.
+            dropped_keys += [file for file in gcs_files if not 
self._transform_file_path(file)]
+            gcs_files = [file for file in gcs_files if 
self._transform_file_path(file)]
+        if dropped_keys:
+            self.log.info(
+                "Skipping %s GCS folder-marker key(s) (omitted from transfer 
and XCom output): %s",
+                len(dropped_keys),
+                dropped_keys,
+            )
+
         s3_hook = S3Hook(
             aws_conn_id=self.dest_aws_conn_id, verify=self.dest_verify, 
extra_args=self.dest_s3_extra_args
         )
diff --git a/providers/amazon/tests/unit/amazon/aws/transfers/test_gcs_to_s3.py 
b/providers/amazon/tests/unit/amazon/aws/transfers/test_gcs_to_s3.py
index 014f26f4943..108eb0debff 100644
--- a/providers/amazon/tests/unit/amazon/aws/transfers/test_gcs_to_s3.py
+++ b/providers/amazon/tests/unit/amazon/aws/transfers/test_gcs_to_s3.py
@@ -187,6 +187,73 @@ class TestGCSToS3Operator:
             assert sorted(MOCK_FILES) == sorted(uploaded_files)
             assert sorted(MOCK_FILES) == sorted(hook.list_keys("bucket", 
delimiter="/"))
 
+    @pytest.mark.parametrize(
+        ("keys", "expected_kept", "expected_dropped"),
+        [
+            ([], [], []),
+            (["a"], ["a"], []),
+            (["a", "b"], ["a", "b"], []),
+            # Non-slash prefix overlaps must NOT be treated as folder markers.
+            (["a", "ax"], ["a", "ax"], []),
+            (["abc", "abcdef"], ["abc", "abcdef"], []),
+            (["foo/", "foo/bar.txt"], ["foo/bar.txt"], ["foo/"]),
+            (
+                ["data/", "data/sub/", "data/sub/file.txt"],
+                ["data/sub/file.txt"],
+                ["data/", "data/sub/"],
+            ),
+            # A lone trailing-slash key with no overlap is a real object and 
stays.
+            (["lonely/"], ["lonely/"], []),
+            (["lonely/", "report.csv"], ["lonely/", "report.csv"], []),
+        ],
+    )
+    def test_strip_overlapping_folder_markers(self, keys, expected_kept, 
expected_dropped):
+        """Folder-marker detection: requires both strict-prefix overlap AND 
trailing slash."""
+        kept, dropped = GCSToS3Operator._strip_overlapping_folder_markers(keys)
+        assert kept == expected_kept
+        assert dropped == expected_dropped
+
+    @mock.patch("airflow.providers.amazon.aws.transfers.gcs_to_s3.GCSHook")
+    def test_execute_skips_overlapping_folder_markers(self, mock_hook):
+        mock_hook.return_value.list.return_value = ["src/", "src/airflow.png", 
"lonely/"]
+        with NamedTemporaryFile() as f:
+            gcs_provide_file = mock_hook.return_value.provide_file
+            gcs_provide_file.return_value.__enter__.return_value.name = f.name
+
+            operator = GCSToS3Operator(
+                task_id=TASK_ID,
+                gcs_bucket=GCS_BUCKET,
+                dest_aws_conn_id="aws_default",
+                dest_s3_key=S3_BUCKET,
+                replace=True,
+            )
+            hook, _ = _create_test_bucket()
+
+            uploaded_files = operator.execute(None)
+            assert uploaded_files == ["lonely/", "src/airflow.png"]
+            assert hook.list_keys("bucket") == ["lonely/", "src/airflow.png"]
+
+    @mock.patch("airflow.providers.amazon.aws.transfers.gcs_to_s3.GCSHook")
+    def test_execute_skips_markers_without_a_basename_when_flattening(self, 
mock_hook):
+        mock_hook.return_value.list.return_value = ["src/", "src/airflow.png", 
"lonely/"]
+        with NamedTemporaryFile() as f:
+            gcs_provide_file = mock_hook.return_value.provide_file
+            gcs_provide_file.return_value.__enter__.return_value.name = f.name
+
+            operator = GCSToS3Operator(
+                task_id=TASK_ID,
+                gcs_bucket=GCS_BUCKET,
+                dest_aws_conn_id="aws_default",
+                dest_s3_key="s3://bucket/dest/",
+                flatten_structure=True,
+                replace=True,
+            )
+            hook, _ = _create_test_bucket()
+
+            uploaded_files = operator.execute(None)
+            assert uploaded_files == ["src/airflow.png"]
+            assert hook.list_keys("bucket") == ["dest/airflow.png"]
+
     @mock.patch("airflow.providers.amazon.aws.transfers.gcs_to_s3.GCSHook")
     def test_execute_with_replace(self, mock_hook):
         mock_hook.return_value.list.return_value = MOCK_FILES

Reply via email to