"""Lambda function module.""" import boto3 import sqlalchemy from vector_utils.aws_utils import sqs import config from src import logger, sql_queries from src.connectors import art_relations def select_and_add_upcs_to_sqs(event_detail_type, log=None): """Select unused upcs and add them to upc_repository-queue.""" sqs_client = boto3.client('sqs', region_name='us-east-1') response = sqs_client.get_queue_attributes( QueueUrl=config.QUEUE_NAME_UPC_PROVISIONER, AttributeNames=['ApproximateNumberOfMessages'] ) current_messages_in_queue = int( response['Attributes']['ApproximateNumberOfMessages'] ) # If there are more upcs in queue then our threshold we want to set to 0 instead of -n needed_upcs = max( config.UPCS_AR_SELECT_LIMIT_ALARM - current_messages_in_queue, 0 ) if needed_upcs == 0: log.info('No upcs needed') return with art_relations.session_scope() as session: result_select = session.execute( sqlalchemy.text(sql_queries.SELECT_UPCS), { 'limit': needed_upcs, 'offset': config.UPCS_AR_SELECT_OFFSET }) upc_list = [row['upc'] for row in result_select] session.execute( sqlalchemy.text(sql_queries.UPDATE_STATUS), { 'list': upc_list }) log.info(f'Number of UPCs claimed: {len(upc_list)}') return sqs.add_messages_to_sqs( config.QUEUE_NAME_UPC_PROVISIONER, upc_list) def handler(event, context): """Lambda entry point.""" log = logger.get_current_logger() event_detail_type = event.get('detail-type') sqs_response = select_and_add_upcs_to_sqs(event_detail_type, log) log.info('SQS response after pushing to queue: ', sqs_response) return sqs_response