feat(comments): coordination + form-letter detection, no LLM required (refs #255)
Some checks failed
CI / lint (push) Successful in 45s
CI / notebooks-smoke (push) Successful in 1m31s
Deploy / notebooks (push) Has been skipped
Deploy / zotero (push) Has been skipped
Deploy / docs (push) Has been skipped
Deploy / api (push) Has been skipped
Deploy / mc (push) Has been skipped
Infra CI / notebooks (push) Successful in 57s
Infra CI / zotero (push) Successful in 14s
Infra CI / docs (push) Successful in 1m43s
CI / test (push) Has been cancelled
Deploy / report (push) Has been cancelled
Infra CI / mc (push) Has been cancelled
Infra CI / api (push) Has been cancelled

rex.comments.coordination: hti5 methodology — Jaccard similarity on
character 5-gram shingles (threshold 0.45) over normalized text,
union-find grouping with stable group ids, form-letter flag at >=3
members. Short texts (<200 chars) are excluded: two short 'I oppose'
notes are agreement, not coordination.

rex.comments.table: shared load path for skin_subs.rulemaking_comments
— outer-joins the classification (#254) and coordination (#255) JSONL
caches onto the comment identity rows, so whichever pass runs first
populates the table and the other enriches it without clobbering.
classify_comments.py refactored onto it; new analyze_coordination.py
driver.

Run against CMS-2025-0304 (the CY2026 OPPS skin-sub docket): of the
384 skin-sub-relevant comments, 203 (53%) are form letters across 31
coordinated groups — largest campaigns 72 and 36 members. Table loaded:
384 rows, coordination columns filled, classification columns NULL
until the #254 LLM run (blocked on API credits).

59 comments tests green.
This commit is contained in:
kert
2026-07-10 22:17:42 -04:00
parent 452f1587b8
commit c9b6a28f07
5 changed files with 457 additions and 112 deletions

View File

@@ -0,0 +1,81 @@
"""Coordination detection over skin-sub-relevant comments (#255).
No LLM involved: near-duplicate detection via Jaccard similarity on
character 5-gram shingles (hti5 methodology, threshold 0.45),
union-find grouping, and form-letter flagging by group size. Results
land in a JSONL cache and the docket's slice of
skin_subs.rulemaking_comments (merging the classification cache from
classify_comments.py when present).
Usage:
uv run python dev/scripts/analyze_coordination.py --docket CMS-2025-0304
"""
from __future__ import annotations
import argparse
import json
from collections import Counter
from pathlib import Path
ROOT = Path(__file__).resolve().parents[2]
BIB_PATH = ROOT / "data" / "bib.sqlite"
DEFAULT_DOCKET = "CMS-2025-0304"
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--docket", default=DEFAULT_DOCKET)
parser.add_argument(
"--no-load", action="store_true", help="analyze but skip the DuckDB load"
)
args = parser.parse_args()
from rex.comments.classify import is_relevant
from rex.comments.coordination import analyze
from rex.comments.table import load_comments, read_cache
comments = load_comments(args.docket, BIB_PATH)
relevant = [c for c in comments if is_relevant(c["text"])]
print(f"{args.docket}: {len(comments)} comments, {len(relevant)} skin-sub relevant")
results = analyze({c["key"]: c["text"] for c in relevant})
by_id = {c["key"]: c["comment_id"] for c in relevant}
groups = Counter(
r["coordination_group"] for r in results.values() if r["coordination_group"]
)
form = sum(1 for r in results.values() if r["is_form_letter"])
print(f"coordinated groups: {len(groups)}, form-letter comments: {form}")
for gid, n in groups.most_common(10):
print(f" group {by_id.get(gid, gid)}: {n} members")
base = ROOT / "data" / "cms"
base.mkdir(parents=True, exist_ok=True)
coordination_path = base / f"comments-coordination-{args.docket}.jsonl"
with coordination_path.open("w") as fh:
for key, rec in sorted(results.items()):
fh.write(json.dumps({"key": key, **rec}) + "\n")
print(f"cache → {coordination_path}")
if args.no_load:
return 0
from conf.connect import duckdb_batch
from rex.comments.table import load_table
classified_path = base / f"comments-classified-{args.docket}.jsonl"
with duckdb_batch("aco") as con:
n = load_table(
con,
args.docket,
relevant,
classified=read_cache(classified_path),
coordination=results,
)
print(f"skin_subs.rulemaking_comments: {n} rows for {args.docket}")
return 0
if __name__ == "__main__":
raise SystemExit(main())

View File

@@ -11,8 +11,9 @@ Pipeline per docket:
stakeholder / provisions / commenter+org) via prisma.llm. stakeholder / provisions / commenter+org) via prisma.llm.
Resumable: results append to a JSONL cache keyed by comment key; Resumable: results append to a JSONL cache keyed by comment key;
re-runs only classify new comments. re-runs only classify new comments.
4. Load the results into skin_subs.rulemaking_comments (replacing 4. Rebuild the docket's slice of skin_subs.rulemaking_comments,
that docket's rows and any DEMO- pilot rows). merging the coordination cache (analyze_coordination.py) when
present.
Usage: Usage:
uv run python dev/scripts/classify_comments.py --docket CMS-2025-0304 uv run python dev/scripts/classify_comments.py --docket CMS-2025-0304
@@ -22,8 +23,6 @@ Usage:
from __future__ import annotations from __future__ import annotations
import argparse import argparse
import json
import sqlite3
from pathlib import Path from pathlib import Path
ROOT = Path(__file__).resolve().parents[2] ROOT = Path(__file__).resolve().parents[2]
@@ -33,49 +32,12 @@ DEFAULT_DOCKET = "CMS-2025-0304"
DEFAULT_MODEL = "claude-haiku-4-5" DEFAULT_MODEL = "claude-haiku-4-5"
def load_comments(docket: str) -> list[dict]: def cache_paths(docket: str) -> tuple[Path, Path]:
"""One record per comment item: key, title, date, text, has_note.""" base = ROOT / "data" / "cms"
con = sqlite3.connect(BIB_PATH) return (
con.row_factory = sqlite3.Row base / f"comments-classified-{docket}.jsonl",
rows = con.execute( base / f"comments-coordination-{docket}.jsonl",
""" )
SELECT i.id, i.key, i.title, i.date_published, i.abstract,
(SELECT content FROM notes n
WHERE n.item_id = i.id AND n.title = 'Comment text') AS note
FROM items i
JOIN item_tags it ON i.id = it.item_id
JOIN tags t ON it.tag_id = t.id
WHERE t.name = ?
""",
(f"reg-docket:{docket}",),
).fetchall()
con.close()
from rex.comments.classify import strip_html
out = []
for r in rows:
text = strip_html(r["note"]) if r["note"] else (r["abstract"] or "")
out.append(
{
"key": r["key"],
"comment_id": r["title"] or r["key"],
"posted_date": r["date_published"] or "",
"text": text,
"has_attachments": bool(r["note"]),
}
)
return out
def load_cache(path: Path) -> dict[str, dict]:
cache: dict[str, dict] = {}
if path.exists():
for line in path.read_text().splitlines():
if line.strip():
rec = json.loads(line)
cache[rec["key"]] = rec
return cache
def main() -> int: def main() -> int:
@@ -93,9 +55,12 @@ def main() -> int:
) )
args = parser.parse_args() args = parser.parse_args()
from rex.comments.classify import is_relevant import json
comments = load_comments(args.docket) from rex.comments.classify import is_relevant
from rex.comments.table import load_comments, read_cache
comments = load_comments(args.docket, BIB_PATH)
relevant = [c for c in comments if is_relevant(c["text"])] relevant = [c for c in comments if is_relevant(c["text"])]
print(f"{args.docket}: {len(comments)} comments, {len(relevant)} skin-sub relevant") print(f"{args.docket}: {len(comments)} comments, {len(relevant)} skin-sub relevant")
if args.dry_run: if args.dry_run:
@@ -103,9 +68,9 @@ def main() -> int:
print(f" {c['comment_id'][:40]:40s} {len(c['text']):>7} chars") print(f" {c['comment_id'][:40]:40s} {len(c['text']):>7} chars")
return 0 return 0
cache_path = ROOT / "data" / "cms" / f"comments-classified-{args.docket}.jsonl" classified_path, coordination_path = cache_paths(args.docket)
cache_path.parent.mkdir(parents=True, exist_ok=True) classified_path.parent.mkdir(parents=True, exist_ok=True)
cache = load_cache(cache_path) cache = read_cache(classified_path)
todo = [c for c in relevant if c["key"] not in cache] todo = [c for c in relevant if c["key"] not in cache]
if args.limit: if args.limit:
todo = todo[: args.limit] todo = todo[: args.limit]
@@ -120,7 +85,7 @@ def main() -> int:
provider = make_provider() provider = make_provider()
errors = 0 errors = 0
with cache_path.open("a") as fh: with classified_path.open("a") as fh:
for i, c in enumerate(todo, 1): for i, c in enumerate(todo, 1):
try: try:
rec = classify(provider, c["text"], title=c["comment_id"]) rec = classify(provider, c["text"], title=c["comment_id"])
@@ -128,14 +93,7 @@ def main() -> int:
errors += 1 errors += 1
print(f" [{i}/{len(todo)}] {c['comment_id'][:36]} ERROR: {e}") print(f" [{i}/{len(todo)}] {c['comment_id'][:36]} ERROR: {e}")
continue continue
rec.update( rec.update(key=c["key"], docket_id=args.docket)
key=c["key"],
comment_id=c["comment_id"],
docket_id=args.docket,
posted_date=c["posted_date"],
has_attachments=c["has_attachments"],
text_length=len(c["text"]),
)
cache[c["key"]] = rec cache[c["key"]] = rec
fh.write(json.dumps(rec) + "\n") fh.write(json.dumps(rec) + "\n")
fh.flush() fh.flush()
@@ -149,62 +107,18 @@ def main() -> int:
if args.no_load: if args.no_load:
return 0 return 0
rows = [r for r in cache.values() if r.get("docket_id") == args.docket]
if not rows:
print("nothing to load")
return 0
from conf.connect import duckdb_batch from conf.connect import duckdb_batch
from rex.comments.table import load_table
with duckdb_batch("aco") as con: with duckdb_batch("aco") as con:
con.execute("CREATE SCHEMA IF NOT EXISTS skin_subs") n = load_table(
con.execute(""" con,
CREATE TABLE IF NOT EXISTS skin_subs.rulemaking_comments ( args.docket,
comment_id VARCHAR, docket_id VARCHAR, commenter_name VARCHAR, relevant,
organization VARCHAR, posted_date VARCHAR, has_attachments BOOLEAN, classified=cache,
text_length BIGINT, position VARCHAR, position_score DOUBLE, coordination=read_cache(coordination_path),
themes VARCHAR, stakeholder_type VARCHAR, coordination_group VARCHAR,
is_form_letter BOOLEAN, provisions VARCHAR
)
""")
# Replace this docket's rows; sweep any DEMO- pilot rows too.
con.execute(
"DELETE FROM skin_subs.rulemaking_comments "
"WHERE docket_id = ? OR comment_id LIKE 'DEMO-%'",
[args.docket],
) )
con.executemany(
"""
INSERT INTO skin_subs.rulemaking_comments
(comment_id, docket_id, commenter_name, organization, posted_date,
has_attachments, text_length, position, position_score, themes,
stakeholder_type, coordination_group, is_form_letter, provisions)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL, ?)
""",
[
(
r["comment_id"],
r["docket_id"],
r.get("commenter_name", ""),
r.get("organization", ""),
r.get("posted_date", ""),
r.get("has_attachments", False),
r.get("text_length", 0),
r["position"],
float(r["position_score"]),
r.get("themes", ""),
r.get("stakeholder_type", "unknown"),
r.get("provisions", ""),
)
for r in rows
],
)
n = con.execute(
"SELECT count(*) FROM skin_subs.rulemaking_comments WHERE docket_id = ?",
[args.docket],
).fetchone()[0]
print(f"skin_subs.rulemaking_comments: {n} rows for {args.docket}") print(f"skin_subs.rulemaking_comments: {n} rows for {args.docket}")
return 0 return 0

View File

@@ -0,0 +1,131 @@
"""Coordination detection for rulemaking comments (#255).
Detects organized comment campaigns without any LLM: near-duplicate
texts are found by Jaccard similarity on character 5-gram shingles
(hti5 methodology, threshold 0.45), connected into groups via
union-find, and groups above a size threshold are flagged as form
letters.
Answers the key question of #255: are manufacturers/distributors
orchestrating form-letter campaigns against the reclassification?
"""
from __future__ import annotations
import re
from collections import defaultdict
JACCARD_THRESHOLD = 0.45
SHINGLE_K = 5
FORM_LETTER_MIN_SIZE = 3
# Below this, shingle sets are too small for Jaccard to mean anything —
# two short "I oppose this rule" notes are agreement, not coordination.
MIN_TEXT_CHARS = 200
def normalize(text: str) -> str:
"""Lowercase, strip non-alphanumerics — punctuation/format noise must
not hide that two letters share a template."""
alnum = re.sub(r"[^a-z0-9]+", " ", (text or "").lower())
return re.sub(r"\s+", " ", alnum).strip()
def shingles(text: str, k: int = SHINGLE_K) -> frozenset[str]:
norm = re.sub(r"\s+", " ", normalize(text))
if len(norm) < k:
return frozenset()
return frozenset(norm[i : i + k] for i in range(len(norm) - k + 1))
def jaccard(a: frozenset[str], b: frozenset[str]) -> float:
if not a or not b:
return 0.0
inter = len(a & b)
return inter / (len(a) + len(b) - inter)
class _UnionFind:
def __init__(self) -> None:
self._parent: dict[str, str] = {}
def find(self, x: str) -> str:
self._parent.setdefault(x, x)
while self._parent[x] != x:
self._parent[x] = self._parent[self._parent[x]]
x = self._parent[x]
return x
def union(self, a: str, b: str) -> None:
ra, rb = self.find(a), self.find(b)
if ra != rb:
self._parent[rb] = ra
def find_groups(
texts: dict[str, str],
*,
threshold: float = JACCARD_THRESHOLD,
min_chars: int = MIN_TEXT_CHARS,
) -> dict[str, str]:
"""Group near-duplicate comments.
Parameters
----------
texts
``{comment_id: text}``.
Returns
-------
dict
``{comment_id: group_id}`` for every comment in a group of two
or more. The group id is the lexicographically first member id,
so it is stable across runs.
"""
sh = {cid: shingles(t) for cid, t in texts.items() if len(t or "") >= min_chars}
ids = sorted(sh)
uf = _UnionFind()
for i, a in enumerate(ids):
for b in ids[i + 1 :]:
if jaccard(sh[a], sh[b]) >= threshold:
uf.union(a, b)
members: dict[str, list[str]] = defaultdict(list)
for cid in ids:
members[uf.find(cid)].append(cid)
out: dict[str, str] = {}
for group in members.values():
if len(group) < 2:
continue
gid = min(group)
for cid in group:
out[cid] = gid
return out
def form_letter_groups(
groups: dict[str, str], *, min_size: int = FORM_LETTER_MIN_SIZE
) -> set[str]:
"""Group ids whose membership is large enough to call a campaign."""
sizes: dict[str, int] = defaultdict(int)
for gid in groups.values():
sizes[gid] += 1
return {gid for gid, n in sizes.items() if n >= min_size}
def analyze(texts: dict[str, str]) -> dict[str, dict]:
"""Full pass: ``{comment_id: {coordination_group, is_form_letter}}``.
Every input id gets a record (group ``""`` when uncoordinated) so
downstream loads can distinguish "analyzed, independent" from
"never analyzed" (NULL).
"""
groups = find_groups(texts)
campaigns = form_letter_groups(groups)
return {
cid: {
"coordination_group": groups.get(cid, ""),
"is_form_letter": groups.get(cid, "") in campaigns,
}
for cid in texts
}

134
src/rex/comments/table.py Normal file
View File

@@ -0,0 +1,134 @@
"""skin_subs.rulemaking_comments — shared load path.
Both analysis passes (LLM classification #254, coordination detection
#255) cache their results as JSONL keyed by bib item key; this module
outer-joins the caches onto the comment identity rows and rebuilds the
docket's slice of the table. Whichever pass runs first populates the
table; the other enriches it without clobbering.
"""
from __future__ import annotations
import json
import sqlite3
from pathlib import Path
from typing import Any
def load_comments(docket: str, bib_path: str | Path) -> list[dict[str, Any]]:
"""One record per bib comment item for *docket*.
Text prefers the "Comment text" note (combined.md — inline body plus
every extracted attachment) and falls back to the inline abstract.
"""
from rex.comments.classify import strip_html
con = sqlite3.connect(str(bib_path))
con.row_factory = sqlite3.Row
rows = con.execute(
"""
SELECT i.id, i.key, i.title, i.date_published, i.abstract,
(SELECT content FROM notes n
WHERE n.item_id = i.id AND n.title = 'Comment text') AS note
FROM items i
JOIN item_tags it ON i.id = it.item_id
JOIN tags t ON it.tag_id = t.id
WHERE t.name = ?
""",
(f"reg-docket:{docket}",),
).fetchall()
con.close()
out = []
for r in rows:
text = strip_html(r["note"]) if r["note"] else (r["abstract"] or "")
out.append(
{
"key": r["key"],
"comment_id": r["title"] or r["key"],
"posted_date": r["date_published"] or "",
"text": text,
"has_attachments": bool(r["note"]),
}
)
return out
def read_cache(path: Path) -> dict[str, dict]:
"""JSONL keyed by bib item key; later lines win."""
cache: dict[str, dict] = {}
if path.exists():
for line in path.read_text().splitlines():
if line.strip():
rec = json.loads(line)
cache[rec["key"]] = rec
return cache
def load_table(
con: Any,
docket: str,
comments: list[dict[str, Any]],
*,
classified: dict[str, dict] | None = None,
coordination: dict[str, dict] | None = None,
) -> int:
"""Rebuild the docket's rows in skin_subs.rulemaking_comments.
*comments* is the identity slice to load (typically the
skin-sub-relevant subset). Classification and coordination fields
come from their caches when present, NULL otherwise. DEMO- pilot
rows are swept. Returns the docket's row count after the load.
"""
classified = classified or {}
coordination = coordination or {}
con.execute("CREATE SCHEMA IF NOT EXISTS skin_subs")
con.execute("""
CREATE TABLE IF NOT EXISTS skin_subs.rulemaking_comments (
comment_id VARCHAR, docket_id VARCHAR, commenter_name VARCHAR,
organization VARCHAR, posted_date VARCHAR, has_attachments BOOLEAN,
text_length BIGINT, position VARCHAR, position_score DOUBLE,
themes VARCHAR, stakeholder_type VARCHAR, coordination_group VARCHAR,
is_form_letter BOOLEAN, provisions VARCHAR
)
""")
con.execute(
"DELETE FROM skin_subs.rulemaking_comments "
"WHERE docket_id = ? OR comment_id LIKE 'DEMO-%'",
[docket],
)
con.executemany(
"""
INSERT INTO skin_subs.rulemaking_comments
(comment_id, docket_id, commenter_name, organization, posted_date,
has_attachments, text_length, position, position_score, themes,
stakeholder_type, coordination_group, is_form_letter, provisions)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
[
(
c["comment_id"],
docket,
cl.get("commenter_name", "") or "",
cl.get("organization", "") or "",
c["posted_date"],
c["has_attachments"],
len(c["text"]),
cl.get("position"),
cl.get("position_score"),
cl.get("themes"),
cl.get("stakeholder_type"),
co.get("coordination_group"),
co.get("is_form_letter"),
cl.get("provisions"),
)
for c in comments
for cl in [classified.get(c["key"], {})]
for co in [coordination.get(c["key"], {})]
],
)
return con.execute(
"SELECT count(*) FROM skin_subs.rulemaking_comments WHERE docket_id = ?",
[docket],
).fetchone()[0]

View File

@@ -0,0 +1,85 @@
"""Tests for rex.comments.coordination — form-letter/campaign detection."""
from __future__ import annotations
from rex.comments.coordination import (
analyze,
find_groups,
form_letter_groups,
jaccard,
normalize,
shingles,
)
TEMPLATE = (
"I am writing to strongly oppose the proposed reclassification of skin "
"substitutes from biologicals to incident-to supplies. The flat rate of "
"$127.28 per square centimeter will devastate patient access to advanced "
"wound care products and harm innovation in cellular tissue-based "
"products. I urge CMS to withdraw this proposal and preserve ASP plus "
"six percent reimbursement for these critical therapies."
)
# Same template, light personalization — the classic form-letter signature.
VARIANT_A = "Dear Administrator, " + TEMPLATE + " Sincerely, Dr. Alice Smith, DPM"
VARIANT_B = "To whom it may concern: " + TEMPLATE + " Regards, Bob Jones, NP"
VARIANT_C = TEMPLATE + " Respectfully submitted, Carol Lee, Wound Care Clinic"
INDEPENDENT = (
"As a health economist who has studied the skin substitute market for a "
"decade, I support the proposed payment change. Average sales price "
"inflation in this sector reflects strategic price-setting rather than "
"clinical value, and the fraud indictments of the past two years show "
"the incentive structure is broken. A packaged rate is overdue, though "
"CMS should phase it in over two years to avoid access cliffs in rural "
"areas where few wound-care alternatives exist for beneficiaries."
)
class TestPrimitives:
def test_normalize_strips_punctuation_and_case(self):
assert normalize("The $127.28/cm² Rate!") == "the 127 28 cm rate"
def test_shingles_short_text_empty(self):
assert shingles("abc") == frozenset()
def test_jaccard_identity_and_disjoint(self):
a = shingles(TEMPLATE)
assert jaccard(a, a) == 1.0
assert jaccard(a, frozenset()) == 0.0
class TestFindGroups:
def test_template_variants_group_independent_stays_out(self):
groups = find_groups(
{
"c-a": VARIANT_A,
"c-b": VARIANT_B,
"c-c": VARIANT_C,
"c-x": INDEPENDENT,
}
)
assert groups.get("c-a") == groups.get("c-b") == groups.get("c-c") == "c-a"
assert "c-x" not in groups
def test_short_texts_never_grouped(self):
groups = find_groups({"s1": "I oppose this rule", "s2": "I oppose this rule"})
assert groups == {}
def test_group_id_stable_regardless_of_order(self):
texts = {"z-late": VARIANT_A, "a-early": VARIANT_B}
assert set(find_groups(texts).values()) == {"a-early"}
class TestFormLetters:
def test_min_size_threshold(self):
groups = {"a": "g1", "b": "g1", "c": "g1", "d": "g2", "e": "g2"}
assert form_letter_groups(groups) == {"g1"}
class TestAnalyze:
def test_every_input_gets_a_record(self):
results = analyze(
{"c-a": VARIANT_A, "c-b": VARIANT_B, "c-c": VARIANT_C, "c-x": INDEPENDENT}
)
assert set(results) == {"c-a", "c-b", "c-c", "c-x"}
assert results["c-a"]["is_form_letter"] is True
assert results["c-x"] == {"coordination_group": "", "is_form_letter": False}