Skip to content

Commit 32e3d7d

Browse files
committed
fix XDEL after XACK
1 parent e723777 commit 32e3d7d

5 files changed

Lines changed: 50 additions & 10 deletions

File tree

api/entrypoints/worker_queues.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,8 @@
3131
from taskiq import AsyncBroker, TaskiqEvents
3232
from taskiq.receiver import Receiver
3333
from taskiq.cli.worker.run import shutdown_broker
34-
from taskiq_redis import RedisStreamBroker
34+
35+
from oss.src.tasks.taskiq.shared.broker import TrimOnAckRedisStreamBroker
3536

3637
from oss.src.core.embeds.service import EmbedsService
3738
from oss.src.core.environments.service import EnvironmentsService
@@ -123,7 +124,7 @@ def _selected_queues() -> List[str]:
123124

124125

125126
def _build_webhooks_broker() -> tuple[AsyncBroker, int]:
126-
broker = RedisStreamBroker(
127+
broker = TrimOnAckRedisStreamBroker(
127128
url=env.redis.uri_durable,
128129
queue_name="queues:webhooks",
129130
consumer_group_name="worker-webhooks",
@@ -135,7 +136,7 @@ def _build_webhooks_broker() -> tuple[AsyncBroker, int]:
135136

136137

137138
def _build_triggers_broker() -> tuple[AsyncBroker, int]:
138-
broker = RedisStreamBroker(
139+
broker = TrimOnAckRedisStreamBroker(
139140
url=env.redis.uri_durable,
140141
queue_name="queues:triggers",
141142
consumer_group_name="worker-triggers",
@@ -175,7 +176,7 @@ def _build_triggers_broker() -> tuple[AsyncBroker, int]:
175176

176177

177178
def _build_interactions_broker() -> tuple[AsyncBroker, int]:
178-
broker = RedisStreamBroker(
179+
broker = TrimOnAckRedisStreamBroker(
179180
url=env.redis.uri_durable,
180181
queue_name="queues:interactions",
181182
consumer_group_name="worker-interactions",

api/entrypoints/worker_streams.py

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,8 @@
1717
from typing import List
1818

1919
from redis.asyncio import Redis
20-
from taskiq_redis import RedisStreamBroker
20+
21+
from oss.src.tasks.taskiq.shared.broker import ProducerOnlyRedisStreamBroker
2122

2223
from oss.src.core.events.service import EventsService
2324
from oss.src.core.secrets.services import VaultService
@@ -85,9 +86,11 @@ async def _build_records_worker(redis_client: Redis) -> StreamConsumer:
8586
async def _build_events_worker(redis_client: Redis) -> StreamConsumer:
8687
events_service = EventsService(events_dao=EventsDAO())
8788

88-
# Webhook dispatch runs inside the events loop as its post-hook.
89+
# Webhook dispatch runs inside the events loop as its post-hook: this broker
90+
# only produces (.kiq), so it must not declare a consumer group it never
91+
# reads — that group sits at 0-0 and reports lag == XLEN forever.
8992
webhooks_dao = WebhooksDAO()
90-
broker = RedisStreamBroker(
93+
broker = ProducerOnlyRedisStreamBroker(
9194
url=env.redis.uri_durable,
9295
queue_name="queues:webhooks",
9396
consumer_group_name="worker-events-webhooks-dispatcher",

api/oss/src/core/evaluations/runtime/broker.py

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@
66
"""
77

88
from taskiq import AsyncBroker
9-
from taskiq_redis import RedisStreamBroker
109

10+
from oss.src.tasks.taskiq.shared.broker import TrimOnAckRedisStreamBroker
1111
from oss.src.core.evaluations.runtime.runner import TaskiqEvaluationTaskRunner
1212
from oss.src.core.evaluations.service import EvaluationsService
1313
from oss.src.core.evaluators.service import SimpleEvaluatorsService
@@ -22,12 +22,13 @@
2222
MAXLEN_QUEUES_EVALUATIONS = 100_000
2323

2424

25-
class NoRedeliveryRedisStreamBroker(RedisStreamBroker):
25+
class NoRedeliveryRedisStreamBroker(TrimOnAckRedisStreamBroker):
2626
"""Stream broker that never redelivers. `listen()` reads only NEW messages
2727
(`>`) and skips the XAUTOCLAIM pending-replay block, so a task that crashed
2828
mid-run is not re-served to later workers. Evaluation tasks are not safely
2929
re-runnable (`retry_on_error=False`); a stuck unacked entry replaying on every
30-
worker restart is worse than dropping it.
30+
worker restart is worse than dropping it. Inherits XDEL-on-ack so completed
31+
entries leave the stream (XLEN = backlog).
3132
"""
3233

3334
async def listen(self):

api/oss/src/tasks/taskiq/shared/__init__.py

Whitespace-only changes.
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
"""Shared TaskIQ RedisStreamBroker mixins.
2+
3+
`TrimOnAckRedisStreamBroker` XDELs each entry right after XACK. Without it, a
4+
TaskIQ queue retains every acked entry until maxlen-trims it, so XLEN reflects
5+
retained history (a rolling maxlen window) rather than backlog — the stream
6+
workers already XACK+XDEL, so this makes the queues match. XDEL fires only on
7+
ack (a terminal state), so retries (re-kicked as new messages) and crash
8+
redelivery (unacked, reclaimed via XAUTOCLAIM from the PEL) are unaffected.
9+
10+
`ProducerOnlyRedisStreamBroker` skips consumer-group declaration on startup, for
11+
a broker used only to register tasks and `.kiq()` — never to `.listen()`. The
12+
base `startup()` unconditionally XGROUP-CREATEs a group that then sits at 0-0
13+
forever (0 consumers), reporting lag == XLEN permanently.
14+
"""
15+
16+
from typing import Awaitable, Callable
17+
18+
from redis.asyncio import Redis
19+
from taskiq_redis import RedisStreamBroker
20+
from taskiq_redis.redis_broker import BaseRedisBroker
21+
22+
23+
class TrimOnAckRedisStreamBroker(RedisStreamBroker):
24+
def _ack_generator(self, id: str, queue_name: str) -> Callable[[], Awaitable[None]]:
25+
async def _ack() -> None:
26+
async with Redis(connection_pool=self.connection_pool) as redis_conn:
27+
await redis_conn.xack(queue_name, self.consumer_group_name, id)
28+
await redis_conn.xdel(queue_name, id)
29+
30+
return _ack
31+
32+
33+
class ProducerOnlyRedisStreamBroker(RedisStreamBroker):
34+
async def startup(self) -> None:
35+
await BaseRedisBroker.startup(self)

0 commit comments

Comments
 (0)