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
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,10 @@
> went on. It needs to be consolidated into what will eventually become the
> final changelog + documentation changes (if any).

* (Minor): The in-progress asynchronous client now exposes `feature_definitions()`. It returns the same dict as `UnleashClient`, keyed by feature name with each toggle's `type` and `project`. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and without any state it returns an empty dict. Every method on the client is now implemented. The client is still unexported. Nothing changes for code using `UnleashClient`.
* (Minor): The in-progress asynchronous client can now resolve variants. `get_variant()` returns the same variant dict as `UnleashClient` and emits the same impression events. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and a toggle the client does not know resolves to the disabled variant. When an initialized client is asked for a toggle it does not know, it logs at `verbose_log_level` that the client does not know the toggle. `feature_definitions()` still raises `NotImplementedError`. The client is still unexported. Nothing changes for code using `UnleashClient`.
* (Minor): The in-progress asynchronous client can now evaluate feature toggles. `is_enabled()` resolves a toggle against the feature state the client holds, with the same results, `fallback_function` handling and impression events as `UnleashClient`. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and a toggle the client does not know resolves to the fallback's answer, or to false without one. `get_variant()` and `feature_definitions()` still raise `NotImplementedError`. The client is still unexported. Nothing changes for code using `UnleashClient`.
* (Minor): The in-progress asynchronous client can now be initialized and shut down. `await initialize_client()` registers with the server, loads the cached feature state, polls every `refresh_interval` seconds and reports metrics every `metrics_interval` seconds, all on the running event loop. It returns without waiting for the server: the first fetch runs one refresh interval later, and until then the client holds the cached state. Like `UnleashClient`, it clears the cached ETag first, so the first fetch is unconditional. Polling is the only mode: `fetch_toggles` is accepted for parity but has no effect, `experimental_mode` is ignored, and there is no bootstrapping. `await destroy()` stops polling, sends the remaining metrics, closes the connection pool and is safe to call more than once. A `destroy()` that runs while registration is still in flight leaves nothing polling. The client also works as an async context manager (`async with AsyncUnleashClient(...) as client:`). `is_enabled()`, `get_variant()` and `feature_definitions()` still raise `NotImplementedError`. The client is still unexported. Nothing changes for code using `UnleashClient`.

* (Minor): The in-progress asynchronous client gains an internal `_AsyncPollingConnector`, in the private `UnleashClient.connectors._async_connector` module. The module is not part of the public API and may change or disappear without notice. It fetches feature state over the asynchronous transport on its own scheduler. Starting it loads the cached state and returns without waiting for the server; the first fetch runs one refresh interval later. With a non-empty cache, READY is therefore emitted when the connector starts rather than after the first fetch. The asynchronous client does not use it yet. Nothing changes for code using `UnleashClient`.
* (Minor): New `UnleashClient.errors` module with a hierarchy of SDK errors. `UnleashClientError` is the base of every error in the hierarchy, so catching it catches all of them without catching `Exception`. Each module gets one error grouping that module's errors, such as `InstanceRegistryError`, and specific errors derive from those, such as `MultipleInstancesNotAllowedError`. A client rejected under `InstanceAllowType.BLOCK` now raises `MultipleInstancesNotAllowedError` instead of a plain `Exception`, with the same message. It still derives from `Exception`, so existing `except Exception` handlers keep catching it. `TransportError` groups the errors raised while talking to the Unleash server. Other errors raised by the SDK are not part of the hierarchy yet.
Expand Down
227 changes: 199 additions & 28 deletions UnleashClient/clients/async_unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,9 @@
``pip install UnleashClient[async]``.
"""

import asyncio
import warnings
from datetime import datetime, timezone
from typing import Callable, Optional

from yggdrasil_engine.engine import UnleashEngine
Expand All @@ -16,24 +19,46 @@
from UnleashClient._headers import _HeaderFactory
from UnleashClient._instance_registry import _get_instance_registry
from UnleashClient._metrics import _AsyncMetricsReporter
from UnleashClient._payloads import _build_register_payload
from UnleashClient.cache import BaseCache, FileCache
from UnleashClient.clients.unleash_client import _RunState
from UnleashClient.config import ExperimentalMode, UnleashConfig
from UnleashClient.constants import REQUEST_RETRIES, REQUEST_TIMEOUT
from UnleashClient.connectors._async_connector import _AsyncPollingConnector
from UnleashClient.constants import (
ETAG,
METRIC_LAST_SENT_TIME,
REQUEST_RETRIES,
REQUEST_TIMEOUT,
)
from UnleashClient.events import BaseEvent
from UnleashClient.impact_metrics import ImpactMetrics
from UnleashClient.utils import InstanceAllowType

_NOT_IMPLEMENTED = (
"AsyncUnleashClient is a work in progress and does not do anything yet. "
"Use UnleashClient."
)
from UnleashClient.utils import LOGGER, InstanceAllowType


class AsyncUnleashClient:
"""
An asyncio-native client for the Unleash feature toggle system.

