Skip to content

Commit 5b2f8d1

Browse files
authored
Merge pull request #5335 from Agenta-AI/chore/fix-redis-exporter-stuff
[chore] Clean up redis exporter
2 parents 655f898 + 4aaa55d commit 5b2f8d1

17 files changed

Lines changed: 52 additions & 176 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/apis/fastapi/sessions/router.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -850,7 +850,6 @@ async def respond_interaction(
850850
interaction_id=str(interaction_id),
851851
answer=answer,
852852
)
853-
log.tick("interactions.enqueued", dims={"queue": "interactions"})
854853
else:
855854
references = (
856855
{

api/oss/src/apis/fastapi/triggers/router.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1602,10 +1602,8 @@ async def ingest_composio_event(
16021602
),
16031603
timeout=_ENQUEUE_TIMEOUT_SECONDS,
16041604
)
1605-
log.tick("triggers.enqueued", dims={"queue": "triggers"})
16061605
except Exception as e:
16071606
log.error("Failed to enqueue trigger event: %s", e)
1608-
log.tick("triggers.enqueue_errors", dims={"queue": "triggers"})
16091607
raise HTTPException(
16101608
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
16111609
detail="Failed to enqueue trigger event",

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/core/evaluations/runtime/runner.py

Lines changed: 0 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -38,11 +38,6 @@ async def process_run_from_source(
3838
kwargs["oldest"] = oldest
3939

4040
result = await self.worker.process_run_from_source.kiq(**kwargs)
41-
log.tick(
42-
"evaluations.enqueued",
43-
dims={"queue": "evaluations"},
44-
task="run_from_source",
45-
)
4641
return result
4742

4843
async def process_run_from_batch(
@@ -76,9 +71,6 @@ async def process_run_from_batch(
7671
kwargs["input_step_key"] = input_step_key
7772

7873
result = await self.worker.process_run_from_batch.kiq(**kwargs)
79-
log.tick(
80-
"evaluations.enqueued", dims={"queue": "evaluations"}, task="run_from_batch"
81-
)
8274
return result
8375

8476
async def process_rerun(
@@ -116,5 +108,4 @@ async def process_rerun(
116108
kwargs["overwrite"] = overwrite
117109

118110
result = await self.worker.process_rerun.kiq(**kwargs)
119-
log.tick("evaluations.enqueued", dims={"queue": "evaluations"}, task="rerun")
120111
return result

api/oss/src/core/events/streaming.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -89,13 +89,7 @@ async def publish_event(
8989
maxlen=MAXLEN_STREAMS_EVENTS,
9090
approximate=True,
9191
)
92-
log.tick(
93-
"events.published",
94-
bytes=len(event_bytes),
95-
dims={"stream": "events"},
96-
)
9792
return True
9893
except Exception as e:
9994
log.error(f"[EVENTS] Failed to publish event: {e}", exc_info=True)
100-
log.tick("events.publish_errors", dims={"stream": "events"})
10195
return False

api/oss/src/core/sessions/records/streaming.py

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -91,13 +91,7 @@ async def publish_record(
9191
maxlen=MAXLEN_STREAMS_RECORDS,
9292
approximate=True,
9393
)
94-
log.tick(
95-
"records.published",
96-
bytes=len(event_bytes),
97-
dims={"stream": "records"},
98-
)
9994
return True
10095
except Exception as e:
10196
log.error(f"[RECORDS] Failed to publish: {e}", exc_info=True)
102-
log.tick("records.publish_errors", dims={"stream": "records"})
10397
return False

api/oss/src/core/tracing/streaming.py

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,4 @@
11
import zlib
2-
from time import perf_counter
32
from typing import List
43
from uuid import UUID
54

@@ -84,7 +83,6 @@ async def publish_spans(
8483

8584
count = 0
8685
total_bytes = 0
87-
started = perf_counter()
8886

8987
# Full span_dtos list is already in hand, so pipeline all XADDs into one
9088
# round-trip instead of one XADD per span on the firehose.
@@ -111,12 +109,4 @@ async def publish_spans(
111109
if count:
112110
await pipe.execute()
113111

114-
log.tick(
115-
"spans.published",
116-
count=count,
117-
bytes=total_bytes,
118-
duration_ms=(perf_counter() - started) * 1000,
119-
dims={"stream": "spans"},
120-
)
121-
122112
return count

api/oss/src/core/triggers/service.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1560,7 +1560,6 @@ async def refresh_schedules(
15601560
timeout=_ENQUEUE_TIMEOUT_SECONDS,
15611561
)
15621562

1563-
log.tick("triggers.enqueued", dims={"queue": "triggers"})
15641563
log.info(
15651564
"[SCHEDULE] Dispatched. ",
15661565
project_id=project_id,

0 commit comments

Comments
 (0)