From a485f5e36a6d26d2334a3566530c3b6f1168d232 Mon Sep 17 00:00:00 2001 From: rcvalerio Date: Thu, 18 Dec 2025 16:25:36 +0000 Subject: [PATCH 1/3] Improve metrics --- docker-compose.lite.yml | 25 +++++++++++++++++++ file-tracker/file_tracker/main.py | 3 +++ file-tracker/file_tracker/metrics.py | 24 ++++++++++++++++++ file-tracker/requirements.txt | 3 ++- prometheus.yml | 16 ++++++++++++ task-runner/requirements.txt | 1 + task-runner/task_runner/main.py | 4 +++ task-runner/task_runner/metrics.py | 17 +++++++++++++ .../task_runner/task_request_handler.py | 15 ++++++++++- 9 files changed, 106 insertions(+), 2 deletions(-) create mode 100644 file-tracker/file_tracker/metrics.py create mode 100644 prometheus.yml create mode 100644 task-runner/task_runner/metrics.py diff --git a/docker-compose.lite.yml b/docker-compose.lite.yml index bd0a488d..851a0711 100644 --- a/docker-compose.lite.yml +++ b/docker-compose.lite.yml @@ -27,6 +27,31 @@ services: network_mode: host volumes: - workdir:/workdir + prometheus: + image: prom/prometheus:latest + network_mode: host + user: "0:0" + volumes: + - ./prometheus.yml:/etc/prometheus/prometheus.yml:ro + - prometheus_data:/prometheus + command: + - --config.file=/etc/prometheus/prometheus.yml + - --storage.tsdb.path=/prometheus + - --storage.tsdb.retention.time=15d + - --web.console.libraries=/usr/share/prometheus/console_libraries + - --web.console.templates=/usr/share/prometheus/consoles + grafana: + image: grafana/grafana:latest + network_mode: host + user: "0:0" + environment: + - GF_SECURITY_ADMIN_USER=admin + - GF_SECURITY_ADMIN_PASSWORD=admin + - GF_SERVER_HTTP_PORT=3500 + volumes: + - grafana_data:/var/lib/grafana volumes: workdir: + prometheus_data: + grafana_data: diff --git a/file-tracker/file_tracker/main.py b/file-tracker/file_tracker/main.py index 720434bd..48f823bd 100644 --- a/file-tracker/file_tracker/main.py +++ b/file-tracker/file_tracker/main.py @@ -4,6 +4,7 @@ from cleanup import TerminationHandler, setup_cleanup_handlers from connection_manager import ConnectionManager +from metrics import start_metrics_server from task_listener import TaskListener @@ -11,6 +12,8 @@ async def main(): workdir = os.getenv("WORKDIR", "/workdir") os.chdir(workdir) + await start_metrics_server() + connection_manager = ConnectionManager.from_env() file_tracker_host = os.getenv("FILE_TRACKER_HOST", "0.0.0.0") file_tracker_port = int(os.getenv("FILE_TRACKER_PORT", "5000")) diff --git a/file-tracker/file_tracker/metrics.py b/file-tracker/file_tracker/metrics.py new file mode 100644 index 00000000..08d7b932 --- /dev/null +++ b/file-tracker/file_tracker/metrics.py @@ -0,0 +1,24 @@ +"""Prometheus metrics server for file-tracker service.""" +import logging +import os + +from aiohttp import web +from prometheus_client import generate_latest + + +async def metrics_handler(request): + """Handler for /metrics endpoint exposing Prometheus client metrics.""" + return web.Response(body=generate_latest(), content_type="text/plain") + + +async def start_metrics_server(host='0.0.0.0', port=None): + """Start a basic Prometheus metrics server exposing client metrics.""" + port = port or int(os.getenv('FILE_TRACKER_METRICS_PORT', '9091')) + app = web.Application() + app.router.add_get('/metrics', metrics_handler) + runner = web.AppRunner(app) + await runner.setup() + site = web.TCPSite(runner, host, port) + await site.start() + logging.info("Metrics server started on %s:%s", host, port) + return runner diff --git a/file-tracker/requirements.txt b/file-tracker/requirements.txt index 5d45f258..0be30042 100644 --- a/file-tracker/requirements.txt +++ b/file-tracker/requirements.txt @@ -1,2 +1,3 @@ aiortc -aiohttp \ No newline at end of file +aiohttp +prometheus-client \ No newline at end of file diff --git a/prometheus.yml b/prometheus.yml new file mode 100644 index 00000000..07c0b4a9 --- /dev/null +++ b/prometheus.yml @@ -0,0 +1,16 @@ +global: + scrape_interval: 15s + evaluation_interval: 15s + +scrape_configs: + - job_name: "file-tracker" + static_configs: + - targets: ["localhost:9091"] + metrics_path: "/metrics" + + - job_name: "task-runner-lite" + static_configs: + - targets: ["localhost:8000"] + metrics_path: "/metrics" + + diff --git a/task-runner/requirements.txt b/task-runner/requirements.txt index d0f0eff9..68a77332 100644 --- a/task-runner/requirements.txt +++ b/task-runner/requirements.txt @@ -10,3 +10,4 @@ gputil pydantic~=2.11.7 tenacity pytest +prometheus-client diff --git a/task-runner/task_runner/main.py b/task-runner/task_runner/main.py index 6883c199..d6849c43 100644 --- a/task-runner/task_runner/main.py +++ b/task-runner/task_runner/main.py @@ -29,6 +29,7 @@ task_execution_loop, utils, ) +from task_runner.metrics import start_metrics_server from task_runner.register_task_runner import register_task_runner from task_runner.task_request_handler import TaskRequestHandler from task_runner.task_status import TaskRunnerTerminationReason @@ -66,6 +67,9 @@ def _set_socks_proxy(): def main(_): _set_socks_proxy() + + start_metrics_server() + workdir = os.getenv("WORKDIR", "/workdir") executer_images_dir = os.getenv("EXECUTER_IMAGES_DIR", "/apptainer") if not executer_images_dir: diff --git a/task-runner/task_runner/metrics.py b/task-runner/task_runner/metrics.py new file mode 100644 index 00000000..7b3361cf --- /dev/null +++ b/task-runner/task_runner/metrics.py @@ -0,0 +1,17 @@ +"""Prometheus metrics server for task-runner service.""" +import logging +import os + +from prometheus_client import Counter, Histogram, Gauge, start_http_server + +# Task metrics +tasks_active = Gauge('tasks_active', 'Number of currently active tasks') +tasks_total = Counter('tasks_total', 'Total number of tasks', ['status']) +task_duration = Histogram('task_duration_seconds', 'Task execution duration in seconds') + + +def start_metrics_server(port=None): + """Start a basic Prometheus metrics server exposing client metrics.""" + port = port or int(os.getenv('TASK_RUNNER_METRICS_PORT', '8000')) + start_http_server(port) + logging.info("Metrics server started on port %s", port) diff --git a/task-runner/task_runner/task_request_handler.py b/task-runner/task_runner/task_request_handler.py index 0a55b799..442fdaba 100644 --- a/task-runner/task_runner/task_request_handler.py +++ b/task-runner/task_runner/task_request_handler.py @@ -32,8 +32,8 @@ task_message_listener, utils, ) +from task_runner.metrics import task_duration, tasks_active, tasks_total from task_runner.operations_logger import OperationName, OperationsLogger -from task_runner.task_status import task_status from task_runner.utils import files KILL_MESSAGE = "kill" @@ -287,6 +287,10 @@ def __call__(self, request: dict[str, str]) -> None: Args: request: Request describing the task to be executed. """ + task_start_time = time.time() + task_status = 'failed' + tasks_active.inc() + # Save the task request to use during output recovery with open(self.request_path, "w", encoding="utf-8") as request_file: json.dump(request, request_file, default=str) @@ -416,12 +420,15 @@ def __call__(self, request: dict[str, str]) -> None: if exit_reason == TaskExitReason.KILLED: new_status = task_status.TaskStatusCode.KILLED.value + task_status = 'killed' elif exit_reason == TaskExitReason.TTL_EXCEEDED: new_status = task_status.TaskStatusCode.TTL_EXCEEDED.value + task_status = 'ttl_exceeded' else: new_status = (task_status.TaskStatusCode.SUCCESS.value if exit_code == 0 else task_status.TaskStatusCode.FAILED.value) + task_status = 'success' if exit_code == 0 else 'failed' safely_delete = self.save_output(new_task_status=new_status) @@ -436,6 +443,7 @@ def __call__(self, request: dict[str, str]) -> None: # Catch all exceptions to ensure that we log the error message except Exception as e: # noqa: BLE001 + task_status = 'error' message = utils.get_exception_root_cause_message(e) try: self._publish_event( @@ -461,6 +469,11 @@ def __call__(self, request: dict[str, str]) -> None: safely_delete = False finally: + # Track metrics + task_duration.observe(time.time() - task_start_time) + tasks_total.labels(status=task_status).inc() + tasks_active.dec() + self.cleaning_up = True self._cleanup(safely_delete) From 45284e3e337e09d201a0f994316545a1bebb6b76 Mon Sep 17 00:00:00 2001 From: rcvalerio Date: Thu, 18 Dec 2025 16:27:53 +0000 Subject: [PATCH 2/3] Lint --- task-runner/task_runner/metrics.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/task-runner/task_runner/metrics.py b/task-runner/task_runner/metrics.py index 7b3361cf..c4f1454c 100644 --- a/task-runner/task_runner/metrics.py +++ b/task-runner/task_runner/metrics.py @@ -2,12 +2,13 @@ import logging import os -from prometheus_client import Counter, Histogram, Gauge, start_http_server +from prometheus_client import Counter, Gauge, Histogram, start_http_server # Task metrics tasks_active = Gauge('tasks_active', 'Number of currently active tasks') tasks_total = Counter('tasks_total', 'Total number of tasks', ['status']) -task_duration = Histogram('task_duration_seconds', 'Task execution duration in seconds') +task_duration = Histogram('task_duration_seconds', + 'Task execution duration in seconds') def start_metrics_server(port=None): From 470b69e920c6cc3ab695639ba143881d69907fe2 Mon Sep 17 00:00:00 2001 From: rcvalerio Date: Thu, 18 Dec 2025 16:32:07 +0000 Subject: [PATCH 3/3] Fix test --- task-runner/task_runner/task_request_handler.py | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/task-runner/task_runner/task_request_handler.py b/task-runner/task_runner/task_request_handler.py index 442fdaba..5224e81a 100644 --- a/task-runner/task_runner/task_request_handler.py +++ b/task-runner/task_runner/task_request_handler.py @@ -30,6 +30,7 @@ executers, observers, task_message_listener, + task_status, utils, ) from task_runner.metrics import task_duration, tasks_active, tasks_total @@ -288,7 +289,7 @@ def __call__(self, request: dict[str, str]) -> None: request: Request describing the task to be executed. """ task_start_time = time.time() - task_status = 'failed' + task_status_str = 'failed' tasks_active.inc() # Save the task request to use during output recovery @@ -420,15 +421,15 @@ def __call__(self, request: dict[str, str]) -> None: if exit_reason == TaskExitReason.KILLED: new_status = task_status.TaskStatusCode.KILLED.value - task_status = 'killed' + task_status_str = 'killed' elif exit_reason == TaskExitReason.TTL_EXCEEDED: new_status = task_status.TaskStatusCode.TTL_EXCEEDED.value - task_status = 'ttl_exceeded' + task_status_str = 'ttl_exceeded' else: new_status = (task_status.TaskStatusCode.SUCCESS.value if exit_code == 0 else task_status.TaskStatusCode.FAILED.value) - task_status = 'success' if exit_code == 0 else 'failed' + task_status_str = 'success' if exit_code == 0 else 'failed' safely_delete = self.save_output(new_task_status=new_status) @@ -443,7 +444,7 @@ def __call__(self, request: dict[str, str]) -> None: # Catch all exceptions to ensure that we log the error message except Exception as e: # noqa: BLE001 - task_status = 'error' + task_status_str = 'error' message = utils.get_exception_root_cause_message(e) try: self._publish_event( @@ -471,7 +472,7 @@ def __call__(self, request: dict[str, str]) -> None: finally: # Track metrics task_duration.observe(time.time() - task_start_time) - tasks_total.labels(status=task_status).inc() + tasks_total.labels(status=task_status_str).inc() tasks_active.dec() self.cleaning_up = True