import time from typing import Any, TYPE_CHECKING import boto3 import botocore.exceptions from aws_testing_utils import config from .logger import log if TYPE_CHECKING: from mypy_boto3_sqs import SQSClient class SQSHandler: """Handles interactions with AWS SQS.""" sqs_client: 'SQSClient' def __init__(self) -> None: self.sqs_client = boto3.client('sqs', region_name=config.AWS_REGION) def get_url(self, sqs_name: str) -> str: """Returns the queue URL for the given queue name.""" log.info(f'Get url for: {sqs_name}') sqs_url = self.sqs_client.get_queue_url(QueueName=sqs_name)['QueueUrl'] return sqs_url def read_message(self, sqs_name: str) -> dict[str, Any]: """Receives and deletes the latest message from the queue. Raises ValueError if empty.""" sqs_url = self.get_url(sqs_name) log.info(f'Read message from: {sqs_url}') response = self.sqs_client.receive_message( QueueUrl=sqs_url, AttributeNames=['SentTimestamp'], MaxNumberOfMessages=1, MessageAttributeNames=['All'], VisibilityTimeout=0, WaitTimeSeconds=0, ) messages = response.get('Messages', []) if not messages: raise ValueError(f'No messages available in queue: {sqs_name}') message: dict[str, Any] = dict(messages[0]) log.info(f'Received message: {message}') receipt_handle = message['ReceiptHandle'] log.info('Delete received message from queue') self.sqs_client.delete_message(QueueUrl=sqs_url, ReceiptHandle=receipt_handle) return message def purge(self, sqs_name: str) -> None: """Purges all messages from the queue, retrying once if a purge is already in progress.""" log.info(f'Purging SQS: {sqs_name}') sqs_url = self.get_url(sqs_name) try: self.sqs_client.purge_queue(QueueUrl=sqs_url) except botocore.exceptions.ClientError as exception: if ( exception.response['Error']['Code'] == 'AWS.SimpleQueueService.PurgeQueueInProgress' ): log.debug(exception.response) log.info('Wait for 60s as PurgeQueueInProgress and retry') time.sleep(60) self.sqs_client.purge_queue(QueueUrl=sqs_url) else: raise exception