"""Cascade-EP trainer — EQUILIBRIUM MODE (the true-EP route). Two-phase (+-beta) relaxation of all layer states to the nudged equilibria via Gauss-Seidel reverse sweeps (solver choice only; readout is taken AT the relaxed states with the standard EP formula), weight grad = (1/2beta)[dF/dtheta|+ - dF/dtheta|-]. Inference = plain forward (standard LLM). Twin of casc_bp_train.py (same seed/data).""" import argparse, math, pickle, time import numpy as np, torch, torch.nn as nn, torch.nn.functional as F from pathlib import Path ap = argparse.ArgumentParser() ap.add_argument('--tag', default='casc_eq6') ap.add_argument('--L', type=int, default=6); ap.add_argument('--C', type=int, default=256) ap.add_argument('--H', type=int, default=8); ap.add_argument('--T', type=int, default=256) ap.add_argument('--B', type=int, default=24); ap.add_argument('--steps', type=int, default=4000) ap.add_argument('--lr', type=float, default=3e-4); ap.add_argument('--warmup', type=int, default=200) ap.add_argument('--beta', type=float, default=0.003); ap.add_argument('--seed', type=int, default=0) ap.add_argument('--K', type=int, default=3) # fb (message-passing) rounds ap.add_argument('--geta', type=float, default=1.0) # fb mixing (1.0 = undamped) ap.add_argument('--save_every', type=int, default=1000); ap.add_argument('--log', type=int, default=100) ap.add_argument('--wandb', default='auto') # ON BY DEFAULT; 'auto' = per-regime project (ept-fineweb-72m / ept-tinystories-42m); --wandb '' to disable ap.add_argument('--wandb_run', default='') ap.add_argument('--kmax', type=int, default=8) # adaptive fb rounds cap ap.add_argument('--noguard', action='store_true') # diagnosis: skip only non-finite grads ap.add_argument('--untie', action='store_true') # separate readout matrix (untied from tok) ap.add_argument('--opt', choices=['adamw', 'muon'], default='adamw') ap.add_argument('--muon_lr', type=float, default=0.02) ap.add_argument('--tok_init', type=float, default=0.0) # >0: init tok/pos with this std (GPT-standard 0.02) ap.add_argument('--compile', action='store_true') # torch.compile each block (free speed where supported) ap.add_argument('--sig_every', type=int, default=25) # tok-sigma refresh interval (amortized) ap.add_argument('--beta_floor', type=float, default=0.0) # >0: floor beta_t (anti finite-beta SNR collapse at depth) ap.add_argument('--beta_fixed', action='store_true') # disable sig^2 schedule, hold beta_t = args.beta constant ap.add_argument('--cosine', action='store_true') # warmup then cosine decay to lr_min_ratio*lr over --steps (long runs) ap.add_argument('--lr_min_ratio', type=float, default=0.1) ap.add_argument('--qk_norm', action='store_true') # RMS-norm q,k per head before scores (OLMo2-style; bounds logits, analog-friendly) ap.add_argument('--final_ln', action='store_true') # final LayerNorm before readout (standard GPT; bounds sig_tok growth -> keeps beta/estimator healthy on long runs) ap.add_argument('--resume', default='') # path to a ckpt (tok/pos/blocks) to continue from; step taken from ckpt ap.add_argument('--sig0', type=float, default=-1.0) # override SIG0 (beta-schedule ref); needed on resume to restore original beta regime ap.add_argument('--olmo2', action='store_true') # OLMo2-standard block: norm-AFTER-sublayer RMSNorm, full-width QK-norm, RoPE(500k), SwiGLU, no-bias, untied head, final RMSNorm, 0.02 init ap.add_argument('--wd', type=float, default=-1.0) # >=0: grouped weight decay (linear weights+head decay; embeddings/norm-gains none). <0 = legacy uniform 1e-4 ap.add_argument('--zloss', type=float, default=0.0) # z-loss coefficient on train objective (OLMo2-style logit regularizer); 0 = off ap.add_argument('--kretry', type=int, default=0) # >0: on drift-reject, RETRY the batch once with this many fb rounds (diag B: K8 converges the marginal batches) instead of dropping it ap.add_argument('--bf_late', type=float, default=0.0) # >0: raise beta_floor to this value from step --bf_late_at (late-training SNR fix; dose-response 2026-07-10) ap.add_argument('--bf_late_at', type=int, default=25000) ap.add_argument('--bsign_rand', action='store_true') # random-sign beta per step (KHS 'random scheme'): averages the O(beta) single-sided bias at single-phase cost ap.add_argument('--bf16', action='store_true') # cast model to bf16 (E-accumulation + tok_sigma stay fp32) — the x0.5 cost lever, GATE before production ap.add_argument('--amp', action='store_true') # PROPER mixed precision: autocast(bf16) matmuls, fp32 params/states/d/E — amp_gate.py PASSED 2026-07-12 (cos 0.9682 vs fp32 0.9687); --bf16 naive-cast stays DEAD (state quantization, RESULT 11) ap.add_argument('--dtop_every', type=int, default=1) # 1 = exact (DEFAULT, BP-parity); 2 = fast mode (~20% cheaper, ~4% CE tax at high lr) ap.add_argument('--gate_every', type=int, default=200) # in-training cos(EP,BP) telemetry; <=0 = fully BP-free (no bp_gate at all) ap.add_argument('--gate_govern', action='store_true') # let gate cos adjust K/bscale (default: observe-only => training control is BP-free) ap.add_argument('--data', default='tinystories_bpe') # dataset dir under ep_run/data (train.bin/val.bin/meta.pkl) ap.add_argument('--sync_check', type=int, default=500) # DDP: verify bitwise param sync every N steps (0=off) ap.add_argument('--ddp_backend', default='nccl', choices=['nccl', 'gloo']) # gloo = correctness tests on shared GPUs ap.add_argument('--ddp_grad_test', action='store_true') # one-step grad equivalence test vs single-GPU big batch, then exit ap.add_argument('--beta_ride', type=float, default=1.0) # cap ceiling: >1 lets the governor RAISE beta # above the schedule, up to ride x schedule ap.add_argument('--beta_ride_up', type=float, default=1.02) # per-step climb rate in the calm branch ap.add_argument('--bpmix', default='') # ABLATION: overwrite EP grads with TRUE BP grads # for selected groups: 'blocks:0-5' | 'blocks:6-11' # | 'attn' | 'ffn' | 'head' (comma-separated) ap.add_argument('--beta_sync', type=int, default=0) # >0: SYNCHRONOUS acceptance — this step's own # relax telemetry gates the commit; on reject, # halve beta and retry same batch (N halvings max) ap.add_argument('--ride_ema', type=int, default=0) # ride-v2(b): rho-EMA before governor decisions ap.add_argument('--ride_cool', type=int, default=0) # ride-v2(c): post-attack climb cooldown (steps) ap.add_argument('--drift_adapt', type=float, default=0.0) # >0: adaptive drift ceiling = this x trailing # accepted-drift EMA (floor 0.05, cap 0.5); # blocks beta-scaled garbage accepts (ride-v2d) ap.add_argument('--beta_cap_rho', type=float, default=0.0) # >0: LOOP-GAIN CAP on beta — if per-sweep residual # ratio rho^ exceeds this, bscale *= 0.8 (beta backs off # under the wall-2 ceiling); recovers x1.02 when rho^ low ap.add_argument('--beta_servo', type=float, default=0.0) # >0: CEILING-HUGGING SERVO. Meter lit (res above the # v2 absolute gate): deadbeat inversion onto the ceiling # — cap *= servo*cap_rho/rho^ (rho ~= G*beta near the # edge, so one step lands beta at servo*ceiling; grows # toward it when under, shrinks when over). Meter dark: # existing ride climb probes upward. Value = safety # fraction of the ceiling to sit at (canonical 0.8). ap.add_argument('--wsync', type=int, default=0) # >0: SYNCHRONOUS WEIGHT-STEP ACCEPTANCE — snapshot # params+momentum before each opt.step; next step's # nudged relax measures the new state through the SAME # _legal gate (res/rho/drift, no new bounds); illegal -> # roll back and re-apply the update at half scale # (p <- (p+snap)/2), up to wsync halvings, then full # revert + skip. The weight trajectory structurally # cannot dwell past the ceiling. Zero new constants. ap.add_argument('--cap_floor', type=float, default=0.05) # hard bottom of the rho-cap; 0 = pure ceiling-tracking # (cap follows the measured ceiling all the way down; a # pinned bottom above the true ceiling = disguised wall-2) ap.add_argument('--relax_tol', type=float, default=0.0) # >0: ADAPTIVE relax — sweep until rel. state change < tol # (or --kmax), geta backtracks x0.6 on residual GROWTH (rho>=1 # signal), then one final graphed round. 0 = legacy fixed-K. ap.add_argument('--muon_mom', type=float, default=0.95) # Muon momentum (late-SNR arm: 0.99 = ~10x noise averaging) ap.add_argument('--adam_b1', type=float, default=0.9) # AdamW beta1 (late-SNR arm companion) ap.add_argument('--est', choices=['single', 'centered', 'richardson'], default='single') # centered: [g(+b)+g(-b)]/2 (O(b^2) bias, 2x relax cost) # richardson: 2g(b)-g(2b) (O(b^2) bias, large-b friendly) ap.add_argument('--est_late', choices=['', 'centered'], default='') ap.add_argument('--est_late_at', type=int, default=0) # switch --est -> --est_late at this step (process-local, # bf_late_at semantics); centered is TAIL medicine ap.add_argument('--qcomp_bits', type=int, default=0) # STAGE-0 T64 scenario: forward/transpose COMPUTE # on grid-snapped weights, fp32 master gets updates # (= word-streaming / shadow accumulation) ap.add_argument('--qup_bits', type=int, default=0) # STAGE-0: quantize weights to an absolute # per-tensor grid after each update (stochastic # rounding); emulates finite analog cell levels ap.add_argument('--centmirror', action='store_true') # centered's -beta pass initialized as the MIRROR # of the +beta solution (d- = -d+ at shared anchor) # + one polish sweep; skips its free pass entirely ap.add_argument('--centfast', action='store_true') # centered via ONE doubled batch [x;x], +beta/-beta halves # (shared kernels; math identical to sequential centered) args = ap.parse_args() if args.olmo2: args.untie = True if args.tok_init <= 0: args.tok_init = 0.02 torch.manual_seed(args.seed) dev = 'cuda' if torch.cuda.is_available() else 'cpu' # ---- DDP (manual: autograd.grad path, guard-synced; torchrun --standalone --nproc_per_node=N) ---- import os import torch.distributed as dist DDP = int(os.environ.get('WORLD_SIZE', '1')) > 1 if DDP: dist.init_process_group(args.ddp_backend) RANK, WORLD = dist.get_rank(), dist.get_world_size() torch.cuda.set_device(int(os.environ['LOCAL_RANK']) % max(torch.cuda.device_count(), 1)) else: RANK, WORLD = 0, 1 DGEN = torch.Generator().manual_seed(args.seed * 7919 + RANK * 104729 + 11) # per-rank DATA stream ONLY # (init/bsign RNGs stay rank-identical) def ddp_avg(gs, params): """average a grad list across ranks; preserves the None pattern (identical graphs => identical pattern) so optimizer skip-semantics match single-GPU exactly.""" if not DDP: return gs none_mask = [g is None for g in gs] filled = [g if g is not None else torch.zeros_like(p) for p, g in zip(params, gs)] flat = torch.cat([g.reshape(-1) for g in filled]) if args.ddp_backend == 'gloo': cf = flat.cpu(); dist.all_reduce(cf, op=dist.ReduceOp.SUM); flat = cf.to(flat.device) else: dist.all_reduce(flat, op=dist.ReduceOp.SUM) flat /= WORLD out, o = [], 0 for p in params: n = p.numel(); out.append(flat[o:o + n].view_as(p)); o += n return [None if m else g for m, g in zip(none_mask, out)] def ddp_max_scalar(v): """global max of a python float (guard decisions must be identical on every rank).""" if not DDP: return v t = torch.tensor([v if math.isfinite(v) else float('inf')], device=dev if dev == 'cuda' else 'cpu') if args.ddp_backend == 'gloo': t = t.cpu() dist.all_reduce(t, op=dist.ReduceOp.MAX) return float(t[0]) def ddp_bcast_scalar(v): if not DDP: return v t = torch.tensor([v], device=dev if dev == 'cuda' else 'cpu') if args.ddp_backend == 'gloo': t = t.cpu() dist.broadcast(t, 0) return float(t[0]) DD = Path('/home/yurenh2/ept/ep_run/data') / args.data vocab = pickle.load(open(DD / 'meta.pkl', 'rb'))['vocab_size'] def get_batch(split): data = np.memmap(DD / ('train.bin' if split == 'train' else 'val.bin'), dtype=np.uint16, mode='r') ix = torch.randint(len(data) - args.T - 1, (args.B,), generator=DGEN) x = torch.stack([torch.from_numpy(data[i:i + args.T].astype(np.int64)) for i in ix]) y = torch.stack([torch.from_numpy(data[i + 1:i + 1 + args.T].astype(np.int64)) for i in ix]) return x.to(dev), y.to(dev) class CausalSelfAttn(nn.Module): """explicit MHA (SDPA-backed) so we can QK-norm q,k per head before the scores.""" def __init__(self, C, H, qk_norm=False): super().__init__() self.H, self.hd, self.qk_norm = H, C // H, qk_norm self.qkv = nn.Linear(C, 3 * C) self.proj = nn.Linear(C, C) if qk_norm: self.q_g = nn.Parameter(torch.ones(self.hd)) self.k_g = nn.Parameter(torch.ones(self.hd)) def forward(self, x): B, T, C = x.shape q, k, v = self.qkv(x).split(C, dim=2) q = q.view(B, T, self.H, self.hd).transpose(1, 2) k = k.view(B, T, self.H, self.hd).transpose(1, 2) v = v.view(B, T, self.H, self.hd).transpose(1, 2) if self.qk_norm: # RMS-norm over head_dim (OLMo2-style), learnable per-dim gain q = q * torch.rsqrt(q.pow(2).mean(-1, keepdim=True) + 1e-6) * self.q_g k = k * torch.rsqrt(k.pow(2).mean(-1, keepdim=True) + 1e-6) * self.k_g y = F.scaled_dot_product_attention(q, k, v, is_causal=True) return self.proj(y.transpose(1, 2).contiguous().view(B, T, C)) class Block(nn.Module): def __init__(self, C, H, qk_norm=False): super().__init__() self.ln1, self.ln2 = nn.LayerNorm(C), nn.LayerNorm(C) self.attn = CausalSelfAttn(C, H, qk_norm) self.ff = nn.Sequential(nn.Linear(C, 4 * C), nn.GELU(), nn.Linear(4 * C, C)) def forward(self, z, mask=None): z = z + self.attn(self.ln1(z)) return z + self.ff(self.ln2(z)) class RMSNorm(nn.Module): def __init__(self, C, eps=1e-6): super().__init__(); self.g = nn.Parameter(torch.ones(C)); self.eps = eps def forward(self, x): return x * torch.rsqrt(x.pow(2).mean(-1, keepdim=True) + self.eps) * self.g class SwiGLU(nn.Module): def __init__(self, C): super().__init__() h = ((8 * C // 3) + 63) // 64 * 64 # ~param-match the 4x-GELU MLP (8C^2) self.w1 = nn.Linear(C, h, bias=False); self.w3 = nn.Linear(C, h, bias=False) self.w2 = nn.Linear(h, C, bias=False) def forward(self, x): return self.w2(F.silu(self.w1(x)) * self.w3(x)) class Olmo2Attn(nn.Module): """OLMo2 attention: no-bias projs, FULL-WIDTH RMS QK-norm (pre-head-split, HF Olmo2 order), then per-head RoPE.""" def __init__(self, C, H, T): super().__init__() self.H, self.hd = H, C // H self.qkv = nn.Linear(C, 3 * C, bias=False); self.proj = nn.Linear(C, C, bias=False) self.qn, self.kn = RMSNorm(C), RMSNorm(C) inv = 1.0 / (500000.0 ** (torch.arange(0, self.hd, 2).float() / self.hd)) fr = torch.outer(torch.arange(T).float(), inv) self.register_buffer('rc', fr.cos(), persistent=False) self.register_buffer('rs', fr.sin(), persistent=False) def rope(self, x): x1, x2 = x[..., ::2], x[..., 1::2] c, s = self.rc[None, None], self.rs[None, None] return torch.stack((x1 * c - x2 * s, x1 * s + x2 * c), dim=-1).flatten(-2) def forward(self, x): B, T, C = x.shape q, k, v = self.qkv(x).split(C, dim=2) q, k = self.qn(q), self.kn(k) q = self.rope(q.view(B, T, self.H, self.hd).transpose(1, 2)) k = self.rope(k.view(B, T, self.H, self.hd).transpose(1, 2)) v = v.view(B, T, self.H, self.hd).transpose(1, 2) y = F.scaled_dot_product_attention(q, k, v, is_causal=True) return self.proj(y.transpose(1, 2).contiguous().view(B, T, C)) class Olmo2Block(nn.Module): """OLMo2 reordered norm (norm AFTER each sublayer, inside the residual) — their training-stability change.""" def __init__(self, C, H, T): super().__init__() self.attn = Olmo2Attn(C, H, T); self.ff = SwiGLU(C) self.na, self.nf = RMSNorm(C), RMSNorm(C) def forward(self, z, mask=None): z = z + self.na(self.attn(z)) return z + self.nf(self.ff(z)) tok = nn.Embedding(vocab, args.C).to(dev) pos = nn.Embedding(args.T, args.C).to(dev) if args.tok_init > 0: with torch.no_grad(): tok.weight.normal_(0, args.tok_init); pos.weight.normal_(0, args.tok_init) blocks = nn.ModuleList([(Olmo2Block(args.C, args.H, args.T) if args.olmo2 else Block(args.C, args.H, args.qk_norm)) for _ in range(args.L)]).to(dev) if args.olmo2: with torch.no_grad(): for m in blocks.modules(): if isinstance(m, nn.Linear): m.weight.normal_(0, 0.02) if args.compile: try: for i in range(args.L): blocks[i] = torch.compile(blocks[i], mode='reduce-overhead') print('[compile] blocks compiled', flush=True) except Exception as e: print(f'[compile] disabled ({e})', flush=True) mask = torch.triu(torch.full((args.T, args.T), float('-inf'), device=dev), 1) W_out = nn.Parameter(torch.randn(vocab, args.C, device=dev) * 0.02) if args.untie else None ln_f = (RMSNorm(args.C) if args.olmo2 else (nn.LayerNorm(args.C) if args.final_ln else nn.Identity())).to(dev) def emb(x): return tok(x) if args.olmo2 else tok(x) + pos(torch.arange(args.T, device=dev))[None] readout = (lambda z: ln_f(z) @ W_out.t()) if args.untie else (lambda z: ln_f(z) @ tok.weight.t()) all_params = list(tok.parameters()) + ([] if args.olmo2 else list(pos.parameters())) + list(blocks.parameters()) + list(ln_f.parameters()) + ([W_out] if args.untie else []) start_step = 0 if args.resume: _ck = torch.load(args.resume, map_location=dev, weights_only=False) tok.load_state_dict(_ck['tok']); pos.load_state_dict(_ck['pos']); blocks.load_state_dict(_ck['blocks']) if _ck.get('wout') is not None and args.untie: with torch.no_grad(): W_out.copy_(_ck['wout'].to(dev)) if _ck.get('lnf') is not None and not isinstance(ln_f, nn.Identity): ln_f.load_state_dict(_ck['lnf']) start_step = int(_ck.get('step', 0)) print(f'[resume] loaded {args.resume} at step {start_step}', flush=True) if args.bf16: for _m in (tok, pos, blocks): _m.to(torch.bfloat16) if not isinstance(ln_f, nn.Identity): ln_f.to(torch.bfloat16) if args.untie: with torch.no_grad(): W_out.data = W_out.data.to(torch.bfloat16) print('[bf16] model cast to bfloat16 (E-accum + sigma stay fp32)', flush=True) if args.opt == 'muon': from muon import build_hybrid opt, sched = build_hybrid(blocks, all_params, args.lr, args.muon_lr, args.warmup, muon_mom=args.muon_mom, adam_b1=args.adam_b1, total_steps=(args.steps if args.cosine else 0), lr_min_ratio=args.lr_min_ratio) else: if args.wd >= 0: # OLMo2-style grouped decay: linear weights + head decay; embeddings/norm-gains none nodecay = {id(p) for p in tok.parameters()} | {id(p) for p in pos.parameters()} | \ {id(p) for p in blocks.parameters() if p.ndim < 2} | {id(p) for p in ln_f.parameters()} opt = torch.optim.AdamW([ {'params': [p for p in all_params if id(p) not in nodecay], 'weight_decay': args.wd}, {'params': [p for p in all_params if id(p) in nodecay], 'weight_decay': 0.0}], lr=args.lr) else: opt = torch.optim.AdamW(all_params, lr=args.lr, weight_decay=1e-4) if args.cosine: def _lrlam(s): if s < args.warmup: return (s + 1) / max(args.warmup, 1) p = min(1.0, (s - args.warmup) / max(1, args.steps - args.warmup)) return args.lr_min_ratio + 0.5 * (1 - args.lr_min_ratio) * (1 + math.cos(math.pi * p)) sched = torch.optim.lr_scheduler.LambdaLR(opt, _lrlam) else: sched = torch.optim.lr_scheduler.LambdaLR(opt, lambda s: min(1.0, (s + 1) / max(args.warmup, 1))) NBT = args.B * args.T def obj_loss(logits2d, y1d): """train objective: CE (+ optional z-loss). Used in the nudge force, theta-readout and bp_gate so EP tracks BP on the SAME objective; evaluate() stays pure CE for comparability.""" l = F.cross_entropy(logits2d, y1d) if args.zloss > 0: l = l + args.zloss * (torch.logsumexp(logits2d.float(), -1) ** 2).mean() return l def free_states_graphed(x): """free forward, keeping per-layer graphs (in_l, out_l) so round-1 backward vjps reuse them.""" with torch.no_grad(): z0 = emb(x) ins, outs, zs = [], [], [] prev = z0 with torch.autocast('cuda', dtype=torch.bfloat16, enabled=args.amp): for b in blocks: i = prev.detach().requires_grad_(True) o = b(i, mask) ins.append(i); outs.append(o); zs.append(o.detach().float()) prev = zs[-1] return z0, zs, ins, outs @torch.no_grad() def tok_sigma(iters=8): """top singular value of tok.weight (power iteration on the raw matrix).""" W = (W_out if args.untie else tok.weight).float() v = torch.randn(W.shape[1], device=dev); v /= v.norm() sig = 1.0 for _ in range(iters): u = W @ v; u /= max(u.norm(), 1e-12) v = W.t() @ u; sig = v.norm(); v /= max(sig, 1e-12) return float(sig) def relax(z0, zs, ins, outs, y, beta, K, x, bmask=None): """K fb rounds with GRAPH REUSE + two dedups: (a) the top CE force d_top is refreshed on even rounds only (states move O(beta) per round -> O(beta^2) error); (b) the LAST rebuild keeps graphs (layer-0 fed a graphed emb) and returns (ins, outs) so the theta-readout reuses them instead of re-running a full graphed chain. bmask (rows,1,1): per-row multiplier on the top force (centfast +/-1 halves); the rows*T-aware scale keeps per-row d identical to the sequential B-sized run.""" d = [None] * args.L geta_l = args.geta adaptive = args.relax_tol > 0 def forces(refresh_top): if refresh_top or d[args.L - 1] is None: zc = zs[args.L - 1].detach().requires_grad_(True) ce = obj_loss(readout(zc).reshape(-1, vocab), y.reshape(-1)) nbt_loc = zc.shape[0] * zc.shape[1] g = torch.autograd.grad(ce, zc)[0] if bmask is not None: g = g * bmask d[args.L - 1] = (-beta * nbt_loc * g).detach() for l in range(args.L - 2, -1, -1): d[l] = torch.autograd.grad(outs[l + 1], ins[l + 1], grad_outputs=d[l + 1].to(outs[l + 1].dtype))[0].detach().float() def rebuild(last): nonlocal ins, outs prev = z0 n_ins, n_outs = [], [] rnum = rden = 0.0 g_eff = 1.0 if last else geta_l # FINAL graphed round is ALWAYS full-step: the theta-read # identity (z - o) = d requires undamped substitution; # mixing there leaks the iteration residual into E (gn 1e5 bug) with torch.autocast('cuda', dtype=torch.bfloat16, enabled=args.amp): for l in range(args.L): if last and l == 0: i = emb(x) # graphed emb for the readout's E-path else: i = prev.detach().requires_grad_(True) o = blocks[l](i, mask) znew = o.detach().float() + d[l] # damped (under-relaxed) mixing: geta<1 restores contraction on stiff operators # (wall-2 toolkit); fixed point unchanged (z = z + geta*(o+d-z) <=> z = o+d) mixed = znew if g_eff >= 1.0 else (zs[l] + g_eff * (znew - zs[l])) with torch.no_grad(): rnum += float((mixed - zs[l]).norm()); rden += float(zs[l].norm()) zs[l] = mixed n_ins.append(i); n_outs.append(o) prev = zs[l] ins, outs = n_ins, n_outs return rnum / max(rden, 1e-9) if not adaptive: # legacy fixed-K path (bit-identical update semantics) rlist = [] for k in range(K): forces(k % args.dtop_every == 0) rlist.append(rebuild(k + 1 == K)) if len(rlist) >= 2 and rlist[-2] > 1e-12: GOV['rho'] = rlist[-1] / rlist[-2] # per-sweep contraction ratio = live loop-gain meter GOV['res'] = rlist[-1] GOV['kuse'] = K GOV['_last_d'] = d return zs, outs prev_res, k = None, 0 while k < args.kmax: forces(k % args.dtop_every == 0) res = rebuild(False) k += 1 if prev_res is not None and prev_res > 1e-12: GOV['rho'] = res / prev_res if prev_res is not None and res > prev_res and res > args.relax_tol: geta_l = max(0.2, geta_l * 0.6) # residual GREW: local rho>=1 -> damp harder prev_res = res if res < args.relax_tol: break forces(True) # final graphed round at the settled state (theta-read) rebuild(True) GOV['kuse'] = k + 1 GOV['_last_d'] = d return zs, outs def dFdtheta(zs, x, y, beta): """theta-readout at FIXED states. Not used by the training loop (relax reuses its own graphs); kept as the INVARIANT-TEST surface for test_bp_free.py. Self-sealing: inputs are detached here so the local-graph property holds for any caller.""" zs = [z.detach() for z in zs] prev = emb(x) E = 0.0 for z, b in zip(zs, blocks): E = E + 0.5 * ((z - b(prev, mask)) ** 2).sum() prev = z # zs detached at entry => blocks l>0 get detached inputs; block 0 gets the graphed emb obj = E / NBT + beta * F.cross_entropy(readout(zs[-1]).reshape(-1, vocab), y.reshape(-1)) gs = torch.autograd.grad(obj, all_params, allow_unused=True) return [g if g is not None else None for g in gs] SIG0 = None BGEN = torch.Generator().manual_seed(args.seed + 990) # separate RNG: sign flips must not shift the data stream GOV = {'K': None, 'bscale': 1.0, 'gema': None, 'drift': 0.0, 'gn': 0.0, 'sig': 0.0} WSNAP = {'p': None, 'o': None} def _clone_state(sd): if torch.is_tensor(sd): return sd.clone() if isinstance(sd, dict): return {k: _clone_state(v) for k, v in sd.items()} if isinstance(sd, list): return [_clone_state(v) for v in sd] return sd def ep_step(x, y): """single-sided EP with a QUALITY-GOVERNED estimator: beta_t = beta0*bscale*sig0^2/sig^2, K = GOV['K'] fb rounds; guard = finiteness + drift + grad-norm sanity only.""" global SIG0 if GOV['K'] is None: GOV['K'] = args.K if GOV.get('step', 0) % args.sig_every == 0 or GOV.get('sig', 0) == 0: GOV['sig'] = ddp_bcast_scalar(tok_sigma()) # all ranks run it (keeps global-RNG lockstep); rank0's value wins GOV['step'] = GOV.get('step', 0) + 1 sig = GOV['sig'] if SIG0 is None: SIG0 = args.sig0 if args.sig0 > 0 else sig beta_t = args.beta * GOV['bscale'] * (SIG0 * SIG0) / max(sig * sig, 1e-9) if args.beta_fixed: beta_t = args.beta * GOV['bscale'] fl = args.beta_floor if args.bf_late > 0.0 and GOV.get('step', 0) >= args.bf_late_at: fl = args.bf_late if fl > 0.0: beta_t = max(beta_t, fl) if args.beta_ride > 1.0: # ride-v2(a): floor jumps must not compose with a pre-charged cap — rescale cap so the # EFFECTIVE beta is continuous across any floor change (the 0.09-at-20k bug, RESULT 37) pf = GOV.get('prev_floor') if pf is not None and fl != pf and pf > 0 and fl > 0: GOV['cap'] = min(max(GOV.get('cap', 1.0) * pf / fl, args.cap_floor), args.beta_ride) GOV['prev_floor'] = fl beta_t = beta_t * GOV.get('cap', 1.0) # wall-2 loop-gain cap OVERRIDES the floor (the ceiling # can sit below the floor near the wall; survival first) if args.bsign_rand and torch.rand((), generator=BGEN).item() < 0.5: beta_t = -beta_t EST = args.est if args.est_late and GOV['step'] >= args.est_late_at: EST = args.est_late CF = (EST == 'centered' and args.centfast) if CF: # doubled batch [x;x]: +beta half / -beta half share every kernel (holofast pattern) x_in, y_in = torch.cat([x, x], 0), torch.cat([y, y], 0) bmask = torch.ones(x_in.shape[0], 1, 1, device=dev); bmask[args.B:] = -1.0 halves = (slice(0, args.B), slice(args.B, None)) else: x_in, y_in, bmask, halves = x, y, None, (slice(None),) z0, zs, ins, outs = free_states_graphed(x_in) zs_free = [z.clone() for z in zs] free_ce = F.cross_entropy(readout(zs_free[-1][:args.B]).reshape(-1, vocab), y.reshape(-1)).item() zp, last_outs = relax(z0, zs, ins, outs, y_in, +beta_t, GOV['K'], x_in, bmask=bmask) def _drift(zp_, zf_): dr = 0.0 for h in halves: # per-half worst drift == sequential guard decisions (max over passes) num = sum(float((a[h] - b[h]).norm()) for a, b in zip(zp_, zf_)) den = sum(float(b[h].norm()) for b in zf_) dr = max(dr, num / max(den, 1e-9)) return dr with torch.no_grad(): drift = _drift(zp, zs_free) gdrift = ddp_max_scalar(drift) # guard DECISIONS on the global worst -> identical on every rank dthr = 0.5 if args.drift_adapt > 0: # ride-v2(d): scale-free adaptive ceiling — reject anything far above the trailing # ACCEPTED drift level (garbage is O(1); healthy drift scales with beta) de = GOV.get('drift_ema') if de is not None: dthr = min(0.5, max(0.05, args.drift_adapt * de)) def _legal(gd): # SYNCHRONOUS acceptance (--beta_sync): judge THIS step by THIS step's own relax # telemetry — an out-of-window beta can be attempted but can never COMMIT. if (not math.isfinite(gd)) or (gd > dthr and not args.noguard): return False if args.beta_sync > 0 or args.wsync > 0 or args.beta_cap_rho > 0: # if you measure the ceiling, illegality COUNTS — crown-3's poison entered # through this gate being wired only to beta_sync (18k accepted garbage steps # with the rho-meter screaming); the meter and the gate are now永久 connected. res_s = ddp_bcast_scalar(GOV.get('res', 0.0)) rho_s = ddp_bcast_scalar(GOV.get('rho') or 0.0) if res_s > 0.02 and rho_s > args.beta_cap_rho: return False return True if not _legal(gdrift): ok_retry = False if args.wsync > 0 and not args.noguard and WSNAP['p'] is not None: # SYNCHRONOUS WEIGHT-STEP ACCEPTANCE: this state (= last opt.step's result) # failed the gate -> the UPDATE was illegal. Halve it in weight space # (p <- (p+snap)/2, delta implicit) and re-measure the same batch; after # wsync halvings, revert fully (momentum too) and skip. Same idiom and same # bounds as beta_sync — no new constants; the trajectory cannot dwell # past the ceiling. for _h in range(args.wsync + 1): with torch.no_grad(): if _h < args.wsync: for p, s in zip(all_params, WSNAP['p']): p.copy_((p + s) * 0.5) else: for p, s in zip(all_params, WSNAP['p']): p.copy_(s) opt.load_state_dict(WSNAP['o']) GOV['skr'] = GOV.get('skr', 0) + 1 z0, zs, ins, outs = free_states_graphed(x_in) zs_free = [z.clone() for z in zs] zp, last_outs = relax(z0, zs, ins, outs, y_in, +beta_t, GOV['K'], x_in, bmask=bmask) with torch.no_grad(): drift = _drift(zp, zs_free) gdrift = ddp_max_scalar(drift) if _legal(gdrift): ok_retry = True break if not ok_retry and args.beta_sync > 0 and not args.noguard: for _h in range(args.beta_sync): # halve beta, retry the SAME batch beta_t = beta_t * 0.5 GOV['skr'] = GOV.get('skr', 0) + 1 z0, zs, ins, outs = free_states_graphed(x_in) zs_free = [z.clone() for z in zs] zp, last_outs = relax(z0, zs, ins, outs, y_in, +beta_t, GOV['K'], x_in, bmask=bmask) with torch.no_grad(): drift = _drift(zp, zs_free) gdrift = ddp_max_scalar(drift) if _legal(gdrift): ok_retry = True if args.beta_ride > 1.0: # the failed trial IS the ceiling measurement GOV['cap'] = max(GOV.get('cap', 1.0) * 0.5 ** (_h + 1), args.cap_floor) break if ok_retry: pass if not ok_retry and args.kretry > 0 and math.isfinite(gdrift) and not args.noguard: GOV['skr'] = GOV.get('skr', 0) + 1 # marginal batch: retry once with deeper relaxation z0, zs, ins, outs = free_states_graphed(x_in) zs_free = [z.clone() for z in zs] zp, last_outs = relax(z0, zs, ins, outs, y_in, +beta_t, args.kretry, x_in, bmask=bmask) with torch.no_grad(): drift = _drift(zp, zs_free) gdrift = ddp_max_scalar(drift) ok_retry = _legal(gdrift) if not ok_retry: GOV['skd'] = GOV.get('skd', 0) + 1 # drift-guard reject (relaxation non-convergence) for p in all_params: p.grad = None return free_ce, beta_t, GOV.get('kuse', GOV['K']), False GOV['drift'] = gdrift if args.drift_adapt > 0: GOV['drift_ema'] = 0.95 * GOV.get('drift_ema', gdrift) + 0.05 * gdrift if args.beta_cap_rho > 0 and GOV.get('rho') is not None: rho_g = ddp_bcast_scalar(GOV['rho']) # rank0's meter rules (identical control on all ranks) res_g = ddp_bcast_scalar(GOV.get('res', 0.0)) # v2: ABSOLUTE-SCALE GATE — rho is only meaningful when the residual is above the noise # floor; at tiny residuals rho ~ noise/noise ~ 1 and v1 starved beta to the cap floor. rho_use = rho_g if args.ride_ema > 0: # ride-v2(b), opt-in: smooth the meter before decisions rho_use = GOV['rho_ema'] = 0.9 * GOV.get('rho_ema', rho_g) + 0.1 * rho_g GOV['cool'] = max(GOV.get('cool', 0) - 1, 0) if args.beta_servo > 0 and res_g > 0.02 and rho_use > 0: # meter lit -> INVERT onto the ceiling (both directions), replacing the AIMD attack GOV['cap'] = min(max(GOV.get('cap', 1.0) * (args.beta_servo * args.beta_cap_rho / rho_use), args.cap_floor), args.beta_ride) elif res_g > 0.02 and rho_use > args.beta_cap_rho: GOV['cap'] = max(GOV.get('cap', 1.0) * 0.85, args.cap_floor) # attack (gentler than v1) if args.ride_cool > 0: GOV['cool'] = args.ride_cool # ride-v2(c), opt-in elif (res_g < 0.01 or rho_use < 0.5 * args.beta_cap_rho) and GOV['cool'] == 0: # recover; with beta_ride > 1 the governor CLIMBS past the schedule — beta finds # its own ceiling and hovers there (ride-the-ceiling; 1.0 = legacy defensive cap) GOV['cap'] = min(GOV.get('cap', 1.0) * args.beta_ride_up, args.beta_ride) if CF: # one-graph centered: [g(+b)+g(-b)]/2 = d[(E+ - E-)/(2b·NBT)]/dtheta; CE-head term from # the +beta half only (matches sequential centered's gsC at the +beta top states). Ec = 0.0 for z, o in zip(zp, last_outs): df = z.detach().float() - o.float() Ec = Ec + 0.5 * (df[:args.B] ** 2).sum() - 0.5 * (df[args.B:] ** 2).sum() obj = Ec / (NBT * 2.0 * beta_t) + obj_loss(readout(zp[-1][:args.B].detach()).reshape(-1, vocab), y.reshape(-1)) gs = torch.autograd.grad(obj, all_params, allow_unused=True) elif EST == 'single': E = 0.0 for z, o in zip(zp, last_outs): E = E + 0.5 * ((z.detach().float() - o.float()) ** 2).sum() # fp32 accumulation (bf16-safe; no-op in fp32) obj = E / (NBT * beta_t) + obj_loss(readout(zp[-1].detach()).reshape(-1, vocab), y.reshape(-1)) gs = torch.autograd.grad(obj, all_params, allow_unused=True) else: # two-pass estimators: g(b) := d[E(b)]/dtheta / (NBT*b) => single-sided bias g_true + c*b. # centered: [g(+b) + g(-b)] / 2 (1/b sign inside => average cancels c*b) # richardson: 2*g(b) - g(2b) (extrapolation cancels c*b at large b) E = 0.0 for z, o in zip(zp, last_outs): E = E + 0.5 * ((z.detach().float() - o.float()) ** 2).sum() gsE = torch.autograd.grad(E / (NBT * beta_t), all_params, allow_unused=True) gsC = torch.autograd.grad(obj_loss(readout(zp[-1].detach()).reshape(-1, vocab), y.reshape(-1)), all_params, allow_unused=True) b2 = -beta_t if EST == 'centered' else 2.0 * beta_t if EST == 'centered' and args.centmirror: # MIRROR WARM-START: d-(free anchor) = -d+ exactly (linear in beta); init the -beta # states as the mirror of the settled +beta solution, then ONE polish sweep corrects # the O(beta^2) even part. Skips the second free pass and K-1 sweeps. dm = [(-di).detach() for di in GOV['_last_d']] zsb_free = zs_free prev = z0 zsb, insb, outsb = [], [], [] with torch.autocast('cuda', dtype=torch.bfloat16, enabled=args.amp): for l in range(args.L): i = prev.detach().requires_grad_(True) o = blocks[l](i, mask) zsb.append(o.detach().float() + dm[l]) insb.append(i); outsb.append(o) prev = zsb[l] _k1 = GOV.get('kuse') zpb, lob = relax(z0, zsb, insb, outsb, y, b2, 1, x) GOV['kuse'] = _k1 # telemetry: report the +beta pass's K, not the mirror polish else: z0b, zsb, insb, outsb = free_states_graphed(x) zsb_free = [z.clone() for z in zsb] zpb, lob = relax(z0b, zsb, insb, outsb, y, b2, GOV['K'], x) with torch.no_grad(): drift2 = sum(float((a - b).norm()) for a, b in zip(zpb, zsb_free)) / max( sum(float(b.norm()) for b in zsb_free), 1e-9) gdrift2 = ddp_max_scalar(drift2) if (not math.isfinite(gdrift2)) or (gdrift2 > 0.5 and not args.noguard): GOV['skd'] = GOV.get('skd', 0) + 1 # second-pass drift reject -> skip step (synced) for p in all_params: p.grad = None return free_ce, beta_t, GOV.get('kuse', GOV['K']), False E2 = 0.0 for z, o in zip(zpb, lob): E2 = E2 + 0.5 * ((z.detach().float() - o.float()) ** 2).sum() gsE2 = torch.autograd.grad(E2 / (NBT * b2), all_params, allow_unused=True) def _comb(a, b): if a is None and b is None: return None a = a if a is not None else torch.zeros_like(b) b = b if b is not None else torch.zeros_like(a) return (a + b) / 2.0 if EST == 'centered' else (2.0 * a - b) gs = [(_comb(e, e2) if (e is not None or e2 is not None) else None) for e, e2 in zip(gsE, gsE2)] gs = [ (g if g is not None else c) if c is None or g is None else g + c for g, c in zip(gs, gsC) ] gs = ddp_avg(gs, all_params) # global-batch gradient; gn/gema/guard below see identical values on all ranks gn = 0.0 for g in gs: if g is not None: gn += float((g ** 2).sum()) gn = gn ** 0.5 if GOV['gema'] is None: GOV['gema'] = gn GOV['gema'] = 0.99 * GOV['gema'] + 0.01 * gn # EMA always updates (frozen-ref bugfix) GOV['gn'] = gn if not math.isfinite(gn) or (gn > 8 * GOV['gema'] and not args.noguard): GOV['skg'] = GOV.get('skg', 0) + 1 # gn-EMA-guard reject (gradient-magnitude spike) for p in all_params: p.grad = None return free_ce, beta_t, GOV.get('kuse', GOV['K']), False for p, g in zip(all_params, gs): p.grad = g return free_ce, beta_t, GOV.get('kuse', GOV['K']), True def bp_gate(x, y): """true BP grads for telemetry cos (called before opt.step; reads p.grad separately).""" z = emb(x) for b in blocks: z = b(z, mask) ce = obj_loss(readout(z).reshape(-1, vocab), y.reshape(-1)) return ddp_avg(list(torch.autograd.grad(ce, all_params, allow_unused=True)), all_params) @torch.no_grad() def evaluate(nb=6): tot = 0.0 for _ in range(nb): x, y = get_batch('val') z = emb(x) for b in blocks: z = b(z, mask) tot += F.cross_entropy(readout(z).reshape(-1, vocab), y.reshape(-1)).item() return tot / nb if args.resume and _ck.get('opt') is not None: try: opt.load_state_dict(_ck['opt']) print('[resume] optimizer state restored (exact chunked-resume)', flush=True) except Exception as e: print(f'[resume] optimizer state NOT restored ({e}) — cold optimizer', flush=True) if DDP: # belt & suspenders on top of identical init seeds: rank0's params are law with torch.no_grad(): for p in all_params: if args.ddp_backend == 'gloo': t = p.data.cpu(); dist.broadcast(t, 0); p.data.copy_(t) else: dist.broadcast(p.data, 0) if RANK == 0: print(f'[ddp] world={WORLD} backend={args.ddp_backend} params broadcast; eff batch {args.B}x{WORLD}={args.B*WORLD}', flush=True) wb = None if args.wandb == 'auto': args.wandb = 'ept-fineweb-72m' if 'fineweb' in args.data else 'ept-tinystories-42m' if args.wandb and RANK == 0: try: import wandb as _w wb = _w.init(entity='eqprop-llm-training', project=args.wandb, name=args.wandb_run or args.tag, id=args.wandb_run or args.tag, resume='allow', config=vars(args)) except Exception as e: print(f'[wandb] disabled ({e})', flush=True) n = sum(p.numel() for p in all_params) if RANK == 0: print(f'[{args.tag}] cascade-EP(EQUILIBRIUM/fb) L{args.L} C{args.C} T{args.T} beta={args.beta} ' f'K={args.K} geta={args.geta} | {n/1e6:.2f}M | {dev}', flush=True) if args.ddp_grad_test: # one-step equivalence: DDP(WORLD ranks x B) averaged grad must equal single-GPU grad on the # SAME WORLD*B batch (exact algebra: per-sample-independent relaxation + mean-linear readout). # Protocol: run WORLD=1 with --B (W*B) first, then torchrun WORLD=N with --B B; both seed 4242. _g = torch.Generator().manual_seed(4242) _data = np.memmap(DD / 'train.bin', dtype=np.uint16, mode='r') _full = torch.randint(len(_data) - args.T - 1, (WORLD * args.B,), generator=_g) _ix = _full[RANK * args.B:(RANK + 1) * args.B] _x = torch.stack([torch.from_numpy(_data[i:i + args.T].astype(np.int64)) for i in _ix]).to(dev) _y = torch.stack([torch.from_numpy(_data[i + 1:i + 1 + args.T].astype(np.int64)) for i in _ix]).to(dev) _ce, _bt, _r, _ok = ep_step(_x, _y) assert _ok, 'grad test: ep_step guarded' _flat = torch.cat([(p.grad if p.grad is not None else torch.zeros_like(p)).reshape(-1).double().cpu() for p in all_params]) if RANK == 0: _f = Path('runs') / f'ddp_grad_w{WORLD}.pt' torch.save({'flat': _flat, 'beta': _bt, 'W': WORLD, 'B': args.B}, _f) print(f'[gradtest] W={WORLD} B/rank={args.B} beta_t={_bt:.3e} ce={_ce:.4f} saved {_f}', flush=True) _ref = Path('runs') / 'ddp_grad_w1.pt' if WORLD > 1 and _ref.exists(): _r1 = torch.load(_ref, weights_only=False) assert _r1['B'] == WORLD * args.B, f"ref B={_r1['B']} != {WORLD*args.B}" _rf = _r1['flat'] _cos = float((_flat @ _rf) / (_flat.norm() * _rf.norm())) _rel = float((_flat - _rf).norm() / _rf.norm()) print(f'[gradtest] VERDICT cos={_cos:.9f} relerr={_rel:.2e} (DDP avg vs single-GPU big-batch)', flush=True) import sys sys.exit(0) best, t0 = 1e9, time.time() skips = 0 for _ in range(start_step): sched.step() # advance LR schedule to the resumed step for step in range(start_step, args.steps + 1): x, y = get_batch('train') if args.qcomp_bits > 0: # STAGE-0 HW GATE #1b (T64 scenario): COMPUTE runs on weights snapped to the DAC # grid (deterministic round-to-nearest); the fp32 master (DDR / shadow accumulator) # receives the update. Equivalent to word-streaming and to resident-cell + shadow. with torch.no_grad(): QSAVE = [p.detach().clone() for p in all_params] for p in all_params: rng = float(p.abs().max()) if rng <= 0: continue g_ = rng / (2 ** (args.qcomp_bits - 1)) p.copy_((p / g_).round() * g_) ce, beta_t, rounds, ok = ep_step(x, y) if args.qcomp_bits > 0: with torch.no_grad(): for p, q in zip(all_params, QSAVE): p.copy_(q) if not ok: skips += 1 gcos = float('nan') if args.gate_every > 0 and step % args.gate_every == 0 and ok: gbp = bp_gate(x, y) num = den1 = den2 = 0.0 for p, g in zip(all_params, gbp): if p.grad is None or g is None: continue num += float((p.grad * g).sum()); den1 += float((p.grad ** 2).sum()); den2 += float((g ** 2).sum()) gcos = num / max((den1 ** 0.5) * (den2 ** 0.5), 1e-12) if args.gate_govern: # opt-in: BP-informed control flow if gcos < 0.97: GOV['K'] = min(GOV['K'] + 2, args.kmax); GOV['bscale'] = max(GOV['bscale'] * 0.7, 0.05) elif gcos > 0.995 and GOV['K'] > args.K: GOV['K'] -= 1; GOV['bscale'] = min(GOV['bscale'] * 1.05, 1.0) if args.bpmix and ok: with torch.enable_grad(): _bg = bp_gate(x, y) # true BP grads, same batch sel = set() for spec in args.bpmix.split(','): spec = spec.strip() if spec.startswith('blocks:'): a_, b_ = spec.split(':')[1].split('-') for bi in range(int(a_), int(b_) + 1): sel.update(id(p) for p in blocks[bi].parameters()) elif spec == 'attn': for blk in blocks: sel.update(id(p) for m in (blk.attn, blk.na) for p in m.parameters()) elif spec == 'ffn': for blk in blocks: sel.update(id(p) for m in (blk.ff, blk.nf) for p in m.parameters()) elif spec == 'head': sel.update(id(p) for p in (list(tok.parameters()) + ([W_out] if isinstance(W_out, torch.nn.Parameter) else []) + list(ln_f.parameters()))) for p, g in zip(all_params, _bg): if id(p) in sel and g is not None: p.grad = g.detach().clone() torch.nn.utils.clip_grad_norm_(all_params, 1.0) if args.wsync > 0: # snapshot the KNOWN-LEGAL pre-step state (this step's relax passed _legal); # next step's relax measures the post-step state and can roll back to here. with torch.no_grad(): WSNAP['p'] = [p.detach().clone() for p in all_params] WSNAP['o'] = _clone_state(opt.state_dict()) opt.step(); sched.step(); opt.zero_grad(set_to_none=True) if args.qup_bits > 0: # STAGE-0 HW GATE: finite conductance levels. Snap every weight to an ABSOLUTE # per-tensor grid (range/2^bits) with stochastic rounding (unbiased) — emulates # analog cell writes; per-step deltas below one level survive only in expectation. with torch.no_grad(): for p in all_params: if p.ndim < 1: continue rng = float(p.abs().max()) if rng <= 0: continue g_ = rng / (2 ** (args.qup_bits - 1)) q = p / g_ fl = q.floor() p.copy_((fl + (torch.rand_like(p) < (q - fl)).float()) * g_) if DDP and args.sync_check > 0 and step % args.sync_check == 0 and step > 0: with torch.no_grad(): h = torch.stack([torch.stack((p.double().sum(), (p.double() ** 2).sum())) for p in all_params]).sum(0) hc = h.cpu() if args.ddp_backend == 'gloo' else h hs = [torch.zeros_like(hc) for _ in range(WORLD)] dist.all_gather(hs, hc) if any(bool((x != hs[0]).any()) for x in hs[1:]): print(f'[ddp] PARAM DESYNC step {step} rank {RANK}: {[x.tolist() for x in hs]}', flush=True) raise RuntimeError('DDP param desync — aborting rather than training garbage') if step % args.log == 0 and RANK == 0: val = evaluate(); best = min(best, val) gtag = '' if math.isnan(gcos) else f' cos={gcos:.4f}' print(f'step {step:5d}/{args.steps} | train {ce:.4f} val {val:.4f} (best {best:.4f}) ' f'| beta={beta_t:.2e} K={rounds} skips={skips}(d{GOV.get("skd",0)}/g{GOV.get("skg",0)}/r{GOV.get("skr",0)}){gtag} ' f'drift={GOV["drift"]:.3f} gn={GOV["gn"]:.2e} sig={GOV["sig"]:.1f} | {step/max(time.time()-t0,1e-9):.3f} it/s', flush=True) if wb is not None: try: wb.log({'train_ce': ce, 'val_ce': val, 'best': best, 'beta_t': beta_t, 'rounds': rounds, 'skips': skips, 'gate_cos': (None if math.isnan(gcos) else gcos)}, step=step) except Exception: pass if (step % args.save_every == 0 or step == args.steps) and step > 0 and RANK == 0: torch.save({'tok': tok.state_dict(), 'pos': pos.state_dict(), 'blocks': blocks.state_dict(), 'wout': (W_out.detach().cpu() if args.untie else None), 'lnf': (ln_f.state_dict() if not isinstance(ln_f, nn.Identity) else None), 'opt': opt.state_dict(), # full optimizer state -> exact resume for chunked HPC jobs 'step': step, 'val': best, 'config': vars(args)}, Path('runs') / f'{args.tag}_s{step}.pt') if RANK == 0: print(f'[{args.tag}] DONE best val CE {best:.4f}', flush=True) if DDP: dist.destroy_process_group() if wb is not None: try: wb.summary['best_val_ce'] = best; wb.finish() except Exception: pass