"""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=''); 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('--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 args = ap.parse_args() torch.manual_seed(args.seed) dev = 'cuda' if torch.cuda.is_available() else 'cpu' DD = Path('/home/yurenh2/ept/ep_run/data/tinystories_bpe') 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,)) 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 Block(nn.Module): def __init__(self, C, H): super().__init__() self.ln1, self.ln2 = nn.LayerNorm(C), nn.LayerNorm(C) self.attn = nn.MultiheadAttention(C, H, batch_first=True) self.ff = nn.Sequential(nn.Linear(C, 4 * C), nn.GELU(), nn.Linear(4 * C, C)) def forward(self, z, mask): h = self.ln1(z); a, _ = self.attn(h, h, h, attn_mask=mask, need_weights=False) z = z + a; return z + self.ff(self.ln2(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([Block(args.C, args.H) for _ in range(args.L)]).to(dev) 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 readout = (lambda z: z @ W_out.t()) if args.untie else (lambda z: z @ tok.weight.t()) all_params = list(tok.parameters()) + list(pos.parameters()) + list(blocks.parameters()) + ([W_out] if args.untie else []) if args.opt == 'muon': from muon import build_hybrid opt, sched = build_hybrid(blocks, all_params, args.lr, args.muon_lr, args.warmup) else: opt = torch.optim.AdamW(all_params, lr=args.lr, weight_decay=1e-4) sched = torch.optim.lr_scheduler.LambdaLR(opt, lambda s: min(1.0, (s + 1) / max(args.warmup, 1))) NBT = args.B * args.T 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 = (tok(x) + pos(torch.arange(args.T, device=dev))[None]) ins, outs, zs = [], [], [] prev = z0 for b in blocks: i = prev.detach().requires_grad_(True) o = b(i, mask) ins.append(i); outs.append(o); zs.append(o.detach()) 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 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): """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.""" d = [None] * args.L for k in range(K): if k % args.dtop_every == 0 or d[args.L - 1] is None: zc = zs[args.L - 1].detach().requires_grad_(True) ce = F.cross_entropy(readout(zc).reshape(-1, vocab), y.reshape(-1)) d[args.L - 1] = (-beta * NBT * torch.autograd.grad(ce, zc)[0]).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])[0].detach() last = (k + 1 == K) prev = z0 n_ins, n_outs = [], [] for l in range(args.L): if last and l == 0: i = tok(x) + pos(torch.arange(args.T, device=dev))[None] # graphed emb for the readout's E-path else: i = prev.detach().requires_grad_(True) o = blocks[l](i, mask) zs[l] = (o.detach() + d[l]) n_ins.append(i); n_outs.append(o) prev = zs[l] ins, outs = n_ins, n_outs return zs, outs def dFdtheta(zs, x, y, beta): """dF/dtheta at fixed relaxed states (z0 rebuilt WITH graph so emb gets its E-path grad).""" prev = tok(x) + pos(torch.arange(args.T, device=dev))[None] E = 0.0 for z, b in zip(zs, blocks): E = E + 0.5 * ((z - b(prev, mask)) ** 2).sum(); prev = z 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 GOV = {'K': None, 'bscale': 1.0, 'gema': None, 'drift': 0.0, 'gn': 0.0, 'sig': 0.0} 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'] = tok_sigma() GOV['step'] = GOV.get('step', 0) + 1 sig = GOV['sig'] if SIG0 is None: SIG0 = sig beta_t = args.beta * GOV['bscale'] * (SIG0 * SIG0) / max(sig * sig, 1e-9) z0, zs, ins, outs = free_states_graphed(x) zs_free = [z.clone() for z in zs] free_ce = F.cross_entropy(readout(zs_free[-1]).reshape(-1, vocab), y.reshape(-1)).item() zp, last_outs = relax(z0, zs, ins, outs, y, +beta_t, GOV['K'], x) with torch.no_grad(): drift = sum(float((a - b).norm()) for a, b in zip(zp, zs_free)) / max( sum(float(b.norm()) for b in zs_free), 1e-9) if (not math.isfinite(drift)) or (drift > 0.5 and not args.noguard): for p in all_params: p.grad = None return free_ce, beta_t, GOV['K'], False GOV['drift'] = drift E = 0.0 for z, o in zip(zp, last_outs): E = E + 0.5 * ((z.detach() - o) ** 2).sum() obj = E / (NBT * beta_t) + F.cross_entropy(readout(zp[-1].detach()).reshape(-1, vocab), y.reshape(-1)) gs = torch.autograd.grad(obj, all_params, allow_unused=True) 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): for p in all_params: p.grad = None return free_ce, beta_t, GOV['K'], False for p, g in zip(all_params, gs): p.grad = g return free_ce, beta_t, GOV['K'], True def bp_gate(x, y): """true BP grads for telemetry cos (called before opt.step; reads p.grad separately).""" z = tok(x) + pos(torch.arange(args.T, device=dev))[None] for b in blocks: z = b(z, mask) ce = F.cross_entropy(readout(z).reshape(-1, vocab), y.reshape(-1)) return list(torch.autograd.grad(ce, all_params, allow_unused=True)) @torch.no_grad() def evaluate(nb=6): tot = 0.0 for _ in range(nb): x, y = get_batch('val') z = tok(x) + pos(torch.arange(args.T, device=dev))[None] for b in blocks: z = b(z, mask) tot += F.cross_entropy(readout(z).reshape(-1, vocab), y.reshape(-1)).item() return tot / nb wb = None if args.wandb: try: import wandb as _w wb = _w.init(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) 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) best, t0 = 1e9, time.time() skips = 0 for step in range(args.steps + 1): x, y = get_batch('train') ce, beta_t, rounds, ok = ep_step(x, y) if not ok: skips += 1 gcos = float('nan') if 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 gcos < 0.97: # estimator governor: spend more 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: # relax back when quality is abundant GOV['K'] -= 1; GOV['bscale'] = min(GOV['bscale'] * 1.05, 1.0) torch.nn.utils.clip_grad_norm_(all_params, 1.0) opt.step(); sched.step(); opt.zero_grad(set_to_none=True) if step % args.log == 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}{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 and step > 0: torch.save({'tok': tok.state_dict(), 'pos': pos.state_dict(), 'blocks': blocks.state_dict(), 'step': step, 'val': best, 'config': vars(args)}, Path('runs') / f'{args.tag}_s{step}.pt') print(f'[{args.tag}] DONE best val CE {best:.4f} (BP twin 2.9746; zil-diagnostic 3.3236)', flush=True) if wb is not None: try: wb.summary['best_val_ce'] = best; wb.finish() except Exception: pass