"""Connector for Amazon Simple Queue Service operations.""" import json from boto3 import resource from moneyhub.constants.constants import CORRELATION_ID_HEADER MAX_BATCH_SIZE = 10 def send_message(queue_name: str, correlation_id: str | None, params: dict) -> dict: """Send a message to SQS. Args: queue_name (str): Name of the queue to send the message to correlation_id (str): correlation_id in the header params (dict[str, str]): The message data to put in the queue Returns: Response: object containing sqs response """ queue = resource('sqs').get_queue_by_name(QueueName=queue_name) if correlation_id: params[CORRELATION_ID_HEADER] = f'{correlation_id}.1' response = queue.send_message(MessageBody=json.dumps(params, sort_keys=True)) return response def send_messages(queue_name: str, correlation_id: str | None, messages: list): """Send several messages to the SQS queue. Args: queue_name (str): Name of the queue to send the message to correlation_id (str): correlation_id in the header messages (list): The message data to put in the queue """ queue = resource('sqs').get_queue_by_name(QueueName=queue_name) chunks = [messages[x:x + MAX_BATCH_SIZE] for x in range(0, len(messages), MAX_BATCH_SIZE)] counter = 1 failures = [] for chunk in chunks: entries = [] for params in chunk: if correlation_id: params[CORRELATION_ID_HEADER] = f'{correlation_id}.1' entry = {'Id': str(counter), 'MessageBody': json.dumps(params, sort_keys=True)} entries.append(entry) counter += 1 response = queue.send_messages(Entries=entries) if 'Failed' in response: for failure in response['Failed']: failures.append(failure['Message']) if failures: message = '\n'.join([ f'Failed to send {len(failures)} (of {len(messages)}) SQS messages:', *failures, ]) raise Exception(message)