Not implemented yet: every method raises :class:`NotImplementedError`.
The client keeps feature state fresh by polling the Unleash server on the
event loop it was initialized on, and reports metrics on the same loop.
Streaming, offline mode and bootstrapping are not supported.

Example::

async with AsyncUnleashClient(
url="https://unleash.example.com/api",
app_name="my-app",
custom_headers={"Authorization": "<API token>"},
) as client:
if client.is_enabled("new-checkout", {"userId": "42"}):
...

variant = client.get_variant("checkout-button", {"userId": "42"})

definitions = client.feature_definitions()

client.impact_metrics.define_counter("purchases", "Number of purchases")
client.impact_metrics.increment_counter("purchases")
"""

def __init__( # noqa: PLR0913, PLR0917
Expand Down Expand Up @@ -115,14 +140,22 @@ def __init__( # noqa: PLR0913, PLR0917
events=self._event_dispatcher,
)
self._transport: _AsyncTransport = _AsyncTransport(self._config, self._headers)
self._scheduler: _AsyncScheduler = _AsyncScheduler()
self._metrics_scheduler: _AsyncScheduler = _AsyncScheduler()
self._metrics: _AsyncMetricsReporter = _AsyncMetricsReporter(
config=self._config,
transport=self._transport,
scheduler=self._scheduler,
scheduler=self._metrics_scheduler,
engine=self._engine,
impact_metrics=self.impact_metrics,
)
self._connector: Optional[_AsyncPollingConnector] = None
self._run_state: _RunState = _RunState.UNINITIALIZED
self._starting: bool = False
self._closed: bool = False

@property
def is_initialized(self) -> bool:
return self._run_state == _RunState.INITIALIZED

def is_enabled(
self,
Expand All @@ -144,7 +177,11 @@ def is_enabled(
:param fallback_function: Allows users to provide a custom function to set default value.
:return: Feature flag result
"""
raise NotImplementedError(_NOT_IMPLEMENTED)
return self._evaluator.is_enabled(
feature_name=feature_name,
context=context,
fallback_function=fallback_function,
)

def get_variant(self, feature_name: str, context: Optional[dict] = None) -> dict:
"""
Expand All @@ -158,33 +195,167 @@ def get_variant(self, feature_name: str, context: Optional[dict] = None) -> dict
:param context: Dictionary with context (e.g. IPs, email) for feature toggle.
:return: Variant and feature flag status.
"""
raise NotImplementedError(_NOT_IMPLEMENTED)
result = self._evaluator.get_variant(feature_name=feature_name, context=context)

if not result.is_found and self.is_initialized:
LOGGER.log(
self._config.verbose_log_level,
"Attempted to get feature flag/variation %s, but the client does not know it.",
feature_name,
)

return result.variant

def feature_definitions(self) -> dict:
"""
Returns a dict containing all feature definitions known to the SDK at the time of calling.
Normally this would be a pared down version of the response from the Unleash API but this
may also be a result from bootstrapping or loading from backup.
Returns a dict containing all feature definitions known to the client at
the time of calling, keyed by feature name. This is a pared down version
of the response from the Unleash server, or of the cached state before
the client has fetched from the server.

Notes:

Example response:
* It is a plain method, not a coroutine, and does not wait for the
server. Before the client has any feature state, it returns an empty
dict.

{
"feature1": {
"project": "default",
"type": "release",
Example response::

{
"feature1": {
"project": "default",
"type": "release",
}
}
}

:return: Feature definitions keyed by feature name.
"""
return self._evaluator.feature_definitions()

async def initialize_client(self, fetch_toggles: bool = True) -> None:
"""
Initializes the client and starts communication with the Unleash server.

This kicks off:

* Client registration
* Loading the cached feature state
* Feature polling, every ``refresh_interval`` seconds
* Metrics reporting, every ``metrics_interval`` seconds

Returns without waiting for the server's feature state. The first fetch
runs one ``refresh_interval`` after this returns, and until then the
client holds the cached state.

Calling it again, or after :meth:`destroy`, warns and does nothing.

This is done automatically when the client is used as an async context
manager:

.. code-block:: python

async with AsyncUnleashClient(
url="https://foo.bar",
app_name="myClient1",
instance_id="myinstanceid",
) as client:
pass

:param fetch_toggles: Accepted for parity with :class:`UnleashClient`. It
has no effect: the client always polls.
:raises aiohttp.InvalidURL: If registration is enabled and the URL is invalid.
:raises ValueError: If a custom strategy is invalid.
"""
raise NotImplementedError(_NOT_IMPLEMENTED)
if self._closed or self._starting or self._run_state > _RunState.UNINITIALIZED:
warnings.warn(
"Attempted to initialize an Unleash Client instance that has already been initialized."
)
return

