"""Lambda function module.""" import json import os import uuid import boto3 from ddex_ingester_common.constants.ddex_providers import ( RISING_88, SME_ANALYTICS_PROVIDER) from ddex_ingester_common.lambda_exceptions import ParseDDEXException from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.release_correction.release_corrections_s3 import ( create_blank_rc_json, delete_all_rc_json_files) from ddex_ingester_common.schemas.ddex_schema import (ParticipantSchema, ProjectSchema) from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.schemas.state_machine_schema import \ StateMachineSchema from ddex_ingester_common.utils.read_ddex import read_ddex import config from constants.rising_88 import (RISING_88_GRID_TO_SUBACCOUNT_ID, RISING_88_VENDOR_ID) logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Parse and deserialize DDEX into a marshmallow schema.""" correlation_id = event.get('correlation_id') or str(uuid.uuid4()) logging_utils.update_logger_correlation_id(logger, correlation_id) logger.info(f'Triggered parse_ddex: {event}') bucket = event.get('bucket') key = event.get('key') execution_name = event.get('execution_name') state_machine_name = event.get('state_machine_name') execution_start_time = event.get('execution_start_time') ddex_provider = event.get('ddex_provider') artwork_ingestion_only = event.get('artwork_ingestion_only') or False ingest_id = event.get('ingest_id') check_key_extension(key, logger) s3_client = boto3.client('s3') data = read_ddex(s3_client, bucket, key, logger) logging_utils.update_logger_with_message_ids( logger, data.message_id, data.message_thread_id, execution_name ) # Save bucket name and key into context data.bucket = bucket data.key = f"{key.rsplit('/', 1)[0]}/" if not data.is_purged_release: format_asset_s3_location(data) # Save state machine execution name into context data.execution_name = execution_name data.state_machine_name = state_machine_name data.execution_start_time = execution_start_time # Save DDEX provider into context data.ddex_provider = ddex_provider # Handle ingesting 88 Rising DDEX with the AWAL sfn if ddex_provider == RISING_88: grid = data.product.grid if grid not in RISING_88_GRID_TO_SUBACCOUNT_ID: raise Exception(f'No subaccount mapping found for GRID: {grid}') data.product.vendor_id = RISING_88_VENDOR_ID data.product.subaccount_id = RISING_88_GRID_TO_SUBACCOUNT_ID[grid] if not data.product.catalog_number: data.product.catalog_number = grid if not data.project: data.project = ProjectSchema() data.project.artist = ParticipantSchema() data.project.name = data.product.product_name data.project.project_code = grid data.project.artist.name = data.product.display_artist_name if data.video and not data.video.product_code: data.video.product_code = grid # save correlation_id to sfn context data.correlation_id = correlation_id # pass on ingest_id data.ingest_id = ingest_id if not data.is_purged_release: # Write json data to s3 json_file_path = f'{data.key}parsed_ddex.json' s3_client.put_object( Bucket=bucket, Key=json_file_path, Body=json.dumps(S3Schema().dump(data)).encode(encoding='UTF-8') ) # Remove existing Release correction JSON files first. delete_all_rc_json_files(event) # Creating a release_correction.json file to be used in subsequent lambdas create_blank_rc_json(event) # Default switchboard deal to False until one is found data.has_switchboard_deal = False data.artwork_ingestion_only = artwork_ingestion_only # appropriate track list will be added to StateMachineSchema in prep_ddex if ddex_provider == SME_ANALYTICS_PROVIDER: data.tracks = [] # Dump parsed DDEX information into context object logger.info('Dumping context') return StateMachineSchema().dump(data) def check_key_extension(key, logger): """Check file extension is valid xml.""" _, extension = os.path.splitext(key) if extension != '.xml': msg = 'S3 event source is not an XML file' logger.warning(msg) raise ParseDDEXException(msg) def format_asset_s3_location(context): """Format assets with s3 filepaths.""" artwork = context.product.artwork tracks = context.tracks video = context.video if artwork: artwork.bucket = context.bucket artwork.key = f'{context.key}{artwork.filepath}{artwork.filename}' for track in tracks: if track.asset: track.asset.bucket = context.bucket track.asset.key = (f'{context.key}{track.asset.filepath}' f'{track.asset.filename}') if video and video.assets: for asset in video.assets: asset.bucket = context.bucket asset.key = f'{context.key}{asset.filepath}{asset.filename}' if len(video.assets) > 1: assets = [] for aspect_ratio in config.PREFERRED_ASPECT_RATIOS: for asset in video.assets: if asset.aspect_ratio == aspect_ratio: assets.append(asset) if assets: for asset_list in assets: if asset_list: context.video.assets = [asset_list] break else: context.video.assets = [video.assets[0]]