Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions doc/changelog.rst
Original file line number Diff line number Diff line change
@@ -1,6 +1,18 @@
Changelog
=========

Changes in Version 4.19.0 (2026/XX/XX)
Comment thread
blink1073 marked this conversation as resolved.
--------------------------------------

Bug fixes
.........

- Fixed a bug where the synchronous client could permanently deadlock under
gevent when a greenlet was killed while checking a connection back into
the pool (`PYTHON-6074`_).

.. _PYTHON-6074: https://jira.mongodb.org/browse/PYTHON-6074

Changes in Version 4.18.0 (2026/09/03)
--------------------------------------

Expand Down
122 changes: 83 additions & 39 deletions pymongo/asynchronous/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -1044,11 +1044,24 @@ async def _get_conn(
if conn:
# We checked out a socket but authentication failed.
await conn.close_conn(ConnectionClosedReason.ERROR)
async with self.size_cond:
self.requests -= 1
if incremented:
self.active_sockets -= 1
self.size_cond.notify()
# Re-apply the accounting if a GreenletExit interrupts
# during the size_cond acquisition; during unwind gevent
# lets the re-acquire complete (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()
Comment thread
blink1073 marked this conversation as resolved.
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(
Expand All @@ -1060,6 +1073,40 @@ async def _get_conn(

return conn

def _checkin_apply(
self, conn: AsyncConnection, txn: bool, cursor: bool, forked: bool
) -> tuple[Optional[str], bool, bool]:
"""Apply checkin accounting; caller holds ``size_cond``.

No cooperative I/O, so safe while a gevent greenlet unwinds. Returns
``(close_conn_reason, emit_closed, appended)`` for outside the lock.
"""
self.active_contexts.discard(conn.cancel_context)
if txn:
self.ntxns -= 1
elif cursor:
self.ncursors -= 1
self.requests -= 1
self.active_sockets -= 1
self.operation_count -= 1
close_conn_reason: Optional[str] = None
emit_closed = False
appended = False
if not forked:
if self.closed:
close_conn_reason = ConnectionClosedReason.POOL_CLOSED
elif conn.closed:
# CMAP requires the closed event be emitted after the check in.
emit_closed = True
elif self.stale_generation(conn.generation, conn.service_id):
close_conn_reason = ConnectionClosedReason.STALE
else:
conn.update_last_checkin_time()
conn.update_is_writable(bool(self.is_writable))
self.conns.appendleft(conn)
appended = True
return close_conn_reason, emit_closed, appended

async def checkin(self, conn: AsyncConnection) -> None:
"""Return the connection to the pool, or if it's closed discard it.

Expand All @@ -1071,44 +1118,41 @@ async def checkin(self, conn: AsyncConnection) -> None:
conn.pinned_txn = False
conn.pinned_cursor = False
self._pinned_sockets.discard(conn)
async with self.lock:
self.active_contexts.discard(conn.cancel_context)
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).
close_conn_reason: Optional[str] = None
emit_closed = False
accounted = False
try:
async with self.size_cond:
close_conn_reason, emit_closed, appended = self._checkin_apply(
conn, txn, cursor, forked
)
accounted = True
if appended:
# Notify any threads waiting to create a connection.
self._max_connecting_cond.notify()
self.size_cond.notify()
finally:
if not accounted:
async with self.size_cond:
close_conn_reason, emit_closed, appended = self._checkin_apply(
conn, txn, cursor, forked
)
if appended:
self._max_connecting_cond.notify()
self.size_cond.notify()
telemetry = self._telemetry
if telemetry._should_publish or (telemetry._log and _is_debug_enabled(_CONNECTION_LOGGER)):
telemetry.checked_in(conn.id)
if self.pid != os.getpid():
if emit_closed:
telemetry.connection_closed(conn.id, ConnectionClosedReason.ERROR)
if forked:
await self.reset_without_pause()
else:
if self.closed:
await conn.close_conn(ConnectionClosedReason.POOL_CLOSED)
elif conn.closed:
# CMAP requires the closed event be emitted after the check in.
self._telemetry.connection_closed(conn.id, ConnectionClosedReason.ERROR)
else:
close_conn = False
async with self.lock:
# Hold the lock to ensure this section does not race with
# Pool.reset().
if self.stale_generation(conn.generation, conn.service_id):
close_conn = True
else:
conn.update_last_checkin_time()
conn.update_is_writable(bool(self.is_writable))
self.conns.appendleft(conn)
# Notify any threads waiting to create a connection.
self._max_connecting_cond.notify()
if close_conn:
await conn.close_conn(ConnectionClosedReason.STALE)

async with self.size_cond:
if txn:
self.ntxns -= 1
elif cursor:
self.ncursors -= 1
self.requests -= 1
self.active_sockets -= 1
self.operation_count -= 1
self.size_cond.notify()
elif close_conn_reason is not None:
await conn.close_conn(close_conn_reason)

async def _perished(self, conn: AsyncConnection) -> bool:
"""Return True and close the connection if it is "perished".
Expand Down
122 changes: 83 additions & 39 deletions pymongo/synchronous/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -1040,11 +1040,24 @@ def _get_conn(
if conn:
# We checked out a socket but authentication failed.
conn.close_conn(ConnectionClosedReason.ERROR)
with self.size_cond:
self.requests -= 1
if incremented:
self.active_sockets -= 1
self.size_cond.notify()
# Re-apply the accounting if a GreenletExit interrupts
# during the size_cond acquisition; during unwind gevent
# lets the re-acquire complete (PYTHON-6074).
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(
Expand All @@ -1056,6 +1069,40 @@ def _get_conn(

return conn

def _checkin_apply(
self, conn: Connection, txn: bool, cursor: bool, forked: bool
) -> tuple[Optional[str], bool, bool]:
"""Apply checkin accounting; caller holds ``size_cond``.

No cooperative I/O, so safe while a gevent greenlet unwinds. Returns
``(close_conn_reason, emit_closed, appended)`` for outside the lock.
"""
self.active_contexts.discard(conn.cancel_context)
if txn:
self.ntxns -= 1
elif cursor:
self.ncursors -= 1
self.requests -= 1
self.active_sockets -= 1
self.operation_count -= 1
close_conn_reason: Optional[str] = None
emit_closed = False
appended = False
if not forked:
if self.closed:
close_conn_reason = ConnectionClosedReason.POOL_CLOSED
elif conn.closed:
# CMAP requires the closed event be emitted after the check in.
emit_closed = True
elif self.stale_generation(conn.generation, conn.service_id):
close_conn_reason = ConnectionClosedReason.STALE
else:
conn.update_last_checkin_time()
conn.update_is_writable(bool(self.is_writable))
self.conns.appendleft(conn)
appended = True
return close_conn_reason, emit_closed, appended

def checkin(self, conn: Connection) -> None:
"""Return the connection to the pool, or if it's closed discard it.

Expand All @@ -1067,44 +1114,41 @@ def checkin(self, conn: Connection) -> None:
conn.pinned_txn = False
conn.pinned_cursor = False
self._pinned_sockets.discard(conn)
with self.lock:
self.active_contexts.discard(conn.cancel_context)
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).
close_conn_reason: Optional[str] = None
emit_closed = False
accounted = False
try:
with self.size_cond:
close_conn_reason, emit_closed, appended = self._checkin_apply(
conn, txn, cursor, forked
)
accounted = True
if appended:
# Notify any threads waiting to create a connection.
self._max_connecting_cond.notify()
self.size_cond.notify()
finally:
if not accounted:
with self.size_cond:
close_conn_reason, emit_closed, appended = self._checkin_apply(
conn, txn, cursor, forked
)
if appended:
self._max_connecting_cond.notify()
self.size_cond.notify()
telemetry = self._telemetry
if telemetry._should_publish or (telemetry._log and _is_debug_enabled(_CONNECTION_LOGGER)):
telemetry.checked_in(conn.id)
if self.pid != os.getpid():
if emit_closed:
telemetry.connection_closed(conn.id, ConnectionClosedReason.ERROR)
if forked:
self.reset_without_pause()
else:
if self.closed:
conn.close_conn(ConnectionClosedReason.POOL_CLOSED)
elif conn.closed:
# CMAP requires the closed event be emitted after the check in.
self._telemetry.connection_closed(conn.id, ConnectionClosedReason.ERROR)
else:
close_conn = False
with self.lock:
# Hold the lock to ensure this section does not race with
# Pool.reset().
if self.stale_generation(conn.generation, conn.service_id):
close_conn = True
else:
conn.update_last_checkin_time()
conn.update_is_writable(bool(self.is_writable))
self.conns.appendleft(conn)
# Notify any threads waiting to create a connection.
self._max_connecting_cond.notify()
if close_conn:
conn.close_conn(ConnectionClosedReason.STALE)

with self.size_cond:
if txn:
self.ntxns -= 1
elif cursor:
self.ncursors -= 1
self.requests -= 1
self.active_sockets -= 1
self.operation_count -= 1
self.size_cond.notify()
elif close_conn_reason is not None:
conn.close_conn(close_conn_reason)

def _perished(self, conn: Connection) -> bool:
"""Return True and close the connection if it is "perished".
Expand Down
Loading
Loading