"""Lambda function module. DDB - DynamoDB, DD DB - direct_delivery database Get an order from DDB streams, load full data from S3, save to DD DB vector order with all the details. """ from collections import defaultdict import json import math import boto3 import sqlalchemy from sentry_sdk import capture_exception from src import config from connectors import direct_delivery from constants import const from constants import fields from src import sql_queries from src import util def get_dict_from_s3(s3_bucket, s3_key): """Get JSON as dict from S3. Args: s3_bucket (str): order's bucket s3_key (str): order's key Returns: dict: full order dict from S3 """ s3 = boto3.resource('s3') s3_obj = s3.Object(s3_bucket, s3_key).get() s3_json = s3_obj['Body'].read().decode() s3_dict = json.loads(s3_json) # json.dumps in another Lambda function converts int dict keys into # strings, and the data is on S3 at this point. So we have to deal with it # here after loading. validated = { int(k): v for k, v in s3_dict[fields.S3_VO_VALIDATED].items()} s3_dict[fields.S3_VO_VALIDATED] = validated return s3_dict def get_order_type_by_encoder_id(encoder_id): """Get order type by encoder id. Args: encoder_id (int): vector encoder type Returns: str: order type """ if encoder_id in const.VO_ENCODER_TO_TYPE: return const.VO_ENCODER_TO_TYPE[encoder_id] raise Exception('Unknown encoder id') def save_order_to_dd_db(order_dict): """Save vector order into encoding_queue. Args: order_dict (dict): set of VO fields from DDB Returns: int: encoding queue id for this VO """ encoder_id = order_dict[fields.VO_ENCODER_ID] order_type = get_order_type_by_encoder_id(encoder_id) with direct_delivery.session_scope() as session: query = session.execute( sqlalchemy.text(sql_queries.INSERT_ENCODING_QUEUE), {'order_id': order_dict[fields.VO_ORDER_ID], 'priority': order_dict[fields.VO_PRIORITY], 'order_type': order_type, 'encoder_id': encoder_id, 'meta_update': order_dict[fields.VO_META_UPDATE]}) encoding_queue_id = query.lastrowid if encoding_queue_id == 0: query = session.execute( sqlalchemy.text(sql_queries.SELECT_ENCODING_QUEUE_ID), {'order_id': order_dict[fields.VO_ORDER_ID], 'order_type': order_type}) encoding_queue_id = query.fetchone()[0] return encoding_queue_id def get_all_order_details(encoding_queue_id): """Get all details for specific encoding queue. Args: encoding_queue_id (int): parent EQ id Returns: dict: key - upc, value - list of store ids """ with direct_delivery.session_scope() as session: query = session.execute( sqlalchemy.text(sql_queries.SELECT_ENCODING_QUEUE_DETAILS), {'encoding_queue_id': encoding_queue_id}) result = defaultdict(list) for detail in query.mappings(): result[detail[fields.EQ_DETAIL_UPC]].append( detail[fields.EQ_DETAIL_STORE_ID]) return dict(result) def get_new_order_details(existing_details, all_details): """Get all new details as (all - existing). Args: existing_details (dict): details that are in dd now, key - upc, value - list of store ids all_details (dict): details that should be in dd, key - upc, value - list of store ids Returns: dict: key - upc, value - list of store ids """ result = defaultdict(set) for upc, all_stores in all_details.items(): if upc not in existing_details: result[upc] = all_stores continue existing_stores = existing_details[upc] result[upc] = list(set(all_stores) - set(existing_stores)) return dict(result) def convert_details_to_list(encoding_queue_id, details): """Convert upc/stores dict to list of columns for insert query. Args: encoding_queue_id (int): parent EQ id details (dict): VO details: key - upc, value - list of store ids Returns: list: of column dicts """ result = [] for upc, stores in details.items(): for store_id in stores: detail_dict = {} detail_dict[fields.EQ_DETAIL_ENCODING_QUEUE_ID] = ( encoding_queue_id) detail_dict[fields.EQ_DETAIL_UPC] = upc detail_dict[fields.EQ_DETAIL_STORE_ID] = store_id result.append(detail_dict) return result def save_details_to_dd_db(details): """Save set of VO details to DD. Args: details (list): values for insert query """ limit = config.BATCH_WRITE_COUNT_LIMIT details_count = len(details) parts_count = math.ceil(details_count / limit) with direct_delivery.session_scope() as session: for i in range(0, parts_count): first_index = limit * i last_index = limit * (i + 1) session.execute( sqlalchemy.text(sql_queries.INSERT_ENCODING_QUEUE_DETAILS), details[first_index:last_index]) def handler(event, context): """Lambda main function.""" try: config.logger.info('Parsing incoming DDB record: %s', event) order_dict = util.parse_aws_dict(event) config.logger.info('Geting full order JSON from S3') s3_dict = get_dict_from_s3( order_dict[fields.DDB_VO_BUCKET_NAME], order_dict[fields.DDB_VO_S3_KEY]) config.logger.info('Saving order to DD') eq_id = save_order_to_dd_db(s3_dict) config.logger.info('Getting encoding queue details from DD') existing_details = get_all_order_details(eq_id) config.logger.info('Selecting only new details') new_details = get_new_order_details( existing_details, s3_dict[fields.S3_VO_VALIDATED]) new_details_list = convert_details_to_list(eq_id, new_details) config.logger.info('Saving order details to DD') save_details_to_dd_db(new_details_list) except Exception as e: capture_exception(e) raise e