"""Throttler Lambda module.""" from vector_utils import queries from vector_utils.aws_utils import sqs from vector_utils.datadog.metrics import call_datadog_with_metric from src.logger import get_current_logger import config 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) dms_id = int(event['dms_id']) limit = int(event.get('limit', 1000)) allowed_encoders = event.get('allowed_encoders') or [17, 18, 19, 23] log.info('Running "{}" for dms_id {}'.format(config.LAMBDA_NAME, dms_id)) dd_tags = [f'dms_id:{dms_id}', f'environment:{config.ENVIRONMENT}'] dd_app = 'throttler' call_datadog_with_metric( f'{dd_app}.attempt', dd_tags, api_key=config.DATADOG_API_KEY, app_key=config.DATADOG_APP_KEY ) open_slots = { 'release': 0, 'ringtone': 0, 'harddrive': 0, 'track': 0 } total_slots = 0 for order_type, open_slot in open_slots.items(): result = queries.get_connection_limit( dms_id, order_type, conn_info=config.RO_DD_MYSQL_CONN_INFO) if result: open_slots[order_type] = result['connection_limit'] - result['current_connections'] # noqa total_slots += open_slots[order_type] log.info('Total slots {}'.format(total_slots)) log.info('Open slots {}'.format(open_slots)) if total_slots < 1: log.info('No more connections left for DMS: {}.'.format(dms_id)) else: results = queries.get_jobs_for_batch_delivery( dms_id, allowed_encoders, conn_info=config.RO_DD_MYSQL_CONN_INFO) if results: log.info('Process batch jobs') process_jobs(dms_id, results, open_slots) results = queries.get_jobs_for_nonbatch_delivery( dms_id, limit, allowed_encoders, conn_info=config.RO_DD_MYSQL_CONN_INFO ) if results: log.info('Process non batch jobs') process_jobs(dms_id, results, open_slots) call_datadog_with_metric( f'{dd_app}.run', dd_tags, api_key=config.DATADOG_API_KEY, app_key=config.DATADOG_APP_KEY ) log.info('Finished "{}" for dms_id {}'.format(config.LAMBDA_NAME, dms_id)) def process_jobs(dms_id, results, open_slots): """Process jobs. Args: dms_id (int): results (dict): open_slots (dict): """ log = get_current_logger(InvocationContext.correlation_id) jobs = process_jobs_data_format(results) for priority, job_details in jobs.items(): for order_type, eqd_ids in job_details.items(): for eqd_id, eqd_details in eqd_ids.items(): log.info('open slots {}'.format(open_slots)) if open_slots[order_type] < 1: log.info('break') break affected_rows = queries.set_job_to_queued_for_delivery( eqd_id, conn_info=config.DD_MYSQL_CONN_INFO) if affected_rows: queue_name = sqs.format_queue_name( sqs.DELIVERY_QUEUE_NAME_PATTERN, encoder_id=eqd_details['encoder_id'], dms_priority=1, priority=priority, env=config.ENVIRONMENT ) log.info('pushing job to queue {}'.format(queue_name)) sqs.add_message_to_sqs(queue_name, { 'eqd_id': eqd_id, 'dms_id': dms_id, 'delivery_batch_id': eqd_details['delivery_batch_id'] }) open_slots[order_type] -= 1 def process_jobs_data_format(results): """Process jobs data formatting helper. Args: results (dict): A result set from SQL Returns: dict """ jobs = {} for result in results: priority = result['priority'] encoding_order_type = result['encoding_order_type'] if priority not in jobs: jobs[priority] = {} if encoding_order_type not in jobs[priority]: jobs[priority][encoding_order_type] = {} jobs[priority][encoding_order_type].update({ result['encoding_queue_detail_id']: { 'encoder_id': result['encoder_id'], 'delivery_batch_id': result['delivery_batch_id'], 'dms_priority': result['dms_priority'] } }) return jobs if __name__ == '__main__': handler(None, None)