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 e15b3a66137 Retry S3Hook.delete_bucket when a late write leaves bucket
non-empty (#73930)
e15b3a66137 is described below
commit e15b3a6613734ca8a6639f52ba687d613d7025c6
Author: Sean Ghaeli <[email protected]>
AuthorDate: Wed Sep 30 05:45:23 2026 -0700
Retry S3Hook.delete_bucket when a late write leaves bucket non-empty
(#73930)
With force_delete, the retries only covered emptying the bucket. An object
that lands after the final empty listing (for example a service still
uploading logs) made DeleteBucket fail with BucketNotEmpty and nothing
retried it.
---
.../src/airflow/providers/amazon/aws/hooks/s3.py | 8 +++++++-
.../amazon/tests/unit/amazon/aws/hooks/test_s3.py | 19 +++++++++++++++++++
2 files changed, 26 insertions(+), 1 deletion(-)
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py
b/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py
index e182afef051..bba4ac7399f 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/hooks/s3.py
@@ -1510,7 +1510,13 @@ class S3Hook(AwsBaseHook):
for retry in range(max_retries):
bucket_keys = self.list_keys(bucket_name=bucket_name)
if not bucket_keys:
- break
+ try:
+ self.conn.delete_bucket(Bucket=bucket_name)
+ return
+ except ClientError as e:
+ if e.response["Error"]["Code"] != "BucketNotEmpty":
+ raise
+ continue
if retry: # Avoid first loop
time.sleep(500)
diff --git a/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py
b/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py
index f7231670e29..5f4da0a8dcc 100644
--- a/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py
+++ b/providers/amazon/tests/unit/amazon/aws/hooks/test_s3.py
@@ -1491,6 +1491,25 @@ class TestAwsS3Hook:
assert
mock_hook.delete_bucket(bucket_name="not-exists-bucket-name", force_delete=True)
assert ctx.value.response["Error"]["Code"] == "NoSuchBucket"
+ @mock_aws
+ @mock.patch("airflow.providers.amazon.aws.hooks.s3.time.sleep")
+ def test_delete_bucket_force_delete_retries_after_late_write(self,
mock_sleep, s3_bucket):
+ hook = S3Hook()
+ hook.load_string("data", key="key", bucket_name=s3_bucket)
+ late_writes = []
+
+ def put_late_object(**kwargs):
+ if not late_writes:
+ late_writes.append("late_key")
+ hook.conn.put_object(Bucket=s3_bucket, Key="late_key",
Body=b"late")
+
+ hook.conn.meta.events.register("before-call.s3.DeleteBucket",
put_late_object)
+ hook.delete_bucket(bucket_name=s3_bucket, force_delete=True)
+
+ assert late_writes == ["late_key"]
+ assert not hook.check_for_bucket(s3_bucket)
+ mock_sleep.assert_called_once_with(500)
+
def test_provide_bucket_name(self):
with mock.patch.object(
S3Hook,