From a1dc6932182157224d3d8ca348273ac155ebd01f Mon Sep 17 00:00:00 2001 From: Abhinav Rastogi Date: Wed, 7 Oct 2026 02:34:56 +0530 Subject: [PATCH 1/3] fix: avoid failover on client-local errors --- src/typesense/async_/api_call.py | 16 ++++++++++++++++ src/typesense/sync/api_call.py | 16 ++++++++++++++++ tests/api_call_test.py | 19 +++++++++++++++++++ 3 files changed, 51 insertions(+) diff --git a/src/typesense/async_/api_call.py b/src/typesense/async_/api_call.py index be1a83d..a85f6a7 100644 --- a/src/typesense/async_/api_call.py +++ b/src/typesense/async_/api_call.py @@ -135,6 +135,20 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): ServiceUnavailable, ) +_CLIENT_ERRORS: typing.Final[ + typing.Tuple[ + typing.Type[httpx.PoolTimeout], + typing.Type[httpx.LocalProtocolError], + typing.Type[httpx.DecodingError], + typing.Type[httpx.TooManyRedirects], + ] +] = ( + httpx.PoolTimeout, + httpx.LocalProtocolError, + httpx.DecodingError, + httpx.TooManyRedirects, +) + class AsyncApiCall: """ @@ -478,6 +492,8 @@ async def _execute_request( as_json, **request_kwargs, ) + except _CLIENT_ERRORS: + raise except _SERVER_ERRORS as server_error: self.node_manager.set_node_health(node, is_healthy=False) if num_retries < self.config.num_retries: diff --git a/src/typesense/sync/api_call.py b/src/typesense/sync/api_call.py index 402a0dc..1290774 100644 --- a/src/typesense/sync/api_call.py +++ b/src/typesense/sync/api_call.py @@ -135,6 +135,20 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): ServiceUnavailable, ) +_CLIENT_ERRORS: typing.Final[ + typing.Tuple[ + typing.Type[httpx.PoolTimeout], + typing.Type[httpx.LocalProtocolError], + typing.Type[httpx.DecodingError], + typing.Type[httpx.TooManyRedirects], + ] +] = ( + httpx.PoolTimeout, + httpx.LocalProtocolError, + httpx.DecodingError, + httpx.TooManyRedirects, +) + class ApiCall: """ @@ -478,6 +492,8 @@ def _execute_request( as_json, **request_kwargs, ) + except _CLIENT_ERRORS: + raise except _SERVER_ERRORS as server_error: self.node_manager.set_node_health(node, is_healthy=False) if num_retries < self.config.num_retries: diff --git a/tests/api_call_test.py b/tests/api_call_test.py index b7c4888..280bb33 100644 --- a/tests/api_call_test.py +++ b/tests/api_call_test.py @@ -461,6 +461,25 @@ def test_selects_next_available_node_on_timeout( assert len(respx.calls) == 3 +def test_client_errors_do_not_mark_nodes_unhealthy( + fake_api_call: ApiCall, + mocker: MockerFixture, +) -> None: + """Pool exhaustion is local to the client and must not trigger failover.""" + node = fake_api_call.node_manager.get_node() + make_request = mocker.patch.object( + fake_api_call.request_handler, + "make_request", + side_effect=httpx.PoolTimeout("No connection available"), + ) + + with pytest.raises(httpx.PoolTimeout): + fake_api_call.get("/test", as_json=True, entity_type=typing.Dict[str, str]) + + assert node.healthy is True + make_request.assert_called_once() + + def test_get_node_no_healthy_nodes( fake_api_call: ApiCall, mocker: MockFixture, From a9be046056efdf85fbde441799934d5f572a7c37 Mon Sep 17 00:00:00 2001 From: Fanis Tharropoulos Date: Wed, 7 Oct 2026 13:08:41 +0300 Subject: [PATCH 2/3] test: cover every client-local error on both clients (#143) --- tests/api_call_test.py | 58 ++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 58 insertions(+) diff --git a/tests/api_call_test.py b/tests/api_call_test.py index 280bb33..ea60b5b 100644 --- a/tests/api_call_test.py +++ b/tests/api_call_test.py @@ -684,3 +684,61 @@ async def test_async_sleeps_retry_interval_between_retries( assert sleep_call == mocker.call( fake_async_api_call.config.retry_interval_seconds, ) + + +@pytest.mark.parametrize( + "client_side_error", + [ + httpx.PoolTimeout("Pool timeout"), + httpx.LocalProtocolError("Local protocol error"), + httpx.DecodingError("Decoding error"), + httpx.TooManyRedirects("Too many redirects"), + ], +) +def test_client_side_error_does_not_mark_node_unhealthy( + fake_api_call: ApiCall, + client_side_error: httpx.HTTPError, +) -> None: + """Test that client-side httpx errors propagate without failing over.""" + with respx.mock: + respx.get("http://nearest:8108/").mock(side_effect=client_side_error) + node0_route = respx.get("http://node0:8108/").mock( + return_value=httpx.Response(200, json={"key": "value"}), + ) + + with pytest.raises(type(client_side_error)): + fake_api_call.get("/", entity_type=typing.Dict[str, str]) + + assert len(respx.calls) == 1 + assert not node0_route.called + + assert fake_api_call.config.nearest_node.healthy is True + + +@pytest.mark.parametrize( + "client_side_error", + [ + httpx.PoolTimeout("Pool timeout"), + httpx.LocalProtocolError("Local protocol error"), + httpx.DecodingError("Decoding error"), + httpx.TooManyRedirects("Too many redirects"), + ], +) +async def test_async_client_side_error_does_not_mark_node_unhealthy( + fake_async_api_call: AsyncApiCall, + client_side_error: httpx.HTTPError, +) -> None: + """Test that client-side httpx errors propagate without failing over (async).""" + with respx.mock: + respx.get("http://nearest:8108/").mock(side_effect=client_side_error) + node0_route = respx.get("http://node0:8108/").mock( + return_value=httpx.Response(200, json={"key": "value"}), + ) + + with pytest.raises(type(client_side_error)): + await fake_async_api_call.get("/", entity_type=typing.Dict[str, str]) + + assert len(respx.calls) == 1 + assert not node0_route.called + + assert fake_async_api_call.config.nearest_node.healthy is True From 3d0d70c03cbfdd48bde619229832eb44bd188b4a Mon Sep 17 00:00:00 2001 From: Fanis Tharropoulos Date: Tue, 6 Oct 2026 13:38:13 +0300 Subject: [PATCH 3/3] fix: mark the node that answered as healthy (#144) --- src/typesense/async_/api_call.py | 9 ++-- src/typesense/sync/api_call.py | 9 ++-- tests/api_call_test.py | 70 ++++++++++++++++++++++++++++++++ 3 files changed, 78 insertions(+), 10 deletions(-) diff --git a/src/typesense/async_/api_call.py b/src/typesense/async_/api_call.py index a85f6a7..8ad68ec 100644 --- a/src/typesense/async_/api_call.py +++ b/src/typesense/async_/api_call.py @@ -487,6 +487,7 @@ async def _execute_request( try: return await self._make_request_and_process_response( method, + node, url, entity_type, as_json, @@ -511,12 +512,13 @@ async def _execute_request( async def _make_request_and_process_response( self, method: str, + node: Node, url: str, entity_type: typing.Type[TEntityDict], as_json: bool, **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: - """Make the async API request and process the response.""" + """Make the async API request to `node` and process the response.""" request_response = await self.request_handler.make_request( method=method, url=url, @@ -525,10 +527,7 @@ async def _make_request_and_process_response( client=self._client, **kwargs, ) - self.node_manager.set_node_health( - self.node_manager.get_node(), - is_healthy=True, - ) + self.node_manager.set_node_health(node, is_healthy=True) return ( typing.cast(TEntityDict, request_response) if as_json diff --git a/src/typesense/sync/api_call.py b/src/typesense/sync/api_call.py index 1290774..f65d320 100644 --- a/src/typesense/sync/api_call.py +++ b/src/typesense/sync/api_call.py @@ -487,6 +487,7 @@ def _execute_request( try: return self._make_request_and_process_response( method, + node, url, entity_type, as_json, @@ -511,12 +512,13 @@ def _execute_request( def _make_request_and_process_response( self, method: str, + node: Node, url: str, entity_type: typing.Type[TEntityDict], as_json: bool, **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: - """Make the async API request and process the response.""" + """Make the async API request to `node` and process the response.""" request_response = self.request_handler.make_request( method=method, url=url, @@ -525,10 +527,7 @@ def _make_request_and_process_response( client=self._client, **kwargs, ) - self.node_manager.set_node_health( - self.node_manager.get_node(), - is_healthy=True, - ) + self.node_manager.set_node_health(node, is_healthy=True) return ( typing.cast(TEntityDict, request_response) if as_json diff --git a/tests/api_call_test.py b/tests/api_call_test.py index ea60b5b..fb4857b 100644 --- a/tests/api_call_test.py +++ b/tests/api_call_test.py @@ -742,3 +742,73 @@ async def test_async_client_side_error_does_not_mark_node_unhealthy( assert not node0_route.called assert fake_async_api_call.config.nearest_node.healthy is True + + +def test_round_robin_visits_each_node_in_turn(fake_api_call: ApiCall) -> None: + """Test that successful requests advance the round-robin by one node each.""" + fake_api_call.config.nearest_node = None + + with respx.mock: + for host in ("node0", "node1", "node2"): + respx.get(f"http://{host}:8108/").mock( + return_value=httpx.Response(200, json={"key": "value"}), + ) + + for _ in range(6): + fake_api_call.get("/", entity_type=typing.Dict[str, str]) + + assert [str(call.request.url) for call in respx.calls] == [ + "http://node0:8108/", + "http://node1:8108/", + "http://node2:8108/", + "http://node0:8108/", + "http://node1:8108/", + "http://node2:8108/", + ] + + +async def test_async_round_robin_visits_each_node_in_turn( + fake_async_api_call: AsyncApiCall, +) -> None: + """Test that successful requests advance the round-robin by one node each (async).""" + fake_async_api_call.config.nearest_node = None + + with respx.mock: + for host in ("node0", "node1", "node2"): + respx.get(f"http://{host}:8108/").mock( + return_value=httpx.Response(200, json={"key": "value"}), + ) + + for _ in range(6): + await fake_async_api_call.get("/", entity_type=typing.Dict[str, str]) + + assert [str(call.request.url) for call in respx.calls] == [ + "http://node0:8108/", + "http://node1:8108/", + "http://node2:8108/", + "http://node0:8108/", + "http://node1:8108/", + "http://node2:8108/", + ] + + +def test_success_marks_only_the_answering_node_healthy( + fake_api_call: ApiCall, +) -> None: + """Test that a success refreshes the node that answered and no other.""" + fake_api_call.config.nearest_node = None + answering_node, unhealthy_node, _ = fake_api_call.node_manager.nodes + answering_node.last_access_ts = 0 + unhealthy_node.healthy = False + unhealthy_node.last_access_ts = int(time.time()) + + with respx.mock: + respx.get("http://node0:8108/").mock( + return_value=httpx.Response(200, json={"key": "value"}), + ) + + fake_api_call.get("/", entity_type=typing.Dict[str, str]) + + assert answering_node.healthy is True + assert answering_node.last_access_ts > 0 + assert unhealthy_node.healthy is False