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
2 changes: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@
* (Minor): Event callbacks are now invoked on a dedicated background thread instead of on whichever thread produced the event. `is_enabled()` and `get_variant()` no longer wait for your callback, so a slow callback can't hold up flag evaluation. Three consequences worth knowing about: callbacks can no longer read thread local state from the caller (Flask `g`, the current Django request, contextvars); they return before the callback has run, so tests asserting straight after the call now need to wait; and reassigning `unleash_event_callback` after construction is no longer honoured.
* (Minor): Connectors take an `EventDispatcher` instead of `ready_callback`/`event_callback`. These classes aren't part of the documented API, so this only affects code importing from `UnleashClient.connectors` directly.
* (Minor): Request headers are now assembled once, by an internal `HeaderFactory`, and passed to each collaborator complete. The headers on the wire are unchanged. `PollingConnector` and `StreamingConnector` no longer add `unleash-interval` and `Accept`/`Content-Type`/`Unleash-Client-Spec` themselves, so code importing from `UnleashClient.connectors` directly must now supply complete headers. Nothing changes for code using `UnleashClient`.
* (Minor): Applying feature state (the cache write, the handover to the engine, and the READY and FETCHED events that follow) now happens in one internal `FeatureStore` instead of being re-implemented by each connector. The cache writes, engine updates and events are exactly what they were. Connectors now take a `store` instead of `engine`, `cache` and `events`, and `BaseConnector.load_features()` is gone, so code importing from `UnleashClient.connectors` directly must build a `FeatureStore` and call `store.load_from_cache()`. Nothing changes for code using `UnleashClient`.
* (Minor): Applying feature state (the cache write, the handover to the engine, and the READY and FETCHED events that follow) now happens in one internal `_FeatureStore`, in the private `UnleashClient._feature_store` module, instead of being re-implemented by each connector. The module is not part of the public API and may change or disappear without notice. The cache writes, engine updates and events are exactly what they were. Connectors now take a `store` instead of `engine`, `cache` and `events`, and `BaseConnector.load_features()` is gone, so code importing from `UnleashClient.connectors` directly must build a `_FeatureStore` and call `store.load_from_cache()`. Nothing changes for code using `UnleashClient`.
* (Minor): The `engine` and `cache` attributes are gone from `UnleashClient`. Both were always configured through the constructor (pass `cache=` to supply your own), and neither appears in the documented API.
* (Bugfix): `refresh_jitter` now reaches the polling job. It was accepted by the constructor, documented, and applied to the offline refresh job, but never passed to the polling connector, so jitter was silently dropped in the default polling mode.
* (Minor): Constructor arguments are now normalized once into an internal `UnleashConfig` object rather than being copied onto the client attribute by attribute. The public `unleash_*` attributes keep their exact values and stay writable, reading and writing through that object, so nothing in calling code needs to change. This is groundwork for an asynchronous client that shares the same configuration handling.
Expand Down
10 changes: 9 additions & 1 deletion UnleashClient/store.py → UnleashClient/_feature_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,22 @@
from UnleashClient.utils import LOGGER


class FeatureStore:
class _FeatureStore:
"""
Owns what happens to feature state once it has arrived: the cache write, the
handover to the engine, and the events that follow.

There is one method per source rather than a single ``apply``, because the
three steps happen in a different order, over different payloads, with
different failure handling depending on where the state came from.

Example::

store = _FeatureStore(engine=UnleashEngine(), cache=cache, events=dispatcher)
store.load_from_cache()

result = transport.fetch_features(etag=store.cached_etag)
store.apply_fetched(raw_state=result.raw_state, etag=result.etag)
"""

