003 Add cdek
This commit is contained in:
@@ -1 +1,5 @@
|
||||
"""CDEK adapter package."""
|
||||
|
||||
from app.adapters.delivery_providers.cdek.client import CDEKProvider
|
||||
|
||||
__all__ = ["CDEKProvider"]
|
||||
|
||||
@@ -1,6 +1,97 @@
|
||||
"""CDEK auth adapter skeleton."""
|
||||
"""CDEK OAuth2 authentication adapter."""
|
||||
|
||||
import asyncio
|
||||
import time
|
||||
from collections.abc import Callable
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
|
||||
class CDEKAuthError(RuntimeError):
|
||||
"""Raised when CDEK OAuth2 flow fails."""
|
||||
|
||||
|
||||
class CDEKAuthClient:
|
||||
def __init__(
|
||||
self,
|
||||
http_client: httpx.AsyncClient,
|
||||
*,
|
||||
base_url: str,
|
||||
client_id: str,
|
||||
client_secret: str,
|
||||
timeout_seconds: float = 10.0,
|
||||
expiry_skew_seconds: float = 30.0,
|
||||
clock: Callable[[], float] | None = None,
|
||||
) -> None:
|
||||
self._http_client = http_client
|
||||
self._token_url = f"{base_url.rstrip('/')}/oauth/token"
|
||||
self._client_id = client_id
|
||||
self._client_secret = client_secret
|
||||
self._timeout_seconds = timeout_seconds
|
||||
self._expiry_skew_seconds = expiry_skew_seconds
|
||||
self._clock = clock or time.monotonic
|
||||
self._token: str | None = None
|
||||
self._token_expires_at: float = 0.0
|
||||
self._refresh_lock = asyncio.Lock()
|
||||
|
||||
async def get_access_token(self) -> str:
|
||||
raise NotImplementedError("CDEK auth implementation is added in task 003.")
|
||||
if self._has_valid_token():
|
||||
return self._cached_token()
|
||||
|
||||
async with self._refresh_lock:
|
||||
if self._has_valid_token():
|
||||
return self._cached_token()
|
||||
return await self._fetch_access_token()
|
||||
|
||||
def _has_valid_token(self) -> bool:
|
||||
return self._token is not None and self._clock() < self._token_expires_at
|
||||
|
||||
def _cached_token(self) -> str:
|
||||
if self._token is None:
|
||||
raise CDEKAuthError("CDEK OAuth token is unexpectedly missing.")
|
||||
return self._token
|
||||
|
||||
async def _fetch_access_token(self) -> str:
|
||||
try:
|
||||
response = await self._http_client.post(
|
||||
self._token_url,
|
||||
data={
|
||||
"grant_type": "client_credentials",
|
||||
"client_id": self._client_id,
|
||||
"client_secret": self._client_secret,
|
||||
},
|
||||
timeout=self._timeout_seconds,
|
||||
)
|
||||
response.raise_for_status()
|
||||
payload = response.json()
|
||||
token = self._extract_access_token(payload)
|
||||
expires_in = self._extract_expires_in(payload)
|
||||
except (httpx.HTTPError, TypeError, ValueError) as exc:
|
||||
raise CDEKAuthError("Failed to obtain CDEK OAuth token.") from exc
|
||||
|
||||
self._token = token
|
||||
ttl_seconds = max(expires_in - self._expiry_skew_seconds, 0.0)
|
||||
self._token_expires_at = self._clock() + ttl_seconds
|
||||
return self._token
|
||||
|
||||
@staticmethod
|
||||
def _extract_access_token(payload: Any) -> str:
|
||||
if not isinstance(payload, dict):
|
||||
raise ValueError("CDEK OAuth payload must be a JSON object.")
|
||||
token = payload.get("access_token")
|
||||
if not isinstance(token, str) or not token.strip():
|
||||
raise ValueError("CDEK OAuth payload has no valid access_token.")
|
||||
return token
|
||||
|
||||
@staticmethod
|
||||
def _extract_expires_in(payload: Any) -> float:
|
||||
if not isinstance(payload, dict):
|
||||
raise ValueError("CDEK OAuth payload must be a JSON object.")
|
||||
expires_in = payload.get("expires_in")
|
||||
if expires_in is None:
|
||||
raise ValueError("CDEK OAuth payload has no expires_in.")
|
||||
parsed_expires_in = float(expires_in)
|
||||
if parsed_expires_in < 0:
|
||||
raise ValueError("CDEK OAuth expires_in cannot be negative.")
|
||||
return parsed_expires_in
|
||||
|
||||
@@ -1,9 +1,195 @@
|
||||
"""CDEK HTTP client adapter skeleton."""
|
||||
"""CDEK HTTP client and provider adapter."""
|
||||
|
||||
import asyncio
|
||||
from collections.abc import Awaitable, Callable
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app.adapters.delivery_providers.base import DeliveryProvider
|
||||
from app.adapters.delivery_providers.cdek.auth import CDEKAuthClient
|
||||
from app.adapters.delivery_providers.cdek.mapper import map_cdek_response
|
||||
from app.config import AdapterConfig
|
||||
from app.schemas.request import DeliveryRequest
|
||||
from app.schemas.response import DeliveryPrice
|
||||
|
||||
|
||||
class CDEKClientError(RuntimeError):
|
||||
"""Raised when CDEK tariff request fails."""
|
||||
|
||||
|
||||
class CDEKClient:
|
||||
async def get_raw_price(self, request: DeliveryRequest) -> dict:
|
||||
_ = request
|
||||
raise NotImplementedError("CDEK client implementation is added in task 003.")
|
||||
def __init__(
|
||||
self,
|
||||
http_client: httpx.AsyncClient,
|
||||
auth_client: CDEKAuthClient,
|
||||
*,
|
||||
base_url: str,
|
||||
timeout_seconds: float = 10.0,
|
||||
retry_attempts: int = 2,
|
||||
retry_backoff_seconds: float = 0.2,
|
||||
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
|
||||
) -> None:
|
||||
self._http_client = http_client
|
||||
self._auth_client = auth_client
|
||||
normalized_base_url = base_url.rstrip("/")
|
||||
self._city_lookup_url = f"{normalized_base_url}/location/cities"
|
||||
self._tariff_url = f"{base_url.rstrip('/')}/calculator/tarifflist"
|
||||
self._timeout_seconds = timeout_seconds
|
||||
self._retry_attempts = retry_attempts
|
||||
self._retry_backoff_seconds = retry_backoff_seconds
|
||||
self._sleep = sleep
|
||||
self._city_code_cache: dict[str, int] = {}
|
||||
|
||||
async def get_raw_price(self, request: DeliveryRequest) -> dict[str, Any]:
|
||||
payload = await self._build_payload(request)
|
||||
for attempt in range(self._retry_attempts + 1):
|
||||
try:
|
||||
token = await self._auth_client.get_access_token()
|
||||
response = await self._http_client.post(
|
||||
self._tariff_url,
|
||||
json=payload,
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=self._timeout_seconds,
|
||||
)
|
||||
except (httpx.TimeoutException, httpx.TransportError) as exc:
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
"CDEK tariff request failed after retry attempts."
|
||||
) from exc
|
||||
|
||||
if self._should_retry(response.status_code):
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
f"CDEK tariff request failed with status {response.status_code}."
|
||||
)
|
||||
|
||||
try:
|
||||
response.raise_for_status()
|
||||
raw_payload = response.json()
|
||||
except (httpx.HTTPError, TypeError, ValueError) as exc:
|
||||
raise CDEKClientError("CDEK tariff request returned invalid payload.") from exc
|
||||
|
||||
if not isinstance(raw_payload, dict):
|
||||
raise CDEKClientError("CDEK tariff payload must be a JSON object.")
|
||||
return raw_payload
|
||||
|
||||
raise CDEKClientError("CDEK tariff request failed unexpectedly.")
|
||||
|
||||
def _retry_delay(self, attempt: int) -> float:
|
||||
return self._retry_backoff_seconds * (attempt + 1)
|
||||
|
||||
@staticmethod
|
||||
def _should_retry(status_code: int) -> bool:
|
||||
return status_code == 429 or status_code >= 500
|
||||
|
||||
async def _build_payload(self, request: DeliveryRequest) -> dict[str, Any]:
|
||||
from_city_code = await self._resolve_city_code(request.from_city)
|
||||
to_city_code = await self._resolve_city_code(request.to_city)
|
||||
return {
|
||||
"from_location": {"code": from_city_code},
|
||||
"to_location": {"code": to_city_code},
|
||||
"packages": [
|
||||
{
|
||||
"weight": int(round(request.weight_kg * 1000)),
|
||||
"length": int(round(request.length_cm)),
|
||||
"width": int(round(request.width_cm)),
|
||||
"height": int(round(request.height_cm)),
|
||||
}
|
||||
],
|
||||
}
|
||||
|
||||
async def _resolve_city_code(self, city: str) -> int:
|
||||
normalized_city = city.strip().casefold()
|
||||
cached_code = self._city_code_cache.get(normalized_city)
|
||||
if cached_code is not None:
|
||||
return cached_code
|
||||
|
||||
for attempt in range(self._retry_attempts + 1):
|
||||
try:
|
||||
token = await self._auth_client.get_access_token()
|
||||
response = await self._http_client.get(
|
||||
self._city_lookup_url,
|
||||
params={"city": city, "country_codes": "RU", "size": 1},
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=self._timeout_seconds,
|
||||
)
|
||||
except (httpx.TimeoutException, httpx.TransportError) as exc:
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
"CDEK city lookup request failed after retry attempts."
|
||||
) from exc
|
||||
|
||||
if self._should_retry(response.status_code):
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
f"CDEK city lookup failed with status {response.status_code}."
|
||||
)
|
||||
try:
|
||||
response.raise_for_status()
|
||||
body = response.json()
|
||||
except (httpx.HTTPError, TypeError, ValueError) as exc:
|
||||
raise CDEKClientError("CDEK city lookup returned invalid payload.") from exc
|
||||
break
|
||||
else:
|
||||
raise CDEKClientError("CDEK city lookup failed unexpectedly.")
|
||||
|
||||
if not isinstance(body, list) or not body:
|
||||
raise CDEKClientError(f"CDEK city lookup returned no matches for '{city}'.")
|
||||
first_item = body[0]
|
||||
if not isinstance(first_item, dict):
|
||||
raise CDEKClientError("CDEK city lookup response entry must be an object.")
|
||||
raw_city_code = first_item.get("code")
|
||||
if raw_city_code is None:
|
||||
raise CDEKClientError("CDEK city lookup response has no city code.")
|
||||
try:
|
||||
city_code = int(raw_city_code)
|
||||
except (TypeError, ValueError) as exc:
|
||||
raise CDEKClientError("CDEK city lookup response city code is invalid.") from exc
|
||||
|
||||
self._city_code_cache[normalized_city] = city_code
|
||||
return city_code
|
||||
|
||||
|
||||
class CDEKProvider(DeliveryProvider):
|
||||
name = "cdek"
|
||||
|
||||
def __init__(self, client: CDEKClient, *, cache_ttl_seconds: int = 900) -> None:
|
||||
self._client = client
|
||||
self.cache_ttl_seconds = cache_ttl_seconds
|
||||
|
||||
@classmethod
|
||||
def from_adapter_config(
|
||||
cls,
|
||||
*,
|
||||
http_client: httpx.AsyncClient,
|
||||
adapter_config: AdapterConfig,
|
||||
) -> "CDEKProvider":
|
||||
auth_client = CDEKAuthClient(
|
||||
http_client=http_client,
|
||||
base_url=adapter_config.cdek_base_url,
|
||||
client_id=adapter_config.cdek_client_id,
|
||||
client_secret=adapter_config.cdek_client_secret,
|
||||
timeout_seconds=adapter_config.cdek_timeout_seconds,
|
||||
)
|
||||
client = CDEKClient(
|
||||
http_client=http_client,
|
||||
auth_client=auth_client,
|
||||
base_url=adapter_config.cdek_base_url,
|
||||
timeout_seconds=adapter_config.cdek_timeout_seconds,
|
||||
retry_attempts=adapter_config.cdek_retry_attempts,
|
||||
retry_backoff_seconds=adapter_config.cdek_retry_backoff_seconds,
|
||||
)
|
||||
return cls(client=client, cache_ttl_seconds=adapter_config.cdek_cache_ttl_seconds)
|
||||
|
||||
async def get_price(self, request: DeliveryRequest) -> DeliveryPrice:
|
||||
raw_payload = await self._client.get_raw_price(request)
|
||||
return map_cdek_response(raw_payload)
|
||||
|
||||
@@ -1,8 +1,51 @@
|
||||
"""CDEK response mapper skeleton."""
|
||||
"""CDEK response mapper to unified delivery schema."""
|
||||
|
||||
from decimal import Decimal
|
||||
from typing import Any
|
||||
|
||||
from app.schemas.response import DeliveryPrice
|
||||
|
||||
|
||||
def map_cdek_response(payload: dict) -> DeliveryPrice:
|
||||
_ = payload
|
||||
raise NotImplementedError("CDEK mapper implementation is added in task 003.")
|
||||
class CDEKMappingError(ValueError):
|
||||
"""Raised when CDEK response cannot be mapped."""
|
||||
|
||||
|
||||
def map_cdek_response(payload: dict[str, Any]) -> DeliveryPrice:
|
||||
tariff = _get_first_tariff(payload)
|
||||
|
||||
service_name = tariff.get("tariff_name") or tariff.get("tariff_code")
|
||||
if service_name is None:
|
||||
raise CDEKMappingError("CDEK tariff_name is missing.")
|
||||
|
||||
raw_price = tariff.get("delivery_sum")
|
||||
if raw_price is None:
|
||||
raise CDEKMappingError("CDEK delivery_sum is missing.")
|
||||
|
||||
period_min = tariff.get("period_min")
|
||||
if period_min is None:
|
||||
raise CDEKMappingError("CDEK period_min is missing.")
|
||||
period_max = tariff.get("period_max", period_min)
|
||||
|
||||
raw_currency = tariff.get("currency") or payload.get("currency") or "RUB"
|
||||
|
||||
try:
|
||||
return DeliveryPrice(
|
||||
provider="cdek",
|
||||
service_name=str(service_name),
|
||||
price=Decimal(str(raw_price)),
|
||||
currency=str(raw_currency).upper(),
|
||||
delivery_days_min=int(period_min),
|
||||
delivery_days_max=int(period_max),
|
||||
)
|
||||
except (ArithmeticError, TypeError, ValueError) as exc:
|
||||
raise CDEKMappingError("CDEK response fields have invalid values.") from exc
|
||||
|
||||
|
||||
def _get_first_tariff(payload: dict[str, Any]) -> dict[str, Any]:
|
||||
tariff_codes = payload.get("tariff_codes")
|
||||
if not isinstance(tariff_codes, list) or not tariff_codes:
|
||||
raise CDEKMappingError("CDEK response must include non-empty tariff_codes.")
|
||||
tariff = tariff_codes[0]
|
||||
if not isinstance(tariff, dict):
|
||||
raise CDEKMappingError("CDEK tariff entry must be an object.")
|
||||
return tariff
|
||||
|
||||
+6
-1
@@ -36,7 +36,12 @@ class RepositoryConfig(BaseModel):
|
||||
|
||||
class AdapterConfig(BaseModel):
|
||||
cdek_base_url: str = "https://api.cdek.ru/v2"
|
||||
cdek_timeout_seconds: float = 10.0
|
||||
cdek_client_id: str = ""
|
||||
cdek_client_secret: str = ""
|
||||
cdek_retry_attempts: int = Field(default=2, ge=0)
|
||||
cdek_retry_backoff_seconds: float = Field(default=0.2, ge=0)
|
||||
cdek_timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
cdek_cache_ttl_seconds: int = Field(default=900, gt=0)
|
||||
|
||||
|
||||
class ObservabilityConfig(BaseModel):
|
||||
|
||||
Reference in New Issue
Block a user