"""Champion loop: fixed held-out test per split generation, Pareto-dominance champions.

Objective: minimise Fail->Pass and Pass->Fail counts on the same working test set.
A checkpoint is the new champion when it Pareto-dominates the current one. When the
working test set looks saturated, the champion is scored once more on the unused
reserve set; if it does not hold up, a new generation (new split, fresh training)
starts. New champions are handed to a separate package command, one at a time.

Every held-out evaluation also keeps its per-song gate margins (`*.scores.jsonl`)
and records the threshold-free `gate_auc` and same-set `gate_curve` in its ledger
row; neither takes part in the champion, dominance or reserve logic.

A student candidate (its overrides hold a `gate_soft_targets` path) waits until that
teacher file exists. With `--teacher-after-round N` the loop builds each generation's
in-sample teacher itself: once every ordinary run of rounds <= N is final, one job
scores the development songs with each trained run's selected checkpoint and averages
the gate margins into the teacher path. Right after a successful build, one more job
scores the ensemble of those same members (mean gate margin, `ensemble_evaluate`) on
the working test set; its eval row competes like any checkpoint and may become the
champion, in which case the package command receives `--manifest`.
"""

import argparse
import copy
import hashlib
import json
import math
import os
import shlex
import signal
import subprocess
import sys
import threading
from collections.abc import Callable, Iterator, Sequence
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Protocol

import yaml
from scripts.resample_loop import (
    NO_CHECKPOINT,
    PIDS_LIMIT,
    PYTHON,
    REPO,
    Candidate,
    Runner,
    SubprocessRunner,
    gpu_free_gib,
    load_candidates,
    move_aside,
    pids_current,
    stamp,
    validation_metrics,
)
from scripts.scheduling_resources import (
    AdmissionPool,
    CompletionSlots,
    Reservation,
    ResourceEstimate,
    read_resource_snapshot,
)
from scripts.training_pool import TrainingWorkerPool, WorkerContext
from scripts.training_workers import FiniteTrainingRunner

from raina_laya.features.grade_model.client import gate_auc

TRAIN_SEED_BASE = 20260928
VALIDATION_FRACTION = "0.10"
HELD_OUT_FRACTION = "0.10"
TEACHER_PATH = "artifacts/teachers/gen-{generation}/insample.jsonl"
MANIFEST_NAME = "teacher-members.json"
# From this round on, every finished round also scores the ensemble of all
# validation-selected checkpoints of the generation so far (members by validation only).
ROUND_ENSEMBLE_FROM = 3
# Member cap of a round ensemble: AUC saturated at 24 members (gen 14: 24/32/42/52
# members gave 0.9234-0.9236) while scoring cost grows with every round. Members
# are ranked by validation only: the selected checkpoint's validation gate AUC from
# its validation-epoch-NNNN.scores.jsonl (ranking by validation Fail->Pass count
# instead kept only Fail-leaning members: 24 of them gave 0.9207). Runs without
# that file fall back to the validation Fail->Pass count.
ROUND_ENSEMBLE_MAX = 24
ENSEMBLE_STATUS = "ensemble"
# Strict: a champion picked on the test set regresses on reserve by selection alone.
ALPHA = 0.01
SUMMARY_KEYS = (
    "run_id",
    "candidate",
    "epoch",
    "checkpoint_name",
    "checkpoint",
    "config_path",
    "eval_json",
    "result",
)

Row = dict[str, Any]


class Handle(Protocol):
    """A launched background process."""

    pid: int

    def poll(self) -> int | None:
        """Return the exit code, or None while it is still running."""
        ...


Launcher = Callable[[Sequence[str], Path], Handle]


def is_ensemble(run: Row) -> bool:
    """Whether a run row is a synthetic ensemble (its config path is the manifest)."""
    return run["status"] == ENSEMBLE_STATUS


def counts(result: Row) -> tuple[int, int]:
    """Return (Fail->Pass count, Pass->Fail count) of one held-out result."""
    return result["fail_to_pass_count"], result["pass_to_fail_count"]


def dominates(new: Row, old: Row) -> bool:
    """Both error counts <= and at least one < (same test set, equal denominators)."""
    (new_f2p, new_p2f), (old_f2p, old_p2f) = counts(new), counts(old)
    return (
        new_f2p <= old_f2p
        and new_p2f <= old_p2f
        and (new_f2p < old_f2p or new_p2f < old_p2f)
    )


def front_champion(front: list[Row], key: str) -> Row:
    """The front member with the fewest Fail->Pass, then Pass->Fail, then rejections.

    Dominance alone leaves the champion to evaluation order (the first point of the
    front is never dominated by its neighbours), so under band decisions the
    champion is this lexicographic best of the front, the same order the reject
    band itself is chosen by. Ties fall to the earliest evaluation.
    """
    return min(
        front,
        key=lambda r: (
            r[key]["fail_to_pass_count"],
            r[key]["pass_to_fail_count"],
            r[key].get("rejected_total", 0),
            r["finished"],
        ),
    )


def pareto_front(evals: list[Row], key: str = "result") -> list[Row]:
    """Eval rows that no other eval row dominates (on `result` or `decided` counts)."""
    return [
        row
        for row in evals
        if not any(dominates(other[key], row[key]) for other in evals)
    ]


def incremental_pareto_front(
    front: list[Row], row: Row, key: str = "result"
) -> list[Row]:
    """Append one point to a frontier, preserving tied points and ledger order."""
    if any(dominates(old[key], row[key]) for old in front):
        return front
    return [old for old in front if not dominates(row[key], old[key])] + [row]


def heldout_result(report: Row) -> Row:
    """Reduce an evaluator report to the counts and rates the loop compares."""
    pass_count, correct = report["pass_count"], report["pass_correct_count"]
    f2p, fail_count = report["fail_to_pass_count"], report["fail_count"]
    p2f = pass_count - correct
    return {
        "fail_to_pass_count": f2p,
        "fail_count": fail_count,
        "fail_to_pass_rate": f2p / fail_count,
        "pass_to_fail_count": p2f,
        "pass_correct_count": correct,
        "pass_count": pass_count,
        "pass_to_fail_rate": p2f / pass_count,
        "macro_f1": report["macro_f1"],
        "s_f1": report["s_f1"],
    }


def worse_z_test(
    base_hits: int, base_total: int, new_hits: int, new_total: int
) -> tuple[float, float]:
    """One-sided two-proportion z-test that the new rate is higher than the base."""
    pooled = (base_hits + new_hits) / (base_total + new_total)
    error = math.sqrt(pooled * (1 - pooled) * (1 / base_total + 1 / new_total))
    if error == 0:
        return 0.0, 1.0
    z = (new_hits / new_total - base_hits / base_total) / error
    return z, 0.5 * math.erfc(z / math.sqrt(2))


def reserve_verdict(test: Row, reserve: Row) -> Row:
    """Hold-up test: fails when either error rate is significantly worse on reserve."""
    f2p_z, f2p_p = worse_z_test(
        test["fail_to_pass_count"],
        test["fail_count"],
        reserve["fail_to_pass_count"],
        reserve["fail_count"],
    )
    p2f_z, p2f_p = worse_z_test(
        test["pass_to_fail_count"],
        test["pass_count"],
        reserve["pass_to_fail_count"],
        reserve["pass_count"],
    )
    held_up = f2p_p >= ALPHA and p2f_p >= ALPHA
    return {
        "f2p_z": f2p_z,
        "f2p_p": f2p_p,
        "p2f_z": p2f_z,
        "p2f_p": p2f_p,
        "verdict": "pass" if held_up else "fail",
    }


def collect_checkpoints(run_dir: Path) -> list[str]:
    """Checkpoint file names to score: every Pareto entry plus the selected one."""
    names: list[str] = []
    pareto = run_dir / "pareto.json"
    if pareto.exists():
        data = json.loads(pareto.read_text(encoding="utf-8"))
        names.extend(item["checkpoint"] for item in data["checkpoints"])
    selection = json.loads((run_dir / "selection.json").read_text(encoding="utf-8"))
    names.append(selection["checkpoint"])
    return list(dict.fromkeys(names))


def summary(row: Row | None) -> Row | None:
    """Compact reference to one eval row for state.json and event rows."""
    return None if row is None else {key: row[key] for key in SUMMARY_KEYS}


def run_id_for(generation: int, round_no: int, name: str) -> str:
    """Run id of one candidate in one generation round."""
    return f"champ-g{generation:03d}-r{round_no:04d}-{name}"


def expand(text: str, generation: int) -> str:
    """Replace `{generation}` by the generation number zero-padded to three digits."""
    return text.replace("{generation}", f"{generation:03d}")


def teacher_target(candidate: Candidate, generation: int) -> Path | None:
    """Teacher file a student candidate distils from; None for an ordinary one."""
    value = candidate.overrides.get("gate_soft_targets")
    return REPO / expand(value, generation) if isinstance(value, str) else None


def rates(result: Row) -> str:
    """Held-out error counts and rates of one result."""
    return (
        f"f2p={result['fail_to_pass_count']}/{result['fail_count']}"
        f" ({result['fail_to_pass_rate']:.2%})"
        f" p2f={result['pass_to_fail_count']}/{result['pass_count']}"
        f" ({result['pass_to_fail_rate']:.2%})"
    )


def describe(row: Row) -> str:
    """One-line held-out error counts and rates of an eval row."""
    members = row.get("members")
    which = f"epoch={row['epoch']}" if members is None else f"members={len(members)}"
    return f"{row['run_id']} {which} {rates(row['result'])}"


