From bb6d53a520ab014689faf8f93e10ef67e3b68cf1 Mon Sep 17 00:00:00 2001 From: Nucleic Date: Thu, 30 Jul 2026 18:53:22 -0700 Subject: [PATCH] Merge nucleic/sleek-ember-seal-uady into dev --- README.md | 41 +++ audit_data.py | 120 +-------- data/curation-review-v1.json | 3 +- data/dataset-v1-manifest.json | 2 +- data/semantic-audit-v1.json | 1 + prepare_data.py | 35 ++- review_contract.py | 493 ++++++++++++++++++++++++++++++++++ review_data.py | 161 +++++++++++ tests/test_prepare_data.py | 62 +++++ tests/test_review_contract.py | 154 +++++++++++ 10 files changed, 962 insertions(+), 110 deletions(-) create mode 100644 review_contract.py create mode 100644 review_data.py create mode 100644 tests/test_review_contract.py diff --git a/README.md b/README.md index cc4f80d..23c8613 100644 --- a/README.md +++ b/README.md @@ -19,6 +19,7 @@ The canonical generated sources are listed in `data/generation-manifest.json`. - fails for review if a high-overlap pair has conflicting labels; - keeps shipped fixtures completely outside source data; - holds every `vague-eval` record out of training; +- optionally applies a completed, versioned human-review ledger before splitting; - stratifies by primary purpose, slice, and primary language; and - verifies that the deterministic test partition still matches the versioned `data/frozen-test-v1.jsonl`. @@ -177,5 +178,45 @@ the frozen split automatically. The 18 word-trigram exclusions and the human-rev completion rule are recorded in `data/curation-review-v1.json`; the semantic report is versioned as `data/semantic-audit-v1.json`. +## Complete the human review + +The deterministic CSV currently contains 1,219 blank review rows. Check progress without +running the embedding audit again: + +```bash +ml/purpose-classifier/.venv/bin/python ml/purpose-classifier/review_data.py +``` + +Mark each row `accept`, `relabel`, or `reject`. `accept` and `reject` leave the four +`reviewed*` fields blank; `reject` requires notes. For `relabel`, blank reviewed fields +retain their generated value, `` clears a secondary purpose, and notes are required. +If a secondary purpose is added or removed, set `reviewedSlice` consistently (`mixed` +when a secondary is present). The validator rejects stale generated columns, missing or +duplicate sample rows, invalid label combinations, and partially completed rows. + +When every row has a human decision, write the versionable ledger: + +```bash +ml/purpose-classifier/.venv/bin/python ml/purpose-classifier/review_data.py --finalize +``` + +Build an isolated candidate split first: + +```bash +ml/purpose-classifier/.venv/bin/python ml/purpose-classifier/prepare_data.py \ + --human-review ml/purpose-classifier/data/human-review-v1.json \ + --output-dir ml/purpose-classifier/.artifacts/reviewed-candidate \ + --frozen-test ml/purpose-classifier/.artifacts/reviewed-frozen-candidate.jsonl \ + --manifest ml/purpose-classifier/.artifacts/reviewed-manifest-candidate.json \ + --refresh-frozen-test +``` + +Inspect the ledger, decision summary, candidate manifest, and split diff. Only then rerun +the same command with the three candidate-path overrides removed to intentionally replace +the versioned frozen dataset and manifest. + +`--regenerate` recreates a blank CSV in the current schema and is only appropriate before +review begins. + The one-time, hardware-bound energy and accelerator-residency procedure is in `ENERGY_AND_RESIDENCY.md`. diff --git a/audit_data.py b/audit_data.py index 2e8a894..1605dfe 100644 --- a/audit_data.py +++ b/audit_data.py @@ -4,12 +4,9 @@ from __future__ import annotations import argparse -import csv -import hashlib import json -import math import sys -from collections import Counter, defaultdict +from collections import Counter from dataclasses import dataclass from pathlib import Path from typing import Any, Sequence @@ -31,10 +28,15 @@ from purpose_data import ( exclude_reviewed_duplicates, load_classifiable_fixtures, load_sources, - normalize_prompt, prompt_hash, write_json, ) +from review_contract import ( + DEFAULT_SAMPLE_FRACTION, + DEFAULT_SAMPLE_SEED, + stratified_review_sample, + write_review_csv, +) from train import ( DEFAULT_MODEL, DEFAULT_MODEL_REVISION, @@ -49,10 +51,6 @@ SCRIPT_DIR = Path(__file__).resolve().parent REPOSITORY_ROOT = SCRIPT_DIR.parent.parent DEFAULT_REPORT = SCRIPT_DIR / "data" / "semantic-audit-v1.json" DEFAULT_REVIEW_CSV = SCRIPT_DIR / ".artifacts" / "human-review-v1.csv" -DEFAULT_SAMPLE_SEED = 0xA11D17 -DEFAULT_SAMPLE_FRACTION = 0.10 - - @dataclass(frozen=True) class AuditRecord: prompt: str @@ -77,103 +75,6 @@ def _relative(path: Path) -> str: return str(path.resolve()) -def _stable_rank(seed: int, record: SourceRecord) -> str: - material = f"{seed}\0{prompt_hash(record.value['prompt'])}".encode("utf-8") - return hashlib.sha256(material).hexdigest() - - -def _review_stratum(record: SourceRecord) -> tuple[str, str, str]: - language = record.value["lang"].split("-", 1)[0].casefold() - return record.value["purpose"], record.value["slice"], language - - -def stratified_review_sample( - records: Sequence[SourceRecord], - *, - fraction: float = DEFAULT_SAMPLE_FRACTION, - seed: int = DEFAULT_SAMPLE_SEED, -) -> list[SourceRecord]: - """Choose exactly round(N*fraction), apportioned across purpose/slice/language.""" - - if not records: - raise DataError("cannot sample an empty review population") - if not 0.0 < fraction <= 1.0: - raise DataError("review fraction must be in (0, 1]") - - target = round(len(records) * fraction) - groups: dict[tuple[str, str, str], list[SourceRecord]] = defaultdict(list) - for record in records: - groups[_review_stratum(record)].append(record) - - # Hamilton apportionment preserves small language/slice strata while still producing - # the exact requested global sample size. - allocations: dict[tuple[str, str, str], int] = {} - remainders: list[tuple[float, str, tuple[str, str, str]]] = [] - allocated = 0 - for key in sorted(groups): - quota = len(groups[key]) * target / len(records) - base = math.floor(quota) - allocations[key] = base - allocated += base - tie_break = hashlib.sha256(f"{seed}\0{key}".encode("utf-8")).hexdigest() - remainders.append((quota - base, tie_break, key)) - for _, _, key in sorted(remainders, reverse=True)[: target - allocated]: - allocations[key] += 1 - - selected: list[SourceRecord] = [] - for key in sorted(groups): - ordered = sorted(groups[key], key=lambda record: _stable_rank(seed, record)) - selected.extend(ordered[: allocations[key]]) - return sorted(selected, key=lambda record: _stable_rank(seed + 1, record)) - - -def write_review_csv(path: Path, records: Sequence[SourceRecord]) -> None: - path.parent.mkdir(parents=True, exist_ok=True) - with path.open("w", encoding="utf-8", newline="") as handle: - writer = csv.DictWriter( - handle, - fieldnames=[ - "promptHash", - "source", - "line", - "prompt", - "generatedPurpose", - "generatedSecondary", - "generatedMixed", - "generatedDifficulty", - "generatedSlice", - "generatedLanguage", - "reviewedPurpose", - "reviewedSecondary", - "reviewedDifficulty", - "reviewStatus", - "reviewNotes", - ], - ) - writer.writeheader() - for record in records: - value = record.value - writer.writerow( - { - "promptHash": prompt_hash(value["prompt"]), - "source": _relative(record.source), - "line": record.line, - "prompt": normalize_prompt(value["prompt"]), - "generatedPurpose": value["purpose"], - "generatedSecondary": value["secondary"] or "", - "generatedMixed": str(value["mixed"]).lower(), - "generatedDifficulty": value["difficulty"], - "generatedSlice": value["slice"], - "generatedLanguage": value["lang"], - "reviewedPurpose": "", - "reviewedSecondary": "", - "reviewedDifficulty": "", - "reviewStatus": "", - "reviewNotes": "", - } - ) - - def semantic_candidates( embeddings: np.ndarray, purposes: Sequence[str], @@ -350,7 +251,11 @@ def audit(args: argparse.Namespace) -> dict[str, Any]: fraction=args.review_fraction, seed=args.review_seed, ) - write_review_csv(args.review_csv.resolve(), sample) + write_review_csv( + args.review_csv.resolve(), + sample, + source_formatter=_relative, + ) audit_records = _audit_records( curated.records, fixtures, args.fixtures.resolve() @@ -476,6 +381,7 @@ def audit(args: argparse.Namespace) -> dict[str, Any]: "reviewedPurpose", "reviewedSecondary", "reviewedDifficulty", + "reviewedSlice", "reviewStatus", ], }, diff --git a/data/curation-review-v1.json b/data/curation-review-v1.json index f686f5f..1823bcd 100644 --- a/data/curation-review-v1.json +++ b/data/curation-review-v1.json @@ -75,10 +75,11 @@ "purpose", "secondary", "difficulty", + "slice", "review status", "notes" ], "generatedArtifact": "ml/purpose-classifier/.artifacts/human-review-v1.csv", - "completionRule": "Every sampled row must be marked accept, relabel, or reject by a human reviewer; relabel/reject decisions are applied to canonical source data before dataset-v1 can be called fully curated." + "completionRule": "Every sampled row must be marked accept, relabel, or reject by a human reviewer. review_data.py validates the exact sample and writes a versioned complete ledger; prepare_data.py applies relabel/reject decisions before splitting. The revised split and frozen evaluation must then be intentionally reviewed and versioned before the dataset can be called fully curated." } } diff --git a/data/dataset-v1-manifest.json b/data/dataset-v1-manifest.json index c8cb287..ff40fde 100644 --- a/data/dataset-v1-manifest.json +++ b/data/dataset-v1-manifest.json @@ -9,7 +9,7 @@ "nearDuplicateThreshold": 0.92, "retainedRecords": 12193, "reviewPath": "ml/purpose-classifier/data/curation-review-v1.json", - "reviewSha256": "2db763df027b545b4fd654108dcadb0717f9c9c0a63555bf8cea26b791cda370", + "reviewSha256": "625251a98bcdcde0bee074e3bba3c564af4347eb7955497a9c3e018aae00b30f", "vagueEvalPolicy": "validation/test only" }, "datasetVersion": "purpose-dataset-v1", diff --git a/data/semantic-audit-v1.json b/data/semantic-audit-v1.json index f13a20f..2c66504 100644 --- a/data/semantic-audit-v1.json +++ b/data/semantic-audit-v1.json @@ -53,6 +53,7 @@ "reviewedPurpose", "reviewedSecondary", "reviewedDifficulty", + "reviewedSlice", "reviewStatus" ], "samplePurposeCounts": { diff --git a/prepare_data.py b/prepare_data.py index 08187d0..4fac65f 100644 --- a/prepare_data.py +++ b/prepare_data.py @@ -28,6 +28,7 @@ from purpose_data import ( write_json, write_jsonl, ) +from review_contract import apply_completed_human_review SCRIPT_DIR = Path(__file__).resolve().parent @@ -121,6 +122,8 @@ def _build_manifest( seed: int, near_duplicate_threshold: float, curation_review_path: Path | None, + human_review_path: Path | None, + human_review_summary: dict[str, Any] | None, ) -> dict[str, Any]: manifest = { "schemaVersion": 1, @@ -181,6 +184,12 @@ def _build_manifest( if curation_review_path is not None: manifest["curation"]["reviewPath"] = _relative(curation_review_path) manifest["curation"]["reviewSha256"] = file_sha256(curation_review_path) + if human_review_path is not None: + manifest["curation"]["humanReview"] = { + "path": _relative(human_review_path), + "sha256": file_sha256(human_review_path), + "summary": human_review_summary, + } return manifest @@ -195,6 +204,7 @@ def prepare( seed: int, near_duplicate_threshold: float, curation_review_path: Path | None = None, + human_review_path: Path | None = None, ) -> dict[str, Any]: records = load_sources(sources) fixtures = load_classifiable_fixtures(fixtures_path) @@ -215,8 +225,18 @@ def prepare( if curation_review_path is not None else CurationResult(records=lexical_curation.records, duplicates=[]) ) + human_review_summary = None + reviewed_records = reviewed_curation.records + if human_review_path is not None: + human_review = apply_completed_human_review( + human_review_path, + reviewed_records, + dataset_version=DATASET_VERSION, + ) + reviewed_records = human_review.records + human_review_summary = human_review.summary curated = CurationResult( - records=reviewed_curation.records, + records=reviewed_records, duplicates=lexical_curation.duplicates + reviewed_curation.duplicates, ) splits = split_records( @@ -283,6 +303,8 @@ def prepare( seed=seed, near_duplicate_threshold=near_duplicate_threshold, curation_review_path=curation_review_path, + human_review_path=human_review_path, + human_review_summary=human_review_summary, ) # The frozen path can be overridden in tests or experiments. manifest["frozenEval"]["syntheticPath"] = _relative(frozen_test_path) @@ -328,6 +350,14 @@ def build_parser() -> argparse.ArgumentParser: default=DEFAULT_CURATION_REVIEW, help="completed semantic duplicate review applied after lexical deduplication", ) + parser.add_argument( + "--human-review", + type=Path, + help=( + "completed label-and-difficulty review ledger from review_data.py; " + "omitted until human review is complete" + ), + ) return parser @@ -349,6 +379,9 @@ def main(argv: Sequence[str] | None = None) -> int: seed=args.seed, near_duplicate_threshold=args.near_duplicate_threshold, curation_review_path=args.curation_review.resolve(), + human_review_path=( + args.human_review.resolve() if args.human_review is not None else None + ), ) except (DataError, OSError, UnicodeError, json.JSONDecodeError) as exc: print(f"error: {exc}", file=sys.stderr) diff --git a/review_contract.py b/review_contract.py new file mode 100644 index 0000000..c5542a1 --- /dev/null +++ b/review_contract.py @@ -0,0 +1,493 @@ +#!/usr/bin/env python3 +"""Fail-closed contract for purpose-classifier human review artifacts.""" + +from __future__ import annotations + +import csv +import hashlib +import json +import math +from collections import Counter, defaultdict +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Callable, Sequence + +from purpose_data import ( + LABEL_SET, + SLICES, + DataError, + SourceRecord, + file_sha256, + normalize_prompt, + prompt_hash, + validate_source_record, + write_json, +) + + +REVIEW_CSV_FIELDS = ( + "promptHash", + "source", + "line", + "prompt", + "generatedPurpose", + "generatedSecondary", + "generatedMixed", + "generatedDifficulty", + "generatedSlice", + "generatedLanguage", + "reviewedPurpose", + "reviewedSecondary", + "reviewedDifficulty", + "reviewedSlice", + "reviewStatus", + "reviewNotes", +) +REVIEW_STATUSES = frozenset({"accept", "relabel", "reject"}) +NONE_SENTINELS = frozenset({"", "none", "null"}) +CHANGE_FIELDS = frozenset({"purpose", "secondary", "difficulty", "slice"}) +DEFAULT_SAMPLE_SEED = 0xA11D17 +DEFAULT_SAMPLE_FRACTION = 0.10 + + +@dataclass(frozen=True) +class ReviewProgress: + records: int + accepted: int + relabeled: int + rejected: int + incomplete: int + decisions: list[dict[str, Any]] + + @property + def completed(self) -> int: + return self.records - self.incomplete + + +@dataclass(frozen=True) +class HumanReviewResult: + records: list[SourceRecord] + summary: dict[str, Any] + + +def _stable_rank(seed: int, record: SourceRecord) -> str: + material = f"{seed}\0{prompt_hash(record.value['prompt'])}".encode("utf-8") + return hashlib.sha256(material).hexdigest() + + +def _review_stratum(record: SourceRecord) -> tuple[str, str, str]: + language = record.value["lang"].split("-", 1)[0].casefold() + return record.value["purpose"], record.value["slice"], language + + +def stratified_review_sample( + records: Sequence[SourceRecord], + *, + fraction: float, + seed: int, +) -> list[SourceRecord]: + """Choose exactly round(N*fraction), apportioned by purpose/slice/language.""" + + if not records: + raise DataError("cannot sample an empty review population") + if not 0.0 < fraction <= 1.0: + raise DataError("review fraction must be in (0, 1]") + + target = round(len(records) * fraction) + groups: dict[tuple[str, str, str], list[SourceRecord]] = defaultdict(list) + for record in records: + groups[_review_stratum(record)].append(record) + + allocations: dict[tuple[str, str, str], int] = {} + remainders: list[tuple[float, str, tuple[str, str, str]]] = [] + allocated = 0 + for key in sorted(groups): + quota = len(groups[key]) * target / len(records) + base = math.floor(quota) + allocations[key] = base + allocated += base + tie_break = hashlib.sha256(f"{seed}\0{key}".encode("utf-8")).hexdigest() + remainders.append((quota - base, tie_break, key)) + for _, _, key in sorted(remainders, reverse=True)[: target - allocated]: + allocations[key] += 1 + + selected: list[SourceRecord] = [] + for key in sorted(groups): + ordered = sorted(groups[key], key=lambda record: _stable_rank(seed, record)) + selected.extend(ordered[: allocations[key]]) + return sorted(selected, key=lambda record: _stable_rank(seed + 1, record)) + + +def _review_row( + record: SourceRecord, + source_formatter: Callable[[Path], str], +) -> dict[str, str]: + value = record.value + return { + "promptHash": prompt_hash(value["prompt"]), + "source": source_formatter(record.source), + "line": str(record.line), + "prompt": normalize_prompt(value["prompt"]), + "generatedPurpose": value["purpose"], + "generatedSecondary": value["secondary"] or "", + "generatedMixed": str(value["mixed"]).lower(), + "generatedDifficulty": str(value["difficulty"]), + "generatedSlice": value["slice"], + "generatedLanguage": value["lang"], + "reviewedPurpose": "", + "reviewedSecondary": "", + "reviewedDifficulty": "", + "reviewedSlice": "", + "reviewStatus": "", + "reviewNotes": "", + } + + +def write_review_csv( + path: Path, + records: Sequence[SourceRecord], + *, + source_formatter: Callable[[Path], str] = str, +) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + with path.open("w", encoding="utf-8", newline="") as handle: + writer = csv.DictWriter(handle, fieldnames=REVIEW_CSV_FIELDS) + writer.writeheader() + for record in records: + writer.writerow(_review_row(record, source_formatter)) + + +def _parse_review_changes( + row: dict[str, str], + record: SourceRecord, + *, + row_number: int, +) -> dict[str, Any]: + value = record.value + changes: dict[str, Any] = {} + + purpose = row["reviewedPurpose"].strip() + if purpose: + if purpose not in LABEL_SET: + raise DataError(f"review CSV row {row_number}: invalid reviewed purpose") + if purpose != value["purpose"]: + changes["purpose"] = purpose + + secondary_text = row["reviewedSecondary"].strip() + secondary = value["secondary"] + if secondary_text: + if secondary_text.casefold() in NONE_SENTINELS: + secondary = None + elif secondary_text in LABEL_SET: + secondary = secondary_text + else: + raise DataError(f"review CSV row {row_number}: invalid reviewed secondary") + if secondary != value["secondary"]: + changes["secondary"] = secondary + + difficulty_text = row["reviewedDifficulty"].strip() + if difficulty_text: + try: + difficulty = float(difficulty_text) + except ValueError as exc: + raise DataError( + f"review CSV row {row_number}: invalid reviewed difficulty" + ) from exc + if not math.isfinite(difficulty) or not 0.0 <= difficulty <= 1.0: + raise DataError( + f"review CSV row {row_number}: reviewed difficulty must be in [0, 1]" + ) + if difficulty != float(value["difficulty"]): + changes["difficulty"] = difficulty + + slice_name = row["reviewedSlice"].strip() + if slice_name: + if slice_name not in SLICES: + raise DataError(f"review CSV row {row_number}: invalid reviewed slice") + if slice_name != value["slice"]: + changes["slice"] = slice_name + + candidate = dict(value) + candidate.update(changes) + candidate["mixed"] = candidate["secondary"] is not None + validate_source_record(candidate, f"review CSV row {row_number}") + return changes + + +def inspect_review_csv( + path: Path, + expected_sample: Sequence[SourceRecord], + *, + source_formatter: Callable[[Path], str] = str, +) -> ReviewProgress: + try: + with path.open("r", encoding="utf-8", newline="") as handle: + reader = csv.DictReader(handle) + if tuple(reader.fieldnames or ()) != REVIEW_CSV_FIELDS: + raise DataError( + f"{path}: review CSV fields do not match the current schema; " + "regenerate the blank artifact before reviewing" + ) + rows = list(reader) + except (OSError, UnicodeError, csv.Error) as exc: + raise DataError(f"{path}: cannot read review CSV: {exc}") from exc + + expected_by_hash = { + prompt_hash(record.value["prompt"]): record for record in expected_sample + } + if len(expected_by_hash) != len(expected_sample): + raise DataError("human-review sample contains duplicate prompt hashes") + if len(rows) != len(expected_sample): + raise DataError( + f"{path}: expected {len(expected_sample)} review rows, found {len(rows)}" + ) + + seen: set[str] = set() + decisions: list[dict[str, Any]] = [] + counts: Counter[str] = Counter() + incomplete = 0 + for row_number, row in enumerate(rows, 2): + if None in row or any(row.get(field) is None for field in REVIEW_CSV_FIELDS): + raise DataError(f"{path}:{row_number}: malformed review CSV column count") + row_hash = row["promptHash"].strip() + if row_hash in seen: + raise DataError(f"{path}:{row_number}: duplicate promptHash {row_hash}") + seen.add(row_hash) + record = expected_by_hash.get(row_hash) + if record is None: + raise DataError( + f"{path}:{row_number}: promptHash is not in the deterministic sample" + ) + + expected = _review_row(record, source_formatter) + for field in REVIEW_CSV_FIELDS[:10]: + actual = row[field] + if field == "generatedDifficulty": + try: + matches = float(actual) == float(expected[field]) + except ValueError: + matches = False + else: + matches = actual == expected[field] + if not matches: + raise DataError( + f"{path}:{row_number}: generated field {field} no longer " + "matches the curated corpus" + ) + + status = row["reviewStatus"].strip().casefold() + reviewer_values = [row[field].strip() for field in REVIEW_CSV_FIELDS[10:14]] + notes = row["reviewNotes"].strip() + if not status: + if any(reviewer_values) or notes: + raise DataError( + f"{path}:{row_number}: reviewer fields require a reviewStatus" + ) + incomplete += 1 + continue + if status not in REVIEW_STATUSES: + raise DataError(f"{path}:{row_number}: invalid reviewStatus {status!r}") + + changes = _parse_review_changes(row, record, row_number=row_number) + if status == "accept": + if any(reviewer_values): + raise DataError( + f"{path}:{row_number}: accept must leave reviewed fields blank" + ) + elif status == "reject": + if any(reviewer_values): + raise DataError( + f"{path}:{row_number}: reject must leave reviewed fields blank" + ) + if not notes: + raise DataError(f"{path}:{row_number}: reject requires reviewNotes") + else: + if not changes: + raise DataError( + f"{path}:{row_number}: relabel must change at least one field" + ) + if not notes: + raise DataError(f"{path}:{row_number}: relabel requires reviewNotes") + + decision: dict[str, Any] = { + "promptHash": row_hash, + "status": status, + } + if changes: + decision["changes"] = changes + if notes: + decision["notes"] = notes + decisions.append(decision) + counts[status] += 1 + + missing = sorted(set(expected_by_hash) - seen) + if missing: + raise DataError(f"{path}: deterministic sample rows are missing") + return ReviewProgress( + records=len(rows), + accepted=counts["accept"], + relabeled=counts["relabel"], + rejected=counts["reject"], + incomplete=incomplete, + decisions=decisions, + ) + + +def finalize_human_review( + csv_path: Path, + output_path: Path, + population: Sequence[SourceRecord], + *, + dataset_version: str, + fraction: float, + seed: int, + source_formatter: Callable[[Path], str] = str, +) -> dict[str, Any]: + sample = stratified_review_sample(population, fraction=fraction, seed=seed) + progress = inspect_review_csv( + csv_path, + sample, + source_formatter=source_formatter, + ) + if progress.incomplete: + raise DataError( + f"{csv_path}: human review is incomplete: {progress.completed}/" + f"{progress.records} rows completed" + ) + + summary = { + "accepted": progress.accepted, + "relabeled": progress.relabeled, + "rejected": progress.rejected, + } + artifact = { + "schemaVersion": 1, + "datasetVersion": dataset_version, + "status": "complete", + "populationRecords": len(population), + "sampleFraction": fraction, + "sampleRecords": len(sample), + "seed": seed, + "sourceCSVSha256": file_sha256(csv_path), + "summary": summary, + "decisions": progress.decisions, + } + if output_path.exists(): + try: + existing = json.loads(output_path.read_text(encoding="utf-8")) + except (OSError, UnicodeError, json.JSONDecodeError) as exc: + raise DataError( + f"{output_path}: cannot verify existing human-review ledger: {exc}" + ) from exc + if existing != artifact: + raise DataError( + f"{output_path}: refusing to replace a different human-review ledger" + ) + return artifact + write_json(output_path, artifact) + return artifact + + +def apply_completed_human_review( + path: Path, + population: Sequence[SourceRecord], + *, + dataset_version: str, +) -> HumanReviewResult: + try: + artifact = json.loads(path.read_text(encoding="utf-8")) + if artifact["schemaVersion"] != 1: + raise DataError(f"{path}: unsupported human-review schema") + if artifact["datasetVersion"] != dataset_version: + raise DataError(f"{path}: human-review dataset version does not match") + if artifact["status"] != "complete": + raise DataError(f"{path}: human review is not complete") + fraction = float(artifact["sampleFraction"]) + seed = int(artifact["seed"]) + decisions = artifact["decisions"] + except DataError: + raise + except (OSError, UnicodeError, json.JSONDecodeError, KeyError, TypeError, ValueError) as exc: + raise DataError(f"{path}: cannot read completed human review: {exc}") from exc + if not isinstance(decisions, list): + raise DataError(f"{path}: human-review decisions must be an array") + if artifact.get("populationRecords") != len(population): + raise DataError(f"{path}: reviewed population no longer matches the corpus") + + sample = stratified_review_sample(population, fraction=fraction, seed=seed) + expected_hashes = {prompt_hash(record.value["prompt"]) for record in sample} + if artifact.get("sampleRecords") != len(sample): + raise DataError(f"{path}: reviewed sample size no longer matches the corpus") + + by_hash: dict[str, dict[str, Any]] = {} + counts: Counter[str] = Counter() + for index, decision in enumerate(decisions, 1): + if not isinstance(decision, dict): + raise DataError(f"{path}: human-review decision {index} must be an object") + row_hash = decision.get("promptHash") + status = decision.get("status") + changes = decision.get("changes", {}) + if not isinstance(row_hash, str) or not isinstance(status, str): + raise DataError(f"{path}: invalid human-review decision {index}") + if row_hash in by_hash: + raise DataError(f"{path}: duplicate human-review promptHash") + if row_hash not in expected_hashes or status not in REVIEW_STATUSES: + raise DataError(f"{path}: invalid human-review decision {index}") + if not isinstance(changes, dict) or not set(changes) <= CHANGE_FIELDS: + raise DataError(f"{path}: invalid changes in human-review decision {index}") + if status != "relabel" and changes: + raise DataError(f"{path}: only relabel decisions may contain changes") + if status == "relabel" and not changes: + raise DataError(f"{path}: relabel decision {index} has no changes") + if "purpose" in changes and ( + not isinstance(changes["purpose"], str) + or changes["purpose"] not in LABEL_SET + ): + raise DataError(f"{path}: invalid purpose in human-review decision {index}") + if "secondary" in changes and ( + changes["secondary"] is not None + and ( + not isinstance(changes["secondary"], str) + or changes["secondary"] not in LABEL_SET + ) + ): + raise DataError(f"{path}: invalid secondary in human-review decision {index}") + if "slice" in changes and ( + not isinstance(changes["slice"], str) or changes["slice"] not in SLICES + ): + raise DataError(f"{path}: invalid slice in human-review decision {index}") + by_hash[row_hash] = decision + counts[status] += 1 + if set(by_hash) != expected_hashes: + raise DataError(f"{path}: decisions do not exactly cover the deterministic sample") + + expected_summary = { + "accepted": counts["accept"], + "relabeled": counts["relabel"], + "rejected": counts["reject"], + } + if artifact.get("summary") != expected_summary: + raise DataError(f"{path}: human-review summary does not match its decisions") + + result: list[SourceRecord] = [] + for record in population: + decision = by_hash.get(prompt_hash(record.value["prompt"])) + if decision is None or decision["status"] == "accept": + result.append(record) + continue + if decision["status"] == "reject": + continue + value = dict(record.value) + value.update(decision["changes"]) + value["mixed"] = value["secondary"] is not None + validate_source_record(value, f"{path}: {decision['promptHash']}") + result.append(SourceRecord(value=value, source=record.source, line=record.line)) + + return HumanReviewResult( + records=result, + summary={ + **expected_summary, + "sampleRecords": len(sample), + "retainedRecords": len(result), + }, + ) diff --git a/review_data.py b/review_data.py new file mode 100644 index 0000000..0e21bb3 --- /dev/null +++ b/review_data.py @@ -0,0 +1,161 @@ +#!/usr/bin/env python3 +"""Check or finalize the deterministic human label-and-difficulty review.""" + +from __future__ import annotations + +import argparse +import json +import sys +from pathlib import Path +from typing import Sequence + +from prepare_data import ( + DATASET_VERSION, + DEFAULT_CURATION_REVIEW, + DEFAULT_FIXTURES, + default_source_paths, + load_reviewed_semantic_exclusions, +) +from purpose_data import ( + DataError, + SourceRecord, + curate_records, + exclude_reviewed_duplicates, + load_classifiable_fixtures, + load_sources, +) +from review_contract import ( + DEFAULT_SAMPLE_FRACTION, + DEFAULT_SAMPLE_SEED, + finalize_human_review, + inspect_review_csv, + stratified_review_sample, + write_review_csv, +) + + +SCRIPT_DIR = Path(__file__).resolve().parent +REPOSITORY_ROOT = SCRIPT_DIR.parent.parent +DEFAULT_REVIEW_CSV = SCRIPT_DIR / ".artifacts" / "human-review-v1.csv" +DEFAULT_OUTPUT = SCRIPT_DIR / "data" / "human-review-v1.json" + + +def _relative(path: Path) -> str: + try: + return str(path.resolve().relative_to(REPOSITORY_ROOT)) + except ValueError: + return str(path.resolve()) + + +def reviewed_population(args: argparse.Namespace) -> list[SourceRecord]: + sources = ( + [path.resolve() for path in args.source] + if args.source + else default_source_paths() + ) + fixtures = load_classifiable_fixtures(args.fixtures.resolve()) + lexical = curate_records( + load_sources(sources), + fixtures, + near_duplicate_threshold=args.near_duplicate_threshold, + ) + return exclude_reviewed_duplicates( + lexical.records, + load_reviewed_semantic_exclusions(args.curation_review.resolve()), + ).records + + +def build_parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--source", action="append", type=Path) + parser.add_argument("--fixtures", type=Path, default=DEFAULT_FIXTURES) + parser.add_argument( + "--curation-review", + type=Path, + default=DEFAULT_CURATION_REVIEW, + help="completed semantic duplicate review used to reconstruct the sample", + ) + parser.add_argument("--review-csv", type=Path, default=DEFAULT_REVIEW_CSV) + parser.add_argument("--output", type=Path, default=DEFAULT_OUTPUT) + parser.add_argument("--review-fraction", type=float, default=DEFAULT_SAMPLE_FRACTION) + parser.add_argument("--review-seed", type=int, default=DEFAULT_SAMPLE_SEED) + parser.add_argument("--near-duplicate-threshold", type=float, default=0.92) + action = parser.add_mutually_exclusive_group() + action.add_argument( + "--finalize", + action="store_true", + help="write the versionable review ledger; fails unless every row is complete", + ) + action.add_argument( + "--regenerate", + action="store_true", + help="replace the blank review CSV from the deterministic current sample", + ) + return parser + + +def main(argv: Sequence[str] | None = None) -> int: + args = build_parser().parse_args(argv) + try: + population = reviewed_population(args) + sample = stratified_review_sample( + population, + fraction=args.review_fraction, + seed=args.review_seed, + ) + if args.regenerate: + if args.review_csv.exists(): + existing = inspect_review_csv( + args.review_csv.resolve(), + sample, + source_formatter=_relative, + ) + if existing.completed: + raise DataError( + f"{args.review_csv}: refusing to replace " + f"{existing.completed} completed review rows" + ) + write_review_csv( + args.review_csv.resolve(), + sample, + source_formatter=_relative, + ) + print( + f"Regenerated {len(sample)} blank human-review rows at " + f"{_relative(args.review_csv)}." + ) + return 0 + progress = inspect_review_csv( + args.review_csv.resolve(), + sample, + source_formatter=_relative, + ) + if args.finalize: + artifact = finalize_human_review( + args.review_csv.resolve(), + args.output.resolve(), + population, + dataset_version=DATASET_VERSION, + fraction=args.review_fraction, + seed=args.review_seed, + source_formatter=_relative, + ) + print( + f"Finalized {artifact['sampleRecords']} human-review decisions at " + f"{_relative(args.output)}." + ) + return 0 + except (DataError, OSError, UnicodeError, json.JSONDecodeError) as exc: + print(f"error: {exc}", file=sys.stderr) + return 1 + + print( + f"Human review: {progress.completed}/{progress.records} complete " + f"({progress.accepted} accept, {progress.relabeled} relabel, " + f"{progress.rejected} reject, {progress.incomplete} remaining)." + ) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_prepare_data.py b/tests/test_prepare_data.py index 5271c5a..9b12cff 100644 --- a/tests/test_prepare_data.py +++ b/tests/test_prepare_data.py @@ -1,3 +1,4 @@ +import csv import json import sys import tempfile @@ -10,6 +11,7 @@ sys.path.insert(0, str(MODULE_DIR)) import prepare_data import purpose_data +import review_contract def example(index: int): @@ -72,6 +74,66 @@ class PrepareIntegrationTests(unittest.TestCase): train = purpose_data.load_jsonl(output / "train.jsonl") self.assertFalse(any(row["slice"] == "vague-eval" for row in train)) + def test_completed_human_review_is_applied_and_recorded(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + source = root / "source.jsonl" + purpose_data.write_jsonl(source, (example(index) for index in range(80))) + fixtures = root / "fixtures.json" + fixtures.write_text( + json.dumps( + [{"prompt": "Plan the cache migration", "purpose": "planning"}] + ), + encoding="utf-8", + ) + source_records = purpose_data.load_sources([source]) + review_population = purpose_data.curate_records( + source_records, + purpose_data.load_classifiable_fixtures(fixtures), + ).records + review_csv = root / "review.csv" + review_contract.write_review_csv(review_csv, review_population) + with review_csv.open("r", encoding="utf-8", newline="") as handle: + rows = list(csv.DictReader(handle)) + for row in rows: + row["reviewStatus"] = "accept" + rows[0]["reviewStatus"] = "reject" + rows[0]["reviewNotes"] = "Not classifiable after human review." + with review_csv.open("w", encoding="utf-8", newline="") as handle: + writer = csv.DictWriter( + handle, + fieldnames=review_contract.REVIEW_CSV_FIELDS, + ) + writer.writeheader() + writer.writerows(rows) + human_review = root / "human-review.json" + review_contract.finalize_human_review( + review_csv, + human_review, + review_population, + dataset_version=prepare_data.DATASET_VERSION, + fraction=1.0, + seed=41, + ) + + manifest = prepare_data.prepare( + sources=[source], + fixtures_path=fixtures, + output_dir=root / "output", + frozen_test_path=root / "frozen.jsonl", + manifest_path=root / "manifest.json", + refresh_frozen_test=True, + seed=23, + near_duplicate_threshold=0.92, + human_review_path=human_review, + ) + + self.assertEqual(79, manifest["curation"]["retainedRecords"]) + self.assertEqual( + 1, + manifest["curation"]["humanReview"]["summary"]["rejected"], + ) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_review_contract.py b/tests/test_review_contract.py new file mode 100644 index 0000000..92e773b --- /dev/null +++ b/tests/test_review_contract.py @@ -0,0 +1,154 @@ +import csv +import sys +import tempfile +import unittest +from pathlib import Path + + +MODULE_DIR = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(MODULE_DIR)) + +import purpose_data +import review_contract + + +def record(index: int) -> purpose_data.SourceRecord: + return purpose_data.SourceRecord( + value={ + "prompt": f"Purpose review prompt {index}", + "purpose": purpose_data.LABELS[index % len(purpose_data.LABELS)], + "secondary": None, + "mixed": False, + "difficulty": 0.4, + "slice": "core", + "lang": "en", + }, + source=Path("source.jsonl"), + line=index + 1, + ) + + +def update_csv(path: Path, update): + with path.open("r", encoding="utf-8", newline="") as handle: + rows = list(csv.DictReader(handle)) + update(rows) + with path.open("w", encoding="utf-8", newline="") as handle: + writer = csv.DictWriter( + handle, + fieldnames=review_contract.REVIEW_CSV_FIELDS, + ) + writer.writeheader() + writer.writerows(rows) + + +class HumanReviewContractTests(unittest.TestCase): + def test_blank_review_reports_progress_without_claiming_completion(self): + population = [record(index) for index in range(10)] + with tempfile.TemporaryDirectory() as directory: + path = Path(directory) / "review.csv" + review_contract.write_review_csv(path, population) + progress = review_contract.inspect_review_csv(path, population) + + self.assertEqual(0, progress.completed) + self.assertEqual(10, progress.incomplete) + self.assertEqual([], progress.decisions) + + def test_finalize_and_apply_accept_relabel_and_reject(self): + population = [record(index) for index in range(10)] + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + csv_path = root / "review.csv" + artifact_path = root / "human-review.json" + review_contract.write_review_csv(csv_path, population) + + def complete(rows): + for row in rows: + row["reviewStatus"] = "accept" + rows[0]["reviewStatus"] = "relabel" + rows[0]["reviewedPurpose"] = "writing" + rows[0]["reviewNotes"] = "Primary intent is prose." + rows[1]["reviewStatus"] = "reject" + rows[1]["reviewNotes"] = "Prompt is not classifiable." + rows[2]["reviewStatus"] = "relabel" + rows[2]["reviewedSecondary"] = "review" + rows[2]["reviewedSlice"] = "mixed" + rows[2]["reviewNotes"] = "Two explicit intents." + + update_csv(csv_path, complete) + artifact = review_contract.finalize_human_review( + csv_path, + artifact_path, + population, + dataset_version="test-v1", + fraction=1.0, + seed=7, + ) + self.assertEqual( + artifact, + review_contract.finalize_human_review( + csv_path, + artifact_path, + population, + dataset_version="test-v1", + fraction=1.0, + seed=7, + ), + ) + result = review_contract.apply_completed_human_review( + artifact_path, + population, + dataset_version="test-v1", + ) + + self.assertEqual( + {"accepted": 7, "relabeled": 2, "rejected": 1}, + artifact["summary"], + ) + self.assertEqual(9, len(result.records)) + by_prompt = {item.value["prompt"]: item.value for item in result.records} + self.assertEqual("writing", by_prompt["Purpose review prompt 0"]["purpose"]) + self.assertNotIn("Purpose review prompt 1", by_prompt) + self.assertEqual("review", by_prompt["Purpose review prompt 2"]["secondary"]) + self.assertTrue(by_prompt["Purpose review prompt 2"]["mixed"]) + self.assertEqual("mixed", by_prompt["Purpose review prompt 2"]["slice"]) + + def test_finalize_fails_closed_when_rows_are_incomplete(self): + population = [record(index) for index in range(5)] + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + csv_path = root / "review.csv" + review_contract.write_review_csv(csv_path, population) + with self.assertRaisesRegex( + purpose_data.DataError, + "human review is incomplete: 0/5", + ): + review_contract.finalize_human_review( + csv_path, + root / "human-review.json", + population, + dataset_version="test-v1", + fraction=1.0, + seed=7, + ) + + def test_relabel_must_preserve_the_source_contract(self): + population = [record(0)] + with tempfile.TemporaryDirectory() as directory: + path = Path(directory) / "review.csv" + review_contract.write_review_csv(path, population) + + def invalidate(rows): + rows[0]["reviewStatus"] = "relabel" + rows[0]["reviewedSecondary"] = "review" + rows[0]["reviewNotes"] = "Two intents." + + update_csv(path, invalidate) + with self.assertRaisesRegex( + purpose_data.DataError, + "mixed slice and mixed field disagree", + ): + review_contract.inspect_review_csv(path, population) + + +if __name__ == "__main__": + unittest.main()