From 235fb91c1697d0be2f3df11f5279134af68e12ce Mon Sep 17 00:00:00 2001 From: Karl Date: Sun, 13 Sep 2026 14:46:55 -0400 Subject: [PATCH 1/2] fix(telemetry): deliver off the caller's thread and stop the thread on close (#105) The opt-in telemetry collector is an advertised public option -- both clients take `enable_telemetry` and `configure_telemetry` is exported -- so it is fixed rather than deleted. - track_request no longer calls _flush on the caller. It buffers (bounded at 1000 events, oldest dropped) and wakes the background thread. Delivery, which is a blocking httpx.post(timeout=5), now only ever runs on that thread, so AsyncOilPriceAPI no longer posts from the event loop. - close() sets enabled=False, signals the flush loop to drain and exit, and joins it with a bounded timeout. It is idempotent and never raises. - configure_telemetry closes the previous global instance instead of leaking its thread. - Both clients route tracking through an identical _track_telemetry guard, so a telemetry failure can no longer change or fail a request result. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015ao5paex73xXvuM424Libo --- oilpriceapi/async_client.py | 23 +- oilpriceapi/client.py | 23 +- oilpriceapi/telemetry.py | 88 +++++-- tests/unit/test_telemetry_lifecycle.py | 308 +++++++++++++++++++++++++ 4 files changed, 420 insertions(+), 22 deletions(-) create mode 100644 tests/unit/test_telemetry_lifecycle.py diff --git a/oilpriceapi/async_client.py b/oilpriceapi/async_client.py index b6cafa6..f2e0698 100644 --- a/oilpriceapi/async_client.py +++ b/oilpriceapi/async_client.py @@ -254,7 +254,7 @@ async def request( logger.debug(f"Async API response: {response.status_code} for {method} {url}") if 200 <= response.status_code < 300: - self._telemetry.track_request( + self._track_telemetry( operation=self._sanitize_path(method, path), duration=_time.time() - start_time, success=True, @@ -352,7 +352,7 @@ async def request( raise last_exception if last_exception: - self._telemetry.track_request( + self._track_telemetry( operation=self._sanitize_path(method, path), duration=_time.time() - start_time, success=False, @@ -362,6 +362,18 @@ async def request( raise OilPriceAPIError("Max retries exceeded") + def _track_telemetry(self, **fields: Any) -> None: + """ + Record a telemetry event without ever affecting the request result. + + Telemetry is opt-in and best effort: a failure inside the collector + must never change, delay, or fail an API call (#105). + """ + try: + self._telemetry.track_request(**fields) + except Exception: # pragma: no cover - telemetry must never surface + logger.debug("Telemetry tracking failed", exc_info=True) + @staticmethod def _sanitize_path(method: str, path: str) -> str: """Strip resource IDs from path for telemetry privacy.""" @@ -402,8 +414,11 @@ async def market_brief( return MarketBrief(**unwrap_data(response)) async def close(self): - """Close the HTTP client and flush telemetry.""" - self._telemetry.close() + """Close the HTTP client and stop telemetry.""" + try: + self._telemetry.close() + except Exception: # pragma: no cover - telemetry must never surface + logger.debug("Telemetry close failed", exc_info=True) if self._client: await self._client.aclose() self._client = None diff --git a/oilpriceapi/client.py b/oilpriceapi/client.py index 57d4ce4..634815e 100644 --- a/oilpriceapi/client.py +++ b/oilpriceapi/client.py @@ -292,7 +292,7 @@ def request( logger.debug(f"API response: {response.status_code} for {method} {url}") if 200 <= response.status_code < 300: - self._telemetry.track_request( + self._track_telemetry( operation=self._sanitize_path_for_telemetry(method, path), duration=time.time() - start_time, success=True, @@ -398,7 +398,7 @@ def request( raise last_exception if last_exception: - self._telemetry.track_request( + self._track_telemetry( operation=f"{method} {path}", duration=time.time() - start_time, success=False, @@ -541,6 +541,18 @@ def request_with_headers( raise OilPriceAPIError("Max retries exceeded") + def _track_telemetry(self, **fields: Any) -> None: + """ + Record a telemetry event without ever affecting the request result. + + Telemetry is opt-in and best effort: a failure inside the collector + must never change, delay, or fail an API call (#105). + """ + try: + self._telemetry.track_request(**fields) + except Exception: # pragma: no cover - telemetry must never surface + logger.debug("Telemetry tracking failed", exc_info=True) + @staticmethod def _sanitize_path_for_telemetry(method: str, path: str) -> str: """Strip resource IDs from path to avoid leaking user data in telemetry.""" @@ -626,8 +638,11 @@ def market_brief( return MarketBrief(**unwrap_data(response)) def close(self): - """Close the HTTP client and flush telemetry.""" - self._telemetry.close() + """Close the HTTP client and stop telemetry.""" + try: + self._telemetry.close() + except Exception: # pragma: no cover - telemetry must never surface + logger.debug("Telemetry close failed", exc_info=True) self._client.close() def __enter__(self): diff --git a/oilpriceapi/telemetry.py b/oilpriceapi/telemetry.py index 23d8a55..f8f4344 100644 --- a/oilpriceapi/telemetry.py +++ b/oilpriceapi/telemetry.py @@ -55,14 +55,23 @@ class Telemetry: Opt-in telemetry collector for SDK health monitoring. Helps detect issues like v1.4.1 timeout bug across user base. + + Delivery is always performed on a background daemon thread; no SDK call + path ever blocks on the telemetry endpoint. """ + #: Buffered events that trigger an early (still background) delivery. + FLUSH_THRESHOLD = 10 + #: Hard cap on the buffer so an unreachable collector cannot leak memory. + MAX_BUFFERED_EVENTS = 1000 + def __init__( self, enabled: bool = False, endpoint: str = "https://telemetry.oilpriceapi.com/v1/events", flush_interval: int = 300, # 5 minutes - debug: bool = False + debug: bool = False, + close_timeout: float = 2.0, ): """ Initialize telemetry. @@ -72,17 +81,26 @@ def __init__( endpoint: Telemetry endpoint URL flush_interval: Seconds between metric flushes debug: Print telemetry events (for testing) + close_timeout: Seconds close() waits for the flush thread to stop """ self.enabled = enabled and HTTPX_AVAILABLE self.endpoint = endpoint self.flush_interval = flush_interval self.debug = debug + self.close_timeout = close_timeout - # Event buffer + # Event buffer. Bounded: telemetry must never grow without limit when + # the collector is unreachable. self._events: List[Dict[str, Any]] = [] self._lock = threading.Lock() self._last_flush = time.time() + # Lifecycle. Delivery happens only on the background thread, which is + # woken either by the flush interval or by a full buffer. + self._flush_thread: Optional[threading.Thread] = None + self._wake = threading.Event() + self._stopping = False + # Session info (collected once) self._session_id = self._generate_session_id() self._sdk_version = self._get_sdk_version() @@ -159,17 +177,29 @@ def track_request( with self._lock: self._events.append(event) + overflow = len(self._events) - self.MAX_BUFFERED_EVENTS + if overflow > 0: + # Drop the oldest events rather than grow without bound. + del self._events[:overflow] + buffered = len(self._events) if self.debug: print(f"[Telemetry] {operation}: {duration*1000:.0f}ms success={success}") - # Flush if buffer is large or time elapsed - if len(self._events) >= 10 or (time.time() - self._last_flush) > self.flush_interval: - self._flush() + # Ask the background thread to deliver. Never send on the caller's + # thread: callers include AsyncOilPriceAPI, running on the event loop. + if buffered >= self.FLUSH_THRESHOLD: + self._wake.set() def _flush(self): - """Flush events to telemetry endpoint.""" - if not self.enabled or not HTTPX_AVAILABLE: + """ + Deliver buffered events to the telemetry endpoint. + + Only ever called from the background flush thread. It is deliberately + not gated on ``self.enabled`` so the final flush during close() can + still drain the buffer after the collector has been disabled. + """ + if not HTTPX_AVAILABLE: return with self._lock: @@ -181,7 +211,7 @@ def _flush(self): self._last_flush = time.time() try: - # Send telemetry in background (non-blocking) + # Runs on the background flush thread, never on a caller's thread. payload = { "events": events, "sdk": "oilpriceapi-python", @@ -205,15 +235,39 @@ def _flush(self): # Silently fail - don't affect SDK operations def _flush_loop(self): - """Background thread to flush telemetry periodically.""" - while self.enabled: - time.sleep(self.flush_interval) + """Background thread that owns every telemetry delivery.""" + while True: + # Wakes early when the buffer fills or when close() is called. + self._wake.wait(self.flush_interval) + self._wake.clear() self._flush() + if self._stopping: + return def close(self): - """Flush remaining events and close telemetry.""" - if self.enabled: - self._flush() + """ + Stop background delivery and drain what is buffered. + + Idempotent: calling it twice (or on a disabled collector) is a no-op, + and it never raises. After close() the collector accepts no further + events and leaves no thread running. + """ + self._stopping = True + was_enabled = self.enabled + self.enabled = False + + thread = self._flush_thread + self._flush_thread = None + self._wake.set() + + if not was_enabled or thread is None: + with self._lock: + self._events.clear() + return + + if thread.is_alive() and thread is not threading.current_thread(): + # Bounded: a stuck collector must not hang the caller's shutdown. + thread.join(timeout=self.close_timeout) # Global telemetry instance (disabled by default) @@ -239,8 +293,14 @@ def configure_telemetry( if endpoint: kwargs["endpoint"] = endpoint + previous = _global_telemetry _global_telemetry = Telemetry(**kwargs) + # Replacing the global config must not leave the previous flush thread + # running. + if previous is not None: + previous.close() + def get_telemetry() -> Optional[Telemetry]: """Get global telemetry instance.""" diff --git a/tests/unit/test_telemetry_lifecycle.py b/tests/unit/test_telemetry_lifecycle.py new file mode 100644 index 0000000..caabce6 --- /dev/null +++ b/tests/unit/test_telemetry_lifecycle.py @@ -0,0 +1,308 @@ +""" +Lifecycle and non-blocking guarantees for the opt-in telemetry transport (#105). + +The telemetry collector is an advertised public option (``enable_telemetry`` on +both clients, plus ``configure_telemetry``), so it is fixed rather than removed. +These tests pin the four defects called out in the issue: + +1. delivery must never happen synchronously on the caller's thread / event loop +2. ``close()`` must stop background activity and be idempotent +3. replacing the global config must not leak the previous thread +4. a telemetry failure must never alter or fail a request result + +No test here is ever allowed to make a real network call: every test replaces +``oilpriceapi.telemetry.httpx.post`` with a local sink. +""" + +import asyncio +import inspect +import threading +import time +from unittest.mock import Mock + +import httpx +import pytest +import respx + +from oilpriceapi import OilPriceAPI +from oilpriceapi.async_client import AsyncOilPriceAPI +from oilpriceapi.telemetry import Telemetry, configure_telemetry, get_telemetry + +API_KEY = "test_api_key_12345" + + +@pytest.fixture(autouse=True) +def no_real_telemetry_network(monkeypatch): + """Fail loudly if telemetry ever reaches the real httpx.post.""" + + def _boom(*args, **kwargs): # pragma: no cover - only runs on regression + raise AssertionError("telemetry attempted a real network call") + + monkeypatch.setattr("oilpriceapi.telemetry.httpx.post", _boom) + yield + + +@pytest.fixture +def sink(monkeypatch): + """Replace the telemetry transport with a local, instrumented sink.""" + + calls = [] + + def _post(url, **kwargs): + calls.append({"url": url, "json": kwargs.get("json")}) + time.sleep(_post.delay) + return Mock(status_code=200) + + _post.delay = 0.0 + _post.calls = calls + monkeypatch.setattr("oilpriceapi.telemetry.httpx.post", _post) + return _post + + +def _track_n(telemetry, n): + for i in range(n): + telemetry.track_request(operation=f"GET /v1/op{i}", duration=0.01, success=True) + + +# --------------------------------------------------------------------------- +# 1. Disabled mode: no thread, no network +# --------------------------------------------------------------------------- + + +def test_disabled_telemetry_starts_no_thread_and_makes_no_network_call(sink): + before = threading.active_count() + telemetry = Telemetry(enabled=False) + try: + _track_n(telemetry, 25) + assert threading.active_count() == before + assert getattr(telemetry, "_flush_thread", None) is None + assert sink.calls == [] + finally: + telemetry.close() + + +def test_client_defaults_start_no_telemetry_thread(sink): + before = threading.active_count() + with OilPriceAPI(api_key=API_KEY) as client: + assert client._telemetry.enabled is False + assert threading.active_count() == before + assert sink.calls == [] + + +# --------------------------------------------------------------------------- +# 2. Enabled mode: ten tracked calls cannot block the caller +# --------------------------------------------------------------------------- + + +def test_ten_tracked_events_do_not_block_the_calling_thread(sink): + sink.delay = 3.0 + telemetry = Telemetry(enabled=True, endpoint="http://localhost:1/telemetry") + try: + start = time.monotonic() + _track_n(telemetry, 10) + elapsed = time.monotonic() - start + assert elapsed < 0.5, ( + f"track_request blocked the caller for {elapsed:.2f}s; " + "delivery must not run on the caller's thread" + ) + finally: + telemetry.enabled = False + telemetry.close() + + +@respx.mock +def test_sync_client_requests_are_not_blocked_by_telemetry_delivery(sink): + sink.delay = 3.0 + respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( + return_value=httpx.Response(200, json={"status": "success", "data": {}}) + ) + client = OilPriceAPI(api_key=API_KEY, enable_telemetry=True) + try: + start = time.monotonic() + for _ in range(10): + client.request("GET", "/v1/prices/latest") + elapsed = time.monotonic() - start + assert elapsed < 1.0, f"ten client calls took {elapsed:.2f}s with a slow telemetry sink" + finally: + client._telemetry.enabled = False + client.close() + + +@pytest.mark.asyncio +@respx.mock +async def test_async_client_event_loop_is_not_blocked_by_telemetry(sink): + sink.delay = 3.0 + respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( + return_value=httpx.Response(200, json={"status": "success", "data": {}}) + ) + client = AsyncOilPriceAPI(api_key=API_KEY, enable_telemetry=True) + + ticks = 0 + + async def heartbeat(): + nonlocal ticks + while True: + await asyncio.sleep(0.01) + ticks += 1 + + beat = asyncio.create_task(heartbeat()) + try: + start = time.monotonic() + for _ in range(10): + await client.request("GET", "/v1/prices/latest") + elapsed = time.monotonic() - start + assert elapsed < 1.0, f"ten async calls took {elapsed:.2f}s with a slow telemetry sink" + ticks_after = ticks + await asyncio.sleep(0.05) + assert ticks > ticks_after, "event loop stalled while telemetry was delivering" + finally: + beat.cancel() + client._telemetry.enabled = False + await client.close() + + +# --------------------------------------------------------------------------- +# 3. close() stops background activity and is idempotent +# --------------------------------------------------------------------------- + + +def test_close_stops_the_background_thread_and_is_idempotent(sink): + telemetry = Telemetry(enabled=True, endpoint="http://localhost:1/telemetry") + thread = telemetry._flush_thread + assert thread is not None and thread.is_alive() + + telemetry.close() + + assert telemetry.enabled is False + thread.join(timeout=2.0) + assert not thread.is_alive(), "close() left the telemetry flush thread running" + + # Idempotent: a second close is a no-op and never raises. + telemetry.close() + telemetry.close() + + # And no further events are buffered or delivered after close. + calls_before = len(sink.calls) + _track_n(telemetry, 25) + assert telemetry._events == [] + assert len(sink.calls) == calls_before + + +def test_client_close_stops_telemetry_thread(sink): + client = OilPriceAPI(api_key=API_KEY, enable_telemetry=True) + thread = client._telemetry._flush_thread + assert thread is not None and thread.is_alive() + client.close() + thread.join(timeout=2.0) + assert not thread.is_alive() + client.close() # idempotent at the client level too + + +@pytest.mark.asyncio +async def test_async_client_close_stops_telemetry_thread(sink): + client = AsyncOilPriceAPI(api_key=API_KEY, enable_telemetry=True) + thread = client._telemetry._flush_thread + assert thread is not None and thread.is_alive() + await client.close() + thread.join(timeout=2.0) + assert not thread.is_alive() + + +# --------------------------------------------------------------------------- +# 4. Replacing the global config leaves no prior thread alive +# --------------------------------------------------------------------------- + + +def test_configure_telemetry_stops_the_previous_instance(sink): + import oilpriceapi.telemetry as telemetry_module + + previous_global = telemetry_module._global_telemetry + try: + configure_telemetry(enabled=True, endpoint="http://localhost:1/telemetry") + first = get_telemetry() + assert first is not None + first_thread = first._flush_thread + assert first_thread is not None and first_thread.is_alive() + + configure_telemetry(enabled=True, endpoint="http://localhost:1/telemetry") + second = get_telemetry() + assert second is not first + + first_thread.join(timeout=2.0) + assert not first_thread.is_alive(), "configure_telemetry leaked the previous flush thread" + assert first.enabled is False + + second_thread = second._flush_thread + configure_telemetry(enabled=False) + if second_thread is not None: + second_thread.join(timeout=2.0) + assert not second_thread.is_alive() + finally: + current = telemetry_module._global_telemetry + if current is not None: + current.close() + telemetry_module._global_telemetry = previous_global + + +# --------------------------------------------------------------------------- +# 5. A telemetry failure cannot alter or fail a request result +# --------------------------------------------------------------------------- + + +class _ExplodingTelemetry: + enabled = True + + def track_request(self, *args, **kwargs): + raise RuntimeError("telemetry backend exploded") + + def close(self): + raise RuntimeError("telemetry close exploded") + + +@respx.mock +def test_sync_request_result_survives_a_telemetry_failure(sink): + payload = {"status": "success", "data": {"code": "BRENT_CRUDE_USD", "price": 75.5}} + respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( + return_value=httpx.Response(200, json=payload) + ) + client = OilPriceAPI(api_key=API_KEY) + client._telemetry = _ExplodingTelemetry() + result = client.request("GET", "/v1/prices/latest") + assert result == payload + + +@pytest.mark.asyncio +@respx.mock +async def test_async_request_result_survives_a_telemetry_failure(sink): + payload = {"status": "success", "data": {"code": "BRENT_CRUDE_USD", "price": 75.5}} + respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( + return_value=httpx.Response(200, json=payload) + ) + client = AsyncOilPriceAPI(api_key=API_KEY) + client._telemetry = _ExplodingTelemetry() + try: + result = await client.request("GET", "/v1/prices/latest") + assert result == payload + finally: + client._telemetry = Telemetry(enabled=False) + await client.close() + + +# --------------------------------------------------------------------------- +# 6. Sync and async clients must handle telemetry identically +# --------------------------------------------------------------------------- + + +def test_sync_and_async_clients_share_identical_telemetry_handling(): + sync_src = inspect.getsource(OilPriceAPI._track_telemetry) + async_src = inspect.getsource(AsyncOilPriceAPI._track_telemetry) + assert sync_src == async_src, "sync and async telemetry guards have diverged" + + for cls in (OilPriceAPI, AsyncOilPriceAPI): + close_src = inspect.getsource(cls.close) + assert "self._telemetry.close()" in close_src + request_src = inspect.getsource(cls.request) + assert "self._telemetry.track_request(" not in request_src, ( + f"{cls.__name__}.request calls track_request directly instead of the guarded helper" + ) + assert "self._track_telemetry(" in request_src From 8473fd40d2ce69ca50107c329b12e5f5fdb01939 Mon Sep 17 00:00:00 2001 From: Karl Date: Sun, 13 Sep 2026 14:51:25 -0400 Subject: [PATCH 2/2] test(telemetry): drop respx so the new tests run on CI's dev extras (#105) respx is not in the [dev] extra, so the new module failed to import on every matrix Python and collection errored. Mock at httpx.Client.request / httpx.AsyncClient.request instead, which is what the rest of tests/unit does. Same assertions, same red against pre-fix code (10 failed, 2 passed). Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015ao5paex73xXvuM424Libo --- tests/unit/test_telemetry_lifecycle.py | 53 +++++++++++++------------- 1 file changed, 27 insertions(+), 26 deletions(-) diff --git a/tests/unit/test_telemetry_lifecycle.py b/tests/unit/test_telemetry_lifecycle.py index caabce6..7968faa 100644 --- a/tests/unit/test_telemetry_lifecycle.py +++ b/tests/unit/test_telemetry_lifecycle.py @@ -11,18 +11,18 @@ 4. a telemetry failure must never alter or fail a request result No test here is ever allowed to make a real network call: every test replaces -``oilpriceapi.telemetry.httpx.post`` with a local sink. +``oilpriceapi.telemetry.httpx.post`` with a local sink, and the API calls are +patched at ``httpx.Client.request`` / ``httpx.AsyncClient.request`` the way the +rest of the unit suite does it. """ import asyncio import inspect import threading import time -from unittest.mock import Mock +from unittest.mock import Mock, patch -import httpx import pytest -import respx from oilpriceapi import OilPriceAPI from oilpriceapi.async_client import AsyncOilPriceAPI @@ -59,6 +59,11 @@ def _post(url, **kwargs): return _post +def _ok_response(payload=None): + """A minimal stand-in for a 200 httpx.Response.""" + return Mock(status_code=200, json=Mock(return_value=payload or {"status": "success", "data": {}})) + + def _track_n(telemetry, n): for i in range(n): telemetry.track_request(operation=f"GET /v1/op{i}", duration=0.01, success=True) @@ -110,12 +115,10 @@ def test_ten_tracked_events_do_not_block_the_calling_thread(sink): telemetry.close() -@respx.mock -def test_sync_client_requests_are_not_blocked_by_telemetry_delivery(sink): +@patch("httpx.Client.request") +def test_sync_client_requests_are_not_blocked_by_telemetry_delivery(mock_request, sink): sink.delay = 3.0 - respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( - return_value=httpx.Response(200, json={"status": "success", "data": {}}) - ) + mock_request.return_value = _ok_response() client = OilPriceAPI(api_key=API_KEY, enable_telemetry=True) try: start = time.monotonic() @@ -129,12 +132,10 @@ def test_sync_client_requests_are_not_blocked_by_telemetry_delivery(sink): @pytest.mark.asyncio -@respx.mock -async def test_async_client_event_loop_is_not_blocked_by_telemetry(sink): +@patch("httpx.AsyncClient.request") +async def test_async_client_event_loop_is_not_blocked_by_telemetry(mock_request, sink): sink.delay = 3.0 - respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( - return_value=httpx.Response(200, json={"status": "success", "data": {}}) - ) + mock_request.return_value = _ok_response() client = AsyncOilPriceAPI(api_key=API_KEY, enable_telemetry=True) ticks = 0 @@ -259,25 +260,25 @@ def close(self): raise RuntimeError("telemetry close exploded") -@respx.mock -def test_sync_request_result_survives_a_telemetry_failure(sink): +@patch("httpx.Client.request") +def test_sync_request_result_survives_a_telemetry_failure(mock_request, sink): payload = {"status": "success", "data": {"code": "BRENT_CRUDE_USD", "price": 75.5}} - respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( - return_value=httpx.Response(200, json=payload) - ) + mock_request.return_value = _ok_response(payload) client = OilPriceAPI(api_key=API_KEY) client._telemetry = _ExplodingTelemetry() - result = client.request("GET", "/v1/prices/latest") - assert result == payload + try: + result = client.request("GET", "/v1/prices/latest") + assert result == payload + finally: + client._telemetry = Telemetry(enabled=False) + client.close() @pytest.mark.asyncio -@respx.mock -async def test_async_request_result_survives_a_telemetry_failure(sink): +@patch("httpx.AsyncClient.request") +async def test_async_request_result_survives_a_telemetry_failure(mock_request, sink): payload = {"status": "success", "data": {"code": "BRENT_CRUDE_USD", "price": 75.5}} - respx.get("https://api.oilpriceapi.com/v1/prices/latest").mock( - return_value=httpx.Response(200, json=payload) - ) + mock_request.return_value = _ok_response(payload) client = AsyncOilPriceAPI(api_key=API_KEY) client._telemetry = _ExplodingTelemetry() try: