diff --git a/app/config.py b/app/config.py index 48d32d1..c4a4098 100644 --- a/app/config.py +++ b/app/config.py @@ -5,7 +5,7 @@ from pathlib import Path import sys from typing import Any -from pydantic import BaseModel, Field, model_validator +from pydantic import BaseModel, Field from pydantic_settings import ( BaseSettings, PydanticBaseSettingsSource, @@ -46,54 +46,6 @@ class AdapterConfig(BaseModel): cdek_cache_ttl_seconds: int = Field(default=900, gt=0) -class ObservabilityConfig(BaseModel): - service_name: str = "g2s-aggregator" - otlp_endpoint: str = "" - log_level: str = "INFO" - - -class Provider5xxAlertConfig(BaseModel): - error_count: int = Field(default=5, ge=1) - window_minutes: int = Field(default=5, ge=1) - - -class ProviderP99LatencyAlertConfig(BaseModel): - threshold_ms: int = Field(default=5000, ge=1) - window_minutes: int = Field(default=10, ge=1) - - -class ProviderUnavailableAlertConfig(BaseModel): - duration_minutes: int = Field(default=5, ge=1) - - -class AlertsConfig(BaseModel): - telegram_enabled: bool = False - telegram_bot_token: str | None = None - telegram_chat_id: str | None = None - provider_5xx: Provider5xxAlertConfig = Field(default_factory=Provider5xxAlertConfig) - provider_p99_latency: ProviderP99LatencyAlertConfig = Field( - default_factory=ProviderP99LatencyAlertConfig - ) - provider_unavailable: ProviderUnavailableAlertConfig = Field( - default_factory=ProviderUnavailableAlertConfig - ) - - @model_validator(mode="after") - def validate_telegram_credentials(self) -> "AlertsConfig": - if not self.telegram_enabled: - return self - - if not self.telegram_bot_token: - raise ValueError( - "alerts.telegram_bot_token is required when alerts.telegram_enabled=true" - ) - if not self.telegram_chat_id: - raise ValueError( - "alerts.telegram_chat_id is required when alerts.telegram_enabled=true" - ) - return self - - class Settings(BaseSettings): model_config = SettingsConfigDict( extra="ignore", @@ -104,8 +56,6 @@ class Settings(BaseSettings): business_logic: BusinessLogicConfig = Field(default_factory=BusinessLogicConfig) repository: RepositoryConfig = Field(default_factory=RepositoryConfig) adapter: AdapterConfig = Field(default_factory=AdapterConfig) - observability: ObservabilityConfig = Field(default_factory=ObservabilityConfig) - alerts: AlertsConfig = Field(default_factory=AlertsConfig) @classmethod def settings_customise_sources( @@ -133,8 +83,6 @@ class _RequiredYamlSections(BaseModel): business_logic: dict[str, Any] repository: dict[str, Any] adapter: dict[str, Any] - observability: dict[str, Any] - alerts: dict[str, Any] def _resolve_runtime_config_file() -> str: diff --git a/app/controllers/middleware.py b/app/controllers/middleware.py deleted file mode 100644 index c72b930..0000000 --- a/app/controllers/middleware.py +++ /dev/null @@ -1,55 +0,0 @@ -"""HTTP middleware registration.""" - -from uuid import uuid4 -from time import perf_counter - -from fastapi import FastAPI, Request -from starlette.middleware.base import BaseHTTPMiddleware -from starlette.responses import Response -from structlog.contextvars import bind_contextvars, clear_contextvars - -from app.observability.metrics import get_metrics - - -class RequestCorrelationMiddleware(BaseHTTPMiddleware): - """Assign and propagate request correlation identifiers.""" - - def __init__(self, app: FastAPI, *, request_id_header: str) -> None: - super().__init__(app) - self._request_id_header = request_id_header - - async def dispatch(self, request: Request, call_next) -> Response: - request_id = str(uuid4()) - request.state.request_id = request_id - bind_contextvars(request_id=request_id) - - started_at = perf_counter() - status_code = 500 - is_error = True - - try: - response = await call_next(request) - status_code = response.status_code - is_error = status_code >= 500 - response.headers[self._request_id_header] = request_id - return response - finally: - duration_ms = (perf_counter() - started_at) * 1000 - get_metrics().record_http_request( - method=request.method, - route=request.url.path, - status_code=status_code, - duration_ms=duration_ms, - is_error=is_error, - ) - clear_contextvars() - - -def install_middleware(app: FastAPI) -> None: - """Register middleware components for the API.""" - - settings = app.state.settings - app.add_middleware( - RequestCorrelationMiddleware, - request_id_header=settings.controller.request_id_header, - ) diff --git a/app/main.py b/app/main.py index 75c39b0..3e70814 100644 --- a/app/main.py +++ b/app/main.py @@ -3,9 +3,7 @@ from fastapi import FastAPI from app.config import Settings, get_settings -from app.controllers.middleware import install_middleware from app.controllers.v1.delivery import router as delivery_router -from app.observability import setup_observability def create_app(settings: Settings | None = None) -> FastAPI: @@ -14,8 +12,6 @@ def create_app(settings: Settings | None = None) -> FastAPI: application = FastAPI(title="G2S Aggregator", version="0.1.0") application.state.settings = resolved_settings - setup_observability(application, resolved_settings.observability) - install_middleware(application) application.include_router(delivery_router, prefix=resolved_settings.controller.api_prefix) return application diff --git a/app/observability/__init__.py b/app/observability/__init__.py deleted file mode 100644 index 0b572ee..0000000 --- a/app/observability/__init__.py +++ /dev/null @@ -1,5 +0,0 @@ -"""Observability setup and helpers.""" - -from app.observability.setup import setup_observability - -__all__ = ["setup_observability"] diff --git a/app/observability/logging.py b/app/observability/logging.py deleted file mode 100644 index 31065be..0000000 --- a/app/observability/logging.py +++ /dev/null @@ -1,48 +0,0 @@ -"""Structured logging setup with request and trace correlation.""" - -import logging -from typing import Any - -import structlog -from opentelemetry import trace -from opentelemetry.trace import SpanContext - - -def configure_structured_logging(log_level: str) -> None: - resolved_level = _resolve_log_level(log_level) - logging.basicConfig(level=resolved_level, format="%(message)s") - - structlog.configure( - processors=[ - structlog.contextvars.merge_contextvars, - structlog.stdlib.add_log_level, - add_trace_context, - structlog.processors.TimeStamper(fmt="iso", utc=True), - structlog.processors.JSONRenderer(), - ], - logger_factory=structlog.stdlib.LoggerFactory(), - wrapper_class=structlog.make_filtering_bound_logger(resolved_level), - cache_logger_on_first_use=True, - ) - - -def add_trace_context( - _: Any, - __: str, - event_dict: dict[str, Any], -) -> dict[str, Any]: - span = trace.get_current_span() - span_context = span.get_span_context() - if _is_valid_span_context(span_context): - event_dict["trace_id"] = format(span_context.trace_id, "032x") - event_dict["span_id"] = format(span_context.span_id, "016x") - return event_dict - - -def _resolve_log_level(log_level: str) -> int: - candidate = log_level.strip().upper() - return logging.getLevelNamesMapping().get(candidate, logging.INFO) - - -def _is_valid_span_context(span_context: SpanContext) -> bool: - return bool(span_context.is_valid and span_context.trace_id and span_context.span_id) diff --git a/app/observability/metrics.py b/app/observability/metrics.py deleted file mode 100644 index 8fb732c..0000000 --- a/app/observability/metrics.py +++ /dev/null @@ -1,129 +0,0 @@ -"""OpenTelemetry metrics setup and recording helpers.""" - -from typing import Sequence - -from opentelemetry import metrics -from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter -from opentelemetry.sdk.metrics import MeterProvider -from opentelemetry.sdk.metrics.export import ( - MetricReader, - PeriodicExportingMetricReader, -) -from opentelemetry.sdk.resources import Resource - -_METER_PROVIDER: MeterProvider | None = None -_METRICS: "ObservabilityMetrics" | None = None -_METER_NAME = "g2s.observability" - - -class ObservabilityMetrics: - def __init__( - self, - meter_name: str, - *, - meter_provider: MeterProvider | None = None, - ) -> None: - if meter_provider is None: - meter = metrics.get_meter(meter_name) - else: - meter = meter_provider.get_meter(meter_name) - self._request_count = meter.create_counter( - name="http.server.request.count", - unit="1", - description="Total number of HTTP requests.", - ) - self._error_count = meter.create_counter( - name="http.server.error.count", - unit="1", - description="Total number of HTTP server errors.", - ) - self._request_latency = meter.create_histogram( - name="http.server.request.latency", - unit="ms", - description="HTTP request latency.", - ) - self._provider_availability = meter.create_histogram( - name="delivery.provider.availability", - unit="1", - description="Provider availability sample (1=available, 0=unavailable).", - ) - - def record_http_request( - self, - *, - method: str, - route: str, - status_code: int, - duration_ms: float, - is_error: bool, - ) -> None: - attributes = { - "http.method": method, - "http.route": route, - "http.status_code": status_code, - } - self._request_count.add(1, attributes) - self._request_latency.record(duration_ms, attributes) - if is_error: - self._error_count.add(1, attributes) - - def record_provider_availability( - self, - *, - provider: str, - is_available: bool, - ) -> None: - self._provider_availability.record( - 1.0 if is_available else 0.0, - {"provider": provider}, - ) - - -def setup_metrics( - *, - service_name: str, - otlp_endpoint: str, - meter_provider: MeterProvider | None = None, - metric_readers: Sequence[MetricReader] | None = None, -) -> MeterProvider: - global _METER_PROVIDER, _METRICS - - if meter_provider is None and _METER_PROVIDER is not None: - return _METER_PROVIDER - - if meter_provider is None: - readers = list(metric_readers or ()) - if otlp_endpoint: - metric_exporter = OTLPMetricExporter( - endpoint=otlp_endpoint, - insecure=otlp_endpoint.startswith("http://"), - ) - readers.append(PeriodicExportingMetricReader(metric_exporter)) - meter_provider = MeterProvider( - resource=Resource.create({"service.name": service_name}), - metric_readers=readers, - ) - - _METER_PROVIDER = meter_provider - _METRICS = ObservabilityMetrics( - _METER_NAME, - meter_provider=meter_provider, - ) - return meter_provider - - -def get_metrics() -> ObservabilityMetrics: - global _METRICS - - if _METRICS is None: - _METRICS = ObservabilityMetrics(_METER_NAME) - return _METRICS - - -def reset_metrics_for_tests() -> None: - global _METER_PROVIDER, _METRICS - - if _METER_PROVIDER is not None: - _METER_PROVIDER.shutdown() - _METER_PROVIDER = None - _METRICS = None diff --git a/app/observability/setup.py b/app/observability/setup.py deleted file mode 100644 index 99cff80..0000000 --- a/app/observability/setup.py +++ /dev/null @@ -1,26 +0,0 @@ -"""High-level observability setup wiring.""" - -from fastapi import FastAPI - -from app.config import ObservabilityConfig -from app.observability.logging import configure_structured_logging -from app.observability.metrics import setup_metrics -from app.observability.tracing import ( - instrument_fastapi_app, - instrument_httpx_client, - setup_tracing, -) - - -def setup_observability(app: FastAPI, config: ObservabilityConfig) -> None: - configure_structured_logging(config.log_level) - setup_tracing( - service_name=config.service_name, - otlp_endpoint=config.otlp_endpoint, - ) - setup_metrics( - service_name=config.service_name, - otlp_endpoint=config.otlp_endpoint, - ) - instrument_fastapi_app(app) - instrument_httpx_client() diff --git a/app/observability/tracing.py b/app/observability/tracing.py deleted file mode 100644 index 5586632..0000000 --- a/app/observability/tracing.py +++ /dev/null @@ -1,80 +0,0 @@ -"""OpenTelemetry tracing setup and instrumentation helpers.""" - -from typing import cast - -from fastapi import FastAPI -from opentelemetry import trace -from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter -from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor -from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor -from opentelemetry.sdk.resources import Resource -from opentelemetry.sdk.trace import TracerProvider -from opentelemetry.sdk.trace.export import BatchSpanProcessor -from opentelemetry.trace import Tracer, TracerProvider as TraceAPIProvider - -_TRACER_PROVIDER: TraceAPIProvider | None = None -_HTTPX_INSTRUMENTED = False - - -def setup_tracing( - *, - service_name: str, - otlp_endpoint: str, - tracer_provider: TraceAPIProvider | None = None, -) -> TraceAPIProvider: - global _TRACER_PROVIDER - - if tracer_provider is None and _TRACER_PROVIDER is not None: - return _TRACER_PROVIDER - - if tracer_provider is None: - tracer_provider = TracerProvider( - resource=Resource.create({"service.name": service_name}), - ) - if otlp_endpoint: - span_exporter = OTLPSpanExporter( - endpoint=otlp_endpoint, - insecure=otlp_endpoint.startswith("http://"), - ) - tracer_provider.add_span_processor(BatchSpanProcessor(span_exporter)) - - _TRACER_PROVIDER = tracer_provider - return tracer_provider - - -def instrument_fastapi_app(app: FastAPI) -> None: - tracer_provider = _TRACER_PROVIDER - if tracer_provider is None: - return - FastAPIInstrumentor.instrument_app(app, tracer_provider=tracer_provider) - - -def instrument_httpx_client() -> None: - global _HTTPX_INSTRUMENTED - - tracer_provider = _TRACER_PROVIDER - if tracer_provider is None: - return - - if _HTTPX_INSTRUMENTED: - return - HTTPXClientInstrumentor().instrument(tracer_provider=tracer_provider) - _HTTPX_INSTRUMENTED = True - - -def get_tracer(name: str) -> Tracer: - tracer_provider = _TRACER_PROVIDER - if tracer_provider is None: - return trace.get_tracer(name) - return cast(Tracer, tracer_provider.get_tracer(name)) - - -def reset_tracing_for_tests() -> None: - global _TRACER_PROVIDER, _HTTPX_INSTRUMENTED - - if _HTTPX_INSTRUMENTED: - HTTPXClientInstrumentor().uninstrument() - if isinstance(_TRACER_PROVIDER, TracerProvider): - _TRACER_PROVIDER.shutdown() - _TRACER_PROVIDER = None - _HTTPX_INSTRUMENTED = False diff --git a/app/services/aggregator.py b/app/services/aggregator.py index fcfa3e7..90e836c 100644 --- a/app/services/aggregator.py +++ b/app/services/aggregator.py @@ -13,8 +13,6 @@ from app.domain.price import ( filter_and_sort_prices, normalize_delivery_request, ) -from app.observability.metrics import get_metrics -from app.observability.tracing import get_tracer from app.schemas.request import DeliveryRequest from app.schemas.response import DeliveryPrice @@ -46,41 +44,33 @@ class AggregatorService: self._filter_and_sort_prices = filter_and_sort_prices_fn async def get_all_prices(self, request: DeliveryRequest) -> list[DeliveryPrice]: - tracer = get_tracer(__name__) - with tracer.start_as_current_span("AggregatorService.get_all_prices") as span: - span.set_attribute("from_city", request.from_city) - span.set_attribute("to_city", request.to_city) - span.set_attribute("weight_kg", float(request.weight_kg)) + normalized_request = normalize_delivery_request( + request, weight_round_scale=self._weight_round_scale + ) + provider_request = self._to_provider_request(normalized_request) - normalized_request = normalize_delivery_request( - request, weight_round_scale=self._weight_round_scale - ) - provider_request = self._to_provider_request(normalized_request) + provider_results = await asyncio.gather( + *( + self._get_provider_price( + provider=provider, + request=provider_request, + cache_key=self._build_cache_key( + provider_name=provider.name, + request=normalized_request, + ), + ) + for provider in self._providers + ), + return_exceptions=True, + ) - provider_results = await asyncio.gather( - *( - self._get_provider_price( - provider=provider, - request=provider_request, - cache_key=self._build_cache_key( - provider_name=provider.name, - request=normalized_request, - ), - ) - for provider in self._providers - ), - return_exceptions=True, - ) - - successful_results = [ - price - for price in provider_results - if isinstance(price, DeliveryPrice) - ] - filtered_and_sorted = self._filter_and_sort_prices(successful_results) - result = [self._coerce_delivery_price(price) for price in filtered_and_sorted] - span.set_attribute("tariffs_found", len(result)) - return result + successful_results = [ + price + for price in provider_results + if isinstance(price, DeliveryPrice) + ] + filtered_and_sorted = self._filter_and_sort_prices(successful_results) + return [self._coerce_delivery_price(price) for price in filtered_and_sorted] async def _get_provider_price( self, @@ -89,59 +79,31 @@ class AggregatorService: request: DeliveryRequest, cache_key: str, ) -> DeliveryPrice: - tracer = get_tracer(__name__) - with tracer.start_as_current_span("DeliveryProvider.get_price") as span: - span.set_attribute("provider", provider.name) - span.set_attribute("from_city", request.from_city) - span.set_attribute("to_city", request.to_city) - span.set_attribute("weight_kg", float(request.weight_kg)) + cached_price = await self._get_cached_price(cache_key) + if cached_price is not None: + return cached_price - cached_price = await self._get_cached_price(cache_key) - span.set_attribute("cache_hit", cached_price is not None) - if cached_price is not None: - return cached_price - - try: - fresh_price = await provider.get_price(request) - except Exception: - get_metrics().record_provider_availability( - provider=provider.name, - is_available=False, - ) - raise - - get_metrics().record_provider_availability( - provider=provider.name, - is_available=True, - ) - await self._set_cached_price( - cache_key, - fresh_price, - ttl=getattr(provider, "cache_ttl_seconds", None), - ) - return fresh_price + fresh_price = await provider.get_price(request) + await self._set_cached_price( + cache_key, + fresh_price, + ttl=getattr(provider, "cache_ttl_seconds", None), + ) + return fresh_price async def _get_cached_price(self, cache_key: str) -> DeliveryPrice | None: - tracer = get_tracer(__name__) - with tracer.start_as_current_span("PriceCache.get") as span: - if self._cache is None: - span.set_attribute("cache_hit", False) - return None - try: - payload = await self._cache.get(cache_key) - except Exception: - span.set_attribute("cache_hit", False) - return None - if payload is None: - span.set_attribute("cache_hit", False) - return None - try: - price = self._coerce_delivery_price(payload) - except Exception: - span.set_attribute("cache_hit", False) - return None - span.set_attribute("cache_hit", True) - return price + if self._cache is None: + return None + try: + payload = await self._cache.get(cache_key) + except Exception: + return None + if payload is None: + return None + try: + return self._coerce_delivery_price(payload) + except Exception: + return None async def _set_cached_price( self, @@ -150,14 +112,12 @@ class AggregatorService: *, ttl: int | None, ) -> None: - tracer = get_tracer(__name__) - with tracer.start_as_current_span("PriceCache.set"): - if self._cache is None: - return - try: - await self._cache.set(cache_key, payload, ttl=ttl) - except Exception: - return + if self._cache is None: + return + try: + await self._cache.set(cache_key, payload, ttl=ttl) + except Exception: + return @staticmethod def _to_provider_request(request: NormalizedDeliveryRequest) -> DeliveryRequest: diff --git a/config.example.yaml b/config.example.yaml index 88712ef..48a1dba 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -21,21 +21,3 @@ adapter: cdek_retry_backoff_seconds: 0.2 cdek_timeout_seconds: 10.0 cdek_cache_ttl_seconds: 900 - -observability: - service_name: "g2s-aggregator" - otlp_endpoint: "http://localhost:4317" - log_level: "INFO" - -alerts: - telegram_enabled: false - telegram_bot_token: null - telegram_chat_id: null - provider_5xx: - error_count: 5 - window_minutes: 5 - provider_p99_latency: - threshold_ms: 5000 - window_minutes: 10 - provider_unavailable: - duration_minutes: 5 diff --git a/config.test.yaml b/config.test.yaml index c288f65..ff37899 100644 --- a/config.test.yaml +++ b/config.test.yaml @@ -21,21 +21,3 @@ adapter: cdek_retry_backoff_seconds: 0.2 cdek_timeout_seconds: 10.0 cdek_cache_ttl_seconds: 900 - -observability: - service_name: "g2s-aggregator-tests" - otlp_endpoint: "" - log_level: "INFO" - -alerts: - telegram_enabled: false - telegram_bot_token: null - telegram_chat_id: null - provider_5xx: - error_count: 5 - window_minutes: 5 - provider_p99_latency: - threshold_ms: 5000 - window_minutes: 10 - provider_unavailable: - duration_minutes: 5 diff --git a/docker-compose.yml b/docker-compose.yml index 7504d89..503d857 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -6,7 +6,6 @@ services: - "8000:8000" depends_on: - redis - - signoz volumes: - ./config.yaml:/config.yaml redis: @@ -15,21 +14,3 @@ services: command: ["redis-server", "--save", "", "--appendonly", "no"] ports: - "6379:6379" - signoz: - image: signoz/signoz:latest - container_name: signoz - environment: - G2S_ALERTS__TELEGRAM_BOT_TOKEN: ${G2S_ALERTS__TELEGRAM_BOT_TOKEN:-} - G2S_ALERTS__TELEGRAM_CHAT_ID: ${G2S_ALERTS__TELEGRAM_CHAT_ID:-} - SIGNOZ_ALERTS_CONTACT_POINT_FILE: /etc/signoz/alerts/signoz_contact_point_telegram.yaml - SIGNOZ_ALERTS_RULES_FILE: /etc/signoz/alerts/signoz_alert_rules.yaml - ports: - - "8080:8080" - - "3301:3301" - - "4317:4317" - volumes: - - signoz_data:/var/lib/signoz - - ./infra/alerts:/etc/signoz/alerts:ro - -volumes: - signoz_data: diff --git a/infra/README.md b/infra/README.md index 32826e9..7daa78b 100644 --- a/infra/README.md +++ b/infra/README.md @@ -4,7 +4,6 @@ - `app`: FastAPI service (`uvicorn app.main:app`) built from root `Dockerfile` - `redis`: price cache backend -- `signoz`: observability backend (UI + OTLP receiver) ## Environment wiring @@ -12,20 +11,14 @@ - `G2S_CONFIG_FILE=/app/config.yaml` - `G2S_REPOSITORY__REDIS_DSN=redis://redis:6379/0` -- `G2S_OBSERVABILITY__OTLP_ENDPOINT=http://signoz:4317` - `PYTHONUNBUFFERED=1` -### `signoz` service - -- `G2S_ALERTS__TELEGRAM_BOT_TOKEN` (optional for Telegram alerts) -- `G2S_ALERTS__TELEGRAM_CHAT_ID` (optional for Telegram alerts) - ## Startup instructions 1. Validate compose syntax: - `docker compose config` 2. Start dependencies: - - `docker compose up -d redis signoz` + - `docker compose up -d redis` 3. Verify running containers: - `docker compose ps` 4. Start application: @@ -37,6 +30,6 @@ ```bash docker compose config -docker compose up -d redis signoz +docker compose up -d redis docker compose ps ``` diff --git a/infra/alerts/README.md b/infra/alerts/README.md deleted file mode 100644 index 28f718b..0000000 --- a/infra/alerts/README.md +++ /dev/null @@ -1,39 +0,0 @@ -# SigNoz -> Telegram alerting setup - -## Alerts config section example (`config.yaml`) - -```yaml -alerts: - telegram_enabled: true - telegram_bot_token: "123456789:telegram-bot-token" - telegram_chat_id: "-1000000000000" - provider_5xx: - error_count: 5 - window_minutes: 5 - provider_p99_latency: - threshold_ms: 5000 - window_minutes: 10 - provider_unavailable: - duration_minutes: 5 -``` - -## Routing artifacts - -- `signoz_contact_point_telegram.yaml`: - - webhook endpoint points to Telegram Bot API `sendMessage` - - uses `G2S_ALERTS__TELEGRAM_BOT_TOKEN` and `G2S_ALERTS__TELEGRAM_CHAT_ID` -- `signoz_alert_rules.yaml`: - - provider 5xx: `>= 5 errors / 5m` - - provider p99 latency: `> 5000ms / 10m` - - provider unavailable: `> 5m` - -## Local verification steps - -1. Export Telegram variables (must match `config.yaml -> alerts` values): - - `export G2S_ALERTS__TELEGRAM_BOT_TOKEN=""` - - `export G2S_ALERTS__TELEGRAM_CHAT_ID=""` -2. Validate compose syntax: - - `docker compose config` -3. In SigNoz UI, create a webhook contact point using `infra/alerts/signoz_contact_point_telegram.yaml`. -4. Create three alert rules from `infra/alerts/signoz_alert_rules.yaml`. -5. Use SigNoz "Test alert" action and verify message delivery in Telegram chat. diff --git a/infra/alerts/signoz_alert_rules.yaml b/infra/alerts/signoz_alert_rules.yaml deleted file mode 100644 index ab0995b..0000000 --- a/infra/alerts/signoz_alert_rules.yaml +++ /dev/null @@ -1,41 +0,0 @@ -rules: - - name: "provider-5xx-errors" - description: "Provider returned >= 5 server errors in 5 minutes." - severity: "critical" - condition: - query: | - sum(increase(g2s_provider_errors_total{status_code=~"5.."}[5m])) by (provider) >= 5 - evaluate_for: "0m" - threshold: - error_count: 5 - window_minutes: 5 - annotations: - summary: "Provider {{ $labels.provider }} returned >= 5 5xx responses during last 5m." - - - name: "provider-p99-latency" - description: "Provider p99 latency is above 5000ms for 10 minutes." - severity: "warning" - condition: - query: | - histogram_quantile( - 0.99, - sum(rate(g2s_provider_latency_ms_bucket[10m])) by (provider, le) - ) > 5000 - evaluate_for: "10m" - threshold: - threshold_ms: 5000 - window_minutes: 10 - annotations: - summary: "Provider {{ $labels.provider }} p99 latency > 5000ms for 10m." - - - name: "provider-unavailable" - description: "Provider is unavailable for more than 5 minutes." - severity: "critical" - condition: - query: | - max_over_time(g2s_provider_available[5m]) == 0 - evaluate_for: "5m" - threshold: - duration_minutes: 5 - annotations: - summary: "Provider {{ $labels.provider }} is unavailable for more than 5m." diff --git a/infra/alerts/signoz_contact_point_telegram.yaml b/infra/alerts/signoz_contact_point_telegram.yaml deleted file mode 100644 index fca8db1..0000000 --- a/infra/alerts/signoz_contact_point_telegram.yaml +++ /dev/null @@ -1,14 +0,0 @@ -contact_point: - name: "telegram-webhook" - type: "webhook" - settings: - method: "POST" - url: "https://api.telegram.org/bot${G2S_ALERTS__TELEGRAM_BOT_TOKEN}/sendMessage" - headers: - Content-Type: "application/json" - body: | - { - "chat_id": "${G2S_ALERTS__TELEGRAM_CHAT_ID}", - "text": "[{{ .Status }}] {{ .RuleName }}\nProvider: {{ index .Labels \"provider\" }}\nSummary: {{ index .Annotations \"summary\" }}", - "disable_notification": false - } diff --git a/spec/index.md b/spec/index.md index 2c088ad..5219a49 100644 --- a/spec/index.md +++ b/spec/index.md @@ -1,7 +1,7 @@ # Spec Tasks Index > ⚠️ This file is generated. Do not edit manually. -> Generated at (UTC): `2026-03-08T19:06:27+00:00` +> Generated at (UTC): `2026-03-08T19:58:23+00:00` ## Tasks @@ -19,9 +19,10 @@ | 009 | DONE | 2026-03-07 | Add local infrastructure stack | `spec/tasks/009_add_local_infra_stack.md` | | 010 | DONE | 2026-03-08 | Migrate to YAML-only configuration loading | `spec/tasks/010_migrate_to_yaml_only_configuration.md` | | 011 | DONE | 2026-03-08 | Migrate Redis repository client to aioredis | `spec/tasks/011_migrate_redis_repository_to_aioredis.md` | +| 012 | DONE | 2026-03-08 | Remove observability and SigNoz stack for phase 1 | `spec/tasks/012_remove_observability_and_signoz_for_phase1.md` | ## Summary -- Total: **12** +- Total: **13** - TODO: **0** -- DONE: **12** +- DONE: **13** diff --git a/spec/overview.md b/spec/overview.md index 42b8728..b25169e 100644 --- a/spec/overview.md +++ b/spec/overview.md @@ -20,8 +20,6 @@ - Возвращать унифицированный список тарифов, отсортированных по цене - Если провайдер вернул ошибку или не ответил вовремя — исключить его из результата, не падая целиком - Кешировать ответы провайдеров для исключения повторных внешних запросов. Время кэширования вынести в конфиг -- Обеспечивать структурированную наблюдаемость: трейсы, метрики, логи — связанные по request ID -- Отправлять алерты в Telegram при аномалиях (ошибки, всплески задержки, недоступность провайдера). token вынести в конфиг --- @@ -113,27 +111,6 @@ delivery_days_max: int | Кеш | Redis | | Конфигурация | pydantic-settings | | Сервер | Uvicorn | -| Инструментация | OpenTelemetry SDK | -| Бэкенд наблюдаемости | SigNoz (self-hosted) | -| Структурные логи | structlog | -| Алерты | SigNoz Alert Rules → Webhook → Telegram Bot API | - ---- - -## Наблюдаемость - -- Каждому запросу присваивается `request_id` (UUID) через middleware -- `request_id` привязывается ко всем логам через `structlog.contextvars` -- Трейсы отправляются через OpenTelemetry; FastAPI и httpx инструментируются автоматически -- Ручные спаны оборачивают: `AggregatorService.get_all_prices`, `get_price` каждого провайдера, обращения к кешу -- Атрибуты спанов: `provider`, `from_city`, `to_city`, `weight_kg`, `cache_hit`, `tariffs_found` -- Отслеживаемые метрики: количество запросов, количество ошибок, время ответа (p50/p99), доступность каждого провайдера -- Все три сигнала (трейсы, метрики, логи) связаны через `trace_id` - -### Условия алертов (Telegram) -- Провайдер вернул 5xx — порог: 5 ошибок за 5 минут -- p99 времени ответа провайдера > 5000ms в течение 10 минут -- Провайдер недоступен более 5 минут --- @@ -143,7 +120,6 @@ delivery_days_max: int app/ ├── controllers/ │ ├── http_client.py -│ ├── middleware.py │ └── v1/ │ └── delivery.py # Controller ├── services/ @@ -174,5 +150,4 @@ app/ # docker-compose сервисы app # FastAPI-приложение redis # Кеш тарифов -signoz # Бэкенд наблюдаемости (трейсы + метрики + логи + алерты) -``` \ No newline at end of file +``` diff --git a/spec/tasks/012_remove_observability_and_signoz_for_phase1.md b/spec/tasks/012_remove_observability_and_signoz_for_phase1.md new file mode 100644 index 0000000..5553cc0 --- /dev/null +++ b/spec/tasks/012_remove_observability_and_signoz_for_phase1.md @@ -0,0 +1,47 @@ +--- +id: 012 +title: Remove observability and SigNoz stack for phase 1 +status: DONE +created: 2026-03-08 +--- + +## Context +На первом этапе принято решение отказаться от логирования, мониторинга, алертинга и инфраструктуры SigNoz. Текущая кодовая база и тесты содержат эти зависимости и сценарии. + +## Goal +Удалить из проекта все runtime- и test-артефакты, связанные с логами, мониторингом, алертами, Telegram alerting и SigNoz, сохранив рабочий API расчета доставки и кеширование. + +## Constraints +- Соблюдать layered architecture из `AGENTS.md`; не переносить бизнес-правила между слоями. +- Scope задачи ограничен удалением observability/alerting/SigNoz и зависимых конфигураций, инфраструктурных и тестовых артефактов. +- Не изменять бизнес-логику расчета тарифов, provider integration и API-контракт `POST /api/v1/delivery/price`. +- Не добавлять новую систему мониторинга или алертинга в рамках этой задачи. +- Не изменять файлы в `spec/`. + +## Acceptance criteria +- Из runtime-приложения удалены middleware, instrumentation и иные механизмы, реализующие request/log correlation, tracing, metrics и alerting. +- В коде и конфигурации приложения отсутствуют секции и параметры `observability` и `alerts`, связанные с OpenTelemetry, structlog, SigNoz и Telegram. +- Из инфраструктурных файлов удалены сервис/настройки SigNoz и маршрутизация алертов. +- Удалены или обновлены тесты observability/alerts; тестовый набор для контроллера, сервиса, репозитория и конфигурации остается зеленым. +- По файлам приложения и инфраструктуры нет упоминаний `signoz`, `opentelemetry`, `structlog`, `telegram`. + +## Definition of Done +- [ ] Удалены runtime-компоненты observability/alerting. +- [ ] Обновлены YAML-конфиги и config schemas без секций observability/alerts. +- [ ] Обновлены инфраструктурные файлы без SigNoz. +- [ ] Удалены/обновлены тесты observability и alerts. +- [ ] Пройдены все команды из раздела Commands. + +## Tests +- Обновить `tests/config/test_config_sections.py` под конфигурацию без observability/alerts. +- Удалить или заменить `tests/config/test_alerts_config.py` и `tests/observability/*` в соответствии с новым scope. +- Обновить `tests/smoke/test_local_infra_stack.py` под инфраструктуру без SigNoz. +- Подтвердить, что `tests/controllers/v1/test_delivery.py`, `tests/services/test_aggregator.py`, `tests/repositories/cache/test_redis_cache.py` проходят без observability-зависимостей. + +## Commands +- `poetry run pytest tests/controllers/v1/test_delivery.py -q` +- `poetry run pytest tests/services/test_aggregator.py -q` +- `poetry run pytest tests/repositories/cache/test_redis_cache.py -q` +- `poetry run pytest tests/config/test_config_sections.py -q` +- `poetry run pytest tests/smoke/test_local_infra_stack.py -q` +- `! rg -n "(signoz|opentelemetry|structlog|telegram)" app tests docker-compose.yml pyproject.toml config.yaml config.example.yaml config.test.yaml infra` diff --git a/tests/config/fixtures/config.default.yaml b/tests/config/fixtures/config.default.yaml index 392ab4b..97478e6 100644 --- a/tests/config/fixtures/config.default.yaml +++ b/tests/config/fixtures/config.default.yaml @@ -13,9 +13,3 @@ repository: adapter: cdek_client_id: "yaml-id" cdek_client_secret: "yaml-secret" - -observability: - service_name: "from-config-yaml" - -alerts: - telegram_enabled: false diff --git a/tests/config/fixtures/config.missing_alerts.yaml b/tests/config/fixtures/config.missing_adapter.yaml similarity index 51% rename from tests/config/fixtures/config.missing_alerts.yaml rename to tests/config/fixtures/config.missing_adapter.yaml index 27232fd..11fd79c 100644 --- a/tests/config/fixtures/config.missing_alerts.yaml +++ b/tests/config/fixtures/config.missing_adapter.yaml @@ -10,13 +10,3 @@ business_logic: repository: redis_dsn: "redis://localhost:6379/0" price_cache_ttl_seconds: 900 - -adapter: - cdek_base_url: "https://api.cdek.ru/v2" - cdek_client_id: "id" - cdek_client_secret: "secret" - -observability: - service_name: "g2s-aggregator" - otlp_endpoint: "" - log_level: "INFO" diff --git a/tests/config/fixtures/config.test.override.yaml b/tests/config/fixtures/config.test.override.yaml index bf4e8be..8243cf2 100644 --- a/tests/config/fixtures/config.test.override.yaml +++ b/tests/config/fixtures/config.test.override.yaml @@ -13,9 +13,3 @@ repository: adapter: cdek_client_id: "test-id" cdek_client_secret: "test-secret" - -observability: - service_name: "from-config-test-yaml" - -alerts: - telegram_enabled: false diff --git a/tests/config/test_alerts_config.py b/tests/config/test_alerts_config.py deleted file mode 100644 index 2cf8aac..0000000 --- a/tests/config/test_alerts_config.py +++ /dev/null @@ -1,93 +0,0 @@ -import pytest -from pydantic import ValidationError - -from app.config import AlertsConfig, Settings - - -def _clear_alert_env(monkeypatch: pytest.MonkeyPatch) -> None: - for name in ( - "G2S_ALERTS__TELEGRAM_ENABLED", - "G2S_ALERTS__TELEGRAM_BOT_TOKEN", - "G2S_ALERTS__TELEGRAM_CHAT_ID", - ): - monkeypatch.delenv(name, raising=False) - - -def test_alerts_settings_are_loaded_from_yaml(tmp_path, monkeypatch) -> None: - _clear_alert_env(monkeypatch) - config_file = tmp_path / "config.yaml" - config_file.write_text( - """ -alerts: - telegram_enabled: true - telegram_bot_token: "bot-token" - telegram_chat_id: "-100777" - provider_5xx: - error_count: 7 - window_minutes: 6 - provider_p99_latency: - threshold_ms: 5200 - window_minutes: 11 - provider_unavailable: - duration_minutes: 8 -""".strip(), - encoding="utf-8", - ) - monkeypatch.setenv("G2S_CONFIG_FILE", str(config_file)) - - settings = Settings() - - assert settings.alerts.telegram_enabled is True - assert settings.alerts.telegram_bot_token == "bot-token" - assert settings.alerts.telegram_chat_id == "-100777" - assert settings.alerts.provider_5xx.error_count == 7 - assert settings.alerts.provider_5xx.window_minutes == 6 - assert settings.alerts.provider_p99_latency.threshold_ms == 5200 - assert settings.alerts.provider_p99_latency.window_minutes == 11 - assert settings.alerts.provider_unavailable.duration_minutes == 8 - - -@pytest.mark.parametrize( - ("yaml_body", "error_message"), - [ - ( - """ -alerts: - telegram_enabled: true - telegram_chat_id: "-100777" -""".strip(), - "alerts.telegram_bot_token is required when alerts.telegram_enabled=true", - ), - ( - """ -alerts: - telegram_enabled: true - telegram_bot_token: "bot-token" -""".strip(), - "alerts.telegram_chat_id is required when alerts.telegram_enabled=true", - ), - ], -) -def test_enabled_telegram_requires_credentials( - tmp_path, - monkeypatch, - yaml_body: str, - error_message: str, -) -> None: - _clear_alert_env(monkeypatch) - config_file = tmp_path / "config.yaml" - config_file.write_text(yaml_body, encoding="utf-8") - monkeypatch.setenv("G2S_CONFIG_FILE", str(config_file)) - - with pytest.raises(ValidationError, match=error_message): - Settings() - - -def test_alert_condition_defaults_match_required_thresholds() -> None: - alerts = AlertsConfig() - - assert alerts.provider_5xx.error_count == 5 - assert alerts.provider_5xx.window_minutes == 5 - assert alerts.provider_p99_latency.threshold_ms == 5000 - assert alerts.provider_p99_latency.window_minutes == 10 - assert alerts.provider_unavailable.duration_minutes == 5 diff --git a/tests/config/test_config_sections.py b/tests/config/test_config_sections.py index c0c7177..90d0a2b 100644 --- a/tests/config/test_config_sections.py +++ b/tests/config/test_config_sections.py @@ -8,8 +8,8 @@ from app.config import Settings, get_settings PROJECT_ROOT = Path(__file__).resolve().parents[2] TEST_CONFIG_FILE = PROJECT_ROOT / "config.test.yaml" -MISSING_ALERTS_CONFIG_FILE = ( - PROJECT_ROOT / "tests" / "config" / "fixtures" / "config.missing_alerts.yaml" +MISSING_REQUIRED_SECTION_CONFIG_FILE = ( + PROJECT_ROOT / "tests" / "config" / "fixtures" / "config.missing_adapter.yaml" ) DEFAULT_CONFIG_FILE_FIXTURE = ( PROJECT_ROOT / "tests" / "config" / "fixtures" / "config.default.yaml" @@ -56,24 +56,21 @@ def test_configuration_sections_are_loaded_from_yaml_file( assert settings.adapter.cdek_retry_backoff_seconds == 0.2 assert settings.adapter.cdek_timeout_seconds == 10.0 assert settings.adapter.cdek_cache_ttl_seconds == 900 - assert settings.observability.service_name == "g2s-aggregator-tests" - assert settings.observability.otlp_endpoint == "" - assert settings.observability.log_level == "INFO" - assert settings.alerts.telegram_enabled is False - assert settings.alerts.telegram_bot_token is None - assert settings.alerts.telegram_chat_id is None def test_get_settings_fails_when_required_yaml_section_is_missing( monkeypatch: pytest.MonkeyPatch, ) -> None: - _use_runtime_config_files(monkeypatch, test_config_file=MISSING_ALERTS_CONFIG_FILE) + _use_runtime_config_files( + monkeypatch, + test_config_file=MISSING_REQUIRED_SECTION_CONFIG_FILE, + ) with pytest.raises(ValidationError) as error: get_settings() locations = {tuple(item["loc"]) for item in error.value.errors()} - assert ("alerts",) in locations + assert ("adapter",) in locations get_settings.cache_clear() @@ -101,5 +98,5 @@ def test_get_settings_uses_config_test_yaml_in_pytest_environment( assert settings.service.provider_timeout_seconds == 17.0 assert settings.controller.api_prefix == "/from-config-test-yaml" - assert settings.observability.service_name == "from-config-test-yaml" + assert settings.adapter.cdek_client_id == "test-id" get_settings.cache_clear() diff --git a/tests/controllers/test_middleware_request_id.py b/tests/controllers/test_middleware_request_id.py deleted file mode 100644 index c2624ee..0000000 --- a/tests/controllers/test_middleware_request_id.py +++ /dev/null @@ -1,101 +0,0 @@ -import asyncio -from uuid import UUID - -import httpx - -from app.config import Settings -from app.controllers.v1.delivery import get_aggregator_service -from app.main import create_app -from app.schemas.request import DeliveryRequest - - -class StubAggregatorService: - def __init__(self) -> None: - self.calls: list[DeliveryRequest] = [] - - async def get_all_prices(self, request: DeliveryRequest) -> list[object]: - self.calls.append(request) - return [] - - -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_app( - *, - request_id_header: str, - service: StubAggregatorService, -) -> tuple[object, StubAggregatorService]: - settings = Settings( - controller={"request_id_header": request_id_header}, - observability={ - "service_name": "g2s-tests", - "otlp_endpoint": "", - "log_level": "INFO", - }, - ) - app = create_app(settings=settings) - - async def override_service() -> StubAggregatorService: - return service - - app.dependency_overrides[get_aggregator_service] = override_service - return app, service - - -def test_request_id_is_generated_per_request_and_returned_in_header() -> None: - app, service = _build_app( - request_id_header="X-Request-ID", - service=StubAggregatorService(), - ) - - async def run_requests() -> tuple[httpx.Response, httpx.Response]: - transport = httpx.ASGITransport(app=app) - async with httpx.AsyncClient( - transport=transport, - base_url="http://testserver", - ) as client: - first = await client.post("/api/v1/delivery/price", json=_build_payload()) - second = await client.post("/api/v1/delivery/price", json=_build_payload()) - return first, second - - first_response, second_response = asyncio.run(run_requests()) - - assert first_response.status_code == 200 - assert second_response.status_code == 200 - first_request_id = first_response.headers["X-Request-ID"] - second_request_id = second_response.headers["X-Request-ID"] - assert UUID(first_request_id).version == 4 - assert UUID(second_request_id).version == 4 - assert first_request_id != second_request_id - assert len(service.calls) == 2 - - -def test_request_id_uses_configured_header_name() -> None: - app, _ = _build_app( - request_id_header="X-Correlation-ID", - service=StubAggregatorService(), - ) - - 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.post("/api/v1/delivery/price", json=_build_payload()) - - response = asyncio.run(run_request()) - - assert response.status_code == 200 - assert "X-Correlation-ID" in response.headers - assert "X-Request-ID" not in response.headers diff --git a/tests/observability/test_logging_context.py b/tests/observability/test_logging_context.py deleted file mode 100644 index d38ba53..0000000 --- a/tests/observability/test_logging_context.py +++ /dev/null @@ -1,50 +0,0 @@ -import structlog -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 structlog.testing import LogCapture -from structlog.contextvars import bind_contextvars, clear_contextvars - -from app.observability.logging import add_trace_context -from app.observability.tracing import get_tracer, reset_tracing_for_tests, setup_tracing - - -def test_logging_context_includes_request_id_and_trace_id() -> 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, - ) - - log_capture = LogCapture() - structlog.configure( - processors=[ - structlog.contextvars.merge_contextvars, - add_trace_context, - log_capture, - ], - logger_factory=structlog.stdlib.LoggerFactory(), - wrapper_class=structlog.make_filtering_bound_logger(20), - cache_logger_on_first_use=False, - ) - logger = structlog.get_logger("test-logger") - - clear_contextvars() - bind_contextvars(request_id="req-test-id") - with get_tracer(__name__).start_as_current_span("test-span"): - logger.info("request_log") - - assert len(log_capture.entries) == 1 - entry = log_capture.entries[0] - assert entry["event"] == "request_log" - assert entry["request_id"] == "req-test-id" - assert len(entry["trace_id"]) == 32 - assert len(entry["span_id"]) == 16 - - clear_contextvars() - reset_tracing_for_tests() diff --git a/tests/observability/test_metrics.py b/tests/observability/test_metrics.py deleted file mode 100644 index 159a262..0000000 --- a/tests/observability/test_metrics.py +++ /dev/null @@ -1,204 +0,0 @@ -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() diff --git a/tests/observability/test_tracing.py b/tests/observability/test_tracing.py deleted file mode 100644 index c79cd3d..0000000 --- a/tests/observability/test_tracing.py +++ /dev/null @@ -1,160 +0,0 @@ -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() diff --git a/tests/smoke/test_local_infra_stack.py b/tests/smoke/test_local_infra_stack.py index b371997..2bf161e 100644 --- a/tests/smoke/test_local_infra_stack.py +++ b/tests/smoke/test_local_infra_stack.py @@ -18,12 +18,12 @@ def test_compose_defines_required_services() -> None: compose = _load_compose() services = compose.get("services") assert isinstance(services, dict) - assert {"app", "redis", "signoz"}.issubset(services.keys()) + assert {"app", "redis"}.issubset(services.keys()) def test_smoke_command_sequence_is_documented() -> None: readme = INFRA_README.read_text(encoding="utf-8") assert "docker compose config" in readme - assert "docker compose up -d redis signoz" in readme + assert "docker compose up -d redis" in readme assert "docker compose ps" in readme