llm-ingestion-okf/tools/okf_corpus_run.py
Kjell Tore Guttormsen aa87eb8818 feat(inbox): walk the drop directory recursively
Door B listed `inbox.iterdir()` and kept only top-level files. A file in a
subdirectory was neither ingested nor refused: it appeared in none of the
result's buckets, so a nested drop produced a bundle that was silently short
of what was dropped and no count said so. That broke the K1b identity for any
inbox with folders in it. Operator decision 2026-09-06.

- `walk_inbox` is the ONE walk rule, shared with `tools/okf_corpus_run.py`:
  the denominator N is now counted over exactly the set of files the door
  ingests, rather than over a second listing that happened to agree.
- Sorted on the whole relative path, not the basename, so the order is a
  function of the tree; that is what keeps rebuild-from-scratch byte-equal to
  an incremental update.
- A concept's `source_file` is the path relative to the inbox root,
  `/`-separated. The concept NAME still comes from the basename, so two
  folders holding one basename hit the existing §3 collision refusal instead
  of one silently claiming the other's concept.
- Dot-directories and a bundle directory inside the inbox are skipped with a
  CODE, in a new `InboxResult.skipped`. Recursion makes the door's own output
  reachable as its own input; a silent skip would be the same
  absence-without-a-denominator defect one level down.
- `--path-prefix` reduces per component and rejoins with `/`, so the caller
  driving a nested corpus can carry the relative directory. Reducing the whole
  string folded the separator into a `-` and flattened `sub/sub2`.

`tests/test_inbox_flow.py::test_subdirectories_are_not_walked` asserted the
opposite and is superseded in place, with the reason written down.

Measured on the K2 corpus (flat, N=43): 39/43 merged, 4 coded, K1b holds. The
bundle digest is
`1472e98aec8643c5beee540f4c42b5e437bd26e7c61d69a91bcff799f06a6d13` over 1108
files -- byte-identical to a run of the same corpus at 190086f WITHOUT this
change (`diff -r` exit 0), so recursion costs a flat inbox nothing. It differs
from the stored 2026-09-03 artifact by one line in `index.md`
(`- [Corpus run history](log.md)`), which 95eb271 added 15 hours after that
bundle was built.

Suite 1113 passed, `ruff` clean, `mypy --strict src/ tools/` clean.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-07 04:11:00 +02:00

497 lines
19 KiB
Python

