mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
fix(prometheus): skip an admissions line a worker could only write part of
This commit is contained in:
parent
ee8600b05b
commit
fed49bc253
2 changed files with 27 additions and 3 deletions
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue