F2's live check needs to index just the 6 CPT bib items, not the whole corpus — --key was wired only into the rules collection (iter_rule_refs); the corpus branch (_refs_for -> iter_corpus_refs) ignored it entirely, so "--collection corpus --key ..." silently indexed everything. iter_corpus_refs takes an optional keys tuple (same Python-side post-filter iter_rule_refs already uses) and _refs_for now passes --key through on the corpus path too.
206 lines
7.2 KiB
Python
206 lines
7.2 KiB
Python
"""stack llm — local RAG over the library (index + fleet + chat serve)."""
|
|
|
|
import typer
|
|
|
|
app = typer.Typer(no_args_is_help=True)
|
|
|
|
_COLLECTIONS = ("comments", "rules", "corpus")
|
|
|
|
|
|
def _refs_for(
|
|
collection: str, store, docket: str, keys: tuple[str, ...], skip_dockets: set[str]
|
|
):
|
|
from llm.source import (
|
|
ZoteroPdfIndex,
|
|
iter_comment_refs,
|
|
iter_corpus_refs,
|
|
iter_rule_refs,
|
|
)
|
|
|
|
if collection == "comments":
|
|
return iter_comment_refs(store, docket=docket, skip_dockets=skip_dockets)
|
|
if collection == "rules":
|
|
return iter_rule_refs(store, keys=keys)
|
|
from conf import ROOT, path
|
|
|
|
zotero = ZoteroPdfIndex.lazy(
|
|
path("db.zotero"), path("storage.zotero"), ROOT / ".state" / "llm"
|
|
)
|
|
return iter_corpus_refs(store, keys=keys, zotero=zotero)
|
|
|
|
|
|
@app.command()
|
|
def index(
|
|
collection: str = typer.Option(
|
|
"comments", help="Which collection: comments | rules | corpus | all."
|
|
),
|
|
docket: str = typer.Option("", help="Limit comments to one docket id."),
|
|
key: list[str] = typer.Option(
|
|
[], "--key", help="Limit rules/corpus to these item keys."
|
|
),
|
|
force: bool = typer.Option(False, help="Re-embed even when unchanged."),
|
|
limit: int = typer.Option(0, help="Stop after N docs per collection (0 = all)."),
|
|
) -> None:
|
|
"""Embed comments/rules/corpus into pgvector (incremental, resumable).
|
|
|
|
Comments are processed newest-posted first; rules come from their FR
|
|
paragraph anchor maps (run `stack bib fr-grab` first for exact links);
|
|
corpus = every non-comment bib item incl. Zotero-only storage PDFs.
|
|
"""
|
|
import itertools
|
|
|
|
from conf.connect import bib
|
|
from llm import config as llm_config
|
|
from llm.index import _engine, docket_complete, index_refs
|
|
from llm.migrate import migrate
|
|
from llm.pool import HostPool
|
|
|
|
targets = _COLLECTIONS if collection == "all" else (collection,)
|
|
if any(t not in _COLLECTIONS for t in targets):
|
|
raise typer.BadParameter("collection must be comments, rules, corpus or all")
|
|
cfg = llm_config.load()
|
|
store = bib()
|
|
sealed_all = {} if force else store.sealed_dockets()
|
|
# One engine for the whole run, migrated *before* anything reads
|
|
# index_docket_state — that table does not exist until migrate() has
|
|
# run, and the first sealed run on an un-migrated DB gets here first.
|
|
engine = None
|
|
if sealed_all and "comments" in targets:
|
|
engine = _engine(cfg)
|
|
migrate(engine)
|
|
for target in targets:
|
|
complete = (
|
|
docket_complete(engine, target)
|
|
if (engine is not None and target == "comments")
|
|
else {}
|
|
)
|
|
skip = {d for d, s in sealed_all.items() if complete.get(d) == s}
|
|
pending_seals = {d: s for d, s in sealed_all.items() if d not in skip}
|
|
refs = _refs_for(target, store, docket, tuple(key), skip)
|
|
if limit:
|
|
refs = itertools.islice(refs, limit)
|
|
stats = index_refs(
|
|
refs,
|
|
collection=target,
|
|
cfg=cfg,
|
|
pool=HostPool.from_config(cfg),
|
|
force=force,
|
|
sealed=pending_seals if target == "comments" else None,
|
|
mark_complete=not limit,
|
|
engine=engine,
|
|
)
|
|
if skip:
|
|
typer.echo(
|
|
f"{target}: {len(skip)} sealed docket(s) already complete — not listed"
|
|
)
|
|
typer.echo(
|
|
f"{target}: indexed={stats['indexed']} skipped={stats['skipped']} chunks={stats['chunks']} "
|
|
f"fp_skipped={stats['fingerprint_skipped']} hash_skipped={stats['hash_skipped']} "
|
|
f"docket_complete={stats['docket_complete']}"
|
|
)
|
|
|
|
|
|
@app.command()
|
|
def restamp(
|
|
collection: str = typer.Option(
|
|
"comments", help="Which collection: comments | rules | corpus | all."
|
|
),
|
|
batch: int = typer.Option(2000, help="Rows per transaction."),
|
|
dry_run: bool = typer.Option(
|
|
False, "--dry-run", help="Scan and report without writing."
|
|
),
|
|
) -> None:
|
|
"""Backfill codes/families/elements onto already-indexed chunks.
|
|
|
|
Metadata-only: computes `pfs.anchors.anchor_metadata` from each row's
|
|
stored document text and merges it into `cmetadata` with a jsonb `||`,
|
|
never touching the embedding — no GPU time, no re-chunking. Safe to
|
|
re-run: rows already stamped with the current values are skipped.
|
|
|
|
F6: never rebuilds the HNSW index (`ensure_hnsw` ALTERs the vector
|
|
column and takes an ACCESS EXCLUSIVE lock for the life of the
|
|
rebuild) — only the cheap metadata GIN indexes, and not even those
|
|
under `--dry-run`.
|
|
"""
|
|
import os
|
|
|
|
from llm import config as llm_config
|
|
from llm.index import _engine
|
|
from llm.migrate import ensure_metadata_indexes, migrate
|
|
from llm.restamp import restamp as restamp_collection
|
|
from pfs.anchors import code_family_index
|
|
from pfs.families import FAMILIES, refresh_from
|
|
|
|
targets = _COLLECTIONS if collection == "all" else (collection,)
|
|
if any(t not in _COLLECTIONS for t in targets):
|
|
raise typer.BadParameter("collection must be comments, rules, corpus or all")
|
|
cfg = llm_config.load()
|
|
engine = _engine(cfg)
|
|
migrate(engine)
|
|
if not dry_run:
|
|
ensure_metadata_indexes(engine)
|
|
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()
|
|
code_index = code_family_index(FAMILIES)
|
|
for target in targets:
|
|
stats = restamp_collection(
|
|
engine,
|
|
collection=target,
|
|
batch=batch,
|
|
dry_run=dry_run,
|
|
code_index=code_index,
|
|
)
|
|
typer.echo(
|
|
f"{target}: scanned={stats['scanned']} updated={stats['updated']} "
|
|
f"seconds={stats['seconds']}"
|
|
)
|
|
|
|
|
|
@app.command()
|
|
def hosts() -> None:
|
|
"""Show the Ollama fleet: declared VRAM, liveness, models, and which
|
|
host + model would answer a chat right now.
|
|
|
|
Source .env first (`set -a; . ./.env; set +a`) so LLM_OLLAMA_HOSTS
|
|
lists the whole fleet; the compose service gets it from its environment.
|
|
"""
|
|
from llm import config as llm_config
|
|
from llm.pool import HostPool, pick_model
|
|
|
|
cfg = llm_config.load()
|
|
pool = HostPool.from_config(cfg)
|
|
declared = pool.hosts
|
|
try:
|
|
live = pool.check(cfg.instruct_model)
|
|
except RuntimeError as e:
|
|
typer.echo(str(e))
|
|
raise typer.Exit(1)
|
|
for row in pool.status():
|
|
typer.echo(
|
|
f"{row['host']:<32} {row['vram_gb']:>5.0f} GB up "
|
|
f"{', '.join(row['models'])}"
|
|
)
|
|
for h in declared:
|
|
if h not in live:
|
|
state = f"up (no {cfg.instruct_model})" if pool.reachable(h) else "DOWN"
|
|
typer.echo(f"{h:<32} {pool.vram(h):>5.0f} GB {state}")
|
|
with pool.acquire_generation() as host:
|
|
typer.echo(f"generation → {pick_model(cfg, pool, host)} @ {host}")
|
|
|
|
|
|
@app.command()
|
|
def serve(
|
|
host: str = typer.Option("127.0.0.1", help="Bind address."),
|
|
port: int = typer.Option(8000, help="Port."),
|
|
) -> None:
|
|
"""Serve the SSO-guarded chat UI (llm.api:app)."""
|
|
import uvicorn
|
|
|
|
uvicorn.run("llm.api:app", host=host, port=port, log_level="info")
|