@@ -18,6 +18,7 @@ from app.config import Settings
|
||||
from app.controllers.http_client import build_controller_http_client
|
||||
from app.repositories.cache.redis_cache import PriceCache
|
||||
from app.repositories.order import OrderRepository
|
||||
from app.runtime.metrics import get_cache_metrics
|
||||
from app.schemas.payment import (
|
||||
InitPaymentRequest,
|
||||
InitPaymentResponse,
|
||||
@@ -75,7 +76,10 @@ def _build_aggregator_service(settings: Settings) -> AggregatorService:
|
||||
config=settings.address_suggestions.tomtom,
|
||||
)
|
||||
providers = (cdek_provider, cse_provider)
|
||||
cache = PriceCache.from_repository_config(settings.repository)
|
||||
cache = PriceCache.from_repository_config(
|
||||
settings.repository,
|
||||
metrics=get_cache_metrics(),
|
||||
)
|
||||
postgres_engine = create_postgres_engine(settings.postgres)
|
||||
postgres_session_factory = create_postgres_session_factory(postgres_engine)
|
||||
order_repository = OrderRepository(session_factory=postgres_session_factory)
|
||||
|
||||
@@ -5,12 +5,14 @@ from fastapi import FastAPI
|
||||
from app.config import Settings, get_settings
|
||||
from app.controllers.v1.delivery import router as delivery_router
|
||||
from app.runtime.logging import configure_logging
|
||||
from app.runtime.metrics import configure_metrics
|
||||
from app.runtime.tracing import configure_tracing
|
||||
|
||||
|
||||
def create_app(settings: Settings | None = None) -> FastAPI:
|
||||
configure_logging()
|
||||
resolved_settings = settings or get_settings()
|
||||
configure_metrics(resolved_settings.observability)
|
||||
|
||||
application = FastAPI(title="G2S Aggregator", version="0.1.0")
|
||||
application.state.settings = resolved_settings
|
||||
|
||||
+52
-2
@@ -7,10 +7,14 @@ import json
|
||||
from typing import Any, Callable, Protocol
|
||||
|
||||
from pydantic import BaseModel
|
||||
import structlog
|
||||
|
||||
from app.config import RepositoryConfig
|
||||
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
class RedisClientProtocol(Protocol):
|
||||
async def get(self, key: str) -> bytes | str | None: ...
|
||||
|
||||
@@ -19,6 +23,14 @@ class RedisClientProtocol(Protocol):
|
||||
async def delete(self, key: str) -> Any: ...
|
||||
|
||||
|
||||
class CacheMetricsProtocol(Protocol):
|
||||
def record_request(self, *, operation: str, outcome: str) -> None: ...
|
||||
|
||||
def record_error(self, *, operation: str, error_type: str) -> None: ...
|
||||
|
||||
def record_hit(self, *, hit: bool) -> None: ...
|
||||
|
||||
|
||||
class PriceCacheRepositoryError(RuntimeError):
|
||||
"""Deterministic repository error raised on Redis operation failures."""
|
||||
|
||||
@@ -26,9 +38,16 @@ class PriceCacheRepositoryError(RuntimeError):
|
||||
class PriceCache:
|
||||
"""Repository for cached provider payloads."""
|
||||
|
||||
def __init__(self, redis_client: RedisClientProtocol, *, ttl_seconds: int) -> None:
|
||||
def __init__(
|
||||
self,
|
||||
redis_client: RedisClientProtocol,
|
||||
*,
|
||||
ttl_seconds: int,
|
||||
metrics: CacheMetricsProtocol | None = None,
|
||||
) -> None:
|
||||
self._redis_client = redis_client
|
||||
self._ttl_seconds = ttl_seconds
|
||||
self._metrics = metrics
|
||||
|
||||
@classmethod
|
||||
def from_repository_config(
|
||||
@@ -36,20 +55,25 @@ class PriceCache:
|
||||
repository_config: RepositoryConfig,
|
||||
*,
|
||||
client_factory: Callable[[str], RedisClientProtocol] | None = None,
|
||||
metrics: CacheMetricsProtocol | None = None,
|
||||
) -> "PriceCache":
|
||||
resolved_factory = client_factory or _default_redis_client_factory
|
||||
redis_client = resolved_factory(repository_config.redis_dsn)
|
||||
return cls(
|
||||
redis_client=redis_client,
|
||||
ttl_seconds=repository_config.price_cache_ttl_seconds,
|
||||
metrics=metrics,
|
||||
)
|
||||
|
||||
async def get(self, key: str) -> Any | None:
|
||||
try:
|
||||
payload = await self._redis_client.get(key)
|
||||
except Exception as exc:
|
||||
self._record_failure(operation="get", error=exc)
|
||||
raise PriceCacheRepositoryError("Redis cache read failed.") from exc
|
||||
if payload is None:
|
||||
self._record_request(operation="get", outcome="miss")
|
||||
self._record_hit(hit=False)
|
||||
return None
|
||||
|
||||
if isinstance(payload, bytes):
|
||||
@@ -59,7 +83,10 @@ class PriceCache:
|
||||
else:
|
||||
raise TypeError("Redis cache payload must be bytes or str.")
|
||||
|
||||
return json.loads(serialized_payload)
|
||||
result = json.loads(serialized_payload)
|
||||
self._record_request(operation="get", outcome="hit")
|
||||
self._record_hit(hit=True)
|
||||
return result
|
||||
|
||||
async def set(self, key: str, value: Any, ttl: int | None = None) -> None:
|
||||
serialized_payload = _serialize(value)
|
||||
@@ -67,13 +94,36 @@ class PriceCache:
|
||||
try:
|
||||
await self._redis_client.set(key, serialized_payload, ex=resolved_ttl)
|
||||
except Exception as exc:
|
||||
self._record_failure(operation="set", error=exc)
|
||||
raise PriceCacheRepositoryError("Redis cache write failed.") from exc
|
||||
self._record_request(operation="set", outcome="success")
|
||||
|
||||
async def invalidate(self, key: str) -> None:
|
||||
try:
|
||||
await self._redis_client.delete(key)
|
||||
except Exception as exc:
|
||||
self._record_failure(operation="invalidate", error=exc)
|
||||
raise PriceCacheRepositoryError("Redis cache invalidate failed.") from exc
|
||||
self._record_request(operation="invalidate", outcome="success")
|
||||
|
||||
def _record_failure(self, *, operation: str, error: Exception) -> None:
|
||||
error_type = type(error).__name__
|
||||
logger.warning(
|
||||
"cache.operation_failed",
|
||||
operation=operation,
|
||||
error_type=error_type,
|
||||
)
|
||||
if self._metrics is not None:
|
||||
self._metrics.record_request(operation=operation, outcome="error")
|
||||
self._metrics.record_error(operation=operation, error_type=error_type)
|
||||
|
||||
def _record_request(self, *, operation: str, outcome: str) -> None:
|
||||
if self._metrics is not None:
|
||||
self._metrics.record_request(operation=operation, outcome=outcome)
|
||||
|
||||
def _record_hit(self, *, hit: bool) -> None:
|
||||
if self._metrics is not None:
|
||||
self._metrics.record_hit(hit=hit)
|
||||
|
||||
|
||||
def _default_redis_client_factory(redis_dsn: str) -> RedisClientProtocol:
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
"""OpenTelemetry metrics bootstrap and application instruments."""
|
||||
|
||||
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 PeriodicExportingMetricReader
|
||||
from opentelemetry.sdk.resources import Resource
|
||||
|
||||
from app.config import ObservabilityConfig
|
||||
|
||||
|
||||
_METER_PROVIDER: MeterProvider | None = None
|
||||
_CACHE_METRICS: "CacheMetrics | None" = None
|
||||
|
||||
|
||||
class CacheMetrics:
|
||||
"""Metrics emitted by the Redis-backed price cache."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
meter = metrics.get_meter("app.repositories.cache")
|
||||
self._requests = meter.create_counter(
|
||||
"cache_requests_total",
|
||||
unit="{request}",
|
||||
description="Number of cache operations by operation and outcome.",
|
||||
)
|
||||
self._errors = meter.create_counter(
|
||||
"cache_errors_total",
|
||||
unit="{error}",
|
||||
description="Number of failed cache operations by error type.",
|
||||
)
|
||||
self._hit_ratio = meter.create_histogram(
|
||||
"cache_hit_ratio",
|
||||
unit="1",
|
||||
description="Cache lookup result: 1 for a hit and 0 for a miss.",
|
||||
)
|
||||
|
||||
def record_request(self, *, operation: str, outcome: str) -> None:
|
||||
self._requests.add(1, {"operation": operation, "outcome": outcome})
|
||||
|
||||
def record_error(self, *, operation: str, error_type: str) -> None:
|
||||
self._errors.add(1, {"operation": operation, "error.type": error_type})
|
||||
|
||||
def record_hit(self, *, hit: bool) -> None:
|
||||
self._hit_ratio.record(1.0 if hit else 0.0)
|
||||
|
||||
|
||||
def configure_metrics(observability: ObservabilityConfig) -> None:
|
||||
"""Configure the process-wide OTLP metrics exporter once."""
|
||||
|
||||
global _METER_PROVIDER
|
||||
|
||||
if not observability.enabled or _METER_PROVIDER is not None:
|
||||
return
|
||||
|
||||
meter_provider = _build_meter_provider(observability)
|
||||
metrics.set_meter_provider(meter_provider)
|
||||
_METER_PROVIDER = meter_provider
|
||||
|
||||
|
||||
def get_cache_metrics() -> CacheMetrics:
|
||||
"""Return the process-wide cache instruments."""
|
||||
|
||||
global _CACHE_METRICS
|
||||
|
||||
if _CACHE_METRICS is None:
|
||||
_CACHE_METRICS = CacheMetrics()
|
||||
return _CACHE_METRICS
|
||||
|
||||
|
||||
def _build_meter_provider(observability: ObservabilityConfig) -> MeterProvider:
|
||||
exporter = OTLPMetricExporter(
|
||||
endpoint=observability.otlp_endpoint,
|
||||
insecure=observability.otlp_insecure,
|
||||
)
|
||||
reader = PeriodicExportingMetricReader(exporter)
|
||||
resource = Resource.create({"service.name": observability.service_name})
|
||||
return MeterProvider(resource=resource, metric_readers=[reader])
|
||||
Reference in New Issue
Block a user