Skip to content

Commit 76fe8dc

Browse files
Route release/drain close sites through _close_best_effort with DEBUG log
The seven release/drain close sites used `contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS)` to absorb a half- torn-down transport's `OSError` / `DqliteConnectionError` / `ProtocolError` / `OperationalError` / `InterfaceError` from `conn.close()` — but `contextlib.suppress` has no logging surface, so the suppressed exception left no trace at any verbosity. The sibling `_initialize_close_unqueued` used `try/except` + `logger.debug(..., exc_info=True)`; the asymmetry was a residual diagnostic gap after the prior `_POOL_CLEANUP_EXCEPTIONS` narrowing pass. Extract `_close_best_effort(conn, site)` matching the initialize- helper shape (shielded close + DEBUG log with `exc_info=True`) and route all seven sites through it. Each site passes a stable identifier so operators can locate the emitting branch in `pool.<site>: close error` log records. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent b6d5207 commit 76fe8dc

2 files changed

Lines changed: 197 additions & 40 deletions

File tree

‎src/dqliteclient/pool.py‎

Lines changed: 74 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -1000,6 +1000,40 @@ async def _initialize_close_unqueued(
10001000
exc_info=True,
10011001
)
10021002

1003+
async def _close_best_effort(self, conn: DqliteConnection, site: str) -> None:
1004+
"""Best-effort shielded close that logs ``_POOL_CLEANUP_EXCEPTIONS``
1005+
at DEBUG with ``exc_info=True`` rather than silently swallowing
1006+
them.
1007+
1008+
``site`` is a short identifier appended to the log key
1009+
(``"pool.<site>: close error"``) so operators can locate which
1010+
call site emitted the record. Mirrors the
1011+
``_initialize_close_unqueued`` discipline — both helpers now log
1012+
through ``exc_info=True`` so the suppressed-exception chain is
1013+
recoverable at DEBUG verbosity instead of being silently
1014+
dropped by ``contextlib.suppress`` (which has no logging
1015+
surface). The asymmetry between the initialize-time helper
1016+
(logged) and the release-time sites (silent) was a residual
1017+
diagnostic gap after the prior ``_POOL_CLEANUP_EXCEPTIONS``
1018+
narrowing pass; routing all release/drain close sites through
1019+
this helper closes it.
1020+
1021+
Cancel-safety: ``asyncio.shield`` runs the close to completion
1022+
even if an outer cancel is delivered during the await; the
1023+
CancelledError still propagates to the caller after the close
1024+
finishes. The CancelledError arm is intentionally NOT caught
1025+
here — preserves the pre-existing behaviour of the sites this
1026+
helper replaces (``contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS)``
1027+
suppressed only the cleanup-exception tuple; CancelledError
1028+
propagated to the surrounding frame). Sites that need to keep
1029+
running through cancel (like ``_initialize_close_unqueued``'s
1030+
per-conn walk) handle that arm explicitly.
1031+
"""
1032+
try:
1033+
await asyncio.shield(conn.close())
1034+
except _POOL_CLEANUP_EXCEPTIONS:
1035+
logger.debug("pool.%s: close error", site, exc_info=True)
1036+
10031037
async def _create_connection(self) -> DqliteConnection:
10041038
"""Create a new connection to the leader.
10051039
@@ -1060,20 +1094,22 @@ async def _put_back_or_release_late_winner(self, conn: DqliteConnection) -> None
10601094
# had happened. Mirrors the ``_drain_idle`` discipline.
10611095
conn._pool_released = False
10621096
try:
1063-
# Suppress the canonical ``_POOL_CLEANUP_EXCEPTIONS``
1064-
# tuple (OSError + DqliteConnectionError + ProtocolError
1065-
# + OperationalError + InterfaceError), not just
1066-
# ``OSError``. A late-winner conn whose transport was
1067-
# broken by a peer reset between checkout and close
1068-
# raises one of the broader categories from
1069-
# ``close()``; pre-fix that exception escaped, the
1097+
# Route through ``_close_best_effort`` so the canonical
1098+
# ``_POOL_CLEANUP_EXCEPTIONS`` tuple (OSError +
1099+
# DqliteConnectionError + ProtocolError +
1100+
# OperationalError + InterfaceError) is absorbed at
1101+
# DEBUG with ``exc_info=True`` (matches the sibling
1102+
# ``_initialize_close_unqueued`` discipline) rather
1103+
# than silently dropped. A late-winner conn whose
1104+
# transport was broken by a peer reset between
1105+
# checkout and close raises one of the broader
1106+
# categories from ``close()``; the helper logs them
1107+
# and proceeds so ``_release_reservation`` still runs
1108+
# below — pre-fix-narrow those exceptions escaped, the
10701109
# ``finally`` restored the flag, and
1071-
# ``_release_reservation`` was NEVER reached — the
1072-
# reservation slot leaked permanently. Matches the
1073-
# canonical ``_release`` discipline at the sibling
1074-
# close-and-track path.
1075-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
1076-
await asyncio.shield(conn.close())
1110+
# ``_release_reservation`` was NEVER reached, leaking
1111+
# the reservation slot permanently.
1112+
await self._close_best_effort(conn, "acquire-late-winner-closed")
10771113
finally:
10781114
# Restore the flag for contract symmetry with
10791115
# ``_drain_idle``: any stale-reference second close()
@@ -1092,18 +1128,19 @@ async def _put_back_or_release_late_winner(self, conn: DqliteConnection) -> None
10921128
# the reference would leak a live reader task and a
10931129
# socket. Close explicitly and adjust the reservation
10941130
# count so the pool shrinks cleanly instead of leaking.
1095-
# Suppress the canonical ``_POOL_CLEANUP_EXCEPTIONS``
1096-
# tuple — a half-torn-down transport raises
1097-
# ``DqliteConnectionError`` / ``InterfaceError`` /
1098-
# ``ProtocolError`` / ``OperationalError`` in addition to
1099-
# plain ``OSError``; pre-fix the broader categories
1100-
# escaped and the slot leaked. Flip ``_pool_released``
1101-
# to ``False`` first so the close actually runs (see
1102-
# the closed-pool arm above for the rationale).
1131+
# Route through ``_close_best_effort`` so the canonical
1132+
# ``_POOL_CLEANUP_EXCEPTIONS`` tuple — a half-torn-down
1133+
# transport raises ``DqliteConnectionError`` /
1134+
# ``InterfaceError`` / ``ProtocolError`` /
1135+
# ``OperationalError`` in addition to plain ``OSError`` —
1136+
# is absorbed at DEBUG with ``exc_info=True``; pre-fix-
1137+
# narrow the broader categories escaped and the slot
1138+
# leaked. Flip ``_pool_released`` to ``False`` first so
1139+
# the close actually runs (see the closed-pool arm above
1140+
# for the rationale).
11031141
conn._pool_released = False
11041142
try:
1105-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
1106-
await asyncio.shield(conn.close())
1143+
await self._close_best_effort(conn, "acquire-late-winner-queuefull")
11071144
finally:
11081145
conn._pool_released = True
11091146
# Route through the helper so the counter stays
@@ -1735,8 +1772,7 @@ async def acquire(self) -> AsyncIterator[DqliteConnection]:
17351772
conn._pool_released = False
17361773
try:
17371774
try:
1738-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
1739-
await asyncio.shield(conn.close())
1775+
await self._close_best_effort(conn, "acquire-drain-stale-conn")
17401776
except RuntimeError:
17411777
# "Event loop is closed" during racing
17421778
# ``engine.dispose()``, or cross-loop
@@ -2173,21 +2209,21 @@ async def _release(self, conn: DqliteConnection) -> None:
21732209
returned_to_queue = False
21742210
try:
21752211
if self._closed:
2176-
# Shield the close so an outer cancel mid-cleanup
2177-
# (e.g. ``asyncio.timeout`` around ``engine.dispose()``)
2178-
# cannot abort the close mid-``wait_closed`` and orphan
2179-
# the StreamReader task. ``_POOL_CLEANUP_EXCEPTIONS``
2180-
# does not include ``CancelledError``; the shield is
2181-
# the load-bearing guard. Symmetric with every other
2182-
# close site in this module.
2183-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
2184-
await asyncio.shield(conn.close())
2212+
# Shield (via ``_close_best_effort``) so an outer cancel
2213+
# mid-cleanup (e.g. ``asyncio.timeout`` around
2214+
# ``engine.dispose()``) cannot abort the close mid-
2215+
# ``wait_closed`` and orphan the StreamReader task.
2216+
# ``_POOL_CLEANUP_EXCEPTIONS`` does not include
2217+
# ``CancelledError``; the shield is the load-bearing
2218+
# guard. ``_POOL_CLEANUP_EXCEPTIONS`` errors are logged
2219+
# at DEBUG with ``exc_info=True`` (mirrors
2220+
# ``_initialize_close_unqueued``).
2221+
await self._close_best_effort(conn, "release-closed")
21852222
conn._pool_released = True
21862223
return
21872224

21882225
if not await self._reset_connection(conn):
2189-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
2190-
await asyncio.shield(conn.close())
2226+
await self._close_best_effort(conn, "release-reset-rolled-back")
21912227
conn._pool_released = True
21922228
return
21932229

@@ -2199,8 +2235,7 @@ async def _release(self, conn: DqliteConnection) -> None:
21992235
# the conn is unreachable). Mirrors the close-vs-initialize
22002236
# symmetry fix.
22012237
if self._closed:
2202-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
2203-
await asyncio.shield(conn.close())
2238+
await self._close_best_effort(conn, "release-post-reset-closed")
22042239
conn._pool_released = True
22052240
return
22062241

@@ -2223,8 +2258,7 @@ async def _release(self, conn: DqliteConnection) -> None:
22232258
try:
22242259
self._pool.put_nowait(conn)
22252260
except asyncio.QueueFull:
2226-
with contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS):
2227-
await asyncio.shield(conn.close())
2261+
await self._close_best_effort(conn, "release-queuefull")
22282262
conn._pool_released = True
22292263
else:
22302264
conn._pool_released = True
Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
"""Pin: the seven release/drain close sites that previously used
2+
``contextlib.suppress(*_POOL_CLEANUP_EXCEPTIONS)`` now route through
3+
``ConnectionPool._close_best_effort``, which logs the absorbed
4+
exception at DEBUG with ``exc_info=True`` rather than silently
5+
dropping it.
6+
7+
Pre-fix the sites suppressed silently — an OSError /
8+
DqliteConnectionError / OperationalError from ``conn.close()`` during
9+
release/drain dissolved with no trace, asymmetric with
10+
``_initialize_close_unqueued`` (which logged via ``exc_info=True``).
11+
An operator investigating "the pool kept hitting close_timeout under
12+
load" had nothing to find at any verbosity.
13+
14+
This test drives the closed-pool release branch (representative of the
15+
seven sites) and asserts a DEBUG record naming the site lands with
16+
``exc_info`` populated. Pairs with
17+
``test_pool_initialize_unqueued_survivor_close_log.py``.
18+
"""
19+
20+
from __future__ import annotations
21+
22+
import logging
23+
from typing import Any
24+
from unittest.mock import MagicMock
25+
26+
import pytest
27+
28+
from dqliteclient.cluster import ClusterClient
29+
from dqliteclient.exceptions import DqliteConnectionError
30+
from dqliteclient.pool import ConnectionPool
31+
32+
33+
class _FakeConn:
34+
def __init__(self, *, close_side_effect: BaseException | None = None) -> None:
35+
self._address = "localhost:9001"
36+
self._in_transaction = False
37+
self._tx_owner = None
38+
self._pool_released = False
39+
self._protocol = MagicMock()
40+
self._protocol.is_wire_coherent = True
41+
self._protocol._writer = MagicMock()
42+
self._protocol._writer.transport = MagicMock()
43+
self._protocol._writer.transport.is_closing = lambda: False
44+
self._protocol._reader = MagicMock()
45+
self._protocol._reader.at_eof = lambda: False
46+
self._close_side_effect = close_side_effect
47+
self.close_calls = 0
48+
49+
@property
50+
def is_connected(self) -> bool:
51+
return self._protocol is not None
52+
53+
async def close(self) -> None:
54+
self.close_calls += 1
55+
if self._close_side_effect is not None:
56+
raise self._close_side_effect
57+
self._pool_released = True
58+
self._protocol = None # type: ignore[assignment]
59+
60+
61+
def _make_pool_with_fake_cluster(fake_conn: _FakeConn) -> ConnectionPool:
62+
async def _connect(**kwargs: Any) -> _FakeConn:
63+
return fake_conn
64+
65+
cluster = MagicMock(spec=ClusterClient)
66+
cluster.connect = _connect
67+
return ConnectionPool(
68+
addresses=["localhost:9001"],
69+
min_size=0,
70+
max_size=1,
71+
timeout=1.0,
72+
cluster=cluster,
73+
)
74+
75+
76+
@pytest.mark.asyncio
77+
async def test_release_closed_branch_logs_close_error_at_debug(
78+
caplog: pytest.LogCaptureFixture,
79+
) -> None:
80+
"""The closed-pool release branch routes through
81+
``_close_best_effort``: a ``DqliteConnectionError`` from
82+
``conn.close()`` is absorbed but emits a DEBUG record naming the
83+
site (``"pool.release-closed: close error"``) with ``exc_info``
84+
populated."""
85+
fake = _FakeConn(close_side_effect=DqliteConnectionError("transport broken"))
86+
pool = _make_pool_with_fake_cluster(fake)
87+
88+
cm = pool.acquire()
89+
conn = await cm.__aenter__()
90+
assert conn is fake
91+
pool._closed = True
92+
93+
with caplog.at_level(logging.DEBUG, logger="dqliteclient.pool"):
94+
# Must NOT raise — the helper absorbs _POOL_CLEANUP_EXCEPTIONS.
95+
await cm.__aexit__(None, None, None)
96+
97+
assert fake.close_calls >= 1
98+
99+
# A DEBUG record named for the site, with exc_info populated.
100+
debug_records = [
101+
r for r in caplog.records if r.levelno == logging.DEBUG and "close error" in r.message
102+
]
103+
assert debug_records, (
104+
"expected a DEBUG record naming the release site; "
105+
f"saw records={[(r.levelname, r.message) for r in caplog.records]}"
106+
)
107+
matched = [r for r in debug_records if "release-closed" in r.message]
108+
assert matched, (
109+
f"expected the release-closed site name in the DEBUG record; "
110+
f"got messages={[r.message for r in debug_records]}"
111+
)
112+
rec = matched[0]
113+
assert rec.exc_info is not None
114+
assert isinstance(rec.exc_info[1], DqliteConnectionError)
115+
116+
117+
def test_close_best_effort_helper_exists_with_logging_shape() -> None:
118+
"""Static pin: ``ConnectionPool`` exposes ``_close_best_effort``, the
119+
DRY helper that replaced the seven ``contextlib.suppress(
120+
*_POOL_CLEANUP_EXCEPTIONS)`` sites. Detects accidental removal of
121+
the helper or rename that would silently drop the logging on the
122+
release path again."""
123+
assert hasattr(ConnectionPool, "_close_best_effort")

0 commit comments

Comments
 (0)