Files
geointel/scripts/run_belgium_building_training_loop.py
T

340 lines
13 KiB
Python

#!/usr/bin/env python3
"""Run checkpointed CUDA train/evaluate iterations until gates pass or a batch yields."""
from __future__ import annotations
import argparse
import hashlib
import json
import shutil
import subprocess
import sys
from datetime import UTC, datetime
from pathlib import Path
from typing import Any
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 write_json(path: Path, value: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
temporary = path.with_suffix(path.suffix + ".tmp")
temporary.write_text(json.dumps(value, indent=2), encoding="utf-8")
temporary.replace(path)
def select_calibration_threshold(report: dict[str, Any]) -> dict[str, Any]:
"""Choose a threshold without consulting test or background evidence."""
eligible = [item for item in report["sweeps"] if item["pure_empty_false_positives"] == 0]
if not eligible:
eligible = report["sweeps"]
return max(
eligible,
key=lambda item: (
min(region["f1"] for region in item["regions"].values()),
item["aggregate"]["f1"],
-item["pure_empty_false_positives"],
),
)
def calibration_failures(
chosen: dict[str, Any],
*,
min_aggregate_f1: float,
min_region_f1: float,
min_region_precision: float,
min_region_recall: float,
max_pure_empty_fp: int,
) -> list[str]:
failures: list[str] = []
if chosen["aggregate"]["f1"] < min_aggregate_f1:
failures.append("calibration_aggregate_f1_below_gate")
for region, values in chosen["regions"].items():
if values["f1"] < min_region_f1:
failures.append(f"calibration_{region}_f1_below_gate")
if values["precision"] < min_region_precision:
failures.append(f"calibration_{region}_precision_below_gate")
if values["recall"] < min_region_recall:
failures.append(f"calibration_{region}_recall_below_gate")
if chosen["pure_empty_false_positives"] > max_pure_empty_fp:
failures.append("calibration_pure_empty_false_positive_gate_failed")
return failures
def training_command(
yolo: str,
*,
model: Path,
data: Path,
project: Path,
name: str,
epochs: int,
seed: int,
batch: int,
workers: int,
max_det: int = 1000,
imgsz: int = 640,
optimizer: str = "auto",
lr0: float | None = None,
mosaic: float = 1.0,
scale: float = 0.5,
translate: float = 0.1,
) -> list[str]:
command = [
yolo,
"train",
f"model={model}",
f"data={data}",
f"epochs={epochs}",
f"imgsz={imgsz}",
f"batch={batch}",
"device=0",
f"workers={workers}",
"patience=35",
"cache=disk",
"close_mosaic=20",
f"max_det={max_det}",
f"optimizer={optimizer}",
f"mosaic={mosaic}",
f"scale={scale}",
f"translate={translate}",
f"seed={seed}",
"deterministic=True",
f"project={project}",
f"name={name}",
"exist_ok=True",
]
if lr0 is not None:
command.append(f"lr0={lr0}")
return command
def run(command: list[str], log_path: Path | None = None, *, allowed: set[int] = {0}) -> int:
if log_path:
log_path.parent.mkdir(parents=True, exist_ok=True)
with log_path.open("a", encoding="utf-8") as log:
completed = subprocess.run(command, stdout=log, stderr=subprocess.STDOUT, check=False)
else:
completed = subprocess.run(command, check=False)
if completed.returncode not in allowed:
raise RuntimeError(f"Command failed ({completed.returncode}): {' '.join(command)}")
return completed.returncode
def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument("--initial-model", type=Path, required=True)
parser.add_argument("--train-yaml", type=Path, required=True)
parser.add_argument("--dataset-audit", type=Path, required=True)
parser.add_argument("--calibration-summary", type=Path, required=True)
parser.add_argument("--test-summary", type=Path, required=True)
parser.add_argument("--background-summary", type=Path, required=True)
parser.add_argument("--corpus-manifest", type=Path, required=True)
parser.add_argument("--output-dir", type=Path, required=True)
parser.add_argument("--iterations", type=int, default=1)
parser.add_argument("--epochs", type=int, default=160)
parser.add_argument("--batch", type=int, default=2)
parser.add_argument("--workers", type=int, default=4)
parser.add_argument("--max-det", type=int, default=1000)
parser.add_argument("--imgsz", type=int, default=640)
parser.add_argument("--optimizer", default="auto")
parser.add_argument("--lr0", type=float)
parser.add_argument("--mosaic", type=float, default=1.0)
parser.add_argument("--scale", type=float, default=0.5)
parser.add_argument("--translate", type=float, default=0.1)
parser.add_argument("--seed", type=int, default=20260731)
parser.add_argument("--yolo", default="yolo")
parser.add_argument("--min-aggregate-f1", type=float, default=0.55)
parser.add_argument("--min-region-f1", type=float, default=0.45)
parser.add_argument("--min-region-precision", type=float, default=0.5)
parser.add_argument("--min-region-recall", type=float, default=0.4)
parser.add_argument("--max-pure-empty-fp", type=int, default=0)
parser.add_argument("--dry-run", action="store_true")
args = parser.parse_args()
if args.iterations < 1:
raise SystemExit("--iterations must be positive")
dataset_audit = json.loads(args.dataset_audit.read_text(encoding="utf-8"))
if dataset_audit.get("status") != "ok":
raise SystemExit(f"Dataset audit is not ok: {args.dataset_audit}")
if int(dataset_audit.get("low_variance_positive_tile_count") or 0) != 0:
raise SystemExit("Dataset audit contains blank/low-variance positive tiles")
state_path = args.output_dir / "training-loop-state.json"
state: dict[str, Any] = {
"schema_version": 1,
"status": "running",
"started_at": datetime.now(UTC).isoformat(),
"initial_model": str(args.initial_model),
"train_yaml": str(args.train_yaml),
"dataset_audit": str(args.dataset_audit),
"corpus_manifest": str(args.corpus_manifest),
"iterations": [],
}
if state_path.is_file():
state = json.loads(state_path.read_text(encoding="utf-8"))
state["status"] = "running"
model = Path(state.get("next_model") or args.initial_model)
first_index = len(state["iterations"]) + 1
scripts_dir = Path(__file__).resolve().parent
for offset in range(args.iterations):
index = first_index + offset
name = f"iteration-{index:03d}"
iteration_dir = args.output_dir / name
train_run = args.output_dir / "runs" / name
command = training_command(
args.yolo,
model=model,
data=args.train_yaml,
project=args.output_dir / "runs",
name=name,
epochs=args.epochs,
seed=args.seed + index,
batch=args.batch,
workers=args.workers,
max_det=args.max_det,
imgsz=args.imgsz,
optimizer=args.optimizer,
lr0=args.lr0,
mosaic=args.mosaic,
scale=args.scale,
translate=args.translate,
)
if args.dry_run:
print(json.dumps({"training_command": command}, indent=2))
return 0
run(command, iteration_dir / "training.log")
best = train_run / "weights" / "best.pt"
if not best.is_file():
raise RuntimeError(f"Training produced no best checkpoint: {best}")
candidate = iteration_dir / "candidate.pt"
shutil.copy2(best, candidate)
reports: dict[str, Path] = {}
for role, summary in (("calibration", args.calibration_summary),):
report = iteration_dir / f"{role}.json"
reports[role] = report
run(
[
sys.executable,
str(scripts_dir / "evaluate_belgium_building_candidate.py"),
"--model",
str(candidate),
"--summary",
str(summary),
"--corpus-manifest",
str(args.corpus_manifest),
"--output",
str(report),
"--device",
"cuda:0",
"--max-det",
str(args.max_det),
"--imgsz",
str(args.imgsz),
],
iteration_dir / f"{role}.log",
)
assessment = iteration_dir / "assessment.json"
calibration = json.loads(reports["calibration"].read_text(encoding="utf-8"))
chosen = select_calibration_threshold(calibration)
failures = calibration_failures(
chosen,
min_aggregate_f1=args.min_aggregate_f1,
min_region_f1=args.min_region_f1,
min_region_precision=args.min_region_precision,
min_region_recall=args.min_region_recall,
max_pure_empty_fp=args.max_pure_empty_fp,
)
if failures:
write_json(
assessment,
{
"schema_version": 1,
"status": "continue_training_loop",
"phase": "calibration_rejected",
"threshold_selection_source": "calibration_only",
"selected_threshold": chosen["threshold"],
"gates": {
"min_aggregate_f1": args.min_aggregate_f1,
"min_region_f1": args.min_region_f1,
"min_region_precision": args.min_region_precision,
"min_region_recall": args.min_region_recall,
"max_pure_empty_false_positives": args.max_pure_empty_fp,
},
"calibration": chosen,
"test": None,
"background": None,
"failures": failures,
},
)
else:
for role, summary in (
("test", args.test_summary),
("background", args.background_summary),
):
report = iteration_dir / f"{role}.json"
reports[role] = report
run(
[
sys.executable,
str(scripts_dir / "evaluate_belgium_building_candidate.py"),
"--model", str(candidate),
"--summary", str(summary),
"--corpus-manifest", str(args.corpus_manifest),
"--output", str(report),
"--device", "cuda:0",
"--max-det", str(args.max_det),
"--imgsz", str(args.imgsz),
],
iteration_dir / f"{role}.log",
)
run(
[
sys.executable,
str(scripts_dir / "assess_belgium_building_training_iteration.py"),
"--calibration", str(reports["calibration"]),
"--test", str(reports["test"]),
"--background", str(reports["background"]),
"--output", str(assessment),
],
iteration_dir / "assessment.log",
allowed={0, 2},
)
decision = json.loads(assessment.read_text(encoding="utf-8"))
record = {
"iteration": index,
"candidate": str(candidate),
"candidate_sha256": sha256(candidate),
"assessment": str(assessment),
"status": decision["status"],
"failures": decision["failures"],
}
state["iterations"].append(record)
state["next_model"] = str(candidate)
if decision["status"] == "training_complete":
state["status"] = "training_complete"
state["completed_at"] = datetime.now(UTC).isoformat()
write_json(state_path, state)
print(json.dumps(state, indent=2))
return 0
model = candidate
write_json(state_path, state)
state["status"] = "continue_training_loop"
state["yielded_at"] = datetime.now(UTC).isoformat()
write_json(state_path, state)
print(json.dumps(state, indent=2))
return 2
if __name__ == "__main__":
raise SystemExit(main())