"""Sqs utils.""" import json import uuid import boto3 from vector_utils.exceptions import ValidationError # make sure its a total of 7 digits and 0 prefix padding with :07d DELIVERY_QUEUE_NAME_PATTERN = '{env}-delivery{encoder_id}_e{priority:07d}_d{dms_priority:07d}' # noqa ENCODING_QUEUE_NAME_PATTERN = '{env}-encoding{encoder_id}_e{priority:07d}_d{dms_priority:07d}' # noqa BATCH_TO_CLOSE_QUEUE_NAME_PATTERN = '{env}-batch_to_close_v2-queue' ASSETS_TO_DOWNLOAD_QUEUE_NAME_PATTERN = '{env}-assets_to_download-queue' # noqa ASSETS_TO_COPY_QUEUE_NAME_PATTERN = '{env}-assets_to_copy_d{loco_id:07d}-queue' GRAS_DELIVERY_QUEUE_NAME_PATTERN = '{env}-gras_delivery-queue' def format_queue_name(queue_pattern, encoder_id=None, dms_priority=None, priority=None, env='dev', loco_id=None): """Format the queue name. This takes in a required constant from above and the corresponding kwargs. Args: queue_pattern (str): This should be a const from above encoder_id (int): The encoder id of the job dms_priority (int): Dms priority in the job priority (int): Priority value of the job env (str): dev/qa/prod loco_id (int): Physical location ID Returns: str: The queue completely formatted Raises: ValidationError: If you send missing arguments or invalid const you will see this. """ if queue_pattern in [DELIVERY_QUEUE_NAME_PATTERN, ENCODING_QUEUE_NAME_PATTERN]: # noqa if not all([encoder_id, dms_priority, env, priority]): raise ValidationError(errors={ 'required_arguments': [ 'encoder_id', 'dms_priority', 'env', 'priority' ] }) return queue_pattern.format( env=env, encoder_id=encoder_id, priority=int(priority), dms_priority=int(dms_priority) ) elif queue_pattern == BATCH_TO_CLOSE_QUEUE_NAME_PATTERN: return queue_pattern.format(env=env) elif queue_pattern == ASSETS_TO_DOWNLOAD_QUEUE_NAME_PATTERN: return queue_pattern.format(env=env) elif queue_pattern == ASSETS_TO_COPY_QUEUE_NAME_PATTERN: if not loco_id: raise ValidationError(errors={ 'required_arguments': [ 'loco_id' ] }) return queue_pattern.format(env=env, loco_id=loco_id) else: raise ValidationError(errors={ 'pattern_must_be_one_of': [ 'DELIVERY_QUEUE_NAME_PATTERN', 'ENCODING_QUEUE_NAME_PATTERN', 'BATCH_TO_CLOSE_QUEUE_NAME_PATTERN', 'ASSETS_TO_DOWNLOAD_QUEUE_NAME_PATTERN', ] }) def add_message_to_sqs(queue_name, payload): """Add message to sqs. Args: queue_name (str): The name of the queue payload (dict|str): The payload to send to SQS Returns: dict """ if type(payload) in [dict, list]: payload = json.dumps(payload) return add_messages_to_sqs( queue_name, [payload], 1) def add_messages_to_sqs(queue_name, payload, batch_size=10): """Add batch message to sqs. Args: queue_name (str): The name of the queue payload (list): List of data to send to SQS batch_size (int): Number of message per batch Returns: dict: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/sqs.html#SQS.Client.send_message_batch # noqa """ if batch_size > 10: raise ValueError('Batch size cannot be more than 10!') sqs_client = boto3.client('sqs', region_name='us-east-1') try: sqs_queue = sqs_client.get_queue_url(QueueName=queue_name) except sqs_client.exceptions.QueueDoesNotExist: sqs_queue = sqs_client.create_queue(QueueName=queue_name) results = { 'Successful': [], 'Failed': [], } chunks = [payload[i:i + batch_size] for i in range( 0, len(payload), batch_size)] for chunk in chunks: messages = [{ 'Id': str(uuid.uuid1()), 'MessageBody': json.dumps( piece) if isinstance(piece, dict) else str(piece) } for piece in chunk] result = sqs_client.send_message_batch( QueueUrl=sqs_queue['QueueUrl'], Entries=messages ) results['Successful'] += result.get('Successful', []) results['Failed'] += result.get('Failed', []) return results def get_messages_from_sqs(queue_name, number_of_messages=1): """Get message to sqs. Args: queue_name (str): The name of the queue number_of_messages (int): The number of messages to receive up to 10 Returns: list: A list of https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/sqs.html#message """ sqs_client = boto3.client('sqs', region_name='us-east-1') try: sqs_queue = sqs_client.get_queue_url(QueueName=queue_name) except sqs_client.exceptions.QueueDoesNotExist: sqs_queue = sqs_client.create_queue(QueueName=queue_name) messages = sqs_client.receive_message( QueueUrl=sqs_queue['QueueUrl'], MaxNumberOfMessages=number_of_messages ) if 'Messages' in messages: return messages['Messages'] return [] def delete_messages_from_sqs(queue_name, receipt_handles): """Add message to sqs. Args: queue_name (str): The name of the queue receipt_handles (list): A list of receipt handles of the messages Returns: dict: https://boto3.amazonaws.com/v1/documentation/api/latest/reference/services/sqs.html#SQS.Queue.delete_messages # noqa """ sqs_client = boto3.client('sqs', region_name='us-east-1') try: sqs_queue = sqs_client.get_queue_url(QueueName=queue_name) except sqs_client.exceptions.QueueDoesNotExist: sqs_queue = sqs_client.create_queue(QueueName=queue_name) deleted = [] for rh in receipt_handles: result = sqs_client.delete_message( QueueUrl=sqs_queue['QueueUrl'], ReceiptHandle=rh) deleted.append(result) return deleted