"""Lambda fan-vendor-weekly-deletions-emails function module.""" import json import logging from typing import DefaultDict from typing import Dict from typing import List import uuid from lambdacommon.common_config import logger import sentry_sdk from sentry_sdk import capture_exception from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from sentry_sdk.integrations.logging import LoggingIntegration import config from src import constants from src.connectors.ses import send_email_with_csv from src.connectors.snowflake import execute_snowflake_query from src.queries import snowflake_queries logging_integration = LoggingIntegration( level=logging.INFO, # Capture info and above as breadcrumbs event_level=logging.CRITICAL # Send only critical as events ) sentry_sdk.init( config.SENTRY_DSN, integrations=[AwsLambdaIntegration(), logging_integration] ) def handler(event, context): """Lambda entry point.""" try: logger.info('Starting lambda fan-vendor-weekly-deletions-emails.') return main_handler(event, context) except Exception as err: logger.error(f'Uncaught exception: {err}') capture_exception(err) raise def main_handler(event, context): """Lambda main logic.""" task_run_id = str(uuid.uuid4()) logger.info(f'Starting vendor deletions export. Task Run ID: {task_run_id}') execute_snowflake_query( snowflake_queries.CREATE_VENDOR_DELETIONS_EXPORT, {'platforms': json.dumps(constants.VENDOR_DELETION_PLATFORMS)}) logger.info('Created VENDOR_DELETIONS_EXPORT table in Snowflake') result = execute_snowflake_query( snowflake_queries.GET_VENDOR_DELETIONS, {}) logger.info('Fetched vendor deletion requests from Snowflake') vendor_emails = group_emails_by_vendor(result) logger.info('Grouped emails by vendor') # Send emails to vendors with CSV attachments for vendor_name, emails_list in vendor_emails.items(): if emails_list: vendor_info = constants.VENDOR_MAPPING.get(vendor_name) if vendor_info: send_email_with_csv( vendor_name=vendor_name, emails_list=emails_list, to_addresses=vendor_info['to_addresses'], cc_addresses=config.CC_ADDRESSES, sender_email=config.SENDER_EMAIL, sender_display_name=constants.SENDER_DISPLAY_NAME ) logger.info(f'Sent email to {vendor_name} with {len(emails_list)} email(s)') else: logger.warning(f'No vendor mapping found for: {vendor_name}') execute_snowflake_query( snowflake_queries.MARK_VENDOR_DELETIONS_AS_EXPORTED, {'task_run_id': task_run_id}) logger.info('Marked vendor deletion requests as exported in Snowflake') # commet out cleanup for debugging purposes # execute_snowflake_query( # snowflake_queries.CLEANUP_VENDOR_DELETIONS_EXPORT, {}) return { 'statusCode': 200, 'body': json.dumps('Completed vendor deletions export process.') } def group_emails_by_vendor(result) -> Dict[str, List[str]]: """Group emails by vendor based on their platforms. Args: result: Iterable of rows with 'EMAIL' and 'PLATFORMS' columns. Returns: Dict mapping vendor names to lists of emails. """ vendor_emails = DefaultDict(list) for row in result: email = row['EMAIL'] platforms = json.loads(row['PLATFORMS']) # Find matching vendor(s) for this email's platforms for vendor, vendor_info in constants.VENDOR_MAPPING.items(): # Check if any of the email's platforms match this vendor's platforms if any(platform in vendor_info['platforms'] for platform in platforms): vendor_emails[vendor].append(email) return vendor_emails