"""Lambda send_success function module. Largely copy/pasted from lambda-gda-account-creation.""" import json from typing import Any from confluent_kafka import Message, Producer from config import ARTIST_DIY_KAFKA_TOPIC, BROKER_SERVERS, logger def delivery_report(err: Exception, msg: Message) -> None: """Log delivery results from callbacks for each message produced, triggered by flush().""" if err is not None: logger.exception(f'Message delivery failed: {err}') else: logger.info(f'Message delivered to {msg.topic()} [{msg.partition()}]') def handler(event: Any, context: Any) -> dict[str, str]: """Lambda entry point.""" try: logger.info('Starting to write success event to Kafka topic.') producer = Producer({'bootstrap.servers': BROKER_SERVERS, 'security.protocol': 'SSL'}) # Asynchronously produce a message, the delivery report callback # will be triggered from poll() above, or flush() below, when the message has # been successfully delivered or failed permanently. producer.produce( ARTIST_DIY_KAFKA_TOPIC, json.dumps(event).encode('utf-8'), callback=delivery_report ) # Wait for any outstanding messages to be delivered and delivery report callback. producer.flush() logger.info(f'Successfully wrote success message to Kafka topic: {ARTIST_DIY_KAFKA_TOPIC}') return {'status': 'success'} except Exception as err: logger.exception(str(err)) raise err