"""Lambda entry point module.""" import logging from typing import Any from typing import Mapping from typing import Optional from lambdacommon.common_config import logger import sentry_sdk from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk.integrations.logging import LoggingIntegration import config from src.logic import dlq_event from src.logic import source_event from src.logic import target_event from src.models.preference_change import InvalidPreferenceMessage logging_integration = LoggingIntegration( level=logging.INFO, event_level=logging.CRITICAL ) sentry_sdk.init( config.settings.sentry_dsn, integrations=[AwsLambdaIntegration(), logging_integration] ) def handler( event: Optional[Mapping[str, Any]], context: Optional[Any], ) -> Mapping[str, str]: """Lambda entry point for processing fan subscription preference updates. Receives MSK events containing fan subscription preference changes, processes them, enriches with Snowflake data, and sends formatted messages to Kafka for Salesforce. Args: event: AWS MSK event payload containing subscription preference changes context: AWS Lambda context object (unused) Returns: Dictionary with status indicator Raises: Exception: Re-raises any exception encountered during processing """ invalid_messages: list[InvalidPreferenceMessage] = [] try: valid_messages, parsed_invalid_messages = source_event.parse(event) invalid_messages = parsed_invalid_messages updated_subscriptions = source_event.flatten_events(valid_messages) target_messages = target_event.prepare_messages(updated_subscriptions) target_event.send_messages(target_messages) return {'status': 'OK'} except Exception as e: invalid_messages.append( InvalidPreferenceMessage( message_key=None, message_value=str(event), error_type=e.__class__.__name__, error_message=str(e), ) ) logger.exception(str(e)) raise e finally: dlq_event.send_invalid_messages(invalid_messages)