kennknowles commented on code in PR #40386:
URL: https://github.com/apache/beam/pull/40386#discussion_r4210464011


##########
.test-infra/tools/build.gradle:
##########
@@ -40,6 +40,15 @@ task removeStaleSpannerResources(type: Exec) {
   commandLine './stale_spanner_cleaner.sh'
 }
 
+task testStaleGcsBucketsCleaner(type: Exec) {
+  commandLine 'bash', '-c', 'python3 -m pip install -r requirements.txt && 
python3 -m unittest test_stale_gcs_buckets_cleaner.py'

Review Comment:
   Nice addition to the pattern here :-)
   
   It would be cool to register this so that the standard `:test` target will 
run it. I'm not sure if that is compatible with `Exec` and it certainly is not 
blocking this review.



##########
.test-infra/tools/stale_gcs_buckets_cleaner.py:
##########
@@ -0,0 +1,112 @@
+#!/usr/bin/env python
+#
+#    Licensed to the Apache Software Foundation (ASF) under one
+#    or more contributor license agreements.  See the NOTICE file
+#    distributed with this work for additional information
+#    regarding copyright ownership.  The ASF licenses this file
+#    to you under the Apache License, Version 2.0 (the
+#    "License"); you may not use this file except in compliance
+#    with the License.  You may obtain a copy of the License at
+#
+#       http://www.apache.org/licenses/LICENSE-2.0
+#
+#    Unless required by applicable law or agreed to in writing,
+#    software distributed under the License is distributed on an
+#    "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+#    KIND, either express or implied.  See the License for the
+#    specific language governing permissions and limitations
+#    under the License.
+#
+"""Deletes stale GCS buckets left behind by GcsUtil integration tests."""
+
+import argparse
+import datetime
+import re
+
+from google.cloud import storage
+
+
+DEFAULT_PROJECT_ID = "apache-beam-testing"
+GCS_TEMP_BUCKET_PREFIX = "apache-beam-temp-bucket-"
+GCS_TEMP_BUCKET_PATTERN = re.compile(
+    rf"^{GCS_TEMP_BUCKET_PREFIX}"
+    r"[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-"
+    r"[0-9a-f]{12}$")
+DEFAULT_MAX_AGE = datetime.timedelta(hours=24)
+UTC = datetime.timezone.utc
+
+
+def _as_utc(value):
+    if value.tzinfo is None:
+        return value.replace(tzinfo=UTC)
+    return value.astimezone(UTC)
+
+
+def clean_stale_gcs_buckets(
+    client,
+    max_age=DEFAULT_MAX_AGE,
+    now=None,
+    dry_run=True,
+):
+    """Deletes test buckets whose creation time is older than ``max_age``.
+
+    Returns the names of buckets selected for deletion. If a deletion fails,
+    the cleaner attempts the remaining buckets before reporting the failures.
+    """
+    current_time = _as_utc(now or datetime.datetime.now(UTC))
+    selected = []
+    failures = []
+
+    buckets = client.list_buckets(
+        project=DEFAULT_PROJECT_ID, prefix=GCS_TEMP_BUCKET_PREFIX)
+    for bucket in buckets:
+        # Keep this check even though list_buckets filters server-side. It is
+        # the final guard before a destructive operation.
+        if not GCS_TEMP_BUCKET_PATTERN.fullmatch(bucket.name):
+            continue
+        if bucket.time_created is None:
+            print(f"Skipping {bucket.name}: creation time is unavailable")
+            continue
+        if current_time - _as_utc(bucket.time_created) <= max_age:
+            continue
+
+        selected.append(bucket.name)
+        if dry_run:
+            print(f"Dry run: would delete gs://{bucket.name}")
+            continue
+
+        print(f"Deleting stale test bucket gs://{bucket.name}")
+        try:
+            # Tests can leave objects behind when they fail before teardown.
+            bucket.delete(force=True)
+        except Exception as error:
+            failures.append(bucket.name)
+            print(f"Failed to delete gs://{bucket.name}: {error}")
+
+    if failures:
+        noun = "bucket" if len(failures) == 1 else "buckets"

Review Comment:
   Let's just have a constant string like `"Bucket failure count: 
{len(failures)}"` both to keep the code more concise and the output more 
obvious, plus you can then search for the error message and find it.



