Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
## [Unreleased]

### Added


- Support additional health check for pool instances (PR #348).

### Changed

### Fixed
Expand Down
81 changes: 61 additions & 20 deletions tarantool/connection_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -398,7 +398,8 @@ def __init__(self,
connection_timeout=CONNECTION_TIMEOUT,
strategy_class=RoundRobinStrategy,
refresh_delay=POOL_REFRESH_DELAY,
fetch_schema=True):
fetch_schema=True,
additional_health_check=None):
"""
:param addrs: List of dictionaries describing server addresses:

Expand Down Expand Up @@ -476,6 +477,17 @@ def __init__(self,
:param fetch_schema: Refer to
:paramref:`~tarantool.Connection.params.fetch_schema`.

:param additional_health_check: An additional health check for each
pool instance. It is called with the instance connection after
the built-in checks succeed.
Return :attr:`Status.UNHEALTHY` to exclude the
instance from request routing until the next successful health
check; return :attr:`Status.HEALTHY` to keep it available. The
check is performed when the pool connects and during periodic
state refreshes. It should be fast and must not raise exceptions;
the pool does not handle callback errors.
:type additional_health_check: :obj:`typing.Callable[[Connection], Status]`, optional

:raise: :exc:`~tarantool.error.ConfigurationError`,
:class:`~tarantool.Connection` exceptions

Expand All @@ -497,6 +509,8 @@ def __init__(self,
new_addrs.append(new_addr)
self.addrs = new_addrs

self.additional_health_check = additional_health_check

# Create connections
self.pool = {}
self.refresh_delay = refresh_delay
Expand Down Expand Up @@ -551,58 +565,85 @@ def _make_key(self, addr):
return f"{addr['host']}:{addr['port']}"
return addr['socket_fd']

def _get_new_state(self, unit):
def _status_check(self, conn, unit):
"""
Get new pool server state.
Check the status of the connection using ``box.info``.

:param unit: Server metainfo.
:param conn: Connection to the instance.
:type conn: :class:`~tarantool.Connection`

:param unit: Instance metainfo.
:type unit: :class:`~tarantool.connection_pool.PoolUnit`

:rtype: :class:`~tarantool.connection_pool.InstanceState`
:return: A pair containing the read-only flag and health status.
:rtype: :obj:`tuple[bool, ~tarantool.connection_pool.Status]`

:meta private:
"""

conn = unit.conn

if conn.is_closed():
try:
conn.connect()
except NetworkError as exc:
msg = (f"Failed to connect to {unit.get_address()}, "
f"reason: {repr(exc)}")
warn(msg, ClusterConnectWarning)
return InstanceState(Status.UNHEALTHY)

try:
resp = conn.call('box.info')
except DatabaseError as exc:
msg = (f"Failed to get box.info for {unit.get_address()}, "
f"reason: {repr(exc)}")
warn(msg, PoolTopologyWarning)
return InstanceState(Status.UNHEALTHY)
return False, Status.UNHEALTHY

try:
read_only = resp.data[0]['ro']
except (IndexError, KeyError) as exc:
msg = (f"Incorrect box.info response from {unit.get_address()}"
f"reason: {repr(exc)}")
warn(msg, PoolTopologyWarning)
return InstanceState(Status.UNHEALTHY)
return False, Status.UNHEALTHY

try:
status = resp.data[0]['status']

if status != 'running':
msg = f"{unit.get_address()} instance status is not 'running'"
warn(msg, PoolTopologyWarning)
return InstanceState(Status.UNHEALTHY)
return read_only, Status.UNHEALTHY
except (IndexError, KeyError) as exc:
msg = (f"Incorrect box.info response from {unit.get_address()}"
f"reason: {repr(exc)}")
warn(msg, PoolTopologyWarning)
return read_only, Status.UNHEALTHY

return read_only, Status.HEALTHY

def _get_new_state(self, unit):
"""
Get new pool server state.

:param unit: Server metainfo.
:type unit: :class:`~tarantool.connection_pool.PoolUnit`

:rtype: :class:`~tarantool.connection_pool.InstanceState`

