diff --git a/docs/bf16-memory.md b/docs/bf16-memory.md new file mode 100644 index 0000000..adf321f --- /dev/null +++ b/docs/bf16-memory.md @@ -0,0 +1,48 @@ +# bf16 memory test and GPU choice (2026-10-06, estimates, nothing was run) + +**No GPU job is started without Kral's go.** Script: `train/hf_train_bf16.py` (`--memory-test --sweep 16000,32000,48000` is the memory test, no data needed, 3 optimizer steps per length, stops at the first OOM, uploads the result as `memtest_*.json` to the output repo). + +## Model (from config.json of Qwen3.8-27B) +64 layers (48 linear attention, 16 full attention), hidden 5120, MLP 17408, vocab 248320, 27.8 B parameters. LoRA rank 16 on q/k/v/o, gate/up/down, in_proj_qkv, in_proj_z, out_proj: **107 M trainable parameters**. + +## Memory by component at 48k tokens (GB; the estimate, not a measurement) +| component | GB | how | +|---|---|---| +| weights bf16 | 51.8 | 27.8 B x 2 bytes | +| LoRA weights + grads + 8-bit Adam | 1.0 | 107 M x 10 bytes | +| layer inputs for the backward pass (checkpoints) | 29.3 on the GPU, about 0 with the Unsloth CPU offload | 48k x 5120 x 2 bytes x 64 layers; the offload needs 29 GB of host RAM (all flavors have 142 GB or more) | +| recompute peak of one layer | 9.2 | MLP tensors 48k x 17408 x 2 bytes x 4 plus 3 GB for attention (assumption) | +| logits and loss | 3 chunked, 67 with full logits | full logits: 48k x 248320 x (bf16 + fp32 upcast + grad); only a chunked or fused cross entropy is realistic | + +## Fit by GPU (92 % of the card counted as usable) +| seq | checkpoints | loss | estimated peak GB | 1x A100 80 GB | 1x RTX PRO 6000 96 GB | 1x H200 141 GB | 2x H200 282 GB (needs model parallel) | +|---|---|---|---|---|---|---|---| +| 16k | CPU offload (Unsloth) | chunked | 61 | fits | fits | fits | fits | +| 16k | CPU offload (Unsloth) | full logits | 80 | no | fits | fits | fits | +| 16k | on the GPU | chunked | 71 | fits | fits | fits | fits | +| 16k | on the GPU | full logits | 90 | no | tight | fits | fits | +| 32k | CPU offload (Unsloth) | chunked | 63 | fits | fits | fits | fits | +| 32k | CPU offload (Unsloth) | full logits | 104 | no | no | fits | fits | +| 32k | on the GPU | chunked | 82 | no | fits | fits | fits | +| 32k | on the GPU | full logits | 124 | no | no | fits | fits | +| 48k | CPU offload (Unsloth) | chunked | 65 | fits | fits | fits | fits | +| 48k | CPU offload (Unsloth) | full logits | 129 | no | no | fits | fits | +| 48k | on the GPU | chunked | 94 | no | tight | fits | fits | +| 48k | on the GPU | full logits | 158 | no | no | no | fits | + +Reading: +- **The loss is the main risk, not the weights.** With full logits only the H200 (141 GB, with the CPU offload) fits 32k and 48k; the A100 and the RTX PRO 6000 do not. The run needs a fused or chunked cross entropy + (Unsloth has one for the architectures it patches; whether it covers `qwen3_5` is not known). The memory test shows it at once: if `seq_len 16000` already fails on an H200, the loss is the cause. + Fallback (not built yet, small): compute the hidden states, then the loss over the 32 % labeled positions only, in chunks of 4k. +- **A100 80 GB (2.50 USD/hour) is only possible with the CPU offload and a chunked loss** (about 65 GB estimated at 48k: no margin). H200 141 GB (5 USD/hour) fits with margin; + RTX PRO 6000 96 GB (2.75 USD/hour) fits with the offload and the chunked loss (est. 65 GB). +- Samples over 32k are a minority (p90 39560 tokens in the current data, p50 24.8k): if 48k does not fit, 32k would drop 15 of 52 samples (29 %; 10 of them DDLS), so test 48k first. +- Multi-GPU (`a100x4`, `h200x2`): possible only with model parallelism (`device_map`), one card works at a time, Unsloth multi-GPU is limited; not recommended before the single card test. + +## Proposed order +1. Memory test on **h200** (about 30 minutes, 2.5 USD): sweep 16000,32000,48000 with the offload. If 48k fits with margin: stay on H200 or try rtx-pro-6000 for the real run. +2. If the loss is the problem: build the chunked loss (about 2 hours of work), repeat the test. +3. Real run after the stage 1 : stage 2 ratio decision (below). + +## Time and cost (estimate from the nf4 run: 221 tokens/s on an A100 at 6.75k tokens per document; bf16 is faster, H200 about 2 to 2.5 times an A100) +Tokens of the proposed run (1 epoch stage 1 and 3 epochs of the 47 stage 2 samples as built today: 2.53 M + 3 x 1.28 M = 6.4 M tokens): A100 about 5.9 h (about 15 USD), H200 about 2.7 h (about 14 USD). The real number comes from the memory test (`step_seconds`). diff --git a/docs/stage2-build.md b/docs/stage2-build.md new file mode 100644 index 0000000..f4295e4 --- /dev/null +++ b/docs/stage2-build.md @@ -0,0 +1,52 @@ +# Stage 2 training set, first build (2026-10-06) + +Builder: `train/build_stage2.py` (output `runs/stage2_data/`, data card `README.md`, report `build_report.json`). Hook for item D: `train/hooks_example.py`. +HF dataset (private): `erhankeseli/abap-stage2-data`. The local Qwen trajectories (series A) are not read. + +## Result +90 accepted DeepSeek trajectories -> 90 after the eval overlap check (0 removed) -> 83 after the 48k limit +(**7 dropped over 48000 tokens: 6 DDLS, 1 INTF**; none cut) -> CLAS cap 35 %: 18 CLAS kept, **31 CLAS in reserve** (`stage2_reserve.jsonl`) +-> **47 train + 5 valid samples** (1.28 M + 0.12 M tokens, 0.41 M loss tokens in train, p50 24771, p95 41829, max 47694). +Loss mask checked on all 47 train samples: no span contains a tool result, the system turn or a user turn; loss share 32 % of the tokens; no token straddles a span boundary. + +## Not good enough yet (kinds under the minimum of 25) +| kind | samples | families | short | +|---|---|---|---| +| CLAS | 18 | 18 | 7 | +| INTF | 2 | 2 | 23 | +| DDLS | 13 | 13 | 12 | +| FUNC | 11 | 11 | 14 | +| PROG | 5 | 5 | 20 | +| TABL | 3 | 3 | 22 | +| STRU | 0 | 0 | 25 | +| MSAG | 0 | 0 | 25 | +| EXC | 0 | 0 | 25 | + +- **STRU, MSAG, exception: 0 trajectories**; INTF 2, TABL 3, PROG 5. Only CLAS, DDLS and FUNC are near their share. Nothing is filled with copies. This is what the restart plan (new kinds first) has to fix. +- **DDLS loses 6 of 19 trajectories to the 48k limit.** The long CDS trajectories (own CDS test class, many reads) are exactly the ones over the limit. Options: raise the limit to 64k (memory test decides), or generate CDS tasks with shorter runs. Not decided. +- Validation: 5 samples ({'CLAS': 2, 'FUNC': 1, 'DDLS': 1, 'PROG': 1}); the new kinds have no validation sample. After the restart the valid set must be rebuilt. +- 2 trajectories of one task and K variants are on the same side (family key); there are no same-task pairs yet in practice. +- Eval overlap: no accepted task overlaps an eval task (spec cosine 0.75, rules 0.60, names 0.60). + +## Stage 1 : stage 2 ratio (proposal, Kral decides) +Stage 1 train: 374 documents, 2.53 M tokens (all tokens carry loss). Stage 2 train today: 47 samples, 0.41 M loss tokens per epoch. Loss tokens per option (the table uses today's data): + +| stage 1 epochs | stage 2 epochs | stage 1 loss tokens | stage 2 loss tokens | stage 2 share | total tokens seen* | +|---|---|---|---|---|---| +| 2 | 3 | 5.05 M | 1.23 M | 20 % | 6.3 M | +| 1 | 3 | 2.53 M | 1.23 M | 33 % | 3.8 M | +| 0.5 | 3 | 1.26 M | 1.23 M | 49 % | 2.5 M | +| 0.33 | 3 | 0.83 M | 1.23 M | 60 % | 2.1 M | +| 1 | 3 | 2.53 M | 3.70 M | 59 % | 6.2 M | + +(*stage 1 tokens plus all stage 2 tokens per epoch; last row: if stage 2 had three times today's data, as expected after the restart.) + +**Proposal: stage 2 for 3 epochs always, stage 1 so that its loss tokens are about two thirds of the stage 2 loss tokens (stage 2 = 60 % of the loss).** +Reasons: (1) stage 2 is the behavior we want (repair after the first error, write after a few reads, the new kinds), and its samples are long and rare; stage 1 is domain knowledge in document form and acts as a +regularizer, so it should not dominate the gradient. (2) The earlier plan of 2 epochs of stage 1 (748 steps) would give 5.1 M stage 1 loss tokens against 1.2 M of stage 2: the model would mostly learn documents again, the format of the tool calls would be a small part. +(3) With 47 samples, more than 3 to 4 epochs of stage 2 risks memorizing them; the valid loss (5 samples today) is too thin to catch it, so keep an eye on the train loss curve and use the checkpoints. +(4) The rule scales with the data: today it means about 0.33 epochs of stage 1 (about 125 documents), with three times the stage 2 data about 1 epoch. `--s1-epochs` takes fractions. +Decision needed: this rule (60 % of the loss on stage 2), or a fixed 1 epoch of stage 1. + +## Also built +`train/hf_train_bf16.py` (bf16, loss mask, mixing, memory test), `docs/bf16-memory.md` (memory table by GPU, estimates), `train/hooks_example.py` (hook for item D, weights). diff --git a/train/build_stage2.py b/train/build_stage2.py new file mode 100644 index 0000000..25712ae --- /dev/null +++ b/train/build_stage2.py @@ -0,0 +1,260 @@ +"""Stage 2 training set builder (Opus item C, 2026-10-06). + + train/.venv/bin/python train/build_stage2.py [--out runs/stage2_data] [--max-tokens 48000] [--clas-cap 0.35] + [--min-per-kind 25] [--valid-frac 0.10] [--hook module:function] + +Input: runs/traj (DeepSeek trajectories; the local Qwen series is never read). Steps: acceptance filter (score, end reason, no +harness text, loop repeats trimmed) -> eval overlap check -> Qwen 3.8 chat template (thinking off) with the tokenizer round trip +check -> samples over --max-tokens are dropped (never cut) -> CLAS capped at --clas-cap of the stage 2 samples (the rest goes to +`stage2_reserve.jsonl`) -> shortage report for kinds under --min-per-kind -> validation split by task family (K variants and all +trajectories of a task on the same side) -> optional hook (own-test score, weights) -> files + data card. + +Loss mask: `assistant_spans` are character spans of the assistant turns (tool calls and the final report, with <|im_end|>); +nothing of the system turn (tool schemas), the user turn or the tool results is in a span. `assistant_tokens` counts the loss tokens. +""" +import argparse +import importlib +import json +import os +import random +import statistics +import sys +import time +from collections import Counter, defaultdict + +ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +sys.path.insert(0, ROOT) +sys.path.insert(0, os.path.join(ROOT, "train")) +import accept as acc # noqa: E402 +import to_qwen as tq # noqa: E402 +from harness import mix, overlap # noqa: E402 +from transformers import AutoTokenizer # noqa: E402 + +POOL = os.path.join(ROOT, "tasks_gen", "train") + + +def pct(v, p): + v = sorted(v) + return v[min(len(v) - 1, int(p * len(v)))] if v else None + + +def family_of(task_id): + try: + t = json.load(open(os.path.join(POOL, task_id, "task.json"))) + except OSError: + return task_id + return t.get("base_task") or task_id + + +def loss_tokens(tok, text, spans): + enc = tok(text, add_special_tokens=False, return_offsets_mapping=True) + n = 0 + for (a, b) in enc["offset_mapping"]: + if any(a >= s and b <= e for s, e in spans): + n += 1 + return len(enc["input_ids"]), n + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--traj", default=os.path.join(ROOT, "runs", "traj")) + ap.add_argument("--out", default=os.path.join(ROOT, "runs", "stage2_data")) + ap.add_argument("--max-tokens", type=int, default=48000) + ap.add_argument("--clas-cap", type=float, default=0.35) + ap.add_argument("--min-per-kind", type=int, default=25) + ap.add_argument("--min-score", type=float, default=80) + ap.add_argument("--valid-frac", type=float, default=0.10) + ap.add_argument("--seed", type=int, default=20261006) + ap.add_argument("--hook", help="module:function; function(row) -> None to drop, or a float weight (repeat factor), or a dict " + "{'keep': bool, 'weight': float, 'extra': {...}} (for example the own-test mutation score)") + a = ap.parse_args() + os.makedirs(a.out, exist_ok=True) + hook = None + if a.hook: + mod, fn = a.hook.split(":") + sys.path.insert(0, os.path.join(ROOT, "train")) + hook = getattr(importlib.import_module(mod), fn) + tok = AutoTokenizer.from_pretrained(os.path.expanduser("~/models/Qwen3.8-27B-4bit")) + evals = overlap.load_pool("eval") + report = {"built": time.strftime("%F %T"), "settings": vars(a), "dropped": defaultdict(Counter), "steps": {}} + + # 1 accepted trajectories + rows = [json.loads(l) for l in open(os.path.join(a.traj, "summary.jsonl"))] + cand = [] + for r in rows: + p = os.path.join(a.traj, r.get("run_dir") or "-", "record.json") + if not os.path.exists(p): + continue + rec = json.load(open(p)) + ok, why = acc.judge(rec, r, a.min_score) + if not ok: + continue + msgs, removed = acc.trim_loops(rec["messages"]) + cand.append({"rec": rec, "msgs": msgs, "removed": removed, "row": r, "repair": acc.has_repair(msgs)}) + report["steps"]["accepted_trajectories"] = len(cand) + + # 2 eval overlap at build time (the task against every eval task) + kept = [] + for c in cand: + tid = c["row"]["task"] + kind = mix.kind_of_task_dir(tid) + try: + hits = overlap.check(overlap.load_task(os.path.join(POOL, tid)), evals) + except (OSError, ValueError): + hits = [] + if hits: + report["dropped"][kind]["eval_overlap"] += 1 + continue + c["kind"] = kind + kept.append(c) + report["steps"]["after_eval_overlap"] = len(kept) + + # 3 Qwen template + tokens + loss tokens; drop over the limit (never cut) + samples = [] + for c in kept: + msgs = tq.to_template_messages(c["msgs"]) + tools = c["rec"]["tools"] + text = tok.apply_chat_template(msgs, tools=tools, tokenize=False, enable_thinking=False) + prefix = tok.apply_chat_template(msgs[:2], tools=tools, tokenize=False, add_generation_prompt=True, enable_thinking=False) + if not text.startswith(prefix) or not tq.check_calls(msgs, text): + report["dropped"][c["kind"]]["template_roundtrip"] += 1 + continue + spans = tq.spans(text) + n, nl = loss_tokens(tok, text, spans) + if n > a.max_tokens: + report["dropped"][c["kind"]]["over_%dk" % (a.max_tokens // 1000)] += 1 + continue + t = c["rec"]["task"] + samples.append({"id": f"{c['row']['task']}_r{c['row']['run']}", "task": c["row"]["task"], "family": family_of(c["row"]["task"]), + "kind": c["kind"], "category": t.get("category"), "object_type": t.get("object_type"), + "attempt": c["row"]["attempt"], "repair": c["repair"], "loop_pairs_removed": c["removed"], + "score": c["rec"]["metadata"].get("score"), "teacher": c["rec"]["teacher"], "weight": 1.0, + "n_tokens": n, "assistant_tokens": nl, "text": text, "assistant_spans": spans}) + report["steps"]["after_token_limit"] = len(samples) + + # 4 CLAS cap: keep the CLAS samples with a repair and the highest score first; the rest is the reserve + by_kind = defaultdict(list) + for s in samples: + by_kind[s["kind"]].append(s) + non_clas = sum(len(v) for k, v in by_kind.items() if k != "CLAS") + clas = sorted(by_kind.get("CLAS", []), key=lambda s: (not s["repair"], -(s["score"] or 0), s["id"])) + limit = int(a.clas_cap / (1 - a.clas_cap) * non_clas) if non_clas else len(clas) + reserve = clas[limit:] + by_kind["CLAS"] = clas[:limit] + report["steps"]["clas_cap"] = {"clas_before": len(clas), "clas_kept": len(by_kind["CLAS"]), "clas_reserve": len(reserve), "limit": limit} + pool = [s for v in by_kind.values() for s in v] + + # 5 shortage report (never filled by copies) + report["kinds"] = {} + for k in mix.TYPE_SHARE: + n = len(by_kind.get(k, [])) + report["kinds"][k] = {"samples": n, "min_required": a.min_per_kind, "short": max(0, a.min_per_kind - n), + "families": len({s["family"] for s in by_kind.get(k, [])})} + + # 6 hook (own-test score, weights) + if hook: + out = [] + for s in pool: + r = hook(s) + if r is None or r is False: + report["dropped"][s["kind"]]["hook"] += 1 + continue + if isinstance(r, dict): + if not r.get("keep", True): + report["dropped"][s["kind"]]["hook"] += 1 + continue + s["weight"] = float(r.get("weight", 1.0)) + s.update(r.get("extra") or {}) + elif isinstance(r, (int, float)) and not isinstance(r, bool): + s["weight"] = float(r) + out.append(s) + pool = out + + # 7 validation split by family, stratified by kind + rnd = random.Random(a.seed) + fam_kind = {} + for s in pool: + fam_kind.setdefault(s["family"], s["kind"]) + valid_fam = set() + for k in mix.TYPE_SHARE: + fams = sorted(f for f, kk in fam_kind.items() if kk == k) + rnd.shuffle(fams) + nv = round(len(fams) * a.valid_frac) + if len(fams) >= 4: + nv = max(nv, 1) + valid_fam |= set(fams[:nv]) + train = [s for s in pool if s["family"] not in valid_fam] + valid = [s for s in pool if s["family"] in valid_fam] + for name, data in (("stage2_train", train), ("stage2_valid", valid), ("stage2_reserve", reserve)): + with open(os.path.join(a.out, name + ".jsonl"), "w") as f: + for s in data: + f.write(json.dumps(s) + "\n") + + # 8 report and data card + def summary(data): + toks = [s["n_tokens"] for s in data] + return {"samples": len(data), "tokens": sum(toks), "loss_tokens": sum(s["assistant_tokens"] for s in data), + "p50": pct(toks, .5), "p90": pct(toks, .9), "p95": pct(toks, .95), "max": max(toks or [0]), + "by_kind": dict(Counter(s["kind"] for s in data)), "by_category": dict(Counter(s["category"] for s in data)), + "repair_share": round(sum(s["repair"] for s in data) / max(len(data), 1), 2)} + report["train"], report["valid"], report["reserve"] = summary(train), summary(valid), summary(reserve) + report["dropped"] = {k: dict(v) for k, v in report["dropped"].items()} + s1 = json.load(open(os.path.join(ROOT, "train", "data", "stats.json"))) if os.path.exists(os.path.join(ROOT, "train", "data", "stats.json")) else {} + report["stage1_stats"] = s1 + json.dump(report, open(os.path.join(a.out, "build_report.json"), "w"), indent=1, default=str) + open(os.path.join(a.out, "README.md"), "w").write(data_card(report)) + print(json.dumps({k: report[k] for k in ("steps", "train", "valid", "reserve", "dropped")}, indent=1, default=str)[:3500]) + print("short kinds:", {k: v["short"] for k, v in report["kinds"].items() if v["short"]}) + + +def data_card(r): + t, v = r["train"], r["valid"] + kinds = "\n".join(f"| {k} | {x['samples']} | {x['families']} | {x['short'] or ''} |" for k, x in r["kinds"].items()) + drop = "\n".join(f"- {k}: {d}" for k, d in r["dropped"].items()) or "- nothing dropped" + return f"""--- +license: mit +task_categories: [text-generation] +tags: [abap, sap, agentic, tool-use, sft] +private: true +--- +# ABAP stage 2 agent trajectories (built {r['built']}) + +Tool-using ABAP development trajectories for supervised fine-tuning of Qwen 3.8 27B. A teacher model solved generated ABAP tasks on a real +SAP ABAP Platform system (A4H, SAP_BASIS 816) through ADT tools; only runs that passed the harness gates and scored at least {r['settings']['min_score']:.0f} of 100 are kept. + +## Sources and licenses +- Teacher: DeepSeek V4.1 Flash (MIT license), via Ollama cloud. No output of Claude or other restricted models. +- Tasks: generated by the same teacher (spec, seed objects, hidden ABAP Unit tests, reference), validated on A4H (oracle 100, null 0, mutation check). Eval tasks are not in this data (overlap check at build time). +- The local Qwen runs (series A) are not in this data. + +## Format +One JSON per line: `text` (the whole conversation in the Qwen 3.8 chat template, thinking off, Qwen XML tool call format, the 20 tool schemas in the system turn), +`assistant_spans` (character spans that carry the loss: assistant turns with tool calls and the final report; none of the system turn, user turn or tool results), +`n_tokens`, `assistant_tokens`, `kind`, `category`, `object_type`, `family` (task family: a K variant has the family of its base task), `repair` (an error followed by a fix), `score`, `weight`. + +## Size +| split | samples | tokens | loss tokens | p50 | p90 | p95 | max | repair share | +|---|---|---|---|---|---|---|---|---| +| train | {t['samples']} | {t['tokens']} | {t['loss_tokens']} | {t['p50']} | {t['p90']} | {t['p95']} | {t['max']} | {t['repair_share']} | +| valid | {v['samples']} | {v['tokens']} | {v['loss_tokens']} | {v['p50']} | {v['p90']} | {v['p95']} | {v['max']} | {v['repair_share']} | + +Samples over {r['settings']['max_tokens']} tokens are dropped, never cut. CLAS is capped at {int(100 * r['settings']['clas_cap'])} % of the samples (the rest is in `stage2_reserve.jsonl`). + +## Kinds (target share: CLAS 28, INTF 7, CDS 25, FUNC 15, PROG 10, TABL 8, STRU 2, MSAG 2.5, exception 2.5 percent) +| kind | samples | families | short of the minimum {r['settings']['min_per_kind']} | +|---|---|---|---| +{kinds} + +## Dropped +{drop} + +## Known limits +- Small data (see the table); kinds marked short have fewer than the minimum samples. +- The proxy added local abaplint messages (`syntaxCheck`) to some failed writes; the EPOD server does not do this itself. +- The teacher sees only the tool results of its own runs; a run that repaired an error is kept with its error (60 to 65 % of the samples). +- Validation split is by task family; do not mix valid samples into the training set. +""" + + +if __name__ == "__main__": + main() diff --git a/train/hf_train_bf16.py b/train/hf_train_bf16.py new file mode 100644 index 0000000..c1a9063 --- /dev/null +++ b/train/hf_train_bf16.py @@ -0,0 +1,163 @@ +# /// script +# requires-python = ">=3.10" +# dependencies = ["unsloth", "datasets", "transformers", "huggingface_hub"] +# /// +"""Stage 1 + stage 2 mixed bf16 LoRA run for Qwen 3.8 27B on Hugging Face Jobs (Opus item E, 2026-10-06). NOT started: no job without Kral's go. + +Memory test first (a few dollars, finds the largest sequence length that fits, no data needed): + hf jobs uv run --flavor h200 --timeout 40m --secrets HF_TOKEN train/hf_train_bf16.py -- --memory-test --sweep 16000,32000,48000 +Real run (after the memory test and Kral's decision on the ratio): + hf jobs uv run --flavor --timeout 10h --secrets HF_TOKEN train/hf_train_bf16.py -- \\ + --stage1 erhankeseli/abap-stage1-data --stage2 erhankeseli/abap-stage2-data --s1-epochs 1 --s2-epochs 3 --out erhankeseli/abap-mixed-adapter + +Data: stage 1 rows have `text` (loss on every token); stage 2 rows have `text` and `assistant_spans` (loss only inside the spans: assistant turns; +none on the system turn with the tool schemas, the user turn or the tool results). A sample longer than --max-seq is skipped, never cut. +Base model in bf16 (no nf4): the adapter then fits the MLX 4-bit base better (stage 1 finding, train/STATE.md). +""" +import argparse +import json +import os +import random +import time + + +def tokenize_masked(tok, text, spans, max_len): + """input_ids and labels; labels are -100 outside the spans (spans=None: loss on every token). None when longer than max_len.""" + enc = tok(text, add_special_tokens=False, return_offsets_mapping=True) + ids = enc["input_ids"] + if len(ids) > max_len: + return None + if spans is None: + return ids, list(ids) + labels = [] + for tid, (a, b) in zip(ids, enc["offset_mapping"]): + labels.append(tid if any(a >= s and b <= e for s, e in spans) else -100) + return ids, labels + + +def build_examples(tok, stage1, stage2, s1_epochs, s2_epochs, max_len, seed): + rnd = random.Random(seed) + rows = [] + for ep in range(int(s1_epochs)): + rows += [("s1", r["text"], None, 1.0) for r in stage1] + frac = s1_epochs - int(s1_epochs) + if frac: + rows += [("s1", r["text"], None, 1.0) for r in rnd.sample(stage1, int(len(stage1) * frac))] + for ep in range(int(s2_epochs)): + for r in stage2: + rows += [("s2", r["text"], r["assistant_spans"], 1.0)] * max(1, round(float(r.get("weight", 1.0)))) + rnd.shuffle(rows) + out, skipped = [], {"s1": 0, "s2": 0} + for src, text, spans, _ in rows: + t = tokenize_masked(tok, text, spans, max_len) + if t is None: + skipped[src] += 1 + continue + out.append({"input_ids": t[0], "labels": t[1], "src": src}) + return out, skipped + + +def main(): + ap = argparse.ArgumentParser() + ap.add_argument("--model", default="Qwen/Qwen3.8-27B") + ap.add_argument("--stage1", default="erhankeseli/abap-stage1-data") + ap.add_argument("--stage2", default="erhankeseli/abap-stage2-data") + ap.add_argument("--s1-file", default="train.jsonl") + ap.add_argument("--s2-file", default="stage2_train.jsonl") + ap.add_argument("--s2-valid", default="stage2_valid.jsonl") + ap.add_argument("--s1-epochs", type=float, default=1.0) + ap.add_argument("--s2-epochs", type=float, default=3.0) + ap.add_argument("--max-seq", type=int, default=48000) + ap.add_argument("--rank", type=int, default=16) + ap.add_argument("--alpha", type=int, default=32) + ap.add_argument("--lr", type=float, default=5e-5) + ap.add_argument("--grad-accum", type=int, default=4) + ap.add_argument("--warmup-steps", type=int, default=20) + ap.add_argument("--no-offload", action="store_true", help="gradient checkpointing on the GPU instead of the Unsloth offload") + ap.add_argument("--out", default="erhankeseli/abap-mixed-adapter") + ap.add_argument("--save-every", type=int, default=0) + ap.add_argument("--memory-test", action="store_true") + ap.add_argument("--sweep", default="16000,32000,48000", help="memory test: sequence lengths, tried in this order, stops at the first OOM") + ap.add_argument("--steps", type=int, default=3, help="memory test: optimizer steps per length") + a = ap.parse_args() + + import torch + from unsloth import FastLanguageModel + from huggingface_hub import HfApi, hf_hub_download + token = os.environ["HF_TOKEN"] + model, tok = FastLanguageModel.from_pretrained(a.model, max_seq_length=a.max_seq, load_in_4bit=False, dtype=torch.bfloat16, token=token) + model = FastLanguageModel.get_peft_model( + model, r=a.rank, lora_alpha=a.alpha, lora_dropout=0.0, bias="none", + target_modules=["q_proj", "k_proj", "v_proj", "o_proj", "gate_proj", "up_proj", "down_proj", "in_proj_qkv", "in_proj_z", "out_proj"], + use_gradient_checkpointing=True if a.no_offload else "unsloth", random_state=20261006) + print("GPU", torch.cuda.get_device_name(0), "x", torch.cuda.device_count(), "| weights allocated GB", round(torch.cuda.memory_allocated() / 2**30, 1), flush=True) + + if a.memory_test: + text = "CLASS zcl_demo DEFINITION PUBLIC. METHODS run IMPORTING iv TYPE i. ENDCLASS.\n" * 4000 + ids_all = tok(text, add_special_tokens=False)["input_ids"] + opt = torch.optim.AdamW([p for p in model.parameters() if p.requires_grad], lr=1e-5) + model.train() + result = {"gpu": torch.cuda.get_device_name(0), "gpus": torch.cuda.device_count(), "offload": not a.no_offload, "sweep": []} + for n in [int(x) for x in a.sweep.split(",")]: + ids = torch.tensor([(ids_all * (n // len(ids_all) + 1))[:n]], device="cuda") + labels = ids.clone() + labels[:, : int(n * 0.68)] = -100 # about 32 % loss tokens, as in the stage 2 data + torch.cuda.reset_peak_memory_stats() + try: + times = [] + for _ in range(a.steps): + t0 = time.time() + loss = model(input_ids=ids, labels=labels).loss + loss.backward() + opt.step() + opt.zero_grad(set_to_none=True) + torch.cuda.synchronize() + times.append(round(time.time() - t0, 1)) + row = {"seq_len": n, "ok": True, "peak_gb": round(torch.cuda.max_memory_allocated() / 2**30, 1), + "reserved_gb": round(torch.cuda.max_memory_reserved() / 2**30, 1), "step_seconds": times, "loss": float(loss)} + except torch.cuda.OutOfMemoryError as e: + row = {"seq_len": n, "ok": False, "error": "OOM", "peak_gb": round(torch.cuda.max_memory_allocated() / 2**30, 1)} + result["sweep"].append(row) + print("MEMTEST", json.dumps(row), flush=True) + break + result["sweep"].append(row) + print("MEMTEST", json.dumps(row), flush=True) + torch.cuda.empty_cache() + print("RESULT", json.dumps(result), flush=True) + HfApi(token=token).upload_file(path_or_fileobj=json.dumps(result, indent=1).encode(), path_in_repo="memtest_%s_%dx.json" % ( + result["gpu"].replace(" ", "_"), result["gpus"]), repo_id=a.out, repo_type="model", create_pr=False) if a.out else None + return + + def rows(repo, fname): + p = hf_hub_download(repo, fname, repo_type="dataset", token=token) + return [json.loads(l) for l in open(p)] + s1, s2 = rows(a.stage1, a.s1_file), rows(a.stage2, a.s2_file) + ex, skipped = build_examples(tok, s1, s2, a.s1_epochs, a.s2_epochs, a.max_seq, 20261006) + loss_tok = sum(sum(1 for x in e["labels"] if x != -100) for e in ex) + print("EXAMPLES", len(ex), "skipped", skipped, "loss tokens", loss_tok, "stage 2 share of loss tokens", + round(sum(sum(1 for x in e["labels"] if x != -100) for e in ex if e["src"] == "s2") / max(loss_tok, 1), 2), flush=True) + from datasets import Dataset + from transformers import Trainer, TrainingArguments + + class Collate: + def __call__(self, batch): # batch size 1; no padding needed + b = batch[0] + return {"input_ids": torch.tensor([b["input_ids"]]), "labels": torch.tensor([b["labels"]]), + "attention_mask": torch.ones(1, len(b["input_ids"]), dtype=torch.long)} + args = TrainingArguments(output_dir="out", per_device_train_batch_size=1, gradient_accumulation_steps=a.grad_accum, num_train_epochs=1, + learning_rate=a.lr, lr_scheduler_type="cosine", warmup_steps=a.warmup_steps, optim="adamw_8bit", bf16=True, + logging_steps=1, save_strategy="no", report_to="none", seed=20261006, remove_unused_columns=False) + tr = Trainer(model=model, args=args, train_dataset=Dataset.from_list(ex), data_collator=Collate()) + t0 = time.time() + tr.train() + info = {"examples": len(ex), "skipped": skipped, "train_seconds": round(time.time() - t0, 1), "rank": a.rank, "alpha": a.alpha, "lr": a.lr, + "peak_gpu_gb": round(torch.cuda.max_memory_allocated() / 2**30, 1), "s1_epochs": a.s1_epochs, "s2_epochs": a.s2_epochs} + print("RESULT", json.dumps(info), flush=True) + model.save_pretrained("adapter") + json.dump(info, open("adapter/job_result.json", "w")) + HfApi(token=token).upload_folder(folder_path="adapter", repo_id=a.out, repo_type="model") + print("PUSHED", a.out) + + +if __name__ == "__main__": + main() diff --git a/train/hooks_example.py b/train/hooks_example.py new file mode 100644 index 0000000..580b544 --- /dev/null +++ b/train/hooks_example.py @@ -0,0 +1,28 @@ +"""Hooks for train/build_stage2.py (--hook hooks_example:own_test_weight). A hook gets one sample (a dict) and returns +None / False (drop), a number (repeat weight), or {"keep": bool, "weight": float, "extra": {...}}.""" +import json +import os + +ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) + + +def identity(row): + return 1.0 + + +def own_test_weight(row): + """Item D (own-test mutation score, metadata only for now): reads runs/traj//own_test_mutation.json when it exists, + stores it as extra data and does NOT drop or reweight (Kral + Opus 2026-10-06: do not change the acceptance yet).""" + run = row["id"].split("_r")[-1] + for d in os.listdir(os.path.join(ROOT, "runs", "traj")): + if d.startswith(run + "_"): + p = os.path.join(ROOT, "runs", "traj", d, "own_test_mutation.json") + if os.path.exists(p): + m = json.load(open(p)) + return {"keep": True, "weight": 1.0, "extra": {"own_test_mutation": m.get("score"), "own_test_mutants": m.get("mutants")}} + return {"keep": True, "weight": 1.0} + + +def repair_up(row): + """Example of a weight: a trajectory with a repair counts twice.""" + return 2.0 if row.get("repair") else 1.0