import asyncio from decimal import Decimal from typing import Any import httpx from opentelemetry.sdk.metrics import MeterProvider from opentelemetry.sdk.metrics.export import InMemoryMetricReader from app.config import Settings from app.controllers.v1.delivery import get_aggregator_service from app.observability.metrics import reset_metrics_for_tests, setup_metrics from app.schemas.request import DeliveryEntity, DeliveryRequest from app.schemas.response import DeliveryPrice from app.services.aggregator import AggregatorService, AggregatorServiceError class ToggleAggregatorService: def __init__(self) -> None: self.should_fail = False async def get_all_prices(self, request: DeliveryRequest) -> list[DeliveryPrice]: _ = request if self.should_fail: raise AggregatorServiceError("failed") return [] class StubProvider: def __init__(self, name: str, *, fail: bool) -> None: self.name = name self.fail = fail self.cache_ttl_seconds = 120 async def get_price(self, request: DeliveryRequest) -> DeliveryPrice: _ = request if self.fail: raise RuntimeError("provider unavailable") return DeliveryPrice( provider=self.name, service_name="economy", price=Decimal("99.90"), currency="RUB", delivery_days_min=2, delivery_days_max=3, ) class StubCache: async def get(self, key: str) -> object | None: _ = key return None async def set(self, key: str, value: object, ttl: int | None = None) -> None: _ = key, value, ttl class StubCacheHit: def __init__(self, payload: object) -> None: self.payload = payload async def get(self, key: str) -> object | None: _ = key return self.payload async def set(self, key: str, value: object, ttl: int | None = None) -> None: _ = key, value, ttl def _build_payload() -> dict[str, object]: return { "entity": "individual", "from_city": "Moscow", "to_city": "Kazan", "weight_kg": 2.5, "length_cm": 30.0, "width_cm": 20.0, "height_cm": 10.0, } def _build_request() -> DeliveryRequest: return DeliveryRequest( entity=DeliveryEntity.INDIVIDUAL, from_city="Moscow", to_city="Kazan", weight_kg=2.5, length_cm=30.0, width_cm=20.0, height_cm=10.0, ) def _iter_data_points(metrics_data: Any, metric_name: str) -> list[Any]: points: list[Any] = [] if metrics_data is None: return points for resource_metric in metrics_data.resource_metrics: for scope_metric in resource_metric.scope_metrics: for metric in scope_metric.metrics: if metric.name == metric_name: points.extend(metric.data.data_points) return points def test_metrics_export_requests_errors_latency_and_provider_availability() -> None: reset_metrics_for_tests() metric_reader = InMemoryMetricReader() meter_provider = MeterProvider(metric_readers=[metric_reader]) setup_metrics( service_name="g2s-tests", otlp_endpoint="", meter_provider=meter_provider, ) from app.main import create_app toggle_service = ToggleAggregatorService() app = create_app( settings=Settings( observability={ "service_name": "g2s-tests", "otlp_endpoint": "", "log_level": "INFO", } ) ) async def override_service() -> ToggleAggregatorService: return toggle_service app.dependency_overrides[get_aggregator_service] = override_service async def run_requests() -> tuple[httpx.Response, httpx.Response]: transport = httpx.ASGITransport(app=app, raise_app_exceptions=False) async with httpx.AsyncClient( transport=transport, base_url="http://testserver", ) as client: success = await client.post("/api/v1/delivery/price", json=_build_payload()) toggle_service.should_fail = True failure = await client.post("/api/v1/delivery/price", json=_build_payload()) return success, failure success_response, failure_response = asyncio.run(run_requests()) assert success_response.status_code == 200 assert failure_response.status_code == 503 availability_service = AggregatorService( providers=[ StubProvider("available-provider", fail=False), StubProvider("unavailable-provider", fail=True), ], cache=StubCache(), ) asyncio.run(availability_service.get_all_prices(_build_request())) metrics_data = metric_reader.get_metrics_data() request_points = _iter_data_points(metrics_data, "http.server.request.count") error_points = _iter_data_points(metrics_data, "http.server.error.count") latency_points = _iter_data_points(metrics_data, "http.server.request.latency") availability_points = _iter_data_points(metrics_data, "delivery.provider.availability") assert sum(point.value for point in request_points) == 2 assert sum(point.value for point in error_points) == 1 assert sum(point.count for point in latency_points) == 2 assert sum(point.count for point in availability_points) == 2 assert sorted(point.sum for point in availability_points) == [0.0, 1.0] reset_metrics_for_tests() def test_provider_availability_not_recorded_on_cache_hit() -> None: reset_metrics_for_tests() metric_reader = InMemoryMetricReader() meter_provider = MeterProvider(metric_readers=[metric_reader]) setup_metrics( service_name="g2s-tests", otlp_endpoint="", meter_provider=meter_provider, ) cached_price = DeliveryPrice( provider="cached-provider", service_name="economy", price=Decimal("95.50"), currency="RUB", delivery_days_min=2, delivery_days_max=3, ).model_dump(mode="json") provider = StubProvider("cached-provider", fail=True) service = AggregatorService( providers=[provider], cache=StubCacheHit(cached_price), ) result = asyncio.run(service.get_all_prices(_build_request())) assert len(result) == 1 assert result[0].provider == "cached-provider" metrics_data = metric_reader.get_metrics_data() availability_points = _iter_data_points(metrics_data, "delivery.provider.availability") assert availability_points == [] reset_metrics_for_tests()