"""Connector for SQS service.""" import copy import decimal import json import uuid from botocore import exceptions from accounting.flows.reserve_payouts import connectors from accounting.flows.reserve_payouts import setting from accounting.flows.reserve_payouts.constants import sqs as sqs_constants from accounting.util import sqs def health_check(): """Perform a simple query to do the health check. Returns: namedtuple: with bool and message attributes, (True, '') if connection is ok, (False, 'Error message') otherwise. """ try: queue = sqs.get_queue_by_name(setting.SQS_QUEUE) sqs.get_queue_number_of_messages(queue) except (exceptions.BotoCoreError, exceptions.ClientError) as e: result = connectors.HealthCheckResult(False, str(e)) else: result = connectors.HealthCheckResult(True, '') return result def get_queue_message_count(): """Get the approximate number of messages in the queue. Returns: int: Approximate number of messages """ queue = sqs.get_queue_by_name(setting.SQS_QUEUE) queue.reload() return sqs.get_queue_number_of_messages(queue) def construct_vendor_message(vendor_data): """Construct a message for the SQS queue. Args: vendor_data (dict): message payload (vendor specific data) Returns: dict: queue.send_message suitable dict """ json_valid_vendor_data = copy.deepcopy(vendor_data) for field_name, field_value in json_valid_vendor_data.items(): if type(field_value) == decimal.Decimal: json_valid_vendor_data[field_name] = float(field_value) msg = { sqs_constants.SQS_MESSAGE_BODY: json.dumps(json_valid_vendor_data), sqs_constants.SQS_MESSAGE_ID: str(uuid.uuid4()), } return msg def batch_vendor_send_messages(message_list): """Send vendor data messages to SQS in batch. Args: message_list (list): vendor transactions sum and contract terms """ queue = sqs.get_queue_by_name(setting.SQS_QUEUE) sublists_generator = sqs.yield_sublists( message_list, setting.SQS_MAX_BATCH_SIZE) for sublist in sublists_generator: res = queue.send_messages(Entries=sublist) # according to boto3 documentation, sqs can contain failed items if 'Failed' in res: response_str = json.dumps(res) msg = 'Failures detected in SQS response: {}'.format(response_str) raise Exception(msg) def get_message(): """Receive a message from the queue. Returns: boto3.resources.factory.sqs.Message: SQS message or None """ queue = sqs.get_queue_by_name(setting.SQS_QUEUE) queue.set_attributes( Attributes={'VisibilityTimeout': str(setting.SQS_VISIBILITY_TIMEOUT)}) return sqs.get_message(queue)