Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions bases/bot_detector/worker_ban_migration/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
)
from bot_detector.event_queue.structs import PlayerBannedStruct
from bot_detector.worker.core import WorkerRunner
from bot_detector.worker.metrics import start_metrics_server
from bot_detector.worker_ban_migration.settings import Settings
from bot_detector.worker_ban_migration.worker import BanMigrationWorker

Expand All @@ -21,6 +22,7 @@


async def main():
start_metrics_server(port=SETTINGS.METRICS_PORT)
session_factory, async_engine = db.get_session_factory(SETTINGS=DBSettings())
stop_event = asyncio.Event()
tasks = []
Expand All @@ -46,6 +48,7 @@ async def main():
model=PlayerBannedStruct,
batch_size=SETTINGS.MAX_BATCH_SIZE,
stop_event=stop_event,
worker_name="ban_migration",
)
task = asyncio.create_task(runner.run())
tasks.append(task)
Expand Down
7 changes: 1 addition & 6 deletions bases/bot_detector/worker_ban_migration/metrics.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,4 @@
import os

from prometheus_client import Counter, start_http_server

if os.environ.get("ENVIRONMENT") != "test":
start_http_server(8000)
from prometheus_client import Counter

# Prometheus metrics
accounts_migrated_counter = Counter(
Expand Down
1 change: 1 addition & 0 deletions bases/bot_detector/worker_ban_migration/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,3 +5,4 @@ class Settings(BaseSettings):
N_WORKERS: int = 1
MAX_BATCH_SIZE: int = 10_000
MAX_INTERVAL_MS: int = 1_000
METRICS_PORT: int = 8000
3 changes: 3 additions & 0 deletions bases/bot_detector/worker_hiscore/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
ScrapedStruct,
)
from bot_detector.worker.core import WorkerRunner
from bot_detector.worker.metrics import start_metrics_server
from bot_detector.worker_hiscore.settings import Settings
from bot_detector.worker_hiscore.worker import HiscoreWorker

Expand Down Expand Up @@ -55,6 +56,7 @@ async def get_data_to_predict_producer() -> QueueProducer[DataToPredictStruct]:


async def main():
start_metrics_server(port=SETTINGS.METRICS_PORT)
session_factory, async_engine = db.get_session_factory(SETTINGS=DBSettings())

player_repo = PlayerRepo()
Expand Down Expand Up @@ -84,6 +86,7 @@ async def main():
)
runner = WorkerRunner(
worker=worker,
worker_name="hiscore",
config=kafka_config,
model=ScrapedStruct,
batch_size=SETTINGS.MAX_BATCH_SIZE,
Expand Down
1 change: 1 addition & 0 deletions bases/bot_detector/worker_hiscore/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,3 +5,4 @@ class Settings(BaseSettings):
N_WORKERS: int = 1
MAX_BATCH_SIZE: int = 10_000
MAX_INTERVAL_MS: int = 5_000
METRICS_PORT: int = 8000
3 changes: 3 additions & 0 deletions bases/bot_detector/worker_report/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
)
from bot_detector.event_queue.structs import ReportsToInsertStruct
from bot_detector.worker.core import WorkerRunner
from bot_detector.worker.metrics import start_metrics_server
from bot_detector.worker_report.settings import Settings
from bot_detector.worker_report.worker import ReportWorker

Expand All @@ -22,6 +23,7 @@


async def main():
start_metrics_server(port=SETTINGS.METRICS_PORT)
session_factory, async_engine = db.get_session_factory(SETTINGS=DBSettings())
report_repo = ReportRepo()
stop_event = asyncio.Event()
Expand Down Expand Up @@ -49,6 +51,7 @@ async def main():
model=ReportsToInsertStruct,
batch_size=SETTINGS.MAX_BATCH_SIZE,
stop_event=stop_event,
worker_name="report",
)
task = asyncio.create_task(runner.run())
tasks.append(task)
Expand Down
1 change: 1 addition & 0 deletions bases/bot_detector/worker_report/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,3 +5,4 @@ class Settings(BaseSettings):
N_WORKERS: int = 1
MAX_BATCH_SIZE: int = 10_000
MAX_INTERVAL_MS: int = 1_000
METRICS_PORT: int = 8000
3 changes: 2 additions & 1 deletion components/bot_detector/worker/__init__.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
from bot_detector.worker.core import Worker, WorkerRunner
from bot_detector.worker.errors import WorkerError
from bot_detector.worker.metrics import start_metrics_server

__all__ = ["Worker", "WorkerError", "WorkerRunner"]
__all__ = ["Worker", "WorkerError", "WorkerRunner", "start_metrics_server"]
46 changes: 45 additions & 1 deletion components/bot_detector/worker/core.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import asyncio
import logging
import time
from abc import ABC, abstractmethod
from typing import Generic, Literal, Type, TypeVar

Expand All @@ -8,6 +9,15 @@
from bot_detector.event_queue.core import Queue
from bot_detector.event_queue.factory import QueueFactory
from bot_detector.worker.errors import WorkerError
from bot_detector.worker.metrics import (
batch_errors_counter,
batch_size_histogram,
handle_duration_histogram,
messages_committed_counter,
messages_consumed_counter,
messages_requeued_counter,
poll_idle_counter,
)
from pydantic import BaseModel

