"""Main Handler.""" from bulk_metadata_ingester_common.constants.catalog_ingestion import FAILURE from bulk_metadata_ingester_common.utils.jitter import jitter from bulk_metadata_ingester_common.utils.logging import get_current_logger import config from config import graphql_gateway from constants.queries import UPDATE_DDEX_INGEST from ddex_ingester_common.helpers.catalog_ingestion import ( CatalogIngestion ) from ddex_ingester_common.helpers.rds import run_rds_query def handler(event, context): """Lambda Entrypoint.""" if 'write_catalog_ingestion' not in event: return event correlation_id = event.get('correlation_id') status = event.get('status', FAILURE) logger = get_current_logger( config.ENVIRONMENT, config.LAMBDA_NAME, logging_level=config.LOGGING_LEVEL, correlation_id=correlation_id) # Jitter calls - Sleep for randomness jitter(logger) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id } ) # Log lambda begins message logger.info(f'write_catalog_ingestion received: {event}') error_message = None if status == FAILURE and event.get('context').get('errors'): error_message = str(event.get('context').get('errors')) write_catalog_ingestion(event, status, error_message, logger) update_data_ingest_status( event.get('ingest_id'), 'ingest_failed' if status == FAILURE else 'ingest_succeeded', logger ) return event def write_catalog_ingestion( event, status, error_message, logger): """Write catalog ingestion to RDS.""" context = event.get('context') key = context.get('key') bucket = context.get('bucket') execution_name = context.get('execution_name') state_machine_name = context.get('state_machine_name') ingest_format = context.get('orig_key').split('.')[-1].lower() config.catalog_ingestion_session.add( CatalogIngestion( state_machine_name=state_machine_name, state_machine_execution_name=execution_name, catalog_ingestion_source_id='bulk-metadata-ingester', s3_bucket_name=bucket, s3_key_name=key, ingest_format=ingest_format, timestamp=context.get('execution_start_time'), vendor_id=context.get('product').get('vendor_id'), subaccount_id=context.get('product').get('subaccount_id'), 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_data_ingest_status(ingest_id, status, logger): """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, UPDATE_DDEX_INGEST, query_args, )