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 e4962bc90d9 [GSoC-273] Implement TestPubsubContext for Python GCP
Integration Tests (#39685)
e4962bc90d9 is described below
commit e4962bc90d974637532e20697ba0a10de2be9ca9
Author: HansMarcus01 <[email protected]>
AuthorDate: Wed Aug 19 10:45:53 2026 -0600
[GSoC-273] Implement TestPubsubContext for Python GCP Integration Tests
(#39685)
* Feat: Implementation of a cleanup handler to oversee the destruction of
Pub/Sub resources generated by Python tests.
* Feat: Enable test mode to avoid accidentally deleting resources.
* Feat: Incorporation of license and log messages
* Fix pytest collection error for Pub/Sub perf tests
Renamed test_pubsub.py to pubsub_test_context.py and added conditional
imports with pylint directives to prevent pytest from crashing on
environments missing GCP dependencies.
* Fix TestPubsubContext initialization and import path
- Added early None validation in TestPubsubContext to prevent AttributeError
when initializing clients in environments without GCP dependencies.
- Fixed the absolute import path in pubsub_io_perf_test.py to use the
standard apache_beam module root.
* Add unit tests for TestPubsubContext to improve patch coverage
Created `test_pubsub_context_unit.py` using unittest.mock to validate
the resource lifecycle logic without requiring actual GCP dependencies.
This covers initialization, resource tracking, successful teardown,
cascading deletions, and the 24-hour debug grace period, resolving
the Codecov patch coverage drop.
* Fix RAT and Formatter CI checks for Pub/Sub test context
Added the required ASF license header to the new unit test file and
ran YAPF in-place to enforce the project's 2-space indentation standard
across the modified Pub/Sub testing utilities.
* test: Rename Pub/Sub context manager and update load_tests/build.gradle
- Renamed the Python lifecycle manager file to `pubsub_test_context.py` for
naming consistency.
- Added the `runPubsubPerfTestWithMonitor` verification task to
`sdks/python/apache_beam/testing/load_tests/build.gradle`.
* Feat: Adding a stress test to the cleanup handler and validating resource
deletion.
Removing the redundant test from build.gradle, as the handler's execution
does not require this module.
* Feat: Correcting the YAPF formatting of the test file
pubsub_test_context_test.py.
Removed local resource cleanup from pubsub_io_perf.py so that the
pubsub_test_context.py handler takes care of it.
* Remove unused code
---
.../apache_beam/io/gcp/pubsub_io_perf_test.py | 13 +-
.../apache_beam/testing/pubsub_test_context.py | 179 +++++++++++++++++++++
.../testing/pubsub_test_context_test.py | 149 +++++++++++++++++
3 files changed, 337 insertions(+), 4 deletions(-)
diff --git a/sdks/python/apache_beam/io/gcp/pubsub_io_perf_test.py
b/sdks/python/apache_beam/io/gcp/pubsub_io_perf_test.py
index 7ca831c980e..e22ab61e852 100644
--- a/sdks/python/apache_beam/io/gcp/pubsub_io_perf_test.py
+++ b/sdks/python/apache_beam/io/gcp/pubsub_io_perf_test.py
@@ -60,6 +60,7 @@ from apache_beam.testing.synthetic_pipeline import
SyntheticSource
from apache_beam.testing.test_pipeline import TestPipeline
from apache_beam.transforms import trigger
from apache_beam.transforms import window
+from apache_beam.testing.pubsub_test_context import TestPubsubContext
# pylint: disable=wrong-import-order, wrong-import-position
try:
@@ -88,6 +89,8 @@ class PubsubIOPerfTest(LoadTest):
'pubsub_namespace_prefix')
self.pubsub_namespace = pubsub_namespace_prefix + unique_id
+ self.pubsub_monitor = TestPubsubContext(project_id=self.project_id)
+
def _setup_pubsub(self):
self.pub_client = pubsub.PublisherClient()
self.topic_name = self.pub_client.topic_path(
@@ -105,6 +108,10 @@ class PubsubIOPerfTest(LoadTest):
self.project_id,
self.pubsub_namespace + '_read_matcher',
)
+ self.pubsub_monitor.register_topic(self.topic_name)
+ self.pubsub_monitor.register_topic(self.matcher_topic_name)
+ self.pubsub_monitor.register_subscription(self.read_sub_name)
+ self.pubsub_monitor.register_subscription(self.read_matcher_sub_name)
class PubsubWritePerfTest(PubsubIOPerfTest):
@@ -205,10 +212,8 @@ class PubsubReadPerfTest(PubsubIOPerfTest):
self.pipeline = TestPipeline(options=PipelineOptions(args))
def cleanup(self):
- self.sub_client.delete_subscription(subscription=self.read_sub_name)
-
self.sub_client.delete_subscription(subscription=self.read_matcher_sub_name)
- self.pub_client.delete_topic(topic=self.topic_name)
- self.pub_client.delete_topic(topic=self.matcher_topic_name)
+ with self.pubsub_monitor:
+ pass
if __name__ == '__main__':
diff --git a/sdks/python/apache_beam/testing/pubsub_test_context.py
b/sdks/python/apache_beam/testing/pubsub_test_context.py
new file mode 100644
index 00000000000..508d8d9b2fb
--- /dev/null
+++ b/sdks/python/apache_beam/testing/pubsub_test_context.py
@@ -0,0 +1,179 @@
+#
+# 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 inspect
+import time
+import logging
+
+logger = logging.getLogger(__name__)
+
+# pylint: disable=wrong-import-order, wrong-import-position
+try:
+ from google.cloud import pubsub_v1
+except ImportError:
+ pubsub_v1 = None
+# pylint: enable=wrong-import-order, wrong-import-position
+
+
+class TestPubsubContext:
+ """A highly advanced Pub/Sub resource lifecycle manager for Python
integration tests.
+ Implements cascading third-party subscription cleanup and selective
+ graceful teardown for debugging on failures.
+
+ Includes a safety 'dry_run' switch for safe deployment and validation of
resources.
+ Any catastrophic leaks are handled independently by the global
'stale_cleaner.py'.
+ """
+ def __init__(
+ self,
+ project_id,
+ dry_run=False
+ ): # Keep dry_run=False to allow actual deletions during testing
+
+ if pubsub_v1 is None:
+ raise ImportError(
+ "The 'google-cloud-pubsub' library is required for
TestPubsubContext. "
+ "Please install it using 'pip install google-cloud-pubsub'.")
+
+ self.project_id = project_id
+ self.dry_run = dry_run
+ self.publisher = pubsub_v1.PublisherClient()
+ self.subscriber = pubsub_v1.SubscriberClient()
+
+ # Lists to track resources created during the test execution
+ self.tracked_topics = []
+ self.tracked_subscriptions = []
+ self.caller_class = "UnknownTestClass"
+ stack = inspect.stack()
+
+ for frame in stack:
+ self_obj = frame[0].f_locals.get('self', None)
+ if self_obj and hasattr(self_obj, '__class__'):
+ self.caller_class = self_obj.__class__.__name__
+ break
+
+ def register_topic(self, topic_path: str):
+ """Registers a topic to be monitored and deleted at the end."""
+ if topic_path not in self.tracked_topics:
+ self.tracked_topics.append(topic_path)
+ logger.info(
+ "[%s] Registering Topic for monitoring: %s",
+ self.caller_class,
+ topic_path)
+
+ def register_subscription(self, subscription_path: str):
+ """Registers a subscription to be monitored and deleted at the end."""
+ if subscription_path not in self.tracked_subscriptions:
+ self.tracked_subscriptions.append(subscription_path)
+ logger.info(
+ "[%s] Registering Subscription for monitoring: %s",
+ self.caller_class,
+ subscription_path)
+
+ def __enter__(self):
+ logger.info(
+ "[START] [%s] Initializing Pub/Sub context (dry_run=%s)",
+ self.caller_class,
+ self.dry_run)
+ return self
+
+ def _delete_cascading_subscriptions(self, topic_path: str):
+ """
+ Finds and deletes from GCP any third-party subscription that is
+ connected to our test topic, preventing loose residual resources.
+ """
+ logger.info(
+ "[%s] Checking for cascading subscriptions on topic: %s",
+ self.caller_class,
+ topic_path)
+ try:
+ # List all subscriptions associated with this specific topic in GCP
+ for sub_path in self.publisher.list_topic_subscriptions(
+ request={"topic": topic_path}):
+ if self.dry_run:
+ logger.info(
+ "[%s] [Cascade] (Dry Run) Would delete subscription: %s",
+ self.caller_class,
+ sub_path)
+ else:
+ logger.info(
+ "[%s] [Teardown - Cascade] Deleting residual third-party
subscription: %s",
+ self.caller_class,
+ sub_path)
+ try:
+ self.subscriber.delete_subscription(
+ request={"subscription": sub_path})
+ except Exception as e:
+ logger.error(
+ "[%s] [Error] Could not delete cascading sub %s: %s",
+ self.caller_class,
+ sub_path,
+ e)
+ except Exception as e:
+ logger.error(
+ "[%s] [Error] Could not list subs for topic %s: %s",
+ self.caller_class,
+ topic_path,
+ e)
+
+ def __exit__(self, exc_type, exc_val, exc_tb):
+ logger.info("Starting teardown of registered resources...")
+ # If the test failed (exc_type is not None), we leave the subscriptions
active for 24 hours
+ # with an automatic TTL in GCP so the developer can debug the backlog.
+ # If the test was successful, we clean up everything immediately to save
100% of the cost.
+ test_failed = exc_type is not None
+
+ if test_failed:
+ logger.warning(
+ "[%s] [ALERT] Failed test detected. Applying Graceful Teardown.",
+ self.caller_class)
+ return False
+
+ logger.info(
+ "[%s] [SUCCESS] Test passed. Proceeding with cleanup.",
+ self.caller_class)
+
+ # 1. Delete registered Subscriptions (Only if the test was successful)
+ for sub_path in list(self.tracked_subscriptions):
+ try:
+ if self.dry_run:
+ logger.info("(Dry Run) Would delete subscription: %s", sub_path)
+ else:
+ logger.info("Deleting temporary subscription: %s", sub_path)
+ self.subscriber.delete_subscription(
+ request={"subscription": sub_path})
+ self.tracked_subscriptions.remove(sub_path)
+ except Exception as e:
+ logger.error(
+ "[%s] [Error] Could not delete subscription %s: %s",
+ self.caller_class,
+ sub_path,
+ e)
+
+ # 2. Cascading Topic Cleanup (Check connected third-party subscriptions)
+ for topic_path in list(self.tracked_topics):
+ # Execute cascading deletion inspired by Java logic
+ self._delete_cascading_subscriptions(topic_path)
+ try:
+ if self.dry_run:
+ logger.info("(Dry Run) Would delete temporary topic: %s", topic_path)
+ else:
+ logger.info("Deleting temporary topic: %s", topic_path)
+ self.publisher.delete_topic(request={"topic": topic_path})
+ self.tracked_topics.remove(topic_path)
+ except Exception as e:
+ logger.error("Could not delete topic %s: %s", topic_path, e)
+ return False
diff --git a/sdks/python/apache_beam/testing/pubsub_test_context_test.py
b/sdks/python/apache_beam/testing/pubsub_test_context_test.py
new file mode 100644
index 00000000000..51fea039f33
--- /dev/null
+++ b/sdks/python/apache_beam/testing/pubsub_test_context_test.py
@@ -0,0 +1,149 @@
+#
+# 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 logging
+import unittest
+from unittest.mock import MagicMock, patch
+
+# Import the renamed class
+from apache_beam.testing.pubsub_test_context import TestPubsubContext
+
+
+class TestPubsubContextUnit(unittest.TestCase):
+
+ # This patch replaces 'pubsub_v1' with a Mock object to avoid the
+ # ImportError we set up in the __init__ method for environments without GCP.
+ @patch('apache_beam.testing.pubsub_test_context.pubsub_v1')
+ def test_context_initialization(self, mock_pubsub):
+ context = TestPubsubContext(project_id="test-project", dry_run=True)
+
+ self.assertEqual(context.project_id, "test-project")
+ self.assertTrue(context.dry_run)
+ self.assertEqual(context.tracked_topics, [])
+ self.assertEqual(context.tracked_subscriptions, [])
+
+ @patch('apache_beam.testing.pubsub_test_context.pubsub_v1')
+ def test_register_topic_and_subscription(self, mock_pubsub):
+ context = TestPubsubContext(project_id="test-project")
+
+ context.register_topic("projects/test-project/topics/test-topic")
+ context.register_subscription(
+ "projects/test-project/subscriptions/test-sub")
+
+ self.assertIn(
+ "projects/test-project/topics/test-topic", context.tracked_topics)
+ self.assertIn(
+ "projects/test-project/subscriptions/test-sub",
+ context.tracked_subscriptions)
+
+ @patch('apache_beam.testing.pubsub_test_context.pubsub_v1')
+ def test_context_manager_success_cleanup(self, mock_pubsub):
+ """Tests that the manager cleans up resources if the test passes
(dry_run=False)."""
+ context = TestPubsubContext(project_id="test-project", dry_run=False)
+
+ # Simulate registering a topic and a subscription
+ context.register_topic("topic-1")
+ context.register_subscription("sub-1")
+
+ # Simulate GCP detecting a cascading subscription
+ context.publisher.list_topic_subscriptions.return_value = ["cascade-sub-1"]
+
+ # Execute the context without errors
+ with context:
+ pass
+
+ # Verify that deletion commands were issued to GCP
+ context.subscriber.delete_subscription.assert_any_call(
+ request={"subscription": "sub-1"})
+ context.subscriber.delete_subscription.assert_any_call(
+ request={"subscription": "cascade-sub-1"})
+ context.publisher.delete_topic.assert_called_with(
+ request={"topic": "topic-1"})
+
+ @patch('apache_beam.testing.pubsub_test_context.pubsub_v1')
+ def test_context_manager_failure_skips_cleanup(self, mock_pubsub):
+ """Tests that resources are NOT deleted if the test fails (exc_type is not
None)."""
+ context = TestPubsubContext(project_id="test-project", dry_run=False)
+ context.register_topic("topic-1")
+
+ try:
+ with context:
+ raise ValueError("Simulated test failure")
+ except ValueError:
+ pass
+
+ # Since there is an error, deletion methods should NOT have been called
+ context.publisher.delete_topic.assert_not_called()
+ context.subscriber.delete_subscription.assert_not_called()
+
+ @patch('apache_beam.testing.pubsub_test_context.pubsub_v1')
+ def test_context_manager_stress_and_scale_cleanup(self, mock_pubsub):
+ """STRESS TEST: Tests that the manager can scale, monitor, and clean up
+ hundreds of concurrent topics and subscriptions safely and without leaks.
+ """
+ # Start the TestPubsubContext in active mode (dry_run=False) to simulate
real GCP interactions
+ context = TestPubsubContext(project_id="test-project", dry_run=False)
+
+ total = 1000
+ expected_topics = []
+ expected_subscriptions = []
+
+ # Bulk register 1000 topics and 1000 simulated parallel test subscriptions.
+ for i in range(total):
+ topic_path = f"projects/test-project/topics/stress-topic-{i}"
+ sub_path = f"projects/test-project/subscriptions/stress-sub-{i}"
+
+ context.register_topic(topic_path)
+ context.register_subscription(sub_path)
+
+ expected_topics.append(topic_path)
+ expected_subscriptions.append(sub_path)
+
+ # Verify that all resources were recorded in the monitor's memory without
omissions.
+ self.assertEqual(len(context.tracked_topics), total)
+ self.assertEqual(len(context.tracked_subscriptions), total)
+ self.assertEqual(context.tracked_topics, expected_topics)
+ self.assertEqual(context.tracked_subscriptions, expected_subscriptions)
+
+ # Configure the mock for the `list_topic_subscriptions` API to return an
empty list.
+ # by default to avoid infinite loops in the cascade simulation
+ context.publisher.list_topic_subscriptions.return_value = []
+
+ # Execute the mass dismantling phase
+ with context:
+ pass
+
+ # VALIDATION OF MASS SUCCESSFUL DELETION IN GCP:
+ # Verify that exactly 1000 unsubscribe calls have been issued.
+ self.assertEqual(context.subscriber.delete_subscription.call_count, total)
+ for sub in expected_subscriptions:
+ context.subscriber.delete_subscription.assert_any_call(
+ request={"subscription": sub})
+
+ # Verify that exactly 1000 topic deletion calls have been issued.
+ self.assertEqual(context.publisher.delete_topic.call_count, total)
+ for topic in expected_topics:
+ context.publisher.delete_topic.assert_any_call(request={"topic": topic})
+
+ # Verify that the manager's memory is completely clean (0 tracked
resources).
+ self.assertEqual(len(context.tracked_topics), 0)
+ self.assertEqual(len(context.tracked_subscriptions), 0)
+
+
+if __name__ == '__main__':
+ logging.basicConfig(level=logging.INFO)
+ unittest.main()