From b2331dbeb13faf42382a7fa2d0903460687ea0a0 Mon Sep 17 00:00:00 2001 From: 9qeklajc Date: Sun, 27 Sep 2026 00:26:15 +0200 Subject: [PATCH] perf: reuse process ssl context for x-cashu clients --- routstr/core/logging.py | 18 +++++++++++++++++- routstr/upstream/http_client.py | 1 + tests/unit/test_queued_logging.py | 21 +++++++++++++++++++++ tests/unit/test_upstream_http_client.py | 16 ++++++++++++++++ 4 files changed, 55 insertions(+), 1 deletion(-) diff --git a/routstr/core/logging.py b/routstr/core/logging.py index 04bbb7f0..8ad91490 100644 --- a/routstr/core/logging.py +++ b/routstr/core/logging.py @@ -155,6 +155,7 @@ class QueuedDailyRotatingFileHandler(logging.Handler): self._kwargs = kwargs self._stopped = True self._next_open_attempt = 0.0 + self._dropped_warned_at = -1.0 self._open() def _open(self) -> None: @@ -175,7 +176,6 @@ class QueuedDailyRotatingFileHandler(logging.Handler): self._target = target self._listener = listener self._stopped = False - self._closed = False with getattr(logging, "_lock"): handler_list = getattr(logging, "_handlerList") # This wrapper owns the target's shutdown and lock ordering. @@ -231,11 +231,27 @@ class QueuedDailyRotatingFileHandler(logging.Handler): try: # Do not acquire the module lock while holding the handler lock. if not self._reopen_locked(): + self._warn_records_dropped() return False except Exception: self.handleError(record) return False + def _warn_records_dropped(self) -> None: + """Report once per backoff window instead of dropping records silently.""" + self.acquire() + try: + if self._dropped_warned_at >= self._next_open_attempt: + return + self._dropped_warned_at = self._next_open_attempt + finally: + self.release() + sys.stderr.write( + f"Logging listener for {self._filename} is unavailable; dropping " + f"records until the next reopen attempt in " + f"{self._reopen_backoff_seconds}s\n" + ) + def _emit_synchronously(self, record: logging.LogRecord) -> None: try: sys.stderr.write(self.format(record) + "\n") diff --git a/routstr/upstream/http_client.py b/routstr/upstream/http_client.py index 36a49e37..fb8cf7f9 100644 --- a/routstr/upstream/http_client.py +++ b/routstr/upstream/http_client.py @@ -500,6 +500,7 @@ def build_x_cashu_client() -> httpx.AsyncClient: """ return httpx.AsyncClient( transport=httpx.AsyncHTTPTransport( + verify=_shared_ssl_context(), retries=UPSTREAM_CONNECT_RETRIES, ), timeout=httpx.Timeout( diff --git a/tests/unit/test_queued_logging.py b/tests/unit/test_queued_logging.py index 83019cef..138bce35 100644 --- a/tests/unit/test_queued_logging.py +++ b/tests/unit/test_queued_logging.py @@ -98,6 +98,27 @@ def test_queued_file_handler_contains_reopen_failures( handler.close() +def test_queued_file_handler_reports_records_dropped_during_backoff( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, capsys: pytest.CaptureFixture[str] +) -> None: + logger, handler = _make_handler(tmp_path, "queued-file-drop-report-test") + handler.close() + + def fail_to_open(*args: object, **kwargs: object) -> None: + raise OSError("disk unavailable") + + monkeypatch.setattr(routstr_logging, "DailyRotatingFileHandler", fail_to_open) + monkeypatch.setattr(type(handler), "handleError", lambda _self, _r: None) + + for _ in range(50): + logger.info("must not vanish without a trace") + + stderr = capsys.readouterr().err + assert stderr.count("dropping records") == 1 + assert "is unavailable" in stderr + handler.close() + + def test_queued_file_handler_emit_does_not_raise_into_caller( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: diff --git a/tests/unit/test_upstream_http_client.py b/tests/unit/test_upstream_http_client.py index bb8bb6b8..3de90aa4 100644 --- a/tests/unit/test_upstream_http_client.py +++ b/tests/unit/test_upstream_http_client.py @@ -2,6 +2,7 @@ import asyncio import concurrent.futures import threading from collections.abc import Callable +from typing import Any, cast from unittest.mock import MagicMock, patch import httpx @@ -12,12 +13,27 @@ from routstr.core.exceptions import UpstreamError from routstr.core.settings import settings from routstr.upstream.http_client import ( acquire_upstream_http_client, + build_x_cashu_client, close_upstream_http_client, get_upstream_http_client, upstream_origin_key, ) +@pytest.mark.asyncio +async def test_x_cashu_client_reuses_the_process_ssl_context() -> None: + """A per-request client must not reload the CA bundle on every call.""" + pooled = get_upstream_http_client("https://api.example.com/v1/chat") + owned = build_x_cashu_client() + try: + pooled_transport = cast(Any, pooled)._transport + owned_transport = cast(Any, owned)._transport + assert owned_transport._pool._ssl_context is pooled_transport._pool._ssl_context + finally: + await owned.aclose() + await close_upstream_http_client() + + @pytest.mark.asyncio async def test_upstream_http_client_is_reused_until_shutdown() -> None: first = get_upstream_http_client("https://api.example.com/v1/chat")