self._starting = True
try:
self._cache.mset(
{METRIC_LAST_SENT_TIME: datetime.now(timezone.utc), ETAG: ""}
)

if self._config.custom_strategies:
self._engine.register_custom_strategies(self._config.custom_strategies)

if not self._config.disable_registration:
await self._transport.register(
_build_register_payload(
self._config, self._config.custom_strategies
)
)

if self._closed:
return

self._connector = _AsyncPollingConnector(
store=self._store,
transport=self._transport,
refresh_interval=self._config.refresh_interval,
refresh_jitter=self._config.refresh_jitter,
)
await self._connector.start()

if not self._config.disable_metrics:
self._metrics.start()
self._metrics_scheduler.start()

async def initialize_client(self) -> None:
raise NotImplementedError(_NOT_IMPLEMENTED)
self._run_state = _RunState.INITIALIZED
except Exception as excep:
LOGGER.warning(
"Exception during AsyncUnleashClient initialization: %s", excep
)
raise
finally:
self._starting = False

async def destroy(self) -> None:
raise NotImplementedError(_NOT_IMPLEMENTED)
"""
Gracefully shuts down the client: stops polling, sends the metrics
collected since the last send and closes the connection to the server.

For cache teardown:

* Default disk-backed FileCache instances are preserved on disk.
* Custom non-FileCache implementations will have ``destroy()`` called.

Calling it more than once does nothing.
"""
if self._closed:
return
self._closed = True
self._run_state = _RunState.SHUTDOWN

if self._connector is not None:
await self._connector.stop()

await self._metrics.stop()

try:
await self._metrics_scheduler.shutdown()
except Exception as exc:
LOGGER.warning("Exception during scheduler teardown: %s", exc)

await self._transport.aclose()

if not isinstance(self._cache, FileCache):
try:
self._cache.destroy()
except Exception as exc:
LOGGER.warning("Exception during cache teardown: %s", exc)

if self._event_dispatcher is not None:
await asyncio.get_running_loop().run_in_executor(
None, self._event_dispatcher.close
)

async def __aenter__(self) -> "AsyncUnleashClient":
raise NotImplementedError(_NOT_IMPLEMENTED)
await self.initialize_client()
return self

async def __aexit__(self, exc_type, exc_val, exc_tb) -> None:
raise NotImplementedError(_NOT_IMPLEMENTED)
async def __aexit__(self, exc_type, exc_val, exc_tb) -> bool:
await self.destroy()
return False
59 changes: 59 additions & 0 deletions tests/specification_tests/test_async_client_specs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
import platform
import sys
import uuid

import pytest

from tests.specification_tests.test_client_specs import TEST_DATA, TEST_NAMES
from tests.utilities.testing_constants import APP_NAME, URL
from UnleashClient.cache import FileCache
from UnleashClient.clients.async_unleash_client import AsyncUnleashClient


async def get_async_client(state, test_context=None, cache_directory=None):
cache_kwargs = {}
if cache_directory is not None:
cache_kwargs["directory"] = str(cache_directory)

cache = FileCache("MOCK_CACHE", **cache_kwargs)
cache.bootstrap_from_dict(state)
env = "default"
if test_context is not None and "environment" in test_context:
env = test_context["environment"]

unleash_client = AsyncUnleashClient(
url=URL,
app_name=APP_NAME,
instance_id="pytest_%s" % uuid.uuid4(),
disable_metrics=True,
disable_registration=True,
cache=cache,
environment=env,
)

await unleash_client.initialize_client(fetch_toggles=False)
return unleash_client


@pytest.mark.skipif(
sys.version_info < (3, 9) and platform.system() == "Windows",
reason="Requires Python >= 3.9 on Windows",
)
@pytest.mark.asyncio
@pytest.mark.parametrize("spec", TEST_DATA, ids=TEST_NAMES)
async def test_spec(spec, tmp_path):
state, test_data, is_variant_test = spec
context = test_data.get("context")
unleash_client = await get_async_client(state, context, tmp_path)
try:
if not is_variant_test:
toggle_name = test_data["toggleName"]
expected = test_data["expectedResult"]
assert unleash_client.is_enabled(toggle_name, context) == expected
else:
toggle_name = test_data["toggleName"]
expected = test_data["expectedResult"]
variant = unleash_client.get_variant(toggle_name, context)
assert variant == expected
finally:
await unleash_client.destroy()
Loading
Loading