"""copy-assets-to-asheville.""" import json from ddex_ingester_common.helpers.rds import run_rds_query from vector_utils.aws_utils import sqs import config from config import secrets_manager_client from constants import constants, graphql_queries, sql_queries logger = config.get_current_logger(config.app_logger) def handler(event, context): """Begin main script execution function.""" logger.info( f'Begin asset copy to Asheville. Event: {event} Context: {context}') # The number of UPCs to process each time the Lambda runs nr_products_per_run = \ event.get('nr_products_per_run') or constants.NR_PRODUCTS_PER_RUN logger.info(f'nr_products_per_run: {nr_products_per_run}') # Get UPCs from RDS logger.info('Reading list of UPCs from RDS') try: rows_to_process = run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, secrets_manager_client.get_cred('rds_read_write_password'), sql_queries.SELECT_UPC_N_ROWS, nr_products_per_run, ) except Exception as err: err_string = str(err) email_payload = f'Failed to execute DB query: {err_string}' email_subject = 'Copy assets to asheville failure' log_error(email_payload, email_subject) return if not rows_to_process: logger.info('No products to process. Ending asset copy to Asheville.') return all_upcs = [row.get('upc') for row in rows_to_process] product_data = [] for upc in all_upcs: try: result = config.graphql_gateway.execute( graphql_queries.GET_PRODUCT_BY_UPC, {'upc': upc} )['data']['productByUpc'] logger.info( f'Ran productByUpc for {upc} and got the result: {result}') if result: product_data.append(json.dumps({ 'project_id': result['project']['projectId'], 'product_id': result['productId'], 'upc': upc, 'physical_location_id': config.PHYSICAL_LOCATION_ID, })) else: logger.info( f'Product was not found. Setting {upc} date to 1971-01-01') run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, secrets_manager_client.get_cred('rds_read_write_password'), sql_queries.UPDATE_ROW_PRODUCT_SET_TO_NOT_FOUND, int(upc), ) except Exception as err: err_string = str(err) email_payload = f'Error retrieving product: {upc}: {err_string}' email_subject = 'Copy assets to asheville failure' log_error(email_payload, email_subject) queue_name = sqs.format_queue_name( sqs.ASSETS_TO_DOWNLOAD_QUEUE_NAME_PATTERN, env=config.ENVIRONMENT ) logger.info(f'Send to SQS queue {queue_name} the messages: {product_data}') try: sqs.add_messages_to_sqs(queue_name, product_data) except Exception as err: err_string = str(err) email_payload = \ f'Error adding to message queue: {upc}: {err_string}' email_subject = 'Copy assets to asheville failure' log_error(email_payload, email_subject) for product in product_data: upc = json.loads(product)['upc'] logger.info(f'Updating RDS row for UPC {upc}') try: run_rds_query( logger, config.RDS_HOST, config.RDS_DB_NAME, config.RDS_RW_USER, secrets_manager_client.get_cred('rds_read_write_password'), sql_queries.UPDATE_ROW_ASSET_COPY_STARTED, upc, ) except Exception as err: err_string = str(err) email_payload = \ f'Error updating DB for product: {upc}: {err_string}' email_subject = 'Copy assets to asheville failure' log_error(email_payload, email_subject) def log_error(email_payload, email_subject): """Log error encountered by lambda.""" logger.info( f' With subject: {email_subject}' f' With payload {email_payload}' )