Repository navigation
PYTHON-6136 Make pool checkout accounting transactional #3089
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
d16d9b2
04dccc5
2191a2f
a909572
dbbc615
40762fa
397e7a4
85e2bc3
e10f2db
8bd3fa5
52f72be
72423b5
1b577fb
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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,28 +991,32 @@ 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) | ||
| 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 | ||
| applied |= _UNDO_REQUESTS | ||
| except BaseException: | ||
| 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 | ||
| 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. | ||
|
|
@@ -1022,6 +1037,7 @@ async def _get_conn( | |
| conn = self.conns.popleft() | ||
| except IndexError: | ||
| self._pending += 1 | ||
| applied |= _UNDO_PENDING | ||
| if conn: # We got a socket from the pool | ||
| if await self._perished(conn): | ||
| conn = None | ||
|
|
@@ -1033,7 +1049,17 @@ async def _get_conn( | |
| finally: | ||
| async with self._max_connecting_cond: | ||
| self._pending -= 1 | ||
| self._max_connecting_cond.notify() | ||
| applied &= ~_UNDO_PENDING | ||
| 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 waiting | ||
| # checkout is not stranded (PYTHON-6136). | ||
| self._max_connecting_cond.notify() | ||
|
|
||
| conn.active = True | ||
| # connect() already adds cancel_context for new connections; only add | ||
|
|
@@ -1043,27 +1069,14 @@ 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 conn is not None: | ||
| # We checked out a socket but authentication failed. | ||
| await conn.close_conn(ConnectionClosedReason.ERROR) | ||
| finally: | ||
| if not accounted: | ||
| async with self.size_cond: | ||
| self.requests -= 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( | ||
|
|
@@ -1075,6 +1088,51 @@ 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). 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 waiting thread. | ||
| self.size_cond.notify() | ||
| notified |= _UNDO_REQUESTS | ||
| if applied & _UNDO_PENDING: | ||
| # A maxConnecting slot was freed; wake the next waiting thread. | ||
| 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: | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Looks like this locks whether it needs to or not, how about something like this instead:
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I verified that the current structure is correct and added a comment. |
||
| if not accounted: | ||
| self._restore_applied(applied) | ||
|
blink1073 marked this conversation as resolved.
|
||
| 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 | ||
| ) -> tuple[Optional[str], bool, bool]: | ||
|
|
@@ -1121,9 +1179,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 | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.