ishaq101's picture
/fix parsing and term extract (#21)
f07443e
Raw History Blame Contribute Delete
26.7 kB
"""Knowledge-pipeline eval runner.
Scores an already-parsed document against `knowledge_gold.yaml` and writes a
timestamped result file. The 2026-08-24 figures were assembled by hand; this is
the reproducible version, so a number can be re-derived instead of trusted.
**Takes a ParsedDocument artifact, never a PDF.** MinerU is not invoked and is
not required — that is the seam doing its job. Point it at the `chunks.json`
written by `python -m src.knowledge_parsing.run`.
# free: E1 recall, E2 compression, the funnel. No API calls, no spend.
./.venv/Scripts/python.exe -m eval.knowledge.run_eval <artifact.json>
# add the paid branches: E3 schema fill and rule recall
./.venv/Scripts/python.exe -m eval.knowledge.run_eval <artifact.json> --extract
# exercise every branch with no spend and no credentials (NOT scoreable)
./.venv/Scripts/python.exe -m eval.knowledge.run_eval <artifact.json> --extract --mock
Module mode only (`-m`). A file-path invocation breaks `sys.path` and the
`src.` imports below fail.
# a BAND instead of a number — repeats the paid stage only. SPENDS N TIMES.
./.venv/Scripts/python.exe -m eval.knowledge.run_eval <artifact.json> --extract --runs 5
Exit codes: `0` clean · `1` below a kill line · `2` bad arguments or artifact ·
`3` no extractor · `4` would overwrite a result · `5` an invariant broke ·
`6` a different model snapshot served the run.
Four things this refuses to do, all for the same reason — a number that cannot
be compared is worse than no number:
* **It will not run on an artifact that does not declare `parser_backend`.**
The same MinerU build emits different text from `pipeline` and `hybrid/high`,
so results measured on different backends are not comparable and nothing
downstream can tell them apart afterwards. `--allow-unlabelled-artifact`
overrides, and stamps the result file so the caveat travels with it.
* **It never overwrites a result file.** Every run writes a new timestamped
file beside the others; the frozen baseline stays frozen.
* **It refuses a run served by a different model snapshot** than the committed
results were measured on (`EXPECTED_MODEL_SNAPSHOT`). X25 recorded what served
a run; this asserts it, so a silent tier move stops the comparison instead of
quietly poisoning the delta column.
* **It refuses a run that broke a structural invariant** (`invariants.py`) —
a span that is not locatable, a duplicate entity id, a repaired rejection.
Those are defects rather than scores, so they fail the run outright.
"""
from __future__ import annotations
import argparse
import json
import statistics
import sys
from collections import Counter
from datetime import datetime
from pathlib import Path
from typing import Any
from src.knowledge_extraction.adapter import parsed_doc_from_artifact
from src.knowledge_extraction.extract import MockExtractor, cacheable, prefix_tokens
from src.knowledge_extraction.genre import classify as classify_genre
from src.knowledge_extraction.service import (
build_clusters,
estimate_cost,
extract_all,
run_filters,
)
from src.knowledge_extraction.settings import EVIDENCE_K
from .invariants import check_invariants
from .score import (
GOLD_PATH,
load_gold,
score_glossary,
score_post_cluster,
score_rules,
score_term_filter,
)
_HERE = Path(__file__).resolve().parent
RESULTS_DIR = _HERE / "results"
# The kill lines. Each is a decision already taken, recorded in
# KNOWLEDGE_PIPELINE_CALIBRATION.md — not a target invented per run.
# E2's kill line was REMOVED 2026-09-10. Compression is now reported as a
# descriptive statistic only, never as pass/fail.
#
# The reason is not that 2.0 was the wrong number — it is that the metric moves
# the WRONG WAY. Over-merging distinct terms into one cluster RAISES compression,
# so the defect that cost the gold term EWH its own entry was pushing this score
# up and through its own pass line. A metric whose pass condition is partly met
# by the defect it should expose is worse than no metric. Retiring the fuzzy pass
# correctly LOWERED compression (BUMA 2.712 -> 2.183), which under a kill line
# would have read as a regression.
#
# E1c (post-cluster recall) is the metric that actually watches this stage.
KILL_LINES = {"E1": 0.70, "E3": 0.80}
# What the last committed run measured, for the delta column. Rule recall has
# no kill line: it is tracked, not gated, and 3/15 is the post-X22 scoring of
# the 2026-08-24 run — NOT the 1/15 in that file, which predates the scorer fix.
# Re-baselined to MISTRAL 2026-09-02 (M1). The delta column compares like with
# like or it lies: `mistral` is now the default backend (knowledge_parsing/
# config.py), so leaving MinerU's numbers here would have printed every future
# run against a parse it cannot be compared to — the same cross-backend
# comparison this harness refuses elsewhere by rejecting unlabelled artifacts.
#
# Previous MinerU baseline, kept for the record, NOT for comparison:
# v2_full_document_2026-08-24_093051.json — E1 0.8049 · E2 2.81 · E3 0.90 · rule 0.2
BASELINE = {
"file": "v2_full_2026-09-02_160326.json",
"E1": 0.8780,
"E2": 2.6970,
"E3": 1.0000,
"rule_recall": 0.5,
}
# ── Step 0: measurement hardening ───────────────────────────────────────
# The snapshot every committed result on this branch was measured against. X25
# started RECORDING what served a run; this makes it an ASSERTION, which is the
# difference between a fact discoverable by diffing two files afterwards and one
# that stops a run at the moment it stops being true.
#
# A deployment name is OUR label and never changes; the snapshot behind it can
# move without notice. When it does, every number in `results/` becomes a
# measurement of a different system and the delta column starts lying. Comparing
# against a moved model is the failure this catches — not the move itself, which
# is Azure's business.
#
# `--expect-model` overrides for a deliberate tier change; `--allow-model-drift`
# downgrades the refusal to a stamped warning when you need the number anyway.
EXPECTED_MODEL_SNAPSHOT = "gpt-5.4-nano-2026-03-17"
# Below this spread, two runs of one document are telling the same story. Above
# it, a single run is not evidence and the band must be quoted instead of the
# point. 0.10 is the handoff's own read of this corpus (identical input scored
# 6/7 one day and 0/12 the next), not a measured tolerance — raise or lower it
# once `--runs` has produced a real band on a real document.
NOISE_BAND = 0.10
def _azure_extractor():
"""Built lazily: importing it requires credentials to be configured."""
from src.knowledge_extraction.extract import AzureExtractor
try:
return AzureExtractor()
except Exception as exc: # noqa: BLE001 - reported, not raised
print(f"cannot build the Azure extractor: {exc!r}", file=sys.stderr)
return None
def _latex_verification_counts(formulas: list[dict]) -> dict[str, int]:
"""How many stored formulas carry which guarantee.
Cheap and free — the marker is set during validation, not by a model call —
and it is the fastest read on whether an artifact's markup is rich enough
to support the guard at all. A run that is all `unverified_operator` was
parsed with a backend that cannot express multiplication.
"""
counts = Counter(f.get("latex_verification") or "none" for f in formulas)
return dict(sorted(counts.items()))
def _verdict(value: float, kill: float | None) -> str:
if kill is None:
return "—"
return "PASS" if value >= kill else "FAIL"
def _delta(value: float, previous: float | None) -> str:
if previous is None:
return "—"
diff = round(value - previous, 4)
return f"{diff:+.4f}"
def _stability(values: list[float]) -> dict[str, Any]:
"""A band, not a number.
The reported value becomes the MEDIAN rather than the mean: one run that
fails schema validation drags a mean into a region no run actually
occupied, and a metric nobody measured is worse than a noisy one.
"""
ordered = sorted(values)
return {
"n": len(ordered),
"min": round(ordered[0], 4),
"median": round(statistics.median(ordered), 4),
"max": round(ordered[-1], 4),
"spread": round(ordered[-1] - ordered[0], 4),
"runs": [round(v, 4) for v in values],
"noisy": bool(ordered[-1] - ordered[0] > NOISE_BAND),
}
def _entity_ids(result) -> dict[str, list[str]]:
"""Every id a run minted, by kind — the thing an expert's approval hangs on.
K11 fixed two ids that did not survive a re-run. Nothing has watched them
since, and a re-run that re-keys its entries silently orphans every approval
recorded against the old ones. Cheap to check, catastrophic to miss.
"""
return {
"glossary": sorted(e.get("term_id") or "" for e in result.glossary),
"rule": sorted(e.get("rule_id") or "" for e in result.rules),
"formula": sorted(e.get("formula_id") or "" for e in result.formulas),
}
def _id_stability(per_run: list[dict[str, list[str]]]) -> dict[str, Any]:
"""Did every run mint the same ids? Anything but `True` is a defect."""
if len(per_run) < 2:
return {"comparable": False, "note": "needs --runs 2 or more"}
first = per_run[0]
drift = {
kind: {
"only_in_first": sorted(set(first[kind]) - set(other[kind]))[:10],
"only_in_later": sorted(set(other[kind]) - set(first[kind]))[:10],
}
for other in per_run[1:]
for kind in first
if set(first[kind]) != set(other[kind])
}
return {"comparable": True, "identical_across_runs": not drift, "drift": drift}
def _served_snapshots(result) -> list[str]:
return sorted({u.model_version for u in result.usages if u.model_version})
def build_report(
doc,
filtered,
clustered,
scores: dict[str, Any],
result,
args,
started: datetime,
seconds: float,
invariants=None,
stability: dict[str, Any] | None = None,
) -> dict:
funnel = {
"pages": doc.n_pages,
"chunks": len(doc.chunks),
"mentions_raw": len(filtered.mentions),
"clusters": clustered.n_clusters,
"compression_ratio": clustered.compression_ratio,
"abbrev_pairs": len(filtered.abbrev_pairs),
"rule_candidates": len(filtered.rule_candidates),
}
report: dict[str, Any] = {
"run": {
"created_at": started.strftime("%Y-%m-%d_%H%M%S"),
"doc_id": doc.doc_id,
"artifact": str(args.artifact),
"content_hash": doc.content_hash,
"parser": f"{doc.parser_name}/{doc.parser_version}",
"parser_backend": doc.parser_backend or "UNLABELLED",
"extracted": bool(args.extract),
"mock": bool(args.mock),
"limit": args.limit,
"runs": args.runs,
"expected_model": args.expect_model,
"seconds": round(seconds, 1),
},
"funnel": funnel,
"genre": classify_genre(doc, filtered).as_dict(),
"scores": scores,
"baseline": BASELINE,
}
if invariants is not None:
report["invariants"] = invariants.as_dict()
if stability:
report["stability"] = stability
if args.allow_unlabelled_artifact and not doc.parser_backend:
report["run"]["comparability"] = (
"NOT COMPARABLE — the artifact does not declare parser_backend and the "
"run was forced. Different MinerU backends emit different text."
)
if args.mock:
report["run"]["comparability"] = (
"NOT SCOREABLE — mock extractor. Branch plumbing only; the content is "
"fabricated, so every quality figure here is meaningless."
)
if result is not None:
# Which model actually served this run. A result file without it is a
# measurement of an unknown: on 2026-08-27 identical input scored 6/7
# one day and 0/12 the next under an unchanged deployment name, and
# nothing recorded said whether the model had moved. See X25.
served = _served_snapshots(result)
report["run"]["served_by"] = {
"model_version": served,
"system_fingerprint": sorted(
{u.system_fingerprint for u in result.usages if u.system_fingerprint}
),
}
# Step 0(a): the recorded snapshot becomes an assertion. A mock run
# reports no snapshot and is exempt — it measures plumbing, not a model.
if not args.mock and served and served != [args.expect_model]:
report["run"]["model_drift"] = {
"expected": args.expect_model,
"served": served,
"forced": bool(args.allow_model_drift),
"note": (
"every committed result was measured on the expected snapshot; "
"a delta against this run compares two different systems"
),
}
prompt, cached, completion = result.total_tokens
report["cost"] = {
"prompt_tokens": prompt,
"cached_tokens": cached,
"completion_tokens": completion,
"cache_hit_rate": round(cached / prompt, 3) if prompt else 0.0,
"calls": len(result.usages),
}
report["quality"] = {
"fields_rejected_by_span_check": len(result.rejected),
"rejections": [r.model_dump(mode="json") for r in result.rejected[:20]],
"glossary_entries": len(result.glossary),
"rule_entries": len(result.rules),
"formula_entries": len(result.formulas),
"abstained_null_definition": sum(
1 for e in result.glossary if not e.get("definition")
),
"latex_verification": _latex_verification_counts(result.formulas),
"links": result.links,
}
return report
def format_summary(report: dict) -> str:
rows = []
for key, score in report["scores"].items():
value = score["value"]
kill = KILL_LINES.get(key)
band = score.get("stability")
suffix = ""
if band:
suffix = (
f" [{band['min']:.4f}..{band['max']:.4f}] n={band['n']}"
+ (" NOISY" if band["noisy"] else "")
)
rows.append(
f" {key:<12} {value:>8.4f} kill {kill if kill else ' — '!s:>5} "
f"{_verdict(value, kill):<5} vs baseline {_delta(value, BASELINE.get(key))}"
f" {score['label']}{suffix}"
)
funnel = report["funnel"]
out = [
"",
f" doc {report['run']['doc_id']} backend {report['run']['parser_backend']}"
f" {funnel['pages']} pages -> {funnel['chunks']} chunks -> "
f"{funnel['mentions_raw']} mentions -> {funnel['clusters']} clusters "
f"({funnel['compression_ratio']}x)",
"",
*rows,
]
if "cost" in report:
cost = report["cost"]
out += [
"",
f" {cost['calls']} calls {cost['prompt_tokens']} prompt tokens, "
f"{cost['cached_tokens']} cached ({cost['cache_hit_rate']:.0%})",
f" {report['quality']['fields_rejected_by_span_check']} fields rejected"
" by the span check",
f" latex_verification: {report['quality']['latex_verification']}",
]
genre = report.get("genre")
if genre:
out += [
"",
f" genre: {genre['genre']} (labels {genre['labels_variant']}, "
f"well-served {genre['well_served']}) {genre['signals']}",
]
for reason in genre["reasons"]:
out += [f" · {reason}"]
for warning in genre["warnings"]:
out += [f" !! {warning}"]
if "invariants" in report:
inv = report["invariants"]
out += [
"",
f" invariants: {'OK' if inv['ok'] else str(inv['n_violations']) + ' VIOLATED'}"
f" checked {inv['checked']}",
]
for violation in inv["violations"][:10]:
out += [f" !! {violation['check']} [{violation['kind']}] {violation['detail']}"]
if "stability" in report:
ids = report["stability"]["entity_ids"]
if ids.get("comparable"):
out += [
f" entity ids identical across runs: {ids['identical_across_runs']}"
]
if "model_drift" in report["run"]:
drift = report["run"]["model_drift"]
out += [
"",
f" !! MODEL DRIFT — expected {drift['expected']}, served "
f"{', '.join(drift['served'])}",
]
if "comparability" in report["run"]:
out += ["", f" !! {report['run']['comparability']}"]
return "\n".join(out)
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="Score the knowledge pipeline against knowledge_gold.yaml"
)
parser.add_argument("artifact", type=Path, help="ParsedDocument JSON (chunks.json)")
parser.add_argument("--doc-id", help="override the artifact's doc_id")
parser.add_argument(
"--extract", action="store_true", help="run the PAID branches (E3, rule recall)"
)
parser.add_argument(
"--mock", action="store_true", help="mock extractor: no network, no spend, NOT scoreable"
)
parser.add_argument("--limit", type=int, help="cap items per branch (pilot)")
parser.add_argument("--active-glossary", type=Path, help="glossary to diff against")
parser.add_argument("--gold", type=Path, default=GOLD_PATH)
parser.add_argument("--label", help="short note recorded in the result file")
parser.add_argument(
"--dry-run", action="store_true", help="token estimate only, zero API calls"
)
parser.add_argument(
"--allow-unlabelled-artifact",
action="store_true",
help="run on an artifact with no parser_backend (results NOT comparable)",
)
parser.add_argument(
"--runs",
type=int,
default=1,
metavar="N",
help=(
"repeat the PAID stage N times and report min/median/max instead of a "
"single number. MULTIPLIES SPEND BY N — the free stages run once"
),
)
parser.add_argument(
"--expect-model",
default=EXPECTED_MODEL_SNAPSHOT,
help=f"model snapshot this run must be served by (default {EXPECTED_MODEL_SNAPSHOT})",
)
parser.add_argument(
"--allow-model-drift",
action="store_true",
help="record a snapshot mismatch as a warning instead of refusing the run",
)
args = parser.parse_args(argv)
if args.runs < 1:
print("--runs must be 1 or more", file=sys.stderr)
return 2
if args.runs > 1 and not args.extract:
print(
"--runs only varies the PAID stage; the free stages are deterministic.\n"
" Add --extract, or drop --runs.",
file=sys.stderr,
)
return 2
if not args.artifact.exists():
print(f"artifact not found: {args.artifact}", file=sys.stderr)
return 2
raw = json.loads(args.artifact.read_text(encoding="utf-8"))
doc = parsed_doc_from_artifact(raw, doc_id=args.doc_id, source_ref=str(args.artifact))
if not doc.parser_backend and not args.allow_unlabelled_artifact:
print(
f"artifact does not declare parser_backend: {args.artifact}\n"
" A score measured on one MinerU backend is not comparable to one\n"
" measured on another, and nothing downstream can tell them apart\n"
" afterwards. Re-export from `src.knowledge_parsing.run`, or pass\n"
" --allow-unlabelled-artifact to record the result as uncomparable.",
file=sys.stderr,
)
return 2
started = datetime.now()
gold = load_gold(args.gold)
filtered = run_filters(doc)
clustered = build_clusters(doc, filtered)
if args.dry_run:
est = estimate_cost(doc, clustered, filtered, args.limit)
print("[dry-run] NO API CALLS MADE")
for key, value in est.items():
print(f" {key}: {value}")
for branch in ("glossary", "rule", "formula", "summary"):
print(
f" prefix[{branch}]: {prefix_tokens(branch)} tokens, "
f"cacheable={cacheable(branch)}"
)
return 0
e1 = score_term_filter(gold, [m.surface for m in filtered.mentions])
canonicals = [c.canonical for c in clustered.clusters]
e1c_loose = score_post_cluster(gold, canonicals)
e1c_strict = score_post_cluster(gold, canonicals, cluster_variants=True)
scores: dict[str, Any] = {
"E1": {
"label": "term-filter recall",
"value": e1.recall,
"detail": e1.as_dict(),
},
# E1c watches the merge stage, which E1 (measured before clustering) and
# E3 (measured only over entries that produced a definition) both miss.
# `E1 - E1c_loose` is the recall clustering destroyed; 0 is the target.
"E1c_loose": {
"label": "post-cluster recall (loose)",
"value": e1c_loose.recall,
"detail": e1c_loose.as_dict(),
},
"E1c_strict": {
"label": "post-cluster recall (owns its own cluster)",
"value": e1c_strict.recall,
"detail": e1c_strict.as_dict(),
},
"E2": {
"label": "review-burden compression (descriptive, NOT gated)",
"value": clustered.compression_ratio,
"detail": {
"mentions": clustered.n_mentions,
"clusters": clustered.n_clusters,
},
},
}
result = None
invariants = None
stability: dict[str, Any] = {}
if args.extract:
extractor = MockExtractor() if args.mock else _azure_extractor()
if extractor is None:
return 3
active = (
json.loads(args.active_glossary.read_text(encoding="utf-8"))
if args.active_glossary
else None
)
# The free stages are deterministic and already done — only the paid
# stage is repeated. Re-running GLiNER N times would cost minutes to
# re-derive an identical clustering.
e3_values: list[float] = []
rule_values: list[float] = []
ids_per_run: list[dict[str, list[str]]] = []
violations_per_run: list[int] = []
for run_index in range(args.runs):
if args.runs > 1:
print(f" run {run_index + 1}/{args.runs} ...", file=sys.stderr)
result = extract_all(doc, clustered, filtered, extractor, args.limit, active)
e3 = score_glossary(gold, result.glossary)
rules = score_rules(gold, result.rules)
e3_values.append(e3.precision)
rule_values.append(rules.recall)
ids_per_run.append(_entity_ids(result))
invariants = check_invariants(result, doc)
violations_per_run.append(len(invariants.violations))
# The LAST run's detail is reported beside the band. Detail from one run
# of several is an example, not a summary — labelled as such below.
scores["E3"] = {
"label": "glossary schema fill (precision)",
"value": round(statistics.median(e3_values), 4),
"detail": e3.as_dict(),
}
scores["rule_recall"] = {
"label": "rule-branch recall",
"value": round(statistics.median(rule_values), 4),
"detail": rules.as_dict(),
}
if args.runs > 1:
scores["E3"]["stability"] = _stability(e3_values)
scores["rule_recall"]["stability"] = _stability(rule_values)
for key in ("E3", "rule_recall"):
scores[key]["detail"]["note"] = (
f"detail is from run {args.runs} of {args.runs}; the value is "
"the median across runs"
)
stability = {
"runs": args.runs,
"entity_ids": _id_stability(ids_per_run),
"violations_per_run": violations_per_run,
}
seconds = (datetime.now() - started).total_seconds()
report = build_report(
doc, filtered, clustered, scores, result, args, started, seconds,
invariants=invariants, stability=stability,
)
if args.label:
report["run"]["label"] = args.label
if not args.extract:
report["run"]["note"] = (
"Free stages only. E3 and rule recall require --extract, which spends."
)
print(format_summary(report))
RESULTS_DIR.mkdir(parents=True, exist_ok=True)
stem = "mock" if args.mock else ("full" if args.extract else "filters")
out_path = RESULTS_DIR / f"v2_{stem}_{started:%Y-%m-%d_%H%M%S}.json"
if out_path.exists(): # never overwrite a committed measurement
print(f"refusing to overwrite {out_path}", file=sys.stderr)
return 4
out_path.write_text(
json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8"
)
print(f"\n-> saved: {out_path.relative_to(_HERE.parent.parent)}")
# Invariants first, and regardless of the tier: a broken structural claim
# means a control leaked, which invalidates the scores rather than joining
# them. Mock runs included — the plumbing is exactly what they test.
if invariants is not None and not invariants.ok:
print(
f"\nINVARIANTS VIOLATED: {len(invariants.violations)} — see the result file",
file=sys.stderr,
)
return 5
if "model_drift" in report["run"] and not args.allow_model_drift:
print(
"\nMODEL SNAPSHOT MISMATCH — the run is recorded but must not be compared\n"
" to the committed results. Re-pin with --expect-model, or pass\n"
" --allow-model-drift to accept the caveat.",
file=sys.stderr,
)
return 6
if EVIDENCE_K and args.extract and not args.mock:
failed = [
k for k, s in scores.items() if k in KILL_LINES and s["value"] < KILL_LINES[k]
]
if failed:
print(f"\nBELOW KILL LINE: {', '.join(failed)}", file=sys.stderr)
return 1
return 0
if __name__ == "__main__":
raise SystemExit(main())