""" Amazon SQS connector ==================== """ import base64 import json import boto3 from ytownership import config from ytownership.connectors.loggly import get_current_logger from ytownership.utils.misc 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): """Stores 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 = 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(base64.b64decode(self._message.body).decode()) 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.AWS_REGION) def get_queue(queue_name, sqs_connection=None): """Get connection to Amazon SQS services""" connection = sqs_connection if sqs_connection else get_connection() queue = connection.get_queue_by_name(QueueName=queue_name) return queue