def __init__(
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/clients/async_unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from UnleashClient._async_scheduler import _AsyncScheduler
from UnleashClient._async_transport import _AsyncTransport
from UnleashClient._evaluator import _Evaluator
from UnleashClient._feature_store import _FeatureStore
from UnleashClient._instance_registry import _get_instance_registry
from UnleashClient.async_metrics_reporter import AsyncMetricsReporter
from UnleashClient.cache import BaseCache, FileCache
Expand All @@ -19,7 +20,6 @@
from UnleashClient.events import BaseEvent, EventDispatcher
from UnleashClient.headers import HeaderFactory
from UnleashClient.impact_metrics import ImpactMetrics
from UnleashClient.store import FeatureStore
from UnleashClient.utils import InstanceAllowType

_NOT_IMPLEMENTED = (
Expand Down Expand Up @@ -104,7 +104,7 @@ def __init__( # noqa: PLR0913, PLR0917
self._cache: BaseCache = cache or FileCache(
self._config.app_name, directory=cache_directory
)
self._store: FeatureStore = FeatureStore(
self._store: _FeatureStore = _FeatureStore(
engine=self._engine, cache=self._cache, events=self._event_dispatcher
)
self._evaluator: _Evaluator = _Evaluator(
Expand Down
6 changes: 3 additions & 3 deletions UnleashClient/clients/unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from yggdrasil_engine.engine import UnleashEngine

from UnleashClient._evaluator import _Evaluator
from UnleashClient._feature_store import _FeatureStore
from UnleashClient._instance_registry import _get_instance_registry
from UnleashClient._scheduler import _ScheduledJob, _Scheduler
from UnleashClient._transport import _Transport
Expand Down Expand Up @@ -43,7 +44,6 @@
from UnleashClient.impact_metrics import ImpactMetrics
from UnleashClient.metrics_reporter import MetricsReporter
from UnleashClient.payloads import build_register_payload
from UnleashClient.store import FeatureStore
from UnleashClient.utils import (
LOGGER,
InstanceAllowType,
Expand Down Expand Up @@ -207,7 +207,7 @@ def __init__( # noqa: PLR0913, PLR0917
self._cache.mset({METRIC_LAST_SENT_TIME: datetime.now(timezone.utc), ETAG: ""})
self.unleash_bootstrapped = self._cache.bootstrapped

self._store = FeatureStore(
self._store = _FeatureStore(
engine=self._engine, cache=self._cache, events=self.__events
)

Expand Down Expand Up @@ -240,7 +240,7 @@ def __init__( # noqa: PLR0913, PLR0917
# move it earlier for bootstrapped clients. See the TODO on
# BootstrapConnector.
BootstrapConnector(
store=FeatureStore(engine=self._engine, cache=self._cache)
store=_FeatureStore(engine=self._engine, cache=self._cache)
).start()

self.connector: BaseConnector = None
Expand Down
6 changes: 3 additions & 3 deletions UnleashClient/connectors/_async_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,11 @@

from UnleashClient._async_scheduler import _AsyncScheduler
from UnleashClient._async_transport import _AsyncTransport
from UnleashClient.store import FeatureStore
from UnleashClient._feature_store import _FeatureStore


class _AsyncBaseConnector(ABC):
def __init__(self, store: FeatureStore) -> None:
def __init__(self, store: _FeatureStore) -> None:
"""
:param store: Applies feature state to the engine and the cache, and
emits the events that follow.
Expand Down Expand Up @@ -42,7 +42,7 @@ class _AsyncPollingConnector(_AsyncBaseConnector):

def __init__(
self,
store: FeatureStore,
store: _FeatureStore,
transport: _AsyncTransport,
refresh_interval: float = 15,
refresh_jitter: Optional[float] = None,
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/connectors/base_connector.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
from abc import ABC, abstractmethod

from UnleashClient.store import FeatureStore
from UnleashClient._feature_store import _FeatureStore


class BaseConnector(ABC):
def __init__(self, store: FeatureStore):
def __init__(self, store: _FeatureStore):
"""
:param store: Applies feature state to the engine and the cache, and
emits the events that follow.
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/connectors/bootstrap_connector.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
from UnleashClient.store import FeatureStore
from UnleashClient._feature_store import _FeatureStore

from .base_connector import BaseConnector


class BootstrapConnector(BaseConnector):
def __init__(
self,
store: FeatureStore,
store: _FeatureStore,
):
super().__init__(store)
self.job = None
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/connectors/offline_connector.py
Original file line number Diff line number Diff line change
@@ -1,15 +1,15 @@
from typing import Optional

from UnleashClient._feature_store import _FeatureStore
from UnleashClient._scheduler import _ScheduledJob, _Scheduler
from UnleashClient.store import FeatureStore

from .base_connector import BaseConnector


class OfflineConnector(BaseConnector):
def __init__(
self,
store: FeatureStore,
store: _FeatureStore,
scheduler: _Scheduler,
refresh_interval: int = 15,
refresh_jitter: Optional[int] = None,
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/connectors/polling_connector.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
from typing import Optional

from UnleashClient._feature_store import _FeatureStore
from UnleashClient._scheduler import _ScheduledJob, _Scheduler
from UnleashClient._transport import _Transport
from UnleashClient.store import FeatureStore

from .base_connector import BaseConnector

Expand All @@ -12,7 +12,7 @@ class PollingConnector(BaseConnector):

def __init__(
self,
store: FeatureStore,
store: _FeatureStore,
scheduler: _Scheduler,
transport: _Transport,
refresh_interval: int = 15,
Expand Down
4 changes: 2 additions & 2 deletions UnleashClient/connectors/streaming_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,16 @@
from ld_eventsource import SSEClient
from ld_eventsource.config import ConnectStrategy, ErrorStrategy, RetryDelayStrategy

from UnleashClient._feature_store import _FeatureStore
from UnleashClient.connectors.base_connector import BaseConnector
from UnleashClient.constants import STREAMING_URL
from UnleashClient.store import FeatureStore
from UnleashClient.utils import LOGGER


class StreamingConnector(BaseConnector):
def __init__(
self,
store: FeatureStore,
store: _FeatureStore,
url: str,
headers: dict,
request_timeout: int,
Expand Down
20 changes: 11 additions & 9 deletions tests/unit_tests/connectors/test_async_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,12 @@
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._feature_store import _FeatureStore
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
Expand All @@ -40,7 +40,7 @@ async def build_connector(server: FakeUnleash):
built = []

def _build_connector(
store: FeatureStore, refresh_interval: float = INTERVAL
store: _FeatureStore, refresh_interval: float = INTERVAL
) -> _AsyncPollingConnector:
config = UnleashConfig(server.base_url, APP_NAME, request_retries=0)
transport = _AsyncTransport(config, HeaderFactory(config))
Expand Down Expand Up @@ -77,7 +77,7 @@ async def test_start_makes_cached_state_evaluable_before_any_fetch(
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
store=_FeatureStore(engine=engine, cache=cache_empty), refresh_interval=NEVER
)

await connector.start()
Expand All @@ -98,7 +98,7 @@ async def test_polling_applies_fetched_state_and_caches_its_etag(
repeat=True,
)
engine = UnleashEngine()
connector = build_connector(store=FeatureStore(engine=engine, cache=cache_empty))
connector = build_connector(store=_FeatureStore(engine=engine, cache=cache_empty))

await connector.start()

Expand All @@ -116,7 +116,7 @@ async def test_polling_sends_the_cached_etag(server, build_connector, cache_empt
)
server.on("GET", FEATURES_PATH, status=304, repeat=True)
connector = build_connector(
store=FeatureStore(engine=UnleashEngine(), cache=cache_empty)
store=_FeatureStore(engine=UnleashEngine(), cache=cache_empty)
)

await connector.start()
Expand All @@ -134,7 +134,7 @@ async def test_failed_poll_keeps_the_last_applied_state(
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))
connector = build_connector(store=_FeatureStore(engine=engine, cache=cache_empty))

await connector.start()

Expand All @@ -152,7 +152,9 @@ async def test_polling_emits_fetched_on_every_fetch_and_ready_once(
):
server.on("GET", FEATURES_PATH, payload=MOCK_FEATURE_RESPONSE, repeat=True)
connector = build_connector(
store=FeatureStore(engine=UnleashEngine(), cache=cache_empty, events=dispatcher)
store=_FeatureStore(
engine=UnleashEngine(), cache=cache_empty, events=dispatcher
)
)

await connector.start()
Expand All @@ -167,7 +169,7 @@ async def test_polling_emits_fetched_on_every_fetch_and_ready_once(
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))
connector = build_connector(store=_FeatureStore(engine=engine, cache=cache_empty))

await connector.start()
await until(lambda: len(server.calls("GET", FEATURES_PATH)) == 1)
Expand All @@ -181,7 +183,7 @@ async def test_stop_interrupts_a_fetch_in_flight(server, build_connector, cache_
@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)
store=_FeatureStore(engine=UnleashEngine(), cache=cache_empty)
)

await connector.stop()
12 changes: 6 additions & 6 deletions tests/unit_tests/connectors/test_offline_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,11 @@

from tests.utilities.events import WAIT_TIMEOUT, EventRecorder
from tests.utilities.mocks.mock_features import MOCK_FEATURE_RESPONSE
from UnleashClient._feature_store import _FeatureStore
from UnleashClient._scheduler import _Scheduler
from UnleashClient.connectors import OfflineConnector
from UnleashClient.constants import FEATURES_URL
from UnleashClient.events import EventDispatcher, UnleashEventType
from UnleashClient.store import FeatureStore


def test_offline_connector_loads_features_on_start(cache_empty):
Expand All @@ -19,7 +19,7 @@ def test_offline_connector_loads_features_on_start(cache_empty):
temp_cache.set(FEATURES_URL, json.dumps(MOCK_FEATURE_RESPONSE))

connector = OfflineConnector(
store=FeatureStore(engine=engine, cache=temp_cache),
store=_FeatureStore(engine=engine, cache=temp_cache),
scheduler=scheduler,
)

Expand All @@ -36,7 +36,7 @@ def test_offline_connector_start_stop(cache_empty):
temp_cache.set(FEATURES_URL, json.dumps(MOCK_FEATURE_RESPONSE))

connector = OfflineConnector(
store=FeatureStore(engine=engine, cache=temp_cache),
store=_FeatureStore(engine=engine, cache=temp_cache),
scheduler=scheduler,
refresh_interval=1,
)
Expand All @@ -59,7 +59,7 @@ def test_offline_connector_emits_ready_event(
temp_cache.set(FEATURES_URL, json.dumps(MOCK_FEATURE_RESPONSE))

connector = OfflineConnector(
store=FeatureStore(engine=engine, cache=temp_cache, events=dispatcher),
store=_FeatureStore(engine=engine, cache=temp_cache, events=dispatcher),
scheduler=scheduler,
)

Expand All @@ -78,7 +78,7 @@ def test_offline_connector_emits_ready_on_an_empty_cache(
scheduler = _Scheduler()

connector = OfflineConnector(
store=FeatureStore(
store=_FeatureStore(
engine=UnleashEngine(), cache=cache_empty, events=dispatcher
),
scheduler=scheduler,
Expand All @@ -100,7 +100,7 @@ def test_offline_connector_without_a_dispatcher_does_not_emit(cache_empty):
temp_cache.set(FEATURES_URL, json.dumps(MOCK_FEATURE_RESPONSE))

connector = OfflineConnector(
store=FeatureStore(engine=engine, cache=temp_cache),
store=_FeatureStore(engine=engine, cache=temp_cache),
scheduler=scheduler,
)

Expand Down
Loading
Loading