From fed49bc253a0957fe221d66ad8c0389b940b8295 Mon Sep 17 00:00:00 2001 From: mateo-berri <277851410+mateo-berri@users.noreply.github.com> Date: Sat, 3 Oct 2026 13:45:21 -0700 Subject: [PATCH] fix(prometheus): skip an admissions line a worker could only write part of --- .../shared_prometheus_series_admissions.py | 16 +++++++++++++--- .../test_prometheus_series_cardinality.py | 14 ++++++++++++++ 2 files changed, 27 insertions(+), 3 deletions(-) diff --git a/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py b/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py index 598dd0cef40..5898b20741a 100644 --- a/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py +++ b/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py @@ -4,13 +4,22 @@ import os from threading import RLock from typing import Final -from pydantic import TypeAdapter +from pydantic import TypeAdapter, ValidationError from litellm.constants import PROMETHEUS_ADMITTED_SERIES_FILE_PREFIX _LABEL_VALUES: Final = TypeAdapter(tuple[str, ...]) +def _parse_admission(line: bytes) -> tuple[str, ...] | None: + """A line a worker could only write part of, which happens when the directory runs out of space, admits + nothing for every worker rather than stopping every worker from reading the lines after it.""" + try: + return _LABEL_VALUES.validate_json(line) + except ValidationError: + return None + + class _MetricAdmissions: __slots__ = ("_label_sets", "_max_series", "_path", "_read_offset") @@ -54,10 +63,11 @@ class _MetricAdmissions: if not newline: return self._read_offset += len(complete_lines) + len(newline) - for line in complete_lines.split(b"\n"): + for label_values in map(_parse_admission, complete_lines.split(b"\n")): if self._is_full(): return - self._label_sets.add(_LABEL_VALUES.validate_json(line)) + if label_values is not None: + self._label_sets.add(label_values) class SharedPrometheusSeriesAdmissions: diff --git a/tests/unit/integrations/test_prometheus_series_cardinality.py b/tests/unit/integrations/test_prometheus_series_cardinality.py index 9894d760788..fba7ca9979b 100644 --- a/tests/unit/integrations/test_prometheus_series_cardinality.py +++ b/tests/unit/integrations/test_prometheus_series_cardinality.py @@ -200,6 +200,20 @@ def test_a_line_another_worker_is_still_writing_is_read_once_it_is_complete(tmp_ assert reader.admit_series("litellm_requests_metric", ("user-b",), max_series=2) +def test_a_line_cut_short_by_a_full_disk_admits_nothing_and_stops_no_worker(tmp_path: Path): + admissions_file: Final = tmp_path / f"{PROMETHEUS_ADMITTED_SERIES_FILE_PREFIX}litellm_requests_metric" + admissions_file.write_bytes(b'["user-a\n') + writer: Final = SharedPrometheusSeriesAdmissions(directory=str(tmp_path)) + reader: Final = SharedPrometheusSeriesAdmissions(directory=str(tmp_path)) + + assert writer.admit_series("litellm_requests_metric", ("user-b",), max_series=2) + assert reader.admit_series("litellm_requests_metric", ("user-b",), max_series=2) + assert reader.admit_series("litellm_requests_metric", ("user-c",), max_series=2) + assert writer.admit_series("litellm_requests_metric", ("user-c",), max_series=2) + assert not writer.admit_series("litellm_requests_metric", ("user-a",), max_series=2) + assert not reader.admit_series("litellm_requests_metric", ("user-a",), max_series=2) + + def test_wiping_the_multiprocess_dir_frees_every_admitted_slot(tmp_path: Path): before_restart: Final = SharedPrometheusSeriesAdmissions(directory=str(tmp_path)) assert before_restart.admit_series("litellm_requests_metric", ("user-a",), max_series=1)