From 632f00bca59a3ca91113808e92662bd6d398a805 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pere=20Pic=C3=B3?= Date: Mon, 28 Sep 2026 11:30:59 +0200 Subject: [PATCH] refactor: make EventDispatcher a private collaborator --- CHANGELOG.md | 2 +- UnleashClient/_evaluator.py | 6 +- UnleashClient/_event_dispatcher.py | 135 ++++++++++++++++++ UnleashClient/_feature_store.py | 4 +- UnleashClient/clients/async_unleash_client.py | 7 +- UnleashClient/clients/unleash_client.py | 8 +- .../connectors/bootstrap_connector.py | 2 +- UnleashClient/events.py | 117 +-------------- tests/conftest.py | 4 +- .../connectors/test_async_connector.py | 5 +- .../connectors/test_offline_connector.py | 7 +- .../connectors/test_polling_connector.py | 5 +- tests/unit_tests/test_events.py | 100 ++++++------- tests/unit_tests/test_feature_store.py | 5 +- tests/utilities/event_callbacks.py | 12 +- 15 files changed, 222 insertions(+), 197 deletions(-) create mode 100644 UnleashClient/_event_dispatcher.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 201c252a..7884994a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -34,7 +34,7 @@ * (Minor): Upgraded to yggdrasil-engine 2.0. Counting toggle and variant evaluations, and deciding whether an impression event is due, now happen inside the engine rather than in the SDK. `is_enabled()` and `get_variant()` return the same types as before, so no calling code needs to change. * (Minor): A `fallback_function` that raises now results in `False` and a logged warning, instead of the exception propagating out of `is_enabled()`. The toggle is not counted in that case. * (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): Connectors take an internal `_EventDispatcher`, in the private `UnleashClient._event_dispatcher` module, which is not part of the public API and may change or disappear without notice, instead of `ready_callback`/`event_callback`. Connectors 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`, in the private `UnleashClient._headers` module, which is not part of the public API and may change or disappear without notice, 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`, 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. diff --git a/UnleashClient/_evaluator.py b/UnleashClient/_evaluator.py index 07587e31..547d9d8f 100644 --- a/UnleashClient/_evaluator.py +++ b/UnleashClient/_evaluator.py @@ -7,9 +7,9 @@ from yggdrasil_engine.engine import UnleashEngine from UnleashClient._context import _ContextEnricher +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient.config import UnleashConfig from UnleashClient.events import ( - EventDispatcher, UnleashEvent, UnleashEventType, ) @@ -31,7 +31,7 @@ def __init__( engine: UnleashEngine, enricher: _ContextEnricher, config: UnleashConfig, - events: Optional[EventDispatcher] = None, + events: Optional[_EventDispatcher] = None, ) -> None: """ :param engine: Feature evaluation engine instance (UnleashEngine). @@ -42,7 +42,7 @@ def __init__( self._engine: UnleashEngine = engine self._enricher: _ContextEnricher = enricher self._config: UnleashConfig = config - self._events: Optional[EventDispatcher] = events + self._events: Optional[_EventDispatcher] = events # pylint: disable=broad-except def is_enabled( diff --git a/UnleashClient/_event_dispatcher.py b/UnleashClient/_event_dispatcher.py new file mode 100644 index 00000000..33cbb471 --- /dev/null +++ b/UnleashClient/_event_dispatcher.py @@ -0,0 +1,135 @@ +"""Background delivery of events to the user's event callback.""" + +import queue +import threading +from typing import Callable, Optional, Union + +from UnleashClient.events import BaseEvent, UnleashEventType +from UnleashClient.utils import LOGGER + +DEFAULT_MAX_QUEUE_SIZE = 100 +DEFAULT_TIMEOUT = 2.0 + + +class _ShutdownMarker: + pass + + +_SHUTDOWN_WAKER: _ShutdownMarker = _ShutdownMarker() + + +class _EventDispatcher: + """ + Delivers events to a user callback on a dedicated background thread, so a + slow or failing callback never holds up flag evaluation. Events beyond the + queue's capacity are dropped and counted, and READY is delivered at most + once. + + Example:: + + def on_event(event: BaseEvent) -> None: + print(event.event_type) + + dispatcher = _EventDispatcher(on_event) + dispatcher.emit_event(UnleashReadyEvent(UnleashEventType.READY, uuid4())) + + dispatcher.close() + """ + + def __init__( + self, + callback: Callable[[BaseEvent], None], + max_size: int = DEFAULT_MAX_QUEUE_SIZE, + ) -> None: + self._callback: Callable[[BaseEvent], None] = callback + self._queue: queue.Queue[Union[_ShutdownMarker, BaseEvent]] = queue.Queue( + maxsize=max_size + ) + self._lock = threading.Lock() + self._thread: Optional[threading.Thread] = None + self._closed = threading.Event() + self._closing = threading.Event() + self._dropped = 0 + self._ready_delivered = False + + def emit_event(self, event: BaseEvent) -> None: + with self._lock: + if self._closed.is_set() or self._closing.is_set(): + return + + self._start_worker() + + is_ready_event = event.event_type == UnleashEventType.READY + + if is_ready_event and self._ready_delivered: + return + + try: + self._queue.put_nowait(event) + if is_ready_event: + self._ready_delivered = True + except queue.Full: + self._dropped += 1 + should_warn = self._dropped == 1 + else: + return + + if should_warn: + LOGGER.warning("Unleash event queue is full; events are being dropped.") + + @property + def dropped_events(self) -> int: + with self._lock: + return self._dropped + + def close(self, timeout: float = DEFAULT_TIMEOUT) -> None: + with self._lock: + if self._closed.is_set(): + return + + self._closing.set() + + thread = self._thread + + if thread is None: + return + + try: + self._queue.put_nowait(_SHUTDOWN_WAKER) + except queue.Full: + pass + + if threading.current_thread() is not thread: + thread.join(timeout) + + self._closed.set() + + def _start_worker(self) -> None: + if self._thread is not None: + return + + self._thread = threading.Thread( + target=self._run, + name="UnleashEventDispatcher", + daemon=True, + ) + self._thread.start() + + def _run(self) -> None: + try: + while not self._closing.is_set(): + item: Union[_ShutdownMarker, BaseEvent] = self._queue.get() + + if isinstance(item, _ShutdownMarker): + return + + try: + self._callback(item) + except Exception: + LOGGER.exception("Error in event callback") + + if self._closed.is_set(): + return + finally: + with self._lock: + self._thread = None diff --git a/UnleashClient/_feature_store.py b/UnleashClient/_feature_store.py index 5fdfd169..82823b01 100644 --- a/UnleashClient/_feature_store.py +++ b/UnleashClient/_feature_store.py @@ -5,11 +5,11 @@ from yggdrasil_engine.engine import UnleashEngine +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient.cache import BaseCache from UnleashClient.constants import ETAG, FEATURES_URL from UnleashClient.events import ( BaseEvent, - EventDispatcher, UnleashEventType, UnleashFetchedEvent, UnleashReadyEvent, @@ -39,7 +39,7 @@ def __init__( self, engine: UnleashEngine, cache: BaseCache, - events: Optional[EventDispatcher] = None, + events: Optional[_EventDispatcher] = None, ) -> None: """ :param engine: Feature evaluation engine instance (UnleashEngine). diff --git a/UnleashClient/clients/async_unleash_client.py b/UnleashClient/clients/async_unleash_client.py index f6b76ceb..be757081 100644 --- a/UnleashClient/clients/async_unleash_client.py +++ b/UnleashClient/clients/async_unleash_client.py @@ -11,6 +11,7 @@ from UnleashClient._async_transport import _AsyncTransport from UnleashClient._context import _ContextEnricher from UnleashClient._evaluator import _Evaluator +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient._feature_store import _FeatureStore from UnleashClient._headers import _HeaderFactory from UnleashClient._instance_registry import _get_instance_registry @@ -18,7 +19,7 @@ from UnleashClient.cache import BaseCache, FileCache from UnleashClient.config import ExperimentalMode, UnleashConfig from UnleashClient.constants import REQUEST_RETRIES, REQUEST_TIMEOUT -from UnleashClient.events import BaseEvent, EventDispatcher +from UnleashClient.events import BaseEvent from UnleashClient.impact_metrics import ImpactMetrics from UnleashClient.utils import InstanceAllowType @@ -87,8 +88,8 @@ def __init__( # noqa: PLR0913, PLR0917 self._enricher: _ContextEnricher = _ContextEnricher(self._config) self._headers: _HeaderFactory = _HeaderFactory(self._config) - self._event_dispatcher: Optional[EventDispatcher] = ( - EventDispatcher(event_callback) if event_callback is not None else None + self._event_dispatcher: Optional[_EventDispatcher] = ( + _EventDispatcher(event_callback) if event_callback is not None else None ) _get_instance_registry().register( diff --git a/UnleashClient/clients/unleash_client.py b/UnleashClient/clients/unleash_client.py index ac7f2e6f..ab6af958 100644 --- a/UnleashClient/clients/unleash_client.py +++ b/UnleashClient/clients/unleash_client.py @@ -11,6 +11,7 @@ from UnleashClient._context import _ContextEnricher from UnleashClient._evaluator import _Evaluator +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient._feature_store import _FeatureStore from UnleashClient._headers import _HeaderFactory from UnleashClient._instance_registry import _get_instance_registry @@ -39,7 +40,6 @@ ) from UnleashClient.events import ( BaseEvent, - EventDispatcher, UnleashEventType, UnleashReadyEvent, ) @@ -65,7 +65,7 @@ def build_ready_callback( Builds a callback function that can be used to notify when the Unleash client is ready. .. deprecated:: - READY is now emitted through :class:`UnleashClient.events.EventDispatcher`, + READY is now emitted through :class:`UnleashClient._event_dispatcher._EventDispatcher`, which deduplicates it itself. This helper is retained for backwards compatibility and is no longer used internally. """ @@ -182,8 +182,8 @@ def __init__( # noqa: PLR0913, PLR0917 self.unleash_event_callback = event_callback # Events are handed to the dispatcher, which delivers them to the user's # callback on its own thread. The callback is never called from here. - self.__events: Optional[EventDispatcher] = ( - EventDispatcher(event_callback) if event_callback is not None else None + self.__events: Optional[_EventDispatcher] = ( + _EventDispatcher(event_callback) if event_callback is not None else None ) self._lifecycle_lock = threading.RLock() self._closed = threading.Event() diff --git a/UnleashClient/connectors/bootstrap_connector.py b/UnleashClient/connectors/bootstrap_connector.py index 191600b6..c7ebca76 100644 --- a/UnleashClient/connectors/bootstrap_connector.py +++ b/UnleashClient/connectors/bootstrap_connector.py @@ -11,7 +11,7 @@ def __init__( super().__init__(store) self.job = None - # TODO: the client hands this connector a store with no EventDispatcher, so + # TODO: the client hands this connector a store with no _EventDispatcher, so # bootstrapping does not emit a READY event. Bootstrapped clients only see # READY once initialize_client() builds a polling, streaming or offline # connector. Passing the dispatcher here would emit READY from diff --git a/UnleashClient/events.py b/UnleashClient/events.py index 8b89b4ef..670489c3 100644 --- a/UnleashClient/events.py +++ b/UnleashClient/events.py @@ -1,13 +1,9 @@ -import queue -import threading from dataclasses import dataclass from enum import Enum from json import loads -from typing import Callable, Optional, Union +from typing import Optional from uuid import UUID -from UnleashClient.utils import LOGGER - class UnleashEventType(Enum): """ @@ -64,114 +60,3 @@ def features(self) -> dict: if not hasattr(self, "_parsed_payload"): self._parsed_payload = loads(self.raw_features)["features"] return self._parsed_payload - - -DEFAULT_MAX_QUEUE_SIZE = 100 -DEFAULT_TIMEOUT = 2.0 - - -class _ShutdownMarker: - pass - - -_SHUTDOWN_WAKER: _ShutdownMarker = _ShutdownMarker() - - -class EventDispatcher: - def __init__( - self, - callback: Callable[[BaseEvent], None], - max_size: int = DEFAULT_MAX_QUEUE_SIZE, - ) -> None: - self._callback: Callable[[BaseEvent], None] = callback - self._queue: queue.Queue[Union[_ShutdownMarker, BaseEvent]] = queue.Queue( - maxsize=max_size - ) - self._lock = threading.Lock() - self._thread: Optional[threading.Thread] = None - self._closed = threading.Event() - self._closing = threading.Event() - self._dropped = 0 - self._ready_delivered = False - - def emit_event(self, event: BaseEvent) -> None: - with self._lock: - if self._closed.is_set() or self._closing.is_set(): - return - - self._start_worker() - - is_ready_event = event.event_type == UnleashEventType.READY - - if is_ready_event and self._ready_delivered: - return - - try: - self._queue.put_nowait(event) - if is_ready_event: - self._ready_delivered = True - except queue.Full: - self._dropped += 1 - should_warn = self._dropped == 1 - else: - return - - if should_warn: - LOGGER.warning("Unleash event queue is full; events are being dropped.") - - @property - def dropped_events(self) -> int: - with self._lock: - return self._dropped - - def close(self, timeout: float = DEFAULT_TIMEOUT) -> None: - with self._lock: - if self._closed.is_set(): - return - - self._closing.set() - - thread = self._thread - - if thread is None: - return - - try: - self._queue.put_nowait(_SHUTDOWN_WAKER) - except queue.Full: - pass - - if threading.current_thread() is not thread: - thread.join(timeout) - - self._closed.set() - - def _start_worker(self) -> None: - if self._thread is not None: - return - - self._thread = threading.Thread( - target=self._run, - name="UnleashEventDispatcher", - daemon=True, - ) - self._thread.start() - - def _run(self) -> None: - try: - while not self._closing.is_set(): - item: Union[_ShutdownMarker, BaseEvent] = self._queue.get() - - if isinstance(item, _ShutdownMarker): - return - - try: - self._callback(item) - except Exception: - LOGGER.exception("Error in event callback") - - if self._closed.is_set(): - return - finally: - with self._lock: - self._thread = None diff --git a/tests/conftest.py b/tests/conftest.py index 7630ab81..6a7c994e 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -6,9 +6,9 @@ from tests.utilities.events import EventRecorder from tests.utilities.mocks import MOCK_ALL_FEATURES, MOCK_CUSTOM_STRATEGY from tests.utilities.mocks.mock_features import MOCK_FEATURES_WITH_SEGMENTS_RESPONSE +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient.cache import FileCache from UnleashClient.constants import ETAG, FEATURES_URL, METRIC_LAST_SENT_TIME -from UnleashClient.events import EventDispatcher @pytest.fixture() @@ -22,7 +22,7 @@ def dispatcher(recorder): A dispatcher wired to the `recorder` fixture, closed on teardown so a wedged worker can't leak into the next test. """ - event_dispatcher = EventDispatcher(recorder) + event_dispatcher = _EventDispatcher(recorder) yield event_dispatcher event_dispatcher.close(timeout=1) diff --git a/tests/unit_tests/connectors/test_async_connector.py b/tests/unit_tests/connectors/test_async_connector.py index b751abc0..56377332 100644 --- a/tests/unit_tests/connectors/test_async_connector.py +++ b/tests/unit_tests/connectors/test_async_connector.py @@ -11,12 +11,13 @@ 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._event_dispatcher import _EventDispatcher from UnleashClient._feature_store import _FeatureStore from UnleashClient._headers import _HeaderFactory 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.events import UnleashEventType API_PREFIX = "/api" FEATURES_PATH = API_PREFIX + FEATURES_URL @@ -147,7 +148,7 @@ async def test_polling_emits_fetched_on_every_fetch_and_ready_once( server, build_connector, cache_empty, - dispatcher: EventDispatcher, + dispatcher: _EventDispatcher, recorder: EventRecorder, ): server.on("GET", FEATURES_PATH, payload=MOCK_FEATURE_RESPONSE, repeat=True) diff --git a/tests/unit_tests/connectors/test_offline_connector.py b/tests/unit_tests/connectors/test_offline_connector.py index f5bd2d58..5db33b4c 100644 --- a/tests/unit_tests/connectors/test_offline_connector.py +++ b/tests/unit_tests/connectors/test_offline_connector.py @@ -4,11 +4,12 @@ from tests.utilities.events import WAIT_TIMEOUT, EventRecorder from tests.utilities.mocks.mock_features import MOCK_FEATURE_RESPONSE +from UnleashClient._event_dispatcher import _EventDispatcher 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.events import UnleashEventType def test_offline_connector_loads_features_on_start(cache_empty): @@ -51,7 +52,7 @@ def test_offline_connector_start_stop(cache_empty): def test_offline_connector_emits_ready_event( - cache_empty, dispatcher: EventDispatcher, recorder: EventRecorder + cache_empty, dispatcher: _EventDispatcher, recorder: EventRecorder ): engine = UnleashEngine() scheduler = _Scheduler() @@ -73,7 +74,7 @@ def test_offline_connector_emits_ready_event( def test_offline_connector_emits_ready_on_an_empty_cache( - cache_empty, dispatcher: EventDispatcher, recorder: EventRecorder + cache_empty, dispatcher: _EventDispatcher, recorder: EventRecorder ): scheduler = _Scheduler() diff --git a/tests/unit_tests/connectors/test_polling_connector.py b/tests/unit_tests/connectors/test_polling_connector.py index b39d6fb0..fad9caa3 100644 --- a/tests/unit_tests/connectors/test_polling_connector.py +++ b/tests/unit_tests/connectors/test_polling_connector.py @@ -18,6 +18,7 @@ REQUEST_TIMEOUT, URL, ) +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient._feature_store import _FeatureStore from UnleashClient._headers import _HeaderFactory from UnleashClient._scheduler import _Scheduler @@ -25,7 +26,7 @@ from UnleashClient.config import UnleashConfig from UnleashClient.connectors import PollingConnector from UnleashClient.constants import ETAG, FEATURES_URL -from UnleashClient.events import EventDispatcher, UnleashEventType +from UnleashClient.events import UnleashEventType FULL_FEATURE_URL = URL + FEATURES_URL @@ -115,7 +116,7 @@ def test_polling_connector_fetch_and_load_failure(cache_empty): @responses.activate def test_polling_connector_emits_fetched_and_ready( - cache_empty, dispatcher: EventDispatcher, recorder: EventRecorder + cache_empty, dispatcher: _EventDispatcher, recorder: EventRecorder ): engine = UnleashEngine() scheduler = _Scheduler() diff --git a/tests/unit_tests/test_events.py b/tests/unit_tests/test_events.py index e38bdc08..d64edb2d 100644 --- a/tests/unit_tests/test_events.py +++ b/tests/unit_tests/test_events.py @@ -20,9 +20,9 @@ variant_event, worker_threads, ) +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient.events import ( BaseEvent, - EventDispatcher, UnleashEventType, ) @@ -58,7 +58,7 @@ def ready_deliveries(callback: RecorderCallback) -> List[BaseEvent]: @pytest.fixture() -def dispatcher_factory() -> Iterator[Callable[..., EventDispatcher]]: +def dispatcher_factory() -> Iterator[Callable[..., _EventDispatcher]]: """ Builds dispatchers and guarantees they're torn down, so a wedged worker thread can't leak into the next test. @@ -66,10 +66,10 @@ def dispatcher_factory() -> Iterator[Callable[..., EventDispatcher]]: Blocking callbacks are released before the close: a test that fails before its own release() would otherwise leave the worker parked inside the callback. """ - created: List[Tuple[object, EventDispatcher]] = [] + created: List[Tuple[object, _EventDispatcher]] = [] - def _build(callback, *args, **kwargs) -> EventDispatcher: - dispatcher = EventDispatcher(callback, *args, **kwargs) + def _build(callback, *args, **kwargs) -> _EventDispatcher: + dispatcher = _EventDispatcher(callback, *args, **kwargs) created.append((callback, dispatcher)) return dispatcher @@ -84,7 +84,7 @@ def _build(callback, *args, **kwargs) -> EventDispatcher: class TestDelivery: def test_an_emitted_event_reaches_the_callback( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -95,7 +95,7 @@ def test_an_emitted_event_reaches_the_callback( assert callback.feature_names == ["testFlag"] def test_the_callback_is_handed_the_event_that_was_emitted( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -107,7 +107,7 @@ def test_the_callback_is_handed_the_event_that_was_emitted( assert callback.events[0] is event def test_events_arrive_in_the_order_they_were_emitted( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -119,7 +119,7 @@ def test_events_arrive_in_the_order_they_were_emitted( assert callback.feature_names == [str(index) for index in range(20)] def test_every_kind_of_event_goes_to_the_same_callback( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -132,7 +132,7 @@ def test_every_kind_of_event_goes_to_the_same_callback( assert callback.events == emitted def test_emitting_does_not_wait_for_the_callback( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback) @@ -147,10 +147,10 @@ def test_emitting_does_not_wait_for_the_callback( assert dispatcher.dropped_events == 0 # in the queue, not on the floor def test_the_emitter_does_not_wait_for_the_whole_callback_duration( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback: SlowCallback = SlowCallback(delay=0.05) - dispatcher: EventDispatcher = dispatcher_factory(callback) + dispatcher: _EventDispatcher = dispatcher_factory(callback) for _ in range(20): dispatcher.emit_event(flag_event()) @@ -160,10 +160,10 @@ def test_the_emitter_does_not_wait_for_the_whole_callback_duration( ) # If emitter had waited for the callback, this would be exactly 20. def test_emitter_does_not_drop_events_despite_slow_callback( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback: SlowCallback = SlowCallback(delay=0.05) - dispatcher: EventDispatcher = dispatcher_factory(callback) + dispatcher: _EventDispatcher = dispatcher_factory(callback) for _ in range(20): dispatcher.emit_event(flag_event()) @@ -171,10 +171,10 @@ def test_emitter_does_not_drop_events_despite_slow_callback( assert dispatcher.dropped_events == 0 def test_emitter_eventually_processes_all_events_despite_slow_callback( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback: SlowCallback = SlowCallback(delay=0.05) - dispatcher: EventDispatcher = dispatcher_factory(callback) + dispatcher: _EventDispatcher = dispatcher_factory(callback) for _ in range(20): dispatcher.emit_event(flag_event()) @@ -184,7 +184,7 @@ def test_emitter_eventually_processes_all_events_despite_slow_callback( class TestWorkerLifecycle: def test_no_worker_runs_until_the_first_event( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): before = len(worker_threads()) dispatcher = dispatcher_factory(RecorderCallback()) @@ -196,7 +196,7 @@ def test_no_worker_runs_until_the_first_event( assert len(worker_threads()) == before + 1 def test_one_worker_serves_every_event( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -209,7 +209,7 @@ def test_one_worker_serves_every_event( assert len(worker_threads()) == before + 1 def test_the_callback_never_runs_on_the_emitting_thread( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -221,7 +221,7 @@ def test_the_callback_never_runs_on_the_emitting_thread( assert threading.current_thread().name != DISPATCHER_THREAD_NAME def test_the_worker_is_a_daemon_thread( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) before = set(worker_threads()) @@ -232,7 +232,7 @@ def test_the_worker_is_a_daemon_thread( assert worker.daemon is True def test_close_leaves_no_worker_behind( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) before = set(worker_threads()) @@ -244,7 +244,7 @@ def test_close_leaves_no_worker_behind( assert not worker.is_alive() def test_a_raising_callback_does_not_kill_the_worker( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RaisingCallback(times=1) dispatcher = dispatcher_factory(callback) @@ -256,7 +256,7 @@ def test_a_raising_callback_does_not_kill_the_worker( assert callback.feature_names == ["explodes", "survives"] def test_the_worker_outlives_a_callback_that_always_raises( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RaisingCallback() dispatcher = dispatcher_factory(callback) @@ -268,7 +268,7 @@ def test_the_worker_outlives_a_callback_that_always_raises( assert callback.feature_names == [str(index) for index in range(50)] def test_a_callback_exception_never_reaches_the_emitter( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RaisingCallback() dispatcher = dispatcher_factory(callback) @@ -289,14 +289,14 @@ class TestBackpressure: """ def test_no_events_are_dropped_before_anything_is_emitted( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) assert dispatcher.dropped_events == 0 def test_nothing_is_dropped_when_the_callback_keeps_up( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback, max_size=100) @@ -308,7 +308,7 @@ def test_nothing_is_dropped_when_the_callback_keeps_up( assert dispatcher.dropped_events == 0 def test_events_beyond_capacity_are_dropped( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=3) @@ -323,7 +323,7 @@ def test_events_beyond_capacity_are_dropped( assert dispatcher.dropped_events == 7 def test_dropped_events_are_never_delivered( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=2) @@ -342,7 +342,7 @@ def test_dropped_events_are_never_delivered( dispatcher.close(timeout=WAIT_TIMEOUT) def test_the_queue_takes_events_again_once_the_callback_drains( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=2) @@ -365,7 +365,7 @@ def test_the_queue_takes_events_again_once_the_callback_drains( assert dispatcher.dropped_events == 3 def test_a_max_size_of_zero_means_an_unbounded_queue( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=0) @@ -380,7 +380,7 @@ def test_a_max_size_of_zero_means_an_unbounded_queue( assert dispatcher.dropped_events == 0 def test_the_drop_count_survives_close( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=1) @@ -407,7 +407,7 @@ class TestReadyDeduplication: """ def test_ready_is_delivered_once_however_many_connectors_emit_it( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -421,7 +421,7 @@ def test_ready_is_delivered_once_however_many_connectors_emit_it( assert len(ready_deliveries(callback)) == 1 def test_ready_carried_on_an_unleash_event_is_deduplicated( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -434,7 +434,7 @@ def test_ready_carried_on_an_unleash_event_is_deduplicated( assert len(ready_deliveries(callback)) == 1 def test_a_ready_event_dropped_by_a_full_queue_is_still_delivered_later( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=1) @@ -462,8 +462,8 @@ def test_a_ready_event_dropped_by_a_full_queue_is_still_delivered_later( class TestFullQueueWarning: def _drop_events( - self, dispatcher_factory: Callable[..., EventDispatcher], count: int - ) -> EventDispatcher: + self, dispatcher_factory: Callable[..., _EventDispatcher], count: int + ) -> _EventDispatcher: """Pins the worker, fills a one-slot queue, then drops ``count`` events on the floor.""" callback = BlockingCallback() dispatcher = dispatcher_factory(callback, max_size=1) @@ -481,7 +481,7 @@ def _drop_events( def test_a_full_queue_is_warned_about( self, - dispatcher_factory: Callable[..., EventDispatcher], + dispatcher_factory: Callable[..., _EventDispatcher], caplog: pytest.LogCaptureFixture, ): caplog.set_level(logging.WARNING, logger="UnleashClient") @@ -494,7 +494,7 @@ def test_a_full_queue_is_warned_about( def test_a_full_queue_is_warned_about_only_once( # line 443 self, - dispatcher_factory: Callable[..., EventDispatcher], + dispatcher_factory: Callable[..., _EventDispatcher], caplog: pytest.LogCaptureFixture, ): caplog.set_level(logging.WARNING, logger="UnleashClient") @@ -511,7 +511,7 @@ def test_a_full_queue_is_warned_about_only_once( # line 443 class TestCallbackErrorLogging: def test_a_callback_exception_is_logged_with_its_message( self, - dispatcher_factory: Callable[..., EventDispatcher], + dispatcher_factory: Callable[..., _EventDispatcher], caplog: pytest.LogCaptureFixture, ): caplog.set_level(logging.ERROR, logger="UnleashClient") @@ -531,7 +531,7 @@ def test_a_callback_exception_is_logged_with_its_message( def test_every_failing_event_is_logged( self, - dispatcher_factory: Callable[..., EventDispatcher], + dispatcher_factory: Callable[..., _EventDispatcher], caplog: pytest.LogCaptureFixture, ): caplog.set_level(logging.ERROR, logger="UnleashClient") @@ -552,7 +552,7 @@ def test_every_failing_event_is_logged( def test_a_clean_run_logs_nothing( self, - dispatcher_factory: Callable[..., EventDispatcher], + dispatcher_factory: Callable[..., _EventDispatcher], caplog: pytest.LogCaptureFixture, ): caplog.set_level(logging.WARNING, logger="UnleashClient") @@ -570,7 +570,7 @@ def test_a_clean_run_logs_nothing( class TestClose: def test_close_drops_events_that_are_still_queued( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback) @@ -590,7 +590,7 @@ def test_close_drops_events_that_are_still_queued( assert callback.feature_names == ["in flight"] def test_close_is_idempotent( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -603,7 +603,7 @@ def test_close_is_idempotent( assert callback.call_count == 1 def test_a_second_close_returns_immediately( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) dispatcher.emit_event(flag_event()) @@ -617,7 +617,7 @@ def test_a_second_close_returns_immediately( assert elapsed < CLOSE_SLACK def test_close_honors_a_zero_timeout( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = BlockingCallback() dispatcher = dispatcher_factory(callback) @@ -632,7 +632,7 @@ def test_close_honors_a_zero_timeout( assert elapsed < CLOSE_SLACK def test_close_returns_promptly_when_nothing_was_ever_emitted( # line 594 - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) @@ -645,7 +645,7 @@ def test_close_returns_promptly_when_nothing_was_ever_emitted( # line 594 assert elapsed < CLOSE_SLACK def test_emit_after_close_is_a_no_op( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) @@ -664,7 +664,7 @@ def test_emit_after_close_is_a_no_op( assert callback.feature_names == ["before"] def test_emit_after_close_starts_no_worker( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): dispatcher = dispatcher_factory(RecorderCallback()) dispatcher.close(timeout=0) @@ -675,7 +675,7 @@ def test_emit_after_close_starts_no_worker( assert len(worker_threads()) == before def test_closing_from_inside_the_callback_does_not_blow_up( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = ClosingCallback(timeout=CLOSE_TIMEOUT) dispatcher = dispatcher_factory(callback) @@ -688,7 +688,7 @@ def test_closing_from_inside_the_callback_does_not_blow_up( assert callback.error is None def test_two_threads_closing_at_once_both_return( - self, dispatcher_factory: Callable[..., EventDispatcher] + self, dispatcher_factory: Callable[..., _EventDispatcher] ): callback = RecorderCallback() dispatcher = dispatcher_factory(callback) diff --git a/tests/unit_tests/test_feature_store.py b/tests/unit_tests/test_feature_store.py index b815fc91..7d8087ef 100644 --- a/tests/unit_tests/test_feature_store.py +++ b/tests/unit_tests/test_feature_store.py @@ -11,10 +11,11 @@ MOCK_FEATURE_RESPONSE_PROJECT, ) from tests.utilities.testing_constants import ETAG_VALUE +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient._feature_store import _FeatureStore from UnleashClient.cache import BaseCache from UnleashClient.constants import ETAG, FEATURES_URL -from UnleashClient.events import EventDispatcher, UnleashEventType +from UnleashClient.events import UnleashEventType FEATURES = json.dumps(MOCK_FEATURE_RESPONSE) OTHER_FEATURES = json.dumps(MOCK_FEATURE_RESPONSE_PROJECT) @@ -77,7 +78,7 @@ def test_load_from_cache_emits_ready(cache_empty, dispatcher, recorder): def test_load_from_cache_on_an_empty_cache_neither_raises_nor_emits( - cache_empty, dispatcher: EventDispatcher, recorder: EventRecorder + cache_empty, dispatcher: _EventDispatcher, recorder: EventRecorder ): engine = UnleashEngine() diff --git a/tests/utilities/event_callbacks.py b/tests/utilities/event_callbacks.py index 1660fa32..9349b655 100644 --- a/tests/utilities/event_callbacks.py +++ b/tests/utilities/event_callbacks.py @@ -1,5 +1,5 @@ """ -Callables and event builders for the EventDispatcher tests. +Callables and event builders for the _EventDispatcher tests. Every callback here records what it received, so a test can assert on delivery whichever behavior it picked. ``RecorderCallback.wait_for`` is how tests wait for the worker thread @@ -10,9 +10,9 @@ import uuid from typing import Callable, List, Optional +from UnleashClient._event_dispatcher import _EventDispatcher from UnleashClient.events import ( BaseEvent, - EventDispatcher, UnleashEvent, UnleashEventType, UnleashFetchedEvent, @@ -149,14 +149,14 @@ def __init__( ) -> None: super().__init__() self._build_event = build_event - self._dispatcher: Optional[EventDispatcher] = None + self._dispatcher: Optional[_EventDispatcher] = None self.entered = threading.Event() self.emitted = threading.Event() self._released = threading.Event() if not gated: self._released.set() - def bind(self, dispatcher: EventDispatcher) -> None: + def bind(self, dispatcher: _EventDispatcher) -> None: self._dispatcher = dispatcher def release(self) -> None: @@ -183,11 +183,11 @@ class ClosingCallback(RecorderCallback): def __init__(self, timeout: float = 0.2) -> None: super().__init__() self._timeout = timeout - self._dispatcher: Optional[EventDispatcher] = None + self._dispatcher: Optional[_EventDispatcher] = None self.returned = threading.Event() self.error: Optional[BaseException] = None - def bind(self, dispatcher: EventDispatcher) -> None: + def bind(self, dispatcher: _EventDispatcher) -> None: self._dispatcher = dispatcher def __call__(self, event: BaseEvent) -> None: