Compare commits

..

11 Commits

Author SHA1 Message Date
Раис Юсупалиев c2728ae268 КОСТЫЛЬ перетираем дату доставки
Deploy / deploy (push) Successful in 53s
2026-06-27 22:52:19 +03:00
Раис Юсупалиев c942885abe fix pickup date
Deploy / deploy (push) Successful in 1m2s
2026-06-27 22:41:10 +03:00
Раис Юсупалиев 7287a1398e еще логи
Deploy / deploy (push) Successful in 56s
2026-06-27 18:42:50 +03:00
Раис Юсупалиев a2d9a243f7 add cse logs
Deploy / deploy (push) Successful in 58s
2026-06-27 18:10:20 +03:00
Раис Юсупалиев 4262b8a200 fix ксе tariffs
Deploy / deploy (push) Successful in 57s
2026-06-27 17:34:06 +03:00
Раис Юсупалиев e4d9b581a6 fix ксе prices
Deploy / deploy (push) Successful in 55s
2026-06-27 16:45:58 +03:00
Раис Юсупалиев facdde00c9 добавлено получение накладной в cse
Deploy / deploy (push) Successful in 2m33s
2026-06-27 07:10:31 +03:00
Раис Юсупалиев 2d8f31d3d9 оценочная стоимость отправляется в cdek
Deploy / deploy (push) Successful in 1m1s
2026-06-26 20:59:27 +03:00
Раис Юсупалиев 6906db7739 убрал packages из cdek payload
Deploy / deploy (push) Successful in 52s
2026-06-26 20:46:21 +03:00
Раис Юсупалиев 843175f12e Рефактор
Deploy / deploy (push) Successful in 18m5s
2026-06-26 20:16:04 +03:00
Раис Юсупалиев c6c37640fd добавлен declaredValue
Deploy / deploy (push) Successful in 53s
2026-06-21 22:55:29 +03:00
48 changed files with 1986 additions and 590 deletions
@@ -0,0 +1,115 @@
"""Generalize provider order state columns."""
from collections.abc import Sequence
from alembic import op
import sqlalchemy as sa
revision: str = "20260626_039"
down_revision: str | None = "20260530_038"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.add_column(
"orders",
sa.Column("provider_order_id", sa.String(length=128), nullable=True),
)
op.add_column(
"orders",
sa.Column("provider_order_status", sa.String(length=64), nullable=True),
)
op.add_column(
"orders",
sa.Column("provider_waybill_id", sa.String(length=128), nullable=True),
)
op.add_column(
"orders",
sa.Column("provider_waybill_url", sa.String(length=2048), nullable=True),
)
op.add_column(
"orders",
sa.Column("provider_polled_at", sa.DateTime(timezone=True), nullable=True),
)
op.execute(
"""
UPDATE orders
SET
provider_order_id = cdek_order_uuid,
provider_order_status = cdek_order_status,
provider_waybill_id = cdek_waybill_uuid,
provider_waybill_url = cdek_waybill_url,
provider_polled_at = cdek_polled_at
WHERE provider = 'cdek'
"""
)
op.execute(
"""
UPDATE orders
SET provider_order_id = cse_order_number
WHERE provider = 'cse'
"""
)
op.drop_column("orders", "cse_order_number")
op.drop_column("orders", "cdek_polled_at")
op.drop_column("orders", "cdek_order_status")
op.drop_column("orders", "cdek_waybill_url")
op.drop_column("orders", "cdek_waybill_uuid")
op.drop_column("orders", "cdek_order_uuid")
def downgrade() -> None:
op.add_column(
"orders",
sa.Column("cdek_order_uuid", sa.String(length=128), nullable=True),
)
op.add_column(
"orders",
sa.Column("cdek_waybill_uuid", sa.String(length=128), nullable=True),
)
op.add_column(
"orders",
sa.Column("cdek_waybill_url", sa.String(length=2048), nullable=True),
)
op.add_column(
"orders",
sa.Column("cdek_order_status", sa.String(length=64), nullable=True),
)
op.add_column(
"orders",
sa.Column("cdek_polled_at", sa.DateTime(timezone=True), nullable=True),
)
op.add_column(
"orders",
sa.Column("cse_order_number", sa.String(length=128), nullable=True),
)
op.execute(
"""
UPDATE orders
SET
cdek_order_uuid = provider_order_id,
cdek_order_status = provider_order_status,
cdek_waybill_uuid = provider_waybill_id,
cdek_waybill_url = provider_waybill_url,
cdek_polled_at = provider_polled_at
WHERE provider = 'cdek'
"""
)
op.execute(
"""
UPDATE orders
SET cse_order_number = provider_order_id
WHERE provider = 'cse'
"""
)
op.drop_column("orders", "provider_polled_at")
op.drop_column("orders", "provider_waybill_url")
op.drop_column("orders", "provider_waybill_id")
op.drop_column("orders", "provider_order_status")
op.drop_column("orders", "provider_order_id")
+23
View File
@@ -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."""
+12 -12
View File
@@ -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
@@ -20,6 +20,7 @@ class CDEKOrderMappingError(ValueError):
_CDEK_WAYBILL_PRINT_TYPE = "WAYBILL"
_CDEK_WAYBILL_RELATED_ENTITY_TYPE = "waybill"
_CDEK_INSURANCE_SERVICE_CODE = "INSURANCE"
@dataclass(frozen=True)
@@ -56,6 +57,9 @@ def map_cdek_order_request(
"to_location": _map_location(request.receiver_address),
"packages": [_map_package(request, order_uuid)],
}
services = _map_services(request)
if services:
payload["services"] = services
if request.content.description:
payload["comment"] = request.content.description
return payload
@@ -126,7 +130,9 @@ def map_cdek_order_info_response(payload: dict[str, Any]) -> CDEKOrderInfo:
"CDEK order info response must include entity.uuid."
)
status_code = _latest_status_code(entity.get("statuses"))
status_code = _invalid_create_request_status(
payload.get("requests")
) or _latest_status_code(entity.get("statuses"))
waybill_uuid, _ = _extract_waybill(payload)
return CDEKOrderInfo(
order_uuid=order_uuid,
@@ -172,6 +178,20 @@ def _latest_status_code(statuses: object) -> str | None:
return dated[-1][1]
def _invalid_create_request_status(requests: object) -> str | None:
if not isinstance(requests, list):
return None
for entry in requests:
if not isinstance(entry, dict):
continue
if entry.get("type") != "CREATE" or entry.get("state") != "INVALID":
continue
errors = entry.get("errors")
if isinstance(errors, list) and errors:
return "INVALID"
return None
def _extract_waybill(payload: dict[str, Any]) -> tuple[str | None, str | None]:
entity = payload.get("entity")
related_entities = (
@@ -312,6 +332,18 @@ def _map_package(request: InitPaymentRequest, order_uuid: str) -> dict[str, Any]
return package
def _map_services(request: InitPaymentRequest) -> list[dict[str, Any]]:
declared_value = request.content.declared_value
if declared_value <= 0:
return []
return [
{
"code": _CDEK_INSURANCE_SERVICE_CODE,
"parameter": declared_value,
}
]
def _map_dimensions(dimensions: Dimensions) -> dict[str, int]:
return {
"length": centimeters_string_to_int(dimensions.length, "length"),
+100 -25
View File
@@ -21,12 +21,17 @@ from app.adapters.delivery_providers.cse.mapper import (
map_cse_calc_response_for_tariff_code,
)
from app.adapters.delivery_providers.cse.order_mapper import (
CSEOrderInfo,
CSEOrderRegistrationParams,
CSEOrderRegistrationResult,
build_calc_body_for_calculation,
build_calc_body_for_payment,
build_tracking_body_for_order,
build_waybill_print_form_body,
map_cse_tracking_response,
map_cse_save_order_request,
map_cse_save_order_response,
map_cse_waybill_print_form_response,
split_tariff_code,
)
from app.adapters.delivery_providers.cse.soap import (
@@ -35,7 +40,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
@@ -113,6 +118,30 @@ class CSEClient:
except CSEMappingError as exc:
raise CSEClientError(str(exc)) from exc
async def get_order(self, order_number: str) -> CSEOrderInfo:
root = await self._post(
"Tracking",
build_tracking_body_for_order(order_number),
request_error_message="CSE Tracking request was rejected with status",
)
try:
return map_cse_tracking_response(root)
except CSEMappingError as exc:
raise CSEClientError(str(exc)) from exc
async def download_waybill_pdf(self, waybill_number: str) -> bytes:
root = await self._post(
"GetFormsForDocuments",
build_waybill_print_form_body(waybill_number),
request_error_message=(
"CSE GetFormsForDocuments request was rejected with status"
),
)
try:
return map_cse_waybill_print_form_response(root)
except CSEMappingError as exc:
raise CSEClientError(str(exc)) from exc
async def _post(
self,
operation: str,
@@ -150,9 +179,11 @@ class CSEClient:
"cse_request_server_error",
operation=operation,
status_code=response.status_code,
response_excerpt=_response_excerpt(response.text),
)
raise CSEClientError(
f"CSE {operation} request failed with status {response.status_code}."
f"CSE {operation} request failed with status "
f"{response.status_code}."
)
if 400 <= response.status_code < 500:
@@ -183,16 +214,15 @@ class CSEClient:
operation: str,
request_error_message: str,
) -> None:
for prop in root.properties:
if prop.key == "Error":
codes = [item.value for item in prop.items if item.value]
error_codes = _response_error_codes(root)
if error_codes:
log.warning(
"cse_response_error",
operation=operation,
error_codes=codes,
error_codes=error_codes,
)
raise CSERequestError(
f"{request_error_message} application error {codes}."
f"{request_error_message} application error {error_codes}."
)
@staticmethod
@@ -203,36 +233,65 @@ class CSEClient:
return self._retry_backoff_seconds * (2**attempt)
def _response_excerpt(text: str, *, limit: int = 1000) -> str:
return " ".join(text.split())[:limit]
def _response_error_codes(root: Element) -> list[str]:
codes: list[str] = []
_collect_response_error_codes(root, codes)
return codes
def _collect_response_error_codes(element: Element, codes: list[str]) -> None:
for prop in element.properties:
if prop.key == "Error":
codes.extend(item.value for item in prop.items if item.value)
for child in (*element.items, *element.tables):
_collect_response_error_codes(child, codes)
class CSEProvider(DeliveryProvider):
name = CSE_PROVIDER_NAME
def __init__(self, client: CSEClient, *, cache_ttl_seconds: int = 900) -> None:
def __init__(
self,
client: CSEClient,
*,
cache_ttl_seconds: int = 900,
delivery_service_guids: tuple[str, ...] = (),
) -> None:
self._client = client
self.cache_ttl_seconds = cache_ttl_seconds
self._delivery_service_guids = delivery_service_guids
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=config.cache_ttl_seconds,
delivery_service_guids=tuple(config.delivery_service_guids),
)
return cls(client=client, cache_ttl_seconds=adapter_config.cse_cache_ttl_seconds)
async def get_prices(
self, request: DeliveryCalculationRequest
@@ -241,8 +300,9 @@ class CSEProvider(DeliveryProvider):
results = await asyncio.gather(
*(
self._prices_for_delivery_type(request, name, label)
self._prices_for_delivery_type(request, name, label, service_guid)
for name, label in delivery_types
for service_guid in self._service_guids_for_price_calculation()
),
return_exceptions=True,
)
@@ -264,18 +324,25 @@ class CSEProvider(DeliveryProvider):
request: DeliveryCalculationRequest,
delivery_type: str,
delivery_type_label: str,
service_guid: str | None,
) -> list[DeliveryPrice]:
body = build_calc_body_for_calculation(request, delivery_type)
body = build_calc_body_for_calculation(request, delivery_type, service_guid)
root = await self._client.calc(body)
try:
return map_cse_calc_response(
root,
delivery_type=delivery_type,
delivery_type_label=delivery_type_label,
service_guid=service_guid,
)
except CSEMappingError as exc:
raise CSEClientError("CSE calc response payload is invalid.") from exc
def _service_guids_for_price_calculation(self) -> tuple[str, ...]:
if not self._delivery_service_guids:
raise CSERequestError("CSE delivery service GUIDs are not configured.")
return self._delivery_service_guids
async def _resolve_delivery_types(self) -> list[tuple[str, str]]:
if self._delivery_types is None:
self._delivery_types = await self._client.get_delivery_types()
@@ -295,14 +362,16 @@ class CSEProvider(DeliveryProvider):
request: InitPaymentRequest,
) -> DeliveryPrice | None:
tariff_code = request.system_data.tariff.tariff_code
delivery_type, _ = split_tariff_code(tariff_code)
body = build_calc_body_for_payment(request, delivery_type)
root = await self._client.calc(body)
try:
delivery_type, service_guid, _ = split_tariff_code(tariff_code)
body = build_calc_body_for_payment(request, delivery_type, service_guid)
root = await self._client.calc(body)
return map_cse_calc_response_for_tariff_code(
root,
tariff_code=tariff_code,
)
except CSERequestError:
raise
except CSEMappingError as exc:
raise CSEClientError(
"CSE payment price validation response payload is invalid."
@@ -312,3 +381,9 @@ class CSEProvider(DeliveryProvider):
self, request: InitPaymentRequest, order_uuid: str
) -> CSEOrderRegistrationResult:
return await self._client.save_order(request, order_uuid)
async def get_order(self, order_number: str) -> CSEOrderInfo:
return await self._client.get_order(order_number)
async def download_waybill_pdf(self, waybill_number: str) -> bytes:
return await self._client.download_waybill_pdf(waybill_number)
+37 -7
View File
@@ -18,11 +18,14 @@ def map_cse_calc_response(
*,
delivery_type: str = "",
delivery_type_label: str = "",
service_guid: str | None = None,
) -> list[DeliveryPrice]:
"""Map a parsed ``Calc`` ``return`` Element into unified delivery prices."""
prices: list[DeliveryPrice] = []
for tariff in _iter_tariffs(root):
if service_guid is not None and tariff.value != service_guid:
continue
price = _map_tariff(tariff, delivery_type, delivery_type_label)
if price is not None:
prices.append(price)
@@ -35,12 +38,14 @@ def map_cse_calc_response_for_tariff_code(
tariff_code: str,
delivery_type_label: str = "",
) -> DeliveryPrice | None:
"""Return the tariff matching ``tariff_code`` (``"<DeliveryType>|<Urgency>"``)."""
"""Return tariff matching ``"<DeliveryType>|<ServiceGuid>|<Urgency>"``."""
delivery_type, _, urgency = tariff_code.partition("|")
delivery_type, tariff_guid, urgency = _split_tariff_code(tariff_code)
for tariff in _iter_tariffs(root):
if _tariff_urgency(tariff) == urgency:
return _map_tariff(tariff, delivery_type, delivery_type_label)
if _tariff_matches(tariff, tariff_guid, urgency):
price = _map_tariff(tariff, delivery_type, delivery_type_label)
if price is not None:
return price
return None
@@ -50,6 +55,19 @@ def _tariff_urgency(tariff: Element) -> str | None:
return tariff.field_value("Urgency") or tariff.value
def _tariff_matches(tariff: Element, tariff_guid: str, urgency: str) -> bool:
if _tariff_urgency(tariff) != urgency:
return False
return tariff.value == tariff_guid
def _split_tariff_code(tariff_code: str) -> tuple[str, str, str]:
parts = tariff_code.split("|")
if len(parts) != 3 or not parts[1] or not parts[2]:
raise CSEMappingError("CSE tariff_code has invalid format.")
return parts[0], parts[1], parts[2]
def _iter_tariffs(root: Element):
for destination in root.items:
for tariff in destination.items:
@@ -62,12 +80,17 @@ def _map_tariff(
delivery_type: str,
delivery_type_label: str,
) -> DeliveryPrice | None:
if _is_additional_service_tariff(tariff):
return None
urgency = _tariff_urgency(tariff)
if not urgency:
return None
# Unified tariff_code carries both the delivery scheme and the urgency so
# registration (SaveDocuments) can set DeliveryOfCargo and Urgency.
tariff_code = f"{delivery_type}|{urgency}"
if not tariff.value:
return None
# Unified tariff_code carries the delivery scheme, CSE service GUID and
# urgency so payment validation can recalculate the exact selected service.
tariff_code = f"{delivery_type}|{tariff.value}|{urgency}"
raw_total = tariff.field_value("Total")
if raw_total is None:
@@ -111,6 +134,13 @@ def _map_tariff(
raise CSEMappingError("CSE tariff fields have invalid values.") from exc
def _is_additional_service_tariff(tariff: Element) -> bool:
value = tariff.field_value("AdditionalService")
if value is None:
return False
return value.strip().lower() == "true"
def _to_int(value: str | None) -> int | None:
if value is None or value == "":
return None
@@ -1,17 +1,21 @@
"""CSE Calc and SaveDocuments payload mappers."""
"""CSE Calc, SaveDocuments and waybill payload mappers."""
import base64
import binascii
from dataclasses import dataclass
from datetime import timedelta
from app.adapters.delivery_providers.cse.constants import (
resolve_cse_geography,
resolve_type_of_cargo,
)
from app.adapters.delivery_providers.cse.errors import CSEMappingError
from app.adapters.delivery_providers.cse.errors import CSEMappingError, CSERequestError
from app.adapters.delivery_providers.cse.soap import Element, make_field
from app.schemas.payment import Address, InitPaymentRequest
from app.schemas.request import DeliveryCalculationRequest
_DATETIME_FORMAT = "%Y-%m-%dT%H:%M:%S"
_DELIVERY_DATE_OFFSET_DAYS = 7
@dataclass(frozen=True)
@@ -32,9 +36,17 @@ class CSEOrderRegistrationResult:
order_number: str
@dataclass(frozen=True)
class CSEOrderInfo:
order_number: str
status_code: str | None
waybill_number: str | None
def build_calc_body_for_calculation(
request: DeliveryCalculationRequest,
delivery_type: str = "",
service_guid: str | None = None,
) -> dict[str, Element]:
fields = [
make_field("SenderGeography", resolve_cse_geography(request.from_city)),
@@ -44,12 +56,15 @@ def build_calc_body_for_calculation(
make_field("Qty", "1", "int"),
]
_append_delivery_type(fields, delivery_type)
if service_guid:
fields.append(make_field("Service", service_guid))
return _calc_body(Element(key="Destination", fields=fields))
def build_calc_body_for_payment(
request: InitPaymentRequest,
delivery_type: str = "",
service_guid: str | None = None,
) -> dict[str, Element]:
system_data = request.system_data
fields = [
@@ -66,6 +81,8 @@ def build_calc_body_for_payment(
make_field("Qty", "1", "int"),
]
_append_delivery_type(fields, delivery_type)
if service_guid:
fields.append(make_field("Service", service_guid))
return _calc_body(Element(key="Destination", fields=fields))
@@ -74,14 +91,17 @@ def _append_delivery_type(fields: list[Element], delivery_type: str) -> None:
fields.append(make_field("DeliveryType", delivery_type))
def split_tariff_code(tariff_code: str) -> tuple[str, str]:
"""Split the CSE ``tariff_code`` into (delivery_type, urgency).
def split_tariff_code(tariff_code: str) -> tuple[str, str, str]:
"""Split the CSE ``tariff_code`` into delivery type, service GUID, urgency.
The unified tariff_code encodes both dimensions as ``"<DeliveryType>|<Urgency>"``.
The expected format is ``"<DeliveryType>|<ServiceGuid>|<Urgency>"``.
``DeliveryType`` may be empty when CSE contract defaults are used.
"""
delivery_type, _, urgency = tariff_code.partition("|")
return delivery_type, urgency
parts = tariff_code.split("|")
if len(parts) != 3 or not parts[1] or not parts[2]:
raise CSERequestError("CSE tariff_code has invalid format.")
return parts[0], parts[1], parts[2]
def _calc_body(destination: Element) -> dict[str, Element]:
@@ -101,7 +121,10 @@ def map_cse_save_order_request(
) -> dict[str, Element]:
system_data = request.system_data
take_date = request.pickup_date.strftime(_DATETIME_FORMAT)
delivery_type, urgency = split_tariff_code(system_data.tariff.tariff_code)
delivery_date = (
request.pickup_date + timedelta(days=_DELIVERY_DATE_OFFSET_DAYS)
).strftime(_DATETIME_FORMAT)
delivery_type, _, urgency = split_tariff_code(system_data.tariff.tariff_code)
fields = [
make_field("TakeDate", take_date, "dateTime"),
@@ -127,18 +150,12 @@ def map_cse_save_order_request(
make_field("TypeOfCargo", resolve_type_of_cargo(system_data.parcel_type)),
make_field("Weight", system_data.weight, "float"),
make_field("CargoPackageQty", "1", "float"),
make_field("CargoCost", str(request.content.declared_value), "float"),
]
if delivery_type:
fields.append(make_field("DeliveryOfCargo", delivery_type))
if request.delivery_date is not None:
fields.append(
make_field(
"DeliveryDate",
request.delivery_date.strftime(_DATETIME_FORMAT),
"dateTime",
)
)
fields.append(make_field("DeliveryDate", delivery_date, "dateTime"))
if request.content.description:
fields.append(make_field("CargoDescription", request.content.description))
fields.append(make_field("Comment", request.content.description))
@@ -175,6 +192,73 @@ def map_cse_save_order_response(root: Element) -> CSEOrderRegistrationResult:
raise CSEMappingError("CSE SaveDocuments response is missing document Number.")
def build_tracking_body_for_order(order_number: str) -> dict[str, Element]:
return {
"documents": Element(
key="Documents",
items=[Element(key=order_number)],
),
"parameters": Element(
key="parameters",
items=[
make_field("DocumentType", "Order"),
make_field("OnlySelectedType", True, "boolean"),
],
),
}
def map_cse_tracking_response(root: Element) -> CSEOrderInfo:
for document in root.items:
order_number = document.property_value("Number") or document.key
if not order_number:
continue
return CSEOrderInfo(
order_number=order_number,
status_code=_latest_tracking_status(document),
waybill_number=_extract_waybill_number(document),
)
raise CSEMappingError("CSE Tracking response is missing document data.")
def build_waybill_print_form_body(waybill_number: str) -> dict[str, Element]:
return {
"documents": Element(
key="Documents",
items=[Element(key=waybill_number)],
),
"parameters": Element(
key="parameters",
items=[
make_field("DocumentType", "waybill"),
make_field("Type", "print"),
make_field(
"Name",
"Универсальная печатная "
"форма документа НАКЛАДНАЯ",
),
make_field("Format", "pdf"),
make_field("OnlySelectedType", True, "boolean"),
],
),
}
def map_cse_waybill_print_form_response(root: Element) -> bytes:
for document in root.items:
bdata = document.bdata
if not bdata:
continue
try:
normalized_bdata = "".join(bdata.split())
return base64.b64decode(normalized_bdata, validate=True)
except (binascii.Error, ValueError) as exc:
raise CSEMappingError(
"CSE GetFormsForDocuments response contains invalid BData."
) from exc
raise CSEMappingError("CSE GetFormsForDocuments response is missing BData.")
def _parcel_type(request: DeliveryCalculationRequest) -> str | None:
return request.parcel_type.value if request.parcel_type is not None else None
@@ -189,3 +273,30 @@ def _compose_address(address: Address) -> str:
def _format_decimal(value: float) -> str:
return format(value, "g")
def _latest_tracking_status(document: Element) -> str | None:
dated: list[tuple[str, str]] = []
for state in document.items:
if not state.key:
continue
date_time = state.property_value("DateTime") or ""
dated.append((date_time, state.key))
if not dated:
return None
dated.sort(key=lambda item: item[0])
return dated[-1][1]
def _extract_waybill_number(document: Element) -> str | None:
for table in document.tables:
if table.key != "Waybills":
continue
for waybill in table.items:
document_type = waybill.property_value("DocumentType")
if document_type is not None and document_type.lower() != "waybill":
continue
number = waybill.property_value("Number") or waybill.key
if number:
return number
return None
+10 -7
View File
@@ -1,9 +1,10 @@
"""SOAP (cargo3) Element serialization and parsing for CSE.
CSE exposes a SOAP/1C web service ("Карго") that transfers data through nested
universal ``Element`` structures (``Key``/``Value``/``ValueType``/``Fields``/
``List``/``Tables``/``Properties``). This module builds request envelopes and
parses responses into a plain Python representation that mappers can navigate.
universal ``Element`` structures (``Key``/``Value``/``ValueType``/
``Properties``/``Fields``/``List``/``Tables``). This module builds request
envelopes and parses responses into a plain Python representation that mappers
can navigate.
"""
from __future__ import annotations
@@ -15,9 +16,8 @@ from xml.etree import ElementTree as ET
CARGO_NS = "http://www.cargo3.ru"
SOAP_NS = "http://www.w3.org/2003/05/soap-envelope"
# Child tags of an Element that hold nested Element lists.
_LIST_TAGS = ("Fields", "List", "Tables", "Properties")
_SCALAR_TAGS = ("Key", "Value", "ValueType")
_BINARY_TAG = "BData"
@dataclass
@@ -27,6 +27,7 @@ class Element:
key: str | None = None
value: str | None = None
value_type: str | None = None
bdata: str | None = None
fields: list["Element"] = field(default_factory=list)
items: list["Element"] = field(default_factory=list) # <List>
tables: list["Element"] = field(default_factory=list)
@@ -100,14 +101,14 @@ def _append_element(parent: ET.Element, tag: str, element: Element) -> None:
_scalar(node, "Value", element.value)
if element.value_type is not None:
_scalar(node, "ValueType", element.value_type)
for child in element.properties:
_append_element(node, "Properties", child)
for child in element.fields:
_append_element(node, "Fields", child)
for child in element.items:
_append_element(node, "List", child)
for child in element.tables:
_append_element(node, "Tables", child)
for child in element.properties:
_append_element(node, "Properties", child)
def _scalar(parent: ET.Element, tag: str, text: str) -> None:
@@ -159,6 +160,8 @@ def _parse_element(node: ET.Element) -> Element:
element.tables.append(_parse_element(child))
elif local == "Properties":
element.properties.append(_parse_element(child))
elif local == _BINARY_TAG:
element.bdata = (child.text or "").strip()
return element
@@ -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,
)
+36 -20
View File
@@ -38,26 +38,40 @@ 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 = ""
# CSE service GUIDs from GetReferenceData: Services that are allowed to be
# shown as delivery tariffs. Additional/non-delivery services must not be
# included here.
delivery_service_guids: list[str] = Field(default_factory=list)
# 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 +171,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 +211,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]
+15 -18
View File
@@ -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(
http_client=http_client,
adapter_config=settings.adapter,
)
cse_provider = CSEProvider.from_adapter_config(
delivery_provider_registry = build_delivery_provider_registry(
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,
+2 -2
View File
@@ -5,7 +5,7 @@ from enum import Enum
class TBankPaymentNotificationAction(str, Enum):
ACKNOWLEDGE_ONLY = "acknowledge_only"
REGISTER_CDEK_ORDER = "register_cdek_order"
REGISTER_PROVIDER_ORDER = "register_provider_order"
def resolve_tbank_payment_notification_action(
@@ -15,7 +15,7 @@ def resolve_tbank_payment_notification_action(
error_code: str,
) -> TBankPaymentNotificationAction:
if status == "CONFIRMED" and success is True and error_code == "0":
return TBankPaymentNotificationAction.REGISTER_CDEK_ORDER
return TBankPaymentNotificationAction.REGISTER_PROVIDER_ORDER
return TBankPaymentNotificationAction.ACKNOWLEDGE_ONLY
+7 -6
View File
@@ -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,
)
+35 -34
View File
@@ -6,7 +6,7 @@ from dataclasses import dataclass
from datetime import datetime
from typing import Any
from sqlalchemy import select
from sqlalchemy import and_, or_, select
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from app.domain.cdek_polling import TERMINAL_ORDER_STATUSES
@@ -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
@@ -109,17 +95,24 @@ class OrderRepository:
*,
limit: int,
) -> Sequence[Order]:
statement = (
select(Order)
.where(
Order.cdek_order_uuid.is_not(None),
Order.cdek_waybill_url.is_(None),
cdek_pending = and_(
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())
cse_pending = and_(
Order.provider == "cse",
Order.provider_order_id.is_not(None),
Order.provider_waybill_id.is_(None),
)
statement = (
select(Order)
.where(or_(cdek_pending, cse_pending))
.order_by(Order.provider_polled_at.asc().nulls_first())
.limit(limit)
.with_for_update(skip_locked=True)
)
@@ -139,10 +132,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 +151,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
@@ -186,10 +179,18 @@ class OrderRepository:
*,
limit: int,
) -> Sequence[Order]:
cdek_ready = and_(
Order.provider == "cdek",
Order.provider_waybill_url.is_not(None),
)
cse_ready = and_(
Order.provider == "cse",
Order.provider_waybill_id.is_not(None),
)
statement = (
select(Order)
.where(
Order.cdek_waybill_url.is_not(None),
or_(cdek_ready, cse_ready),
Order.waybill_email_sent_at.is_(None),
)
.order_by(Order.created_at.asc())
+1
View File
@@ -75,6 +75,7 @@ class Contact(_CamelModel):
class Content(_CamelModel):
description: str | None = None
declared_value: int = Field(alias="declared_value", strict=True)
class SystemDataTariff(_CamelModel):
+7 -25
View File
@@ -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,14 +635,7 @@ 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(
order = await self._order_repository.mark_provider_order_registered(
session,
order_uuid,
registration_id,
+27 -13
View File
@@ -1,7 +1,7 @@
"""Background service that e-mails CDEK waybill PDFs to customers."""
"""Background service that e-mails provider waybill PDFs to customers."""
import asyncio
from collections.abc import Callable, Sequence
from collections.abc import Callable, Mapping, Sequence
from contextlib import AbstractAsyncContextManager
from dataclasses import dataclass
from datetime import datetime, timezone
@@ -15,15 +15,18 @@ logger = structlog.get_logger(__name__)
_EMAIL_SUBJECT_TEMPLATE = "Накладная по заказу {order_uuid}"
_EMAIL_BODY_TEMPLATE = (
"Здравствуйте!\n\n"
"По вашему заказу {order_uuid} сформирована транспортная накладная CDEK.\n"
"По вашему заказу {order_uuid} сформирована "
"транспортная накладная.\n"
"PDF-файл накладной приложен к этому письму.\n"
)
_EMAIL_BODY_URL_LINE_TEMPLATE = (
"Также накладная доступна по ссылке: {waybill_url}\n"
)
_ATTACHMENT_FILENAME_TEMPLATE = "waybill_{order_uuid}.pdf"
class WaybillPDFDownloaderProtocol(Protocol):
async def download_waybill_pdf(self, url: str) -> bytes: ...
async def download_waybill_pdf(self, order: "OrderRecord") -> bytes: ...
class EmailSenderProtocol(Protocol):
@@ -40,8 +43,10 @@ class EmailSenderProtocol(Protocol):
class OrderRecord(Protocol):
order_uuid: str
provider: str
account_email: str
cdek_waybill_url: str | None
provider_waybill_id: str | None
provider_waybill_url: str | None
class WaybillEmailSenderRepositoryProtocol(Protocol):
@@ -72,13 +77,13 @@ class WaybillEmailSenderService:
self,
*,
order_repository: WaybillEmailSenderRepositoryProtocol,
waybill_downloader: WaybillPDFDownloaderProtocol,
waybill_downloaders: Mapping[str, WaybillPDFDownloaderProtocol],
email_sender: EmailSenderProtocol,
batch_size: int,
datetime_now: Callable[[], datetime] = lambda: datetime.now(timezone.utc),
) -> None:
self._repository = order_repository
self._waybill_downloader = waybill_downloader
self._waybill_downloaders = dict(waybill_downloaders)
self._email_sender = email_sender
self._batch_size = batch_size
self._datetime_now = datetime_now
@@ -99,6 +104,7 @@ class WaybillEmailSenderService:
logger.exception(
"waybill_email_order_failed",
order_uuid=order.order_uuid,
provider=order.provider,
account_email=order.account_email,
)
return SendBatchSummary(
@@ -133,15 +139,16 @@ class WaybillEmailSenderService:
continue
async def _handle_order(self, session: object, order: OrderRecord) -> None:
waybill_url = order.cdek_waybill_url
if waybill_url is None:
if order.provider_waybill_url is None and order.provider_waybill_id is None:
return
pdf_bytes = await self._waybill_downloader.download_waybill_pdf(waybill_url)
downloader = self._resolve_downloader(order.provider)
pdf_bytes = await downloader.download_waybill_pdf(order)
subject = _EMAIL_SUBJECT_TEMPLATE.format(order_uuid=order.order_uuid)
body = _EMAIL_BODY_TEMPLATE.format(
order_uuid=order.order_uuid,
waybill_url=waybill_url,
body = _EMAIL_BODY_TEMPLATE.format(order_uuid=order.order_uuid)
if order.provider_waybill_url is not None:
body += _EMAIL_BODY_URL_LINE_TEMPLATE.format(
waybill_url=order.provider_waybill_url,
)
filename = _ATTACHMENT_FILENAME_TEMPLATE.format(order_uuid=order.order_uuid)
@@ -162,6 +169,13 @@ class WaybillEmailSenderService:
logger.info(
"waybill_email_sent",
order_uuid=order.order_uuid,
provider=order.provider,
account_email=order.account_email,
sent_at=sent_at,
)
def _resolve_downloader(self, provider: str) -> WaybillPDFDownloaderProtocol:
downloader = self._waybill_downloaders.get(provider)
if downloader is None:
raise RuntimeError(f"Waybill downloader is not configured for {provider}.")
return downloader
+70 -32
View File
@@ -1,7 +1,7 @@
"""Background service that polls CDEK for waybill updates."""
"""Background service that polls providers for waybill updates."""
import asyncio
from collections.abc import Callable, Sequence
from collections.abc import Callable, Mapping, Sequence
from contextlib import AbstractAsyncContextManager
from dataclasses import dataclass
from datetime import datetime, timezone
@@ -9,26 +9,22 @@ from typing import Protocol
import structlog
from app.adapters.delivery_providers.cdek.order_mapper import (
CDEKOrderInfo,
CDEKWaybillInfo,
)
logger = structlog.get_logger(__name__)
class CDEKOrderInfoAdapterProtocol(Protocol):
async def get_order(self, cdek_order_uuid: str) -> CDEKOrderInfo: ...
class OrderInfoAdapterProtocol(Protocol):
async def get_order(self, provider_order_id: str) -> object: ...
class CDEKWaybillInfoAdapterProtocol(Protocol):
async def get_waybill(self, cdek_waybill_uuid: str) -> CDEKWaybillInfo: ...
class WaybillInfoAdapterProtocol(Protocol):
async def get_waybill(self, provider_waybill_id: str) -> object: ...
class OrderRecord(Protocol):
order_uuid: str
cdek_order_uuid: str | None
cdek_waybill_uuid: str | None
provider: str
provider_order_id: str | None
provider_waybill_id: str | None
class WaybillPollerRepositoryProtocol(Protocol):
@@ -70,14 +66,14 @@ class WaybillPollerService:
self,
*,
order_repository: WaybillPollerRepositoryProtocol,
order_info_adapter: CDEKOrderInfoAdapterProtocol,
waybill_info_adapter: CDEKWaybillInfoAdapterProtocol,
order_info_adapters: Mapping[str, OrderInfoAdapterProtocol],
waybill_info_adapters: Mapping[str, WaybillInfoAdapterProtocol],
batch_size: int,
datetime_now: Callable[[], datetime] = lambda: datetime.now(timezone.utc),
) -> None:
self._repository = order_repository
self._order_info_adapter = order_info_adapter
self._waybill_info_adapter = waybill_info_adapter
self._order_info_adapters = dict(order_info_adapters)
self._waybill_info_adapters = dict(waybill_info_adapters)
self._batch_size = batch_size
self._datetime_now = datetime_now
@@ -97,8 +93,9 @@ 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.provider,
provider_order_id=order.provider_order_id,
provider_waybill_id=order.provider_waybill_id,
)
return PollBatchSummary(
processed=len(orders),
@@ -133,37 +130,78 @@ 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)
adapter = self._resolve_order_info_adapter(order.provider)
info = await adapter.get_order(provider_order_id)
status_code = _extract_status_code(info)
waybill_id = _extract_waybill_id(info)
await self._repository.record_order_poll(
session,
order_uuid=order.order_uuid,
order_status=info.status_code,
waybill_uuid=info.waybill_uuid,
order_status=status_code,
waybill_uuid=waybill_id,
polled_at=polled_at,
)
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.provider,
provider_order_id=provider_order_id,
provider_order_status=status_code,
provider_waybill_id=waybill_id,
)
return
waybill = await self._waybill_info_adapter.get_waybill(order.cdek_waybill_uuid)
adapter = self._resolve_waybill_info_adapter(order.provider)
waybill = await adapter.get_waybill(order.provider_waybill_id)
waybill_url = _extract_waybill_url(waybill)
await self._repository.record_waybill_poll(
session,
order_uuid=order.order_uuid,
waybill_url=waybill.url,
waybill_url=waybill_url,
polled_at=polled_at,
)
logger.info(
"waybill_poll_waybill_result",
order_uuid=order.order_uuid,
cdek_waybill_uuid=order.cdek_waybill_uuid,
cdek_waybill_url=waybill.url,
provider=order.provider,
provider_waybill_id=order.provider_waybill_id,
provider_waybill_url=waybill_url,
)
def _resolve_order_info_adapter(self, provider: str) -> OrderInfoAdapterProtocol:
adapter = self._order_info_adapters.get(provider)
if adapter is None:
raise RuntimeError(f"Order info adapter is not configured for {provider}.")
return adapter
def _resolve_waybill_info_adapter(
self, provider: str
) -> WaybillInfoAdapterProtocol:
adapter = self._waybill_info_adapters.get(provider)
if adapter is None:
raise RuntimeError(
f"Waybill info adapter is not configured for {provider}."
)
return adapter
def _extract_status_code(info: object) -> str | None:
value = getattr(info, "status_code", None)
return value if isinstance(value, str) and value else None
def _extract_waybill_id(info: object) -> str | None:
for attr in ("waybill_id", "waybill_uuid", "waybill_number"):
value = getattr(info, attr, None)
if isinstance(value, str) and value:
return value
return None
def _extract_waybill_url(info: object) -> str | None:
value = getattr(info, "url", None)
return value if isinstance(value, str) and value else None
+67 -11
View File
@@ -1,13 +1,21 @@
"""Background worker that e-mails CDEK waybill PDFs."""
"""Background worker that e-mails provider waybill PDFs."""
import asyncio
import signal
from dataclasses import dataclass
import httpx
import structlog
from app.adapters.delivery_providers.registry import (
resolve_delivery_provider_timeout_seconds,
)
from app.adapters.delivery_providers.cdek.auth import CDEKAuthClient
from app.adapters.delivery_providers.cdek.client import CDEKClient
from app.adapters.delivery_providers.cse.client import CSEClient
from app.adapters.delivery_providers.cse.order_mapper import (
CSEOrderRegistrationParams,
)
from app.adapters.email import SMTPEmailSender
from app.adapters.postgres.engine import (
create_postgres_engine,
@@ -21,23 +29,71 @@ from app.services.waybill_email_sender import WaybillEmailSenderService
logger = structlog.get_logger(__name__)
@dataclass(frozen=True)
class CDEKWaybillPDFDownloader:
client: CDEKClient
async def download_waybill_pdf(self, order: object) -> bytes:
url = getattr(order, "provider_waybill_url", None)
if not isinstance(url, str) or not url:
raise RuntimeError("CDEK waybill URL is missing.")
return await self.client.download_waybill_pdf(url)
@dataclass(frozen=True)
class CSEWaybillPDFDownloader:
client: CSEClient
async def download_waybill_pdf(self, order: object) -> bytes:
waybill_number = getattr(order, "provider_waybill_id", None)
if not isinstance(waybill_number, str) or not waybill_number:
raise RuntimeError("CSE waybill number is missing.")
return await self.client.download_waybill_pdf(waybill_number)
async def _run(settings: Settings, stop_event: asyncio.Event) -> None:
http_client = httpx.AsyncClient(timeout=settings.adapter.cdek_timeout_seconds)
http_client = httpx.AsyncClient(
timeout=resolve_delivery_provider_timeout_seconds(settings.delivery_providers)
)
waybill_downloaders: dict[str, object] = {}
cdek_config = settings.delivery_providers.cdek
if cdek_config.enabled:
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,
)
waybill_downloaders["cdek"] = CDEKWaybillPDFDownloader(cdek_client)
cse_config = settings.delivery_providers.cse
if cse_config.enabled:
cse_client = CSEClient(
http_client=http_client,
base_url=cse_config.base_url,
login=cse_config.login,
password=cse_config.password,
registration_params=CSEOrderRegistrationParams(
payer=cse_config.payer,
payment_method=cse_config.payment_method,
shipping_method=cse_config.shipping_method,
),
timeout_seconds=cse_config.timeout_seconds,
retry_attempts=cse_config.retry_attempts,
retry_backoff_seconds=cse_config.retry_backoff_seconds,
)
waybill_downloaders["cse"] = CSEWaybillPDFDownloader(cse_client)
email_sender = SMTPEmailSender(
smtp_host=settings.email.smtp_host,
smtp_port=settings.email.smtp_port,
@@ -52,7 +108,7 @@ async def _run(settings: Settings, stop_event: asyncio.Event) -> None:
repository = OrderRepository(session_factory=session_factory)
service = WaybillEmailSenderService(
order_repository=repository,
waybill_downloader=cdek_client,
waybill_downloaders=waybill_downloaders,
email_sender=email_sender,
batch_size=settings.waybill_email_sender.batch_size,
)
+47 -12
View File
@@ -1,4 +1,4 @@
"""Background worker that polls CDEK for waybill updates."""
"""Background worker that polls delivery providers for waybill updates."""
import asyncio
import signal
@@ -6,8 +6,15 @@ import signal
import httpx
import structlog
from app.adapters.delivery_providers.registry import (
resolve_delivery_provider_timeout_seconds,
)
from app.adapters.delivery_providers.cdek.auth import CDEKAuthClient
from app.adapters.delivery_providers.cdek.client import CDEKClient
from app.adapters.delivery_providers.cse.client import CSEClient
from app.adapters.delivery_providers.cse.order_mapper import (
CSEOrderRegistrationParams,
)
from app.adapters.postgres.engine import (
create_postgres_engine,
create_postgres_session_factory,
@@ -21,29 +28,57 @@ 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)
http_client = httpx.AsyncClient(
timeout=resolve_delivery_provider_timeout_seconds(settings.delivery_providers)
)
order_info_adapters: dict[str, object] = {}
waybill_info_adapters: dict[str, object] = {}
cdek_config = settings.delivery_providers.cdek
if cdek_config.enabled:
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,
)
order_info_adapters["cdek"] = cdek_client
waybill_info_adapters["cdek"] = cdek_client
cse_config = settings.delivery_providers.cse
if cse_config.enabled:
cse_client = CSEClient(
http_client=http_client,
base_url=cse_config.base_url,
login=cse_config.login,
password=cse_config.password,
registration_params=CSEOrderRegistrationParams(
payer=cse_config.payer,
payment_method=cse_config.payment_method,
shipping_method=cse_config.shipping_method,
),
timeout_seconds=cse_config.timeout_seconds,
retry_attempts=cse_config.retry_attempts,
retry_backoff_seconds=cse_config.retry_backoff_seconds,
)
order_info_adapters["cse"] = cse_client
engine = create_postgres_engine(settings.postgres)
session_factory = create_postgres_session_factory(engine)
repository = OrderRepository(session_factory=session_factory)
service = WaybillPollerService(
order_repository=repository,
order_info_adapter=cdek_client,
waybill_info_adapter=cdek_client,
order_info_adapters=order_info_adapters,
waybill_info_adapters=waybill_info_adapters,
batch_size=settings.waybill_poller.batch_size,
)
logger.info(
+24 -18
View File
@@ -14,24 +14,30 @@ repository:
redis_dsn: "redis://redis:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.edu.cdek.ru/v2"
cdek_client_id: "${CDEK_CLIENT_ID}"
cdek_client_secret: "${CDEK_CLIENT_SECRET}"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
cse_base_url: "http://lk-test.cse.ru/1c/ws/web1c.1cws"
cse_login: "${CSE_LOGIN}"
cse_password: "${CSE_PASSWORD}"
cse_retry_attempts: 2
cse_retry_backoff_seconds: 0.2
cse_timeout_seconds: 10.0
cse_cache_ttl_seconds: 900
cse_payer: "0" # Заказчик
cse_payment_method: "1" # Безналичный расчёт
cse_shipping_method: "5052d0b3-5ea3-46f2-823f-1472686a51dd" # авто
delivery_providers:
cdek:
enabled: true
base_url: "https://api.edu.cdek.ru/v2"
client_id: "${CDEK_CLIENT_ID}"
client_secret: "${CDEK_CLIENT_SECRET}"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: true
base_url: "http://lk-test.cse.ru/1c/ws/web1c.1cws"
login: "${CSE_LOGIN}"
password: "${CSE_PASSWORD}"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
delivery_service_guids:
- "6da21fe8-4f13-11dc-bda1-0015170f8c09" # Россия доставка
payer: "0" # Заказчик
payment_method: "1" # Безналичный расчёт
shipping_method: "5052d0b3-5ea3-46f2-823f-1472686a51dd" # авто
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
+24 -18
View File
@@ -14,24 +14,30 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
cse_base_url: "http://lk-test.cse.ru/1c/ws/web1c.1cws"
cse_login: "test"
cse_password: "2016"
cse_retry_attempts: 2
cse_retry_backoff_seconds: 0.2
cse_timeout_seconds: 10.0
cse_cache_ttl_seconds: 900
cse_payer: "0"
cse_payment_method: "1"
cse_shipping_method: "5052d0b3-5ea3-46f2-823f-1472686a51dd"
delivery_providers:
cdek:
enabled: true
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: true
base_url: "http://lk-test.cse.ru/1c/ws/web1c.1cws"
login: "test"
password: "2016"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
delivery_service_guids:
- "tariff-guid-2"
payer: "0"
payment_method: "1"
shipping_method: "5052d0b3-5ea3-46f2-823f-1472686a51dd"
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -13,7 +13,7 @@ from app.adapters.delivery_providers.cdek.client import (
)
from app.cities import cities_map
from app import config as config_module
from app.config import AdapterConfig, Settings
from app.config import CDEKDeliveryProviderConfig, Settings
from app.schemas.request import DeliveryCalculationRequest, DeliveryEntity
@@ -287,14 +287,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.test/v2"
cdek_client_id: "yaml-id"
cdek_client_secret: "yaml-secret"
cdek_retry_attempts: 0
cdek_retry_backoff_seconds: 0.1
cdek_timeout_seconds: 7.5
cdek_cache_ttl_seconds: 777
delivery_providers:
cdek:
base_url: "https://api.cdek.test/v2"
client_id: "yaml-id"
client_secret: "yaml-secret"
retry_attempts: 0
retry_backoff_seconds: 0.1
timeout_seconds: 7.5
cache_ttl_seconds: 777
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -323,9 +326,9 @@ email:
monkeypatch.setattr(config_module, "_resolve_runtime_config_file", lambda: str(config_file))
settings = Settings()
http_client = RecordingHTTPClient()
provider = CDEKProvider.from_adapter_config(
provider = CDEKProvider.from_config(
http_client=http_client, # type: ignore[arg-type]
adapter_config=settings.adapter,
config=settings.delivery_providers.cdek,
)
result = asyncio.run(provider.get_prices(_make_request()))
@@ -349,14 +352,14 @@ email:
def test_provider_uses_default_adapter_timeout_and_cache_ttl() -> None:
adapter_config = AdapterConfig(
cdek_client_id="default-id",
cdek_client_secret="default-secret",
config = CDEKDeliveryProviderConfig(
client_id="default-id",
client_secret="default-secret",
)
http_client = RecordingHTTPClient()
provider = CDEKProvider.from_adapter_config(
provider = CDEKProvider.from_config(
http_client=http_client, # type: ignore[arg-type]
adapter_config=adapter_config,
config=config,
)
result = asyncio.run(provider.get_prices(_make_request()))
@@ -29,6 +29,23 @@ def test_order_payload_requests_waybill_print() -> None:
assert payload["print"] == "WAYBILL"
def test_order_payload_maps_declared_value_to_insurance_service() -> None:
payload = map_cdek_order_request(make_init_payment_request(), "order-uuid-1")
assert payload["services"] == [{"code": "INSURANCE", "parameter": 50000}]
def test_order_payload_omits_insurance_service_for_zero_declared_value() -> None:
payload = map_cdek_order_request(
make_init_payment_request(
content={"description": "Docs", "declared_value": 0}
),
"order-uuid-1",
)
assert "services" not in payload
def test_map_cdek_order_response_extracts_order_uuid_without_waybill() -> None:
result = map_cdek_order_response({"entity": {"uuid": "cdek-order-uuid"}})
@@ -143,7 +160,7 @@ def test_order_payload_for_parcel_includes_dimensions() -> None:
def test_order_payload_for_doc_omits_dimensions() -> None:
request = make_init_payment_request(
content={"description": "Docs"},
content={"description": "Docs", "declared_value": 25000},
systemData={
"tariff": {
"provider": "СДЭК",
@@ -198,7 +215,7 @@ def test_order_payload_includes_company_requisites_for_legal_entity() -> None:
assert sender["kpp"] == "770701001"
def test_order_payload_propagates_description_to_comments_without_items() -> None:
def test_order_payload_maps_content_to_comments_without_items() -> None:
payload = map_cdek_order_request(make_init_payment_request(), "order-uuid-1")
assert payload["comment"] == "Headphones"
@@ -208,10 +225,14 @@ def test_order_payload_propagates_description_to_comments_without_items() -> Non
def test_order_payload_falls_back_package_comment_to_order_uuid() -> None:
payload = map_cdek_order_request(
make_init_payment_request(content={"description": None}), "order-uuid-1"
make_init_payment_request(
content={"description": None, "declared_value": 50000}
),
"order-uuid-1",
)
assert payload["packages"][0]["comment"] == "order-uuid-1"
assert "items" not in payload["packages"][0]
def test_order_payload_includes_company_for_individual_sender_as_full_name() -> None:
@@ -272,6 +293,37 @@ def test_map_cdek_order_info_response_returns_none_status_when_statuses_empty()
)
def test_map_cdek_order_info_response_maps_invalid_create_request_to_invalid_status() -> None:
result = map_cdek_order_info_response(
{
"entity": {
"uuid": "cdek-order-uuid",
"statuses": [
{"code": "ACCEPTED", "date_time": "2026-05-24T10:00:00+0000"}
],
},
"requests": [
{
"type": "CREATE",
"state": "INVALID",
"errors": [
{
"code": "ve_delivery_can_not_has_goods",
"message": "Заказ типа доставка не может содержать товары",
}
],
}
],
}
)
assert result == CDEKOrderInfo(
order_uuid="cdek-order-uuid",
status_code="INVALID",
waybill_uuid=None,
)
def test_map_cdek_order_info_response_raises_for_missing_entity() -> None:
with pytest.raises(CDEKOrderMappingError):
map_cdek_order_info_response({"requests": []})
@@ -124,7 +124,7 @@ def test_provider_get_payment_price_omits_dimensions_for_doc() -> None:
)
)
request = make_init_payment_request(
content={"description": "Docs"},
content={"description": "Docs", "declared_value": 25000},
systemData={
"tariff": {
"provider": "СДЭК",
@@ -7,6 +7,7 @@ import pytest
from app.adapters.delivery_providers.cse import (
CSEClient,
CSEClientError,
CSEProvider,
CSERequestError,
)
@@ -70,6 +71,162 @@ _SAVE_RESPONSE = """<?xml version="1.0" encoding="UTF-8"?>
</soap:Body>
</soap:Envelope>"""
_SAVE_RESPONSE_DOCUMENT_ERROR = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
<m:SaveDocumentsResponse xmlns:m="http://www.cargo3.ru">
<m:return>
<m:Key>SaveDocuments</m:Key>
<m:List>
<m:Key>Order</m:Key>
<m:Properties>
<m:Key>Error</m:Key>
<m:Value>true</m:Value>
<m:ValueType>boolean</m:ValueType>
<m:List>
<m:Key>Description</m:Key>
<m:Value>SenderAddress is invalid</m:Value>
</m:List>
</m:Properties>
</m:List>
</m:return>
</m:SaveDocumentsResponse>
</soap:Body>
</soap:Envelope>"""
_TRACKING_RESPONSE = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
<m:TrackingResponse xmlns:m="http://www.cargo3.ru">
<m:return>
<m:Key>Tracking</m:Key>
<m:List>
<m:Key>CSE-000123</m:Key>
<m:Value>Order</m:Value>
<m:Properties><m:Key>Number</m:Key><m:Value>CSE-000123</m:Value></m:Properties>
<m:List>
<m:Key>Заказ принят, идет обработка заказа.</m:Key>
<m:Properties>
<m:Key>DateTime</m:Key><m:Value>2026-06-01T10:00:00</m:Value>
</m:Properties>
</m:List>
<m:List>
<m:Key>Накладная оформлена.</m:Key>
<m:Properties>
<m:Key>DateTime</m:Key><m:Value>2026-06-01T10:02:00</m:Value>
</m:Properties>
</m:List>
<m:Tables>
<m:Key>Waybills</m:Key>
<m:List>
<m:Key>496-AA-1676378</m:Key>
<m:Properties>
<m:Key>DocumentType</m:Key><m:Value>Waybill</m:Value>
</m:Properties>
<m:Properties>
<m:Key>Number</m:Key><m:Value>496-AA-1676378</m:Value>
</m:Properties>
</m:List>
</m:Tables>
</m:List>
</m:return>
</m:TrackingResponse>
</soap:Body>
</soap:Envelope>"""
_FORM_RESPONSE = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
<m:GetFormsForDocumentsResponse xmlns:m="http://www.cargo3.ru">
<m:return>
<m:Key>GetPrintForms</m:Key>
<m:List>
<m:Key>496-AA-1676378</m:Key>
<m:Properties><m:Key>FormFormat</m:Key><m:Value>PDF</m:Value></m:Properties>
<m:BData>JVBERg==</m:BData>
</m:List>
</m:return>
</m:GetFormsForDocumentsResponse>
</soap:Body>
</soap:Envelope>"""
_CALC_RESPONSE_WITH_ADDITIONAL_SERVICE_FIRST = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
<m:CalcResponse xmlns:m="http://www.cargo3.ru">
<m:return>
<m:Key>Calc</m:Key>
<m:List>
<m:Key>Destination</m:Key>
<m:List>
<m:Key>Tariff</m:Key>
<m:Value>additional-service-guid</m:Value>
<m:Fields><m:Key>Total</m:Key><m:Value>100.00</m:Value></m:Fields>
<m:Fields><m:Key>CurrencyName</m:Key><m:Value>RUR</m:Value></m:Fields>
<m:Fields>
<m:Key>Service</m:Key><m:Value>Доп. услуга</m:Value>
</m:Fields>
<m:Fields><m:Key>Urgency</m:Key><m:Value>urg-exp</m:Value></m:Fields>
<m:Fields>
<m:Key>AdditionalService</m:Key><m:Value>true</m:Value>
</m:Fields>
</m:List>
<m:List>
<m:Key>Tariff</m:Key>
<m:Value>delivery-guid</m:Value>
<m:Fields><m:Key>Total</m:Key><m:Value>793.00</m:Value></m:Fields>
<m:Fields><m:Key>CurrencyName</m:Key><m:Value>RUR</m:Value></m:Fields>
<m:Fields><m:Key>Service</m:Key><m:Value>Экспресс</m:Value></m:Fields>
<m:Fields><m:Key>Urgency</m:Key><m:Value>urg-exp</m:Value></m:Fields>
<m:Fields>
<m:Key>AdditionalService</m:Key><m:Value>false</m:Value>
</m:Fields>
<m:Fields><m:Key>MinPeriod</m:Key><m:Value>1</m:Value></m:Fields>
<m:Fields><m:Key>MaxPeriod</m:Key><m:Value>2</m:Value></m:Fields>
</m:List>
</m:List>
</m:return>
</m:CalcResponse>
</soap:Body>
</soap:Envelope>"""
_CALC_RESPONSE_WITH_BUYOUT_SERVICE_FIRST = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
<m:CalcResponse xmlns:m="http://www.cargo3.ru">
<m:return>
<m:Key>Calc</m:Key>
<m:List>
<m:Key>Destination</m:Key>
<m:List>
<m:Key>Tariff</m:Key>
<m:Value>buyout-service-guid</m:Value>
<m:Fields><m:Key>Total</m:Key><m:Value>100.00</m:Value></m:Fields>
<m:Fields><m:Key>CurrencyName</m:Key><m:Value>RUR</m:Value></m:Fields>
<m:Fields>
<m:Key>Service</m:Key><m:Value>Частичный выкуп</m:Value>
</m:Fields>
<m:Fields><m:Key>Urgency</m:Key><m:Value>urg-exp</m:Value></m:Fields>
</m:List>
<m:List>
<m:Key>Tariff</m:Key>
<m:Value>delivery-guid</m:Value>
<m:Fields><m:Key>Total</m:Key><m:Value>793.00</m:Value></m:Fields>
<m:Fields><m:Key>CurrencyName</m:Key><m:Value>RUR</m:Value></m:Fields>
<m:Fields><m:Key>Service</m:Key><m:Value>Экспресс</m:Value></m:Fields>
<m:Fields><m:Key>Urgency</m:Key><m:Value>urg-exp</m:Value></m:Fields>
<m:Fields>
<m:Key>AdditionalService</m:Key><m:Value>false</m:Value>
</m:Fields>
<m:Fields><m:Key>MinPeriod</m:Key><m:Value>1</m:Value></m:Fields>
<m:Fields><m:Key>MaxPeriod</m:Key><m:Value>2</m:Value></m:Fields>
</m:List>
</m:List>
</m:return>
</m:CalcResponse>
</soap:Body>
</soap:Envelope>"""
_ERROR_RESPONSE = """<?xml version="1.0" encoding="UTF-8"?>
<soap:Envelope xmlns:soap="http://www.w3.org/2003/05/soap-envelope">
<soap:Body>
@@ -181,6 +338,37 @@ def test_build_envelope_contains_credentials_and_payload() -> None:
assert "geo-2" in text
def test_save_order_envelope_serializes_properties_before_fields() -> None:
body = order_mapper.map_cse_save_order_request(
make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {"length": "10", "width": "10", "height": "10"},
}
),
"order-uuid-1",
_params(),
)
text = build_envelope(
"SaveDocuments", login="test", password="2016", body=body
).decode("utf-8")
order_start = text.index("<m:Key>Order</m:Key>")
properties_index = text.index("<m:Properties>", order_start)
fields_index = text.index("<m:Fields>", order_start)
assert properties_index < fields_index
def test_parse_calc_response_round_trip() -> None:
root = parse_response(_CALC_RESPONSE, "Calc")
assert root.key == "Calc"
@@ -192,19 +380,23 @@ def test_provider_get_prices_maps_all_tariffs() -> None:
[
httpx.Response(200, text=_DELIVERY_TYPES_RESPONSE),
httpx.Response(200, text=_CALC_RESPONSE),
httpx.Response(200, text=_CALC_RESPONSE),
]
)
provider = CSEProvider(client)
provider = CSEProvider(
client,
delivery_service_guids=("tariff-guid-1", "tariff-guid-2"),
)
prices = asyncio.run(provider.get_prices(_calc_request()))
# PVZ-requiring scheme (Склад-Дверь) is excluded: only delivery types + one
# Calc (for door-to-door) are requested.
assert len(http_client.calls) == 2
# Calc per configured CSE delivery service are requested.
assert len(http_client.calls) == 3
assert [price.provider for price in prices] == ["cse", "cse"]
assert [price.tariff_code for price in prices] == [
"ДоставкаДоДверей|urg-std",
"ДоставкаДоДверей|urg-exp",
"ДоставкаДоДверей|tariff-guid-1|urg-std",
"ДоставкаДоДверей|tariff-guid-2|urg-exp",
]
assert prices[0].price == Decimal("1114.92")
assert prices[0].currency == "RUB"
@@ -215,7 +407,7 @@ def test_provider_get_prices_maps_all_tariffs() -> None:
def test_provider_get_payment_price_selects_matching_tariff() -> None:
response = httpx.Response(200, text=_CALC_RESPONSE)
client, _ = _build_client([response])
client, http_client = _build_client([response])
provider = CSEProvider(client)
request = make_init_payment_request(
systemData={
@@ -225,7 +417,7 @@ def test_provider_get_payment_price_selects_matching_tariff() -> None:
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|urg-exp",
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
@@ -236,8 +428,110 @@ def test_provider_get_payment_price_selects_matching_tariff() -> None:
price = asyncio.run(provider.get_payment_price(request))
assert price is not None
assert price.tariff_code == "ДоставкаДоДверей|urg-exp"
assert price.tariff_code == "ДоставкаДоДверей|tariff-guid-2|urg-exp"
assert price.price == Decimal("2000.00")
content = http_client.calls[0]["content"].decode("utf-8")
assert "Service" in content
assert "tariff-guid-2" in content
def test_provider_get_payment_price_ignores_matching_additional_service() -> None:
response = httpx.Response(200, text=_CALC_RESPONSE_WITH_ADDITIONAL_SERVICE_FIRST)
client, _ = _build_client([response])
provider = CSEProvider(client)
request = make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 79300,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|delivery-guid|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {"length": "10", "width": "10", "height": "10"},
}
)
price = asyncio.run(provider.get_payment_price(request))
assert price is not None
assert price.tariff_code == "ДоставкаДоДверей|delivery-guid|urg-exp"
assert price.price == Decimal("793.00")
def test_provider_get_payment_price_selects_matching_service_guid() -> None:
response = httpx.Response(200, text=_CALC_RESPONSE_WITH_BUYOUT_SERVICE_FIRST)
client, _ = _build_client([response])
provider = CSEProvider(client)
request = make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 79300,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|delivery-guid|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {"length": "10", "width": "10", "height": "10"},
}
)
price = asyncio.run(provider.get_payment_price(request))
assert price is not None
assert price.tariff_code == "ДоставкаДоДверей|delivery-guid|urg-exp"
assert price.price == Decimal("793.00")
def test_provider_get_prices_excludes_additional_service_tariffs() -> None:
client, _ = _build_client(
[
httpx.Response(200, text=_DELIVERY_TYPES_RESPONSE),
httpx.Response(200, text=_CALC_RESPONSE_WITH_ADDITIONAL_SERVICE_FIRST),
]
)
provider = CSEProvider(client, delivery_service_guids=("delivery-guid",))
prices = asyncio.run(provider.get_prices(_calc_request()))
assert [price.price for price in prices] == [Decimal("793.00")]
assert [price.service_name for price in prices] == [
"Экспресс — Дверь-Дверь"
]
def test_provider_get_prices_uses_configured_service_guid_filter() -> None:
client, http_client = _build_client(
[
httpx.Response(200, text=_DELIVERY_TYPES_RESPONSE),
httpx.Response(200, text=_CALC_RESPONSE_WITH_BUYOUT_SERVICE_FIRST),
]
)
provider = CSEProvider(client, delivery_service_guids=("delivery-guid",))
prices = asyncio.run(provider.get_prices(_calc_request()))
assert [price.price for price in prices] == [Decimal("793.00")]
assert [price.service_name for price in prices] == [
"Экспресс — Дверь-Дверь"
]
content = http_client.calls[1]["content"].decode("utf-8")
assert "Service" in content
assert "delivery-guid" in content
def test_provider_get_prices_requires_delivery_service_guid_configuration() -> None:
client, _ = _build_client([httpx.Response(200, text=_DELIVERY_TYPES_RESPONSE)])
provider = CSEProvider(client)
with pytest.raises(CSERequestError):
asyncio.run(provider.get_prices(_calc_request()))
def test_provider_register_order_returns_document_number() -> None:
@@ -246,13 +540,96 @@ def test_provider_register_order_returns_document_number() -> None:
provider = CSEProvider(client)
result = asyncio.run(
provider.register_order(make_init_payment_request(), "order-uuid-1")
provider.register_order(
make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {"length": "10", "width": "10", "height": "10"},
}
),
"order-uuid-1",
)
)
assert result.order_number == "CSE-000123"
assert http_client.calls[0]["url"] == "http://lk-test.cse.ru/1c/ws/web1c.1cws"
def test_provider_register_order_raises_document_level_error() -> None:
response = httpx.Response(200, text=_SAVE_RESPONSE_DOCUMENT_ERROR)
client, _ = _build_client([response])
provider = CSEProvider(client)
with pytest.raises(CSERequestError, match="SenderAddress is invalid"):
asyncio.run(
provider.register_order(
make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {
"length": "10",
"width": "10",
"height": "10",
},
}
),
"order-uuid-1",
)
)
def test_provider_get_order_extracts_waybill_number_from_tracking() -> None:
response = httpx.Response(200, text=_TRACKING_RESPONSE)
client, http_client = _build_client([response])
provider = CSEProvider(client)
result = asyncio.run(provider.get_order("CSE-000123"))
assert result.order_number == "CSE-000123"
assert result.status_code == "Накладная оформлена."
assert result.waybill_number == "496-AA-1676378"
content = http_client.calls[0]["content"].decode("utf-8")
assert "Tracking" in content
assert "DocumentType" in content
assert "Order" in content
def test_provider_download_waybill_pdf_uses_print_form() -> None:
response = httpx.Response(200, text=_FORM_RESPONSE)
client, http_client = _build_client([response])
provider = CSEProvider(client)
pdf = asyncio.run(provider.download_waybill_pdf("496-AA-1676378"))
assert pdf == b"%PDF"
content = http_client.calls[0]["content"].decode("utf-8")
assert "GetFormsForDocuments" in content
assert "DocumentType" in content
assert "waybill" in content
assert "Type" in content
assert "print" in content
assert "Format" in content
assert "pdf" in content
def test_save_order_request_uses_selected_tariff_urgency() -> None:
request = make_init_payment_request(
systemData={
@@ -262,7 +639,7 @@ def test_save_order_request_uses_selected_tariff_urgency() -> None:
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|urg-exp",
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
@@ -276,6 +653,32 @@ def test_save_order_request_uses_selected_tariff_urgency() -> None:
assert order.field_value("DeliveryOfCargo") == "ДоставкаДоДверей"
assert order.field_value("Payer") == "payer-1"
assert order.field_value("ShippingMethod") == "sm-1"
assert order.field_value("CargoCost") == "50000"
def test_save_order_request_uses_pickup_date_plus_seven_days_as_delivery_date() -> None:
request = make_init_payment_request(
deliveryDate="2026-05-18T18:00:00.000Z",
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {"length": "10", "width": "10", "height": "10"},
},
)
body = order_mapper.map_cse_save_order_request(request, "order-uuid-1", _params())
order = body["data"].items[0]
assert order.field_value("TakeDate") == "2026-05-15T10:00:00"
assert order.field_value("DeliveryDate") == "2026-05-22T10:00:00"
def test_application_error_in_response_raises_request_error() -> None:
@@ -285,7 +688,63 @@ def test_application_error_in_response_raises_request_error() -> None:
httpx.Response(200, text=_ERROR_RESPONSE),
]
)
provider = CSEProvider(client)
provider = CSEProvider(client, delivery_service_guids=("delivery-guid",))
with pytest.raises(CSERequestError):
asyncio.run(provider.get_prices(_calc_request()))
def test_server_error_logs_response_excerpt(monkeypatch: pytest.MonkeyPatch) -> None:
warning_calls: list[dict[str, Any]] = []
def capture_warning(event: str, **kwargs: Any) -> None:
warning_calls.append({"event": event, **kwargs})
monkeypatch.setattr(
"app.adapters.delivery_providers.cse.client.log.warning",
capture_warning,
)
response = httpx.Response(
500,
text="<fault>\n CSE rejected SaveDocuments because field X is invalid \n</fault>",
)
client, _ = _build_client([response])
provider = CSEProvider(client)
with pytest.raises(CSEClientError):
asyncio.run(
provider.register_order(
make_init_payment_request(
systemData={
"tariff": {
"provider": "cse",
"serviceName": "Экспресс",
"price": 200000,
"deliveryDaysMin": 1,
"deliveryDaysMax": 2,
"tariffCode": "ДоставкаДоДверей|tariff-guid-2|urg-exp",
},
"parcelType": "parcel",
"weight": "1.0",
"dimensions": {
"length": "10",
"width": "10",
"height": "10",
},
}
),
"order-uuid-1",
)
)
assert warning_calls == [
{
"event": "cse_request_server_error",
"operation": "SaveDocuments",
"status_code": 500,
"response_excerpt": (
"<fault> CSE rejected SaveDocuments because field X is invalid "
"</fault>"
),
}
]
@@ -0,0 +1,50 @@
import asyncio
import httpx
from app.adapters.delivery_providers.registry import (
build_delivery_provider_registry,
resolve_delivery_provider_timeout_seconds,
)
from app.config import (
CDEKDeliveryProviderConfig,
CSEDeliveryProviderConfig,
DeliveryProvidersConfig,
)
def test_registry_builds_enabled_provider_capability_maps() -> None:
config = DeliveryProvidersConfig(
cdek=CDEKDeliveryProviderConfig(
base_url="https://cdek.test/v2",
client_id="id",
client_secret="secret",
cache_ttl_seconds=111,
),
cse=CSEDeliveryProviderConfig(
enabled=False,
),
)
http_client = httpx.AsyncClient()
try:
registry = build_delivery_provider_registry(
http_client=http_client,
config=config,
)
finally:
asyncio.run(http_client.aclose())
assert [provider.name for provider in registry.providers] == ["cdek"]
assert sorted(registry.payment_price_validation_adapters) == ["cdek"]
assert sorted(registry.order_registration_adapters) == ["cdek"]
assert registry.providers[0].cache_ttl_seconds == 111
def test_resolve_delivery_provider_timeout_uses_max_enabled_timeout() -> None:
config = DeliveryProvidersConfig(
cdek=CDEKDeliveryProviderConfig(timeout_seconds=3.0),
cse=CSEDeliveryProviderConfig(timeout_seconds=7.5),
)
assert resolve_delivery_provider_timeout_seconds(config) == 7.5
+6 -3
View File
@@ -11,9 +11,12 @@ business_logic:
repository:
redis_dsn: "redis://localhost:6379/0"
adapter:
cdek_client_id: "yaml-id"
cdek_client_secret: "yaml-secret"
delivery_providers:
cdek:
client_id: "yaml-id"
client_secret: "yaml-secret"
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -14,14 +14,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -14,14 +14,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -14,14 +14,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -14,14 +14,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -13,14 +13,17 @@ repository:
redis_dsn: "redis://localhost:6379/0"
price_cache_ttl_seconds: 900
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -11,9 +11,12 @@ business_logic:
repository:
redis_dsn: "redis://localhost:6379/0"
adapter:
cdek_client_id: "test-id"
cdek_client_secret: "test-secret"
delivery_providers:
cdek:
client_id: "test-id"
client_secret: "test-secret"
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
+29 -10
View File
@@ -132,7 +132,7 @@ business_logic:
repository: {{}}
adapter: {{}}
delivery_providers: {{}}
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
@@ -174,13 +174,32 @@ def test_configuration_sections_are_loaded_from_yaml_file(
assert settings.business_logic.provider_price_multiplier == Decimal("1.0")
assert settings.repository.redis_dsn == "redis://localhost:6379/0"
assert settings.repository.price_cache_ttl_seconds == 900
assert settings.adapter.cdek_base_url == "https://api.cdek.ru/v2"
assert settings.adapter.cdek_client_id == "test-client-id"
assert settings.adapter.cdek_client_secret == "test-client-secret"
assert settings.adapter.cdek_retry_attempts == 2
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.delivery_providers.cdek.enabled is True
assert settings.delivery_providers.cdek.base_url == "https://api.cdek.ru/v2"
assert settings.delivery_providers.cdek.client_id == "test-client-id"
assert settings.delivery_providers.cdek.client_secret == "test-client-secret"
assert settings.delivery_providers.cdek.retry_attempts == 2
assert settings.delivery_providers.cdek.retry_backoff_seconds == 0.2
assert settings.delivery_providers.cdek.timeout_seconds == 10.0
assert settings.delivery_providers.cdek.cache_ttl_seconds == 900
assert settings.delivery_providers.cse.enabled is True
assert (
settings.delivery_providers.cse.base_url
== "http://lk-test.cse.ru/1c/ws/web1c.1cws"
)
assert settings.delivery_providers.cse.login == "test"
assert settings.delivery_providers.cse.password == "2016"
assert settings.delivery_providers.cse.retry_attempts == 2
assert settings.delivery_providers.cse.retry_backoff_seconds == 0.2
assert settings.delivery_providers.cse.timeout_seconds == 10.0
assert settings.delivery_providers.cse.cache_ttl_seconds == 900
assert settings.delivery_providers.cse.delivery_service_guids == ["tariff-guid-2"]
assert settings.delivery_providers.cse.payer == "0"
assert settings.delivery_providers.cse.payment_method == "1"
assert (
settings.delivery_providers.cse.shipping_method
== "5052d0b3-5ea3-46f2-823f-1472686a51dd"
)
assert settings.tbank_payment.init_url == "https://securepay.tinkoff.ru/v2/Init"
assert (
settings.tbank_payment.notification_url
@@ -233,7 +252,7 @@ def test_get_settings_fails_when_required_yaml_section_is_missing(
get_settings()
locations = {tuple(item["loc"]) for item in error.value.errors()}
assert ("adapter",) in locations
assert ("delivery_providers",) in locations
get_settings.cache_clear()
@@ -287,7 +306,7 @@ 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.business_logic.provider_price_multiplier == Decimal("1.25")
assert settings.adapter.cdek_client_id == "test-id"
assert settings.delivery_providers.cdek.client_id == "test-id"
assert settings.tbank_payment.auth.terminal_key == "override-terminal-key"
assert settings.tbank_payment.auth.password == "override-password"
assert (
@@ -5,6 +5,7 @@ import pytest
from app.controllers.v1 import delivery as delivery_controller
from app.controllers.v1.delivery import get_aggregator_service
from app.adapters.delivery_providers.registry import DeliveryProviderRegistry
from app.main import create_app
from app.schemas.request import AddressSuggestRequest, SuggestAddressRequest
from app.schemas.response import AddressSuggestion
@@ -52,16 +53,6 @@ def test_post_address_suggest_uses_registered_provider_in_default_dependency(
async def aclose(self) -> None:
return None
class StubCDEKProvider:
name = "stub-cdek"
cache_ttl_seconds = 900
@classmethod
def from_adapter_config(cls, *, http_client, adapter_config) -> "StubCDEKProvider":
assert http_client is stub_http_client
_ = adapter_config
return cls()
class StubAddressProvider:
def __init__(self, name: str, response: list[AddressSuggestion]) -> None:
self.name = name
@@ -169,12 +160,38 @@ def test_post_address_suggest_uses_registered_provider_in_default_dependency(
http_client_timeouts.append(timeout_seconds)
return stub_http_client
def fake_resolve_delivery_provider_timeout_seconds(config: object) -> float:
_ = config
return 10.0
def fake_build_delivery_provider_registry(
*,
http_client: object,
config: object,
) -> DeliveryProviderRegistry:
assert http_client is stub_http_client
_ = config
return DeliveryProviderRegistry(
providers=(),
payment_price_validation_adapters={},
order_registration_adapters={},
)
monkeypatch.setattr(
delivery_controller,
"build_controller_http_client",
fake_build_controller_http_client,
)
monkeypatch.setattr(delivery_controller, "CDEKProvider", StubCDEKProvider)
monkeypatch.setattr(
delivery_controller,
"resolve_delivery_provider_timeout_seconds",
fake_resolve_delivery_provider_timeout_seconds,
)
monkeypatch.setattr(
delivery_controller,
"build_delivery_provider_registry",
fake_build_delivery_provider_registry,
)
monkeypatch.setattr(
delivery_controller,
"DadataAddressSuggestionProvider",
+29 -9
View File
@@ -9,6 +9,7 @@ from app.controllers.v1.delivery import get_aggregator_service
from app.main import create_app
from app.schemas.request import DeliveryCalculationRequest, DeliveryEntity, ParcelType
from app.schemas.response import DeliveryPrice
from app.adapters.delivery_providers.registry import DeliveryProviderRegistry
from app.services.aggregator import (
AggregatorService,
AggregatorServiceError,
@@ -118,13 +119,6 @@ def test_post_delivery_price_uses_registered_provider_in_default_dependency(
),
]
class StubCDEKProvider:
@classmethod
def from_adapter_config(cls, *, http_client, adapter_config) -> StubProvider:
assert http_client is stub_http_client
_ = adapter_config
return stub_provider
class StubCache:
def __init__(self) -> None:
self.storage: dict[str, object] = {}
@@ -154,12 +148,38 @@ def test_post_delivery_price_uses_registered_provider_in_default_dependency(
http_client_timeouts.append(timeout_seconds)
return stub_http_client
def fake_resolve_delivery_provider_timeout_seconds(config: object) -> float:
_ = config
return 12.5
def fake_build_delivery_provider_registry(
*,
http_client: object,
config: object,
) -> DeliveryProviderRegistry:
assert http_client is stub_http_client
_ = config
return DeliveryProviderRegistry(
providers=(stub_provider,),
payment_price_validation_adapters={},
order_registration_adapters={},
)
monkeypatch.setattr(
delivery_controller,
"build_controller_http_client",
fake_build_controller_http_client,
)
monkeypatch.setattr(delivery_controller, "CDEKProvider", StubCDEKProvider)
monkeypatch.setattr(
delivery_controller,
"resolve_delivery_provider_timeout_seconds",
fake_resolve_delivery_provider_timeout_seconds,
)
monkeypatch.setattr(
delivery_controller,
"build_delivery_provider_registry",
fake_build_delivery_provider_registry,
)
monkeypatch.setattr(delivery_controller, "PriceCache", StubPriceCache)
app = create_app()
@@ -199,7 +219,7 @@ def test_post_delivery_price_uses_registered_provider_in_default_dependency(
}
]
assert second_response.json() == first_response.json()
assert len(http_client_timeouts) == 1
assert http_client_timeouts == [12.5]
assert len(stub_provider.calls) == 1
+41
View File
@@ -164,6 +164,47 @@ def test_post_init_payment_rejects_non_integer_price() -> None:
assert service.calls == []
def test_post_init_payment_rejects_missing_declared_value() -> None:
service = StubAggregatorService(response=None)
app = create_app()
_install_service_override(app, service)
invalid_payload = make_init_payment_payload()
del invalid_payload["content"]["declared_value"]
response = _post(app, invalid_payload)
assert response.status_code == 422
assert service.calls == []
def test_post_init_payment_rejects_non_integer_declared_value() -> None:
service = StubAggregatorService(response=None)
app = create_app()
_install_service_override(app, service)
invalid_payload = make_init_payment_payload()
invalid_payload["content"]["declared_value"] = "50000"
response = _post(app, invalid_payload)
assert response.status_code == 422
assert service.calls == []
def test_post_init_payment_rejects_camel_case_declared_value() -> None:
service = StubAggregatorService(response=None)
app = create_app()
_install_service_override(app, service)
invalid_payload = make_init_payment_payload()
invalid_payload["content"]["declaredValue"] = invalid_payload["content"].pop(
"declared_value"
)
response = _post(app, invalid_payload)
assert response.status_code == 422
assert service.calls == []
def test_post_init_payment_requires_company_fields_when_is_company_true() -> None:
service = StubAggregatorService(response=None)
app = create_app()
+2 -2
View File
@@ -5,14 +5,14 @@ from app.domain.payment_notifications import (
)
def test_confirmed_success_zero_error_code_registers_cdek_order() -> None:
def test_confirmed_success_zero_error_code_registers_provider_order() -> None:
result = resolve_tbank_payment_notification_action(
status="CONFIRMED",
success=True,
error_code="0",
)
assert result is TBankPaymentNotificationAction.REGISTER_CDEK_ORDER
assert result is TBankPaymentNotificationAction.REGISTER_PROVIDER_ORDER
def test_confirmed_without_success_acknowledges_only() -> None:
+1
View File
@@ -49,6 +49,7 @@ def make_init_payment_payload(**overrides: Any) -> dict[str, Any]:
},
"content": {
"description": "Headphones",
"declared_value": 50000,
},
"pickupDate": "2026-05-15T10:00:00.000Z",
"deliveryDate": "2026-05-18T18:00:00.000Z",
+11 -8
View File
@@ -186,14 +186,17 @@ repository:
redis_dsn: "redis://redis.internal:6380/5"
price_cache_ttl_seconds: 123
adapter:
cdek_base_url: "https://api.cdek.ru/v2"
cdek_client_id: "test-client-id"
cdek_client_secret: "test-client-secret"
cdek_retry_attempts: 2
cdek_retry_backoff_seconds: 0.2
cdek_timeout_seconds: 10.0
cdek_cache_ttl_seconds: 900
delivery_providers:
cdek:
base_url: "https://api.cdek.ru/v2"
client_id: "test-client-id"
client_secret: "test-client-secret"
retry_attempts: 2
retry_backoff_seconds: 0.2
timeout_seconds: 10.0
cache_ttl_seconds: 900
cse:
enabled: false
tbank_payment:
init_url: "https://securepay.tinkoff.ru/v2/Init"
+108 -82
View File
@@ -69,9 +69,9 @@ def test_create_order_persists_all_required_fields() -> None:
assert persisted_order.payload == order_data.payload
assert persisted_order.payment_status is None
assert persisted_order.tbank_payment_id is None
assert persisted_order.cdek_order_uuid is None
assert persisted_order.cdek_waybill_uuid is None
assert persisted_order.cdek_waybill_url is None
assert persisted_order.provider_order_id is None
assert persisted_order.provider_waybill_id is None
assert persisted_order.provider_waybill_url is None
assert persisted_order.created_at is not None
assert persisted_order.updated_at is not None
@@ -175,7 +175,7 @@ def test_mark_payment_status_returns_none_for_missing_order() -> None:
asyncio.run(_with_repository(run))
def test_mark_cdek_order_registered_persists_cdek_order_uuid_only() -> None:
def test_mark_provider_order_registered_persists_provider_order_id_only() -> None:
async def run(
repository: OrderRepository,
session_factory: async_sessionmaker[AsyncSession],
@@ -184,7 +184,7 @@ def test_mark_cdek_order_registered_persists_cdek_order_uuid_only() -> None:
await repository.create_order(session, _make_order_data())
async with repository.session() as session:
order = await repository.mark_cdek_order_registered(
order = await repository.mark_provider_order_registered(
session,
"order-uuid-1",
"cdek-order-uuid-1",
@@ -197,14 +197,14 @@ def test_mark_cdek_order_registered_persists_cdek_order_uuid_only() -> None:
persisted_order = result.scalar_one()
assert order is not None
assert persisted_order.cdek_order_uuid == "cdek-order-uuid-1"
assert persisted_order.cdek_waybill_uuid is None
assert persisted_order.cdek_waybill_url is None
assert persisted_order.provider_order_id == "cdek-order-uuid-1"
assert persisted_order.provider_waybill_id is None
assert persisted_order.provider_waybill_url is None
asyncio.run(_with_repository(run))
def test_mark_cdek_order_registered_is_idempotent_for_same_uuid() -> None:
def test_mark_provider_order_registered_is_idempotent_for_same_uuid() -> None:
async def run(
repository: OrderRepository,
session_factory: async_sessionmaker[AsyncSession],
@@ -213,14 +213,14 @@ def test_mark_cdek_order_registered_is_idempotent_for_same_uuid() -> None:
await repository.create_order(session, _make_order_data())
async with repository.session() as session:
await repository.mark_cdek_order_registered(
await repository.mark_provider_order_registered(
session,
"order-uuid-1",
"cdek-order-uuid-1",
)
async with repository.session() as session:
await repository.mark_cdek_order_registered(
await repository.mark_provider_order_registered(
session,
"order-uuid-1",
"cdek-order-uuid-1",
@@ -231,18 +231,18 @@ def test_mark_cdek_order_registered_is_idempotent_for_same_uuid() -> None:
orders = result.scalars().all()
assert len(orders) == 1
assert orders[0].cdek_order_uuid == "cdek-order-uuid-1"
assert orders[0].provider_order_id == "cdek-order-uuid-1"
asyncio.run(_with_repository(run))
def test_mark_cdek_order_registered_returns_none_for_missing_order() -> None:
def test_mark_provider_order_registered_returns_none_for_missing_order() -> None:
async def run(
repository: OrderRepository,
_session_factory: async_sessionmaker[AsyncSession],
) -> None:
async with repository.session() as session:
order = await repository.mark_cdek_order_registered(
order = await repository.mark_provider_order_registered(
session,
"missing-order",
"cdek-order-uuid-1",
@@ -257,37 +257,38 @@ async def _seed_order(
repository: OrderRepository,
*,
order_uuid: str,
cdek_order_uuid: str | None,
cdek_order_status: str | None = None,
cdek_waybill_uuid: str | None = None,
cdek_waybill_url: str | None = None,
cdek_polled_at: datetime | None = None,
provider_order_id: str | None,
provider: str = "cdek",
provider_order_status: str | None = None,
provider_waybill_id: str | None = None,
provider_waybill_url: str | None = None,
provider_polled_at: datetime | None = None,
) -> None:
async with repository.session() as session:
await repository.create_order(
session, _make_order_data(order_uuid=order_uuid)
session, _make_order_data(order_uuid=order_uuid, provider=provider)
)
if cdek_order_uuid is not None:
await repository.mark_cdek_order_registered(
session, order_uuid, cdek_order_uuid
if provider_order_id is not None:
await repository.mark_provider_order_registered(
session, order_uuid, provider_order_id
)
if (
cdek_order_status is not None
or cdek_waybill_uuid is not None
or cdek_polled_at is not None
provider_order_status is not None
or provider_waybill_id is not None
or provider_polled_at is not None
):
order = await repository.get_order_by_order_uuid(session, order_uuid)
assert order is not None
if cdek_order_status is not None:
order.cdek_order_status = cdek_order_status
if cdek_waybill_uuid is not None:
order.cdek_waybill_uuid = cdek_waybill_uuid
if cdek_polled_at is not None:
order.cdek_polled_at = cdek_polled_at
if cdek_waybill_url is not None:
if provider_order_status is not None:
order.provider_order_status = provider_order_status
if provider_waybill_id is not None:
order.provider_waybill_id = provider_waybill_id
if provider_polled_at is not None:
order.provider_polled_at = provider_polled_at
if provider_waybill_url is not None:
order = await repository.get_order_by_order_uuid(session, order_uuid)
assert order is not None
order.cdek_waybill_url = cdek_waybill_url
order.provider_waybill_url = provider_waybill_url
def test_list_orders_pending_waybill_returns_orders_without_url() -> None:
@@ -295,26 +296,38 @@ def test_list_orders_pending_waybill_returns_orders_without_url() -> None:
repository: OrderRepository,
_session_factory: async_sessionmaker[AsyncSession],
) -> None:
await _seed_order(repository, order_uuid="pending", cdek_order_uuid="o1")
await _seed_order(repository, order_uuid="pending", provider_order_id="o1")
await _seed_order(
repository,
order_uuid="done",
cdek_order_uuid="o2",
cdek_waybill_uuid="w2",
cdek_waybill_url="https://cdek.test/2.pdf",
provider_order_id="o2",
provider_waybill_id="w2",
provider_waybill_url="https://cdek.test/2.pdf",
)
await _seed_order(
repository,
order_uuid="invalid",
cdek_order_uuid="o3",
cdek_order_status="INVALID",
provider_order_id="o3",
provider_order_status="INVALID",
)
await _seed_order(
repository,
order_uuid="cse-pending",
provider="cse",
provider_order_id="cse-1",
)
await _seed_order(
repository,
order_uuid="cse-with-waybill",
provider="cse",
provider_order_id="cse-2",
provider_waybill_id="496-AA-1676378",
)
await _seed_order(repository, order_uuid="no-cdek", cdek_order_uuid=None)
async with repository.session() as session:
orders = await repository.list_orders_pending_waybill(session, limit=10)
assert [order.order_uuid for order in orders] == ["pending"]
assert [order.order_uuid for order in orders] == ["pending", "cse-pending"]
asyncio.run(_with_repository(run))
@@ -330,16 +343,16 @@ def test_list_orders_pending_waybill_orders_polled_at_nulls_first() -> None:
await _seed_order(
repository,
order_uuid="late",
cdek_order_uuid="o-late",
cdek_polled_at=later,
provider_order_id="o-late",
provider_polled_at=later,
)
await _seed_order(
repository,
order_uuid="early",
cdek_order_uuid="o-early",
cdek_polled_at=earlier,
provider_order_id="o-early",
provider_polled_at=earlier,
)
await _seed_order(repository, order_uuid="never", cdek_order_uuid="o-never")
await _seed_order(repository, order_uuid="never", provider_order_id="o-never")
async with repository.session() as session:
orders = await repository.list_orders_pending_waybill(session, limit=10)
@@ -356,7 +369,7 @@ def test_record_order_poll_sets_status_and_waybill_uuid() -> None:
repository: OrderRepository,
_session_factory: async_sessionmaker[AsyncSession],
) -> None:
await _seed_order(repository, order_uuid="o", cdek_order_uuid="cdek-o")
await _seed_order(repository, order_uuid="o", provider_order_id="cdek-o")
async with repository.session() as session:
order = await repository.record_order_poll(
@@ -368,9 +381,9 @@ def test_record_order_poll_sets_status_and_waybill_uuid() -> None:
)
assert order is not None
assert order.cdek_order_status == "ACCEPTED"
assert order.cdek_waybill_uuid == "waybill-1"
assert order.cdek_polled_at == polled
assert order.provider_order_status == "ACCEPTED"
assert order.provider_waybill_id == "waybill-1"
assert order.provider_polled_at == polled
asyncio.run(_with_repository(run))
@@ -385,8 +398,8 @@ def test_record_order_poll_does_not_overwrite_existing_waybill_uuid() -> None:
await _seed_order(
repository,
order_uuid="o",
cdek_order_uuid="cdek-o",
cdek_waybill_uuid="existing-waybill",
provider_order_id="cdek-o",
provider_waybill_id="existing-waybill",
)
async with repository.session() as session:
@@ -399,7 +412,7 @@ def test_record_order_poll_does_not_overwrite_existing_waybill_uuid() -> None:
)
assert order is not None
assert order.cdek_waybill_uuid == "existing-waybill"
assert order.provider_waybill_id == "existing-waybill"
asyncio.run(_with_repository(run))
@@ -414,8 +427,8 @@ def test_record_waybill_poll_sets_url_only_when_previously_null() -> None:
await _seed_order(
repository,
order_uuid="o",
cdek_order_uuid="cdek-o",
cdek_waybill_uuid="waybill-1",
provider_order_id="cdek-o",
provider_waybill_id="waybill-1",
)
async with repository.session() as session:
@@ -427,8 +440,8 @@ def test_record_waybill_poll_sets_url_only_when_previously_null() -> None:
)
assert order is not None
assert order.cdek_waybill_url == "https://cdek.test/1.pdf"
assert order.cdek_polled_at == polled
assert order.provider_waybill_url == "https://cdek.test/1.pdf"
assert order.provider_polled_at == polled
async with repository.session() as session:
order = await repository.record_waybill_poll(
@@ -438,12 +451,12 @@ def test_record_waybill_poll_sets_url_only_when_previously_null() -> None:
polled_at=polled,
)
assert order is not None
assert order.cdek_waybill_url == "https://cdek.test/1.pdf"
assert order.provider_waybill_url == "https://cdek.test/1.pdf"
asyncio.run(_with_repository(run))
def test_list_orders_pending_waybill_email_returns_orders_with_url_and_no_sent_at() -> None:
def test_list_orders_pending_waybill_email_returns_ready_orders() -> None:
async def run(
repository: OrderRepository,
_session_factory: async_sessionmaker[AsyncSession],
@@ -451,22 +464,35 @@ def test_list_orders_pending_waybill_email_returns_orders_with_url_and_no_sent_a
await _seed_order(
repository,
order_uuid="ready",
cdek_order_uuid="o1",
cdek_waybill_uuid="w1",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_order_id="o1",
provider_waybill_id="w1",
provider_waybill_url="https://cdek.test/1.pdf",
)
await _seed_order(
repository,
order_uuid="no-url",
cdek_order_uuid="o2",
cdek_waybill_uuid="w2",
provider_order_id="o2",
provider_waybill_id="w2",
)
await _seed_order(
repository,
order_uuid="already-sent",
cdek_order_uuid="o3",
cdek_waybill_uuid="w3",
cdek_waybill_url="https://cdek.test/3.pdf",
provider_order_id="o3",
provider_waybill_id="w3",
provider_waybill_url="https://cdek.test/3.pdf",
)
await _seed_order(
repository,
order_uuid="cse-ready",
provider="cse",
provider_order_id="cse-o1",
provider_waybill_id="496-AA-1676378",
)
await _seed_order(
repository,
order_uuid="cse-no-waybill",
provider="cse",
provider_order_id="cse-o2",
)
async with repository.session() as session:
sent = await repository.get_order_by_order_uuid(session, "already-sent")
@@ -480,7 +506,7 @@ def test_list_orders_pending_waybill_email_returns_orders_with_url_and_no_sent_a
session, limit=10
)
assert [order.order_uuid for order in orders] == ["ready"]
assert [order.order_uuid for order in orders] == ["ready", "cse-ready"]
asyncio.run(_with_repository(run))
@@ -493,16 +519,16 @@ def test_list_orders_pending_waybill_email_orders_by_created_at_asc() -> None:
await _seed_order(
repository,
order_uuid="first",
cdek_order_uuid="o1",
cdek_waybill_uuid="w1",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_order_id="o1",
provider_waybill_id="w1",
provider_waybill_url="https://cdek.test/1.pdf",
)
await _seed_order(
repository,
order_uuid="second",
cdek_order_uuid="o2",
cdek_waybill_uuid="w2",
cdek_waybill_url="https://cdek.test/2.pdf",
provider_order_id="o2",
provider_waybill_id="w2",
provider_waybill_url="https://cdek.test/2.pdf",
)
async with repository.session() as session:
@@ -526,9 +552,9 @@ def test_record_waybill_email_sent_sets_timestamp_once() -> None:
await _seed_order(
repository,
order_uuid="o",
cdek_order_uuid="cdek-o",
cdek_waybill_uuid="w",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_order_id="cdek-o",
provider_waybill_id="w",
provider_waybill_url="https://cdek.test/1.pdf",
)
async with repository.session() as session:
@@ -580,8 +606,8 @@ def test_record_waybill_poll_updates_polled_at_when_url_is_none() -> None:
await _seed_order(
repository,
order_uuid="o",
cdek_order_uuid="cdek-o",
cdek_waybill_uuid="waybill-1",
provider_order_id="cdek-o",
provider_waybill_id="waybill-1",
)
async with repository.session() as session:
@@ -593,7 +619,7 @@ def test_record_waybill_poll_updates_polled_at_when_url_is_none() -> None:
)
assert order is not None
assert order.cdek_waybill_url is None
assert order.cdek_polled_at == polled
assert order.provider_waybill_url is None
assert order.provider_polled_at == polled
asyncio.run(_with_repository(run))
+5 -11
View File
@@ -115,8 +115,7 @@ class StoredOrder:
)
payment_status: str | None = None
tbank_payment_id: int | None = None
cdek_order_uuid: str | None = None
cse_order_number: str | None = None
provider_order_id: str | None = None
payment_email_sent_at: object | None = None
account_email: str | None = None
@@ -148,17 +147,12 @@ class StubOrderRepository:
self._order.tbank_payment_id = payment_id
return self._order
async def mark_cse_order_registered(
self, session: object, order_uuid: str, cse_order_number: str
async def mark_provider_order_registered(
self, session: object, order_uuid: str, provider_order_id: str
) -> StoredOrder:
self._order.cse_order_number = cse_order_number
self._order.provider_order_id = provider_order_id
return self._order
async def mark_cdek_order_registered(
self, session: object, order_uuid: str, cdek_order_uuid: str
) -> StoredOrder:
raise AssertionError("CDEK persistence must not be used for a CSE order.")
def _notification() -> TBankPaymentNotification:
return TBankPaymentNotification(
@@ -190,4 +184,4 @@ def test_notification_routes_registration_to_order_provider() -> None:
assert result == "OK"
assert len(cse_registration.calls) == 1
assert order.cse_order_number == "CSE-000123"
assert order.provider_order_id == "CSE-000123"
+21 -22
View File
@@ -32,10 +32,9 @@ class StoredOrder:
payload: dict[str, Any] = field(default_factory=_default_payload)
payment_status: str | None = None
tbank_payment_id: int | None = None
cdek_order_uuid: str | None = None
cse_order_number: str | None = None
cdek_waybill_uuid: str | None = None
cdek_waybill_url: str | None = None
provider_order_id: str | None = None
provider_waybill_id: str | None = None
provider_waybill_url: str | None = None
payment_email_sent_at: object | None = None
@@ -72,10 +71,10 @@ class StubOrderRepository:
self,
*,
orders: list[StoredOrder] | None = None,
mark_cdek_errors: list[Exception | None] | None = None,
mark_provider_errors: list[Exception | None] | None = None,
) -> None:
self._orders = {order.order_uuid: order for order in orders or []}
self._mark_cdek_errors = mark_cdek_errors or []
self._mark_provider_errors = mark_provider_errors or []
self.session_value = object()
self.calls: list[tuple[str, tuple[object, ...]]] = []
@@ -113,24 +112,24 @@ class StubOrderRepository:
order.tbank_payment_id = payment_id
return order
async def mark_cdek_order_registered(
async def mark_provider_order_registered(
self,
session: object,
order_uuid: str,
cdek_order_uuid: str,
provider_order_id: str,
) -> StoredOrder | None:
self.calls.append(
("mark_cdek_order_registered", (session, order_uuid, cdek_order_uuid))
("mark_provider_order_registered", (session, order_uuid, provider_order_id))
)
if self._mark_cdek_errors:
error = self._mark_cdek_errors.pop(0)
if self._mark_provider_errors:
error = self._mark_provider_errors.pop(0)
if error is not None:
raise error
order = self._orders.get(order_uuid)
if order is None:
return None
order.cdek_order_uuid = cdek_order_uuid
order.provider_order_id = provider_order_id
return order
async def record_payment_email_sent(
@@ -230,15 +229,15 @@ def test_confirmed_notification_registers_cdek_order_and_saves_uuid() -> None:
assert payment_adapter.notifications == [notification]
assert order.payment_status == "CONFIRMED"
assert order.tbank_payment_id == 8347568144
assert order.cdek_order_uuid == "cdek-order-uuid-1"
assert order.cdek_waybill_uuid is None
assert order.cdek_waybill_url is None
assert order.provider_order_id == "cdek-order-uuid-1"
assert order.provider_waybill_id is None
assert order.provider_waybill_url is None
assert len(cdek_adapter.calls) == 1
assert cdek_adapter.calls[0][1] == "order-uuid-1"
def test_duplicate_confirmed_notification_does_not_call_cdek() -> None:
order = StoredOrder(cdek_order_uuid="existing-cdek-order-uuid")
order = StoredOrder(provider_order_id="existing-cdek-order-uuid")
cdek_adapter = StubCDEKOrderAdapter()
service = AggregatorService(
providers=[],
@@ -253,7 +252,7 @@ def test_duplicate_confirmed_notification_does_not_call_cdek() -> None:
assert result == "OK"
assert cdek_adapter.calls == []
assert order.cdek_order_uuid == "existing-cdek-order-uuid"
assert order.provider_order_id == "existing-cdek-order-uuid"
def test_non_confirmed_notification_acknowledges_without_cdek() -> None:
@@ -330,7 +329,7 @@ def test_repeated_confirmed_after_cdek_uuid_save_failure_uses_same_external_id()
order = StoredOrder()
order_repository = StubOrderRepository(
orders=[order],
mark_cdek_errors=[RuntimeError("db down"), None],
mark_provider_errors=[RuntimeError("db down"), None],
)
cdek_adapter = StubCDEKOrderAdapter(
responses=[
@@ -357,12 +356,12 @@ def test_repeated_confirmed_after_cdek_uuid_save_failure_uses_same_external_id()
with pytest.raises(TBankPaymentNotificationProcessingError):
asyncio.run(service.handle_tbank_payment_notification(notification))
assert order.cdek_order_uuid is None
assert order.provider_order_id is None
result = asyncio.run(service.handle_tbank_payment_notification(notification))
assert result == "OK"
assert order.cdek_order_uuid == "same-cdek-order-uuid"
assert order.provider_order_id == "same-cdek-order-uuid"
assert [order_uuid for _, order_uuid in cdek_adapter.calls] == [
"order-uuid-1",
"order-uuid-1",
@@ -417,7 +416,7 @@ def test_duplicate_notification_does_not_resend_payment_email() -> None:
from datetime import datetime, timezone
order = StoredOrder(
cdek_order_uuid="existing-cdek-order-uuid",
provider_order_id="existing-cdek-order-uuid",
payment_email_sent_at=datetime(2026, 1, 1, tzinfo=timezone.utc),
)
email_sender = StubEmailSender()
@@ -454,4 +453,4 @@ def test_email_failure_does_not_break_notification_handling() -> None:
assert result == "OK"
assert len(email_sender.calls) == 1
assert order.cdek_order_uuid == "cdek-order-uuid-1"
assert order.provider_order_id == "cdek-order-uuid-1"
+47 -11
View File
@@ -10,7 +10,9 @@ from app.services.waybill_email_sender import WaybillEmailSenderService
class StoredOrder:
order_uuid: str
account_email: str
cdek_waybill_url: str | None = None
provider: str = "cdek"
provider_waybill_id: str | None = None
provider_waybill_url: str | None = None
waybill_email_sent_at: datetime | None = None
@@ -41,7 +43,10 @@ class StubRepository:
return [
order
for order in self._orders.values()
if order.cdek_waybill_url is not None
if (
order.provider_waybill_url is not None
or order.provider_waybill_id is not None
)
and order.waybill_email_sent_at is None
]
@@ -68,9 +73,11 @@ class StubDownloader:
self._results = results
self.calls: list[str] = []
async def download_waybill_pdf(self, url: str) -> bytes:
self.calls.append(url)
result = self._results[url]
async def download_waybill_pdf(self, order: StoredOrder) -> bytes:
locator = order.provider_waybill_url or order.provider_waybill_id
assert locator is not None
self.calls.append(locator)
result = self._results[locator]
if isinstance(result, Exception):
raise result
return result
@@ -114,7 +121,7 @@ def _make_service(
) -> WaybillEmailSenderService:
return WaybillEmailSenderService(
order_repository=repository,
waybill_downloader=downloader or StubDownloader({}),
waybill_downloaders={"cdek": downloader or StubDownloader({})},
email_sender=email_sender or StubEmailSender(),
batch_size=10,
datetime_now=lambda: _SENT_AT,
@@ -125,7 +132,7 @@ def test_poll_once_downloads_pdf_sends_email_and_marks_sent() -> None:
order = StoredOrder(
order_uuid="o-1",
account_email="client@example.com",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_waybill_url="https://cdek.test/1.pdf",
)
repo = StubRepository([order])
downloader = StubDownloader({"https://cdek.test/1.pdf": b"%PDF"})
@@ -151,11 +158,40 @@ def test_poll_once_downloads_pdf_sends_email_and_marks_sent() -> None:
assert order.waybill_email_sent_at == _SENT_AT
def test_poll_once_sends_cse_waybill_pdf_without_url() -> None:
order = StoredOrder(
order_uuid="o-1",
provider="cse",
account_email="client@example.com",
provider_waybill_id="496-AA-1676378",
)
repo = StubRepository([order])
downloader = StubDownloader({"496-AA-1676378": b"%PDF"})
email_sender = StubEmailSender()
service = WaybillEmailSenderService(
order_repository=repo,
waybill_downloaders={"cse": downloader},
email_sender=email_sender,
batch_size=10,
datetime_now=lambda: _SENT_AT,
)
summary = asyncio.run(service.poll_once())
assert summary.processed == 1
assert summary.succeeded == 1
assert downloader.calls == ["496-AA-1676378"]
assert len(email_sender.calls) == 1
assert "https://" not in email_sender.calls[0]["body"]
assert email_sender.calls[0]["attachment_bytes"] == b"%PDF"
assert order.waybill_email_sent_at == _SENT_AT
def test_poll_once_download_error_keeps_order_pending_and_skips_send() -> None:
order = StoredOrder(
order_uuid="o-1",
account_email="client@example.com",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_waybill_url="https://cdek.test/1.pdf",
)
repo = StubRepository([order])
downloader = StubDownloader({"https://cdek.test/1.pdf": RuntimeError("cdek 500")})
@@ -177,7 +213,7 @@ def test_poll_once_smtp_error_keeps_order_pending() -> None:
order = StoredOrder(
order_uuid="o-1",
account_email="bad@example.com",
cdek_waybill_url="https://cdek.test/1.pdf",
provider_waybill_url="https://cdek.test/1.pdf",
)
repo = StubRepository([order])
downloader = StubDownloader({"https://cdek.test/1.pdf": b"%PDF"})
@@ -197,12 +233,12 @@ def test_poll_once_failure_on_one_order_does_not_break_batch() -> None:
bad = StoredOrder(
order_uuid="bad",
account_email="bad@example.com",
cdek_waybill_url="https://cdek.test/bad.pdf",
provider_waybill_url="https://cdek.test/bad.pdf",
)
good = StoredOrder(
order_uuid="good",
account_email="good@example.com",
cdek_waybill_url="https://cdek.test/good.pdf",
provider_waybill_url="https://cdek.test/good.pdf",
)
repo = StubRepository([bad, good])
downloader = StubDownloader(
+73 -35
View File
@@ -7,17 +7,19 @@ from app.adapters.delivery_providers.cdek.order_mapper import (
CDEKOrderInfo,
CDEKWaybillInfo,
)
from app.adapters.delivery_providers.cse.order_mapper import CSEOrderInfo
from app.services.waybill_poller import WaybillPollerService
@dataclass
class StoredOrder:
order_uuid: str
cdek_order_uuid: str | None = None
cdek_order_status: str | None = None
cdek_waybill_uuid: str | None = None
cdek_waybill_url: str | None = None
cdek_polled_at: datetime | None = None
provider: str = "cdek"
provider_order_id: str | None = None
provider_order_status: str | None = None
provider_waybill_id: str | None = None
provider_waybill_url: str | None = None
provider_polled_at: datetime | None = None
class StubSessionContext:
@@ -69,10 +71,10 @@ class StubRepository:
order = self._orders.get(order_uuid)
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
return order
async def record_waybill_poll(
@@ -96,9 +98,9 @@ class StubRepository:
order = self._orders.get(order_uuid)
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
return order
@@ -107,9 +109,9 @@ class StubOrderInfoAdapter:
self._results = results
self.calls: list[str] = []
async def get_order(self, cdek_order_uuid: str) -> CDEKOrderInfo:
self.calls.append(cdek_order_uuid)
result = self._results[cdek_order_uuid]
async def get_order(self, provider_order_id: str) -> CDEKOrderInfo:
self.calls.append(provider_order_id)
result = self._results[provider_order_id]
if isinstance(result, Exception):
raise result
return result
@@ -120,15 +122,18 @@ class StubWaybillInfoAdapter:
self._results = results
self.calls: list[str] = []
async def get_waybill(self, cdek_waybill_uuid: str) -> CDEKWaybillInfo:
self.calls.append(cdek_waybill_uuid)
result = self._results[cdek_waybill_uuid]
async def get_waybill(self, provider_waybill_id: str) -> CDEKWaybillInfo:
self.calls.append(provider_waybill_id)
result = self._results[provider_waybill_id]
if isinstance(result, Exception):
raise result
return result
_POLLED_AT = datetime(2026, 5, 24, 12, 0, tzinfo=timezone.utc)
_CSE_WAYBILL_CREATED_STATUS = (
"На основании заказа оформлена накладная."
)
def _make_service(
@@ -139,15 +144,15 @@ def _make_service(
) -> WaybillPollerService:
return WaybillPollerService(
order_repository=repository,
order_info_adapter=order_info or StubOrderInfoAdapter({}),
waybill_info_adapter=waybill_info or StubWaybillInfoAdapter({}),
order_info_adapters={"cdek": order_info or StubOrderInfoAdapter({})},
waybill_info_adapters={"cdek": waybill_info or StubWaybillInfoAdapter({})},
batch_size=10,
datetime_now=lambda: _POLLED_AT,
)
def test_poll_once_fetches_order_info_when_waybill_uuid_is_missing() -> None:
order = StoredOrder(order_uuid="o", cdek_order_uuid="cdek-o")
order = StoredOrder(order_uuid="o", provider_order_id="cdek-o")
repo = StubRepository([order])
order_info = StubOrderInfoAdapter(
{
@@ -164,16 +169,16 @@ def test_poll_once_fetches_order_info_when_waybill_uuid_is_missing() -> None:
assert summary.processed == 1 and summary.succeeded == 1 and summary.failed == 0
assert order_info.calls == ["cdek-o"]
assert order.cdek_order_status == "ACCEPTED"
assert order.cdek_waybill_uuid == "waybill-1"
assert order.cdek_polled_at == _POLLED_AT
assert order.provider_order_status == "ACCEPTED"
assert order.provider_waybill_id == "waybill-1"
assert order.provider_polled_at == _POLLED_AT
def test_poll_once_fetches_waybill_info_when_waybill_uuid_is_present() -> None:
order = StoredOrder(
order_uuid="o",
cdek_order_uuid="cdek-o",
cdek_waybill_uuid="waybill-1",
provider_order_id="cdek-o",
provider_waybill_id="waybill-1",
)
repo = StubRepository([order])
waybill_info = StubWaybillInfoAdapter(
@@ -190,12 +195,45 @@ def test_poll_once_fetches_waybill_info_when_waybill_uuid_is_present() -> None:
assert summary.processed == 1 and summary.succeeded == 1 and summary.failed == 0
assert waybill_info.calls == ["waybill-1"]
assert order.cdek_waybill_url == "https://cdek.test/1.pdf"
assert order.cdek_polled_at == _POLLED_AT
assert order.provider_waybill_url == "https://cdek.test/1.pdf"
assert order.provider_polled_at == _POLLED_AT
def test_poll_once_fetches_cse_waybill_number_from_order_tracking() -> None:
order = StoredOrder(
order_uuid="o",
provider="cse",
provider_order_id="CSE-000123",
)
repo = StubRepository([order])
order_info = StubOrderInfoAdapter(
{
"CSE-000123": CSEOrderInfo(
order_number="CSE-000123",
status_code=_CSE_WAYBILL_CREATED_STATUS,
waybill_number="496-AA-1676378",
)
}
)
service = WaybillPollerService(
order_repository=repo,
order_info_adapters={"cse": order_info},
waybill_info_adapters={},
batch_size=10,
datetime_now=lambda: _POLLED_AT,
)
summary = asyncio.run(service.poll_once())
assert summary.processed == 1 and summary.succeeded == 1 and summary.failed == 0
assert order_info.calls == ["CSE-000123"]
assert order.provider_order_status == _CSE_WAYBILL_CREATED_STATUS
assert order.provider_waybill_id == "496-AA-1676378"
assert order.provider_polled_at == _POLLED_AT
def test_poll_once_records_terminal_status_without_waybill() -> None:
order = StoredOrder(order_uuid="o", cdek_order_uuid="cdek-o")
order = StoredOrder(order_uuid="o", provider_order_id="cdek-o")
repo = StubRepository([order])
order_info = StubOrderInfoAdapter(
{
@@ -211,16 +249,16 @@ def test_poll_once_records_terminal_status_without_waybill() -> None:
summary = asyncio.run(service.poll_once())
assert summary.succeeded == 1
assert order.cdek_order_status == "INVALID"
assert order.cdek_waybill_uuid is None
assert order.provider_order_status == "INVALID"
assert order.provider_waybill_id is None
def test_poll_once_failure_on_one_order_does_not_break_batch() -> None:
failing = StoredOrder(order_uuid="bad", cdek_order_uuid="cdek-bad")
failing = StoredOrder(order_uuid="bad", provider_order_id="cdek-bad")
good = StoredOrder(
order_uuid="good",
cdek_order_uuid="cdek-good",
cdek_waybill_uuid="waybill-good",
provider_order_id="cdek-good",
provider_waybill_id="waybill-good",
)
repo = StubRepository([failing, good])
order_info = StubOrderInfoAdapter({"cdek-bad": RuntimeError("cdek down")})
@@ -241,7 +279,7 @@ def test_poll_once_failure_on_one_order_does_not_break_batch() -> None:
assert summary.processed == 2
assert summary.succeeded == 1
assert summary.failed == 1
assert good.cdek_waybill_url == "https://cdek.test/good.pdf"
assert good.provider_waybill_url == "https://cdek.test/good.pdf"
def test_run_forever_exits_when_stop_event_is_set() -> None:
@@ -3,7 +3,8 @@ from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
from app.config import (
AdapterConfig,
CDEKDeliveryProviderConfig,
DeliveryProvidersConfig,
EmailAdapterConfig,
ObservabilityConfig,
PostgresConfig,
@@ -17,10 +18,12 @@ from app.workers.waybill_email_sender import _run
def _make_settings() -> Settings:
return Settings(
adapter=AdapterConfig(
cdek_base_url="https://api.cdek.test/v2",
cdek_client_id="id",
cdek_client_secret="secret",
delivery_providers=DeliveryProvidersConfig(
cdek=CDEKDeliveryProviderConfig(
base_url="https://api.cdek.test/v2",
client_id="id",
client_secret="secret",
)
),
tbank_payment=TBankPaymentConfig(
init_url="https://pay.test/init",
@@ -61,6 +64,7 @@ def test_run_exits_when_stop_event_is_set() -> None:
),
patch("app.workers.waybill_email_sender.CDEKAuthClient"),
patch("app.workers.waybill_email_sender.CDEKClient"),
patch("app.workers.waybill_email_sender.CSEClient"),
patch("app.workers.waybill_email_sender.SMTPEmailSender"),
patch(
"app.workers.waybill_email_sender.create_postgres_engine",
+9 -5
View File
@@ -10,7 +10,8 @@ from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
from app.config import (
AdapterConfig,
CDEKDeliveryProviderConfig,
DeliveryProvidersConfig,
EmailAdapterConfig,
ObservabilityConfig,
PostgresConfig,
@@ -24,10 +25,12 @@ from app.workers.waybill_poller import _run
def _make_settings() -> Settings:
return Settings(
adapter=AdapterConfig(
cdek_base_url="https://api.cdek.test/v2",
cdek_client_id="id",
cdek_client_secret="secret",
delivery_providers=DeliveryProvidersConfig(
cdek=CDEKDeliveryProviderConfig(
base_url="https://api.cdek.test/v2",
client_id="id",
client_secret="secret",
)
),
tbank_payment=TBankPaymentConfig(
init_url="https://pay.test/init",
@@ -66,6 +69,7 @@ def test_run_exits_when_stop_event_is_set() -> None:
),
patch("app.workers.waybill_poller.CDEKAuthClient"),
patch("app.workers.waybill_poller.CDEKClient"),
patch("app.workers.waybill_poller.CSEClient"),
patch(
"app.workers.waybill_poller.create_postgres_engine",
return_value=engine_instance,