Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 46 additions & 18 deletions tests/test_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import asyncio
import gc
import inspect
import logging
import os
import pathlib
import platform
Expand Down Expand Up @@ -1593,28 +1594,55 @@ async def test_full_reconnect_on_node_change_role(self):
# Force a new connection to be created
await pool.fetchval('SELECT 1')

await self.simulate_cluster_recovery_mode()
logs = []

# current pool connection info cache is expired,
# but we don't know it yet
with self.assertRaises(asyncpg.TargetServerAttributeNotMatched) as cm:
await pool.execute('SELECT 1')
class CaptureLogs(logging.Handler):
def emit(self, record):
logs.append(record)

self.assertEqual(
cm.exception.args[0],
"None of the hosts match the target attribute requirement "
"<SessionAttribute.primary: 'primary'>"
)
# Maintenance may race foreground acquisition, so retries are
# optional. Capture their expected warnings through pool shutdown.
with mock.patch.object(pg_pool.logger, 'handlers', [CaptureLogs()]), \
mock.patch.object(pg_pool.logger, 'propagate', False):
try:
await self.simulate_cluster_recovery_mode()

# force reconnect
with self.assertRaises(asyncpg.TargetServerAttributeNotMatched) as cm:
await pool.execute('SELECT 1')
# current pool connection info cache is expired,
# but we don't know it yet
with self.assertRaises(
asyncpg.TargetServerAttributeNotMatched) as cm:
await pool.execute('SELECT 1')

self.assertEqual(
cm.exception.args[0],
"None of the hosts match the target attribute requirement "
"<SessionAttribute.primary: 'primary'>"
)
self.assertEqual(
cm.exception.args[0],
"None of the hosts match the target attribute requirement "
"<SessionAttribute.primary: 'primary'>"
)

# force reconnect
with self.assertRaises(
asyncpg.TargetServerAttributeNotMatched) as cm:
await pool.execute('SELECT 1')

self.assertEqual(
cm.exception.args[0],
"None of the hosts match the target attribute requirement "
"<SessionAttribute.primary: 'primary'>"
)
finally:
await pool.close()

for record in logs:
self.assertEqual(record.levelno, logging.WARNING)
self.assertEqual(record.getMessage(),
'Failed to restore the pool connection floor; '
'retrying')
self.assertIsNotNone(record.exc_info)
self.assertIsInstance(record.exc_info[1], (
OSError,
asyncpg.CannotConnectNowError,
asyncpg.TargetServerAttributeNotMatched,
))


@unittest.skipIf(os.environ.get('PGHOST'), 'using remote cluster for testing')
Expand Down
Loading