"""Run a corpus through the whole path and report numbers, never a claim.
The instrument behind K1 and K2. It exists because the alternative -- a person
running the door by hand and writing down what they saw -- has already produced
a number that was right about a directory that no longer existed.
Three rules it enforces rather than describes:
**K1b is a command.** The conservation identity `merged + Sigma(coded
rejections) == N` is CHECKED here, and a run where it does not hold EXITS
NON-ZERO. Asserted in prose it would be something a reader has to trust; as an
exit status it fails the run that produced it. When it fails, the unaccounted
files are NAMED -- "some file went missing" is not actionable.
**`N` is computed, never typed.** It is the file count of the corpus
directory, read at run time. A literal would keep passing after the corpus
changed and would then report a fact about a directory that no longer exists.
**Three counts, never one.** The guard sits between extraction and persist, so
a healthy persisted count can hide a pile of quarantines. Extracted, gated and
persisted are separate numbers for that reason.
**The degenerate-merge rule is a DEFINITION, not a threshold: a merge is
degenerate when the extracted text is zero characters after stripping
whitespace.** A concept with an empty body cannot carry one unit of knowledge,
so counting it as a merge would report extraction failure as success.
The resolved converter path and version are printed in the output, because the
vendored binary is bypassed silently otherwise -- measured three times, wheel
3.9 against host 3.10.2.
"""
from __future__ import annotations
import argparse
import json
import sys
import time
from collections.abc import Mapping
from dataclasses import dataclass, replace
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
from llm_ingestion_okf.errors import IngestError # noqa: E402
from llm_ingestion_okf.extract import extract_text # noqa: E402
from llm_ingestion_okf.inbox import ( # noqa: E402
GateDecision,
InboxResult,
process_inbox,
relative_source,
walk_inbox,
)
from llm_ingestion_okf.profiles import ( # noqa: E402
SEGMENTED_OKF_V0_2,
STRUCTURED_V1,
BundleProfile,
)
from llm_ingestion_okf.segmentation import ( # noqa: E402
SegmentationPlan,
parse_segmentation_plan,
)
__all__ = [
"CorpusReport",
"converter_identity",
"load_plans",
"is_degenerate",
"main",
"measure",
"replace",
"unaccounted_names",
]
HARNESS_ID = "okf-corpus-run"
# The log's name and title in ONE place, because two of them now read it: the
# file's own frontmatter and the root index entry that points at it. Two
# literals would let the link's label drift away from the thing it labels.
LOG_NAME = "log.md"
LOG_TITLE = "Corpus run history"
def is_degenerate(text: str) -> bool:
"""Zero characters after stripping whitespace. The whole rule, in one line.
A definition rather than a threshold on purpose: a threshold invites a
later argument about where it should sit, and every such argument has to be
had again the next time the corpus changes.
"""
return not text.strip()
def converter_identity() -> tuple[str, str]:
"""The converter this run would use, resolved by path, and its version.
Reported rather than assumed. `pypandoc` prefers the HIGHEST version it can
find over the one this package vendored, so a run that did not say which
binary produced its text would be unattributable.
"""
from llm_ingestion_okf._pandoc import PANDOC_VERSION, resolve_pandoc
try:
return (str(resolve_pandoc()), PANDOC_VERSION)
except IngestError as exc:
return (f"unresolved ({exc.code})", PANDOC_VERSION)
def unaccounted_names(
*, dropped: tuple[str, ...], merged: tuple[str, ...], coded: tuple[str, ...]
) -> tuple[str, ...]:
"""Every dropped file that is in neither column, in sorted order.
The conservation check, isolated so it can be driven with an inventory the
door could not produce. A harness whose failure path is unreachable is a
harness that proves nothing when it passes.
"""
return tuple(sorted(set(dropped) - set(merged) - set(coded)))
@dataclass(frozen=True)
class CorpusReport:
"""One corpus run's numbers, every one of them with its denominator."""
corpus: str
ingested_at: str
n: int
extracted: int
gated: int
persisted: int
substantive: int
degenerate: int
rejected: int
seconds_total: float
converter_path: str
converter_version: str
codes: tuple[tuple[str, int], ...]
unaccounted: tuple[str, ...]
@property
def merged(self) -> int:
return self.substantive + self.degenerate
def render(self) -> str:
per_file = self.seconds_total / self.n if self.n else 0.0
lines = [
f"# Corpus run: {self.corpus}",
"",
f"N (denominator, the directory's file count) = {self.n}",
"",
"## Three counts, never one",
"",
"The guard sits between extraction and persist, so a healthy persisted",
"count can hide a pile of quarantines.",
"",
f"- extracted: {self.extracted}/{self.n}",
f"- gated: {self.gated}/{self.n}",
f"- persisted: {self.persisted}/{self.n}",
"",
"## The numerator, split",
"",
"A merge is degenerate when the extracted text is zero characters after",
"stripping whitespace -- a definition, not a threshold.",
"",
f"- substantive: {self.substantive}/{self.n}",
f"- degenerate: {self.degenerate}/{self.n}",
f"- rejected (coded): {self.rejected}/{self.n}",
"",
f"merged + coded rejections = {self.merged + self.rejected}; N = {self.n}",
"",
"## Converter",
"",
f"- resolved converter path: {self.converter_path}",
f"- pinned converter version: {self.converter_version}",
"",
"## Wall time",
"",
f"- total: {self.seconds_total:.2f} s",
f"- per file: {per_file:.3f} s",
"",
"## Rejection codes",
"",
]
lines.extend(
f"- `{code}`: {count}/{self.n}" for code, count in self.codes or (("(none)", 0),)
)
if self.unaccounted:
lines += ["", "## UNACCOUNTED", ""]
lines.extend(f"- {name}" for name in self.unaccounted)
return "\n".join(lines) + "\n"
def render_log(self) -> str:
"""The bundle's own `log.md`, in SPEC section 9 form.
Written because a consumer measured that K1b was NOT checkable from the
bundle: `merged` is countable from the concepts, `N` is not, so the
conservation identity could only be taken on trust from a report that
does not travel with the artifact. Section 9 already reserves this file
for the history of a scope, and the denominator is the one fact about
this run that the bundle cannot otherwise recover.
Dated from `ingested_at`, never the wall clock: determinism here is
bit-exact, and a date that moved between two replays of the same corpus
would put a changing byte in an artifact that must not change.
"""
codes = self.codes or (("(none)", 0),)
rejections = ", ".join(f"`{code}`: {count}" for code, count in codes)
lines = [
"---",
"type: Log",
f"title: {LOG_TITLE}",
"---",
"",
f"# {LOG_TITLE}",
"",
f"## {self.ingested_at[:10]}",
"",
f"* **Ingested**: {self.corpus} — N = {self.n} "
f"(the corpus directory's file count, computed at run time), "
f"merged = {self.merged} ({self.substantive} substantive, "
f"{self.degenerate} degenerate), coded rejections = {self.rejected}.",
f"* **Rejected**: {rejections}.",
f"* **Conservation (K1b)**: merged + coded rejections = "
f"{self.merged} + {self.rejected} = {self.merged + self.rejected}; "
f"N = {self.n}. The run exits non-zero when these differ.",
f"* **Converter**: {self.converter_path}, version {self.converter_version}.",
]
if self.unaccounted:
lines.append("* **Unaccounted**: " + ", ".join(self.unaccounted) + " — K1b FAILED.")
return "\n".join(lines) + "\n"
def load_plans(plans_dir: Path) -> dict[str, SegmentationPlan]:
"""Every proposal artifact in a directory, keyed by filename.
The key is for the operator, never for selection: `process_inbox` matches a
plan to a drop by the source content hash, so a renamed document still finds
its plan and a plan filed under the wrong name still cannot be applied to
the wrong bytes.
A directory with no artifacts raises rather than returning an empty mapping.
An empty mapping is indistinguishable from "no plans were asked for", and
the run would then report a flat bundle as a success -- the exact silent
skip that produced a corpus with zero `adjudication` keys.
"""
files = sorted(plans_dir.glob("*.json"))
if not files:
raise IngestError(
f"no segmentation plans in {plans_dir} -- a run asked to replay plans and "
"given none would build a flat bundle and report it as a success",
code="segmentation_plan_invalid",
)
return {
path.name: parse_segmentation_plan(json.loads(path.read_text(encoding="utf-8")))
for path in files
}
def _gate(text: str) -> GateDecision:
return GateDecision(sanitized_text=text, disposition="warn")
def _split_merges(corpus: Path, result: InboxResult) -> tuple[int, int]:
"""Merged files split into substantive and degenerate, by the stated rule.
Re-extracted here rather than read back off the bundle: the rule is about
the EXTRACTED text, and a concept body has already been through the gate.
"""
substantive = 0
degenerate = 0
for item in result.persisted:
source = corpus / item.source_file
try:
text = extract_text(source.name, source.read_bytes())
except (IngestError, OSError):
continue
if is_degenerate(text):
degenerate += 1
else:
substantive += 1
return (substantive, degenerate)
def measure(
corpus: Path,
bundle: Path,
*,
ingested_at: str,
plans: Mapping[str, SegmentationPlan] | None = None,
profile: BundleProfile = STRUCTURED_V1,
root_frontmatter_values: Mapping[str, str] | None = None,
) -> CorpusReport:
"""Run the corpus through the door and count what happened.
Keyword-only with defaults, so the flat call that produced the published
K1/K2 numbers stays source-compatible and byte-identical.
"""
# ONE walk rule, imported rather than restated: the denominator has to be
# counted over exactly the set of files the door ingests, or the
# conservation identity would hold over a different N than the run did.
walked, _ = walk_inbox(corpus, exclude=bundle)
dropped = tuple(relative_source(path, corpus) for path in walked)
started = time.monotonic()
result = process_inbox(
corpus,
bundle,
ingested_at,
okf_type="reference",
gate=_gate,
profile=profile,
root_frontmatter_values=root_frontmatter_values,
segmentations=plans,
)
elapsed = time.monotonic() - started
merged_names = tuple(item.source_file for item in result.persisted)
blocked = result.quarantined + result.rejected
coded_names = tuple(item.source_file for item in result.failed) + tuple(
item.source_file for item in blocked
)
counts: dict[str, int] = {}
for failure in result.failed:
counts[failure.error.code] = counts.get(failure.error.code, 0) + 1
for item in blocked:
counts[item.disposition] = counts.get(item.disposition, 0) + 1
substantive, degenerate = _split_merges(corpus, result)
path, version = converter_identity()
return CorpusReport(
corpus=str(corpus),
ingested_at=ingested_at,
n=len(dropped),
# A file that reached the gate was extracted; the gate here persists
# everything it sees, so the two differ only when a gate refuses.
extracted=len(merged_names) + len(blocked),
gated=len(merged_names) + len(blocked),
persisted=len(merged_names),
substantive=substantive,
degenerate=degenerate,
rejected=len(coded_names),
seconds_total=elapsed,
converter_path=path,
converter_version=version,
codes=tuple(sorted(counts.items())),
unaccounted=unaccounted_names(dropped=dropped, merged=merged_names, coded=coded_names),
)
def parse_args(argv: list[str] | None) -> argparse.Namespace:
parser = argparse.ArgumentParser(
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
)
parser.add_argument("--corpus", type=Path, required=True, help="the directory to run")
parser.add_argument("--report", type=Path, required=True, help="where to write the report")
parser.add_argument("--bundle", type=Path, default=None, help="where to build the bundle")
parser.add_argument(
"--ingested-at", default="2026-09-02T00:00:00Z", help="stamped verbatim, as everywhere"
)
parser.add_argument(
"--plans-dir",
type=Path,
default=None,
help=(
"directory of per-document segmentation proposals to REPLAY. Produced by "
"tools/okf_propose_segments.py first, one per document; this harness never "
"proposes a split of its own, because the split is a judgement and the run "
"path is a deterministic replay of one"
),
)
parser.add_argument(
"--bundle-id",
default=None,
help="required with --plans-dir: what a consumer joins the bundle's concepts on",
)
parser.add_argument(
"--okf-version",
default=None,
help=(
"required with --plans-dir: the upstream OKF version this bundle declares. "
"An argument and never a constant -- the VALUE belongs to the catalog "
"(decision E1), and a literal here would claim a decision this repository "
"does not own"
),
)
return parser.parse_args(argv)
def link_log_in_root_index(bundle: Path, profile: BundleProfile) -> None:
"""Point the root index at the log, so the walk section 8 supports reaches it.
Measured on the K2 artifact: the bundle carried a conformant root `log.md`
that no index named, so a consumer entering at `index.md` never reached the
one file carrying `N`.
A LOCAL choice, not a conformance requirement, and the difference is worth
stating rather than implying. Section 9 puts `log.md` at any level and
section 8 has an index enumerate its directory's contents, but upstream's
own bundles do not link it: measured at `9a15b13`, 0 of the 24 shipped
`index.md` files name the single `log.md` in the set. Upstream therefore
shows the link is not REQUIRED -- not that it is disallowed.
It belongs to the harness and not the library. The log's content IS the
run's outcome, so it cannot exist when the indexes are projected; an index
that enumerated it off the directory would gain the link only from the
second run onward and break rebuild-equals-incremental, the property the
segmented bundle is built on. Writing it after the log instead keeps both
runs identical.
THE MEMBERSHIP TEST IS LOAD-BEARING, and measured rather than assumed: the
two reprojections do not treat this line the same way. The per-directory
one drops every managed line before re-emitting its block, so the link is
gone by the time this runs. The flat one keeps a managed line whose target
is not an owned concept -- deliberately, because claiming somebody else's
link on the strength of a regex would delete curated content -- so `log.md`
survives there. Appending unconditionally therefore doubled the entry on
the second unsegmented run. Re-writing the line only when it is absent is
idempotent under both, without either side having to know about the other.
"""
index_path = bundle / profile.index.name
if not index_path.is_file():
return
body = index_path.read_text(encoding="utf-8")
link = profile.index.render_link(LOG_TITLE, LOG_NAME) + "\n"
if link in body.splitlines(keepends=True):
return
index_path.write_text(body + link, encoding="utf-8", newline="")
def main(argv: list[str] | None = None) -> int:
args = parse_args(argv)
if not args.corpus.is_dir():
print(f"{HARNESS_ID}: FAILED - no corpus directory at {args.corpus}", file=sys.stderr)
return 2
bundle = args.bundle or args.report.parent / f"{args.corpus.name}-bundle"
# Both root values or neither, checked BEFORE anything is read or written.
# A segmented run that discovered a missing `bundle_id` half way through
# would leave a partial bundle behind, and this library refuses half-built
# bundles at every other door.
plans: dict[str, SegmentationPlan] | None = None
profile = STRUCTURED_V1
root_values: dict[str, str] | None = None
if args.plans_dir is not None:
missing = [
flag
for flag, value in (
("--bundle-id", args.bundle_id),
("--okf-version", args.okf_version),
)
if value is None
]
if missing:
print(
f"{HARNESS_ID}: FAILED - {', '.join(missing)} is required with --plans-dir; "
"a profile names a key and the caller owns its value",
file=sys.stderr,
)
return 2
try:
plans = load_plans(args.plans_dir)
except (IngestError, OSError, ValueError) as exc:
print(f"{HARNESS_ID}: FAILED - {exc}", file=sys.stderr)
return 2
profile = SEGMENTED_OKF_V0_2
root_values = {"okf_version": args.okf_version, "bundle_id": args.bundle_id}
report = measure(
args.corpus,
bundle,
ingested_at=args.ingested_at,
plans=plans,
profile=profile,
root_frontmatter_values=root_values,
)
args.report.parent.mkdir(parents=True, exist_ok=True)
args.report.write_text(report.render(), encoding="utf-8", newline="")
# Into the BUNDLE, not next to the report: section 9's `log.md` is part of
# the artifact a consumer receives, and a log that stayed behind in the
# harness's output directory would leave the bundle exactly as unverifiable
# as it was before.
bundle.mkdir(parents=True, exist_ok=True)
(bundle / LOG_NAME).write_text(report.render_log(), encoding="utf-8", newline="")
link_log_in_root_index(bundle, profile)
print(report.render())
if report.unaccounted or report.merged + report.rejected != report.n:
print(
f"{HARNESS_ID}: K1b FAILED - merged ({report.merged}) + coded rejections "
f"({report.rejected}) != N ({report.n}). Unaccounted: "
f"{', '.join(report.unaccounted) or '(none named)'}",
file=sys.stderr,
)
return 1
return 0
if __name__ == "__main__":
raise SystemExit(main())