:meta private:
"""

conn = unit.conn

if conn.is_closed():
try:
conn.connect()
except NetworkError as exc:
msg = (f"Failed to connect to {unit.get_address()}, "
f"reason: {repr(exc)}")
warn(msg, ClusterConnectWarning)
return InstanceState(Status.UNHEALTHY)

read_only, status = self._status_check(conn, unit)
if status == Status.UNHEALTHY:
return InstanceState(Status.UNHEALTHY)

if self.additional_health_check is not None:
status = self.additional_health_check(conn)

if status == Status.UNHEALTHY:
return InstanceState(Status.UNHEALTHY)

return InstanceState(Status.HEALTHY, read_only)

def _refresh_state(self, key):
Expand Down
79 changes: 78 additions & 1 deletion test/suites/test_pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -594,14 +594,24 @@ def test_17_instance_bootstrap_error_does_not_kill_refresh(self):
warnings.simplefilter('ignore', category=PoolTolopogyWarning)

self.set_cluster_ro([False, True, True, True, True])
callback_calls = 0
target_port = self.addrs[0]['port']

def additional_health_check(conn):
nonlocal callback_calls
if conn.port == target_port:
callback_calls += 1
return Status.HEALTHY

self.pool = tarantool.ConnectionPool(
addrs=self.addrs,
user='test',
password='test',
refresh_delay=0.2)
refresh_delay=0.2,
additional_health_check=additional_health_check)

self.pool.ping(mode=tarantool.Mode.RW)
callback_calls = 0

unit = self.pool.pool[f"{self.addrs[0]['host']}:{self.addrs[0]['port']}"]

Expand All @@ -622,6 +632,7 @@ def expect_instance_unhealthy_and_refresh_alive():
self.assertTrue(unit.thread.is_alive(),
'refresh thread died on a DatabaseError')
self.assertEqual(unit.state.status, Status.UNHEALTHY)
self.assertEqual(callback_calls, 0)

self.retry(func=expect_instance_unhealthy_and_refresh_alive)

Expand All @@ -636,6 +647,72 @@ def expect_rw_request_succeed():

self.retry(func=expect_rw_request_succeed)

def test_18_additional_health_check_excludes_and_recovers_instance(self):
self.set_cluster_ro([False, False, True, False, False])
target_addr = self.addrs[2]
target_key = f"{target_addr['host']}:{target_addr['port']}"
unhealthy = True

def additional_health_check(conn):
if unhealthy and conn.port == target_addr['port']:
return Status.UNHEALTHY
return Status.HEALTHY

self.pool = tarantool.ConnectionPool(
addrs=self.addrs,
user='test',
password='test',
refresh_delay=0.2,
additional_health_check=additional_health_check)

self.assertEqual(self.pool.pool[target_key].state.status,
Status.UNHEALTHY)
with self.assertRaises(PoolTopologyError):
self.pool.ping(mode=tarantool.Mode.RO)

unhealthy = False

def expect_instance_recovered():
self.assertEqual(self.pool.pool[target_key].state.status,
Status.HEALTHY)
self.pool.ping(mode=tarantool.Mode.RO)

self.retry(func=expect_instance_recovered)

def test_19_additional_health_check_not_called_before_running(self):
warnings.simplefilter('ignore', category=PoolTopologyWarning)

self.set_cluster_ro([False, True, True, True, True])
target_addr = self.addrs[0]
callback_calls = 0

resp = self.servers[0].admin(r"""
rawset(_G, 'box_info_backup', box.info)
box.info = function()
local info = box_info_backup()
return {ro = info.ro, status = 'loading'}
end
return true
""")
assert_admin_success(resp)

def additional_health_check(conn):
nonlocal callback_calls
if conn.port == target_addr['port']:
callback_calls += 1
return Status.HEALTHY

self.pool = tarantool.ConnectionPool(
addrs=self.addrs,
user='test',
password='test',
additional_health_check=additional_health_check)

self.assertEqual(callback_calls, 0)
self.assertEqual(
self.pool.pool[f"{target_addr['host']}:{target_addr['port']}"].state.status,
Status.UNHEALTHY)

def tearDown(self):
if self.pool:
self.pool.close()
Expand Down
Loading