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 @@ -4,6 +4,7 @@
> went on. It needs to be consolidated into what will eventually become the
> final changelog + documentation changes (if any).

* (Bugfix): `UnleashClient` can be imported again without the optional `aiohttp` dependency. Importing the package used to fail with an `ImportError` asking for `UnleashClient[async]`, even for code that never touches the asynchronous client. Only the asynchronous client requires `aiohttp` now, and importing it without `aiohttp` still raises that `ImportError`.
* (Minor): The in-progress asynchronous client now exposes `feature_definitions()`. It returns the same dict as `UnleashClient`, keyed by feature name with each toggle's `type` and `project`. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and without any state it returns an empty dict. Every method on the client is now implemented. The client is still unexported. Nothing changes for code using `UnleashClient`.
* (Minor): The in-progress asynchronous client can now resolve variants. `get_variant()` returns the same variant dict as `UnleashClient` and emits the same impression events. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and a toggle the client does not know resolves to the disabled variant. When an initialized client is asked for a toggle it does not know, it logs at `verbose_log_level` that the client does not know the toggle. `feature_definitions()` still raises `NotImplementedError`. The client is still unexported. Nothing changes for code using `UnleashClient`.
* (Minor): The in-progress asynchronous client can now evaluate feature toggles. `is_enabled()` resolves a toggle against the feature state the client holds, with the same results, `fallback_function` handling and impression events as `UnleashClient`. It is a plain method, not a coroutine, and it does not wait for the server: before the first fetch it answers from the cached state, and a toggle the client does not know resolves to the fallback's answer, or to false without one. `get_variant()` and `feature_definitions()` still raise `NotImplementedError`. The client is still unexported. Nothing changes for code using `UnleashClient`.
Expand Down
96 changes: 96 additions & 0 deletions UnleashClient/_async_metrics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
"""Metrics reporting, for the async Unleash client."""

from typing import Optional

from yggdrasil_engine.engine import UnleashEngine

from UnleashClient._async_scheduler import _AsyncJob, _AsyncScheduler
from UnleashClient._async_transport import _AsyncTransport
from UnleashClient._payloads import _build_metrics_payload
from UnleashClient.config import UnleashConfig
from UnleashClient.impact_metrics import ImpactMetrics
from UnleashClient.utils import LOGGER


class _AsyncMetricsReporter:
"""
Sends feature and impact metrics to Unleash on a recurring interval.

:meth:`start` must be called, and :meth:`stop` awaited, from the event loop the client
runs on, and the loop must stay open for as long as metrics are being reported.

Example::

reporter = _AsyncMetricsReporter(
config=config,
transport=transport,
scheduler=scheduler,
engine=engine,
impact_metrics=impact_metrics,
)
reporter.start()

await reporter.flush()

await reporter.stop()
"""

def __init__(
self,
config: UnleashConfig,
transport: _AsyncTransport,
scheduler: _AsyncScheduler,
engine: UnleashEngine,
impact_metrics: ImpactMetrics,
) -> None:
self._config: UnleashConfig = config
self._transport: _AsyncTransport = transport
self._scheduler: _AsyncScheduler = scheduler
self._engine: UnleashEngine = engine
self._impact_metrics: ImpactMetrics = impact_metrics
self._job: Optional[_AsyncJob] = None

def start(self) -> None:
"""Schedules a send every ``metrics_interval`` seconds, with ``metrics_jitter`` of jitter."""
self._job = self._scheduler.every(
interval_seconds=int(self._config.metrics_interval),
jitter_seconds=self._config.metrics_jitter,
fn=self.flush,
)

async def flush(self) -> None:
"""
Sends one bucket of feature and impact metrics.

Sends nothing when neither has anything to report. When a send fails or is
cancelled, its impact metrics are restored so the next send carries them.
"""
bucket = self._engine.get_metrics()
impact_metrics = self._impact_metrics.collect()

if not (bucket or impact_metrics):
LOGGER.debug("No feature flags with metrics, skipping metrics submission.")
return

payload = _build_metrics_payload(self._config, bucket, impact_metrics)
sent = False
try:
sent = await self._transport.send_metrics(payload)
finally:
if not sent and impact_metrics:
self._impact_metrics.restore(impact_metrics)

async def stop(self) -> None:
"""
Stops the recurring send and flushes whatever is left.

Does nothing when :meth:`start` was never called. A send still in flight is
cancelled: its impact metrics go out with the final flush, and its feature
metrics are lost.
"""
if self._job is None:
return

