Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion kale/publisher.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,7 @@ def publish_messages_to_dead_letter_queue(self, dlq_name, messages):
:raises: SendMessagesException: SQS responded with a partial success. Some
messages were not delivered.
"""
sqs_dead_letter_queue = self._get_or_create_queue(dlq_name)
sqs_dead_letter_queue = self._get_or_create_queue(dlq_name, is_dlq=True)

response = sqs_dead_letter_queue.send_messages(
Entries=[{
Expand Down
7 changes: 4 additions & 3 deletions kale/sqs.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,12 +55,13 @@ def __init__(self, *args, **kwargs):
self._client = self._session.client('sqs', endpoint_url=endpoint_url)
self._sqs = self._session.resource('sqs', endpoint_url=endpoint_url)

def _get_or_create_queue(self, queue_name):
def _get_or_create_queue(self, queue_name, is_dlq=False):
"""Fetch or create a queue.

:param str queue_name: string for queue name.
:param bool is_dlq:True iff the queue is a dead letter queue and should be tagged as such
:return: Queue
:rtype: boto3.resources.factory.sqs.Queue
:rtype: boto3.resources.factory.sqs.Queue
"""

# Check local cache first.
Expand All @@ -76,7 +77,7 @@ def _get_or_create_queue(self, queue_name):
raise e

logger.info('Creating new SQS queue: %s' % queue_name)
queue = self._client.create_queue(QueueName=queue_name)
queue = self._client.create_queue(QueueName=queue_name, tags={"dlq":str(is_dlq)})

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

queue_url = queue.get('QueueUrl')

# create queue object
Expand Down