| """Validator loop — manifest → eval → KOTH decision → weights. |
| |
| The validator never trains. Each round it: |
| |
| 1. Reads the current :class:`TrainingManifest` from the owner dataset repo and |
| verifies its signature + that king and challenger share the contract digest |
| (the controlled-experiment guarantee). |
| 2. Pulls the king's and challenger's trained checkpoints and scores both on the |
| *same* held-out eval windows. |
| 3. Runs the paired-bootstrap KOTH verdict and folds it into the sticky |
| champion state (``dethrone_cp`` consecutive wins to take the throne; |
| ``dethrone_cp = 1`` makes it single-round). |
| 4. Sets weights: equal share across the current king plus up to |
| ``[scoring] reward_prior_kings`` registered prior kings (teutonic-style), |
| collapsing to winner-take-all when ``reward_prior_kings = 0``. |
| |
| The pure orchestration in :meth:`ValidatorRunner.process_round` is testable by |
| injecting ``evaluate_fn`` and ``windows``; HF + torch + chain are isolated |
| behind the defaults. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import contextlib |
| import hashlib |
| import json |
| import logging |
| from collections.abc import Callable |
| from dataclasses import dataclass, field, replace |
| from datetime import UTC |
| from pathlib import Path |
| from typing import TYPE_CHECKING |
|
|
| from ..eval.koth import RoundResult, evaluate_round |
| from ..eval.scoring import WindowScore |
| from ..eval.window import EvalWindow |
| from ..shared.config import ChainConfig |
| from ..shared.manifest import ( |
| TrainedEntry, |
| TrainingManifest, |
| contract_digest, |
| parse_trained_pointer, |
| verify_signature, |
| ) |
| from ..shared.receipt import ( |
| EntryScores, |
| EvalContext, |
| Participant, |
| RoundReceipt, |
| VerdictRecord, |
| WindowScoreRecord, |
| build_receipt, |
| ) |
| from . import state as state_mod |
| from .state import ChampionState, StateTransition |
|
|
| if TYPE_CHECKING: |
| from ..trainer.remote import RemoteHost |
| from .cascade import CascadeController |
|
|
| log = logging.getLogger("cascade.validator") |
|
|
| |
| EvaluateFn = Callable[[TrainedEntry, list[EvalWindow]], list[WindowScore]] |
| |
| |
| |
| GiftRowsFn = Callable[[TrainedEntry], dict | None] |
|
|
|
|
| def participants_from_commitments(commitments: list, cutoff_block: int, |
| floor_block: int = 0) -> tuple[Participant, ...]: |
| """The round's eligible participant set, for the public receipt. |
| |
| Mirrors the trainer's eligibility rule (``trainer.loop.resolve_commitments``): |
| parseable generator pointers revealed STRICTLY BEFORE the epoch boundary and |
| at/after the go-live ``floor_block``, latest commit per hotkey among the |
| eligible ones — but keeps ``commit_block`` so an auditor can re-check every |
| entrant against the cutoff. Sorted by UID for a deterministic receipt body. |
| |
| ``commitments`` should carry each hotkey's FULL reveal history |
| (``poll_commitments(include_history=True)``): with only the latest reveal, a |
| miner who re-committed after the boundary vanishes from the record even |
| though their eligible pre-cutoff entry fielded the round. |
| """ |
| from ..interface.validation import parse_commit |
|
|
| best: dict[str, Participant] = {} |
| for c in commitments: |
| if floor_block and c.commit_block < floor_block: |
| continue |
| if c.commit_block >= cutoff_block: |
| continue |
| parsed = parse_commit(c.payload) |
| if parsed is None: |
| continue |
| prev = best.get(c.hotkey) |
| if prev is None or c.commit_block >= prev.commit_block: |
| best[c.hotkey] = Participant( |
| hotkey=c.hotkey, uid=c.uid, gen_ref=parsed.ref, commit_block=c.commit_block |
| ) |
| return tuple(sorted(best.values(), key=lambda p: p.uid)) |
|
|
|
|
| @dataclass(frozen=True) |
| class RoundOutcome: |
| result: RoundResult |
| transition: StateTransition |
| |
| |
| |
| entry_scores: tuple[EntryScores, ...] = () |
| |
| |
| king_tenure_rounds: int = 0 |
|
|
|
|
| |
| |
| |
| |
| |
| |
| POOL_PIN_READ_GRACE_SECONDS = 1800.0 |
|
|
|
|
| @dataclass |
| class ValidatorRunner: |
| cfg: ChainConfig |
| state: ChampionState = field(default_factory=ChampionState) |
| evaluate_fn: EvaluateFn | None = None |
| gift_rows_fn: GiftRowsFn | None = None |
| cache_dir: Path | None = None |
| device: str = "cpu" |
| |
| |
| |
| |
| |
| |
| |
| eval_host_fn: Callable[[], RemoteHost | None] | None = None |
| verify_signatures: bool = True |
| |
| |
| |
| |
| |
| cascade: CascadeController | None = None |
| |
| |
| |
| |
| _last_weight_block: int | None = None |
| |
| |
| |
| |
| |
| _pin_read_first_failure: dict[str, float] = field(default_factory=dict, repr=False) |
|
|
| |
|
|
| def check_manifest(self, manifest: TrainingManifest) -> str | None: |
| """Return a rejection reason string, or None if the manifest is usable. |
| |
| Enforces (1) the trainer-hotkey signature, (2) the contract-digest match |
| (king and challenger trained under the same terms), and (3) that the |
| manifest targets our configured base architecture and eval dataset. |
| """ |
| if self.verify_signatures and not verify_signature(manifest, self.cfg.manifest.trainer_hotkey): |
| return "signature_invalid" |
| want_contract = contract_digest(self.cfg.training) |
| if manifest.contract_digest != want_contract: |
| return f"contract_digest_mismatch: {manifest.contract_digest} != {want_contract}" |
| if manifest.base_arch_digest != self.cfg.training.base_arch_digest: |
| return "base_arch_digest_mismatch" |
| if manifest.eval_dataset != self.cfg.eval.eval_dataset: |
| return "eval_dataset_mismatch" |
| gpu_reason = self._check_gpu(manifest) |
| if gpu_reason is not None: |
| return gpu_reason |
| ws_reason = self._check_warm_start(manifest) |
| if ws_reason is not None: |
| return ws_reason |
| return None |
|
|
| def _check_warm_start(self, manifest: TrainingManifest) -> str | None: |
| """Warm-start pin gate (Cascade, DEC-CA-0005): the manifest's signed |
| ``warm_start_ckpt`` must equal the init THIS validator's deterministic |
| promotion installed (its own ``warm_start_init_path`` file; "" before any |
| promotion). Every validator computes the same promotion (block-anchored |
| clock + trainer-signed bench scores), so agreement is fleet-wide. A |
| mismatch — trainer trained from random when a promotion is live, or from |
| a stale/foreign init — rejects the round rather than silently scoring |
| runs trained off-baseline; the trainer re-syncs by the next round. Only |
| enforced when Cascade is wired (off ⇒ pure KOTH, field ignored).""" |
| if self.cascade is None: |
| return None |
| expected = "" |
| p = Path(self.cfg.validator.warm_start_init_path) |
| if p.is_file(): |
| try: |
| expected = str(json.loads(p.read_text(encoding="utf-8")).get("checkpoint_id") or "") |
| except Exception as e: |
| return f"warm_start_state_unreadable: {p}: {e}" |
| if manifest.warm_start_ckpt != expected: |
| return (f"warm_start_mismatch: manifest trained from " |
| f"{manifest.warm_start_ckpt or '<random init>'!r}, this validator " |
| f"expects {expected or '<random init>'!r}") |
| return None |
|
|
| @staticmethod |
| def check_pool_pin( |
| manifest: TrainingManifest, window_source: object, *, block: int | None |
| ) -> str | None: |
| """Verify this round's eval pool against the trainer-signed pin. |
| |
| A pinned manifest carries the ``(key, sha256)`` of the snapshot the |
| trainer screened on; the validator's own deterministic selection (same |
| epoch block, same rule) must resolve to the identical pair. The pin is |
| inside the signed body, so pool integrity descends from the trainer |
| signature rather than the unsigned ``pool/index.json`` — a poisoned |
| index or tampered tar surfaces here as a loud reject instead of scoring |
| on attacker-chosen data. Unpinned manifests (older trainers) keep the |
| legacy index-trust behaviour. Returns a reject reason or ``None``. |
| |
| Raises :class:`~cascade.shared.hippius.StorageError` when the pool |
| index could not be READ (auth/network/5xx) — that is a transient, not a |
| verdict, and the caller decides retry-vs-reject (the live loop retries |
| within :data:`POOL_PIN_READ_GRACE_SECONDS`). Every other failure is a |
| reject reason. |
| """ |
| if not (manifest.eval_pool_key and manifest.eval_pool_sha256): |
| return None |
| prov_fn = getattr(window_source, "provenance_for_round", None) |
| if prov_fn is None: |
| return ("pool_pin_unverifiable: manifest pins the eval pool but this " |
| "validator's pool source reports no provenance") |
| try: |
| key, sha = prov_fn(int(manifest.round_id), block=block) |
| except Exception as e: |
| from ..shared.hippius import StorageError |
|
|
| if isinstance(e, StorageError): |
| raise |
| return f"pool_pin_unverifiable: provenance lookup failed: {e}" |
| if not key or not sha: |
| return ("pool_pin_unverifiable: manifest pins the eval pool but this " |
| "validator resolved no snapshot for the round") |
| if (key, sha) != (manifest.eval_pool_key, manifest.eval_pool_sha256): |
| return (f"pool_pin_mismatch: manifest signed {manifest.eval_pool_key}@" |
| f"{manifest.eval_pool_sha256[:12]}…, this validator resolved " |
| f"{key}@{sha[:12]}…") |
| return None |
|
|
| def _pool_pin_read_failed( |
| self, round_id: str, err: Exception, *, now: float | None = None |
| ) -> str | None: |
| """Grace bookkeeping for an UNREADABLE pool index at the pin gate. |
| |
| Returns ``None`` while ``round_id``'s read failures span less than |
| :data:`POOL_PIN_READ_GRACE_SECONDS` — the caller then skips the cycle |
| with no latch and no receipt, so the next manifest poll re-attempts the |
| round from scratch. Once the grace expires, returns the terminal reject |
| reason (and forgets the round, so a later re-publish starts a fresh |
| grace window). |
| """ |
| import time |
|
|
| now = time.monotonic() if now is None else now |
| first = self._pin_read_first_failure.setdefault(str(round_id), now) |
| waited = now - first |
| if waited < POOL_PIN_READ_GRACE_SECONDS: |
| log.warning( |
| "pool index unreadable at pin gate for round=%s (%s); retrying " |
| "next poll (%.0fs into %.0fs grace, no reject latched)", |
| round_id, err, waited, POOL_PIN_READ_GRACE_SECONDS, |
| ) |
| return None |
| self._pin_read_first_failure.pop(str(round_id), None) |
| return ( |
| f"pool_pin_unverifiable: provenance lookup failed persistently " |
| f"({waited:.0f}s > {POOL_PIN_READ_GRACE_SECONDS:.0f}s grace): {err}" |
| ) |
|
|
| def _check_gpu(self, manifest: TrainingManifest) -> str | None: |
| """Matched-hardware gate for byte-exact re-derivation. |
| |
| If ``[training] expected_gpu`` is pinned, every entry must report that GPU. |
| Otherwise require only that king and challenger ran the same GPU (when both |
| report one) — equal compute is already guaranteed by the token budget, but |
| a byte-exact audit needs the comparison run on one SKU. |
| """ |
| pinned = self.cfg.training.expected_gpu |
| gpus = {e.gpu_name for e in manifest.entries if e.gpu_name} |
| if pinned: |
| bad = sorted(g for g in gpus if g != pinned) |
| if bad or any(not e.gpu_name for e in manifest.entries): |
| return f"gpu_mismatch: expected {pinned!r}, manifest has {sorted(gpus)!r}" |
| elif len(gpus) > 1: |
| return f"gpu_mismatch: king/challenger on different GPUs {sorted(gpus)!r}" |
| return None |
|
|
| |
|
|
| def _fetch_checkpoint_dir(self, entry: TrainedEntry) -> Path: |
| """Fetch a trained checkpoint from the Hippius Hub registry to a local |
| dir and return it. The OCI digest in the ref pins the bytes, so the |
| fetch is self-verifying; repeated fetches of the same ref land in the |
| same digest-named dir (cheap to reuse).""" |
| from ..shared.hippius import HubConfig, HubRef, fetch_from_hub |
|
|
| ref = parse_trained_pointer(entry.trained_pointer) |
| if ref is None: |
| raise ValueError(f"malformed trained_pointer: {entry.trained_pointer!r}") |
| hub = HubConfig.from_storage(self.cfg.storage) |
| dest = Path(self.cache_dir or "./_eval_ckpts") / HubRef.parse(ref).digest.replace(":", "-") |
| fetch_from_hub(ref, dest, hub) |
| return dest |
|
|
| def _evaluate(self, entry: TrainedEntry, windows: list[EvalWindow]) -> list[WindowScore]: |
| if self.evaluate_fn is not None: |
| return self.evaluate_fn(entry, windows) |
| |
| |
| from .evaluator import evaluate_checkpoint |
|
|
| dest = self._fetch_checkpoint_dir(entry) |
| return evaluate_checkpoint( |
| dest, windows, num_samples=self.cfg.eval.num_samples, device=self.device |
| ) |
|
|
| |
|
|
| def _eval_host(self) -> RemoteHost | None: |
| """The offload pod for THIS eval — resolved fresh every time so an |
| elastic pod (rented per round manifest, torn down on the receipt) |
| appears and disappears without a validator restart.""" |
| return self.eval_host_fn() if self.eval_host_fn is not None else None |
|
|
| def _gift_rows(self, entry: TrainedEntry) -> dict | None: |
| """Gift-eval ratio rows for one entry — injected in tests, else the |
| sidecar bridge on the fetched checkpoint dir.""" |
| if self.gift_rows_fn is not None: |
| return self.gift_rows_fn(entry) |
|
|
| ec = self.cfg.eval |
| dest = self._fetch_checkpoint_dir(entry) |
| num_samples = ec.gift_gate_num_samples or ec.num_samples |
| eval_host = self._eval_host() |
| if eval_host is not None: |
| |
| |
| from .eval_offload import gift_rows_via_host |
|
|
| return gift_rows_via_host( |
| eval_host, dest, |
| datasets=ec.gift_gate_datasets, |
| num_samples=num_samples, |
| data_dir=(ec.gift_gate_data_dir or None), |
| device="cuda", |
| timeout_s=ec.gift_gate_timeout_s, |
| ) |
| from ..eval.benchmarks import run_gift_rows |
|
|
| return run_gift_rows( |
| dest, |
| project_dir=ec.benchmark_project_dir, |
| datasets=ec.gift_gate_datasets, |
| num_samples=num_samples, |
| device=self.device, |
| data_dir=(ec.gift_gate_data_dir or None), |
| timeout_s=ec.gift_gate_timeout_s, |
| ) |
|
|
| def _run_gift_gate( |
| self, |
| result: RoundResult, |
| king_entry: TrainedEntry, |
| chal_entry: TrainedEntry, |
| *, |
| seed: int | str, |
| round_id: str, |
| ) -> RoundResult: |
| """Fold the public-benchmark gate into a *winning* round result. |
| |
| Scores both sides on gift-eval (via the sidecar bridge) and runs the |
| paired no-regression bootstrap. The gate is uncomputable — and the |
| round therefore inconclusive under ``enforce`` — when either sidecar run |
| fails, gift-eval was skipped/errored, or the two runs scored against |
| different pinned data revisions (a consensus-safety check: king and |
| challenger must be judged on identical public data). |
| """ |
| from ..eval.gift_gate import evaluate_gift_gate, uncomputable_gate |
| from ..eval.koth import apply_gift_gate |
|
|
| p = self.cfg.koth_params() |
| mode = p.gift_gate_mode |
| king_run = self._gift_rows(king_entry) |
| chal_run = self._gift_rows(chal_entry) |
| if ( |
| king_run is None or chal_run is None |
| or king_run.get("status") != "ok" or chal_run.get("status") != "ok" |
| ): |
| gate = uncomputable_gate(p.gift_gate_tolerance, "gift-eval sidecar unavailable/errored") |
| elif king_run.get("revision") != chal_run.get("revision"): |
| gate = uncomputable_gate( |
| p.gift_gate_tolerance, |
| f"data-revision mismatch: king {king_run.get('revision')} != " |
| f"chal {chal_run.get('revision')}", |
| ) |
| else: |
| gate = evaluate_gift_gate( |
| king_run["rows"], chal_run["rows"], |
| tolerance=p.gift_gate_tolerance, |
| alpha=p.bootstrap_alpha, |
| B=p.bootstrap_B, |
| seed=seed, |
| min_configs=p.gift_gate_min_configs, |
| ) |
| log.info( |
| "gift-gate round=%s mode=%s computed=%s passed=%s lcb=%s tol=%.4f " |
| "n_configs=%d king_agg=%.5f chal_agg=%.5f%s", |
| round_id, mode, gate.computed, gate.passed, |
| f"{gate.lcb:.5f}" if gate.computed else "n/a", gate.tolerance, |
| gate.n_configs, gate.king_agg, gate.chal_agg, |
| "" if gate.computed else f" reason={gate.reason!r}", |
| ) |
| return apply_gift_gate(result, gate, mode=mode) |
|
|
| def _maybe_run_benchmarks( |
| self, manifest: TrainingManifest, outcome: RoundOutcome | None |
| ) -> None: |
| """Log public-benchmark numbers for a newly crowned king (log-only). |
| |
| Best-effort and strictly off the consensus path: it runs only when a |
| challenger just dethroned the king, scores that new king's checkpoint via |
| the isolated sidecar, and logs whatever comes back. Any failure is |
| swallowed — a benchmark hiccup must never disturb weights or KOTH state. |
| """ |
| ec = self.cfg.eval |
| if not ec.run_benchmarks: |
| return |
| if outcome is None or not outcome.transition.dethroned: |
| return |
| new_king = manifest.entry_for_role("challenger") |
| if new_king is None: |
| return |
| try: |
| from ..eval.benchmarks import format_report, run_benchmarks |
|
|
| ckpt = self._fetch_checkpoint_dir(new_king) |
| report = run_benchmarks( |
| ckpt, |
| project_dir=ec.benchmark_project_dir, |
| suites=ec.benchmark_suites or ("gift-eval", "boom", "time"), |
| num_samples=ec.benchmark_num_samples or ec.num_samples, |
| max_series=ec.benchmark_max_series, |
| device=self.device, |
| ) |
| if report is not None: |
| log.info( |
| "benchmarks round=%s king=%s %s", |
| manifest.round_id, self.state.king_hotkey, format_report(report), |
| ) |
| except Exception as e: |
| log.warning("benchmark hook failed for round=%s: %s", manifest.round_id, e) |
|
|
| |
|
|
| def _current_king_entry(self, manifest: TrainingManifest) -> TrainedEntry | None: |
| """The manifest checkpoint the reigning champion produced this round. |
| |
| Cascade times the *validator's champion*, not the manifest's (lagging) |
| king role — so the checkpoint to score is the entry whose miner hotkey is |
| the champion. Prefers the primary throne size (what the benchmark sidecar |
| scores) and falls back to any size that hotkey trained.""" |
| hk = self.state.king_hotkey |
| if hk is None: |
| return None |
| matches = [e for e in manifest.entries if e.miner_hotkey == hk] |
| if not matches: |
| return None |
| primary = self.cfg.throne_contracts()[0].arch_preset |
| return next((e for e in matches if e.size == primary), matches[0]) |
|
|
| @staticmethod |
| def _bench_scores_dict(entry: TrainedEntry) -> dict | None: |
| """The six Cascade numbers off a manifest entry's trainer-signed |
| ``bench_scores`` (GIFT-Eval / BOOM / TIME CRPS+MASE), or ``None`` when the |
| entry carries none. This is the authoritative, consensus-safe source: every |
| validator reads the identical signed numbers.""" |
| bs = entry.bench_scores |
| if bs is None: |
| return None |
| return { |
| "gifteval_crps": bs.gifteval_crps, "gifteval_mase": bs.gifteval_mase, |
| "boom_crps": bs.boom_crps, "boom_mase": bs.boom_mase, |
| "time_crps": bs.time_crps, "time_mase": bs.time_mase, |
| } |
|
|
| def _bench_metrics_via_sidecar(self, entry: TrainedEntry) -> dict | None: |
| """Fallback: score one checkpoint on GIFT-Eval, BOOM, and TIME via the |
| out-of-process sidecar, returning the six numbers or ``None`` when any suite |
| is missing/errored. Used only when the manifest carries no ``bench_scores`` |
| (e.g. a trainer that predates the Cascade hook). NOTE: independently-run GPU |
| sweeps are not bit-reproducible, so this path is not consensus-safe across |
| validators — prefer the trainer-signed numbers.""" |
| from ..eval.benchmarks import extract_bench_scores, run_benchmarks |
|
|
| ec = self.cfg.eval |
| ckpt = self._fetch_checkpoint_dir(entry) |
| num_samples = ec.benchmark_num_samples or ec.num_samples |
| eval_host = self._eval_host() |
| if eval_host is not None: |
| |
| |
| from .eval_offload import bench_scores_via_host |
|
|
| metrics = bench_scores_via_host( |
| eval_host, ckpt, |
| num_samples=num_samples, |
| max_series=ec.cascade_bench_max_series, |
| data_dir=(ec.gift_gate_data_dir or None), |
| device="cuda", |
| timeout_s=ec.gift_gate_timeout_s, |
| ) |
| else: |
| report = run_benchmarks( |
| ckpt, |
| project_dir=ec.benchmark_project_dir, |
| suites=("gift-eval", "boom", "time"), |
| num_samples=num_samples, |
| max_series=ec.cascade_bench_max_series, |
| device=self.device, |
| ) |
| metrics = extract_bench_scores(report) |
| if metrics is None: |
| log.warning("cascade: incomplete GIFT-Eval/BOOM/TIME metrics for king checkpoint %s; " |
| "not recording this round", entry.trained_pointer) |
| return metrics |
|
|
| def _record_king_checkpoint( |
| self, manifest: TrainingManifest, now: float |
| ) -> None: |
| """Add the reigning king's checkpoint to the reign log so a later Cascade |
| selection is a lookup, not a re-eval. Prefers the trainer's signed |
| ``bench_scores`` on the manifest (consensus-safe); falls back to scoring via |
| the local sidecar only when the manifest carries none. Best-effort: a miss |
| just means this round's checkpoint isn't a promotion candidate.""" |
| if self.cascade is None or self.state.king_hotkey is None: |
| return |
| entry = self._current_king_entry(manifest) |
| if entry is None: |
| return |
| metrics = self._bench_scores_dict(entry) or self._bench_metrics_via_sidecar(entry) |
| if metrics is None: |
| return |
| self.cascade.record_checkpoint(entry.trained_pointer, now=now, size=entry.size, **metrics) |
|
|
| def _cascade_round( |
| self, manifest: TrainingManifest, outcome: RoundOutcome | None |
| ) -> None: |
| """One Cascade step, run at the end of a round (after weights/receipts). |
| Resets the reign clock on a dethrone, records the reigning king's |
| checkpoint, then checks the clock — a fired Cascade installs the promoted |
| init and re-crowns the SAME king (DEC-CA-0004: the champion throne is |
| never touched). Fully guarded: Cascade never disturbs KOTH.""" |
| if self.cascade is None: |
| return |
| import time |
|
|
| now = time.time() |
| try: |
| |
| |
| block = self._epoch_start_block(manifest) |
| |
| |
| if outcome is not None and outcome.transition.dethroned and outcome.transition.new_king_hotkey: |
| self.cascade.note_dethrone(outcome.transition.new_king_hotkey, block=block) |
| elif self.cascade.state.king_hotkey is None and self.state.king_hotkey is not None: |
| self.cascade.note_dethrone(self.state.king_hotkey, block=block) |
| self._record_king_checkpoint(manifest, now) |
| event = self.cascade.cascade_check(block=block, now=now) |
| if event is not None: |
| self._apply_cascade(event) |
| except Exception as e: |
| log.warning("cascade step failed for round=%s: %s", manifest.round_id, e) |
|
|
| def _apply_cascade(self, event: object) -> None: |
| """Log a fired Cascade. The champion throne is deliberately untouched |
| (DEC-CA-0004): the king persists — vacating had no benefit (both roles |
| train from the shared init) and a vacant throne refillable only via the |
| dethrone branch froze the reign clock when the incumbent kept winning.""" |
| winner = getattr(event, "winner", None) |
| king = getattr(event, "old_king", None) |
| log.info( |
| "cascade: promotion installed (king %s persists); field trains from " |
| "checkpoint %s next round", |
| (king or "?")[:12], |
| getattr(winner, "checkpoint_id", "?"), |
| ) |
|
|
| def process_round( |
| self, |
| manifest: TrainingManifest, |
| windows: list[EvalWindow], |
| base_seed: int | str, |
| ) -> RoundOutcome | None: |
| """Evaluate one manifest against the eval windows and update state. |
| |
| A round carries one (king, challenger) pair PER trained size (the primary |
| plus any ``[[training.sizes]]``). Each size's pair is scored on the SAME |
| windows, then the per-size scores are POOLED — king's across sizes vs |
| challenger's across sizes, in identical order — and a single paired |
| bootstrap decides ONE throne on the combined score (scaling-aware KOTH). |
| Pooling preserves pairing because each size's king and challenger share |
| the window ``abs_target``. |
| |
| Returns None (king holds, no state change) when the manifest carries no |
| size with both a king and a challenger, or fails the contract gate. |
| Otherwise returns the round outcome with the (already-applied) transition. |
| """ |
| reason = self.check_manifest(manifest) |
| if reason is not None: |
| log.warning("rejecting manifest round=%s: %s", manifest.round_id, reason) |
| return None |
|
|
| king_by_size = {e.size: e for e in manifest.entries_for_role("king")} |
| chal_by_size = {e.size: e for e in manifest.entries_for_role("challenger")} |
| paired_sizes = [s for s in manifest.sizes() if s in king_by_size and s in chal_by_size] |
| if not paired_sizes: |
| log.info("manifest round=%s has no king/challenger pair; king holds", manifest.round_id) |
| return None |
|
|
| king_scores: list[WindowScore] = [] |
| chal_scores: list[WindowScore] = [] |
| score_records: list[EntryScores] = [] |
| for size in paired_sizes: |
| import time as _time |
|
|
| _t0 = _time.perf_counter() |
| ks = self._evaluate(king_by_size[size], windows) |
| _t_king = _time.perf_counter() - _t0 |
| _t1 = _time.perf_counter() |
| cs = self._evaluate(chal_by_size[size], windows) |
| _t_chal = _time.perf_counter() - _t1 |
| log.info( |
| "round=%s eval-timing size=%s device=%s n_windows=%d num_samples=%d " |
| "king=%.1fs challenger=%.1fs total=%.1fs", |
| manifest.round_id, size, self.device, len(windows), |
| self.cfg.eval.num_samples, _t_king, _t_chal, _t_king + _t_chal, |
| ) |
| king_scores += ks |
| chal_scores += cs |
| for entry, scores in ((king_by_size[size], ks), (chal_by_size[size], cs)): |
| score_records.append(EntryScores( |
| role=entry.role, size=size, |
| hotkey=entry.miner_hotkey, uid=entry.miner_uid, |
| scores=tuple(WindowScoreRecord.from_score(s) for s in scores), |
| )) |
| |
| |
| chal_entry = chal_by_size[paired_sizes[0]] |
|
|
| tenure_at_decision = self.state.tenure_rounds |
| result = evaluate_round( |
| king_scores, |
| chal_scores, |
| self.cfg.koth_params(), |
| seed=base_seed, |
| king_tenure_rounds=tenure_at_decision, |
| ) |
| |
| |
| |
| |
| if self.cfg.scoring.gift_gate_mode != "off" and result.challenger_wins_round: |
| gate_size = ( |
| self.cfg.training.arch_preset |
| if self.cfg.training.arch_preset in king_by_size and |
| self.cfg.training.arch_preset in chal_by_size |
| else paired_sizes[0] |
| ) |
| result = self._run_gift_gate( |
| result, king_by_size[gate_size], chal_by_size[gate_size], |
| seed=base_seed, round_id=manifest.round_id, |
| ) |
| transition = state_mod.apply_round( |
| self.state, |
| challenger_hotkey=chal_entry.miner_hotkey, |
| challenger_uid=chal_entry.miner_uid, |
| result=result, |
| dethrone_cp=self.cfg.scoring.dethrone_cp, |
| keep_former_kings=self.cfg.scoring.reward_prior_kings, |
| ) |
| self.state = transition.state |
| log.info( |
| "round=%s lcb=%.4f margin=%.4f win=%s %s king=%s tenure=%d", |
| manifest.round_id, result.lcb, result.margin, result.challenger_wins_round, |
| transition.note, self.state.king_hotkey, self.state.tenure_rounds, |
| ) |
| |
| |
| |
| if result.win_rate is not None: |
| log.info( |
| "round=%s diag n_clusters=%d win_rate=%.3f wilcoxon_p=%s per_domain=%s", |
| manifest.round_id, result.n_clusters, result.win_rate, |
| f"{result.wilcoxon_p:.4g}" if result.wilcoxon_p is not None else "n/a", |
| {d: f"{wr:.2f}/n{n}" for d, (wr, n) in (result.per_domain_win_rate or {}).items()}, |
| ) |
| return RoundOutcome( |
| result=result, transition=transition, entry_scores=tuple(score_records), |
| king_tenure_rounds=tenure_at_decision, |
| ) |
|
|
| def _epoch_start_block(self, manifest: TrainingManifest) -> int: |
| """The round's epoch-boundary block: ``created_block`` floored to the |
| epoch grid. Monotonic and identical for every validator (from the shared |
| manifest), so it is the consensus key for daily eval-pool snapshot |
| selection — unlike the round id, which is a block *hash* (non-monotonic). |
| """ |
| epoch_blocks = max(1, self.cfg.round.epoch_blocks) |
| return (manifest.created_block // epoch_blocks) * epoch_blocks |
|
|
| |
|
|
| def build_round_receipt( |
| self, |
| manifest: TrainingManifest, |
| *, |
| base_seed: int, |
| epoch_start_block: int, |
| epoch_block_hash: str, |
| outcome: RoundOutcome | None = None, |
| windows: list[EvalWindow] | None = None, |
| reject_reason: str | None = None, |
| participants: tuple[Participant, ...] = (), |
| pool_provenance: tuple[str, str] = ("", ""), |
| reward_uids: tuple[int, ...] = (), |
| weights: tuple[float, ...] = (), |
| validator_hotkey: str = "", |
| ) -> RoundReceipt: |
| """Assemble the round's public receipt (pure — no I/O, no signing). |
| |
| A gated-out manifest — or one with no (king, challenger) pair to score — |
| yields a ``rejected`` receipt carrying the reason; a scored round yields |
| the full record: chain context, embedded manifest, participant set, the |
| eval slice, every per-window score, the verdict, and the weight vector. |
| """ |
| from ..trainer.contract import RoundSeeds |
|
|
| seeds = RoundSeeds.derive(base_seed, self.cfg.training) |
| if reject_reason is not None: |
| return build_receipt( |
| round_id=manifest.round_id, status="rejected", |
| epoch_start_block=epoch_start_block, epoch_block_hash=epoch_block_hash, |
| base_seed=base_seed, seeds=seeds, manifest=manifest, |
| participants=participants, reject_reason=reject_reason, |
| reward_uids=reward_uids, weights=weights, |
| validator_hotkey=validator_hotkey, |
| ) |
| if outcome is None or windows is None: |
| raise ValueError("a scored receipt needs both outcome and windows") |
| eval_context = EvalContext( |
| pool_ref=pool_provenance[0], |
| pool_digest=pool_provenance[1], |
| window_ids=tuple(w.series_id for w in windows), |
| n_windows=len(windows), |
| num_samples=self.cfg.eval.num_samples, |
| ) |
| verdict = VerdictRecord.from_round( |
| outcome.result, outcome.transition, |
| params=self.cfg.koth_params(), bootstrap_seed=base_seed, |
| king_tenure_rounds=outcome.king_tenure_rounds, |
| ) |
| return build_receipt( |
| round_id=manifest.round_id, status="scored", |
| epoch_start_block=epoch_start_block, epoch_block_hash=epoch_block_hash, |
| base_seed=base_seed, seeds=seeds, manifest=manifest, |
| participants=participants, eval_context=eval_context, |
| entry_scores=outcome.entry_scores, verdict=verdict, |
| reward_uids=reward_uids, weights=weights, |
| validator_hotkey=validator_hotkey, |
| ) |
|
|
| def _publish_round_receipt( |
| self, |
| client: object, |
| manifest: TrainingManifest, |
| base_seed: int, |
| *, |
| outcome: RoundOutcome | None = None, |
| windows: list[EvalWindow] | None = None, |
| reject_reason: str | None = None, |
| window_source: object = None, |
| reward_uids: tuple[int, ...] = (), |
| weights: tuple[float, ...] = (), |
| ) -> None: |
| """Gather chain context, sign, and publish the round receipt. |
| |
| Best-effort by design: a receipt failure must never disturb weights or |
| KOTH state (they are already committed), so every chain lookup degrades |
| to an empty field and any publish error is logged and swallowed. The |
| audit CLI treats missing context as WARN, not PASS. |
| """ |
| from datetime import datetime |
|
|
| from ..shared.hippius import ( |
| open_manifest_store, |
| publish_receipt, |
| update_receipt_index, |
| ) |
| from ..shared.receipt import dump_receipt, sign_receipt, summarize_receipt |
|
|
| try: |
| epoch_start = self._epoch_start_block(manifest) |
| epoch_hash = "" |
| current_block: int | None = None |
| participants: tuple[Participant, ...] = () |
| try: |
| epoch_hash = client.block_hash(epoch_start) |
| participants = participants_from_commitments( |
| client.poll_commitments(include_history=True), |
| cutoff_block=epoch_start, |
| floor_block=self.cfg.round.commit_floor_block, |
| ) |
| |
| |
| current_block = int(client.current_block()) |
| except Exception as e: |
| log.warning("receipt chain context unavailable for round=%s: %s", |
| manifest.round_id, e) |
| provenance = ("", "") |
| prov_fn = getattr(window_source, "provenance_for_round", None) |
| if prov_fn is not None: |
| try: |
| |
| |
| provenance = tuple(prov_fn(base_seed, block=epoch_start)) |
| except Exception as e: |
| log.warning("pool provenance unavailable for round=%s: %s", |
| manifest.round_id, e) |
| wallet = getattr(client, "wallet", lambda: None)() |
| hotkey_ss58 = str(getattr(getattr(wallet, "hotkey", None), "ss58_address", "") or "") |
|
|
| receipt = self.build_round_receipt( |
| manifest, |
| base_seed=base_seed, |
| epoch_start_block=epoch_start, |
| epoch_block_hash=str(epoch_hash), |
| outcome=outcome, |
| windows=windows, |
| reject_reason=reject_reason, |
| participants=participants, |
| pool_provenance=(provenance[0], provenance[1]), |
| reward_uids=reward_uids, |
| weights=weights, |
| validator_hotkey=hotkey_ss58, |
| ) |
| if wallet is not None: |
| receipt = sign_receipt(receipt, wallet) |
| else: |
| log.warning("publishing an UNSIGNED receipt (no wallet) for round=%s", |
| manifest.round_id) |
| store = open_manifest_store(self.cfg.storage) |
| key = publish_receipt(store, dump_receipt(receipt), manifest.round_id, |
| validator_hotkey=hotkey_ss58) |
| log.info("published %s receipt round=%s signed=%s → s3://%s/%s", |
| receipt.status, manifest.round_id, receipt.signature is not None, |
| self.cfg.storage.manifest_bucket, key) |
| |
| |
| try: |
| now_iso = datetime.now(UTC).isoformat(timespec="seconds") |
| |
| |
| |
| |
| |
| chain = { |
| "as_of": now_iso, |
| "current_block": current_block, |
| "epoch_start_block": epoch_start, |
| "epoch_blocks": int(self.cfg.round.epoch_blocks), |
| "block_time_s": 12.0, |
| } |
| update_receipt_index( |
| store, summarize_receipt(receipt), |
| updated_at=now_iso, |
| subnet={"netuid": self.cfg.subnet.netuid, "name": self.cfg.subnet.name}, |
| chain=chain, |
| ) |
| except Exception as e: |
| log.warning("receipt index update failed for round=%s: %s", |
| manifest.round_id, e) |
| except Exception as e: |
| log.warning("receipt publication failed for round=%s: %s", manifest.round_id, e) |
|
|
| |
|
|
| def _publish_chain_status(self, client: object, store: object) -> None: |
| """Publish the dashboard's live ``status/chain.json`` (current block, |
| epoch grid, stage windows, revealed submissions) on the poll cadence. |
| |
| Purely presentational and best-effort — it feeds the web dashboard's |
| round-stage strip and live submissions panel between receipts, and any |
| failure (chain flake, storage outage) is swallowed: status telemetry |
| must never disturb a round. |
| """ |
| from datetime import datetime |
|
|
| from ..shared.chain_status import build_chain_status, publish_chain_status |
|
|
| try: |
| status = build_chain_status( |
| self.cfg, |
| current_block=int(client.current_block()), |
| commitments=client.poll_commitments(), |
| network=str(getattr(client, "network", "")), |
| as_of=datetime.now(UTC).isoformat(timespec="seconds"), |
| ) |
| publish_chain_status(store, status) |
| except Exception as e: |
| log.debug("chain status publish skipped: %s", e) |
|
|
| def run_forever(self, client: object, *, window_source: object) -> None: |
| """Poll the manifest bucket → evaluate → set weights, once per round. |
| |
| ``window_source`` is a :class:`cascade.validator.windows.WindowSource` |
| (the loaded private pool). Each new manifest's ``round_id`` is the base |
| seed; the same seed drives the rotating window slice so every validator |
| scores the identical set. |
| """ |
| import time |
|
|
| from ..shared.hippius import StorageError, open_manifest_store, read_latest_manifest |
| from ..shared.manifest import load_manifest |
|
|
| store = open_manifest_store(self.cfg.storage) |
| poll = self.cfg.manifest.poll_seconds |
| |
| |
| |
| |
| last_round: str | None = None |
| last_digest: str | None = None |
| while True: |
| try: |
| |
| |
| |
| self._publish_chain_status(client, store) |
| raw = read_latest_manifest(store) |
| digest = hashlib.sha256(raw.encode("utf-8")).hexdigest() |
| if digest == last_digest: |
| log.debug("manifest unchanged (round=%s sha=%s…); skipping", |
| last_round, digest[:12]) |
| else: |
| manifest = load_manifest(raw) |
| base_seed = int(manifest.round_id) |
| if manifest.round_id == last_round: |
| log.warning( |
| "manifest for already-handled round=%s RE-PUBLISHED " |
| "with different content (sha %s… -> %s…); re-judging", |
| manifest.round_id, |
| (last_digest or "")[:12], digest[:12], |
| ) |
| log.info( |
| "new manifest round=%s entries=%d (%s); gating + scoring …", |
| manifest.round_id, len(manifest.entries), |
| ",".join(f"{e.role}:uid{e.miner_uid}" for e in manifest.entries), |
| ) |
| |
| reason = self.check_manifest(manifest) |
| retry_pin = False |
| if reason is None: |
| |
| |
| try: |
| reason = self.check_pool_pin( |
| manifest, window_source, |
| block=self._epoch_start_block(manifest), |
| ) |
| except StorageError as e: |
| |
| |
| |
| reason = self._pool_pin_read_failed(manifest.round_id, e) |
| retry_pin = reason is None |
| else: |
| self._pin_read_first_failure.pop(str(manifest.round_id), None) |
| if retry_pin: |
| |
| |
| |
| |
| pass |
| elif reason is not None: |
| log.warning("rejecting manifest round=%s: %s", manifest.round_id, reason) |
| last_round, last_digest = manifest.round_id, digest |
| |
| |
| self._publish_round_receipt( |
| client, manifest, base_seed, |
| reject_reason=reason, window_source=window_source, |
| ) |
| elif not self.king_synced(manifest): |
| |
| |
| |
| |
| |
| last_round, last_digest = manifest.round_id, digest |
| self.state, reject_reason = self._resync_step(manifest) |
| self._persist_state() |
| reward_uids = self._reward_uids(manifest, None, client) |
| weights_vec = self._apply_weights(client, manifest.round_id, reward_uids) |
| self._publish_round_receipt( |
| client, manifest, base_seed, |
| reject_reason=reject_reason, |
| window_source=window_source, |
| reward_uids=tuple(reward_uids), weights=weights_vec, |
| ) |
| else: |
| |
| |
| windows = window_source.windows_for_round( |
| base_seed, self.cfg.eval.n_windows, |
| block=self._epoch_start_block(manifest), |
| ) |
| |
| |
| |
| |
| |
| outcome = self.process_round(manifest, windows, base_seed) |
| |
| |
| if self.state.resync_holds or self.state.last_resync_round_id: |
| self.state = replace( |
| self.state, resync_holds=0, last_resync_round_id=None) |
| last_round, last_digest = manifest.round_id, digest |
| self._persist_state() |
| reward_uids = self._reward_uids(manifest, outcome, client) |
| weights_vec = self._apply_weights(client, manifest.round_id, reward_uids) |
| |
| |
| |
| if outcome is not None: |
| self._publish_round_receipt( |
| client, manifest, base_seed, |
| outcome=outcome, windows=windows, |
| window_source=window_source, |
| reward_uids=tuple(reward_uids), weights=weights_vec, |
| ) |
| else: |
| |
| |
| self._publish_round_receipt( |
| client, manifest, base_seed, |
| reject_reason="no_king_challenger_pair", |
| window_source=window_source, |
| reward_uids=tuple(reward_uids), weights=weights_vec, |
| ) |
| |
| |
| self._maybe_run_benchmarks(manifest, outcome) |
| |
| |
| |
| |
| self._cascade_round(manifest, outcome) |
| except Exception as e: |
| log.exception("round processing failed; retrying after poll: %s", e) |
| try: |
| self._maybe_reassert_weights(client) |
| except Exception as e: |
| log.warning("weight re-assert check failed: %s", e) |
| time.sleep(poll) |
|
|
| def _maybe_reassert_weights(self, client: object) -> None: |
| """Re-push the standing weight vector every ``weight_set_interval_blocks``. |
| |
| Subtensor treats a validator whose ``last_update`` is older than the |
| subnet's ``activity_cutoff`` (5000 blocks ≈ 16.7 h on netuid 91) as |
| inactive in Yuma consensus. A mainnet round is 7200 blocks, so voting |
| only when a manifest lands blows through the cutoff every round — and a |
| stalled trainer silences the validator entirely. The vector needs no |
| manifest: it is recomputed from the persisted champion state (king + |
| registered prior kings; no champion burns to ``burn_uid``), so this is |
| a pure freshness signal that can never move the throne. The interval |
| must stay ≥ the subnet's ``weights_rate_limit`` (100 blocks on netuid |
| 91) or the chain silently no-ops the extrinsic; ≤ 0 disables. |
| """ |
| interval = int(self.cfg.validator.weight_set_interval_blocks) |
| if interval <= 0: |
| return |
| cur = int(client.current_block()) |
| if self._last_weight_block is not None and cur - self._last_weight_block < interval: |
| return |
| uids: list[int] = [] |
| if self.state.king_hotkey is not None: |
| uid = client.uid_for_hotkey(self.state.king_hotkey) |
| if uid is not None: |
| uids.append(uid) |
| for hk in self.state.former_kings: |
| uid = client.uid_for_hotkey(hk) |
| if uid is not None and uid not in uids: |
| uids.append(uid) |
| log.info("re-asserting weights (last set %s, block %d): reward_uids=%s", |
| "never" if self._last_weight_block is None else |
| f"{cur - self._last_weight_block} blocks ago", cur, |
| uids or [self.cfg.scoring.burn_uid]) |
| self._apply_weights(client, "weight-reassert", uids) |
| |
| |
| self._last_weight_block = cur |
|
|
| @staticmethod |
| def _manifest_king_hotkey(manifest: TrainingManifest) -> str | None: |
| e = manifest.entry_for_role("king") |
| return e.miner_hotkey if e is not None else None |
|
|
| @staticmethod |
| def _manifest_king_uid(manifest: TrainingManifest) -> int | None: |
| e = manifest.entry_for_role("king") |
| return e.miner_uid if e is not None else None |
|
|
| def king_synced(self, manifest: TrainingManifest) -> bool: |
| """Whether the round's *trained* king matches the validator's champion. |
| |
| The trainer picks the king it trains from on-chain incentive, which lags |
| the validator's dethrone verdicts. Until the champion |
| the validator crowned actually becomes the highest-incentive UID — and so |
| the king the trainer trains — the two disagree, and a round trained |
| against the *old* king must not have its verdict applied to the *new* |
| champion. Synced when the champion is unset (bootstrap) or the trained |
| king is the champion. |
| """ |
| if self.state.king_hotkey is None: |
| return True |
| return self._manifest_king_hotkey(manifest) == self.state.king_hotkey |
|
|
| def _resync_step(self, manifest: TrainingManifest) -> tuple[ChampionState, str]: |
| """Next champion state + receipt reason for a king-resync round. |
| |
| Called when the trained king != champion (``king_synced`` is False). By |
| default it holds the throne, bumping the consecutive-hold counter, and the |
| caller keeps voting the champion so incentive migrates and the trainer |
| re-syncs. SAFETY VALVE: once the champion has stayed un-synced for |
| ``scoring.king_resync_max_rounds`` consecutive rounds it can never be the |
| king the trainer trains (e.g. it has no usable commitment), so holding |
| forever would wedge the subnet — the valve abandons it and adopts the |
| trainer's trained king (:func:`state.demote_to_trained`), and normal |
| scoring resumes next round. ``king_resync_max_rounds <= 0`` disables the |
| valve (hold indefinitely). The counter advances once per DISTINCT |
| un-synced round: a restart re-gates the same stale manifest, and |
| counting those re-gates let five restarts during a pause trip the |
| valve and demote a healthy champion (2026-07-22). Pure: returns the |
| new state; the caller persists, votes, and publishes. |
| """ |
| champ = self.state.king_hotkey |
| trained = self._manifest_king_hotkey(manifest) |
| round_id = str(manifest.round_id) |
| same_round = self.state.last_resync_round_id == round_id |
| holds = self.state.resync_holds if same_round else self.state.resync_holds + 1 |
| cap = self.cfg.scoring.king_resync_max_rounds |
| if 0 < cap <= holds and trained is not None: |
| log.warning( |
| "round=%s king_resync SAFETY VALVE: champion %s un-synced %d rounds " |
| "(cap=%d) — demoting to trained king %s, resuming normal scoring next round", |
| manifest.round_id, (champ or "?")[:12], holds, cap, (trained or "?")[:12], |
| ) |
| return ( |
| state_mod.demote_to_trained( |
| self.state, trained_hotkey=trained, |
| trained_uid=self._manifest_king_uid(manifest), |
| ), |
| f"king_resync_demoted: champion {champ} un-synced {holds} rounds " |
| f"(cap={cap}); adopted trained king {trained}", |
| ) |
| log.warning( |
| "round=%s trainer king %s != champion %s; voting champion to re-sync " |
| "incentive, KOTH state held (%d/%s)", |
| manifest.round_id, (trained or "?")[:12], (champ or "?")[:12], |
| holds, cap if cap > 0 else "∞", |
| ) |
| return ( |
| replace(self.state, resync_holds=holds, last_resync_round_id=round_id), |
| f"king_resyncing: champion {champ} != trained king {trained}", |
| ) |
|
|
| def _king_uid_to_vote(self, manifest: TrainingManifest, *, client: object | None = None) -> int | None: |
| """The UID to put the king's weight on this round. |
| |
| The **validator's champion state** is the authority on who holds the |
| throne, so vote *that* king every round — not the (lagging) king the |
| trainer happened to train. This is what makes a dethrone STICK: the new |
| champion keeps the weight, incentive migrates to it, and next round the |
| trainer trains it as king (they re-sync). Voting the trained/manifest |
| king instead — the old behaviour — reverted a dethrone the moment the |
| trainer lagged one round, orphaning the champion. The champion hotkey is |
| resolved to its current UID via the metagraph (robust to re-registration); |
| the manifest king is used only to bootstrap when there is no champion yet. |
| """ |
| if self.state.king_hotkey is not None: |
| if client is not None: |
| resolved = client.uid_for_hotkey(self.state.king_hotkey) |
| if resolved is not None: |
| return resolved |
| return self.state.king_uid |
| king_entry = manifest.entry_for_role("king") |
| return king_entry.miner_uid if king_entry is not None else None |
|
|
| def _reward_uids( |
| self, manifest: TrainingManifest, outcome: RoundOutcome | None, client: object |
| ) -> list[int]: |
| """UIDs that share this round's weight: the current king plus any |
| ``former_kings`` still registered (teutonic-style equal-share payout). |
| |
| Returns an empty list when there is no king to vote for at all (no |
| champion and no manifest king); the loop hands that to |
| ``set_equal_share_weights``, which burns to ``burn_uid`` rather than |
| reverting. The list is otherwise deduped/range-checked there too. |
| |
| ``[validator] force_burn`` empties the list HERE — not just at the |
| weight push — so the published receipt's ``reward_uids`` agree with the |
| burn vector actually set (``cascade-audit`` recomputes one from the |
| other and fails on a mismatch). |
| """ |
| if self.cfg.validator.force_burn: |
| return [] |
| uids: list[int] = [] |
| king_uid = self._king_uid_to_vote(manifest, client=client) |
| if king_uid is not None: |
| uids.append(king_uid) |
| for hk in self.state.former_kings: |
| uid = client.uid_for_hotkey(hk) |
| if uid is not None: |
| uids.append(uid) |
| return uids |
|
|
| def _apply_weights(self, client: object, round_id: str, reward_uids: list[int]) -> tuple[float, ...]: |
| """Set the equal-share weight vector on chain; return it (empty on failure). |
| |
| Shared by the scored path and the king-resync path. Always sets weights — |
| an empty ``reward_uids`` burns to ``burn_uid`` so emission still leaves the |
| network. A failed extrinsic is logged and retried next round (the empty |
| vector is recorded truthfully in the receipt). |
| |
| ``[validator] force_burn`` overrides the vector to a burn HERE — the |
| single choke point every push flows through (scored, resync, re-assert) — |
| so the receipt records the burn that was actually set. Champion state is |
| never touched by this override.""" |
| from ..shared.chain import decayed_share_vector |
|
|
| if self.cfg.validator.force_burn: |
| log.warning( |
| "FORCE-BURN active (round=%s): %sburning to uid %d (champion state " |
| "untouched — unset [validator] force_burn and restart to resume voting)", |
| round_id, |
| f"dropping reward_uids={reward_uids}; " if reward_uids else "", |
| self.cfg.scoring.burn_uid, |
| ) |
| reward_uids = [] |
|
|
| decay = self.cfg.scoring.king_decay |
| try: |
| n_uids = client.n_uids() |
| client.set_equal_share_weights( |
| reward_uids, n_uids, decay=decay, burn_uid=self.cfg.scoring.burn_uid, |
| ) |
| vec = tuple(decayed_share_vector( |
| reward_uids, n_uids, decay=decay, burn_uid=self.cfg.scoring.burn_uid)) |
| log.info("round=%s weights set: reward_uids=%s (n_uids=%d, burn_uid=%d)", |
| round_id, reward_uids or [self.cfg.scoring.burn_uid], n_uids, |
| self.cfg.scoring.burn_uid) |
| |
| with contextlib.suppress(Exception): |
| self._last_weight_block = int(client.current_block()) |
| return vec |
| except Exception as e: |
| log.warning("weight set failed for round=%s (king holds, retried next round): %s", |
| round_id, e) |
| return () |
|
|
| def _persist_state(self) -> None: |
| from . import state as state_mod |
|
|
| try: |
| Path(self.cfg.validator.state_db_path).write_text( |
| state_mod.dumps(self.state), encoding="utf-8" |
| ) |
| except Exception as e: |
| log.warning("failed to persist validator state: %s", e) |
|
|
|
|
| def _load_state(path: str) -> ChampionState: |
| """Load persisted champion state from ``state_db_path`` (JSON), or a fresh |
| state if the file is absent/unreadable.""" |
| p = Path(path) |
| if not p.is_file(): |
| return ChampionState() |
| try: |
| return state_mod.loads(p.read_text(encoding="utf-8")) |
| except Exception as e: |
| log.warning("could not load validator state from %s (%s); starting fresh", path, e) |
| return ChampionState() |
|
|
|
|
| def _bootstrap_state_from_receipts(store: object, anchor: str) -> ChampionState | None: |
| """Champion inherited from the signed public receipt trail, or ``None``. |
| |
| First-boot inheritance for a validator with no local state: the throne |
| otherwise lives only in each validator's private state DB, so a validator |
| joining mid-reign would judge the next manifest blind (``king_synced`` |
| treats an unset champion as synced) and crown whichever king it happened |
| to see win first — a different champion than every validator that |
| witnessed the real dethrone (OPSLOG 2026-07-17). |
| |
| ``anchor`` is the pinned receipt-signing ss58 (``[manifest] |
| validator_hotkey``, falling back to ``trainer_hotkey``) — the same trust |
| anchor the validator already applies to manifests, extended once, at |
| first boot, to the receipt trail. Reads the anchor's |
| ``receipts/<anchor>/latest.json`` (legacy shared pointer as fallback) and |
| adopts the throne recorded by a *scored* receipt whose signature |
| verifies. When ``latest.json`` is a hold/rejected receipt (verdict-less |
| by construction), the receipt *index* is consulted — but only as an |
| UNTRUSTED pointer to candidate round ids: nothing is adopted except from |
| a per-round receipt whose signature verifies against the anchor. |
| |
| Anything short of that — missing objects, unreadable JSON, a bad |
| signature, a genesis throne (``king_hotkey`` unset) — returns ``None`` |
| and the caller proceeds with the stock blank-slate behaviour. Storage |
| faults must never block validator startup. |
| """ |
| if not anchor: |
| return None |
| from ..shared.hippius import ( |
| RECEIPT_INDEX_KEY, |
| RECEIPT_LATEST_KEY, |
| receipt_latest_key, |
| receipt_round_key, |
| ) |
| from ..shared.receipt import load_receipt, verify_receipt_signature |
|
|
| def _adopt_from(key: str) -> ChampionState | None: |
| try: |
| text = store.get_text(key) |
| except Exception: |
| return None |
| try: |
| receipt = load_receipt(text) |
| except Exception as e: |
| log.warning("receipt bootstrap: unreadable receipt at %s (%s); skipped", key, e) |
| return None |
| if not verify_receipt_signature(receipt, anchor): |
| log.warning("receipt bootstrap: receipt at %s is not signed by the pinned " |
| "hotkey %s…; skipped", key, anchor[:8]) |
| return None |
| v = receipt.verdict |
| if receipt.status != "scored" or v is None or not v.king_hotkey or v.king_uid is None: |
| return None |
| log.info("receipt bootstrap: adopting champion %s (uid %d) from signed scored " |
| "receipt round=%s", v.king_hotkey, int(v.king_uid), receipt.round_id) |
| return ChampionState(king_hotkey=str(v.king_hotkey), king_uid=int(v.king_uid)) |
|
|
| adopted = _adopt_from(receipt_latest_key(anchor)) or _adopt_from(RECEIPT_LATEST_KEY) |
| if adopted is not None: |
| return adopted |
| try: |
| rows = json.loads(store.get_text(RECEIPT_INDEX_KEY)).get("rounds", []) |
| except Exception: |
| rows = [] |
| |
| |
| for row in reversed(rows): |
| if str(row.get("status")) != "scored": |
| continue |
| adopted = _adopt_from(receipt_round_key(str(row.get("round_id", "")), anchor)) |
| if adopted is not None: |
| return adopted |
| return None |
|
|
|
|
| def _warm_start_installer(path: Path) -> Callable[[object], None]: |
| """The default Cascade installer: promote the winning checkpoint by writing its |
| pointer (and its eval numbers) to ``warm_start_init_path`` — the seam the |
| trainer reads to warm-start every subsequent round from. Promotes AS-IS; no |
| retrain/fine-tune.""" |
|
|
| def _install(winner: object) -> None: |
| import time |
|
|
| path.parent.mkdir(parents=True, exist_ok=True) |
| path.write_text( |
| json.dumps( |
| { |
| "checkpoint_id": getattr(winner, "checkpoint_id", None), |
| "size": getattr(winner, "size", ""), |
| "score": getattr(winner, "score", None), |
| "gifteval_crps": getattr(winner, "gifteval_crps", None), |
| "gifteval_mase": getattr(winner, "gifteval_mase", None), |
| "boom_crps": getattr(winner, "boom_crps", None), |
| "boom_mase": getattr(winner, "boom_mase", None), |
| "time_crps": getattr(winner, "time_crps", None), |
| "time_mase": getattr(winner, "time_mase", None), |
| "installed_at": time.time(), |
| }, |
| sort_keys=True, |
| ), |
| encoding="utf-8", |
| ) |
| log.info("cascade: warm-start init written to %s (checkpoint %s)", |
| path, getattr(winner, "checkpoint_id", "?")) |
|
|
| return _install |
|
|
|
|
| def _build_cascade(cfg: ChainConfig) -> CascadeController: |
| """Construct the Cascade controller from config, restoring the persisted reign |
| clock + checkpoint log so it resumes across restarts.""" |
| from .cascade import CascadeController, load_state |
|
|
| state_path = Path(cfg.validator.cascade_state_db_path) |
| return CascadeController( |
| reign_days=cfg.scoring.cascade_reign_days, |
| state=load_state(state_path), |
| install_fn=_warm_start_installer(Path(cfg.validator.warm_start_init_path)), |
| state_path=state_path, |
| ) |
|
|
|
|
| def build_runner( |
| *, |
| chain_toml: Path | None = None, |
| cache_dir: Path | None = None, |
| device: str = "cpu", |
| eval_host_fn: Callable[[], RemoteHost | None] | None = None, |
| ) -> ValidatorRunner: |
| """Construct a runner from ``chain.toml``, restoring persisted champion |
| state. Wallet/chain wiring for live weight-setting is attached by |
| ``cascade-validator`` (see main.py). ``eval_host_fn`` (optional) resolves |
| the GPU pod to offload heavy evals to — re-invoked per eval, so an elastic |
| provisioner-rented pod is picked up lazily; the wallet stays on this box.""" |
| from ..shared.config import load_chain_config |
|
|
| cfg = load_chain_config(chain_toml) |
| state = _load_state(cfg.validator.state_db_path) |
| if (cfg.validator.bootstrap_from_receipts and state.king_hotkey is None |
| and state.rounds_seen == 0 and not state.former_kings): |
| |
| |
| from ..shared.hippius import open_manifest_store |
|
|
| anchor = cfg.manifest.validator_hotkey or cfg.manifest.trainer_hotkey |
| try: |
| adopted = _bootstrap_state_from_receipts( |
| open_manifest_store(cfg.storage), anchor) |
| except Exception as e: |
| log.warning("receipt bootstrap skipped (%s); starting blank", e) |
| adopted = None |
| if adopted is not None: |
| state = adopted |
| |
| |
| cascade = _build_cascade(cfg) if cfg.scoring.cascade_enabled else None |
| return ValidatorRunner( |
| cfg=cfg, state=state, |
| cache_dir=cache_dir, device=device, cascade=cascade, eval_host_fn=eval_host_fn, |
| ) |
|
|