From d16d9b2b83cd1cb9504a1fc90a790ea640b310d3 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 16:05:49 -0500 Subject: [PATCH 01/13] PYTHON-6136 Fix pool checkout accounting rollback on BaseException A BaseException (KeyboardInterrupt, CancelledError, GreenletExit) landing inside Pool.checkout could leak checkout accounting: the size gate (self.requests) was only rolled back after a socket was acquired, and the maxConnecting _pending counter was only rolled back on the happy path. Roll back both counters, with notify(), when a BaseException interrupts the gate or the connection attempt. Rework the gevent killall race test to amplify the checkout unwind windows and assert the counters fully drain after the pool settles. --- pymongo/asynchronous/pool.py | 66 ++++++++++++++++++++++++-------- pymongo/synchronous/pool.py | 66 ++++++++++++++++++++++++-------- test/asynchronous/test_client.py | 57 +++++++++++++++++++-------- test/test_client.py | 57 +++++++++++++++++++-------- 4 files changed, 182 insertions(+), 64 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 69ac14e7e4..7d34193d73 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -980,24 +980,44 @@ async def _get_conn( else: deadline = None - async with self.size_cond: - self._raise_if_not_ready(checkout_started_time, emit_event=True) - while not (self.requests < self.max_pool_size): - timeout = deadline - time.monotonic() if deadline else None - if not await _async_cond_wait(self.size_cond, timeout): - # Timed out, notify the next thread to ensure a - # timeout doesn't consume the condition. - if self.requests < self.max_pool_size: - self.size_cond.notify() - self._raise_wait_queue_timeout(checkout_started_time) + # Roll back the size gate if a BaseException lands here (PYTHON-6136). + requests_incremented = False + try: + async with self.size_cond: self._raise_if_not_ready(checkout_started_time, emit_event=True) - self.requests += 1 + while not (self.requests < self.max_pool_size): + timeout = deadline - time.monotonic() if deadline else None + if not await _async_cond_wait(self.size_cond, timeout): + # Timed out, notify the next thread to ensure a + # timeout doesn't consume the condition. + if self.requests < self.max_pool_size: + self.size_cond.notify() + self._raise_wait_queue_timeout(checkout_started_time) + self._raise_if_not_ready(checkout_started_time, emit_event=True) + self.requests += 1 + requests_incremented = True + except BaseException: + if requests_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + accounted = False + try: + async with self.size_cond: + self.requests -= 1 + accounted = True + self.size_cond.notify() + finally: + if not accounted: + async with self.size_cond: + self.requests -= 1 + self.size_cond.notify() + raise # We've now acquired the semaphore and must release it on error. conn = None incremented = False emitted_event = False is_new_conn = False + pending_incremented = False try: async with self.lock: self.active_sockets += 1 @@ -1022,6 +1042,7 @@ async def _get_conn( conn = self.conns.popleft() except IndexError: self._pending += 1 + pending_incremented = True if conn: # We got a socket from the pool if await self._perished(conn): conn = None @@ -1033,6 +1054,7 @@ async def _get_conn( finally: async with self._max_connecting_cond: self._pending -= 1 + pending_incremented = False self._max_connecting_cond.notify() conn.active = True @@ -1043,12 +1065,24 @@ async def _get_conn( self.active_contexts.add(conn.cancel_context) # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: + if pending_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + pending_accounted = False + try: + async with self._max_connecting_cond: + self._pending -= 1 + pending_accounted = True + self._max_connecting_cond.notify() + finally: + if not pending_accounted: + async with self._max_connecting_cond: + self._pending -= 1 + self._max_connecting_cond.notify() + if conn: # We checked out a socket but authentication failed. await conn.close_conn(ConnectionClosedReason.ERROR) - # Re-apply the accounting if a GreenletExit interrupts - # during the size_cond acquisition; during unwind gevent - # lets the re-acquire complete (PYTHON-6074). + # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). accounted = False try: async with self.size_cond: @@ -1121,9 +1155,7 @@ async def checkin(self, conn: AsyncConnection) -> None: conn.pinned_cursor = False self._pinned_sockets.discard(conn) forked = self.pid != os.getpid() - # Re-apply the accounting if a gevent GreenletExit interrupts during - # the size_cond acquisition; gevent lets the re-acquire complete while - # unwinding (PYTHON-6074). + # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). close_conn_reason: Optional[str] = None emit_closed = False accounted = False diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index c977729ba8..56713aeae5 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -976,24 +976,44 @@ def _get_conn( else: deadline = None - with self.size_cond: - self._raise_if_not_ready(checkout_started_time, emit_event=True) - while not (self.requests < self.max_pool_size): - timeout = deadline - time.monotonic() if deadline else None - if not _cond_wait(self.size_cond, timeout): - # Timed out, notify the next thread to ensure a - # timeout doesn't consume the condition. - if self.requests < self.max_pool_size: - self.size_cond.notify() - self._raise_wait_queue_timeout(checkout_started_time) + # Roll back the size gate if a BaseException lands here (PYTHON-6136). + requests_incremented = False + try: + with self.size_cond: self._raise_if_not_ready(checkout_started_time, emit_event=True) - self.requests += 1 + while not (self.requests < self.max_pool_size): + timeout = deadline - time.monotonic() if deadline else None + if not _cond_wait(self.size_cond, timeout): + # Timed out, notify the next thread to ensure a + # timeout doesn't consume the condition. + if self.requests < self.max_pool_size: + self.size_cond.notify() + self._raise_wait_queue_timeout(checkout_started_time) + self._raise_if_not_ready(checkout_started_time, emit_event=True) + self.requests += 1 + requests_incremented = True + except BaseException: + if requests_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + accounted = False + try: + with self.size_cond: + self.requests -= 1 + accounted = True + self.size_cond.notify() + finally: + if not accounted: + with self.size_cond: + self.requests -= 1 + self.size_cond.notify() + raise # We've now acquired the semaphore and must release it on error. conn = None incremented = False emitted_event = False is_new_conn = False + pending_incremented = False try: with self.lock: self.active_sockets += 1 @@ -1018,6 +1038,7 @@ def _get_conn( conn = self.conns.popleft() except IndexError: self._pending += 1 + pending_incremented = True if conn: # We got a socket from the pool if self._perished(conn): conn = None @@ -1029,6 +1050,7 @@ def _get_conn( finally: with self._max_connecting_cond: self._pending -= 1 + pending_incremented = False self._max_connecting_cond.notify() conn.active = True @@ -1039,12 +1061,24 @@ def _get_conn( self.active_contexts.add(conn.cancel_context) # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: + if pending_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + pending_accounted = False + try: + with self._max_connecting_cond: + self._pending -= 1 + pending_accounted = True + self._max_connecting_cond.notify() + finally: + if not pending_accounted: + with self._max_connecting_cond: + self._pending -= 1 + self._max_connecting_cond.notify() + if conn: # We checked out a socket but authentication failed. conn.close_conn(ConnectionClosedReason.ERROR) - # Re-apply the accounting if a GreenletExit interrupts - # during the size_cond acquisition; during unwind gevent - # lets the re-acquire complete (PYTHON-6074). + # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). accounted = False try: with self.size_cond: @@ -1117,9 +1151,7 @@ def checkin(self, conn: Connection) -> None: conn.pinned_cursor = False self._pinned_sockets.discard(conn) forked = self.pid != os.getpid() - # Re-apply the accounting if a gevent GreenletExit interrupts during - # the size_cond acquisition; gevent lets the re-acquire complete while - # unwinding (PYTHON-6074). + # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). close_conn_reason: Optional[str] = None emit_closed = False accounted = False diff --git a/test/asynchronous/test_client.py b/test/asynchronous/test_client.py index 56810b1433..b8ec8479d4 100644 --- a/test/asynchronous/test_client.py +++ b/test/asynchronous/test_client.py @@ -2744,20 +2744,18 @@ def test_gevent_kill_churn_deadlock(self): import gevent.thread as _gthread from gevent import Timeout, spawn - # AMPLIFY_RACE=1 widens gevent's brief sleep on a contended lock - # so a kill lands there reliably. Only the bare sleep() is - # widened; timed sleeps pass through. The unfixed test then - # deadlocks within seconds. + # AMPLIFY_RACE=1 widens the kill windows below; keep AMPLIFY_SECONDS + # small or ops starve. + amplify_seconds = 0.0 if os.environ.get("AMPLIFY_RACE", "0") == "1": - _AMPLIFY_SECONDS = float(os.environ.get("AMPLIFY_SECONDS", "0.02")) + amplify_seconds = float(os.environ.get("AMPLIFY_SECONDS", "0.002")) _orig_thread_sleep = _gthread.sleep def _amplified_sleep(*args): - if not args: # bare sleep(): the courtesy yield on a failed - # non-blocking acquire (Condition.notify -> _is_owned -> acquire(False)) - _orig_thread_sleep(_AMPLIFY_SECONDS) - else: # sleep(0.001), sleep(2), etc.: passthrough + if not args: # bare sleep(): checkin courtesy yield + _orig_thread_sleep(amplify_seconds) + else: # timed sleeps pass through _orig_thread_sleep(*args) _gthread.sleep = _amplified_sleep @@ -2767,6 +2765,31 @@ def _amplified_sleep(*args): coll = client.pymongo_test.coll coll.insert_one({}) + # Widen the post-gate checkout windows (PYTHON-6136). + if amplify_seconds: + pool = async_get_pool(client) # type:ignore + + class _AmplifiedCondition: + def __init__(self, cond, seconds): + self._cond = cond + self._seconds = seconds + + def __enter__(self): + self._cond.__enter__() + return self + + def __exit__(self, *args): + self._cond.__exit__(*args) + time.sleep(self._seconds) + + def __getattr__(self, name): + return getattr(self._cond, name) + + pool.size_cond = _AmplifiedCondition(pool.size_cond, amplify_seconds) + pool._max_connecting_cond = _AmplifiedCondition( + pool._max_connecting_cond, amplify_seconds + ) + op_count = [0] running = [True] workers: list = [] @@ -2780,9 +2803,12 @@ def worker(): except Exception: return + # Scale the kill cadence or ops starve. + reaper_interval = max(0.003, amplify_seconds * 8) + def reaper(): while running[0]: - time.sleep(0.003) + time.sleep(reaper_interval) if not workers: continue idx = random.randrange(len(workers)) @@ -2823,12 +2849,8 @@ def reaper(): coll.find_one({}) except Timeout: self.fail("Pool gate saturated (PYTHON-6074)") - # Deterministic check: a saturated size gate pins the pool's - # checkout counters at maxPoolSize (PYTHON-6074). - pool = async_get_pool(client) # type:ignore - self.assertLess(pool.requests, pool.max_pool_size) - self.assertLess(pool.active_sockets, pool.max_pool_size) self.assertGreater(op_count[0], 0) + pool = async_get_pool(client) # type:ignore finally: running[0] = False gevent.killall(workers, block=False) @@ -2845,6 +2867,11 @@ def reaper(): client.close() except Timeout: pass + # Post-settle: the counters must have fully drained. + time.sleep(1.0) + self.assertEqual(pool.requests, 0) + self.assertEqual(pool._pending, 0) + self.assertEqual(pool.active_sockets, 0) class TestClientLazyConnect(AsyncIntegrationTest): diff --git a/test/test_client.py b/test/test_client.py index ebd1e670ee..9e7a41d506 100644 --- a/test/test_client.py +++ b/test/test_client.py @@ -2695,20 +2695,18 @@ def test_gevent_kill_churn_deadlock(self): import gevent.thread as _gthread from gevent import Timeout, spawn - # AMPLIFY_RACE=1 widens gevent's brief sleep on a contended lock - # so a kill lands there reliably. Only the bare sleep() is - # widened; timed sleeps pass through. The unfixed test then - # deadlocks within seconds. + # AMPLIFY_RACE=1 widens the kill windows below; keep AMPLIFY_SECONDS + # small or ops starve. + amplify_seconds = 0.0 if os.environ.get("AMPLIFY_RACE", "0") == "1": - _AMPLIFY_SECONDS = float(os.environ.get("AMPLIFY_SECONDS", "0.02")) + amplify_seconds = float(os.environ.get("AMPLIFY_SECONDS", "0.002")) _orig_thread_sleep = _gthread.sleep def _amplified_sleep(*args): - if not args: # bare sleep(): the courtesy yield on a failed - # non-blocking acquire (Condition.notify -> _is_owned -> acquire(False)) - _orig_thread_sleep(_AMPLIFY_SECONDS) - else: # sleep(0.001), sleep(2), etc.: passthrough + if not args: # bare sleep(): checkin courtesy yield + _orig_thread_sleep(amplify_seconds) + else: # timed sleeps pass through _orig_thread_sleep(*args) _gthread.sleep = _amplified_sleep @@ -2718,6 +2716,31 @@ def _amplified_sleep(*args): coll = client.pymongo_test.coll coll.insert_one({}) + # Widen the post-gate checkout windows (PYTHON-6136). + if amplify_seconds: + pool = get_pool(client) # type:ignore + + class _AmplifiedCondition: + def __init__(self, cond, seconds): + self._cond = cond + self._seconds = seconds + + def __enter__(self): + self._cond.__enter__() + return self + + def __exit__(self, *args): + self._cond.__exit__(*args) + time.sleep(self._seconds) + + def __getattr__(self, name): + return getattr(self._cond, name) + + pool.size_cond = _AmplifiedCondition(pool.size_cond, amplify_seconds) + pool._max_connecting_cond = _AmplifiedCondition( + pool._max_connecting_cond, amplify_seconds + ) + op_count = [0] running = [True] workers: list = [] @@ -2731,9 +2754,12 @@ def worker(): except Exception: return + # Scale the kill cadence or ops starve. + reaper_interval = max(0.003, amplify_seconds * 8) + def reaper(): while running[0]: - time.sleep(0.003) + time.sleep(reaper_interval) if not workers: continue idx = random.randrange(len(workers)) @@ -2774,12 +2800,8 @@ def reaper(): coll.find_one({}) except Timeout: self.fail("Pool gate saturated (PYTHON-6074)") - # Deterministic check: a saturated size gate pins the pool's - # checkout counters at maxPoolSize (PYTHON-6074). - pool = get_pool(client) # type:ignore - self.assertLess(pool.requests, pool.max_pool_size) - self.assertLess(pool.active_sockets, pool.max_pool_size) self.assertGreater(op_count[0], 0) + pool = get_pool(client) # type:ignore finally: running[0] = False gevent.killall(workers, block=False) @@ -2796,6 +2818,11 @@ def reaper(): client.close() except Timeout: pass + # Post-settle: the counters must have fully drained. + time.sleep(1.0) + self.assertEqual(pool.requests, 0) + self.assertEqual(pool._pending, 0) + self.assertEqual(pool.active_sockets, 0) class TestClientLazyConnect(IntegrationTest): From 04dccc5f7fda56c0b1a0afc027a0c59f7bcc2f13 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 17:27:24 -0500 Subject: [PATCH 02/13] PYTHON-6136 Always roll back the size gate on checkout unwind --- pymongo/asynchronous/pool.py | 49 +++++++++++++++++++----------------- pymongo/synchronous/pool.py | 49 +++++++++++++++++++----------------- 2 files changed, 52 insertions(+), 46 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 7d34193d73..ec1c88a7e7 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1065,39 +1065,42 @@ async def _get_conn( self.active_contexts.add(conn.cancel_context) # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: - if pending_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - pending_accounted = False - try: - async with self._max_connecting_cond: - self._pending -= 1 - pending_accounted = True - self._max_connecting_cond.notify() - finally: - if not pending_accounted: + try: + if pending_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + pending_accounted = False + try: async with self._max_connecting_cond: self._pending -= 1 + pending_accounted = True self._max_connecting_cond.notify() + finally: + if not pending_accounted: + async with self._max_connecting_cond: + self._pending -= 1 + self._max_connecting_cond.notify() - if conn: - # We checked out a socket but authentication failed. - await conn.close_conn(ConnectionClosedReason.ERROR) - # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). - accounted = False - try: - async with self.size_cond: - self.requests -= 1 - if incremented: - self.active_sockets -= 1 - accounted = True - self.size_cond.notify() + if conn: + # We checked out a socket but authentication failed. + await conn.close_conn(ConnectionClosedReason.ERROR) finally: - if not accounted: + # Always roll back the size gate, even if the cleanups above + # were interrupted (PYTHON-6136). + accounted = False + try: async with self.size_cond: self.requests -= 1 if incremented: self.active_sockets -= 1 + accounted = True self.size_cond.notify() + finally: + if not accounted: + async with self.size_cond: + self.requests -= 1 + if incremented: + self.active_sockets -= 1 + self.size_cond.notify() if not emitted_event: self._telemetry.checkout_failed( diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index 56713aeae5..a851f6a9b7 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1061,39 +1061,42 @@ def _get_conn( self.active_contexts.add(conn.cancel_context) # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: - if pending_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - pending_accounted = False - try: - with self._max_connecting_cond: - self._pending -= 1 - pending_accounted = True - self._max_connecting_cond.notify() - finally: - if not pending_accounted: + try: + if pending_incremented: + # Gevent grants the re-acquire during unwind (PYTHON-6074). + pending_accounted = False + try: with self._max_connecting_cond: self._pending -= 1 + pending_accounted = True self._max_connecting_cond.notify() + finally: + if not pending_accounted: + with self._max_connecting_cond: + self._pending -= 1 + self._max_connecting_cond.notify() - if conn: - # We checked out a socket but authentication failed. - conn.close_conn(ConnectionClosedReason.ERROR) - # Re-apply the accounting if a BaseException interrupts here (PYTHON-6074). - accounted = False - try: - with self.size_cond: - self.requests -= 1 - if incremented: - self.active_sockets -= 1 - accounted = True - self.size_cond.notify() + if conn: + # We checked out a socket but authentication failed. + conn.close_conn(ConnectionClosedReason.ERROR) finally: - if not accounted: + # Always roll back the size gate, even if the cleanups above + # were interrupted (PYTHON-6136). + accounted = False + try: with self.size_cond: self.requests -= 1 if incremented: self.active_sockets -= 1 + accounted = True self.size_cond.notify() + finally: + if not accounted: + with self.size_cond: + self.requests -= 1 + if incremented: + self.active_sockets -= 1 + self.size_cond.notify() if not emitted_event: self._telemetry.checkout_failed( From 2191a2fe0ff69ddd4cd1bbb6efe7e0d4d8f9029f Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 19:51:13 -0500 Subject: [PATCH 03/13] PYTHON-6136 Add regression tests for interrupted checkout cleanup --- test/asynchronous/test_pooling.py | 50 +++++++++++++++++++++++++++++++ test/test_pooling.py | 50 +++++++++++++++++++++++++++++++ 2 files changed, 100 insertions(+) diff --git a/test/asynchronous/test_pooling.py b/test/asynchronous/test_pooling.py index 661bd4e3d2..8452d18b61 100644 --- a/test/asynchronous/test_pooling.py +++ b/test/asynchronous/test_pooling.py @@ -342,6 +342,56 @@ async def __aexit__(self, *args): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + async def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): + # PYTHON-6136: a KeyboardInterrupt from connect() must roll back the + # size gate, the maxConnecting gate, and active_sockets. + cx_pool = await self.create_pool(max_pool_size=1) + + with patch.object(cx_pool, "connect", side_effect=KeyboardInterrupt()): + with self.assertRaises(KeyboardInterrupt): + async with cx_pool.checkout(): + pass + + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + + async def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): + # PYTHON-6136: an interruption during the pending-gate cleanup must + # not skip the size-gate rollback below it. + cx_pool = await self.create_pool(max_pool_size=1) + + class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): + def __init__(self, lock): + super().__init__(lock) + self.enters = 0 + + async def __aenter__(self): + self.enters += 1 + if self.enters == 2: + # First enter is the checkout wait, second is connect()'s + # cleanup. Simulate a kill delivered while waiting there. + raise KeyboardInterrupt() + return await super().__aenter__() + + async def __aexit__(self, *args): + return await super().__aexit__(*args) + + def notify(self, n=1): + # The handler's pending-gate rollback is itself killed. + raise KeyboardInterrupt() + + cx_pool._max_connecting_cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + + with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): + with self.assertRaises(KeyboardInterrupt): + async with cx_pool.checkout(): + pass + + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + async def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. cx_pool = await self.create_pool() diff --git a/test/test_pooling.py b/test/test_pooling.py index a3f0eaf589..2e8e243cbf 100644 --- a/test/test_pooling.py +++ b/test/test_pooling.py @@ -342,6 +342,56 @@ def __exit__(self, *args): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): + # PYTHON-6136: a KeyboardInterrupt from connect() must roll back the + # size gate, the maxConnecting gate, and active_sockets. + cx_pool = self.create_pool(max_pool_size=1) + + with patch.object(cx_pool, "connect", side_effect=KeyboardInterrupt()): + with self.assertRaises(KeyboardInterrupt): + with cx_pool.checkout(): + pass + + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + + def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): + # PYTHON-6136: an interruption during the pending-gate cleanup must + # not skip the size-gate rollback below it. + cx_pool = self.create_pool(max_pool_size=1) + + class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): + def __init__(self, lock): + super().__init__(lock) + self.enters = 0 + + def __enter__(self): + self.enters += 1 + if self.enters == 2: + # First enter is the checkout wait, second is connect()'s + # cleanup. Simulate a kill delivered while waiting there. + raise KeyboardInterrupt() + return super().__enter__() + + def __exit__(self, *args): + return super().__exit__(*args) + + def notify(self, n=1): + # The handler's pending-gate rollback is itself killed. + raise KeyboardInterrupt() + + cx_pool._max_connecting_cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + + with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): + with self.assertRaises(KeyboardInterrupt): + with cx_pool.checkout(): + pass + + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. cx_pool = self.create_pool() From a909572aad3199c59e30d4bf7025fc6e9dcd5fcd Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 20:05:04 -0500 Subject: [PATCH 04/13] PYTHON-6136 Bind pool before the race test's try block --- test/asynchronous/test_client.py | 3 +-- test/test_client.py | 3 +-- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/test/asynchronous/test_client.py b/test/asynchronous/test_client.py index b8ec8479d4..62133b814f 100644 --- a/test/asynchronous/test_client.py +++ b/test/asynchronous/test_client.py @@ -2765,9 +2765,9 @@ def _amplified_sleep(*args): coll = client.pymongo_test.coll coll.insert_one({}) + pool = async_get_pool(client) # type:ignore # Widen the post-gate checkout windows (PYTHON-6136). if amplify_seconds: - pool = async_get_pool(client) # type:ignore class _AmplifiedCondition: def __init__(self, cond, seconds): @@ -2850,7 +2850,6 @@ def reaper(): except Timeout: self.fail("Pool gate saturated (PYTHON-6074)") self.assertGreater(op_count[0], 0) - pool = async_get_pool(client) # type:ignore finally: running[0] = False gevent.killall(workers, block=False) diff --git a/test/test_client.py b/test/test_client.py index 9e7a41d506..d6e004fa75 100644 --- a/test/test_client.py +++ b/test/test_client.py @@ -2716,9 +2716,9 @@ def _amplified_sleep(*args): coll = client.pymongo_test.coll coll.insert_one({}) + pool = get_pool(client) # type:ignore # Widen the post-gate checkout windows (PYTHON-6136). if amplify_seconds: - pool = get_pool(client) # type:ignore class _AmplifiedCondition: def __init__(self, cond, seconds): @@ -2801,7 +2801,6 @@ def reaper(): except Timeout: self.fail("Pool gate saturated (PYTHON-6074)") self.assertGreater(op_count[0], 0) - pool = get_pool(client) # type:ignore finally: running[0] = False gevent.killall(workers, block=False) From dbbc6150ebbde413be862b19040ea5fb804c17eb Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 20:29:05 -0500 Subject: [PATCH 05/13] PYTHON-6136 Roll back operation_count on every failed checkout --- pymongo/asynchronous/pool.py | 7 +++++++ pymongo/synchronous/pool.py | 7 +++++++ test/asynchronous/test_pooling.py | 6 ++++++ test/test_pooling.py | 6 ++++++ 4 files changed, 26 insertions(+) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index ec1c88a7e7..3e54c78b95 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1003,13 +1003,18 @@ async def _get_conn( try: async with self.size_cond: self.requests -= 1 + self.operation_count -= 1 accounted = True self.size_cond.notify() finally: if not accounted: async with self.size_cond: self.requests -= 1 + self.operation_count -= 1 self.size_cond.notify() + else: + # The gate never admitted; still undo the load increment above. + self.operation_count -= 1 raise # We've now acquired the semaphore and must release it on error. @@ -1090,6 +1095,7 @@ async def _get_conn( try: async with self.size_cond: self.requests -= 1 + self.operation_count -= 1 if incremented: self.active_sockets -= 1 accounted = True @@ -1098,6 +1104,7 @@ async def _get_conn( if not accounted: async with self.size_cond: self.requests -= 1 + self.operation_count -= 1 if incremented: self.active_sockets -= 1 self.size_cond.notify() diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index a851f6a9b7..9ae94b8365 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -999,13 +999,18 @@ def _get_conn( try: with self.size_cond: self.requests -= 1 + self.operation_count -= 1 accounted = True self.size_cond.notify() finally: if not accounted: with self.size_cond: self.requests -= 1 + self.operation_count -= 1 self.size_cond.notify() + else: + # The gate never admitted; still undo the load increment above. + self.operation_count -= 1 raise # We've now acquired the semaphore and must release it on error. @@ -1086,6 +1091,7 @@ def _get_conn( try: with self.size_cond: self.requests -= 1 + self.operation_count -= 1 if incremented: self.active_sockets -= 1 accounted = True @@ -1094,6 +1100,7 @@ def _get_conn( if not accounted: with self.size_cond: self.requests -= 1 + self.operation_count -= 1 if incremented: self.active_sockets -= 1 self.size_cond.notify() diff --git a/test/asynchronous/test_pooling.py b/test/asynchronous/test_pooling.py index 8452d18b61..948f0648f4 100644 --- a/test/asynchronous/test_pooling.py +++ b/test/asynchronous/test_pooling.py @@ -306,6 +306,7 @@ def notify(): # Accounting was applied exactly once. self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool.operation_count) async def test_checkout_error_accounting_on_kill_during_acquire(self): # PYTHON-6074: an exception delivered while the checkout error @@ -341,6 +342,7 @@ async def __aexit__(self, *args): # The fallback applied the accounting exactly once. self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool.operation_count) async def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): # PYTHON-6136: a KeyboardInterrupt from connect() must roll back the @@ -355,6 +357,7 @@ async def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) async def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must @@ -391,6 +394,7 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) async def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. @@ -513,6 +517,8 @@ async def test_wait_queue_timeout(self): 1, f"Waited {duration:.2f} seconds for a socket, expected {wait_queue_timeout:f}", ) + # The load metric must not be inflated by the failed checkout. + self.assertEqual(0, pool.operation_count) async def test_no_wait_queue_timeout(self): # Verify get_socket() with no wait_queue_timeout blocks forever. diff --git a/test/test_pooling.py b/test/test_pooling.py index 2e8e243cbf..fead75b6f8 100644 --- a/test/test_pooling.py +++ b/test/test_pooling.py @@ -306,6 +306,7 @@ def notify(): # Accounting was applied exactly once. self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool.operation_count) def test_checkout_error_accounting_on_kill_during_acquire(self): # PYTHON-6074: an exception delivered while the checkout error @@ -341,6 +342,7 @@ def __exit__(self, *args): # The fallback applied the accounting exactly once. self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool.operation_count) def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): # PYTHON-6136: a KeyboardInterrupt from connect() must roll back the @@ -355,6 +357,7 @@ def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must @@ -391,6 +394,7 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.requests) self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. @@ -513,6 +517,8 @@ def test_wait_queue_timeout(self): 1, f"Waited {duration:.2f} seconds for a socket, expected {wait_queue_timeout:f}", ) + # The load metric must not be inflated by the failed checkout. + self.assertEqual(0, pool.operation_count) def test_no_wait_queue_timeout(self): # Verify get_socket() with no wait_queue_timeout blocks forever. From 40762fa1f3beffef02139ef4c7f3a2ec742bb168 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 20:42:36 -0500 Subject: [PATCH 06/13] PYTHON-6136 Lock-protect the operation_count rollback on the gate path --- pymongo/asynchronous/pool.py | 10 +++++++++- pymongo/synchronous/pool.py | 10 +++++++++- 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 3e54c78b95..3c58cb924f 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1014,7 +1014,15 @@ async def _get_conn( self.size_cond.notify() else: # The gate never admitted; still undo the load increment above. - self.operation_count -= 1 + accounted = False + try: + async with self.size_cond: + self.operation_count -= 1 + accounted = True + finally: + if not accounted: + async with self.size_cond: + self.operation_count -= 1 raise # We've now acquired the semaphore and must release it on error. diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index 9ae94b8365..7dd5722df8 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1010,7 +1010,15 @@ def _get_conn( self.size_cond.notify() else: # The gate never admitted; still undo the load increment above. - self.operation_count -= 1 + accounted = False + try: + with self.size_cond: + self.operation_count -= 1 + accounted = True + finally: + if not accounted: + with self.size_cond: + self.operation_count -= 1 raise # We've now acquired the semaphore and must release it on error. From 397e7a419d15a43702f7060da2de30fba72728ef Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Thu, 1 Oct 2026 20:53:14 -0500 Subject: [PATCH 07/13] PYTHON-6136 Replace checkout rollback flags with an undo bitmask Every failed checkout now restores its counters through one accounted-guarded replay, replacing the per-mutation flags. Behavior, lock regions, and telemetry ordering are unchanged. --- pymongo/asynchronous/pool.py | 119 +++++++++++++++-------------------- pymongo/synchronous/pool.py | 119 +++++++++++++++-------------------- 2 files changed, 102 insertions(+), 136 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 3c58cb924f..e4aea1131b 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -106,6 +106,13 @@ _IS_SYNC = False +# Flags recording the counters a checkout has incremented, so a failed +# checkout can restore them exactly once (PYTHON-6136). +_UNDO_OPERATION_COUNT = 1 +_UNDO_REQUESTS = 2 +_UNDO_SOCKETS = 4 +_UNDO_PENDING = 8 + class AsyncConnection(_ConnectionTelemetryInfo): """Store a connection with some metadata. @@ -969,8 +976,12 @@ async def _get_conn( "Attempted to check out a connection from closed connection pool" ) + # Every counter increment sets a flag in ``applied`` so a failed + # checkout can restore them exactly once (PYTHON-6136). + applied = 0 async with self.lock: self.operation_count += 1 + applied |= _UNDO_OPERATION_COUNT # Get a free socket or create one. if _csot.get_timeout(): @@ -980,8 +991,6 @@ async def _get_conn( else: deadline = None - # Roll back the size gate if a BaseException lands here (PYTHON-6136). - requests_incremented = False try: async with self.size_cond: self._raise_if_not_ready(checkout_started_time, emit_event=True) @@ -995,46 +1004,19 @@ async def _get_conn( self._raise_wait_queue_timeout(checkout_started_time) self._raise_if_not_ready(checkout_started_time, emit_event=True) self.requests += 1 - requests_incremented = True + applied |= _UNDO_REQUESTS except BaseException: - if requests_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - accounted = False - try: - async with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - accounted = True - self.size_cond.notify() - finally: - if not accounted: - async with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - self.size_cond.notify() - else: - # The gate never admitted; still undo the load increment above. - accounted = False - try: - async with self.size_cond: - self.operation_count -= 1 - accounted = True - finally: - if not accounted: - async with self.size_cond: - self.operation_count -= 1 + await self._restore_counters(applied) raise # We've now acquired the semaphore and must release it on error. conn = None - incremented = False emitted_event = False is_new_conn = False - pending_incremented = False try: async with self.lock: self.active_sockets += 1 - incremented = True + applied |= _UNDO_SOCKETS while conn is None: # CMAP: we MUST wait for either maxConnecting OR for a socket # to be checked back into the pool. @@ -1055,7 +1037,7 @@ async def _get_conn( conn = self.conns.popleft() except IndexError: self._pending += 1 - pending_incremented = True + applied |= _UNDO_PENDING if conn: # We got a socket from the pool if await self._perished(conn): conn = None @@ -1067,7 +1049,7 @@ async def _get_conn( finally: async with self._max_connecting_cond: self._pending -= 1 - pending_incremented = False + applied &= ~_UNDO_PENDING self._max_connecting_cond.notify() conn.active = True @@ -1079,43 +1061,13 @@ async def _get_conn( # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: try: - if pending_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - pending_accounted = False - try: - async with self._max_connecting_cond: - self._pending -= 1 - pending_accounted = True - self._max_connecting_cond.notify() - finally: - if not pending_accounted: - async with self._max_connecting_cond: - self._pending -= 1 - self._max_connecting_cond.notify() - - if conn: + if conn is not None: # We checked out a socket but authentication failed. await conn.close_conn(ConnectionClosedReason.ERROR) finally: - # Always roll back the size gate, even if the cleanups above - # were interrupted (PYTHON-6136). - accounted = False - try: - async with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - if incremented: - self.active_sockets -= 1 - accounted = True - self.size_cond.notify() - finally: - if not accounted: - async with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - if incremented: - self.active_sockets -= 1 - self.size_cond.notify() + # Restore the counters even if the cleanup above was + # interrupted (PYTHON-6136). + await self._restore_counters(applied) if not emitted_event: self._telemetry.checkout_failed( @@ -1127,6 +1079,37 @@ async def _get_conn( return conn + def _restore_applied(self, applied: int) -> None: + """Restore the counters flagged in ``applied``. Caller holds ``size_cond``.""" + if applied & _UNDO_OPERATION_COUNT: + self.operation_count -= 1 + if applied & _UNDO_REQUESTS: + self.requests -= 1 + if applied & _UNDO_SOCKETS: + self.active_sockets -= 1 + if applied & _UNDO_PENDING: + self._pending -= 1 + + async def _restore_counters(self, applied: int) -> None: + """Restore the counters a failed checkout incremented (PYTHON-6136). + + Gevent grants the re-acquire during unwind (PYTHON-6074). + """ + accounted = False + try: + async with self.size_cond: + self._restore_applied(applied) + accounted = True + if applied & _UNDO_REQUESTS: + # A pool slot was freed; wake the next waiter. + self.size_cond.notify() + finally: + if not accounted: + async with self.size_cond: + self._restore_applied(applied) + if applied & _UNDO_REQUESTS: + self.size_cond.notify() + def _checkin_apply( self, conn: AsyncConnection, txn: bool, cursor: bool, forked: bool ) -> tuple[Optional[str], bool, bool]: diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index 7dd5722df8..ff0e139b10 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -106,6 +106,13 @@ _IS_SYNC = True +# Flags recording the counters a checkout has incremented, so a failed +# checkout can restore them exactly once (PYTHON-6136). +_UNDO_OPERATION_COUNT = 1 +_UNDO_REQUESTS = 2 +_UNDO_SOCKETS = 4 +_UNDO_PENDING = 8 + class Connection(_ConnectionTelemetryInfo): """Store a connection with some metadata. @@ -965,8 +972,12 @@ def _get_conn( "Attempted to check out a connection from closed connection pool" ) + # Every counter increment sets a flag in ``applied`` so a failed + # checkout can restore them exactly once (PYTHON-6136). + applied = 0 with self.lock: self.operation_count += 1 + applied |= _UNDO_OPERATION_COUNT # Get a free socket or create one. if _csot.get_timeout(): @@ -976,8 +987,6 @@ def _get_conn( else: deadline = None - # Roll back the size gate if a BaseException lands here (PYTHON-6136). - requests_incremented = False try: with self.size_cond: self._raise_if_not_ready(checkout_started_time, emit_event=True) @@ -991,46 +1000,19 @@ def _get_conn( self._raise_wait_queue_timeout(checkout_started_time) self._raise_if_not_ready(checkout_started_time, emit_event=True) self.requests += 1 - requests_incremented = True + applied |= _UNDO_REQUESTS except BaseException: - if requests_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - accounted = False - try: - with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - accounted = True - self.size_cond.notify() - finally: - if not accounted: - with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - self.size_cond.notify() - else: - # The gate never admitted; still undo the load increment above. - accounted = False - try: - with self.size_cond: - self.operation_count -= 1 - accounted = True - finally: - if not accounted: - with self.size_cond: - self.operation_count -= 1 + self._restore_counters(applied) raise # We've now acquired the semaphore and must release it on error. conn = None - incremented = False emitted_event = False is_new_conn = False - pending_incremented = False try: with self.lock: self.active_sockets += 1 - incremented = True + applied |= _UNDO_SOCKETS while conn is None: # CMAP: we MUST wait for either maxConnecting OR for a socket # to be checked back into the pool. @@ -1051,7 +1033,7 @@ def _get_conn( conn = self.conns.popleft() except IndexError: self._pending += 1 - pending_incremented = True + applied |= _UNDO_PENDING if conn: # We got a socket from the pool if self._perished(conn): conn = None @@ -1063,7 +1045,7 @@ def _get_conn( finally: with self._max_connecting_cond: self._pending -= 1 - pending_incremented = False + applied &= ~_UNDO_PENDING self._max_connecting_cond.notify() conn.active = True @@ -1075,43 +1057,13 @@ def _get_conn( # Catch KeyboardInterrupt, CancelledError, etc. and cleanup. except BaseException: try: - if pending_incremented: - # Gevent grants the re-acquire during unwind (PYTHON-6074). - pending_accounted = False - try: - with self._max_connecting_cond: - self._pending -= 1 - pending_accounted = True - self._max_connecting_cond.notify() - finally: - if not pending_accounted: - with self._max_connecting_cond: - self._pending -= 1 - self._max_connecting_cond.notify() - - if conn: + if conn is not None: # We checked out a socket but authentication failed. conn.close_conn(ConnectionClosedReason.ERROR) finally: - # Always roll back the size gate, even if the cleanups above - # were interrupted (PYTHON-6136). - accounted = False - try: - with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - if incremented: - self.active_sockets -= 1 - accounted = True - self.size_cond.notify() - finally: - if not accounted: - with self.size_cond: - self.requests -= 1 - self.operation_count -= 1 - if incremented: - self.active_sockets -= 1 - self.size_cond.notify() + # Restore the counters even if the cleanup above was + # interrupted (PYTHON-6136). + self._restore_counters(applied) if not emitted_event: self._telemetry.checkout_failed( @@ -1123,6 +1075,37 @@ def _get_conn( return conn + def _restore_applied(self, applied: int) -> None: + """Restore the counters flagged in ``applied``. Caller holds ``size_cond``.""" + if applied & _UNDO_OPERATION_COUNT: + self.operation_count -= 1 + if applied & _UNDO_REQUESTS: + self.requests -= 1 + if applied & _UNDO_SOCKETS: + self.active_sockets -= 1 + if applied & _UNDO_PENDING: + self._pending -= 1 + + def _restore_counters(self, applied: int) -> None: + """Restore the counters a failed checkout incremented (PYTHON-6136). + + Gevent grants the re-acquire during unwind (PYTHON-6074). + """ + accounted = False + try: + with self.size_cond: + self._restore_applied(applied) + accounted = True + if applied & _UNDO_REQUESTS: + # A pool slot was freed; wake the next witer. + self.size_cond.notify() + finally: + if not accounted: + with self.size_cond: + self._restore_applied(applied) + if applied & _UNDO_REQUESTS: + self.size_cond.notify() + def _checkin_apply( self, conn: Connection, txn: bool, cursor: bool, forked: bool ) -> tuple[Optional[str], bool, bool]: From 85e2bc3458f7533d9ba7fd82375af532613f4a39 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 05:07:49 -0500 Subject: [PATCH 08/13] PYTHON-6136 Wake maxConnecting waiters when a failed checkout frees a pending slot --- pymongo/asynchronous/pool.py | 5 +++++ pymongo/synchronous/pool.py | 5 +++++ test/asynchronous/test_pooling.py | 12 +++++++++--- test/test_pooling.py | 12 +++++++++--- 4 files changed, 28 insertions(+), 6 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index e4aea1131b..34495dbef3 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1103,12 +1103,17 @@ async def _restore_counters(self, applied: int) -> None: if applied & _UNDO_REQUESTS: # A pool slot was freed; wake the next waiter. self.size_cond.notify() + if applied & _UNDO_PENDING: + # A maxConnecting slot was freed; wake the next waiter. + self._max_connecting_cond.notify() finally: if not accounted: async with self.size_cond: self._restore_applied(applied) if applied & _UNDO_REQUESTS: self.size_cond.notify() + if applied & _UNDO_PENDING: + self._max_connecting_cond.notify() def _checkin_apply( self, conn: AsyncConnection, txn: bool, cursor: bool, forked: bool diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index ff0e139b10..a68812cb9a 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1099,12 +1099,17 @@ def _restore_counters(self, applied: int) -> None: if applied & _UNDO_REQUESTS: # A pool slot was freed; wake the next witer. self.size_cond.notify() + if applied & _UNDO_PENDING: + # A maxConnecting slot was freed; wake the next witer. + self._max_connecting_cond.notify() finally: if not accounted: with self.size_cond: self._restore_applied(applied) if applied & _UNDO_REQUESTS: self.size_cond.notify() + if applied & _UNDO_PENDING: + self._max_connecting_cond.notify() def _checkin_apply( self, conn: Connection, txn: bool, cursor: bool, forked: bool diff --git a/test/asynchronous/test_pooling.py b/test/asynchronous/test_pooling.py index 948f0648f4..a69ca15111 100644 --- a/test/asynchronous/test_pooling.py +++ b/test/asynchronous/test_pooling.py @@ -361,13 +361,14 @@ async def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): async def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must - # not skip the size-gate rollback below it. + # not skip the counter restore, which must wake maxConnecting waiters. cx_pool = await self.create_pool(max_pool_size=1) class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): def __init__(self, lock): super().__init__(lock) self.enters = 0 + self.notifies = 0 async def __aenter__(self): self.enters += 1 @@ -381,10 +382,13 @@ async def __aexit__(self, *args): return await super().__aexit__(*args) def notify(self, n=1): - # The handler's pending-gate rollback is itself killed. + # The counter restore wakes maxConnecting waiters, then is + # itself killed. + self.notifies += 1 raise KeyboardInterrupt() - cx_pool._max_connecting_cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + cx_pool._max_connecting_cond = cond with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): with self.assertRaises(KeyboardInterrupt): @@ -395,6 +399,8 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) self.assertEqual(0, cx_pool.operation_count) + # The restore notified the maxConnecting gate before the kill. + self.assertEqual(1, cond.notifies) async def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. diff --git a/test/test_pooling.py b/test/test_pooling.py index fead75b6f8..debea79fe4 100644 --- a/test/test_pooling.py +++ b/test/test_pooling.py @@ -361,13 +361,14 @@ def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must - # not skip the size-gate rollback below it. + # not skip the counter restore, which must wake maxConnecting witers. cx_pool = self.create_pool(max_pool_size=1) class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): def __init__(self, lock): super().__init__(lock) self.enters = 0 + self.notifies = 0 def __enter__(self): self.enters += 1 @@ -381,10 +382,13 @@ def __exit__(self, *args): return super().__exit__(*args) def notify(self, n=1): - # The handler's pending-gate rollback is itself killed. + # The counter restore wakes maxConnecting witers, then is + # itself killed. + self.notifies += 1 raise KeyboardInterrupt() - cx_pool._max_connecting_cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + cond = _InterruptOnSecondEnter(cx_pool._max_connecting_cond._lock) + cx_pool._max_connecting_cond = cond with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): with self.assertRaises(KeyboardInterrupt): @@ -395,6 +399,8 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) self.assertEqual(0, cx_pool.operation_count) + # The restore notified the maxConnecting gate before the kill. + self.assertEqual(1, cond.notifies) def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. From e10f2dbe09215f5e0ee48cd286f67e68fb9855c2 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 05:23:28 -0500 Subject: [PATCH 09/13] PYTHON-6136 Retry interrupted gate notifications and fix comment mangling Track notification completion separately from counter restoration so a kill landing inside notify() retries the wake-up instead of stranding waiters. Reword comments containing 'waiter', which the synchro replacement table mangles ('aiter' is a substring). --- pymongo/asynchronous/pool.py | 24 +++++++++++++++--------- pymongo/synchronous/pool.py | 24 +++++++++++++++--------- test/asynchronous/test_pooling.py | 10 ++++++---- test/test_pooling.py | 10 ++++++---- 4 files changed, 42 insertions(+), 26 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 34495dbef3..a03ebc6f30 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1093,27 +1093,33 @@ def _restore_applied(self, applied: int) -> None: async def _restore_counters(self, applied: int) -> None: """Restore the counters a failed checkout incremented (PYTHON-6136). - Gevent grants the re-acquire during unwind (PYTHON-6074). + Gevent grants the re-acquire during unwind (PYTHON-6074). A kill can + also land inside notify() (a yield point), so notifications are + tracked and retried separately from the counter restore. """ accounted = False + notified = 0 try: async with self.size_cond: self._restore_applied(applied) accounted = True if applied & _UNDO_REQUESTS: - # A pool slot was freed; wake the next waiter. + # A pool slot was freed; wake the next waiting thread. self.size_cond.notify() + notified |= _UNDO_REQUESTS if applied & _UNDO_PENDING: - # A maxConnecting slot was freed; wake the next waiter. + # A maxConnecting slot was freed; wake the next waiting thread. self._max_connecting_cond.notify() + notified |= _UNDO_PENDING finally: - if not accounted: - async with self.size_cond: + async with self.size_cond: + if not accounted: self._restore_applied(applied) - if applied & _UNDO_REQUESTS: - self.size_cond.notify() - if applied & _UNDO_PENDING: - self._max_connecting_cond.notify() + missing = applied & ~notified + if missing & _UNDO_REQUESTS: + self.size_cond.notify() + if missing & _UNDO_PENDING: + self._max_connecting_cond.notify() def _checkin_apply( self, conn: AsyncConnection, txn: bool, cursor: bool, forked: bool diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index a68812cb9a..1af592cf5a 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1089,27 +1089,33 @@ def _restore_applied(self, applied: int) -> None: def _restore_counters(self, applied: int) -> None: """Restore the counters a failed checkout incremented (PYTHON-6136). - Gevent grants the re-acquire during unwind (PYTHON-6074). + Gevent grants the re-acquire during unwind (PYTHON-6074). A kill can + also land inside notify() (a yield point), so notifications are + tracked and retried separately from the counter restore. """ accounted = False + notified = 0 try: with self.size_cond: self._restore_applied(applied) accounted = True if applied & _UNDO_REQUESTS: - # A pool slot was freed; wake the next witer. + # A pool slot was freed; wake the next waiting thread. self.size_cond.notify() + notified |= _UNDO_REQUESTS if applied & _UNDO_PENDING: - # A maxConnecting slot was freed; wake the next witer. + # A maxConnecting slot was freed; wake the next waiting thread. self._max_connecting_cond.notify() + notified |= _UNDO_PENDING finally: - if not accounted: - with self.size_cond: + with self.size_cond: + if not accounted: self._restore_applied(applied) - if applied & _UNDO_REQUESTS: - self.size_cond.notify() - if applied & _UNDO_PENDING: - self._max_connecting_cond.notify() + missing = applied & ~notified + if missing & _UNDO_REQUESTS: + self.size_cond.notify() + if missing & _UNDO_PENDING: + self._max_connecting_cond.notify() def _checkin_apply( self, conn: Connection, txn: bool, cursor: bool, forked: bool diff --git a/test/asynchronous/test_pooling.py b/test/asynchronous/test_pooling.py index a69ca15111..dfad69966b 100644 --- a/test/asynchronous/test_pooling.py +++ b/test/asynchronous/test_pooling.py @@ -361,7 +361,8 @@ async def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): async def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must - # not skip the counter restore, which must wake maxConnecting waiters. + # not skip the counter restore, which must wake threads waiting at + # the maxConnecting gate. cx_pool = await self.create_pool(max_pool_size=1) class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): @@ -382,7 +383,7 @@ async def __aexit__(self, *args): return await super().__aexit__(*args) def notify(self, n=1): - # The counter restore wakes maxConnecting waiters, then is + # The counter restore wakes the maxConnecting gate, then is # itself killed. self.notifies += 1 raise KeyboardInterrupt() @@ -399,8 +400,9 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) self.assertEqual(0, cx_pool.operation_count) - # The restore notified the maxConnecting gate before the kill. - self.assertEqual(1, cond.notifies) + # The restore notified the maxConnecting gate; the notify interrupted + # by the kill was retried. + self.assertEqual(2, cond.notifies) async def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. diff --git a/test/test_pooling.py b/test/test_pooling.py index debea79fe4..585882791d 100644 --- a/test/test_pooling.py +++ b/test/test_pooling.py @@ -361,7 +361,8 @@ def test_checkout_error_accounting_on_connect_keyboard_interrupt(self): def test_checkout_error_accounting_on_kill_during_pending_cleanup(self): # PYTHON-6136: an interruption during the pending-gate cleanup must - # not skip the counter restore, which must wake maxConnecting witers. + # not skip the counter restore, which must wake threads waiting at + # the maxConnecting gate. cx_pool = self.create_pool(max_pool_size=1) class _InterruptOnSecondEnter(type(cx_pool._max_connecting_cond)): @@ -382,7 +383,7 @@ def __exit__(self, *args): return super().__exit__(*args) def notify(self, n=1): - # The counter restore wakes maxConnecting witers, then is + # The counter restore wakes the maxConnecting gate, then is # itself killed. self.notifies += 1 raise KeyboardInterrupt() @@ -399,8 +400,9 @@ def notify(self, n=1): self.assertEqual(0, cx_pool.active_sockets) self.assertEqual(0, cx_pool._pending) self.assertEqual(0, cx_pool.operation_count) - # The restore notified the maxConnecting gate before the kill. - self.assertEqual(1, cond.notifies) + # The restore notified the maxConnecting gate; the notify interrupted + # by the kill was retried. + self.assertEqual(2, cond.notifies) def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. From 8bd3fa5d2a6f9ee9b83d417669f65d4f159231ff Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 06:00:25 -0500 Subject: [PATCH 10/13] PYTHON-6136 Retry the maxConnecting cleanup notify after a kill A kill landing inside the connect cleanup's notify() (a gevent yield point) could strand a checkout waiting at the maxConnecting gate. Retry the notify with the accounted idiom so the wake-up survives. --- pymongo/asynchronous/pool.py | 11 ++++++++++- pymongo/synchronous/pool.py | 11 ++++++++++- test/asynchronous/test_pooling.py | 31 +++++++++++++++++++++++++++++++ test/test_pooling.py | 31 +++++++++++++++++++++++++++++++ 4 files changed, 82 insertions(+), 2 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index a03ebc6f30..737d4068a2 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1050,7 +1050,16 @@ async def _get_conn( async with self._max_connecting_cond: self._pending -= 1 applied &= ~_UNDO_PENDING - self._max_connecting_cond.notify() + notified = False + try: + self._max_connecting_cond.notify() + notified = True + finally: + if not notified: + # A kill landed inside notify() (a gevent + # yield point); retry so a waiter is not + # stranded (PYTHON-6136). + self._max_connecting_cond.notify() conn.active = True # connect() already adds cancel_context for new connections; only add diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index 1af592cf5a..7a304086dc 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1046,7 +1046,16 @@ def _get_conn( with self._max_connecting_cond: self._pending -= 1 applied &= ~_UNDO_PENDING - self._max_connecting_cond.notify() + notified = False + try: + self._max_connecting_cond.notify() + notified = True + finally: + if not notified: + # A kill landed inside notify() (a gevent + # yield point); retry so a witer is not + # stranded (PYTHON-6136). + self._max_connecting_cond.notify() conn.active = True # connect() already adds cancel_context for new connections; only add diff --git a/test/asynchronous/test_pooling.py b/test/asynchronous/test_pooling.py index dfad69966b..25da5134b9 100644 --- a/test/asynchronous/test_pooling.py +++ b/test/asynchronous/test_pooling.py @@ -404,6 +404,37 @@ def notify(self, n=1): # by the kill was retried. self.assertEqual(2, cond.notifies) + async def test_checkout_error_accounting_on_kill_during_pending_notify(self): + # PYTHON-6136: a kill landing inside the cleanup notify must not + # strand a checkout waiting at the maxConnecting gate. + cx_pool = await self.create_pool(max_pool_size=1) + + class _InterruptOnFirstNotify(type(cx_pool._max_connecting_cond)): + def __init__(self, lock): + super().__init__(lock) + self.notifies = 0 + + def notify(self, n=1): + self.notifies += 1 + if self.notifies == 1: + # Simulate a kill delivered inside notify(). + raise KeyboardInterrupt() + + cond = _InterruptOnFirstNotify(cx_pool._max_connecting_cond._lock) + cx_pool._max_connecting_cond = cond + + with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): + with self.assertRaises(KeyboardInterrupt): + async with cx_pool.checkout(): + pass + + # The cleanup notify was interrupted, then retried. + self.assertEqual(2, cond.notifies) + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) + async def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. cx_pool = await self.create_pool() diff --git a/test/test_pooling.py b/test/test_pooling.py index 585882791d..3a81bc796f 100644 --- a/test/test_pooling.py +++ b/test/test_pooling.py @@ -404,6 +404,37 @@ def notify(self, n=1): # by the kill was retried. self.assertEqual(2, cond.notifies) + def test_checkout_error_accounting_on_kill_during_pending_notify(self): + # PYTHON-6136: a kill landing inside the cleanup notify must not + # strand a checkout waiting at the maxConnecting gate. + cx_pool = self.create_pool(max_pool_size=1) + + class _InterruptOnFirstNotify(type(cx_pool._max_connecting_cond)): + def __init__(self, lock): + super().__init__(lock) + self.notifies = 0 + + def notify(self, n=1): + self.notifies += 1 + if self.notifies == 1: + # Simulate a kill delivered inside notify(). + raise KeyboardInterrupt() + + cond = _InterruptOnFirstNotify(cx_pool._max_connecting_cond._lock) + cx_pool._max_connecting_cond = cond + + with patch.object(cx_pool, "connect", side_effect=asyncio.CancelledError()): + with self.assertRaises(KeyboardInterrupt): + with cx_pool.checkout(): + pass + + # The cleanup notify was interrupted, then retried. + self.assertEqual(2, cond.notifies) + self.assertEqual(0, cx_pool.requests) + self.assertEqual(0, cx_pool.active_sockets) + self.assertEqual(0, cx_pool._pending) + self.assertEqual(0, cx_pool.operation_count) + def test_pool_removes_closed_socket(self): # Test that Pool removes explicitly closed socket. cx_pool = self.create_pool() From 52f72bedca07e5ae72bd3f47246eaefa882c0bfc Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 06:13:16 -0500 Subject: [PATCH 11/13] PYTHON-6136 Reword retry comment to avoid the synchro mangler --- pymongo/asynchronous/pool.py | 4 ++-- pymongo/synchronous/pool.py | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index 737d4068a2..e75013b2a6 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1057,8 +1057,8 @@ async def _get_conn( finally: if not notified: # A kill landed inside notify() (a gevent - # yield point); retry so a waiter is not - # stranded (PYTHON-6136). + # yield point); retry so a waiting + # checkout is not stranded (PYTHON-6136). self._max_connecting_cond.notify() conn.active = True diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index 7a304086dc..f629d342f5 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1053,8 +1053,8 @@ def _get_conn( finally: if not notified: # A kill landed inside notify() (a gevent - # yield point); retry so a witer is not - # stranded (PYTHON-6136). + # yield point); retry so a waiting + # checkout is not stranded (PYTHON-6136). self._max_connecting_cond.notify() conn.active = True From 72423b5da112913e3ddba75990de9fc6ba6f3da1 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 06:26:43 -0500 Subject: [PATCH 12/13] PYTHON-6136 Assert operation_count drains in the stress test --- test/asynchronous/test_client.py | 5 ++++- test/test_client.py | 5 ++++- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/test/asynchronous/test_client.py b/test/asynchronous/test_client.py index 62133b814f..563d29cbad 100644 --- a/test/asynchronous/test_client.py +++ b/test/asynchronous/test_client.py @@ -2866,11 +2866,14 @@ def reaper(): client.close() except Timeout: pass - # Post-settle: the counters must have fully drained. + # Post-settle: the counters must have fully drained. Note close() + # does not zero operation_count (only a fork does), so a leaked + # increment survives and is caught here. time.sleep(1.0) self.assertEqual(pool.requests, 0) self.assertEqual(pool._pending, 0) self.assertEqual(pool.active_sockets, 0) + self.assertEqual(pool.operation_count, 0) class TestClientLazyConnect(AsyncIntegrationTest): diff --git a/test/test_client.py b/test/test_client.py index d6e004fa75..0c8d0e7595 100644 --- a/test/test_client.py +++ b/test/test_client.py @@ -2817,11 +2817,14 @@ def reaper(): client.close() except Timeout: pass - # Post-settle: the counters must have fully drained. + # Post-settle: the counters must have fully drained. Note close() + # does not zero operation_count (only a fork does), so a leaked + # increment survives and is caught here. time.sleep(1.0) self.assertEqual(pool.requests, 0) self.assertEqual(pool._pending, 0) self.assertEqual(pool.active_sockets, 0) + self.assertEqual(pool.operation_count, 0) class TestClientLazyConnect(IntegrationTest): From 1b577fb984257481f485d2d0979e7912ed4e1b76 Mon Sep 17 00:00:00 2001 From: Steven Silvester Date: Fri, 2 Oct 2026 13:01:47 -0500 Subject: [PATCH 13/13] PYTHON-6136 Document why _restore_counters always reacquires the lock --- pymongo/asynchronous/pool.py | 3 +++ pymongo/synchronous/pool.py | 3 +++ 2 files changed, 6 insertions(+) diff --git a/pymongo/asynchronous/pool.py b/pymongo/asynchronous/pool.py index e75013b2a6..b7f7a1f6a0 100644 --- a/pymongo/asynchronous/pool.py +++ b/pymongo/asynchronous/pool.py @@ -1121,6 +1121,9 @@ async def _restore_counters(self, applied: int) -> None: self._max_connecting_cond.notify() notified |= _UNDO_PENDING finally: + # Always reacquired: `applied` keeps restore-only flags + # `notified` never tracks, so skipping needs a mask synced + # to the notify sites; uncontended acquires don't yield (PYTHON-6136). async with self.size_cond: if not accounted: self._restore_applied(applied) diff --git a/pymongo/synchronous/pool.py b/pymongo/synchronous/pool.py index f629d342f5..dd562ac68e 100644 --- a/pymongo/synchronous/pool.py +++ b/pymongo/synchronous/pool.py @@ -1117,6 +1117,9 @@ def _restore_counters(self, applied: int) -> None: self._max_connecting_cond.notify() notified |= _UNDO_PENDING finally: + # Always reacquired: `applied` keeps restore-only flags + # `notified` never tracks, so skipping needs a mask synced + # to the notify sites; uncontended acquires don't yield (PYTHON-6136). with self.size_cond: if not accounted: self._restore_applied(applied)