"""In-memory delivery acknowledgements for Reportable messages.
The send-to-server plugins queue messages rather than deliver them, so a
producer that keeps its own copy of the payload cannot tell a delivered
message from one a lost send round dropped. Every send path reports the ids
the transport accepted here, and producers register the ids they care about.
Acknowledgements are deliberately not persisted: one that never arrives
leaves the message unconfirmed and makes the producer send it again. That
costs a duplicate the server may well store twice, whereas a silently
dropped payload cannot be recovered at all.
"""
import logging
from typing import Callable, Optional
logger = logging.getLogger(__name__)
class DeliveryAckRegistry:
def __init__(self) -> None:
self._callbacks: dict[str, Callable[[], None]] = {}
def watch(self, message_id: str, on_delivered: Callable[[], None]) -> None:
if not message_id:
return
self._callbacks[message_id] = on_delivered
def unwatch(self, message_id: str) -> None:
self._callbacks.pop(message_id, None)
def confirm(self, message_id: Optional[str]) -> None:
on_delivered = self._callbacks.pop(message_id, None)
if on_delivered is None:
return
try:
on_delivered()
except Exception:
# a producer's bookkeeping must never break a send round, but it
# failing means the producer will re-send forever: log the
# traceback, this is the only place that sees it
logger.exception(
"Delivery acknowledgement for %s failed", message_id
)
registry = DeliveryAckRegistry()