"""Lambda sr-record-event function module.""" import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration import config import json from datetime import datetime from datetime import timezone from src.common.connectors import kafka_producer from src.common import logger from src.constants import delivery_event # initialize sentry sentry_dsn = config.secrets_manager_client.get_cred('SENTRY_DSN') if sentry_dsn: logger.info('Initializing with sentry') sentry_sdk.init( sentry_dsn, integrations=[AwsLambdaIntegration()] ) else: logger.info('Initializing without sentry') def _determine_execution_type(event): if 'delivery_type' in event and event['delivery_type'] is not None: # noqa:E501 return event['delivery_type'].upper() if 'execution_type' in event and event['execution_type'] is not None: return event['execution_type'] def _determine_message(event): status = 'ok' details = event.get('details') if 'error' in event: status = 'error' elif details and details.get('undelivered', False): status = 'undelivered' details = event.get('error') if status == 'error' else details return { 'status': status, 'details': details } def _format_kafka_message(event): details = event.get('details') return { 'sound_recording_id': event['sound_recording']['id'], 'sound_recording_version': event['sound_recording']['version'], 'step_function_execution_id': event['execution_name'], 'service': ( event['service'] if 'service' in event and event['service'] else delivery_event.SERVICE_TIKTOK ), 'execution_type': _determine_execution_type(event), 'timestamp': datetime.now(timezone.utc).isoformat(), 'message': _determine_message(event), 'event_type': ( delivery_event.UNDELIVERED_EVENT_TYPE if details and details.get('undelivered', False) else event['event_type'] ) } def handler(event, context): """Lambda entry point.""" try: data = _format_kafka_message(event) kafka_producer.produce_message( json.dumps(data), config.KAFKA_TOPIC) return event except Exception as e: logger.exception(str(e)) raise e