mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-11 03:38:38 +00:00
fix(prometheus): frame each admissions record with newlines so a cut-off record cannot swallow the next
A record a worker could only write part of used to merge with the next worker's record, and both were skipped for one request. Each record is now written between two newlines, so the fragment is a line of its own. The clock fixture in the series tests starts from a constant instead of reading the real clock
This commit is contained in:
parent
fed49bc253
commit
2aefd9c76e
2 changed files with 7 additions and 8 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue