diff --git a/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py b/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py index 5898b20741a..1df27175445 100644 --- a/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py +++ b/litellm/integrations/prometheus_helpers/shared_prometheus_series_admissions.py @@ -12,8 +12,6 @@ _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: @@ -48,7 +46,7 @@ class _MetricAdmissions: def _append(self, label_values: tuple[str, ...]) -> None: descriptor: Final = os.open(self._path, os.O_WRONLY | os.O_APPEND | os.O_CREAT, 0o600) try: - os.write(descriptor, _LABEL_VALUES.dump_json(label_values) + b"\n") + os.write(descriptor, b"\n" + _LABEL_VALUES.dump_json(label_values) + b"\n") finally: os.close(descriptor) @@ -75,7 +73,9 @@ class SharedPrometheusSeriesAdmissions: ``PROMETHEUS_MULTIPROC_DIR``. Each metric has one append-only file there, and its first ``max_series`` distinct lines are the admitted label sets. Every worker reads the same lines in the same order, so all of them, including a worker that replaces an exited one, admit the same label sets and a scrape that merges - the workers stays at the cap.""" + the workers stays at the cap. Each record sits between two newlines, so a record a worker could only write + part of (the directory ran out of space) is a line of its own that admits nothing for every worker, and + it neither hides the records after it nor runs into the next worker's record.""" def __init__(self, directory: str) -> None: self._directory = directory diff --git a/tests/unit/integrations/test_prometheus_series_cardinality.py b/tests/unit/integrations/test_prometheus_series_cardinality.py index fba7ca9979b..4af0d91dfa1 100644 --- a/tests/unit/integrations/test_prometheus_series_cardinality.py +++ b/tests/unit/integrations/test_prometheus_series_cardinality.py @@ -1,7 +1,6 @@ import re from pathlib import Path from threading import Thread -from time import monotonic from typing import Final import pytest @@ -56,7 +55,7 @@ def isolated_registry_and_settings(monkeypatch): @pytest.fixture def clock(monkeypatch): - now: Final = [monotonic()] + now: Final = [1_000.0] monkeypatch.setattr(bounded_prometheus_series_tracker.time, "monotonic", lambda: now[0]) return now @@ -200,9 +199,9 @@ 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): +def test_a_record_cut_short_by_a_full_disk_admits_nothing_and_hides_no_other_record(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') + admissions_file.write_bytes(b'\n["user-a') writer: Final = SharedPrometheusSeriesAdmissions(directory=str(tmp_path)) reader: Final = SharedPrometheusSeriesAdmissions(directory=str(tmp_path))