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

1"""Retry worker for deferred mail deliveries. 

2 

3Executes what :class:`~lexigram.notification.delivery.retry.RetryingMailer` 

4schedules: re-sends deliveries whose backoff window has elapsed and records 

5outcomes via the store. 

6 

7Wire it to the task executor:: 

8 

9 executor.register_handler( 

10 "notification.flush_retries", 

11 lambda: flush_retries(store, backend), 

12 ) 

13""" 

14 

15from __future__ import annotations 

16 

17from typing import Any 

18 

19from lexigram.contracts.mailer import EmailMessage 

20from lexigram.logging import get_logger 

21 

22logger = get_logger(__name__) 

23 

24_MAX_ATTEMPTS = 5 

25"""Abandon a delivery after this many attempts.""" 

26 

27 

28async def flush_retries(store: Any, backend: Any, limit: int = 50) -> int: 

29 """Re-send deferred deliveries whose backoff window has elapsed. 

30 

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. 

36 

37 Returns: 

38 The number of deliveries completed during this flush. 

39 """ 

40 

41 due = await store.due_deliveries(limit) 

42 delivered = 0 

43 

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 

52 

53 message = EmailMessage(to=recipient.split(","), subject=subject, body=body) 

54 result = await backend.send(message) 

55 

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 

61 

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 ) 

82 

83 return delivered 

84 

85 

86__all__ = ["flush_retries"]