"""E2 stage A (item 21): dense short-CoT supervision through the carry whiteboard on GSM8K. train_carry.py skeleton, one change that matters: CE targets are the model's own VERIFIED terse scratchpad + answer (gsm_cot_data.json, ~30-60 tokens) instead of the ~3-token bare answer — the dense-output ingredient the latent GSM arms structurally lacked. Control: --feedforward = same supervision, adapter(e,e), no recurrence. """ import argparse import json import os import math import random import sys import time from pathlib import Path import torch import torch.nn.functional as F from carry_common import PAUSE_ID, carry_logits from loop_common import BandLooper, MergeAdapter, chat_prompt, DIRECT_SUFFIX sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) from jlens.core import load_model # noqa: E402 OUT = Path(os.environ.get("LOOP_OUT", Path(__file__).resolve().parent.parent / "results-loop")) STEPS = 600 BATCH = 4 LR = 1e-3 WARMUP = 20 K_PREFILL = 2 P_BY_LABEL = {"easy": 2, "hard": 6} ap = argparse.ArgumentParser() ap.add_argument("--feedforward", action="store_true") ap.add_argument("--seed", type=int, default=0) ap.add_argument("--lensnoise", default=None, metavar="RANK,SCALE", help="E2-N1: inject noise into the carried state during " "training, shaped by the top-RANK sensitivity " "directions of jbar at the band entrance; noise norm " "= SCALE * per-position state norm (e.g. 32,0.05)") ap.add_argument("--drop-steps", type=int, default=0, metavar="D", help="E2-L rung B: delete the first D scratchpad steps, " "each replaced by --pause-per-step extra pauses") ap.add_argument("--pause-per-step", type=int, default=10, help="pauses per deleted step (median step = 10 tokens)") ap.add_argument("--warm-start", default=None, metavar="ADAPTER_PT") ap.add_argument("--steps", type=int, default=STEPS) ap.add_argument("--lr", type=float, default=LR) ARGS = ap.parse_args() STEPS, LR = ARGS.steps, ARGS.lr TAG = ("carrycot_ff" if ARGS.feedforward else "carrycot") + ( f"_b{ARGS.drop_steps}" if ARGS.drop_steps else "") + ( f"_ln{ARGS.lensnoise.replace(',', '_')}" if ARGS.lensnoise else "") + ( f"_s{ARGS.seed}" if ARGS.seed else "") def drop_cot_steps(cot, d): """Delete the first d scratchpad lines (front-first: the deleted computation must ride the pause-chain before the visible remainder). Returns (new_cot, n_deleted); unparseable cots pass through intact.""" lines = [l for l in cot.split("\n") if l.strip()] ans = [l for l in lines if l.startswith("Answer")] steps = [l for l in lines if not l.startswith("Answer")] if len(ans) != 1 or not steps: return cot, 0 n = min(d, len(steps)) return "\n".join(steps[n:] + ans), n class LensNoiseWrapper(torch.nn.Module): """Perturb the carried state s (not the anchor e) before the merge, within the span of the lens's top-r readout-sensitive directions. Train-time only (noise_on flag); eval and checkpoints use .base.""" def __init__(self, base, jbar_layer, rank, scale): super().__init__() self.base = base self.scale = scale self.noise_on = True J = jbar_layer.float() _, _, Vt = torch.linalg.svd(J, full_matrices=False) self.register_buffer("V", Vt[:rank].T.contiguous()) # (d, r) def forward(self, e, s): if self.noise_on and self.scale > 0: z = torch.randn(*s.shape[:-1], self.V.shape[1], device=s.device, dtype=torch.float32) n = z @ self.V.T n = n * (s.float().norm(dim=-1, keepdim=True) * self.scale / (n.norm(dim=-1, keepdim=True) + 1e-6)) s = (s.float() + n).to(s.dtype) return self.base(e, s) def lr_at(step): if step < WARMUP: return LR * (step + 1) / WARMUP t = (step - WARMUP) / max(1, STEPS - WARMUP) return 1e-4 + 0.5 * (LR - 1e-4) * (1 + math.cos(math.pi * t)) def build_batch(tok, items, p, device="cuda"): seqs, labs, plens = [], [], [] for it in items: pr = tok(chat_prompt(tok, it["question"], DIRECT_SUFFIX), add_special_tokens=False)["input_ids"] a = tok(it["cot"] + "", add_special_tokens=False)["input_ids"] pi = p + it.get("extra_pauses", 0) seqs.append(pr + [PAUSE_ID] * pi + a) labs.append([-100] * (len(pr) + pi) + a) plens.append(len(pr)) T = max(len(s) for s in seqs) pad = tok.pad_token_id or 0 ids = torch.full((len(seqs), T), pad, dtype=torch.long) lab = torch.full((len(seqs), T), -100, dtype=torch.long) msk = torch.zeros((len(seqs), T), dtype=torch.long) for i, (s, l) in enumerate(zip(seqs, labs)): ids[i, : len(s)] = torch.tensor(s) lab[i, : len(s)] = torch.tensor(l) msk[i, : len(s)] = 1 return (ids.to(device), msk.to(device), lab.to(device), torch.tensor(plens, device=device)) @torch.no_grad() def val_loss(looper, adapter, tok, items, p): tot, n = 0.0, 0 for i in range(0, len(items), BATCH): ids, msk, lab, plens = build_batch(tok, items[i : i + BATCH], p) logits = carry_logits(looper, adapter, ids, msk, plens, K_PREFILL, feedforward=ARGS.feedforward) loss = F.cross_entropy(logits[:, :-1].flatten(0, 1).float(), lab[:, 1:].flatten(), ignore_index=-100) tot += loss.item() * len(ids) n += len(ids) return tot / n def main(): rng = random.Random(ARGS.seed) torch.manual_seed(ARGS.seed) star = {it["idx"]: it for it in json.load(open(OUT / "star_data.json")) if it["split"] == "train"} cots = json.load(open(OUT / "gsm_cot_data.json")) data = [] n_dropped = 0 for r in cots: it = star.get(r["idx"]) if it is None: continue cot, ndel = (drop_cot_steps(r["cot"], ARGS.drop_steps) if ARGS.drop_steps else (r["cot"], 0)) n_dropped += ndel data.append({"question": it["question"], "label": r["label"], "cot": cot, "extra_pauses": ndel * ARGS.pause_per_step}) if ARGS.drop_steps: print(f"rung B d={ARGS.drop_steps}: {n_dropped} steps deleted " f"across {len(data)} items " f"({ARGS.pause_per_step} pauses each)", flush=True) model, tok = load_model(dtype=torch.bfloat16) for pp in model.parameters(): pp.requires_grad_(False) looper = BandLooper(model) adapter = MergeAdapter( d=model.config.get_text_config().hidden_size).cuda() if ARGS.lensnoise: r, sc = ARGS.lensnoise.split(",") jbar = torch.load(Path(__file__).resolve().parent.parent / "results/jbar.pt", map_location="cpu")["Jbar"] from loop_common import BAND adapter = LensNoiseWrapper(adapter, jbar[BAND[0]], int(r), float(sc)).cuda() print(f"lens-noise: rank={r} scale={sc} on jbar L{BAND[0]}", flush=True) if ARGS.warm_start: (adapter.base if ARGS.lensnoise else adapter).load_state_dict( torch.load(ARGS.warm_start, map_location="cuda")) print(f"warm-started from {ARGS.warm_start}", flush=True) opt = torch.optim.AdamW(adapter.parameters(), lr=LR, weight_decay=0.01) keep = [it for it in data if len(tok(it["question"])["input_ids"]) + len(tok(it["cot"])["input_ids"]) + it.get("extra_pauses", 0) + 16 <= 460] rng.shuffle(keep) pool = {l: [it for it in keep if it["label"] == l] for l in ("easy", "hard")} val = {l: pool[l][:12] for l in pool} pool = {l: pool[l][12:] for l in pool} print(f"pool: easy={len(pool['easy'])} hard={len(pool['hard'])}", flush=True) if min(len(v) for v in pool.values()) < BATCH: print("INSUFFICIENT POOL — aborting", flush=True) return log = [] t0 = time.time() for step in range(STEPS): lbl = ("easy", "hard")[step % 2] batch = rng.sample(pool[lbl], BATCH) p = P_BY_LABEL[lbl] ids, msk, lab, plens = build_batch(tok, batch, p) for g in opt.param_groups: g["lr"] = lr_at(step) logits = carry_logits(looper, adapter, ids, msk, plens, K_PREFILL, use_checkpoint=True, feedforward=ARGS.feedforward) loss = F.cross_entropy(logits[:, :-1].flatten(0, 1).float(), lab[:, 1:].flatten(), ignore_index=-100) opt.zero_grad(set_to_none=True) loss.backward() torch.nn.utils.clip_grad_norm_(adapter.parameters(), 1.0) opt.step() log.append({"step": step, "loss": loss.item()}) if step % 10 == 0: print(f"step {step:4d} {lbl:4s} loss={loss.item():.4f} " f"({(time.time()-t0)/(step+1):.1f}s/step)", flush=True) if step % 200 == 199 or step == STEPS - 1: if ARGS.lensnoise: adapter.noise_on = False for l in ("easy", "hard"): v = val_loss(looper, adapter, tok, val[l], P_BY_LABEL[l]) print(f" val@{step}: {l}={v:.3f}", flush=True) if ARGS.lensnoise: adapter.noise_on = True sd = (adapter.base if ARGS.lensnoise else adapter).state_dict() torch.save(sd, OUT / f"adapter_{TAG}_e{step+1}.pt") json.dump(log, open(OUT / f"train_{TAG}_log.json", "w")) print("done", flush=True) if __name__ == "__main__": main()