##########
.test-infra/tools/test_stale_gcs_buckets_cleaner.py:
##########
@@ -0,0 +1,152 @@
+#!/usr/bin/env python
+#
+#    Licensed to the Apache Software Foundation (ASF) under one
+#    or more contributor license agreements.  See the NOTICE file
+#    distributed with this work for additional information
+#    regarding copyright ownership.  The ASF licenses this file
+#    to you under the Apache License, Version 2.0 (the
+#    "License"); you may not use this file except in compliance
+#    with the License.  You may obtain a copy of the License at
+#
+#       http://www.apache.org/licenses/LICENSE-2.0
+#
+#    Unless required by applicable law or agreed to in writing,
+#    software distributed under the License is distributed on an
+#    "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+#    KIND, either express or implied.  See the License for the
+#    specific language governing permissions and limitations
+#    under the License.
+#
+
+import datetime
+import unittest
+from unittest import mock
+
+import stale_gcs_buckets_cleaner as cleaner
+
+
+UTC = datetime.timezone.utc
+NOW = datetime.datetime(2026, 9, 30, 12, tzinfo=UTC)
+STALE_BUCKET = "apache-beam-temp-bucket-123e4567-e89b-42d3-a456-426614174000"
+FRESH_BUCKET = "apache-beam-temp-bucket-123e4567-e89b-42d3-a456-426614174001"
+
+
+def bucket(name, created):
+    value = mock.Mock()

Review Comment:
   Instead of mocking, a fake is more robust.
   
   The official suggestion would be 
https://github.com/googleapis/storage-testbench since it will have high 
fidelity. But here, you could probably use a lightweight fake like so
   
   ```
   from unittest import mock
   import pytest
   
   class FakeGcsStorage:
       def __init__(self):
           self.buckets = {}
   
       def upload(self, bucket_name: str, blob_name: str, data: bytes):
           self.buckets.setdefault(bucket_name, {})[blob_name] = data
   
       def download(self, bucket_name: str, blob_name: str) -> bytes:
           return self.buckets[bucket_name][blob_name]
   
   def test_service_with_fake_storage():
       fake_storage = FakeGcsStorage()
       fake_storage.upload("my-bucket", "data.txt", b"my content")
       assert fake_storage.download("my-bucket", "data.txt") == b"my content"
   ```
   
   (just example code - fill in the methods you need with in-process 
trivial-but-real implementations)
   



##########
.test-infra/tools/stale_gcs_buckets_cleaner.py:
##########
@@ -0,0 +1,112 @@
+#!/usr/bin/env python
+#
+#    Licensed to the Apache Software Foundation (ASF) under one
+#    or more contributor license agreements.  See the NOTICE file
+#    distributed with this work for additional information
+#    regarding copyright ownership.  The ASF licenses this file
+#    to you under the Apache License, Version 2.0 (the
+#    "License"); you may not use this file except in compliance
+#    with the License.  You may obtain a copy of the License at
+#
+#       http://www.apache.org/licenses/LICENSE-2.0
+#
+#    Unless required by applicable law or agreed to in writing,
+#    software distributed under the License is distributed on an
+#    "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+#    KIND, either express or implied.  See the License for the
+#    specific language governing permissions and limitations
+#    under the License.
+#
+"""Deletes stale GCS buckets left behind by GcsUtil integration tests."""
+
+import argparse
+import datetime
+import re
+
+from google.cloud import storage
+
+
+DEFAULT_PROJECT_ID = "apache-beam-testing"
+GCS_TEMP_BUCKET_PREFIX = "apache-beam-temp-bucket-"
+GCS_TEMP_BUCKET_PATTERN = re.compile(
+    rf"^{GCS_TEMP_BUCKET_PREFIX}"
+    r"[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-"
+    r"[0-9a-f]{12}$")
+DEFAULT_MAX_AGE = datetime.timedelta(hours=24)
+UTC = datetime.timezone.utc
+
+
+def _as_utc(value):
+    if value.tzinfo is None:
+        return value.replace(tzinfo=UTC)
+    return value.astimezone(UTC)
+
+
+def clean_stale_gcs_buckets(
+    client,
+    max_age=DEFAULT_MAX_AGE,
+    now=None,
+    dry_run=True,
+):
+    """Deletes test buckets whose creation time is older than ``max_age``.
+
+    Returns the names of buckets selected for deletion. If a deletion fails,
+    the cleaner attempts the remaining buckets before reporting the failures.
+    """
+    current_time = _as_utc(now or datetime.datetime.now(UTC))
+    selected = []
+    failures = []
+
+    buckets = client.list_buckets(
+        project=DEFAULT_PROJECT_ID, prefix=GCS_TEMP_BUCKET_PREFIX)
+    for bucket in buckets:
+        # Keep this check even though list_buckets filters server-side. It is
+        # the final guard before a destructive operation.
+        if not GCS_TEMP_BUCKET_PATTERN.fullmatch(bucket.name):
+            continue
+        if bucket.time_created is None:
+            print(f"Skipping {bucket.name}: creation time is unavailable")
+            continue
+        if current_time - _as_utc(bucket.time_created) <= max_age:
+            continue
+
+        selected.append(bucket.name)
+        if dry_run:
+            print(f"Dry run: would delete gs://{bucket.name}")
+            continue
+
+        print(f"Deleting stale test bucket gs://{bucket.name}")
+        try:
+            # Tests can leave objects behind when they fail before teardown.
+            bucket.delete(force=True)

Review Comment:
   I don't know if the summary is accurate, but I checked on this and I 
understand it is necessary if the bucket has any contents, but will still fail 
if it has more than 256 objects, which seems pretty likely.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to