@@ -1,7 +1,9 @@
|
||||
"""Base interface for delivery providers."""
|
||||
|
||||
from abc import ABC, abstractmethod
|
||||
from typing import Protocol, runtime_checkable
|
||||
|
||||
from app.schemas.payment import InitPaymentRequest
|
||||
from app.schemas.request import DeliveryCalculationRequest
|
||||
from app.schemas.response import DeliveryPrice
|
||||
|
||||
@@ -18,6 +20,27 @@ class DeliveryProvider(ABC):
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class PaymentPriceValidationProvider(Protocol):
|
||||
name: str
|
||||
|
||||
async def get_payment_price(
|
||||
self,
|
||||
request: InitPaymentRequest,
|
||||
) -> DeliveryPrice | None: ...
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class OrderRegistrationProvider(Protocol):
|
||||
name: str
|
||||
|
||||
async def register_order(
|
||||
self,
|
||||
request: InitPaymentRequest,
|
||||
order_uuid: str,
|
||||
) -> object: ...
|
||||
|
||||
|
||||
class ProviderClientError(RuntimeError):
|
||||
"""Raised when a provider call fails for temporary or provider-side reasons."""
|
||||
|
||||
|
||||
@@ -33,7 +33,7 @@ from app.adapters.delivery_providers.cdek.order_mapper import (
|
||||
resolve_cdek_city_code,
|
||||
)
|
||||
from app.cities import cities_map
|
||||
from app.config import AdapterConfig
|
||||
from app.config import CDEKDeliveryProviderConfig
|
||||
from app.schemas.payment import InitPaymentRequest
|
||||
from app.schemas.request import DeliveryCalculationRequest
|
||||
from app.schemas.response import DeliveryPrice
|
||||
@@ -452,28 +452,28 @@ class CDEKProvider(DeliveryProvider):
|
||||
self.cache_ttl_seconds = cache_ttl_seconds
|
||||
|
||||
@classmethod
|
||||
def from_adapter_config(
|
||||
def from_config(
|
||||
cls,
|
||||
*,
|
||||
http_client: httpx.AsyncClient,
|
||||
adapter_config: AdapterConfig,
|
||||
config: CDEKDeliveryProviderConfig,
|
||||
) -> "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,
|
||||
base_url=config.base_url,
|
||||
client_id=config.client_id,
|
||||
client_secret=config.client_secret,
|
||||
timeout_seconds=config.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,
|
||||
base_url=config.base_url,
|
||||
timeout_seconds=config.timeout_seconds,
|
||||
retry_attempts=config.retry_attempts,
|
||||
retry_backoff_seconds=config.retry_backoff_seconds,
|
||||
)
|
||||
return cls(client=client, cache_ttl_seconds=adapter_config.cdek_cache_ttl_seconds)
|
||||
return cls(client=client, cache_ttl_seconds=config.cache_ttl_seconds)
|
||||
|
||||
async def get_prices(
|
||||
self, request: DeliveryCalculationRequest
|
||||
|
||||
@@ -35,7 +35,7 @@ from app.adapters.delivery_providers.cse.soap import (
|
||||
make_field,
|
||||
parse_response,
|
||||
)
|
||||
from app.config import AdapterConfig
|
||||
from app.config import CSEDeliveryProviderConfig
|
||||
from app.schemas.payment import InitPaymentRequest
|
||||
from app.schemas.request import DeliveryCalculationRequest
|
||||
from app.schemas.response import DeliveryPrice
|
||||
@@ -212,27 +212,27 @@ class CSEProvider(DeliveryProvider):
|
||||
self._delivery_types: list[tuple[str, str]] | None = None
|
||||
|
||||
@classmethod
|
||||
def from_adapter_config(
|
||||
def from_config(
|
||||
cls,
|
||||
*,
|
||||
http_client: httpx.AsyncClient,
|
||||
adapter_config: AdapterConfig,
|
||||
config: CSEDeliveryProviderConfig,
|
||||
) -> "CSEProvider":
|
||||
client = CSEClient(
|
||||
http_client=http_client,
|
||||
base_url=adapter_config.cse_base_url,
|
||||
login=adapter_config.cse_login,
|
||||
password=adapter_config.cse_password,
|
||||
base_url=config.base_url,
|
||||
login=config.login,
|
||||
password=config.password,
|
||||
registration_params=CSEOrderRegistrationParams(
|
||||
payer=adapter_config.cse_payer,
|
||||
payment_method=adapter_config.cse_payment_method,
|
||||
shipping_method=adapter_config.cse_shipping_method,
|
||||
payer=config.payer,
|
||||
payment_method=config.payment_method,
|
||||
shipping_method=config.shipping_method,
|
||||
),
|
||||
timeout_seconds=adapter_config.cse_timeout_seconds,
|
||||
retry_attempts=adapter_config.cse_retry_attempts,
|
||||
retry_backoff_seconds=adapter_config.cse_retry_backoff_seconds,
|
||||
timeout_seconds=config.timeout_seconds,
|
||||
retry_attempts=config.retry_attempts,
|
||||
retry_backoff_seconds=config.retry_backoff_seconds,
|
||||
)
|
||||
return cls(client=client, cache_ttl_seconds=adapter_config.cse_cache_ttl_seconds)
|
||||
return cls(client=client, cache_ttl_seconds=config.cache_ttl_seconds)
|
||||
|
||||
async def get_prices(
|
||||
self, request: DeliveryCalculationRequest
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
"""Delivery provider registry and wiring."""
|
||||
|
||||
from dataclasses import dataclass
|
||||
from typing import Mapping, Sequence
|
||||
|
||||
import httpx
|
||||
|
||||
from app.adapters.delivery_providers.base import (
|
||||
DeliveryProvider,
|
||||
OrderRegistrationProvider,
|
||||
PaymentPriceValidationProvider,
|
||||
)
|
||||
from app.adapters.delivery_providers.cdek import CDEKProvider
|
||||
from app.adapters.delivery_providers.cse import CSEProvider
|
||||
from app.config import DeliveryProvidersConfig
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class DeliveryProviderRegistry:
|
||||
providers: Sequence[DeliveryProvider]
|
||||
payment_price_validation_adapters: Mapping[str, PaymentPriceValidationProvider]
|
||||
order_registration_adapters: Mapping[str, OrderRegistrationProvider]
|
||||
|
||||
|
||||
def resolve_delivery_provider_timeout_seconds(
|
||||
config: DeliveryProvidersConfig,
|
||||
) -> float:
|
||||
timeouts = [
|
||||
provider_config.timeout_seconds
|
||||
for provider_config in (config.cdek, config.cse)
|
||||
if provider_config.enabled
|
||||
]
|
||||
return max(timeouts, default=10.0)
|
||||
|
||||
|
||||
def build_delivery_provider_registry(
|
||||
*,
|
||||
http_client: httpx.AsyncClient,
|
||||
config: DeliveryProvidersConfig,
|
||||
) -> DeliveryProviderRegistry:
|
||||
providers: list[DeliveryProvider] = []
|
||||
payment_price_validation_adapters: dict[str, PaymentPriceValidationProvider] = {}
|
||||
order_registration_adapters: dict[str, OrderRegistrationProvider] = {}
|
||||
|
||||
if config.cdek.enabled:
|
||||
cdek_provider = CDEKProvider.from_config(
|
||||
http_client=http_client,
|
||||
config=config.cdek,
|
||||
)
|
||||
providers.append(cdek_provider)
|
||||
payment_price_validation_adapters[cdek_provider.name] = cdek_provider
|
||||
order_registration_adapters[cdek_provider.name] = cdek_provider
|
||||
|
||||
if config.cse.enabled:
|
||||
cse_provider = CSEProvider.from_config(
|
||||
http_client=http_client,
|
||||
config=config.cse,
|
||||
)
|
||||
providers.append(cse_provider)
|
||||
payment_price_validation_adapters[cse_provider.name] = cse_provider
|
||||
order_registration_adapters[cse_provider.name] = cse_provider
|
||||
|
||||
return DeliveryProviderRegistry(
|
||||
providers=tuple(providers),
|
||||
payment_price_validation_adapters=payment_price_validation_adapters,
|
||||
order_registration_adapters=order_registration_adapters,
|
||||
)
|
||||
+33
-21
@@ -30,7 +30,7 @@ class ServiceConfig(BaseModel):
|
||||
|
||||
class BusinessLogicConfig(BaseModel):
|
||||
weight_round_scale: int = 2
|
||||
provider_price_multiplier: Decimal = Field(..., gt=0)
|
||||
provider_price_multiplier: Decimal = Field(default=Decimal("1"), gt=0)
|
||||
|
||||
|
||||
class RepositoryConfig(BaseModel):
|
||||
@@ -38,26 +38,36 @@ class RepositoryConfig(BaseModel):
|
||||
price_cache_ttl_seconds: int = 900
|
||||
|
||||
|
||||
class AdapterConfig(BaseModel):
|
||||
cdek_base_url: str = "https://api.cdek.ru/v2"
|
||||
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)
|
||||
cse_base_url: str = "https://web.cse.ru/1c/ws/Web1C.1cws"
|
||||
cse_login: str = ""
|
||||
cse_password: str = ""
|
||||
cse_retry_attempts: int = Field(default=2, ge=0)
|
||||
cse_retry_backoff_seconds: float = Field(default=0.2, ge=0)
|
||||
cse_timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
cse_cache_ttl_seconds: int = Field(default=900, gt=0)
|
||||
class DeliveryProviderBaseConfig(BaseModel):
|
||||
enabled: bool = True
|
||||
timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
retry_attempts: int = Field(default=2, ge=0)
|
||||
retry_backoff_seconds: float = Field(default=0.2, ge=0)
|
||||
cache_ttl_seconds: int = Field(default=900, gt=0)
|
||||
|
||||
|
||||
class CDEKDeliveryProviderConfig(DeliveryProviderBaseConfig):
|
||||
base_url: str = "https://api.cdek.ru/v2"
|
||||
client_id: str = ""
|
||||
client_secret: str = ""
|
||||
|
||||
|
||||
class CSEDeliveryProviderConfig(DeliveryProviderBaseConfig):
|
||||
base_url: str = "https://web.cse.ru/1c/ws/Web1C.1cws"
|
||||
login: str = ""
|
||||
password: str = ""
|
||||
# Contract-specific required parameters for SaveDocuments (order registration).
|
||||
# Urgency is not here: it comes from the tariff selected by the client.
|
||||
cse_payer: str = ""
|
||||
cse_payment_method: str = ""
|
||||
cse_shipping_method: str = ""
|
||||
payer: str = ""
|
||||
payment_method: str = ""
|
||||
shipping_method: str = ""
|
||||
|
||||
|
||||
class DeliveryProvidersConfig(BaseModel):
|
||||
cdek: CDEKDeliveryProviderConfig = Field(
|
||||
default_factory=CDEKDeliveryProviderConfig
|
||||
)
|
||||
cse: CSEDeliveryProviderConfig = Field(default_factory=CSEDeliveryProviderConfig)
|
||||
|
||||
|
||||
class TBankPaymentAuthConfig(BaseModel):
|
||||
@@ -157,7 +167,9 @@ class Settings(BaseSettings):
|
||||
service: ServiceConfig = Field(default_factory=ServiceConfig)
|
||||
business_logic: BusinessLogicConfig = Field(default_factory=BusinessLogicConfig)
|
||||
repository: RepositoryConfig = Field(default_factory=RepositoryConfig)
|
||||
adapter: AdapterConfig = Field(default_factory=AdapterConfig)
|
||||
delivery_providers: DeliveryProvidersConfig = Field(
|
||||
default_factory=DeliveryProvidersConfig
|
||||
)
|
||||
tbank_payment: TBankPaymentConfig
|
||||
postgres: PostgresConfig
|
||||
address_suggestions: AddressSuggestionsConfig = Field(
|
||||
@@ -195,7 +207,7 @@ class _RequiredYamlSections(BaseModel):
|
||||
service: dict[str, Any]
|
||||
business_logic: dict[str, Any]
|
||||
repository: dict[str, Any]
|
||||
adapter: dict[str, Any]
|
||||
delivery_providers: dict[str, Any]
|
||||
tbank_payment: dict[str, Any]
|
||||
postgres: dict[str, Any]
|
||||
address_suggestions: dict[str, Any]
|
||||
|
||||
@@ -10,8 +10,10 @@ from app.adapters.address_suggestions.tomtom import TomTomAddressSuggestionProvi
|
||||
from app.adapters.address_suggestions.yandex_geosuggest import (
|
||||
YandexGeosuggestAddressSuggestionProvider,
|
||||
)
|
||||
from app.adapters.delivery_providers.cdek import CDEKProvider
|
||||
from app.adapters.delivery_providers.cse import CSEProvider
|
||||
from app.adapters.delivery_providers.registry import (
|
||||
build_delivery_provider_registry,
|
||||
resolve_delivery_provider_timeout_seconds,
|
||||
)
|
||||
from app.adapters.email import SMTPEmailSender
|
||||
from app.adapters.tbank import TBankAdapter
|
||||
from app.config import Settings
|
||||
@@ -49,15 +51,13 @@ _AGGREGATOR_SERVICE_STATE_KEY = "aggregator_service"
|
||||
|
||||
def _build_aggregator_service(settings: Settings) -> AggregatorService:
|
||||
http_client = build_controller_http_client(
|
||||
timeout_seconds=settings.adapter.cdek_timeout_seconds
|
||||
timeout_seconds=resolve_delivery_provider_timeout_seconds(
|
||||
settings.delivery_providers
|
||||
)
|
||||
)
|
||||
cdek_provider = CDEKProvider.from_adapter_config(
|
||||
delivery_provider_registry = build_delivery_provider_registry(
|
||||
http_client=http_client,
|
||||
adapter_config=settings.adapter,
|
||||
)
|
||||
cse_provider = CSEProvider.from_adapter_config(
|
||||
http_client=http_client,
|
||||
adapter_config=settings.adapter,
|
||||
config=settings.delivery_providers,
|
||||
)
|
||||
payment_adapter = TBankAdapter.from_config(
|
||||
http_client=http_client,
|
||||
@@ -75,7 +75,6 @@ def _build_aggregator_service(settings: Settings) -> AggregatorService:
|
||||
http_client=http_client,
|
||||
config=settings.address_suggestions.tomtom,
|
||||
)
|
||||
providers = (cdek_provider, cse_provider)
|
||||
cache = PriceCache.from_repository_config(
|
||||
settings.repository,
|
||||
metrics=get_cache_metrics(),
|
||||
@@ -93,18 +92,16 @@ def _build_aggregator_service(settings: Settings) -> AggregatorService:
|
||||
timeout_seconds=settings.email.timeout_seconds,
|
||||
)
|
||||
service = AggregatorService(
|
||||
providers=providers,
|
||||
providers=delivery_provider_registry.providers,
|
||||
cache=cache,
|
||||
payment_adapter=payment_adapter,
|
||||
payment_price_validation_adapters={
|
||||
cdek_provider.name: cdek_provider,
|
||||
cse_provider.name: cse_provider,
|
||||
},
|
||||
payment_price_validation_adapters=(
|
||||
delivery_provider_registry.payment_price_validation_adapters
|
||||
),
|
||||
order_repository=order_repository,
|
||||
order_registration_adapters={
|
||||
cdek_provider.name: cdek_provider,
|
||||
cse_provider.name: cse_provider,
|
||||
},
|
||||
order_registration_adapters=(
|
||||
delivery_provider_registry.order_registration_adapters
|
||||
),
|
||||
email_sender=email_sender,
|
||||
address_suggestion_providers=(
|
||||
dadata_provider,
|
||||
|
||||
@@ -45,12 +45,13 @@ class Order(Base):
|
||||
)
|
||||
payment_status: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
tbank_payment_id: Mapped[int | None] = mapped_column(BigInteger, nullable=True)
|
||||
cdek_order_uuid: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
cse_order_number: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
cdek_order_status: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
cdek_waybill_uuid: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
cdek_waybill_url: Mapped[str | None] = mapped_column(String(2048), nullable=True)
|
||||
cdek_polled_at: Mapped[datetime | None] = mapped_column(
|
||||
provider_order_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
provider_order_status: Mapped[str | None] = mapped_column(String(64), nullable=True)
|
||||
provider_waybill_id: Mapped[str | None] = mapped_column(String(128), nullable=True)
|
||||
provider_waybill_url: Mapped[str | None] = mapped_column(
|
||||
String(2048), nullable=True
|
||||
)
|
||||
provider_polled_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
nullable=True,
|
||||
)
|
||||
|
||||
@@ -75,31 +75,17 @@ class OrderRepository:
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
async def mark_cdek_order_registered(
|
||||
async def mark_provider_order_registered(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
order_uuid: str,
|
||||
cdek_order_uuid: str,
|
||||
provider_order_id: str,
|
||||
) -> Order | None:
|
||||
order = await self.get_order_by_order_uuid(session, order_uuid)
|
||||
if order is None:
|
||||
return None
|
||||
|
||||
order.cdek_order_uuid = cdek_order_uuid
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
async def mark_cse_order_registered(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
order_uuid: str,
|
||||
cse_order_number: str,
|
||||
) -> Order | None:
|
||||
order = await self.get_order_by_order_uuid(session, order_uuid)
|
||||
if order is None:
|
||||
return None
|
||||
|
||||
order.cse_order_number = cse_order_number
|
||||
order.provider_order_id = provider_order_id
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
@@ -112,14 +98,15 @@ class OrderRepository:
|
||||
statement = (
|
||||
select(Order)
|
||||
.where(
|
||||
Order.cdek_order_uuid.is_not(None),
|
||||
Order.cdek_waybill_url.is_(None),
|
||||
Order.provider == "cdek",
|
||||
Order.provider_order_id.is_not(None),
|
||||
Order.provider_waybill_url.is_(None),
|
||||
(
|
||||
Order.cdek_order_status.is_(None)
|
||||
| Order.cdek_order_status.not_in(TERMINAL_ORDER_STATUSES)
|
||||
Order.provider_order_status.is_(None)
|
||||
| Order.provider_order_status.not_in(TERMINAL_ORDER_STATUSES)
|
||||
),
|
||||
)
|
||||
.order_by(Order.cdek_polled_at.asc().nulls_first())
|
||||
.order_by(Order.provider_polled_at.asc().nulls_first())
|
||||
.limit(limit)
|
||||
.with_for_update(skip_locked=True)
|
||||
)
|
||||
@@ -139,10 +126,10 @@ class OrderRepository:
|
||||
if order is None:
|
||||
return None
|
||||
|
||||
order.cdek_order_status = order_status
|
||||
if waybill_uuid is not None and order.cdek_waybill_uuid is None:
|
||||
order.cdek_waybill_uuid = waybill_uuid
|
||||
order.cdek_polled_at = polled_at
|
||||
order.provider_order_status = order_status
|
||||
if waybill_uuid is not None and order.provider_waybill_id is None:
|
||||
order.provider_waybill_id = waybill_uuid
|
||||
order.provider_polled_at = polled_at
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
@@ -158,9 +145,9 @@ class OrderRepository:
|
||||
if order is None:
|
||||
return None
|
||||
|
||||
if waybill_url is not None and order.cdek_waybill_url is None:
|
||||
order.cdek_waybill_url = waybill_url
|
||||
order.cdek_polled_at = polled_at
|
||||
if waybill_url is not None and order.provider_waybill_url is None:
|
||||
order.provider_waybill_url = waybill_url
|
||||
order.provider_polled_at = polled_at
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
@@ -189,7 +176,8 @@ class OrderRepository:
|
||||
statement = (
|
||||
select(Order)
|
||||
.where(
|
||||
Order.cdek_waybill_url.is_not(None),
|
||||
Order.provider == "cdek",
|
||||
Order.provider_waybill_url.is_not(None),
|
||||
Order.waybill_email_sent_at.is_(None),
|
||||
)
|
||||
.order_by(Order.created_at.asc())
|
||||
|
||||
+11
-29
@@ -141,18 +141,11 @@ class OrderRepositoryProtocol(Protocol):
|
||||
payment_id: int,
|
||||
) -> object | None: ...
|
||||
|
||||
async def mark_cdek_order_registered(
|
||||
async def mark_provider_order_registered(
|
||||
self,
|
||||
session: object,
|
||||
order_uuid: str,
|
||||
cdek_order_uuid: str,
|
||||
) -> object | None: ...
|
||||
|
||||
async def mark_cse_order_registered(
|
||||
self,
|
||||
session: object,
|
||||
order_uuid: str,
|
||||
cse_order_number: str,
|
||||
provider_order_id: str,
|
||||
) -> object | None: ...
|
||||
|
||||
|
||||
@@ -357,7 +350,7 @@ class AggregatorService:
|
||||
error=str(exc),
|
||||
)
|
||||
raise InvalidInitPaymentRequestError(
|
||||
"Payment init request is invalid for CDEK price validation."
|
||||
"Payment init request is invalid for provider price validation."
|
||||
) from exc
|
||||
except ProviderClientError as exc:
|
||||
logger.warning(
|
||||
@@ -389,7 +382,7 @@ class AggregatorService:
|
||||
requested_price_kopecks=requested_price,
|
||||
)
|
||||
raise InvalidInitPaymentRequestError(
|
||||
"CDEK did not return the requested tariff for payment validation."
|
||||
"Provider did not return the requested tariff for payment validation."
|
||||
)
|
||||
|
||||
expected_amount_kopecks = calculate_expected_payment_amount_kopecks(
|
||||
@@ -411,7 +404,7 @@ class AggregatorService:
|
||||
provider_price=str(getattr(provider_price, "price", None)),
|
||||
)
|
||||
raise InvalidInitPaymentRequestError(
|
||||
"Payment amount does not match CDEK validated delivery price."
|
||||
"Payment amount does not match provider validated delivery price."
|
||||
)
|
||||
|
||||
async def handle_tbank_payment_notification(
|
||||
@@ -593,11 +586,7 @@ class AggregatorService:
|
||||
|
||||
@staticmethod
|
||||
def _has_existing_registration(order: object, provider: str) -> bool:
|
||||
if provider == "cdek":
|
||||
return bool(getattr(order, "cdek_order_uuid", None))
|
||||
if provider == "cse":
|
||||
return bool(getattr(order, "cse_order_number", None))
|
||||
return False
|
||||
return bool(getattr(order, "provider_order_id", None))
|
||||
|
||||
@staticmethod
|
||||
def _extract_registration_id(result: object) -> str:
|
||||
@@ -646,18 +635,11 @@ class AggregatorService:
|
||||
registration_id = self._extract_registration_id(result)
|
||||
try:
|
||||
async with self._order_repository.session() as session:
|
||||
if provider == "cse":
|
||||
order = await self._order_repository.mark_cse_order_registered(
|
||||
session,
|
||||
order_uuid,
|
||||
registration_id,
|
||||
)
|
||||
else:
|
||||
order = await self._order_repository.mark_cdek_order_registered(
|
||||
session,
|
||||
order_uuid,
|
||||
registration_id,
|
||||
)
|
||||
order = await self._order_repository.mark_provider_order_registered(
|
||||
session,
|
||||
order_uuid,
|
||||
registration_id,
|
||||
)
|
||||
if order is None:
|
||||
logger.warning(
|
||||
"order_registration_order_not_found",
|
||||
|
||||
@@ -41,7 +41,7 @@ class EmailSenderProtocol(Protocol):
|
||||
class OrderRecord(Protocol):
|
||||
order_uuid: str
|
||||
account_email: str
|
||||
cdek_waybill_url: str | None
|
||||
provider_waybill_url: str | None
|
||||
|
||||
|
||||
class WaybillEmailSenderRepositoryProtocol(Protocol):
|
||||
@@ -133,7 +133,7 @@ class WaybillEmailSenderService:
|
||||
continue
|
||||
|
||||
async def _handle_order(self, session: object, order: OrderRecord) -> None:
|
||||
waybill_url = order.cdek_waybill_url
|
||||
waybill_url = order.provider_waybill_url
|
||||
if waybill_url is None:
|
||||
return
|
||||
|
||||
|
||||
@@ -27,8 +27,8 @@ class CDEKWaybillInfoAdapterProtocol(Protocol):
|
||||
|
||||
class OrderRecord(Protocol):
|
||||
order_uuid: str
|
||||
cdek_order_uuid: str | None
|
||||
cdek_waybill_uuid: str | None
|
||||
provider_order_id: str | None
|
||||
provider_waybill_id: str | None
|
||||
|
||||
|
||||
class WaybillPollerRepositoryProtocol(Protocol):
|
||||
@@ -97,8 +97,8 @@ class WaybillPollerService:
|
||||
logger.exception(
|
||||
"waybill_poll_order_failed",
|
||||
order_uuid=order.order_uuid,
|
||||
cdek_order_uuid=order.cdek_order_uuid,
|
||||
cdek_waybill_uuid=order.cdek_waybill_uuid,
|
||||
provider_order_id=order.provider_order_id,
|
||||
provider_waybill_id=order.provider_waybill_id,
|
||||
)
|
||||
return PollBatchSummary(
|
||||
processed=len(orders),
|
||||
@@ -133,11 +133,11 @@ class WaybillPollerService:
|
||||
|
||||
async def _handle_order(self, session: object, order: OrderRecord) -> None:
|
||||
polled_at = self._datetime_now()
|
||||
if order.cdek_waybill_uuid is None:
|
||||
cdek_order_uuid = order.cdek_order_uuid
|
||||
if cdek_order_uuid is None:
|
||||
if order.provider_waybill_id is None:
|
||||
provider_order_id = order.provider_order_id
|
||||
if provider_order_id is None:
|
||||
return
|
||||
info = await self._order_info_adapter.get_order(cdek_order_uuid)
|
||||
info = await self._order_info_adapter.get_order(provider_order_id)
|
||||
await self._repository.record_order_poll(
|
||||
session,
|
||||
order_uuid=order.order_uuid,
|
||||
@@ -148,13 +148,15 @@ class WaybillPollerService:
|
||||
logger.info(
|
||||
"waybill_poll_order_result",
|
||||
order_uuid=order.order_uuid,
|
||||
cdek_order_uuid=cdek_order_uuid,
|
||||
cdek_order_status=info.status_code,
|
||||
cdek_waybill_uuid=info.waybill_uuid,
|
||||
provider_order_id=provider_order_id,
|
||||
provider_order_status=info.status_code,
|
||||
provider_waybill_id=info.waybill_uuid,
|
||||
)
|
||||
return
|
||||
|
||||
waybill = await self._waybill_info_adapter.get_waybill(order.cdek_waybill_uuid)
|
||||
waybill = await self._waybill_info_adapter.get_waybill(
|
||||
order.provider_waybill_id
|
||||
)
|
||||
await self._repository.record_waybill_poll(
|
||||
session,
|
||||
order_uuid=order.order_uuid,
|
||||
@@ -164,6 +166,6 @@ class WaybillPollerService:
|
||||
logger.info(
|
||||
"waybill_poll_waybill_result",
|
||||
order_uuid=order.order_uuid,
|
||||
cdek_waybill_uuid=order.cdek_waybill_uuid,
|
||||
cdek_waybill_url=waybill.url,
|
||||
provider_waybill_id=order.provider_waybill_id,
|
||||
provider_waybill_url=waybill.url,
|
||||
)
|
||||
|
||||
@@ -22,21 +22,22 @@ logger = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
async def _run(settings: Settings, stop_event: asyncio.Event) -> None:
|
||||
http_client = httpx.AsyncClient(timeout=settings.adapter.cdek_timeout_seconds)
|
||||
cdek_config = settings.delivery_providers.cdek
|
||||
http_client = httpx.AsyncClient(timeout=cdek_config.timeout_seconds)
|
||||
auth_client = CDEKAuthClient(
|
||||
http_client=http_client,
|
||||
base_url=settings.adapter.cdek_base_url,
|
||||
client_id=settings.adapter.cdek_client_id,
|
||||
client_secret=settings.adapter.cdek_client_secret,
|
||||
timeout_seconds=settings.adapter.cdek_timeout_seconds,
|
||||
base_url=cdek_config.base_url,
|
||||
client_id=cdek_config.client_id,
|
||||
client_secret=cdek_config.client_secret,
|
||||
timeout_seconds=cdek_config.timeout_seconds,
|
||||
)
|
||||
cdek_client = CDEKClient(
|
||||
http_client=http_client,
|
||||
auth_client=auth_client,
|
||||
base_url=settings.adapter.cdek_base_url,
|
||||
timeout_seconds=settings.adapter.cdek_timeout_seconds,
|
||||
retry_attempts=settings.adapter.cdek_retry_attempts,
|
||||
retry_backoff_seconds=settings.adapter.cdek_retry_backoff_seconds,
|
||||
base_url=cdek_config.base_url,
|
||||
timeout_seconds=cdek_config.timeout_seconds,
|
||||
retry_attempts=cdek_config.retry_attempts,
|
||||
retry_backoff_seconds=cdek_config.retry_backoff_seconds,
|
||||
)
|
||||
email_sender = SMTPEmailSender(
|
||||
smtp_host=settings.email.smtp_host,
|
||||
|
||||
@@ -21,21 +21,22 @@ logger = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
async def _run(settings: Settings, stop_event: asyncio.Event) -> None:
|
||||
http_client = httpx.AsyncClient(timeout=settings.adapter.cdek_timeout_seconds)
|
||||
cdek_config = settings.delivery_providers.cdek
|
||||
http_client = httpx.AsyncClient(timeout=cdek_config.timeout_seconds)
|
||||
auth_client = CDEKAuthClient(
|
||||
http_client=http_client,
|
||||
base_url=settings.adapter.cdek_base_url,
|
||||
client_id=settings.adapter.cdek_client_id,
|
||||
client_secret=settings.adapter.cdek_client_secret,
|
||||
timeout_seconds=settings.adapter.cdek_timeout_seconds,
|
||||
base_url=cdek_config.base_url,
|
||||
client_id=cdek_config.client_id,
|
||||
client_secret=cdek_config.client_secret,
|
||||
timeout_seconds=cdek_config.timeout_seconds,
|
||||
)
|
||||
cdek_client = CDEKClient(
|
||||
http_client=http_client,
|
||||
auth_client=auth_client,
|
||||
base_url=settings.adapter.cdek_base_url,
|
||||
timeout_seconds=settings.adapter.cdek_timeout_seconds,
|
||||
retry_attempts=settings.adapter.cdek_retry_attempts,
|
||||
retry_backoff_seconds=settings.adapter.cdek_retry_backoff_seconds,
|
||||
base_url=cdek_config.base_url,
|
||||
timeout_seconds=cdek_config.timeout_seconds,
|
||||
retry_attempts=cdek_config.retry_attempts,
|
||||
retry_backoff_seconds=cdek_config.retry_backoff_seconds,
|
||||
)
|
||||
engine = create_postgres_engine(settings.postgres)
|
||||
session_factory = create_postgres_session_factory(engine)
|
||||
|
||||
Reference in New Issue
Block a user