Skip to content
Closed
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
106 changes: 78 additions & 28 deletions pymongo/asynchronous/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -980,24 +980,57 @@ 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
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
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
Expand All @@ -1022,6 +1055,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
Expand All @@ -1033,6 +1067,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
Expand All @@ -1043,27 +1078,44 @@ async def _get_conn(
self.active_contexts.add(conn.cancel_context)
# Catch KeyboardInterrupt, CancelledError, etc. and cleanup.
except BaseException:
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).
accounted = False
try:
async with self.size_cond:
self.requests -= 1
if incremented:
self.active_sockets -= 1
accounted = True
self.size_cond.notify()
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)
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
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()

if not emitted_event:
self._telemetry.checkout_failed(
Expand Down Expand Up @@ -1121,9 +1173,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
Expand Down
106 changes: 78 additions & 28 deletions pymongo/synchronous/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -976,24 +976,57 @@ 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
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
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
Expand All @@ -1018,6 +1051,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
Expand All @@ -1029,6 +1063,7 @@ def _get_conn(
finally:
with self._max_connecting_cond:
self._pending -= 1
pending_incremented = False
self._max_connecting_cond.notify()

conn.active = True
Expand All @@ -1039,27 +1074,44 @@ def _get_conn(
self.active_contexts.add(conn.cancel_context)
# Catch KeyboardInterrupt, CancelledError, etc. and cleanup.
except BaseException:
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).
accounted = False
try:
with self.size_cond:
self.requests -= 1
if incremented:
self.active_sockets -= 1
accounted = True
self.size_cond.notify()
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)
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
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()

if not emitted_event:
self._telemetry.checkout_failed(
Expand Down Expand Up @@ -1117,9 +1169,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
Expand Down
Loading
Loading