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 86aef85e35c Keep S3 Dag bundle downloads within the configured 
directory (#73756)
86aef85e35c is described below

commit 86aef85e35cdd458dd2388fe68b88e6c2366c0ce
Author: jingi723 <[email protected]>
AuthorDate: Tue Oct 6 02:10:47 2026 +0900

    Keep S3 Dag bundle downloads within the configured directory (#73756)
    
    * Keep S3 directory sync within the requested prefix
    
    S3 lists keys using string prefixes, so a directory prefix without a
    trailing slash also selects sibling files and directories. Those keys
    cannot be mapped below the requested local directory and prevent a Dag
    bundle from refreshing.
    
    * Clarify the S3 sync key-equals-prefix test case
    
    The exact-prefix case fails differently from sibling keys, so its name 
should make that boundary clear in CI failures.
    
    * Limit S3 Dag bundle downloads to the configured directory
    
    The bundle validates a directory prefix during initialization, but 
downloads with a raw string prefix. Matching objects outside that directory can 
prevent initialization and refresh.
---
 providers/amazon/docs/bundles/index.rst            |  5 ++++
 providers/amazon/docs/changelog.rst                |  6 ++++
 .../src/airflow/providers/amazon/aws/bundles/s3.py |  6 +++-
 .../tests/unit/amazon/aws/bundles/test_s3.py       | 33 ++++++++++++++++++++++
 4 files changed, 49 insertions(+), 1 deletion(-)

diff --git a/providers/amazon/docs/bundles/index.rst 
b/providers/amazon/docs/bundles/index.rst
index 0c1d3ea6302..f332b2f3b73 100644
--- a/providers/amazon/docs/bundles/index.rst
+++ b/providers/amazon/docs/bundles/index.rst
@@ -27,6 +27,11 @@ S3DagBundle
 Use the :class:`~airflow.providers.amazon.aws.bundles.s3.S3DagBundle` to 
configure an S3 bundle in your Airflow's
 ``[dag_processor] dag_bundle_config_list``.
 
+The ``prefix`` selects a subdirectory, with or without a trailing slash. For 
example,
+``dags`` and ``dags/`` both download objects under ``dags/``. An empty prefix 
downloads
+the whole bucket. IAM policies that restrict listing with the ``s3:prefix`` 
condition
+must allow the directory prefix including its trailing slash.
+
 Example of using the S3DagBundle:
 
 **JSON format example**:
diff --git a/providers/amazon/docs/changelog.rst 
b/providers/amazon/docs/changelog.rst
index 3c64bc19c85..1e8473af276 100644
--- a/providers/amazon/docs/changelog.rst
+++ b/providers/amazon/docs/changelog.rst
@@ -26,6 +26,12 @@
 Changelog
 ---------
 
+.. warning::
+  ``S3DagBundle`` now appends ``/`` to non-empty directory prefixes when 
downloading Dags.
+  For a configured prefix of ``dags``, listing requests now use ``dags/``. IAM 
policies
+  that restrict ``s3:prefix`` by exact value must permit the directory prefix 
with its
+  trailing slash. Empty prefixes and prefixes that already end in ``/`` are 
unchanged.
+
 9.37.0
 ......
 
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/bundles/s3.py 
b/providers/amazon/src/airflow/providers/amazon/aws/bundles/s3.py
index 65bacb4b388..c7f5eb77789 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/bundles/s3.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/bundles/s3.py
@@ -36,6 +36,7 @@ class S3DagBundle(BaseDagBundle):
     :param aws_conn_id: Airflow connection ID for AWS.  Defaults to 
AwsBaseHook.default_conn_name.
     :param bucket_name: The name of the S3 bucket containing the Dag files.
     :param prefix:  Optional subdirectory within the S3 bucket where the Dags 
are stored.
+                    A trailing slash is optional.
                     If None, Dags are assumed to be at the root of the bucket 
(Optional).
     """
 
@@ -129,9 +130,12 @@ class S3DagBundle(BaseDagBundle):
             self._log.debug(
                 "Downloading Dags from s3://%s/%s to %s", self.bucket_name, 
self.prefix, self.s3_dags_dir
             )
+            sync_prefix = self.prefix
+            if sync_prefix and not sync_prefix.endswith("/"):
+                sync_prefix += "/"
             self.s3_hook.sync_to_local_dir(
                 bucket_name=self.bucket_name,
-                s3_prefix=self.prefix,
+                s3_prefix=sync_prefix,
                 local_dir=self.s3_dags_dir,
                 delete_stale=True,
             )
diff --git a/providers/amazon/tests/unit/amazon/aws/bundles/test_s3.py 
b/providers/amazon/tests/unit/amazon/aws/bundles/test_s3.py
index f886ebe4367..9df5202fe71 100644
--- a/providers/amazon/tests/unit/amazon/aws/bundles/test_s3.py
+++ b/providers/amazon/tests/unit/amazon/aws/bundles/test_s3.py
@@ -221,3 +221,36 @@ class TestS3DagBundle:
         bundle.refresh()
         assert bundle._log.debug.call_count == 2
         assert bundle._log.debug.call_args_list == [download_log_call, 
download_log_call]
+
+    @pytest.mark.parametrize("prefix", ["dags", "dags/", "project/dags", 
"project/dags/"])
+    @pytest.mark.parametrize("extra_suffix", ["_archive/other.py", 
pytest.param("", id="key_equals_prefix")])
+    def test_refresh_uses_directory_prefix(self, s3_client, prefix, 
extra_suffix):
+        s3_client.create_bucket(Bucket=S3_BUCKET_NAME)
+        directory = prefix.rstrip("/")
+        old_key = f"{directory}/nested/old.py"
+        extra_key = directory + extra_suffix
+        s3_client.put_object(Bucket=S3_BUCKET_NAME, Key=old_key, Body=b"old")
+        s3_client.put_object(Bucket=S3_BUCKET_NAME, Key=extra_key, 
Body=b"outside")
+
+        bundle = S3DagBundle(name="test", bucket_name=S3_BUCKET_NAME, 
prefix=prefix)
+        original_url = bundle.view_url_template()
+        original_repr = repr(bundle)
+        bundle.initialize()
+        assert bundle.is_initialized
+        assert (bundle.path / "nested/old.py").read_bytes() == b"old"
+        assert {p.relative_to(bundle.path).as_posix() for p in 
bundle.path.rglob("*") if p.is_file()} == {
+            "nested/old.py"
+        }
+
+        s3_client.delete_object(Bucket=S3_BUCKET_NAME, Key=old_key)
+        s3_client.put_object(Bucket=S3_BUCKET_NAME, Key=f"{directory}/new.py", 
Body=b"new-content")
+        bundle.refresh()
+        assert not (bundle.path / "nested").exists()
+        assert (bundle.path / "new.py").read_bytes() == b"new-content"
+        assert {p.relative_to(bundle.path).as_posix() for p in 
bundle.path.rglob("*") if p.is_file()} == {
+            "new.py"
+        }
+        assert bundle.prefix == prefix
+        assert repr(bundle) == original_repr
+        assert bundle.view_url_template() == original_url
+        assert s3_client.get_object(Bucket=S3_BUCKET_NAME, 
Key=extra_key)["Body"].read() == b"outside"

Reply via email to