"""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, 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__) @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=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=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=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, 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_downloaders=waybill_downloaders, 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())