"""Utility helpers for preparing and sending target Kafka messages.""" from typing import Any from typing import Literal from typing import Optional from typing import cast 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 from src.logic import subscription_info from src.models.preference_change import SalesforcePreferenceMessage from src.models.preference_change import UpdatedSubscription def prepare_messages( batch_items: dict[str, UpdatedSubscription], ) -> list[tuple[str, SalesforcePreferenceMessage]]: """Prepare Salesforce messages from batch items. Fetches additional data from Snowflake and maps to SalesforcePreferenceMessage format. Args: batch_items: Dictionary mapping subscription IDs to UpdatedSubscription objects Returns: List of (message_key, message_value) tuples for Kafka publishing, where: - message_key: Salesforce Subscription ID - message_value: SalesforcePreferenceMessage instance """ result = [] if not batch_items: logger.info('No batch items to process, returning empty list') return [] logger.info(f'Preparing messages for {len(batch_items)} subscription updates') # Extract subscription IDs for Snowflake query subscription_ids = list(batch_items.keys()) logger.debug(f'Fetching subscription data from Snowflake for IDs: {subscription_ids}') # Fetch enriched data from Snowflake snowflake_rows = subscription_info.get_subscriptions(subscription_ids) logger.info(f'Retrieved {len(snowflake_rows)} rows from Snowflake for {len(subscription_ids)} subscription IDs') for row in snowflake_rows: subscription_id = str(row['FAN_SUBSCRIPTION_ID']) action = cast(Literal['Sub', 'Unsub'], 'Sub' if batch_items[subscription_id].new_value else 'Unsub') result.append( ( str(row['SUBSCRIPTION_ID__C']), SalesforcePreferenceMessage( Subscription_ID__c=str(row['SUBSCRIPTION_ID__C']), Action__c=action, Mailing_List_ID__c=str(row['MAILING_LIST_ID__C']), Source__c='New Preference Center', Fan_ID__c=str(row['FAN_ID__C']), ) ) ) logger.info(f'Prepared {len(result)} messages') return result def send_messages(messages: list[tuple[str, SalesforcePreferenceMessage]]) -> None: """Send messages to Kafka topic. Publishes formatted preference update messages to the configured Kafka topic for consumption by Salesforce integration. Args: messages: List of (message_key, message_value) tuples to publish Raises: Exception: If message delivery fails (raised in callback) """ logger.info(f'Starting to send {len(messages)} messages to Kafka topic: {config.settings.kafka_target_topic}') if not messages: logger.info('No messages to send, skipping Kafka producer initialization') return event_producer_params = dict( bootstrap_servers=config.settings.kafka_bootstrap_servers, key_serializer=StringSerializer(), value_serializer=SimpleJSONSerializer(), security_protocol=config.settings.kafka_security_protocol ) with EventProducer(**event_producer_params) as producer: for message_key, message_value in messages: producer.produce( topic=config.settings.kafka_target_topic, event_key=message_key, event_value=message_value.model_dump(by_alias=True, exclude_none=True), headers={config.settings.kafka_message_header_key: ''.encode()}, auto_flush=False, callback=msg_delivery_callback ) logger.info(f'Successfully sent {len(messages)} messages to Kafka') def msg_delivery_callback(err: Optional[Exception], msg: Any) -> None: """Handle Kafka message delivery result. Callback invoked after Kafka producer attempts to deliver a message. Logs errors and raises exceptions for failed deliveries. Args: err: Exception if delivery failed, None if successful msg: Kafka message object containing key and metadata Raises: Exception: Re-raises the error if message delivery failed """ if err: logger.error(f'There was an error producing message for {msg.key()} to Kafka') raise err