Files
g2s-aggregator/tests/observability/test_tracing.py
T
Раис Юсупалиев 8d47917ba2 007 add telemetry
2026-03-08 11:43:15 +03:00

161 lines
5.2 KiB
Python

import asyncio
from decimal import Decimal
import httpx
from fastapi import FastAPI
from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor
from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter
from opentelemetry.trace import SpanKind
from app.observability.tracing import (
instrument_fastapi_app,
instrument_httpx_client,
reset_tracing_for_tests,
setup_tracing,
)
from app.schemas.request import DeliveryEntity, DeliveryRequest
from app.schemas.response import DeliveryPrice
from app.services.aggregator import AggregatorService
class StubProvider:
name = "stub-provider"
cache_ttl_seconds = 120
async def get_price(self, request: DeliveryRequest) -> DeliveryPrice:
_ = request
return DeliveryPrice(
provider=self.name,
service_name="economy",
price=Decimal("100.10"),
currency="RUB",
delivery_days_min=2,
delivery_days_max=4,
)
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
def _build_request() -> DeliveryRequest:
return DeliveryRequest(
entity=DeliveryEntity.INDIVIDUAL,
from_city="Moscow",
to_city="Kazan",
weight_kg=1.5,
length_cm=10.0,
width_cm=20.0,
height_cm=30.0,
)
def test_manual_spans_capture_required_attributes() -> None:
reset_tracing_for_tests()
span_exporter = InMemorySpanExporter()
tracer_provider = TracerProvider(resource=Resource.create({"service.name": "tests"}))
tracer_provider.add_span_processor(SimpleSpanProcessor(span_exporter))
setup_tracing(
service_name="tests",
otlp_endpoint="",
tracer_provider=tracer_provider,
)
service = AggregatorService(providers=[StubProvider()], cache=StubCache())
asyncio.run(service.get_all_prices(_build_request()))
spans = span_exporter.get_finished_spans()
span_by_name = {span.name: span for span in spans}
assert "AggregatorService.get_all_prices" in span_by_name
assert "DeliveryProvider.get_price" in span_by_name
assert "PriceCache.get" in span_by_name
assert "PriceCache.set" in span_by_name
service_span = span_by_name["AggregatorService.get_all_prices"]
assert service_span.attributes["from_city"] == "Moscow"
assert service_span.attributes["to_city"] == "Kazan"
assert service_span.attributes["weight_kg"] == 1.5
assert service_span.attributes["tariffs_found"] == 1
provider_span = span_by_name["DeliveryProvider.get_price"]
assert provider_span.attributes["provider"] == "stub-provider"
assert provider_span.attributes["cache_hit"] is False
assert provider_span.attributes["from_city"] == "Moscow"
assert provider_span.attributes["to_city"] == "Kazan"
assert provider_span.attributes["weight_kg"] == 1.5
cache_span = span_by_name["PriceCache.get"]
assert cache_span.attributes["cache_hit"] is False
reset_tracing_for_tests()
def test_fastapi_and_httpx_instrumentation_emit_server_and_client_spans() -> None:
reset_tracing_for_tests()
span_exporter = InMemorySpanExporter()
tracer_provider = TracerProvider(resource=Resource.create({"service.name": "tests"}))
tracer_provider.add_span_processor(SimpleSpanProcessor(span_exporter))
setup_tracing(
service_name="tests",
otlp_endpoint="",
tracer_provider=tracer_provider,
)
app = FastAPI()
@app.get("/proxy")
async def proxy() -> dict[str, int]:
def handler(request: httpx.Request) -> httpx.Response:
_ = request
return httpx.Response(status_code=200, json={"ok": True})
transport = httpx.MockTransport(handler)
async with httpx.AsyncClient(
transport=transport,
base_url="https://provider.test",
) as client:
HTTPXClientInstrumentor.instrument_client(client, tracer_provider=tracer_provider)
response = await client.get("/quote")
return {"status_code": response.status_code}
instrument_fastapi_app(app)
instrument_httpx_client()
async def run_request() -> httpx.Response:
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(
transport=transport,
base_url="http://testserver",
) as client:
return await client.get("/proxy")
response = asyncio.run(run_request())
assert response.status_code == 200
spans = span_exporter.get_finished_spans()
assert any(span.kind == SpanKind.SERVER for span in spans)
assert any(span.kind == SpanKind.CLIENT for span in spans)
assert any(
(
"provider.test/quote"
in (
span.attributes.get("http.url")
or span.attributes.get("url.full")
or ""
)
)
for span in spans
if span.kind == SpanKind.CLIENT
)
reset_tracing_for_tests()