"""Lambda function module.""" from typing import Dict import uuid import config from ddex_ingester_common.constants.catalog_ingestion import ( DDEX_FORMAT, DDEX_PROVIDER_TO_SOURCE_ID, FAILURE, SKIP, SUCCESS ) from ddex_ingester_common.constants.ddex_providers import \ SME_ANALYTICS_PROVIDER from ddex_ingester_common.helpers.catalog_ingestion import ( CatalogIngestion ) from ddex_ingester_common.helpers.rds import run_rds_query from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.models.state_machine.body import ( Body as StateMachineContext) from ddex_ingester_common.schemas.state_machine_schema import ( StateMachineSchema ) import queries logger = logging_utils.get_logger(config.app_logger) def handler(event: Dict, context: object) -> Dict: """Lambda entrypoint.""" # If the execution fails on parse_ddex we have no context data # Without context we don't have the required data to write to the table if not event['context'].get('key') and not event['context'].get('bucket'): logger.info('Context not found in event payload.') return event status = event.get('status', FAILURE) context = event['context'] logger.info( f'UPC: {context.get("product", {}).get("upc", "undefined")}, ' f'DDEX_PROVIDER: {context.get("ddex_provider", "undefined")}, ' f'STATUS: {status}') state_machine_data = StateMachineSchema().load(event.get('context')) correlation_id = state_machine_data.correlation_id or str(uuid.uuid4()) state_machine_data.correlation_id = correlation_id logging_utils.update_logger_correlation_id(logger, correlation_id) logging_utils.update_logger_with_message_ids( logger, state_machine_data.message_id, state_machine_data.message_thread_id, state_machine_data.execution_name ) logger.info(f'Writing catalog ingestion data for event: {event}') error_message = None if status == FAILURE and state_machine_data.errors: error_message = str(state_machine_data.errors) write_catalog_ingestion(state_machine_data, status, error_message) db_status = 'ingest_succeeded' if status == FAILURE: db_status = 'ingest_failed' elif status == SKIP: db_status = 'skipped' if context.get('ddex_provider') == SME_ANALYTICS_PROVIDER and \ context.get('is_purged_release'): if status == SUCCESS: db_status = 'soft_deleted' elif status == FAILURE: db_status = 'soft_delete_failed' update_ddex_ingest_status( state_machine_data.ingest_id, db_status ) return StateMachineSchema().dump(state_machine_data) def write_catalog_ingestion( state_machine_data: StateMachineContext, status: str, error_message: str): """Write catalog ingestion to RDS.""" config.catalog_ingestion_session.add( CatalogIngestion( state_machine_name=state_machine_data.state_machine_name, state_machine_execution_name=state_machine_data.execution_name, catalog_ingestion_source_id=DDEX_PROVIDER_TO_SOURCE_ID.get( state_machine_data.ddex_provider), s3_bucket_name=state_machine_data.bucket, s3_key_name=state_machine_data.key, ingest_format=DDEX_FORMAT, timestamp=state_machine_data.execution_start_time, vendor_id=( state_machine_data.product.vendor_id if state_machine_data.product else None), subaccount_id=( state_machine_data.product.subaccount_id if state_machine_data.product else None), status=status, error_message=error_message, ) ) logger.info( 'Saving Catalog Ingestion Session with data: ' f'{config.catalog_ingestion_session.data}') config.catalog_ingestion_session.save() def update_ddex_ingest_status(ingest_id, status): """Update DDEX ingest status. Args: ingest_id (int): Ingest ID status (str): Status """ if not ingest_id or not status: return query_args = (status, ingest_id) run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_USER, config.RDS_PASSWORD, queries.UPDATE_DDEX_INGEST, query_args, )