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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
1 change: 1 addition & 0 deletions setup.cfg
Original file line number Diff line number Diff line change
Expand Up @@ -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
8 changes: 3 additions & 5 deletions src/typesense/async_/analytics_rule_v1.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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",
)
Expand Down
18 changes: 7 additions & 11 deletions src/typesense/async_/analytics_rules_v1.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
)
Expand All @@ -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",
)
Expand All @@ -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",
)
Expand Down
106 changes: 70 additions & 36 deletions src/typesense/async_/api_call.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@

import httpx

from typesense.concurrency_limit import AsyncConcurrencyLimit
from typesense.configuration import Configuration, Node
from typesense.exceptions import (
HTTPStatus0Error,
Expand All @@ -50,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

Expand All @@ -59,7 +61,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]]
)
Expand Down Expand Up @@ -94,7 +96,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]

Expand All @@ -115,26 +117,25 @@ 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,
)

# 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"),
)


class AsyncApiCall:
"""
Expand All @@ -146,21 +147,51 @@ 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=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,
),
)

async def __aenter__(self) -> "AsyncApiCall":
Expand All @@ -174,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(
Expand Down Expand Up @@ -473,11 +505,14 @@ async def _execute_request(
try:
return await self._make_request_and_process_response(
method,
node,
url,
entity_type,
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:
Expand All @@ -495,24 +530,23 @@ 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."""
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(
self.node_manager.get_node(),
is_healthy=True,
)
"""Make the async API request to `node` and process the response."""
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)
if as_json
Expand Down
19 changes: 15 additions & 4 deletions src/typesense/async_/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -86,14 +87,24 @@ 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.

Args:
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",
Expand All @@ -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
)
Expand Down Expand Up @@ -164,5 +175,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])
Loading
Loading