Files
stack/src/llm/source.py
kert e154451896 fix(llm): batch attachment-path lookups for rule/corpus fingerprints (refs #680)
source.py: iter_rule_refs and iter_corpus_refs used to run one
_attachment_paths query per item just to build the cheap fingerprint,
before any change was even known. _attachment_paths_by_key replaces
that with one scan of the whole attachments table, used by
iter_rule_refs (only when at least one row lacks an anchor sha) and
iter_corpus_refs (every run). The per-item _attachment_paths calls on
the actual load path (_attachment_text, _rule_text,
_attachment_sections) are untouched — those only run for items that
already changed.

Also: the misplaced "rules" section marker now sits above the
rule-only helpers it should have covered from the start
(_attachment_text/_rule_text/rule_paragraphs/_anchor_doc, previously
stranded in the "comments" section), and iter_comment_refs's inner
_load no longer shadows the enclosing loop's row variable (renamed to
item_row).
2026-09-11 17:32:19 -04:00

616 lines
22 KiB
Python

"""Doc iterators over the bib store + comment extraction tree.
Every source is listed as a :class:`DocRef` first — key, collection,
docket and a cheap *fingerprint* (file stats + ``updated_at``, or the FR
anchor sha) — so the indexer can decide what changed before any file is
read. ``ref.load()`` then does the expensive part (extraction, per-item
SQL) and returns ``None`` for a text-less item. ``iter_*_docs`` are thin
wrappers that load every ref.
Comments: text prefers the #253 extraction output
(``.state/comments/<docket>/<comment_id>/combined.md`` body via
``rex.comments.combine.parse_combined``); falls back to the bib
``abstract`` for attachment-less comments. Corpus: every non-comment item;
attachment text extracted with the same rex extractors, joined with the
abstract.
Docket is derived from the comment's URL (``.../comment/<docket>-<seq>``)
rather than a ``docket:`` tag — the live store tags comments with
``reg-docket:<docket>`` (and a separate, unrelated empty ``docket:`` tag),
so parsing the comment id is the reliable path to a comment's docket.
"""
from __future__ import annotations
import logging
import os
import shutil
import sqlite3
from dataclasses import dataclass
from pathlib import Path
from typing import Callable, Iterator
from bib.dockets import fingerprint_files
from bib.store import Store
from llm.chunk import Doc, Paragraph
log = logging.getLogger(__name__)
_COMMENT_URL_PREFIX = "https://www.regulations.gov/comment/"
_ATTACHMENT_EXT = (".pdf", ".docx", ".doc", ".txt")
@dataclass(frozen=True)
class DocRef:
"""A document the indexer *may* need: enough to decide (key +
fingerprint) without building it. ``load()`` does the expensive part
and returns ``None`` when the item has no text."""
key: str
collection: str
docket: str | None
fingerprint: str
load: Callable[[], "Doc | None"]
def _default_root() -> Path:
from conf import ROOT
return ROOT / ".state" / "comments"
def _year_of(store: Store, item_key: str) -> str:
row = (
store._con()
.execute(
"SELECT t.name FROM tags t "
"JOIN item_tags it ON it.tag_id = t.id "
"JOIN items i ON i.id = it.item_id "
"WHERE i.key = ? AND t.name LIKE 'year:%'",
(item_key,),
)
.fetchone()
)
return row[0].split(":", 1)[1] if row else ""
def _attachment_paths(store: Store, item_key: str) -> list[Path]:
rows = (
store._con()
.execute(
"SELECT a.storage_path FROM attachments a "
"JOIN items i ON i.id = a.item_id WHERE i.key = ?",
(item_key,),
)
.fetchall()
)
return [Path(r[0]) for r in rows]
def _attachment_paths_by_key(store: Store) -> dict[str, list[Path]]:
"""Every item's attachment paths, grouped by key, in one scan.
``iter_rule_refs``/``iter_corpus_refs`` used to call
:func:`_attachment_paths` once per item just to build the
fingerprint — one query per rule/corpus item on every run, before
any change is even known. One table scan replaces that N+1."""
rows = (
store._con()
.execute(
"SELECT i.key, a.storage_path FROM attachments a JOIN items i ON i.id = a.item_id"
)
.fetchall()
)
by_key: dict[str, list[Path]] = {}
for key, path in rows:
by_key.setdefault(key, []).append(Path(path))
return by_key
# ── comments ──
def _comment_listing(store: Store, docket: str = "") -> list[sqlite3.Row]:
"""One query: key, comment url, date, title, updated_at, year — newest
posted first so incremental index runs surface the latest comments
before older backlog.
*docket* filters via the URL (``.../comment/<docket>-%``); an empty
docket matches every comment and the docket is recovered from each
comment_id's ``<docket>-<seq>`` suffix.
"""
pattern = (
f"{_COMMENT_URL_PREFIX}{docket}-%" if docket else f"{_COMMENT_URL_PREFIX}%"
)
return (
store._con()
.execute(
"SELECT i.key, i.url, COALESCE(i.date_published,'') AS date, "
"COALESCE(i.title,'') AS title, COALESCE(i.updated_at,'') AS updated_at, "
"COALESCE((SELECT t.name FROM tags t JOIN item_tags it ON it.tag_id = t.id "
" WHERE it.item_id = i.id AND t.name LIKE 'year:%' LIMIT 1), '') AS year "
"FROM items i WHERE i.url LIKE ? "
"ORDER BY i.date_published DESC, i.key",
(pattern,),
)
.fetchall()
)
def comment_key_map(store: Store, docket: str) -> dict[str, tuple[str, str]]:
"""comment_id -> (bib item key, year) for one docket."""
return {
row["url"].rsplit("/", 1)[-1]: (row["key"], row["year"].split(":", 1)[-1])
for row in _comment_listing(store, docket)
}
def _comment_files(comment_dir: Path) -> tuple[tuple[str, str], ...]:
if not comment_dir.is_dir():
return ()
return tuple(
(p.name, str(p))
for p in sorted(comment_dir.iterdir())
if p.is_file() and p.suffix.lower() in _ATTACHMENT_EXT
)
def iter_comment_refs(
store: Store,
*,
docket: str = "",
root: Path | None = None,
skip_dockets: set[str] | frozenset[str] = frozenset(),
) -> Iterator[DocRef]:
"""One DocRef per comment, newest first. Fingerprint = file stats of
combined.md + attachments + the item's updated_at; no file is read
and the extraction machinery is not even imported."""
root = root if root is not None else _default_root()
for row in _comment_listing(store, docket):
comment_id = row["url"].rsplit("/", 1)[-1]
dk = docket or comment_id.rsplit("-", 1)[0]
if dk in skip_dockets:
continue
key, date, title, updated_at = (
row["key"],
row["date"],
row["title"],
row["updated_at"],
)
year = row["year"].split(":", 1)[-1]
comment_dir = root / dk / comment_id
combined = comment_dir / "combined.md"
files = _comment_files(comment_dir)
fp = (
fingerprint_files([combined, *(Path(p) for _, p in files)])
+ "|"
+ updated_at
)
meta = {
"docket": dk,
"comment_id": comment_id,
"doctype": "comment",
"kind": "comment",
"year": year,
"date": date[:10],
"title": title,
}
def _load(key=key, meta=meta, combined=combined, files=files) -> Doc | None:
from rex.comments.combine import parse_combined
if combined.exists():
_, body = parse_combined(combined.read_text())
if body.strip():
return Doc(key=key, text=body, metadata=meta, files=files)
item_row = (
store._con()
.execute("SELECT abstract FROM items WHERE key = ?", (key,))
.fetchone()
)
abstract = (item_row["abstract"] if item_row else "") or ""
return (
Doc(key=key, text=abstract, metadata=meta) if abstract.strip() else None
)
yield DocRef(
key=key, collection="comments", docket=dk, fingerprint=fp, load=_load
)
def iter_comment_docs(
store: Store, *, docket: str = "", root: Path | None = None
) -> Iterator[Doc]:
"""One Doc per comment, newest first: extraction body, else abstract."""
for ref in iter_comment_refs(store, docket=docket, root=root):
doc = ref.load()
if doc is not None:
yield doc
# ── rules ──
def _attachment_text(store: Store, item_key: str) -> str:
from rex.comments.combine import extract_attachment
parts = []
for path in _attachment_paths(store, item_key):
if path.exists():
result = extract_attachment(path)
if result.status == "ok":
parts.append(result.text)
return "\n\n".join(parts)
def _rule_text(store: Store, item_key: str) -> str:
"""TXT attachment (HTML wrapper stripped) preferred; else PDF-extract.
The FR ``.txt`` attachment is normalized once at capture time (see
``dev/scripts/fetch_fr_attachments.py`` and the one-off
``dev/scripts/migrate_fr_txt.py`` migration) via the shared
``rex.frtext.clean_fr_text`` cleaner. It's applied again here as a
cheap, idempotent guard — not a load-bearing transform — so
not-yet-migrated or freshly-downloaded files still come out clean.
"""
from rex.frtext import clean_fr_text
for path in _attachment_paths(store, item_key):
if path.suffix.lower() == ".txt" and path.exists():
return clean_fr_text(path.read_text(encoding="utf-8", errors="replace"))
return _attachment_text(store, item_key)
def rule_paragraphs(store: Store, item_key: str) -> tuple[Paragraph, ...]:
"""The rule's FR paragraph anchors in document order (empty when
``stack bib fr-grab`` has not run for it)."""
rows = (
store._con()
.execute(
"SELECT p_id, page, ordinal, text FROM fr_anchors "
"WHERE item_key = ? ORDER BY p_id",
(item_key,),
)
.fetchall()
)
return tuple(Paragraph(r[0], r[1], r[2], r[3]) for r in rows)
def _anchor_doc(store: Store, item_key: str) -> tuple[str, str]:
row = (
store._con()
.execute(
"SELECT html_url, fr_volume FROM fr_anchor_docs WHERE item_key = ?",
(item_key,),
)
.fetchone()
)
return (row[0], str(row[1])) if row else ("", "")
def iter_rule_refs(
store: Store, *, keys: tuple[str, ...] = (), tag: str = ""
) -> Iterator[DocRef]:
"""One DocRef per FR rule item. Fingerprint = the grabbed anchor
document's sha256 when present, else the attachment file stats +
updated_at."""
rows = (
store._con()
.execute(
"SELECT i.key, COALESCE(i.updated_at,'') AS updated_at, "
"(SELECT d.sha256 FROM fr_anchor_docs d WHERE d.item_key = i.key) AS anchor_sha "
"FROM items i WHERE i.item_type = 'rule'"
+ (
" AND i.id IN (SELECT item_id FROM item_tags WHERE tag_id IN "
"(SELECT id FROM tags WHERE name = ?))"
if tag
else ""
)
+ " ORDER BY i.id",
(tag,) if tag else (),
)
.fetchall()
)
# Attachment paths are only needed for rows without an anchor sha;
# one batched scan replaces one _attachment_paths query per such row.
paths_by_key = (
_attachment_paths_by_key(store)
if any(not r["anchor_sha"] for r in rows)
else {}
)
for row in rows:
key = row["key"]
if keys and key not in keys:
continue
if row["anchor_sha"]:
# updated_at too: the anchor sha covers the FR body, not the
# item's own metadata (title, cms-rule: tag, date).
fp = f"anchors:{row['anchor_sha']}|{row['updated_at']}"
else:
fp = fingerprint_files(paths_by_key.get(key, [])) + "|" + row["updated_at"]
def _load(key=key) -> Doc | None:
return _build_rule_doc(store, store.get(key))
yield DocRef(
key=key, collection="rules", docket=None, fingerprint=fp, load=_load
)
def _build_rule_doc(store: Store, item) -> Doc | None:
"""Anchor paragraphs when grabbed (exact ``#p-N`` provenance per
chunk), else TXT attachment, else PDF-extract; None when empty."""
paragraphs = rule_paragraphs(store, item.key)
html_url, volume = _anchor_doc(store, item.key)
if paragraphs:
text = "\n\n".join(p.text for p in paragraphs if p.text.strip())
else:
text = _rule_text(store, item.key)
if not text.strip():
return None
cms_rule = next(
(t.split(":", 1)[1] for t in item.tags if t.startswith("cms-rule:")), ""
)
return Doc(
key=item.key,
text=text,
metadata={
"doctype": "rule",
"kind": "rule",
"cms_rule_id": cms_rule,
"fr_document_number": item.document_number or "",
"year": _year_of(store, item.key),
"date": (item.date_published or "")[:10],
"title": item.title,
"item_key": item.key,
"html_url": html_url,
"fr_volume": volume,
},
paragraphs=paragraphs,
)
def iter_rule_docs(
store: Store, *, keys: tuple[str, ...] = (), tag: str = ""
) -> Iterator[Doc]:
"""One Doc per FR rule item: anchor paragraphs when grabbed (exact
``#p-N`` provenance per chunk), else TXT attachment, else PDF-extract."""
for ref in iter_rule_refs(store, keys=keys, tag=tag):
doc = ref.load()
if doc is not None:
yield doc
def _source_mtime_ns(src: Path) -> int:
"""Newest mtime across zotero.sqlite and its ``-wal`` sidecar.
Zotero commits into the WAL and only touches the main database at a
checkpoint, so ``src.stat()`` alone reports a database that has not
changed for weeks while the library is being edited daily.
"""
newest = src.stat().st_mtime_ns if src.exists() else 0
wal = src.with_name(src.name + "-wal")
if wal.exists():
newest = max(newest, wal.stat().st_mtime_ns)
return newest
def _copy_snapshot(src: Path, snap: Path) -> bool:
"""Copy the Zotero database to *snap*. False (logged) on failure —
the caller must not treat whatever is at *snap* as fresh."""
try:
shutil.copy2(src, snap)
except OSError as e:
log.warning("zotero snapshot failed (%s) — no Zotero PDF fallback", e)
return False
return True
class ZoteroPdfIndex:
"""bib/Zotero item key → storage PDFs, read from a *snapshot copy* of
zotero.sqlite (the live file is locked by the Zotero desktop and its
WAL must never be read in place)."""
def __init__(self, by_key: dict[str, list[Path]]) -> None:
self._by_key = by_key
self._loader: Callable[[], dict[str, list[Path]]] | None = None
@classmethod
def lazy(
cls, sqlite_path: Path, storage_dir: Path, tmp_dir: Path
) -> "ZoteroPdfIndex":
"""Snapshot on first ``pdfs_for``; skip the copy when the existing
snapshot is already as new as the source (the 1.95 GB copy is pure
waste on a run that never reaches the Zotero fallback)."""
inst = cls({})
def _load() -> dict[str, list[Path]]:
snap = Path(tmp_dir) / "zotero.sqlite"
src = Path(sqlite_path)
# Zotero writes to zotero.sqlite-wal and only touches the main
# file at a checkpoint, so the main mtime alone would leave a
# stale snapshot in place indefinitely — take the newer of the two.
newest = _source_mtime_ns(src)
if snap.exists() and src.exists() and snap.stat().st_mtime_ns >= newest:
return cls._read(snap, Path(storage_dir))
if not src.exists():
log.warning("zotero db %s missing — no Zotero PDF fallback", src)
return {}
Path(tmp_dir).mkdir(parents=True, exist_ok=True)
if not _copy_snapshot(src, snap):
# Only a *successful* copy may be stamped: stamping the
# stale snapshot left behind by a failed one would pass
# it off as current on every future run.
return {}
if newest:
# copy2 carries the *main* file's mtime, which is older
# than the WAL; stamp what we actually captured so the
# next run doesn't re-copy 1.95 GB for nothing.
os.utime(snap, ns=(newest, newest))
return cls._read(snap, Path(storage_dir))
inst._loader = _load
return inst
@classmethod
def snapshot(
cls, sqlite_path: Path, storage_dir: Path, tmp_dir: Path
) -> "ZoteroPdfIndex":
if not Path(sqlite_path).exists():
log.warning("zotero db %s missing — no Zotero PDF fallback", sqlite_path)
return cls({})
tmp_dir.mkdir(parents=True, exist_ok=True)
snap = tmp_dir / "zotero.sqlite"
if not _copy_snapshot(Path(sqlite_path), snap):
return cls({})
return cls(cls._read(snap, Path(storage_dir)))
@classmethod
def _read(cls, snap: Path, storage_dir: Path) -> dict[str, list[Path]]:
"""parent item key → existing storage PDFs, from a snapshot copy."""
try:
con = sqlite3.connect(f"file:{snap}?mode=ro", uri=True)
rows = con.execute(
"SELECT p.key, a.key, ia.path FROM itemAttachments ia "
"JOIN items a ON a.itemID = ia.itemID "
"JOIN items p ON p.itemID = ia.parentItemID "
"WHERE ia.path LIKE 'storage:%.pdf' "
"AND a.itemID NOT IN (SELECT itemID FROM deletedItems)"
).fetchall()
con.close()
except (OSError, sqlite3.Error) as e:
log.warning("zotero snapshot failed (%s) — no Zotero PDF fallback", e)
return {}
by_key: dict[str, list[Path]] = {}
for parent_key, att_key, path in rows:
pdf = Path(storage_dir) / att_key / path[len("storage:") :]
if pdf.exists():
by_key.setdefault(parent_key, []).append(pdf)
return by_key
def pdfs_for(self, key: str) -> list[Path]:
if self._loader is not None:
self._by_key = self._loader()
self._loader = None
return list(self._by_key.get(key, []))
def _attachment_sections(
store: Store, item_key: str
) -> tuple[list[str], list[tuple[str, str]]]:
"""(markdown sections, files) for an item's bib attachments."""
from rex.comments.combine import extract_attachment
sections: list[str] = []
files: list[tuple[str, str]] = []
for path in _attachment_paths(store, item_key):
if path.exists():
result = extract_attachment(path)
if result.status == "ok" and result.text.strip():
sections.append(f"## {path.name}\n\n{result.text.strip()}")
files.append((path.name, str(path)))
return sections, files
# ── corpus ──
def iter_corpus_refs(
store: Store,
*,
tag: str = "",
keys: tuple[str, ...] = (),
zotero: "ZoteroPdfIndex | None" = None,
) -> Iterator[DocRef]:
"""One DocRef per non-comment, non-skipped item. Fingerprint =
updated_at + the bib attachment file stats; the Zotero fallback is
only consulted by ``load()``.
Items tagged ``llm:skip`` are excluded — a generic opt-out for
material that should not be embedded (e.g. a scratch export, or a
duplicate scan of something already in the corpus). It's never
applied automatically; a caller opts an item in by tagging it.
*keys*, when given, limits the yield to those item keys (same
Python-side post-filter ``iter_rule_refs`` uses for its own
``keys`` — the corpus scan is cheap enough that this doesn't need
to be pushed into the SQL)."""
sql = (
"SELECT i.key, COALESCE(i.updated_at,'') AS updated_at FROM items i "
"WHERE i.id NOT IN (SELECT item_id FROM item_tags WHERE tag_id IN "
"(SELECT id FROM tags WHERE name IN ('doctype:comment', 'llm:skip')))"
+ (
" AND i.id IN (SELECT item_id FROM item_tags WHERE tag_id IN "
"(SELECT id FROM tags WHERE name = ?))"
if tag
else ""
)
+ " ORDER BY i.id"
)
rows = store._con().execute(sql, (tag,) if tag else ()).fetchall()
# One batched scan replaces one _attachment_paths query per item —
# the fingerprint step used to cost a query per corpus item before
# any change was even known.
paths_by_key = _attachment_paths_by_key(store)
for row in rows:
key = row["key"]
if keys and key not in keys:
continue
fp = row["updated_at"] + "|" + fingerprint_files(paths_by_key.get(key, []))
def _load(key=key) -> Doc | None:
return _build_corpus_doc(store, store.get(key), zotero)
yield DocRef(
key=key, collection="corpus", docket=None, fingerprint=fp, load=_load
)
def _build_corpus_doc(
store: Store, item, zotero: "ZoteroPdfIndex | None"
) -> Doc | None:
"""Attachment sections (bib, else Zotero-only storage PDFs) +
abstract; None when empty."""
from rex.comments.combine import extract_attachment
sections, files = _attachment_sections(store, item.key)
if not sections and zotero is not None:
for pdf in zotero.pdfs_for(item.key):
result = extract_attachment(pdf)
if result.status == "ok" and result.text.strip():
sections.append(f"## {pdf.name}\n\n{result.text.strip()}")
files.append((pdf.name, str(pdf)))
abstract = item.abstract.strip()
parts = sections + ([abstract] if abstract else [])
text = "\n\n".join(parts)
if not text.strip():
return None
project = next(
(t.split(":", 1)[1] for t in item.tags if t.startswith("project:")), ""
)
return Doc(
key=item.key,
text=text,
metadata={
"doctype": item.item_type,
"kind": "corpus",
"year": _year_of(store, item.key),
"date": (item.date_published or "")[:10],
"title": item.title,
"url": item.url or "",
"project": project,
},
files=tuple(files),
)
def iter_corpus_docs(
store: Store, *, tag: str = "", zotero: ZoteroPdfIndex | None = None
) -> Iterator[Doc]:
"""Every non-comment item: attachment sections (bib, else Zotero-only
storage PDFs) + abstract."""
for ref in iter_corpus_refs(store, tag=tag, zotero=zotero):
doc = ref.load()
if doc is not None:
yield doc