feat(llm): lazy DocRefs with fingerprints, single-query listings, lazy Zotero snapshot (refs #615)
This commit is contained in:
@@ -1,5 +1,12 @@
|
|||||||
"""Doc iterators over the bib store + comment extraction tree.
|
"""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
|
Comments: text prefers the #253 extraction output
|
||||||
(``.state/comments/<docket>/<comment_id>/combined.md`` body via
|
(``.state/comments/<docket>/<comment_id>/combined.md`` body via
|
||||||
``rex.comments.combine.parse_combined``); falls back to the bib
|
``rex.comments.combine.parse_combined``); falls back to the bib
|
||||||
@@ -18,9 +25,11 @@ from __future__ import annotations
|
|||||||
import logging
|
import logging
|
||||||
import shutil
|
import shutil
|
||||||
import sqlite3
|
import sqlite3
|
||||||
|
from dataclasses import dataclass
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Iterator
|
from typing import Callable, Iterator
|
||||||
|
|
||||||
|
from bib.dockets import fingerprint_files
|
||||||
from bib.store import Store
|
from bib.store import Store
|
||||||
from llm.chunk import Doc, Paragraph
|
from llm.chunk import Doc, Paragraph
|
||||||
|
|
||||||
@@ -30,6 +39,19 @@ _COMMENT_URL_PREFIX = "https://www.regulations.gov/comment/"
|
|||||||
_ATTACHMENT_EXT = (".pdf", ".docx", ".doc", ".txt")
|
_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:
|
def _default_root() -> Path:
|
||||||
from conf import ROOT
|
from conf import ROOT
|
||||||
|
|
||||||
@@ -51,12 +73,26 @@ def _year_of(store: Store, item_key: str) -> str:
|
|||||||
return row[0].split(":", 1)[1] if row else ""
|
return row[0].split(":", 1)[1] if row else ""
|
||||||
|
|
||||||
|
|
||||||
def _comment_rows(
|
def _attachment_paths(store: Store, item_key: str) -> list[Path]:
|
||||||
store: Store, docket: str = ""
|
rows = (
|
||||||
) -> list[tuple[str, str, str, str, str]]:
|
store._con()
|
||||||
"""(docket, comment_id, item key, date_published, title) for every
|
.execute(
|
||||||
comment matching *docket*, newest posted first so incremental index
|
"SELECT a.storage_path FROM attachments a "
|
||||||
runs surface the latest comments before older backlog.
|
"JOIN items i ON i.id = a.item_id WHERE i.key = ?",
|
||||||
|
(item_key,),
|
||||||
|
)
|
||||||
|
.fetchall()
|
||||||
|
)
|
||||||
|
return [Path(r[0]) for r in rows]
|
||||||
|
|
||||||
|
|
||||||
|
# ── 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* filters via the URL (``.../comment/<docket>-%``); an empty
|
||||||
docket matches every comment and the docket is recovered from each
|
docket matches every comment and the docket is recovered from each
|
||||||
@@ -65,29 +101,26 @@ def _comment_rows(
|
|||||||
pattern = (
|
pattern = (
|
||||||
f"{_COMMENT_URL_PREFIX}{docket}-%" if docket else f"{_COMMENT_URL_PREFIX}%"
|
f"{_COMMENT_URL_PREFIX}{docket}-%" if docket else f"{_COMMENT_URL_PREFIX}%"
|
||||||
)
|
)
|
||||||
rows = (
|
return (
|
||||||
store._con()
|
store._con()
|
||||||
.execute(
|
.execute(
|
||||||
"SELECT i.key, i.url, COALESCE(i.date_published, ''), "
|
"SELECT i.key, i.url, COALESCE(i.date_published,'') AS date, "
|
||||||
"COALESCE(i.title, '') FROM items i WHERE i.url LIKE ? "
|
"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",
|
"ORDER BY i.date_published DESC, i.key",
|
||||||
(pattern,),
|
(pattern,),
|
||||||
)
|
)
|
||||||
.fetchall()
|
.fetchall()
|
||||||
)
|
)
|
||||||
out = []
|
|
||||||
for key, url, date, title in rows:
|
|
||||||
comment_id = url.rsplit("/", 1)[-1]
|
|
||||||
dk = docket or comment_id.rsplit("-", 1)[0]
|
|
||||||
out.append((dk, comment_id, key, date, title))
|
|
||||||
return out
|
|
||||||
|
|
||||||
|
|
||||||
def comment_key_map(store: Store, docket: str) -> dict[str, tuple[str, str]]:
|
def comment_key_map(store: Store, docket: str) -> dict[str, tuple[str, str]]:
|
||||||
"""comment_id -> (bib item key, year) for one docket."""
|
"""comment_id -> (bib item key, year) for one docket."""
|
||||||
return {
|
return {
|
||||||
comment_id: (key, _year_of(store, key))
|
row["url"].rsplit("/", 1)[-1]: (row["key"], row["year"].split(":", 1)[-1])
|
||||||
for _, comment_id, key, _, _ in _comment_rows(store, docket)
|
for row in _comment_listing(store, docket)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -101,52 +134,84 @@ def _comment_files(comment_dir: Path) -> tuple[tuple[str, str], ...]:
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def iter_comment_docs(
|
def iter_comment_refs(
|
||||||
store: Store, *, docket: str = "", root: Path | None = None
|
store: Store,
|
||||||
) -> Iterator[Doc]:
|
*,
|
||||||
"""One Doc per comment, newest first: extraction body, else abstract."""
|
docket: str = "",
|
||||||
from rex.comments.combine import parse_combined
|
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()
|
root = root if root is not None else _default_root()
|
||||||
for dk, comment_id, key, date, title in _comment_rows(store, docket):
|
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 = {
|
meta = {
|
||||||
"docket": dk,
|
"docket": dk,
|
||||||
"comment_id": comment_id,
|
"comment_id": comment_id,
|
||||||
"doctype": "comment",
|
"doctype": "comment",
|
||||||
"kind": "comment",
|
"kind": "comment",
|
||||||
"year": _year_of(store, key),
|
"year": year,
|
||||||
"date": date[:10],
|
"date": date[:10],
|
||||||
"title": title,
|
"title": title,
|
||||||
}
|
}
|
||||||
comment_dir = root / dk / comment_id
|
|
||||||
combined = comment_dir / "combined.md"
|
def _load(key=key, meta=meta, combined=combined, files=files) -> Doc | None:
|
||||||
|
from rex.comments.combine import parse_combined
|
||||||
|
|
||||||
if combined.exists():
|
if combined.exists():
|
||||||
_, body = parse_combined(combined.read_text())
|
_, body = parse_combined(combined.read_text())
|
||||||
if body.strip():
|
if body.strip():
|
||||||
yield Doc(
|
return Doc(key=key, text=body, metadata=meta, files=files)
|
||||||
key=key, text=body, metadata=meta, files=_comment_files(comment_dir)
|
row = (
|
||||||
|
store._con()
|
||||||
|
.execute("SELECT abstract FROM items WHERE key = ?", (key,))
|
||||||
|
.fetchone()
|
||||||
)
|
)
|
||||||
continue
|
abstract = (row["abstract"] if row else "") or ""
|
||||||
item = store.get(key)
|
return (
|
||||||
if item.abstract.strip():
|
Doc(key=key, text=abstract, metadata=meta) if abstract.strip() else None
|
||||||
yield Doc(key=key, text=item.abstract, metadata=meta)
|
)
|
||||||
|
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
def _attachment_text(store: Store, item_key: str) -> str:
|
def _attachment_text(store: Store, item_key: str) -> str:
|
||||||
from rex.comments.combine import extract_attachment
|
from rex.comments.combine import extract_attachment
|
||||||
|
|
||||||
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()
|
|
||||||
)
|
|
||||||
parts = []
|
parts = []
|
||||||
for (storage_path,) in rows:
|
for path in _attachment_paths(store, item_key):
|
||||||
path = Path(storage_path)
|
|
||||||
if path.exists():
|
if path.exists():
|
||||||
result = extract_attachment(path)
|
result = extract_attachment(path)
|
||||||
if result.status == "ok":
|
if result.status == "ok":
|
||||||
@@ -166,17 +231,7 @@ def _rule_text(store: Store, item_key: str) -> str:
|
|||||||
"""
|
"""
|
||||||
from rex.frtext import clean_fr_text
|
from rex.frtext import clean_fr_text
|
||||||
|
|
||||||
rows = (
|
for path in _attachment_paths(store, item_key):
|
||||||
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()
|
|
||||||
)
|
|
||||||
for (storage_path,) in rows:
|
|
||||||
path = Path(storage_path)
|
|
||||||
if path.suffix.lower() == ".txt" and path.exists():
|
if path.suffix.lower() == ".txt" and path.exists():
|
||||||
return clean_fr_text(path.read_text(encoding="utf-8", errors="replace"))
|
return clean_fr_text(path.read_text(encoding="utf-8", errors="replace"))
|
||||||
return _attachment_text(store, item_key)
|
return _attachment_text(store, item_key)
|
||||||
@@ -209,14 +264,56 @@ def _anchor_doc(store: Store, item_key: str) -> tuple[str, str]:
|
|||||||
return (row[0], str(row[1])) if row else ("", "")
|
return (row[0], str(row[1])) if row else ("", "")
|
||||||
|
|
||||||
|
|
||||||
def iter_rule_docs(
|
# ── rules ──
|
||||||
|
|
||||||
|
|
||||||
|
def iter_rule_refs(
|
||||||
store: Store, *, keys: tuple[str, ...] = (), tag: str = ""
|
store: Store, *, keys: tuple[str, ...] = (), tag: str = ""
|
||||||
) -> Iterator[Doc]:
|
) -> Iterator[DocRef]:
|
||||||
"""One Doc per FR rule item: anchor paragraphs when grabbed (exact
|
"""One DocRef per FR rule item. Fingerprint = the grabbed anchor
|
||||||
``#p-N`` provenance per chunk), else TXT attachment, else PDF-extract."""
|
document's sha256 when present, else the attachment file stats +
|
||||||
for item in store.list_items(item_type="rule", tag=tag):
|
updated_at."""
|
||||||
if keys and item.key not in keys:
|
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()
|
||||||
|
)
|
||||||
|
for row in rows:
|
||||||
|
key = row["key"]
|
||||||
|
if keys and key not in keys:
|
||||||
continue
|
continue
|
||||||
|
if row["anchor_sha"]:
|
||||||
|
fp = f"anchors:{row['anchor_sha']}"
|
||||||
|
else:
|
||||||
|
fp = (
|
||||||
|
fingerprint_files(_attachment_paths(store, 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)
|
paragraphs = rule_paragraphs(store, item.key)
|
||||||
html_url, volume = _anchor_doc(store, item.key)
|
html_url, volume = _anchor_doc(store, item.key)
|
||||||
if paragraphs:
|
if paragraphs:
|
||||||
@@ -224,11 +321,11 @@ def iter_rule_docs(
|
|||||||
else:
|
else:
|
||||||
text = _rule_text(store, item.key)
|
text = _rule_text(store, item.key)
|
||||||
if not text.strip():
|
if not text.strip():
|
||||||
continue
|
return None
|
||||||
cms_rule = next(
|
cms_rule = next(
|
||||||
(t.split(":", 1)[1] for t in item.tags if t.startswith("cms-rule:")), ""
|
(t.split(":", 1)[1] for t in item.tags if t.startswith("cms-rule:")), ""
|
||||||
)
|
)
|
||||||
yield Doc(
|
return Doc(
|
||||||
key=item.key,
|
key=item.key,
|
||||||
text=text,
|
text=text,
|
||||||
metadata={
|
metadata={
|
||||||
@@ -247,6 +344,17 @@ def iter_rule_docs(
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
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
|
||||||
|
|
||||||
|
|
||||||
class ZoteroPdfIndex:
|
class ZoteroPdfIndex:
|
||||||
"""bib/Zotero item key → storage PDFs, read from a *snapshot copy* of
|
"""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
|
zotero.sqlite (the live file is locked by the Zotero desktop and its
|
||||||
@@ -254,6 +362,30 @@ class ZoteroPdfIndex:
|
|||||||
|
|
||||||
def __init__(self, by_key: dict[str, list[Path]]) -> None:
|
def __init__(self, by_key: dict[str, list[Path]]) -> None:
|
||||||
self._by_key = by_key
|
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)
|
||||||
|
if (
|
||||||
|
snap.exists()
|
||||||
|
and src.exists()
|
||||||
|
and snap.stat().st_mtime_ns >= src.stat().st_mtime_ns
|
||||||
|
):
|
||||||
|
return cls._read(snap, Path(storage_dir))
|
||||||
|
return cls.snapshot(src, storage_dir, Path(tmp_dir))._by_key
|
||||||
|
|
||||||
|
inst._loader = _load
|
||||||
|
return inst
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def snapshot(
|
def snapshot(
|
||||||
@@ -265,7 +397,18 @@ class ZoteroPdfIndex:
|
|||||||
tmp_dir.mkdir(parents=True, exist_ok=True)
|
tmp_dir.mkdir(parents=True, exist_ok=True)
|
||||||
snap = tmp_dir / "zotero.sqlite"
|
snap = tmp_dir / "zotero.sqlite"
|
||||||
try:
|
try:
|
||||||
|
# copy2 preserves the source mtime, so a later lazy() sees
|
||||||
|
# snap.mtime >= src.mtime until Zotero writes again.
|
||||||
shutil.copy2(sqlite_path, snap)
|
shutil.copy2(sqlite_path, snap)
|
||||||
|
except OSError as e:
|
||||||
|
log.warning("zotero snapshot failed (%s) — no Zotero PDF fallback", e)
|
||||||
|
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)
|
con = sqlite3.connect(f"file:{snap}?mode=ro", uri=True)
|
||||||
rows = con.execute(
|
rows = con.execute(
|
||||||
"SELECT p.key, a.key, ia.path FROM itemAttachments ia "
|
"SELECT p.key, a.key, ia.path FROM itemAttachments ia "
|
||||||
@@ -277,15 +420,18 @@ class ZoteroPdfIndex:
|
|||||||
con.close()
|
con.close()
|
||||||
except (OSError, sqlite3.Error) as e:
|
except (OSError, sqlite3.Error) as e:
|
||||||
log.warning("zotero snapshot failed (%s) — no Zotero PDF fallback", e)
|
log.warning("zotero snapshot failed (%s) — no Zotero PDF fallback", e)
|
||||||
return cls({})
|
return {}
|
||||||
by_key: dict[str, list[Path]] = {}
|
by_key: dict[str, list[Path]] = {}
|
||||||
for parent_key, att_key, path in rows:
|
for parent_key, att_key, path in rows:
|
||||||
pdf = Path(storage_dir) / att_key / path[len("storage:") :]
|
pdf = Path(storage_dir) / att_key / path[len("storage:") :]
|
||||||
if pdf.exists():
|
if pdf.exists():
|
||||||
by_key.setdefault(parent_key, []).append(pdf)
|
by_key.setdefault(parent_key, []).append(pdf)
|
||||||
return cls(by_key)
|
return by_key
|
||||||
|
|
||||||
def pdfs_for(self, key: str) -> list[Path]:
|
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, []))
|
return list(self._by_key.get(key, []))
|
||||||
|
|
||||||
|
|
||||||
@@ -295,19 +441,9 @@ def _attachment_sections(
|
|||||||
"""(markdown sections, files) for an item's bib attachments."""
|
"""(markdown sections, files) for an item's bib attachments."""
|
||||||
from rex.comments.combine import extract_attachment
|
from rex.comments.combine import extract_attachment
|
||||||
|
|
||||||
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()
|
|
||||||
)
|
|
||||||
sections: list[str] = []
|
sections: list[str] = []
|
||||||
files: list[tuple[str, str]] = []
|
files: list[tuple[str, str]] = []
|
||||||
for (storage_path,) in rows:
|
for path in _attachment_paths(store, item_key):
|
||||||
path = Path(storage_path)
|
|
||||||
if path.exists():
|
if path.exists():
|
||||||
result = extract_attachment(path)
|
result = extract_attachment(path)
|
||||||
if result.status == "ok" and result.text.strip():
|
if result.status == "ok" and result.text.strip():
|
||||||
@@ -316,16 +452,46 @@ def _attachment_sections(
|
|||||||
return sections, files
|
return sections, files
|
||||||
|
|
||||||
|
|
||||||
def iter_corpus_docs(
|
# ── corpus ──
|
||||||
store: Store, *, tag: str = "", zotero: ZoteroPdfIndex | None = None
|
|
||||||
) -> Iterator[Doc]:
|
|
||||||
"""Every non-comment item: attachment sections (bib, else Zotero-only
|
def iter_corpus_refs(
|
||||||
storage PDFs) + abstract."""
|
store: Store, *, tag: str = "", zotero: "ZoteroPdfIndex | None" = None
|
||||||
|
) -> Iterator[DocRef]:
|
||||||
|
"""One DocRef per non-comment item. Fingerprint = updated_at + the
|
||||||
|
bib attachment file stats; the Zotero fallback is only consulted by
|
||||||
|
``load()``."""
|
||||||
|
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 = 'doctype:comment'))"
|
||||||
|
+ (
|
||||||
|
" 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"
|
||||||
|
)
|
||||||
|
for row in store._con().execute(sql, (tag,) if tag else ()).fetchall():
|
||||||
|
key = row["key"]
|
||||||
|
fp = row["updated_at"] + "|" + fingerprint_files(_attachment_paths(store, 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
|
from rex.comments.combine import extract_attachment
|
||||||
|
|
||||||
for item in store.list_items(tag=tag):
|
|
||||||
if "doctype:comment" in item.tags:
|
|
||||||
continue
|
|
||||||
sections, files = _attachment_sections(store, item.key)
|
sections, files = _attachment_sections(store, item.key)
|
||||||
if not sections and zotero is not None:
|
if not sections and zotero is not None:
|
||||||
for pdf in zotero.pdfs_for(item.key):
|
for pdf in zotero.pdfs_for(item.key):
|
||||||
@@ -337,11 +503,11 @@ def iter_corpus_docs(
|
|||||||
parts = sections + ([abstract] if abstract else [])
|
parts = sections + ([abstract] if abstract else [])
|
||||||
text = "\n\n".join(parts)
|
text = "\n\n".join(parts)
|
||||||
if not text.strip():
|
if not text.strip():
|
||||||
continue
|
return None
|
||||||
project = next(
|
project = next(
|
||||||
(t.split(":", 1)[1] for t in item.tags if t.startswith("project:")), ""
|
(t.split(":", 1)[1] for t in item.tags if t.startswith("project:")), ""
|
||||||
)
|
)
|
||||||
yield Doc(
|
return Doc(
|
||||||
key=item.key,
|
key=item.key,
|
||||||
text=text,
|
text=text,
|
||||||
metadata={
|
metadata={
|
||||||
@@ -355,3 +521,14 @@ def iter_corpus_docs(
|
|||||||
},
|
},
|
||||||
files=tuple(files),
|
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
|
||||||
|
|||||||
192
tests/llm/test_source_refs.py
Normal file
192
tests/llm/test_source_refs.py
Normal file
@@ -0,0 +1,192 @@
|
|||||||
|
"""llm.source — lazy DocRefs: fingerprints without loads, docket skips,
|
||||||
|
lazy Zotero snapshot."""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from bib.item import Item, Rule
|
||||||
|
from bib.store import Store
|
||||||
|
from llm.source import (
|
||||||
|
DocRef,
|
||||||
|
ZoteroPdfIndex,
|
||||||
|
iter_comment_refs,
|
||||||
|
iter_corpus_refs,
|
||||||
|
iter_rule_refs,
|
||||||
|
)
|
||||||
|
|
||||||
|
DOCKET = "CMS-2019-0111"
|
||||||
|
CID = f"{DOCKET}-0042"
|
||||||
|
COMBINED = f"---\ncomment_id: {CID}\ndocket_id: {DOCKET}\n---\n\nWe object.\n"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def store(tmp_path):
|
||||||
|
s = Store(":memory:", storage_dir=tmp_path / "storage")
|
||||||
|
key = s.create(
|
||||||
|
Item(
|
||||||
|
item_type="report",
|
||||||
|
title="A comment",
|
||||||
|
url=f"https://www.regulations.gov/comment/{CID}",
|
||||||
|
abstract="Inline.",
|
||||||
|
date_published="2019-09-27",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
for tag in ("doctype:comment", "year:2019", f"reg-docket:{DOCKET}"):
|
||||||
|
s.add_tag(key, tag)
|
||||||
|
s._comment_key = key
|
||||||
|
return s
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.fixture
|
||||||
|
def root(tmp_path):
|
||||||
|
d = tmp_path / DOCKET / CID
|
||||||
|
d.mkdir(parents=True)
|
||||||
|
(d / "combined.md").write_text(COMBINED)
|
||||||
|
return tmp_path
|
||||||
|
|
||||||
|
|
||||||
|
class TestCommentRefs:
|
||||||
|
def test_ref_has_fingerprint_and_lazy_load(self, store, root, monkeypatch):
|
||||||
|
reads = []
|
||||||
|
real = Path.read_text
|
||||||
|
monkeypatch.setattr(
|
||||||
|
Path,
|
||||||
|
"read_text",
|
||||||
|
lambda self, *a, **k: reads.append(self) or real(self, *a, **k),
|
||||||
|
)
|
||||||
|
refs = list(iter_comment_refs(store, docket=DOCKET, root=root))
|
||||||
|
assert len(refs) == 1
|
||||||
|
r = refs[0]
|
||||||
|
assert isinstance(r, DocRef)
|
||||||
|
assert (r.key, r.collection, r.docket) == (
|
||||||
|
store._comment_key,
|
||||||
|
"comments",
|
||||||
|
DOCKET,
|
||||||
|
)
|
||||||
|
assert len(r.fingerprint) > 64 and "|" in r.fingerprint
|
||||||
|
assert reads == [] # nothing read yet
|
||||||
|
doc = r.load()
|
||||||
|
assert "We object" in doc.text
|
||||||
|
assert doc.metadata["year"] == "2019" and doc.metadata["docket"] == DOCKET
|
||||||
|
assert reads # load() read combined.md
|
||||||
|
|
||||||
|
def test_fingerprint_changes_with_new_attachment(self, store, root):
|
||||||
|
f1 = next(iter_comment_refs(store, docket=DOCKET, root=root)).fingerprint
|
||||||
|
(root / DOCKET / CID / "attachment_1.pdf").write_bytes(b"%PDF")
|
||||||
|
f2 = next(iter_comment_refs(store, docket=DOCKET, root=root)).fingerprint
|
||||||
|
assert f1 != f2
|
||||||
|
|
||||||
|
def test_fingerprint_changes_with_updated_at(self, store, root):
|
||||||
|
f1 = next(iter_comment_refs(store, docket=DOCKET, root=root)).fingerprint
|
||||||
|
store._con().execute(
|
||||||
|
"UPDATE items SET updated_at='2030-01-01T00:00:00Z' WHERE key=?",
|
||||||
|
(store._comment_key,),
|
||||||
|
)
|
||||||
|
f2 = next(iter_comment_refs(store, docket=DOCKET, root=root)).fingerprint
|
||||||
|
assert f1 != f2
|
||||||
|
|
||||||
|
def test_skip_dockets_excludes_rows(self, store, root):
|
||||||
|
assert list(iter_comment_refs(store, root=root, skip_dockets={DOCKET})) == []
|
||||||
|
|
||||||
|
def test_unextracted_load_falls_back_to_abstract(self, store, tmp_path):
|
||||||
|
r = next(iter_comment_refs(store, docket=DOCKET, root=tmp_path))
|
||||||
|
assert r.load().text == "Inline."
|
||||||
|
|
||||||
|
def test_empty_comment_load_returns_none(self, store, tmp_path):
|
||||||
|
store._con().execute(
|
||||||
|
"UPDATE items SET abstract='' WHERE key=?", (store._comment_key,)
|
||||||
|
)
|
||||||
|
r = next(iter_comment_refs(store, docket=DOCKET, root=tmp_path))
|
||||||
|
assert r.load() is None
|
||||||
|
|
||||||
|
def test_single_query_for_year(self, store, root, monkeypatch):
|
||||||
|
"""No per-comment year lookups: exactly one SELECT on items for the listing."""
|
||||||
|
calls: list[str] = []
|
||||||
|
store._con().set_trace_callback(calls.append)
|
||||||
|
list(iter_comment_refs(store, docket=DOCKET, root=root))
|
||||||
|
store._con().set_trace_callback(None)
|
||||||
|
selects = [c for c in calls if c.lstrip().upper().startswith("SELECT")]
|
||||||
|
assert len(selects) == 1
|
||||||
|
|
||||||
|
|
||||||
|
class TestRuleRefs:
|
||||||
|
def test_fingerprint_from_anchor_sha(self, store):
|
||||||
|
key = store.create(
|
||||||
|
Rule(title="R", url="https://fr/1", document_number="2019-1")
|
||||||
|
)
|
||||||
|
store._con().execute(
|
||||||
|
"INSERT INTO fr_anchor_docs (item_key, document_number, html_url, start_page, end_page, fr_volume, sha256) VALUES (?,?,?,?,?,?,?)",
|
||||||
|
(key, "2019-1", "https://fr/1", 1, 2, 84, "abc"),
|
||||||
|
)
|
||||||
|
store._con().execute(
|
||||||
|
"INSERT INTO fr_anchors (item_key, p_id, page, ordinal, text) VALUES (?,?,?,?,?)",
|
||||||
|
(key, 1, 1, 1, "Para one."),
|
||||||
|
)
|
||||||
|
refs = list(iter_rule_refs(store))
|
||||||
|
assert [r.key for r in refs] == [key]
|
||||||
|
assert refs[0].fingerprint == "anchors:abc"
|
||||||
|
assert refs[0].collection == "rules" and refs[0].docket is None
|
||||||
|
assert refs[0].load().text == "Para one."
|
||||||
|
|
||||||
|
|
||||||
|
class TestCorpusRefs:
|
||||||
|
def test_excludes_comments_and_is_lazy(self, store):
|
||||||
|
key = store.create(
|
||||||
|
Item(
|
||||||
|
item_type="report", title="Report", url="https://x/r", abstract="Body."
|
||||||
|
)
|
||||||
|
)
|
||||||
|
refs = list(iter_corpus_refs(store))
|
||||||
|
assert [r.key for r in refs] == [key]
|
||||||
|
assert refs[0].collection == "corpus"
|
||||||
|
assert refs[0].load().metadata["kind"] == "corpus"
|
||||||
|
|
||||||
|
def test_zotero_not_consulted_until_load(self, store, tmp_path):
|
||||||
|
store.create(
|
||||||
|
Item(
|
||||||
|
item_type="report", title="Report", url="https://x/r", abstract="Body."
|
||||||
|
)
|
||||||
|
)
|
||||||
|
calls = []
|
||||||
|
|
||||||
|
class Z(ZoteroPdfIndex):
|
||||||
|
def pdfs_for(self, key):
|
||||||
|
calls.append(key)
|
||||||
|
return []
|
||||||
|
|
||||||
|
refs = list(iter_corpus_refs(store, zotero=Z({})))
|
||||||
|
assert calls == []
|
||||||
|
refs[0].load()
|
||||||
|
assert calls
|
||||||
|
|
||||||
|
|
||||||
|
class TestLazyZotero:
|
||||||
|
def test_no_copy_until_first_lookup(self, tmp_path):
|
||||||
|
src = tmp_path / "zotero.sqlite"
|
||||||
|
import sqlite3
|
||||||
|
|
||||||
|
con = sqlite3.connect(src)
|
||||||
|
con.executescript(
|
||||||
|
"CREATE TABLE items(itemID INTEGER, key TEXT); CREATE TABLE itemAttachments(itemID INTEGER, parentItemID INTEGER, path TEXT); CREATE TABLE deletedItems(itemID INTEGER);"
|
||||||
|
)
|
||||||
|
con.close()
|
||||||
|
snap_dir = tmp_path / "snap"
|
||||||
|
z = ZoteroPdfIndex.lazy(src, tmp_path / "storage", snap_dir)
|
||||||
|
assert not (snap_dir / "zotero.sqlite").exists()
|
||||||
|
assert z.pdfs_for("ABCD1234") == []
|
||||||
|
assert (snap_dir / "zotero.sqlite").exists()
|
||||||
|
m1 = (snap_dir / "zotero.sqlite").stat().st_mtime_ns
|
||||||
|
z2 = ZoteroPdfIndex.lazy(src, tmp_path / "storage", snap_dir)
|
||||||
|
z2.pdfs_for("ABCD1234")
|
||||||
|
assert (
|
||||||
|
snap_dir / "zotero.sqlite"
|
||||||
|
).stat().st_mtime_ns == m1 # source unchanged → no recopy
|
||||||
|
bump = src.stat().st_mtime_ns + 10**9
|
||||||
|
os.utime(src, ns=(bump, bump)) # source now newer than the snapshot
|
||||||
|
z3 = ZoteroPdfIndex.lazy(src, tmp_path / "storage", snap_dir)
|
||||||
|
z3.pdfs_for("ABCD1234")
|
||||||
|
assert (snap_dir / "zotero.sqlite").stat().st_mtime_ns != m1
|
||||||
Reference in New Issue
Block a user