import json from typing import Dict, List, Tuple from notification_service.utils import chunks, send_metric_by_dimension import boto3 from boto3_type_annotations.sns import Client as SNSClient from boto3_type_annotations.sqs import Client as SQSClient from smelog.factory import SmeBoundLogger class SQSService: sqs: SQSClient = boto3.client('sqs') def __init__(self, queue: str, logger: SmeBoundLogger, max_msg_num: int = 10) -> None: self._queue = queue self._logger = logger self.max_msg_num = max_msg_num def receive_msgs(self) -> Dict: bodies = {} try: msgs = self.sqs.receive_message( QueueUrl=self._queue, MaxNumberOfMessages=self.max_msg_num, WaitTimeSeconds=1 ) except Exception as exc: # pylint: disable=broad-except self._logger.exception('SQS Service receive exception.', exc=str(exc)) else: for msg in msgs.get('Messages', []): bodies[msg['ReceiptHandle']] = json.loads(msg['Body']) return bodies def remove_msg(self, receipt_handle: str) -> None: try: self.sqs.delete_message(QueueUrl=self._queue, ReceiptHandle=receipt_handle) except Exception as exc: # pylint: disable=broad-except self._logger.exception('SQS Service delete msg exception.', exc=str(exc)) class SNSService: sns: SNSClient = boto3.client('sns') max_msgs_count = 10 def __init__(self, sns_topic_arn: str, logger: SmeBoundLogger) -> None: self.sns_topic_arn = sns_topic_arn self._logger = logger def send_msgs(self, msgs: List[Dict]) -> Tuple: success_sent = [] failed_sent = [] for msg in msgs: try: self.sns.publish( TopicArn=self.sns_topic_arn, MessageStructure='string', Message=json.dumps(msg) ) except Exception as exc: # pylint: disable=broad-except self._logger.exception('SNS Service send exception.', exc=str(exc), msg=msg) failed_sent.append(msg) else: success_sent.append(msg) return success_sent, failed_sent def send_batch_msgs(self, msgs: List, job_index: int) -> Tuple: success_sent = 0 failed_sent = 0 self._logger.info( 'The process of sending data to SNS was successfully started.', job_index=job_index, ) for chunk in chunks(msgs, self.max_msgs_count): batch = [] for msg in chunk: send_metric_by_dimension(msg=msg) batch.append( { 'Id': str(msgs.index(msg)), 'MessageStructure': 'string', 'Message': json.dumps(msg) } ) try: self.sns.publish_batch( TopicArn=self.sns_topic_arn, PublishBatchRequestEntries=batch ) except Exception as exc: # pylint: disable=broad-except self._logger.exception(exc=exc) failed_sent += len(batch) else: success_sent += len(batch) self._logger.info( 'The process of sending data to SNS was successfully finished.', job_index=job_index, ) return success_sent, failed_sent