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 @@ -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.
Expand Down
6 changes: 3 additions & 3 deletions UnleashClient/_evaluator.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)
Expand All @@ -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).
Expand All @@ -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(
Expand Down
135 changes: 135 additions & 0 deletions UnleashClient/_event_dispatcher.py
Original file line number Diff line number Diff line change
@@ -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
4 changes: 2 additions & 2 deletions UnleashClient/_feature_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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).
Expand Down
7 changes: 4 additions & 3 deletions UnleashClient/clients/async_unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,15 @@
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
from UnleashClient._metrics import _AsyncMetricsReporter
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

Expand Down Expand Up @@ -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(
Expand Down
8 changes: 4 additions & 4 deletions UnleashClient/clients/unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -39,7 +40,6 @@
)
from UnleashClient.events import (
BaseEvent,
EventDispatcher,
UnleashEventType,
UnleashReadyEvent,
)
Expand All @@ -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.
"""
Expand Down Expand Up @@ -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()
Expand Down
2 changes: 1 addition & 1 deletion UnleashClient/connectors/bootstrap_connector.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
117 changes: 1 addition & 116 deletions UnleashClient/events.py
Original file line number Diff line number Diff line change
@@ -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):
"""
Expand Down Expand Up @@ -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
Loading
Loading