From 35141ad9c5e75bc103fe01e8d860878b58050130 Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 09:07:11 +0400 Subject: [PATCH 1/9] Remove deprecated pkg_resources --- pyproject.toml | 3 ++- setup.py | 7 +++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index e2b18388..e6e285cb 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -62,7 +62,8 @@ docs = [ [build-system] requires = [ "setuptools>=77.0.3", - "Cython(>=3.2.1,<4.0.0)" + "Cython(>=3.2.1,<4.0.0)", + "packaging", ] build-backend = "setuptools.build_meta" diff --git a/setup.py b/setup.py index f9fafadf..247fe596 100644 --- a/setup.py +++ b/setup.py @@ -188,8 +188,6 @@ def finalize_options(self): need_cythonize = True if need_cythonize: - import pkg_resources - # Double check Cython presence in case setup_requires # didn't go into effect (most likely because someone # imported Cython before setup_requires injected the @@ -201,8 +199,9 @@ def finalize_options(self): 'please install {} to compile asyncpg from source'.format( CYTHON_DEPENDENCY)) - cython_dep = pkg_resources.Requirement.parse(CYTHON_DEPENDENCY) - if Cython.__version__ not in cython_dep: + from packaging.requirements import Requirement + cython_dep = Requirement(CYTHON_DEPENDENCY) + if Cython.__version__ not in cython_dep.specifier: raise RuntimeError( 'asyncpg requires {}, got Cython=={}'.format( CYTHON_DEPENDENCY, Cython.__version__ From af1afbfe98e31b6f2214507ad067c42928e3aee8 Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 08:54:30 +0400 Subject: [PATCH 2/9] Rename Pool.min_size into Pool.init_size --- asyncpg/_testbase/__init__.py | 4 +- asyncpg/pool.py | 34 ++++++------ tests/test_adversity.py | 6 +-- tests/test_cache_invalidation.py | 6 +-- tests/test_connect.py | 4 +- tests/test_pool.py | 90 ++++++++++++++++---------------- 6 files changed, 72 insertions(+), 72 deletions(-) diff --git a/asyncpg/_testbase/__init__.py b/asyncpg/_testbase/__init__.py index 95775e11..33e2a975 100644 --- a/asyncpg/_testbase/__init__.py +++ b/asyncpg/_testbase/__init__.py @@ -267,7 +267,7 @@ def _shutdown_cluster(cluster): def create_pool(dsn=None, *, - min_size=10, + init_size=10, max_size=10, max_queries=50000, max_inactive_connection_lifetime=60.0, @@ -281,7 +281,7 @@ def create_pool(dsn=None, *, **connect_kwargs): return pool_class( dsn, - min_size=min_size, + init_size=init_size, max_size=max_size, max_queries=max_queries, loop=loop, diff --git a/asyncpg/pool.py b/asyncpg/pool.py index 5c7ea9ca..5fbce013 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -338,7 +338,7 @@ class Pool: """ __slots__ = ( - '_queue', '_loop', '_minsize', '_maxsize', + '_queue', '_loop', '_initsize', '_maxsize', '_init', '_connect', '_reset', '_connect_args', '_connect_kwargs', '_holders', '_initialized', '_initializing', '_closing', '_closed', '_connection_class', '_record_class', '_generation', @@ -346,7 +346,7 @@ class Pool: ) def __init__(self, *connect_args, - min_size, + init_size, max_size, max_queries, max_inactive_connection_lifetime, @@ -374,12 +374,12 @@ def __init__(self, *connect_args, if max_size <= 0: raise ValueError('max_size is expected to be greater than zero') - if min_size < 0: + if init_size < 0: raise ValueError( - 'min_size is expected to be greater or equal to zero') + 'init_size is expected to be greater or equal to zero') - if min_size > max_size: - raise ValueError('min_size is greater than max_size') + if init_size > max_size: + raise ValueError('init_size is greater than max_size') if max_queries <= 0: raise ValueError('max_queries is expected to be greater than zero') @@ -399,7 +399,7 @@ def __init__(self, *connect_args, 'record_class is expected to be a subclass of ' 'asyncpg.Record, got {!r}'.format(record_class)) - self._minsize = min_size + self._initsize = init_size self._maxsize = max_size self._holders = [] @@ -454,7 +454,7 @@ async def _initialize(self): self._holders.append(ch) self._queue.put_nowait(ch) - if self._minsize: + if self._initsize: # Since we use a LIFO queue, the first items in the queue will be # the last ones in `self._holders`. We want to pre-connect the # first few connections in the queue, therefore we want to walk @@ -465,11 +465,11 @@ async def _initialize(self): first_ch = self._holders[-1] # type: PoolConnectionHolder await first_ch.connect() - if self._minsize > 1: + if self._initsize > 1: connect_tasks = [] for i, ch in enumerate(reversed(self._holders[:-1])): - # `minsize - 1` because we already have first_ch - if i >= self._minsize - 1: + # `initsize - 1` because we already have first_ch + if i >= self._initsize - 1: break connect_tasks.append(ch.connect()) @@ -489,12 +489,12 @@ def get_size(self): """ return sum(h.is_connected() for h in self._holders) - def get_min_size(self): - """Return the minimum number of connections in this pool. + def get_init_size(self): + """Return the initial number of connections in this pool. .. versionadded:: 0.25.0 """ - return self._minsize + return self._initsize def get_max_size(self): """Return the maximum allowed number of connections in this pool. @@ -1073,7 +1073,7 @@ def __await__(self): def create_pool(dsn=None, *, - min_size=10, + init_size=10, max_size=10, max_queries=50000, max_inactive_connection_lifetime=300.0, @@ -1147,7 +1147,7 @@ def create_pool(dsn=None, *, the connections in this pool. Must be a subclass of :class:`~asyncpg.Record`. - :param int min_size: + :param int init_size: Number of connection the pool will be initialized with. :param int max_size: @@ -1235,7 +1235,7 @@ def create_pool(dsn=None, *, dsn, connection_class=connection_class, record_class=record_class, - min_size=min_size, + init_size=init_size, max_size=max_size, max_queries=max_queries, loop=loop, diff --git a/tests/test_adversity.py b/tests/test_adversity.py index a6e03feb..6c05e6a3 100644 --- a/tests/test_adversity.py +++ b/tests/test_adversity.py @@ -29,7 +29,7 @@ async def test_connection_close_timeout(self): @tb.with_timeout(30.0) async def test_pool_acquire_timeout(self): pool = await self.create_pool( - database='postgres', min_size=2, max_size=2) + database='postgres', init_size=2, max_size=2) try: self.proxy.trigger_connectivity_loss() for _ in range(2): @@ -46,7 +46,7 @@ async def test_pool_acquire_timeout(self): @tb.with_timeout(30.0) async def test_pool_release_timeout(self): pool = await self.create_pool( - database='postgres', min_size=2, max_size=2) + database='postgres', init_size=2, max_size=2) try: with self.assertRaises(asyncio.TimeoutError): async with pool.acquire(timeout=0.5): @@ -74,7 +74,7 @@ def kill_connectivity(): self.proxy.trigger_connectivity_loss() new_pool = self.create_pool( - database='postgres', min_size=pool_size, max_size=pool_size, + database='postgres', init_size=pool_size, max_size=pool_size, timeout=cmd_timeout, command_timeout=cmd_timeout) with self.assertRunUnder(worst_runtime): diff --git a/tests/test_cache_invalidation.py b/tests/test_cache_invalidation.py index 5cab2d92..49ab04fb 100644 --- a/tests/test_cache_invalidation.py +++ b/tests/test_cache_invalidation.py @@ -77,7 +77,7 @@ async def test_prepare_cache_invalidation_in_transaction(self): async def test_prepare_cache_invalidation_in_pool(self): pool = await self.create_pool(database='postgres', - min_size=2, max_size=2) + init_size=2, max_size=2) await self.con.execute('CREATE TABLE tab1(a int, b int)') @@ -309,10 +309,10 @@ async def test_type_cache_invalidation_on_change_attr(self): async def test_type_cache_invalidation_in_pool(self): await self.con.execute('CREATE DATABASE testdb') pool = await self.create_pool(database='postgres', - min_size=2, max_size=2) + init_size=2, max_size=2) pool_chk = await self.create_pool(database='testdb', - min_size=2, max_size=2) + init_size=2, max_size=2) await self.con.execute('CREATE TYPE typ1 AS (x int, y int)') await self.con.execute('CREATE TABLE tab1(a int, b typ1)') diff --git a/tests/test_connect.py b/tests/test_connect.py index 955fb825..441376f2 100644 --- a/tests/test_connect.py +++ b/tests/test_connect.py @@ -1961,7 +1961,7 @@ async def test_ssl_connection_pool(self): host='localhost', user='ssl_user', database='postgres', - min_size=5, + init_size=5, max_size=10, ssl=ssl_context) @@ -2221,7 +2221,7 @@ async def test_nossl_connection_pool(self): host='localhost', user='ssl_user', database='postgres', - min_size=5, + init_size=5, max_size=10, ssl='prefer') diff --git a/tests/test_pool.py b/tests/test_pool.py index 695363b7..3c5a43d0 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -47,7 +47,7 @@ async def test_pool_01(self): for n in {1, 5, 10, 20, 100}: with self.subTest(tasksnum=n): pool = await self.create_pool(database='postgres', - min_size=5, max_size=10) + init_size=5, max_size=10) async def worker(): con = await pool.acquire() @@ -62,7 +62,7 @@ async def test_pool_02(self): for n in {1, 3, 5, 10, 20, 100}: with self.subTest(tasksnum=n): async with self.create_pool(database='postgres', - min_size=5, max_size=5) as pool: + init_size=5, max_size=5) as pool: async def worker(): con = await pool.acquire(timeout=5) @@ -74,7 +74,7 @@ async def worker(): async def test_pool_03(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire(timeout=1) with self.assertRaises(asyncio.TimeoutError): @@ -85,7 +85,7 @@ async def test_pool_03(self): async def test_pool_04(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -110,7 +110,7 @@ async def test_pool_05(self): for n in {1, 3, 5, 10, 20, 100}: with self.subTest(tasksnum=n): pool = await self.create_pool(database='postgres', - min_size=5, max_size=10) + init_size=5, max_size=10) async def worker(): async with pool.acquire() as con: @@ -127,7 +127,7 @@ async def setup(con): fut.set_result(con) async with self.create_pool(database='postgres', - min_size=5, max_size=5, + init_size=5, max_size=5, setup=setup) as pool: async with pool.acquire() as con: pass @@ -169,7 +169,7 @@ async def user(pool): raise RuntimeError('init was not called') async with self.create_pool(database='postgres', - min_size=2, + init_size=2, max_size=5, connect=connect, init=init, @@ -196,7 +196,7 @@ async def bad_connect(*args, **kwargs): async def test_pool_08(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) with self.assertRaisesRegex(asyncpg.InterfaceError, 'is not a member'): @@ -204,10 +204,10 @@ async def test_pool_08(self): async def test_pool_09(self): pool1 = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) pool2 = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) try: con = await pool1.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -222,7 +222,7 @@ async def test_pool_09(self): async def test_pool_10(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire() await pool.release(con) @@ -232,7 +232,7 @@ async def test_pool_10(self): async def test_pool_11(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) async with pool.acquire() as con: self.assertIn(repr(con._con), repr(con)) # Test __repr__. @@ -289,7 +289,7 @@ async def test_pool_11(self): async def test_pool_12(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) async with pool.acquire() as con: self.assertTrue(isinstance(con, pg_connection.Connection)) @@ -299,7 +299,7 @@ async def test_pool_12(self): async def test_pool_13(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) async with pool.acquire() as con: self.assertIn('Execute an SQL command', con.execute.__doc__) @@ -335,7 +335,7 @@ async def setup(con): last_con = None cons = [] async with self.create_pool(database='postgres', - min_size=1, max_size=1, + init_size=1, max_size=1, setup=setup) as pool: with self.assertRaises(Error): await pool.acquire() @@ -349,7 +349,7 @@ async def setup(con): last_con = None cons = [] async with self.create_pool(database='postgres', - min_size=0, max_size=1, + init_size=0, max_size=1, init=setup) as pool: with self.assertRaises(Error): await pool.acquire() @@ -391,7 +391,7 @@ async def test_pool_auth(self): pool = await self.create_pool(database='postgres', user='pooluser', password='poolpassword', - min_size=5, max_size=10) + init_size=5, max_size=10) async def worker(): con = await pool.acquire() @@ -412,7 +412,7 @@ async def worker(): async def test_pool_handles_task_cancel_in_acquire_with_timeout(self): # See https://github.com/MagicStack/asyncpg/issues/547 pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) async def worker(): async with pool.acquire(timeout=100): @@ -433,7 +433,7 @@ async def test_pool_handles_task_cancel_in_release(self): # Use SlowResetConnectionPool to simulate # the Task.cancel() and __aexit__ race. pool = await self.create_pool(database='postgres', - min_size=1, max_size=1, + init_size=1, max_size=1, connection_class=SlowResetConnection) async def worker(): @@ -454,7 +454,7 @@ async def test_pool_handles_query_cancel_in_release(self): # Use SlowResetConnectionPool to simulate # the Task.cancel() and __aexit__ race. pool = await self.create_pool(database='postgres', - min_size=1, max_size=1, + init_size=1, max_size=1, connection_class=SlowCancelConnection) async def worker(): @@ -473,7 +473,7 @@ async def worker(): async def test_pool_no_acquire_deadlock(self): async with self.create_pool(database='postgres', - min_size=1, max_size=1, + init_size=1, max_size=1, max_queries=1) as pool: async def sleep_and_release(): @@ -507,7 +507,7 @@ async def test(pool): cons.add(con) async with self.create_pool( - database='postgres', min_size=10, max_size=10, + database='postgres', init_size=10, max_size=10, max_queries=1, connection_class=MyConnection, statement_cache_size=3) as pool: @@ -518,7 +518,7 @@ async def test(pool): async def test_pool_release_in_xact(self): """Test that Connection.reset() closes any open transaction.""" async with self.create_pool(database='postgres', - min_size=1, max_size=1) as pool: + init_size=1, max_size=1) as pool: async def get_xact_id(con): return await con.fetchval('select txid_current()') @@ -581,7 +581,7 @@ async def test_execute_with_arg(pool): async def run(N, meth): async with self.create_pool(database='postgres', - min_size=5, max_size=10) as pool: + init_size=5, max_size=10) as pool: coros = [meth(pool) for _ in range(N)] res = await asyncio.gather(*coros) @@ -608,7 +608,7 @@ async def worker(pool): N = 200 async with self.create_pool(database='postgres', - min_size=5, max_size=10) as pool: + init_size=5, max_size=10) as pool: await pool.execute('CREATE TABLE exmany (a text, b int)') try: @@ -625,7 +625,7 @@ async def worker(pool): async def test_pool_max_inactive_time_01(self): async with self.create_pool( - database='postgres', min_size=1, max_size=1, + database='postgres', init_size=1, max_size=1, max_inactive_connection_lifetime=0.1) as pool: # Test that it's OK if a query takes longer time to execute @@ -644,7 +644,7 @@ async def test_pool_max_inactive_time_01(self): async def test_pool_max_inactive_time_02(self): async with self.create_pool( - database='postgres', min_size=1, max_size=1, + database='postgres', init_size=1, max_size=1, max_inactive_connection_lifetime=0.5) as pool: # Test that we have a new connection after pool not @@ -667,7 +667,7 @@ async def test_pool_max_inactive_time_02(self): async def test_pool_max_inactive_time_03(self): async with self.create_pool( - database='postgres', min_size=1, max_size=1, + database='postgres', init_size=1, max_size=1, max_inactive_connection_lifetime=1) as pool: # Test that we start counting inactive time *after* @@ -708,7 +708,7 @@ async def worker(pool): N += 1 async with self.create_pool( - database='postgres', min_size=10, max_size=30, + database='postgres', init_size=10, max_size=30, max_inactive_connection_lifetime=0.1) as pool: workers = [worker(pool) for _ in range(50)] @@ -720,7 +720,7 @@ async def test_pool_max_inactive_time_05(self): # Test that idle never-acquired connections abide by # the max inactive lifetime. async with self.create_pool( - database='postgres', min_size=2, max_size=2, + database='postgres', init_size=2, max_size=2, max_inactive_connection_lifetime=0.2) as pool: self.assertIsNotNone(pool._holders[0]._con) @@ -736,7 +736,7 @@ async def test_pool_max_inactive_time_05(self): async def test_pool_handles_inactive_connection_errors(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -757,10 +757,10 @@ async def test_pool_handles_inactive_connection_errors(self): async def test_pool_size_and_capacity(self): async with self.create_pool( database='postgres', - min_size=2, + init_size=2, max_size=3, ) as pool: - self.assertEqual(pool.get_min_size(), 2) + self.assertEqual(pool.get_init_size(), 2) self.assertEqual(pool.get_max_size(), 3) self.assertEqual(pool.get_size(), 2) self.assertEqual(pool.get_idle_size(), 2) @@ -788,7 +788,7 @@ async def test_pool_closing(self): async def test_pool_handles_transaction_exit_in_asyncgen_1(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -809,7 +809,7 @@ class MyException(Exception): async def test_pool_handles_transaction_exit_in_asyncgen_2(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -833,7 +833,7 @@ class MyException(Exception): async def test_pool_handles_asyncgen_finalization(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -854,7 +854,7 @@ class MyException(Exception): async def test_pool_close_waits_for_release(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) flag = self.loop.create_future() conn_released = False @@ -877,7 +877,7 @@ async def worker(): async def test_pool_close_timeout(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) flag = self.loop.create_future() @@ -896,7 +896,7 @@ async def worker(): async def test_pool_expire_connections(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) con = await pool.acquire() try: @@ -909,7 +909,7 @@ async def test_pool_expire_connections(self): async def test_pool_set_connection_args(self): pool = await self.create_pool(database='postgres', - min_size=1, max_size=1) + init_size=1, max_size=1) # Test that connection is expired on release. con = await pool.acquire() @@ -951,7 +951,7 @@ async def test_pool_set_connection_args(self): await pool.close() async def test_pool_init_race(self): - pool = self.create_pool(database='postgres', min_size=1, max_size=1) + pool = self.create_pool(database='postgres', init_size=1, max_size=1) t1 = asyncio.ensure_future(pool) t2 = asyncio.ensure_future(pool) @@ -965,7 +965,7 @@ async def test_pool_init_race(self): await pool.close() async def test_pool_init_and_use_race(self): - pool = self.create_pool(database='postgres', min_size=1, max_size=1) + pool = self.create_pool(database='postgres', init_size=1, max_size=1) pool_task = asyncio.ensure_future(pool) await asyncio.sleep(0) @@ -980,7 +980,7 @@ async def test_pool_init_and_use_race(self): await pool.close() async def test_pool_remote_close(self): - pool = await self.create_pool(min_size=1, max_size=1) + pool = await self.create_pool(init_size=1, max_size=1) backend_pid_fut = self.loop.create_future() async def worker(): @@ -1037,7 +1037,7 @@ async def test_full_reconnect_on_node_change_role(self): return pool = await self.create_pool( - min_size=1, + init_size=1, max_size=1, target_session_attrs='primary' ) @@ -1081,7 +1081,7 @@ async def test_standby_pool_01(self): with self.subTest(tasksnum=n): pool = await self.create_pool( database='postgres', user='postgres', - min_size=5, max_size=10) + init_size=5, max_size=10) async def worker(): con = await pool.acquire() From 7b1fcde339fdd1d00d7f581c07146c8277fe1b9a Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 09:31:20 +0400 Subject: [PATCH 3/9] Support min_size in the pool --- asyncpg/_testbase/__init__.py | 2 + asyncpg/pool.py | 41 +++++++++++- tests/test_adversity.py | 7 +- tests/test_cache_invalidation.py | 6 +- tests/test_connect.py | 2 + tests/test_pool.py | 109 +++++++++++++++++++------------ 6 files changed, 116 insertions(+), 51 deletions(-) diff --git a/asyncpg/_testbase/__init__.py b/asyncpg/_testbase/__init__.py index 33e2a975..96d43403 100644 --- a/asyncpg/_testbase/__init__.py +++ b/asyncpg/_testbase/__init__.py @@ -268,6 +268,7 @@ def _shutdown_cluster(cluster): def create_pool(dsn=None, *, init_size=10, + min_size=10, max_size=10, max_queries=50000, max_inactive_connection_lifetime=60.0, @@ -282,6 +283,7 @@ def create_pool(dsn=None, *, return pool_class( dsn, init_size=init_size, + min_size=min_size, max_size=max_size, max_queries=max_queries, loop=loop, diff --git a/asyncpg/pool.py b/asyncpg/pool.py index 5fbce013..055ce6e5 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -293,6 +293,15 @@ def _deactivate_inactive_connection(self) -> None: 'attempting to deactivate an acquired connection') if self._con is not None: + # The connection is idle and not in use, + # but we have min size limitation. So keep it alive for a while. + if self._pool.get_size() <= self._pool.get_min_size(): + # We already in the callback. Clean the field + self._inactive_callback = None + # But next time it can be the case when we have to terminate it + self._setup_inactive_callback() + return + # The connection is idle and not in use, so it's fine to # use terminate() instead of close(). self._con.terminate() @@ -338,7 +347,7 @@ class Pool: """ __slots__ = ( - '_queue', '_loop', '_initsize', '_maxsize', + '_queue', '_loop', '_initsize', '_minsize', '_maxsize', '_init', '_connect', '_reset', '_connect_args', '_connect_kwargs', '_holders', '_initialized', '_initializing', '_closing', '_closed', '_connection_class', '_record_class', '_generation', @@ -347,6 +356,7 @@ class Pool: def __init__(self, *connect_args, init_size, + min_size, max_size, max_queries, max_inactive_connection_lifetime, @@ -374,6 +384,13 @@ def __init__(self, *connect_args, if max_size <= 0: raise ValueError('max_size is expected to be greater than zero') + if min_size < 0: + raise ValueError( + 'min_size is expected to be greater or equal to zero') + + if min_size > max_size: + raise ValueError('min_size is greater than max_size') + if init_size < 0: raise ValueError( 'init_size is expected to be greater or equal to zero') @@ -381,6 +398,9 @@ def __init__(self, *connect_args, if init_size > max_size: raise ValueError('init_size is greater than max_size') + if init_size < min_size: + raise ValueError('init_size is smaller than min_size') + if max_queries <= 0: raise ValueError('max_queries is expected to be greater than zero') @@ -400,6 +420,7 @@ def __init__(self, *connect_args, 'asyncpg.Record, got {!r}'.format(record_class)) self._initsize = init_size + self._minsize = min_size self._maxsize = max_size self._holders = [] @@ -492,10 +513,17 @@ def get_size(self): def get_init_size(self): """Return the initial number of connections in this pool. - .. versionadded:: 0.25.0 + .. versionadded:: 0.4.0 """ return self._initsize + def get_min_size(self): + """Return the minimum number of connections in this pool. + + .. versionadded:: 0.4.0 + """ + return self._minsize + def get_max_size(self): """Return the maximum allowed number of connections in this pool. @@ -1073,7 +1101,10 @@ def __await__(self): def create_pool(dsn=None, *, - init_size=10, + # Trigger min_size exception everywhere in the world + # to pay attention to the breaking changes + init_size=1, + min_size=10, max_size=10, max_queries=50000, max_inactive_connection_lifetime=300.0, @@ -1150,6 +1181,9 @@ def create_pool(dsn=None, *, :param int init_size: Number of connection the pool will be initialized with. + :param int minsize: + Min number of connections in the pool. + :param int max_size: Max number of connections in the pool. @@ -1236,6 +1270,7 @@ def create_pool(dsn=None, *, connection_class=connection_class, record_class=record_class, init_size=init_size, + min_size=min_size, max_size=max_size, max_queries=max_queries, loop=loop, diff --git a/tests/test_adversity.py b/tests/test_adversity.py index 6c05e6a3..6263be8a 100644 --- a/tests/test_adversity.py +++ b/tests/test_adversity.py @@ -29,7 +29,7 @@ async def test_connection_close_timeout(self): @tb.with_timeout(30.0) async def test_pool_acquire_timeout(self): pool = await self.create_pool( - database='postgres', init_size=2, max_size=2) + database='postgres', init_size=2, min_size=0, max_size=2) try: self.proxy.trigger_connectivity_loss() for _ in range(2): @@ -46,7 +46,7 @@ async def test_pool_acquire_timeout(self): @tb.with_timeout(30.0) async def test_pool_release_timeout(self): pool = await self.create_pool( - database='postgres', init_size=2, max_size=2) + database='postgres', init_size=2, min_size=0, max_size=2) try: with self.assertRaises(asyncio.TimeoutError): async with pool.acquire(timeout=0.5): @@ -74,7 +74,8 @@ def kill_connectivity(): self.proxy.trigger_connectivity_loss() new_pool = self.create_pool( - database='postgres', init_size=pool_size, max_size=pool_size, + database='postgres', + init_size=pool_size, min_size=0, max_size=pool_size, timeout=cmd_timeout, command_timeout=cmd_timeout) with self.assertRunUnder(worst_runtime): diff --git a/tests/test_cache_invalidation.py b/tests/test_cache_invalidation.py index 49ab04fb..f2901c9e 100644 --- a/tests/test_cache_invalidation.py +++ b/tests/test_cache_invalidation.py @@ -77,7 +77,7 @@ async def test_prepare_cache_invalidation_in_transaction(self): async def test_prepare_cache_invalidation_in_pool(self): pool = await self.create_pool(database='postgres', - init_size=2, max_size=2) + init_size=2, min_size=0, max_size=2) await self.con.execute('CREATE TABLE tab1(a int, b int)') @@ -309,10 +309,10 @@ async def test_type_cache_invalidation_on_change_attr(self): async def test_type_cache_invalidation_in_pool(self): await self.con.execute('CREATE DATABASE testdb') pool = await self.create_pool(database='postgres', - init_size=2, max_size=2) + init_size=2, min_size=0, max_size=2) pool_chk = await self.create_pool(database='testdb', - init_size=2, max_size=2) + init_size=2, min_size=0, max_size=2) await self.con.execute('CREATE TYPE typ1 AS (x int, y int)') await self.con.execute('CREATE TABLE tab1(a int, b typ1)') diff --git a/tests/test_connect.py b/tests/test_connect.py index 441376f2..247f4696 100644 --- a/tests/test_connect.py +++ b/tests/test_connect.py @@ -1962,6 +1962,7 @@ async def test_ssl_connection_pool(self): user='ssl_user', database='postgres', init_size=5, + min_size=5, max_size=10, ssl=ssl_context) @@ -2222,6 +2223,7 @@ async def test_nossl_connection_pool(self): user='ssl_user', database='postgres', init_size=5, + min_size=5, max_size=10, ssl='prefer') diff --git a/tests/test_pool.py b/tests/test_pool.py index 3c5a43d0..1c77a5c9 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -47,7 +47,9 @@ async def test_pool_01(self): for n in {1, 5, 10, 20, 100}: with self.subTest(tasksnum=n): pool = await self.create_pool(database='postgres', - init_size=5, max_size=10) + init_size=5, + min_size=1, + max_size=10) async def worker(): con = await pool.acquire() @@ -62,7 +64,9 @@ async def test_pool_02(self): for n in {1, 3, 5, 10, 20, 100}: with self.subTest(tasksnum=n): async with self.create_pool(database='postgres', - init_size=5, max_size=5) as pool: + init_size=5, + min_size=1, + max_size=5) as pool: async def worker(): con = await pool.acquire(timeout=5) @@ -74,7 +78,7 @@ async def worker(): async def test_pool_03(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) con = await pool.acquire(timeout=1) with self.assertRaises(asyncio.TimeoutError): @@ -85,7 +89,7 @@ async def test_pool_03(self): async def test_pool_04(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -110,7 +114,9 @@ async def test_pool_05(self): for n in {1, 3, 5, 10, 20, 100}: with self.subTest(tasksnum=n): pool = await self.create_pool(database='postgres', - init_size=5, max_size=10) + init_size=5, + min_size=1, + max_size=10) async def worker(): async with pool.acquire() as con: @@ -127,7 +133,7 @@ async def setup(con): fut.set_result(con) async with self.create_pool(database='postgres', - init_size=5, max_size=5, + init_size=5, min_size=1, max_size=5, setup=setup) as pool: async with pool.acquire() as con: pass @@ -170,7 +176,7 @@ async def user(pool): async with self.create_pool(database='postgres', init_size=2, - max_size=5, + min_size=1, max_size=5, connect=connect, init=init, setup=setup, @@ -196,7 +202,7 @@ async def bad_connect(*args, **kwargs): async def test_pool_08(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) with self.assertRaisesRegex(asyncpg.InterfaceError, 'is not a member'): @@ -204,10 +210,10 @@ async def test_pool_08(self): async def test_pool_09(self): pool1 = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) pool2 = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) try: con = await pool1.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -222,7 +228,7 @@ async def test_pool_09(self): async def test_pool_10(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) con = await pool.acquire() await pool.release(con) @@ -232,7 +238,7 @@ async def test_pool_10(self): async def test_pool_11(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) async with pool.acquire() as con: self.assertIn(repr(con._con), repr(con)) # Test __repr__. @@ -289,7 +295,7 @@ async def test_pool_11(self): async def test_pool_12(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) async with pool.acquire() as con: self.assertTrue(isinstance(con, pg_connection.Connection)) @@ -299,7 +305,7 @@ async def test_pool_12(self): async def test_pool_13(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) async with pool.acquire() as con: self.assertIn('Execute an SQL command', con.execute.__doc__) @@ -335,7 +341,7 @@ async def setup(con): last_con = None cons = [] async with self.create_pool(database='postgres', - init_size=1, max_size=1, + init_size=1, min_size=1, max_size=1, setup=setup) as pool: with self.assertRaises(Error): await pool.acquire() @@ -349,7 +355,7 @@ async def setup(con): last_con = None cons = [] async with self.create_pool(database='postgres', - init_size=0, max_size=1, + init_size=0, min_size=0, max_size=1, init=setup) as pool: with self.assertRaises(Error): await pool.acquire() @@ -391,7 +397,7 @@ async def test_pool_auth(self): pool = await self.create_pool(database='postgres', user='pooluser', password='poolpassword', - init_size=5, max_size=10) + init_size=5, min_size=1, max_size=10) async def worker(): con = await pool.acquire() @@ -412,7 +418,7 @@ async def worker(): async def test_pool_handles_task_cancel_in_acquire_with_timeout(self): # See https://github.com/MagicStack/asyncpg/issues/547 pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) async def worker(): async with pool.acquire(timeout=100): @@ -433,7 +439,7 @@ async def test_pool_handles_task_cancel_in_release(self): # Use SlowResetConnectionPool to simulate # the Task.cancel() and __aexit__ race. pool = await self.create_pool(database='postgres', - init_size=1, max_size=1, + init_size=1, min_size=1, max_size=1, connection_class=SlowResetConnection) async def worker(): @@ -454,7 +460,7 @@ async def test_pool_handles_query_cancel_in_release(self): # Use SlowResetConnectionPool to simulate # the Task.cancel() and __aexit__ race. pool = await self.create_pool(database='postgres', - init_size=1, max_size=1, + init_size=1, min_size=1, max_size=1, connection_class=SlowCancelConnection) async def worker(): @@ -473,7 +479,7 @@ async def worker(): async def test_pool_no_acquire_deadlock(self): async with self.create_pool(database='postgres', - init_size=1, max_size=1, + init_size=1, min_size=1, max_size=1, max_queries=1) as pool: async def sleep_and_release(): @@ -507,7 +513,7 @@ async def test(pool): cons.add(con) async with self.create_pool( - database='postgres', init_size=10, max_size=10, + database='postgres', init_size=10, min_size=1, max_size=10, max_queries=1, connection_class=MyConnection, statement_cache_size=3) as pool: @@ -518,7 +524,9 @@ async def test(pool): async def test_pool_release_in_xact(self): """Test that Connection.reset() closes any open transaction.""" async with self.create_pool(database='postgres', - init_size=1, max_size=1) as pool: + init_size=1, + min_size=1, + max_size=1) as pool: async def get_xact_id(con): return await con.fetchval('select txid_current()') @@ -581,7 +589,9 @@ async def test_execute_with_arg(pool): async def run(N, meth): async with self.create_pool(database='postgres', - init_size=5, max_size=10) as pool: + init_size=5, + min_size=0, + max_size=10) as pool: coros = [meth(pool) for _ in range(N)] res = await asyncio.gather(*coros) @@ -608,7 +618,9 @@ async def worker(pool): N = 200 async with self.create_pool(database='postgres', - init_size=5, max_size=10) as pool: + init_size=5, + min_size=0, + max_size=10) as pool: await pool.execute('CREATE TABLE exmany (a text, b int)') try: @@ -625,7 +637,7 @@ async def worker(pool): async def test_pool_max_inactive_time_01(self): async with self.create_pool( - database='postgres', init_size=1, max_size=1, + database='postgres', init_size=1, min_size=0, max_size=1, max_inactive_connection_lifetime=0.1) as pool: # Test that it's OK if a query takes longer time to execute @@ -644,7 +656,7 @@ async def test_pool_max_inactive_time_01(self): async def test_pool_max_inactive_time_02(self): async with self.create_pool( - database='postgres', init_size=1, max_size=1, + database='postgres', init_size=1, min_size=0, max_size=1, max_inactive_connection_lifetime=0.5) as pool: # Test that we have a new connection after pool not @@ -667,7 +679,7 @@ async def test_pool_max_inactive_time_02(self): async def test_pool_max_inactive_time_03(self): async with self.create_pool( - database='postgres', init_size=1, max_size=1, + database='postgres', init_size=1, min_size=0, max_size=1, max_inactive_connection_lifetime=1) as pool: # Test that we start counting inactive time *after* @@ -708,7 +720,7 @@ async def worker(pool): N += 1 async with self.create_pool( - database='postgres', init_size=10, max_size=30, + database='postgres', init_size=10, min_size=0, max_size=30, max_inactive_connection_lifetime=0.1) as pool: workers = [worker(pool) for _ in range(50)] @@ -720,7 +732,7 @@ async def test_pool_max_inactive_time_05(self): # Test that idle never-acquired connections abide by # the max inactive lifetime. async with self.create_pool( - database='postgres', init_size=2, max_size=2, + database='postgres', init_size=2, min_size=0, max_size=2, max_inactive_connection_lifetime=0.2) as pool: self.assertIsNotNone(pool._holders[0]._con) @@ -736,7 +748,7 @@ async def test_pool_max_inactive_time_05(self): async def test_pool_handles_inactive_connection_errors(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=0, max_size=1) con = await pool.acquire(timeout=POOL_NOMINAL_TIMEOUT) @@ -758,9 +770,11 @@ async def test_pool_size_and_capacity(self): async with self.create_pool( database='postgres', init_size=2, + min_size=1, max_size=3, ) as pool: self.assertEqual(pool.get_init_size(), 2) + self.assertEqual(pool.get_min_size(), 1) self.assertEqual(pool.get_max_size(), 3) self.assertEqual(pool.get_size(), 2) self.assertEqual(pool.get_idle_size(), 2) @@ -788,7 +802,7 @@ async def test_pool_closing(self): async def test_pool_handles_transaction_exit_in_asyncgen_1(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -809,7 +823,7 @@ class MyException(Exception): async def test_pool_handles_transaction_exit_in_asyncgen_2(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -833,7 +847,7 @@ class MyException(Exception): async def test_pool_handles_asyncgen_finalization(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) locals_ = {} exec(textwrap.dedent('''\ @@ -854,7 +868,7 @@ class MyException(Exception): async def test_pool_close_waits_for_release(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) flag = self.loop.create_future() conn_released = False @@ -877,7 +891,7 @@ async def worker(): async def test_pool_close_timeout(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) flag = self.loop.create_future() @@ -896,7 +910,7 @@ async def worker(): async def test_pool_expire_connections(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) con = await pool.acquire() try: @@ -909,7 +923,7 @@ async def test_pool_expire_connections(self): async def test_pool_set_connection_args(self): pool = await self.create_pool(database='postgres', - init_size=1, max_size=1) + init_size=1, min_size=1, max_size=1) # Test that connection is expired on release. con = await pool.acquire() @@ -951,7 +965,12 @@ async def test_pool_set_connection_args(self): await pool.close() async def test_pool_init_race(self): - pool = self.create_pool(database='postgres', init_size=1, max_size=1) + pool = self.create_pool( + database='postgres', + init_size=1, + min_size=1, + max_size=1, + ) t1 = asyncio.ensure_future(pool) t2 = asyncio.ensure_future(pool) @@ -965,7 +984,12 @@ async def test_pool_init_race(self): await pool.close() async def test_pool_init_and_use_race(self): - pool = self.create_pool(database='postgres', init_size=1, max_size=1) + pool = self.create_pool( + database='postgres', + init_size=1, + min_size=1, + max_size=1, + ) pool_task = asyncio.ensure_future(pool) await asyncio.sleep(0) @@ -980,7 +1004,7 @@ async def test_pool_init_and_use_race(self): await pool.close() async def test_pool_remote_close(self): - pool = await self.create_pool(init_size=1, max_size=1) + pool = await self.create_pool(init_size=1, min_size=1, max_size=1) backend_pid_fut = self.loop.create_future() async def worker(): @@ -1038,6 +1062,7 @@ async def test_full_reconnect_on_node_change_role(self): pool = await self.create_pool( init_size=1, + min_size=1, max_size=1, target_session_attrs='primary' ) @@ -1081,7 +1106,7 @@ async def test_standby_pool_01(self): with self.subTest(tasksnum=n): pool = await self.create_pool( database='postgres', user='postgres', - init_size=5, max_size=10) + init_size=5, min_size=0, max_size=10) async def worker(): con = await pool.acquire() From 199358a22e396d3b391537b7edf28303c461fa3a Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 10:14:53 +0400 Subject: [PATCH 4/9] Add tests for check min_size --- tests/test_pool.py | 104 +++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 104 insertions(+) diff --git a/tests/test_pool.py b/tests/test_pool.py index 1c77a5c9..70a54a24 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -746,6 +746,110 @@ async def test_pool_max_inactive_time_05(self): # but should be closed nonetheless. self.assertIs(pool._holders[1]._con, None) + async def test_pool_min_size_keeps_connections_alive(self): + # Test that min_size prevents idle connections from being closed. + async with self.create_pool( + database='postgres', init_size=2, min_size=2, max_size=2, + max_inactive_connection_lifetime=0.2) as pool: + + con0 = pool._holders[0]._con + con1 = pool._holders[1]._con + self.assertIsNotNone(con0) + self.assertIsNotNone(con1) + + await asyncio.sleep(0.5) + + # Connections should be kept alive because pool size == min_size. + self.assertIs(pool._holders[0]._con, con0) + self.assertIs(pool._holders[1]._con, con1) + + async def test_pool_min_size_partial_keep(self): + # When the pool has more connections than min_size, only the excess + # connections should be allowed to expire; the min_size ones are kept. + async with self.create_pool( + database='postgres', init_size=3, min_size=1, max_size=3, + max_inactive_connection_lifetime=0.2) as pool: + + # Force all 3 connections to be created by acquiring them all. + c1 = await pool.acquire() + c2 = await pool.acquire() + c3 = await pool.acquire() + await pool.release(c1) + await pool.release(c2) + await pool.release(c3) + + self.assertEqual(pool.get_size(), 3) + + await asyncio.sleep(0.5) + + # At least min_size (1) connection must survive. + alive = sum( + 1 for h in pool._holders if h._con is not None + ) + self.assertGreaterEqual(alive, 1) + + async def test_pool_min_size_zero_allows_full_expiry(self): + # When min_size=0, all idle connections are allowed to expire. + async with self.create_pool( + database='postgres', init_size=2, min_size=0, max_size=2, + max_inactive_connection_lifetime=0.2) as pool: + + self.assertIsNotNone(pool._holders[0]._con) + self.assertIsNotNone(pool._holders[1]._con) + + await asyncio.sleep(0.5) + + self.assertIs(pool._holders[0]._con, None) + self.assertIs(pool._holders[1]._con, None) + + async def test_pool_min_size_validation(self): + # init_size < min_size should raise. + with self.assertRaisesRegex(ValueError, 'init_size is smaller than min_size'): + await self.create_pool( + database='postgres', init_size=1, min_size=2, max_size=5) + + # min_size > max_size should raise. + with self.assertRaisesRegex(ValueError, 'min_size is greater than max_size'): + await self.create_pool( + database='postgres', init_size=3, min_size=3, max_size=2) + + # init_size > max_size should raise. + with self.assertRaisesRegex(ValueError, 'init_size is greater than max_size'): + await self.create_pool( + database='postgres', init_size=5, min_size=1, max_size=3) + + # init_size < 0 should raise. + with self.assertRaisesRegex(ValueError, + 'init_size is expected to be greater or equal to zero'): + await self.create_pool( + database='postgres', init_size=-1, min_size=0, max_size=3) + + async def test_pool_init_size_and_min_size_getters(self): + async with self.create_pool( + database='postgres', init_size=3, min_size=2, max_size=5) as pool: + self.assertEqual(pool.get_init_size(), 3) + self.assertEqual(pool.get_min_size(), 2) + self.assertEqual(pool.get_max_size(), 5) + self.assertEqual(pool.get_size(), 3) + + async def test_pool_min_size_reconnect_after_expiry(self): + # Connections kept alive by min_size should still be functional. + async with self.create_pool( + database='postgres', init_size=1, min_size=1, max_size=1, + max_inactive_connection_lifetime=0.2) as pool: + + con_before = pool._holders[0]._con + self.assertIsNotNone(con_before) + + await asyncio.sleep(0.5) + + # Connection must still be alive due to min_size=1. + self.assertIs(pool._holders[0]._con, con_before) + + # And it must still work. + result = await pool.fetchval('SELECT 42::int') + self.assertEqual(result, 42) + async def test_pool_handles_inactive_connection_errors(self): pool = await self.create_pool(database='postgres', init_size=1, min_size=0, max_size=1) From 76a8a7091fb11c0034badeb79214df4e7a46843c Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 10:24:14 +0400 Subject: [PATCH 5/9] Cleanup docstrings --- asyncpg/pool.py | 20 +++++++++++++++----- 1 file changed, 15 insertions(+), 5 deletions(-) diff --git a/asyncpg/pool.py b/asyncpg/pool.py index 055ce6e5..0389d792 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -513,14 +513,18 @@ def get_size(self): def get_init_size(self): """Return the initial number of connections in this pool. - .. versionadded:: 0.4.0 + .. versionadded:: 0.32.0 """ return self._initsize def get_min_size(self): """Return the minimum number of connections in this pool. - .. versionadded:: 0.4.0 + .. versionadded:: 0.25.0 + + .. versionchanged:: 0.32.0 + The parameter now controls the connection floor rather than the + initial pool size (see ``init_size``). """ return self._minsize @@ -1179,10 +1183,10 @@ def create_pool(dsn=None, *, :class:`~asyncpg.Record`. :param int init_size: - Number of connection the pool will be initialized with. + Number of connections the pool will be initialized with. - :param int minsize: - Min number of connections in the pool. + :param int min_size: + Minimum number of connections the pool will keep alive at all times. :param int max_size: Max number of connections in the pool. @@ -1264,6 +1268,12 @@ def create_pool(dsn=None, *, .. versionchanged:: 0.30.0 Added the *connect* and *reset* parameters. + + .. versionchanged:: 0.32.0 + The *min_size* parameter now defines the connection floor (minimum + number of live connections kept at all times). The former role of + *min_size* — setting the initial pool size — is now handled by the + new *init_size* parameter. """ return Pool( dsn, From 8672b958768ca37c9a59a67352eade0c48e7a4b1 Mon Sep 17 00:00:00 2001 From: Dmitry Shatilov Date: Mon, 18 May 2026 23:18:08 +0400 Subject: [PATCH 6/9] Remove the trigger: default pool should work --- asyncpg/pool.py | 4 +--- tests/test_pool.py | 19 +++++++++++++------ 2 files changed, 14 insertions(+), 9 deletions(-) diff --git a/asyncpg/pool.py b/asyncpg/pool.py index 0389d792..a118cced 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -1105,9 +1105,7 @@ def __await__(self): def create_pool(dsn=None, *, - # Trigger min_size exception everywhere in the world - # to pay attention to the breaking changes - init_size=1, + init_size=10, min_size=10, max_size=10, max_queries=50000, diff --git a/tests/test_pool.py b/tests/test_pool.py index 70a54a24..27e90123 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -804,29 +804,36 @@ async def test_pool_min_size_zero_allows_full_expiry(self): async def test_pool_min_size_validation(self): # init_size < min_size should raise. - with self.assertRaisesRegex(ValueError, 'init_size is smaller than min_size'): + with self.assertRaisesRegex(ValueError, + 'init_size is smaller than min_size'): await self.create_pool( database='postgres', init_size=1, min_size=2, max_size=5) # min_size > max_size should raise. - with self.assertRaisesRegex(ValueError, 'min_size is greater than max_size'): + with self.assertRaisesRegex(ValueError, + 'min_size is greater than max_size'): await self.create_pool( database='postgres', init_size=3, min_size=3, max_size=2) # init_size > max_size should raise. - with self.assertRaisesRegex(ValueError, 'init_size is greater than max_size'): + with self.assertRaisesRegex(ValueError, + 'init_size is greater than max_size'): await self.create_pool( database='postgres', init_size=5, min_size=1, max_size=3) # init_size < 0 should raise. - with self.assertRaisesRegex(ValueError, - 'init_size is expected to be greater or equal to zero'): + with self.assertRaisesRegex( + ValueError, + 'init_size is expected to be greater or equal to zero'): await self.create_pool( database='postgres', init_size=-1, min_size=0, max_size=3) async def test_pool_init_size_and_min_size_getters(self): async with self.create_pool( - database='postgres', init_size=3, min_size=2, max_size=5) as pool: + database='postgres', + init_size=3, + min_size=2, + max_size=5) as pool: self.assertEqual(pool.get_init_size(), 3) self.assertEqual(pool.get_min_size(), 2) self.assertEqual(pool.get_max_size(), 5) From 9c476fd57111937b8a65edf455227529b85c0076 Mon Sep 17 00:00:00 2001 From: Elvis Pranskevichus Date: Wed, 30 Sep 2026 09:42:17 -0700 Subject: [PATCH 7/9] Fix pool minimum connection maintenance and shutdown cleanup --- asyncpg/_testbase/__init__.py | 10 +- asyncpg/connect_utils.py | 11 +- asyncpg/connection.py | 11 +- asyncpg/pool.py | 139 ++++++++++++++++---- tests/test_connect.py | 22 ++++ tests/test_pool.py | 240 +++++++++++++++++++++++++++++++++- 6 files changed, 395 insertions(+), 38 deletions(-) diff --git a/asyncpg/_testbase/__init__.py b/asyncpg/_testbase/__init__.py index 96d43403..326564cf 100644 --- a/asyncpg/_testbase/__init__.py +++ b/asyncpg/_testbase/__init__.py @@ -267,7 +267,7 @@ def _shutdown_cluster(cluster): def create_pool(dsn=None, *, - init_size=10, + init_size=None, min_size=10, max_size=10, max_queries=50000, @@ -368,10 +368,16 @@ def setUp(self): self._pools = [] def tearDown(self): - super().tearDown() + maintenance_tasks = [] for pool in self._pools: pool.terminate() + if pool._maintenance_task is not None: + maintenance_tasks.append(pool._maintenance_task) + if maintenance_tasks: + self.loop.run_until_complete(asyncio.gather( + *maintenance_tasks, return_exceptions=True)) self._pools = [] + super().tearDown() def create_pool(self, pool_class=pg_pool.Pool, connection_class=pg_connection.Connection, **kwargs): diff --git a/asyncpg/connect_utils.py b/asyncpg/connect_utils.py index 07c4fdde..297fc006 100644 --- a/asyncpg/connect_utils.py +++ b/asyncpg/connect_utils.py @@ -1096,7 +1096,16 @@ async def __connect_addr( else: connector = loop.create_connection(proto_factory, *addr) - tr, pr = await connector + try: + tr, pr = await connector + except (Exception, asyncio.CancelledError): + # The protocol can exist before create_connection() returns. If + # that operation is cancelled, nobody will await authentication. + if not connected.done(): + connected.cancel() + elif not connected.cancelled(): + connected.exception() + raise try: await connected diff --git a/asyncpg/connection.py b/asyncpg/connection.py index 71fb04f8..a917f23f 100644 --- a/asyncpg/connection.py +++ b/asyncpg/connection.py @@ -54,6 +54,7 @@ class Connection(metaclass=ConnectionMeta): '_intro_query', '_reset_query', '_proxy', '_stmt_exclusive_section', '_config', '_params', '_addr', '_log_listeners', '_termination_listeners', '_cancellations', + '_pool_holder', '_source_traceback', '_query_loggers', '__weakref__') def __init__(self, protocol, transport, loop, @@ -105,6 +106,7 @@ def __init__(self, protocol, transport, loop, self._reset_query = None self._proxy = None + self._pool_holder = None # Used to serialize operations that might involve anonymous # statements. Specifically, we want to make the following @@ -1581,10 +1583,11 @@ def _cleanup(self): # Free the resources associated with this connection. # This must be called when a connection is terminated. - if self._proxy is not None: - # Connection is a member of a pool, so let the pool - # know that this connection is dead. - self._proxy._holder._release_on_close() + if self._pool_holder is not None: + # Idle connections have no proxy, but still belong to a holder. + holder, self._pool_holder = self._pool_holder, None + if holder._con is self: + holder._release_on_close() self._mark_stmts_as_closed() self._listeners.clear() diff --git a/asyncpg/pool.py b/asyncpg/pool.py index a118cced..6980f663 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -25,6 +25,24 @@ logger = logging.getLogger(__name__) +class _PoolConnectionHolderQueue(asyncio.LifoQueue): + """Prefer live connections, with LIFO ordering within each group.""" + + def _get(self): + for i in range(len(self._queue) - 1, -1, -1): + if self._queue[i].is_connected(): + return self._queue.pop(i) + return super()._get() + + def get_disconnected_nowait(self): + for i in range(len(self._queue) - 1, -1, -1): + if not self._queue[i].is_connected(): + holder = self._queue.pop(i) + self._wakeup_next(self._putters) + return holder + raise asyncio.QueueEmpty + + class PoolConnectionProxyMeta(type): def __new__( @@ -150,8 +168,17 @@ async def connect(self) -> None: 'PoolConnectionHolder.connect() called while another ' 'connection already exists') - self._con = await self._pool._get_new_connection() - self._generation = self._pool._generation + generation = self._pool._generation + con = await self._pool._get_new_connection() + if self._pool.is_closing(): + await con.close() + raise exceptions.InterfaceError('pool is closing') + if con.is_closed(): + raise exceptions.ConnectionDoesNotExistError( + 'connection was closed during pool initialization') + self._con = con + con._pool_holder = self + self._generation = generation self._maybe_cancel_inactive_callback() self._setup_inactive_callback() @@ -292,28 +319,23 @@ def _deactivate_inactive_connection(self) -> None: raise exceptions.InternalClientError( 'attempting to deactivate an acquired connection') + self._inactive_callback = None if self._con is not None: - # The connection is idle and not in use, - # but we have min size limitation. So keep it alive for a while. - if self._pool.get_size() <= self._pool.get_min_size(): - # We already in the callback. Clean the field - self._inactive_callback = None - # But next time it can be the case when we have to terminate it - self._setup_inactive_callback() + if (self.is_connected() and not self._pool.is_closing() and + self._pool.get_size() <= self._pool.get_min_size()): + # A floor connection needs no further inactivity checks. + # Acquiring and releasing it will arm a new timer. return # The connection is idle and not in use, so it's fine to # use terminate() instead of close(). self._con.terminate() - # Must call clear_connection, because _deactivate_connection - # is called when the connection is *not* checked out, and - # so terminate() above will not call the below. - self._release_on_close() def _release_on_close(self) -> None: self._maybe_cancel_inactive_callback() self._release() self._con = None + self._pool._schedule_min_size_maintenance() def _release(self) -> None: """Release this connection holder.""" @@ -351,11 +373,12 @@ class Pool: '_init', '_connect', '_reset', '_connect_args', '_connect_kwargs', '_holders', '_initialized', '_initializing', '_closing', '_closed', '_connection_class', '_record_class', '_generation', - '_setup', '_max_queries', '_max_inactive_connection_lifetime' + '_setup', '_max_queries', '_max_inactive_connection_lifetime', + '_maintenance_task', ) def __init__(self, *connect_args, - init_size, + init_size=None, min_size, max_size, max_queries, @@ -391,6 +414,9 @@ def __init__(self, *connect_args, if min_size > max_size: raise ValueError('min_size is greater than max_size') + if init_size is None: + init_size = min_size + if init_size < 0: raise ValueError( 'init_size is expected to be greater or equal to zero') @@ -434,6 +460,7 @@ def __init__(self, *connect_args, self._closing = False self._closed = False self._generation = 0 + self._maintenance_task = None self._connect = connect if connect is not None else connection.connect self._connect_args = connect_args @@ -459,12 +486,19 @@ async def _async__init__(self): try: await self._initialize() return self + except (Exception, asyncio.CancelledError): + # Failed initialization must not leave warm connections behind. + self._closed = True + for holder in self._holders: + holder.terminate() + raise finally: self._initializing = False self._initialized = True + self._schedule_min_size_maintenance() async def _initialize(self): - self._queue = asyncio.LifoQueue(maxsize=self._maxsize) + self._queue = _PoolConnectionHolderQueue(maxsize=self._maxsize) for _ in range(self._maxsize): ch = PoolConnectionHolder( self, @@ -496,6 +530,48 @@ async def _initialize(self): await asyncio.gather(*connect_tasks) + def _schedule_min_size_maintenance(self): + if (not self._initialized or self._initializing or self.is_closing() + or not self._minsize or self._maintenance_task is not None): + return + if self.get_size() < self._minsize: + self._maintenance_task = self._loop.create_task( + self._maintain_min_size()) + + async def _maintain_min_size(self): + retry_delay = 1.0 + try: + while not self.is_closing() and self.get_size() < self._minsize: + try: + holder = self._queue.get_disconnected_nowait() + except asyncio.QueueEmpty: + # Acquirers are already connecting the remaining holders. + return + + failed = False + try: + if holder._con is not None: + holder.terminate() + await holder.connect() + except asyncio.CancelledError: + raise + except Exception: + if self.is_closing(): + return + failed = True + logger.warning('Failed to restore the pool connection ' + 'floor; retrying', exc_info=True) + finally: + self._queue.put_nowait(holder) + + if failed: + await asyncio.sleep(retry_delay) + retry_delay = min(retry_delay * 2, 60.0) + else: + retry_delay = 1.0 + finally: + self._maintenance_task = None + def is_closing(self): """Return ``True`` if the pool is closing or is closed. @@ -913,6 +989,7 @@ async def _acquire_impl(): proxy = await ch.acquire() # type: PoolConnectionProxy except (Exception, asyncio.CancelledError): self._queue.put_nowait(ch) + self._schedule_min_size_maintenance() raise else: # Record the timeout, as we will apply it by default @@ -992,6 +1069,12 @@ async def close(self): warning_callback = None try: + if self._maintenance_task is not None: + self._maintenance_task.cancel() + await asyncio.gather( + self._maintenance_task, return_exceptions=True) + self._maintenance_task = None + warning_callback = self._loop.call_later( 60, self._warn_on_long_close) @@ -1024,9 +1107,11 @@ def terminate(self): if self._closed: return self._check_init() + self._closed = True + if self._maintenance_task is not None: + self._maintenance_task.cancel() for ch in self._holders: ch.terminate() - self._closed = True async def expire_connections(self): """Expire all currently open connections. @@ -1105,7 +1190,7 @@ def __await__(self): def create_pool(dsn=None, *, - init_size=10, + init_size=None, min_size=10, max_size=10, max_queries=50000, @@ -1181,10 +1266,14 @@ def create_pool(dsn=None, *, :class:`~asyncpg.Record`. :param int init_size: - Number of connections the pool will be initialized with. + Number of connections the pool will be initialized with. Defaults + to *min_size*. Must be between *min_size* and *max_size*. :param int min_size: - Minimum number of connections the pool will keep alive at all times. + Minimum number of connections retained during idle periods. Closed + connections are replaced in the background to restore this floor. + The pool may temporarily fall below it while reconnecting or while + the server is unavailable. Pass ``0`` to allow the pool to drain. :param int max_size: Max number of connections in the pool. @@ -1195,7 +1284,8 @@ def create_pool(dsn=None, *, :param float max_inactive_connection_lifetime: Number of seconds after which inactive connections in the - pool will be closed. Pass ``0`` to disable this mechanism. + pool above *min_size* will be closed. Pass ``0`` to disable this + mechanism. :param coroutine connect: A coroutine that is called instead of @@ -1268,10 +1358,9 @@ def create_pool(dsn=None, *, Added the *connect* and *reset* parameters. .. versionchanged:: 0.32.0 - The *min_size* parameter now defines the connection floor (minimum - number of live connections kept at all times). The former role of - *min_size* — setting the initial pool size — is now handled by the - new *init_size* parameter. + The *min_size* parameter now defines the connection floor. The former + role of *min_size* — setting the initial pool size — is now handled by + the new *init_size* parameter, which defaults to *min_size*. """ return Pool( dsn, diff --git a/tests/test_connect.py b/tests/test_connect.py index 247f4696..4ec818b9 100644 --- a/tests/test_connect.py +++ b/tests/test_connect.py @@ -1658,6 +1658,28 @@ async def test_connect_args_validation(self): class TestConnection(tb.ConnectedTestCase): + async def test_connection_cancelled_during_transport_setup(self): + for connection_lost_first in (False, True): + with self.subTest(connection_lost_first=connection_lost_first): + async def interrupted_connector(factory, *args, **kwargs): + proto = factory() + if connection_lost_first: + proto.connection_lost(None) + else: + self.loop.call_soon(proto.connection_lost, None) + raise asyncio.CancelledError + + with unittest.mock.patch.object( + self.loop, 'create_connection', interrupted_connector, + ): + with self.assertRaises(asyncio.CancelledError): + await self.connect(host='127.0.0.1', ssl=False) + + # Authentication must not report an unobserved failure after + # cancellation, even if the transport closes afterwards. + await asyncio.sleep(0) + gc.collect() + async def test_connection_isinstance(self): self.assertTrue(isinstance(self.con, pg_connection.Connection)) self.assertTrue(isinstance(self.con, object)) diff --git a/tests/test_pool.py b/tests/test_pool.py index 27e90123..7b141918 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -6,6 +6,7 @@ import asyncio +import gc import inspect import os import pathlib @@ -14,6 +15,7 @@ import textwrap import time import unittest +import weakref import asyncpg from asyncpg import _testbase as tb @@ -43,6 +45,13 @@ async def _cancel(self, waiter): class TestPool(tb.ConnectedTestCase): + async def wait_for_pool_size(self, pool, size): + async def wait(): + while pool.get_size() != size: + await asyncio.sleep(0.01) + + await asyncio.wait_for(wait(), 5) + async def test_pool_01(self): for n in {1, 5, 10, 20, 100}: with self.subTest(tasksnum=n): @@ -770,7 +779,7 @@ async def test_pool_min_size_partial_keep(self): database='postgres', init_size=3, min_size=1, max_size=3, max_inactive_connection_lifetime=0.2) as pool: - # Force all 3 connections to be created by acquiring them all. + # Exercise release timers for all three connections. c1 = await pool.acquire() c2 = await pool.acquire() c3 = await pool.acquire() @@ -782,11 +791,230 @@ async def test_pool_min_size_partial_keep(self): await asyncio.sleep(0.5) - # At least min_size (1) connection must survive. - alive = sum( - 1 for h in pool._holders if h._con is not None - ) - self.assertGreaterEqual(alive, 1) + # Only min_size (1) connection should survive. + self.assertEqual(pool.get_size(), 1) + + async def test_pool_implicit_init_size(self): + # Existing small, large, and lazy pool configurations still work. + for size in (0, 1, 20): + with self.subTest(size=size): + async with asyncpg.create_pool( + **self.get_connection_spec(), + min_size=size, max_size=max(1, size), + ) as pool: + self.assertEqual(pool.get_init_size(), size) + self.assertEqual(pool.get_size(), size) + + async def test_pool_min_size_reuses_retained_connections(self): + async with self.create_pool( + init_size=3, min_size=1, max_size=3, + max_inactive_connection_lifetime=0.05, + ) as pool: + await self.wait_for_pool_size(pool, 1) + retained = next(h._con for h in pool._holders if h.is_connected()) + + for _ in range(2): + async with pool.acquire() as con: + self.assertIs(con._con, retained) + self.assertEqual(await con.fetchval('SELECT 42'), 42) + await asyncio.sleep(0.1) + self.assertEqual(pool.get_size(), 1) + + async def test_pool_shutdown_cleans_idle_connections(self): + class WeakPool(pg_pool.Pool): + pass + + for action in ('close', 'terminate'): + for floor in (0, 1): + with self.subTest(action=action, floor=floor): + pool = await tb.create_pool( + **self.get_connection_spec(), pool_class=WeakPool, + init_size=1, min_size=floor, max_size=1, + max_inactive_connection_lifetime=0.05, + ) + try: + if action == 'close': + await pool.close() + else: + pool.terminate() + self.assertIsNone(pool._holders[0]._con) + self.assertIsNone(pool._holders[0]._inactive_callback) + finally: + pool.terminate() + + ref = weakref.ref(pool) + del pool + gc.collect() + self.assertIsNone(ref()) + + async def test_pool_min_size_restored_after_recycling(self): + for cause in ('max_queries', 'expire', 'close', 'terminate'): + with self.subTest(cause=cause): + async with self.create_pool( + init_size=1, min_size=1, max_size=1, + max_queries=1 if cause == 'max_queries' else 50000, + max_inactive_connection_lifetime=0, + ) as pool: + async with pool.acquire() as con: + old_con = con._con + await con.fetchval('SELECT 42') + if cause == 'expire': + await pool.expire_connections() + elif cause == 'close': + await con.close() + elif cause == 'terminate': + con.terminate() + + await self.wait_for_pool_size(pool, 1) + self.assertIsNot(pool._holders[0]._con, old_con) + self.assertTrue(old_con.is_closed()) + + async def test_pool_min_size_restored_after_idle_connection_loss(self): + async with self.create_pool( + init_size=1, min_size=1, max_size=1, + max_inactive_connection_lifetime=0, + ) as pool: + old_con = pool._holders[0]._con + terminated = asyncio.Event() + old_con.add_termination_listener(lambda con: terminated.set()) + + await self.con.execute( + 'SELECT pg_terminate_backend($1)', old_con.get_server_pid()) + await asyncio.wait_for(terminated.wait(), 5) + await self.wait_for_pool_size(pool, 1) + self.assertIsNot(pool._holders[0]._con, old_con) + self.assertEqual(await pool.fetchval('SELECT 42'), 42) + + async def test_pool_min_size_retries_failed_reconnect(self): + offline = False + attempted = asyncio.Event() + + async def connect(*args, **kwargs): + if offline: + attempted.set() + raise OSError('server temporarily unavailable') + return await pg_connection.connect(*args, **kwargs) + + async with self.create_pool( + init_size=1, min_size=1, max_size=1, connect=connect, + ) as pool: + offline = True + with self.assertLogs('asyncpg.pool', level='WARNING'): + pool._holders[0]._con.terminate() + await asyncio.wait_for(attempted.wait(), 5) + self.assertEqual(pool.get_size(), 0) + + offline = False + await self.wait_for_pool_size(pool, 1) + self.assertEqual(await pool.fetchval('SELECT 42'), 42) + + async def test_pool_floor_reconnect_honors_expired_generation(self): + started = asyncio.Event() + resume = asyncio.Event() + attempts = 0 + + async def init(con): + nonlocal attempts + attempts += 1 + if attempts == 2: + started.set() + await resume.wait() + + async with self.create_pool( + init_size=1, min_size=1, max_size=1, init=init, + server_settings={'application_name': 'old_pool_args'}, + ) as pool: + pool._holders[0]._con.terminate() + await asyncio.wait_for(started.wait(), 5) + pool.set_connect_args(**self.get_connection_spec({ + 'server_settings': {'application_name': 'new_pool_args'}, + })) + await pool.expire_connections() + resume.set() + await self.wait_for_pool_size(pool, 1) + + async with pool.acquire() as con: + self.assertEqual(con.get_settings().application_name, + 'new_pool_args') + + async def test_pool_min_size_retries_closed_connection(self): + offline = False + attempted = asyncio.Event() + + async def init(con): + if offline: + await con.close() + attempted.set() + + async with self.create_pool( + init_size=1, min_size=1, max_size=1, init=init, + ) as pool: + offline = True + with self.assertLogs('asyncpg.pool', level='WARNING'): + pool._holders[0]._con.terminate() + await asyncio.wait_for(attempted.wait(), 5) + self.assertEqual(pool.get_size(), 0) + + offline = False + await self.wait_for_pool_size(pool, 1) + self.assertEqual(await pool.fetchval('SELECT 42'), 42) + + async def test_pool_min_size_restored_after_multiple_losses(self): + async with self.create_pool( + init_size=3, min_size=3, max_size=5, + ) as pool: + old_connections = { + h._con for h in pool._holders if h.is_connected() + } + for con in old_connections: + con.terminate() + + async def worker(): + async with pool.acquire() as con: + self.assertLessEqual(pool.get_size(), pool.get_max_size()) + self.assertEqual(await con.fetchval('SELECT 42'), 42) + + await asyncio.gather(*(worker() for _ in range(10))) + self.assertGreaterEqual(pool.get_size(), pool.get_min_size()) + self.assertTrue(all(h._con not in old_connections + for h in pool._holders if h.is_connected())) + + async def test_pool_shutdown_cancels_floor_reconnect(self): + for action in ('close', 'terminate'): + with self.subTest(action=action): + started = asyncio.Event() + cancelled = asyncio.Event() + + async def connect(*args, **kwargs): + if started.is_set(): + self.fail('duplicate background connection attempt') + if pool is not None: + started.set() + try: + await asyncio.Future() + finally: + cancelled.set() + return await pg_connection.connect(*args, **kwargs) + + pool = None + pool = await self.create_pool( + init_size=1, min_size=1, max_size=1, connect=connect, + ) + pool._holders[0]._con.terminate() + await asyncio.wait_for(started.wait(), 5) + + # The background connector reserves the holder, so an + # acquirer must wait rather than exceed max_size. + with self.assertRaises(asyncio.TimeoutError): + await pool.acquire(timeout=0.05) + + if action == 'close': + await pool.close() + else: + pool.terminate() + await asyncio.wait_for(cancelled.wait(), 5) + self.assertTrue(pool.is_closing()) + self.assertEqual(pool.get_size(), 0) async def test_pool_min_size_zero_allows_full_expiry(self): # When min_size=0, all idle connections are allowed to expire. From b284a92226f7032f6b3f737bc7e0c52bc6dd87e3 Mon Sep 17 00:00:00 2001 From: Elvis Pranskevichus Date: Wed, 30 Sep 2026 13:39:24 -0700 Subject: [PATCH 8/9] Prevent garbage collection from resurrecting connection pools --- asyncpg/connection.py | 4 ++-- asyncpg/pool.py | 6 ++++-- tests/test_pool.py | 47 +++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 53 insertions(+), 4 deletions(-) diff --git a/asyncpg/connection.py b/asyncpg/connection.py index 3a70585d..3b69ad0f 100644 --- a/asyncpg/connection.py +++ b/asyncpg/connection.py @@ -1603,8 +1603,8 @@ def _cleanup(self): if self._pool_holder is not None: # Idle connections have no proxy, but still belong to a holder. - holder, self._pool_holder = self._pool_holder, None - if holder._con is self: + holder, self._pool_holder = self._pool_holder(), None + if holder is not None and holder._con is self: holder._release_on_close() self._mark_stmts_as_closed() diff --git a/asyncpg/pool.py b/asyncpg/pool.py index 02ba2e9d..b0bce0ff 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -14,6 +14,7 @@ from types import TracebackType from typing import Any, Optional, Type import warnings +import weakref from . import compat from . import connection @@ -132,7 +133,7 @@ class PoolConnectionHolder: '_max_queries', '_setup', '_max_inactive_time', '_in_use', '_inactive_callback', '_timeout', - '_generation') + '_generation', '__weakref__') def __init__( self, @@ -176,7 +177,8 @@ async def connect(self) -> None: raise exceptions.ConnectionDoesNotExistError( 'connection was closed during pool initialization') self._con = con - con._pool_holder = self + # A collected pool must not be resurrected by connection cleanup. + con._pool_holder = weakref.ref(self) self._generation = generation self._maybe_cancel_inactive_callback() self._setup_inactive_callback() diff --git a/tests/test_pool.py b/tests/test_pool.py index ea6c4ac9..bc6180e0 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -901,6 +901,53 @@ class WeakPool(pg_pool.Pool): gc.collect() self.assertIsNone(ref()) + async def test_pool_gc_does_not_restore_min_size(self): + for inactive_timeout in (0, 0.05): + with self.subTest(inactive_timeout=inactive_timeout): + connections = [] + maintenance_pools = [] + + class WeakPool(pg_pool.Pool): + async def _maintain_min_size(self): + maintenance_pools.append(self) + await super()._maintain_min_size() + + async def connect(*args, **kwargs): + con = await pg_connection.connect(*args, **kwargs) + connections.append(weakref.ref(con)) + return con + + pool = await tb.create_pool( + **self.get_connection_spec(), pool_class=WeakPool, + init_size=2, min_size=2, max_size=2, connect=connect, + max_inactive_connection_lifetime=inactive_timeout, + ) + ref = weakref.ref(pool) + try: + if inactive_timeout: + await asyncio.sleep(inactive_timeout * 2) + self.assertTrue(all(h._inactive_callback is None + for h in pool._holders)) + del pool + with self.assertWarnsRegex(ResourceWarning, + 'unclosed connection'): + gc.collect() + + # Give any incorrectly scheduled maintenance task time + # to reconnect, then ensure no new connections appeared. + await asyncio.sleep(0.1) + self.assertFalse(maintenance_pools) + self.assertEqual(len(connections), 2) + self.assertIsNone(ref()) + self.assertTrue(all(con() is None for con in connections)) + finally: + remaining = ref() + if remaining is not None: + await remaining.close() + for remaining in maintenance_pools: + await remaining.close() + maintenance_pools.clear() + async def test_pool_min_size_restored_after_recycling(self): for cause in ('max_queries', 'expire', 'close', 'terminate'): with self.subTest(cause=cause): From b47f3e22968d42c8532fd30cb6abc117358d46e7 Mon Sep 17 00:00:00 2001 From: Elvis Pranskevichus Date: Wed, 30 Sep 2026 14:10:48 -0700 Subject: [PATCH 9/9] Prevent pool maintenance during connection finalization --- asyncpg/connection.py | 4 ++++ asyncpg/pool.py | 2 +- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/asyncpg/connection.py b/asyncpg/connection.py index 3b69ad0f..c4067fc7 100644 --- a/asyncpg/connection.py +++ b/asyncpg/connection.py @@ -136,6 +136,10 @@ def __del__(self): warnings.warn(msg, ResourceWarning) if not self._loop.is_closed(): + # A weak holder reference may still be live during GC. + # Finalization must not notify the pool and restart + # maintenance. + self._pool_holder = None self.terminate() async def add_listener(self, channel, callback): diff --git a/asyncpg/pool.py b/asyncpg/pool.py index b0bce0ff..29e9cb9d 100644 --- a/asyncpg/pool.py +++ b/asyncpg/pool.py @@ -177,7 +177,7 @@ async def connect(self) -> None: raise exceptions.ConnectionDoesNotExistError( 'connection was closed during pool initialization') self._con = con - # A collected pool must not be resurrected by connection cleanup. + # Notify live holders without keeping an abandoned pool alive. con._pool_holder = weakref.ref(self) self._generation = generation self._maybe_cancel_inactive_callback()