From a1dc6932182157224d3d8ca348273ac155ebd01f Mon Sep 17 00:00:00 2001 From: Abhinav Rastogi Date: Wed, 7 Oct 2026 02:34:56 +0530 Subject: [PATCH 1/6] 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/6] 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/6] 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 From 6aae5005d66b7a6c64bdd771fffaa9b0025cac38 Mon Sep 17 00:00:00 2001 From: Fanis Tharropoulos Date: Tue, 6 Oct 2026 13:58:47 +0300 Subject: [PATCH 4/6] feat: configure the connection pool and cap requests in flight (#146) --- src/typesense/async_/api_call.py | 30 ++++-- src/typesense/concurrency_limit.py | 89 ++++++++++++++++ src/typesense/configuration.py | 63 ++++++++++++ src/typesense/sync/api_call.py | 30 ++++-- tests/api_call_test.py | 131 ++++++++++++++++++++++++ tests/configuration_test.py | 43 ++++++++ tests/configuration_validations_test.py | 37 +++++++ utils/run-unasync.py | 2 + 8 files changed, 407 insertions(+), 18 deletions(-) create mode 100644 src/typesense/concurrency_limit.py diff --git a/src/typesense/async_/api_call.py b/src/typesense/async_/api_call.py index 8ad68ec..916b9bb 100644 --- a/src/typesense/async_/api_call.py +++ b/src/typesense/async_/api_call.py @@ -37,6 +37,7 @@ import httpx +from typesense.concurrency_limit import AsyncConcurrencyLimit from typesense.configuration import Configuration, Node from typesense.exceptions import ( HTTPStatus0Error, @@ -174,7 +175,17 @@ def __init__(self, config: Configuration): self.node_manager = NodeManager(config) self.request_handler = RequestHandler(config) self._client = httpx.AsyncClient( - timeout=config.connection_timeout_seconds, + timeout=httpx.Timeout( + config.connection_timeout_seconds, + pool=config.pool_timeout_seconds, + ), + limits=httpx.Limits( + max_connections=config.max_connections, + max_keepalive_connections=config.max_keepalive_connections, + ), + ) + self._concurrency_limit = AsyncConcurrencyLimit( + config.max_concurrent_requests, ) async def __aenter__(self) -> "AsyncApiCall": @@ -519,14 +530,15 @@ async def _make_request_and_process_response( **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: """Make the async API request to `node` and process the response.""" - request_response = await self.request_handler.make_request( - method=method, - url=url, - as_json=as_json, - entity_type=entity_type, - client=self._client, - **kwargs, - ) + async with self._concurrency_limit: + request_response = await self.request_handler.make_request( + method=method, + url=url, + as_json=as_json, + entity_type=entity_type, + client=self._client, + **kwargs, + ) self.node_manager.set_node_health(node, is_healthy=True) return ( typing.cast(TEntityDict, request_response) diff --git a/src/typesense/concurrency_limit.py b/src/typesense/concurrency_limit.py new file mode 100644 index 0000000..34ebdc5 --- /dev/null +++ b/src/typesense/concurrency_limit.py @@ -0,0 +1,89 @@ +""" +Optional caps on the number of requests a client sends at once. + +``AsyncConcurrencyLimit`` is used by the async client and ``ConcurrencyLimit`` by the +sync client (``utils/run-unasync.py`` maps one name to the other). Both are no-ops +when ``max_concurrent_requests`` is ``None``. + +Keeping the cap below the httpx pool's ``max_connections`` means requests queue here +instead of in the pool, so a burst of slow requests cannot exhaust the pool and +raise ``httpx.PoolTimeout``. +""" + +import asyncio +import sys +import threading +from types import TracebackType + +if sys.version_info >= (3, 11): + import typing +else: + import typing_extensions as typing + + +class AsyncConcurrencyLimit: + """Async context manager that holds a slot for the duration of a request.""" + + def __init__(self, max_concurrent_requests: typing.Optional[int]) -> None: + """ + Initialize the limit. + + Args: + max_concurrent_requests (Optional[int]): The maximum number of requests + in flight at once, or ``None`` for no limit. + """ + self._max_concurrent_requests = max_concurrent_requests + # Created on first use, inside the running event loop. On Python < 3.10 a + # semaphore binds to the loop that is current when it is constructed. + self._semaphore: typing.Optional[asyncio.Semaphore] = None + + async def __aenter__(self) -> None: + """Wait for a free slot.""" + if self._max_concurrent_requests is None: + return + if self._semaphore is None: + self._semaphore = asyncio.Semaphore(self._max_concurrent_requests) + await self._semaphore.acquire() + + async def __aexit__( + self, + exc_type: typing.Optional[typing.Type[BaseException]], + exc_val: typing.Optional[BaseException], + exc_tb: typing.Optional[TracebackType], + ) -> None: + """Release the slot.""" + if self._semaphore is not None: + self._semaphore.release() + + +class ConcurrencyLimit: + """Context manager that holds a slot for the duration of a request.""" + + def __init__(self, max_concurrent_requests: typing.Optional[int]) -> None: + """ + Initialize the limit. + + Args: + max_concurrent_requests (Optional[int]): The maximum number of requests + in flight at once, or ``None`` for no limit. + """ + self._semaphore: typing.Optional[threading.Semaphore] = ( + None + if max_concurrent_requests is None + else threading.Semaphore(max_concurrent_requests) + ) + + def __enter__(self) -> None: + """Wait for a free slot.""" + if self._semaphore is not None: + self._semaphore.acquire() + + def __exit__( + self, + exc_type: typing.Optional[typing.Type[BaseException]], + exc_val: typing.Optional[BaseException], + exc_tb: typing.Optional[TracebackType], + ) -> None: + """Release the slot.""" + if self._semaphore is not None: + self._semaphore.release() diff --git a/src/typesense/configuration.py b/src/typesense/configuration.py index aaa741e..4f9144a 100644 --- a/src/typesense/configuration.py +++ b/src/typesense/configuration.py @@ -82,6 +82,23 @@ class ConfigDict(typing.TypedDict): connection_timeout_seconds (float): The connection timeout in seconds. suppress_deprecation_warnings (bool): Whether to suppress deprecation warnings. + + pool_timeout_seconds (float): How long a request waits for a free connection + in the pool before raising ``httpx.PoolTimeout``. Defaults to + ``connection_timeout_seconds``. Setting it lower than + ``connection_timeout_seconds`` makes the httpcore connection leak + (encode/httpcore#1093) more likely under load. + + max_connections (int): The maximum number of connections in the pool. + Defaults to 100. + + max_keepalive_connections (int): The maximum number of idle connections + kept alive in the pool. Defaults to 20. + + max_concurrent_requests (int): The maximum number of requests in flight at + once; further requests wait for a slot. Keep it below + ``max_connections`` so a burst of slow requests cannot exhaust the pool. + Defaults to no limit. """ nodes: typing.List[typing.Union[str, NodeConfigDict]] @@ -100,6 +117,10 @@ class ConfigDict(typing.TypedDict): ] # deprecated connection_timeout_seconds: typing.NotRequired[float] suppress_deprecation_warnings: typing.NotRequired[bool] + pool_timeout_seconds: typing.NotRequired[float] + max_connections: typing.NotRequired[int] + max_keepalive_connections: typing.NotRequired[int] + max_concurrent_requests: typing.NotRequired[int] class Node: @@ -188,6 +209,10 @@ class Configuration: retry_interval_seconds (float): The interval in seconds between retries. healthcheck_interval_seconds (int): The interval in seconds between health checks. verify (bool): Whether to verify the SSL certificate. + pool_timeout_seconds (float): How long to wait for a free pooled connection. + max_connections (int): The maximum number of connections in the pool. + max_keepalive_connections (int): The maximum number of idle pooled connections. + max_concurrent_requests (int | None): The maximum number of requests in flight. """ def __init__( @@ -232,6 +257,18 @@ def __init__( self.suppress_deprecation_warnings = config_dict.get( "suppress_deprecation_warnings", False ) + self.pool_timeout_seconds = config_dict.get( + "pool_timeout_seconds", + self.connection_timeout_seconds, + ) + self.max_connections = config_dict.get("max_connections", 100) + self.max_keepalive_connections = config_dict.get( + "max_keepalive_connections", + 20, + ) + self.max_concurrent_requests: typing.Optional[int] = config_dict.get( + "max_concurrent_requests", + ) def _handle_nearest_node( self, @@ -295,6 +332,32 @@ def validate_config_dict(config_dict: ConfigDict) -> None: if nearest_node: ConfigurationValidations.validate_nearest_node(nearest_node) + ConfigurationValidations.validate_connection_pool(config_dict) + + @staticmethod + def validate_connection_pool(config_dict: ConfigDict) -> None: + """ + Validate the connection pool and concurrency settings. + + Args: + config_dict (ConfigDict): The configuration dictionary to validate. + + Raises: + ConfigError: If a pool or concurrency setting is out of range. + """ + positive_settings: typing.Dict[str, typing.Optional[float]] = { + "pool_timeout_seconds": config_dict.get("pool_timeout_seconds"), + "max_connections": config_dict.get("max_connections"), + "max_concurrent_requests": config_dict.get("max_concurrent_requests"), + } + for key, config_value in positive_settings.items(): + if config_value is not None and config_value <= 0: + raise ConfigError(f"`{key}` must be greater than 0.") + + max_keepalive_connections = config_dict.get("max_keepalive_connections") + if max_keepalive_connections is not None and max_keepalive_connections < 0: + raise ConfigError("`max_keepalive_connections` must not be negative.") + @staticmethod def validate_required_config_fields(config_dict: ConfigDict) -> None: """ diff --git a/src/typesense/sync/api_call.py b/src/typesense/sync/api_call.py index f65d320..4184599 100644 --- a/src/typesense/sync/api_call.py +++ b/src/typesense/sync/api_call.py @@ -37,6 +37,7 @@ import httpx +from typesense.concurrency_limit import ConcurrencyLimit from typesense.configuration import Configuration, Node from typesense.exceptions import ( HTTPStatus0Error, @@ -174,7 +175,17 @@ def __init__(self, config: Configuration): self.node_manager = NodeManager(config) self.request_handler = RequestHandler(config) self._client = httpx.Client( - timeout=config.connection_timeout_seconds, + timeout=httpx.Timeout( + config.connection_timeout_seconds, + pool=config.pool_timeout_seconds, + ), + limits=httpx.Limits( + max_connections=config.max_connections, + max_keepalive_connections=config.max_keepalive_connections, + ), + ) + self._concurrency_limit = ConcurrencyLimit( + config.max_concurrent_requests, ) def __enter__(self) -> "ApiCall": @@ -519,14 +530,15 @@ def _make_request_and_process_response( **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: """Make the async API request to `node` and process the response.""" - request_response = self.request_handler.make_request( - method=method, - url=url, - as_json=as_json, - entity_type=entity_type, - client=self._client, - **kwargs, - ) + with self._concurrency_limit: + request_response = self.request_handler.make_request( + method=method, + url=url, + as_json=as_json, + entity_type=entity_type, + client=self._client, + **kwargs, + ) self.node_manager.set_node_health(node, is_healthy=True) return ( typing.cast(TEntityDict, request_response) diff --git a/tests/api_call_test.py b/tests/api_call_test.py index fb4857b..9c77d03 100644 --- a/tests/api_call_test.py +++ b/tests/api_call_test.py @@ -1,8 +1,11 @@ """Unit Tests for the ApiCall class.""" +import asyncio import logging import sys +import threading import time +from concurrent.futures import ThreadPoolExecutor from pytest_mock import MockFixture @@ -812,3 +815,131 @@ def test_success_marks_only_the_answering_node_healthy( assert answering_node.healthy is True assert answering_node.last_access_ts > 0 assert unhealthy_node.healthy is False + + +def test_client_uses_connection_pool_settings( + fake_config: Configuration, + mocker: MockerFixture, +) -> None: + """Test that the httpx client is built from the connection pool settings.""" + client_mock = mocker.patch("typesense.sync.api_call.httpx.Client") + fake_config.connection_timeout_seconds = 3.0 + fake_config.pool_timeout_seconds = 1.5 + fake_config.max_connections = 200 + fake_config.max_keepalive_connections = 50 + + ApiCall(fake_config) + + client_mock.assert_called_once_with( + timeout=httpx.Timeout(3.0, pool=1.5), + limits=httpx.Limits(max_connections=200, max_keepalive_connections=50), + ) + + +def test_async_client_uses_connection_pool_settings( + fake_config: Configuration, + mocker: MockerFixture, +) -> None: + """Test that the httpx async client is built from the connection pool settings.""" + client_mock = mocker.patch("typesense.async_.api_call.httpx.AsyncClient") + fake_config.connection_timeout_seconds = 3.0 + fake_config.pool_timeout_seconds = 1.5 + fake_config.max_connections = 200 + fake_config.max_keepalive_connections = 50 + + AsyncApiCall(fake_config) + + client_mock.assert_called_once_with( + timeout=httpx.Timeout(3.0, pool=1.5), + limits=httpx.Limits(max_connections=200, max_keepalive_connections=50), + ) + + +def _count_requests_in_flight( + concurrent_requests: int, + max_concurrent_requests: typing.Optional[int], + fake_config: Configuration, +) -> int: + """Send requests from several threads and return the peak number in flight.""" + fake_config.max_concurrent_requests = max_concurrent_requests + api_call = ApiCall(fake_config) + lock = threading.Lock() + in_flight = 0 + peak = 0 + + def slow_response(request: httpx.Request) -> httpx.Response: + nonlocal in_flight, peak + with lock: + in_flight += 1 + peak = max(peak, in_flight) + time.sleep(0.05) + with lock: + in_flight -= 1 + return httpx.Response(200, json={"key": "value"}) + + with respx.mock: + respx.get("http://nearest:8108/").mock(side_effect=slow_response) + with ThreadPoolExecutor(max_workers=concurrent_requests) as executor: + for _ in range(concurrent_requests): + executor.submit(api_call.get, "/", entity_type=typing.Dict[str, str]) + + return peak + + +def test_max_concurrent_requests_caps_requests_in_flight( + fake_config: Configuration, +) -> None: + """Test that no more than ``max_concurrent_requests`` requests are in flight.""" + assert _count_requests_in_flight(6, 2, fake_config) == 2 + + +def test_requests_in_flight_are_unlimited_by_default( + fake_config: Configuration, +) -> None: + """Test that requests are not capped when ``max_concurrent_requests`` is unset.""" + assert _count_requests_in_flight(6, None, fake_config) == 6 + + +async def _async_count_requests_in_flight( + concurrent_requests: int, + max_concurrent_requests: typing.Optional[int], + fake_config: Configuration, +) -> int: + """Send concurrent async requests and return the peak number in flight.""" + fake_config.max_concurrent_requests = max_concurrent_requests + api_call = AsyncApiCall(fake_config) + in_flight = 0 + peak = 0 + + async def slow_response(request: httpx.Request) -> httpx.Response: + nonlocal in_flight, peak + in_flight += 1 + peak = max(peak, in_flight) + await asyncio.sleep(0.01) + in_flight -= 1 + return httpx.Response(200, json={"key": "value"}) + + with respx.mock: + respx.get("http://nearest:8108/").mock(side_effect=slow_response) + await asyncio.gather( + *( + api_call.get("/", entity_type=typing.Dict[str, str]) + for _ in range(concurrent_requests) + ), + ) + + return peak + + +async def test_async_max_concurrent_requests_caps_requests_in_flight( + fake_config: Configuration, +) -> None: + """Test that no more than ``max_concurrent_requests`` requests are in flight (async).""" + assert await _async_count_requests_in_flight(6, 2, fake_config) == 2 + + +async def test_async_requests_in_flight_are_unlimited_by_default( + fake_config: Configuration, +) -> None: + """Test that async requests are not capped when ``max_concurrent_requests`` is unset.""" + assert await _async_count_requests_in_flight(6, None, fake_config) == 6 diff --git a/tests/configuration_test.py b/tests/configuration_test.py index 626c477..092c93b 100644 --- a/tests/configuration_test.py +++ b/tests/configuration_test.py @@ -207,3 +207,46 @@ def test_configuration_invalid_nearest_node_url() -> None: match="Node URL does not contain the port.", ): Configuration(config) + + +def test_configuration_connection_pool_defaults() -> None: + """Test the connection pool defaults, with the pool timeout following the connection timeout.""" + configuration = Configuration( + { + "nodes": [DEFAULT_NODE], + "api_key": "xyz", + "connection_timeout_seconds": 7.0, + }, + ) + + expected = { + "pool_timeout_seconds": 7.0, + "max_connections": 100, + "max_keepalive_connections": 20, + "max_concurrent_requests": None, + } + + assert_to_contain_object(configuration, expected) + + +def test_configuration_connection_pool_explicit() -> None: + """Test the connection pool settings with explicit values.""" + configuration = Configuration( + { + "nodes": [DEFAULT_NODE], + "api_key": "xyz", + "pool_timeout_seconds": 1.5, + "max_connections": 200, + "max_keepalive_connections": 50, + "max_concurrent_requests": 150, + }, + ) + + expected = { + "pool_timeout_seconds": 1.5, + "max_connections": 200, + "max_keepalive_connections": 50, + "max_concurrent_requests": 150, + } + + assert_to_contain_object(configuration, expected) diff --git a/tests/configuration_validations_test.py b/tests/configuration_validations_test.py index d408e05..8cf8061 100644 --- a/tests/configuration_validations_test.py +++ b/tests/configuration_validations_test.py @@ -1,7 +1,13 @@ """Tests for the ConfigurationValidations class.""" +import sys import types +if sys.version_info >= (3, 11): + import typing +else: + import typing_extensions as typing + import pytest from typesense.configuration import ConfigDict, ConfigurationValidations @@ -199,3 +205,34 @@ def test_validate_config_dict_with_wrong_nearest_node() -> None: "api_key": "xyz", }, ) + + +@pytest.mark.parametrize( + ("key", "config_value", "message"), + [ + ("pool_timeout_seconds", 0, "`pool_timeout_seconds` must be greater than 0."), + ("max_connections", 0, "`max_connections` must be greater than 0."), + ( + "max_concurrent_requests", + -1, + "`max_concurrent_requests` must be greater than 0.", + ), + ( + "max_keepalive_connections", + -1, + "`max_keepalive_connections` must not be negative.", + ), + ], +) +def test_validate_config_dict_with_invalid_connection_pool( + key: str, + config_value: float, + message: str, +) -> None: + """Test validate_config_dict with out-of-range connection pool settings.""" + config_dict = {"nodes": [DEFAULT_NODE], "api_key": "xyz", key: config_value} + + with pytest.raises(ConfigError, match=message): + ConfigurationValidations.validate_config_dict( + typing.cast(ConfigDict, config_dict), + ) diff --git a/utils/run-unasync.py b/utils/run-unasync.py index aa4dcbd..7d836d4 100644 --- a/utils/run-unasync.py +++ b/utils/run-unasync.py @@ -27,6 +27,8 @@ def collect_class_replacements(source_dir: Path) -> dict[str, str]: # client (unasync strips ``await``); map the module token so the import and call # are rewritten too. replacements["asyncio"] = "time" + # Defined in the shared ``typesense.concurrency_limit`` module, outside async_. + replacements["AsyncConcurrencyLimit"] = "ConcurrencyLimit" return replacements From 977e8e2042cbe77450faba074f81e1d352fd1c22 Mon Sep 17 00:00:00 2001 From: Fanis Tharropoulos Date: Tue, 6 Oct 2026 14:27:12 +0300 Subject: [PATCH 5/6] fix: resolve typesense imports in mypy and fix the type errors it surfaces --- setup.cfg | 1 + src/typesense/async_/analytics_rule_v1.py | 8 +-- src/typesense/async_/analytics_rules_v1.py | 18 +++--- src/typesense/async_/api_call.py | 4 +- src/typesense/async_/client.py | 4 +- src/typesense/async_/documents.py | 68 +++++++++++++++------ src/typesense/async_/keys.py | 5 +- src/typesense/async_/operations.py | 51 ++++++++-------- src/typesense/async_/override.py | 2 +- src/typesense/async_/overrides.py | 2 +- src/typesense/async_/synonym.py | 2 +- src/typesense/async_/synonyms.py | 2 +- src/typesense/preprocess.py | 7 ++- src/typesense/request_handler.py | 69 +++++++++++++++++----- src/typesense/sync/analytics_rule_v1.py | 8 +-- src/typesense/sync/analytics_rules_v1.py | 18 +++--- src/typesense/sync/api_call.py | 4 +- src/typesense/sync/client.py | 4 +- src/typesense/sync/documents.py | 68 +++++++++++++++------ src/typesense/sync/keys.py | 5 +- src/typesense/sync/operations.py | 51 ++++++++-------- src/typesense/sync/override.py | 2 +- src/typesense/sync/overrides.py | 2 +- src/typesense/sync/synonym.py | 2 +- src/typesense/sync/synonyms.py | 2 +- 25 files changed, 255 insertions(+), 154 deletions(-) diff --git a/setup.cfg b/setup.cfg index 088736f..72beb40 100644 --- a/setup.cfg +++ b/setup.cfg @@ -56,6 +56,7 @@ enable_error_code = redundant-self, explicit_package_bases = true +mypy_path = src ignore_missing_imports = true strict = true warn_unreachable = true diff --git a/src/typesense/async_/analytics_rule_v1.py b/src/typesense/async_/analytics_rule_v1.py index d640853..5584623 100644 --- a/src/typesense/async_/analytics_rule_v1.py +++ b/src/typesense/async_/analytics_rule_v1.py @@ -74,11 +74,9 @@ async def retrieve( Union[RuleSchemaForQueries, RuleSchemaForCounters]: The schema containing the rule details. """ - response: typing.Union[ - RuleSchemaForQueries, RuleSchemaForCounters - ] = await self.api_call.get( + response = await self.api_call.get( self._endpoint_path, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], as_json=True, ) return typing.cast( @@ -101,7 +99,7 @@ async def delete(self) -> RuleDeleteSchema: return response @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "AsyncAnalyticsRuleV1 is deprecated on v30+. Use client.analytics.rules[rule_id] instead.", flag_name="analytics_rules_v1_deprecation", ) diff --git a/src/typesense/async_/analytics_rules_v1.py b/src/typesense/async_/analytics_rules_v1.py index 1aac207..2e905e4 100644 --- a/src/typesense/async_/analytics_rules_v1.py +++ b/src/typesense/async_/analytics_rules_v1.py @@ -89,7 +89,7 @@ def __getitem__(self, rule_id: str) -> AsyncAnalyticsRuleV1: self.rules[rule_id] = AsyncAnalyticsRuleV1(self.api_call, rule_id) return self.rules[rule_id] - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "AsyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) @@ -115,21 +115,19 @@ async def create( The created rule. Returns RuleSchemaForCounters for counter rules and RuleSchemaForQueries for query rules. """ - response: typing.Union[ - RuleSchemaForCounters, RuleSchemaForQueries - ] = await self.api_call.post( + response = await self.api_call.post( AsyncAnalyticsRulesV1.resource_path, body=rule, params=rule_parameters, as_json=True, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], ) return typing.cast( typing.Union[RuleSchemaForCounters, RuleSchemaForQueries], response, ) - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "AsyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) @@ -148,19 +146,17 @@ async def upsert( Returns: Union[RuleSchemaForCounters, RuleCreateSchemaForQueries]: The upserted rule. """ - response: typing.Union[ - RuleSchemaForCounters, RuleCreateSchemaForQueries - ] = await self.api_call.put( + response = await self.api_call.put( "/".join([self.resource_path, rule_id]), body=rule, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], ) return typing.cast( typing.Union[RuleSchemaForCounters, RuleCreateSchemaForQueries], response, ) - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "AsyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) diff --git a/src/typesense/async_/api_call.py b/src/typesense/async_/api_call.py index 916b9bb..9bc98ec 100644 --- a/src/typesense/async_/api_call.py +++ b/src/typesense/async_/api_call.py @@ -60,7 +60,7 @@ import typing_extensions as typing TEntityDict = typing.TypeVar("TEntityDict") -TParams = typing.TypeVar("TParams", bound=typing.Dict[str, typing.Any]) +TParams = typing.TypeVar("TParams", bound=typing.Mapping[str, object]) TBody = typing.TypeVar( "TBody", bound=typing.Union[str, bytes, typing.Mapping[str, typing.Any]] ) @@ -95,7 +95,7 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): params: typing.NotRequired[typing.Union[TParams, None]] data: typing.NotRequired[typing.Union[TBody, None]] - content: typing.NotRequired[typing.Union[TBody, str, None]] + content: typing.NotRequired[typing.Union[str, bytes, None]] headers: typing.NotRequired[typing.Dict[str, str]] timeout: typing.NotRequired[float] diff --git a/src/typesense/async_/client.py b/src/typesense/async_/client.py index 1ecb807..8175bcb 100644 --- a/src/typesense/async_/client.py +++ b/src/typesense/async_/client.py @@ -164,5 +164,5 @@ def typed_collection( """ if name is None: name = model.__name__.lower() - collection: AsyncCollection[TDoc] = self.collections[name] - return collection + # ``collections`` is typed for the default DocumentSchema; narrow it to the model. + return typing.cast(AsyncCollection[TDoc], self.collections[name]) diff --git a/src/typesense/async_/documents.py b/src/typesense/async_/documents.py index 8228762..399c82d 100644 --- a/src/typesense/async_/documents.py +++ b/src/typesense/async_/documents.py @@ -61,6 +61,17 @@ None, ] +# One line of an import response. ``ImportResponse`` is a union of lists, one per +# return mode, so the helpers below build a list of these and ``import_`` casts it +# to the list type its overloads promise. +_ImportResponseItem = typing.Union[ + ImportResponseWithDoc[TDoc], + ImportResponseWithId, + ImportResponseWithDocAndId[TDoc], + ImportResponseSuccess, + ImportResponseFail[TDoc], +] + class AsyncDocuments(typing.Generic[TDoc]): """ @@ -125,12 +136,14 @@ async def create( Returns: TDoc: The created document. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "create" + write_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "create", + } response = await self.api_call.post( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=write_parameters, as_json=True, entity_type=typing.Dict[str, str], ) @@ -154,7 +167,14 @@ async def create_many( The list of import responses. """ logger.warn("`create_many` is deprecated: please use `import_`.") - return await self.import_(documents, dirty_values_parameters) + # Dirty values parameters are a subset of the write parameters. + return await self.import_( + documents, + typing.cast( + typing.Optional[DocumentWriteParameters], + dirty_values_parameters, + ), + ) async def upsert( self, @@ -172,12 +192,14 @@ async def upsert( Returns: TDoc: The upserted document. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "upsert" + write_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "upsert", + } response = await self.api_call.post( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=write_parameters, as_json=True, entity_type=typing.Dict[str, str], ) @@ -199,12 +221,14 @@ async def update( Returns: UpdateByFilterResponse: The response containing information about the update. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "update" + update_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "update", + } response: UpdateByFilterResponse = await self.api_call.patch( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=update_parameters, entity_type=UpdateByFilterResponse, ) return response @@ -301,9 +325,14 @@ async def import_( return await self._import_raw(documents, import_parameters) if batch_size: - return await self._batch_import(documents, import_parameters, batch_size) - - return await self._bulk_import(documents, import_parameters) + response_objs = await self._batch_import( + documents, + import_parameters, + batch_size, + ) + else: + response_objs = await self._bulk_import(documents, import_parameters) + return typing.cast(ImportResponse[TDoc], response_objs) async def export( self, @@ -410,9 +439,9 @@ async def _batch_import( documents: typing.List[TDoc], import_parameters: _ImportParameters, batch_size: int, - ) -> ImportResponse[TDoc]: + ) -> typing.List[_ImportResponseItem[TDoc]]: """Import documents in batches.""" - response_objs: ImportResponse[TDoc] = [] + response_objs: typing.List[_ImportResponseItem[TDoc]] = [] for batch_index in range(0, len(documents), batch_size): batch = documents[batch_index : batch_index + batch_size] api_response = await self._bulk_import(batch, import_parameters) @@ -423,7 +452,7 @@ async def _bulk_import( self, documents: typing.List[TDoc], import_parameters: _ImportParameters, - ) -> ImportResponse[TDoc]: + ) -> typing.List[_ImportResponseItem[TDoc]]: """Import a list of documents in bulk.""" document_strs = [json.dumps(doc) for doc in documents] if not document_strs: @@ -439,9 +468,12 @@ async def _bulk_import( ) return self._parse_import_response(res) - def _parse_import_response(self, response: str) -> ImportResponse[TDoc]: + def _parse_import_response( + self, + response: str, + ) -> typing.List[_ImportResponseItem[TDoc]]: """Parse the import response string into a list of response objects.""" - response_objs: typing.List[ImportResponse] = [] + response_objs: typing.List[_ImportResponseItem[TDoc]] = [] for res_obj_str in response.split("\n"): try: res_obj_json = json.loads(res_obj_str) diff --git a/src/typesense/async_/keys.py b/src/typesense/async_/keys.py index 0dd8d94..4639113 100644 --- a/src/typesense/async_/keys.py +++ b/src/typesense/async_/keys.py @@ -29,7 +29,6 @@ ApiKeyCreateResponseSchema, ApiKeyCreateSchema, ApiKeyRetrieveSchema, - ApiKeySchema, ) if sys.version_info >= (3, 11): @@ -103,11 +102,11 @@ async def create(self, schema: ApiKeyCreateSchema) -> ApiKeyCreateResponseSchema ... } ... ) """ - response: ApiKeySchema = await self.api_call.post( + response: ApiKeyCreateResponseSchema = await self.api_call.post( AsyncKeys.resource_path, as_json=True, body=schema, - entity_type=ApiKeySchema, + entity_type=ApiKeyCreateResponseSchema, ) return response diff --git a/src/typesense/async_/operations.py b/src/typesense/async_/operations.py index ca61a1f..4a36608 100644 --- a/src/typesense/async_/operations.py +++ b/src/typesense/async_/operations.py @@ -60,8 +60,10 @@ def __init__(self, api_call: AsyncApiCall): """ self.api_call = api_call + # The generic ``str`` overload below also matches "schema_changes"; overloads are + # tried in order, so this one wins. @typing.overload - async def perform( + async def perform( # type: ignore[overload-overlap] self, operation_name: typing.Literal["schema_changes"], query_params: None = None, @@ -132,36 +134,36 @@ async def perform( @typing.overload async def perform( self, - operation_name: str, - query_params: typing.Union[typing.Dict[str, str], None] = None, + operation_name: typing.Literal["snapshot"], + query_params: SnapshotParameters, ) -> OperationResponse: """ - Perform a generic operation. + Perform a snapshot operation. Args: - operation_name (str): The name of the operation. - query_params (Union[Dict[str, str], None], optional): - Query parameters for the operation. + operation_name (Literal["snapshot"]): The name of the operation. + query_params (SnapshotParameters): Query parameters for the snapshot operation. Returns: - OperationResponse: The response from the operation. + OperationResponse: The response from the snapshot operation. """ @typing.overload async def perform( self, - operation_name: typing.Literal["snapshot"], - query_params: SnapshotParameters, + operation_name: str, + query_params: typing.Union[typing.Dict[str, str], None] = None, ) -> OperationResponse: """ - Perform a snapshot operation. + Perform a generic operation. Args: - operation_name (Literal["snapshot"]): The name of the operation. - query_params (SnapshotParameters): Query parameters for the snapshot operation. + operation_name (str): The name of the operation. + query_params (Union[Dict[str, str], None], optional): + Query parameters for the operation. Returns: - OperationResponse: The response from the snapshot operation. + OperationResponse: The response from the operation. """ async def perform( @@ -181,7 +183,7 @@ async def perform( typing.Dict[str, str], None, ] = None, - ) -> OperationResponse: + ) -> typing.Union[OperationResponse, typing.List[SchemaChangesResponse]]: """ Perform an operation on the Typesense API. @@ -202,13 +204,16 @@ async def perform( >>> response = await operations.perform("vote") >>> health = await operations.is_healthy() """ - response: OperationResponse = await self.api_call.post( + response = await self.api_call.post( self._endpoint_path(operation_name), params=query_params, as_json=True, - entity_type=OperationResponse, + entity_type=object, + ) + return typing.cast( + typing.Union[OperationResponse, typing.List[SchemaChangesResponse]], + response, ) - return response async def is_healthy(self) -> bool: """ @@ -222,16 +227,14 @@ async def is_healthy(self) -> bool: >>> healthy = await operations.is_healthy() >>> print(healthy) """ - call_resp: HealthCheckResponse = await self.api_call.get( + call_resp: object = await self.api_call.get( AsyncOperations.health_path, as_json=True, entity_type=HealthCheckResponse, ) - if isinstance(call_resp, typing.Dict): - is_ok: bool = call_resp.get("ok", False) - else: - is_ok = False - return is_ok + if isinstance(call_resp, dict): + return bool(call_resp.get("ok", False)) + return False async def toggle_slow_request_log( self, diff --git a/src/typesense/async_/override.py b/src/typesense/async_/override.py index 58e5a26..3e3e0f8 100644 --- a/src/typesense/async_/override.py +++ b/src/typesense/async_/override.py @@ -87,7 +87,7 @@ async def delete(self) -> OverrideDeleteSchema: return response @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The override API (collections/{collection}/overrides/{override_id}) is deprecated is removed on v30+. " "Use curation sets (curation_sets) instead.", flag_name="overrides_deprecation", diff --git a/src/typesense/async_/overrides.py b/src/typesense/async_/overrides.py index b8e725b..99d1190 100644 --- a/src/typesense/async_/overrides.py +++ b/src/typesense/async_/overrides.py @@ -129,7 +129,7 @@ async def retrieve(self) -> OverrideRetrieveSchema: ) return response - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "AsyncOverrides is deprecated on v30+. Use client.curation_sets instead.", flag_name="overrides_deprecation", ) diff --git a/src/typesense/async_/synonym.py b/src/typesense/async_/synonym.py index 3ad6bc2..73cd46c 100644 --- a/src/typesense/async_/synonym.py +++ b/src/typesense/async_/synonym.py @@ -79,7 +79,7 @@ async def delete(self) -> SynonymDeleteSchema: ) @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The synonym API (collections/{collection}/synonyms/{synonym_id}) is deprecated is removed on v30+. " "Use synonym sets (synonym_sets) instead.", flag_name="synonyms_deprecation", diff --git a/src/typesense/async_/synonyms.py b/src/typesense/async_/synonyms.py index 027172e..ea2b3aa 100644 --- a/src/typesense/async_/synonyms.py +++ b/src/typesense/async_/synonyms.py @@ -124,7 +124,7 @@ async def retrieve(self) -> SynonymsRetrieveSchema: ) return response - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The synonyms API (collections/{collection}/synonyms) is deprecated is removed on v30+. " "Use synonym sets (synonym_sets) instead.", flag_name="synonyms_deprecation", diff --git a/src/typesense/preprocess.py b/src/typesense/preprocess.py index b45db0c..15e13d7 100644 --- a/src/typesense/preprocess.py +++ b/src/typesense/preprocess.py @@ -110,7 +110,9 @@ def process_param_list( return ",".join(stringified_list) -def stringify_search_params(parameter_dict: ParamSchema) -> StringifiedParamSchema: +def stringify_search_params( + parameter_dict: typing.Mapping[str, object], +) -> StringifiedParamSchema: """ Convert the search parameters to strings. @@ -118,7 +120,8 @@ def stringify_search_params(parameter_dict: ParamSchema) -> StringifiedParamSche to their string representations. List values are converted to comma-separated strings. Args: - parameter_dict (ParamSchema): The search parameters. + parameter_dict (Mapping[str, object]): The search parameters, e.g. a + ``SearchParameters`` TypedDict or a ``ParamSchema`` dictionary. Returns: StringifiedParamSchema: The search parameters as strings. diff --git a/src/typesense/request_handler.py b/src/typesense/request_handler.py index 38e6c24..1e8d81e 100644 --- a/src/typesense/request_handler.py +++ b/src/typesense/request_handler.py @@ -48,8 +48,23 @@ ) TEntityDict = typing.TypeVar("TEntityDict") -TParams = typing.TypeVar("TParams", bound=typing.Dict[str, typing.Any]) -TBody = typing.TypeVar("TBody", bound=typing.Union[str, bytes]) +TParams = typing.TypeVar("TParams", bound=typing.Mapping[str, object]) +TBody = typing.TypeVar( + "TBody", bound=typing.Union[str, bytes, typing.Mapping[str, typing.Any]] +) + +# The query parameter values httpx accepts, once booleans are normalized to strings. +_QueryParams = typing.Mapping[ + str, + typing.Union[ + str, + int, + float, + bool, + None, + typing.Sequence[typing.Union[str, int, float, bool, None]], + ], +] _ERROR_CODE_MAP: typing.Mapping[str, typing.Type[TypesenseClientError]] = ( MappingProxyType( @@ -99,7 +114,7 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): data: typing.NotRequired[ typing.Union[TBody, str, typing.Dict[str, typing.Any], None] ] - content: typing.NotRequired[typing.Union[TBody, str, None]] + content: typing.NotRequired[typing.Union[str, bytes, None]] headers: typing.NotRequired[typing.Dict[str, str]] timeout: typing.NotRequired[float] @@ -128,6 +143,30 @@ def __init__(self, config: Configuration): """ self.config = config + @typing.overload + def make_request( + self, + *, + method: str, + url: str, + entity_type: typing.Type[TEntityDict], + as_json: typing.Union[typing.Literal[True], typing.Literal[False]] = True, + client: httpx.AsyncClient, + **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], + ) -> typing.Coroutine[typing.Any, typing.Any, typing.Union[TEntityDict, str]]: ... + + @typing.overload + def make_request( + self, + *, + method: str, + url: str, + entity_type: typing.Type[TEntityDict], + as_json: typing.Union[typing.Literal[True], typing.Literal[False]] = True, + client: httpx.Client, + **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], + ) -> typing.Union[TEntityDict, str]: ... + def make_request( self, *, @@ -183,9 +222,9 @@ def make_request( request_kwargs["params"] = params if body := kwargs.get("data"): - if not isinstance(body, (str, bytes)): - body = json.dumps(body) - request_kwargs["content"] = typing.cast(TBody, body) + request_kwargs["content"] = ( + body if isinstance(body, (str, bytes)) else json.dumps(body) + ) if isinstance(client, httpx.AsyncClient): return self._make_async_request( @@ -207,13 +246,13 @@ def _make_sync_request( ) -> typing.Union[TEntityDict, str]: """Make a synchronous HTTP request using httpx.Client.""" params: typing.Union[TParams, None] = request_kwargs.get("params") - content: typing.Union[TBody, str, None] = request_kwargs.get("content") + content: typing.Union[str, bytes, None] = request_kwargs.get("content") headers: typing.Dict[str, str] = request_kwargs.get("headers", {}) response = client.request( method, url, - params=params, + params=typing.cast(typing.Optional[_QueryParams], params), content=content, headers=headers, ) @@ -242,13 +281,13 @@ async def _make_async_request( ) -> typing.Union[TEntityDict, str]: """Make an asynchronous HTTP request using httpx.AsyncClient.""" params: typing.Union[TParams, None] = request_kwargs.get("params") - content: typing.Union[TBody, str, None] = request_kwargs.get("content") + content: typing.Union[str, bytes, None] = request_kwargs.get("content") headers: typing.Dict[str, str] = request_kwargs.get("headers", {}) response = await client.request( method, url, - params=params, + params=typing.cast(typing.Optional[_QueryParams], params), content=content, headers=headers, ) @@ -267,17 +306,19 @@ async def _make_async_request( return response.text @staticmethod - def normalize_params(params: typing.Dict[str, typing.Any]) -> None: + def normalize_params(params: typing.Mapping[str, object]) -> None: """ - Normalize boolean parameters in the request. + Normalize boolean parameters in the request, in place. Args: - params (Dict[str, Any]): The parameters to normalize. + params (Mapping[str, object]): The parameters to normalize. They are + typed as read-only so TypedDict parameters are accepted, but must + be a ``dict`` at runtime. Raises: ValueError: If params is not a dictionary. """ - if not isinstance(params, typing.Dict): + if not isinstance(params, dict): raise ValueError("Params must be a dictionary.") for key, parameter_value in params.items(): if isinstance(parameter_value, bool): diff --git a/src/typesense/sync/analytics_rule_v1.py b/src/typesense/sync/analytics_rule_v1.py index 38e8f41..9eec662 100644 --- a/src/typesense/sync/analytics_rule_v1.py +++ b/src/typesense/sync/analytics_rule_v1.py @@ -74,11 +74,9 @@ def retrieve( Union[RuleSchemaForQueries, RuleSchemaForCounters]: The schema containing the rule details. """ - response: typing.Union[ - RuleSchemaForQueries, RuleSchemaForCounters - ] = self.api_call.get( + response = self.api_call.get( self._endpoint_path, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], as_json=True, ) return typing.cast( @@ -101,7 +99,7 @@ def delete(self) -> RuleDeleteSchema: return response @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "SyncAnalyticsRuleV1 is deprecated on v30+. Use client.analytics.rules[rule_id] instead.", flag_name="analytics_rules_v1_deprecation", ) diff --git a/src/typesense/sync/analytics_rules_v1.py b/src/typesense/sync/analytics_rules_v1.py index e63f802..30edf75 100644 --- a/src/typesense/sync/analytics_rules_v1.py +++ b/src/typesense/sync/analytics_rules_v1.py @@ -89,7 +89,7 @@ def __getitem__(self, rule_id: str) -> AnalyticsRuleV1: self.rules[rule_id] = AnalyticsRuleV1(self.api_call, rule_id) return self.rules[rule_id] - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "SyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) @@ -115,21 +115,19 @@ def create( The created rule. Returns RuleSchemaForCounters for counter rules and RuleSchemaForQueries for query rules. """ - response: typing.Union[ - RuleSchemaForCounters, RuleSchemaForQueries - ] = self.api_call.post( + response = self.api_call.post( AnalyticsRulesV1.resource_path, body=rule, params=rule_parameters, as_json=True, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], ) return typing.cast( typing.Union[RuleSchemaForCounters, RuleSchemaForQueries], response, ) - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "SyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) @@ -148,19 +146,17 @@ def upsert( Returns: Union[RuleSchemaForCounters, RuleCreateSchemaForQueries]: The upserted rule. """ - response: typing.Union[ - RuleSchemaForCounters, RuleCreateSchemaForQueries - ] = self.api_call.put( + response = self.api_call.put( "/".join([self.resource_path, rule_id]), body=rule, - entity_type=dict, + entity_type=typing.Dict[str, typing.Any], ) return typing.cast( typing.Union[RuleSchemaForCounters, RuleCreateSchemaForQueries], response, ) - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "SyncAnalyticsRulesV1 is deprecated on v30+. Use client.analytics instead.", flag_name="analytics_rules_v1_deprecation", ) diff --git a/src/typesense/sync/api_call.py b/src/typesense/sync/api_call.py index 4184599..09c6107 100644 --- a/src/typesense/sync/api_call.py +++ b/src/typesense/sync/api_call.py @@ -60,7 +60,7 @@ import typing_extensions as typing TEntityDict = typing.TypeVar("TEntityDict") -TParams = typing.TypeVar("TParams", bound=typing.Dict[str, typing.Any]) +TParams = typing.TypeVar("TParams", bound=typing.Mapping[str, object]) TBody = typing.TypeVar( "TBody", bound=typing.Union[str, bytes, typing.Mapping[str, typing.Any]] ) @@ -95,7 +95,7 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): params: typing.NotRequired[typing.Union[TParams, None]] data: typing.NotRequired[typing.Union[TBody, None]] - content: typing.NotRequired[typing.Union[TBody, str, None]] + content: typing.NotRequired[typing.Union[str, bytes, None]] headers: typing.NotRequired[typing.Dict[str, str]] timeout: typing.NotRequired[float] diff --git a/src/typesense/sync/client.py b/src/typesense/sync/client.py index ef1afb0..b11e542 100644 --- a/src/typesense/sync/client.py +++ b/src/typesense/sync/client.py @@ -164,5 +164,5 @@ def typed_collection( """ if name is None: name = model.__name__.lower() - collection: Collection[TDoc] = self.collections[name] - return collection + # ``collections`` is typed for the default DocumentSchema; narrow it to the model. + return typing.cast(Collection[TDoc], self.collections[name]) diff --git a/src/typesense/sync/documents.py b/src/typesense/sync/documents.py index b22ef69..0c7d7f7 100644 --- a/src/typesense/sync/documents.py +++ b/src/typesense/sync/documents.py @@ -61,6 +61,17 @@ None, ] +# One line of an import response. ``ImportResponse`` is a union of lists, one per +# return mode, so the helpers below build a list of these and ``import_`` casts it +# to the list type its overloads promise. +_ImportResponseItem = typing.Union[ + ImportResponseWithDoc[TDoc], + ImportResponseWithId, + ImportResponseWithDocAndId[TDoc], + ImportResponseSuccess, + ImportResponseFail[TDoc], +] + class Documents(typing.Generic[TDoc]): """ @@ -125,12 +136,14 @@ def create( Returns: TDoc: The created document. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "create" + write_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "create", + } response = self.api_call.post( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=write_parameters, as_json=True, entity_type=typing.Dict[str, str], ) @@ -154,7 +167,14 @@ def create_many( The list of import responses. """ logger.warn("`create_many` is deprecated: please use `import_`.") - return self.import_(documents, dirty_values_parameters) + # Dirty values parameters are a subset of the write parameters. + return self.import_( + documents, + typing.cast( + typing.Optional[DocumentWriteParameters], + dirty_values_parameters, + ), + ) def upsert( self, @@ -172,12 +192,14 @@ def upsert( Returns: TDoc: The upserted document. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "upsert" + write_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "upsert", + } response = self.api_call.post( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=write_parameters, as_json=True, entity_type=typing.Dict[str, str], ) @@ -199,12 +221,14 @@ def update( Returns: UpdateByFilterResponse: The response containing information about the update. """ - dirty_values_parameters = dirty_values_parameters or {} - dirty_values_parameters["action"] = "update" + update_parameters: typing.Dict[str, object] = { + **(dirty_values_parameters or {}), + "action": "update", + } response: UpdateByFilterResponse = self.api_call.patch( self._endpoint_path(), body=document, - params=dirty_values_parameters, + params=update_parameters, entity_type=UpdateByFilterResponse, ) return response @@ -301,9 +325,14 @@ def import_( return self._import_raw(documents, import_parameters) if batch_size: - return self._batch_import(documents, import_parameters, batch_size) - - return self._bulk_import(documents, import_parameters) + response_objs = self._batch_import( + documents, + import_parameters, + batch_size, + ) + else: + response_objs = self._bulk_import(documents, import_parameters) + return typing.cast(ImportResponse[TDoc], response_objs) def export( self, @@ -410,9 +439,9 @@ def _batch_import( documents: typing.List[TDoc], import_parameters: _ImportParameters, batch_size: int, - ) -> ImportResponse[TDoc]: + ) -> typing.List[_ImportResponseItem[TDoc]]: """Import documents in batches.""" - response_objs: ImportResponse[TDoc] = [] + response_objs: typing.List[_ImportResponseItem[TDoc]] = [] for batch_index in range(0, len(documents), batch_size): batch = documents[batch_index : batch_index + batch_size] api_response = self._bulk_import(batch, import_parameters) @@ -423,7 +452,7 @@ def _bulk_import( self, documents: typing.List[TDoc], import_parameters: _ImportParameters, - ) -> ImportResponse[TDoc]: + ) -> typing.List[_ImportResponseItem[TDoc]]: """Import a list of documents in bulk.""" document_strs = [json.dumps(doc) for doc in documents] if not document_strs: @@ -439,9 +468,12 @@ def _bulk_import( ) return self._parse_import_response(res) - def _parse_import_response(self, response: str) -> ImportResponse[TDoc]: + def _parse_import_response( + self, + response: str, + ) -> typing.List[_ImportResponseItem[TDoc]]: """Parse the import response string into a list of response objects.""" - response_objs: typing.List[ImportResponse] = [] + response_objs: typing.List[_ImportResponseItem[TDoc]] = [] for res_obj_str in response.split("\n"): try: res_obj_json = json.loads(res_obj_str) diff --git a/src/typesense/sync/keys.py b/src/typesense/sync/keys.py index b70ec5e..9bd494e 100644 --- a/src/typesense/sync/keys.py +++ b/src/typesense/sync/keys.py @@ -29,7 +29,6 @@ ApiKeyCreateResponseSchema, ApiKeyCreateSchema, ApiKeyRetrieveSchema, - ApiKeySchema, ) if sys.version_info >= (3, 11): @@ -103,11 +102,11 @@ def create(self, schema: ApiKeyCreateSchema) -> ApiKeyCreateResponseSchema: ... } ... ) """ - response: ApiKeySchema = self.api_call.post( + response: ApiKeyCreateResponseSchema = self.api_call.post( Keys.resource_path, as_json=True, body=schema, - entity_type=ApiKeySchema, + entity_type=ApiKeyCreateResponseSchema, ) return response diff --git a/src/typesense/sync/operations.py b/src/typesense/sync/operations.py index e560b76..450c345 100644 --- a/src/typesense/sync/operations.py +++ b/src/typesense/sync/operations.py @@ -60,8 +60,10 @@ def __init__(self, api_call: ApiCall): """ self.api_call = api_call + # The generic ``str`` overload below also matches "schema_changes"; overloads are + # tried in order, so this one wins. @typing.overload - def perform( + def perform( # type: ignore[overload-overlap] self, operation_name: typing.Literal["schema_changes"], query_params: None = None, @@ -132,36 +134,36 @@ def perform( @typing.overload def perform( self, - operation_name: str, - query_params: typing.Union[typing.Dict[str, str], None] = None, + operation_name: typing.Literal["snapshot"], + query_params: SnapshotParameters, ) -> OperationResponse: """ - Perform a generic operation. + Perform a snapshot operation. Args: - operation_name (str): The name of the operation. - query_params (Union[Dict[str, str], None], optional): - Query parameters for the operation. + operation_name (Literal["snapshot"]): The name of the operation. + query_params (SnapshotParameters): Query parameters for the snapshot operation. Returns: - OperationResponse: The response from the operation. + OperationResponse: The response from the snapshot operation. """ @typing.overload def perform( self, - operation_name: typing.Literal["snapshot"], - query_params: SnapshotParameters, + operation_name: str, + query_params: typing.Union[typing.Dict[str, str], None] = None, ) -> OperationResponse: """ - Perform a snapshot operation. + Perform a generic operation. Args: - operation_name (Literal["snapshot"]): The name of the operation. - query_params (SnapshotParameters): Query parameters for the snapshot operation. + operation_name (str): The name of the operation. + query_params (Union[Dict[str, str], None], optional): + Query parameters for the operation. Returns: - OperationResponse: The response from the snapshot operation. + OperationResponse: The response from the operation. """ def perform( @@ -181,7 +183,7 @@ def perform( typing.Dict[str, str], None, ] = None, - ) -> OperationResponse: + ) -> typing.Union[OperationResponse, typing.List[SchemaChangesResponse]]: """ Perform an operation on the Typesense API. @@ -202,13 +204,16 @@ def perform( >>> response = await operations.perform("vote") >>> health = await operations.is_healthy() """ - response: OperationResponse = self.api_call.post( + response = self.api_call.post( self._endpoint_path(operation_name), params=query_params, as_json=True, - entity_type=OperationResponse, + entity_type=object, + ) + return typing.cast( + typing.Union[OperationResponse, typing.List[SchemaChangesResponse]], + response, ) - return response def is_healthy(self) -> bool: """ @@ -222,16 +227,14 @@ def is_healthy(self) -> bool: >>> healthy = await operations.is_healthy() >>> print(healthy) """ - call_resp: HealthCheckResponse = self.api_call.get( + call_resp: object = self.api_call.get( Operations.health_path, as_json=True, entity_type=HealthCheckResponse, ) - if isinstance(call_resp, typing.Dict): - is_ok: bool = call_resp.get("ok", False) - else: - is_ok = False - return is_ok + if isinstance(call_resp, dict): + return bool(call_resp.get("ok", False)) + return False def toggle_slow_request_log( self, diff --git a/src/typesense/sync/override.py b/src/typesense/sync/override.py index 8a24e9e..78aff10 100644 --- a/src/typesense/sync/override.py +++ b/src/typesense/sync/override.py @@ -87,7 +87,7 @@ def delete(self) -> OverrideDeleteSchema: return response @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The override API (collections/{collection}/overrides/{override_id}) is deprecated is removed on v30+. " "Use curation sets (curation_sets) instead.", flag_name="overrides_deprecation", diff --git a/src/typesense/sync/overrides.py b/src/typesense/sync/overrides.py index 7682ff5..99c4667 100644 --- a/src/typesense/sync/overrides.py +++ b/src/typesense/sync/overrides.py @@ -129,7 +129,7 @@ def retrieve(self) -> OverrideRetrieveSchema: ) return response - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "SyncOverrides is deprecated on v30+. Use client.curation_sets instead.", flag_name="overrides_deprecation", ) diff --git a/src/typesense/sync/synonym.py b/src/typesense/sync/synonym.py index d091fdd..27ff6cf 100644 --- a/src/typesense/sync/synonym.py +++ b/src/typesense/sync/synonym.py @@ -79,7 +79,7 @@ def delete(self) -> SynonymDeleteSchema: ) @property - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The synonym API (collections/{collection}/synonyms/{synonym_id}) is deprecated is removed on v30+. " "Use synonym sets (synonym_sets) instead.", flag_name="synonyms_deprecation", diff --git a/src/typesense/sync/synonyms.py b/src/typesense/sync/synonyms.py index d6e055b..f3c4e45 100644 --- a/src/typesense/sync/synonyms.py +++ b/src/typesense/sync/synonyms.py @@ -124,7 +124,7 @@ def retrieve(self) -> SynonymsRetrieveSchema: ) return response - @warn_deprecation( # type: ignore[untyped-decorator] + @warn_deprecation( "The synonyms API (collections/{collection}/synonyms) is deprecated is removed on v30+. " "Use synonym sets (synonym_sets) instead.", flag_name="synonyms_deprecation", From c2f3eb5400756a7988b55c8c390249ab50f8fdf3 Mon Sep 17 00:00:00 2001 From: Fanis Tharropoulos Date: Tue, 6 Oct 2026 14:17:56 +0300 Subject: [PATCH 6/6] feat: accept a user-supplied httpx or httpx2 client (#146) --- README.md | 27 ++++++ pyproject.toml | 4 + src/typesense/async_/api_call.py | 77 +++++++++------- src/typesense/async_/client.py | 15 ++- src/typesense/http_backend.py | 77 ++++++++++++++++ src/typesense/request_handler.py | 41 ++++++--- src/typesense/sync/api_call.py | 77 +++++++++------- src/typesense/sync/client.py | 15 ++- tests/http_backend_test.py | 151 ++++++++++++++++++++++++++++++ utils/run-unasync.py | 2 + uv.lock | 153 +++++++++++++++++++++++-------- 11 files changed, 515 insertions(+), 124 deletions(-) create mode 100644 src/typesense/http_backend.py create mode 100644 tests/http_backend_test.py diff --git a/README.md b/README.md index 9537087..fd919b6 100644 --- a/README.md +++ b/README.md @@ -42,6 +42,33 @@ if __name__ == "__main__": See `examples/async_collection_operations.py` for a fuller async walkthrough. +## Using httpx2 + +The client sends requests with [httpx](https://www.python-httpx.org/) by default. On Python 3.10+ you can pass an [httpx2](https://github.com/pydantic/httpx2) client instead. httpx2 is Pydantic's maintained continuation of httpx, and it fixes a connection pool leak in httpcore ([encode/httpcore#1093](https://github.com/encode/httpcore/issues/1093)) that can leave an `AsyncClient` failing every request with `PoolTimeout` under load. + +``` +$ pip install "typesense[httpx2]" +``` + +```python +import httpx2 +import typesense + +http_client = httpx2.AsyncClient( + timeout=httpx2.Timeout(2.0), + limits=httpx2.Limits(max_connections=100, max_keepalive_connections=20), +) +client = typesense.AsyncClient( + { + "api_key": "abcd", + "nodes": [{"host": "localhost", "port": "8108", "protocol": "http"}], + }, + http_client=http_client, +) +``` + +`typesense.Client` takes an `httpx2.Client` the same way. The connection pool settings in the config (`pool_timeout_seconds`, `max_connections`, `max_keepalive_connections`) only apply to the default client, so set them on your own client instead. The Typesense client does not close a client you pass in. + ## Compatibility | Typesense Server | typesense-python | diff --git a/pyproject.toml b/pyproject.toml index d918e31..bd2bf29 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -27,6 +27,10 @@ dependencies = [ ] dynamic = ["version"] +[project.optional-dependencies] +# Maintained continuation of httpx; fixes the httpcore pool leak in encode/httpcore#1093. +httpx2 = ["httpx2>=2.6.0; python_version >= '3.10'"] + [project.urls] Documentation = "https://typesense.org/" Source = "https://github.com/typesense/typesense-python" diff --git a/src/typesense/async_/api_call.py b/src/typesense/async_/api_call.py index 9bc98ec..6ded2ed 100644 --- a/src/typesense/async_/api_call.py +++ b/src/typesense/async_/api_call.py @@ -51,6 +51,7 @@ ServiceUnavailable, TypesenseClientError, ) +from typesense.http_backend import ASYNC_CLIENT_TYPES, AsyncClientType, backend_errors from typesense.node_manager import NodeManager from typesense.request_handler import RequestHandler @@ -116,38 +117,23 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): }, ) -_SERVER_ERRORS: typing.Final[ - typing.Tuple[ - typing.Type[httpx.TimeoutException], - typing.Type[httpx.ConnectError], - typing.Type[httpx.HTTPError], - typing.Type[httpx.RequestError], - typing.Type[HTTPStatus0Error], - typing.Type[ServerError], - typing.Type[ServiceUnavailable], - ] -] = ( - httpx.TimeoutException, - httpx.ConnectError, - httpx.HTTPError, - httpx.RequestError, +_SERVER_ERRORS: typing.Final[typing.Tuple[typing.Type[Exception], ...]] = ( + *backend_errors("TimeoutException"), + *backend_errors("ConnectError"), + *backend_errors("HTTPError"), + *backend_errors("RequestError"), HTTPStatus0Error, ServerError, 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, +# Raised by httpx inside the client, so they say nothing about the node's +# health. They subclass entries of _SERVER_ERRORS and must be caught first. +_CLIENT_ERRORS: typing.Final[typing.Tuple[typing.Type[Exception], ...]] = ( + *backend_errors("PoolTimeout"), + *backend_errors("LocalProtocolError"), + *backend_errors("DecodingError"), + *backend_errors("TooManyRedirects"), ) @@ -161,19 +147,42 @@ class AsyncApiCall: Attributes: config (Configuration): The configuration object for the Typesense client. node_manager (NodeManager): Manages the nodes in the Typesense cluster. - _client (httpx.AsyncClient): The httpx async client for making requests. + _client (httpx.AsyncClient | httpx2.AsyncClient): The async client for + making requests. """ - def __init__(self, config: Configuration): + def __init__( + self, + config: Configuration, + http_client: typing.Optional[AsyncClientType] = None, + ): """ Initialize the AsyncApiCall instance. Args: config (Configuration): The configuration object for the Typesense client. + http_client (httpx.AsyncClient | httpx2.AsyncClient, optional): A client + to send requests with instead of the default httpx client. The + connection pool settings in ``config`` are not applied to it, and it + is not closed by ``aclose``. + + Raises: + TypeError: If ``http_client`` is not an httpx or httpx2 async client. """ self.config = config self.node_manager = NodeManager(config) self.request_handler = RequestHandler(config) + self._concurrency_limit = AsyncConcurrencyLimit( + config.max_concurrent_requests, + ) + self._owns_client = http_client is None + if http_client is not None: + if not isinstance(http_client, ASYNC_CLIENT_TYPES): + raise TypeError( + "`http_client` must be an httpx.AsyncClient or httpx2.AsyncClient.", + ) + self._client: AsyncClientType = http_client + return self._client = httpx.AsyncClient( timeout=httpx.Timeout( config.connection_timeout_seconds, @@ -184,9 +193,6 @@ def __init__(self, config: Configuration): max_keepalive_connections=config.max_keepalive_connections, ), ) - self._concurrency_limit = AsyncConcurrencyLimit( - config.max_concurrent_requests, - ) async def __aenter__(self) -> "AsyncApiCall": """Async context manager entry.""" @@ -199,11 +205,12 @@ async def __aexit__( exc_tb: typing.Optional[TracebackType], ) -> None: """Async context manager exit.""" - await self._client.aclose() + await self.aclose() async def aclose(self) -> None: - """Close the httpx client.""" - await self._client.aclose() + """Close the httpx client, unless it was passed in by the caller.""" + if self._owns_client: + await self._client.aclose() @typing.overload async def get( diff --git a/src/typesense/async_/client.py b/src/typesense/async_/client.py index 8175bcb..a156040 100644 --- a/src/typesense/async_/client.py +++ b/src/typesense/async_/client.py @@ -55,6 +55,7 @@ from .stopwords import AsyncStopwords from .synonym_sets import AsyncSynonymSets from typesense.configuration import ConfigDict, Configuration +from typesense.http_backend import AsyncClientType TDoc = typing.TypeVar("TDoc", bound=DocumentSchema) @@ -86,7 +87,11 @@ class AsyncClient: conversations_models (ConversationsModels): Instance for managing conversation models. """ - def __init__(self, config_dict: ConfigDict) -> None: + def __init__( + self, + config_dict: ConfigDict, + http_client: typing.Optional[AsyncClientType] = None, + ) -> None: """ Initialize the Client instance. @@ -94,6 +99,12 @@ def __init__(self, config_dict: ConfigDict) -> None: config_dict (ConfigDict): A dictionary containing the configuration for the Typesense client. + http_client (httpx.AsyncClient | httpx2.AsyncClient, optional): + A client to send requests with instead of the default httpx client, + e.g. an ``httpx2.AsyncClient`` (``pip install typesense[httpx2]``). + The connection pool settings in ``config_dict`` are not applied to + it, and the Typesense client does not close it. + Example: >>> config = { ... "api_key": "your_api_key", @@ -105,7 +116,7 @@ def __init__(self, config_dict: ConfigDict) -> None: >>> client = Client(config) """ self.config = Configuration(config_dict) - self.api_call = AsyncApiCall(self.config) + self.api_call = AsyncApiCall(self.config, http_client) self.collections: AsyncCollections[DocumentSchema] = AsyncCollections( self.api_call ) diff --git a/src/typesense/http_backend.py b/src/typesense/http_backend.py new file mode 100644 index 0000000..69e5a40 --- /dev/null +++ b/src/typesense/http_backend.py @@ -0,0 +1,77 @@ +""" +HTTP backends the Typesense client can send requests with. + +The client builds an ``httpx`` client by default. Users can pass their own client +instead, either from ``httpx`` or from ``httpx2`` (install ``typesense[httpx2]``, +Python 3.10+). ``httpx2`` is Pydantic's maintained continuation of ``httpx``. It +fixes a connection pool leak in ``httpcore`` (encode/httpcore#1093) that can leave +an async client failing every request with ``PoolTimeout``. + +``httpx2`` has the same API as ``httpx`` but its own classes, so every +``isinstance`` check and ``except`` clause has to cover both packages. This module +builds those type tuples once, including ``httpx2`` only when it is installed. +""" + +import importlib +import sys +from types import ModuleType + +import httpx + +if sys.version_info >= (3, 11): + import typing +else: + import typing_extensions as typing + +if typing.TYPE_CHECKING: + import httpx2 # noqa: F401 (used in the string annotations in this module) + + +def _import_httpx2() -> typing.Optional[ModuleType]: + """Return the ``httpx2`` module, or ``None`` when it is not installed.""" + # ``import_module`` is typed as returning a module whether or not httpx2 is + # installed; a plain ``import`` is ``Any`` to mypy when it is missing (Python 3.9). + try: + return importlib.import_module("httpx2") + except ImportError: + return None + + +_httpx2 = _import_httpx2() + +_BACKENDS: typing.Final[typing.Tuple[ModuleType, ...]] = tuple( + backend for backend in (httpx, _httpx2) if backend is not None +) + +SyncClientType = typing.Union[httpx.Client, "httpx2.Client"] +AsyncClientType = typing.Union[httpx.AsyncClient, "httpx2.AsyncClient"] +ResponseType = typing.Union[httpx.Response, "httpx2.Response"] + + +def backend_errors(name: str) -> typing.Tuple[typing.Type[Exception], ...]: + """ + Return the exception class called ``name`` from every installed backend. + + Args: + name (str): The exception name shared by ``httpx`` and ``httpx2``. + + Returns: + Tuple[Type[Exception], ...]: The matching classes, for ``except`` clauses. + """ + return tuple(getattr(backend, name) for backend in _BACKENDS) + + +# Declared precisely for type checkers so ``isinstance`` narrows to the client +# unions above; at runtime they only contain the backends that are installed. +if typing.TYPE_CHECKING: + CLIENT_TYPES: typing.Tuple[ + typing.Type[httpx.Client], + typing.Type["httpx2.Client"], + ] + ASYNC_CLIENT_TYPES: typing.Tuple[ + typing.Type[httpx.AsyncClient], + typing.Type["httpx2.AsyncClient"], + ] +else: + CLIENT_TYPES = tuple(backend.Client for backend in _BACKENDS) + ASYNC_CLIENT_TYPES = tuple(backend.AsyncClient for backend in _BACKENDS) diff --git a/src/typesense/request_handler.py b/src/typesense/request_handler.py index 1e8d81e..22457f6 100644 --- a/src/typesense/request_handler.py +++ b/src/typesense/request_handler.py @@ -19,14 +19,13 @@ - Normalizes boolean parameters for API requests - Supports both sync (httpx.Client) and async (httpx.AsyncClient) HTTP clients -Note: This module relies on the 'httpx' library for both sync and async operations. +Note: This module sends requests with an httpx or httpx2 client, sync or async. """ import json import sys from types import MappingProxyType -import httpx if sys.version_info >= (3, 11): import typing @@ -46,6 +45,19 @@ ServiceUnavailable, TypesenseClientError, ) +from typesense.http_backend import ( + ASYNC_CLIENT_TYPES, + CLIENT_TYPES, + AsyncClientType, + SyncClientType, + ResponseType, + backend_errors, +) + +_DECODING_ERRORS: typing.Final[typing.Tuple[typing.Type[Exception], ...]] = ( + json.JSONDecodeError, + *backend_errors("DecodingError"), +) TEntityDict = typing.TypeVar("TEntityDict") TParams = typing.TypeVar("TParams", bound=typing.Mapping[str, object]) @@ -151,7 +163,7 @@ def make_request( url: str, entity_type: typing.Type[TEntityDict], as_json: typing.Union[typing.Literal[True], typing.Literal[False]] = True, - client: httpx.AsyncClient, + client: AsyncClientType, **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Coroutine[typing.Any, typing.Any, typing.Union[TEntityDict, str]]: ... @@ -163,7 +175,7 @@ def make_request( url: str, entity_type: typing.Type[TEntityDict], as_json: typing.Union[typing.Literal[True], typing.Literal[False]] = True, - client: httpx.Client, + client: SyncClientType, **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: ... @@ -174,7 +186,7 @@ def make_request( url: str, entity_type: typing.Type[TEntityDict], as_json: typing.Union[typing.Literal[True], typing.Literal[False]] = True, - client: typing.Union[httpx.Client, httpx.AsyncClient], + client: typing.Union[SyncClientType, AsyncClientType], **kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[ TEntityDict, @@ -226,14 +238,15 @@ def make_request( body if isinstance(body, (str, bytes)) else json.dumps(body) ) - if isinstance(client, httpx.AsyncClient): + if isinstance(client, ASYNC_CLIENT_TYPES): return self._make_async_request( method, url, entity_type, as_json, client, **request_kwargs ) - else: + if isinstance(client, CLIENT_TYPES): return self._make_sync_request( method, url, entity_type, as_json, client, **request_kwargs ) + raise TypeError("`client` must be an httpx or httpx2 client.") def _make_sync_request( self, @@ -241,7 +254,7 @@ def _make_sync_request( url: str, entity_type: typing.Type[TEntityDict], as_json: bool, - client: httpx.Client, + client: SyncClientType, **request_kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: """Make a synchronous HTTP request using httpx.Client.""" @@ -249,7 +262,7 @@ def _make_sync_request( content: typing.Union[str, bytes, None] = request_kwargs.get("content") headers: typing.Dict[str, str] = request_kwargs.get("headers", {}) - response = client.request( + response: ResponseType = client.request( method, url, params=typing.cast(typing.Optional[_QueryParams], params), @@ -276,7 +289,7 @@ async def _make_async_request( url: str, entity_type: typing.Type[TEntityDict], as_json: bool, - client: httpx.AsyncClient, + client: AsyncClientType, **request_kwargs: typing.Unpack[SessionFunctionKwargs[TParams, TBody]], ) -> typing.Union[TEntityDict, str]: """Make an asynchronous HTTP request using httpx.AsyncClient.""" @@ -284,7 +297,7 @@ async def _make_async_request( content: typing.Union[str, bytes, None] = request_kwargs.get("content") headers: typing.Dict[str, str] = request_kwargs.get("headers", {}) - response = await client.request( + response: ResponseType = await client.request( method, url, params=typing.cast(typing.Optional[_QueryParams], params), @@ -325,12 +338,12 @@ def normalize_params(params: typing.Mapping[str, object]) -> None: params[key] = str(parameter_value).lower() @staticmethod - def _get_error_message(response: httpx.Response) -> str: + def _get_error_message(response: ResponseType) -> str: """ Extract the error message from an API response. Args: - response (httpx.Response): The API response. + response (httpx.Response | httpx2.Response): The API response. Returns: str: The extracted error message or a default message. @@ -339,7 +352,7 @@ def _get_error_message(response: httpx.Response) -> str: if content_type.startswith("application/json"): try: return typing.cast(str, response.json().get("message", "API error.")) - except (json.JSONDecodeError, httpx.DecodingError): + except _DECODING_ERRORS: return f"API error: Invalid JSON response: {response.text}" if response.text: return f"API error. {response.text}" diff --git a/src/typesense/sync/api_call.py b/src/typesense/sync/api_call.py index 09c6107..42c58a0 100644 --- a/src/typesense/sync/api_call.py +++ b/src/typesense/sync/api_call.py @@ -51,6 +51,7 @@ ServiceUnavailable, TypesenseClientError, ) +from typesense.http_backend import CLIENT_TYPES, SyncClientType, backend_errors from typesense.node_manager import NodeManager from typesense.request_handler import RequestHandler @@ -116,38 +117,23 @@ class SessionFunctionKwargs(typing.Generic[TParams, TBody], typing.TypedDict): }, ) -_SERVER_ERRORS: typing.Final[ - typing.Tuple[ - typing.Type[httpx.TimeoutException], - typing.Type[httpx.ConnectError], - typing.Type[httpx.HTTPError], - typing.Type[httpx.RequestError], - typing.Type[HTTPStatus0Error], - typing.Type[ServerError], - typing.Type[ServiceUnavailable], - ] -] = ( - httpx.TimeoutException, - httpx.ConnectError, - httpx.HTTPError, - httpx.RequestError, +_SERVER_ERRORS: typing.Final[typing.Tuple[typing.Type[Exception], ...]] = ( + *backend_errors("TimeoutException"), + *backend_errors("ConnectError"), + *backend_errors("HTTPError"), + *backend_errors("RequestError"), HTTPStatus0Error, ServerError, 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, +# Raised by httpx inside the client, so they say nothing about the node's +# health. They subclass entries of _SERVER_ERRORS and must be caught first. +_CLIENT_ERRORS: typing.Final[typing.Tuple[typing.Type[Exception], ...]] = ( + *backend_errors("PoolTimeout"), + *backend_errors("LocalProtocolError"), + *backend_errors("DecodingError"), + *backend_errors("TooManyRedirects"), ) @@ -161,19 +147,42 @@ class ApiCall: Attributes: config (Configuration): The configuration object for the Typesense client. node_manager (NodeManager): Manages the nodes in the Typesense cluster. - _client (httpx.Client): The httpx async client for making requests. + _client (httpx.Client | httpx2.Client): The async client for + making requests. """ - def __init__(self, config: Configuration): + def __init__( + self, + config: Configuration, + http_client: typing.Optional[SyncClientType] = None, + ): """ Initialize the ApiCall instance. Args: config (Configuration): The configuration object for the Typesense client. + http_client (httpx.Client | httpx2.Client, optional): A client + to send requests with instead of the default httpx client. The + connection pool settings in ``config`` are not applied to it, and it + is not closed by ``close``. + + Raises: + TypeError: If ``http_client`` is not an httpx or httpx2 async client. """ self.config = config self.node_manager = NodeManager(config) self.request_handler = RequestHandler(config) + self._concurrency_limit = ConcurrencyLimit( + config.max_concurrent_requests, + ) + self._owns_client = http_client is None + if http_client is not None: + if not isinstance(http_client, CLIENT_TYPES): + raise TypeError( + "`http_client` must be an httpx.Client or httpx2.Client.", + ) + self._client: SyncClientType = http_client + return self._client = httpx.Client( timeout=httpx.Timeout( config.connection_timeout_seconds, @@ -184,9 +193,6 @@ def __init__(self, config: Configuration): max_keepalive_connections=config.max_keepalive_connections, ), ) - self._concurrency_limit = ConcurrencyLimit( - config.max_concurrent_requests, - ) def __enter__(self) -> "ApiCall": """Async context manager entry.""" @@ -199,11 +205,12 @@ def __exit__( exc_tb: typing.Optional[TracebackType], ) -> None: """Async context manager exit.""" - self._client.close() + self.close() def close(self) -> None: - """Close the httpx client.""" - self._client.close() + """Close the httpx client, unless it was passed in by the caller.""" + if self._owns_client: + self._client.close() @typing.overload def get( diff --git a/src/typesense/sync/client.py b/src/typesense/sync/client.py index b11e542..3e01d8a 100644 --- a/src/typesense/sync/client.py +++ b/src/typesense/sync/client.py @@ -55,6 +55,7 @@ from .stopwords import Stopwords from .synonym_sets import SynonymSets from typesense.configuration import ConfigDict, Configuration +from typesense.http_backend import SyncClientType TDoc = typing.TypeVar("TDoc", bound=DocumentSchema) @@ -86,7 +87,11 @@ class Client: conversations_models (ConversationsModels): Instance for managing conversation models. """ - def __init__(self, config_dict: ConfigDict) -> None: + def __init__( + self, + config_dict: ConfigDict, + http_client: typing.Optional[SyncClientType] = None, + ) -> None: """ Initialize the Client instance. @@ -94,6 +99,12 @@ def __init__(self, config_dict: ConfigDict) -> None: config_dict (ConfigDict): A dictionary containing the configuration for the Typesense client. + http_client (httpx.Client | httpx2.Client, optional): + A client to send requests with instead of the default httpx client, + e.g. an ``httpx2.Client`` (``pip install typesense[httpx2]``). + The connection pool settings in ``config_dict`` are not applied to + it, and the Typesense client does not close it. + Example: >>> config = { ... "api_key": "your_api_key", @@ -105,7 +116,7 @@ def __init__(self, config_dict: ConfigDict) -> None: >>> client = Client(config) """ self.config = Configuration(config_dict) - self.api_call = ApiCall(self.config) + self.api_call = ApiCall(self.config, http_client) self.collections: Collections[DocumentSchema] = Collections( self.api_call ) diff --git a/tests/http_backend_test.py b/tests/http_backend_test.py new file mode 100644 index 0000000..690ffc3 --- /dev/null +++ b/tests/http_backend_test.py @@ -0,0 +1,151 @@ +"""Tests for sending requests with a user-supplied httpx or httpx2 client.""" + +import sys + +if sys.version_info >= (3, 11): + import typing +else: + import typing_extensions as typing + +import httpx +import pytest + +from typesense.async_.api_call import AsyncApiCall +from typesense.configuration import Configuration +from typesense.http_backend import ASYNC_CLIENT_TYPES, CLIENT_TYPES +from typesense.sync.api_call import ApiCall + +httpx2 = pytest.importorskip("httpx2") + + +def _ok_response(request: typing.Any) -> typing.Any: + """Return a successful JSON response.""" + return httpx2.Response(200, json={"key": "value"}) + + +def test_backend_errors_include_httpx2() -> None: + """Test that httpx2 clients are recognised when httpx2 is installed.""" + assert httpx2.Client in CLIENT_TYPES + assert httpx2.AsyncClient in ASYNC_CLIENT_TYPES + assert httpx.Client in CLIENT_TYPES + assert httpx.AsyncClient in ASYNC_CLIENT_TYPES + + +def test_sends_requests_with_httpx2_client(fake_config: Configuration) -> None: + """Test that requests go through a user-supplied httpx2 client.""" + http_client = httpx2.Client(transport=httpx2.MockTransport(_ok_response)) + api_call = ApiCall(fake_config, http_client) + + response = api_call.get("/", entity_type=typing.Dict[str, str]) + + assert response == {"key": "value"} + + +async def test_async_sends_requests_with_httpx2_client( + fake_config: Configuration, +) -> None: + """Test that requests go through a user-supplied httpx2 async client.""" + http_client = httpx2.AsyncClient(transport=httpx2.MockTransport(_ok_response)) + api_call = AsyncApiCall(fake_config, http_client) + + response = await api_call.get("/", entity_type=typing.Dict[str, str]) + + assert response == {"key": "value"} + + +def test_httpx2_connect_error_fails_over(fake_config: Configuration) -> None: + """Test that an httpx2 connection error marks the node unhealthy and fails over.""" + + def handler(request: typing.Any) -> typing.Any: + if request.url.host == "nearest": + raise httpx2.ConnectError("Connection refused", request=request) + return httpx2.Response(200, json={"key": "value"}) + + http_client = httpx2.Client(transport=httpx2.MockTransport(handler)) + api_call = ApiCall(fake_config, http_client) + + response = api_call.get("/", entity_type=typing.Dict[str, str]) + + assert response == {"key": "value"} + assert fake_config.nearest_node.healthy is False + + +def test_httpx2_client_side_error_does_not_fail_over( + fake_config: Configuration, +) -> None: + """Test that an httpx2 PoolTimeout propagates without marking the node unhealthy.""" + requested_hosts: typing.List[str] = [] + + def handler(request: typing.Any) -> typing.Any: + requested_hosts.append(request.url.host) + raise httpx2.PoolTimeout("Pool timeout", request=request) + + http_client = httpx2.Client(transport=httpx2.MockTransport(handler)) + api_call = ApiCall(fake_config, http_client) + + with pytest.raises(httpx2.PoolTimeout): + api_call.get("/", entity_type=typing.Dict[str, str]) + + assert requested_hosts == ["nearest"] + assert fake_config.nearest_node.healthy is True + + +def test_httpx2_server_error_response_fails_over(fake_config: Configuration) -> None: + """Test that a 503 through an httpx2 client fails over to the next node.""" + + def handler(request: typing.Any) -> typing.Any: + if request.url.host == "nearest": + return httpx2.Response(503, json={"message": "unavailable"}) + return httpx2.Response(200, json={"key": "value"}) + + http_client = httpx2.Client(transport=httpx2.MockTransport(handler)) + api_call = ApiCall(fake_config, http_client) + + response = api_call.get("/", entity_type=typing.Dict[str, str]) + + assert response == {"key": "value"} + assert fake_config.nearest_node.healthy is False + + +def test_rejects_async_client_for_sync_api_call(fake_config: Configuration) -> None: + """Test that the sync client refuses an async http client.""" + with pytest.raises(TypeError, match="`http_client` must be"): + ApiCall(fake_config, typing.cast(typing.Any, httpx2.AsyncClient())) + + +def test_rejects_sync_client_for_async_api_call(fake_config: Configuration) -> None: + """Test that the async client refuses a sync http client.""" + with pytest.raises(TypeError, match="`http_client` must be"): + AsyncApiCall(fake_config, typing.cast(typing.Any, httpx2.Client())) + + +def test_close_leaves_user_supplied_client_open(fake_config: Configuration) -> None: + """Test that closing the api call does not close a client the caller owns.""" + http_client = httpx2.Client(transport=httpx2.MockTransport(_ok_response)) + api_call = ApiCall(fake_config, http_client) + + api_call.close() + + assert http_client.is_closed is False + + +def test_close_closes_default_client(fake_config: Configuration) -> None: + """Test that closing the api call closes the client it created.""" + api_call = ApiCall(fake_config) + + api_call.close() + + assert api_call._client.is_closed is True + + +async def test_async_close_leaves_user_supplied_client_open( + fake_config: Configuration, +) -> None: + """Test that closing the async api call does not close a caller-owned client.""" + http_client = httpx2.AsyncClient(transport=httpx2.MockTransport(_ok_response)) + api_call = AsyncApiCall(fake_config, http_client) + + await api_call.aclose() + + assert http_client.is_closed is False + await http_client.aclose() diff --git a/utils/run-unasync.py b/utils/run-unasync.py index 7d836d4..5dd8816 100644 --- a/utils/run-unasync.py +++ b/utils/run-unasync.py @@ -29,6 +29,8 @@ def collect_class_replacements(source_dir: Path) -> dict[str, str]: replacements["asyncio"] = "time" # Defined in the shared ``typesense.concurrency_limit`` module, outside async_. replacements["AsyncConcurrencyLimit"] = "ConcurrencyLimit" + # Defined in the shared ``typesense.http_backend`` module, outside async_. + replacements["ASYNC_CLIENT_TYPES"] = "CLIENT_TYPES" return replacements diff --git a/uv.lock b/uv.lock index a1ae574..f017369 100644 --- a/uv.lock +++ b/uv.lock @@ -1,8 +1,9 @@ version = 1 -revision = 3 +revision = 5 requires-python = ">=3.9" resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", "python_full_version < '3.10'", ] @@ -12,7 +13,8 @@ version = "4.12.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "exceptiongroup", marker = "python_full_version < '3.11'" }, - { name = "idna" }, + { name = "idna", version = "3.11", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, + { name = "idna", version = "3.20", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.10'" }, { name = "typing-extensions", marker = "python_full_version < '3.13'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/16/ce/8a777047513153587e5434fd752e89334ac33e379aa3497db860eeb60377/anyio-4.12.0.tar.gz", hash = "sha256:73c693b567b0c55130c104d0b43a9baf3aa6a31fc6110116509f27bf75e21ec0", size = 228266, upload-time = "2025-11-28T23:37:38.911Z" } @@ -271,7 +273,8 @@ name = "coverage" version = "7.13.0" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] sdist = { url = "https://files.pythonhosted.org/packages/b6/45/2c665ca77ec32ad67e25c77daf1cee28ee4558f3bc571cdbaf88a00b9f23/coverage-7.13.0.tar.gz", hash = "sha256:a394aa27f2d7ff9bc04cf703817773a59ad6dfbd577032e690f961d2460ee936", size = 820905, upload-time = "2025-12-08T13:14:38.055Z" } wheels = [ @@ -373,7 +376,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions", marker = "python_full_version < '3.13'" }, + { name = "typing-extensions" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -388,7 +391,7 @@ resolution-markers = [ "python_full_version < '3.10'", ] dependencies = [ - { name = "tzdata", marker = "python_full_version < '3.10'" }, + { name = "tzdata" }, ] sdist = { url = "https://files.pythonhosted.org/packages/3d/84/e95acaa848b855e15c83331d0401ee5f84b2f60889255c2e055cb4fb6bdf/faker-37.12.0.tar.gz", hash = "sha256:7505e59a7e02fa9010f06c3e1e92f8250d4cfbb30632296140c2d6dbef09b0fa", size = 1935741, upload-time = "2025-10-24T15:19:58.764Z" } wheels = [ @@ -400,10 +403,11 @@ name = "faker" version = "38.2.0" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] dependencies = [ - { name = "tzdata", marker = "python_full_version >= '3.10'" }, + { name = "tzdata" }, ] sdist = { url = "https://files.pythonhosted.org/packages/64/27/022d4dbd4c20567b4c294f79a133cc2f05240ea61e0d515ead18c995c249/faker-38.2.0.tar.gz", hash = "sha256:20672803db9c7cb97f9b56c18c54b915b6f1d8991f63d1d673642dc43f5ce7ab", size = 1941469, upload-time = "2025-11-19T16:37:31.892Z" } wheels = [ @@ -432,6 +436,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/7e/f5/f66802a942d491edb555dd61e3a9961140fd64c90bce1eafd741609d334d/httpcore-1.0.9-py3-none-any.whl", hash = "sha256:2d400746a40668fc9dec9810239072b40b4484b640a8c38fd654a024c7a1bf55", size = 78784, upload-time = "2025-04-24T22:06:20.566Z" }, ] +[[package]] +name = "httpcore2" +version = "2.13.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "h11" }, + { name = "truststore" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/cb/f3/1db7aa2bc2524062192bb0e0323969492d1883152a232fe36eea65f4e35c/httpcore2-2.13.1.tar.gz", hash = "sha256:e0aa977abe17e69a3b820a24542a6fa88702676d83880b8d194dcd18408e5103", size = 68071, upload-time = "2026-09-23T07:47:22.372Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/09/ba/a4568248771ce81957bfb7cc600264a40fbcda092391ee1c415c50be4bea/httpcore2-2.13.1-py3-none-any.whl", hash = "sha256:e1e05d4f25f7d7d496bfb96748f6f4b67657b03da069b3a68c36069f3db73d0a", size = 83423, upload-time = "2026-09-23T07:47:19.365Z" }, +] + [[package]] name = "httpx" version = "0.28.1" @@ -440,28 +457,71 @@ dependencies = [ { name = "anyio" }, { name = "certifi" }, { name = "httpcore" }, - { name = "idna" }, + { name = "idna", version = "3.11", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, + { name = "idna", version = "3.20", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.10'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/b1/df/48c586a5fe32a0f01324ee087459e112ebb7224f646c0b5023f5e79e9956/httpx-0.28.1.tar.gz", hash = "sha256:75e98c5f16b0f35b567856f597f06ff2270a374470a5c2392242528e3e3e42fc", size = 141406, upload-time = "2024-12-06T15:37:23.222Z" } wheels = [ { url = "https://files.pythonhosted.org/packages/2a/39/e50c7c3a983047577ee07d2a9e53faf5a69493943ec3f6a384bdc792deb2/httpx-0.28.1-py3-none-any.whl", hash = "sha256:d909fcccc110f8c7faf814ca82a9a4d816bc5a6dbfea25d6591d6985b8ba59ad", size = 73517, upload-time = "2024-12-06T15:37:21.509Z" }, ] +[[package]] +name = "httpx2" +version = "2.13.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "anyio", marker = "sys_platform != 'emscripten'" }, + { name = "httpcore2", marker = "sys_platform != 'emscripten'" }, + { name = "httpx2-jsfetch", marker = "python_full_version >= '3.12' and sys_platform == 'emscripten'" }, + { name = "idna", version = "3.20", source = { registry = "https://pypi.org/simple" } }, + { name = "truststore", marker = "sys_platform != 'emscripten'" }, + { name = "typing-extensions", marker = "python_full_version < '3.13'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/d5/44/474bef2a0e9d90f1715d32cb98b0738695ca17ba324095fb2497ed7fbd59/httpx2-2.13.1.tar.gz", hash = "sha256:e48744a19e3af5ee48313d0ce5fe941d5422fae5705ea922a4aabf94d7800dfa", size = 100405, upload-time = "2026-09-23T07:47:23.052Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/d8/9c/6fe8931fd9f381042a9e4c7d5a7b4cbf7016b252bec0c99a49fce42c3326/httpx2-2.13.1-py3-none-any.whl", hash = "sha256:6dff50fabc270ee5fd25d845d0b078ed20564579744d6d962850975996d2f9a4", size = 95597, upload-time = "2026-09-23T07:47:20.995Z" }, +] + +[[package]] +name = "httpx2-jsfetch" +version = "1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/cd/c4/0e5636363151a2a1795e0a77617168b9ca438e1748ec05fc9b5687f93d64/httpx2_jsfetch-1.0.tar.gz", hash = "sha256:70a0e3eabfef7cce5ad9c629f7d01ca05e418f586646f4ddf14782e4c1454c60", size = 6872, upload-time = "2026-08-07T00:13:07.492Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/9b/43/832f631d32e4f1211caa2ba368317739fe71f0b8530e4c9d15dc454bac2a/httpx2_jsfetch-1.0-py3-none-any.whl", hash = "sha256:cb916b707601e69a07721aabc8f3f6659be3a6893bc1ff5c6f9e02241df2da32", size = 6382, upload-time = "2026-08-07T00:13:06.567Z" }, +] + [[package]] name = "idna" version = "3.11" source = { registry = "https://pypi.org/simple" } +resolution-markers = [ + "python_full_version < '3.10'", +] sdist = { url = "https://files.pythonhosted.org/packages/6f/6d/0703ccc57f3a7233505399edb88de3cbd678da106337b9fcde432b65ed60/idna-3.11.tar.gz", hash = "sha256:795dafcc9c04ed0c1fb032c2aa73654d8e8c5023a7df64a53f39190ada629902", size = 194582, upload-time = "2025-10-12T14:55:20.501Z" } wheels = [ { url = "https://files.pythonhosted.org/packages/0e/61/66938bbb5fc52dbdf84594873d5b51fb1f7c7794e9c0f5bd885f30bc507b/idna-3.11-py3-none-any.whl", hash = "sha256:771a87f49d9defaf64091e6e6fe9c18d4833f140bd19464795bc32d966ca37ea", size = 71008, upload-time = "2025-10-12T14:55:18.883Z" }, ] +[[package]] +name = "idna" +version = "3.20" +source = { registry = "https://pypi.org/simple" } +resolution-markers = [ + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", +] +sdist = { url = "https://files.pythonhosted.org/packages/f5/08/8eea9d4b8302028f3abb2c0813953f7aec26d33b7a8960ed760e65ff29fa/idna-3.20.tar.gz", hash = "sha256:a7db850025b95ded1eae8a46181a1a6c56c92c96f0e2b005d9ff8dc0210cab44", size = 216463, upload-time = "2026-09-17T14:11:04.752Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/58/a2/bb081bab032533a855d44de1d56f8e8426114ff1ba5d1f07a438a0a654f8/idna-3.20-py3-none-any.whl", hash = "sha256:ab7ae7122974553370f0bdb919e1a960b2cd1bc1ef0276416d896db81c14582c", size = 69583, upload-time = "2026-09-17T14:11:03.168Z" }, +] + [[package]] name = "importlib-metadata" version = "8.7.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "zipp", marker = "python_full_version < '3.10'" }, + { name = "zipp" }, ] sdist = { url = "https://files.pythonhosted.org/packages/76/66/650a33bd90f786193e4de4b3ad86ea60b53c89b669a5c7be931fac31cdb0/importlib_metadata-8.7.0.tar.gz", hash = "sha256:d13b81ad223b890aa16c5471f2ac3056cf76c5f10f82d6f9292f0b415f389000", size = 56641, upload-time = "2025-04-27T15:29:01.736Z" } wheels = [ @@ -485,7 +545,8 @@ name = "iniconfig" version = "2.3.0" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] sdist = { url = "https://files.pythonhosted.org/packages/72/34/14ca021ce8e5dfedc35312d08ba8bf51fdd999c576889fc2c24cb97f4f10/iniconfig-2.3.0.tar.gz", hash = "sha256:c76315c77db068650d49c5b56314774a7804df16fee4402c1f19d6d15d8c4730", size = 20503, upload-time = "2025-10-18T21:55:43.219Z" } wheels = [ @@ -500,7 +561,7 @@ resolution-markers = [ "python_full_version < '3.10'", ] dependencies = [ - { name = "importlib-metadata", marker = "python_full_version < '3.10'" }, + { name = "importlib-metadata" }, ] sdist = { url = "https://files.pythonhosted.org/packages/1e/82/fa43935523efdfcce6abbae9da7f372b627b27142c3419fcf13bf5b0c397/isort-6.1.0.tar.gz", hash = "sha256:9b8f96a14cfee0677e78e941ff62f03769a06d412aabb9e2a90487b3b7e8d481", size = 824325, upload-time = "2025-10-01T16:26:45.027Z" } wheels = [ @@ -512,7 +573,8 @@ name = "isort" version = "7.0.0" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] sdist = { url = "https://files.pythonhosted.org/packages/63/53/4f3c058e3bace40282876f9b553343376ee687f3c35a525dc79dbd450f88/isort-7.0.0.tar.gz", hash = "sha256:5513527951aadb3ac4292a41a16cbc50dd1642432f5e8c20057d414bdafb4187", size = 805049, upload-time = "2025-10-11T13:30:59.107Z" } wheels = [ @@ -707,13 +769,13 @@ resolution-markers = [ "python_full_version < '3.10'", ] dependencies = [ - { name = "colorama", marker = "python_full_version < '3.10' and sys_platform == 'win32'" }, - { name = "exceptiongroup", marker = "python_full_version < '3.10'" }, - { name = "iniconfig", version = "2.1.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, - { name = "packaging", marker = "python_full_version < '3.10'" }, - { name = "pluggy", marker = "python_full_version < '3.10'" }, - { name = "pygments", marker = "python_full_version < '3.10'" }, - { name = "tomli", marker = "python_full_version < '3.10'" }, + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "exceptiongroup" }, + { name = "iniconfig", version = "2.1.0", source = { registry = "https://pypi.org/simple" } }, + { name = "packaging" }, + { name = "pluggy" }, + { name = "pygments" }, + { name = "tomli" }, ] sdist = { url = "https://files.pythonhosted.org/packages/a3/5c/00a0e072241553e1a7496d638deababa67c5058571567b92a7eaa258397c/pytest-8.4.2.tar.gz", hash = "sha256:86c0d0b93306b961d58d62a4db4879f27fe25513d4b969df351abdddb3c30e01", size = 1519618, upload-time = "2025-09-04T14:34:22.711Z" } wheels = [ @@ -725,16 +787,17 @@ name = "pytest" version = "9.0.2" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] dependencies = [ - { name = "colorama", marker = "python_full_version >= '3.10' and sys_platform == 'win32'" }, - { name = "exceptiongroup", marker = "python_full_version == '3.10.*'" }, - { name = "iniconfig", version = "2.3.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.10'" }, - { name = "packaging", marker = "python_full_version >= '3.10'" }, - { name = "pluggy", marker = "python_full_version >= '3.10'" }, - { name = "pygments", marker = "python_full_version >= '3.10'" }, - { name = "tomli", marker = "python_full_version == '3.10.*'" }, + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "exceptiongroup", marker = "python_full_version < '3.11'" }, + { name = "iniconfig", version = "2.3.0", source = { registry = "https://pypi.org/simple" } }, + { name = "packaging" }, + { name = "pluggy" }, + { name = "pygments" }, + { name = "tomli", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/d1/db/7ef3487e0fb0049ddb5ce41d3a49c235bf9ad299b6a25d5780a89f19230f/pytest-9.0.2.tar.gz", hash = "sha256:75186651a92bd89611d1d9fc20f0b4345fd827c41ccd5c299a868a05d70edf11", size = 1568901, upload-time = "2025-12-06T21:30:51.014Z" } wheels = [ @@ -749,9 +812,9 @@ resolution-markers = [ "python_full_version < '3.10'", ] dependencies = [ - { name = "backports-asyncio-runner", marker = "python_full_version < '3.10'" }, - { name = "pytest", version = "8.4.2", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, - { name = "typing-extensions", marker = "python_full_version < '3.10'" }, + { name = "backports-asyncio-runner" }, + { name = "pytest", version = "8.4.2", source = { registry = "https://pypi.org/simple" } }, + { name = "typing-extensions" }, ] sdist = { url = "https://files.pythonhosted.org/packages/42/86/9e3c5f48f7b7b638b216e4b9e645f54d199d7abbbab7a64a13b4e12ba10f/pytest_asyncio-1.2.0.tar.gz", hash = "sha256:c609a64a2a8768462d0c99811ddb8bd2583c33fd33cf7f21af1c142e824ffb57", size = 50119, upload-time = "2025-09-12T07:33:53.816Z" } wheels = [ @@ -763,12 +826,13 @@ name = "pytest-asyncio" version = "1.3.0" source = { registry = "https://pypi.org/simple" } resolution-markers = [ - "python_full_version >= '3.10'", + "python_full_version >= '3.12' and sys_platform == 'emscripten'", + "(python_full_version >= '3.10' and python_full_version < '3.12') or (python_full_version >= '3.10' and sys_platform != 'emscripten')", ] dependencies = [ - { name = "backports-asyncio-runner", marker = "python_full_version == '3.10.*'" }, - { name = "pytest", version = "9.0.2", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.10'" }, - { name = "typing-extensions", marker = "python_full_version >= '3.10' and python_full_version < '3.13'" }, + { name = "backports-asyncio-runner", marker = "python_full_version < '3.11'" }, + { name = "pytest", version = "9.0.2", source = { registry = "https://pypi.org/simple" } }, + { name = "typing-extensions", marker = "python_full_version < '3.13'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/90/2c/8af215c0f776415f3590cac4f9086ccefd6fd463befeae41cd4d3f193e5a/pytest_asyncio-1.3.0.tar.gz", hash = "sha256:d7f52f36d231b80ee124cd216ffb19369aa168fc10095013c6b014a34d3ee9e5", size = 50087, upload-time = "2025-11-10T16:07:47.256Z" } wheels = [ @@ -804,7 +868,8 @@ source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "certifi" }, { name = "charset-normalizer" }, - { name = "idna" }, + { name = "idna", version = "3.11", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, + { name = "idna", version = "3.20", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.10'" }, { name = "urllib3" }, ] sdist = { url = "https://files.pythonhosted.org/packages/c9/74/b3ff8e6c8446842c3f5c837e9c3dfcfe2018ea6ecef224c710c85ef728f4/requests-2.32.5.tar.gz", hash = "sha256:dbba0bac56e100853db0ea71b82b4dfd5fe2bf6d3754a8893c3af500cec7d7cf", size = 134517, upload-time = "2025-08-18T20:46:02.573Z" } @@ -917,6 +982,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/77/b8/0135fadc89e73be292b473cb820b4f5a08197779206b33191e801feeae40/tomli-2.3.0-py3-none-any.whl", hash = "sha256:e95b1af3c5b07d9e643909b5abbec77cd9f1217e6d0bca72b0234736b9fb1f1b", size = 14408, upload-time = "2025-10-08T22:01:46.04Z" }, ] +[[package]] +name = "truststore" +version = "0.10.4" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/53/a3/1585216310e344e8102c22482f6060c7a6ea0322b63e026372e6dcefcfd6/truststore-0.10.4.tar.gz", hash = "sha256:9d91bd436463ad5e4ee4aba766628dd6cd7010cf3e2461756b3303710eebc301", size = 26169, upload-time = "2025-08-12T18:49:02.73Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/19/97/56608b2249fe206a67cd573bc93cd9896e1efb9e98bce9c163bcdc704b88/truststore-0.10.4-py3-none-any.whl", hash = "sha256:adaeaecf1cbb5f4de3b1959b42d41f6fab57b2b1666adb59e89cb0b53361d981", size = 18660, upload-time = "2025-08-12T18:49:01.46Z" }, +] + [[package]] name = "typesense" source = { virtual = "." } @@ -925,6 +999,11 @@ dependencies = [ { name = "typing-extensions" }, ] +[package.optional-dependencies] +httpx2 = [ + { name = "httpx2", marker = "python_full_version >= '3.10'" }, +] + [package.dev-dependencies] dev = [ { name = "coverage", version = "7.10.7", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" }, @@ -949,8 +1028,10 @@ dev = [ [package.metadata] requires-dist = [ { name = "httpx", specifier = ">=0.28.1" }, + { name = "httpx2", marker = "python_full_version >= '3.10' and extra == 'httpx2'", specifier = ">=2.6.0" }, { name = "typing-extensions" }, ] +provides-extras = ["httpx2"] [package.metadata.requires-dev] dev = [