def ensemble_run_id(generation: int) -> str:
    """Synthetic run id of one generation's teacher-member ensemble."""
    return f"ens-g{generation:03d}-teacher"


def round_ensemble_run_id(generation: int, round_no: int) -> str:
    """Synthetic run id of the all-selected-checkpoints ensemble after one round."""
    return f"ens-g{generation:03d}-r{round_no:04d}-all"


def validation_gate_auc(run_dir: Path, checkpoint_name: str) -> float | None:
    """Validation gate AUC of one checkpoint from the run's per-epoch scores file.

    None when the run predates per-epoch validation scores or the file is not the
    evaluator's JSONL (the caller then ranks by the selection's Fail->Pass count).
    """
    epoch = int(Path(checkpoint_name).stem.rsplit("-", 1)[1])
    path = run_dir / f"validation-epoch-{epoch:04d}.scores.jsonl"
    if not path.exists():
        return None
    try:
        rows = [
            json.loads(line)
            for line in path.read_text(encoding="utf-8").splitlines()
            if line.strip()
        ]
        margins = [float(r["margin"]) for r in rows]
        return gate_auc(margins, [int(r["target"]) for r in rows])
    except (ValueError, KeyError, TypeError):
        return None


def round_manifest_name(round_no: int) -> str:
    """Manifest file name of one round's all-members ensemble."""
    return f"round-{round_no:04d}-all.json"


def popen_detached(
    command: Sequence[str], log_path: Path, env: dict[str, str]
) -> subprocess.Popen[bytes]:
    """Start a background command in its own session, logging to a file."""
    log_path.parent.mkdir(parents=True, exist_ok=True)
    with log_path.open("ab") as log:
        return subprocess.Popen(  # noqa: S603
            command,
            cwd=REPO,
            env=env,
            stdout=log,
            stderr=subprocess.STDOUT,
            start_new_session=True,
        )


@dataclass(frozen=True)
class Job:
    """One training run of one candidate in one generation round."""

    generation: int
    round_no: int
    seed: int
    candidate: str
    config_path: Path
    config_sha256: str
    run_id: str


@dataclass(frozen=True)
class TeacherJob:
    """The in-sample teacher build of one generation from its trained early runs."""

    generation: int
    runs: tuple[Row, ...]


@dataclass(frozen=True)
class EnsembleJob:
    """The test evaluation of one generation's teacher-member ensemble."""

    generation: int


@dataclass(frozen=True)
class RoundEnsembleJob:
    """The test evaluation of the ensemble of every selected checkpoint so far."""

    generation: int
    round_no: int


@dataclass(frozen=True)
class ReserveJob:
    """The exact champion and saturation trigger awaiting resource admission."""

    generation: int
    champion: Row
    test_evals: int


type LoopJob = Job | TeacherJob | EnsembleJob | RoundEnsembleJob | ReserveJob


