From 6719169c796e15ee95e5a4290e43ab92705f8bf2 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 12:21:24 -0400 Subject: [PATCH 1/8] fix(bigquery): add gc.collect before socket leak test to reduce flakiness --- packages/google-cloud-bigquery/tests/system/test_client.py | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index 9ddec48428b1..8e7ed8c47e56 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2182,6 +2182,11 @@ def test_dbapi_dry_run_query(self): def test_dbapi_connection_does_not_leak_sockets(self): pytest.importorskip("google.cloud.bigquery_storage") + + # Ensure any garbage collection from previous tests is done to avoid false positives. + import gc + gc.collect() + current_process = psutil.Process() conn_start = current_process.net_connections() conn_count_start = len(conn_start) From 7279ba1fa290a0f10c9daf11cd3766d2cc154a69 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 14:31:01 -0400 Subject: [PATCH 2/8] feat(bigquery): rewrite socket leak test to use whitebox object tracking --- .../tests/system/test_client.py | 106 ++++++++++++------ 1 file changed, 72 insertions(+), 34 deletions(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index 8e7ed8c47e56..e977db4e3f1a 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2183,15 +2183,49 @@ def test_dbapi_dry_run_query(self): def test_dbapi_connection_does_not_leak_sockets(self): pytest.importorskip("google.cloud.bigquery_storage") - # Ensure any garbage collection from previous tests is done to avoid false positives. - import gc - gc.collect() + import google.api_core.grpc_helpers + import requests - current_process = psutil.Process() - conn_start = current_process.net_connections() - conn_count_start = len(conn_start) + # Track HTTP Sessions + created_sessions = [] + closed_sessions = [] - with helpers.patch_tracked_requests(): + original_session_init = requests.Session.__init__ + original_session_close = requests.Session.close + + def patched_session_init(self, *args, **kwargs): + original_session_init(self, *args, **kwargs) + created_sessions.append(self) + + def patched_session_close(self): + original_session_close(self) + closed_sessions.append(self) + + # Track gRPC Channels + created_channels = [] + closed_channels = [] + + original_create_channel = google.api_core.grpc_helpers.create_channel + + def patched_create_channel(*args, **kwargs): + channel = original_create_channel(*args, **kwargs) + created_channels.append(channel) + + original_channel_close = channel.close + + def patched_channel_close(*c_args, **c_kwargs): + original_channel_close(*c_args, **c_kwargs) + closed_channels.append(channel) + + channel.close = patched_channel_close + return channel + + # Apply patches + requests.Session.__init__ = patched_session_init + requests.Session.close = patched_session_close + google.api_core.grpc_helpers.create_channel = patched_create_channel + + try: # Provide no explicit clients, so that the connection will create and own them. connection = dbapi.connect() cursor = connection.cursor() @@ -2208,37 +2242,41 @@ def test_dbapi_connection_does_not_leak_sockets(self): self.assertEqual(len(rows), 100000) connection.close() - import gc - gc.collect() - for _ in range(30): # Wait up to 3 seconds - conn_end = current_process.net_connections() - conn_count_end = len(conn_end) - if conn_count_end <= conn_count_start: - break - time.sleep(0.1) + # Assertions + self.assertEqual( + len(created_sessions), + len(closed_sessions), + f"HTTP Sessions leak detected! Created: {len(created_sessions)}, Closed: {len(closed_sessions)}", + ) - try: - self.assertLessEqual(conn_count_end, conn_count_start) - except AssertionError as e: - # Due to flakiness in this test (likely caused by OS cleanup delays or - # non-deterministic garbage collection of sockets), we want to capture - # the detailed state of connections in future failing runs to help - # decrease false positives and identify the root cause. - conn_debug = [ - f"Status: {c.status}, Laddr: {c.laddr}, Raddr: {c.raddr}" - for c in current_process.net_connections() - ] - debug_msg = "\n".join(conn_debug) - - raise AssertionError( - f"{e}\n\n" - f"--- Socket Leak Debug Info ---\n" - f"Start Count: {conn_count_start}\n" - f"End Count: {conn_count_end}\n" - f"Current Connections:\n{debug_msg}" + self.assertEqual( + len(created_channels), + len(closed_channels), + f"gRPC Channels leak detected! Created: {len(created_channels)}, Closed: {len(closed_channels)}", ) + finally: + # Revert patches + requests.Session.__init__ = original_session_init + requests.Session.close = original_session_close + google.api_core.grpc_helpers.create_channel = original_create_channel + + # Clean up any unclosed sessions/channels to avoid leaking in the test runner + for s in created_sessions: + if s not in closed_sessions: + try: + s.close() + except Exception: + pass + + for c in created_channels: + if c not in closed_channels: + try: + c.close() + except Exception: + pass + def _load_table_for_dml(self, rows, dataset_id, table_id): from google.cloud._testing import _NamedTemporaryFile from google.cloud.bigquery.job import ( From 6d31f0ca498814d7dc23657a7dd18f586814cf9d Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 16:02:20 -0400 Subject: [PATCH 3/8] feat(bigquery): use set-based comparisons in socket leak test --- .../tests/system/test_client.py | 24 ++++++++++++++----- 1 file changed, 18 insertions(+), 6 deletions(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index e977db4e3f1a..b155cff8fac2 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2244,16 +2244,28 @@ def patched_channel_close(*c_args, **c_kwargs): connection.close() # Assertions + created_session_ids = {id(s) for s in created_sessions} + closed_session_ids = {id(s) for s in closed_sessions} + leaked_session_ids = created_session_ids - closed_session_ids + self.assertEqual( - len(created_sessions), - len(closed_sessions), - f"HTTP Sessions leak detected! Created: {len(created_sessions)}, Closed: {len(closed_sessions)}", + len(leaked_session_ids), + 0, + f"HTTP Sessions leak detected! Leaked: {len(leaked_session_ids)}. " + f"Unique Created: {len(created_session_ids)}, Unique Closed: {len(closed_session_ids)}. " + f"Total Created: {len(created_sessions)}, Total Closed: {len(closed_sessions)}", ) + created_channel_ids = {id(c) for c in created_channels} + closed_channel_ids = {id(c) for c in closed_channels} + leaked_channel_ids = created_channel_ids - closed_channel_ids + self.assertEqual( - len(created_channels), - len(closed_channels), - f"gRPC Channels leak detected! Created: {len(created_channels)}, Closed: {len(closed_channels)}", + len(leaked_channel_ids), + 0, + f"gRPC Channels leak detected! Leaked: {len(leaked_channel_ids)}. " + f"Unique Created: {len(created_channel_ids)}, Unique Closed: {len(closed_channel_ids)}. " + f"Total Created: {len(created_channels)}, Total Closed: {len(closed_channels)}", ) finally: From 914a1c7ddf98ee3cfa7b54d3eace9f8f040ca177 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 16:22:04 -0400 Subject: [PATCH 4/8] feat(bigquery): add stack trace tracking to identify leaked sessions --- .../tests/system/test_client.py | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index b155cff8fac2..5e292cc4296f 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2189,13 +2189,17 @@ def test_dbapi_connection_does_not_leak_sockets(self): # Track HTTP Sessions created_sessions = [] closed_sessions = [] + session_creation_stacks = {} original_session_init = requests.Session.__init__ original_session_close = requests.Session.close + import traceback + def patched_session_init(self, *args, **kwargs): original_session_init(self, *args, **kwargs) created_sessions.append(self) + session_creation_stacks[id(self)] = traceback.format_stack(limit=10) def patched_session_close(self): original_session_close(self) @@ -2248,12 +2252,23 @@ def patched_channel_close(*c_args, **c_kwargs): closed_session_ids = {id(s) for s in closed_sessions} leaked_session_ids = created_session_ids - closed_session_ids + debug_info = "" + if leaked_session_ids: + debug_info += "\n--- Leaked Sessions Creation Stacks ---\n" + for s_id in leaked_session_ids: + debug_info += f"\nSession ID: {s_id}\n" + debug_info += "".join( + session_creation_stacks.get(s_id, ["No stack trace available"]) + ) + debug_info += "----------------------------------------\n" + self.assertEqual( len(leaked_session_ids), 0, f"HTTP Sessions leak detected! Leaked: {len(leaked_session_ids)}. " f"Unique Created: {len(created_session_ids)}, Unique Closed: {len(closed_session_ids)}. " - f"Total Created: {len(created_sessions)}, Total Closed: {len(closed_sessions)}", + f"Total Created: {len(created_sessions)}, Total Closed: {len(closed_sessions)}" + f"{debug_info}", ) created_channel_ids = {id(c) for c in created_channels} From 3c53d69ae060b95c9a7111b7809f71180c226450 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 16:41:04 -0400 Subject: [PATCH 5/8] feat(bigquery): restore patch_tracked_requests to fix known leak --- .../tests/system/test_client.py | 29 ++++++++++--------- 1 file changed, 15 insertions(+), 14 deletions(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index 5e292cc4296f..caad63133aba 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2230,22 +2230,23 @@ def patched_channel_close(*c_args, **c_kwargs): google.api_core.grpc_helpers.create_channel = patched_create_channel try: - # Provide no explicit clients, so that the connection will create and own them. - connection = dbapi.connect() - cursor = connection.cursor() - - cursor.execute( + with helpers.patch_tracked_requests(): + # Provide no explicit clients, so that the connection will create and own them. + connection = dbapi.connect() + cursor = connection.cursor() + + cursor.execute( + """ + SELECT id, `by`, timestamp + FROM `bigquery-public-data.hacker_news.full` + ORDER BY `id` ASC + LIMIT 100000 """ - SELECT id, `by`, timestamp - FROM `bigquery-public-data.hacker_news.full` - ORDER BY `id` ASC - LIMIT 100000 - """ - ) - rows = cursor.fetchall() - self.assertEqual(len(rows), 100000) + ) + rows = cursor.fetchall() + self.assertEqual(len(rows), 100000) - connection.close() + connection.close() # Assertions created_session_ids = {id(s) for s in created_sessions} From f8442b7ffa7ac1ef3b8f0f804e0e209e55eaa497 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 20:34:25 -0400 Subject: [PATCH 6/8] refactor(bigquery): use sets for tracking leaked sessions/channels per PR comments --- .../tests/system/test_client.py | 79 ++++++------------- 1 file changed, 25 insertions(+), 54 deletions(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index caad63133aba..a2620ca55746 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2187,39 +2187,35 @@ def test_dbapi_connection_does_not_leak_sockets(self): import requests # Track HTTP Sessions - created_sessions = [] - closed_sessions = [] - session_creation_stacks = {} + created_sessions = set() + closed_sessions = set() original_session_init = requests.Session.__init__ original_session_close = requests.Session.close - import traceback - def patched_session_init(self, *args, **kwargs): original_session_init(self, *args, **kwargs) - created_sessions.append(self) - session_creation_stacks[id(self)] = traceback.format_stack(limit=10) + created_sessions.add(self) def patched_session_close(self): original_session_close(self) - closed_sessions.append(self) + closed_sessions.add(self) # Track gRPC Channels - created_channels = [] - closed_channels = [] + created_channels = set() + closed_channels = set() original_create_channel = google.api_core.grpc_helpers.create_channel def patched_create_channel(*args, **kwargs): channel = original_create_channel(*args, **kwargs) - created_channels.append(channel) + created_channels.add(channel) original_channel_close = channel.close def patched_channel_close(*c_args, **c_kwargs): original_channel_close(*c_args, **c_kwargs) - closed_channels.append(channel) + closed_channels.add(channel) channel.close = patched_channel_close return channel @@ -2249,39 +2245,16 @@ def patched_channel_close(*c_args, **c_kwargs): connection.close() # Assertions - created_session_ids = {id(s) for s in created_sessions} - closed_session_ids = {id(s) for s in closed_sessions} - leaked_session_ids = created_session_ids - closed_session_ids - - debug_info = "" - if leaked_session_ids: - debug_info += "\n--- Leaked Sessions Creation Stacks ---\n" - for s_id in leaked_session_ids: - debug_info += f"\nSession ID: {s_id}\n" - debug_info += "".join( - session_creation_stacks.get(s_id, ["No stack trace available"]) - ) - debug_info += "----------------------------------------\n" - self.assertEqual( - len(leaked_session_ids), - 0, - f"HTTP Sessions leak detected! Leaked: {len(leaked_session_ids)}. " - f"Unique Created: {len(created_session_ids)}, Unique Closed: {len(closed_session_ids)}. " - f"Total Created: {len(created_sessions)}, Total Closed: {len(closed_sessions)}" - f"{debug_info}", + created_sessions, + closed_sessions, + f"HTTP Sessions leak detected! Leaked: {created_sessions - closed_sessions}", ) - created_channel_ids = {id(c) for c in created_channels} - closed_channel_ids = {id(c) for c in closed_channels} - leaked_channel_ids = created_channel_ids - closed_channel_ids - self.assertEqual( - len(leaked_channel_ids), - 0, - f"gRPC Channels leak detected! Leaked: {len(leaked_channel_ids)}. " - f"Unique Created: {len(created_channel_ids)}, Unique Closed: {len(closed_channel_ids)}. " - f"Total Created: {len(created_channels)}, Total Closed: {len(closed_channels)}", + created_channels, + closed_channels, + f"gRPC Channels leak detected! Leaked: {created_channels - closed_channels}", ) finally: @@ -2291,19 +2264,17 @@ def patched_channel_close(*c_args, **c_kwargs): google.api_core.grpc_helpers.create_channel = original_create_channel # Clean up any unclosed sessions/channels to avoid leaking in the test runner - for s in created_sessions: - if s not in closed_sessions: - try: - s.close() - except Exception: - pass - - for c in created_channels: - if c not in closed_channels: - try: - c.close() - except Exception: - pass + for s in created_sessions - closed_sessions: + try: + s.close() + except Exception: + pass + + for c in created_channels - closed_channels: + try: + c.close() + except Exception: + pass def _load_table_for_dml(self, rows, dataset_id, table_id): from google.cloud._testing import _NamedTemporaryFile From c021da45b700a586991555f2bdf1703988acbeaa Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 21:07:50 -0400 Subject: [PATCH 7/8] refactor(bigquery): use instance-level patching for Session.close() for symmetry and isolation --- .../tests/system/test_client.py | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/packages/google-cloud-bigquery/tests/system/test_client.py b/packages/google-cloud-bigquery/tests/system/test_client.py index a2620ca55746..6f14cc1ed6a0 100644 --- a/packages/google-cloud-bigquery/tests/system/test_client.py +++ b/packages/google-cloud-bigquery/tests/system/test_client.py @@ -2191,15 +2191,18 @@ def test_dbapi_connection_does_not_leak_sockets(self): closed_sessions = set() original_session_init = requests.Session.__init__ - original_session_close = requests.Session.close def patched_session_init(self, *args, **kwargs): original_session_init(self, *args, **kwargs) created_sessions.add(self) - def patched_session_close(self): - original_session_close(self) - closed_sessions.add(self) + original_close = self.close + + def patched_close(*s_args, **s_kwargs): + original_close(*s_args, **s_kwargs) + closed_sessions.add(self) + + self.close = patched_close # Track gRPC Channels created_channels = set() @@ -2222,7 +2225,6 @@ def patched_channel_close(*c_args, **c_kwargs): # Apply patches requests.Session.__init__ = patched_session_init - requests.Session.close = patched_session_close google.api_core.grpc_helpers.create_channel = patched_create_channel try: @@ -2260,7 +2262,6 @@ def patched_channel_close(*c_args, **c_kwargs): finally: # Revert patches requests.Session.__init__ = original_session_init - requests.Session.close = original_session_close google.api_core.grpc_helpers.create_channel = original_create_channel # Clean up any unclosed sessions/channels to avoid leaking in the test runner From 4d49f175ff0e4197ddae4bb3b1dee0eef7e817d6 Mon Sep 17 00:00:00 2001 From: chalmer lowe Date: Thu, 30 Jul 2026 21:38:05 -0400 Subject: [PATCH 8/8] fix(bigquery): close cursors before clients and add defensive None checks in Connection.close() --- .../google/cloud/bigquery/dbapi/connection.py | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/packages/google-cloud-bigquery/google/cloud/bigquery/dbapi/connection.py b/packages/google-cloud-bigquery/google/cloud/bigquery/dbapi/connection.py index 53f32ce387ce..450e4188beef 100644 --- a/packages/google-cloud-bigquery/google/cloud/bigquery/dbapi/connection.py +++ b/packages/google-cloud-bigquery/google/cloud/bigquery/dbapi/connection.py @@ -78,10 +78,14 @@ def close(self): """ self._closed = True - if self._owns_client: + for cursor_ in list(self._cursors_created): + if not cursor_._closed: + cursor_.close() + + if self._owns_client and self._client is not None: self._client.close() - if self._owns_bqstorage_client: + if self._owns_bqstorage_client and self._bqstorage_client is not None: # There is no close() on the BQ Storage client itself. transport = self._bqstorage_client.transport transport.close() @@ -95,10 +99,6 @@ def close(self): if channel is not None: channel.close() - for cursor_ in self._cursors_created: - if not cursor_._closed: - cursor_.close() - def commit(self): """No-op, but for consistency raise an error if connection is closed."""