T = TypeVar("T", bound=BaseModel)
Expand Down Expand Up @@ -60,10 +70,12 @@ def __init__(
config: KafkaConfig | InMemoryConfig,
model: Type[T],
worker: Worker[T],
worker_name: str,
stop_event: asyncio.Event,
batch_size: int = 1000,
) -> None:
self._worker = worker
self._worker_name = worker_name
self._batch_size = batch_size
self._config = config
self._model = model
Expand Down Expand Up @@ -115,6 +127,7 @@ async def _requeue(self, batch: list[T]) -> None:

async def _consume(self) -> None:
"""Main loop: get_many → handle → commit. Requeues what handle returns."""
name = self._worker_name
while not self._stop_event.is_set():
batch: list[T] = []
try:
Expand All @@ -125,12 +138,22 @@ async def _consume(self) -> None:

batch = result
if not batch:
poll_idle_counter.labels(worker=name).inc()
await asyncio.sleep(0.1)
continue

batch_size_histogram.labels(worker=name).observe(len(batch))
messages_consumed_counter.labels(worker=name).inc(len(batch))
logger.info(f"Consumed {len(batch)} messages")

result = await self._worker.handle(batch)
start = time.monotonic()
try:
result = await self._worker.handle(batch)
finally:
handle_duration_histogram.labels(worker=name).observe(
time.monotonic() - start
)

if isinstance(result, WorkerError):
if result.error_batch:
logger.error(
Expand All @@ -139,6 +162,9 @@ async def _consume(self) -> None:
f"requeuing={len(result.error_batch)} "
f"of {len(batch)}"
)
messages_requeued_counter.labels(
worker=name, reason="worker_error"
).inc(len(result.error_batch))
await self._requeue(result.error_batch)
await asyncio.sleep(1)
else:
Expand All @@ -149,20 +175,38 @@ async def _consume(self) -> None:
commit_err = await self._queue.commit()
if isinstance(commit_err, Exception):
logger.error(f"Failed to commit batch: {commit_err}")
else:
messages_committed_counter.labels(worker=name).inc(
len(batch)
)
elif result:
logger.info(f"Requeuing {len(result)} of {len(batch)} messages")
messages_requeued_counter.labels(worker=name, reason="requeue").inc(
len(result)
)
await self._requeue(result)
else:
commit_err = await self._queue.commit()
if isinstance(commit_err, Exception):
logger.error(f"Failed to commit batch: {commit_err}")
else:
messages_committed_counter.labels(worker=name).inc(len(batch))
except asyncio.CancelledError:
if batch:
messages_requeued_counter.labels(
worker=name, reason="cancelled"
).inc(len(batch))
await self._requeue(batch)
raise
except Exception as e:
batch_errors_counter.labels(
worker=name, error_type=type(e).__name__
).inc()
logger.error(f"Error processing batch: {e}", exc_info=True)
if batch:
messages_requeued_counter.labels(worker=name, reason="error").inc(
len(batch)
)
await self._requeue(batch)
await asyncio.sleep(1)

Expand Down
63 changes: 63 additions & 0 deletions components/bot_detector/worker/metrics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import logging
import os

from prometheus_client import Counter, Histogram, start_http_server

logger = logging.getLogger(__name__)

_started = False

WORKER_LABELS = ["worker"]


def start_metrics_server(port: int = 8000) -> None:
"""Start the Prometheus metrics HTTP server.

Idempotent: subsequent calls are no-ops. No-op when ENVIRONMENT == "test"
so unit tests don't bind a port.
"""
global _started
if _started:
return
if os.environ.get("ENVIRONMENT") == "test":
return
start_http_server(port)
_started = True
logger.info(f"Metrics server listening on :{port}")


messages_consumed_counter = Counter(
name="worker_messages_consumed_total",
documentation="Total messages pulled from the queue by WorkerRunner",
labelnames=WORKER_LABELS,
)
messages_committed_counter = Counter(
name="worker_messages_committed_total",
documentation="Total messages successfully committed (processed)",
labelnames=WORKER_LABELS,
)
messages_requeued_counter = Counter(
name="worker_messages_requeued_total",
documentation="Total messages requeued, by reason",
labelnames=["worker", "reason"],
)
batch_size_histogram = Histogram(
name="worker_batch_size",
documentation="Distribution of consumed batch sizes",
labelnames=WORKER_LABELS,
)
handle_duration_histogram = Histogram(
name="worker_handle_duration_seconds",
documentation="Wall-clock time spent in Worker.handle() per batch",
labelnames=WORKER_LABELS,
)
batch_errors_counter = Counter(
name="worker_batch_errors_total",
documentation="Exceptions raised inside handle(), by exception type",
labelnames=["worker", "error_type"],
)
poll_idle_counter = Counter(
name="worker_poll_idle_total",
documentation="Empty polls where no messages were available",
labelnames=WORKER_LABELS,
)
1 change: 1 addition & 0 deletions projects/worker_hiscore/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ dependencies = [
"pydantic==2.10.4",
"sqlalchemy==2.0.44",
"orjson==3.11.5",
"prometheus-client==0.22.0",
]

[project.scripts]
Expand Down
Loading
Loading