"""Full-weight cross-entropy training with periodic, fixed-fold evaluation. A simplified port of autojev's trainer: same objective, schedule, optimizer and checkpoint selection, without the streamed tranche/audit machinery. """ import argparse import copy import hashlib import json import math import os from pathlib import Path import random import shutil import subprocess import time import tomllib from collections.abc import Mapping, Sequence from typing import cast import torch import torch.nn.functional as F from kev.evaluate import ( Metrics, Prediction, by_panel, calibration_ok, evaluate_logits, fit_temperature, hard_label, label_index, metrics, options, read_predictions, read_rows, selection_key, validate_coverage, write_json, ) from kev.events import record from kev.model import BASE_MODEL, DecisionModel from kev.optim import CPUOffloadAdamW from kev.tracking import Tracker from kev.types import Example, JSONValue # Settings that must match for an exact resume; everything else may change (e.g. stop_after). RESUME_INVARIANT = ("train", "development", "temperature", "reference", "public", "base_model", "epochs", "batch_size", "effective_batch_size", "token_budget", "max_length", "lr", "weight_decay", "warmup_fraction", "min_lr_ratio", "seed", "extend_from") class Arguments(argparse.Namespace): config: str | None train: str development: str temperature: str reference: str | None public: str | None run: str base_model: str device: str | None epochs: int batch_size: int effective_batch_size: int token_budget: int max_length: int lr: float weight_decay: float warmup_fraction: float min_lr_ratio: float seed: int eval_every: int public_eval_every: int resume_every: int keep_checkpoints: int stop_after: int | None resume: bool extend_from: str | None cpu_threads: int eval_batch_size: int wandb_project: str | None wandb_mode: str eval_token_budget: int def digest(path: str | Path) -> str: with Path(path).open("rb") as stream: return hashlib.file_digest(stream, "sha256").hexdigest() def append(path: Path, value: object) -> None: with path.open("a") as stream: stream.write(json.dumps(value, ensure_ascii=False) + "\n") def length_estimate(row: Example) -> int: measured = row["source"].get("input_tokens") if isinstance(measured, int) and not isinstance(measured, bool): return measured + 16 return len(json.dumps([row["state"], row["question"]], ensure_ascii=False)) // 3 + 192 + 512 * len(row.get("images", [])) def microbatches(rows: Sequence[Example], batch_size: int, token_budget: int) -> list[list[Example]]: result: list[list[Example]] = [] pending: list[Example] = [] longest = 0 for row in rows: length = length_estimate(row) if pending and (len(pending) == batch_size or max(longest, length) * (len(pending) + 1) > token_budget): result.append(pending) pending, longest = [], 0 pending.append(row) longest = max(longest, length) if pending: result.append(pending) return result def targets(rows: Sequence[Example], device: torch.device) -> torch.Tensor: values = torch.zeros((len(rows), 255), dtype=torch.float32, device=device) for index, row in enumerate(rows): labels, target = options(row["question"]), row["target"] if isinstance(target, list): distribution = target elif row["question"]["type"] == "noul": positive = float(cast(float, target)) distribution = [1.0 - positive, positive] else: distribution = [float(label == target) for label in labels] if len(distribution) != len(labels) or any(not math.isfinite(p) or p < 0 for p in distribution) or abs(sum(distribution) - 1) > 1e-6: raise ValueError(f"Invalid training target: {row['id']}") values[index, :len(distribution)] = torch.tensor(distribution, device=device) return values def augment(rows: Sequence[Example], rng: random.Random) -> list[Example]: """Shuffle choice option order so the readout cannot learn positional priors.""" result = copy.deepcopy(list(rows)) for row in result: if row["question"]["type"] == "choice": criteria = row["question"]["criteria"] target = row["target"] weights = dict(zip(criteria, target, strict=True)) if isinstance(target, list) else None items = list(criteria.items()) rng.shuffle(items) row["question"]["criteria"] = dict(items) if weights is not None: row["target"] = [weights[key] for key, _ in items] return result def learning_rate_factor(step: int, total_steps: int, warmup_fraction: float, min_lr_ratio: float) -> float: """Linear warmup, then cosine decay to min_lr_ratio of the peak.""" warmup = max(1, int(warmup_fraction * total_steps)) if step <= warmup: return step / warmup progress = (step - warmup) / max(1, total_steps - warmup) return min_lr_ratio + (1 - min_lr_ratio) * 0.5 * (1 + math.cos(math.pi * progress)) def check_partitions(partitions: Mapping[str, Sequence[Example]]) -> None: """Reject duplicate IDs and any ID or (dataset, family) shared across folds.""" families: dict[str, set[tuple[str, str]]] = {} identifiers: dict[str, set[str]] = {} for name, rows in partitions.items(): identifiers[name] = {row["id"] for row in rows} if len(identifiers[name]) != len(rows): raise ValueError(f"Duplicate IDs in {name}") families[name] = {(str(row["source"].get("dataset", row["suite"])), row["family"]) for row in rows} for other in identifiers: if name != other and (identifiers[name] & identifiers[other] or families[name] & families[other]): raise ValueError(f"Partitions overlap: {name}/{other}") def synchronize() -> None: if torch.cuda.is_available(): torch.cuda.synchronize() @torch.inference_mode() def infer(model: DecisionModel, rows: Sequence[Example], args: Arguments) -> list[list[float]]: was_training = model.training model.eval() result: list[list[float]] = [] for batch in microbatches(rows, args.batch_size, args.token_budget): logits = model(model.prepare(batch, max_length=args.max_length)).detach().cpu() for row, values in zip(batch, logits, strict=True): result.append(cast(list[float], values[:len(options(row["question"]))].tolist())) model.train(was_training) return result def save_predictions(path: Path, predictions: Sequence[Prediction]) -> None: with path.open("w") as stream: for prediction in predictions: stream.write(json.dumps(prediction, ensure_ascii=False) + "\n") def sync_directory(path: Path) -> None: descriptor = os.open(path, os.O_RDONLY | os.O_DIRECTORY) try: os.fsync(descriptor) finally: os.close(descriptor) def select_checkpoint(root: Path, step: int) -> None: pending = root / "selected.pending" pending.unlink(missing_ok=True) pending.symlink_to(f"step-{step:05d}", target_is_directory=True) os.replace(pending, root / "selected") sync_directory(root) def prune_checkpoints(root: Path, keep: int, protected: int | None) -> None: """Each checkpoint is ~54 GB; keep the newest `keep` plus the one the saved resume state selects.""" steps = sorted((path for path in root.glob("step-*") if path.is_dir()), key=lambda path: path.name) for path in steps[:-keep]: if protected is not None and path.name == f"step-{protected:05d}": continue shutil.rmtree(path) record("checkpoint_pruned", path=str(path)) def save_selected(model: DecisionModel, root: Path, temperature: float, step: int, provenance: dict[str, str], keep: int, protected: int | None) -> None: destination = root / f"step-{step:05d}" if destination.exists(): # left over from an interrupted attempt past the resume point shutil.rmtree(destination) model.save(destination, temperature=temperature, step=step, provenance=cast(JSONValue, provenance)) for file in destination.rglob("*"): if file.is_file(): with file.open("rb") as stream: os.fsync(stream.fileno()) sync_directory(destination) select_checkpoint(root, step) prune_checkpoints(root, keep, protected) def truncate_logs(run: Path, step: int) -> None: """On resume, drop log lines written after the saved step so the history stays one trajectory.""" for name in ("training.jsonl", "evaluations.jsonl", "public-evaluations.jsonl"): path = run / name if path.exists(): lines = [line for line in path.read_text().splitlines() if json.loads(line)["step"] <= step] path.write_text("".join(line + "\n" for line in lines)) def parse_arguments(argv: Sequence[str] | None = None) -> Arguments: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--config", help="TOML file of defaults; command-line flags override it") parser.add_argument("--train", required=True) parser.add_argument("--development", required=True) parser.add_argument("--temperature", required=True, help="Calibration fold used only to fit the temperature") parser.add_argument("--reference", help="Reference (e.g. Jev) predictions on --development; gates selection on calibration") parser.add_argument("--public", help="Extra diagnostic fold, evaluated every --public-eval-every steps") parser.add_argument("--run", required=True, help="Run directory: logs, predictions, checkpoints/, resume.pt") parser.add_argument("--base-model", default=BASE_MODEL) parser.add_argument("--device") parser.add_argument("--epochs", type=int, default=1) parser.add_argument("--batch-size", type=int, default=32) parser.add_argument("--effective-batch-size", type=int, default=256) parser.add_argument("--token-budget", type=int, default=8192) parser.add_argument("--max-length", type=int, default=8192) parser.add_argument("--lr", type=float, default=2e-6) parser.add_argument("--weight-decay", type=float, default=0.01) parser.add_argument("--warmup-fraction", type=float, default=0.05) parser.add_argument("--min-lr-ratio", type=float, default=0.1) parser.add_argument("--seed", type=int, default=20260920) parser.add_argument("--eval-every", type=int, default=50) parser.add_argument("--public-eval-every", type=int, default=50) parser.add_argument("--resume-every", type=int, default=50) parser.add_argument("--keep-checkpoints", type=int, default=2) parser.add_argument("--stop-after", type=int, help="Pause after this global step (a pilot or a planned break)") parser.add_argument("--resume", action="store_true", help="Continue exactly from /resume.pt") parser.add_argument("--extend-from", help="FSDP only: continue a completed run with a new epoch schedule in a separate directory") parser.add_argument("--cpu-threads", type=int, default=32) parser.add_argument("--wandb-project", help="Log to this W&B project (off when unset)") parser.add_argument("--wandb-mode", default="online", choices=("online", "offline", "disabled")) parser.add_argument("--eval-batch-size", type=int, default=64, help="Rows per inference batch (FSDP trainer)") parser.add_argument("--eval-token-budget", type=int, default=65536, help="Padded tokens per inference batch (FSDP trainer)") # A config file supplies defaults, so its values also satisfy required flags. preliminary = argparse.ArgumentParser(add_help=False) preliminary.add_argument("--config") config_path = preliminary.parse_known_args(argv)[0].config if config_path: defaults = tomllib.loads(Path(config_path).read_text()) unknown = set(defaults) - {action.dest for action in parser._actions} if unknown: parser.error(f"Unknown config keys: {sorted(unknown)}") parser.set_defaults(**defaults) for action in parser._actions: if action.dest in defaults: action.required = False args = parser.parse_args(argv, namespace=Arguments()) if min(args.epochs, args.batch_size, args.effective_batch_size, args.token_budget, args.eval_every, args.public_eval_every, args.resume_every, args.keep_checkpoints) < 1: parser.error("Batch, epoch, interval and retention settings must be positive") if args.stop_after is not None and args.stop_after < 1: parser.error("--stop-after must be at least one step") return args def main(argv: Sequence[str] | None = None) -> None: args = parse_arguments(argv) if args.extend_from: raise ValueError("--extend-from is supported by kev.train_fsdp only") run = Path(args.run) output = run / "checkpoints" if (run / "config.json").exists() and not args.resume: raise ValueError("Run already exists; pass --resume or choose a new run directory") if args.resume and not (run / "resume.pt").exists(): raise ValueError("--resume needs an existing /resume.pt") output.mkdir(parents=True, exist_ok=True) os.environ.setdefault("KEV_EVENTS", str(run / "events.jsonl")) train = read_rows(Path(args.train)) development, temperature_rows = read_rows(Path(args.development)), read_rows(Path(args.temperature)) public_rows = read_rows(Path(args.public)) if args.public else [] if not train or not development or not temperature_rows: raise ValueError("Training, development and temperature folds must be nonempty") check_partitions({"train": train, "development": development, "temperature": temperature_rows, "public": public_rows}) reference: Metrics | None = None if args.reference: reference_predictions = read_predictions(Path(args.reference)) validate_coverage(development, reference_predictions) reference = metrics(reference_predictions) inputs = {"train": args.train, "development": args.development, "temperature": args.temperature, "reference": args.reference, "public": args.public} hashes = {name: digest(path) for name, path in inputs.items() if path} package = Path(__file__).resolve().parent code_hashes = {name: digest(package / name) for name in ("train.py", "model.py", "optim.py", "evaluate.py", "types.py")} lock = package.parents[1] / "uv.lock" if lock.exists(): code_hashes["uv.lock"] = digest(lock) try: git_commit: str | None = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=package, text=True, stderr=subprocess.DEVNULL).strip() except (subprocess.CalledProcessError, FileNotFoundError): git_commit = None config = {**vars(args), "data_sha256": hashes, "code_sha256": code_hashes, "git_commit": git_commit, "train_rows": len(train), "development_rows": len(development), "temperature_rows": len(temperature_rows), "public_rows": len(public_rows)} if not args.resume: write_json(run / "config.json", config) tracker = Tracker(run, args.wandb_project, args.wandb_mode, config) random.seed(args.seed) torch.manual_seed(args.seed) torch.cuda.manual_seed_all(args.seed) started = time.monotonic() model = DecisionModel(train=True, base_model=args.base_model, device=args.device, gradient_checkpointing=True, cpu_threads=args.cpu_threads) optimizer = CPUOffloadAdamW(model.named_parameters(), lr=args.lr, weight_decay=args.weight_decay) parameters = [parameter for parameter in model.parameters() if parameter.requires_grad] groups: list[list[Example]] = [] for epoch in range(args.epochs): ordered = list(train) random.Random(args.seed + epoch).shuffle(ordered) for offset in range(0, len(ordered), args.effective_batch_size): groups.append(sorted(ordered[offset:offset + args.effective_batch_size], key=length_estimate)) total_steps = len(groups) step, examples_seen = 0, 0 best: Metrics | None = None best_step: int | None = None resumable_best_step: int | None = None # the checkpoint resume.pt would reselect; never pruned selected_temperature: float | None = None if args.resume: state = torch.load(run / "resume.pt", map_location="cpu", weights_only=True) if state["data_sha256"] != hashes or state["total_steps"] != total_steps: raise ValueError("Resume data differs from the saved run") if state["config"]["code_sha256"] != code_hashes: raise ValueError("Training implementation differs from the saved run") for key in RESUME_INVARIANT: if state["config"][key] != vars(args)[key]: raise ValueError(f"Resume configuration differs: {key}") step, examples_seen = state["step"], state["examples_seen"] optimizer.load_state_dict(state["optimizer"]) if args.stop_after is not None and args.stop_after <= step: raise ValueError("The requested stopping step must follow the saved step") best, best_step, selected_temperature = state["best"], state["best_step"], state["selected_temperature"] resumable_best_step = best_step random.setstate(state["python_rng"]) torch.set_rng_state(state["torch_rng"]) if state["cuda_rng"]: torch.cuda.set_rng_state_all(state["cuda_rng"]) del state truncate_logs(run, step) for path in output.glob("step-*"): if path.is_dir() and int(path.name.rsplit("-", 1)[1]) > step: shutil.rmtree(path) if best_step is not None: select_checkpoint(output, best_step) record("training_started", run=run.name, git_commit=git_commit, data_sha256=hashes, initialization="exact_resume" if args.resume else "base_fresh_optimizer", step=step, total_steps=total_steps, device=model.device_name, trainable_parameters=sum(p.numel() for p in parameters)) def evaluate(include_public: bool = False) -> None: nonlocal best, best_step, selected_temperature began = time.monotonic() temperature_logits = infer(model, temperature_rows, args) fitted = fit_temperature(temperature_logits, [label_index(options(row["question"]), hard_label(row)) for row in temperature_rows]) logits = infer(model, development, args) raw_predictions, fitted_predictions = evaluate_logits(development, logits), evaluate_logits(development, logits, fitted) raw, calibrated = metrics(raw_predictions), metrics(fitted_predictions) eligible = reference is None or calibration_ok(calibrated, reference) improved = eligible and (best is None or selection_key(calibrated, step) < selection_key(best, cast(int, best_step))) if improved: save_selected(model, output, fitted, step, {"run": run.name, "git_commit": git_commit or "", **hashes, **code_hashes}, args.keep_checkpoints, resumable_best_step) best, best_step, selected_temperature = calibrated, step, fitted save_predictions(run / f"development-{step:05d}.jsonl", fitted_predictions) save_predictions(run / f"temperature-{step:05d}.jsonl", evaluate_logits(temperature_rows, temperature_logits, fitted)) panels = by_panel(development, fitted_predictions) tracker.log_evaluation(step, raw, calibrated, panels, fitted, time.monotonic() - began) value = record("evaluation", run=run.name, step=step, examples_seen=examples_seen, raw=raw, fitted=calibrated, panels=panels, temperature=fitted, reference=reference, eligible=eligible, selected_step=best_step, elapsed_seconds=time.monotonic() - started, evaluation_seconds=time.monotonic() - began) append(run / "evaluations.jsonl", value) print(json.dumps({key: value[key] for key in ("step", "temperature", "eligible", "selected_step")} | {"accuracy": calibrated["accuracy"], "ece": calibrated["ece"], "brier": calibrated["brier"]}), flush=True) if include_public and public_rows: predictions = evaluate_logits(public_rows, infer(model, public_rows, args), fitted) save_predictions(run / f"public-{step:05d}.jsonl", predictions) value = record("public_evaluation", run=run.name, step=step, temperature=fitted, metrics=metrics(predictions)) append(run / "public-evaluations.jsonl", value) def save_resume() -> None: nonlocal resumable_best_step began = time.monotonic() state = {"optimizer": optimizer.state_dict(), "step": step, "examples_seen": examples_seen, "total_steps": total_steps, "data_sha256": hashes, "config": config, "best": best, "best_step": best_step, "selected_temperature": selected_temperature, "python_rng": random.getstate(), "torch_rng": torch.get_rng_state(), "cuda_rng": torch.cuda.get_rng_state_all() if torch.cuda.is_available() else []} pending = run / "resume.pt.pending" torch.save(state, pending) with pending.open("rb") as stream: os.fsync(stream.fileno()) os.replace(pending, run / "resume.pt") sync_directory(run) resumable_best_step = best_step prune_checkpoints(output, args.keep_checkpoints, resumable_best_step) record("resume_saved", run=run.name, step=step, bytes=(run / "resume.pt").stat().st_size, seconds=time.monotonic() - began) if step == 0: evaluate(include_public=bool(public_rows)) stop = min(total_steps, args.stop_after) if args.stop_after is not None else total_steps while step < stop: rows = groups[step] began = time.monotonic() model.train() group = augment(rows, random.Random(args.seed + 100003 * (step + 1))) optimizer.zero_grad() total_loss = 0.0 input_tokens = 0 for batch in microbatches(group, args.batch_size, args.token_budget): prepared = model.prepare(batch, max_length=args.max_length) logits = model(prepared) target = targets(batch, logits.device) loss = -(target * F.log_softmax(logits, dim=-1)).sum(-1).mean() if not torch.isfinite(loss): raise RuntimeError(f"Nonfinite loss at step {step + 1}") (loss * len(batch) / len(group)).backward() total_loss += float(loss.detach()) * len(batch) input_tokens += prepared.input_tokens del logits, loss, target, prepared gradient_norm = float(torch.nn.utils.clip_grad_norm_(parameters, 1.0)) if not math.isfinite(gradient_norm): raise RuntimeError(f"Nonfinite gradient at step {step + 1}") step += 1 factor = learning_rate_factor(step, total_steps, args.warmup_fraction, args.min_lr_ratio) for group_parameters in optimizer.param_groups: group_parameters["lr"] = args.lr * factor optimizer_started = time.monotonic() optimizer.step() synchronize() examples_seen += len(group) value = record("training_step", run=run.name, step=step, loss=total_loss / len(group), learning_rate=args.lr * factor, examples_seen=examples_seen, group_examples=len(group), input_tokens=input_tokens, gradient_norm=gradient_norm, step_seconds=time.monotonic() - began, optimizer_seconds=time.monotonic() - optimizer_started, elapsed_seconds=time.monotonic() - started, gpu_peak_gb=torch.cuda.max_memory_allocated() / 1e9 if torch.cuda.is_available() else None) append(run / "training.jsonl", value) tracker.log_training(value) print(json.dumps(value), flush=True) if step % args.eval_every == 0 or step == stop: evaluate(include_public=step % args.public_eval_every == 0 or step == total_steps) if step % args.resume_every == 0 or step == stop: save_resume() summary = {"run": run.name, "steps": step, "planned_steps": total_steps, "complete": step == total_steps, "examples_seen": examples_seen, "best_step": best_step, "selected_temperature": selected_temperature, "selected_metrics": best, "reference": reference, "data_sha256": hashes, "git_commit": git_commit, "elapsed_seconds": time.monotonic() - started, "checkpoint": str(output / "selected") if best else None} write_json(run / "summary.json", summary) record("training_finished" if step == total_steps else "training_paused", **summary) tracker.summary({"best_step": best_step, "selected_temperature": selected_temperature, "steps": step, **({f"best/{key}": best[key] for key in ("accuracy", "ece", "brier", "nll")} if best else {})}) tracker.finish() print(json.dumps(summary, indent=2), flush=True) if __name__ == "__main__": main()