From 5e89c2a1fec1afbff40931fa7d927b4c2381ce95 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pere=20Pic=C3=B3?= Date: Tue, 29 Sep 2026 13:59:13 +0200 Subject: [PATCH] fix: prevent aiohttp from becoming mandatory to use the sync SDK to do so, I've moved the async metrics reporter to its own file. --- CHANGELOG.md | 1 + UnleashClient/_async_metrics.py | 96 ++++ UnleashClient/_metrics.py | 90 +--- UnleashClient/clients/async_unleash_client.py | 2 +- .../clients/test_async_unleash_client.py | 3 +- tests/unit_tests/test_async_metrics.py | 453 ++++++++++++++++++ tests/unit_tests/test_metrics.py | 425 +--------------- 7 files changed, 555 insertions(+), 515 deletions(-) create mode 100644 UnleashClient/_async_metrics.py create mode 100644 tests/unit_tests/test_async_metrics.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 82759e21..c2bf974f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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`. diff --git a/UnleashClient/_async_metrics.py b/UnleashClient/_async_metrics.py new file mode 100644 index 00000000..502dd48c --- /dev/null +++ b/UnleashClient/_async_metrics.py @@ -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() diff --git a/UnleashClient/_metrics.py b/UnleashClient/_metrics.py index 5c9cd720..5262216b 100644 --- a/UnleashClient/_metrics.py +++ b/UnleashClient/_metrics.py @@ -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 @@ -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() diff --git a/UnleashClient/clients/async_unleash_client.py b/UnleashClient/clients/async_unleash_client.py index 47abfc86..6fb38900 100644 --- a/UnleashClient/clients/async_unleash_client.py +++ b/UnleashClient/clients/async_unleash_client.py @@ -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 @@ -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 diff --git a/tests/unit_tests/clients/test_async_unleash_client.py b/tests/unit_tests/clients/test_async_unleash_client.py index 7b6cf190..2d8f37ad 100644 --- a/tests/unit_tests/clients/test_async_unleash_client.py +++ b/tests/unit_tests/clients/test_async_unleash_client.py @@ -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 ( diff --git a/tests/unit_tests/test_async_metrics.py b/tests/unit_tests/test_async_metrics.py new file mode 100644 index 00000000..d2526b39 --- /dev/null +++ b/tests/unit_tests/test_async_metrics.py @@ -0,0 +1,453 @@ +import asyncio +import json +from typing import Callable, List, Optional, TypedDict + +import pytest +import pytest_asyncio +from pytest import mark +from yggdrasil_engine.engine import UnleashEngine + +from tests.utilities.fake_unleash_server import FakeUnleash +from UnleashClient._async_metrics import _AsyncMetricsReporter +from UnleashClient._async_scheduler import _AsyncJob, _AsyncJobFn, _AsyncScheduler +from UnleashClient._async_transport import _AsyncTransport +from UnleashClient._headers import _HeaderFactory +from UnleashClient.config import UnleashConfig +from UnleashClient.constants import CLIENT_SPEC_VERSION, METRICS_URL +from UnleashClient.impact_metrics import ImpactMetrics + +APP_NAME = "pytest" + +# The fake server mounts Unleash under /api, so the async transport builds the URL +# shape it builds against a deployment. +API_PREFIX = "/api" +METRICS_PATH = API_PREFIX + METRICS_URL + +# Any flag will do; the bucket only has to be non-empty for a send to happen. +COUNTED_FLAG = "something-to-make-sure-metrics-get-sent" + +# Long enough that the scheduled flush never runs during a test, so a recorded request +# can only have come from the call the test made. +NEVER = 3600 + + +class Registration(TypedDict): + interval_seconds: float + jitter_seconds: Optional[float] + fn: _AsyncJobFn + job: _AsyncJob + + +class RecordingAsyncScheduler(_AsyncScheduler): + """Captures what reached every() and cancel_and_wait().""" + + def __init__(self) -> None: + super().__init__() + self.registered: List[Registration] = [] + self.cancelled: List[_AsyncJob] = [] + + def every(self, interval_seconds, jitter_seconds, fn, kwargs=None): + job = super().every(interval_seconds, jitter_seconds, fn, kwargs) + self.registered.append( + { + "interval_seconds": interval_seconds, + "jitter_seconds": jitter_seconds, + "fn": fn, + "job": job, + } + ) + return job + + async def cancel_and_wait(self, job): + self.cancelled.append(job) + await super().cancel_and_wait(job) + + +class SilentImpactMetrics: + """ + Impact metrics that never yield anything: what a client that records none looks + like, and what a failed collection degrades to. + """ + + def __init__(self): + self.restored = [] + + def collect(self): + return None + + def restore(self, metrics): + self.restored.append(metrics) + + +@pytest_asyncio.fixture +async def server(): + """A real Unleash server on an ephemeral port, stopped on teardown.""" + fake = FakeUnleash() + await fake.start(API_PREFIX) + try: + yield fake + finally: + await fake.close() + + +def metrics_body(server: FakeUnleash, index: int = 0) -> dict: + return json.loads(server.calls("POST", METRICS_PATH)[index].body) + + +async def until_requested(server: FakeUnleash, count: int = 1) -> None: + while len(server.calls("POST", METRICS_PATH)) < count: + await asyncio.sleep(0.01) + + +class TestAsyncMetricsReporter: + @pytest_asyncio.fixture + async def build_reporter(self, server: FakeUnleash): + """ + Factory pointed at the fake server. Keyword arguments override the defaults on the + config. + + Every reporter it builds gets its own RecordingScheduler, and both are shut down on + teardown: a job left running outlives the test, and an unclosed ClientSession is + reported by aiohttp on garbage collection. + """ + built = [] + + def _build_reporter(impact_metrics=None, **kwargs) -> _AsyncMetricsReporter: + defaults = { + "metrics_interval": NEVER, + } + defaults.update(kwargs) + config = UnleashConfig(server.base_url, APP_NAME, **defaults) + engine = UnleashEngine() + reporter = _AsyncMetricsReporter( + config=config, + transport=_AsyncTransport(config, _HeaderFactory(config)), + scheduler=RecordingAsyncScheduler(), + engine=engine, + impact_metrics=( + impact_metrics + if impact_metrics is not None + else ImpactMetrics( + engine, config.app_name, config.impact_metrics_environment + ) + ), + ) + built.append(reporter) + return reporter + + try: + yield _build_reporter + finally: + for reporter in built: + await reporter._scheduler.shutdown() + await reporter._transport.aclose() + + @pytest_asyncio.fixture + async def reporter( + self, + build_reporter: Callable[..., _AsyncMetricsReporter], + ) -> _AsyncMetricsReporter: + """The reporter the tests that need no config override share.""" + return build_reporter() + + # flush + + @mark.asyncio + async def test_flush_sends_nothing_when_nothing_was_recorded( + self, server, reporter + ): + server.on("POST", METRICS_PATH, status=202, payload={}) + + await reporter.flush() + + assert len(server.calls("POST", METRICS_PATH)) == 0 + + @mark.asyncio + async def test_flush_sends_the_bucket_it_collected(self, server, reporter): + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + assert metrics_body(server)["bucket"]["toggles"][COUNTED_FLAG]["yes"] == 1 + + @mark.asyncio + async def test_flush_identifies_the_client(self, server, build_reporter): + reporter = build_reporter(instance_id="123") + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + body = metrics_body(server) + assert body["appName"] == APP_NAME + assert body["instanceId"] == "123" + assert body["connectionId"] == reporter._config.connection_id + + @mark.asyncio + async def test_flush_sends_the_platform_metadata(self, server, reporter): + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + body = metrics_body(server) + assert body["yggdrasilVersion"] is not None + assert body["specVersion"] == CLIENT_SPEC_VERSION + assert body["platformName"] is not None + assert body["platformVersion"] is not None + + @mark.asyncio + async def test_flush_includes_sdk_flavor_when_set(self, server, build_reporter): + reporter = build_reporter( + sdk_flavor="unleash-openfeature-python-provider", sdk_flavor_version="1.2.3" + ) + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + body = metrics_body(server) + assert body["sdkFlavor"] == "unleash-openfeature-python-provider" + assert body["sdkFlavorVersion"] == "1.2.3" + + @mark.asyncio + async def test_flush_omits_sdk_flavor_when_unset(self, server, reporter): + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + body = metrics_body(server) + assert "sdkFlavor" not in body + assert "sdkFlavorVersion" not in body + + @mark.asyncio + async def test_the_config_is_read_on_every_flush(self, server, reporter): + # UnleashClient.unleash_app_name has a setter, so a client can change it after the + # reporter was constructed. + server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) + + reporter._engine.count_toggle(COUNTED_FLAG, True) + await reporter.flush() + + reporter._config.app_name = "renamed" + reporter._engine.count_toggle(COUNTED_FLAG, True) + await reporter.flush() + + assert metrics_body(server, 0)["appName"] == APP_NAME + assert metrics_body(server, 1)["appName"] == "renamed" + + @mark.asyncio + async def test_the_flush_goes_through_the_async_transport(self, reporter): + # The flush runs on the client's loop, so a blocking transport would hold it up for + # the length of every POST. + assert isinstance(reporter._transport, _AsyncTransport) + assert asyncio.iscoroutinefunction(reporter._transport.send_metrics) + assert asyncio.iscoroutinefunction(reporter.flush) + + # start and stop + + @mark.asyncio + async def test_start_registers_the_flush_with_the_metrics_interval_and_jitter( + self, + build_reporter, + ): + reporter = build_reporter(metrics_interval=30, metrics_jitter=10) + + reporter.start() + + (call,) = reporter._scheduler.registered + assert call["interval_seconds"] == 30 + assert call["jitter_seconds"] == 10 + assert call["fn"] == reporter.flush + + @mark.asyncio + async def test_start_coerces_the_metrics_interval(self, build_reporter): + # A client can be built with metrics_interval="30": the constructor does not + # validate types. + reporter = build_reporter() + reporter._config.metrics_interval = "30" + + reporter.start() + + (call,) = reporter._scheduler.registered + assert call["interval_seconds"] == 30 + + @mark.asyncio + async def test_stop_flushes_what_is_left_and_cancels_the_job( + self, server, build_reporter + ): + # A short-lived client can be destroyed before its first interval elapses, so the + # bucket has to go out on the way down. + reporter = build_reporter() + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter.start() + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.stop() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + (call,) = reporter._scheduler.registered + assert reporter._scheduler.cancelled == [call["job"]] + + @mark.asyncio + async def test_stop_sends_nothing_when_start_was_never_called( + self, server, reporter + ): + # Metrics disabled, or destroy() on a client that was never initialized. + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.stop() + + assert len(server.calls("POST", METRICS_PATH)) == 0 + + @mark.asyncio + async def test_stop_is_idempotent(self, server, build_reporter): + reporter = build_reporter() + server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) + reporter.start() + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.stop() + await reporter.stop() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + + @mark.asyncio + async def test_constructing_the_reporter_registers_no_job_and_opens_no_session( + self, + build_reporter, + ): + # __init__ is synchronous and runs with no loop: the client's constructor has to + # work before anything is awaited. + reporter = build_reporter() + + assert reporter._scheduler.registered == [] + assert reporter._transport._session is None + + # impact metrics + + @mark.asyncio + async def test_impact_metrics_go_out_with_the_bucket(self, server, reporter): + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._impact_metrics.define_counter("purchases", "Number of purchases") + reporter._impact_metrics.increment_counter("purchases", 1) + + await reporter.flush() + + assert metrics_body(server)["impactMetrics"][0]["name"] == "purchases" + + @mark.asyncio + async def test_impact_metrics_alone_are_enough_to_trigger_a_send( + self, server, reporter + ): + # Nothing was evaluated, so there is no toggle bucket at all. + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._impact_metrics.define_counter("purchases", "Number of purchases") + reporter._impact_metrics.increment_counter("purchases", 1) + + await reporter.flush() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + assert metrics_body(server)["bucket"] is None + + @mark.asyncio + async def test_impact_metrics_are_restored_when_the_send_fails( + self, server, reporter + ): + server.on("POST", METRICS_PATH, status=500, payload={}) + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._impact_metrics.define_counter("my_counter", "Test counter") + reporter._impact_metrics.increment_counter("my_counter", 5) + + await reporter.flush() + await reporter.flush() + + resent = metrics_body(server, 1)["impactMetrics"][0] + assert resent["name"] == "my_counter" + assert resent["samples"][0]["value"] == 5 + + @mark.asyncio + async def test_impact_metrics_are_restored_when_the_send_is_cancelled( + self, server, reporter + ): + server.on("POST", METRICS_PATH, hang=True) + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._impact_metrics.define_counter("my_counter", "Test counter") + reporter._impact_metrics.increment_counter("my_counter", 5) + flush = asyncio.create_task(reporter.flush()) + await until_requested(server) + + flush.cancel() + with pytest.raises(asyncio.CancelledError): + await flush + await reporter.flush() + + resent = metrics_body(server, 1)["impactMetrics"][0] + assert resent["name"] == "my_counter" + assert resent["samples"][0]["value"] == 5 + + @mark.asyncio + async def test_stop_during_a_send_resends_its_impact_metrics( + self, server, build_reporter + ): + reporter = build_reporter(metrics_interval=1) + server.on("POST", METRICS_PATH, hang=True) + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._impact_metrics.define_counter("my_counter", "Test counter") + reporter._impact_metrics.increment_counter("my_counter", 5) + reporter.start() + reporter._scheduler.start() + await until_requested(server) + + await asyncio.wait_for(reporter.stop(), timeout=5) + + assert len(server.calls("POST", METRICS_PATH)) == 2 + assert metrics_body(server, 1)["impactMetrics"][0]["samples"][0]["value"] == 5 + + @mark.asyncio + async def test_impact_metrics_are_not_restored_when_the_send_succeeds( + self, server, reporter + ): + server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) + reporter._impact_metrics.define_counter("my_counter", "Test counter") + reporter._impact_metrics.increment_counter("my_counter", 5) + + await reporter.flush() + await reporter.flush() + + # The engine keeps the counter definition, so a second send still describes it, + # back at zero, rather than replaying the 5 the server already took. + assert metrics_body(server, 1)["impactMetrics"][0]["samples"][0]["value"] == 0 + + @mark.asyncio + async def test_nothing_is_restored_when_there_were_no_impact_metrics( + self, server, build_reporter + ): + impact_metrics = SilentImpactMetrics() + reporter = build_reporter(impact_metrics=impact_metrics) + server.on("POST", METRICS_PATH, status=500, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + assert impact_metrics.restored == [] + + @mark.asyncio + async def test_the_bucket_still_goes_out_when_impact_collection_yields_nothing( + self, server, build_reporter + ): + # ImpactMetrics.collect() returns None on a broken engine; that must not take the + # toggle metrics down with it. + reporter = build_reporter(impact_metrics=SilentImpactMetrics()) + server.on("POST", METRICS_PATH, status=202, payload={}) + reporter._engine.count_toggle(COUNTED_FLAG, True) + + await reporter.flush() + + assert len(server.calls("POST", METRICS_PATH)) == 1 + assert "impactMetrics" not in metrics_body(server) diff --git a/tests/unit_tests/test_metrics.py b/tests/unit_tests/test_metrics.py index c238580b..1c132321 100644 --- a/tests/unit_tests/test_metrics.py +++ b/tests/unit_tests/test_metrics.py @@ -1,19 +1,11 @@ -import asyncio import json -from typing import Callable, List, Optional, TypedDict -import pytest -import pytest_asyncio import responses from apscheduler.schedulers.background import BackgroundScheduler -from pytest import mark from yggdrasil_engine.engine import UnleashEngine -from tests.utilities.fake_unleash_server import FakeUnleash -from UnleashClient._async_scheduler import _AsyncJob, _AsyncJobFn, _AsyncScheduler -from UnleashClient._async_transport import _AsyncTransport from UnleashClient._headers import _HeaderFactory -from UnleashClient._metrics import _AsyncMetricsReporter, _MetricsReporter +from UnleashClient._metrics import _MetricsReporter from UnleashClient._scheduler import _Scheduler from UnleashClient._transport import _Transport from UnleashClient.config import UnleashConfig @@ -25,18 +17,9 @@ FULL_METRICS_URL = URL + METRICS_URL -# The fake server mounts Unleash under /api, so the async transport builds the URL -# shape it builds against a deployment. -API_PREFIX = "/api" -METRICS_PATH = API_PREFIX + METRICS_URL - # Any flag will do; the bucket only has to be non-empty for a send to happen. COUNTED_FLAG = "something-to-make-sure-metrics-get-sent" -# Long enough that the scheduled flush never runs during a test, so a recorded request -# can only have come from the call the test made. -NEVER = 3600 - class RecordingScheduler(BackgroundScheduler): """Captures what reached add_job.""" @@ -66,38 +49,6 @@ def remove_all_jobs(self, *args, **kwargs): pass -class Registration(TypedDict): - interval_seconds: float - jitter_seconds: Optional[float] - fn: _AsyncJobFn - job: _AsyncJob - - -class RecordingAsyncScheduler(_AsyncScheduler): - """Captures what reached every() and cancel_and_wait().""" - - def __init__(self) -> None: - super().__init__() - self.registered: List[Registration] = [] - self.cancelled: List[_AsyncJob] = [] - - def every(self, interval_seconds, jitter_seconds, fn, kwargs=None): - job = super().every(interval_seconds, jitter_seconds, fn, kwargs) - self.registered.append( - { - "interval_seconds": interval_seconds, - "jitter_seconds": jitter_seconds, - "fn": fn, - "job": job, - } - ) - return job - - async def cancel_and_wait(self, job): - self.cancelled.append(job) - await super().cancel_and_wait(job) - - class SilentImpactMetrics: """ Impact metrics that never yield anything: what a client that records none looks @@ -134,26 +85,6 @@ def build_sync_reporter( ) -@pytest_asyncio.fixture -async def server(): - """A real Unleash server on an ephemeral port, stopped on teardown.""" - fake = FakeUnleash() - await fake.start(API_PREFIX) - try: - yield fake - finally: - await fake.close() - - -def metrics_body(server: FakeUnleash, index: int = 0) -> dict: - return json.loads(server.calls("POST", METRICS_PATH)[index].body) - - -async def until_requested(server: FakeUnleash, count: int = 1) -> None: - while len(server.calls("POST", METRICS_PATH)) < count: - await asyncio.sleep(0.01) - - class TestMetricsReporter: # flush @@ -412,357 +343,3 @@ def test_the_bucket_still_goes_out_when_impact_collection_yields_nothing( assert len(responses.calls) == 1 assert "impactMetrics" not in json.loads(responses.calls[0].request.body) - - -class TestAsyncMetricsReporter: - @pytest_asyncio.fixture - async def build_reporter(self, server: FakeUnleash): - """ - Factory pointed at the fake server. Keyword arguments override the defaults on the - config. - - Every reporter it builds gets its own RecordingScheduler, and both are shut down on - teardown: a job left running outlives the test, and an unclosed ClientSession is - reported by aiohttp on garbage collection. - """ - built = [] - - def _build_reporter(impact_metrics=None, **kwargs) -> _AsyncMetricsReporter: - defaults = { - "metrics_interval": NEVER, - } - defaults.update(kwargs) - config = UnleashConfig(server.base_url, APP_NAME, **defaults) - engine = UnleashEngine() - reporter = _AsyncMetricsReporter( - config=config, - transport=_AsyncTransport(config, _HeaderFactory(config)), - scheduler=RecordingAsyncScheduler(), - engine=engine, - impact_metrics=( - impact_metrics - if impact_metrics is not None - else ImpactMetrics( - engine, config.app_name, config.impact_metrics_environment - ) - ), - ) - built.append(reporter) - return reporter - - try: - yield _build_reporter - finally: - for reporter in built: - await reporter._scheduler.shutdown() - await reporter._transport.aclose() - - @pytest_asyncio.fixture - async def reporter( - self, - build_reporter: Callable[..., _AsyncMetricsReporter], - ) -> _AsyncMetricsReporter: - """The reporter the tests that need no config override share.""" - return build_reporter() - - # flush - - @mark.asyncio - async def test_flush_sends_nothing_when_nothing_was_recorded( - self, server, reporter - ): - server.on("POST", METRICS_PATH, status=202, payload={}) - - await reporter.flush() - - assert len(server.calls("POST", METRICS_PATH)) == 0 - - @mark.asyncio - async def test_flush_sends_the_bucket_it_collected(self, server, reporter): - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - assert metrics_body(server)["bucket"]["toggles"][COUNTED_FLAG]["yes"] == 1 - - @mark.asyncio - async def test_flush_identifies_the_client(self, server, build_reporter): - reporter = build_reporter(instance_id="123") - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - body = metrics_body(server) - assert body["appName"] == APP_NAME - assert body["instanceId"] == "123" - assert body["connectionId"] == reporter._config.connection_id - - @mark.asyncio - async def test_flush_sends_the_platform_metadata(self, server, reporter): - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - body = metrics_body(server) - assert body["yggdrasilVersion"] is not None - assert body["specVersion"] == CLIENT_SPEC_VERSION - assert body["platformName"] is not None - assert body["platformVersion"] is not None - - @mark.asyncio - async def test_flush_includes_sdk_flavor_when_set(self, server, build_reporter): - reporter = build_reporter( - sdk_flavor="unleash-openfeature-python-provider", sdk_flavor_version="1.2.3" - ) - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - body = metrics_body(server) - assert body["sdkFlavor"] == "unleash-openfeature-python-provider" - assert body["sdkFlavorVersion"] == "1.2.3" - - @mark.asyncio - async def test_flush_omits_sdk_flavor_when_unset(self, server, reporter): - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - body = metrics_body(server) - assert "sdkFlavor" not in body - assert "sdkFlavorVersion" not in body - - @mark.asyncio - async def test_the_config_is_read_on_every_flush(self, server, reporter): - # UnleashClient.unleash_app_name has a setter, so a client can change it after the - # reporter was constructed. - server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) - - reporter._engine.count_toggle(COUNTED_FLAG, True) - await reporter.flush() - - reporter._config.app_name = "renamed" - reporter._engine.count_toggle(COUNTED_FLAG, True) - await reporter.flush() - - assert metrics_body(server, 0)["appName"] == APP_NAME - assert metrics_body(server, 1)["appName"] == "renamed" - - @mark.asyncio - async def test_the_flush_goes_through_the_async_transport(self, reporter): - # The flush runs on the client's loop, so a blocking transport would hold it up for - # the length of every POST. - assert isinstance(reporter._transport, _AsyncTransport) - assert asyncio.iscoroutinefunction(reporter._transport.send_metrics) - assert asyncio.iscoroutinefunction(reporter.flush) - - # start and stop - - @mark.asyncio - async def test_start_registers_the_flush_with_the_metrics_interval_and_jitter( - self, - build_reporter, - ): - reporter = build_reporter(metrics_interval=30, metrics_jitter=10) - - reporter.start() - - (call,) = reporter._scheduler.registered - assert call["interval_seconds"] == 30 - assert call["jitter_seconds"] == 10 - assert call["fn"] == reporter.flush - - @mark.asyncio - async def test_start_coerces_the_metrics_interval(self, build_reporter): - # A client can be built with metrics_interval="30": the constructor does not - # validate types. - reporter = build_reporter() - reporter._config.metrics_interval = "30" - - reporter.start() - - (call,) = reporter._scheduler.registered - assert call["interval_seconds"] == 30 - - @mark.asyncio - async def test_stop_flushes_what_is_left_and_cancels_the_job( - self, server, build_reporter - ): - # A short-lived client can be destroyed before its first interval elapses, so the - # bucket has to go out on the way down. - reporter = build_reporter() - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter.start() - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.stop() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - (call,) = reporter._scheduler.registered - assert reporter._scheduler.cancelled == [call["job"]] - - @mark.asyncio - async def test_stop_sends_nothing_when_start_was_never_called( - self, server, reporter - ): - # Metrics disabled, or destroy() on a client that was never initialized. - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.stop() - - assert len(server.calls("POST", METRICS_PATH)) == 0 - - @mark.asyncio - async def test_stop_is_idempotent(self, server, build_reporter): - reporter = build_reporter() - server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) - reporter.start() - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.stop() - await reporter.stop() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - - @mark.asyncio - async def test_constructing_the_reporter_registers_no_job_and_opens_no_session( - self, - build_reporter, - ): - # __init__ is synchronous and runs with no loop: the client's constructor has to - # work before anything is awaited. - reporter = build_reporter() - - assert reporter._scheduler.registered == [] - assert reporter._transport._session is None - - # impact metrics - - @mark.asyncio - async def test_impact_metrics_go_out_with_the_bucket(self, server, reporter): - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._impact_metrics.define_counter("purchases", "Number of purchases") - reporter._impact_metrics.increment_counter("purchases", 1) - - await reporter.flush() - - assert metrics_body(server)["impactMetrics"][0]["name"] == "purchases" - - @mark.asyncio - async def test_impact_metrics_alone_are_enough_to_trigger_a_send( - self, server, reporter - ): - # Nothing was evaluated, so there is no toggle bucket at all. - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._impact_metrics.define_counter("purchases", "Number of purchases") - reporter._impact_metrics.increment_counter("purchases", 1) - - await reporter.flush() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - assert metrics_body(server)["bucket"] is None - - @mark.asyncio - async def test_impact_metrics_are_restored_when_the_send_fails( - self, server, reporter - ): - server.on("POST", METRICS_PATH, status=500, payload={}) - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._impact_metrics.define_counter("my_counter", "Test counter") - reporter._impact_metrics.increment_counter("my_counter", 5) - - await reporter.flush() - await reporter.flush() - - resent = metrics_body(server, 1)["impactMetrics"][0] - assert resent["name"] == "my_counter" - assert resent["samples"][0]["value"] == 5 - - @mark.asyncio - async def test_impact_metrics_are_restored_when_the_send_is_cancelled( - self, server, reporter - ): - server.on("POST", METRICS_PATH, hang=True) - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._impact_metrics.define_counter("my_counter", "Test counter") - reporter._impact_metrics.increment_counter("my_counter", 5) - flush = asyncio.create_task(reporter.flush()) - await until_requested(server) - - flush.cancel() - with pytest.raises(asyncio.CancelledError): - await flush - await reporter.flush() - - resent = metrics_body(server, 1)["impactMetrics"][0] - assert resent["name"] == "my_counter" - assert resent["samples"][0]["value"] == 5 - - @mark.asyncio - async def test_stop_during_a_send_resends_its_impact_metrics( - self, server, build_reporter - ): - reporter = build_reporter(metrics_interval=1) - server.on("POST", METRICS_PATH, hang=True) - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._impact_metrics.define_counter("my_counter", "Test counter") - reporter._impact_metrics.increment_counter("my_counter", 5) - reporter.start() - reporter._scheduler.start() - await until_requested(server) - - await asyncio.wait_for(reporter.stop(), timeout=5) - - assert len(server.calls("POST", METRICS_PATH)) == 2 - assert metrics_body(server, 1)["impactMetrics"][0]["samples"][0]["value"] == 5 - - @mark.asyncio - async def test_impact_metrics_are_not_restored_when_the_send_succeeds( - self, server, reporter - ): - server.on("POST", METRICS_PATH, status=202, payload={}, repeat=True) - reporter._impact_metrics.define_counter("my_counter", "Test counter") - reporter._impact_metrics.increment_counter("my_counter", 5) - - await reporter.flush() - await reporter.flush() - - # The engine keeps the counter definition, so a second send still describes it, - # back at zero, rather than replaying the 5 the server already took. - assert metrics_body(server, 1)["impactMetrics"][0]["samples"][0]["value"] == 0 - - @mark.asyncio - async def test_nothing_is_restored_when_there_were_no_impact_metrics( - self, server, build_reporter - ): - impact_metrics = SilentImpactMetrics() - reporter = build_reporter(impact_metrics=impact_metrics) - server.on("POST", METRICS_PATH, status=500, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - assert impact_metrics.restored == [] - - @mark.asyncio - async def test_the_bucket_still_goes_out_when_impact_collection_yields_nothing( - self, server, build_reporter - ): - # ImpactMetrics.collect() returns None on a broken engine; that must not take the - # toggle metrics down with it. - reporter = build_reporter(impact_metrics=SilentImpactMetrics()) - server.on("POST", METRICS_PATH, status=202, payload={}) - reporter._engine.count_toggle(COUNTED_FLAG, True) - - await reporter.flush() - - assert len(server.calls("POST", METRICS_PATH)) == 1 - assert "impactMetrics" not in metrics_body(server)