diff --git a/kale/publisher.py b/kale/publisher.py index 94e8f0c..cb601fb 100644 --- a/kale/publisher.py +++ b/kale/publisher.py @@ -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=[{ diff --git a/kale/sqs.py b/kale/sqs.py index bc1b718..4e40f8e 100644 --- a/kale/sqs.py +++ b/kale/sqs.py @@ -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. @@ -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)}) queue_url = queue.get('QueueUrl') # create queue object diff --git a/kale/version.py b/kale/version.py index 383582f..b9c2f2a 100644 --- a/kale/version.py +++ b/kale/version.py @@ -12,4 +12,4 @@ # See the License for the specific language governing permissions and # limitations under the License. -__version__ = '2.2.3' # http://semver.org/ +__version__ = '2.2.4' # http://semver.org/