"""Kafka Producer for Triggered Send Events.""" from typing import Optional from confluent_kafka import KafkaError from confluent_kafka import Message from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer from lambdacommon.common_config import logger import config # noqa from src.models import TriggeredSendFailedEvent from src.models import TriggeredSendSuccessfulEvent class TriggeredSendEventProducer: """Singleton wrapper for Kafka EventProducer.""" _instance = None def __new__(cls, *args, **kwargs): """Create a new instance of TriggeredSendEventProducer.""" if not cls._instance: cls._instance = super(TriggeredSendEventProducer, cls).__new__(cls) cls._instance._initialize(*args, **kwargs) return cls._instance def _initialize(self) -> None: """Initialize the EventProducer instance. Sets up the Kafka producer with appropriate serializers and connection configuration from the app config. """ self.producer = EventProducer( bootstrap_servers=config.KAFKA_BOOTSTRAP_SERVERS, key_serializer=StringSerializer(), value_serializer=SimpleJSONSerializer(), security_protocol=config.KAFKA_SECURITY_PROTOCOL ) def produce_successful_event(self, event: TriggeredSendSuccessfulEvent) -> None: """Produce a successful triggered send event to the target Kafka topic. Args: event (TriggeredSendSuccessfulEvent): The successful event to be published. """ self.producer.produce( topic=config.SUCCESSFUL_TRIGGERED_SENDS_TOPIC, event_key=None, event_value=event.model_dump(), callback=_msg_delivery_callback ) def produce_failed_event(self, event: TriggeredSendFailedEvent) -> None: """Produce a failed triggered send event to the target Kafka topic. Args: event (TriggeredSendFailedEvent): The failed event to be published. """ self.producer.produce( topic=config.FAILED_TRIGGERED_SENDS_TOPIC, event_key=None, event_value=event.model_dump(), callback=_msg_delivery_callback ) def _msg_delivery_callback(err: Optional[KafkaError], msg: Message) -> None: """Handle Kafka message delivery result. This callback is invoked by the Kafka producer when a message has been successfully delivered or when delivery failed. Args: err (Optional[KafkaError]): The error that occurred during message delivery, if any. msg (Message): The message that was delivered or failed. Raises: Exception: If there was an error producing the message to Kafka. """ if err: logger.error(f'Error producing message to Kafka topic {msg.topic()}: {err}') raise Exception(f'Kafka producer error: {err}') else: logger.debug(f'Message delivered to {msg.topic()} [{msg.partition()}] at offset {msg.offset()}')