"""Lambda function module.""" import json import os import uuid from ddex_ingester_common.helpers.asset import copy_asset, create_asset_token from ddex_ingester_common.helpers.s3_ddex import load_ddex_json from ddex_ingester_common.lambda_exceptions import (AudioException, LambdaException) from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.release_correction.release_correction import \ is_release_correction from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.schemas.state_machine_schema import ( StateMachineSchema, TrackSchema) import config from config import graphql_gateway from constants import queries logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Lambda entrypoint.""" logger.info(f'Triggered ddex-ingester-handle-audio-asset: {event}') try: return handle_track_asset(event) except Exception as e: logger.error(f'handle_track_asset encountered an error: {e}') raise def handle_track_asset(event): """Handle uploading track audio asset.""" s3_data = S3Schema().load(load_ddex_json(event.get('context'))) state_machine_data = StateMachineSchema().load(event.get('context')) state_machine_track = TrackSchema().load(event.get('track')) s3_track = find_s3_track(state_machine_track.release_reference, s3_data) 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) graphql_gateway.set_headers( { 'Orchard-User-Id': config.OA_USER, 'Correlation-Id': correlation_id, # Will be used by some ows-assets routes for feature flag checks 'Grass-Account-Id': state_machine_data.product.vendor_id, 'Grass-Account-Type': 'vendor' } ) logging_utils.update_logger_with_message_ids( logger, state_machine_data.message_id, state_machine_data.message_thread_id, state_machine_data.execution_name ) if not state_machine_track.tuid: logger.error(f'TUID not found for ISRC: {state_machine_track.isrc}') raise AudioException( f'TUID not found for ISRC: {state_machine_track.isrc}' ) logger.info(f'Asset source bucket: {state_machine_track.asset.bucket}') logger.info(f'Asset key: {state_machine_track.asset.key}') # Ask for a token from ows-assets asset_token = create_asset_token(graphql_gateway, 'audio') filename = asset_token['filename'] + os.path.splitext( state_machine_track.asset.filename)[1] target_bucket = asset_token['bucket'] # Check if we need to construct a release correction item is_correction = is_release_correction(state_machine_data) if is_correction: if not state_machine_data.error_correction or \ not state_machine_data.error_correction.release_correction_id: raise LambdaException( 'context error correction is empty or missing required data') logger.info('Creating error correction details for context') error_correction_items = construct_error_correction_items( state_machine_data.product.product_id, state_machine_track ) logger.info( f'Error correction items constructed: {error_correction_items}' ) product_correction_detail_data = \ format_correction_detail_data( state_machine_data.product.product_id, state_machine_data.error_correction.release_correction_id, error_correction_items ) results = graphql_gateway.execute( queries.create_product_correction_detail, {'data': product_correction_detail_data} )['data']['createProductCorrectionDetail'] correction_items = [] for correction_detail in results: item = format_correction_detail_response( correction_detail) correction_items.append(item) state_machine_data.error_correction.items.extend(correction_items) # Construct metadata and copy our asset to target bucket from ows-assets s3_metadata = construct_s3_metadata(state_machine_data.product, state_machine_track, '1' if is_correction else '0') logger.info(f'S3 metadata: {s3_metadata}') logger.info(f'Target bucket: {target_bucket}') logger.info(f'Target filename/key: {filename}') copy_asset(s3_metadata, s3_track.asset, target_bucket, filename) state_machine_track.asset.ows_assets_filename = filename return { 'context': StateMachineSchema().dump(state_machine_data), 'track': TrackSchema().dump(state_machine_track) } def construct_s3_metadata(product: object, state_machine_track: object, is_correction: str) -> dict: """Construct assets metadata for s3. Args: product (object): Product object track (object): Track object is_correction (str): Correction value Returns: dict """ return { 'asset_type': format_asset_type(state_machine_track.asset), 'product_id': str(product.product_id), 'upc': product.upc, 'track_unique_id': str(state_machine_track.tuid), 'original_filename': state_machine_track.asset.filename, 'is_correction': is_correction } def construct_error_correction_items(product_id, state_machine_track) -> dict: """Construct error correction object.""" return [ { 'field_name': 'track', 'key_value': True, 'key_id': state_machine_track.tuid, 'table_name': 'track', }, { 'field_name': 'originalFileName', 'key_value': state_machine_track.asset.filename, 'key_id': state_machine_track.tuid, 'table_name': 'track', }, ] def format_asset_type(asset: object): """Format asset type from asset object. Args: asset (object): Asset object Returns: dict """ return os.path.splitext(asset.filename)[1].strip('.').upper() def format_correction_detail_data(product_id, release_correction_id, items): """Format correction detail payload. Args: product_id (int): Product ID release_correction_id (str): Release Correction ID items (list): List of correction detail items Returns: dict """ corrections = [] for item in items: corrections.append({ 'fieldName': item['field_name'], 'keyValue': json.dumps(item['key_value']), 'keyId': item['key_id'], 'tableName': item['table_name'] }) return { 'productId': product_id, 'releaseCorrectionId': release_correction_id, 'corrections': corrections } def format_correction_detail_response(correction_detail): """Format correction detail creation response. Args: correction_detail (dict): Item in correction detail creation response Returns: dict """ return { 'release_correction_detail_id': correction_detail['correctionDetailId'], 'key_id': correction_detail['keyId'], 'key_value': correction_detail['keyValue'], 'table_name': correction_detail['tableName'], 'field_name': correction_detail['fieldName'], } def find_s3_track(state_machine_track_reference, s3_data): """Find the s3_track that matches state_machine_track_reference. Args: state_machine_track_reference(ID): track release_reference to find s3_data (S3Schema): Parsed DDEX object Returns: track """ for s3_track in s3_data.tracks: if state_machine_track_reference == s3_track.release_reference: return s3_track