"""Lambda earnings_transfer function module.""" from __future__ import annotations from datetime import datetime from typing import Any, Mapping import sentry_sdk from lambdacommon.common_config import logger from sentry_sdk.integrations.aws_lambda import AwsLambdaIntegration import config from src.connectors.ows_client import get_ows_client from src.connectors.snowflake import SnowflakeConnectionFactory, build_snowflake_config from src.errors import InputValidationError from src.parsers.earnings_transfer_parser import parse_earnings_transfers from src.processor import EarningsTransferProcessor if config.SENTRY_DSN: sentry_sdk.init( dsn=config.SENTRY_DSN, environment=config.ENVIRONMENT, integrations=[AwsLambdaIntegration(timeout_warning=True)], ) def handler(event: Mapping[str, Any] | None, context: Any) -> dict[str, Any]: """Lambda entry point. Invoked manually. Loads earnings transfer configurations from ows-royalties, enriches them with Snowflake data, and runs the calculation pipeline. """ logger.info( 'Function ARN', extra={'arn': getattr(context, 'invoked_function_arn', 'unknown')}, ) try: snow_factory = SnowflakeConnectionFactory(build_snowflake_config()) with snow_factory.connection() as conn: processor = EarningsTransferProcessor(config.S3_BUCKET_NAME) start_time = datetime.now() records = parse_earnings_transfers(get_ows_client(), conn) result = processor.process_records(records) elapsed = datetime.now() - start_time logger.info('Finished processing', extra={'elapsed': str(elapsed)}) except InputValidationError as e: # Permanent failure — do not retry logger.error('Validation error: %s', e) return {'status': 'ERROR', 'error': str(e)} # TransientError and other exceptions propagate for Lambda retry return {'status': 'OK', **result.model_dump()}