"""Main Handler.""" from typing import Dict, List from bulk_metadata_ingester_common.models.bulk_release import BulkRelease from bulk_metadata_ingester_common.utils.error import graphql_execute from bulk_metadata_ingester_common.utils.jitter import jitter from bulk_metadata_ingester_common.utils.logging import get_current_logger from bulk_metadata_ingester_common.utils.release_utils import \ get_release_track_list import config from config import graphql_gateway from constants import queries from constants.constants import WORKSTATION_VALIDATION_RULE_ID from ddex_ingester_common.constants.catalog_ingestion import WARNING from ddex_ingester_common.helpers.catalog_ingestion import ( CatalogIngestionValidationResult ) from ddex_ingester_common.lambda_exceptions import ValidateProductException from lambdacommon.graphql import graphql def handler(event: Dict, context: object) -> Dict: """Lambda Entrypoint.""" # TODO remove when Lambda is implemented and working # This key won't be in the event unless we explicitly add it # It skips some Lambdas to allow testing of other parts of the sfn if 'validate_product' not in event: return event key = event.get('key') bucket = event.get('bucket') item_index = event.get('item_index') state_machine_execution_name = event.get('execution_name') state_machine_name = event.get('state_machine_name') correlation_id = event.get('correlation_id') release = event.get('release', {}) tracks = event.get('tracks', {}) 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'validate_product received: {event}') # Get release from JSON doc release_track_list = get_release_track_list(bucket, key, item_index) # Move JSON dict to model modeled_release = BulkRelease(json_release_rows=release_track_list) # Rehydrate context modeled_release.rehydrate(release, tracks) # Execute lambda logic errors = [] try: validate_result = graphql_execute( graphql_gateway, queries.validate_product, {'productId': str(modeled_release.product_id)}, logger )['data']['product'] if not validate_result['validation']['isValid']: product_errors = validate_result['validation']['errors'] if product_errors: logger.info(f'Product errors: {product_errors}') state_machine_name = 'hardcoded_machine_name' state_machine_execution_name = 'hardcoded_execution_name' results = [ CatalogIngestionValidationResult( state_machine_name=state_machine_name, state_machine_execution_name=state_machine_execution_name, # noqa validation_rule_id=WORKSTATION_VALIDATION_RULE_ID, response=WARNING, message=error['reason'], category=error['code'] ) for error in product_errors ] config.catalog_ingestion_session.add(results) product_errors = flatten_errors(product_errors) errors.append(product_errors) for track in validate_result['tracks']: isrc = track['isrc'] if not track['validation']['isValid']: track_errors = track['validation']['errors'] # no warnings at the time of this implementation # track_warnings = track['validation']['warnings'] if track_errors: logger.info(f'Track errors: {track_errors}') results = [ CatalogIngestionValidationResult( state_machine_name=state_machine_name, state_machine_execution_name=state_machine_execution_name, # noqa validation_rule_id=WORKSTATION_VALIDATION_RULE_ID, response=WARNING, isrc=isrc, message=error['reason'], category=error['code'] ) for error in track_errors ] config.catalog_ingestion_session.add(results) track_errors = flatten_errors(track_errors, isrc) errors.append(track_errors) config.catalog_ingestion_session.save() except graphql.GraphQLError as err: raise ValidateProductException('Graphql error') from err # Get all passed and updated fields updates = modeled_release.get_modified() result = { **event, 'release': { **updates['release'], }, 'tracks': { **updates['tracks'] } } return result def flatten_errors(errors: List[Dict], isrc: str = None) -> List[str]: """Convert nested list to one dimension error list.""" flattened_errors = [] for key, value in enumerate(errors): if type(value) == list: flattened_errors += flatten_errors(value, isrc) else: message = value['reason'] if isrc: message = f'ISRC: {isrc}: {message}' flattened_errors.append(message) return flattened_errors