331 lines
14 KiB
Python
331 lines
14 KiB
Python
"""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)
|
|
ap.add_argument("--tag-suffix", default="",
|
|
help="appended to TAG (distinguish control variants)")
|
|
ap.add_argument("--bandlora", type=int, default=0, metavar="RANK",
|
|
help="item 24: loop-only LoRA (lora_band.LoopLoRA) on every "
|
|
"band layer, uniform scale 1.0 — active only during "
|
|
"band re-runs, k=0 stays bit-exact")
|
|
ap.add_argument("--lora-lr", type=float, default=1e-3)
|
|
ap.add_argument("--lensteach", type=float, default=0.0, metavar="LAMBDA",
|
|
help="item 25: latent process supervision — lens-CE at the "
|
|
"replacement pauses against the DELETED step's tokens "
|
|
"(1:1 pause j <-> step token j), mixed at LAMBDA into "
|
|
"the output CE")
|
|
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"_blr{ARGS.bandlora}" if ARGS.bandlora else "") + (
|
|
f"_lt{str(ARGS.lensteach).replace('.', '')}" if ARGS.lensteach else "") + (
|
|
f"_ln{ARGS.lensnoise.replace(',', '_')}" if ARGS.lensnoise else "") + (
|
|
f"_s{ARGS.seed}" if ARGS.seed else "") + ARGS.tag_suffix
|
|
|
|
|
|
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, deleted_text); 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, "\n".join(steps[: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"] + "<end_of_turn>",
|
|
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, deleted = (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, "deleted": deleted,
|
|
"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)
|
|
lora_params, lora_band_layers = [], []
|
|
if ARGS.bandlora:
|
|
from lora_band import inject_band_lora
|
|
from loop_common import BAND
|
|
lora_band_layers = list(range(BAND[0], BAND[1] + 1))
|
|
scales = {l: 1.0 for l in lora_band_layers}
|
|
lora_params = inject_band_lora(looper.tm, BAND[0], scales,
|
|
rank=ARGS.bandlora)
|
|
for p in lora_params:
|
|
p.data = p.data.cuda()
|
|
print(f"band-lora r={ARGS.bandlora}: "
|
|
f"{sum(p.numel() for p in lora_params)/1e6:.1f}M params, "
|
|
f"layers {lora_band_layers[0]}-{lora_band_layers[-1]}, "
|
|
f"lr={ARGS.lora_lr}", flush=True)
|
|
lens_teach = None
|
|
if ARGS.lensteach:
|
|
from loop_common import BAND
|
|
J30 = torch.load(Path(__file__).resolve().parent.parent
|
|
/ "results/jbar.pt",
|
|
map_location="cuda")["Jbar"][BAND[1]].float()
|
|
tm = looper.tm
|
|
softcap = model.config.get_text_config().final_logit_softcapping
|
|
|
|
def lens_teach(h):
|
|
proj = h.float() @ J30.T
|
|
x = tm.norm(proj.to(tm.norm.weight.dtype))
|
|
lg = model.lm_head(x)
|
|
if softcap:
|
|
lg = softcap * torch.tanh(lg / softcap)
|
|
return lg
|
|
|
|
n_tgt = 0
|
|
for it in data:
|
|
if it["deleted"]:
|
|
it["lens_targets"] = tok(
|
|
it["deleted"], add_special_tokens=False
|
|
)["input_ids"][: it["extra_pauses"]]
|
|
n_tgt += bool(it["lens_targets"])
|
|
print(f"lens-teach λ={ARGS.lensteach}: targets on {n_tgt} items "
|
|
f"(pause j <-> deleted-step token j, lens at L{BAND[1]})",
|
|
flush=True)
|
|
groups = [{"params": list(adapter.parameters()), "lr": LR, "base": LR}]
|
|
if lora_params:
|
|
groups.append({"params": lora_params, "lr": ARGS.lora_lr,
|
|
"base": ARGS.lora_lr})
|
|
opt = torch.optim.AdamW(groups, 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"] = g["base"] * lr_at(step) / LR
|
|
if ARGS.lensteach:
|
|
logits, S = carry_logits(looper, adapter, ids, msk, plens,
|
|
K_PREFILL, use_checkpoint=True,
|
|
feedforward=ARGS.feedforward,
|
|
return_states=True)
|
|
else:
|
|
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)
|
|
lce_val = 0.0
|
|
if ARGS.lensteach:
|
|
terms = []
|
|
for b, it in enumerate(batch):
|
|
tgt = it.get("lens_targets")
|
|
if not tgt:
|
|
continue
|
|
s0 = int(plens[b]) + p
|
|
hs = S[b, s0 : s0 + len(tgt)]
|
|
terms.append(F.cross_entropy(
|
|
lens_teach(hs).float(),
|
|
torch.tensor(tgt, device=hs.device)))
|
|
if terms:
|
|
lce = torch.stack(terms).mean()
|
|
lce_val = lce.item()
|
|
loss = loss + ARGS.lensteach * lce
|
|
opt.zero_grad(set_to_none=True)
|
|
loss.backward()
|
|
torch.nn.utils.clip_grad_norm_(
|
|
list(adapter.parameters()) + lora_params, 1.0)
|
|
opt.step()
|
|
log.append({"step": step, "loss": loss.item(), "lce": lce_val})
|
|
if step % 10 == 0:
|
|
print(f"step {step:4d} {lbl:4s} loss={loss.item():.4f} "
|
|
f"lce={lce_val:.3f} "
|
|
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")
|
|
if lora_params:
|
|
torch.save({"rank": ARGS.bandlora, "band": lora_band_layers,
|
|
"tensors": [p.detach().cpu()
|
|
for p in lora_params]},
|
|
OUT / f"lora_{TAG}_e{step+1}.pt")
|
|
json.dump(log, open(OUT / f"train_{TAG}_log.json", "w"))
|
|
print("done", flush=True)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|