Добавлена отправка накладных на почту
This commit is contained in:
@@ -253,6 +253,55 @@ class CDEKClient:
|
||||
"CDEK get waybill response payload is invalid."
|
||||
) from exc
|
||||
|
||||
async def download_waybill_pdf(self, url: str) -> bytes:
|
||||
failure_message = "CDEK waybill PDF download failed"
|
||||
for attempt in range(self._retry_attempts + 1):
|
||||
try:
|
||||
token = await self._auth_client.get_access_token()
|
||||
response = await self._http_client.get(
|
||||
url,
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
timeout=self._timeout_seconds,
|
||||
)
|
||||
except (httpx.TimeoutException, httpx.TransportError) as exc:
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
f"{failure_message} after retry attempts."
|
||||
) from exc
|
||||
|
||||
if self._should_retry(response.status_code):
|
||||
if attempt < self._retry_attempts:
|
||||
await self._sleep(self._retry_delay(attempt))
|
||||
continue
|
||||
raise CDEKClientError(
|
||||
f"{failure_message} with retriable status "
|
||||
f"{response.status_code}."
|
||||
)
|
||||
|
||||
if 400 <= response.status_code < 500:
|
||||
log.warning(
|
||||
"cdek_waybill_pdf_download_rejected",
|
||||
url=url,
|
||||
status_code=response.status_code,
|
||||
response_body=_response_text_or_none(response),
|
||||
)
|
||||
raise CDEKRequestError(
|
||||
f"{failure_message} with status {response.status_code}."
|
||||
)
|
||||
|
||||
try:
|
||||
response.raise_for_status()
|
||||
except httpx.HTTPError as exc:
|
||||
raise CDEKClientError(
|
||||
f"{failure_message}: invalid response."
|
||||
) from exc
|
||||
|
||||
return response.content
|
||||
|
||||
raise CDEKClientError(f"{failure_message} unexpectedly.")
|
||||
|
||||
async def _get_json(self, url: str, *, failure_message: str) -> dict[str, Any]:
|
||||
for attempt in range(self._retry_attempts + 1):
|
||||
try:
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
"""E-mail adapters."""
|
||||
|
||||
from app.adapters.email.smtp_client import SMTPEmailSender, SMTPEmailSenderError
|
||||
|
||||
__all__ = ["SMTPEmailSender", "SMTPEmailSenderError"]
|
||||
@@ -0,0 +1,135 @@
|
||||
"""SMTP e-mail adapter built on top of aiosmtplib."""
|
||||
|
||||
from email.message import EmailMessage
|
||||
from typing import Protocol
|
||||
|
||||
import aiosmtplib
|
||||
import structlog
|
||||
|
||||
log = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
class SMTPEmailSenderError(Exception):
|
||||
"""Raised when the SMTP delivery fails."""
|
||||
|
||||
|
||||
class SMTPSendFn(Protocol):
|
||||
async def __call__(
|
||||
self,
|
||||
message: EmailMessage,
|
||||
*,
|
||||
hostname: str,
|
||||
port: int,
|
||||
username: str | None,
|
||||
password: str | None,
|
||||
use_tls: bool,
|
||||
start_tls: bool,
|
||||
timeout: float,
|
||||
) -> object: ...
|
||||
|
||||
|
||||
class SMTPEmailSender:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
smtp_host: str,
|
||||
smtp_port: int,
|
||||
username: str,
|
||||
password: str,
|
||||
from_address: str,
|
||||
use_tls: bool = True,
|
||||
timeout_seconds: float = 10.0,
|
||||
send: SMTPSendFn | None = None,
|
||||
) -> None:
|
||||
self._smtp_host = smtp_host
|
||||
self._smtp_port = smtp_port
|
||||
self._username = username
|
||||
self._password = password
|
||||
self._from_address = from_address
|
||||
self._use_tls = use_tls
|
||||
self._timeout_seconds = timeout_seconds
|
||||
self._send = send or _default_send
|
||||
|
||||
async def send_email(
|
||||
self,
|
||||
*,
|
||||
to: str,
|
||||
subject: str,
|
||||
body: str,
|
||||
attachment_bytes: bytes,
|
||||
attachment_filename: str,
|
||||
) -> None:
|
||||
message = self._build_message(
|
||||
to=to,
|
||||
subject=subject,
|
||||
body=body,
|
||||
attachment_bytes=attachment_bytes,
|
||||
attachment_filename=attachment_filename,
|
||||
)
|
||||
try:
|
||||
await self._send(
|
||||
message,
|
||||
hostname=self._smtp_host,
|
||||
port=self._smtp_port,
|
||||
username=self._username or None,
|
||||
password=self._password or None,
|
||||
use_tls=self._use_tls and self._smtp_port == 465,
|
||||
start_tls=self._use_tls and self._smtp_port != 465,
|
||||
timeout=self._timeout_seconds,
|
||||
)
|
||||
except Exception as exc:
|
||||
log.warning(
|
||||
"smtp_send_failed",
|
||||
smtp_host=self._smtp_host,
|
||||
smtp_port=self._smtp_port,
|
||||
to=to,
|
||||
error=str(exc),
|
||||
)
|
||||
raise SMTPEmailSenderError(
|
||||
f"SMTP send to {to} failed: {exc}"
|
||||
) from exc
|
||||
|
||||
def _build_message(
|
||||
self,
|
||||
*,
|
||||
to: str,
|
||||
subject: str,
|
||||
body: str,
|
||||
attachment_bytes: bytes,
|
||||
attachment_filename: str,
|
||||
) -> EmailMessage:
|
||||
message = EmailMessage()
|
||||
message["From"] = self._from_address
|
||||
message["To"] = to
|
||||
message["Subject"] = subject
|
||||
message.set_content(body)
|
||||
message.add_attachment(
|
||||
attachment_bytes,
|
||||
maintype="application",
|
||||
subtype="pdf",
|
||||
filename=attachment_filename,
|
||||
)
|
||||
return message
|
||||
|
||||
|
||||
async def _default_send(
|
||||
message: EmailMessage,
|
||||
*,
|
||||
hostname: str,
|
||||
port: int,
|
||||
username: str | None,
|
||||
password: str | None,
|
||||
use_tls: bool,
|
||||
start_tls: bool,
|
||||
timeout: float,
|
||||
) -> object:
|
||||
return await aiosmtplib.send(
|
||||
message,
|
||||
hostname=hostname,
|
||||
port=port,
|
||||
username=username,
|
||||
password=password,
|
||||
use_tls=use_tls,
|
||||
start_tls=start_tls,
|
||||
timeout=timeout,
|
||||
)
|
||||
@@ -121,6 +121,21 @@ class WaybillPollerConfig(BaseModel):
|
||||
batch_size: int = Field(default=50, gt=0)
|
||||
|
||||
|
||||
class EmailAdapterConfig(BaseModel):
|
||||
smtp_host: str = Field(..., min_length=1)
|
||||
smtp_port: int = Field(..., gt=0, le=65535)
|
||||
username: str = ""
|
||||
password: str = ""
|
||||
from_address: str = Field(..., min_length=1)
|
||||
use_tls: bool = True
|
||||
timeout_seconds: float = Field(default=10.0, gt=0)
|
||||
|
||||
|
||||
class WaybillEmailSenderConfig(BaseModel):
|
||||
interval_seconds: float = Field(default=30.0, gt=0)
|
||||
batch_size: int = Field(default=50, gt=0)
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
model_config = SettingsConfigDict(
|
||||
extra="ignore",
|
||||
@@ -138,6 +153,10 @@ class Settings(BaseSettings):
|
||||
)
|
||||
observability: ObservabilityConfig
|
||||
waybill_poller: WaybillPollerConfig = Field(default_factory=WaybillPollerConfig)
|
||||
email: EmailAdapterConfig
|
||||
waybill_email_sender: WaybillEmailSenderConfig = Field(
|
||||
default_factory=WaybillEmailSenderConfig
|
||||
)
|
||||
|
||||
@classmethod
|
||||
def settings_customise_sources(
|
||||
@@ -169,6 +188,7 @@ class _RequiredYamlSections(BaseModel):
|
||||
postgres: dict[str, Any]
|
||||
address_suggestions: dict[str, Any]
|
||||
observability: dict[str, Any]
|
||||
email: dict[str, Any]
|
||||
|
||||
|
||||
def _resolve_runtime_config_file() -> str:
|
||||
|
||||
@@ -48,6 +48,10 @@ class Order(Base):
|
||||
DateTime(timezone=True),
|
||||
nullable=True,
|
||||
)
|
||||
waybill_email_sent_at: Mapped[datetime | None] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
nullable=True,
|
||||
)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True),
|
||||
server_default=func.now(),
|
||||
|
||||
@@ -144,3 +144,38 @@ class OrderRepository:
|
||||
order.cdek_polled_at = polled_at
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
async def list_orders_pending_waybill_email(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
*,
|
||||
limit: int,
|
||||
) -> Sequence[Order]:
|
||||
statement = (
|
||||
select(Order)
|
||||
.where(
|
||||
Order.cdek_waybill_url.is_not(None),
|
||||
Order.waybill_email_sent_at.is_(None),
|
||||
)
|
||||
.order_by(Order.created_at.asc())
|
||||
.limit(limit)
|
||||
.with_for_update(skip_locked=True)
|
||||
)
|
||||
result = await session.execute(statement)
|
||||
return result.scalars().all()
|
||||
|
||||
async def record_waybill_email_sent(
|
||||
self,
|
||||
session: AsyncSession,
|
||||
*,
|
||||
order_uuid: str,
|
||||
sent_at: datetime,
|
||||
) -> Order | None:
|
||||
order = await self.get_order_by_order_uuid(session, order_uuid)
|
||||
if order is None:
|
||||
return None
|
||||
|
||||
if order.waybill_email_sent_at is None:
|
||||
order.waybill_email_sent_at = sent_at
|
||||
await session.flush()
|
||||
return order
|
||||
|
||||
@@ -4,7 +4,14 @@ from datetime import datetime
|
||||
from decimal import Decimal, InvalidOperation
|
||||
from typing import Literal
|
||||
|
||||
from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
|
||||
from pydantic import (
|
||||
BaseModel,
|
||||
ConfigDict,
|
||||
EmailStr,
|
||||
Field,
|
||||
field_validator,
|
||||
model_validator,
|
||||
)
|
||||
from pydantic.alias_generators import to_camel
|
||||
|
||||
|
||||
@@ -115,7 +122,7 @@ class InitPaymentRequest(_CamelModel):
|
||||
content: Content
|
||||
pickup_date: datetime
|
||||
delivery_date: datetime | None = None
|
||||
account_email: str = Field(min_length=1)
|
||||
account_email: EmailStr
|
||||
system_data: SystemData
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,167 @@
|
||||
"""Background service that e-mails CDEK waybill PDFs to customers."""
|
||||
|
||||
import asyncio
|
||||
from collections.abc import Callable, Sequence
|
||||
from contextlib import AbstractAsyncContextManager
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
from typing import Protocol
|
||||
|
||||
import structlog
|
||||
|
||||
logger = structlog.get_logger(__name__)
|
||||
|
||||
|
||||
_EMAIL_SUBJECT_TEMPLATE = "Накладная по заказу {order_uuid}"
|
||||
_EMAIL_BODY_TEMPLATE = (
|
||||
"Здравствуйте!\n\n"
|
||||
"По вашему заказу {order_uuid} сформирована транспортная накладная CDEK.\n"
|
||||
"PDF-файл накладной приложен к этому письму.\n"
|
||||
"Также накладная доступна по ссылке: {waybill_url}\n"
|
||||
)
|
||||
_ATTACHMENT_FILENAME_TEMPLATE = "waybill_{order_uuid}.pdf"
|
||||
|
||||
|
||||
class WaybillPDFDownloaderProtocol(Protocol):
|
||||
async def download_waybill_pdf(self, url: str) -> bytes: ...
|
||||
|
||||
|
||||
class EmailSenderProtocol(Protocol):
|
||||
async def send_email(
|
||||
self,
|
||||
*,
|
||||
to: str,
|
||||
subject: str,
|
||||
body: str,
|
||||
attachment_bytes: bytes,
|
||||
attachment_filename: str,
|
||||
) -> None: ...
|
||||
|
||||
|
||||
class OrderRecord(Protocol):
|
||||
order_uuid: str
|
||||
account_email: str
|
||||
cdek_waybill_url: str | None
|
||||
|
||||
|
||||
class WaybillEmailSenderRepositoryProtocol(Protocol):
|
||||
def session(self) -> AbstractAsyncContextManager[object]: ...
|
||||
|
||||
async def list_orders_pending_waybill_email(
|
||||
self, session: object, *, limit: int
|
||||
) -> Sequence[OrderRecord]: ...
|
||||
|
||||
async def record_waybill_email_sent(
|
||||
self,
|
||||
session: object,
|
||||
*,
|
||||
order_uuid: str,
|
||||
sent_at: datetime,
|
||||
) -> object | None: ...
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SendBatchSummary:
|
||||
processed: int
|
||||
succeeded: int
|
||||
failed: int
|
||||
|
||||
|
||||
class WaybillEmailSenderService:
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
order_repository: WaybillEmailSenderRepositoryProtocol,
|
||||
waybill_downloader: 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._email_sender = email_sender
|
||||
self._batch_size = batch_size
|
||||
self._datetime_now = datetime_now
|
||||
|
||||
async def poll_once(self) -> SendBatchSummary:
|
||||
async with self._repository.session() as session:
|
||||
orders = await self._repository.list_orders_pending_waybill_email(
|
||||
session, limit=self._batch_size
|
||||
)
|
||||
succeeded = 0
|
||||
failed = 0
|
||||
for order in orders:
|
||||
try:
|
||||
await self._handle_order(session, order)
|
||||
succeeded += 1
|
||||
except Exception:
|
||||
failed += 1
|
||||
logger.exception(
|
||||
"waybill_email_order_failed",
|
||||
order_uuid=order.order_uuid,
|
||||
account_email=order.account_email,
|
||||
)
|
||||
return SendBatchSummary(
|
||||
processed=len(orders),
|
||||
succeeded=succeeded,
|
||||
failed=failed,
|
||||
)
|
||||
|
||||
async def run_forever(
|
||||
self,
|
||||
*,
|
||||
interval_seconds: float,
|
||||
stop_event: asyncio.Event,
|
||||
) -> None:
|
||||
while not stop_event.is_set():
|
||||
try:
|
||||
summary = await self.poll_once()
|
||||
logger.info(
|
||||
"waybill_email_tick",
|
||||
processed=summary.processed,
|
||||
succeeded=summary.succeeded,
|
||||
failed=summary.failed,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("waybill_email_tick_failed")
|
||||
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
stop_event.wait(), timeout=interval_seconds
|
||||
)
|
||||
except asyncio.TimeoutError:
|
||||
continue
|
||||
|
||||
async def _handle_order(self, session: object, order: OrderRecord) -> None:
|
||||
waybill_url = order.cdek_waybill_url
|
||||
if waybill_url is None:
|
||||
return
|
||||
|
||||
pdf_bytes = await self._waybill_downloader.download_waybill_pdf(waybill_url)
|
||||
subject = _EMAIL_SUBJECT_TEMPLATE.format(order_uuid=order.order_uuid)
|
||||
body = _EMAIL_BODY_TEMPLATE.format(
|
||||
order_uuid=order.order_uuid,
|
||||
waybill_url=waybill_url,
|
||||
)
|
||||
filename = _ATTACHMENT_FILENAME_TEMPLATE.format(order_uuid=order.order_uuid)
|
||||
|
||||
await self._email_sender.send_email(
|
||||
to=order.account_email,
|
||||
subject=subject,
|
||||
body=body,
|
||||
attachment_bytes=pdf_bytes,
|
||||
attachment_filename=filename,
|
||||
)
|
||||
|
||||
sent_at = self._datetime_now()
|
||||
await self._repository.record_waybill_email_sent(
|
||||
session,
|
||||
order_uuid=order.order_uuid,
|
||||
sent_at=sent_at,
|
||||
)
|
||||
logger.info(
|
||||
"waybill_email_sent",
|
||||
order_uuid=order.order_uuid,
|
||||
account_email=order.account_email,
|
||||
sent_at=sent_at,
|
||||
)
|
||||
@@ -0,0 +1,91 @@
|
||||
"""Background worker that e-mails CDEK waybill PDFs."""
|
||||
|
||||
import asyncio
|
||||
import signal
|
||||
|
||||
import httpx
|
||||
import structlog
|
||||
|
||||
from app.adapters.delivery_providers.cdek.auth import CDEKAuthClient
|
||||
from app.adapters.delivery_providers.cdek.client import CDEKClient
|
||||
from app.adapters.email import SMTPEmailSender
|
||||
from app.adapters.postgres.engine import (
|
||||
create_postgres_engine,
|
||||
create_postgres_session_factory,
|
||||
)
|
||||
from app.config import Settings, get_settings
|
||||
from app.repositories.order import OrderRepository
|
||||
from app.runtime.logging import configure_logging
|
||||
from app.services.waybill_email_sender import WaybillEmailSenderService
|
||||
|
||||
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)
|
||||
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,
|
||||
)
|
||||
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,
|
||||
)
|
||||
email_sender = SMTPEmailSender(
|
||||
smtp_host=settings.email.smtp_host,
|
||||
smtp_port=settings.email.smtp_port,
|
||||
username=settings.email.username,
|
||||
password=settings.email.password,
|
||||
from_address=settings.email.from_address,
|
||||
use_tls=settings.email.use_tls,
|
||||
timeout_seconds=settings.email.timeout_seconds,
|
||||
)
|
||||
engine = create_postgres_engine(settings.postgres)
|
||||
session_factory = create_postgres_session_factory(engine)
|
||||
repository = OrderRepository(session_factory=session_factory)
|
||||
service = WaybillEmailSenderService(
|
||||
order_repository=repository,
|
||||
waybill_downloader=cdek_client,
|
||||
email_sender=email_sender,
|
||||
batch_size=settings.waybill_email_sender.batch_size,
|
||||
)
|
||||
logger.info(
|
||||
"waybill_email_sender_started",
|
||||
interval_seconds=settings.waybill_email_sender.interval_seconds,
|
||||
batch_size=settings.waybill_email_sender.batch_size,
|
||||
)
|
||||
try:
|
||||
await service.run_forever(
|
||||
interval_seconds=settings.waybill_email_sender.interval_seconds,
|
||||
stop_event=stop_event,
|
||||
)
|
||||
finally:
|
||||
await http_client.aclose()
|
||||
await engine.dispose()
|
||||
logger.info("waybill_email_sender_stopped")
|
||||
|
||||
|
||||
async def main() -> None:
|
||||
configure_logging()
|
||||
settings = get_settings()
|
||||
stop_event = asyncio.Event()
|
||||
|
||||
loop = asyncio.get_running_loop()
|
||||
for sig in (signal.SIGTERM, signal.SIGINT):
|
||||
try:
|
||||
loop.add_signal_handler(sig, stop_event.set)
|
||||
except NotImplementedError:
|
||||
pass
|
||||
|
||||
await _run(settings, stop_event)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user