"""AWS Lambda handler for adjustment file preparation triggered by EventBridge. Validates and prepares adjustment files using DuckDB for in-memory processing, loads validated data to MySQL staging, and publishes results to EventBridge for downstream processing. Flow: EventBridge → Lambda → DuckDB (validation) → MySQL (staging) → EventBridge (results) Error Handling: - TransientError: Triggers Lambda retry (DB connection, S3 throttling) - PermanentError: Halts processing, updates batch status to ERROR - Transaction rollback and DuckDB cleanup on all errors """ from __future__ import annotations import os from typing import Any, Mapping import sentry_sdk from pydantic import ValidationError from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration from config import config from src.connectors.duckdb import DuckDBConfig, DuckDBConnectionFactory from src.connectors.duckdb.utils import ( add_s3_secret, add_snowflake_secret, ) from src.connectors.mysql import MySQLConnectionFactory from src.connectors.snowflake import SnowflakeConnectionFactory from src.errors import ( PermanentError, TransientError, ) from src.gateways import SnowflakeGateway from src.infra.log import logger from src.infra.resources import ResourceManager from src.repositories.royalty_accounting import RoyaltyAccountingClient from src.schemas import AdjustmentFilePrepareEvent from src.services.file_loader import AdjustmentFileLoader from src.services.processor import AdjustmentFilePrepareProcessor from src.services.reference_loader import ReferenceLoader from src.services.s3_downloader import S3Downloader from src.services.validator import DuckDBValidator from src.sql import DuckDBQuery, load_sql from src.utils.file_utils import create_temp_file if config.sentry_dsn: sentry_sdk.init( dsn=config.sentry_dsn, environment=config.env, integrations=[AwsLambdaIntegration(timeout_warning=True)], ) def handler(event: Mapping[str, Any] | None, context: Any) -> dict[str, Any]: """Lambda handler for adjustment file validation and preparation. Downloads file from S3, validates via DuckDB, stages to MySQL, and returns results. Args: event: EventBridge event with batch_id, s3_bucket, s3_key, correlation_id. context: Lambda runtime context (ARN, memory limits). Returns: Response dict with validation results (row counts, amounts) and prepared file location. Raises: TransientError: Retriable errors triggering Lambda retry. PermanentError: Non-retriable errors halting processing. ValidationError: Invalid event structure. """ duck_path = create_temp_file(config.duckdb.TEMP_DIR, suffix='.duckdb') try: logger.info(f'Function ARN: {context.invoked_function_arn}') validated_event = AdjustmentFilePrepareEvent(**(event or {})) batch_id = validated_event.detail.metadata.target_id correlation_id = validated_event.detail.metadata.correlation_id logger.info(f'Correlation ID: {correlation_id}') logger.info(f'Batch ID: {batch_id}') logger.info('Initializing') duck_config = DuckDBConfig( file_path=duck_path, max_memory_mb=int(config.duckdb.MEM_PCT * int(context.memory_limit_in_mb)), ) duck_factory = DuckDBConnectionFactory(duck_config) dst_factory = MySQLConnectionFactory(config.mysql) snow_factory = SnowflakeConnectionFactory(config.snowflake) with ( dst_factory.connection() as destination_conn, duck_factory.connection() as duck_conn, snow_factory.connection() as snow_conn, ): try: with duck_conn.cursor() as cursor: add_s3_secret(cursor, config.aws_region) snowflake_secret = add_snowflake_secret(cursor, config.snowflake) cursor.execute(load_sql(DuckDBQuery.RegisterMacros)) adjustment_file_loader = AdjustmentFileLoader( duck_conn, max_rows=config.policy.MAX_FILE_ROWS, max_str_len=config.policy.MAX_COMMENT_LEN, ) validator = DuckDBValidator(duck_conn) gateway = SnowflakeGateway( duck_conn, snow_conn, snowflake_secret, config.env ) data_service = ReferenceLoader(duck_conn, gateway) ra_client = RoyaltyAccountingClient(destination_conn) s3_file_downloader = S3Downloader( ResourceManager.get_s3_connection(), config.policy.MAX_FILE_BYTES ) logger.info('Processing') processor = AdjustmentFilePrepareProcessor( adjustment_file_loader=adjustment_file_loader, adjustment_file_validator=validator, reference_data_service=data_service, duck_conn=duck_conn, royalty_accounting_client=ra_client, s3_file_downloader=s3_file_downloader, ) result = processor.process(validated_event) destination_conn.commit() # Return result return result.model_dump() # type: ignore except Exception as e: try: destination_conn.rollback() except Exception as db_err: logger.error(f'Rollback failed: {db_err}') raise e from db_err else: raise e except ValidationError as e: # Invalid event structure is a permanent error logger.error(f'Invalid event structure: {e}') raise PermanentError(str(e)) from e except TransientError as e: # Retriable errors (DB connection, IntegrityError, etc) logger.warning(f'Transient error encountered: {e}') raise TransientError(str(e)) from e except PermanentError as e: # Non-retriable errors (validation failures, not found, etc) logger.error(f'Permanent error encountered: {e}') raise PermanentError(str(e)) from e except Exception as e: # Unexpected errors logger.exception(f'Unexpected error encountered: {e}') raise finally: # Cleanup DuckDB try: logger.info(f'Deleting DuckDB: {duck_path}') os.remove(duck_path) except FileNotFoundError: pass except OSError as e: logger.error(f'Failed to cleanup file {duck_path}: {e}')