feat(comments): mtime-based extract currency; skipped dirs are never read; skip_dockets + reattach (refs #615)

This commit is contained in:
kert
2026-09-08 13:10:20 -04:00
parent 5ff3a5e5bf
commit e7287c827f
4 changed files with 192 additions and 35 deletions

View File

@@ -120,6 +120,34 @@ def derive_siblings_from_combined(comment_dir: Path) -> list[Path]:
return written return written
def source_paths(comment_dir: Path) -> list[Path]:
"""Files that feed extraction: attachments and bodies, never our own
outputs (``combined.md``, ``*.md`` siblings, ``*.tmp``) or dotfiles."""
out: list[Path] = []
for p in comment_dir.iterdir():
n = p.name
if not p.is_file() or n.startswith(".") or n == _FILENAME:
continue
if n.endswith(".md") or n.endswith(".tmp"):
continue
out.append(p)
return sorted(out)
def is_current(comment_dir: Path) -> bool:
"""True when ``combined.md`` exists and is at least as new as every
source file — a stat-only check, no reads."""
out = comment_dir / _FILENAME
try:
out_m = out.stat().st_mtime_ns
except FileNotFoundError:
return False
for p in source_paths(comment_dir):
if p.stat().st_mtime_ns > out_m:
return False
return True
def parse_combined(text: str) -> tuple[dict[str, Any], str]: def parse_combined(text: str) -> tuple[dict[str, Any], str]:
"""Split a combined.md into (frontmatter dict, body markdown).""" """Split a combined.md into (frontmatter dict, body markdown)."""
if not text.startswith("---\n"): if not text.startswith("---\n"):
@@ -136,11 +164,7 @@ def parse_combined(text: str) -> tuple[dict[str, Any], str]:
def _attachment_paths(comment_dir: Path) -> list[Path]: def _attachment_paths(comment_dir: Path) -> list[Path]:
return [ return source_paths(comment_dir)
p
for p in comment_dir.iterdir()
if p.is_file() and p.name != _FILENAME and not p.name.startswith(".")
]
def _render_section(filename: str, result: ExtractResult) -> str: def _render_section(filename: str, result: ExtractResult) -> str:

View File

@@ -12,12 +12,14 @@ from collections.abc import Callable
from concurrent.futures import ThreadPoolExecutor, as_completed from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path from pathlib import Path
from rex.comments.combine import derive_siblings_from_combined, extract_comment from rex.comments.combine import (
derive_siblings_from_combined,
extract_comment,
is_current,
)
log = logging.getLogger(__name__) log = logging.getLogger(__name__)
_COMBINED = "combined.md"
def walk_and_extract( def walk_and_extract(
root: Path, root: Path,
@@ -28,6 +30,8 @@ def walk_and_extract(
workers: int = 8, workers: int = 8,
force: bool = False, force: bool = False,
on_extracted: Callable[[str, Path], None] | None = None, on_extracted: Callable[[str, Path], None] | None = None,
skip_dockets: set[str] | frozenset[str] = frozenset(),
reattach: bool = False,
) -> dict[str, int]: ) -> dict[str, int]:
"""Process every comment dir under *root*. """Process every comment dir under *root*.
@@ -36,21 +40,22 @@ def walk_and_extract(
``comment_id`` and returns the inline body text from bib.sqlite (or ``comment_id`` and returns the inline body text from bib.sqlite (or
"" if missing). "" if missing).
*on_extracted*, if given, is called on the main thread with A dir is *current* when its combined.md is newer than every source
``(comment_id, comment_dir)`` for every dir that ends up with a file (``combine.is_current``); current dirs are skipped without being
combined.md — both newly-written ones and pre-existing skips. The read. *on_extracted* fires for newly written dirs; with *reattach*
callback is responsible for discovering whatever markdown files it it also fires for skipped dirs (after deriving any missing sibling
wants to consume in the dir (combined.md plus per-attachment MDs), which is the repair path for bib notes. *skip_dockets* names
siblings). Skipped dirs that pre-date the sibling-write feature get docket dirs to ignore entirely (sealed dockets), unless *force*.
siblings derived from combined.md before the callback fires, so the
callback always sees a complete set.
Returns counts: ``{"written": N, "skipped": N, "failed": N}``. Returns counts: ``{"written": N, "skipped": N, "failed": N}``, plus
``"skipped_sealed": N`` when any docket was skipped via *skip_dockets*.
""" """
candidate_dirs, skipped_dirs = _collect_dirs( candidate_dirs, skipped_dirs, sealed = _collect_dirs(
root, docket=docket, limit=limit, force=force root, docket=docket, limit=limit, force=force, skip_dockets=skip_dockets
) )
stats = {"written": 0, "skipped": len(skipped_dirs), "failed": 0} stats = {"written": 0, "skipped": len(skipped_dirs), "failed": 0}
if sealed:
stats["skipped_sealed"] = sealed
if candidate_dirs: if candidate_dirs:
with ThreadPoolExecutor(max_workers=workers) as pool: with ThreadPoolExecutor(max_workers=workers) as pool:
@@ -65,7 +70,7 @@ def walk_and_extract(
cdir = futures[fut] cdir = futures[fut]
on_extracted(cdir.name, cdir) on_extracted(cdir.name, cdir)
if on_extracted: if on_extracted and reattach:
for cdir in skipped_dirs: for cdir in skipped_dirs:
try: try:
derive_siblings_from_combined(cdir) derive_siblings_from_combined(cdir)
@@ -82,31 +87,39 @@ def _collect_dirs(
docket: str | None, docket: str | None,
limit: int | None, limit: int | None,
force: bool, force: bool,
) -> tuple[list[Path], list[Path]]: skip_dockets: set[str] | frozenset[str] = frozenset(),
"""Return ``(dirs_to_process, dirs_skipped)``. ) -> tuple[list[Path], list[Path], int]:
"""Return ``(dirs_to_process, dirs_skipped, sealed_docket_count)``.
``dirs_skipped`` are dirs that already have combined.md and would not ``dirs_skipped`` are dirs that are already current (``is_current``)
be re-processed. We return them as a list (not a count) so callers and would not be re-processed. We return them as a list (not a count)
can still drive per-dir post-processing — e.g. attaching the existing so callers can still drive per-dir post-processing — e.g. attaching
combined.md to bib on a re-run. the existing combined.md to bib on a re-run (``reattach``).
``sealed_docket_count`` counts docket dirs skipped entirely because
their name is in *skip_dockets* (and *force* is False).
""" """
todo: list[Path] = [] todo: list[Path] = []
skipped: list[Path] = [] skipped: list[Path] = []
sealed = 0
for docket_dir in sorted(root.iterdir()): for docket_dir in sorted(root.iterdir()):
if not docket_dir.is_dir(): if not docket_dir.is_dir():
continue continue
if docket and docket_dir.name != docket: if docket and docket_dir.name != docket:
continue continue
if docket_dir.name in skip_dockets and not force:
sealed += 1
continue
for cdir in sorted(docket_dir.iterdir()): for cdir in sorted(docket_dir.iterdir()):
if not cdir.is_dir(): if not cdir.is_dir():
continue continue
if (cdir / _COMBINED).is_file() and not force: if not force and is_current(cdir):
skipped.append(cdir) skipped.append(cdir)
continue continue
todo.append(cdir) todo.append(cdir)
if limit and len(todo) >= limit: if limit and len(todo) >= limit:
return todo, skipped return todo, skipped, sealed
return todo, skipped return todo, skipped, sealed
def _one( def _one(
@@ -115,7 +128,7 @@ def _one(
force: bool, force: bool,
) -> str: ) -> str:
"""Process one comment dir. Returns "written" / "skipped" / "failed".""" """Process one comment dir. Returns "written" / "skipped" / "failed"."""
if (comment_dir / _COMBINED).is_file() and not force: if not force and is_current(comment_dir):
return "skipped" return "skipped"
try: try:
body = inline_body_lookup(comment_dir.name) body = inline_body_lookup(comment_dir.name)
@@ -123,7 +136,7 @@ def _one(
log.warning("inline body lookup failed for %s: %s", comment_dir.name, e) log.warning("inline body lookup failed for %s: %s", comment_dir.name, e)
body = "" body = ""
try: try:
extract_comment(comment_dir, inline_body=body, force=force) extract_comment(comment_dir, inline_body=body, force=True)
except Exception as e: # noqa: BLE001 except Exception as e: # noqa: BLE001
log.warning("extract_comment failed for %s: %s", comment_dir, e) log.warning("extract_comment failed for %s: %s", comment_dir, e)
return "failed" return "failed"

View File

@@ -0,0 +1,117 @@
from __future__ import annotations
import os
from pathlib import Path
from rex.comments.combine import is_current, source_paths
from rex.comments.walker import walk_and_extract
def _dir(tmp_path: Path, docket="CMS-2024-0001", cid="CMS-2024-0001-0001") -> Path:
d = tmp_path / docket / cid
d.mkdir(parents=True)
return d
def test_source_paths_excludes_outputs(tmp_path: Path):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"x")
(d / "attachment_1.pdf.md").write_text("sibling")
(d / "combined.md").write_text("c")
(d / "combined.md.tmp").write_text("t")
(d / ".hidden").write_text("h")
assert [p.name for p in source_paths(d)] == ["attachment_1.pdf"]
def test_is_current_false_without_combined(tmp_path: Path):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"x")
assert is_current(d) is False
def test_is_current_true_when_combined_newer(tmp_path: Path):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"x")
os.utime(d / "attachment_1.pdf", ns=(1_000, 1_000))
(d / "combined.md").write_text("c")
assert is_current(d) is True
def test_is_current_false_when_source_newer(tmp_path: Path):
d = _dir(tmp_path)
(d / "combined.md").write_text("c")
os.utime(d / "combined.md", ns=(1_000, 1_000))
(d / "attachment_2.pdf").write_bytes(b"new")
assert is_current(d) is False
def test_walker_reextracts_stale_dir(tmp_path: Path, monkeypatch):
d = _dir(tmp_path)
(d / "combined.md").write_text("---\ncomment_id: x\n---\n\nold\n")
os.utime(d / "combined.md", ns=(1_000, 1_000))
(d / "attachment_1.pdf").write_bytes(b"%PDF") # newer than combined.md
called = []
monkeypatch.setattr(
"rex.comments.walker.extract_comment",
lambda cdir, **kw: called.append(cdir) or (cdir / "combined.md"),
)
stats = walk_and_extract(tmp_path, inline_body_lookup=lambda _c: "", workers=1)
assert stats["written"] == 1 and called == [d]
def test_walker_skips_current_dir_without_reading(tmp_path: Path, monkeypatch):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"%PDF")
os.utime(d / "attachment_1.pdf", ns=(1_000, 1_000))
(d / "combined.md").write_text("---\ncomment_id: x\n---\n\nbody\n")
reads = []
real_read_text = Path.read_text
def spy(self, *a, **kw):
if self.name == "combined.md":
reads.append(self)
return real_read_text(self, *a, **kw)
monkeypatch.setattr(Path, "read_text", spy)
seen = []
stats = walk_and_extract(
tmp_path,
inline_body_lookup=lambda _c: "",
workers=1,
on_extracted=lambda cid, _d: seen.append(cid),
)
assert stats == {"written": 0, "skipped": 1, "failed": 0}
assert reads == [] # skipped dirs are not read
assert seen == [] # and the callback does not fire without --reattach
def test_walker_reattach_fires_callback_for_skipped(tmp_path: Path):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"%PDF")
os.utime(d / "attachment_1.pdf", ns=(1_000, 1_000))
(d / "combined.md").write_text(
"---\ncomment_id: x\ndocket_id: y\n---\n\n## attachment_1.pdf\n\nbody\n"
)
seen = []
walk_and_extract(
tmp_path,
inline_body_lookup=lambda _c: "",
workers=1,
reattach=True,
on_extracted=lambda cid, _d: seen.append(cid),
)
assert seen == [d.name]
assert (d / "attachment_1.pdf.md").is_file() # siblings derived on reattach
def test_walker_skip_dockets(tmp_path: Path):
d = _dir(tmp_path)
(d / "attachment_1.pdf").write_bytes(b"%PDF")
stats = walk_and_extract(
tmp_path,
inline_body_lookup=lambda _c: "",
workers=1,
skip_dockets={"CMS-2024-0001"},
)
assert stats == {"written": 0, "skipped": 0, "failed": 0, "skipped_sealed": 1}
assert not (d / "combined.md").exists()

View File

@@ -2,6 +2,7 @@
from __future__ import annotations from __future__ import annotations
import os
from pathlib import Path from pathlib import Path
import fitz import fitz
@@ -44,6 +45,7 @@ def test_walk_and_extract_writes_all_combined(tmp_path: Path):
def test_walk_and_extract_skips_existing(tmp_path: Path): def test_walk_and_extract_skips_existing(tmp_path: Path):
a, _b = _setup_two_comments(tmp_path) a, _b = _setup_two_comments(tmp_path)
(a / "combined.md").write_text("---\ncomment_id: x\n---\n\npre-existing\n") (a / "combined.md").write_text("---\ncomment_id: x\n---\n\npre-existing\n")
os.utime(a / "attachment_1.pdf", ns=(1_000, 1_000))
stats = walk_and_extract(tmp_path, inline_body_lookup=_stub_inline_body, workers=1) stats = walk_and_extract(tmp_path, inline_body_lookup=_stub_inline_body, workers=1)
@@ -96,10 +98,10 @@ def test_on_extracted_called_for_written_dirs(tmp_path: Path):
def test_on_extracted_called_for_skipped_dirs(tmp_path: Path): def test_on_extracted_called_for_skipped_dirs(tmp_path: Path):
"""Pre-existing combined.md still gets the callback, and any missing """With reattach=True, pre-existing combined.md still gets the callback,
sibling MDs are derived from it before the callback fires — so and any missing sibling MDs are derived from it before the callback
previously-extracted dirs end up with the same set of MDs as fires — so previously-extracted dirs end up with the same set of MDs
freshly-written ones.""" as freshly-written ones."""
a, b = _setup_two_comments(tmp_path) a, b = _setup_two_comments(tmp_path)
# Pre-write a stale combined.md with one section so derive_siblings # Pre-write a stale combined.md with one section so derive_siblings
# has something to split. # has something to split.
@@ -112,6 +114,7 @@ def test_on_extracted_called_for_skipped_dirs(tmp_path: Path):
tmp_path, tmp_path,
inline_body_lookup=_stub_inline_body, inline_body_lookup=_stub_inline_body,
workers=1, workers=1,
reattach=True,
on_extracted=lambda cid, _d: seen.append(cid), on_extracted=lambda cid, _d: seen.append(cid),
) )