"""Lambda phf_sales_data_process function module.""" import csv from itertools import islice import os import re from sqlalchemy import exc from constants import errors from constants import general from constants import publishing from constants import sales_file_structure from models import period from models import phf_publishing_escrow from warmer_util import catch_warmer_event import file_processing import lambda_exceptions import s3 import status as common_status import util @catch_warmer_event() @util.handle_lambda_result( lambda_name=general.LAMBDA_NAME, default_error_code=errors.FILE_PROCESSING_ERROR_CODE) def handler(event, context): """Lambda entry point.""" status = 200 s3_object = event.get('s3_object', {'bucket': None, 'key': None}) additional_info = event.get('additional_info', {}) start_index = additional_info.get('start_index', 0) batch_start_index = start_index stop_index = ( start_index + sales_file_structure.ROWS_TO_PROCESS_PER_ITERATION) error_code = None error_description = event.get('error_description', []) # This variable is used for monitoring original file number before all # processing and validations, so user will be able to find rows with errors # default line number is 1 because header line is skipped during processing line_number = additional_info.get('line_number', 1) batch_start_line_number = line_number retries_count = additional_info.get('retries_count', 0) s3_downloaded_file = None try: file_path = '/tmp/{}'.format(s3_object['key']) s3_downloaded_file = s3.download_object( s3_object['bucket'], s3_object['key'], file_path) if start_index == 0: (phf_publishing_escrow. set_all_publishing_escrows_from_file_inactive(s3_object['key'])) periods = period.get_all_periods() file_encoding = file_processing.guess_file_segment_encoding( s3_downloaded_file, start_index, stop_index) with open( s3_downloaded_file, encoding=file_encoding, errors='strict', newline='\r\n') as sales_file: headers = re.sub(r'(?= publishing.DEFAULT_BATCH_SIZE: file_processing.bulk_insert_new_phf_metadata( processed_tracks, processed_escrows) processed_escrows = [] processed_tracks = {} batch_start_index = current_index + 1 batch_start_line_number = line_number + 1 if (context.get_remaining_time_in_millis() <= general.ABORT_IF_TIME_LIMIT_LEFT_LESS_THAN_MS): start_next_processing_on = current_index + 1 if collected_errors: error_code = ( errors.FILE_PARTIALLY_PROCESSED_ERROR_CODE) error_description = collected_errors return { 'status': 200, 's3_object': s3_object, 'error_code': error_code, 'error_description': error_description or None, 'additional_info': { 'start_index': start_next_processing_on, 'line_number': line_number}} except lambda_exceptions.DataValidationFailed as e: error_code = getattr(e, 'error_code', error_code) collected_errors = collect_errors(collected_errors, e) except exc.SQLAlchemyError as e: (phf_publishing_escrow. set_all_publishing_escrows_from_file_inactive( s3_object['key'])) return { 'status': 500, 's3_object': s3_object, 'error_code': errors.DATABASE_ERROR_CODE, 'error_description': str(e)} except Exception as e: error_code = getattr(e, 'error_code', error_code) if error_code == errors.DATA_VALIDATION_ERROR_CODE: line_number = None collected_errors = collect_errors( collected_errors, e, line_number) continue total_processed = start_index + rows_processed if not rows_processed: raise lambda_exceptions.ContentIsMissing if processed_tracks or processed_escrows: file_processing.bulk_insert_new_phf_metadata( processed_tracks, processed_escrows) if collected_errors: error_code = errors.FILE_PARTIALLY_PROCESSED_ERROR_CODE error_description = collected_errors if total_processed < stop_index: status = 201 if status == 201 and error_description: common_status.send_failure_status( general.LAMBDA_NAME, status, s3_object['key'], error_code, error_description, s3_object['bucket']) return { 'status': status, 's3_object': s3_object, 'error_code': error_code, 'error_description': error_description or None, 'additional_info': { 'start_index': stop_index, 'line_number': line_number}} except csv.Error: raise lambda_exceptions.FileEncodingNotIdentified except exc.IntegrityError as e: if retries_count < 3: # Restart processing from records after last ingested batch return { 'status': 200, 's3_object': s3_object, 'error_code': error_code, 'error_description': error_description or None, 'additional_info': { 'start_index': batch_start_index, 'line_number': batch_start_line_number, 'retries_count': retries_count + 1}} else: raise e except exc.SQLAlchemyError as e: (phf_publishing_escrow. set_all_publishing_escrows_from_file_inactive( s3_object['key'])) return { 'status': 500, 's3_object': s3_object, 'error_code': errors.DATABASE_ERROR_CODE, 'error_description': str(e)} except lambda_exceptions.LoggedException as e: error_description = [str(e)] status = 400 return { 'status': status, 's3_object': s3_object, 'error_code': error_code, 'error_description': error_description or None} finally: if s3_downloaded_file: os.remove(s3_downloaded_file) def collect_errors(collected_errors, recent_error, line_number=None): """Prepare e-mail message body. Args: collected_errors (dict): errors collected during processing. recent_error (Exception): error instance. line_number (int): number where an error occurred. Returns: collected_errors (dict): errors collected during processing including the last one. """ error_code = getattr( recent_error, 'error_code', errors.FILE_PROCESSING_ERROR_CODE) if error_code == errors.CONTENT_DOES_NOT_MATCH_HEADERS_ERROR_CODE: collected_errors.setdefault(error_code, {}) collected_errors[error_code].setdefault('line_numbers', []).append( line_number) elif error_code == errors.MISSING_MANDATORY_FIELDS_ERROR_CODE: missing_fields = recent_error.missing_fields collected_errors.setdefault(error_code, {}) collected_errors[error_code][line_number] = missing_fields elif error_code == errors.DATA_VALIDATION_ERROR_CODE: multiple_rows_invalid_fields = recent_error.invalid_fields collected_errors.setdefault(error_code, {}) for field, value in multiple_rows_invalid_fields.items(): collected_errors[error_code].setdefault(field, []) collected_errors[error_code][field] = value else: recent_error = str(recent_error) collected_errors.setdefault( errors.PARTIAL_PROCESSING_ERRORS, {}) collected_errors[errors.PARTIAL_PROCESSING_ERRORS].setdefault( recent_error, []).append(line_number) return collected_errors