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"