This is an automated email from the ASF dual-hosted git repository.
pabloem pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new f0734081aca Feat: new cleaning rule to orphaned subscriptions (#39538)
f0734081aca is described below
commit f0734081aca456e6a3f4cd76e26a1e275b9afa0a
Author: HansMarcus01 <[email protected]>
AuthorDate: Mon Aug 3 10:46:54 2026 -0600
Feat: new cleaning rule to orphaned subscriptions (#39538)
* Feat: A new prefix was added to the cleaner to clean up newly discovered
orphaned subscriptions of taxirides topic
* Fix: Modifications were made to manage the cleanup of disconnected and
active subscriptions under prefixes.
* Fix: Deleting the hardcode about taxi prefix
* Fix: Delete comment in spanish
* Fix: resolving the issue where active subscriptions are not deleted
---
.test-infra/tools/stale_cleaner.py | 15 +++++---
.test-infra/tools/test_stale_cleaner.py | 61 ++++++++++++++++++++++++---------
2 files changed, 54 insertions(+), 22 deletions(-)
diff --git a/.test-infra/tools/stale_cleaner.py
b/.test-infra/tools/stale_cleaner.py
index 2c702cfa0fd..a00bb3ad7ab 100644
--- a/.test-infra/tools/stale_cleaner.py
+++ b/.test-infra/tools/stale_cleaner.py
@@ -336,10 +336,13 @@ class PubSubSubscriptionCleaner(StaleCleaner):
for subscription in
self.client.list_subscriptions(request={"project": self.project_path}):
subscription_name = subscription.name
# Apply prefix filtering if prefixes are defined
- if not self.prefixes or
any(subscription_name.startswith(f"{self.project_path}/subscriptions/{prefix}")
for prefix in self.prefixes):
- # Check if the subscription has a topic associated with it
- if subscription.detached:
+ if subscription.detached:
d[subscription_name] =
GoogleCloudResource(resource_name=subscription_name, clock=self.clock)
+ #Only attached subscriptions with the NYC taxi prefix are
eligible.
+ elif any(
+
subscription_name.startswith(f"{self.project_path}/subscriptions/{prefix}") for
prefix in self.prefixes
+ ):
+ d[subscription_name] =
GoogleCloudResource(resource_name=subscription_name, clock=self.clock)
return d
@@ -416,8 +419,10 @@ def clean_pubsub_subscriptions():
project_id = DEFAULT_PROJECT_ID
bucket_name = DEFAULT_BUCKET_NAME
- # No prefixes are defined for subscriptions so we will delete all stale
subscriptions
- prefixes = []
+ # Restrict subscription cleanup to the NYC taxi prefix only.
+ prefixes = [
+ "taxirides-realtime_beam_",
+ ]
# Create a PubSubSubscriptionCleaner instance
cleaner = PubSubSubscriptionCleaner(project_id=project_id,
bucket_name=bucket_name,
diff --git a/.test-infra/tools/test_stale_cleaner.py
b/.test-infra/tools/test_stale_cleaner.py
index c53fbc1a44d..08cdab39b85 100644
--- a/.test-infra/tools/test_stale_cleaner.py
+++ b/.test-infra/tools/test_stale_cleaner.py
@@ -431,30 +431,57 @@ class PubSubSubscriptionCleanerTest(unittest.TestCase):
self.assertEqual(self.cleaner.time_threshold, self.time_threshold)
self.assertIsInstance(self.cleaner.clock, FakeClock)
- def test_active_resources(self):
- """Test _active_resources method."""
- # Mock subscriptions
- sub1 = mock.Mock()
- sub1.name = "projects/test-project/subscriptions/test-prefix-sub1"
- sub1.topic = "projects/test-project/topics/some-topic"
+ def test_active_resources_active_subscriptions(self):
+ """Verify that active subscriptions with the 'taxirides' prefix are
identified."""
+ self.cleaner.prefixes = ["taxirides-realtime_beam_"]
- sub2 = mock.Mock()
- sub2.name =
"projects/test-project/subscriptions/test-prefix-sub2-detached"
- sub2.topic = "_deleted-topic_"
+ # Active suscription with the correct taxi prefix
+ sub_taxi_active = mock.Mock()
+ sub_taxi_active.name =
f"projects/{self.project_id}/subscriptions/taxirides-realtime_beam_-12345"
+ sub_taxi_active.topic =
"projects/pubsub-public-data/topics/taxirides-realtime"
+ sub_taxi_active.detached = False
- sub3 = mock.Mock()
- sub3.name = "projects/test-project/subscriptions/other-prefix-sub3"
- sub3.topic = "projects/test-project/topics/another-topic"
+ # Active subscription with a different prefix
+ sub_other_active = mock.Mock()
+ sub_other_active.name =
f"projects/{self.project_id}/subscriptions/other-prefix-sub"
+ sub_other_active.topic =
f"projects/{self.project_id}/topics/another-topic"
+ sub_other_active.detached = False
- self.mock_subscriber_client.list_subscriptions.return_value = [sub1,
sub2, sub3]
+ self.mock_subscriber_client.list_subscriptions.return_value =
[sub_taxi_active, sub_other_active]
with SilencePrint():
active = self.cleaner._active_resources()
- self.assertIn("projects/test-project/subscriptions/test-prefix-sub1",
active)
-
self.assertIn("projects/test-project/subscriptions/test-prefix-sub2-detached",
active)
-
self.assertNotIn("projects/test-project/subscriptions/other-prefix-sub3",
active)
- self.assertEqual(len(active), 2)
+ # Verify that only the taxi subscription is captured, discarding the
other one
+ self.assertIn(sub_taxi_active.name, active)
+ self.assertNotIn(sub_other_active.name, active)
+ self.assertEqual(len(active), 1)
+
+ def test_active_resources_detached_subscriptions(self):
+ """Verify that detached subscriptions with the 'test-prefix'
prefix are identified."""
+ self.cleaner.prefixes = ["test-prefix"]
+
+ # Standar suscription with a detached topic (should be included)
+ sub_detached = mock.Mock()
+ sub_detached.name =
f"projects/{self.project_id}/subscriptions/test-prefix-detached"
+ sub_detached.topic = "_deleted-topic_"
+ sub_detached.detached = True
+
+ # Standard connected subscription (should ignore)
+ sub_attached = mock.Mock()
+ sub_attached.name =
f"projects/{self.project_id}/subscriptions/other-prefix-attached"
+ sub_attached.topic =
f"projects/{self.project_id}/topics/some-topic"
+ sub_attached.detached = False
+
+ self.mock_subscriber_client.list_subscriptions.return_value =
[sub_detached, sub_attached]
+
+ with SilencePrint():
+ active = self.cleaner._active_resources()
+
+ # Only the detached subscription should be included in the active
resources
+ self.assertIn(sub_detached.name, active)
+ self.assertNotIn(sub_attached.name, active)
+ self.assertEqual(len(active), 1)
def test_delete_resource(self):
"""Test _delete_resource method."""