Files
jspace/scripts/train_merge_unified.py
T
2026-07-14 12:32:36 +02:00

189 lines
6.9 KiB
Python

"""Unified merge adapter: GSM8K + MBPP, prompt-only looping, fresh init.
One adapter, one regime (loop_mask = prompt span; generated/answer tokens
never loop), mixed-task batches. Hardened protocol elements: per-task val
holdouts, checkpoint every 100 steps kept separately (best-val selection and
k chosen on val, never test).
"""
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 loop_common import (AdaptiveMergeAdapter, BandLooper, MergeAdapter,
chat_prompt, DIRECT_SUFFIX)
from prep_mbpp import DIRECT_SUFFIX as MBPP_SUFFIX
from prep_mbpp import mbpp_prompt
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 = 800
BATCH = 6
LR = 1e-3
WARMUP = 20
MAX_TOK = 512
K_BUCKETS = [(1, ("easy",)), (2, ("easy", "hard")), (4, ("hard",))]
SEED = 0
VAL_N = 16 # per task per bucket
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 item_texts(tok, it):
if it["task"] == "gsm":
return (chat_prompt(tok, it["question"], DIRECT_SUFFIX),
it["gold"] + "<end_of_turn>")
return (mbpp_prompt(tok, it, MBPP_SUFFIX),
"```python\n" + it["sol_code"] + "\n```<end_of_turn>")
def build_batch(tok, items, device="cuda"):
seqs, labs = [], []
for it in items:
ptxt, atxt = item_texts(tok, it)
p = tok(ptxt, add_special_tokens=False)["input_ids"]
a = tok(atxt, add_special_tokens=False)["input_ids"]
seqs.append(p + a)
labs.append([-100] * len(p) + a)
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
lmask = (lab == -100) & (msk == 1)
return ids.to(device), msk.to(device), lab.to(device), lmask.to(device)
@torch.no_grad()
def val_loss(looper, adapter, tok, items, k, feedforward=False):
tot, n = 0.0, 0
for i in range(0, len(items), BATCH):
ids, msk, lab, lmask = build_batch(tok, items[i : i + BATCH])
logits = looper.loop_logits(adapter, ids, k, attention_mask=msk,
loop_mask=lmask, feedforward=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():
ap = argparse.ArgumentParser()
ap.add_argument("--noloop", action="store_true",
help="feedforward control: adapter(e,e), no recurrence")
ap.add_argument("--tasks", default="gsm,mbpp")
ap.add_argument("--tag", default=None)
ap.add_argument("--adaptive", action="store_true")
args = ap.parse_args()
tasks = args.tasks.split(",")
tag = args.tag or ("ff" if args.noloop else "uni")
rng = random.Random(SEED)
torch.manual_seed(SEED)
gsm, mbpp = [], []
if "gsm" in tasks:
gsm = [dict(it, task="gsm")
for it in json.load(open(OUT / "star_data.json"))
if it["split"] == "train" and it["label"] != "drop"]
if "mbpp" in tasks:
mbpp = [dict(it, task="mbpp")
for it in json.load(open(OUT / "mbpp_data.json"))
if it["split"] == "train" and it["label"] != "drop"
and it.get("sol_code")]
model, tok = load_model(dtype=torch.bfloat16)
for p in model.parameters():
p.requires_grad_(False)
looper = BandLooper(model)
cls = AdaptiveMergeAdapter if args.adaptive else MergeAdapter
adapter = cls(
d=model.config.get_text_config().hidden_size).cuda()
opt = torch.optim.AdamW(adapter.parameters(), lr=LR, weight_decay=0.01)
def fits(it):
ptxt, atxt = item_texts(tok, it)
return (len(tok(ptxt)["input_ids"]) + len(tok(atxt)["input_ids"])
<= MAX_TOK)
items = [it for it in gsm + mbpp if fits(it)]
rng.shuffle(items)
pool, val = {"easy": [], "hard": []}, {}
for task in tasks:
for lbl in ("easy", "hard"):
sub = [it for it in items if it["task"] == task
and it["label"] == lbl]
val[(task, lbl)] = sub[:VAL_N]
pool.setdefault(lbl, []).extend(sub[VAL_N:])
print("pool: easy={} hard={} (gsm {} / mbpp {})".format(
len(pool["easy"]), len(pool["hard"]),
sum(it["task"] == "gsm" for l in pool.values() for it in l),
sum(it["task"] == "mbpp" for l in pool.values() for it in l)),
flush=True)
log = []
t0 = time.time()
for step in range(STEPS):
k, labels = K_BUCKETS[step % len(K_BUCKETS)]
cand = [it for lbl in labels for it in pool[lbl]]
batch = rng.sample(cand, BATCH)
ids, msk, lab, lmask = build_batch(tok, batch)
for g in opt.param_groups:
g["lr"] = lr_at(step)
logits = looper.loop_logits(adapter, ids, k, attention_mask=msk,
use_checkpoint=True, loop_mask=lmask,
feedforward=args.noloop)
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, "k": k, "loss": loss.item()})
if step % 10 == 0:
print(f"step {step:4d} k={k} loss={loss.item():.4f} "
f"({(time.time()-t0)/(step+1):.1f}s/step)", flush=True)
if step % 100 == 99 or step == STEPS - 1:
vals = {}
kk_grid = (0, 1) if args.noloop else (0, 1, 2, 4)
for (task, lbl), vitems in val.items():
for kk in kk_grid:
vals[f"{task}_{lbl}_k{kk}"] = val_loss(
looper, adapter, tok, vitems, kk,
feedforward=args.noloop)
print(f" val@{step}: " + " ".join(
f"{n}={v:.3f}" for n, v in sorted(vals.items())), flush=True)
log.append({"step": step, "val": vals})
torch.save(adapter.state_dict(),
OUT / f"adapter_{tag}_e{step+1}.pt")
json.dump(log, open(OUT / f"train_{tag}_log.json", "w"), indent=1)
print("done", flush=True)
if __name__ == "__main__":
main()