Coverage for src/lexigram/notification/delivery/worker.py: 69%
36 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-26 02:31 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-26 02:31 +0800
1"""Retry worker for deferred mail deliveries.
3Executes what :class:`~lexigram.notification.delivery.retry.RetryingMailer`
4schedules: re-sends deliveries whose backoff window has elapsed and records
5outcomes via the store.
7Wire it to the task executor::
9 executor.register_handler(
10 "notification.flush_retries",
11 lambda: flush_retries(store, backend),
12 )
13"""
15from __future__ import annotations
17from typing import Any
19from lexigram.contracts.mailer import EmailMessage
20from lexigram.logging import get_logger
22logger = get_logger(__name__)
24_MAX_ATTEMPTS = 5
25"""Abandon a delivery after this many attempts."""
28async def flush_retries(store: Any, backend: Any, limit: int = 50) -> int:
29 """Re-send deferred deliveries whose backoff window has elapsed.
31 Args:
32 store: Delivery-state store exposing ``due_deliveries`` plus the
33 ``DeliveryStoreProtocol`` mutation methods.
34 backend: Raw mail backend used for the re-send attempts.
35 limit: Maximum deliveries processed per flush.
37 Returns:
38 The number of deliveries completed during this flush.
39 """
41 due = await store.due_deliveries(limit)
42 delivered = 0
44 for entry in due:
45 message_payload = entry.get("message") or {}
46 recipient = entry.get("recipient", "")
47 subject = str(message_payload.get("subject", ""))
48 body = str(message_payload.get("body", ""))
49 if not recipient or not subject:
50 await store.mark_failed(entry["delivery_id"], "empty payload")
51 continue
53 message = EmailMessage(to=recipient.split(","), subject=subject, body=body)
54 result = await backend.send(message)
56 if result.is_ok():
57 await store.mark_delivered(entry["delivery_id"])
58 delivered += 1
59 logger.info("retry_mail_delivered", delivery_id=entry["delivery_id"])
60 continue
62 attempts = int(entry.get("attempts", 0)) + 1
63 error = str(result.unwrap_err())
64 if attempts >= _MAX_ATTEMPTS:
65 await store.mark_failed(entry["delivery_id"], reason=error)
66 logger.error(
67 "retry_mail_abandoned",
68 delivery_id=entry["delivery_id"],
69 attempts=attempts,
70 error=error,
71 )
72 else:
73 delay = min(60 * (2 ** max(0, attempts - 1)), 3600)
74 await store.increment_retry(entry["delivery_id"])
75 await store.schedule_retry(entry["delivery_id"], delay)
76 logger.warning(
77 "retry_mail_scheduled",
78 delivery_id=entry["delivery_id"],
79 attempts=attempts,
80 delay_seconds=delay,
81 )
83 return delivered
86__all__ = ["flush_retries"]