"""Lambda function module. Load all closed and not processed EO from AR, save to S3 and DynamoDB, update EO status in AR. """ import collections import copy import json import boto3 from botocore import config as botocore_config from lambdacommon import util from src.logger import get_current_logger import src.sql_queries as sql_queries import config from constants import fields class InvocationContext: """This is for passing around the invoke context to logger.""" correlation_id = None def get_encoder_ids(): """Parse and return encoder ids list from environment variable. Returns: list: list of encoder ids (int) """ encoder_ids = config.ORDER_ENCODER_ID_LIST.split(',') return [int(el.strip()) for el in encoder_ids] def get_orders(): """Get set of encoding order details from AR. Returns: dict: key (str) EO id, value (list) details as dict """ # parse env variable with encoder ids results = get_orders_from_ar() orders = collections.OrderedDict() unique_products = collections.OrderedDict() for detail in results: order_id = detail[fields.AR_EO_ORDER_ID] if order_id not in orders: orders[order_id] = { fields.AR_EO_ORDER_ID: order_id, fields.AR_EO_ENTRY_DATE: detail[ fields.AR_EO_ENTRY_DATE], fields.AR_EO_PRIORITY: detail[fields.AR_EO_PRIORITY], fields.AR_EO_ENCODER_ID: detail[fields.AR_EO_ENCODER_ID], fields.AR_EO_META_UPDATE: detail[ fields.AR_EO_META_UPDATE], fields.AR_EO_USER_ID: detail[fields.AR_EO_USER_ID], fields.EO_DETAILS: [], } orders[order_id][fields.S3_VO_PRODUCTS] = [] orders[order_id][fields.EO_DETAILS].append({ fields.AR_EO_UPC: detail[fields.AR_EO_UPC], fields.AR_EO_STORE_ID: detail[fields.AR_EO_STORE_ID]}) orders[order_id][fields.S3_VO_PRODUCTS].append( { 'project_id': detail[fields.AR_EO_PROJECT_ID], 'product_id': detail[fields.AR_EO_PRODUCT_ID], 'upc': detail[fields.AR_EO_UPC], } ) if detail[fields.AR_EO_META_UPDATE] == 'N': unique_products[detail[fields.AR_EO_UPC]] = { 'project_id': detail[fields.AR_EO_PROJECT_ID], 'product_id': detail[fields.AR_EO_PRODUCT_ID], 'upc': detail[fields.AR_EO_UPC], } physical_locations = config.PHYSICAL_LOCATIONS.split(',') products = [] for upc, product in unique_products.items(): for physical_location_id in physical_locations: tmp_product = copy.deepcopy(product) tmp_product['physical_location_id'] = int(physical_location_id) products.append(tmp_product) return {'orders': orders, 'products': products} def get_orders_from_ar(): """Query art_relations to get unprocessed orders. Returns: list """ with util.ar_connection(config.AR_DB_CREDENTIALS) as conn: with conn.cursor() as cursor: cursor.execute(sql_queries.GET_ENCODING_ORDERS.format( limit=config.ORDERS_AR_SELECT_LIMIT, encoder_ids=config.ORDER_ENCODER_ID_LIST)) results = cursor.fetchall() order_ids = [row[fields.AR_EO_ORDER_ID] for row in results] if order_ids: cursor.execute( sql_queries.GET_ENCODING_ORDER_DETAILS.format( order_ids=','.join( [str(order_id) for order_id in order_ids]))) return cursor.fetchall() return [] def update_order_status(order_id, status): """ Update AR EO status. Args: order_id (int): EO id status (str): new status """ with util.ar_connection(config.AR_DB_CREDENTIALS) as conn: with conn.cursor() as cursor: cursor.execute(sql_queries.UPDATE_ENCODING_ORDER_STATUS.format( status=status, order_id=order_id)) conn.commit() def order_to_base_dict(order): """Generate a dict with base common attributes. Args: order (dict): encoding order Returns: dict: dictionary with some common EO fields """ base_dict = dict() base_dict[fields.VO_ORDER_ID] = ( str(order[fields.AR_EO_ORDER_ID])) base_dict[fields.VO_CREATED_AT] = ( order[fields.AR_EO_ENTRY_DATE].strftime('%Y-%m-%dT%H:%M:%S')) base_dict[fields.VO_PRIORITY] = ( order[fields.AR_EO_PRIORITY]) base_dict[fields.VO_ENCODER_ID] = ( order[fields.AR_EO_ENCODER_ID]) base_dict[fields.VO_META_UPDATE] = ( order[fields.AR_EO_META_UPDATE]) base_dict[fields.VO_USER_ID] = ( order[fields.AR_EO_USER_ID]) return base_dict def order_to_dynamo_dict(order): """Build dict for DynamoDB. Args: order (dict): encoding order with all EO fields included Returns: dict: DynamoDB EO information """ ddb_dict = order_to_base_dict(order) ddb_dict[fields.DDB_VO_BUCKET_NAME] = ( config.ORDERS_S3_BUCKET) ddb_dict[fields.DDB_VO_S3_KEY] = ( get_s3_key(order[fields.AR_EO_ORDER_ID])) return ddb_dict def order_to_full_dict(order): """Build a dict with all EO information including details. Args: order (dict): EO with details Returns: dict: full EO information including details """ full_dict = order_to_base_dict(order) products = list() stores = list() validated = {} for detail in order[fields.EO_DETAILS]: upc = detail[fields.AR_EO_UPC] if upc not in products: products.append(upc) store_id = detail[fields.AR_EO_STORE_ID] if store_id not in stores: stores.append(store_id) if upc not in validated: validated[upc] = list() validated[upc].append(store_id) full_dict[fields.S3_VO_PRODUCTS] = products full_dict[fields.S3_VO_STORES] = stores full_dict[fields.S3_VO_VALIDATED] = validated return full_dict def get_s3_key(order_id): """Generate s3 order path. Args: order_id (int): EO id Returns: str: EO s3 key """ return '{:s}/{:s}.json'.format( config.ORDERS_S3_PREFIX.strip('/'), str(order_id)) def put_to_s3(s3_bucket, s3_key, s3_json): """Put JSON to s3 with key. Args: s3_bucket (str): s3 bucket name s3_key (str): unique s3 order key s3_json (str): EO full JSON """ s3 = boto3.resource('s3') s3.Object(s3_bucket, s3_key).put( Body=s3_json, ContentType='application/json') def put_to_dynamo(order_dict): """Put EO to DynamoDB. Args: order_dict (dict): EO base information """ dynamodb = get_dynamodb() table = dynamodb.Table(config.ORDERS_DDB_TABLE) table.put_item(Item=order_dict) def handler(event, context): """Lambda main function.""" InvocationContext.correlation_id = context.aws_request_id log = get_current_logger(InvocationContext.correlation_id) log.info('Retrieving orders from AR') orders, products = get_orders().values() log.info('Got orders from AR: {}'.format(orders.keys())) for order_id, order in orders.items(): log.info('Processing order ID {}'.format(order_id)) ddb_dict = order_to_dynamo_dict(order) s3_dict = order_to_full_dict(order) log.info('Encoding order ID {} contains {} details '.format( order_id, len(order[fields.EO_DETAILS]))) log.info('Writing data to S3') put_to_s3( ddb_dict[fields.DDB_VO_BUCKET_NAME], ddb_dict[fields.DDB_VO_S3_KEY], json.dumps(s3_dict)) log.info('Writing data to DynamoDB: {}'.format(ddb_dict)) put_to_dynamo(ddb_dict) log.info( 'Updating AR status for order ID {}'.format(order_id)) update_order_status(order_id, fields.AR_EO_PROCESSED_DONE) def get_dynamodb(): """Get Boto DynamoDB resource.""" config_instance = botocore_config.Config( retries={'max_attempts': config.DDB_WRITE_MAX_ATTEMPTS}) return boto3.resource('dynamodb', config=config_instance)