From bddac60965a1fa302fc93692c3fa5485c4407f1e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=D0=A0=D0=B0=D0=B8=D1=81=20=D0=AE=D1=81=D1=83=D0=BF=D0=B0?= =?UTF-8?q?=D0=BB=D0=B8=D0=B5=D0=B2?= Date: Sat, 18 Apr 2026 00:33:45 +0300 Subject: [PATCH] =?UTF-8?q?=D0=94=D0=BE=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5?= =?UTF-8?q?=D0=BD=D0=BE=20=D1=81=D0=BE=D1=85=D1=80=D0=B0=D0=BD=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D0=B5=20=D0=B7=D0=B0=D0=BA=D0=B0=D0=B7=D0=BE=D0=B2=20?= =?UTF-8?q?=D0=B2=20postgres?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- Dockerfile | 2 + alembic.ini | 41 +++++ alembic/env.py | 66 ++++++++ alembic/script.py.mako | 20 +++ .../20260412_028_create_orders_table.py | 51 ++++++ app/adapters/postgres/__init__.py | 0 app/adapters/postgres/engine.py | 16 ++ app/adapters/tbank/base.py | 15 ++ app/adapters/tbank/client.py | 80 +++++++++- app/config.py | 6 + app/controllers/v1/delivery.py | 6 + app/repositories/order/__init__.py | 5 + app/repositories/order/models.py | 63 ++++++++ app/repositories/order/repository.py | 53 +++++++ app/services/aggregator.py | 66 ++++++++ config.example.yaml | 3 + config.test.yaml | 3 + docker-compose.yml | 33 +++- poetry.lock | 17 +- pyproject.toml | 3 + spec/index.md | 5 +- spec/overview.md | 31 +++- .../028_add_postgresql_order_persistence.md | 75 +++++++++ .../delivery_providers/cdek/test_client.py | 3 + tests/adapters/tbank/test_client.py | 36 ++++- tests/config/fixtures/config.default.yaml | 3 + .../config.invalid_price_multiplier.yaml | 3 + .../fixtures/config.missing_adapter.yaml | 3 + .../config.missing_address_suggestions.yaml | 3 + .../config.missing_observability.yaml | 3 + ...config.missing_observability_endpoint.yaml | 3 + .../config.missing_price_multiplier.yaml | 3 + .../config/fixtures/config.test.override.yaml | 3 + tests/config/test_config_sections.py | 9 ++ tests/repositories/cache/test_redis_cache.py | 3 + tests/repositories/order/test_repository.py | 145 ++++++++++++++++++ tests/services/test_init_payment.py | 105 +++++++++++++ tests/smoke/test_local_infra_stack.py | 3 +- 38 files changed, 971 insertions(+), 17 deletions(-) create mode 100644 alembic.ini create mode 100644 alembic/env.py create mode 100644 alembic/script.py.mako create mode 100644 alembic/versions/20260412_028_create_orders_table.py create mode 100644 app/adapters/postgres/__init__.py create mode 100644 app/adapters/postgres/engine.py create mode 100644 app/repositories/order/__init__.py create mode 100644 app/repositories/order/models.py create mode 100644 app/repositories/order/repository.py create mode 100644 spec/tasks/028_add_postgresql_order_persistence.md create mode 100644 tests/repositories/order/test_repository.py diff --git a/Dockerfile b/Dockerfile index f111240..613dbbe 100644 --- a/Dockerfile +++ b/Dockerfile @@ -12,5 +12,7 @@ RUN poetry config virtualenvs.create false \ COPY app ./app COPY config.yaml ./config.yaml +COPY alembic.ini ./alembic.ini +COPY alembic ./alembic CMD ["poetry", "run", "uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/alembic.ini b/alembic.ini new file mode 100644 index 0000000..8de810c --- /dev/null +++ b/alembic.ini @@ -0,0 +1,41 @@ +[alembic] +script_location = alembic +prepend_sys_path = . +path_separator = os +sqlalchemy.url = + +[post_write_hooks] + +[loggers] +keys = root,sqlalchemy,alembic + +[handlers] +keys = console + +[formatters] +keys = generic + +[logger_root] +level = WARNING +handlers = console +qualname = + +[logger_sqlalchemy] +level = WARNING +handlers = +qualname = sqlalchemy.engine + +[logger_alembic] +level = INFO +handlers = +qualname = alembic + +[handler_console] +class = StreamHandler +args = (sys.stderr,) +level = NOTSET +formatter = generic + +[formatter_generic] +format = %(levelname)-5.5s [%(name)s] %(message)s +datefmt = %H:%M:%S diff --git a/alembic/env.py b/alembic/env.py new file mode 100644 index 0000000..1d80b43 --- /dev/null +++ b/alembic/env.py @@ -0,0 +1,66 @@ +"""Alembic async migration environment.""" + +from asyncio import run +from logging.config import fileConfig + +from alembic import context +from sqlalchemy import pool +from sqlalchemy.engine import Connection +from sqlalchemy.ext.asyncio import async_engine_from_config + +from app.config import get_settings +from app.repositories.order.models import Base + +config = context.config + +if config.config_file_name is not None: + fileConfig(config.config_file_name) + +target_metadata = Base.metadata + + +def _database_url() -> str: + configured_url = config.get_main_option("sqlalchemy.url") + if configured_url: + return configured_url + return get_settings().postgres.dsn + + +def run_migrations_offline() -> None: + context.configure( + url=_database_url(), + target_metadata=target_metadata, + literal_binds=True, + dialect_opts={"paramstyle": "named"}, + ) + + with context.begin_transaction(): + context.run_migrations() + + +def do_run_migrations(connection: Connection) -> None: + context.configure(connection=connection, target_metadata=target_metadata) + + with context.begin_transaction(): + context.run_migrations() + + +async def run_async_migrations() -> None: + configuration = config.get_section(config.config_ini_section, {}) + configuration["sqlalchemy.url"] = _database_url() + connectable = async_engine_from_config( + configuration, + prefix="sqlalchemy.", + poolclass=pool.NullPool, + ) + + async with connectable.connect() as connection: + await connection.run_sync(do_run_migrations) + + await connectable.dispose() + + +if context.is_offline_mode(): + run_migrations_offline() +else: + run(run_async_migrations()) diff --git a/alembic/script.py.mako b/alembic/script.py.mako new file mode 100644 index 0000000..2d50c24 --- /dev/null +++ b/alembic/script.py.mako @@ -0,0 +1,20 @@ +"""${message}""" + +from collections.abc import Sequence + +from alembic import op +import sqlalchemy as sa +${imports if imports else ""} + +revision: str = ${repr(up_revision)} +down_revision: str | None = ${repr(down_revision)} +branch_labels: str | Sequence[str] | None = ${repr(branch_labels)} +depends_on: str | Sequence[str] | None = ${repr(depends_on)} + + +def upgrade() -> None: + ${upgrades if upgrades else "pass"} + + +def downgrade() -> None: + ${downgrades if downgrades else "pass"} diff --git a/alembic/versions/20260412_028_create_orders_table.py b/alembic/versions/20260412_028_create_orders_table.py new file mode 100644 index 0000000..d92907a --- /dev/null +++ b/alembic/versions/20260412_028_create_orders_table.py @@ -0,0 +1,51 @@ +"""Create orders table.""" + +from collections.abc import Sequence + +from alembic import op +import sqlalchemy as sa +from sqlalchemy.dialects import postgresql + +revision: str = "20260412_028" +down_revision: str | None = None +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.create_table( + "orders", + sa.Column("id", postgresql.UUID(as_uuid=True), nullable=False), + sa.Column("order_uuid", sa.String(length=128), nullable=False), + sa.Column("payment_url", sa.String(length=2048), nullable=False), + sa.Column("price", sa.Integer(), nullable=False), + sa.Column("delivery_type", sa.Integer(), nullable=False), + sa.Column("tariff_code", sa.Integer(), nullable=False), + sa.Column("sender", postgresql.JSONB(astext_type=sa.Text()), nullable=False), + sa.Column("recipient", postgresql.JSONB(astext_type=sa.Text()), nullable=False), + sa.Column( + "from_location", + postgresql.JSONB(astext_type=sa.Text()), + nullable=False, + ), + sa.Column( + "to_location", + postgresql.JSONB(astext_type=sa.Text()), + nullable=False, + ), + sa.Column("packages", postgresql.JSONB(astext_type=sa.Text()), nullable=False), + sa.Column("services", postgresql.JSONB(astext_type=sa.Text()), nullable=True), + sa.Column("comment", sa.String(length=1024), nullable=True), + sa.Column( + "created_at", + sa.DateTime(timezone=True), + server_default=sa.text("now()"), + nullable=False, + ), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint("order_uuid", name="uq_orders_order_uuid"), + ) + + +def downgrade() -> None: + op.drop_table("orders") diff --git a/app/adapters/postgres/__init__.py b/app/adapters/postgres/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/adapters/postgres/engine.py b/app/adapters/postgres/engine.py new file mode 100644 index 0000000..a4b32f5 --- /dev/null +++ b/app/adapters/postgres/engine.py @@ -0,0 +1,16 @@ +"""PostgreSQL async engine and session factory management.""" + +from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker +from sqlalchemy.ext.asyncio import create_async_engine + +from app.config import PostgresConfig + + +def create_postgres_engine(config: PostgresConfig) -> AsyncEngine: + return create_async_engine(config.dsn, pool_pre_ping=True) + + +def create_postgres_session_factory( + engine: AsyncEngine, +) -> async_sessionmaker[AsyncSession]: + return async_sessionmaker(engine, expire_on_commit=False) diff --git a/app/adapters/tbank/base.py b/app/adapters/tbank/base.py index dba0a0f..29d7bcc 100644 --- a/app/adapters/tbank/base.py +++ b/app/adapters/tbank/base.py @@ -7,3 +7,18 @@ class TBankPaymentAdapterError(RuntimeError): class TBankPaymentRequestError(TBankPaymentAdapterError): """Raised when TBank rejects payment request data.""" + + def __init__( + self, + message: str, + *, + status_code: int | None = None, + error_code: str | None = None, + provider_message: str | None = None, + details: str | None = None, + ) -> None: + super().__init__(message) + self.status_code = status_code + self.error_code = error_code + self.provider_message = provider_message + self.details = details diff --git a/app/adapters/tbank/client.py b/app/adapters/tbank/client.py index 21832b4..e411760 100644 --- a/app/adapters/tbank/client.py +++ b/app/adapters/tbank/client.py @@ -88,9 +88,10 @@ class TBankAdapter: ) if 400 <= response.status_code < 500: - raise TBankPaymentRequestError( - "TBank payment init request was rejected with status " - f"{response.status_code}." + raise _build_tbank_request_error( + "TBank payment init request was rejected.", + status_code=response.status_code, + payload=_response_json_or_none(response), ) try: @@ -134,7 +135,10 @@ class TBankAdapter: ) if payload.get("Success") is False: - raise TBankPaymentRequestError("TBank payment init request was rejected.") + raise _build_tbank_request_error( + "TBank payment init request was rejected.", + payload=payload, + ) payment_url = payload.get("PaymentURL") if not isinstance(payment_url, str) or not payment_url.strip(): @@ -155,3 +159,71 @@ def _build_tbank_token(payload: dict[str, Any], *, password: str) -> str: str(token_payload[key]) for key in sorted(token_payload) ) return hashlib.sha256(token_source.encode("utf-8")).hexdigest() + + +def _response_json_or_none(response: httpx.Response) -> object | None: + try: + return response.json() + except (TypeError, ValueError): + return None + + +def _build_tbank_request_error( + message: str, + *, + status_code: int | None = None, + payload: object | None = None, +) -> TBankPaymentRequestError: + error_code = _payload_text_value(payload, "ErrorCode") + provider_message = _payload_text_value(payload, "Message") + details = _payload_text_value(payload, "Details") + return TBankPaymentRequestError( + _format_tbank_request_error_message( + message, + status_code=status_code, + error_code=error_code, + provider_message=provider_message, + details=details, + ), + status_code=status_code, + error_code=error_code, + provider_message=provider_message, + details=details, + ) + + +def _payload_text_value(payload: object | None, key: str) -> str | None: + if not isinstance(payload, dict): + return None + + value = payload.get(key) + if value is None: + return None + + text = str(value).strip() + if not text: + return None + return text + + +def _format_tbank_request_error_message( + message: str, + *, + status_code: int | None, + error_code: str | None, + provider_message: str | None, + details: str | None, +) -> str: + fields: list[str] = [] + if status_code is not None: + fields.append(f"status_code={status_code}") + if error_code is not None: + fields.append(f"error_code={error_code}") + if provider_message is not None: + fields.append(f"message={provider_message}") + if details is not None: + fields.append(f"details={details}") + + if not fields: + return message + return f"{message} {' '.join(fields)}" diff --git a/app/config.py b/app/config.py index a3fabe5..1075662 100644 --- a/app/config.py +++ b/app/config.py @@ -61,6 +61,10 @@ class TBankPaymentConfig(BaseModel): retry_backoff_seconds: float = Field(default=0.2, ge=0) +class PostgresConfig(BaseModel): + dsn: str = Field(..., min_length=1) + + class DadataAddressSuggestionsConfig(BaseModel): url: str = "https://suggestions.dadata.ru/suggestions/api/4_1/rs/suggest/address" api_key: str = "" @@ -121,6 +125,7 @@ class Settings(BaseSettings): repository: RepositoryConfig = Field(default_factory=RepositoryConfig) adapter: AdapterConfig = Field(default_factory=AdapterConfig) tbank_payment: TBankPaymentConfig + postgres: PostgresConfig address_suggestions: AddressSuggestionsConfig = Field( default_factory=AddressSuggestionsConfig ) @@ -153,6 +158,7 @@ class _RequiredYamlSections(BaseModel): repository: dict[str, Any] adapter: dict[str, Any] tbank_payment: dict[str, Any] + postgres: dict[str, Any] address_suggestions: dict[str, Any] observability: dict[str, Any] diff --git a/app/controllers/v1/delivery.py b/app/controllers/v1/delivery.py index 6863aa7..3780183 100644 --- a/app/controllers/v1/delivery.py +++ b/app/controllers/v1/delivery.py @@ -2,6 +2,7 @@ from fastapi import APIRouter, Depends, HTTPException, Request, status +from app.adapters.postgres.engine import create_postgres_engine, create_postgres_session_factory from app.adapters.address_suggestions.dadata import DadataAddressSuggestionProvider from app.adapters.address_suggestions.tomtom import TomTomAddressSuggestionProvider from app.adapters.address_suggestions.yandex_geosuggest import ( @@ -12,6 +13,7 @@ from app.adapters.tbank import TBankAdapter from app.config import Settings from app.controllers.http_client import build_controller_http_client from app.repositories.cache.redis_cache import PriceCache +from app.repositories.order import OrderRepository from app.schemas.payment import InitPaymentRequest, InitPaymentResponse from app.schemas.request import AddressSuggestRequest, DeliveryCalculationRequest from app.schemas.response import AddressSuggestion, DeliveryPrice @@ -58,10 +60,14 @@ def _build_aggregator_service(settings: Settings) -> AggregatorService: ) providers = (cdek_provider,) cache = PriceCache.from_repository_config(settings.repository) + postgres_engine = create_postgres_engine(settings.postgres) + postgres_session_factory = create_postgres_session_factory(postgres_engine) + order_repository = OrderRepository(session_factory=postgres_session_factory) service = AggregatorService( providers=providers, cache=cache, payment_adapter=payment_adapter, + order_repository=order_repository, address_suggestion_providers=( dadata_provider, yandex_geosuggest_provider, diff --git a/app/repositories/order/__init__.py b/app/repositories/order/__init__.py new file mode 100644 index 0000000..6a4ee04 --- /dev/null +++ b/app/repositories/order/__init__.py @@ -0,0 +1,5 @@ +"""Order repository exports.""" + +from app.repositories.order.repository import OrderData, OrderRepository + +__all__ = ("OrderData", "OrderRepository") diff --git a/app/repositories/order/models.py b/app/repositories/order/models.py new file mode 100644 index 0000000..0d9c7dd --- /dev/null +++ b/app/repositories/order/models.py @@ -0,0 +1,63 @@ +"""SQLAlchemy models for order persistence.""" + +from datetime import datetime +from typing import Any +from uuid import UUID, uuid4 + +from sqlalchemy import DateTime, Integer, String, UniqueConstraint, Uuid, func +from sqlalchemy.dialects.postgresql import JSONB +from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column +from sqlalchemy.types import JSON + + +def _json_payload_type() -> JSON: + return JSON().with_variant(JSONB, "postgresql") + + +class Base(DeclarativeBase): + pass + + +class Order(Base): + __tablename__ = "orders" + __table_args__ = ( + UniqueConstraint("order_uuid", name="uq_orders_order_uuid"), + ) + + id: Mapped[UUID] = mapped_column( + Uuid(as_uuid=True), + primary_key=True, + default=uuid4, + ) + order_uuid: Mapped[str] = mapped_column(String(128), nullable=False) + payment_url: Mapped[str] = mapped_column(String(2048), nullable=False) + price: Mapped[int] = mapped_column(Integer, nullable=False) + delivery_type: Mapped[int] = mapped_column(Integer, nullable=False) + tariff_code: Mapped[int] = mapped_column(Integer, nullable=False) + sender: Mapped[dict[str, Any]] = mapped_column(_json_payload_type(), nullable=False) + recipient: Mapped[dict[str, Any]] = mapped_column( + _json_payload_type(), + nullable=False, + ) + from_location: Mapped[dict[str, Any]] = mapped_column( + _json_payload_type(), + nullable=False, + ) + to_location: Mapped[dict[str, Any]] = mapped_column( + _json_payload_type(), + nullable=False, + ) + packages: Mapped[list[dict[str, Any]]] = mapped_column( + _json_payload_type(), + nullable=False, + ) + services: Mapped[list[dict[str, Any]] | None] = mapped_column( + _json_payload_type(), + nullable=True, + ) + comment: Mapped[str | None] = mapped_column(String(1024), nullable=True) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + server_default=func.now(), + nullable=False, + ) diff --git a/app/repositories/order/repository.py b/app/repositories/order/repository.py new file mode 100644 index 0000000..112a3ef --- /dev/null +++ b/app/repositories/order/repository.py @@ -0,0 +1,53 @@ +"""PostgreSQL order repository.""" + +from contextlib import AbstractAsyncContextManager +from dataclasses import dataclass +from typing import Any + +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from app.repositories.order.models import Order + + +@dataclass(frozen=True) +class OrderData: + order_uuid: str + payment_url: str + price: int + delivery_type: int + tariff_code: int + sender: dict[str, Any] + recipient: dict[str, Any] + from_location: dict[str, Any] + to_location: dict[str, Any] + packages: list[dict[str, Any]] + services: list[dict[str, Any]] | None + comment: str | None + + +class OrderRepository: + def __init__(self, session_factory: async_sessionmaker[AsyncSession]) -> None: + self._session_factory = session_factory + + def session(self) -> AbstractAsyncContextManager[AsyncSession]: + return self._session_factory.begin() + + async def create_order(self, session: AsyncSession, order_data: OrderData) -> Order: + order = Order( + order_uuid=order_data.order_uuid, + payment_url=order_data.payment_url, + price=order_data.price, + delivery_type=order_data.delivery_type, + tariff_code=order_data.tariff_code, + sender=order_data.sender, + recipient=order_data.recipient, + from_location=order_data.from_location, + to_location=order_data.to_location, + packages=order_data.packages, + services=order_data.services, + comment=order_data.comment, + ) + session.add(order) + await session.flush() + await session.refresh(order) + return order diff --git a/app/services/aggregator.py b/app/services/aggregator.py index fd7c830..c680891 100644 --- a/app/services/aggregator.py +++ b/app/services/aggregator.py @@ -4,9 +4,12 @@ import asyncio import hashlib import json from collections.abc import Iterable, Mapping, Sequence +from contextlib import AbstractAsyncContextManager from decimal import Decimal from typing import Protocol +import structlog + from app.adapters.address_suggestions.base import ( AddressSuggestionClientError, AddressSuggestionProvider, @@ -24,10 +27,13 @@ from app.domain.price import ( filter_and_sort_prices, normalize_delivery_request, ) +from app.repositories.order import OrderData from app.schemas.payment import InitPaymentRequest, InitPaymentResponse from app.schemas.request import AddressSuggestRequest, DeliveryCalculationRequest from app.schemas.response import AddressSuggestion, DeliveryPrice +logger = structlog.get_logger(__name__) + class AggregatorServiceError(RuntimeError): """Base exception for AggregatorService failures.""" @@ -67,6 +73,12 @@ class PaymentAdapterProtocol(Protocol): async def create_payment_link(self, order_uuid: str, amount_kopecks: int) -> str: ... +class OrderRepositoryProtocol(Protocol): + def session(self) -> AbstractAsyncContextManager[object]: ... + + async def create_order(self, session: object, order_data: OrderData) -> object: ... + + class FilterAndSortPricesFn(Protocol): def __call__( self, @@ -83,6 +95,7 @@ class AggregatorService: providers: Sequence[DeliveryProvider], cache: PriceCacheProtocol | None = None, payment_adapter: PaymentAdapterProtocol | None = None, + order_repository: OrderRepositoryProtocol | None = None, address_suggestion_providers: Sequence[AddressSuggestionProvider] = (), address_suggestion_country_to_provider: Mapping[str, str] | None = None, *, @@ -93,6 +106,7 @@ class AggregatorService: self._providers = tuple(providers) self._cache = cache self._payment_adapter = payment_adapter + self._order_repository = order_repository self._weight_round_scale = weight_round_scale self._provider_price_multiplier = provider_price_multiplier self._filter_and_sort_prices = filter_and_sort_prices_fn @@ -180,6 +194,14 @@ class AggregatorService: amount_kopecks=request.price, ) except TBankPaymentRequestError as exc: + logger.exception( + "payment_init_rejected", + order_uuid=request.order_uuid, + provider_status_code=exc.status_code, + provider_error_code=exc.error_code, + provider_error_message=exc.provider_message, + provider_error_details=exc.details, + ) raise InvalidInitPaymentRequestError( "Payment init request is invalid for the configured provider." ) from exc @@ -192,8 +214,52 @@ class AggregatorService: "Payment initialization is temporarily unavailable." ) from exc + await self._persist_order(request=request, payment_url=payment_url) return InitPaymentResponse(payment_url=payment_url) + async def _persist_order( + self, + *, + request: InitPaymentRequest, + payment_url: str, + ) -> None: + if self._order_repository is None: + return + + try: + async with self._order_repository.session() as session: + await self._order_repository.create_order( + session, + self._to_order_data(request=request, payment_url=payment_url), + ) + except Exception: + logger.exception( + "order_persistence_failed", + order_uuid=request.order_uuid, + ) + + @staticmethod + def _to_order_data( + *, + request: InitPaymentRequest, + payment_url: str, + ) -> OrderData: + payload = request.model_dump(mode="json") + return OrderData( + order_uuid=request.order_uuid, + payment_url=payment_url, + price=request.price, + delivery_type=request.type, + tariff_code=request.tariff_code, + sender=payload["sender"], + recipient=payload["recipient"], + from_location=payload["from_location"], + to_location=payload["to_location"], + packages=payload["packages"], + services=payload["services"], + comment=request.comment, + ) + async def _get_provider_prices( self, *, diff --git a/config.example.yaml b/config.example.yaml index 3d1e9ab..9b44a38 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -32,6 +32,9 @@ tbank_payment: retry_attempts: 2 retry_backoff_seconds: 0.2 +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@postgres:5432/g2s_aggregator" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/config.test.yaml b/config.test.yaml index db4af8c..890ce96 100644 --- a/config.test.yaml +++ b/config.test.yaml @@ -32,6 +32,9 @@ tbank_payment: retry_attempts: 2 retry_backoff_seconds: 0.2 +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/g2s_aggregator_test" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/docker-compose.yml b/docker-compose.yml index 3bc202a..5f3e7a3 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -5,12 +5,43 @@ services: ports: - "8000:8000" depends_on: - - redis + redis: + condition: service_started + postgres: + condition: service_healthy + migrations: + condition: service_completed_successfully volumes: - ./config.yaml:/config.yaml + migrations: + image: yusupal1ev/g2s-aggregator:0.0.2 + container_name: g2s-aggregator-migrations + depends_on: + postgres: + condition: service_healthy + volumes: + - ./config.yaml:/config.yaml + command: ["poetry", "run", "alembic", "upgrade", "head"] + restart: "no" redis: image: redis:7-alpine container_name: redis command: ["redis-server", "--save", "", "--appendonly", "no"] ports: - "6379:6379" + postgres: + image: postgres:16-alpine + container_name: postgres + environment: + POSTGRES_DB: g2s_aggregator + ports: + - "5432:5432" + volumes: + - postgres:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U postgres -d g2s_aggregator"] + interval: 5s + timeout: 3s + retries: 10 +volumes: + postgres: diff --git a/poetry.lock b/poetry.lock index a71cd81..3b0cae1 100644 --- a/poetry.lock +++ b/poetry.lock @@ -18,6 +18,21 @@ typing-extensions = "*" [package.extras] hiredis = ["hiredis (>=1.0)"] +[[package]] +name = "aiosqlite" +version = "0.22.1" +description = "asyncio bridge to the standard sqlite3 module" +optional = false +python-versions = ">=3.9" +files = [ + {file = "aiosqlite-0.22.1-py3-none-any.whl", hash = "sha256:21c002eb13823fad740196c5a2e9d8e62f6243bd9e7e4a1f87fb5e44ecb4fceb"}, + {file = "aiosqlite-0.22.1.tar.gz", hash = "sha256:043e0bd78d32888c0a9ca90fc788b38796843360c855a7262a532813133a0650"}, +] + +[package.extras] +dev = ["attribution (==1.8.0)", "black (==25.11.0)", "build (>=1.2)", "coverage[toml] (==7.10.7)", "flake8 (==7.3.0)", "flake8-bugbear (==24.12.12)", "flit (==3.12.0)", "mypy (==1.19.0)", "ufmt (==2.8.0)", "usort (==1.0.8.post1)"] +docs = ["sphinx (==8.1.3)", "sphinx-mdinclude (==0.6.2)"] + [[package]] name = "alembic" version = "1.18.4" @@ -1990,4 +2005,4 @@ type = ["pytest-mypy"] [metadata] lock-version = "2.0" python-versions = "^3.14" -content-hash = "a6d02bc869931fda634adb1c542cf314710f59e72d9962cb54b55571a42be2f2" +content-hash = "d3da307072ea8d39a768a338440f91ceaf72b09eab699bd06875714de74c8c90" diff --git a/pyproject.toml b/pyproject.toml index 5ac8e9a..92a7b3a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,6 +4,7 @@ version = "0.1.0" description = "" authors = ["Your Name "] readme = "README.md" +package-mode = false [tool.poetry.dependencies] python = "^3.14" @@ -25,9 +26,11 @@ redis = "^7.3.0" sqlalchemy = "^2.0.49" asyncpg = "^0.31.0" alembic = "^1.18.4" +aiosqlite = "^0.22.1" [tool.pytest.ini_options] pythonpath = ["."] +addopts = ["--import-mode=importlib"] [build-system] diff --git a/spec/index.md b/spec/index.md index a12c1c2..c562142 100644 --- a/spec/index.md +++ b/spec/index.md @@ -34,9 +34,10 @@ | 025 | DONE | 2026-04-03 | Remove company from Create Delivery Order parties | `spec/tasks/025_remove_company_from_create_delivery_order.md` | | 026 | DONE | 2026-04-05 | Align CDEK order contract with single phone and kilogram package weight | `spec/tasks/026_align_cdek_order_contract_single_phone_and_weight_units.md` | | 027 | DONE | 2026-04-11 | Add TBank payment adapter, init_payment endpoint and rename order flow | `spec/tasks/027_add_tbank_payment_adapter_and_order_payment_link.md` | +| 028 | DONE | 2026-04-12 | Add PostgreSQL adapter, order repository and persist order after payment link creation | `spec/tasks/028_add_postgresql_order_persistence.md` | ## Summary -- Total: **28** +- Total: **29** - TODO: **0** -- DONE: **28** +- DONE: **29** diff --git a/spec/overview.md b/spec/overview.md index ce37dd0..a28049a 100644 --- a/spec/overview.md +++ b/spec/overview.md @@ -20,6 +20,7 @@ - Предоставлять отдельный endpoint подсказок адреса, чтобы frontend мог получить точное значение для `from_location.address` и `to_location.address` перед инициализацией оплаты доставки - Принимать запрос на инициализацию оплаты доставки через TBank и возвращать ссылку на оплату без регистрации заказа в CDEK - Принимать сумму оплаты в поле `price` в копейках +- После успешного получения ссылки на оплату от TBank сохранять все данные заявки вместе со ссылкой на оплату в PostgreSQL; ошибка сохранения не блокирует возврат ссылки клиенту - Выбирать сервис подсказок адреса по `country_code` через маппинг стран в конфиге - Для `RU`, `BY` и `KZ`, сопоставленных с provider id `dadata`, использовать `dadata.ru` - Для `AM`, `AZ`, `KG`, `MD`, `TJ`, `TM` и `UZ`, сопоставленных с provider id `yandex_geosuggest`, использовать Yandex Geosuggest @@ -69,7 +70,7 @@ - `AggregatorService.suggest_addresses(request: AddressSuggestRequest) -> list[AddressSuggestion]` - `AggregatorService.suggest_addresses()` выбирает address suggestion provider по `country_code` через injected config mapping и оркестрирует ровно один adapter call - `AggregatorService.init_payment(request: InitPaymentRequest) -> InitPaymentResponse` -- `AggregatorService.init_payment()` оркестрирует инициализацию платежа через injected TBank adapter dependency +- `AggregatorService.init_payment()` оркестрирует инициализацию платежа через injected TBank adapter dependency, затем сохраняет данные заявки вместе с `payment_url` через injected order repository; ошибка сохранения логируется, но не блокирует возврат `payment_url` - Service не содержит бизнес-логики и provider HTTP-деталей ### Business Logic (`app/domain/`) @@ -86,6 +87,12 @@ - Операции: `get(key)`, `set(key, value, ttl)`, `invalidate(key)` - Без бизнес-решений; формирование ключа — ответственность Adapter +### Repository (`app/repositories/order/`) +- `OrderRepository` — сохранение данных заявки в PostgreSQL +- Операции: `create_order(session, order_data)` — сохраняет запись заявки с данными `InitPaymentRequest` и `payment_url` +- SQLAlchemy ORM model таблицы `orders` в `models.py` +- Без бизнес-решений; только CRUD-примитивы + ### Adapter (`app/adapters/delivery_providers`) - `base.py` — абстрактный интерфейс `DeliveryProvider`: ```python @@ -109,6 +116,12 @@ - TBank adapter инкапсулирует HTTP-взаимодействие с TBank Init API, auth token, retries, timeout, serialization и error handling - TBank adapter владеет собственной секцией конфигурации `tbank_payment` с полями `init_url`, `auth.terminal_key`, `auth.password`, `timeout_seconds`, `retry_attempts`, `retry_backoff_seconds` +### Adapter (`app/adapters/postgres`) +- Управление подключением к PostgreSQL: создание `AsyncEngine` и `async_sessionmaker` из конфигурации +- Инкапсулирует инфраструктурные детали SQLAlchemy async engine +- Владеет собственной секцией конфигурации `postgres` с обязательным полем `dsn` +- Без SQL-запросов, без бизнес-логики + ### Adapter (`app/adapters/address_suggestions`) - `base.py` — абстрактный интерфейс `AddressSuggestionProvider`: ```python @@ -223,6 +236,8 @@ payment_url: str | HTTP-клиент | httpx (async) | | Валидация | Pydantic v2 | | Кеш | Redis | +| База данных | PostgreSQL (asyncpg + SQLAlchemy 2.0 async) | +| Миграции | Alembic (async) | | Конфигурация | pydantic-settings | | Сервер | Uvicorn | @@ -259,15 +274,22 @@ app/ │ └── tbank/ │ ├── base.py │ └── client.py +│ └── postgres/ +│ └── engine.py # AsyncEngine & session factory ├── repositories/ -│ └── cache/ -│ └── redis_cache.py # Repository +│ ├── cache/ +│ │ └── redis_cache.py # Repository (Redis) +│ └── order/ +│ ├── models.py # SQLAlchemy ORM model +│ └── repository.py # Repository (PostgreSQL) ├── schemas/ │ ├── address.py │ ├── request.py │ ├── response.py │ └── payment.py -└── config.py +├── config.py +alembic/ # Alembic migrations +alembic.ini ``` --- @@ -278,4 +300,5 @@ app/ # docker-compose сервисы app # FastAPI-приложение redis # Кеш тарифов +postgres # База данных заявок ``` diff --git a/spec/tasks/028_add_postgresql_order_persistence.md b/spec/tasks/028_add_postgresql_order_persistence.md new file mode 100644 index 0000000..3a3cea7 --- /dev/null +++ b/spec/tasks/028_add_postgresql_order_persistence.md @@ -0,0 +1,75 @@ +--- +id: 028 +title: Add PostgreSQL adapter, order repository and persist order after payment link creation +status: DONE +created: 2026-04-12 +--- + +## Context +После создания ссылки на оплату через TBank adapter данные заявки нигде не сохраняются. Необходимо добавить PostgreSQL-адаптер для управления подключением к базе данных, репозиторий для сохранения данных заявки и интегрировать сохранение в существующий flow `AggregatorService.init_payment()`. Зависимости `sqlalchemy` (2.0.49) и `asyncpg` (0.31.0) уже присутствуют в `pyproject.toml`. Alembic (1.18.4) также доступен для миграций. + +## Goal +1. Добавить PostgreSQL-адаптер (`app/adapters/postgres/`) для управления async-сессиями SQLAlchemy (`AsyncEngine`, `async_sessionmaker`). +2. Добавить репозиторий заявок (`app/repositories/order/`) с методом `create_order()` для сохранения данных заявки вместе со ссылкой на оплату в PostgreSQL. +3. Определить SQLAlchemy model для таблицы заявок в `app/repositories/order/models.py`. +4. Создать Alembic-миграцию для создания таблицы заявок. +5. Расширить `AggregatorService.init_payment()`: после успешного получения `payment_url` от TBank adapter сохранять данные `InitPaymentRequest` вместе с `payment_url` через order repository. +6. Добавить секцию конфигурации PostgreSQL (`PostgresConfig`) в `app/config.py` и пример в `config.yaml`. +7. Добавить сервис PostgreSQL в `docker-compose.yml`. + +## Constraints +- PostgreSQL adapter (`app/adapters/postgres/`) MUST содержать только управление подключением (engine, session factory). Без бизнес-логики, без SQL-запросов. +- Repository (`app/repositories/order/`) MUST содержать только операции с базой данных. Без бизнес-решений, без workflow-логики. +- Service MUST оркестрировать вызовы TBank adapter и order repository. Если сохранение в БД завершается ошибкой после успешного получения `payment_url`, Service MUST всё равно вернуть `payment_url` клиенту (сохранение не должно блокировать ответ); ошибку сохранения логировать. +- SQLAlchemy model MUST использовать `sqlalchemy.orm.DeclarativeBase` (SQLAlchemy 2.0 style). +- Для миграций использовать Alembic с async-конфигурацией (`asyncpg`). +- Конфигурация PostgreSQL MUST быть в отдельной секции `postgres` в `config.yaml` с обязательным полем `dsn`; `config.yaml` уже в `.gitignore`. +- `_RequiredYamlSections` в `app/config.py` MUST быть обновлён для включения секции `postgres`. +- Order repository передаётся в `AggregatorService` через dependency injection (новый параметр конструктора). +- PostgreSQL adapter создаётся в wiring (`_build_aggregator_service`) в controller и передаёт session factory в order repository. +- Таблица заявок MUST содержать как минимум: `id` (UUID, PK), `order_uuid` (str, unique), `payment_url` (str), `price` (int, копейки), `tariff_code` (int), `sender` (JSONB), `recipient` (JSONB), `from_location` (JSONB), `to_location` (JSONB), `packages` (JSONB), `services` (JSONB, nullable), `comment` (str, nullable), `created_at` (timestamp with timezone, server default). +- Scope НЕ включает: чтение/обновление/удаление заявок, API-endpoint для списка заявок, webhook-обработку платёжных уведомлений, изменения price flow, address suggestion flow. +- НЕ изменять существующие тесты TBank adapter, не изменять поведение price и address suggestion endpoints. + +## Acceptance criteria +- В `app/adapters/postgres/` существует модуль с функцией создания `AsyncEngine` и `async_sessionmaker` из конфигурации. +- В `app/repositories/order/` существует `OrderRepository` с async-методом `create_order(session, order_data)`, сохраняющим запись заявки. +- В `app/repositories/order/models.py` определена SQLAlchemy ORM model таблицы `orders` со всеми обязательными полями. +- Alembic инициализирован с async-конфигурацией; существует миграция для создания таблицы `orders`. +- `AggregatorService.__init__()` принимает опциональный `order_repository` через DI. +- `AggregatorService.init_payment()` после успешного получения `payment_url` вызывает `order_repository.create_order()` с данными из `InitPaymentRequest` и `payment_url`. +- Если `order_repository.create_order()` выбрасывает исключение, `init_payment()` логирует ошибку и возвращает `InitPaymentResponse(payment_url=...)` без ошибки клиенту. +- В `app/config.py` добавлена `PostgresConfig` с полем `dsn: str`. +- Секция `postgres` присутствует в `_RequiredYamlSections`. +- В `docker-compose.yml` добавлен сервис `postgres` и `app` зависит от него. +- Wiring в `_build_aggregator_service` создаёт PostgreSQL engine, session factory, `OrderRepository` и передаёт его в `AggregatorService`. +- Запросы к эндпоинту `POST /api/v1/delivery/init-payment` продолжают возвращать `InitPaymentResponse` с `payment_url`. + +## Definition of Done +- [ ] Создан модуль `app/adapters/postgres/` с engine/session factory. +- [ ] Создан `app/repositories/order/repository.py` с `OrderRepository.create_order()`. +- [ ] Создан `app/repositories/order/models.py` с ORM model таблицы `orders`. +- [ ] Alembic инициализирован (`alembic.ini`, `alembic/`), создана миграция для таблицы `orders`. +- [ ] Добавлена `PostgresConfig` в `app/config.py`; `_RequiredYamlSections` обновлён. +- [ ] `AggregatorService` принимает `order_repository` через DI и использует его в `init_payment()`. +- [ ] Ошибки сохранения заявки не блокируют возврат `payment_url` клиенту. +- [ ] Обновлён wiring в `app/controllers/v1/delivery.py`. +- [ ] Добавлен сервис `postgres` в `docker-compose.yml`. +- [ ] `config.yaml` пример содержит секцию `postgres`. +- [ ] `config.test.yaml` содержит секцию `postgres` (может использовать sqlite или тестовый DSN). +- [ ] Все существующие тесты продолжают проходить. +- [ ] Добавлены новые тесты. + +## Tests +- Добавить `tests/repositories/order/test_repository.py`: проверка `create_order()` с in-memory SQLite async engine (SQLAlchemy async); проверка, что все обязательные поля сохраняются; проверка обработки дублирования `order_uuid` (unique constraint). +- Обновить `tests/services/test_init_payment.py`: добавить test case, где `order_repository.create_order()` вызывается после успешного создания payment link; добавить test case, где `order_repository.create_order()` выбрасывает исключение, а `init_payment()` всё равно возвращает `payment_url`. +- Обновить `tests/config/test_config_sections.py` для проверки наличия секции `postgres` в yaml. +- При необходимости обновить `tests/smoke/test_app_import.py` для проверки wiring order repository. + +## Commands +- `poetry run pytest tests/repositories/order/test_repository.py -q` +- `poetry run pytest tests/services/test_init_payment.py -q` +- `poetry run pytest tests/config/test_config_sections.py -q` +- `poetry run pytest tests/smoke/test_app_import.py -q` +- `poetry run pytest -q` +- `python3 spec/gen_spec_index.py --check` diff --git a/tests/adapters/delivery_providers/cdek/test_client.py b/tests/adapters/delivery_providers/cdek/test_client.py index 1a94ab3..1b5b9a0 100644 --- a/tests/adapters/delivery_providers/cdek/test_client.py +++ b/tests/adapters/delivery_providers/cdek/test_client.py @@ -302,6 +302,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/cdek_adapter_test" + observability: enabled: false service_name: "cdek-adapter-test-service" diff --git a/tests/adapters/tbank/test_client.py b/tests/adapters/tbank/test_client.py index b1ac022..68907c6 100644 --- a/tests/adapters/tbank/test_client.py +++ b/tests/adapters/tbank/test_client.py @@ -111,28 +111,56 @@ def test_create_payment_link_posts_signed_payload_and_maps_payment_url() -> None def test_create_payment_link_maps_4xx_to_request_error() -> None: response = httpx.Response( 400, - json={"Success": False, "Message": "bad request"}, + json={ + "Success": False, + "ErrorCode": "101", + "Message": "bad request", + "Details": "invalid token", + }, request=httpx.Request("POST", "https://securepay.tinkoff.ru/v2/Init"), ) http_client = SequenceHTTPClient([response]) adapter = _make_adapter(http_client) - with pytest.raises(TBankPaymentRequestError, match="status 400"): + with pytest.raises(TBankPaymentRequestError) as exc_info: asyncio.run(adapter.create_payment_link("order-uuid-1", 125000)) + assert exc_info.value.status_code == 400 + assert exc_info.value.error_code == "101" + assert exc_info.value.provider_message == "bad request" + assert exc_info.value.details == "invalid token" + assert str(exc_info.value) == ( + "TBank payment init request was rejected. " + "status_code=400 error_code=101 message=bad request details=invalid token" + ) + def test_create_payment_link_maps_unsuccessful_payload_to_request_error() -> None: response = httpx.Response( 200, - json={"Success": False, "Message": "bad request"}, + json={ + "Success": False, + "ErrorCode": "102", + "Message": "duplicate order id", + "Details": "OrderId must be unique", + }, request=httpx.Request("POST", "https://securepay.tinkoff.ru/v2/Init"), ) http_client = SequenceHTTPClient([response]) adapter = _make_adapter(http_client) - with pytest.raises(TBankPaymentRequestError): + with pytest.raises(TBankPaymentRequestError) as exc_info: asyncio.run(adapter.create_payment_link("order-uuid-1", 125000)) + assert exc_info.value.status_code is None + assert exc_info.value.error_code == "102" + assert exc_info.value.provider_message == "duplicate order id" + assert exc_info.value.details == "OrderId must be unique" + assert str(exc_info.value) == ( + "TBank payment init request was rejected. " + "error_code=102 message=duplicate order id details=OrderId must be unique" + ) + def test_create_payment_link_maps_transport_errors_to_client_error() -> None: request = httpx.Request("POST", "https://securepay.tinkoff.ru/v2/Init") diff --git a/tests/config/fixtures/config.default.yaml b/tests/config/fixtures/config.default.yaml index 1aee977..4f1d9b3 100644 --- a/tests/config/fixtures/config.default.yaml +++ b/tests/config/fixtures/config.default.yaml @@ -21,6 +21,9 @@ tbank_payment: terminal_key: "yaml-terminal-key" password: "yaml-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/default" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.invalid_price_multiplier.yaml b/tests/config/fixtures/config.invalid_price_multiplier.yaml index 9fae52a..3441b80 100644 --- a/tests/config/fixtures/config.invalid_price_multiplier.yaml +++ b/tests/config/fixtures/config.invalid_price_multiplier.yaml @@ -29,6 +29,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/invalid_price" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.missing_adapter.yaml b/tests/config/fixtures/config.missing_adapter.yaml index 2bdc089..ff5ea38 100644 --- a/tests/config/fixtures/config.missing_adapter.yaml +++ b/tests/config/fixtures/config.missing_adapter.yaml @@ -18,6 +18,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/missing_adapter" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.missing_address_suggestions.yaml b/tests/config/fixtures/config.missing_address_suggestions.yaml index 9a52a80..70632d0 100644 --- a/tests/config/fixtures/config.missing_address_suggestions.yaml +++ b/tests/config/fixtures/config.missing_address_suggestions.yaml @@ -29,6 +29,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/missing_address" + observability: enabled: false service_name: "missing-address-suggestions-service" diff --git a/tests/config/fixtures/config.missing_observability.yaml b/tests/config/fixtures/config.missing_observability.yaml index cb0f103..3f6a7ce 100644 --- a/tests/config/fixtures/config.missing_observability.yaml +++ b/tests/config/fixtures/config.missing_observability.yaml @@ -29,6 +29,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/missing_observability" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.missing_observability_endpoint.yaml b/tests/config/fixtures/config.missing_observability_endpoint.yaml index a6a5940..d26120e 100644 --- a/tests/config/fixtures/config.missing_observability_endpoint.yaml +++ b/tests/config/fixtures/config.missing_observability_endpoint.yaml @@ -29,6 +29,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/missing_observability_endpoint" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.missing_price_multiplier.yaml b/tests/config/fixtures/config.missing_price_multiplier.yaml index 4921d42..e140dc0 100644 --- a/tests/config/fixtures/config.missing_price_multiplier.yaml +++ b/tests/config/fixtures/config.missing_price_multiplier.yaml @@ -28,6 +28,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/missing_price" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/fixtures/config.test.override.yaml b/tests/config/fixtures/config.test.override.yaml index 796591d..974be73 100644 --- a/tests/config/fixtures/config.test.override.yaml +++ b/tests/config/fixtures/config.test.override.yaml @@ -24,6 +24,9 @@ tbank_payment: retry_attempts: 1 retry_backoff_seconds: 0.05 +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/override" + address_suggestions: country_to_provider: RU: "dadata" diff --git a/tests/config/test_config_sections.py b/tests/config/test_config_sections.py index 30ee9fb..b67a1ed 100644 --- a/tests/config/test_config_sections.py +++ b/tests/config/test_config_sections.py @@ -132,6 +132,10 @@ def test_configuration_sections_are_loaded_from_yaml_file( assert settings.tbank_payment.timeout_seconds == 10.0 assert settings.tbank_payment.retry_attempts == 2 assert settings.tbank_payment.retry_backoff_seconds == 0.2 + assert ( + settings.postgres.dsn + == "postgresql+asyncpg://postgres:postgres@localhost:5432/g2s_aggregator_test" + ) assert settings.address_suggestions.country_to_provider == _expected_country_mapping() assert ( settings.address_suggestions.dadata.url @@ -199,6 +203,7 @@ def test_get_settings_returns_cached_instance(monkeypatch: pytest.MonkeyPatch) - assert first.service.provider_timeout_seconds == 10.0 assert first.business_logic.provider_price_multiplier == Decimal("1.0") assert first.tbank_payment.auth.terminal_key == "test-terminal-key" + assert first.postgres.dsn.endswith("/g2s_aggregator_test") assert first.address_suggestions.country_to_provider["RU"] == "dadata" assert first.observability.service_name == "g2s-aggregator-test" get_settings.cache_clear() @@ -224,6 +229,10 @@ def test_get_settings_uses_config_test_yaml_in_pytest_environment( assert settings.tbank_payment.timeout_seconds == 4.25 assert settings.tbank_payment.retry_attempts == 1 assert settings.tbank_payment.retry_backoff_seconds == 0.05 + assert ( + settings.postgres.dsn + == "postgresql+asyncpg://postgres:postgres@localhost:5432/override" + ) assert settings.address_suggestions.country_to_provider == { "RU": "dadata", "AM": "yandex_geosuggest", diff --git a/tests/repositories/cache/test_redis_cache.py b/tests/repositories/cache/test_redis_cache.py index e5bc51a..43fba2a 100644 --- a/tests/repositories/cache/test_redis_cache.py +++ b/tests/repositories/cache/test_redis_cache.py @@ -185,6 +185,9 @@ tbank_payment: terminal_key: "test-terminal-key" password: "test-password" +postgres: + dsn: "postgresql+asyncpg://postgres:postgres@localhost:5432/repository_test" + observability: enabled: false service_name: "repository-test-service" diff --git a/tests/repositories/order/test_repository.py b/tests/repositories/order/test_repository.py new file mode 100644 index 0000000..84867f4 --- /dev/null +++ b/tests/repositories/order/test_repository.py @@ -0,0 +1,145 @@ +import asyncio +from collections.abc import Awaitable, Callable +from typing import Any + +import pytest +from sqlalchemy import select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker +from sqlalchemy.ext.asyncio import create_async_engine + +from app.repositories.order import OrderData, OrderRepository +from app.repositories.order.models import Base, Order + + +def _make_order_data(**overrides: object) -> OrderData: + payload: dict[str, Any] = { + "order_uuid": "order-uuid-1", + "payment_url": "https://pay.test/payment/1", + "price": 125000, + "delivery_type": 2, + "tariff_code": 535, + "comment": "Test payment", + "sender": { + "name": "Petr Petrov", + "email": "sender@example.com", + "phone": {"number": "+79009876543"}, + }, + "recipient": { + "name": "Ivan Ivanov", + "email": "ivan@example.com", + "phone": {"number": "+79001234567"}, + }, + "from_location": { + "address": "Lenina 1", + "city": "Moscow", + "country_code": "RU", + }, + "to_location": { + "address": "Pushkina 10", + "city": "Novosibirsk", + "country_code": "RU", + }, + "services": [{"code": "INSURANCE", "parameter": "1000"}], + "packages": [ + { + "number": "1", + "weight": 1, + "length": 20, + "width": 15, + "height": 10, + "comment": "Package 1", + } + ], + } + payload.update(overrides) + return OrderData(**payload) + + +async def _with_repository( + test_fn: Callable[ + [OrderRepository, async_sessionmaker[AsyncSession]], + Awaitable[None], + ], +) -> None: + engine = create_async_engine("sqlite+aiosqlite:///:memory:") + try: + async with engine.begin() as connection: + await connection.run_sync(Base.metadata.create_all) + session_factory = async_sessionmaker(engine, expire_on_commit=False) + await test_fn(OrderRepository(session_factory=session_factory), session_factory) + finally: + await engine.dispose() + + +def test_create_order_persists_all_required_fields() -> None: + async def run( + repository: OrderRepository, + session_factory: async_sessionmaker[AsyncSession], + ) -> None: + order_data = _make_order_data() + + async with repository.session() as session: + order = await repository.create_order(session, order_data) + + async with session_factory() as session: + result = await session.execute( + select(Order).where(Order.order_uuid == "order-uuid-1") + ) + persisted_order = result.scalar_one() + + assert order.id == persisted_order.id + assert persisted_order.order_uuid == "order-uuid-1" + assert persisted_order.payment_url == "https://pay.test/payment/1" + assert persisted_order.price == 125000 + assert persisted_order.delivery_type == 2 + assert persisted_order.tariff_code == 535 + assert persisted_order.sender == order_data.sender + assert persisted_order.recipient == order_data.recipient + assert persisted_order.from_location == order_data.from_location + assert persisted_order.to_location == order_data.to_location + assert persisted_order.packages == order_data.packages + assert persisted_order.services == order_data.services + assert persisted_order.comment == "Test payment" + assert persisted_order.created_at is not None + + asyncio.run(_with_repository(run)) + + +def test_create_order_persists_nullable_services_and_comment() -> None: + async def run( + repository: OrderRepository, + session_factory: async_sessionmaker[AsyncSession], + ) -> None: + order_data = _make_order_data(services=None, comment=None) + + async with repository.session() as session: + await repository.create_order(session, order_data) + + async with session_factory() as session: + result = await session.execute( + select(Order).where(Order.order_uuid == "order-uuid-1") + ) + persisted_order = result.scalar_one() + + assert persisted_order.services is None + assert persisted_order.comment is None + + asyncio.run(_with_repository(run)) + + +def test_create_order_rejects_duplicate_order_uuid() -> None: + async def run( + repository: OrderRepository, + _session_factory: async_sessionmaker[AsyncSession], + ) -> None: + order_data = _make_order_data() + + async with repository.session() as session: + await repository.create_order(session, order_data) + + with pytest.raises(IntegrityError): + async with repository.session() as session: + await repository.create_order(session, order_data) + + asyncio.run(_with_repository(run)) diff --git a/tests/services/test_init_payment.py b/tests/services/test_init_payment.py index 446c467..44fff7c 100644 --- a/tests/services/test_init_payment.py +++ b/tests/services/test_init_payment.py @@ -6,6 +6,7 @@ from app.adapters.tbank.base import ( TBankPaymentAdapterError, TBankPaymentRequestError, ) +from app.repositories.order import OrderData from app.schemas.payment import InitPaymentRequest, InitPaymentResponse from app.services.aggregator import ( AggregatorService, @@ -34,6 +35,33 @@ class StubPaymentAdapter: return self._response +class StubOrderSessionContext: + def __init__(self, session: object) -> None: + self._session = session + + async def __aenter__(self) -> object: + return self._session + + async def __aexit__(self, exc_type, exc, traceback) -> None: + return None + + +class StubOrderRepository: + def __init__(self, *, error: Exception | None = None) -> None: + self._error = error + self.session_value = object() + self.calls: list[tuple[object, OrderData]] = [] + + def session(self) -> StubOrderSessionContext: + return StubOrderSessionContext(self.session_value) + + async def create_order(self, session: object, order_data: OrderData) -> object: + self.calls.append((session, order_data)) + if self._error is not None: + raise self._error + return object() + + def _make_init_payment_request(**overrides: object) -> InitPaymentRequest: payload: dict[str, object] = { "order_uuid": "order-uuid-1", @@ -88,6 +116,83 @@ def test_init_payment_delegates_to_adapter_and_returns_payment_url() -> None: assert adapter.calls == [("order-uuid-1", 125000)] +def test_init_payment_persists_order_after_successful_payment_link() -> None: + request = _make_init_payment_request() + adapter = StubPaymentAdapter(response="https://pay.test/payment/1") + order_repository = StubOrderRepository() + service = AggregatorService( + providers=[], + payment_adapter=adapter, + order_repository=order_repository, + ) + + result = asyncio.run(service.init_payment(request)) + + assert result == InitPaymentResponse(payment_url="https://pay.test/payment/1") + assert adapter.calls == [("order-uuid-1", 125000)] + assert order_repository.calls == [ + ( + order_repository.session_value, + OrderData( + order_uuid="order-uuid-1", + payment_url="https://pay.test/payment/1", + price=125000, + delivery_type=2, + tariff_code=535, + sender={ + "name": "Petr Petrov", + "email": "sender@example.com", + "phone": {"number": "+79009876543"}, + }, + recipient={ + "name": "Ivan Ivanov", + "email": "ivan@example.com", + "phone": {"number": "+79001234567"}, + }, + from_location={ + "address": "Lenina 1", + "city": "Moscow", + "country_code": "RU", + }, + to_location={ + "address": "Pushkina 10", + "city": "Novosibirsk", + "country_code": "RU", + }, + packages=[ + { + "number": "1", + "weight": 1, + "length": 20, + "width": 15, + "height": 10, + "comment": "Package 1", + } + ], + services=[{"code": "INSURANCE", "parameter": "1000"}], + comment="Test payment", + ), + ) + ] + + +def test_init_payment_returns_payment_url_when_order_persistence_fails() -> None: + request = _make_init_payment_request() + adapter = StubPaymentAdapter(response="https://pay.test/payment/1") + order_repository = StubOrderRepository(error=RuntimeError("database down")) + service = AggregatorService( + providers=[], + payment_adapter=adapter, + order_repository=order_repository, + ) + + result = asyncio.run(service.init_payment(request)) + + assert result == InitPaymentResponse(payment_url="https://pay.test/payment/1") + assert adapter.calls == [("order-uuid-1", 125000)] + assert len(order_repository.calls) == 1 + + def test_init_payment_maps_provider_request_errors_to_invalid_payment_error() -> None: request = _make_init_payment_request() adapter = StubPaymentAdapter(error=TBankPaymentRequestError("bad payload")) diff --git a/tests/smoke/test_local_infra_stack.py b/tests/smoke/test_local_infra_stack.py index 2bf161e..7108674 100644 --- a/tests/smoke/test_local_infra_stack.py +++ b/tests/smoke/test_local_infra_stack.py @@ -18,7 +18,8 @@ def test_compose_defines_required_services() -> None: compose = _load_compose() services = compose.get("services") assert isinstance(services, dict) - assert {"app", "redis"}.issubset(services.keys()) + assert {"app", "redis", "postgres"}.issubset(services.keys()) + assert "postgres" in services["app"]["depends_on"] def test_smoke_command_sequence_is_documented() -> None: