feat(bib): discover_docket + walk_docket — watermark, counts, auto-seal (refs #615)
This commit is contained in:
@@ -29,12 +29,14 @@ import os
|
|||||||
import re
|
import re
|
||||||
import time
|
import time
|
||||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||||
from dataclasses import dataclass, field
|
from dataclasses import dataclass, field, replace
|
||||||
|
from datetime import date
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import TYPE_CHECKING, Callable, Iterator
|
from typing import TYPE_CHECKING, Callable, Iterator
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
|
||||||
|
from bib.dockets import Docket
|
||||||
from bib.item import Source
|
from bib.item import Source
|
||||||
from bib.tag import Tag
|
from bib.tag import Tag
|
||||||
|
|
||||||
@@ -798,6 +800,19 @@ def upsert_comment(
|
|||||||
extra_tags: list[str] | None = None,
|
extra_tags: list[str] | None = None,
|
||||||
) -> str:
|
) -> str:
|
||||||
"""Upsert the comment as a Source item, return bib key."""
|
"""Upsert the comment as a Source item, return bib key."""
|
||||||
|
return upsert_comment_status(store, comment, cms_id=cms_id, extra_tags=extra_tags)[
|
||||||
|
0
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def upsert_comment_status(
|
||||||
|
store: Store,
|
||||||
|
comment: Comment,
|
||||||
|
*,
|
||||||
|
cms_id: str = "",
|
||||||
|
extra_tags: list[str] | None = None,
|
||||||
|
) -> tuple[str, str]:
|
||||||
|
"""Like :func:`upsert_comment`, also returning created/updated/unchanged."""
|
||||||
# Title = the comment's own CMS-DOCKET-XXXX-NNNN id, optionally
|
# Title = the comment's own CMS-DOCKET-XXXX-NNNN id, optionally
|
||||||
# prefixed with the byline (org or first/last name) when the
|
# prefixed with the byline (org or first/last name) when the
|
||||||
# commenter filled those fields in. The previous "Comment on
|
# commenter filled those fields in. The previous "Comment on
|
||||||
@@ -830,7 +845,138 @@ def upsert_comment(
|
|||||||
for t in tags:
|
for t in tags:
|
||||||
item.add_tag(t)
|
item.add_tag(t)
|
||||||
|
|
||||||
return store.upsert(item)
|
return store.upsert_status(item)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Docket walks ───────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def discover_docket(client: Client, docket_id: str, *, rule_cms_id: str = "") -> Docket:
|
||||||
|
"""One ``/documents`` listing → the docket's commentable document.
|
||||||
|
|
||||||
|
Picks the document that has an ``objectId`` and a ``commentEndDate``
|
||||||
|
(latest close date wins when several qualify). Called once per
|
||||||
|
docket; the result is persisted so later runs make no API call.
|
||||||
|
"""
|
||||||
|
best: dict | None = None
|
||||||
|
for fr_doc in client.find_documents_in_docket(docket_id):
|
||||||
|
attrs = fr_doc.get("attributes") or {}
|
||||||
|
if not attrs.get("objectId") or not attrs.get("commentEndDate"):
|
||||||
|
continue
|
||||||
|
if (
|
||||||
|
best is None
|
||||||
|
or attrs["commentEndDate"]
|
||||||
|
> (best.get("attributes") or {})["commentEndDate"]
|
||||||
|
):
|
||||||
|
best = fr_doc
|
||||||
|
attrs = (best or {}).get("attributes") or {}
|
||||||
|
return Docket(
|
||||||
|
id=docket_id,
|
||||||
|
rule_cms_id=rule_cms_id,
|
||||||
|
fr_document_id=(best or {}).get("id", ""),
|
||||||
|
fr_object_id=attrs.get("objectId", ""),
|
||||||
|
comment_end_date=(attrs.get("commentEndDate") or "")[:10],
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass
|
||||||
|
class WalkResult:
|
||||||
|
created: int = 0
|
||||||
|
updated: int = 0
|
||||||
|
unchanged: int = 0
|
||||||
|
watermark: str = ""
|
||||||
|
clean: bool = True
|
||||||
|
sealed: bool = False
|
||||||
|
|
||||||
|
|
||||||
|
def walk_docket(
|
||||||
|
store: Store,
|
||||||
|
client: Client,
|
||||||
|
docket: Docket,
|
||||||
|
*,
|
||||||
|
attachments: bool = False,
|
||||||
|
limit: int = 0,
|
||||||
|
force: bool = False,
|
||||||
|
quiet_days: int = 30,
|
||||||
|
today: date | None = None,
|
||||||
|
scratch_root: Path = Path(".state/comments"),
|
||||||
|
echo: Callable[[str], None] | None = None,
|
||||||
|
) -> WalkResult:
|
||||||
|
"""Pull *docket*'s comments from the stored watermark, upsert them,
|
||||||
|
advance the watermark on a clean walk, and auto-seal when
|
||||||
|
:func:`should_seal` says so. ``force`` walks from page 1 and never
|
||||||
|
seals. Returns per-status counts.
|
||||||
|
"""
|
||||||
|
from bib.dockets import should_seal
|
||||||
|
|
||||||
|
res = WalkResult()
|
||||||
|
since = "" if force else docket.pull_watermark
|
||||||
|
max_lm = docket.pull_watermark if not force else ""
|
||||||
|
n = 0
|
||||||
|
|
||||||
|
def _err(_e: Exception) -> None:
|
||||||
|
res.clean = False
|
||||||
|
|
||||||
|
scratch = scratch_root / docket.id
|
||||||
|
for c in client.iter_comments(docket.fr_object_id, since=since, on_error=_err):
|
||||||
|
if limit and n >= limit:
|
||||||
|
res.clean = False # a capped walk is not a complete one
|
||||||
|
break
|
||||||
|
key, status = upsert_comment_status(
|
||||||
|
store, c, cms_id=docket.rule_cms_id, extra_tags=[f"reg-docket:{docket.id}"]
|
||||||
|
)
|
||||||
|
setattr(res, status, getattr(res, status) + 1)
|
||||||
|
if attachments and c.attachment_count and status != "unchanged":
|
||||||
|
for att in client.attachments_for(c.id):
|
||||||
|
path = client.download_attachment(att.url, scratch / c.id)
|
||||||
|
if path:
|
||||||
|
store.attach_file(key, path, title=att.filename)
|
||||||
|
lm = (c.raw.get("attributes") or {}).get("lastModifiedDate") or ""
|
||||||
|
if lm > max_lm:
|
||||||
|
max_lm = lm
|
||||||
|
n += 1
|
||||||
|
if n % 50 == 0:
|
||||||
|
store._con().commit() # noqa: SLF001
|
||||||
|
if echo:
|
||||||
|
echo(f" {n} comments")
|
||||||
|
store._con().commit() # noqa: SLF001
|
||||||
|
res.watermark = max_lm
|
||||||
|
|
||||||
|
if res.clean and not force:
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
|
||||||
|
# ``today`` doubles as a deterministic override of "now" for
|
||||||
|
# tests exercising the auto-seal quiet-window boundary; real
|
||||||
|
# callers leave it unset and get the actual wall-clock time.
|
||||||
|
if today is not None:
|
||||||
|
now = f"{today.isoformat()}T00:00:00Z"
|
||||||
|
else:
|
||||||
|
now = datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||||
|
updated = replace(
|
||||||
|
docket, pull_watermark=max_lm, last_pull_at=now, last_pull_new=res.created
|
||||||
|
)
|
||||||
|
store.docket_upsert(updated)
|
||||||
|
if should_seal(updated, today or date.today(), quiet_days):
|
||||||
|
store.docket_seal(
|
||||||
|
docket.id, reason="auto", counts=docket_counts(store, docket.id)
|
||||||
|
)
|
||||||
|
res.sealed = True
|
||||||
|
return res
|
||||||
|
|
||||||
|
|
||||||
|
def docket_counts(store: Store, docket_id: str) -> dict[str, int]:
|
||||||
|
"""``{"comments": n, "enriched": n}`` from SQL only (no filesystem)."""
|
||||||
|
con = store._con() # noqa: SLF001
|
||||||
|
like = f"https://www.regulations.gov/comment/{docket_id}-%"
|
||||||
|
comments = con.execute(
|
||||||
|
"SELECT count(*) FROM items WHERE url LIKE ?", (like,)
|
||||||
|
).fetchone()[0]
|
||||||
|
enriched = con.execute(
|
||||||
|
"""SELECT count(*) FROM items i WHERE i.url LIKE ? AND i.id IN (
|
||||||
|
SELECT item_id FROM item_tags WHERE tag_id IN (SELECT id FROM tags WHERE name='enriched:ok'))""",
|
||||||
|
(like,),
|
||||||
|
).fetchone()[0]
|
||||||
|
return {"comments": comments, "enriched": enriched}
|
||||||
|
|
||||||
|
|
||||||
# ── Internals ──────────────────────────────────────────────────
|
# ── Internals ──────────────────────────────────────────────────
|
||||||
|
|||||||
@@ -156,7 +156,7 @@ class TestParseComment:
|
|||||||
class TestUpsertComment:
|
class TestUpsertComment:
|
||||||
def test_creates_item(self):
|
def test_creates_item(self):
|
||||||
store = MagicMock()
|
store = MagicMock()
|
||||||
store.upsert.return_value = "KEY1"
|
store.upsert_status.return_value = ("KEY1", "created")
|
||||||
c = Comment(
|
c = Comment(
|
||||||
id="CMS-2017-0092-0002",
|
id="CMS-2017-0092-0002",
|
||||||
title="Test",
|
title="Test",
|
||||||
@@ -168,8 +168,8 @@ class TestUpsertComment:
|
|||||||
)
|
)
|
||||||
key = upsert_comment(store, c, cms_id="CMS-1676-P")
|
key = upsert_comment(store, c, cms_id="CMS-1676-P")
|
||||||
assert key == "KEY1"
|
assert key == "KEY1"
|
||||||
store.upsert.assert_called_once()
|
store.upsert_status.assert_called_once()
|
||||||
item = store.upsert.call_args[0][0]
|
item = store.upsert_status.call_args[0][0]
|
||||||
assert "AMA" in item.title
|
assert "AMA" in item.title
|
||||||
assert item.url == "https://www.regulations.gov/comment/CMS-2017-0092-0002"
|
assert item.url == "https://www.regulations.gov/comment/CMS-2017-0092-0002"
|
||||||
|
|
||||||
|
|||||||
@@ -114,7 +114,7 @@ class TestByline:
|
|||||||
class TestUpsertComment:
|
class TestUpsertComment:
|
||||||
def test_full(self):
|
def test_full(self):
|
||||||
store = MagicMock()
|
store = MagicMock()
|
||||||
store.upsert.return_value = "KEY1"
|
store.upsert_status.return_value = ("KEY1", "created")
|
||||||
c = Comment(
|
c = Comment(
|
||||||
id="CMS-2023-0001-0001",
|
id="CMS-2023-0001-0001",
|
||||||
title="My Comment",
|
title="My Comment",
|
||||||
@@ -130,7 +130,7 @@ class TestUpsertComment:
|
|||||||
|
|
||||||
def test_no_title_uses_comment_on_id(self):
|
def test_no_title_uses_comment_on_id(self):
|
||||||
store = MagicMock()
|
store = MagicMock()
|
||||||
store.upsert.return_value = "KEY2"
|
store.upsert_status.return_value = ("KEY2", "created")
|
||||||
c = Comment(
|
c = Comment(
|
||||||
id="C2",
|
id="C2",
|
||||||
title="",
|
title="",
|
||||||
|
|||||||
136
tests/bib/test_walk_docket.py
Normal file
136
tests/bib/test_walk_docket.py
Normal file
@@ -0,0 +1,136 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from datetime import date
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
from bib.dockets import Docket
|
||||||
|
from bib.regulations_gov import Comment, WalkResult, discover_docket, walk_docket
|
||||||
|
from bib.store import Store
|
||||||
|
|
||||||
|
D = "CMS-2026-2377"
|
||||||
|
|
||||||
|
|
||||||
|
def _c(n: int, lm: str) -> Comment:
|
||||||
|
return Comment(
|
||||||
|
id=f"{D}-{n}",
|
||||||
|
title="",
|
||||||
|
posted_date=lm[:10],
|
||||||
|
received_date=lm[:10],
|
||||||
|
docket_id=D,
|
||||||
|
comment_on_id="x",
|
||||||
|
raw={"attributes": {"lastModifiedDate": lm}},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _store() -> Store:
|
||||||
|
return Store(":memory:", storage_dir="/tmp/nope")
|
||||||
|
|
||||||
|
|
||||||
|
def _docket(**kw) -> Docket:
|
||||||
|
base = dict(
|
||||||
|
id=D,
|
||||||
|
rule_cms_id="CMS-1848-P",
|
||||||
|
fr_object_id="obj",
|
||||||
|
comment_end_date="2026-09-14",
|
||||||
|
)
|
||||||
|
base.update(kw)
|
||||||
|
return Docket(**base)
|
||||||
|
|
||||||
|
|
||||||
|
def test_discover_docket_picks_commentable_doc():
|
||||||
|
api = MagicMock()
|
||||||
|
api.find_documents_in_docket.return_value = [
|
||||||
|
{"id": "X-1", "attributes": {"objectId": "o1"}}, # no comment window
|
||||||
|
{
|
||||||
|
"id": "X-2",
|
||||||
|
"attributes": {"objectId": "o2", "commentEndDate": "2026-09-14T03:59:59Z"},
|
||||||
|
},
|
||||||
|
]
|
||||||
|
d = discover_docket(api, D, rule_cms_id="CMS-1848-P")
|
||||||
|
assert d == Docket(
|
||||||
|
id=D,
|
||||||
|
rule_cms_id="CMS-1848-P",
|
||||||
|
fr_document_id="X-2",
|
||||||
|
fr_object_id="o2",
|
||||||
|
comment_end_date="2026-09-14",
|
||||||
|
)
|
||||||
|
api.find_documents_in_docket.assert_called_once_with(D)
|
||||||
|
|
||||||
|
|
||||||
|
def test_walk_creates_counts_and_advances_watermark():
|
||||||
|
s = _store()
|
||||||
|
s.docket_upsert(_docket())
|
||||||
|
api = MagicMock()
|
||||||
|
api.iter_comments.return_value = [
|
||||||
|
_c(1, "2026-09-01T00:00:00Z"),
|
||||||
|
_c(2, "2026-09-02T00:00:00Z"),
|
||||||
|
]
|
||||||
|
r = walk_docket(s, api, s.docket_get(D), today=date(2026, 9, 8))
|
||||||
|
assert r == WalkResult(
|
||||||
|
created=2,
|
||||||
|
updated=0,
|
||||||
|
unchanged=0,
|
||||||
|
watermark="2026-09-02T00:00:00Z",
|
||||||
|
clean=True,
|
||||||
|
sealed=False,
|
||||||
|
)
|
||||||
|
d = s.docket_get(D)
|
||||||
|
assert d.pull_watermark == "2026-09-02T00:00:00Z"
|
||||||
|
assert d.last_pull_new == 2 and d.last_pull_at
|
||||||
|
api.iter_comments.assert_called_once()
|
||||||
|
assert api.iter_comments.call_args.kwargs["since"] == ""
|
||||||
|
|
||||||
|
|
||||||
|
def test_walk_passes_watermark_and_counts_unchanged():
|
||||||
|
s = _store()
|
||||||
|
s.docket_upsert(_docket())
|
||||||
|
api = MagicMock()
|
||||||
|
api.iter_comments.return_value = [_c(1, "2026-09-01T00:00:00Z")]
|
||||||
|
walk_docket(s, api, s.docket_get(D), today=date(2026, 9, 8))
|
||||||
|
api.iter_comments.return_value = [
|
||||||
|
_c(1, "2026-09-01T00:00:00Z")
|
||||||
|
] # boundary re-yield
|
||||||
|
r = walk_docket(s, api, s.docket_get(D), today=date(2026, 9, 8))
|
||||||
|
assert api.iter_comments.call_args.kwargs["since"] == "2026-09-01T00:00:00Z"
|
||||||
|
assert (r.created, r.unchanged) == (0, 1)
|
||||||
|
assert s.docket_get(D).last_pull_new == 0
|
||||||
|
|
||||||
|
|
||||||
|
def test_unclean_walk_keeps_old_watermark():
|
||||||
|
s = _store()
|
||||||
|
s.docket_upsert(_docket(pull_watermark="2026-08-01T00:00:00Z"))
|
||||||
|
api = MagicMock()
|
||||||
|
|
||||||
|
def _iter(_obj, *, since="", on_error=None):
|
||||||
|
yield _c(9, "2026-09-05T00:00:00Z")
|
||||||
|
on_error(
|
||||||
|
httpx.HTTPStatusError("boom", request=MagicMock(), response=MagicMock())
|
||||||
|
)
|
||||||
|
|
||||||
|
api.iter_comments.side_effect = _iter
|
||||||
|
r = walk_docket(s, api, s.docket_get(D), today=date(2026, 9, 8))
|
||||||
|
assert r.clean is False and r.created == 1
|
||||||
|
assert s.docket_get(D).pull_watermark == "2026-08-01T00:00:00Z"
|
||||||
|
|
||||||
|
|
||||||
|
def test_force_walks_from_scratch_and_never_seals():
|
||||||
|
s = _store()
|
||||||
|
s.docket_upsert(_docket(pull_watermark="2026-08-01T00:00:00Z"))
|
||||||
|
api = MagicMock()
|
||||||
|
api.iter_comments.return_value = []
|
||||||
|
r = walk_docket(s, api, s.docket_get(D), force=True, today=date(2027, 1, 1))
|
||||||
|
assert api.iter_comments.call_args.kwargs["since"] == ""
|
||||||
|
assert r.sealed is False and not s.docket_get(D).sealed
|
||||||
|
|
||||||
|
|
||||||
|
def test_auto_seal_after_quiet_empty_pull():
|
||||||
|
s = _store()
|
||||||
|
s.docket_upsert(_docket())
|
||||||
|
api = MagicMock()
|
||||||
|
api.iter_comments.return_value = []
|
||||||
|
r = walk_docket(s, api, s.docket_get(D), today=date(2026, 10, 20), quiet_days=30)
|
||||||
|
assert r.sealed is True
|
||||||
|
d = s.docket_get(D)
|
||||||
|
assert d.sealed and d.seal_reason == "auto"
|
||||||
Reference in New Issue
Block a user