#!/usr/bin/env python3 """Audit a frozen Belgian corpus and emit a deterministic review queue.""" from __future__ import annotations import argparse import hashlib import json import sys from collections import Counter from datetime import datetime from pathlib import Path from typing import Any SCRIPT_DIR = Path(__file__).resolve().parent if str(SCRIPT_DIR) not in sys.path: sys.path.insert(0, str(SCRIPT_DIR)) from training_dataset_eligibility import frozen_manifest_training_eligibility_failures # noqa: E402 REQUIRED_SPLITS = ("train", "val", "calibration", "test", "background-test") REGIONS = ("flanders", "wallonia", "brussels") SHA256_LENGTH = 64 def sha256(path: Path) -> str: digest = hashlib.sha256() with path.open("rb") as stream: for chunk in iter(lambda: stream.read(1024 * 1024), b""): digest.update(chunk) return digest.hexdigest() def valid_review_timestamp(value: Any) -> bool: if not isinstance(value, str) or not value.strip(): return False try: parsed = datetime.fromisoformat(value.strip().replace("Z", "+00:00")) except ValueError: return False return parsed.tzinfo is not None def valid_sha256(value: Any) -> bool: return ( isinstance(value, str) and len(value) == SHA256_LENGTH and all(character in "0123456789abcdefABCDEF" for character in value) ) def main() -> int: parser = argparse.ArgumentParser() parser.add_argument("--corpus-dir", type=Path, required=True) parser.add_argument("--output-dir", type=Path, required=True) parser.add_argument("--review-decisions", type=Path) args = parser.parse_args() manifest_path = args.corpus_dir / "operator_samples_manifest.json" manifest = json.loads(manifest_path.read_text(encoding="utf-8")) leakage = json.loads((args.corpus_dir / "spatial-leakage-audit.json").read_text(encoding="utf-8")) manifest_policy = manifest.get("training_eligibility") fixture_mode = bool(manifest_policy.get("fixture_mode")) if isinstance(manifest_policy, dict) else False eligibility_failures = frozen_manifest_training_eligibility_failures( manifest_path, fixture_mode=fixture_mode, ) samples = manifest["samples"] split_counts = Counter((sample["region"], sample["split"]) for sample in samples) decision_counts: Counter[str] = Counter() review_queue: list[dict[str, Any]] = [] total_input = 0 total_accepted = 0 temporal_unknown = 0 failures: list[str] = list(eligibility_failures) for region in REGIONS: for split in REQUIRED_SPLITS: minimum = 4 if split == "train" else 2 if split_counts[(region, split)] < minimum: failures.append(f"{region}/{split} has {split_counts[(region, split)]}, requires {minimum}") for sample in samples: audit_path = args.corpus_dir / "pairs" / sample["sample_slug"] / "label-audit.json" audit = json.loads(audit_path.read_text(encoding="utf-8")) total_input += int(audit["input_feature_count"]) total_accepted += int(audit["accepted_feature_count"]) decision_counts.update(audit["decision_counts"]) if audit.get("temporal_alignment_status") == "unknown": temporal_unknown += 1 expected_empty = bool(sample.get("sample_role") == "background_candidate" and sample.get("require_empty")) if expected_empty and audit["accepted_feature_count"] != 0: failures.append(f"{sample['sample_slug']} is not pure empty after normalization") priority = "high" if audit["accepted_feature_count"] >= 200 or sample["split"] in {"test", "background-test"} else "normal" review_queue.append( { "sample_slug": sample["sample_slug"], "region": sample["region"], "context": sample.get("context"), "split": sample["split"], "accepted_feature_count": audit["accepted_feature_count"], "temporal_alignment_status": audit.get("temporal_alignment_status"), "priority": priority, "decision": "pending_human_review", } ) if leakage.get("status") != "ok": failures.append("spatial leakage audit failed") reviewed = 0 review_complete = False review_evidence_failures: list[str] = [] accepted_review_evidence: dict[str, dict[str, Any]] = {} review_decisions_path: Path | None = None if args.review_decisions and args.review_decisions.is_file(): review_decisions_path = args.review_decisions.resolve(strict=False) decisions = json.loads(review_decisions_path.read_text(encoding="utf-8")) by_slug = {item["sample_slug"]: item for item in decisions.get("decisions", [])} for item in review_queue: decision = by_slug.get(item["sample_slug"]) if decision: item["decision"] = decision.get("decision") item["reviewer"] = decision.get("reviewer") item["notes"] = decision.get("notes") item["reviewed_at"] = decision.get("reviewed_at") item["reviewed_artifact_path"] = decision.get("reviewed_artifact_path") item["reviewed_artifact_sha256"] = decision.get("reviewed_artifact_sha256") reviewed_artifact = ( Path(item["reviewed_artifact_path"]).expanduser().resolve(strict=False) if isinstance(item.get("reviewed_artifact_path"), str) and item["reviewed_artifact_path"].strip() else None ) has_evidence = ( item["decision"] == "accepted" and isinstance(item.get("reviewer"), str) and bool(item["reviewer"].strip()) and valid_review_timestamp(item.get("reviewed_at")) and isinstance(item.get("reviewed_artifact_path"), str) and bool(item["reviewed_artifact_path"].strip()) and valid_sha256(item.get("reviewed_artifact_sha256")) and reviewed_artifact is not None and reviewed_artifact.is_file() and item["reviewed_artifact_sha256"] == sha256(reviewed_artifact) ) if has_evidence: reviewed += 1 accepted_review_evidence[str(item["sample_slug"])] = { "reviewer": item["reviewer"].strip(), "reviewed_at": item["reviewed_at"], "reviewed_artifact_path": item["reviewed_artifact_path"], "reviewed_artifact_sha256": item["reviewed_artifact_sha256"], } elif item["decision"] in {"accepted", "rejected"}: review_evidence_failures.append( f"{item['sample_slug']}:accepted review lacks reviewer/timestamp/artifact evidence" ) review_complete = reviewed == len(review_queue) and all(item["decision"] == "accepted" for item in review_queue) human_review_evidence: dict[str, Any] | None = None if review_decisions_path is not None: human_review_evidence = { "review_decisions_path": str(review_decisions_path), "review_decisions_sha256": sha256(review_decisions_path), "required_sample_count": len(review_queue), "accepted_sample_count": reviewed, "accepted_sample_slugs": sorted(accepted_review_evidence), "reviewer_ids": sorted( {item["reviewer"] for item in accepted_review_evidence.values()} ), "reviewed_at_by_sample": { slug: accepted_review_evidence[slug]["reviewed_at"] for slug in sorted(accepted_review_evidence) }, "reviewed_artifact_path_by_sample": { slug: accepted_review_evidence[slug]["reviewed_artifact_path"] for slug in sorted(accepted_review_evidence) }, "reviewed_artifact_sha256_by_sample": { slug: accepted_review_evidence[slug]["reviewed_artifact_sha256"] for slug in sorted(accepted_review_evidence) }, } status = "failed" if failures else ("ok" if review_complete else "needs_human_review") report = { "status": status, "dataset_version": manifest["dataset_version"], "corpus_manifest_path": str(manifest_path.resolve(strict=False)), "corpus_manifest_sha256": sha256(manifest_path), "manifest_immutable": manifest["immutable"], "sample_count": len(samples), "split_counts": {f"{region}/{split}": split_counts[(region, split)] for region in REGIONS for split in REQUIRED_SPLITS}, "input_feature_count": total_input, "accepted_feature_count": total_accepted, "decision_counts": dict(sorted(decision_counts.items())), "temporal_unknown_sample_count": temporal_unknown, "spatial_leakage_status": leakage.get("status"), "training_eligibility_status": manifest_policy.get("status") if isinstance(manifest_policy, dict) else None, "training_eligibility_fixture_mode": fixture_mode, "training_eligibility_failures": eligibility_failures, "reviewed_sample_count": reviewed, "review_complete": review_complete, "human_review_evidence": human_review_evidence, "review_evidence_failures": sorted(set(review_evidence_failures)), "failures": failures, "review_queue": review_queue, } args.output_dir.mkdir(parents=True, exist_ok=True) (args.output_dir / "belgium-building-corpus-audit.json").write_text( json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8" ) lines = [ f"# Belgian building corpus audit: {manifest['dataset_version']}", "", f"Status: `{status}`", f"Samples: {len(samples)}; accepted labels: {total_accepted}/{total_input}.", f"Spatial leakage: `{leakage.get('status')}`; human reviewed: {reviewed}/{len(samples)}.", "", "## Review queue", "", "| Sample | Region | Context | Split | Labels | Priority | Decision |", "| --- | --- | --- | --- | ---: | --- | --- |", ] lines.extend( f"| {item['sample_slug']} | {item['region']} | {item['context']} | {item['split']} | " f"{item['accepted_feature_count']} | {item['priority']} | {item['decision']} |" for item in review_queue ) (args.output_dir / "belgium-building-corpus-audit.md").write_text("\n".join(lines) + "\n", encoding="utf-8") print(json.dumps({key: value for key, value in report.items() if key != "review_queue"}, ensure_ascii=False, indent=2)) return 1 if failures else 0 if __name__ == "__main__": raise SystemExit(main())