From e052c69a71cc7fef00fac494162dc54c706dfab4 Mon Sep 17 00:00:00 2001 From: Dima Anfimov Date: Wed, 16 Sep 2026 15:23:36 +0200 Subject: [PATCH] feat: support x-delivery-limit for quorum queues --- README.md | 22 +++++++ pyproject.toml | 2 +- taskiq_aio_pika/broker.py | 29 ++++++--- taskiq_aio_pika/queue.py | 2 + tests/test_poison_message.py | 112 +++++++++++++++++++++++++++++++++++ uv.lock | 18 ++---- 6 files changed, 163 insertions(+), 22 deletions(-) create mode 100644 tests/test_poison_message.py diff --git a/README.md b/README.md index 8393c93..6511b3c 100644 --- a/README.md +++ b/README.md @@ -122,6 +122,28 @@ async def main(): await prio_task.kicker().with_labels(priority=None).kiq() ``` +## Poison message handling + +Quorum queues (the default queue type used by this broker) keep track of how many times a message has been redelivered. This is useful for preventing "poison messages" — messages that crash the consumer (e.g. via an OOM) before it can ack, nack, or update a retry count — from being redelivered forever. + +Set `delivery_limit` on a `Queue` to have RabbitMQ automatically dead-letter a message once it has been redelivered too many times, instead of requeuing it indefinitely: + +```python +from taskiq_aio_pika import AioPikaBroker, Queue, QueueType + +broker = AioPikaBroker( + task_queues=[ + Queue( + name="taskiq", + type=QueueType.QUORUM, + delivery_limit=5, + ), + ], +) +``` + +Once a message has been redelivered more than `delivery_limit` times, RabbitMQ dead-letters it to the broker's dead-letter queue instead of redelivering it again — no application code involved. `delivery_limit` is only supported by quorum queues. See the [RabbitMQ docs](https://www.rabbitmq.com/docs/quorum-queues#poison-message-handling) for details. + ## Custom Queue and Exchange arguments You can pass custom arguments to the underlying RabbitMQ queues and exchange declaration by using the `Queue`/`Exchange` classes from `taskiq_aio_pika`. If you used `faststream` before you are probably familiar with this concept. diff --git a/pyproject.toml b/pyproject.toml index 17e8a4a..52db81e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -31,7 +31,7 @@ repository = "https://github.com/taskiq-python/taskiq-aio-pika" keywords = ["taskiq", "tasks", "distributed", "async", "aio-pika"] requires-python = ">=3.10,<4" dependencies = [ - "taskiq>=0.12.0,<1", + "taskiq>=0.12.6,<1", "aio-pika>=9.0.0", "typing-extensions>=4.14.0", ] diff --git a/taskiq_aio_pika/broker.py b/taskiq_aio_pika/broker.py index 764ccdc..e113f48 100644 --- a/taskiq_aio_pika/broker.py +++ b/taskiq_aio_pika/broker.py @@ -18,7 +18,7 @@ QueueNotDeclaredError, ) from taskiq_aio_pika.exchange import Exchange -from taskiq_aio_pika.queue import Queue +from taskiq_aio_pika.queue import Queue, QueueType from taskiq_aio_pika.utils import merge_async_iterables _T = TypeVar("_T") @@ -201,10 +201,9 @@ async def _declare_dead_letter_queue( ) -> None: if self._dead_letter_queue.declare: dead_letter_queue_arguments = self._dead_letter_queue.arguments.copy() - if self._dead_letter_queue.max_priority is not None: - dead_letter_queue_arguments["x-max-priority"] = ( - self._dead_letter_queue.max_priority - ) + dead_letter_queue_arguments.update( + self._optional_queue_arguments(self._dead_letter_queue), + ) dead_letter_queue_arguments["x-queue-type"] = ( self._dead_letter_queue.type.value ) @@ -229,6 +228,23 @@ async def _declare_dead_letter_queue( f"was not declared and does not exist.", ) from error + @staticmethod + def _optional_queue_arguments(queue: Queue) -> FieldTable: + """Build queue arguments that are only set when explicitly configured.""" + arguments: FieldTable = {} + if queue.max_priority is not None: + arguments["x-max-priority"] = queue.max_priority + if queue.delivery_limit is not None: + if queue.type != QueueType.QUORUM: + logger.warning( + "delivery_limit is set for queue '%s', but it is only supported by quorum queues " + "(queue type is '%s'); it will be ignored by RabbitMQ.", + queue.name, + queue.type.value, + ) + arguments["x-delivery-limit"] = queue.delivery_limit + return arguments + async def _declare_queues( self, channel: AbstractChannel, @@ -259,8 +275,7 @@ async def _declare_queues( for queue in filter(lambda queue: queue.declare, queues): per_queue_arguments: FieldTable = queue_default_arguments.copy() - if queue.max_priority is not None: - per_queue_arguments["x-max-priority"] = queue.max_priority + per_queue_arguments.update(self._optional_queue_arguments(queue)) per_queue_arguments["x-queue-type"] = queue.type.value if self._delay_queue and queue.name == self._delay_queue.name: per_queue_arguments["x-dead-letter-exchange"] = self._exchange.name diff --git a/taskiq_aio_pika/queue.py b/taskiq_aio_pika/queue.py index e3beafc..6db285a 100644 --- a/taskiq_aio_pika/queue.py +++ b/taskiq_aio_pika/queue.py @@ -26,6 +26,7 @@ class Queue: passive: Whether to check if the queue exists without creating it. auto_delete: Whether the queue should be auto-deleted. max_priority: The maximum priority for the queue. + delivery_limit: The maximum number of delivery attempts for a message before it is dead-lettered. arguments: Additional arguments for the queue declaration. timeout: Timeout for queue declaration. routing_key: The routing key for the queue. @@ -46,6 +47,7 @@ class Queue: passive: bool = False auto_delete: bool = False max_priority: int | None = None + delivery_limit: int | None = None arguments: FieldTable = field(default_factory=dict) timeout: int | float | None = None diff --git a/tests/test_poison_message.py b/tests/test_poison_message.py new file mode 100644 index 0000000..7affbf3 --- /dev/null +++ b/tests/test_poison_message.py @@ -0,0 +1,112 @@ +import aiormq +import pytest +from aio_pika import Channel +from aio_pika.exceptions import QueueEmpty +from taskiq import BrokerMessage + +from taskiq_aio_pika import AioPikaBroker +from taskiq_aio_pika.exchange import Exchange +from taskiq_aio_pika.queue import Queue, QueueType +from tests.conftest import _cleanup_amqp_resources + + +async def test_when_delivery_limit_exceeded__message_is_dead_lettered( + amqp_url: str, + test_channel: Channel, + queue_name: str, + dead_queue_name: str, + exchange_name: str, +) -> None: + # given + broker = AioPikaBroker( + url=amqp_url, + exchange=Exchange(name=exchange_name, declare=True), + dead_letter_queue=Queue(name=dead_queue_name, declare=True), + task_queues=[ + Queue( + name=queue_name, + declare=True, + type=QueueType.QUORUM, + delivery_limit=2, + ), + ], + ) + + try: + await broker.startup() + main_queue = await test_channel.get_queue(queue_name) + dead_letter_queue = await test_channel.get_queue(dead_queue_name) + + await broker.kick( + BrokerMessage( + task_id="1", + task_name="name", + message=b"poison", + labels={}, + ), + ) + + # when + for _ in range( + 3, + ): # simulate a consumer that never acks the message, so it gets requeued + message = await main_queue.get() + await message.nack(requeue=True) + + # then + with pytest.raises(QueueEmpty): + await main_queue.get() + + dead_lettered_message = await dead_letter_queue.get() + assert dead_lettered_message.body == b"poison" + finally: + await broker.shutdown() + await _cleanup_amqp_resources( + amqp_url, + [exchange_name], + [queue_name, dead_queue_name], + ) + + +async def test_when_delivery_limit_set_on_classic_queue__warning_is_logged( + amqp_url: str, + queue_name: str, + dead_queue_name: str, + exchange_name: str, + caplog: pytest.LogCaptureFixture, +) -> None: + # given + broker = AioPikaBroker( + url=amqp_url, + exchange=Exchange(name=exchange_name, declare=True), + dead_letter_queue=Queue(name=dead_queue_name, declare=True), + task_queues=[ + Queue( + name=queue_name, + declare=True, + type=QueueType.CLASSIC, + delivery_limit=2, + ), + ], + ) + + try: + # when + with ( + caplog.at_level("WARNING", logger="taskiq.aio_pika_broker"), + pytest.raises(aiormq.exceptions.ChannelPreconditionFailed), + ): + await broker.startup() + + # then + assert any( + "delivery_limit" in record.message and queue_name in record.message + for record in caplog.records + ) + finally: + await broker.shutdown() + await _cleanup_amqp_resources( + amqp_url, + [exchange_name], + [queue_name, dead_queue_name], + ) diff --git a/uv.lock b/uv.lock index 2c403a2..ac9c6ac 100644 --- a/uv.lock +++ b/uv.lock @@ -537,15 +537,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/cb/b1/3846dd7f199d53cb17f49cba7e651e9ce294d8497c8c150530ed11865bb8/iniconfig-2.3.0-py3-none-any.whl", hash = "sha256:f631c04d2c48c52b84d0d0549c99ff3859c98df65b3101406327ecc7d53fbf12", size = 7484, upload-time = "2025-10-18T21:55:41.639Z" }, ] -[[package]] -name = "izulu" -version = "0.50.0" -source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/d0/58/6d6335c78b7ade54d8a6c6dbaa589e5c21b3fd916341d5a16f774c72652a/izulu-0.50.0.tar.gz", hash = "sha256:cc8e252d5e8560c70b95380295008eeb0786f7b745a405a40d3556ab3252d5f5", size = 48558, upload-time = "2025-03-24T15:52:21.51Z" } -wheels = [ - { url = "https://files.pythonhosted.org/packages/4a/9f/bf9d33546bbb6e5e80ebafe46f90b7d8b4a77410b7b05160b0ca8978c15a/izulu-0.50.0-py3-none-any.whl", hash = "sha256:4e9ae2508844e7c5f62c468a8b9e2deba2f60325ef63f01e65b39fd9a6b3fab4", size = 18095, upload-time = "2025-03-24T15:52:19.667Z" }, -] - [[package]] name = "multidict" version = "6.7.0" @@ -1161,21 +1152,20 @@ wheels = [ [[package]] name = "taskiq" -version = "0.12.0" +version = "0.12.6" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "aiohttp" }, { name = "anyio" }, - { name = "izulu" }, { name = "packaging" }, { name = "pycron" }, { name = "pydantic" }, { name = "taskiq-dependencies" }, { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/67/9b/bb9b3ab5024051e80013170950bf5acb7729636918bab0ceb91a3900c815/taskiq-0.12.0.tar.gz", hash = "sha256:722d64b8176affb146635c7ac356d3e44efa446a3dc7a373694c0eb8852672b6", size = 60099, upload-time = "2025-11-26T18:38:54.618Z" } +sdist = { url = "https://files.pythonhosted.org/packages/dd/6b/a7df3a5d1e7bbfe75a41f378f71c16e89ac19fbab931f8e8d76a8e339357/taskiq-0.12.6.tar.gz", hash = "sha256:4667db6366cfda24d23f0e34151b118d71b0907229f8ff00db9dbc2465f7baf8", size = 408574, upload-time = "2026-08-29T08:07:44.052Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/37/88/e0bb05fcca198d313a50c5c461711a6f36d3a8d29010b65960796bdb03cd/taskiq-0.12.0-py3-none-any.whl", hash = "sha256:8fea577bbf72ceabd77338f643510c64787521f43e7907e605ebeead660c1a74", size = 90388, upload-time = "2025-11-26T18:38:53.4Z" }, + { url = "https://files.pythonhosted.org/packages/68/89/949d2d26ae24a5282b37d39e7e05b714f960840166458c7bc78cfade83b5/taskiq-0.12.6-py3-none-any.whl", hash = "sha256:fa8c5071ecafc6a1891e5ede1dc580093d99bccaf0d258b1848d2e5f35d22250", size = 94895, upload-time = "2026-08-29T08:07:42.591Z" }, ] [[package]] @@ -1222,7 +1212,7 @@ typecheck = [ [package.metadata] requires-dist = [ { name = "aio-pika", specifier = ">=9.0.0" }, - { name = "taskiq", specifier = ">=0.12.0,<1" }, + { name = "taskiq", specifier = ">=0.12.6,<1" }, { name = "typing-extensions", specifier = ">=4.14.0" }, ]