class Loop:
    """Global job pool with an append-only ledger; all state is replayed from it."""

    def __init__(  # noqa: PLR0915 - explicit replay and owned dispatch configuration
        self,
        *,
        state_dir: Path,
        candidates_path: Path,
        runs_root: Path,
        data_root: Path,
        parallel: int,
        base_seed: int,
        min_free_gib: float,
        saturation_evals: int,
        max_reserve_checks: int,
        runner: Runner,
        package_command: Sequence[str] | None = None,
        package_launcher: Launcher | None = None,
        gpu_free: Callable[[], float] = gpu_free_gib,
        pids: Callable[[], int] = pids_current,
        poll_seconds: float = 30.0,
        launch_gap_seconds: float = 30.0,
        teacher_after_round: int = 0,
        teacher_path: str = TEACHER_PATH,
        band_decisions: bool = False,
        incremental_pareto: bool = True,
        training_runner: FiniteTrainingRunner | None = None,
        score_cache_dir: Path | None = None,
        reuse_teacher_preparation: bool = False,
        worker_pool: TrainingWorkerPool | None = None,
        resource_estimator: Callable[[LoopJob], ResourceEstimate] | None = None,
    ) -> None:
        self.state_dir = state_dir
        self.candidates_path = candidates_path
        self.runs_root = runs_root
        self.data_root = data_root
        self.parallel = parallel
        self.base_seed = base_seed
        self.min_free_gib = min_free_gib
        self.saturation_evals = saturation_evals
        self.max_reserve_checks = max_reserve_checks
        self._execution = threading.local()
        self._default_runner = runner
        self.package_command = package_command
        self.package_launcher = package_launcher
        self.gpu_free = gpu_free
        self.pids = pids
        self.poll_seconds = poll_seconds
        self.launch_gap_seconds = launch_gap_seconds
        self.teacher_after_round = teacher_after_round
        self.teacher_path = teacher_path
        # Compare champions on reject-band decided counts instead of boundary 0.
        self.band_decisions = band_decisions
        self.compare_key = "decided" if band_decisions else "result"
        self.incremental_pareto = incremental_pareto
        (
            self._default_training_runner,
            self.score_cache_dir,
            self.reuse_teacher_preparation,
        ) = (
            training_runner,
            score_cache_dir,
            reuse_teacher_preparation,
        )
        self.ledger = state_dir / "ledger.jsonl"
        self.stop_file = state_dir / "STOP"
        self.terminating = False
        self.wake = threading.Event()
        self.lock = threading.RLock()
        self.waits = 0
        self.reserve_running = False
        self.package_proc: tuple[Handle, Row] | None = None
        self.package_queued: Row | None = None
        # Run ids of the ordinary candidates of rounds <= teacher_after_round, per
        # generation, as the job generator saw them in this process.
        self.teacher_plan: dict[int, set[str]] = {}
        self.teacher_tried: set[int] = set()
        self.ensemble_tried: set[int] = set()
        # (generation, round) -> ordinary run ids of that round (students included)
        self.round_plan: dict[tuple[int, int], set[str]] = {}
        self.round_ensemble_tried: set[tuple[int, int]] = set()
        self.teacher_waits: set[tuple[int, Path]] = set()
        # State below is rebuilt from the ledger by apply().
        self.generation = 0
        self.worker_pool, self.resource_estimator = worker_pool, resource_estimator
        if (worker_pool is None) != (resource_estimator is None):
            message = "worker pool requires an explicit actual-input resource estimator"
            raise ValueError(message)
        self.champion: Row | None = None
        self.evals: list[Row] = []
        self.pareto: list[Row] = []
        self.checked_evals = 0
        self.reserve_checks = 0
        self.advance_reason: str | None = None
        self.runs: dict[str, Row] = {}
        self.evaluated: set[tuple[str, str, str]] = set()
        self.packaged: set[tuple[str, str]] = set()
        self.teachers: dict[int, Row] = {}
        if self.ledger.exists():
            for line in self.ledger.read_text(encoding="utf-8").splitlines():
                self.apply(json.loads(line))

    @property
    def runner(self) -> Runner:
        """Route each admitted job through its fixed device's native runner."""
        context: WorkerContext | None = getattr(self._execution, "context", None)
        return context.native if context is not None else self._default_runner

    @runner.setter
    def runner(self, runner: Runner) -> None:
        self._default_runner = runner

    @property
    def training_runner(self) -> FiniteTrainingRunner | None:
        """Keep a finite training process private to one admitted execution slot."""
        context: WorkerContext | None = getattr(self._execution, "context", None)
        return (
            context.training if context is not None else self._default_training_runner
        )

    @training_runner.setter
    def training_runner(self, runner: FiniteTrainingRunner | None) -> None:
        self._default_training_runner = runner

    # ---- ledger and state -------------------------------------------------

    def log(self, message: str) -> None:
        """Append one event line to loop.log and stdout."""
        line = f"[{stamp()}] {message}"
        with self.lock:
            self.state_dir.mkdir(parents=True, exist_ok=True)
            with (self.state_dir / "loop.log").open("a", encoding="utf-8") as stream:
                _ = stream.write(line + "\n")
            print(line, flush=True)

    def apply(self, row: Row) -> None:
        """Fold one ledger row into the in-memory state (live and on replay)."""
        match row["type"]:
            case "generation":
                self.generation = row["generation"]
                self.champion, self.evals, self.pareto = None, [], []
                self.checked_evals = self.reserve_checks = 0
                self.advance_reason = None
            case "run":
                self.runs[row["run_id"]] = row
            case "eval":
                self.apply_eval(row)
            case "reserve_check" if row["generation"] == self.generation:
                self.apply_reserve_check(row)
            case "teacher" if row["status"] == "ok":
                self.teachers[row["generation"]] = row
            case "package" if row["event"] in {"launched", "dropped"}:
                self.packaged.add((row["run_id"], row["checkpoint_name"]))

    def apply_eval(self, row: Row) -> None:
        """Track evaluated checkpoints, the champion and the Pareto set."""
        self.evaluated.add((row["run_id"], row["checkpoint_name"], row["role"]))
        if (
            row["role"] != "test"
            or row["status"] != "ok"
            or row["generation"] != self.generation
        ):
            return
        key = self.compare_key
        if row.get(key) is None:  # evaluated before the band existed: not comparable
            return
        self.evals.append(row)
        self.pareto = (
            incremental_pareto_front(self.pareto, row, key)
            if self.incremental_pareto
            else pareto_front(self.evals, key)
        )
        if self.band_decisions:
            self.champion = front_champion(self.pareto, key)
        elif self.champion is None or dominates(row[key], self.champion[key]):
            self.champion = row

    def apply_reserve_check(self, row: Row) -> None:
        """Reset the saturation counter; a failed or last check ends the generation."""
        if row.get("stale_champion", False):
            return
        self.checked_evals = row["test_evals"]
        if row["status"] != "ok":
            return
        self.reserve_checks += 1
        if row["verdict"] == "fail":
            self.advance_reason = "reserve_failed"
        elif self.reserve_checks >= self.max_reserve_checks:
            self.advance_reason = "max_reserve_checks"

    def record(self, row: Row) -> None:
        """Append one ledger row (fsynced), apply it, and refresh state.json."""
        with self.lock:
            self.state_dir.mkdir(parents=True, exist_ok=True)
            with self.ledger.open("a", encoding="utf-8") as stream:
                _ = stream.write(json.dumps(row, sort_keys=False) + "\n")
                stream.flush()
                os.fsync(stream.fileno())
            self.apply(row)
            state = {
                "generation": self.generation,
                "split": str(
                    self.state_dir / "splits" / f"gen-{self.generation:03d}.yaml"
                ),
                "champion": summary(self.champion),
                "pareto_set": [summary(item) for item in self.pareto],
                "test_evals": len(self.evals),
                "tests_since_reserve_check": len(self.evals) - self.checked_evals,
                "reserve_checks": self.reserve_checks,
                "advance_reason": self.advance_reason,
                "updated": stamp(),
            }
            tmp = self.state_dir / "state.json.tmp"
            _ = tmp.write_text(json.dumps(state, indent=2) + "\n", encoding="utf-8")
            _ = tmp.replace(self.state_dir / "state.json")

    def start_generation(self, reason: str) -> None:
        """Open the next generation (the split itself is built lazily)."""
        with self.lock:
            number = self.generation + 1
            self.record(
                {
                    "type": "generation",
                    "event": "start" if self.generation == 0 else "advance",
                    "generation": number,
                    "reason": reason,
                    "split_seed": self.base_seed + number,
                    "previous_champion": summary(self.champion),
                    "previous_pareto_set": [summary(row) for row in self.pareto],
                    "started": stamp(),
                }
            )
        self.log(f"generation {number} start reason={reason}")

    # ---- splits, configs, candidates --------------------------------------

    def split_for(self, generation: int) -> Path:
        """Return the generation's final split, building it if it is missing."""
        splits = self.state_dir / "splits"
        stem = f"gen-{generation:03d}"
        final = splits / f"{stem}.yaml"
        if final.exists():
            return final
        devval = splits / f"{stem}-devval.yaml"
        tested = splits / f"{stem}-test.yaml"
        tmp = splits / f"{stem}.yaml.tmp"
        for partial in (devval, tested, tmp):
            if partial.exists():
                _ = move_aside(partial)
        seed = str(self.base_seed + generation)
        scripts = REPO / "scripts"
        carve = str(scripts / "carve_test_split.py")
        commands = [
            [
                PYTHON,
                str(scripts / "generate_folder_grade_split.py"),
                *("--data-root", str(self.data_root)),
                *("--output", str(devval)),
                *("--seed", seed),
                *("--validation-fraction", VALIDATION_FRACTION),
            ],
            [
                PYTHON,
                carve,
                *("--split", str(devval)),
                *("--output", str(tested)),
                *("--role", "test"),
                *("--test-fraction", HELD_OUT_FRACTION),
                *("--seed", seed),
            ],
            [
                PYTHON,
                carve,
                *("--split", str(tested)),
                *("--output", str(tmp)),
                *("--role", "reserve"),
                *("--test-fraction", HELD_OUT_FRACTION),
                *("--seed", seed),
            ],
        ]
        log = self.state_dir / "logs" / f"split-{stem}.log"
        for command in commands:
            rc = self.runner(command, log)
            if rc != 0:
                message = f"split generation failed rc={rc}: {command[1]}"
                raise RuntimeError(message)
        _ = tmp.rename(final)
        return final

    def config_for(
        self,
        candidate: Candidate,
        generation: int,
        round_no: int,
        split: Path,
    ) -> tuple[Path, str]:
        """Write the job config for one candidate and return (path, sha256).

        String overrides may contain `{generation}`, replaced by the generation number
        zero-padded to three digits; other values and other braces stay untouched.
        """
        config = yaml.safe_load(candidate.base_config.read_text(encoding="utf-8"))
        config.update(
            {
                key: expand(value, generation) if isinstance(value, str) else value
                for key, value in candidate.overrides.items()
            }
        )
        config["split"] = str(split)
        config["seed"] = TRAIN_SEED_BASE + round_no
        text = yaml.safe_dump(config, sort_keys=False)
        path = (
            self.state_dir
            / "configs"
            / f"gen-{generation:03d}"
            / f"round-{round_no:04d}"
            / f"{candidate.name}.yaml"
        )
        if path.exists() and path.read_text(encoding="utf-8") != text:
            _ = move_aside(path)
        if not path.exists():
            path.parent.mkdir(parents=True, exist_ok=True)
            with path.open("x", encoding="utf-8") as stream:
                _ = stream.write(text)
        return path, hashlib.sha256(text.encode()).hexdigest()

    def read_candidates(self, previous: list[Candidate] | None) -> list[Candidate]:
        """Re-read the candidates file; keep the previous list if it is invalid."""
        try:
            return load_candidates(self.candidates_path)
        except (OSError, TypeError, ValueError, yaml.YAMLError) as error:
            if previous is None:
                raise
            self.log(f"candidates file invalid, keeping previous list: {error}")
            return previous

    def request_advance(self, reason: str) -> None:
        """Open the next generation before any new job (e.g. after a relabelled corpus)."""
        if self.generation == 0:
            message = "no generation to advance from"
            raise ValueError(message)
        self.advance_reason = reason

    def stopping(self) -> bool:
        """Whether new launches are disabled."""
        return self.terminating or self.stop_file.exists()

    def terminate(self) -> None:
        """Stop scheduling and make workers skip bookkeeping."""
        self.terminating = True
        self.wake.set()
        if self.worker_pool is not None:
            self.worker_pool.terminate()

    # ---- job generation ---------------------------------------------------

    def missing_evals(self, run: Row) -> list[str]:
        """Checkpoints of a trained run that have no test eval row yet."""
        return [
            name
            for name in run["checkpoints"]
            if (run["run_id"], name, "test") not in self.evaluated
        ]

    def pending(self, run_id: str) -> bool:
        """Whether a run still needs training or test evaluations."""
        run = self.runs.get(run_id)
        return run is None or bool(self.missing_evals(run))

    def make_job(
        self, candidate: Candidate, generation: int, round_no: int, split: Path
    ) -> Job:
        """Build the job; a resumed run keeps the config it was trained with."""
        run_id = run_id_for(generation, round_no, candidate.name)
        run = self.runs.get(run_id)
        if run is not None:
            path, sha = Path(run["config_path"]), run["config_sha256"]
        else:
            path, sha = self.config_for(candidate, generation, round_no, split)
        seed = TRAIN_SEED_BASE + round_no
        return Job(generation, round_no, seed, candidate.name, path, sha, run_id)

    def jobs(self) -> Iterator[Job | TeacherJob | EnsembleJob | RoundEnsembleJob]:
        """Yield pending jobs round after round within the current generation.

        A due teacher build comes first; students whose teacher file is missing sit
        the round out.
        """
        candidates: list[Candidate] | None = None
        generation, round_no = self.generation, 0
        idle_logged = False
        while not self.stopping():
            if self.advance_reason:
                self.start_generation(self.advance_reason)
            if self.generation != generation:
                generation, round_no = self.generation, 0
            candidates = self.read_candidates(candidates)
            yield from self.due_builds(generation, round_no)
            waiting = self.waiting_students(candidates, generation)
            if len(waiting) == len(candidates):  # empty, or only students waiting
                if not idle_logged:
                    self.log("no active candidates; waiting")
                    idle_logged = True
                _ = self.wake.wait(self.poll_seconds)
                continue
            idle_logged = False
            round_no += 1
            self.plan_teacher(generation, round_no, candidates)
            self.round_plan.setdefault((generation, round_no), set()).update(
                run_id_for(generation, round_no, c.name)
                for c in candidates
                if c not in waiting
            )
            todo = [
                c
                for c in candidates
                if c not in waiting
                and self.pending(run_id_for(generation, round_no, c.name))
            ]
            if not todo:
                continue
            split = self.split_for(generation)
            self.log(
                f"generation {generation} round {round_no} start "
                f"seed={TRAIN_SEED_BASE + round_no} "
                f"candidates={','.join(c.name for c in todo)}"
            )
            for candidate in todo:
                if self.advance_reason or self.stopping():
                    break
                yield from self.due_builds(generation, round_no)
                job = self.make_job(candidate, generation, round_no, split)
                if (
                    job.run_id not in self.runs
                    and (self.runs_root / job.run_id).exists()
                ):
                    now = stamp()
                    self.record(self.run_row(job, "stale_run_dir", now, now))
                    self.log(f"job skip {job.run_id} status=stale_run_dir")
                    continue
                yield job

    # ---- training and evaluation ------------------------------------------

    @staticmethod
    def run_row(
        job: Job, status: str, started: str, finished: str, **extra: object
    ) -> Row:
        """Build one `run` ledger row with the full fixed key set."""
        row: Row = {
            "type": "run",
            "generation": job.generation,
            "round": job.round_no,
            "seed": job.seed,
            "candidate": job.candidate,
            "run_id": job.run_id,
            "config_path": str(job.config_path),
            "config_sha256": job.config_sha256,
            "status": status,
            "train_rc": None,
            "started": started,
            "finished": finished,
            "epochs_run": None,
            "checkpoints": [],
            "error": None,
        }
        row.update(extra)
        return row

    def train(self, job: Job) -> None:
        """Train one job and record its `run` row (unless terminating)."""
        started = stamp()
        train_log = self.state_dir / "logs" / f"{job.run_id}.train.log"
        rc = (
            self.training_runner.run(
                config=job.config_path,
                config_sha256=job.config_sha256,
                output_root=self.runs_root,
                run_id=job.run_id,
                log_path=train_log,
                generation=job.generation,
                score_cache_dir=self.score_cache_dir,
            )
            if self.training_runner is not None
            else self.runner(
                [
                    PYTHON,
                    *("-m", "raina_laya.choice_train"),
                    *("--config", str(job.config_path)),
                    *("--output-root", str(self.runs_root)),
                    *("--run-id", job.run_id),
                    *(
                        ()
                        if self.score_cache_dir is None
                        else ("--score-cache-dir", str(self.score_cache_dir))
                    ),
                ],
                train_log,
            )
        )
        if self.terminating:
            return
        run_dir = self.runs_root / job.run_id
        fields: Row = {
            "train_rc": rc,
            "epochs_run": len(list(run_dir.glob("validation-epoch-*.json"))),
        }
        status = "train_failed"
        if rc == 0:
            try:
                fields["checkpoints"] = collect_checkpoints(run_dir)
                status = "trained"
            except (OSError, KeyError, ValueError) as error:
                fields["error"] = repr(error)
        elif rc == 1 and NO_CHECKPOINT in train_log.read_text(
            encoding="utf-8", errors="replace"
        ):
            status = "no_eligible_checkpoint"
        self.record(self.run_row(job, status, started, stamp(), **fields))

    def eval_path(self, run_id: str, name: str, role: str) -> Path:
        """Report path of one checkpoint evaluation."""
        return self.state_dir / "evals" / run_id / f"{Path(name).stem}.{role}.json"

    def scores_path(self, run_id: str, name: str, role: str) -> Path:
        """Per-song gate margin path of one checkpoint evaluation."""
        stem = Path(name).stem
        return self.state_dir / "evals" / run_id / f"{stem}.{role}.scores.jsonl"

    def evaluate(  # noqa: PLR0917 - positional report/scores paths predate `band`
        self,
        run: Row,
        name: str,
        role: str,
        output: Path | None = None,
        scores: Path | None = None,
        band: Path | None = None,
    ) -> Row | None:
        """Score one checkpoint on a held-out role; reuse a finished report.

        The report and scores go to the loop's evals folder unless given. `band` is
        the validation scores file that fixes the reject band of this report.
        """
        output = output or self.eval_path(run["run_id"], name, role)
        if not output.exists():
            output.parent.mkdir(parents=True, exist_ok=True)
            scores = scores or self.scores_path(run["run_id"], name, role)
            if scores.exists():  # a killed evaluator left scores but no report
                _ = move_aside(scores)
            if is_ensemble(run):
                target = [
                    "raina_laya.ensemble_evaluate",
                    "--manifest",
                    run["config_path"],
                ]
            else:
                target = [
                    "raina_laya.choice_evaluate",
                    *("--config", run["config_path"]),
                    *("--checkpoint", str(self.runs_root / run["run_id"] / name)),
                ]
            rc = self.runner(
                [
                    PYTHON,
                    "-m",
                    *target,
                    *("--output", str(output)),
                    *("--role", role),
                    *("--scores-output", str(scores)),
                    *(() if band is None else ("--band-scores", str(band))),
                    *(
                        ()
                        if self.score_cache_dir is None
                        else ("--score-cache-dir", str(self.score_cache_dir))
                    ),
                ],
                self.state_dir
                / "logs"
                / f"{run['run_id']}.{Path(name).stem}.{role}.eval.log",
            )
            if rc != 0 or self.terminating:
                return None
        return json.loads(output.read_text(encoding="utf-8"))

    def band_scores(self, run: Row, name: str) -> Path | None:
        """Validation scores that fix one checkpoint's reject band.

        Training writes `validation-epoch-NNNN.scores.jsonl` into the run directory;
        runs from before that (and ensembles) are scored once on validation instead.
        """
        if not is_ensemble(run):
            epoch = int(Path(name).stem.rsplit("-", 1)[1])
            local = (
                self.runs_root
                / run["run_id"]
                / f"validation-epoch-{epoch:04d}.scores.jsonl"
            )
            if local.exists():
                return local
        if self.evaluate(run, name, "validation") is None:
            return None
        scores = self.scores_path(run["run_id"], name, "validation")
        return scores if scores.exists() else None

    def evaluate_many(
        self, run: Row, names: Sequence[str], role: str, bands: dict[str, Path | None]
    ) -> dict[str, Row | None]:
        """Score several checkpoints of one run in one evaluator process.

        Finished reports are reused; an ensemble (one manifest) goes through
        `evaluate`. The dataset is loaded once per process, which is what makes
        this cheaper than one process per checkpoint.
        """
        if is_ensemble(run):
            return {
                name: self.evaluate(run, name, role, band=bands.get(name))
                for name in names
            }
        jobs: list[Row] = []
        for name in names:
            output = self.eval_path(run["run_id"], name, role)
            if output.exists():
                continue
            output.parent.mkdir(parents=True, exist_ok=True)
            scores = self.scores_path(run["run_id"], name, role)
            if scores.exists():  # a killed evaluator left scores but no report
                _ = move_aside(scores)
            job: Row = {
                "checkpoint": str(self.runs_root / run["run_id"] / name),
                "output": str(output),
                "scores_output": str(scores),
            }
            if bands.get(name) is not None:
                job["band_scores"] = str(bands[name])
            jobs.append(job)
        if jobs:
            logs = self.state_dir / "logs"
            logs.mkdir(parents=True, exist_ok=True)
            jobs_path = logs / f"{run['run_id']}.{role}.jobs.json"
            _ = jobs_path.write_text(json.dumps(jobs, indent=1), encoding="utf-8")
            _ = self.runner(
                [
                    PYTHON,
                    "-m",
                    "raina_laya.choice_evaluate",
                    *("--config", run["config_path"]),
                    *("--role", role),
                    *("--jobs", str(jobs_path)),
                    *(
                        ()
                        if self.score_cache_dir is None
                        else ("--score-cache-dir", str(self.score_cache_dir))
                    ),
                ],
                logs / f"{run['run_id']}.{role}.eval.log",
            )
        reports: dict[str, Row | None] = {}
        for name in names:
            output = self.eval_path(run["run_id"], name, role)
            reports[name] = (
                json.loads(output.read_text(encoding="utf-8"))
                if output.exists()
                else None
            )
        return reports

    def eval_row(self, run: Row, name: str, role: str, report: Row | None) -> Row:
        """Build one `eval` ledger row (status failed when there is no report)."""
        ensemble = is_ensemble(run)
        epoch = None if ensemble else int(Path(name).stem.rsplit("-", 1)[1])
        run_dir = self.runs_root / run["run_id"]
        row: Row = {
            "type": "eval",
            "generation": run["generation"],
            "role": role,
            "run_id": run["run_id"],
            "candidate": run["candidate"],
            "round": run["round"],
            "checkpoint": run["config_path"] if ensemble else str(run_dir / name),
            "checkpoint_name": name,
            "checkpoint_sha256": None,
            "epoch": epoch,
            "config_path": run["config_path"],
            "status": "failed",
            "validation": None
            if epoch is None
            else validation_metrics(run_dir / f"validation-epoch-{epoch:04d}.json"),
            "result": None,
            "gate_auc": None,
            "gate_curve": None,
            "reject_band": None,
            "decided": None,
            "eval_json": str(self.eval_path(run["run_id"], name, role)),
            "finished": stamp(),
        }
        if ensemble:
            row.update(manifest=run["config_path"], members=run["members"])
        if report is not None:
            row.update(
                status="ok",
                checkpoint_sha256=report["checkpoint_sha256"],
                result=heldout_result(report),
                # Reports finished before these keys existed simply lack them.
                gate_auc=report.get("gate_auc"),
                gate_curve=report.get("gate_curve"),
                reject_band=report.get("reject_band"),
                decided=report.get("decided"),
            )
        return row

    def record_eval(self, row: Row) -> None:
        """Record a test eval row, then champion and Pareto-set events and packaging."""
        new_champion = None
        with self.lock:
            old_champion = self.champion
            old_front = {(r["run_id"], r["checkpoint_name"]) for r in self.pareto}
            self.record(row)
            if row["status"] == "ok":
                auc = row.get("gate_auc")
                suffix = "" if auc is None else f" auc={auc:.4f}"
                if "members" in row:
                    self.log(
                        f"ensemble eval g={row['generation']} "
                        f"members={len(row['members'])} {rates(row['result'])}{suffix}"
                    )
                else:
                    self.log(f"eval {row['role']} {describe(row)}{suffix}")
            else:
                self.log(f"eval failed {row['run_id']} {row['checkpoint_name']}")
            if self.champion is not old_champion and self.champion is not None:
                new_champion = self.champion
                self.record(
                    {
                        "type": "champion",
                        "generation": new_champion["generation"],
                        "previous": summary(old_champion),
                        "champion": summary(new_champion),
                        "finished": stamp(),
                    }
                )
                self.log(f"new champion g={self.generation} {describe(new_champion)}")
            front = {(r["run_id"], r["checkpoint_name"]) for r in self.pareto}
            if front != old_front:
                self.record(
                    {
                        "type": "pareto_set",
                        "generation": self.generation,
                        "members": [summary(r) for r in self.pareto],
                        "finished": stamp(),
                    }
                )
        if new_champion is not None:
            self.request_package(new_champion)

    def evaluate_run(self, run_id: str) -> None:
        """Score every not yet scored checkpoint of a trained run on the test role."""
        run = self.runs.get(run_id)
        if run is None:
            return
        names = self.missing_evals(run)
        if (
            self.reuse_teacher_preparation
            and self.band_decisions
            and not is_ensemble(run)
        ):
            missing_validation = [
                name
                for name in names
                if not (
                    self.runs_root
                    / run_id
                    / f"validation-epoch-{int(Path(name).stem.rsplit('-', 1)[1]):04d}.scores.jsonl"
                ).exists()
            ]
            if missing_validation:
                self.evaluate_many(
                    run,
                    missing_validation,
                    "validation",
                    dict.fromkeys(missing_validation),
                )
        bands: dict[str, Path | None] = {}
        for name in names:
            bands[name] = self.band_scores(run, name)
            if self.terminating:
                return
        if not names:
            return
        reports = self.evaluate_many(run, names, "test", bands)
        if self.terminating:
            return
        for name in names:
            self.record_eval(self.eval_row(run, name, "test", reports[name]))
            self.maybe_reserve_check(run["generation"])

    def execute(self, job: LoopJob) -> None:
        """Train (unless already trained) and evaluate one job, or build a teacher."""
        if isinstance(job, ReserveJob):
            try:
                self.reserve_check(job.champion, job.test_evals, require_current=True)
            finally:
                self.reserve_running = False
            return
        if isinstance(job, TeacherJob):
            self.build_teacher(job)
            return
        if isinstance(job, EnsembleJob):
            self.evaluate_ensemble(job)
            return
        if isinstance(job, RoundEnsembleJob):
            self.evaluate_round_ensemble(job)
            return
        self.log(f"job start {job.run_id}")
        try:
            if job.run_id not in self.runs:
                self.train(job)
            self.evaluate_run(job.run_id)
        except Exception as error:  # noqa: BLE001
            self.log(f"job error {job.run_id}: {error!r}")
        run = self.runs.get(job.run_id)
        if run is not None and not self.terminating:
            scored = sum(
                1 for key in self.evaluated if key[0] == job.run_id and key[2] == "test"
            )
            self.log(
                f"job end {job.run_id} status={run['status']} "
                f"epochs={run['epochs_run']} test_evals={scored}"
            )

    # ---- in-sample teacher ------------------------------------------------

    def teacher_file(self, generation: int) -> Path:
        """Path of one generation's in-sample teacher file."""
        return REPO / expand(self.teacher_path, generation)

    def waiting_students(
        self, candidates: list[Candidate], generation: int
    ) -> list[Candidate]:
        """Students whose teacher file does not exist yet (each path logged once)."""
        waiting: list[Candidate] = []
        for candidate in candidates:
            path = teacher_target(candidate, generation)
            if path is None or path.exists():
                continue
            if (generation, path) not in self.teacher_waits:
                self.teacher_waits.add((generation, path))
                self.log(f"waiting for teacher {path}")
            waiting.append(candidate)
        return waiting

    def plan_teacher(
        self, generation: int, round_no: int, candidates: list[Candidate]
    ) -> None:
        """Remember the ordinary runs of a round the teacher build has to wait for."""
        if round_no <= self.teacher_after_round:
            self.teacher_plan.setdefault(generation, set()).update(
                run_id_for(generation, round_no, c.name)
                for c in candidates
                if teacher_target(c, generation) is None
            )

    def teacher_builds(self, generation: int, round_no: int) -> Iterator[TeacherJob]:
        """Yield the generation's teacher build once its early ordinary runs are final.

        The plan holds every round <= teacher_after_round once `round_no` reached it.
        A build is offered once per process; a restart offers it again.
        """
        ids = sorted(self.teacher_plan.get(generation, ()))
        if (
            round_no < self.teacher_after_round
            or not ids
            or generation in self.teacher_tried
            or any(self.pending(run_id) for run_id in ids)
            or self.teacher_file(generation).exists()
        ):
            return
        self.teacher_tried.add(generation)
        rows = (self.runs[run_id] for run_id in ids)
        yield TeacherJob(generation, tuple(r for r in rows if r["status"] == "trained"))

    def build_teacher(self, job: TeacherJob) -> None:
        """Build one teacher file and record the outcome (not when terminating)."""
        started, path = stamp(), self.teacher_file(job.generation)
        self.log(f"teacher start g={job.generation} runs={len(job.runs)}")
        members: list[Row] = []
        digest = None
        try:
            digest = self.write_teacher(job, path, members)
        except Exception as error:  # noqa: BLE001
            self.log(f"teacher error g={job.generation}: {error!r}")
        if self.terminating and digest is None:
            return
        status = "failed" if digest is None else "ok"
        extra = {} if digest is None else {"sha256": digest}
        self.record(
            {
                "type": "teacher",
                "generation": job.generation,
                "status": status,
                "path": str(path),
                **extra,
                "members": members,
                "started": started,
                "finished": stamp(),
            }
        )
        self.log(f"teacher {status} g={job.generation} members={len(members)} {path}")

    def write_teacher(
        self, job: TeacherJob, path: Path, members: list[Row]
    ) -> str | None:
        """Score every member on development songs and average the gate margins.

        Appends each member to `members`; returns the sha256 of the new teacher file,
        or None when any step failed.
        """
        folder = path.parent / "members"
        inputs: list[str] = []
        for run in job.runs:
            run_id = run["run_id"]
            selection = self.runs_root / run_id / "selection.json"
            name = json.loads(selection.read_text(encoding="utf-8"))["checkpoint"]
            members.append({"run_id": run_id, "checkpoint_name": name})
            scores = folder / f"{run_id}.dev.scores.jsonl"
            inputs.append(str(scores))
            report = folder / f"{run_id}.dev.json"
            if (
                not self.reuse_teacher_preparation
                and self.evaluate(run, name, "development", report, scores) is None
            ):
                return None
        if self.reuse_teacher_preparation:
            jobs = [
                (
                    run,
                    member["checkpoint_name"],
                    folder / f"{run['run_id']}.dev.json",
                    folder / f"{run['run_id']}.dev.scores.jsonl",
                )
                for run, member in zip(job.runs, members, strict=True)
            ]
            if not self.evaluate_compatible_members(
                jobs, "development", f"teacher-g{job.generation:03d}"
            ):
                return None
        tmp = path.with_name(f"{path.name}.tmp")
        if tmp.exists():  # a killed run left a partial file
            _ = move_aside(tmp)
        rc = self.runner(
            [
                PYTHON,
                str(REPO / "scripts" / "teacher_scores.py"),
                *("--inputs", *inputs),
                *("--output", str(tmp)),
            ],
            self.state_dir / "logs" / f"teacher-g{job.generation:03d}.log",
        )
        if rc != 0 or self.terminating:
            return None
        _ = tmp.rename(path)
        return hashlib.sha256(path.read_bytes()).hexdigest()

    def evaluate_compatible_members(
        self,
        entries: Sequence[tuple[Row, str, Path, Path]],
        role: str,
        namespace: str,
    ) -> bool:
        """Bundle verified input lineages, passing each member's actual config."""
        groups: dict[str, list[tuple[Row, str, Path, Path]]] = {}
        hashes: dict[tuple[Path, int, int, int, int, int], str] = {}
        for entry in entries:
            run, _name, report, _scores = entry
            if report.exists():
                continue
            payload = Path(run["config_path"]).read_bytes()
            if hashlib.sha256(payload).hexdigest() != run["config_sha256"]:
                message = "teacher member config SHA256 changed"
                raise ValueError(message)
            key = self.teacher_input_key(payload, hashes)
            groups.setdefault(key, []).append(entry)
        for index, group in enumerate(groups.values()):
            jobs: list[Row] = []
            for run, name, report, scores in group:
                report.parent.mkdir(parents=True, exist_ok=True)
                if scores.exists():
                    _ = move_aside(scores)
                jobs.append(
                    {
                        "config": run["config_path"],
                        "config_sha256": run["config_sha256"],
                        "checkpoint": str(self.runs_root / run["run_id"] / name),
                        "output": str(report),
                        "scores_output": str(scores),
                    }
                )
            logs = self.state_dir / "logs"
            logs.mkdir(parents=True, exist_ok=True)
            jobs_path = logs / f"{namespace}.{role}.compatible-{index:04d}.jobs.json"
            if jobs_path.exists():
                _ = move_aside(jobs_path)
            with jobs_path.open("x", encoding="utf-8") as stream:
                _ = stream.write(json.dumps(jobs, indent=1) + "\n")
            code = self.runner(
                [
                    PYTHON,
                    "-m",
                    "raina_laya.choice_evaluate",
                    "--config",
                    group[0][0]["config_path"],
                    "--role",
                    role,
                    "--jobs",
                    str(jobs_path),
                    *(
                        ()
                        if self.score_cache_dir is None
                        else ("--score-cache-dir", str(self.score_cache_dir))
                    ),
                ],
                logs / f"{namespace}.{role}.compatible-{index:04d}.log",
            )
            if code != 0 or self.terminating:
                return False
        return all(report.exists() for _run, _name, report, _scores in entries)

    @staticmethod
    def teacher_input_key(
        payload: bytes, hashes: dict[tuple[Path, int, int, int, int, int], str]
    ) -> str:
        """Identify explicit input paths; evaluator validates every job independently."""
        settings = yaml.safe_load(payload)
        if not {"data_root", "split", "cache_inventory"}.issubset(settings):
            # Incomplete legacy fixtures retain their conservative complete-config key.
            return json.dumps(settings, sort_keys=True, separators=(",", ":"))
        lineage: Row = {"embedding_layers": settings.get("embedding_layers")}
        for key in (
            "data_root",
            "split",
            "cache_inventory",
            "embedding_database",
            "embedding_layer_dir",
        ):
            if not settings.get(key):
                lineage[key] = None
                continue
            path = (REPO / settings[key]).resolve()
            value = path.stat()
            identity = (
                path,
                value.st_dev,
                value.st_ino,
                value.st_size,
                value.st_mtime_ns,
                value.st_ctime_ns,
            )
            content = None
            if key in {"split", "cache_inventory"}:
                if identity not in hashes:
                    hashes[identity] = hashlib.sha256(path.read_bytes()).hexdigest()
                content = hashes[identity]
            lineage[key] = (str(path), identity[1:], content)
        return json.dumps(lineage, sort_keys=True, separators=(",", ":"))

    def due_builds(
        self, generation: int, round_no: int
    ) -> Iterator[TeacherJob | EnsembleJob | RoundEnsembleJob]:
        """Yield a due teacher build, then due ensemble evaluations."""
        yield from self.teacher_builds(generation, round_no)
        yield from self.ensemble_builds(generation)
        yield from self.round_ensemble_builds(generation, round_no)

    def round_ensemble_builds(
        self, generation: int, round_no: int
    ) -> Iterator[RoundEnsembleJob]:
        """Yield the all-members ensemble of each finished round >= ROUND_ENSEMBLE_FROM.

        A round is finished once the next round has started (all its jobs were
        yielded) and none of its runs, nor those of earlier rounds, is pending.
        Offered once per process; a restart offers unscored ones again.
        """
        if generation != self.generation:
            return
        for (g, r), _ids in sorted(self.round_plan.items()):
            key = (round_ensemble_run_id(g, r), round_manifest_name(r), "test")
            if (
                g != generation
                or r < ROUND_ENSEMBLE_FROM
                or r >= round_no
                or (g, r) in self.round_ensemble_tried
                or key in self.evaluated
                or any(
                    self.pending(run_id)
                    for (gg, rr), ids in self.round_plan.items()
                    if gg == g and rr <= r
                    for run_id in ids
                )
            ):
                continue
            self.round_ensemble_tried.add((g, r))
            if len(self.round_members(g, r)) >= 2:  # noqa: PLR2004 - an ensemble
                yield RoundEnsembleJob(g, r)

    def round_members(self, generation: int, round_no: int) -> list[Row]:
        """Selected checkpoints of trained ordinary runs of rounds <= round_no.

        At most ROUND_ENSEMBLE_MAX members, the best by the selected checkpoint's
        validation gate AUC (fallback: fewest validation Fail->Pass); ties keep
        run-id order.
        """
        ranked: list[tuple[tuple[float, int], Row]] = []
        for (g, r), ids in sorted(self.round_plan.items()):
            if g != generation or r > round_no:
                continue
            for run_id in sorted(ids):
                run = self.runs.get(run_id)
                selection = self.runs_root / run_id / "selection.json"
                if run is None or run["status"] != "trained" or not selection.exists():
                    continue
                chosen = json.loads(selection.read_text(encoding="utf-8"))
                name = chosen["checkpoint"]
                auc = validation_gate_auc(self.runs_root / run_id, name)
                rank = (
                    -auc if auc is not None else 0.0,
                    int(chosen.get("false_negatives", 0)),
                )
                ranked.append(
                    (
                        rank,
                        {
                            "run_id": run_id,
                            "checkpoint_name": name,
                            "config": run["config_path"],
                            "checkpoint": str(self.runs_root / run_id / name),
                            "validation_gate_auc": auc,
                            "validation_fail_to_pass": rank[1],
                        },
                    )
                )
        ranked.sort(key=lambda item: item[0])
        return [member for _rank, member in ranked[:ROUND_ENSEMBLE_MAX]]

    def round_ensemble_run(self, generation: int, round_no: int) -> Row:
        """Return one round's all-members ensemble run, writing its manifest first."""
        run_id = round_ensemble_run_id(generation, round_no)
        if run_id in self.runs:
            return self.runs[run_id]
        members = self.round_members(generation, round_no)
        name = round_manifest_name(round_no)
        path = self.state_dir / "ensembles" / f"gen-{generation:03d}" / name
        if not path.exists():  # reuse a manifest left by an earlier attempt
            path.parent.mkdir(parents=True, exist_ok=True)
            entries = [
                {"config": m["config"], "checkpoint": m["checkpoint"]} for m in members
            ]
            with path.open("x", encoding="utf-8") as stream:
                _ = stream.write(json.dumps({"members": entries}, indent=2) + "\n")
        now = stamp()
        run: Row = {
            "type": "run",
            "generation": generation,
            "round": round_no,
            "seed": None,
            "candidate": "ensemble_round_members",
            "run_id": run_id,
            "config_path": str(path),
            "config_sha256": hashlib.sha256(path.read_bytes()).hexdigest(),
            "status": ENSEMBLE_STATUS,
            "train_rc": None,
            "started": now,
            "finished": now,
            "epochs_run": None,
            "checkpoints": [name],
            "error": None,
            "members": members,
        }
        self.record(run)
        return run

    def evaluate_round_ensemble(self, job: RoundEnsembleJob) -> None:
        """Score one round's all-members ensemble on the test role like a checkpoint."""
        try:
            run = self.round_ensemble_run(job.generation, job.round_no)
            self.log(
                f"round ensemble g={job.generation} r={job.round_no} "
                f"members={len(run['members'])}"
            )
            self.evaluate_run(run["run_id"])
        except Exception as error:  # noqa: BLE001
            self.log(
                f"round ensemble error g={job.generation} r={job.round_no}: {error!r}"
            )

    def ensemble_builds(self, generation: int) -> Iterator[EnsembleJob]:
        """Yield the ensemble evaluation once the teacher is built and it is unscored.

        A failed evaluation row also counts as scored, like for any checkpoint.
        """
        key = (ensemble_run_id(generation), MANIFEST_NAME, "test")
        if (
            generation != self.generation
            or generation not in self.teachers
            or generation in self.ensemble_tried
            or key in self.evaluated
        ):
            return
        self.ensemble_tried.add(generation)
        yield EnsembleJob(generation)

    def ensemble_run(self, generation: int) -> Row:
        """Return the generation's synthetic ensemble run, writing its manifest first."""
        run_id = ensemble_run_id(generation)
        if run_id in self.runs:
            return self.runs[run_id]
        members = self.teachers[generation]["members"]
        path = self.state_dir / "ensembles" / f"gen-{generation:03d}" / MANIFEST_NAME
        entries = [
            {
                "config": self.runs[m["run_id"]]["config_path"],
                "checkpoint": str(self.runs_root / m["run_id"] / m["checkpoint_name"]),
            }
            for m in members
        ]
        if not path.exists():  # reuse a manifest left by an earlier attempt
            path.parent.mkdir(parents=True, exist_ok=True)
            with path.open("x", encoding="utf-8") as stream:
                _ = stream.write(json.dumps({"members": entries}, indent=2) + "\n")
        now = stamp()
        run: Row = {
            "type": "run",
            "generation": generation,
            "round": self.teacher_after_round,
            "seed": None,
            "candidate": "ensemble_teacher_members",
            "run_id": run_id,
            "config_path": str(path),
            "config_sha256": hashlib.sha256(path.read_bytes()).hexdigest(),
            "status": ENSEMBLE_STATUS,
            "train_rc": None,
            "started": now,
            "finished": now,
            "epochs_run": None,
            "checkpoints": [MANIFEST_NAME],
            "error": None,
            "members": entries,
        }
        self.record(run)
        return run

    def evaluate_ensemble(self, job: EnsembleJob) -> None:
        """Score the teacher-member ensemble on the test role like a checkpoint."""
        try:
            run = self.ensemble_run(job.generation)
            self.evaluate_run(run["run_id"])
        except Exception as error:  # noqa: BLE001
            self.log(f"ensemble error g={job.generation}: {error!r}")

    # ---- saturation and reserve check -------------------------------------

    def maybe_reserve_check(self, generation: int) -> None:
        """Score the champion on the reserve set once enough test evals piled up."""
        if self.worker_pool is not None:
            self.wake.set()
            return
        with self.lock:
            if (
                generation != self.generation
                or self.advance_reason
                or self.reserve_running
                or self.champion is None
                or len(self.evals) - self.checked_evals < self.saturation_evals
            ):
                return
            self.reserve_running = True
            champion, test_evals = self.champion, len(self.evals)
        try:
            self.reserve_check(champion, test_evals)
        finally:
            self.reserve_running = False

    def reserve_check(
        self, champion: Row, test_evals: int, *, require_current: bool = False
    ) -> None:
        """Evaluate the champion on reserve and record the hold-up verdict."""
        run, name = self.runs[champion["run_id"]], champion["checkpoint_name"]
        band = self.band_scores(run, name)
        if self.terminating:
            return
        report = self.evaluate(run, name, "reserve", band=band)
        if self.terminating:
            return
        if (run["run_id"], name, "reserve") not in self.evaluated:
            self.record(self.eval_row(run, name, "reserve", report))
        row: Row = {
            "type": "reserve_check",
            "generation": champion["generation"],
            "run_id": champion["run_id"],
            "candidate": champion["candidate"],
            "checkpoint_name": name,
            "epoch": champion["epoch"],
            "test_evals": test_evals,
            "status": "failed",
            "verdict": "error",
            "test": champion["result"],
            "reserve": None,
            "test_decided": champion.get("decided"),
            "reserve_decided": None,
            "compared": self.compare_key,
            "eval_json": str(self.eval_path(run["run_id"], name, "reserve")),
            "finished": stamp(),
        }
        if report is None:
            self._record_reserve_check(row, champion, require_current)
            self.log(f"reserve check failed to run {describe(champion)}")
            return
        reserve = heldout_result(report)
        row.update(reserve=reserve, reserve_decided=report.get("decided"))
        if self.band_decisions:
            test, held = champion["decided"], report.get("decided")
        else:
            test, held = champion["result"], reserve
        if held is None:  # the reserve report has no band to compare on
            self._record_reserve_check(row, champion, require_current)
            self.log(f"reserve check has no decided counts {describe(champion)}")
            return
        row.update(status="ok", **reserve_verdict(test, held))
        self._record_reserve_check(row, champion, require_current)
        reserve = held
        self.log(
            f"reserve check g={champion['generation']} verdict={row['verdict']} "
            f"{champion['run_id']} epoch={champion['epoch']}"
            f" f2p test={test['fail_to_pass_rate']:.2%}"
            f" reserve={reserve['fail_to_pass_rate']:.2%} p={row['f2p_p']:.3f}"
            f" p2f test={test['pass_to_fail_rate']:.2%}"
            f" reserve={reserve['pass_to_fail_rate']:.2%} p={row['p2f_p']:.3f}"
            f" checks={self.reserve_checks}/{self.max_reserve_checks}"
        )

    def _record_reserve_check(
        self, row: Row, champion: Row, require_current: bool
    ) -> None:
        with self.lock:
            if require_current:
                row["stale_champion"] = (
                    self.generation != champion["generation"]
                    or self.champion is None
                    or (self.champion["run_id"], self.champion["checkpoint_name"])
                    != (champion["run_id"], champion["checkpoint_name"])
                )
            self.record(row)

    # ---- packaging --------------------------------------------------------

    def package_row(self, event: str, champion: Row, **extra: object) -> Row:
        """Build one `package` ledger row."""
        return {
            "type": "package",
            "event": event,
            "generation": champion["generation"],
            "run_id": champion["run_id"],
            "epoch": champion["epoch"],
            "checkpoint_name": champion["checkpoint_name"],
            "time": stamp(),
            **extra,
        }

    def request_package(self, champion: Row) -> None:
        """Launch the package command now, or queue the champion (newest wins)."""
        if self.package_command is None or self.package_launcher is None:
            return
        with self.lock:
            if self.package_proc is None:
                self.launch_package(champion)
                return
            if self.package_queued is not None:
                self.record(self.package_row("dropped", self.package_queued))
                self.log(f"package dropped {describe(self.package_queued)}")
            self.package_queued = champion
            self.log(f"package queued {describe(champion)}")

    def launch_package(self, champion: Row) -> None:
        """Start the package command for one champion without waiting for it."""
        if self.package_command is None or self.package_launcher is None:
            return
        manifest = champion.get("manifest")
        suffix = "" if manifest else f"-{champion['epoch']:04d}"
        log = self.state_dir / "logs" / f"package-{champion['run_id']}{suffix}.log"
        source = (
            ["--manifest", manifest]
            if manifest
            else [
                "--checkpoint",
                champion["checkpoint"],
                "--config",
                champion["config_path"],
            ]
        )
        command = [
            *self.package_command,
            *source,
            *("--test-eval", champion["eval_json"]),
            *("--generation", str(champion["generation"])),
            *("--state-dir", str(self.state_dir)),
        ]
        try:
            proc = self.package_launcher(command, log)
        except OSError as error:
            self.record(self.package_row("launch_failed", champion, error=repr(error)))
            self.log(f"package launch failed {describe(champion)}: {error!r}")
            return
        self.package_proc = (proc, champion)
        self.record(
            self.package_row(
                "launched", champion, command=command, pid=proc.pid, log=str(log)
            )
        )
        self.log(f"package launched pid={proc.pid} {describe(champion)}")

    def poll_package(self) -> None:
        """Record a finished package command and start the queued champion."""
        with self.lock:
            if self.package_proc is None:
                return
            proc, champion = self.package_proc
            code = proc.poll()
            if code is None:
                return
            self.package_proc = None
            self.record(self.package_row("finished", champion, exit_code=code))
            self.log(f"package finished exit={code} {describe(champion)}")
            queued, self.package_queued = self.package_queued, None
            if queued is not None:
                self.launch_package(queued)

    def drain_package(self) -> None:
        """Wait for the running package command (not after SIGTERM)."""
        while True:
            self.poll_package()
            if self.package_proc is None or self.terminating:
                break
            _ = self.wake.wait(self.poll_seconds)
        if self.package_proc is not None:
            self.log(f"package pid={self.package_proc[0].pid} left running")

    # ---- scheduling -------------------------------------------------------

    def _next_reserve_job(self) -> ReserveJob | None:
        with self.lock:
            if (
                self.advance_reason
                or self.reserve_running
                or self.champion is None
                or len(self.evals) - self.checked_evals < self.saturation_evals
            ):
                return None
            self.reserve_running = True
            return ReserveJob(
                self.generation, copy.deepcopy(self.champion), len(self.evals)
            )

    @staticmethod
    def dispatch_identity(job: LoopJob) -> str:
        """Give each dependency a stable ledger identity before native launch."""
        match job:
            case Job():
                return job.run_id
            case TeacherJob():
                return f"teacher-g{job.generation:03d}"
            case EnsembleJob():
                return ensemble_run_id(job.generation)
            case RoundEnsembleJob():
                return round_ensemble_run_id(job.generation, job.round_no)
            case ReserveJob():
                return (
                    f"reserve-{job.champion['run_id']}-"
                    f"{job.champion['checkpoint_name']}-n{job.test_evals}"
                )

    def _record_dispatch(
        self, reservation: Reservation, job: LoopJob, event: str, **extra: object
    ) -> None:
        self.record(
            {
                "type": "dispatch",
                "event": event,
                "generation": reservation.generation,
                "job": reservation.job_id,
                "kind": type(job).__name__,
                "slot": reservation.sequence,
                "worker": reservation.worker_id,
                "physical_device": reservation.device,
                "peak_usage": reservation.estimate.peak.values(),
                "retained_usage": reservation.estimate.retained.values(),
                "time": stamp(),
                **extra,
            }
        )

    def _execute_admitted(
        self, job: LoopJob, reservation: Reservation, completed: CompletionSlots
    ) -> None:
        pool = self.worker_pool
        if pool is None:
            message = "admitted execution requires its owned worker pool"
            raise RuntimeError(message)
        error: BaseException | None = None
        started = False
        cancelled = False
        try:
            self._execution.context = pool.context(reservation)
            started = True
            self.execute(job)
        except BaseException as caught:  # noqa: BLE001 - completion carries exact error
            cancelled = (
                not started
                and reservation not in pool.admission.active
                and (
                    self.stopping()
                    or reservation.generation != pool.admission.generation
                )
            )
            error = None if cancelled else caught
        finally:
            self._execution.context = None
            if isinstance(job, ReserveJob):
                self.reserve_running = False
            if started:
                run = self.runs.get(job.run_id) if isinstance(job, Job) else None
                failed = error is not None or (
                    isinstance(job, Job) and (run is None or run["status"] != "trained")
                )
                try:
                    pool.finish(reservation, failed=failed)
                except BaseException as caught:  # noqa: BLE001 - preserve cleanup failure
                    error = caught
            try:
                self._record_dispatch(
                    reservation,
                    job,
                    "cancelled"
                    if cancelled
                    else "finished"
                    if error is None
                    else "failed",
                    error=repr(error) if error is not None else None,
                )
            except BaseException as caught:  # noqa: BLE001 - never lose completion
                error = caught
            finally:
                # Resource ownership and native bookkeeping precede notification.
                completed.complete(reservation.sequence, error)

    def _launch_admitted(
        self,
        job: LoopJob,
        completed: CompletionSlots,
        threads: dict[int, threading.Thread],
    ) -> bool:
        pool, estimate = self.worker_pool, self.resource_estimator
        if pool is None or estimate is None:
            message = "admission requires its explicit estimator and worker owner"
            raise RuntimeError(message)
        reservation = pool.reserve(
            self.dispatch_identity(job), job.generation, estimate(job)
        )
        if reservation is None:
            return False
        completed.register(reservation.sequence)
        thread = threading.Thread(
            target=self._execute_admitted,
            args=(job, reservation, completed),
            daemon=True,
        )
        self._record_dispatch(reservation, job, "reserved")
        threads[reservation.sequence] = thread
        try:
            thread.start()
        except BaseException:
            del threads[reservation.sequence]
            completed.cancel(reservation.sequence)
            pool.cancel(reservation)
            self._record_dispatch(reservation, job, "launch_failed")
            raise
        return True

    def _run_admitted(self) -> int:  # noqa: C901, PLR0912, PLR0915 - owned dispatch lifecycle
        pool = self.worker_pool
        if pool is None:
            message = "admitted scheduling requires its owned worker pool"
            raise RuntimeError(message)
        if pool.admission.generation != self.generation:
            pool.advance_generation(self.generation)
        completed = CompletionSlots(self.wake)
        threads: dict[int, threading.Thread] = {}
        pending = self.jobs()
        job: LoopJob | None = None
        reserve: ReserveJob | None = None
        exhausted = False
        try:
            while not self.terminating:
                for finished in completed.arm():
                    threads.pop(finished.slot).join()
                    if finished.error is not None:
                        raise finished.error
                self.poll_package()
                if pool.admission.generation != self.generation:
                    pool.advance_generation(self.generation)
                stopping = self.stopping()
                if stopping:
                    pool.stop()
                if job and (self.advance_reason or job.generation != self.generation):
                    job = None
                if reserve and reserve.generation != self.generation:
                    self.reserve_running = False
                    reserve = None
                free_slot = len(threads) < self.parallel
                if not stopping and free_slot:
                    reserve = reserve or self._next_reserve_job()
                    if job is None and reserve is None and not exhausted:
                        job = next(pending, None)
                        exhausted = job is None
                    candidate = reserve or job
                    if candidate is not None and self._launch_admitted(
                        candidate, completed, threads
                    ):
                        if candidate is reserve:
                            reserve = None
                        else:
                            job = None
                        _ = completed.wait(self.launch_gap_seconds)
                        continue
                if (exhausted or stopping) and not threads:
                    break
                _ = completed.wait(self.poll_seconds)
        except BaseException:
            pool.terminate()
            raise
        finally:
            if self.terminating:
                pool.terminate()
            for thread in threads.values():
                thread.join()
            for reservation in pool.admission.active:
                pool.finish(reservation, failed=True)
            pool.close()
            self.reserve_running = False
        self.drain_package()
        return 0

    def resources_ok(self) -> bool:
        """Check GPU memory and cgroup task headroom before a launch."""
        free, pids = self.gpu_free(), self.pids()
        if free >= self.min_free_gib and pids < PIDS_LIMIT:
            self.waits = 0
            return True
        if self.waits % 20 == 0:
            self.log(f"waiting for resources: gpu_free={free:.1f}GiB pids={pids}")
        self.waits += 1
        return False

    def run(self) -> int:  # noqa: C901 - preserve baseline scheduling defaults
        """Keep up to `parallel` jobs running until STOP or SIGTERM."""
        if self.generation == 0:
            self.start_generation("initial")
        champion = self.champion
        if (
            champion
            and (champion["run_id"], champion["checkpoint_name"]) not in self.packaged
        ):
            self.request_package(champion)
        if self.worker_pool is not None:
            return self._run_admitted()
        pending = self.jobs()
        threads: list[threading.Thread] = []
        job: Job | TeacherJob | EnsembleJob | RoundEnsembleJob | None = None
        exhausted = False
        stop_logged = False
        try:
            while not self.terminating:
                self.poll_package()
                threads = [t for t in threads if t.is_alive()]
                free_slot = len(threads) < self.parallel
                if self.stopping() and not stop_logged:
                    self.log("stop requested; launching nothing new")
                    stop_logged = True
                if job and (self.advance_reason or job.generation != self.generation):
                    job = None
                if job is None and not exhausted and free_slot and not self.stopping():
                    job = next(pending, None)
                    exhausted = job is None
                if job and free_slot and not self.stopping() and self.resources_ok():
                    thread = threading.Thread(
                        target=self.execute, args=(job,), daemon=True
                    )
                    thread.start()
                    threads.append(thread)
                    job = None
                    _ = self.wake.wait(self.launch_gap_seconds)
                    continue
                if (exhausted or self.stopping()) and not threads:
                    break
                _ = self.wake.wait(self.poll_seconds)
        finally:
            for thread in threads:
                thread.join()
            self.close_training_worker()
        self.drain_package()
        return 0

    def close_training_worker(self) -> None:
        """Release the finite training subprocess after loop threads drain."""
        if self.training_runner is not None:
            self.training_runner.close()


