fix(llm,pfs): thread one code_index into chunk_doc so index and restamp agree (F1)
llm.chunk.chunk_doc scanned the import-time pfs.families.FAMILIES (five hand families) at chunk time, while `stack llm restamp` refreshed from the DuckDB replica first — an indexed chunk's `families` metadata could disagree with what a restamp of the same text would compute. chunk_doc now takes an optional `code_index` (pfs.anchors.code_family_ index(families)) and threads it into every anchor_metadata call, on both the section-chunking and paragraph-chunking paths. llm.index. index_refs builds one such index per run — hand families plus pfs.families.refresh_from against the DuckDB replica when it exists, the same mechanism `stack llm restamp` already used — and passes it to every chunk_doc call for the run. Corrected pfs/anchors.py's docstring, which claimed the two paths already wrote identical values.
This commit is contained in:
@@ -13,6 +13,7 @@ from __future__ import annotations
|
|||||||
import hashlib
|
import hashlib
|
||||||
import re
|
import re
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
from typing import Mapping
|
||||||
|
|
||||||
from pfs.anchors import anchor_metadata
|
from pfs.anchors import anchor_metadata
|
||||||
|
|
||||||
@@ -111,7 +112,11 @@ def _pack(section: str, target: int, overlap: int) -> list[str]:
|
|||||||
|
|
||||||
|
|
||||||
def chunk_doc(
|
def chunk_doc(
|
||||||
doc: Doc, *, target_chars: int = 2000, overlap_chars: int = 200
|
doc: Doc,
|
||||||
|
*,
|
||||||
|
target_chars: int = 2000,
|
||||||
|
overlap_chars: int = 200,
|
||||||
|
code_index: Mapping[str, tuple[str, ...]] | None = None,
|
||||||
) -> list[Chunk]:
|
) -> list[Chunk]:
|
||||||
"""Chunk *doc* into <= target_chars windows with overlap between them.
|
"""Chunk *doc* into <= target_chars windows with overlap between them.
|
||||||
|
|
||||||
@@ -120,13 +125,21 @@ def chunk_doc(
|
|||||||
Every chunk carries ``section`` — the markdown heading it sits under
|
Every chunk carries ``section`` — the markdown heading it sits under
|
||||||
(``""`` when none), which is how comment chunks know their attachment.
|
(``""`` when none), which is how comment chunks know their attachment.
|
||||||
|
|
||||||
|
*code_index* is ``pfs.anchors.code_family_index(families)`` — pass
|
||||||
|
the same one the restamp backfill builds (hand families plus
|
||||||
|
``pfs.families.refresh_from`` against the DuckDB replica) so a chunk
|
||||||
|
stamped at index time and the same chunk restamped later carry
|
||||||
|
identical ``families`` metadata (Ruling F1). ``None`` (the default)
|
||||||
|
falls back to a live scan of the import-time ``pfs.families.FAMILIES``
|
||||||
|
registry — only hand families, no CPT-derived ones.
|
||||||
|
|
||||||
Raises ``ValueError`` when ``overlap_chars >= target_chars`` (an
|
Raises ``ValueError`` when ``overlap_chars >= target_chars`` (an
|
||||||
overlap that large or larger would never let the window advance).
|
overlap that large or larger would never let the window advance).
|
||||||
"""
|
"""
|
||||||
if overlap_chars >= target_chars:
|
if overlap_chars >= target_chars:
|
||||||
raise ValueError("overlap_chars must be smaller than target_chars")
|
raise ValueError("overlap_chars must be smaller than target_chars")
|
||||||
if doc.paragraphs:
|
if doc.paragraphs:
|
||||||
return _chunk_paragraphs(doc, target_chars, overlap_chars)
|
return _chunk_paragraphs(doc, target_chars, overlap_chars, code_index)
|
||||||
body = _CONTROL.sub("", _FRONTMATTER.sub("", doc.text)).strip()
|
body = _CONTROL.sub("", _FRONTMATTER.sub("", doc.text)).strip()
|
||||||
if not body:
|
if not body:
|
||||||
return []
|
return []
|
||||||
@@ -148,14 +161,19 @@ def chunk_doc(
|
|||||||
# The heading is scanned with the piece: a section titled
|
# The heading is scanned with the piece: a section titled
|
||||||
# "## G0556 — APCM" is about G0556 all the way down, but
|
# "## G0556 — APCM" is about G0556 all the way down, but
|
||||||
# only its first chunk repeats the code.
|
# only its first chunk repeats the code.
|
||||||
**anchor_metadata(f"{heading} {piece}"),
|
**anchor_metadata(f"{heading} {piece}", code_index=code_index),
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
for seq, (heading, piece) in enumerate(pieces)
|
for seq, (heading, piece) in enumerate(pieces)
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
def _chunk_paragraphs(doc: Doc, target: int, overlap: int) -> list[Chunk]:
|
def _chunk_paragraphs(
|
||||||
|
doc: Doc,
|
||||||
|
target: int,
|
||||||
|
overlap: int,
|
||||||
|
code_index: Mapping[str, tuple[str, ...]] | None = None,
|
||||||
|
) -> list[Chunk]:
|
||||||
"""Greedy pack of whole paragraphs; an oversized paragraph is
|
"""Greedy pack of whole paragraphs; an oversized paragraph is
|
||||||
hard-wrapped with overlap, every piece keeping its own anchor."""
|
hard-wrapped with overlap, every piece keeping its own anchor."""
|
||||||
packed: list[tuple[str, Paragraph, Paragraph]] = [] # text, first, last
|
packed: list[tuple[str, Paragraph, Paragraph]] = [] # text, first, last
|
||||||
@@ -190,7 +208,7 @@ def _chunk_paragraphs(doc: Doc, target: int, overlap: int) -> list[Chunk]:
|
|||||||
"item_key": doc.key,
|
"item_key": doc.key,
|
||||||
"seq": str(seq),
|
"seq": str(seq),
|
||||||
"section": "",
|
"section": "",
|
||||||
**anchor_metadata(text),
|
**anchor_metadata(text, code_index=code_index),
|
||||||
"p_id": str(first.p_id),
|
"p_id": str(first.p_id),
|
||||||
"p_id_last": str(last.p_id),
|
"p_id_last": str(last.p_id),
|
||||||
"page": str(first.page),
|
"page": str(first.page),
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ table creation to first ``add_embeddings()``.
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
from typing import Iterable
|
from typing import Iterable
|
||||||
|
|
||||||
from sqlalchemy import create_engine, text
|
from sqlalchemy import create_engine, text
|
||||||
@@ -36,6 +37,27 @@ from llm.source import DocRef
|
|||||||
log = logging.getLogger(__name__)
|
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
|
def _engine(cfg: LlmConfig) -> Engine: # pragma: no cover — needs a live DB
|
||||||
return create_engine(pg_url(cfg))
|
return create_engine(pg_url(cfg))
|
||||||
|
|
||||||
@@ -162,6 +184,7 @@ def index_refs(
|
|||||||
pool.check(cfg.embed_model)
|
pool.check(cfg.embed_model)
|
||||||
store = vectorstore(collection, cfg, pool)
|
store = vectorstore(collection, cfg, pool)
|
||||||
seen = _state(engine, collection)
|
seen = _state(engine, collection)
|
||||||
|
code_index = _code_index(cfg)
|
||||||
stats = {
|
stats = {
|
||||||
"indexed": 0,
|
"indexed": 0,
|
||||||
"skipped": 0,
|
"skipped": 0,
|
||||||
@@ -202,7 +225,7 @@ def index_refs(
|
|||||||
)
|
)
|
||||||
stats["hash_skipped"] += 1
|
stats["hash_skipped"] += 1
|
||||||
continue
|
continue
|
||||||
chunks = enrich_pdf_pages(doc, chunk_doc(doc))
|
chunks = enrich_pdf_pages(doc, chunk_doc(doc, code_index=code_index))
|
||||||
if not chunks:
|
if not chunks:
|
||||||
_unindexable(ref)
|
_unindexable(ref)
|
||||||
continue
|
continue
|
||||||
|
|||||||
@@ -1,7 +1,12 @@
|
|||||||
"""Anchor metadata for a piece of text: the codes it names, the families
|
"""Anchor metadata for a piece of text: the codes it names, the families
|
||||||
those codes belong to, and the element slugs its phrases match. One
|
those codes belong to, and the element slugs its phrases match. One
|
||||||
function, used by the indexer at chunk time and by the restamp backfill,
|
function, used by the indexer at chunk time (``llm.chunk.chunk_doc``,
|
||||||
so both write identical values. Pure."""
|
called from ``llm.index.index_refs``) and by the restamp backfill
|
||||||
|
(``llm.restamp.restamp``). Both callers build one ``code_family_index``
|
||||||
|
per run the same way — hand families, plus ``pfs.families.refresh_from``
|
||||||
|
against the DuckDB replica when it exists — and pass it through as
|
||||||
|
``code_index``, so the two paths compute identical values from the same
|
||||||
|
family snapshot by construction, not by coincidence (Ruling F1). Pure."""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
|||||||
@@ -213,3 +213,41 @@ class TestAnchorMetadata:
|
|||||||
assert "elements" in chunk.metadata
|
assert "elements" in chunk.metadata
|
||||||
assert chunk.metadata["codes"] == "G0556"
|
assert chunk.metadata["codes"] == "G0556"
|
||||||
assert chunk.metadata["families"] == "APCM"
|
assert chunk.metadata["families"] == "APCM"
|
||||||
|
|
||||||
|
|
||||||
|
class TestCodeIndexParam:
|
||||||
|
"""F1: `chunk_doc(..., code_index=...)` must thread the caller's
|
||||||
|
prebuilt index into every `anchor_metadata` call — the same index
|
||||||
|
`stack llm restamp` builds — instead of scanning the import-time
|
||||||
|
`pfs.families.FAMILIES` (hand families only)."""
|
||||||
|
|
||||||
|
def test_code_index_stamps_a_family_absent_from_hand_families(self):
|
||||||
|
# "90001" belongs to no hand family — without an index, the
|
||||||
|
# chunk's families metadata is empty.
|
||||||
|
text = "Sample text about 90001."
|
||||||
|
without = chunk_doc(_doc(text))
|
||||||
|
assert without[0].metadata["families"] == ""
|
||||||
|
|
||||||
|
idx = {"90001": ("CPT-HEADING-X",)}
|
||||||
|
with_idx = chunk_doc(_doc(text), code_index=idx)
|
||||||
|
assert with_idx[0].metadata["families"] == "CPT-HEADING-X"
|
||||||
|
|
||||||
|
def test_default_none_behaviour_is_unchanged(self):
|
||||||
|
text = "Use 99439; consent; per calendar month."
|
||||||
|
chunks = chunk_doc(_doc(text))
|
||||||
|
assert chunks[0].metadata["codes"] == "99439"
|
||||||
|
assert chunks[0].metadata["families"] == "CCM"
|
||||||
|
|
||||||
|
def test_code_index_threads_through_paragraph_chunks(self):
|
||||||
|
from llm.chunk import Paragraph
|
||||||
|
|
||||||
|
paras = (Paragraph(1, 100, 1, "90001 needs a heading family."),)
|
||||||
|
doc = Doc(
|
||||||
|
key="R", text="90001 needs a heading family.", metadata={}, paragraphs=paras
|
||||||
|
)
|
||||||
|
without = chunk_doc(doc)
|
||||||
|
assert without[0].metadata["families"] == ""
|
||||||
|
|
||||||
|
idx = {"90001": ("CPT-HEADING-X",)}
|
||||||
|
with_idx = chunk_doc(doc, code_index=idx)
|
||||||
|
assert with_idx[0].metadata["families"] == "CPT-HEADING-X"
|
||||||
|
|||||||
@@ -48,6 +48,7 @@ def _run(
|
|||||||
mark_complete=True,
|
mark_complete=True,
|
||||||
complete_rows=(),
|
complete_rows=(),
|
||||||
chunk_doc=None,
|
chunk_doc=None,
|
||||||
|
code_index=None,
|
||||||
):
|
):
|
||||||
store = MagicMock()
|
store = MagicMock()
|
||||||
engine = MagicMock()
|
engine = MagicMock()
|
||||||
@@ -72,6 +73,10 @@ def _run(
|
|||||||
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=lambda d, c: c) as enrich,
|
||||||
patch("llm.index.HostPool") as MockPool,
|
patch("llm.index.HostPool") as MockPool,
|
||||||
|
# Real `_code_index` opens the live DuckDB replica — never in a
|
||||||
|
# unit test; every test here gets a fixed (default empty) index
|
||||||
|
# unless it asks for one.
|
||||||
|
patch("llm.index._code_index", return_value=code_index or {}),
|
||||||
_maybe_chunk_doc(chunk_doc),
|
_maybe_chunk_doc(chunk_doc),
|
||||||
):
|
):
|
||||||
MockPool.return_value.check.return_value = ["http://h1:11434"]
|
MockPool.return_value.check.return_value = ["http://h1:11434"]
|
||||||
@@ -211,8 +216,81 @@ def test_sealed_docket_not_marked_when_a_ref_yields_no_chunks():
|
|||||||
"""Text that chunks to nothing is a real failure — it blocks the seal
|
"""Text that chunks to nothing is a real failure — it blocks the seal
|
||||||
and is not stamped."""
|
and is not stamped."""
|
||||||
stats, _, conn, _ = _run(
|
stats, _, conn, _ = _run(
|
||||||
[_ref("fp1")], [], sealed={"D": "s"}, chunk_doc=lambda d: []
|
[_ref("fp1")], [], sealed={"D": "s"}, chunk_doc=lambda d, **k: []
|
||||||
)
|
)
|
||||||
assert stats["skipped"] == 1 and stats["docket_complete"] == 0
|
assert stats["skipped"] == 1 and stats["docket_complete"] == 0
|
||||||
assert _state_inserts(conn) == []
|
assert _state_inserts(conn) == []
|
||||||
assert _docket_inserts(conn) == []
|
assert _docket_inserts(conn) == []
|
||||||
|
|
||||||
|
|
||||||
|
# ── F1: indexer and restamp must write the same `families` ──────────
|
||||||
|
|
||||||
|
|
||||||
|
def test_index_refs_threads_code_index_into_chunk_doc():
|
||||||
|
"""`index_refs` passes the same `code_index` it built to every
|
||||||
|
`chunk_doc` call — a code that maps to a CPT-derived family in the
|
||||||
|
index built for this run must show up in the chunk's metadata."""
|
||||||
|
captured = []
|
||||||
|
|
||||||
|
def fake_chunk_doc(doc, **kwargs):
|
||||||
|
from llm.chunk import Chunk
|
||||||
|
|
||||||
|
captured.append(kwargs.get("code_index"))
|
||||||
|
return [Chunk(id="c1", text=doc.text, metadata={})]
|
||||||
|
|
||||||
|
fixed_index = {"99490": ("CCM",)}
|
||||||
|
stats, store, conn, enrich = _run(
|
||||||
|
[_ref("fp2")],
|
||||||
|
[("K1", "stale", "fp1")],
|
||||||
|
chunk_doc=fake_chunk_doc,
|
||||||
|
code_index=fixed_index,
|
||||||
|
)
|
||||||
|
assert stats["indexed"] == 1
|
||||||
|
assert captured == [fixed_index]
|
||||||
|
|
||||||
|
|
||||||
|
def test_code_index_builds_once_per_run_via_refresh_from_replica():
|
||||||
|
"""`llm.index._code_index` refreshes from the DuckDB replica when it
|
||||||
|
exists — the same mechanism `stack llm restamp` uses — before
|
||||||
|
building `code_family_index`, so both paths compute `families` from
|
||||||
|
the same snapshot (Ruling F1)."""
|
||||||
|
import llm.index as index_mod
|
||||||
|
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
class FakeCon:
|
||||||
|
def close(self):
|
||||||
|
calls.append("closed")
|
||||||
|
|
||||||
|
def fake_connect(path, read_only):
|
||||||
|
calls.append(("connect", path, read_only))
|
||||||
|
return FakeCon()
|
||||||
|
|
||||||
|
def fake_refresh_from(con):
|
||||||
|
calls.append("refreshed")
|
||||||
|
return 0
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("llm.index.os.path.exists", return_value=True),
|
||||||
|
patch("duckdb.connect", side_effect=fake_connect),
|
||||||
|
patch("pfs.families.refresh_from", side_effect=fake_refresh_from),
|
||||||
|
):
|
||||||
|
idx = index_mod._code_index(CFG)
|
||||||
|
|
||||||
|
assert calls == [("connect", CFG.duckdb_replica, True), "refreshed", "closed"]
|
||||||
|
assert isinstance(idx, dict)
|
||||||
|
|
||||||
|
|
||||||
|
def test_code_index_skips_refresh_when_no_replica_file():
|
||||||
|
"""No replica on disk (a fresh checkout, or replica=False) — the
|
||||||
|
index falls back to hand families only, no DuckDB touched."""
|
||||||
|
import llm.index as index_mod
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("llm.index.os.path.exists", return_value=False),
|
||||||
|
patch("pfs.families.refresh_from") as refresh,
|
||||||
|
):
|
||||||
|
idx = index_mod._code_index(CFG)
|
||||||
|
|
||||||
|
refresh.assert_not_called()
|
||||||
|
assert isinstance(idx, dict)
|
||||||
|
|||||||
Reference in New Issue
Block a user