job, self._job = self._job, None
await self._scheduler.cancel_and_wait(job)
await self.flush()
90 changes: 1 addition & 89 deletions UnleashClient/_metrics.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,7 @@
"""Metrics reporting, for the sync and async Unleash clients."""

from typing import Optional
"""Metrics reporting, for the sync Unleash client."""

from yggdrasil_engine.engine import UnleashEngine

from UnleashClient._async_scheduler import _AsyncJob, _AsyncScheduler
from UnleashClient._async_transport import _AsyncTransport
from UnleashClient._payloads import _build_metrics_payload
from UnleashClient._scheduler import _ScheduledJob, _Scheduler
from UnleashClient._transport import _Transport
Expand Down Expand Up @@ -98,87 +94,3 @@ def stop(self) -> None:
self.flush()
self._scheduler.cancel(self._job)
self._job = None


class _AsyncMetricsReporter:
"""
Sends feature and impact metrics to Unleash on a recurring interval.

:meth:`start` must be called, and :meth:`stop` awaited, from the event loop the client
runs on, and the loop must stay open for as long as metrics are being reported.

Example::

reporter = _AsyncMetricsReporter(
config=config,
transport=transport,
scheduler=scheduler,
engine=engine,
impact_metrics=impact_metrics,
)
reporter.start()

await reporter.flush()

await reporter.stop()
"""

def __init__(
self,
config: UnleashConfig,
transport: _AsyncTransport,
scheduler: _AsyncScheduler,
engine: UnleashEngine,
impact_metrics: ImpactMetrics,
) -> None:
self._config: UnleashConfig = config
self._transport: _AsyncTransport = transport
self._scheduler: _AsyncScheduler = scheduler
self._engine: UnleashEngine = engine
self._impact_metrics: ImpactMetrics = impact_metrics
self._job: Optional[_AsyncJob] = None

def start(self) -> None:
"""Schedules a send every ``metrics_interval`` seconds, with ``metrics_jitter`` of jitter."""
self._job = self._scheduler.every(
interval_seconds=int(self._config.metrics_interval),
jitter_seconds=self._config.metrics_jitter,
fn=self.flush,
)

async def flush(self) -> None:
"""
Sends one bucket of feature and impact metrics.

Sends nothing when neither has anything to report. When a send fails or is
cancelled, its impact metrics are restored so the next send carries them.
"""
bucket = self._engine.get_metrics()
impact_metrics = self._impact_metrics.collect()

if not (bucket or impact_metrics):
LOGGER.debug("No feature flags with metrics, skipping metrics submission.")
return

payload = _build_metrics_payload(self._config, bucket, impact_metrics)
sent = False
try:
sent = await self._transport.send_metrics(payload)
finally:
if not sent and impact_metrics:
self._impact_metrics.restore(impact_metrics)

async def stop(self) -> None:
"""
Stops the recurring send and flushes whatever is left.

Does nothing when :meth:`start` was never called. A send still in flight is
cancelled: its impact metrics go out with the final flush, and its feature
metrics are lost.
"""
if self._job is None:
return

job, self._job = self._job, None
await self._scheduler.cancel_and_wait(job)
await self.flush()
2 changes: 1 addition & 1 deletion UnleashClient/clients/async_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._async_metrics import _AsyncMetricsReporter
from UnleashClient._async_scheduler import _AsyncScheduler
from UnleashClient._async_transport import _AsyncTransport
from UnleashClient._context import _ContextEnricher
Expand All @@ -18,7 +19,6 @@
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._payloads import _build_register_payload
from UnleashClient.cache import BaseCache, FileCache
from UnleashClient.clients.unleash_client import _RunState
Expand Down
3 changes: 2 additions & 1 deletion tests/unit_tests/clients/test_async_unleash_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,8 @@
)
from tests.utilities.testing_constants import APP_NAME, URL
from UnleashClient import INSTANCES, UnleashClient
from UnleashClient._metrics import _AsyncMetricsReporter, _MetricsReporter
from UnleashClient._async_metrics import _AsyncMetricsReporter
from UnleashClient._metrics import _MetricsReporter
from UnleashClient.cache import FileCache
from UnleashClient.clients.async_unleash_client import AsyncUnleashClient
from UnleashClient.constants import (
Expand Down
Loading
Loading