Skip to content

Commit a6c2e4f

Browse files
committed
refactor: rename delay label to x_delay to avoid conflict with SmartRetryMiddleware
1 parent cf89301 commit a6c2e4f

5 files changed

Lines changed: 15 additions & 15 deletions

File tree

README.md

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -39,14 +39,14 @@ broker = AioPikaBroker(
3939
)
4040
```
4141

42-
After that you have to specify delay label. You can do it with `task` decorator, or by using kicker.
42+
After that you have to specify x_delay label. You can do it with `task` decorator, or by using kicker.
4343

4444
In this type of delay we are using additional queue with `expiration` parameter. After declared time message will be deleted from `delay` queue and sent to the main queue. For example:
4545

4646
```python
4747
broker = AioPikaBroker(...)
4848

49-
@broker.task(delay=3)
49+
@broker.task(x_delay=3)
5050
async def delayed_task() -> int:
5151
return 1
5252

@@ -57,11 +57,11 @@ async def main():
5757
await delayed_task.kiq()
5858

5959
# This message is going to be received after the delay in 4 seconds.
60-
# Since we overridden the `delay` label using kicker.
61-
await delayed_task.kicker().with_labels(delay=4).kiq()
60+
# Since we overridden the `x_delay` label using kicker.
61+
await delayed_task.kicker().with_labels(x_delay=4).kiq()
6262

6363
# This message is going to be send immediately. Since we deleted the label.
64-
await delayed_task.kicker().with_labels(delay=None).kiq()
64+
await delayed_task.kicker().with_labels(x_delay=None).kiq()
6565

6666
# Of course the delay is managed by rabbitmq, so you don't
6767
# have to wait delay period before message is going to be sent.
@@ -80,7 +80,7 @@ broker = AioPikaBroker(
8080
delayed_message_exchange_plugin=True,
8181
)
8282

83-
@broker.task(delay=3)
83+
@broker.task(x_delay=3)
8484
async def delayed_task() -> int:
8585
return 1
8686

@@ -91,8 +91,8 @@ async def main():
9191
await delayed_task.kiq()
9292

9393
# This message is going to be received after the delay in 4 seconds.
94-
# Since we overridden the `delay` label using kicker.
95-
await delayed_task.kicker().with_labels(delay=4).kiq()
94+
# Since we overridden the `x_delay` label using kicker.
95+
await delayed_task.kicker().with_labels(x_delay=4).kiq()
9696
```
9797

9898
## Priorities

examples/delayed_task.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ async def add_one(value: int) -> int:
2525
async def main() -> None:
2626
await broker.startup()
2727
# Send the task to the broker.
28-
task = await add_one.kicker().with_labels(delay=2).kiq(1)
28+
task = await add_one.kicker().with_labels(x_delay=2).kiq(1)
2929
print("Task sent with 2 seconds delay.")
3030
# Wait for the result.
3131
result = await task.wait_result(timeout=3)

taskiq_aio_pika/broker.py

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -370,7 +370,7 @@ async def kick(self, message: BrokerMessage) -> None:
370370
delivery_mode=DeliveryMode.PERSISTENT,
371371
priority=priority,
372372
)
373-
delay = parse_val(float, message.labels.get("delay"))
373+
x_delay = parse_val(float, message.labels.get("x_delay"))
374374

375375
if len(self._task_queues) == 1:
376376
routing_key_name = (
@@ -392,20 +392,20 @@ async def kick(self, message: BrokerMessage) -> None:
392392
f"Check routing keys and queue names in broker queues.",
393393
)
394394

395-
if delay is None:
395+
if x_delay is None:
396396
exchange = await self.write_channel.get_exchange(
397397
self._exchange.name,
398398
ensure=False,
399399
)
400400
await exchange.publish(rmq_message, routing_key=routing_key_name)
401401
elif self._delayed_message_exchange_plugin:
402-
rmq_message.headers["x-delay"] = int(delay * 1000)
402+
rmq_message.headers["x-delay"] = int(x_delay * 1000)
403403
exchange = await self.write_channel.get_exchange(
404404
self._delayed_message_exchange.name,
405405
)
406406
await exchange.publish(rmq_message, routing_key=routing_key_name)
407407
elif self._delay_queue:
408-
rmq_message.expiration = timedelta(seconds=delay)
408+
rmq_message.expiration = timedelta(seconds=x_delay)
409409
await self.write_channel.default_exchange.publish(
410410
rmq_message,
411411
routing_key=self._delay_queue.routing_key or self._delay_queue.name,

tests/test_delay.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ async def test_when_delayed_message_queue_exists__then_send_with_delay_must_work
2020
task_id="1",
2121
task_name="name",
2222
message=b"message",
23-
labels={"delay": "2"},
23+
labels={"x_delay": "2"},
2424
)
2525
await broker.kick(broker_msg)
2626

tests/test_delay_with_plugin.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ async def test_when_delayed_message_plugin_enabled__then_send_with_delay_must_wo
1919
task_id="1",
2020
task_name="name",
2121
message=b"message",
22-
labels={"delay": "2"},
22+
labels={"x_delay": "2"},
2323
)
2424

2525
# when & then

0 commit comments

Comments
 (0)