"""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 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"] + "") return (mbpp_prompt(tok, it, MBPP_SUFFIX), "```python\n" + it["sol_code"] + "\n```") 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) 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) adapter = MergeAdapter( 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()