New `llm.metrics`: one lazily-built instrument set off `perf.meter("stack.llm")`,
so every call site is a no-op mock when telemetry is disabled (as it is in CI and,
today, in the container — see below). A build lock guards first use because
`embed_texts` fans batches across a thread pool.
Instruments, all `stack_llm_…` per the P26 convention:
stack_llm_dispatch_total{host,kind} pool slots handed out
stack_llm_inflight{host} slots held right now
stack_llm_dispatch_seconds{host,kind} how long a slot was held
stack_llm_embedded_texts_total{host} texts embedded
stack_llm_embed_batch_seconds{host} one /api/embed round trip
stack_llm_chat_seconds{stage} retrieve | generate
stack_llm_chat_total{outcome} ok | error
stack_llm_chat_tokens_total{direction} prompt | completion (Ollama's
final frame; skipped when absent)
stack_llm_indexed_chunks_total{collection}
stack_llm_indexed_items_total{collection,outcome} indexed | skipped
Attribute values are closed sets — host, kind, stage, outcome, collection — never
item keys, dockets or question text, so the series count stays bounded.
Call sites: `HostPool.acquire`/`acquire_generation` wrap the held slot (kind embed
vs generate); `embed_texts` times each batch around the POST. `stream_answer` is
now a thin wrapper that counts the outcome once (ok, or error and re-raise) around
`_stream_answer`, which records the retrieve/generate stage split and the token
counts. The indexer counts one item at each of its five existing decision points;
no behaviour changed — every added line is a metric call.
The #575 tagging chain does not exist yet, so tag throughput and abstain rate from
the issue's scope are deliberately left out; they belong with that module.
Dashboard `infra/grafana/dashboards/llm.json` (uid `llm`, "LLM service",
schemaVersion 39): request rate by status filtered to `service_name="llm"`, chat
outcomes, dispatch rate and in-flight per host, p50/p95 dispatch latency by
host+kind, chat latency by stage, embedding throughput, chat tokens, indexed
chunks and items per collection, plus a Loki tail.
Export path — the metrics do NOT reach Prometheus yet, and this commit does not
change that. `perf._meter.setup_meter_provider` installs only a
`PrometheusMetricReader` (an in-process prometheus_client registry); the
`OTEL_EXPORTER_OTLP_ENDPOINT` the llm container sets is read by `_tracer.py` for
spans only, so nothing is pushed to otel-collector:4317 and nothing is exposed for
a scrape (`llm` has no /metrics route and no target in
infra/prometheus/targets/services.yml). On top of that `[telemetry] enabled` is
false in stack.toml and the llm service sets no STACK_TELEMETRY, so `perf.meter`
hands back the no-op mock today. Making #579's "metrics scraped" true needs a
follow-up that (a) enables telemetry for llm, (b) gives metrics a real exporter,
and (c) accounts for the collector's `namespace: stack`, which would re-prefix
OTLP-delivered names to `stack_stack_llm_…`. Filed as a note on #579.
Tests: tests/llm/test_metrics.py patches `perf.meter` with a recording fake and
asserts every instrument name, kind and attribute set, drives the real pool, chat
and indexer paths through it, and checks the helpers stay inert against the real
`MockMeter`; tests/test_observability.py gains the llm dashboard to
EXPECTED_DASHBOARDS plus panel/datasource/quantile assertions.
409 lines
14 KiB
Python
409 lines
14 KiB
Python
"""llm.metrics — OTel instruments for dispatch, embeddings, chat, indexing.
|
|
|
|
Two halves: with ``perf.meter`` patched to a recording fake (so the
|
|
instrument names and attributes are asserted exactly), and against the
|
|
real ``perf`` no-op ``MockMeter`` that every test run gets with
|
|
telemetry disabled — the call sites must be inert, not merely quiet.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from unittest.mock import MagicMock, patch
|
|
|
|
import pytest
|
|
|
|
from llm import metrics
|
|
from llm.config import LlmConfig
|
|
from llm.pool import HostPool, embed_texts
|
|
|
|
CFG = LlmConfig(
|
|
ollama_hosts=("http://h1:11434",),
|
|
host_vram={"http://h1:11434": 24},
|
|
embed_model="embed",
|
|
instruct_model="chat",
|
|
instruct_model_large="big",
|
|
embed_dim=768,
|
|
pg_host="x",
|
|
pg_port=5432,
|
|
pg_db="llm",
|
|
pg_user="llm",
|
|
build_ann_index=False,
|
|
top_n=3,
|
|
)
|
|
H1 = "http://h1:11434"
|
|
H2 = "http://h2:11434"
|
|
|
|
|
|
class FakeInstrument:
|
|
"""Records every ``add``/``record`` call as (amount, attributes)."""
|
|
|
|
def __init__(self, name: str, kind: str) -> None:
|
|
self.name = name
|
|
self.kind = kind
|
|
self.calls: list[tuple[float, dict]] = []
|
|
|
|
def add(self, amount, attributes=None):
|
|
self.calls.append((amount, dict(attributes or {})))
|
|
|
|
def record(self, amount, attributes=None):
|
|
self.calls.append((amount, dict(attributes or {})))
|
|
|
|
|
|
class FakeMeter:
|
|
def __init__(self) -> None:
|
|
self.name = ""
|
|
self.instruments: dict[str, FakeInstrument] = {}
|
|
|
|
def _make(self, name, kind):
|
|
self.instruments[name] = FakeInstrument(name, kind)
|
|
return self.instruments[name]
|
|
|
|
def create_counter(self, name, **kw):
|
|
return self._make(name, "counter")
|
|
|
|
def create_histogram(self, name, **kw):
|
|
return self._make(name, "histogram")
|
|
|
|
def create_up_down_counter(self, name, **kw):
|
|
return self._make(name, "up_down_counter")
|
|
|
|
|
|
@pytest.fixture()
|
|
def fake_meter():
|
|
"""Build the instrument set against a recording meter."""
|
|
meter = FakeMeter()
|
|
|
|
def _meter(name):
|
|
meter.name = name
|
|
return meter
|
|
|
|
metrics.reset()
|
|
with patch("perf.meter", _meter):
|
|
metrics.instruments()
|
|
yield meter
|
|
metrics.reset()
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _reset_instruments():
|
|
"""No test may inherit another's cached instrument set."""
|
|
metrics.reset()
|
|
yield
|
|
metrics.reset()
|
|
|
|
|
|
def _calls(meter: FakeMeter, name: str) -> list[tuple[float, dict]]:
|
|
return meter.instruments[name].calls
|
|
|
|
|
|
class TestInstrumentSet:
|
|
def test_names_and_kinds(self, fake_meter):
|
|
assert {n: i.kind for n, i in fake_meter.instruments.items()} == {
|
|
"stack_llm_dispatch_total": "counter",
|
|
"stack_llm_inflight": "up_down_counter",
|
|
"stack_llm_dispatch_seconds": "histogram",
|
|
"stack_llm_embedded_texts_total": "counter",
|
|
"stack_llm_embed_batch_seconds": "histogram",
|
|
"stack_llm_chat_seconds": "histogram",
|
|
"stack_llm_chat_total": "counter",
|
|
"stack_llm_chat_tokens_total": "counter",
|
|
"stack_llm_indexed_chunks_total": "counter",
|
|
"stack_llm_indexed_items_total": "counter",
|
|
}
|
|
|
|
def test_meter_scope(self, fake_meter):
|
|
assert fake_meter.name == "stack.llm"
|
|
|
|
def test_instruments_built_once(self, fake_meter):
|
|
first = metrics.instruments()
|
|
assert metrics.instruments() is first
|
|
assert len(fake_meter.instruments) == 10
|
|
|
|
def test_reset_rebuilds(self, fake_meter):
|
|
first = metrics.instruments()
|
|
metrics.reset()
|
|
with patch("perf.meter", lambda name: fake_meter):
|
|
assert metrics.instruments() is not first
|
|
|
|
|
|
class TestHelpers:
|
|
def test_dispatch_counts_times_and_balances_inflight(self, fake_meter):
|
|
with metrics.dispatch(H1, "embed"):
|
|
held = list(_calls(fake_meter, "stack_llm_inflight"))
|
|
assert held == [(1, {"host": H1})]
|
|
assert _calls(fake_meter, "stack_llm_dispatch_total") == [
|
|
(1, {"host": H1, "kind": "embed"})
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_inflight") == [
|
|
(1, {"host": H1}),
|
|
(-1, {"host": H1}),
|
|
]
|
|
(seconds, attrs) = _calls(fake_meter, "stack_llm_dispatch_seconds")[0]
|
|
assert seconds >= 0
|
|
assert attrs == {"host": H1, "kind": "embed"}
|
|
|
|
def test_dispatch_releases_inflight_on_error(self, fake_meter):
|
|
with pytest.raises(RuntimeError):
|
|
with metrics.dispatch(H2, "generate"):
|
|
raise RuntimeError("boom")
|
|
assert _calls(fake_meter, "stack_llm_inflight")[-1] == (-1, {"host": H2})
|
|
assert len(_calls(fake_meter, "stack_llm_dispatch_seconds")) == 1
|
|
|
|
def test_embedded(self, fake_meter):
|
|
metrics.embedded(H1, 64, 1.5)
|
|
assert _calls(fake_meter, "stack_llm_embedded_texts_total") == [
|
|
(64, {"host": H1})
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_embed_batch_seconds") == [
|
|
(1.5, {"host": H1})
|
|
]
|
|
|
|
def test_chat_stage_and_outcome(self, fake_meter):
|
|
metrics.chat_stage("retrieve", 0.25)
|
|
metrics.chat_outcome("ok")
|
|
assert _calls(fake_meter, "stack_llm_chat_seconds") == [
|
|
(0.25, {"stage": "retrieve"})
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_chat_total") == [(1, {"outcome": "ok"})]
|
|
|
|
def test_chat_tokens_both_directions(self, fake_meter):
|
|
metrics.chat_tokens(120, 40)
|
|
assert _calls(fake_meter, "stack_llm_chat_tokens_total") == [
|
|
(120, {"direction": "prompt"}),
|
|
(40, {"direction": "completion"}),
|
|
]
|
|
|
|
def test_chat_tokens_skips_zero(self, fake_meter):
|
|
metrics.chat_tokens(0, 0)
|
|
assert _calls(fake_meter, "stack_llm_chat_tokens_total") == []
|
|
|
|
def test_indexed_item_with_chunks(self, fake_meter):
|
|
metrics.indexed_item("comments", "indexed", 7)
|
|
assert _calls(fake_meter, "stack_llm_indexed_items_total") == [
|
|
(1, {"collection": "comments", "outcome": "indexed"})
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_indexed_chunks_total") == [
|
|
(7, {"collection": "comments"})
|
|
]
|
|
|
|
def test_indexed_item_skipped_records_no_chunks(self, fake_meter):
|
|
metrics.indexed_item("rules", "skipped")
|
|
assert _calls(fake_meter, "stack_llm_indexed_items_total") == [
|
|
(1, {"collection": "rules", "outcome": "skipped"})
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_indexed_chunks_total") == []
|
|
|
|
|
|
class TestPool:
|
|
def test_acquire_dispatches_as_embed(self, fake_meter):
|
|
with HostPool([H1]).acquire() as host:
|
|
assert host == H1
|
|
assert _calls(fake_meter, "stack_llm_dispatch_total") == [
|
|
(1, {"host": H1, "kind": "embed"})
|
|
]
|
|
|
|
def test_acquire_generation_dispatches_as_generate(self, fake_meter):
|
|
with HostPool([H1, H2], vram_gb={H2: 24.0}).acquire_generation() as host:
|
|
assert host == H2
|
|
assert _calls(fake_meter, "stack_llm_dispatch_total") == [
|
|
(1, {"host": H2, "kind": "generate"})
|
|
]
|
|
|
|
@patch("llm.pool.httpx.Client")
|
|
def test_embed_texts_counts_texts_per_host(self, MockClient, fake_meter):
|
|
resp = MagicMock()
|
|
resp.json.return_value = {"embeddings": [[0.1], [0.2]]}
|
|
MockClient.return_value.__enter__.return_value.post.return_value = resp
|
|
|
|
out = embed_texts(HostPool([H1]), "embed", ["a", "b"], batch_size=2)
|
|
|
|
assert out == [[0.1], [0.2]]
|
|
assert _calls(fake_meter, "stack_llm_embedded_texts_total") == [
|
|
(2, {"host": H1})
|
|
]
|
|
(seconds, attrs) = _calls(fake_meter, "stack_llm_embed_batch_seconds")[0]
|
|
assert seconds >= 0
|
|
assert attrs == {"host": H1}
|
|
|
|
|
|
def _pool_mock():
|
|
pool = MagicMock()
|
|
pool.acquire_generation.return_value.__enter__.return_value = H1
|
|
pool.vram.return_value = 24.0
|
|
pool.serves.return_value = True
|
|
return pool
|
|
|
|
|
|
def _stream(MockClient, lines):
|
|
client = MockClient.return_value.__enter__.return_value
|
|
resp = client.stream.return_value.__enter__.return_value
|
|
resp.iter_lines.return_value = iter(lines)
|
|
return resp
|
|
|
|
|
|
class TestChat:
|
|
@patch("llm.rag.valuation_evidence", return_value=None)
|
|
@patch("llm.rag.lineage_evidence", return_value=None)
|
|
@patch("llm.rag.httpx.Client")
|
|
@patch("llm.rag.retrieve", return_value=[])
|
|
def test_ok_records_stages_outcome_and_tokens(
|
|
self, _retrieve, MockClient, _lin, _ev, fake_meter
|
|
):
|
|
from llm.rag import stream_answer
|
|
|
|
_stream(
|
|
MockClient,
|
|
[
|
|
'{"message":{"content":"hi"},"done":false}',
|
|
'{"message":{"content":""},"done":true,'
|
|
'"prompt_eval_count":120,"eval_count":40}',
|
|
],
|
|
)
|
|
|
|
events = list(stream_answer("q", cfg=CFG, pool=_pool_mock()))
|
|
|
|
assert events[-1] == {"type": "done"}
|
|
assert [
|
|
a[1]["stage"] for a in _calls(fake_meter, "stack_llm_chat_seconds")
|
|
] == [
|
|
"retrieve",
|
|
"generate",
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_chat_total") == [(1, {"outcome": "ok"})]
|
|
assert _calls(fake_meter, "stack_llm_chat_tokens_total") == [
|
|
(120, {"direction": "prompt"}),
|
|
(40, {"direction": "completion"}),
|
|
]
|
|
|
|
@patch("llm.rag.valuation_evidence", return_value=None)
|
|
@patch("llm.rag.lineage_evidence", return_value=None)
|
|
@patch("llm.rag.httpx.Client")
|
|
@patch("llm.rag.retrieve", return_value=[])
|
|
def test_missing_token_counts_are_not_recorded(
|
|
self, _retrieve, MockClient, _lin, _ev, fake_meter
|
|
):
|
|
from llm.rag import stream_answer
|
|
|
|
_stream(MockClient, ['{"message":{"content":""},"done":true}'])
|
|
list(stream_answer("q", cfg=CFG, pool=_pool_mock()))
|
|
assert _calls(fake_meter, "stack_llm_chat_tokens_total") == []
|
|
|
|
@patch("llm.rag.valuation_evidence", return_value=None)
|
|
@patch("llm.rag.lineage_evidence", return_value=None)
|
|
@patch("llm.rag.httpx.Client")
|
|
@patch("llm.rag.retrieve", return_value=[])
|
|
def test_error_outcome_and_no_generate_stage(
|
|
self, _retrieve, MockClient, _lin, _ev, fake_meter
|
|
):
|
|
from llm.rag import stream_answer
|
|
|
|
resp = _stream(MockClient, [])
|
|
resp.raise_for_status.side_effect = RuntimeError("ollama down")
|
|
|
|
with pytest.raises(RuntimeError, match="ollama down"):
|
|
list(stream_answer("q", cfg=CFG, pool=_pool_mock()))
|
|
|
|
assert _calls(fake_meter, "stack_llm_chat_total") == [(1, {"outcome": "error"})]
|
|
stages = [a[1]["stage"] for a in _calls(fake_meter, "stack_llm_chat_seconds")]
|
|
assert stages == ["retrieve"]
|
|
|
|
|
|
class TestIndexer:
|
|
def _run(self, refs, state_rows, fake_meter):
|
|
from llm.index import index_refs
|
|
|
|
engine = MagicMock()
|
|
conn = engine.begin.return_value.__enter__.return_value
|
|
|
|
def fake_execute(clause, *a, **k):
|
|
r = MagicMock()
|
|
sql = str(clause)
|
|
r.fetchall.return_value = state_rows if "FROM index_state" in sql else []
|
|
return r
|
|
|
|
conn.execute.side_effect = fake_execute
|
|
with (
|
|
patch("llm.index.migrate"),
|
|
patch("llm.index.vectorstore", return_value=MagicMock()),
|
|
patch("llm.index._code_index", return_value={}),
|
|
patch("llm.index.embed_texts", return_value=[[0.1]]),
|
|
patch("llm.index.enrich_pdf_pages", side_effect=lambda doc, chunks: chunks),
|
|
):
|
|
return index_refs(
|
|
refs,
|
|
collection="comments",
|
|
cfg=CFG,
|
|
pool=MagicMock(),
|
|
engine=engine,
|
|
)
|
|
|
|
def _ref(self, key, fingerprint, doc):
|
|
from llm.source import DocRef
|
|
|
|
return DocRef(
|
|
key=key,
|
|
collection="comments",
|
|
docket="D",
|
|
fingerprint=fingerprint,
|
|
load=lambda: doc,
|
|
)
|
|
|
|
def test_indexed_and_skipped_items(self, fake_meter):
|
|
from llm.chunk import Doc, content_hash
|
|
|
|
body = Doc(key="K1", text="Some body text.", metadata={"docket": "D"})
|
|
same = Doc(key="K2", text="Unchanged text.", metadata={"docket": "D"})
|
|
blank = Doc(key="K3", text=" ", metadata={"docket": "D"})
|
|
stats = self._run(
|
|
[
|
|
self._ref("K1", "fp-new", body),
|
|
self._ref("K2", "fp-changed", same),
|
|
self._ref("K3", "fp-blank", blank),
|
|
self._ref("K4", "fp-old", body),
|
|
],
|
|
[("K2", content_hash(same.text), "fp-stale"), ("K4", "h", "fp-old")],
|
|
fake_meter,
|
|
)
|
|
|
|
assert (stats["indexed"], stats["chunks"]) == (1, 1)
|
|
assert _calls(fake_meter, "stack_llm_indexed_items_total") == [
|
|
(1, {"collection": "comments", "outcome": "indexed"}),
|
|
(1, {"collection": "comments", "outcome": "skipped"}),
|
|
(1, {"collection": "comments", "outcome": "skipped"}),
|
|
(1, {"collection": "comments", "outcome": "skipped"}),
|
|
]
|
|
assert _calls(fake_meter, "stack_llm_indexed_chunks_total") == [
|
|
(1, {"collection": "comments"})
|
|
]
|
|
|
|
def test_unindexable_ref_counts_skipped(self, fake_meter):
|
|
from llm.chunk import Doc
|
|
|
|
doc = Doc(key="K9", text="text that chunks to nothing", metadata={})
|
|
with patch("llm.index.chunk_doc", return_value=[]):
|
|
stats = self._run([self._ref("K9", "fp", doc)], [], fake_meter)
|
|
assert stats["skipped"] == 1
|
|
assert _calls(fake_meter, "stack_llm_indexed_items_total") == [
|
|
(1, {"collection": "comments", "outcome": "skipped"})
|
|
]
|
|
|
|
|
|
class TestNoopMeter:
|
|
"""Telemetry is disabled in CI — the real ``perf`` mock must absorb
|
|
every call without the caller noticing."""
|
|
|
|
def test_every_helper_is_inert(self, monkeypatch):
|
|
monkeypatch.setenv("STACK_TELEMETRY", "false")
|
|
from perf import meter
|
|
from perf._meter import MockMeter, _MockInstrument
|
|
|
|
assert isinstance(meter(metrics.METER_NAME), MockMeter)
|
|
assert isinstance(metrics.instruments().dispatch_total, _MockInstrument)
|
|
with metrics.dispatch(H1, "embed"):
|
|
pass
|
|
metrics.embedded(H1, 3, 0.1)
|
|
metrics.chat_stage("generate", 0.2)
|
|
metrics.chat_outcome("ok")
|
|
metrics.chat_tokens(1, 2)
|
|
metrics.indexed_item("rules", "indexed", 4)
|