"""register-delivery-to-ingest.""" from datetime import datetime, timedelta import json import os import re import traceback from typing import List import uuid from xml.dom.minidom import parseString import boto3 import config from constants import queries from ddex_ingester_common.constants.catalog_ingestion import ( DDEX_PROVIDER_TO_SOURCE_ID, REJECT ) from ddex_ingester_common.constants.ddex_providers import ( ALTAFONTE, ALTAFONTE_FOLDER_NAME, AWAL, RISING_88, RISING_88_FOLDER_NAME, SME, SME_ANALYTICS_PROVIDER, SME_ANALYTICS_PROVIDER_FOLDER_NAME ) from ddex_ingester_common.constants.send_acknowledgement import \ SME_ERROR_STATUS, SME_SUCCESS_STATUS from ddex_ingester_common.helpers.rds import run_rds_query from ddex_ingester_common.logging import utils as logging_utils from ddex_ingester_common.schemas.s3_schema import S3Schema from ddex_ingester_common.send_acknowledgement import ddex_scsm from ddex_ingester_common.utils import ses from ddex_ingester_common.utils.prep_ddex import prep_s3_context from ddex_ingester_common.utils.read_ddex import read_ddex from ddex_ingester_common.validation.rules import get_validation_rules from ddex_ingester_common.validation.utils import save_validation_result from mako.lookup import TemplateLookup import requests from requests import auth logger = logging_utils.get_logger(config.app_logger) def handler(event, context): """Lambda entrypoint.""" logger.info(f'Lambda triggered with event {event} and context {context}') # Check event key is for an xml file correlation_id = event.get('correlation_id') or str(uuid.uuid4()) logging_utils.update_logger_correlation_id(logger, correlation_id) delivered_time = event.get('time', '') event_detail = event.get('detail', {}) if not event_detail.get('requestParameters', {}): event_detail['requestParameters'] = {} event_request_params = event_detail['requestParameters'] if event.get('detail-type', '') == 'Object Created': delivery_key = event_detail.get('object', {}).get('key') delivery_bucket = event_detail.get('bucket', {}).get('name') # Simulate old object structure for backward compatibility event_request_params['bucketName'] = delivery_bucket event_request_params['key'] = delivery_key else: delivery_bucket = event_request_params.get('bucketName') delivery_key = event_request_params.get('key') delivery_provider = event_request_params.get('ddex_provider') if not delivery_provider: if RISING_88_FOLDER_NAME in delivery_key: delivery_provider = RISING_88 elif ALTAFONTE_FOLDER_NAME in delivery_key: delivery_provider = ALTAFONTE # SME_ANALYTICS_PROVIDER must be before SME elif SME_ANALYTICS_PROVIDER_FOLDER_NAME in delivery_key: delivery_provider = SME_ANALYTICS_PROVIDER elif SME.lower() in delivery_key \ and SME_ANALYTICS_PROVIDER_FOLDER_NAME not in delivery_key: delivery_provider = SME elif AWAL.lower() in delivery_key: delivery_provider = AWAL else: logger.error('Missing delivery provider') return all_ddex_providers = [SME, AWAL, RISING_88, ALTAFONTE, SME_ANALYTICS_PROVIDER] if delivery_provider.upper() not in all_ddex_providers: logger.error(f'Unknown delivery provider: {delivery_provider}') return if not delivery_bucket or not delivery_key: logger.error('Missing request parameter') return if not delivery_key.endswith('.xml'): logger.error('Delivered file is not of type XML') return if 'resources' in os.path.dirname(delivery_key): logger.error('Ignoring "resources" folder') return provider_source_id = DDEX_PROVIDER_TO_SOURCE_ID.get( delivery_provider.upper() ) if not provider_source_id: logger.error(f'Cannot map delivery provider: {delivery_provider}') return if not delivered_time: # Turn 2021-09-14 13:01:14.524565 into 2021-09-09T13:39:00Z delivered_time = datetime.utcnow().strftime('%Y-%m-%dT%H:%M:%S') + 'Z' logger.debug(f'Missing event time, defaulting to: {delivered_time}') logger.info( f'Execution data: ' f'time - {delivered_time} | ' f'provider - {delivery_provider} | ' f'bucket - {delivery_bucket} | ' f'key - {delivery_key}') if delivery_provider == SME_ANALYTICS_PROVIDER: try: s3_client = boto3.client('s3') doc = read_ddex(s3_client, delivery_bucket, delivery_key, logger) if doc.is_purged_release: grid = delivery_key.split('/')[-1].split('.')[0] logger.info( 'Found purged_release : Bucket:' f' {delivery_bucket} Key: {delivery_key}') _store_purged_ddex_delivery(provider_source_id, delivery_bucket, delivery_key, grid, doc.product.upc) send_acknowledgement(grid, SME_SUCCESS_STATUS) return is_confidential = doc.product.is_confidential has_artwork = get_has_artwork(doc) # Note earliest column not used in flow control since DS-9968 earliest_start_datetime = \ get_earliest_start_time(doc, 'start_date_time', False) earliest_preorder_datetime = \ get_earliest_start_time(doc, 'start_date_time', True) earliest_start_date = \ get_earliest_start_time(doc, 'start_date', False) earliest_preorder_date = \ get_earliest_start_time(doc, 'start_date', True) release_type = doc.product.release_type doc.bucket = delivery_bucket doc.key = delivery_key upc, grid, catalog_number = get_ids_from_ddex(delivery_bucket, delivery_key) if grid is None: raise ValueError('Can not get grid') s3_schema = S3Schema().load(S3Schema().dump(doc)) prep_ddex = prep_s3_context(delivery_provider, release_type, s3_schema, False, is_ccm_som_livre_grps_ddex_ingestion_enabled(), # noqa logger) results = validate_and_save_results(delivery_provider, prep_ddex, delivery_key) rejected_rules = [ result.message for result in results if result.response == REJECT ] db_status = queries.INGEST_READY_STATUS acknowledgment_status = SME_SUCCESS_STATUS if rejected_rules: db_status = queries.VALIDATION_FAILED_STATUS acknowledgment_status = SME_ERROR_STATUS for date in _get_ingestion_dates( S3Schema().load(S3Schema().dump(doc)), is_confidential): _store_ddex_delivery(provider_source_id, delivery_bucket, delivery_key, is_confidential, has_artwork, earliest_start_datetime, earliest_preorder_datetime, earliest_start_date, earliest_preorder_date, db_status, upc, grid, date.strftime('%Y-%m-%d %H:%M:%S')) send_acknowledgement(grid, acknowledgment_status) if acknowledgment_status == SME_ERROR_STATUS: send_error_notification( delivery_bucket, delivery_key, delivery_provider, None, rejected_rules) except Exception as exception: logger.error(exception) send_error_notification( delivery_bucket, delivery_key, delivery_provider, exception, None) else: query_args = ( provider_source_id, delivery_bucket, delivery_key, queries.INGEST_READY_STATUS ) run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.secrets_manager_client.get_cred('rds_read_write_password'), queries.INSERT_QUERY, query_args, ) return event def get_has_artwork(doc): """Return true if product has artwork.""" has_artwork = False if doc.product.artwork: has_artwork = True return has_artwork def get_earliest_start_time(doc, start_time_attribute, is_pre_order_deal): """Get the earliest time for pre_order or not pre_order deals. Args: doc: DDEXSchema data start_time_attribute: deal field to get start_date_time or start_date is_pre_order_deal: get earliest_start_time for pre_order_deals or not Returns: min start_time """ earliest_start_time = None for release_deal in doc.deals: for deal_term in release_deal.deal_terms: if bool(deal_term.pre_order) == is_pre_order_deal: date = getattr(deal_term, start_time_attribute) if date: date_format = config.date_formats[start_time_attribute] if date.endswith('Z'): date_format = date_format + 'Z' date = datetime.strptime(date, date_format) else: continue if not earliest_start_time: earliest_start_time = date elif earliest_start_time > date: earliest_start_time = date if earliest_start_time: return earliest_start_time.strftime( config.date_formats[start_time_attribute]) def get_ids_from_ddex(bucket, key): """Get ids directly from the DDEX file if the schema failed to parse.""" logger.info(f'Getting ids from ddex file : Bucket: {bucket} Key: {key}') s3_client = boto3.client('s3') file = s3_client.get_object(Bucket=bucket, Key=key) upc = None grid = None catalog_number = None try: doc = parseString(file.get('Body').read()) release_list = doc.getElementsByTagName('Release') for release in release_list: if release.getAttribute('IsMainRelease'): release_id = release.getElementsByTagName('ReleaseId')[0] grid_node = release_id.getElementsByTagName('GRid')[0] upc_node = release_id.getElementsByTagName('ICPN')[0] c_nr_node = release_id.getElementsByTagName('CatalogNumber')[0] grid = grid_node.firstChild.data upc = upc_node.firstChild.data catalog_number = c_nr_node.firstChild.data logger.info( f'UPC: {upc} GRID: {grid} Catalog Number: {catalog_number}' ) break except Exception: logger.info(f'{key} is a malformed XML or not an XML file.') if not grid: logger.info('Getting grid using filename') _, filename = os.path.split(key) matches = re.match('A10[0-9A-Z]{15}', filename) grid = matches.group(0) return upc, grid, catalog_number def send_acknowledgement(grid, status): """Send acknowledge with XML attachment. Args: status (str): grid (str): """ sme_ddexws_user = config.register_delivery_secrets_manager_client \ .get_cred('SME_ANALYTICS_PROVIDER_USER') sme_ddexws_passwd = config.register_delivery_secrets_manager_client \ .get_cred('SME_ANALYTICS_PROVIDER_PASSWD') scsm_generator = ddex_scsm. \ DDEXSupplyChainStatusMessage(grid, status) xml_data = scsm_generator.output_xml() session = requests.Session() request = requests.Request( method='post', url=config.SME_DDEXWS_HOST, auth=auth.HTTPBasicAuth( username=sme_ddexws_user, password=sme_ddexws_passwd ), headers={'Content-Type': 'application/xml'}, data=xml_data.decode('UTF-8') ).prepare() response = session.send(request) if response.status_code != 200: err_msg = ( 'Call to SME DDEX WS failed. ' f'Response: {response.content}') logger.error(err_msg) raise Exception(err_msg) logger.info(f'Call to SME DDEX WS with status {status}.') def send_error_notification( bucket, key, ddex_provider, exception: Exception, rejected_rules): """Send email with error message.""" recipients = json.loads(config.REPORT_ERRORS_SME_ANALYTICS_DUMMY_EMAIL) upc, grid, catalog_number = None, None, None try: upc, grid, catalog_number = get_ids_from_ddex(bucket, key) except Exception: logger.error(f'Failed to get_ids from {bucket}/{key}') if ddex_provider == SME_ANALYTICS_PROVIDER: sme_product = 'SME Analytics Product' else: sme_product = 'SME Interop Product' email_subject = f'{sme_product} s3 location: {key}' if grid: email_subject = f'{sme_product} {grid},' \ f's3 location: {key}' if upc: email_subject = f'{sme_product} {upc}, ' \ f's3 location: {key}' if grid and upc: email_subject = \ f'{sme_product} {grid} {upc}, ' \ f's3 location: {key}' if exception: error_message = f'Failed to process product s3:{bucket}/{key} .' \ f'exception_info : {exception.args} ' \ f'exception_traceback : ' \ f'{traceback.format_exception(exception)}' email_subject = 'Failure to ingest ' + email_subject else: email_subject = 'Validation failed ' + email_subject rejected_validations = '\n'.join(rejected_rules) error_message = f'Failed validation rules : {rejected_validations} ' mylookup = TemplateLookup(directories=['constants/templates']) base_template = mylookup.get_template('no_context_failure.mak') email_payload = base_template.render( upc=upc, grid=grid, sony_product_id=catalog_number, error_message=error_message, # Rest of the table fields additional_details=None, title=None, artist=None, format=None, owner=None, label=None, status=None, nfd=None, ).replace('\n', '') send_email(recipients, email_payload, email_subject, []) logger.info('Successfully sent email.') def send_email( recipients: List[str], email_payload: str, email_subject: str, cc_addresses: List): """Send email to given contact. Args: recipients (list): Recipient list email_payload (string): Email text content containing errors email_subject (string): Subject of email to be sent cc_addresses (list): CC list recipients Returns: list """ logger.info( f'Sending an email to: {recipients}' f' With CC: {cc_addresses}' f' With subject: {email_subject}' f' With payload {email_payload}' ) ses.Email().send( email_subject, config.EMAIL_FROM_ADDRESS, recipients, email_payload, None, cc_addresses ) def validate_and_save_results(delivery_provider, doc, key): """Validate and save results.""" rules = get_validation_rules(delivery_provider, config.APPLICATION_NAME, None, config.graphql_gateway) data = S3Schema().dump(doc) validation_results = [] for rule in rules: logger.info(f'Executing validation rule: {rule.__name__}') validation_results.append(rule(data)) results = save_validation_result( validation_results, f'Lambda_{config.LAMBDA_NAME}', f's3 file key : {key}', config.catalog_ingestion_session ) return results def is_ccm_som_livre_grps_ddex_ingestion_enabled(): """Lookup if ccm_som_livre_grps_ddex_ingestion split is enabled.""" split = config.split_client is_enabled = split and split.get_treatment( config.APPLICATION_NAME, 'ccm_som_livre_grps_ddex_ingestion', {'service': config.APPLICATION_NAME} ) == 'on' return is_enabled def _store_purged_ddex_delivery(provider_source_id, delivery_bucket, delivery_key, grid, upc): query_args = ( provider_source_id, delivery_bucket, delivery_key, False, False, None, None, None, None, queries.PURGED_RELEASE, upc, grid, None ) run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.RDS_PASSWORD, queries.INSERT_QUERY_SME_ANALYTICS, query_args, ) def _store_ddex_delivery(provider_source_id, delivery_bucket, delivery_key, is_confidential, has_artwork, earliest_start_datetime, earliest_preorder_datetime, earliest_start_date, earliest_preorder_date, db_status, upc, grid, original_release_datetime): query_args = ( provider_source_id, delivery_bucket, delivery_key, is_confidential, has_artwork, earliest_start_datetime, earliest_preorder_datetime, earliest_start_date, earliest_preorder_date, db_status, upc, grid, original_release_datetime ) run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, config.RDS_PASSWORD, queries.INSERT_QUERY_SME_ANALYTICS, query_args, ) def _get_ingestion_dates(s3_schema, is_confidential): earliest_start_dates = sorted(set( [datetime.strptime(track.original_release_date, '%Y-%m-%d') for track in s3_schema.tracks if track.original_release_date])) logger.info(f'earliest_start_dates : {earliest_start_dates}') if _r0_deal_exists(s3_schema): r0_start_date_time = _get_r0_deal_earliest_start_datetime(s3_schema) else: r0_start_date_time = None filter_out_outdated_dates = [date for date in earliest_start_dates if date >= datetime.now() - timedelta(days=1) ] logger.info(f'filter_out_outdated_dates : {filter_out_outdated_dates}') if is_confidential: if r0_start_date_time is None: raise Exception('R0 deal is missing for confidential product') filter_out_dates_before_r0_release_date = [date for date in filter_out_outdated_dates if date >= r0_start_date_time ] # for confidential r0.start_date is the first processing event filter_out_dates_before_r0_release_date.append(r0_start_date_time) scheduling_dates = filter_out_dates_before_r0_release_date else: scheduling_dates = filter_out_outdated_dates # for not confidential in case if there are release dates in the past # get latest earliest_start_dates in order to schedule ingestion # for all release dates in past if len(filter_out_outdated_dates) != len(earliest_start_dates): scheduling_dates = filter_out_outdated_dates outdated_dates = sorted(set(set(earliest_start_dates) - set( filter_out_outdated_dates))) scheduling_dates.append(outdated_dates[-1]) if r0_start_date_time is not None \ and r0_start_date_time > datetime.now(): scheduling_dates.append(r0_start_date_time) logger.info(f'scheduling_dates : {len(set(scheduling_dates))}') return sorted(set(scheduling_dates)) def _get_r0_deal_earliest_start_datetime(s3_schema): r0_deal_terms = None for deal in s3_schema.deals: if 'R0' in deal.release_references: r0_deal_terms = deal.deal_terms[0] start_date = None start_date_time = None if r0_deal_terms.start_date: start_date = datetime.strptime(r0_deal_terms.start_date, '%Y-%m-%d') if r0_deal_terms.start_date_time: date_format = '%Y-%m-%dT%H:%M:%S' if r0_deal_terms.start_date_time.endswith('Z'): date_format = date_format + 'Z' start_date_time = datetime.strptime(r0_deal_terms.start_date_time, date_format) return min(x for x in [start_date, start_date_time] if x is not None) def _r0_deal_exists(s3_schema): for deal in s3_schema.deals: if 'R0' in deal.release_references: return True return False