mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 08:08:13 +08:00
In some environments permissions for an IAM role or access key might not include "SQS create". The broker should first try to get the queue and if this fails try to create it as fall back to the current default behaviour.
79 lines
2.3 KiB
Python
79 lines
2.3 KiB
Python
from boto3 import Session
|
|
from botocore.client import ClientError
|
|
|
|
from django_q.brokers import Broker
|
|
from django_q.conf import Conf
|
|
|
|
|
|
QUEUE_DOES_NOT_EXIST = "AWS.SimpleQueueService.NonExistentQueue"
|
|
|
|
|
|
class Sqs(Broker):
|
|
def __init__(self, list_key: str = Conf.PREFIX):
|
|
self.sqs = None
|
|
super(Sqs, self).__init__(list_key)
|
|
self.queue = self.get_queue()
|
|
|
|
def enqueue(self, task):
|
|
response = self.queue.send_message(MessageBody=task)
|
|
return response.get("MessageId")
|
|
|
|
def dequeue(self):
|
|
# sqs supports max 10 messages in bulk
|
|
if Conf.BULK > 10:
|
|
Conf.BULK = 10
|
|
tasks = self.queue.receive_messages(
|
|
MaxNumberOfMessages=Conf.BULK, VisibilityTimeout=Conf.RETRY
|
|
)
|
|
if tasks:
|
|
return [(t.receipt_handle, t.body) for t in tasks]
|
|
|
|
def acknowledge(self, task_id):
|
|
return self.delete(task_id)
|
|
|
|
def queue_size(self) -> int:
|
|
return int(self.queue.attributes["ApproximateNumberOfMessages"])
|
|
|
|
def lock_size(self) -> int:
|
|
return int(self.queue.attributes["ApproximateNumberOfMessagesNotVisible"])
|
|
|
|
def delete(self, task_id):
|
|
message = self.sqs.Message(self.queue.url, task_id)
|
|
message.delete()
|
|
|
|
def fail(self, task_id):
|
|
self.delete(task_id)
|
|
|
|
def delete_queue(self):
|
|
self.queue.delete()
|
|
|
|
def purge_queue(self):
|
|
self.queue.purge()
|
|
|
|
def ping(self) -> bool:
|
|
return "sqs" in self.connection.get_available_resources()
|
|
|
|
def info(self) -> str:
|
|
return "AWS SQS"
|
|
|
|
@staticmethod
|
|
def get_connection(list_key: str = Conf.PREFIX) -> Session:
|
|
config = Conf.SQS
|
|
if "aws_region" in config:
|
|
config["region_name"] = config["aws_region"]
|
|
del config["aws_region"]
|
|
return Session(**config)
|
|
|
|
def get_queue(self):
|
|
self.sqs = self.connection.resource("sqs")
|
|
|
|
try:
|
|
# try to return an existing queue by name. If the queue does not
|
|
# exist try to create it.
|
|
return self.sqs.get_queue_by_name(QueueName=self.list_key)
|
|
except ClientError as exp:
|
|
if not exp.response["Error"]["Code"] == QUEUE_DOES_NOT_EXIST:
|
|
raise exp
|
|
|
|
return self.sqs.create_queue(QueueName=self.list_key)
|