diff --git a/CHANGELOG.md b/CHANGELOG.md index bdcd3686..5dd8a07b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ > final changelog + documentation changes (if any). +* (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. Other errors raised by the SDK are not part of the hierarchy yet. * (Bugfix): The in-progress asynchronous metrics reporter keeps impact metrics whose send timed out or was cancelled, and sends them with the next flush. A timed-out send now counts as a failed send, like any status other than 202. Cancelling the reporter mid-send, as `stop()` does, used to lose the impact metrics of that send; they now go out with the final flush. Feature metrics from a failed send are still dropped, as before. `UnleashClient` is unaffected. * (Minor): The in-progress asynchronous client gains an internal `_AsyncScheduler`, in the private `UnleashClient._async_scheduler` module, which runs recurring jobs as tasks on the event loop with the same interval, first-run delay and one-sided jitter as the synchronous `_Scheduler`. The module is not part of the public API and may change or disappear without notice. The asynchronous metrics reporter schedules its sends through it. Nothing changes for code using `UnleashClient`. diff --git a/UnleashClient/connectors/_async_connector.py b/UnleashClient/connectors/_async_connector.py new file mode 100644 index 00000000..399e1827 --- /dev/null +++ b/UnleashClient/connectors/_async_connector.py @@ -0,0 +1,87 @@ +from abc import ABC, abstractmethod +from typing import Optional + +from UnleashClient._async_scheduler import _AsyncScheduler +from UnleashClient.async_transport import AsyncTransport +from UnleashClient.store import FeatureStore + + +class _AsyncBaseConnector(ABC): + def __init__(self, store: FeatureStore) -> None: + """ + :param store: Applies feature state to the engine and the cache, and + emits the events that follow. + """ + self._store = store + + @abstractmethod + async def start(self) -> None: + pass + + @abstractmethod + async def stop(self) -> None: + pass + + +class _AsyncPollingConnector(_AsyncBaseConnector): + """ + Keeps feature state fresh by fetching it on a fixed interval. Starting loads + the cached state and schedules the fetch, without waiting for it. + + Example:: + + connector = _AsyncPollingConnector( + store=store, + transport=transport, + refresh_interval=15, + ) + await connector.start() + + await connector.stop() + """ + + def __init__( + self, + store: FeatureStore, + transport: AsyncTransport, + refresh_interval: float = 15, + refresh_jitter: Optional[float] = None, + ) -> None: + """ + :param store: Applies feature state to the engine and the cache. + :param transport: Performs the fetch against the Unleash server. + :param refresh_interval: Seconds between fetches. + :param refresh_jitter: Maximum seconds to randomly offset each fetch by, or + None for no jitter. + """ + super().__init__(store) + self._transport: AsyncTransport = transport + self._refresh_interval = refresh_interval + self._refresh_jitter = refresh_jitter + self._scheduler: _AsyncScheduler = _AsyncScheduler() + + async def _fetch_and_load(self) -> None: + result = await self._transport.fetch_features(etag=self._store.cached_etag) + + self._store.apply_fetched(raw_state=result.raw_state, etag=result.etag) + + async def start(self) -> None: + """ + Loads the cached feature state, then fetches every ``refresh_interval`` + seconds. The first fetch runs one interval after this returns. + """ + self._store.load_from_cache() + + _ = self._scheduler.every( + interval_seconds=self._refresh_interval, + jitter_seconds=self._refresh_jitter, + fn=self._fetch_and_load, + ) + self._scheduler.start() + + async def stop(self) -> None: + """ + Stops fetching, and returns once a fetch in flight has been interrupted. + Safe to call when :meth:`start` was never called. + """ + await self._scheduler.shutdown() diff --git a/tests/conftest.py b/tests/conftest.py index 57c359ff..7630ab81 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -28,18 +28,18 @@ def dispatcher(recorder): @pytest.fixture() -def cache_empty(): +def cache_empty(tmp_path): cache_name = "pytest_%s" % uuid.uuid4() - temporary_cache = FileCache(cache_name) + temporary_cache = FileCache(cache_name, directory=str(tmp_path)) temporary_cache.mset({METRIC_LAST_SENT_TIME: datetime.now(timezone.utc), ETAG: ""}) yield temporary_cache temporary_cache.destroy() @pytest.fixture() -def cache_full(): +def cache_full(tmp_path): cache_name = "pytest_%s" % uuid.uuid4() - temporary_cache = FileCache(cache_name) + temporary_cache = FileCache(cache_name, directory=str(tmp_path)) temporary_cache.mset( { FEATURES_URL: MOCK_ALL_FEATURES, @@ -52,9 +52,9 @@ def cache_full(): @pytest.fixture() -def cache_custom(): +def cache_custom(tmp_path): cache_name = "pytest_%s" % uuid.uuid4() - temporary_cache = FileCache(cache_name) + temporary_cache = FileCache(cache_name, directory=str(tmp_path)) temporary_cache.mset( { FEATURES_URL: MOCK_CUSTOM_STRATEGY, @@ -67,9 +67,9 @@ def cache_custom(): @pytest.fixture() -def cache_segments(): +def cache_segments(tmp_path): cache_name = "pytest_%s" % uuid.uuid4() - temporary_cache = FileCache(cache_name) + temporary_cache = FileCache(cache_name, directory=str(tmp_path)) temporary_cache.mset( { FEATURES_URL: MOCK_FEATURES_WITH_SEGMENTS_RESPONSE, diff --git a/tests/unit_tests/connectors/test_async_connector.py b/tests/unit_tests/connectors/test_async_connector.py new file mode 100644 index 00000000..5f39dbc3 --- /dev/null +++ b/tests/unit_tests/connectors/test_async_connector.py @@ -0,0 +1,187 @@ +import asyncio +import json +from typing import Callable + +import pytest_asyncio +from pytest import mark +from yggdrasil_engine.engine import UnleashEngine + +from tests.utilities.events import WAIT_TIMEOUT, EventRecorder +from tests.utilities.fake_unleash_server import FakeUnleash +from tests.utilities.mocks.mock_features import MOCK_FEATURE_RESPONSE +from tests.utilities.testing_constants import APP_NAME, ETAG_VALUE +from UnleashClient.async_transport import AsyncTransport +from UnleashClient.config import UnleashConfig +from UnleashClient.connectors._async_connector import _AsyncPollingConnector +from UnleashClient.constants import ETAG, FEATURES_URL +from UnleashClient.events import EventDispatcher, UnleashEventType +from UnleashClient.headers import HeaderFactory +from UnleashClient.store import FeatureStore + +API_PREFIX = "/api" +FEATURES_PATH = API_PREFIX + FEATURES_URL + +INTERVAL = 0.01 +NEVER = 3600 + + +@pytest_asyncio.fixture +async def server(): + fake = FakeUnleash() + await fake.start(API_PREFIX) + try: + yield fake + finally: + await fake.close() + + +@pytest_asyncio.fixture +async def build_connector(server: FakeUnleash): + built = [] + + def _build_connector( + store: FeatureStore, refresh_interval: float = INTERVAL + ) -> _AsyncPollingConnector: + config = UnleashConfig(server.base_url, APP_NAME, request_retries=0) + transport = AsyncTransport(config, HeaderFactory(config)) + connector = _AsyncPollingConnector( + store=store, transport=transport, refresh_interval=refresh_interval + ) + built.append((connector, transport)) + return connector + + try: + yield _build_connector + finally: + for connector, transport in built: + await connector.stop() + await transport.aclose() + + +async def until(predicate: Callable[[], bool]) -> None: + async def poll() -> None: + while not predicate(): + await asyncio.sleep(INTERVAL) + + await asyncio.wait_for(poll(), WAIT_TIMEOUT) + + +def is_enabled(engine: UnleashEngine, flag: str) -> bool: + return bool(engine.is_enabled(flag, {}).is_enabled) + + +@mark.asyncio +async def test_start_makes_cached_state_evaluable_before_any_fetch( + server, build_connector, cache_empty +): + cache_empty.set(FEATURES_URL, json.dumps(MOCK_FEATURE_RESPONSE)) + engine = UnleashEngine() + connector = build_connector( + store=FeatureStore(engine=engine, cache=cache_empty), refresh_interval=NEVER + ) + + await connector.start() + + assert is_enabled(engine, "testFlag") + assert server.calls("GET", FEATURES_PATH) == [] + + +@mark.asyncio +async def test_polling_applies_fetched_state_and_caches_its_etag( + server, build_connector, cache_empty +): + server.on( + "GET", + FEATURES_PATH, + payload=MOCK_FEATURE_RESPONSE, + headers={"etag": ETAG_VALUE}, + repeat=True, + ) + engine = UnleashEngine() + connector = build_connector(store=FeatureStore(engine=engine, cache=cache_empty)) + + await connector.start() + + await until(lambda: is_enabled(engine, "testFlag")) + assert cache_empty.get(ETAG) == ETAG_VALUE + + +@mark.asyncio +async def test_polling_sends_the_cached_etag(server, build_connector, cache_empty): + server.on( + "GET", + FEATURES_PATH, + payload=MOCK_FEATURE_RESPONSE, + headers={"etag": ETAG_VALUE}, + ) + server.on("GET", FEATURES_PATH, status=304, repeat=True) + connector = build_connector( + store=FeatureStore(engine=UnleashEngine(), cache=cache_empty) + ) + + await connector.start() + + await until(lambda: len(server.calls("GET", FEATURES_PATH)) >= 2) + first, second = server.calls("GET", FEATURES_PATH)[:2] + assert "If-None-Match" not in first.headers + assert second.headers["If-None-Match"] == ETAG_VALUE + + +@mark.asyncio +async def test_failed_poll_keeps_the_last_applied_state( + server, build_connector, cache_empty +): + server.on("GET", FEATURES_PATH, payload=MOCK_FEATURE_RESPONSE) + server.on("GET", FEATURES_PATH, status=500, repeat=True) + engine = UnleashEngine() + connector = build_connector(store=FeatureStore(engine=engine, cache=cache_empty)) + + await connector.start() + + await until(lambda: len(server.calls("GET", FEATURES_PATH)) >= 3) + assert is_enabled(engine, "testFlag") + + +@mark.asyncio +async def test_polling_emits_fetched_on_every_fetch_and_ready_once( + server, + build_connector, + cache_empty, + dispatcher: EventDispatcher, + recorder: EventRecorder, +): + server.on("GET", FEATURES_PATH, payload=MOCK_FEATURE_RESPONSE, repeat=True) + connector = build_connector( + store=FeatureStore(engine=UnleashEngine(), cache=cache_empty, events=dispatcher) + ) + + await connector.start() + await until(lambda: len(recorder.of_type(UnleashEventType.FETCHED)) >= 2) + await connector.stop() + dispatcher.close(timeout=WAIT_TIMEOUT) + + assert len(recorder.of_type(UnleashEventType.READY)) == 1 + + +@mark.asyncio +async def test_stop_interrupts_a_fetch_in_flight(server, build_connector, cache_empty): + server.on("GET", FEATURES_PATH, payload=MOCK_FEATURE_RESPONSE, hang=True) + engine = UnleashEngine() + connector = build_connector(store=FeatureStore(engine=engine, cache=cache_empty)) + + await connector.start() + await until(lambda: len(server.calls("GET", FEATURES_PATH)) == 1) + await asyncio.wait_for(connector.stop(), WAIT_TIMEOUT) + await server.close() + + assert not is_enabled(engine, "testFlag") + assert len(server.calls("GET", FEATURES_PATH)) == 1 + + +@mark.asyncio +async def test_stop_is_safe_when_never_started(build_connector, cache_empty): + connector = build_connector( + store=FeatureStore(engine=UnleashEngine(), cache=cache_empty) + ) + + await connector.stop()