def main(argv: Sequence[str] | None = None) -> int:  # noqa: PLR0915 - CLI dispatch policies
    """Parse arguments and run the loop."""
    parser = argparse.ArgumentParser(description=__doc__)
    _ = parser.add_argument("--state-dir", default="artifacts/champion-loop")
    _ = parser.add_argument(
        "--candidates", default="configs/research/champion/candidates.yaml"
    )
    _ = parser.add_argument("--parallel", type=int, default=5)
    _ = parser.add_argument("--base-seed", type=int, default=20261200)
    _ = parser.add_argument("--min-free-gib", type=float, default=30.0)
    _ = parser.add_argument("--saturation-evals", type=int, default=200)
    _ = parser.add_argument("--max-reserve-checks", type=int, default=3)
    _ = parser.add_argument(
        "--package-command", default=".venv/bin/python scripts/package_champion.py"
    )
    _ = parser.add_argument("--no-package", action="store_true")
    _ = parser.add_argument("--runs-root", default="artifacts/training-runs")
    _ = parser.add_argument("--data-root", default="data/20260912_for_training")
    _ = parser.add_argument("--poll-seconds", type=float, default=30.0)
    _ = parser.add_argument("--launch-gap-seconds", type=float, default=30.0)
    _ = parser.add_argument(
        "--advance-generation",
        metavar="REASON",
        help="open the next generation on start, e.g. after the corpus labels changed",
    )
    _ = parser.add_argument(
        "--teacher-after-round",
        type=int,
        default=0,
        metavar="N",
        help="build each generation's in-sample teacher once its ordinary runs of "
        "rounds <= N are final (0: never)",
    )
    _ = parser.add_argument(
        "--teacher-path",
        default=TEACHER_PATH,
        help="teacher file, relative to the repo root; {generation} is zero-padded",
    )
    _ = parser.add_argument(
        "--band-decisions",
        action="store_true",
        help="compare champions on reject-band decided counts (validation-chosen "
        "band) instead of boundary 0; start a new generation when switching",
    )
    _ = parser.add_argument(
        "--no-incremental-pareto",
        action="store_true",
        help="recompute the frontier from all evaluations on every append and replay",
    )
    _ = parser.add_argument(
        "--training-worker-max-jobs",
        type=int,
        default=0,
        help="opt in to one finite preparation-reusing training worker",
    )
    _ = parser.add_argument("--training-worker-idle-seconds", type=float, default=30.0)
    _ = parser.add_argument("--reuse-teacher-preparation", action="store_true")
    _ = parser.add_argument("--score-cache-dir", type=Path)
    _ = parser.add_argument("--resource-worker-pool", action="store_true")
    _ = parser.add_argument("--worker-devices", default="0")
    _ = parser.add_argument("--worker-cuda-context-gib", type=float)
    options = parser.parse_args(argv)
    runner = SubprocessRunner()
    loop = Loop(
        state_dir=REPO / options.state_dir,
        candidates_path=REPO / options.candidates,
        runs_root=REPO / options.runs_root,
        data_root=Path(options.data_root),
        parallel=options.parallel,
        base_seed=options.base_seed,
        min_free_gib=options.min_free_gib,
        saturation_evals=options.saturation_evals,
        max_reserve_checks=options.max_reserve_checks,
        runner=runner,
        package_command=(
            None if options.no_package else shlex.split(options.package_command)
        ),
        package_launcher=lambda command, log: popen_detached(command, log, runner.env),
        poll_seconds=options.poll_seconds,
        launch_gap_seconds=options.launch_gap_seconds,
        teacher_after_round=options.teacher_after_round,
        teacher_path=options.teacher_path,
        band_decisions=options.band_decisions,
        incremental_pareto=not options.no_incremental_pareto,
        score_cache_dir=options.score_cache_dir,
        reuse_teacher_preparation=options.reuse_teacher_preparation,
    )
    if options.resource_worker_pool:
        from scripts.champion_admission import (  # noqa: PLC0415 - explicit opt-in path
            ChampionResourceEstimator,
        )

        try:
            devices = tuple(int(item) for item in options.worker_devices.split(","))
        except ValueError:
            parser.error("worker devices must name physical GPU0 or3")
        if (
            options.training_worker_max_jobs <= 0
            or options.worker_cuda_context_gib is None
            or not math.isfinite(options.worker_cuda_context_gib)
            or options.worker_cuda_context_gib < 0
            or options.parallel <= 0
            or not devices
            or len(set(devices)) != len(devices)
            or any(device not in (0, 3) for device in devices)
        ):
            parser.error(
                "resource pool requires finite jobs, valid devices and CUDA allowance"
            )
        admission = AdmissionPool(
            read_resource_snapshot,
            parallel=min(options.parallel, len(devices)),
            generation=loop.generation,
            allowed_devices=devices,
        )
        loop.worker_pool = TrainingWorkerPool(
            admission,
            REPO,
            PYTHON,
            runner.env,
            max_jobs=options.training_worker_max_jobs,
            idle_seconds=options.training_worker_idle_seconds,
            stopping=loop.stopping,
        )
        loop.resource_estimator = ChampionResourceEstimator(
            loop, cuda_context_bytes=int(options.worker_cuda_context_gib * 1024**3)
        )
    elif options.training_worker_max_jobs:
        loop.training_runner = FiniteTrainingRunner(
            PYTHON,
            REPO,
            runner.env,
            max_jobs=options.training_worker_max_jobs,
            idle_seconds=options.training_worker_idle_seconds,
            stopping=loop.stopping,
        )

    def on_sigterm(_signum: int, _frame: object) -> None:
        loop.terminate()
        runner.terminate()
        if loop.training_runner is not None:
            loop.training_runner.terminate()

    _ = signal.signal(signal.SIGTERM, on_sigterm)
    if options.advance_generation:
        loop.request_advance(options.advance_generation)
    return loop.run()


if __name__ == "__main__":
    sys.exit(main())
