"""Job Priority Rule Match Lambda module.""" from lambdacommon import common_config from lambdacommon import util from vector_utils.datadog.metrics import call_datadog_with_metric from src import config from src import exceptions from src import query from src.logger import get_current_logger class InvocationContext: """This is for passing around the invoke context to logger.""" correlation_id = None def handler(event, context): """Entrypoint to Lambda function. Args: event (optional): AWS Lambda event dependent structure with metadata. context (LambdaContext): AWS Lambda context. """ InvocationContext.correlation_id = context.aws_request_id log = get_current_logger(InvocationContext.correlation_id) log.info('Running "{}"'.format(config.LAMBDA_NAME)) with util.dd_connection(config.DD_MYSQL_CONN_INFO) as conn: with conn.cursor() as cursor: cursor.execute( query.SELECT_JOBS_TO_PROCESS, { 'encoder_ids': config.ENCODER_IDS.split(','), 'limit': config.DD_MYSQL_LIMIT } ) jobs = cursor.fetchall() jobs = get_data_from_ar(jobs) result = process_jobs(jobs) log.info('Finished "{}"'.format(config.LAMBDA_NAME)) return result def get_data_from_ar(jobs): """Get data from art relations. This will query art relations with applicable queries to enrich existing job pull from direct delivery. This function so far merges two lists of dictionaries based on UPC. Args: jobs (list(dict)): List of jobs from database Returns: list(dict): List of jobs from database but potentially enriched """ log = get_current_logger(InvocationContext.correlation_id) upc_list = [x['upc'] for x in jobs] if not upc_list: return jobs log.info('Retrieve data from AR') with util.ar_connection(config.AR_MYSQL_CONN_INFO) as conn: with conn.cursor() as cursor: cursor.execute( query.GET_RELEASE_BY_UPC, { 'upcs': upc_list } ) results = cursor.fetchall() for job in jobs: in_results = next( (x for x in results if int(x['upc']) == int(job['upc'])), None) if in_results: job.update(**in_results) return jobs def calculate_priority(job): """Prioritize job. Args: dict: Enriched job from database Returns: int: priority """ job_id = job['encoding_queue_detail_id'] try: date_diff = job['release_date_diff'] except KeyError: return 6 order_priority = job['order_priority'] if order_priority == 1: return 1 elif order_priority == 2: default_op_2 = 5 if date_diff is None: return default_op_2 date_range_priorities = { 14: 2, 60: 3, 120: 4 } return next( ( p for x, p in date_range_priorities.items() if -x <= date_diff <= x ), default_op_2 ) elif order_priority == 3: return 6 raise exceptions.PriorityCalculationFailed(job_id) def save_priorities(jobs): """Save job priorities. Args list[tuple]: List of jobs from database with priority """ log = get_current_logger(InvocationContext.correlation_id) log.info('Store matched rules') with util.dd_connection(config.DD_MYSQL_CONN_INFO) as conn: with conn.cursor() as cursor: cursor.executemany(query.INSERT_INTO_JOB_PRIORITY, jobs) conn.commit() def process_jobs(jobs): """Insert priority for jobs. Args: jobs (list): List of jobs from database Returns: dict: Dictionary with the number of jobs processed """ log = get_current_logger(InvocationContext.correlation_id) log.info('Start processing jobs') processed = [ ( x['encoding_queue_detail_id'], calculate_priority(x) ) for x in jobs ] save_priorities(processed) if config.DATADOG_API_KEY and config.DATADOG_APP_KEY: log.info('Send data to datadog') app_name = 'job_priority_rule_match' tags = [f'environment:{common_config.ENVIRONMENT}'] call_datadog_with_metric( f'{app_name}.jobs_to_process', tags, api_key=config.DATADOG_API_KEY, app_key=config.DATADOG_APP_KEY, value=len(jobs) ) call_datadog_with_metric( f'{app_name}.jobs_processed', tags, api_key=config.DATADOG_API_KEY, app_key=config.DATADOG_APP_KEY, value=len(processed) ) return { 'jobs': len(jobs), 'processed': len(processed) } if __name__ == '__main__': handler(None, None)