diff --git a/src/cli/llm.py b/src/cli/llm.py index 307da6d..fa6cdfb 100644 --- a/src/cli/llm.py +++ b/src/cli/llm.py @@ -50,6 +50,7 @@ def index( 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,) @@ -58,10 +59,17 @@ def index( 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(cfg), target) - if (sealed_all and target == "comments") + 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} @@ -77,6 +85,7 @@ def index( force=force, sealed=pending_seals if target == "comments" else None, mark_complete=not limit, + engine=engine, ) if skip: typer.echo( diff --git a/src/llm/index.py b/src/llm/index.py index 79cb93b..9d4e46e 100644 --- a/src/llm/index.py +++ b/src/llm/index.py @@ -110,18 +110,23 @@ def index_refs( 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). After the loop, every - docket in *sealed* that was iterated is recorded in - ``index_docket_state`` so the next run does not list it at all — - unless *mark_complete* is False (a ``--limit`` run is never complete). + write. PDFs are opened only on path (3). 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 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) pool.check(cfg.embed_model) store = vectorstore(collection, cfg, pool) @@ -135,6 +140,12 @@ def index_refs( "docket_complete": 0, } 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: if ref.docket: @@ -145,7 +156,7 @@ def index_refs( continue doc = ref.load() if doc is None or not doc.text.strip(): - stats["skipped"] += 1 + _unindexable(ref) continue h = content_hash(doc.text) if not force and prev_hash == h: @@ -161,7 +172,7 @@ def index_refs( continue chunks = enrich_pdf_pages(doc, chunk_doc(doc)) if not chunks: - stats["skipped"] += 1 + _unindexable(ref) continue vectors = embed_texts(pool, cfg.embed_model, [c.text for c in chunks]) _delete_old_chunks(engine, ref.key, collection) @@ -197,6 +208,13 @@ def index_refs( 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( diff --git a/tests/cli/test_llm_index_sealed.py b/tests/cli/test_llm_index_sealed.py index 9d0b502..c934f89 100644 --- a/tests/cli/test_llm_index_sealed.py +++ b/tests/cli/test_llm_index_sealed.py @@ -76,3 +76,76 @@ def test_force_and_limit_disable_skips_and_completion( assert result.exit_code == 0, result.output assert mock_iter.call_args.kwargs["skip_dockets"] == set() 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 diff --git a/tests/llm/test_index_refs.py b/tests/llm/test_index_refs.py index e9ad1a0..81ea9c9 100644 --- a/tests/llm/test_index_refs.py +++ b/tests/llm/test_index_refs.py @@ -33,7 +33,14 @@ def _ref(fp="fp1", docket="D", loads=None, doc=DOC): 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() engine = MagicMock() @@ -56,7 +63,7 @@ def _run( patch("llm.index.vectorstore", return_value=store), patch("llm.index.embed_texts", return_value=[[0.0] * 3]), 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, ): 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 stats, _, conn, _ = _run([_ref("fp1")], [], sealed={}) 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]) + ]