"""Amazon SQS connector.""" import json import boto3 from sentry_sdk import capture_exception from availability import config from availability.connectors import loggly from availability.utils import DotDict CORRELATION_ID_ATTRIBUTE = 'Correlation-Id' APPROXIMATE_RECEIVE_COUNT_ATTRIBUTE = 'ApproximateReceiveCount' class JSONMessageExt(): """Wrapper for SQS Messages that provides a logger with Correlation-Id.""" def __init__(self, message=None): """Create new instance of JSONMessageExt.""" self._message = message self._context = DotDict() @property def context(self): """Store an additional context for message (e.g. logging adapter). Returns: DotDict: an object that holds message context """ return self._context @property def correlation_id(self): """Get Correlation-Id extracted from message attributes. Returns: str: correlation Id """ return self._message.message_attributes.get( CORRELATION_ID_ATTRIBUTE, {}).get('StringValue') @property def logger(self): """Get logging adapter for current message. Returns: OwsLoggingAdapter: Logging adapter with Correlation-Id, extracted from SQS message attributes """ if not self.context.logger: self.context.logger = loggly.get_current_logger( self.correlation_id) return self.context.logger @property def receive_count(self): """Get approximate message receive count. Returns: int: approximate number of times the message has been received """ count = self._message.attributes.get( APPROXIMATE_RECEIVE_COUNT_ATTRIBUTE, 0) return int(count) @property def message(self): """Get the message object.""" return self._message def get_body(self): """Get the message body.""" return json.loads(self._message.body) def delete(self): """Delete the message from SQS queue.""" return self._message.delete() def change_visibility(self, queue_timeout): """Set visibility of message in SQS queue.""" return self._message.change_visibility(VisibilityTimeout=queue_timeout) def get_connection(): """Get connection to Amazon SQS services.""" return boto3.resource('sqs', region_name=config.SQS_DEFAULT_REGION) def get_queue(queue_name, sqs_connection=None): """Get Amazon SQS queue. Args: queue_name (str): SQS queue name sqs_connection (obj): SQS connection object Returns: Queue: Instance of Amazon SQS queue """ try: connection = sqs_connection if sqs_connection else get_connection() queue = connection.get_queue_by_name(QueueName=queue_name) if queue: queue.set_attributes( Attributes={ 'MessageRetentionPeriod': str( config.SQS_MESSAGE_RETENTION_PERIOD_SECONDS) } ) return queue except Exception as e: capture_exception(e) return None