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.
308 lines
12 KiB
Python
308 lines
12 KiB
Python
"""Incremental, resumable embedding indexer.
|
|
|
|
Per doc: compare ``content_hash`` against ``index_state``; skip when
|
|
unchanged, else delete the item's old chunks in this collection (by
|
|
``cmetadata->>'item_key'`` + ``collection_id``), embed via the host
|
|
pool, upsert with deterministic ids, and record the new hash — three
|
|
commits per doc, ordered embed → delete → add → record-state, so a
|
|
killed run reprocesses any half-done doc on the next pass (its hash was
|
|
never recorded) rather than losing work.
|
|
|
|
The metadata delete runs after ``vectorstore()`` has already constructed
|
|
the ``PGVector`` store: the installed ``langchain-postgres==0.0.17``
|
|
creates ``langchain_pg_embedding``/``langchain_pg_collection`` and the
|
|
collection row synchronously in ``PGVector.__post_init__`` (sync mode),
|
|
so the table always exists by the time any ``DELETE`` runs — but the
|
|
delete is still wrapped defensively in case a future version defers
|
|
table creation to first ``add_embeddings()``.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import os
|
|
from typing import Iterable
|
|
|
|
from sqlalchemy import create_engine, text
|
|
from sqlalchemy.engine import Engine
|
|
from sqlalchemy.exc import ProgrammingError
|
|
|
|
from llm import metrics
|
|
from llm.chunk import Doc, chunk_doc, content_hash
|
|
from llm.config import LlmConfig, pg_url
|
|
from llm.migrate import ensure_hnsw, migrate
|
|
from llm.pages import enrich_pdf_pages
|
|
from llm.pool import HostPool, PoolEmbeddings, embed_texts
|
|
from llm.source import DocRef
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
|
|
def _code_index(cfg: LlmConfig) -> dict[str, tuple[str, ...]]:
|
|
"""Build ``pfs.anchors.code_family_index`` the same way ``stack llm
|
|
restamp`` does — hand families, plus ``pfs.families.refresh_from``
|
|
against the DuckDB replica when it exists — so a chunk stamped here
|
|
at index time carries the same ``families`` metadata a later restamp
|
|
would compute for it (Ruling F1: the two paths must agree by
|
|
construction, not by luck)."""
|
|
from pfs.anchors import code_family_index
|
|
from pfs.families import FAMILIES, refresh_from
|
|
|
|
if os.path.exists(cfg.duckdb_replica):
|
|
import duckdb
|
|
|
|
con = duckdb.connect(cfg.duckdb_replica, read_only=True)
|
|
try:
|
|
refresh_from(con)
|
|
finally:
|
|
con.close()
|
|
return code_family_index(FAMILIES)
|
|
|
|
|
|
def _engine(cfg: LlmConfig) -> Engine: # pragma: no cover — needs a live DB
|
|
return create_engine(pg_url(cfg))
|
|
|
|
|
|
def vectorstore(collection: str, cfg: LlmConfig, pool: HostPool):
|
|
from langchain_postgres import PGVector # pragma: no cover — needs live DB
|
|
|
|
return PGVector( # pragma: no cover — needs a live pgvector DB
|
|
embeddings=PoolEmbeddings(pool, cfg.embed_model),
|
|
collection_name=collection,
|
|
connection=pg_url(cfg),
|
|
embedding_length=cfg.embed_dim,
|
|
use_jsonb=True,
|
|
)
|
|
|
|
|
|
def _state(engine: Engine, collection: str) -> dict[str, tuple[str, str]]:
|
|
with engine.begin() as conn:
|
|
rows = conn.execute(
|
|
text(
|
|
"SELECT item_key, content_hash, fingerprint FROM index_state "
|
|
"WHERE collection = :c"
|
|
),
|
|
{"c": collection},
|
|
).fetchall()
|
|
return {r[0]: (r[1], r[2] or "") for r in rows}
|
|
|
|
|
|
def docket_complete(engine: Engine, collection: str) -> dict[str, str]:
|
|
"""docket → sealed_at for dockets fully indexed under that seal."""
|
|
with engine.begin() as conn:
|
|
rows = conn.execute(
|
|
text(
|
|
"SELECT docket, sealed_at FROM index_docket_state WHERE collection = :c"
|
|
),
|
|
{"c": collection},
|
|
).fetchall()
|
|
return {r[0]: r[1] for r in rows}
|
|
|
|
|
|
def _delete_old_chunks(engine: Engine, item_key: str, collection: str) -> None:
|
|
"""Delete *item_key*'s chunks in *collection* only.
|
|
|
|
``index_state``'s PK is (item_key, collection): the same item may be
|
|
indexed into several collections, so an unscoped delete-by-item_key
|
|
would silently drop another collection's chunks. Column names match
|
|
langchain-postgres 0.0.17: ``langchain_pg_collection.uuid`` is the PK
|
|
that ``langchain_pg_embedding.collection_id`` references.
|
|
"""
|
|
try:
|
|
with engine.begin() as conn:
|
|
conn.execute(
|
|
text(
|
|
"DELETE FROM langchain_pg_embedding "
|
|
"WHERE cmetadata->>'item_key' = :k "
|
|
"AND collection_id = (SELECT uuid "
|
|
"FROM langchain_pg_collection WHERE name = :c)"
|
|
),
|
|
{"k": item_key, "c": collection},
|
|
)
|
|
except ProgrammingError:
|
|
# Fresh database, store tables not created yet — nothing to delete.
|
|
log.debug("langchain_pg_embedding not present yet; skipping delete")
|
|
|
|
|
|
def _record_state(
|
|
engine: Engine, key: str, collection: str, h: str, n: int, fingerprint: str
|
|
) -> None:
|
|
"""Upsert this item's row in ``index_state``.
|
|
|
|
``h == "" and n == 0`` is the *stamped-final* form: an item with no
|
|
text at all. It is a finished state, not a failure — no chunks exist
|
|
to delete or embed — and the empty hash can never collide with a real
|
|
``content_hash`` (a sha256 hexdigest), so the next run recognises it
|
|
by fingerprint alone and never loads it again. If text does appear
|
|
later the fingerprint changes and the item is embedded normally.
|
|
"""
|
|
with engine.begin() as conn:
|
|
conn.execute(
|
|
text(
|
|
"INSERT INTO index_state "
|
|
"(item_key, collection, content_hash, chunk_count, fingerprint) "
|
|
"VALUES (:k, :c, :h, :n, :f) "
|
|
"ON CONFLICT (item_key, collection) DO UPDATE SET "
|
|
"content_hash = :h, chunk_count = :n, fingerprint = :f, "
|
|
"indexed_at = now()"
|
|
),
|
|
{"k": key, "c": collection, "h": h, "n": n, "f": fingerprint},
|
|
)
|
|
|
|
|
|
def index_refs(
|
|
refs: Iterable[DocRef],
|
|
*,
|
|
collection: str,
|
|
cfg: LlmConfig,
|
|
pool: HostPool,
|
|
force: bool = False,
|
|
sealed: dict[str, str] | None = None,
|
|
mark_complete: bool = True,
|
|
engine: Engine | None = None,
|
|
) -> dict:
|
|
"""Fingerprint-first incremental indexing.
|
|
|
|
Per ref: (1) fingerprint equal to the stored one → skip without
|
|
loading; (2) load, hash the text; hash equal → record the new
|
|
fingerprint, skip embedding; (3) chunk, locate PDF pages, embed,
|
|
write. PDFs are opened only on path (3). A ref that loads blank is
|
|
*stamped final* (see :func:`_record_state`) rather than treated as a
|
|
failure — plenty of real comments carry no text at all.
|
|
|
|
After the loop, a docket in *sealed* is recorded in
|
|
``index_docket_state`` — so the next run does not list it at all —
|
|
only when every one of its refs was accounted for. A ref that loaded
|
|
*with* text but chunked to nothing is *unindexable* and blocks
|
|
completion, because a seal must never claim coverage this run did not
|
|
achieve (an exception on the embed/write path aborts the run and so
|
|
blocks it too). *mark_complete* False (a ``--limit`` run is never
|
|
complete) blocks it as well. Pass *engine* to reuse a caller's engine
|
|
instead of building a second one.
|
|
"""
|
|
engine = engine if engine is not None else _engine(cfg)
|
|
migrate(engine)
|
|
pool.check(cfg.embed_model)
|
|
store = vectorstore(collection, cfg, pool)
|
|
seen = _state(engine, collection)
|
|
code_index = _code_index(cfg)
|
|
stats = {
|
|
"indexed": 0,
|
|
"skipped": 0,
|
|
"chunks": 0,
|
|
"fingerprint_skipped": 0,
|
|
"hash_skipped": 0,
|
|
"docket_complete": 0,
|
|
}
|
|
dockets_seen: set[str] = set()
|
|
unindexable: dict[str, int] = {}
|
|
|
|
def _unindexable(ref: DocRef) -> None:
|
|
stats["skipped"] += 1
|
|
metrics.indexed_item(collection, "skipped")
|
|
if ref.docket:
|
|
unindexable[ref.docket] = unindexable.get(ref.docket, 0) + 1
|
|
|
|
for ref in refs:
|
|
if ref.docket:
|
|
dockets_seen.add(ref.docket)
|
|
prev_hash, prev_fp = seen.get(ref.key, ("", ""))
|
|
if not force and ref.fingerprint and ref.fingerprint == prev_fp:
|
|
stats["fingerprint_skipped"] += 1
|
|
metrics.indexed_item(collection, "skipped")
|
|
continue
|
|
doc = ref.load()
|
|
if doc is None or not doc.text.strip():
|
|
_record_state(engine, ref.key, collection, "", 0, ref.fingerprint)
|
|
stats["skipped"] += 1
|
|
metrics.indexed_item(collection, "skipped")
|
|
continue
|
|
h = content_hash(doc.text)
|
|
if not force and prev_hash == h:
|
|
with engine.begin() as conn:
|
|
conn.execute(
|
|
text(
|
|
"UPDATE index_state SET fingerprint = :f "
|
|
"WHERE item_key = :k AND collection = :c"
|
|
),
|
|
{"f": ref.fingerprint, "k": ref.key, "c": collection},
|
|
)
|
|
stats["hash_skipped"] += 1
|
|
metrics.indexed_item(collection, "skipped")
|
|
continue
|
|
chunks = enrich_pdf_pages(doc, chunk_doc(doc, code_index=code_index))
|
|
if not chunks:
|
|
_unindexable(ref)
|
|
continue
|
|
vectors = embed_texts(pool, cfg.embed_model, [c.text for c in chunks])
|
|
_delete_old_chunks(engine, ref.key, collection)
|
|
store.add_embeddings(
|
|
texts=[c.text for c in chunks],
|
|
embeddings=vectors,
|
|
metadatas=[c.metadata for c in chunks],
|
|
ids=[c.id for c in chunks],
|
|
)
|
|
_record_state(engine, ref.key, collection, h, len(chunks), ref.fingerprint)
|
|
stats["indexed"] += 1
|
|
stats["chunks"] += len(chunks)
|
|
metrics.indexed_item(collection, "indexed", len(chunks))
|
|
# needs 100+ docs in one run to hit this line
|
|
if stats["indexed"] % 100 == 0:
|
|
log.info("indexed %(indexed)s (+%(chunks)s)", stats) # pragma: no cover
|
|
|
|
if mark_complete and sealed:
|
|
for d in sorted(dockets_seen & set(sealed)):
|
|
if unindexable.get(d):
|
|
log.info(
|
|
"docket %s not marked complete: %s unindexable item(s)",
|
|
d,
|
|
unindexable[d],
|
|
)
|
|
continue
|
|
with engine.begin() as conn:
|
|
conn.execute(
|
|
text(
|
|
"INSERT INTO index_docket_state (collection, docket, sealed_at) "
|
|
"VALUES (:c, :d, :s) "
|
|
"ON CONFLICT (collection, docket) DO UPDATE SET "
|
|
"sealed_at = :s, indexed_at = now()"
|
|
),
|
|
{"c": collection, "d": d, "s": sealed[d]},
|
|
)
|
|
stats["docket_complete"] += 1
|
|
|
|
if cfg.build_ann_index:
|
|
ensure_hnsw(engine, cfg.embed_dim)
|
|
else:
|
|
log.info("ANN index skipped (build_ann_index=false); using exact search")
|
|
return stats
|
|
|
|
|
|
def index_docs(
|
|
docs: Iterable[Doc],
|
|
*,
|
|
collection: str,
|
|
cfg: LlmConfig,
|
|
pool: HostPool,
|
|
force: bool = False,
|
|
) -> dict:
|
|
"""Back-compat wrapper: Docs without fingerprints (never fingerprint-skipped)."""
|
|
refs = (
|
|
DocRef(
|
|
key=d.key,
|
|
collection=collection,
|
|
docket=d.metadata.get("docket") or None,
|
|
fingerprint="",
|
|
load=(lambda d=d: d),
|
|
)
|
|
for d in docs
|
|
)
|
|
stats = index_refs(
|
|
refs, collection=collection, cfg=cfg, pool=pool, force=force, sealed=None
|
|
)
|
|
# Fold the fingerprint/hash skip counters into "skipped" so the 3-key
|
|
# contract still means "not embedded" — index_refs itself keeps them
|
|
# separate; this is only a local projection.
|
|
skipped = stats["skipped"] + stats["fingerprint_skipped"] + stats["hash_skipped"]
|
|
return {"indexed": stats["indexed"], "skipped": skipped, "chunks": stats["chunks"]}
|