Files
g2s-aggregator/app/workers/waybill_poller.py
T
Раис Юсупалиев 843175f12e
Deploy / deploy (push) Successful in 18m5s
Рефактор
2026-06-26 20:16:04 +03:00

83 lines
2.6 KiB
Python

"""Background worker that polls CDEK for waybill updates."""
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.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_poller import WaybillPollerService
logger = structlog.get_logger(__name__)
async def _run(settings: Settings, stop_event: asyncio.Event) -> None:
cdek_config = settings.delivery_providers.cdek
http_client = httpx.AsyncClient(timeout=cdek_config.timeout_seconds)
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,
)
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,
batch_size=settings.waybill_poller.batch_size,
)
logger.info(
"waybill_poller_started",
interval_seconds=settings.waybill_poller.interval_seconds,
batch_size=settings.waybill_poller.batch_size,
)
try:
await service.run_forever(
interval_seconds=settings.waybill_poller.interval_seconds,
stop_event=stop_event,
)
finally:
await http_client.aclose()
await engine.dispose()
logger.info("waybill_poller_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())