mirror of
https://github.com/Routstr/routstr-core.git
synced 2026-10-05 12:28:22 +00:00
perf: reuse process ssl context for x-cashu clients
This commit is contained in:
+17
-1
@@ -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")
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user