fix(llm): migrate before docket_complete; completion requires zero unindexable items (refs #615)

`stack llm index` read index_docket_state before anything had created
it: migrate() only ran inside index_refs, so the first sealed run on an
un-migrated database died with UndefinedTable. Build one engine in the
CLI, migrate it, and hand it down (index_refs takes an `engine` kwarg);
--force still builds nothing.

index_refs also marked a sealed docket complete when some of its refs
loaded blank or produced no chunks — a seal claiming coverage the run
never achieved. Count unindexable refs per docket and write
index_docket_state only for dockets with none.
This commit is contained in:
kert
2026-09-08 15:31:22 -04:00
parent 3df7a3ff0e
commit 066b304111
4 changed files with 147 additions and 11 deletions

View File

@@ -50,6 +50,7 @@ def index(
from conf.connect import bib from conf.connect import bib
from llm import config as llm_config from llm import config as llm_config
from llm.index import _engine, docket_complete, index_refs from llm.index import _engine, docket_complete, index_refs
from llm.migrate import migrate
from llm.pool import HostPool from llm.pool import HostPool
targets = _COLLECTIONS if collection == "all" else (collection,) targets = _COLLECTIONS if collection == "all" else (collection,)
@@ -58,10 +59,17 @@ def index(
cfg = llm_config.load() cfg = llm_config.load()
store = bib() store = bib()
sealed_all = {} if force else store.sealed_dockets() 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: for target in targets:
complete = ( complete = (
docket_complete(_engine(cfg), target) docket_complete(engine, target)
if (sealed_all and target == "comments") if (engine is not None and target == "comments")
else {} else {}
) )
skip = {d for d, s in sealed_all.items() if complete.get(d) == s} skip = {d for d, s in sealed_all.items() if complete.get(d) == s}
@@ -77,6 +85,7 @@ def index(
force=force, force=force,
sealed=pending_seals if target == "comments" else None, sealed=pending_seals if target == "comments" else None,
mark_complete=not limit, mark_complete=not limit,
engine=engine,
) )
if skip: if skip:
typer.echo( typer.echo(

View File

@@ -110,18 +110,23 @@ def index_refs(
force: bool = False, force: bool = False,
sealed: dict[str, str] | None = None, sealed: dict[str, str] | None = None,
mark_complete: bool = True, mark_complete: bool = True,
engine: Engine | None = None,
) -> dict: ) -> dict:
"""Fingerprint-first incremental indexing. """Fingerprint-first incremental indexing.
Per ref: (1) fingerprint equal to the stored one → skip without Per ref: (1) fingerprint equal to the stored one → skip without
loading; (2) load, hash the text; hash equal → record the new loading; (2) load, hash the text; hash equal → record the new
fingerprint, skip embedding; (3) chunk, locate PDF pages, embed, fingerprint, skip embedding; (3) chunk, locate PDF pages, embed,
write. PDFs are opened only on path (3). After the loop, every write. PDFs are opened only on path (3). After the loop, a docket in
docket in *sealed* that was iterated is recorded in *sealed* is recorded in ``index_docket_state`` — so the next run does
``index_docket_state`` so the next run does not list it at all — not list it at all — only when every one of its refs was accounted
unless *mark_complete* is False (a ``--limit`` run is never complete). for: a ref that loaded blank or produced no chunks is *unindexable*
and blocks completion, because a seal must never claim coverage this
run did not achieve. *mark_complete* False (a ``--limit`` run is
never complete) blocks it too. Pass *engine* to reuse a caller's
engine instead of building a second one.
""" """
engine = _engine(cfg) engine = engine if engine is not None else _engine(cfg)
migrate(engine) migrate(engine)
pool.check(cfg.embed_model) pool.check(cfg.embed_model)
store = vectorstore(collection, cfg, pool) store = vectorstore(collection, cfg, pool)
@@ -135,6 +140,12 @@ def index_refs(
"docket_complete": 0, "docket_complete": 0,
} }
dockets_seen: set[str] = set() dockets_seen: set[str] = set()
unindexable: dict[str, int] = {}
def _unindexable(ref: DocRef) -> None:
stats["skipped"] += 1
if ref.docket:
unindexable[ref.docket] = unindexable.get(ref.docket, 0) + 1
for ref in refs: for ref in refs:
if ref.docket: if ref.docket:
@@ -145,7 +156,7 @@ def index_refs(
continue continue
doc = ref.load() doc = ref.load()
if doc is None or not doc.text.strip(): if doc is None or not doc.text.strip():
stats["skipped"] += 1 _unindexable(ref)
continue continue
h = content_hash(doc.text) h = content_hash(doc.text)
if not force and prev_hash == h: if not force and prev_hash == h:
@@ -161,7 +172,7 @@ def index_refs(
continue continue
chunks = enrich_pdf_pages(doc, chunk_doc(doc)) chunks = enrich_pdf_pages(doc, chunk_doc(doc))
if not chunks: if not chunks:
stats["skipped"] += 1 _unindexable(ref)
continue continue
vectors = embed_texts(pool, cfg.embed_model, [c.text for c in chunks]) vectors = embed_texts(pool, cfg.embed_model, [c.text for c in chunks])
_delete_old_chunks(engine, ref.key, collection) _delete_old_chunks(engine, ref.key, collection)
@@ -197,6 +208,13 @@ def index_refs(
if mark_complete and sealed: if mark_complete and sealed:
for d in sorted(dockets_seen & set(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: with engine.begin() as conn:
conn.execute( conn.execute(
text( text(

View File

@@ -76,3 +76,76 @@ def test_force_and_limit_disable_skips_and_completion(
assert result.exit_code == 0, result.output assert result.exit_code == 0, result.output
assert mock_iter.call_args.kwargs["skip_dockets"] == set() assert mock_iter.call_args.kwargs["skip_dockets"] == set()
assert mock_index.call_args.kwargs["mark_complete"] is False assert mock_index.call_args.kwargs["mark_complete"] is False
@patch("llm.index._engine")
@patch("llm.migrate.migrate")
@patch("llm.index.docket_complete", return_value={})
@patch("llm.source.iter_comment_refs")
@patch("llm.pool.HostPool.from_config")
@patch("llm.index.index_refs")
@patch("conf.connect.bib")
@patch("llm.config.load")
def test_migrate_runs_before_docket_complete_and_engine_is_reused(
mock_load,
mock_bib,
mock_index,
mock_pool,
mock_iter,
mock_complete,
mock_migrate,
mock_engine,
):
"""index_docket_state may not exist yet on the first sealed run."""
order: list[str] = []
mock_migrate.side_effect = lambda e: order.append("migrate")
mock_complete.side_effect = lambda e, c: (order.append("complete"), {})[1]
mock_load.return_value = MagicMock()
store = MagicMock()
store.sealed_dockets.return_value = {"CMS-2019-0111": "s1"}
mock_bib.return_value = store
mock_iter.return_value = iter([])
mock_index.return_value = _STATS
result = runner.invoke(app, ["index"])
assert result.exit_code == 0, result.output
assert order == ["migrate", "complete"]
mock_engine.assert_called_once() # one engine, shared
assert mock_migrate.call_args.args[0] is mock_engine.return_value
assert mock_complete.call_args.args[0] is mock_engine.return_value
assert mock_index.call_args.kwargs["engine"] is mock_engine.return_value
@patch("llm.index._engine")
@patch("llm.migrate.migrate")
@patch("llm.index.docket_complete", return_value={})
@patch("llm.source.iter_comment_refs")
@patch("llm.pool.HostPool.from_config")
@patch("llm.index.index_refs")
@patch("conf.connect.bib")
@patch("llm.config.load")
def test_force_touches_neither_engine_nor_migrate(
mock_load,
mock_bib,
mock_index,
mock_pool,
mock_iter,
mock_complete,
mock_migrate,
mock_engine,
):
mock_load.return_value = MagicMock()
store = MagicMock()
store.sealed_dockets.return_value = {"CMS-2019-0111": "s1"}
mock_bib.return_value = store
mock_iter.return_value = iter([])
mock_index.return_value = _STATS
result = runner.invoke(app, ["index", "--force"])
assert result.exit_code == 0, result.output
mock_engine.assert_not_called()
mock_migrate.assert_not_called()
mock_complete.assert_not_called()
assert mock_index.call_args.kwargs["engine"] is None

View File

@@ -33,7 +33,14 @@ def _ref(fp="fp1", docket="D", loads=None, doc=DOC):
def _run( def _run(
refs, state_rows, *, force=False, sealed=None, mark_complete=True, complete_rows=() refs,
state_rows,
*,
force=False,
sealed=None,
mark_complete=True,
complete_rows=(),
chunks_for=lambda d, c: c,
): ):
store = MagicMock() store = MagicMock()
engine = MagicMock() engine = MagicMock()
@@ -56,7 +63,7 @@ def _run(
patch("llm.index.vectorstore", return_value=store), patch("llm.index.vectorstore", return_value=store),
patch("llm.index.embed_texts", return_value=[[0.0] * 3]), patch("llm.index.embed_texts", return_value=[[0.0] * 3]),
patch("llm.index.ensure_hnsw"), patch("llm.index.ensure_hnsw"),
patch("llm.index.enrich_pdf_pages", side_effect=lambda d, c: c) as enrich, patch("llm.index.enrich_pdf_pages", side_effect=chunks_for) as enrich,
patch("llm.index.HostPool") as MockPool, patch("llm.index.HostPool") as MockPool,
): ):
MockPool.return_value.check.return_value = ["http://h1:11434"] MockPool.return_value.check.return_value = ["http://h1:11434"]
@@ -143,3 +150,32 @@ def test_not_marked_when_mark_complete_false_or_unsealed():
assert stats["docket_complete"] == 0 assert stats["docket_complete"] == 0
stats, _, conn, _ = _run([_ref("fp1")], [], sealed={}) stats, _, conn, _ = _run([_ref("fp1")], [], sealed={})
assert stats["docket_complete"] == 0 assert stats["docket_complete"] == 0
def test_sealed_docket_not_marked_when_a_ref_is_unindexable():
"""A docket with an item that has no indexable text is not complete."""
blank = DocRef(
key="K9", collection="comments", docket="D", fingerprint="x", load=lambda: None
)
stats, _, conn, _ = _run(
[_ref("fp1"), blank], [], sealed={"D": "2026-10-20T00:00:00Z"}
)
assert stats["skipped"] == 1 and stats["indexed"] == 1
assert stats["docket_complete"] == 0
assert not [
c
for c in conn.execute.call_args_list
if "INSERT INTO index_docket_state" in str(c.args[0])
]
def test_sealed_docket_not_marked_when_a_ref_yields_no_chunks():
stats, _, conn, _ = _run(
[_ref("fp1")], [], sealed={"D": "s"}, chunks_for=lambda d, c: []
)
assert stats["skipped"] == 1 and stats["docket_complete"] == 0
assert not [
c
for c in conn.execute.call_args_list
if "INSERT INTO index_docket_state" in str(c.args[0])
]