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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
87 changes: 87 additions & 0 deletions UnleashClient/connectors/_async_connector.py
Original file line number Diff line number Diff line change
@@ -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()
16 changes: 8 additions & 8 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down
187 changes: 187 additions & 0 deletions tests/unit_tests/connectors/test_async_connector.py
Original file line number Diff line number Diff line change
@@ -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()
Loading