"""Lambda function module. Load all closed and non processed EO from AR, save to S3 and DynamoDB, update EO status in AR. """ import collections import json import boto3 from botocore import config as botocore_config import sqlalchemy import common_config import config from connectors import art_relations from constants import common_const from constants import common_fields from constants import fields import dynamodb_config import exceptions import s3_config import sql_queries 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 vaiable with encoder ids encoder_ids = get_encoder_ids() with art_relations.session_scope() as session: # execute query query = session.execute( sqlalchemy.text(sql_queries.OBTAIN_ENCODING_ORDERS), {'limit': config.ORDERS_AR_SELECT_LIMIT, 'in_values': tuple(encoder_ids)}) # build resulting dict orders = collections.OrderedDict() for detail in query: 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: list() } 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] }) return orders def update_order_status(order_id, status): """ Update AR EO status. Args: order_id (int): EO id status (str): new status """ with art_relations.session_scope() as session: session.execute( sqlalchemy.text(sql_queries.UPDATE_ENCODING_ORDER_STATUS), {'status': status, 'order_id': order_id}) 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 = {} base_dict[common_fields.VO_ORDER_ID] = ( str(order[fields.AR_EO_ORDER_ID])) base_dict[common_fields.VO_CREATED_AT] = ( order[fields.AR_EO_ENTRY_DATE].strftime( common_const.JSON_DATETIME_FORMAT)) base_dict[common_fields.VO_PRIORITY] = ( order[fields.AR_EO_PRIORITY]) base_dict[common_fields.VO_ENCODER_ID] = ( order[fields.AR_EO_ENCODER_ID]) base_dict[common_fields.VO_META_UPDATE] = ( order[fields.AR_EO_META_UPDATE]) base_dict[common_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[common_fields.DDB_VO_BUCKET_NAME] = ( s3_config.ORDERS_S3_BUCKET) ddb_dict[common_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 (list): 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[common_fields.S3_VO_PRODUCTS] = products full_dict[common_fields.S3_VO_STORES] = stores full_dict[common_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( s3_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(dynamodb_config.ORDERS_DDB_TABLE) table.put_item(Item=order_dict) @exceptions.sentry_capture_exception def handler(event, context): """Lambda main function.""" common_config.logger.info('Retrieving orders form AR') orders = get_orders() common_config.logger.info('Got orders from AR: %s', orders.keys()) for order_id, order in orders.items(): common_config.logger.debug('Processing order ID %d', order_id) ddb_dict = order_to_dynamo_dict(order) s3_dict = order_to_full_dict(order) common_config.logger.debug('Writing data to S3') put_to_s3( ddb_dict[common_fields.DDB_VO_BUCKET_NAME], ddb_dict[common_fields.DDB_VO_S3_KEY], json.dumps(s3_dict)) common_config.logger.debug('Writing data to DynamoDB: %s', ddb_dict) put_to_dynamo(ddb_dict) common_config.logger.debug( 'Updating AR status for order ID %d', 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)