ArushBuilds commited on
Commit
2735668
·
1 Parent(s): 742b752

Upload 4 files

Browse files
Files changed (4) hide show
  1. config.py +136 -0
  2. inference.py +361 -0
  3. model.py +1813 -0
  4. tokenizer.py +389 -0
config.py ADDED
@@ -0,0 +1,136 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+
2
+ use_liger=True
3
+ xla_on_gpu=True
4
+ # ---------------------------------------------------------------------------
5
+ # Core shape -- chosen for GPU efficiency, not just "a number that fits"
6
+ # ---------------------------------------------------------------------------
7
+ CONTEXT = 512
8
+ max_gen_tokens = 256
9
+ yarn_scale = 1.0 # neutralized YaRN extension ramp for clean RoPE sanity checks
10
+ gen_headroom = 128
11
+ vocab_size = 16384
12
+
13
+ # Missing config fields expected by the shared GPTConfig dataclass and the
14
+ # model factory path. Keeping them explicit here removes the fallback-warning
15
+ # cascade shown in the traceback and stabilizes object construction.
16
+ bias = False
17
+
18
+
19
+
20
+ gradient_checkpointing = False
21
+ num_kv_heads = 4
22
+ pattern = "dense"
23
+ sliding_window_size = 512
24
+ # tensor_parallel_size is the knob the config dataclass reads.
25
+ tensor_parallel_size = 1
26
+
27
+ # YaRN config defaults from the model constructor.
28
+
29
+ mod_capacity = 0.25
30
+ use_mtp = False # set True to activate
31
+ mtp_depth = 1 # 1 = predict 2 tokens ahead; 2 = also predict 3 ahead
32
+ mtp_lambda = 0.3
33
+ use_nvfp4 = False
34
+ use_int8 = False
35
+ int8_group_size = 128
36
+ mod_gate_entropy_coeff = 0.01
37
+ # Experimental hardware-specific paths. These are deliberately opt-in and
38
+ # are read live by model.py/prefetch.py so post-import CLI overrides work.
39
+
40
+ numberoflayers = 12 # 4 complete pyramid_swa groups (4x4) -- no odd
41
+ # trailing layer forced into an incomplete group
42
+ numberofheads = 12
43
+ D_MODEL = 192 # multiple of 64 -> head_dim = 192/12 = 16, a clean
44
+ # Tensor Core tile-friendly size (fp16 wants
45
+ # multiples of 16; 16 divides evenly with no
46
+ # padding waste on T4/A100 GEMM tiling)
47
+ # AdamW-scale LR for a model this size; only
48
+ # takes effect if use_static_lr=False below
49
+
50
+ # ---------------------------------------------------------------------------
51
+ # Attention
52
+ # ---------------------------------------------------------------------------
53
+ num_kv_heads = 4 # 2x GQA compression on global layers (8->4 kv
54
+ # heads); pyramid_swa steps kv heads [1,2,4]
55
+ # across each group's 3 SWA layers (MQA->GQA),
56
+ # verified against the real build_attention_layers
57
+ # logic in model.py, not assumed
58
+
59
+ sliding_window_size = 512 # explicit (was previously left to fall back to
60
+ # the default with a warning) -- 1/4 of CONTEXT,
61
+ # a reasonable local-attention window for this
62
+ # sequence length
63
+ use_sdpa=False
64
+ optimizer_8bit = True
65
+ # ---------------------------------------------------------------------------
66
+ # FFN / MoE
67
+ # ---------------------------------------------------------------------------
68
+ ffn_mult = 8 / 3 # SwiGLU hidden_dim multiplier. model.py rounds
69
+ # d_model*ffn_mult UP to the nearest multiple of
70
+ # 64 (256 * 8/3 = 682.67 -> 704), avoiding the
71
+ # prime-dimension GEMM padding waste a plain
72
+ # round() would produce (683 is prime).
73
+ use_moe = False # plain dense SwiGLU_FFN every layer. AryaSparseMoE
74
+ # currently has NO auxiliary load-balancing loss
75
+ # (expert_bias is trainable but nothing pushes it
76
+ # toward balanced routing) -- leave this False
77
+ # unless/until that's added, especially at this
78
+ # small a scale where routing collapse compounds
79
+ # fastest across depth.
80
+ n_experts = 8
81
+ n_shared = 1
82
+ # T4 path
83
+
84
+ # (use_int8 ignored on SM100+ -- nvFP4 takes priority)
85
+ # ---------------------------------------------------------------------------
86
+ # Precision / batching
87
+ # ---------------------------------------------------------------------------
88
+ PRECISION = 'fp16' # T4 has no bf16 Tensor Cores -- fp16 is correct
89
+ # here. On A100+, bf16 is the better choice (same
90
+ # speed, wider dynamic range, no GradScaler
91
+ # needed) -- change this if/when you move hardware.
92
+ GLOBAL_DTYPE = PRECISION
93
+ USE_ASYNC_LAYER_PREFETCH = True
94
+ DISABLE_MMAP = False
95
+ use_flash_outproj_add_rmsnorm = False
96
+ use_flash_rope = False
97
+ use_flash_swiglu = False
98
+ use_fused_add_rmsnorm = False
99
+ batch_size = 64 # starting point, not a measured optimum -- train.py
100
+ # now logs `vram X/Y GB` on every step (added
101
+ # recently). Watch that number on your first run
102
+ # and raise batch_size toward your actual VRAM
103
+ # headroom; this value has NOT been validated
104
+ # against a real GPU run for this exact config.
105
+ weight_decay = 0.01 # train.py now splits params into decay/no-decay
106
+ # groups automatically (biases, RMSNorm gains,
107
+ # expert_bias get weight_decay=0.0 regardless of
108
+ # this value) -- this number only applies to
109
+ # actual weight matrices + the tied embedding.
110
+
111
+ # ---------------------------------------------------------------------------
112
+ # Schedule
113
+ # ---------------------------------------------------------------------------
114
+ # Defaults to warmup+cosine (NOT static) for this fresh config -- unlike
115
+ # your existing config.py, this file has no prior "it works, don't touch it"
116
+ # history behind a flat LR, so there's no reason to ship it with the known
117
+ # no-warmup risk baked in. WARMUP_STEPS=500 and MAX_STEPS=20000 are
118
+ # hardcoded in train.py's main() (not read from config.py at all) -- change
119
+ # them there directly if this run needs a different horizon.
120
+ use_static_lr =True
121
+ static_lr = 4e-3 # only used if you flip use_static_lr back to True
122
+ learning_rate = 3e-3
123
+ # ---------------------------------------------------------------------------
124
+ # Parallelism / misc
125
+ # ---------------------------------------------------------------------------
126
+
127
+
128
+ data_parallel_only = True
129
+ dropout = 0.0 # explicit (was previously left to the
130
+ # default with a warning)
131
+ use_muon = True
132
+ # ---------------------------------------------------------------------------
133
+ # Data
134
+ # ---------------------------------------------------------------------------
135
+ DATA_PATH = './training_data'
136
+ data_path = DATA_PATH
inference.py ADDED
@@ -0,0 +1,361 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """
2
+ inference.py -- load a train.py checkpoint and generate text with it.
3
+
4
+ python inference.py --checkpoint ./checkpoints/latest.pt --prompt "Once upon a time"
5
+ python inference.py --checkpoint ./checkpoints/latest.pt --interactive
6
+
7
+ The checkpoint carries its own architecture config (see train.py's
8
+ save_checkpoint), so you don't need config.py to match the training run --
9
+ the model is reconstructed exactly as it was trained, every time.
10
+
11
+ KNOWN LIMITATION: model.py's attention modules have no KV cache. Every
12
+ generated token re-runs a full forward pass over the whole sequence so far
13
+ (O(T) attention passes, each O(T) or O(T*window) itself -- no incremental
14
+ state is carried between steps). Fine for short completions and for
15
+ sanity-checking a checkpoint; a real KV cache is the natural next piece of
16
+ work if you want fast interactive generation, and would need to be threaded
17
+ through CompressedSparseAttention / HeavilyCompressedAttention /
18
+ SlidingWindowAttention individually since each keeps different per-branch
19
+ state (compressed blocks vs. raw recent tokens).
20
+ """
21
+
22
+ import argparse
23
+ import contextlib
24
+ import sys
25
+
26
+ import torch
27
+
28
+ from model import GPT, GPTConfig
29
+ from tokenizer import TokenizerWrapper
30
+
31
+
32
+ # ─────────────────────────────────────────────────────────────────────────────
33
+ # Device selection
34
+ # ─────────────────────────────────────────────────────────────────────────────
35
+
36
+ def pick_device(requested: str | None) -> torch.device:
37
+ if requested is not None:
38
+ return torch.device(requested)
39
+ if torch.cuda.is_available():
40
+ return torch.device("cuda")
41
+ try:
42
+ import torch_xla.core.xla_model as xm # noqa: F401
43
+ return torch.device("xla")
44
+ except ImportError:
45
+ pass
46
+ if torch.backends.mps.is_available():
47
+ return torch.device("mps")
48
+ return torch.device("cpu")
49
+
50
+
51
+ # ─────────────────────────────────────────────────────────────────────────────
52
+ # Checkpoint loading
53
+ # ─────────────────────────────────────────────────────────────────────────────
54
+
55
+ def load_model(checkpoint_path: str, device: torch.device) -> tuple[GPT, GPTConfig]:
56
+ checkpoint = torch.load(checkpoint_path, map_location=device, weights_only=False)
57
+
58
+ if "model_config" not in checkpoint:
59
+ raise RuntimeError(
60
+ f"'{checkpoint_path}' has no saved model_config -- it was written "
61
+ "before save_checkpoint() started saving the architecture config "
62
+ "alongside the weights. Either retrain (new checkpoints save it "
63
+ "automatically), or construct GPTConfig(...) by hand here to "
64
+ "match whatever config.py looked like when this checkpoint was "
65
+ "produced, and call GPT(that_config) instead of this function."
66
+ )
67
+
68
+ model_config = GPTConfig(**checkpoint["model_config"])
69
+ model = GPT(model_config)
70
+
71
+ state_dict = checkpoint["model_state_dict"]
72
+
73
+
74
+ if any(k.startswith("_orig_mod.") for k in state_dict):
75
+ state_dict = {
76
+ (k[len("_orig_mod."):] if k.startswith("_orig_mod.") else k): v
77
+ for k, v in state_dict.items()
78
+ }
79
+ print(
80
+ "[Inference] Checkpoint was saved from a torch.compile()'d model "
81
+ "(keys prefixed with '_orig_mod.') -- stripped the prefix to "
82
+ "load onto this plain GPT instance.",
83
+ file=sys.stderr,
84
+ )
85
+
86
+ model.load_state_dict(state_dict)
87
+ model.to(device)
88
+ model.eval()
89
+
90
+ n_params = model.get_num_params() / 1e6
91
+ step = checkpoint.get("step", "?")
92
+ print(f"[Inference] Loaded checkpoint (step {step}, {n_params:.2f}M params) onto {device}", file=sys.stderr)
93
+ return model, model_config
94
+
95
+
96
+ # ─────────────────────────────────────────────────────────────────────────────
97
+ # Sampling
98
+ # ─────────────────────────────────────────────────────────────────────────────
99
+
100
+ def _apply_repetition_penalty(logits: torch.Tensor, generated_ids: list[int], penalty: float) -> torch.Tensor:
101
+ """CTRL-style penalty, COUNT-SCALED: a token seen N times gets divided
102
+ by penalty**N (if positive) / multiplied by penalty**N (if negative),
103
+ not a single flat divide regardless of occurrence count.
104
+
105
+ The original flat version (divide once no matter how many times a
106
+ token has recurred) is too weak on small/repetitive-domain models
107
+ (e.g. TinyStories-style data where "so so so happy!" is a genuinely
108
+ common training pattern): once the model locks onto a token, the
109
+ penalty never gets any stronger no matter how many times it fires,
110
+ so a token that's only mildly favored over the runner-up never gets
111
+ pushed low enough to escape the loop. Scaling by occurrence count
112
+ means each additional repeat compounds the penalty, so a loop that's
113
+ 3-4 tokens deep gets meaningfully suppressed even if a single
114
+ application wouldn't have been enough.
115
+ """
116
+ if penalty == 1.0 or not generated_ids:
117
+ return logits
118
+ from collections import Counter
119
+ counts = Counter(generated_ids)
120
+ seen = torch.tensor(sorted(counts.keys()), device=logits.device, dtype=torch.long)
121
+ reps = torch.tensor([counts[t] for t in sorted(counts.keys())], device=logits.device, dtype=logits.dtype)
122
+ vals = logits[0, seen]
123
+ factor = penalty ** reps
124
+ logits[0, seen] = torch.where(vals > 0, vals / factor, vals * factor)
125
+ return logits
126
+
127
+
128
+ def _block_repeated_ngrams(logits: torch.Tensor, generated_ids: list[int], ngram_size: int) -> torch.Tensor:
129
+ """Hard n-gram repeat block (HF's `no_repeat_ngram_size`, same idea).
130
+
131
+ Look at the last (ngram_size - 1) generated tokens. If that same
132
+ prefix has appeared earlier in generated_ids, whatever token followed
133
+ it there is banned from being chosen again right now (logit -> -inf).
134
+ This is a HARD constraint, unlike repetition_penalty which only
135
+ nudges probabilities down -- it's what actually stops "so so so so
136
+ so..." dead rather than just making it less likely each step. Only
137
+ kicks in once at least ngram_size-1 tokens have been generated.
138
+ """
139
+ if ngram_size <= 0 or len(generated_ids) < ngram_size - 1:
140
+ return logits
141
+ prefix_len = ngram_size - 1
142
+ if prefix_len == 0:
143
+ return logits
144
+ current_prefix = tuple(generated_ids[-prefix_len:])
145
+ banned = set()
146
+ for i in range(len(generated_ids) - prefix_len):
147
+ if tuple(generated_ids[i:i + prefix_len]) == current_prefix:
148
+ banned.add(generated_ids[i + prefix_len])
149
+ if banned:
150
+ banned_idx = torch.tensor(sorted(banned), device=logits.device, dtype=torch.long)
151
+ logits[0, banned_idx] = float("-inf")
152
+ return logits
153
+
154
+
155
+ def _top_k_filter(logits: torch.Tensor, top_k: int) -> torch.Tensor:
156
+ k = min(top_k, logits.size(-1))
157
+ values, _ = torch.topk(logits, k)
158
+ logits[logits < values[:, [-1]]] = float("-inf")
159
+ return logits
160
+
161
+
162
+ def _top_p_filter(logits: torch.Tensor, top_p: float) -> torch.Tensor:
163
+ """Nucleus sampling: keep the smallest prefix of sorted tokens whose
164
+ cumulative probability crosses top_p, drop the rest."""
165
+ sorted_logits, sorted_idx = torch.sort(logits, descending=True)
166
+ probs = torch.softmax(sorted_logits, dim=-1)
167
+ cum_probs = torch.cumsum(probs, dim=-1)
168
+ # remove token i if the cumulative prob BEFORE it already exceeds top_p
169
+ # (i.e. it wasn't needed to cross the threshold) -- matches the standard
170
+ # HF TopPLogitsWarper convention.
171
+ remove = (cum_probs - probs) > top_p
172
+ sorted_logits = sorted_logits.masked_fill(remove, float("-inf"))
173
+ return torch.full_like(logits, float("-inf")).scatter(1, sorted_idx, sorted_logits)
174
+
175
+
176
+ @torch.no_grad()
177
+ def generate(
178
+ model: GPT,
179
+ tokenizer: TokenizerWrapper,
180
+ prompt: str,
181
+ device: torch.device,
182
+ max_new_tokens: int = 200,
183
+ temperature: float = 1.0,
184
+ top_k: int | None = None,
185
+ top_p: float | None = None,
186
+ repetition_penalty: float = 1.0,
187
+ no_repeat_ngram_size: int = 0,
188
+ stream: bool = True,
189
+ seed: int | None = None,
190
+ num_samples: int = 1,
191
+ ) -> list[str]:
192
+ if seed is not None:
193
+ torch.manual_seed(seed)
194
+
195
+ prompt_ids = tokenizer.encode(prompt, add_bos=True, add_eos=False)
196
+ idx = torch.tensor([prompt_ids] * num_samples, dtype=torch.long, device=device)
197
+ block_size = model.config.block_size
198
+ eos_id = getattr(tokenizer, "eos_id", None)
199
+
200
+ generated_ids = [list(prompt_ids) for _ in range(num_samples)]
201
+ finished = [False] * num_samples
202
+ # Decode the FULL sequence each step and print only the new suffix,
203
+ # rather than decoding one new token at a time: BPE/subword tokens don't
204
+ # each cleanly map to standalone text (a token can be half a multi-byte
205
+ # character or mid-word piece), so decoding incrementally token-by-token
206
+ # can emit garbled text at token boundaries. Re-decoding the whole
207
+ # sequence and diffing against what's already been printed sidesteps
208
+ # that regardless of the tokenizer's internals. O(T) decode cost per
209
+ # step, which is fine at CLI-generation scale.
210
+ prev_text = [tokenizer.decode(g) for g in generated_ids] if stream else None
211
+ if stream and num_samples == 1:
212
+ print(prev_text[0], end="", flush=True)
213
+
214
+ for _ in range(max_new_tokens):
215
+ if all(finished):
216
+ break
217
+ idx_cond = idx if idx.size(1) <= block_size else idx[:, -block_size:]
218
+ logits, _ = model(idx_cond)
219
+ logits = logits[:, -1, :].float() # last position, fp32 for stable sampling math regardless of training dtype
220
+
221
+ for b in range(num_samples):
222
+ if finished[b]:
223
+ continue
224
+ row = logits[b:b + 1]
225
+ row = _apply_repetition_penalty(row, generated_ids[b], repetition_penalty)
226
+ row = _block_repeated_ngrams(row, generated_ids[b], no_repeat_ngram_size)
227
+ row = row / max(temperature, 1e-5)
228
+ if top_k is not None:
229
+ row = _top_k_filter(row, top_k)
230
+ if top_p is not None:
231
+ row = _top_p_filter(row, top_p)
232
+ logits[b:b + 1] = row
233
+
234
+ probs = torch.softmax(logits, dim=-1)
235
+ next_ids = torch.multinomial(probs, num_samples=1) # (num_samples, 1)
236
+
237
+ # Once a sequence hits EOS, keep feeding it its own last token so the
238
+ # batch stays rectangular, but stop appending to its recorded output.
239
+ for b in range(num_samples):
240
+ if finished[b]:
241
+ next_ids[b, 0] = idx[b, -1]
242
+ continue
243
+ tid = next_ids[b, 0].item()
244
+ generated_ids[b].append(tid)
245
+ if eos_id is not None and tid == eos_id:
246
+ finished[b] = True
247
+
248
+ idx = torch.cat([idx, next_ids], dim=1)
249
+
250
+ if stream and num_samples == 1:
251
+ new_text = tokenizer.decode(generated_ids[0])
252
+ print(new_text[len(prev_text[0]):], end="", flush=True)
253
+ prev_text[0] = new_text
254
+
255
+ if stream and num_samples == 1:
256
+ print()
257
+
258
+ return [tokenizer.decode(g) for g in generated_ids]
259
+
260
+
261
+ # ─────────────────────────────────────────────────────────────────────────────
262
+ # CLI
263
+ # ─────────────────────────────────────────────────────────────────────────────
264
+
265
+ def build_arg_parser() -> argparse.ArgumentParser:
266
+ p = argparse.ArgumentParser(description="Generate text from a trained Arya/Veylon checkpoint.")
267
+ p.add_argument("--checkpoint", required=True, help="Path to a checkpoint saved by train.py")
268
+ p.add_argument("--tokenizer", default="./tokenizer.model", help="Path passed to TokenizerWrapper")
269
+ p.add_argument("--prompt", default=None, help="Prompt text (omit if using --interactive)")
270
+ p.add_argument("--interactive", action="store_true", help="Drop into a REPL: one prompt per line")
271
+ p.add_argument("--max-new-tokens", type=int, default=200)
272
+ p.add_argument("--temperature", type=float, default=0.8)
273
+ p.add_argument("--top-k", type=int, default=50)
274
+ p.add_argument("--top-p", type=float, default=None, help="Nucleus sampling threshold, e.g. 0.9. Combine freely with --top-k.")
275
+ p.add_argument("--repetition-penalty", type=float, default=1.3,
276
+ help="CTRL-style penalty on already-generated tokens, now count-scaled "
277
+ "(penalty**occurrences). 1.0 disables it -- small models loop into "
278
+ "degenerate repetition ('so so so so...') without some penalty here, "
279
+ "so this defaults ON.")
280
+ p.add_argument("--no-repeat-ngram-size", type=int, default=3,
281
+ help="Hard-ban repeating any n-gram of this size (HF-style). 0 disables it. "
282
+ "This is what actually stops infinite word loops -- repetition-penalty "
283
+ "alone only makes them less likely, this makes them impossible.")
284
+ p.add_argument("--num-samples", type=int, default=1, help="Generate N completions per prompt (batched)")
285
+ p.add_argument("--seed", type=int, default=None)
286
+ p.add_argument("--device", default=None, help="cuda / cpu / mps / xla -- auto-detected if omitted")
287
+ p.add_argument("--dtype", default="auto", choices=["auto", "fp32", "fp16", "bf16"],
288
+ help="Inference precision. 'auto' uses fp16 on CUDA, fp32 elsewhere.")
289
+ return p
290
+
291
+
292
+ def _resolve_dtype(dtype_arg: str, device: torch.device) -> torch.dtype | None:
293
+ if device.type not in ("cuda", "cpu"):
294
+ # torch.autocast is only exercised here against cuda/cpu; MPS and XLA
295
+ # handle precision differently (XLA in particular usually wants
296
+ # XLA_USE_BF16 or torch_xla's own autocast, not torch.autocast).
297
+ # Guessing at that without real hardware to test against would be
298
+ # worse than just running those in fp32 and saying so.
299
+ if dtype_arg != "auto":
300
+ print(f"[Inference] --dtype={dtype_arg} requested but autocast isn't wired up for device "
301
+ f"'{device.type}' -- running in fp32.", file=sys.stderr)
302
+ return None
303
+ if dtype_arg == "fp32":
304
+ return None # no autocast
305
+ if dtype_arg == "fp16":
306
+ return torch.float16
307
+ if dtype_arg == "bf16":
308
+ return torch.bfloat16
309
+ # auto
310
+ return torch.float16 if device.type == "cuda" else None
311
+
312
+
313
+ def main() -> None:
314
+ args = build_arg_parser().parse_args()
315
+ if not args.prompt and not args.interactive:
316
+ print("Provide --prompt \"...\" or pass --interactive.", file=sys.stderr)
317
+ sys.exit(1)
318
+
319
+ device = pick_device(args.device)
320
+ model, model_config = load_model(args.checkpoint, device)
321
+ tokenizer = TokenizerWrapper(args.tokenizer)
322
+ autocast_dtype = _resolve_dtype(args.dtype, device)
323
+
324
+ def run(prompt: str, stream: bool) -> list[str]:
325
+ ctx = torch.autocast(device_type=device.type, dtype=autocast_dtype) if autocast_dtype else contextlib.nullcontext()
326
+ with ctx:
327
+ return generate(
328
+ model, tokenizer, prompt, device,
329
+ max_new_tokens=args.max_new_tokens,
330
+ temperature=args.temperature,
331
+ top_k=args.top_k,
332
+ top_p=args.top_p,
333
+ repetition_penalty=args.repetition_penalty,
334
+ no_repeat_ngram_size=args.no_repeat_ngram_size,
335
+ stream=stream,
336
+ seed=args.seed,
337
+ num_samples=args.num_samples,
338
+ )
339
+
340
+ if args.interactive:
341
+ print("[Inference] Interactive mode. Ctrl-C or empty line to exit.", file=sys.stderr)
342
+ while True:
343
+ try:
344
+ prompt = input("\n>>> ")
345
+ except (EOFError, KeyboardInterrupt):
346
+ break
347
+ if not prompt.strip():
348
+ break
349
+ outputs = run(prompt, stream=(args.num_samples == 1))
350
+ if args.num_samples > 1:
351
+ for i, text in enumerate(outputs):
352
+ print(f"\n--- sample {i} ---\n{text}")
353
+ else:
354
+ outputs = run(args.prompt, stream=(args.num_samples == 1))
355
+ if args.num_samples > 1:
356
+ for i, text in enumerate(outputs):
357
+ print(f"\n--- sample {i} ---\n{text}")
358
+
359
+
360
+ if __name__ == "__main__":
361
+ main()
model.py ADDED
@@ -0,0 +1,1813 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ import functools
2
+ import math
3
+ import warnings
4
+ from dataclasses import dataclass
5
+ from typing import Optional, Tuple
6
+
7
+ import torch
8
+ import torch.nn as nn
9
+ import torch.utils.checkpoint
10
+ from torch.nn import functional as F
11
+
12
+ try:
13
+ import config as alpha_config
14
+ except ImportError: # pragma: no cover - fallback for package-style imports
15
+ from . import config as alpha_config
16
+
17
+ try:
18
+ import kernel
19
+ except ImportError: # pragma: no cover - fallback for package-style imports
20
+ from . import kernel
21
+
22
+ try:
23
+ import fused_kernels
24
+ except ImportError: # pragma: no cover - fallback for package-style imports
25
+ try:
26
+ from . import fused_kernels
27
+ except ImportError:
28
+ fused_kernels = None
29
+
30
+ try:
31
+ import flash_ops
32
+ except ImportError: # pragma: no cover - fallback for package-style imports
33
+ try:
34
+ from . import flash_ops
35
+ except ImportError:
36
+ flash_ops = None
37
+
38
+
39
+
40
+
41
+ # ── nvFP4 via torchao RHT (Blackwell SM100+ only) ───────────────────────────
42
+ # Uses NVFP4DynamicActivationNVFP4WeightConfig from torchao.prototype.mx_formats.
43
+ # This replaces the old manual block-quantization implementation.
44
+ # Dynamic activation scaling: per-tensor FP8 scale computed at runtime.
45
+ # Weight scaling: per-block-16 FP4 scale (RHT = Reciprocal-Half-Tile layout).
46
+ # Both scales are fused into a single Blackwell GEMM kernel via torch.compile.
47
+ # No manual weight surgery needed -- torchao rewrites nn.Linear in-place via quantize_().
48
+
49
+ _TORCHAO_AVAILABLE: bool = False
50
+ _NVFP4_RHT_AVAILABLE: bool = False
51
+ _INT8_AVAILABLE: bool = False
52
+ try:
53
+ import torchao as _torchao_probe # noqa: F401
54
+ _TORCHAO_AVAILABLE = True
55
+ from torchao.quantization import quantize_ as _torchao_quantize
56
+ from torchao.quantization import (
57
+ Int8DynamicActivationInt8WeightConfig as _Int8Config,
58
+ )
59
+ # Per-group INT8 weight quant: available as Int8WeightOnlyConfig with
60
+ # group_size, paired with dynamic int8 activations via a composed config.
61
+ # Stable API: try the grouped variant; fall back to per-channel if absent.
62
+ try:
63
+ from torchao.quantization import Int8WeightOnlyConfig as _Int8WOConfig
64
+ _INT8_GROUPED_AVAILABLE = True
65
+ except ImportError:
66
+ _Int8WOConfig = None # type: ignore[assignment]
67
+ _INT8_GROUPED_AVAILABLE = False
68
+ _INT8_AVAILABLE = True
69
+ try:
70
+ from torchao.prototype.mx_formats import (
71
+ NVFP4DynamicActivationNVFP4WeightConfig as _NVFP4Config,
72
+ )
73
+ _NVFP4_RHT_AVAILABLE = True
74
+ except ImportError:
75
+ _NVFP4Config = None # type: ignore[assignment]
76
+ except ImportError:
77
+ _torchao_quantize = None # type: ignore[assignment]
78
+ _Int8Config = None # type: ignore[assignment]
79
+ _NVFP4Config = None # type: ignore[assignment]
80
+
81
+ _NVFP4_DISABLED_REASON: "str | None" = None
82
+
83
+ # Layers to keep in full precision.
84
+ # lm_head / wte: embedding + output projection -- high vocab sensitivity.
85
+ # norm / ln_: RMSNorm scale weights -- tiny tensors, quantizing them is pointless.
86
+ # router / gate: MoE router logits drive routing decisions; FP4 routing = collapsed experts.
87
+ # expert_bias: per-expert scalar bias, not a Linear weight.
88
+ _NVFP4_SKIP_KEYWORDS: tuple[str, ...] = (
89
+ "wte", "lm_head", "ln_", "ln_f", "norm",
90
+ "router", "gate", "expert_bias",
91
+ )
92
+
93
+
94
+ def _nvfp4_filter(module: nn.Module, fqn: str) -> bool:
95
+ """torchao filter_fn: True = quantize this module."""
96
+ if not isinstance(module, nn.Linear):
97
+ return False
98
+ for kw in _NVFP4_SKIP_KEYWORDS:
99
+ if kw in fqn:
100
+ return False
101
+ # Minimum size guard: torchao's RHT kernel requires in_features % 16 == 0
102
+ # and out_features >= 64. Skip tiny linears (e.g. MoD router D->1).
103
+ out_f, in_f = module.weight.shape
104
+ if in_f < 32 or out_f < 64 or in_f % 16 != 0:
105
+ return False
106
+ return True
107
+
108
+
109
+ def _nvfp4_enabled() -> bool:
110
+ global _NVFP4_DISABLED_REASON
111
+ if not getattr(alpha_config, 'use_nvfp4', False):
112
+ return False
113
+ if not torch.cuda.is_available():
114
+ _NVFP4_DISABLED_REASON = "nvFP4: no CUDA device"
115
+ return False
116
+ major, _minor = torch.cuda.get_device_capability()
117
+ if major < 10:
118
+ msg = (
119
+ f"nvFP4 requires Blackwell (sm_100+); found sm_{major}{_minor} "
120
+ f"({torch.cuda.get_device_name()}) -- disabled."
121
+ )
122
+ if msg not in _alpha_config_warned:
123
+ _alpha_config_warned.add(msg)
124
+ warnings.warn(msg)
125
+ _NVFP4_DISABLED_REASON = msg
126
+ return False
127
+ if not _NVFP4_RHT_AVAILABLE:
128
+ _NVFP4_DISABLED_REASON = (
129
+ "nvFP4: torchao or torchao.prototype.mx_formats not available. "
130
+ "Install: pip install torchao --pre"
131
+ )
132
+ warnings.warn(_NVFP4_DISABLED_REASON)
133
+ return False
134
+ return True
135
+
136
+
137
+ def apply_nvfp4_to_model(model: nn.Module) -> int:
138
+ """Quantize eligible Linear layers to NVFP4 (RHT, dynamic activations).
139
+
140
+ Uses torchao.quantization.quantize_() with
141
+ NVFP4DynamicActivationNVFP4WeightConfig, which applies:
142
+ - Weight: FP4 e2m1 with per-block-16 FP8 scales (RHT layout)
143
+ - Activation: FP4 with per-tensor FP32 dynamic scale (computed at runtime)
144
+
145
+ The fused GEMM kernel is emitted by torch.compile on Blackwell SM100+.
146
+ Must be called BEFORE torch.compile(); quantize_() rewrites modules in-place.
147
+
148
+ Returns the number of quantized Linear layers.
149
+ """
150
+ before = {fqn for fqn, m in model.named_modules() if _nvfp4_filter(m, fqn)}
151
+ _torchao_quantize(model, _NVFP4Config(), filter_fn=_nvfp4_filter)
152
+ count = 0
153
+ for fqn, m in model.named_modules():
154
+ if fqn in before and isinstance(m, nn.Module):
155
+ if type(m).__name__ != "Linear":
156
+ count += 1
157
+ return count
158
+
159
+
160
+ # ── Jetfire-style INT8 via torchao (T4 / Turing sm_75 and any CUDA device) ──
161
+ # Int8DynamicActivationInt8WeightConfig applies:
162
+ # Weight: INT8 symmetric per-group quantization (group_size from config).
163
+ # group_size=128 is the Jetfire default; smaller (32) tightens
164
+ # outlier containment at the cost of more scale overhead.
165
+ # Activation: INT8 dynamic per-tensor scale computed at runtime (no calibration).
166
+ # Inductor fuses the INT8 GEMM + dequant into a single Triton or CUTLASS kernel.
167
+ # Works on T4 (sm_75) and any device with INT8 tensor core support (Turing+).
168
+ # Much cheaper than FP16 GEMM on T4 where INT8 tensor cores run at 2× the rate.
169
+
170
+ # Reuse same skip list as FP4 -- same sensitive layers apply.
171
+ _INT8_SKIP_KEYWORDS: tuple[str, ...] = _NVFP4_SKIP_KEYWORDS
172
+
173
+ _INT8_DISABLED_REASON: "str | None" = None
174
+
175
+
176
+ def _int8_filter(module: nn.Module, fqn: str) -> bool:
177
+ """torchao filter_fn for INT8: True = quantize this module."""
178
+ if not isinstance(module, nn.Linear):
179
+ return False
180
+ for kw in _INT8_SKIP_KEYWORDS:
181
+ if kw in fqn:
182
+ return False
183
+ # group_size constraint: in_features must be divisible by group_size.
184
+ # We can't check the configured group_size here without a closure, so
185
+ # enforce the stricter of the two defaults (128). apply_int8_to_model
186
+ # re-checks with the actual group_size before calling quantize_().
187
+ out_f, in_f = module.weight.shape
188
+ if in_f < 64 or out_f < 32:
189
+ return False
190
+ return True
191
+
192
+
193
+ def _int8_enabled() -> bool:
194
+ global _INT8_DISABLED_REASON
195
+ if not getattr(alpha_config, 'use_int8', False):
196
+ return False
197
+ if not torch.cuda.is_available():
198
+ _INT8_DISABLED_REASON = "INT8: no CUDA device"
199
+ return False
200
+ if not _INT8_AVAILABLE:
201
+ _INT8_DISABLED_REASON = (
202
+ "INT8: torchao not installed or Int8DynamicActivationInt8WeightConfig "
203
+ "not found. Install: pip install torchao"
204
+ )
205
+ warnings.warn(_INT8_DISABLED_REASON)
206
+ return False
207
+ # Mutual exclusion: nvFP4 and INT8 cannot both be active.
208
+ # nvFP4 only runs on SM100+; if the user enabled both on a Blackwell machine
209
+ # nvFP4 takes priority (higher throughput).
210
+ if _nvfp4_enabled():
211
+ msg = (
212
+ "INT8 disabled: use_nvfp4=True also set and this is SM100+ hardware -- "
213
+ "nvFP4 takes priority. Set use_nvfp4=False to use INT8 instead."
214
+ )
215
+ if msg not in _alpha_config_warned:
216
+ _alpha_config_warned.add(msg)
217
+ warnings.warn(msg)
218
+ _INT8_DISABLED_REASON = msg
219
+ return False
220
+ return True
221
+
222
+
223
+ def apply_int8_to_model(model: nn.Module, group_size: int = 128) -> int:
224
+ """Quantize eligible Linear layers to INT8 (Jetfire-style, dynamic activations).
225
+
226
+ Config selection priority:
227
+ 1. Int8DynamicActivationInt8WeightConfig() -- stable torchao, per-channel
228
+ weights + dynamic per-tensor INT8 activations. Best for T4 throughput.
229
+ 2. Int8WeightOnlyConfig(group_size=N) -- weight-only, activations stay fp16.
230
+ Only used as fallback if (1) is unavailable (shouldn't happen on stable).
231
+
232
+ Works on T4 (sm_75+). Fused kernel emitted by torch.compile via Inductor.
233
+ Must be called BEFORE torch.compile().
234
+ """
235
+ def _filter_with_gs(module: nn.Module, fqn: str) -> bool:
236
+ if not _int8_filter(module, fqn):
237
+ return False
238
+ # group_size divisibility only matters for grouped weight configs;
239
+ # Int8DynamicActivationInt8WeightConfig is per-channel so no constraint.
240
+ # Keep the check anyway to skip genuinely misaligned tiny layers.
241
+ in_f = module.weight.shape[1]
242
+ if in_f % 32 != 0: # 32 = minimum alignment for any INT8 kernel
243
+ return False
244
+ return True
245
+
246
+ # Snapshot module ids before -- torchao replaces nn.Linear with a subclass
247
+ # in-place; the fqn stays the same but id() changes. Class-name check is
248
+ # unreliable because torchao's subclass is also called "Linear" in some versions.
249
+ before_ids = {fqn: id(m) for fqn, m in model.named_modules() if _filter_with_gs(m, fqn)}
250
+
251
+ # Prefer Int8DynamicActivationInt8WeightConfig: dynamic INT8 on both weight
252
+ # and activation, fused into a single CUTLASS/Triton INT8 GEMM by Inductor.
253
+ # This is the correct stable API -- no constructor args needed.
254
+ if _INT8_AVAILABLE and _Int8Config is not None:
255
+ config_obj = _Int8Config()
256
+ elif _INT8_GROUPED_AVAILABLE and _Int8WOConfig is not None:
257
+ warnings.warn(
258
+ f"[INT8] Int8DynamicActivationInt8WeightConfig unavailable -- "
259
+ f"falling back to Int8WeightOnlyConfig(group_size={group_size}). "
260
+ f"Activations will stay fp16 (weight-only quant).",
261
+ )
262
+ config_obj = _Int8WOConfig(group_size=group_size)
263
+ else:
264
+ raise RuntimeError("[INT8] No usable INT8 config found in torchao -- pip install torchao")
265
+
266
+ _torchao_quantize(model, config_obj, filter_fn=_filter_with_gs)
267
+
268
+ # Count by id() change: torchao replaces the module object in-place on the
269
+ # parent, so the same fqn now resolves to a different object.
270
+ count = sum(
271
+ 1 for fqn, old_id in before_ids.items()
272
+ if id(dict(model.named_modules()).get(fqn)) != old_id
273
+ )
274
+ return count
275
+
276
+
277
+
278
+
279
+
280
+ _alpha_config_warned: set[str] = set()
281
+
282
+
283
+ def _alpha_config_attr(name: str, default):
284
+
285
+ if not hasattr(alpha_config, name):
286
+ if name not in _alpha_config_warned:
287
+ _alpha_config_warned.add(name)
288
+ warnings.warn(
289
+ f"config module has no attribute '{name}' -- falling back to "
290
+ f"default {default!r}. If this is unexpected, check for a "
291
+ f"casing mismatch or rename in your config.py.",
292
+ stacklevel=2,
293
+ )
294
+ return getattr(alpha_config, name, default)
295
+
296
+
297
+ try:
298
+ from torch.nn.functional import rms_norm as _native_rms_norm
299
+ except (ImportError, AttributeError):
300
+ _native_rms_norm = None
301
+
302
+
303
+
304
+ try:
305
+ from liger_kernel.ops.rms_norm import LigerRMSNormFunction as _LigerRMSNormFn
306
+ except ImportError:
307
+ _LigerRMSNormFn = None
308
+
309
+ try:
310
+ from liger_kernel.ops.rope import LigerRopeFunction as _LigerRopeFn
311
+ except ImportError:
312
+ _LigerRopeFn = None
313
+
314
+ try:
315
+ from liger_kernel.transformers.fused_linear_cross_entropy import (
316
+ LigerFusedLinearCrossEntropyLoss as _LigerFusedLinearCrossEntropyLoss,
317
+ )
318
+ except ImportError:
319
+ _LigerFusedLinearCrossEntropyLoss = None
320
+
321
+ _liger_warned: set[str] = set()
322
+
323
+
324
+ def _liger_warn_once(key: str, msg: str) -> None:
325
+ if key not in _liger_warned:
326
+ _liger_warned.add(key)
327
+ warnings.warn(f"[model.py/liger] {msg}", stacklevel=3)
328
+
329
+
330
+ def _liger_cuda_gate(x: torch.Tensor) -> bool:
331
+
332
+ return x.device.type == "cuda"
333
+
334
+
335
+ def _use_liger_config() -> bool:
336
+
337
+ return bool(_alpha_config_attr('use_liger', False))
338
+
339
+
340
+ def _use_flash_ops_master() -> bool:
341
+
342
+ return bool(_alpha_config_attr('use_flash_ops', True))
343
+
344
+
345
+ def _use_flash_rope_config() -> bool:
346
+
347
+ return _use_flash_ops_master() and bool(_alpha_config_attr('use_flash_rope', True))
348
+
349
+
350
+ def _use_flash_outproj_config() -> bool:
351
+
352
+ return _use_flash_ops_master() and bool(_alpha_config_attr('use_flash_outproj_add_rmsnorm', True))
353
+
354
+
355
+ def _use_flash_swiglu_config() -> bool:
356
+
357
+ return _use_flash_ops_master() and bool(_alpha_config_attr('use_flash_swiglu', True))
358
+
359
+
360
+
361
+ def _use_fused_add_rmsnorm_config() -> bool:
362
+
363
+
364
+ return bool(_alpha_config_attr('use_fused_add_rmsnorm', True))
365
+
366
+
367
+ @torch.compiler.disable
368
+ def _liger_rmsnorm_attempt(x: torch.Tensor, weight: torch.Tensor, eps: float):
369
+
370
+ try:
371
+ out = _LigerRMSNormFn.apply(x, weight, eps)
372
+ kernel._record_backend("liger_rmsnorm_success")
373
+ return out
374
+ except Exception as e: # noqa: BLE001 -- must never crash training
375
+ kernel._record_backend(f"liger_rmsnorm_fallback:{type(e).__name__}")
376
+ _liger_warn_once(
377
+ f"rmsnorm_fail_{type(e).__name__}",
378
+ f"Liger RMSNorm failed on {x.device} ({e!r}) -- falling back to native/eager RMSNorm.",
379
+ )
380
+ return None
381
+
382
+
383
+ class RMSNorm(nn.Module):
384
+ """Root Mean Square Layer Normalization (used in modern transformers)"""
385
+ def __init__(self, dim, eps=1e-5):
386
+ super().__init__()
387
+ self.eps = eps
388
+ self.weight = nn.Parameter(torch.ones(dim))
389
+ self.dim = dim
390
+
391
+ def forward(self, x):
392
+
393
+ if _use_liger_config() and _LigerRMSNormFn is not None and _liger_cuda_gate(x):
394
+ out = _liger_rmsnorm_attempt(x, self.weight, self.eps)
395
+ if out is not None:
396
+ return out
397
+ if _native_rms_norm is not None:
398
+
399
+ return _native_rms_norm(x, (self.dim,), self.weight, self.eps)
400
+
401
+ x_fp32 = x.float()
402
+ rms = torch.sqrt(x_fp32.pow(2).mean(-1, keepdim=True) + self.eps)
403
+ weight = self.weight.to(x.dtype)
404
+ return ((x_fp32 / rms).to(x.dtype)) * weight
405
+ def rotate_half(x):
406
+ x1, x2 = x.chunk(2, dim=-1)
407
+ return torch.cat((-x2, x1), dim=-1)
408
+
409
+
410
+ @dataclass
411
+ class AttentionOutput:
412
+ """Uniform return type so callers can optionally inspect weights / cache cost."""
413
+
414
+ output: torch.Tensor
415
+ attention_weights: Optional[torch.Tensor] = None
416
+ kv_cache_bytes: Optional[int] = None
417
+ selected_indices: Optional[torch.Tensor] = None
418
+
419
+
420
+ def module_output_tensor(x):
421
+
422
+ return x.output if isinstance(x, AttentionOutput) else x
423
+
424
+
425
+ def reshape_heads(x: torch.Tensor, num_heads: int, head_dim: int) -> torch.Tensor:
426
+
427
+ b, t, _ = x.shape
428
+ return x.view(b, t, num_heads, head_dim).transpose(1, 2)
429
+
430
+
431
+ def merge_heads(x: torch.Tensor) -> torch.Tensor:
432
+
433
+ b, h, t, d = x.shape
434
+ return x.transpose(1, 2).contiguous().view(b, t, h * d)
435
+
436
+
437
+ def expand_kv_heads(x: torch.Tensor, num_heads: int) -> torch.Tensor:
438
+
439
+ b, h_kv, t, d = x.shape
440
+ if h_kv == num_heads:
441
+ return x
442
+ if num_heads % h_kv != 0:
443
+ raise ValueError(f"num_heads ({num_heads}) must be divisible by kv heads ({h_kv})")
444
+ reps = num_heads // h_kv
445
+ return x.repeat_interleave(reps, dim=1)
446
+
447
+
448
+ def _yarn_find_correction_dim(num_rotations, dim, base=10000.0, orig_max_pos=2048):
449
+
450
+ return (dim * math.log(orig_max_pos / (num_rotations * 2 * math.pi))) / (2 * math.log(base))
451
+
452
+
453
+ def _yarn_find_correction_range(low_rot, high_rot, dim, base=10000.0, orig_max_pos=2048):
454
+ low = math.floor(_yarn_find_correction_dim(low_rot, dim, base, orig_max_pos))
455
+ high = math.ceil(_yarn_find_correction_dim(high_rot, dim, base, orig_max_pos))
456
+ return max(low, 0), min(high, dim - 1)
457
+
458
+
459
+ def _yarn_linear_ramp_mask(low, high, dim):
460
+
461
+ if low > high:
462
+ low, high = high, low
463
+ if low == high:
464
+ high += 0.001 # avoid a divide-by-zero when the ramp is degenerate
465
+ ramp = (torch.arange(dim, dtype=torch.float32) - low) / (high - low)
466
+ return torch.clamp(ramp, 0, 1)
467
+
468
+
469
+ def yarn_attention_temperature(scale: float) -> float:
470
+
471
+ if scale <= 1.0:
472
+ return 1.0
473
+ t = 0.1 * math.log(scale) + 1.0
474
+ return 1.0 / t
475
+
476
+
477
+ @functools.lru_cache(maxsize=64)
478
+ def _rope_cache(seq_len: int, dim: int, device, dtype, base: float = 10000.0):
479
+ inv_freq = 1.0 / (base ** (torch.arange(0, dim, 2, device=device, dtype=torch.float32) / dim))
480
+ t = torch.arange(seq_len, device=device, dtype=torch.float32)
481
+ freqs = torch.outer(t, inv_freq)
482
+ emb = torch.cat((freqs, freqs), dim=-1)
483
+ return emb.cos().to(dtype), emb.sin().to(dtype)
484
+
485
+
486
+ def build_rope_cache(
487
+ seq_len: int,
488
+ dim: int,
489
+ base: float = 10000.0,
490
+ yarn_scale: float = 1.0,
491
+ yarn_orig_max_pos: int | None = None,
492
+ yarn_alpha: float = 1.0,
493
+ yarn_beta: float = 32.0,
494
+ ) -> tuple[torch.Tensor, torch.Tensor]:
495
+ if yarn_scale <= 1.0:
496
+ inv_freq = 1.0 / (base ** (torch.arange(0, dim, 2, dtype=torch.float32) / dim))
497
+ else:
498
+ if yarn_orig_max_pos is None:
499
+ raise ValueError(
500
+ "build_rope_cache: yarn_scale > 1.0 requires yarn_orig_max_pos "
501
+ "(the ORIGINAL trained context length, i.e. CONTEXT/block_size "
502
+ "from config.py before scaling) to compute the NTK-by-parts "
503
+ "ramp correctly -- it can't be inferred from seq_len alone, "
504
+ "since seq_len here is already the SCALED target length."
505
+ )
506
+ pos_freqs = base ** (torch.arange(0, dim, 2, dtype=torch.float32) / dim)
507
+ inv_freq_extrapolation = 1.0 / pos_freqs # gamma=1 region: untouched
508
+ inv_freq_interpolation = 1.0 / (yarn_scale * pos_freqs) # gamma=0 region: NTK/linear-interpolated
509
+ low, high = _yarn_find_correction_range(
510
+ yarn_alpha, yarn_beta, dim, base, yarn_orig_max_pos
511
+ )
512
+ if low > high:
513
+ warnings.warn(
514
+ f"YaRN correction range inverted (low={low} > high={high}) "
515
+ f"for dim={dim} with yarn_alpha={yarn_alpha}, yarn_beta="
516
+ f"{yarn_beta} -- these defaults were tuned for LLaMA-family "
517
+ f"head_dim=128 (per Peng et al. 2023), and can invert at "
518
+ f"smaller head_dim. Auto-corrected by swapping low/high so "
519
+ f"the ramp stays monotonic, but the resulting correction "
520
+ f"window differs from LLaMA's -- if generation quality at "
521
+ f"the extended length looks off, try tuning yarn_alpha/"
522
+ f"yarn_beta in config.py rather than trusting these "
523
+ f"defaults blindly for this model's shape."
524
+ )
525
+
526
+ inv_freq_mask = 1.0 - _yarn_linear_ramp_mask(low, high, dim // 2)
527
+ inv_freq = (
528
+ inv_freq_interpolation * (1 - inv_freq_mask)
529
+ + inv_freq_extrapolation * inv_freq_mask
530
+ )
531
+ t = torch.arange(seq_len, dtype=torch.float32)
532
+ freqs = torch.outer(t, inv_freq)
533
+ emb = torch.cat((freqs, freqs), dim=-1)
534
+ return emb.cos(), emb.sin()
535
+
536
+
537
+ def apply_rope(
538
+ x: torch.Tensor,
539
+ positions: torch.Tensor | None = None,
540
+ rope_dim: int | None = None,
541
+ max_position: int | None = None,
542
+ cos_full: torch.Tensor | None = None,
543
+ sin_full: torch.Tensor | None = None,
544
+ ) -> torch.Tensor:
545
+
546
+ b, h, t, d = x.shape
547
+ rope_dim = d if rope_dim is None else rope_dim
548
+
549
+ is_sequential = positions is None
550
+ if positions is None:
551
+ positions = torch.arange(t, device=x.device)
552
+
553
+ if cos_full is None or sin_full is None:
554
+ if max_position is None:
555
+ max_position = t
556
+ max_pos = int(max_position) if max_position is not None else int(positions.max().item()) + 1 if positions.numel() > 0 else 1
557
+ cos_full, sin_full = _rope_cache(max_pos, rope_dim, x.device, x.dtype)
558
+
559
+
560
+ if is_sequential:
561
+ if t > cos_full.shape[0]:
562
+ raise ValueError(
563
+ f"RoPE table too small: highest requested position="
564
+ f"{t - 1}, table size={cos_full.shape[0]}. If this "
565
+ f"happened during generate(), set gen_headroom in config.py "
566
+ f"to however many extra positions past CONTEXT (block_size) "
567
+ f"generation needs -- see GPT.__init__'s gen_headroom "
568
+ f"wiring. This is separate from max_gen_tokens (a hard cap "
569
+ f"on tokens produced per generate() call, unrelated to "
570
+ f"RoPE table size). Default gen_headroom is 0 if unset."
571
+ )
572
+ elif positions is not None and positions.numel() > 0:
573
+ highest_pos = int(positions.max().item())
574
+ if highest_pos >= cos_full.shape[0]:
575
+ raise ValueError(
576
+ f"RoPE table too small: highest requested position="
577
+ f"{highest_pos}, table size={cos_full.shape[0]}. If this "
578
+ f"happened during generate(), set gen_headroom in config.py "
579
+ f"to however many extra positions past CONTEXT (block_size) "
580
+ f"generation needs -- see GPT.__init__'s gen_headroom "
581
+ f"wiring. This is separate from max_gen_tokens (a hard cap "
582
+ f"on tokens produced per generate() call, unrelated to "
583
+ f"RoPE table size). Default gen_headroom is 0 if unset."
584
+ )
585
+
586
+
587
+ if is_sequential:
588
+ cos = cos_full[:t].to(x.dtype).unsqueeze(0).unsqueeze(0)
589
+ sin = sin_full[:t].to(x.dtype).unsqueeze(0).unsqueeze(0)
590
+ else:
591
+ cos = cos_full[positions].to(x.dtype).unsqueeze(0).unsqueeze(0)
592
+ sin = sin_full[positions].to(x.dtype).unsqueeze(0).unsqueeze(0)
593
+
594
+ x_rot, x_pass = x[..., :rope_dim], x[..., rope_dim:]
595
+ x_rot = (x_rot * cos) + (rotate_half(x_rot) * sin)
596
+ return torch.cat([x_rot, x_pass], dim=-1) if x_pass.shape[-1] > 0 else x_rot
597
+
598
+
599
+ def apply_rope_qk(
600
+ q: torch.Tensor,
601
+ k: torch.Tensor,
602
+ cos_full: torch.Tensor,
603
+ sin_full: torch.Tensor,
604
+ rope_dim: int,
605
+ positions: torch.Tensor | None = None,
606
+ ) -> tuple[torch.Tensor, torch.Tensor]:
607
+
608
+ b, h, t, d = q.shape
609
+ hk = k.shape[1]
610
+
611
+
612
+ if _use_flash_rope_config() and flash_ops is not None and q.is_cuda:
613
+ cos_t = cos_full[:t].to(q.dtype)
614
+ sin_t = sin_full[:t].to(q.dtype)
615
+ result = flash_ops.fused_rope_qk(q, k, cos_t, sin_t, rope_dim)
616
+ if result is not None:
617
+ return result
618
+
619
+ if (
620
+ _use_liger_config()
621
+ and _LigerRopeFn is not None
622
+ and _liger_cuda_gate(q)
623
+ and rope_dim == d
624
+ and h == hk
625
+ ):
626
+ result = _liger_rope_attempt(q, k, cos_full[:t], sin_full[:t])
627
+ if result is not None:
628
+ return result
629
+ q_rot = apply_rope(q, positions=positions, rope_dim=rope_dim, cos_full=cos_full, sin_full=sin_full)
630
+ k_rot = apply_rope(k, positions=positions, rope_dim=rope_dim, cos_full=cos_full, sin_full=sin_full)
631
+ return q_rot, k_rot
632
+
633
+
634
+ @torch.compiler.disable
635
+ def _liger_rope_attempt(q: torch.Tensor, k: torch.Tensor, cos: torch.Tensor, sin: torch.Tensor):
636
+
637
+ try:
638
+ cos = cos.to(q.dtype)
639
+ sin = sin.to(q.dtype)
640
+ q_rot, k_rot = _LigerRopeFn.apply(q, k, cos, sin)
641
+ kernel._record_backend("liger_rope_success")
642
+ return q_rot, k_rot
643
+ except Exception as e: # noqa: BLE001 -- must never crash training
644
+ kernel._record_backend(f"liger_rope_fallback:{type(e).__name__}")
645
+ _liger_warn_once(
646
+ f"rope_fail_{type(e).__name__}",
647
+ f"Liger RoPE failed on {q.device} ({e!r}) -- falling back to eager RoPE.",
648
+ )
649
+ return None
650
+
651
+
652
+ @torch.compiler.disable
653
+ def _liger_fused_ce_attempt(hidden: torch.Tensor, weight: torch.Tensor, targets: torch.Tensor):
654
+
655
+ try:
656
+ loss_fn = _LigerFusedLinearCrossEntropyLoss(ignore_index=-1)
657
+ # Liger's fused linear+CE kernel expects the activation tensor first
658
+ # and the output projection weight second; the previous order was
659
+ # swapped and would pass an invalid tensor layout to the kernel.
660
+ loss = loss_fn(hidden.reshape(-1, hidden.size(-1)), weight, targets.reshape(-1))
661
+ kernel._record_backend("liger_fused_ce_success")
662
+ return loss
663
+ except Exception as e: # noqa: BLE001 -- must never crash training
664
+ kernel._record_backend(f"liger_fused_ce_fallback:{type(e).__name__}")
665
+ _liger_warn_once(
666
+ f"fused_ce_fail_{type(e).__name__}",
667
+ f"Liger fused linear CE failed on {hidden.device} ({e!r}) -- "
668
+ f"falling back to eager lm_head + F.cross_entropy.",
669
+ )
670
+ return None
671
+
672
+
673
+ def make_causal_mask(q_len: int, k_len: int, device) -> torch.Tensor:
674
+
675
+ q_idx = torch.arange(q_len, device=device).view(q_len, 1)
676
+ k_idx = torch.arange(k_len, device=device).view(1, k_len)
677
+ return torch.where(k_idx <= q_idx, torch.zeros(1, device=device), torch.full((1,), -1e9, device=device))
678
+
679
+
680
+ def make_sliding_window_causal_mask(q_len: int, k_len: int, window_size: int, device) -> torch.Tensor:
681
+
682
+ q_idx = torch.arange(q_len, device=device).view(q_len, 1)
683
+ k_idx = torch.arange(k_len, device=device).view(1, k_len)
684
+ visible = (k_idx <= q_idx) & (k_idx > q_idx - window_size)
685
+ return torch.where(visible, torch.zeros(1, device=device), torch.full((1,), -1e9, device=device))
686
+
687
+
688
+
689
+ def masked_softmax(scores: torch.Tensor, mask: torch.Tensor | None, dim: int = -1) -> torch.Tensor:
690
+ if mask is not None:
691
+ scores = scores + mask.to(scores.dtype)
692
+
693
+ return torch.softmax(scores.float(), dim=dim).to(scores.dtype)
694
+
695
+
696
+ def estimate_kv_cache_bytes(
697
+ batch_size, seq_len, num_kv_heads, head_dim, dtype,
698
+ compression_ratio: int = 1, window_size: int | None = None,
699
+ ) -> int:
700
+
701
+ elem_size = torch.zeros(1, dtype=dtype).element_size()
702
+ effective_len = max(1, -(-seq_len // compression_ratio)) # ceil div
703
+ if window_size is not None:
704
+ effective_len = min(effective_len, window_size)
705
+ return int(batch_size * num_kv_heads * effective_len * head_dim * elem_size * 2)
706
+
707
+
708
+
709
+
710
+ class DenseMHA(nn.Module):
711
+
712
+
713
+ def __init__(
714
+ self,
715
+ hidden_size: int,
716
+ num_heads: int,
717
+ head_dim: int | None = None,
718
+ num_kv_heads: int | None = None,
719
+ dropout: float = 0.0,
720
+ use_rope: bool = True,
721
+ causal: bool = True,
722
+ bias: bool = True,
723
+ max_seq_len: int = 4096,
724
+ window_size: int | None = None,
725
+ tp_size: int = 1,
726
+ rope_dim: int | None = None,
727
+ yarn_scale: float = 1.0,
728
+ yarn_orig_max_pos: int | None = None,
729
+ yarn_alpha: float = 1.0,
730
+ yarn_beta: float = 32.0,
731
+ ) -> None:
732
+ super().__init__()
733
+ if head_dim is None:
734
+ if hidden_size % num_heads != 0:
735
+ raise ValueError("hidden_size must be divisible by num_heads when head_dim is omitted")
736
+ head_dim = hidden_size // num_heads
737
+ num_kv_heads = num_heads if num_kv_heads is None else num_kv_heads
738
+ if num_heads % num_kv_heads != 0:
739
+ raise ValueError(f"num_heads ({num_heads}) must be divisible by num_kv_heads ({num_kv_heads})")
740
+ if num_heads % tp_size != 0:
741
+ raise ValueError(
742
+ f"tensor_parallel_size ({tp_size}) must divide num_heads ({num_heads}) "
743
+ f"for this layer -- reduce tensor_parallel_size or increase n_head."
744
+ )
745
+
746
+ local_num_heads = num_heads // tp_size
747
+
748
+ if num_kv_heads < tp_size:
749
+ local_num_kv_heads = num_kv_heads
750
+ else:
751
+ local_num_kv_heads = num_kv_heads // tp_size if num_kv_heads % tp_size == 0 else num_kv_heads
752
+
753
+ self.hidden_size = hidden_size
754
+ self.num_heads = local_num_heads
755
+ self.num_kv_heads = local_num_kv_heads
756
+ self.head_dim = head_dim
757
+ self.inner_dim = local_num_heads * head_dim
758
+ self.use_rope = use_rope
759
+ self.causal = causal
760
+ self.tp_size = tp_size
761
+ self.tp_group = None
762
+ self.window_size = window_size
763
+
764
+ self.rope_dim = head_dim if rope_dim is None else rope_dim
765
+
766
+
767
+
768
+ self._q_width = self.inner_dim
769
+ self._k_width = self.num_kv_heads * head_dim
770
+ self._v_width = self.num_kv_heads * head_dim
771
+ self.qkv_proj = nn.Linear(hidden_size, self._q_width + self._k_width + self._v_width, bias=bias)
772
+ self.out_proj = nn.Linear(self.inner_dim, hidden_size, bias=bias)
773
+ self.dropout = nn.Dropout(dropout)
774
+
775
+ if use_rope:
776
+
777
+ rope_cos, rope_sin = build_rope_cache(
778
+ max_seq_len, self.rope_dim,
779
+ yarn_scale=yarn_scale, yarn_orig_max_pos=yarn_orig_max_pos,
780
+ yarn_alpha=yarn_alpha, yarn_beta=yarn_beta,
781
+ )
782
+ self.register_buffer("rope_cos", rope_cos, persistent=False)
783
+ self.register_buffer("rope_sin", rope_sin, persistent=False)
784
+
785
+ self.yarn_temp_factor = yarn_attention_temperature(yarn_scale)
786
+
787
+ def forward(
788
+ self,
789
+ hidden_states: torch.Tensor,
790
+ output_attentions: bool = False,
791
+ return_kv_cache_estimate: bool = False,
792
+ skip_out_proj: bool = False,
793
+ positions: torch.Tensor | None = None,
794
+ ) -> torch.Tensor | AttentionOutput:
795
+ batch_size, seq_len, _ = hidden_states.shape
796
+
797
+ qkv = self.qkv_proj(hidden_states)
798
+ q_flat, k_flat, v_flat = qkv.split([self._q_width, self._k_width, self._v_width], dim=-1)
799
+ q = reshape_heads(q_flat, self.num_heads, self.head_dim)
800
+ k = reshape_heads(k_flat, self.num_kv_heads, self.head_dim)
801
+ v = reshape_heads(v_flat, self.num_kv_heads, self.head_dim)
802
+
803
+ if self.use_rope:
804
+ q, k = apply_rope_qk(
805
+ q,
806
+ k,
807
+ cos_full=self.rope_cos,
808
+ sin_full=self.rope_sin,
809
+ rope_dim=self.rope_dim,
810
+ positions=positions,
811
+ )
812
+ if self.yarn_temp_factor != 1.0:
813
+
814
+ q = q * self.yarn_temp_factor
815
+
816
+ if output_attentions:
817
+
818
+ if self.num_kv_heads != self.num_heads:
819
+ reps = self.num_heads // self.num_kv_heads
820
+ k_eager = k.repeat_interleave(reps, dim=1)
821
+ v_eager = v.repeat_interleave(reps, dim=1)
822
+ else:
823
+ k_eager, v_eager = k, v
824
+ scores = torch.matmul(q, k_eager.transpose(-2, -1)) / math.sqrt(self.head_dim)
825
+ mask = None
826
+ if self.window_size is not None:
827
+ mask = make_sliding_window_causal_mask(
828
+ seq_len, seq_len, self.window_size, hidden_states.device
829
+ ).view(1, 1, seq_len, seq_len)
830
+ elif self.causal:
831
+ mask = make_causal_mask(seq_len, seq_len, hidden_states.device).view(1, 1, seq_len, seq_len)
832
+ weights = masked_softmax(scores, mask, dim=-1)
833
+
834
+ weights = self.dropout(weights)
835
+ context = torch.matmul(weights, v_eager)
836
+ else:
837
+ context = kernel.fused_attention(
838
+ q,
839
+ k,
840
+ v,
841
+ causal=self.causal,
842
+ window_size=self.window_size,
843
+ dropout_p=self.dropout.p,
844
+ training=self.training,
845
+ )
846
+ weights = None
847
+
848
+ merged_context = merge_heads(context)
849
+
850
+
851
+ if skip_out_proj and self.tp_size == 1 and not output_attentions and not return_kv_cache_estimate:
852
+ return merged_context
853
+
854
+ output = self.out_proj(merged_context)
855
+
856
+
857
+ if self.tp_size > 1:
858
+ if self.tp_group is None:
859
+ raise RuntimeError(
860
+ "DenseMHA was built with tp_size>1 but tp_group was never "
861
+ "set -- call GPT.set_tp_group(group) after construction, "
862
+ "before the first forward pass."
863
+ )
864
+ torch.distributed.all_reduce(output, op=torch.distributed.ReduceOp.SUM, group=self.tp_group)
865
+
866
+ if output_attentions or return_kv_cache_estimate:
867
+ cache_bytes = None
868
+ if return_kv_cache_estimate:
869
+ cache_bytes = estimate_kv_cache_bytes(
870
+ batch_size, seq_len, self.num_kv_heads, self.head_dim,
871
+ hidden_states.dtype, window_size=self.window_size,
872
+ )
873
+ return AttentionOutput(output=output, attention_weights=weights if output_attentions else None, kv_cache_bytes=cache_bytes)
874
+ return output
875
+
876
+
877
+
878
+ class SwiGLU_FFN(nn.Module):
879
+ """
880
+ SwiGLU Feed Forward Network.
881
+
882
+
883
+ """
884
+
885
+ def __init__(self, d_model, ffn_mult=8/3, dropout=0.0, tp_size: int = 1):
886
+ super().__init__()
887
+
888
+ raw_hidden = d_model * ffn_mult
889
+ hidden_dim = int(((raw_hidden + 63) // 64) * 64)
890
+ if hidden_dim % tp_size != 0:
891
+ raise ValueError(
892
+ f"tensor_parallel_size ({tp_size}) must divide the FFN hidden_dim "
893
+ f"({hidden_dim} = round({d_model} * {ffn_mult}) -> next mult of 64) "
894
+ f"-- adjust ffn_mult or tensor_parallel_size."
895
+ )
896
+ local_hidden_dim = hidden_dim // tp_size
897
+ self.tp_size = tp_size
898
+ self.tp_group = None # set post-construction via GPT.set_tp_group()
899
+ self.W_gate = nn.Linear(d_model, local_hidden_dim, bias=False)
900
+ self.W_value = nn.Linear(d_model, local_hidden_dim, bias=False)
901
+ self.W_out = nn.Linear(local_hidden_dim, d_model, bias=False)
902
+ self.dropout = nn.Dropout(dropout)
903
+
904
+ def forward(self, x):
905
+ gate_raw = self.W_gate(x)
906
+ value = self.W_value(x)
907
+
908
+
909
+ hidden = None
910
+ if _use_flash_swiglu_config() and flash_ops is not None and gate_raw.is_cuda:
911
+ hidden = flash_ops.fused_swiglu(gate_raw, value)
912
+ if hidden is None:
913
+ hidden = F.silu(gate_raw) * value
914
+
915
+ hidden = self.dropout(hidden)
916
+ out = self.W_out(hidden)
917
+ if self.tp_size > 1:
918
+ if self.tp_group is None:
919
+ raise RuntimeError(
920
+ "SwiGLU_FFN was built with tp_size>1 but tp_group was never "
921
+ "set -- call GPT.set_tp_group(group) after construction, "
922
+ "before the first forward pass."
923
+ )
924
+ torch.distributed.all_reduce(out, op=torch.distributed.ReduceOp.SUM, group=self.tp_group)
925
+ return out # Return delta only
926
+ class AryaSparseMoE(nn.Module):
927
+ """
928
+ DeepSeek-style Sparse Mixture of Experts.
929
+ Expects input to be ALREADY normalized (Pre-LN).
930
+ Returns the combined FFN delta.
931
+ """
932
+
933
+ def __init__(self, d_model, n_experts=8, n_shared=1, top_k=2, ffn_mult=8/3, dropout=0.0,
934
+ use_capacity_routing=False, capacity_factor=1.25, fixed_capacity: int | None = None):
935
+ super().__init__()
936
+ self.n_experts = n_experts
937
+ self.n_routed = n_experts - n_shared
938
+ self.top_k = top_k
939
+ self.d_model = d_model
940
+ self.use_capacity_routing = use_capacity_routing
941
+ self.capacity_factor = capacity_factor
942
+ self.fixed_capacity = fixed_capacity
943
+
944
+ if self.n_routed < top_k:
945
+ raise ValueError(
946
+ f"AryaSparseMoE config error: n_experts={n_experts} - n_shared={n_shared} "
947
+ f"= {self.n_routed} routed experts, but top_k={top_k} routing needs at "
948
+ f"least {top_k} routed experts to choose from."
949
+ )
950
+
951
+ self.shared_expert = SwiGLU_FFN(d_model, ffn_mult, dropout=dropout)
952
+ self.routed_experts = nn.ModuleList([
953
+ SwiGLU_FFN(d_model, ffn_mult, dropout=dropout)
954
+ for _ in range(self.n_routed)
955
+ ])
956
+
957
+ self.router = nn.Linear(d_model, self.n_routed, bias=False)
958
+ self.expert_bias = nn.Parameter(torch.zeros(self.n_routed))
959
+
960
+ def forward(self, x):
961
+ B, T, D = x.shape
962
+
963
+
964
+ shared_out = self.shared_expert(x)
965
+
966
+
967
+ router_logits = self.router(x)
968
+ selection_scores = router_logits + self.expert_bias
969
+ full_probs = F.softmax(selection_scores.float(), dim=-1)
970
+ topk_probs, topk_indices = full_probs.topk(self.top_k, dim=-1)
971
+ gate_weights = topk_probs.to(router_logits.dtype)
972
+
973
+ x_flat = x.view(B * T, D)
974
+ gate_flat = gate_weights.view(B * T, self.top_k)
975
+ idx_flat = topk_indices.view(B * T, self.top_k)
976
+
977
+ if self.use_capacity_routing:
978
+ routed_out = self._forward_capacity_routed(x_flat, gate_flat, idx_flat)
979
+ routed_out = routed_out.view(B, T, D)
980
+ out = shared_out + routed_out
981
+ return out
982
+
983
+ routed_out = torch.zeros_like(x_flat)
984
+ for expert_id in range(self.n_routed):
985
+ expert = self.routed_experts[expert_id]
986
+ token_mask = (idx_flat == expert_id).any(dim=-1)
987
+ if not token_mask.any():
988
+ continue
989
+
990
+ selected_tokens = x_flat[token_mask]
991
+ expert_out = expert(selected_tokens) # Returns delta only
992
+ weight = (
993
+ (idx_flat[token_mask] == expert_id).float() * gate_flat[token_mask]
994
+ ).sum(dim=-1)
995
+
996
+ routed_out[token_mask] += weight.unsqueeze(-1) * expert_out
997
+
998
+ routed_out = routed_out.view(B, T, D)
999
+ out = shared_out + routed_out
1000
+ return out
1001
+
1002
+ def _forward_capacity_routed(self, x_flat, gate_flat, idx_flat):
1003
+ n_tokens, D = x_flat.shape
1004
+ n_routed, capacity = self.n_routed, self._capacity(n_tokens)
1005
+
1006
+ flat_expert_ids = idx_flat.reshape(-1)
1007
+ flat_gates = gate_flat.reshape(-1)
1008
+ flat_token_ids = torch.arange(n_tokens, device=x_flat.device).unsqueeze(1).expand(-1, self.top_k).reshape(-1)
1009
+
1010
+
1011
+ one_hot = F.one_hot(flat_expert_ids, num_classes=n_routed).to(torch.float32)
1012
+ position_in_expert = (one_hot.cumsum(dim=0) * one_hot).sum(dim=1).long() - 1
1013
+ keep = position_in_expert < capacity
1014
+ safe_position = position_in_expert.clamp(min=0, max=capacity - 1)
1015
+
1016
+ flat_slot = flat_expert_ids * capacity + safe_position
1017
+
1018
+ gathered_x = x_flat[flat_token_ids] * keep.unsqueeze(-1).to(x_flat.dtype)
1019
+ dispatch = torch.zeros(n_routed * capacity, D, device=x_flat.device, dtype=x_flat.dtype)
1020
+ dispatch.index_add_(0, flat_slot, gathered_x)
1021
+ dispatch = dispatch.view(n_routed, capacity, D)
1022
+
1023
+ expert_out = torch.stack(
1024
+ [expert(dispatch[e]) for e, expert in enumerate(self.routed_experts)],
1025
+ dim=0,
1026
+ ).view(n_routed * capacity, D)
1027
+
1028
+ contrib = expert_out[flat_slot] * keep.unsqueeze(-1).to(x_flat.dtype) * flat_gates.unsqueeze(-1)
1029
+ routed_out = torch.zeros(n_tokens, D, device=x_flat.device, dtype=x_flat.dtype)
1030
+ routed_out.index_add_(0, flat_token_ids, contrib)
1031
+ return routed_out
1032
+
1033
+ def _capacity(self, n_tokens) -> int:
1034
+ if self.fixed_capacity is not None:
1035
+ return self.fixed_capacity
1036
+ total_assignments = n_tokens * self.top_k
1037
+ avg_load = -(-total_assignments // self.n_routed)
1038
+ numer = int(round(self.capacity_factor * 1000))
1039
+ capacity = -(-(avg_load * numer) // 1000)
1040
+ return capacity + 1
1041
+
1042
+
1043
+ def build_attention_layers(
1044
+ num_layers: int,
1045
+ hidden_size: int,
1046
+ num_heads: int,
1047
+ dropout: float = 0.0,
1048
+ bias: bool = True,
1049
+ num_kv_heads: int | None = None,
1050
+ pattern: str = "dense",
1051
+ max_seq_len: int = 4096,
1052
+ sliding_window_size: int = 256,
1053
+ tp_size: int = 1,
1054
+ rope_dim: int | None = None,
1055
+ yarn_scale: float = 1.0,
1056
+ yarn_orig_max_pos: int | None = None,
1057
+ yarn_alpha: float = 1.0,
1058
+ yarn_beta: float = 32.0,
1059
+ ):
1060
+
1061
+ if num_kv_heads is None:
1062
+ num_kv_heads = num_heads
1063
+
1064
+ layers = []
1065
+ layer_types = []
1066
+
1067
+ if pattern == "dense":
1068
+ for _ in range(num_layers):
1069
+ layers.append(DenseMHA(
1070
+ hidden_size=hidden_size, num_heads=num_heads, num_kv_heads=num_kv_heads,
1071
+ dropout=dropout, bias=bias, max_seq_len=max_seq_len,
1072
+ window_size=None, tp_size=tp_size, rope_dim=rope_dim,
1073
+ yarn_scale=yarn_scale, yarn_orig_max_pos=yarn_orig_max_pos,
1074
+ yarn_alpha=yarn_alpha, yarn_beta=yarn_beta,
1075
+ ))
1076
+ layer_types.append("dense")
1077
+
1078
+ elif pattern == "pyramid_swa":
1079
+
1080
+ if num_kv_heads <= 1:
1081
+ pyramid_kv_heads = [1, 1, 1]
1082
+ else:
1083
+ import math as _math
1084
+
1085
+ mqa_floor = max(2, _math.ceil(num_kv_heads / 4))
1086
+ mqa_floor = min(mqa_floor, num_kv_heads) # never exceed full heads
1087
+ pyramid_kv_heads = sorted(set([
1088
+ mqa_floor,
1089
+ max(mqa_floor, _math.ceil(num_kv_heads / 2)),
1090
+ num_kv_heads,
1091
+ ]))
1092
+
1093
+ while len(pyramid_kv_heads) < 3:
1094
+ pyramid_kv_heads.append(num_kv_heads)
1095
+ pyramid_kv_heads = pyramid_kv_heads[:3]
1096
+
1097
+ for i in range(num_layers):
1098
+ pos_in_group = i % 4
1099
+ is_last_layer = (i == num_layers - 1)
1100
+
1101
+ if pos_in_group == 3 or is_last_layer:
1102
+ layers.append(DenseMHA(
1103
+ hidden_size=hidden_size, num_heads=num_heads, num_kv_heads=num_kv_heads,
1104
+ dropout=dropout, bias=bias, max_seq_len=max_seq_len,
1105
+ window_size=None, tp_size=tp_size, rope_dim=rope_dim,
1106
+ yarn_scale=yarn_scale, yarn_orig_max_pos=yarn_orig_max_pos,
1107
+ yarn_alpha=yarn_alpha, yarn_beta=yarn_beta,
1108
+ ))
1109
+ layer_types.append("global")
1110
+ else:
1111
+ kv_heads_here = pyramid_kv_heads[pos_in_group]
1112
+ layers.append(DenseMHA(
1113
+ hidden_size=hidden_size, num_heads=num_heads, num_kv_heads=kv_heads_here,
1114
+ dropout=dropout, bias=bias, max_seq_len=max_seq_len,
1115
+ window_size=sliding_window_size, tp_size=tp_size, rope_dim=rope_dim,
1116
+ yarn_scale=yarn_scale, yarn_orig_max_pos=yarn_orig_max_pos,
1117
+ yarn_alpha=yarn_alpha, yarn_beta=yarn_beta,
1118
+ ))
1119
+ layer_types.append("swa")
1120
+ elif pattern == "mod":
1121
+
1122
+ for i in range(num_layers):
1123
+ layers.append(DenseMHA(
1124
+ hidden_size=hidden_size, num_heads=num_heads, num_kv_heads=num_kv_heads,
1125
+ dropout=dropout, bias=bias, max_seq_len=max_seq_len,
1126
+ window_size=None, tp_size=tp_size, rope_dim=rope_dim,
1127
+ yarn_scale=yarn_scale, yarn_orig_max_pos=yarn_orig_max_pos,
1128
+ yarn_alpha=yarn_alpha, yarn_beta=yarn_beta,
1129
+ ))
1130
+ is_full_dense = (i % 4 == 3) or (i == num_layers - 1)
1131
+ layer_types.append("mod_dense" if is_full_dense else "mod")
1132
+ else:
1133
+ raise ValueError(
1134
+ f"Unknown attention pattern {pattern!r} -- supported patterns: "
1135
+ f"'dense', 'pyramid_swa', 'mod'. "
1136
+ f"(BigBird/conv-hybrid/LSA/HCA patterns have been removed.)"
1137
+ )
1138
+
1139
+ return nn.ModuleList(layers), layer_types
1140
+
1141
+
1142
+ class AryaBlock(nn.Module):
1143
+ """Single transformer block: Layer Norm → Attention → MLP (with residual connections).
1144
+
1145
+ `attn_module` is DenseMHA,
1146
+ pre-built by GPT via build_attention_layers. Every layer is fully
1147
+ standalone -- no cross-layer attention/KV sharing of any kind.
1148
+ `role` is purely a label ("dense" | "lsa") for inspection/debugging.
1149
+ """
1150
+
1151
+ def __init__(self, config, attn_module: nn.Module, role: str = "dense"):
1152
+ super().__init__()
1153
+ self.ln_1 = RMSNorm(config.n_embd, eps=1e-5)
1154
+ # Pre-LN for the MLP branch: ensure router and experts see the
1155
+ # identical normalized tensor the MLPs expect.
1156
+ self.ln_2 = RMSNorm(config.n_embd, eps=1e-5)
1157
+ self.attn = attn_module
1158
+ self.role = role # "dense" | "lsa"
1159
+
1160
+
1161
+ self.use_moe = getattr(config, 'use_moe', False)
1162
+ if self.use_moe:
1163
+ if getattr(config, 'tp_size', 1) > 1:
1164
+ raise ValueError(
1165
+ "tensor_parallel_size > 1 is not supported together with use_moe=True -- "
1166
+ "MoE needs its own (unimplemented here) expert-parallel sharding scheme, "
1167
+ "column/row-parallel head splitting does not apply to expert routing. "
1168
+ "Set use_moe=False or tensor_parallel_size=1."
1169
+ )
1170
+ self.ffn = AryaSparseMoE(
1171
+ d_model=config.n_embd,
1172
+ n_experts=config.n_experts,
1173
+ n_shared=config.n_shared,
1174
+ top_k=getattr(config, 'moe_top_k', 2),
1175
+ ffn_mult=config.ffn_mult,
1176
+ dropout=config.dropout,
1177
+ use_capacity_routing=True,
1178
+ fixed_capacity=getattr(config, 'moe_fixed_capacity', None),
1179
+ )
1180
+ else:
1181
+ self.ffn = SwiGLU_FFN(config.n_embd, config.ffn_mult, dropout=config.dropout, tp_size=getattr(config, 'tp_size', 1))
1182
+
1183
+ def forward(self, x):
1184
+ """Every attention layer (DenseMHA, LSA) is standalone and
1185
+ computes its own K/V from `x` alone --
1186
+ no external_kv, no cross-layer sharing, no produced_kv to pass on.
1187
+ """
1188
+ normed = self.ln_1(x)
1189
+
1190
+
1191
+ fused_out = None
1192
+ raw_ctx = None
1193
+ attn_delta = None
1194
+
1195
+ use_fused_outproj = (
1196
+ _use_flash_outproj_config()
1197
+ and flash_ops is not None
1198
+ and self.attn.out_proj.bias is None
1199
+ )
1200
+ if use_fused_outproj:
1201
+ attn_raw = self.attn(normed, skip_out_proj=True)
1202
+ candidate = module_output_tensor(attn_raw)
1203
+ if candidate.is_cuda and candidate.shape[-1] == self.attn.out_proj.weight.shape[1]:
1204
+ raw_ctx = candidate
1205
+
1206
+ gemm_dtype = raw_ctx.dtype
1207
+ fused_out = flash_ops.fused_outproj_add_rmsnorm(
1208
+ raw_ctx.reshape(-1, raw_ctx.shape[-1]),
1209
+ self.attn.out_proj.weight.to(dtype=gemm_dtype),
1210
+ x.reshape(-1, x.shape[-1]).to(dtype=gemm_dtype),
1211
+ self.ln_2.weight.to(dtype=gemm_dtype),
1212
+ self.ln_2.eps,
1213
+ )
1214
+ if fused_out is not None:
1215
+ residual_out, normed_out = fused_out
1216
+
1217
+ x = residual_out.to(dtype=x.dtype).reshape(x.shape)
1218
+ ffn_in = normed_out.reshape(x.shape)
1219
+ else:
1220
+
1221
+ skip_out_proj_honored = self.attn.tp_size == 1
1222
+ if skip_out_proj_honored:
1223
+ raw_ctx = candidate
1224
+ attn_delta = self.attn.out_proj(raw_ctx)
1225
+ else:
1226
+ attn_delta = candidate
1227
+
1228
+ if fused_out is None:
1229
+ if attn_delta is None:
1230
+ if raw_ctx is not None:
1231
+
1232
+ attn_delta = self.attn.out_proj(raw_ctx)
1233
+ else:
1234
+
1235
+ attn_raw = self.attn(normed)
1236
+ attn_delta = module_output_tensor(attn_raw)
1237
+
1238
+
1239
+
1240
+ tier2_out = None
1241
+ if _use_fused_add_rmsnorm_config() and fused_kernels is not None:
1242
+ tier2_out = fused_kernels.fused_add_rmsnorm(
1243
+ x, attn_delta, self.ln_2.weight.to(dtype=x.dtype), self.ln_2.eps
1244
+ )
1245
+ if tier2_out is not None:
1246
+ x, ffn_in = tier2_out
1247
+ else:
1248
+ # Tier 3: fully eager, byte-for-byte what this block
1249
+ # always ran before any fusion work this session.
1250
+ x = x + attn_delta
1251
+ ffn_in = self.ln_2(x)
1252
+
1253
+
1254
+ ffn_delta = self.ffn(ffn_in)
1255
+ x = x + ffn_delta
1256
+ return x
1257
+
1258
+
1259
+ class MoDBlock(nn.Module):
1260
+ """Mixture-of-Depths via a continuous sigmoid gate -- no topk, no gather, no scatter.
1261
+
1262
+ WILL THE GATE ACTUALLY SELECT TOKENS?
1263
+ ───���──────────────────────────────────
1264
+ Short answer: not without explicit sparsity pressure.
1265
+
1266
+ A plain sigmoid gate starts near 0.5 (for zero-mean inputs after RMSNorm)
1267
+ and stays fractional -- it applies 40% of the delta to some tokens, 70%
1268
+ to others, 30% to others. The gradient never forces it toward 0 or 1
1269
+ unless there is a cost to using non-binary values. Two mechanisms push
1270
+ the gate toward genuine token selection:
1271
+
1272
+ 1. mod_gate_entropy_coeff (config.py, default 0.01):
1273
+ Adds a per-step auxiliary loss term:
1274
+ L_entropy = -coeff * mean( gate*log(gate) + (1-gate)*log(1-gate) )
1275
+ This is the binary entropy H(gate), negated -- maximizing binary entropy
1276
+ means pushing gate toward 0 or 1, where H=0 (the gate is certain).
1277
+ A gate at 0.5 has H=1.0 (maximum uncertainty); a gate at 0.05 or 0.95
1278
+ has H≈0.29. The coefficient controls how hard this regularizer pushes:
1279
+ 0.0 -> pure sigmoid, stays fractional, no selection
1280
+ 0.01 -> gentle push; gates typically reach 0.1/0.9 range by step 2000
1281
+ 0.1 -> strong push; gates near 0/1 within ~500 steps but can destabilize
1282
+ This term is computed in forward() and returned alongside the block output.
1283
+ train.py adds it to the main loss, multiplied by the coefficient.
1284
+
1285
+ 2. Natural gradient pressure from the main loss:
1286
+ If one token position's delta reliably INCREASES the loss regardless of
1287
+ context, the model learns to set gate≈0 there (the identity is better).
1288
+ If another position's delta reliably DECREASES the loss, gate≈1 emerges.
1289
+ This happens without any regularization, just more slowly -- typically
1290
+ ~5000-10000 steps before meaningful gate bimodality appears.
1291
+
1292
+ The entropy loss (mechanism 1) dramatically speeds up mechanism 2 by making
1293
+ the gate pay a cost for ambiguity. Use mod_gate_entropy_coeff=0.01 as a
1294
+ starting point and inspect a histogram of gate values at step 1000 and 5000.
1295
+ If the distribution is bimodal (peaks near 0 and 1), the gate is selecting.
1296
+ If it's unimodal near 0.5, increase the coefficient.
1297
+
1298
+
1299
+
1300
+ GATE DESIGN
1301
+ ───────────
1302
+ gate = sigmoid(router(x)) # (B, T, 1) -- D->1 matmul + sigmoid
1303
+ delta = inner(x) - x # (B, T, D) -- net block update
1304
+ out = x + gate * delta # elementwise, fixed shape, fused
1305
+ aux = entropy_coeff * H_binary(gate) # scalar, added to training loss
1306
+
1307
+ Parameter cost: D floats (router.weight). At D=256 that's 256 params total.
1308
+ """
1309
+
1310
+ def __init__(self, config, inner: nn.Module, capacity: float = 0.125):
1311
+ super().__init__()
1312
+ self.inner = inner
1313
+ self.capacity = capacity
1314
+ self._entropy_coeff = float(getattr(config, 'mod_gate_entropy_coeff', 0.01))
1315
+
1316
+
1317
+ self.router = nn.Linear(config.n_embd, 1, bias=False)
1318
+
1319
+ def forward(self, x: torch.Tensor) -> tuple[torch.Tensor, torch.Tensor]:
1320
+
1321
+
1322
+ gate = torch.sigmoid(self.router(x)) # (B, T, 1)
1323
+
1324
+
1325
+ inner_out = self.inner(x) # (B,T,D)
1326
+ out = x + gate * (inner_out - x) # (B,T,D); fuses delta into out
1327
+
1328
+
1329
+ if self.training and self._entropy_coeff > 0.0:
1330
+ g = gate.float().clamp(1e-6, 1.0 - 1e-6) # fp32 for log stability
1331
+
1332
+ ent = -(g * torch.log(g) + (1.0 - g) * torch.log(1.0 - g)) # H per element
1333
+ gate_aux = self._entropy_coeff * ent.mean()
1334
+ else:
1335
+ gate_aux = torch.zeros((), device=x.device, dtype=torch.float32)
1336
+
1337
+ return out, gate_aux
1338
+
1339
+ def extra_repr(self) -> str:
1340
+ return f"gate_based=True,entropy_coeff={self._entropy_coeff}"
1341
+
1342
+
1343
+ class MultiTokenPrediction(nn.Module):
1344
+ """Multi-Token Prediction auxiliary heads (Gloeckle et al. / DeepSeek-V3
1345
+ style): predict tokens at offsets 2, 3, ... beyond the main next-token
1346
+ prediction, using the same pre-lm_head hidden state.
1347
+
1348
+ Each head k predicts token t+k+1 from hidden state h_t via a separate
1349
+ RMSNorm (so each depth can learn its own scale) followed by the shared
1350
+ lm_head weight (weight tying: same projection used for the main head),
1351
+ keeping parameter count increase minimal.
1352
+
1353
+ Loss contribution: mtp_lambda * mean(CE_1, CE_2, ..., CE_depth).
1354
+ The main model's loss is unaffected; MTP loss is added on top.
1355
+
1356
+
1357
+ """
1358
+
1359
+ def __init__(self, config, depth: int = 1):
1360
+ super().__init__()
1361
+ self.depth = depth
1362
+ # One RMSNorm per MTP depth; no separate linear (share lm_head.weight).
1363
+ self.norms = nn.ModuleList([
1364
+ RMSNorm(config.n_embd, eps=1e-5) for _ in range(depth)
1365
+ ])
1366
+
1367
+ def compute_loss(
1368
+ self,
1369
+ x: torch.Tensor, # (B, T, D) -- pre-lm_head hidden states
1370
+ targets: torch.Tensor, # (B, T) -- already the +1 shifted targets
1371
+ lm_head_weight: torch.Tensor, # (vocab_size, D) -- tied weight
1372
+ mtp_lambda: float = 0.3,
1373
+ ) -> torch.Tensor:
1374
+
1375
+
1376
+ total = torch.zeros((), device=x.device, dtype=torch.float32)
1377
+ valid_depths = 0
1378
+ for k, norm in enumerate(self.norms):
1379
+ offset = k + 1 # additional token offset beyond the main +1 shift
1380
+ Tk = x.size(1) - offset
1381
+ if Tk <= 0:
1382
+ continue
1383
+
1384
+ h = norm(x[:, :Tk]) # (B, Tk, D)
1385
+
1386
+ logits_k = F.linear(h, lm_head_weight) # (B, Tk, vocab_size)
1387
+ tgt_k = targets[:, offset:] # (B, Tk)
1388
+ loss_k = F.cross_entropy(
1389
+ logits_k.reshape(-1, logits_k.size(-1)).float(),
1390
+ tgt_k.reshape(-1).long(),
1391
+ ignore_index=-1,
1392
+ )
1393
+ total = total + loss_k # fp32 + fp32 -- no precision loss
1394
+ valid_depths += 1
1395
+
1396
+ if valid_depths == 0:
1397
+
1398
+ return torch.zeros((), device=x.device, dtype=torch.float32)
1399
+ return mtp_lambda * (total / valid_depths) # fp32 scalar
1400
+
1401
+
1402
+ @dataclass
1403
+ class GPTConfig:
1404
+ """Configuration for GPT model"""
1405
+ block_size: int = _alpha_config_attr('CONTEXT', 1024)
1406
+ vocab_size: int = _alpha_config_attr('vocab_size', 50304)
1407
+ n_layer: int = _alpha_config_attr('numberoflayers', 12)
1408
+ n_head: int = _alpha_config_attr('numberofheads', 12)
1409
+ n_embd: int = _alpha_config_attr('D_MODEL', 768)
1410
+ dropout: float = _alpha_config_attr('dropout', 0.0)
1411
+
1412
+ bias: bool = _alpha_config_attr('bias', False)
1413
+ ffn_mult: float = _alpha_config_attr('ffn_mult', 4.0)
1414
+
1415
+ d_rope: int | None = _alpha_config_attr('d_rope', None)
1416
+
1417
+ top_k: int = _alpha_config_attr('top_k', 64)
1418
+ n_experts: int = _alpha_config_attr('n_experts', 8)
1419
+ n_shared: int = _alpha_config_attr('n_shared', 1)
1420
+ use_moe: bool = _alpha_config_attr('use_moe', False)
1421
+
1422
+ use_liger: bool = _alpha_config_attr('use_liger', False)
1423
+
1424
+ moe_top_k: int = _alpha_config_attr('moe_top_k', 2)
1425
+
1426
+ gradient_checkpointing: bool = _alpha_config_attr('gradient_checkpointing', False) # trade ~20-30% more compute for substantially less activation memory (grows with n_layer -- see GPT.forward)
1427
+ num_kv_heads: int | None = _alpha_config_attr('num_kv_heads', None)
1428
+
1429
+ pattern: str = _alpha_config_attr('pattern', 'pyramid_swa')
1430
+
1431
+ sliding_window_size: int = _alpha_config_attr('sliding_window_size', 256)
1432
+
1433
+ tp_size: int = _alpha_config_attr('tensor_parallel_size', 1)
1434
+
1435
+
1436
+ gen_headroom: int = _alpha_config_attr('gen_headroom', 0)
1437
+
1438
+ max_gen_tokens: int | None = _alpha_config_attr('max_gen_tokens', None)
1439
+
1440
+
1441
+ yarn_scale: float = _alpha_config_attr('yarn_scale', 1.0)
1442
+ yarn_orig_max_pos: int | None = _alpha_config_attr('yarn_orig_max_pos', None)
1443
+ yarn_alpha: float = _alpha_config_attr('yarn_alpha', 1.0)
1444
+ yarn_beta: float = _alpha_config_attr('yarn_beta', 32.0)
1445
+
1446
+
1447
+ mod_capacity: float = _alpha_config_attr('mod_capacity', 0.125)
1448
+
1449
+
1450
+ use_mtp: bool = _alpha_config_attr('use_mtp', False)
1451
+ mtp_depth: int = _alpha_config_attr('mtp_depth', 1)
1452
+ mtp_lambda: float = _alpha_config_attr('mtp_lambda', 0.3)
1453
+
1454
+
1455
+ use_nvfp4: bool = _alpha_config_attr('use_nvfp4', False)
1456
+
1457
+
1458
+ use_int8: bool = _alpha_config_attr('use_int8', False)
1459
+ int8_group_size: int = _alpha_config_attr('int8_group_size', 128)
1460
+
1461
+ mod_gate_entropy_coeff: float = _alpha_config_attr('mod_gate_entropy_coeff', 0.01)
1462
+
1463
+ class GPT(nn.Module):
1464
+ """Full GPT language model"""
1465
+
1466
+ def __init__(self, config):
1467
+ super().__init__()
1468
+ if config.vocab_size is None:
1469
+ raise ValueError("config.vocab_size must be set")
1470
+ if config.block_size is None:
1471
+ raise ValueError("config.block_size must be set")
1472
+ self.config = config
1473
+
1474
+
1475
+
1476
+ head_dim = config.n_embd // config.n_head
1477
+ rope_dim = head_dim if getattr(config, 'd_rope', None) is None else config.d_rope
1478
+
1479
+ gen_headroom = max(0, int(getattr(config, 'gen_headroom', 0)))
1480
+ rope_table_len = config.block_size + gen_headroom
1481
+
1482
+
1483
+ yarn_scale = float(getattr(config, 'yarn_scale', 1.0))
1484
+ yarn_orig_max_pos_cfg = getattr(config, 'yarn_orig_max_pos', None)
1485
+ yarn_orig_max_pos = int(yarn_orig_max_pos_cfg) if yarn_orig_max_pos_cfg is not None else config.block_size
1486
+ yarn_alpha = float(getattr(config, 'yarn_alpha', 1.0))
1487
+ yarn_beta = float(getattr(config, 'yarn_beta', 32.0))
1488
+ effective_max_seq_len = int(rope_table_len * yarn_scale) if yarn_scale > 1.0 else rope_table_len
1489
+
1490
+ attn_layers, self.layer_types = build_attention_layers(
1491
+ num_layers=config.n_layer,
1492
+ hidden_size=config.n_embd,
1493
+ num_heads=config.n_head,
1494
+ dropout=config.dropout,
1495
+ bias=config.bias,
1496
+ num_kv_heads=config.num_kv_heads,
1497
+ pattern=config.pattern,
1498
+ max_seq_len=effective_max_seq_len,
1499
+ sliding_window_size=config.sliding_window_size,
1500
+ tp_size=config.tp_size,
1501
+ rope_dim=rope_dim,
1502
+ yarn_scale=yarn_scale,
1503
+ yarn_orig_max_pos=yarn_orig_max_pos,
1504
+ yarn_alpha=yarn_alpha,
1505
+ yarn_beta=yarn_beta,
1506
+ )
1507
+
1508
+
1509
+ mod_capacity = float(getattr(config, 'mod_capacity', 0.125))
1510
+ blocks = []
1511
+ for i in range(config.n_layer):
1512
+ block = AryaBlock(config, attn_layers[i], role=self.layer_types[i])
1513
+ if self.layer_types[i] == "mod":
1514
+ # Wrap in MoDBlock: 3/4 of layers route only top-capacity
1515
+ # fraction of tokens through the full AryaBlock.
1516
+ block = MoDBlock(config, block, capacity=mod_capacity)
1517
+ blocks.append(block)
1518
+
1519
+ self.transformer = nn.ModuleDict(dict(
1520
+ wte=nn.Embedding(config.vocab_size, config.n_embd),
1521
+ drop=nn.Dropout(config.dropout),
1522
+ h=nn.ModuleList(blocks),
1523
+ ln_f=RMSNorm(config.n_embd, eps=1e-5),
1524
+ ))
1525
+ self.lm_head = nn.Linear(config.n_embd, config.vocab_size, bias=False)
1526
+
1527
+
1528
+ use_mtp = bool(getattr(config, 'use_mtp', False))
1529
+ mtp_depth = int(getattr(config, 'mtp_depth', 1))
1530
+ self.mtp = MultiTokenPrediction(config, depth=mtp_depth) if use_mtp else None
1531
+
1532
+ self._mtp_lambda = float(getattr(config, 'mtp_lambda', 0.3))
1533
+
1534
+
1535
+ self.apply(self._init_weights)
1536
+
1537
+
1538
+ self.transformer.wte.weight = self.lm_head.weight
1539
+
1540
+
1541
+ residual_std = 0.02 / math.sqrt(config.n_layer)
1542
+ for name, p in self.named_parameters():
1543
+ if name.endswith("out_proj.weight") or name.endswith("W_out.weight"):
1544
+ torch.nn.init.normal_(p, mean=0.0, std=residual_std)
1545
+
1546
+ # ── Quantization (mutually exclusive; nvFP4 takes priority on SM100+) ──
1547
+ # Both quantize_() calls must happen BEFORE torch.compile() in train.py.
1548
+ # _int8_enabled() already checks _nvfp4_enabled() and refuses if both are set.
1549
+ self._use_nvfp4: bool = False
1550
+ self._use_int8: bool = False
1551
+
1552
+ if _nvfp4_enabled():
1553
+ n_quantized = apply_nvfp4_to_model(self)
1554
+ self._use_nvfp4 = True
1555
+ warnings.warn(
1556
+ f"[nvFP4/RHT] torchao quantized {n_quantized} Linear layers "
1557
+ f"(Blackwell SM100+). Weight: FP4 e2m1 per-block-16 RHT. "
1558
+ f"Activation: FP4 dynamic per-tensor scale (FP8 intermediate). "
1559
+ f"Fused GEMM kernel emitted by torch.compile. "
1560
+ f"Skipped: wte, lm_head, norms, router, gate, small/misaligned linears. "
1561
+ f"Verify val_loss vs bf16 baseline before committing to a long run.",
1562
+ stacklevel=2,
1563
+ )
1564
+ elif _int8_enabled():
1565
+ gs = int(getattr(config, 'int8_group_size', 128))
1566
+ n_quantized = apply_int8_to_model(self, group_size=gs)
1567
+ self._use_int8 = True
1568
+ warnings.warn(
1569
+ f"[INT8/Jetfire] torchao quantized {n_quantized} Linear layers "
1570
+ f"(T4/Turing+). Weight: INT8 per-group-{gs} symmetric. "
1571
+ f"Activation: INT8 dynamic per-tensor scale (runtime, no calibration). "
1572
+ f"Fused kernel emitted by torch.compile via Inductor. "
1573
+ f"Skipped: wte, lm_head, norms, router, gate, small linears. "
1574
+ f"Verify val_loss vs fp16 baseline before a long run.",
1575
+ stacklevel=2,
1576
+ )
1577
+
1578
+ # Weight tying can be broken by quantized module replacement.
1579
+ # Reassert the canonical tie after either quantization pass.
1580
+ self.transformer.wte.weight = self.lm_head.weight
1581
+ assert self.transformer.wte.weight is self.lm_head.weight
1582
+
1583
+ def set_tp_group(self, group) -> None:
1584
+
1585
+ if self.config.tp_size <= 1:
1586
+ return
1587
+ for block in self.transformer.h:
1588
+
1589
+ real_block = block.inner if isinstance(block, MoDBlock) else block
1590
+ real_block.attn.tp_group = group
1591
+ if hasattr(real_block, "ffn") and real_block.ffn is not None:
1592
+ real_block.ffn.tp_group = group
1593
+
1594
+
1595
+ if torch.distributed.is_initialized() and group is not None:
1596
+ for block in self.transformer.h:
1597
+ real_block = block.inner if isinstance(block, MoDBlock) else block
1598
+ attn = real_block.attn
1599
+ if not isinstance(attn, DenseMHA):
1600
+ continue
1601
+ if attn.num_kv_heads >= attn.tp_size:
1602
+ continue
1603
+ with torch.no_grad():
1604
+ # Full qkv projection tensor is the canonical object to align.
1605
+ torch.distributed.broadcast(attn.qkv_proj.weight.data, src=0, group=group)
1606
+ if attn.qkv_proj.bias is not None:
1607
+ torch.distributed.broadcast(attn.qkv_proj.bias.data, src=0, group=group)
1608
+
1609
+ def _init_weights(self, module):
1610
+ if isinstance(module, nn.Linear):
1611
+ torch.nn.init.normal_(module.weight, mean=0.0, std=0.02)
1612
+ if module.bias is not None:
1613
+ torch.nn.init.zeros_(module.bias)
1614
+ elif isinstance(module, nn.Embedding):
1615
+ torch.nn.init.normal_(module.weight, mean=0.0, std=0.02)
1616
+
1617
+ def get_num_params(self, non_embedding=True):
1618
+ """Count number of parameters"""
1619
+ n_params = sum(p.numel() for p in self.parameters())
1620
+ if non_embedding:
1621
+ n_params -= self.transformer.wte.weight.numel()
1622
+ return n_params
1623
+
1624
+ def forward(self, idx, targets=None, need_logits=True):
1625
+
1626
+ idx = idx.long()
1627
+ if targets is not None:
1628
+ targets = targets.long()
1629
+ b, t = idx.size()
1630
+ torch._check(
1631
+ t <= self.config.block_size,
1632
+
1633
+ lambda: f"Sequence too long: {t} > {self.config.block_size}",
1634
+ )
1635
+
1636
+ tok_emb = self.transformer.wte(idx) # (B, T, n_embd)
1637
+ x = self.transformer.drop(tok_emb)
1638
+
1639
+
1640
+ mod_gate_loss = torch.zeros((), device=x.device, dtype=torch.float32)
1641
+
1642
+ for block in self.transformer.h:
1643
+ is_mod = isinstance(block, MoDBlock)
1644
+ if self.config.gradient_checkpointing and self.training:
1645
+
1646
+ out = torch.utils.checkpoint.checkpoint(block, x, use_reentrant=False)
1647
+ if is_mod:
1648
+ x, gate_aux = out
1649
+ mod_gate_loss = mod_gate_loss + gate_aux.float()
1650
+ else:
1651
+ x = out
1652
+ else:
1653
+ if is_mod:
1654
+ x, gate_aux = block(x)
1655
+ mod_gate_loss = mod_gate_loss + gate_aux.float()
1656
+ else:
1657
+ x = block(x)
1658
+
1659
+ x = self.transformer.ln_f(x)
1660
+
1661
+ if targets is not None:
1662
+ mtp_loss = torch.zeros((), device=x.device, dtype=torch.float32)
1663
+ if self.mtp is not None:
1664
+ mtp_loss = self.mtp.compute_loss(
1665
+ x, targets, self.lm_head.weight, mtp_lambda=self._mtp_lambda
1666
+ )
1667
+
1668
+ # Keep the primary language-model CE loss separate from auxiliary MTP/MoD penalties.
1669
+ # This lets caller-side logging and validation track the true CE signal without
1670
+ # the auxiliary terms masking the baseline perplexity curve.
1671
+ if not need_logits and _use_liger_config() and _LigerFusedLinearCrossEntropyLoss is not None and _liger_cuda_gate(x):
1672
+ liger_ce = _liger_fused_ce_attempt(x, self.lm_head.weight, targets)
1673
+ if liger_ce is not None:
1674
+ total_loss = liger_ce + mtp_loss + mod_gate_loss
1675
+ self.last_ce_loss = liger_ce.detach()
1676
+ self.last_total_loss = total_loss.detach()
1677
+ return None, total_loss
1678
+
1679
+ logits = F.linear(x, self.lm_head.weight)
1680
+ ce_loss = F.cross_entropy(
1681
+ logits.reshape(-1, logits.size(-1)).float(),
1682
+ targets.reshape(-1),
1683
+ ignore_index=-1
1684
+ )
1685
+
1686
+ # Keep the loss scalar finite by avoiding a patched post-hoc
1687
+ # nan_to_num that stamps out the gradient entirely. Repair the
1688
+ # overflow sources (logit clamp/initialization) out front instead.
1689
+ total_loss = ce_loss + mtp_loss + mod_gate_loss
1690
+ self.last_ce_loss = ce_loss.detach()
1691
+ self.last_total_loss = total_loss.detach()
1692
+
1693
+ if not need_logits:
1694
+ logits = None
1695
+ else:
1696
+ logits = F.linear(x[:, [-1], :], self.lm_head.weight) # (B, 1, vocab_size)
1697
+ ce_loss = None
1698
+ total_loss = None
1699
+
1700
+ return logits, total_loss
1701
+
1702
+ @torch.no_grad()
1703
+ def generate(self, idx, max_new_tokens=None, temperature=1.0, top_k=None):
1704
+
1705
+ if max_new_tokens is None:
1706
+ max_new_tokens = getattr(self.config, 'max_gen_tokens', None)
1707
+ if max_new_tokens is None:
1708
+ raise ValueError(
1709
+ "generate() needs max_new_tokens, either passed "
1710
+ "directly or set as max_gen_tokens in config.py."
1711
+ )
1712
+ for _ in range(max_new_tokens):
1713
+ # Crop to block size if needed
1714
+ idx_cond = idx if idx.size(1) <= self.config.block_size else idx[:, -self.config.block_size:]
1715
+
1716
+ # Forward pass
1717
+ logits, _ = self(idx_cond) # (B, T, vocab_size)
1718
+
1719
+ # Get logits for next token (only last position)
1720
+ logits = logits[:, -1, :] / temperature # (B, vocab_size)
1721
+
1722
+ # Top-k filtering
1723
+ if top_k is not None:
1724
+ v, _ = torch.topk(logits, min(top_k, logits.size(-1)))
1725
+ logits[logits < v[:, [-1]]] = -float('Inf')
1726
+
1727
+
1728
+ probs = F.softmax(logits.float(), dim=-1) # (B, vocab_size)
1729
+
1730
+ # Sample next token
1731
+ idx_next = torch.multinomial(probs, num_samples=1) # (B, 1)
1732
+
1733
+ # Append to sequence
1734
+ idx = torch.cat((idx, idx_next), dim=1) # (B, T+1)
1735
+
1736
+ return idx
1737
+
1738
+
1739
+ def audit_fused_kernels(model, device, dtype=torch.float16, verbose=True):
1740
+
1741
+ config = model.config
1742
+ head_dim = config.n_embd // config.n_head
1743
+ rope_dim = head_dim if getattr(config, 'd_rope', None) is None else config.d_rope
1744
+ b, t = 2, 8
1745
+ results = {}
1746
+
1747
+ if flash_ops is None:
1748
+ if verbose:
1749
+ print(
1750
+ "[FusionAudit] flash_ops extension failed to load entirely -- "
1751
+ "all three flash_* fusions (rope/outproj_add_rmsnorm/swiglu) "
1752
+ "are running eager fallback for the WHOLE run, not just under "
1753
+ "some condition. Check flash.cu compile output / nvcc "
1754
+ "availability above for the real cause."
1755
+ )
1756
+ return {"flash_ops_loaded": False}
1757
+ results["flash_ops_loaded"] = True
1758
+
1759
+ # -- fused RoPE --
1760
+ try:
1761
+ q = torch.randn(b, t, config.n_head, head_dim, device=device, dtype=dtype)
1762
+ k = torch.randn(b, t, config.num_kv_heads or config.n_head, head_dim, device=device, dtype=dtype)
1763
+ cos, sin = build_rope_cache(config.block_size, rope_dim)
1764
+ cos, sin = cos.to(device), sin.to(device)
1765
+ out = flash_ops.fused_rope_qk(q, k, cos[:t].to(dtype), sin[:t].to(dtype), rope_dim)
1766
+ results["fused_rope"] = out is not None
1767
+ except Exception as e: # noqa: BLE001 -- diagnostic only, must not crash
1768
+ results["fused_rope"] = False
1769
+ results["fused_rope_error"] = repr(e)
1770
+
1771
+ # -- fused out_proj + residual-add + RMSNorm --
1772
+ try:
1773
+ x = torch.randn(b * t, config.n_embd, device=device, dtype=dtype)
1774
+ w = torch.randn(config.n_embd, config.n_embd, device=device, dtype=dtype)
1775
+ residual = torch.randn(b * t, config.n_embd, device=device, dtype=dtype)
1776
+ norm_w = torch.ones(config.n_embd, device=device, dtype=dtype)
1777
+ out = flash_ops.fused_outproj_add_rmsnorm(x, w, residual, norm_w, 1e-5)
1778
+ results["fused_outproj_add_rmsnorm"] = out is not None
1779
+ except Exception as e: # noqa: BLE001
1780
+ results["fused_outproj_add_rmsnorm"] = False
1781
+ results["fused_outproj_add_rmsnorm_error"] = repr(e)
1782
+
1783
+ # -- fused SwiGLU --
1784
+ try:
1785
+ hidden = int(config.n_embd * config.ffn_mult)
1786
+ gate = torch.randn(b * t, hidden, device=device, dtype=dtype)
1787
+ value = torch.randn(b * t, hidden, device=device, dtype=dtype)
1788
+ out = flash_ops.fused_swiglu(gate, value)
1789
+ results["fused_swiglu"] = out is not None
1790
+ except Exception as e: # noqa: BLE001
1791
+ results["fused_swiglu"] = False
1792
+ results["fused_swiglu_error"] = repr(e)
1793
+
1794
+ if verbose:
1795
+ print("[FusionAudit] one-time probe of fused kernels against real on-device tensors:")
1796
+ for name in ("fused_rope", "fused_outproj_add_rmsnorm", "fused_swiglu"):
1797
+ ok = results.get(name, False)
1798
+ status = "ENGAGED" if ok else "FELL BACK TO EAGER"
1799
+ print(f" {name:28s} -> {status}")
1800
+ if not ok and f"{name}_error" in results:
1801
+ print(f" reason: {results[f'{name}_error']}")
1802
+ n_ok = sum(results.get(n, False) for n in ("fused_rope", "fused_outproj_add_rmsnorm", "fused_swiglu"))
1803
+ if n_ok < 3:
1804
+ print(
1805
+ f" [FusionAudit] {3 - n_ok}/3 fusions NOT engaging -- given "
1806
+ f"this model is memory-bandwidth-bound (per the roofline "
1807
+ f"analysis), a missed fusion means real intermediate-tensor "
1808
+ f"HBM traffic that shouldn't be there. Worth fixing before "
1809
+ f"chasing anything else."
1810
+ )
1811
+ else:
1812
+ print(" [FusionAudit] all 3/3 fusions engaged -- fusion is not the bottleneck here.")
1813
+ return results
tokenizer.py ADDED
@@ -0,0 +1,389 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ from __future__ import annotations
2
+
3
+ import json
4
+ import logging
5
+ from pathlib import Path
6
+ from typing import Sequence, List, Optional
7
+
8
+ import sentencepiece as spm
9
+
10
+ logger = logging.getLogger(__name__)
11
+
12
+ # 32K is the well-established baseline vocab size for BPE/SentencePiece
13
+ # LLM tokenizers (Llama-1/2, T5, Gopher, Chinchilla all use exactly this).
14
+ # 128K+ only pays off for heavy multilingual/code coverage; for a small,
15
+ # largely-English, narrow-domain model, 32K is the standard, safe default.
16
+ DEFAULT_VOCAB_SIZE = 32000
17
+
18
+
19
+ def train_sentencepiece(
20
+ data_files: Sequence[str],
21
+ model_prefix: str = 'tokenizer',
22
+ vocab_size: int = DEFAULT_VOCAB_SIZE,
23
+ model_type: str = 'bpe',
24
+ character_coverage: float = 0.9995,
25
+ byte_fallback: bool = True,
26
+ pad_id: int = 1,
27
+ unk_id: int = 0,
28
+ bos_id: int = 2,
29
+ eos_id: int = 3,
30
+ add_dummy_prefix: bool = True,
31
+ num_threads: int = 8,
32
+ input_sentence_size: int = 5_000_000,
33
+ shuffle_input_sentence: bool = True,
34
+ max_sentence_length: int = 16384,
35
+ split_digits: bool = True,
36
+ allow_whitespace_only_pieces: bool = True,
37
+ train_extremely_large_corpus: bool = False,
38
+ ) -> str:
39
+ """
40
+ Train a SentencePiece BPE tokenizer with byte-fallback — the same
41
+ scheme used by Llama-2, Mistral, and EuroLLM (BPE + byte_fallback via
42
+ SentencePiece specifically, not a hand-rolled BPE implementation).
43
+
44
+ Why SentencePiece and not a hand-written tiktoken export: SentencePiece's
45
+ C++ core does encode/decode and merge-rank bookkeeping internally and
46
+ natively — there is no manual ID-renumbering or rank-export step for
47
+ calling code to get wrong. (A prior tiktoken-based rewrite of this
48
+ tokenizer had exactly that class of bug: hand-exported merge ranks were
49
+ non-contiguous because special tokens occupied ids 0-3 in the source
50
+ vocab, silently corrupting merge-priority order and decode() mappings —
51
+ manifesting as repetitive garbage output like "to to to" despite a
52
+ healthy training loss. Delegating to SentencePiece's own encode/decode
53
+ removes that entire class of bug by construction.)
54
+
55
+ Notes on defaults:
56
+ - character_coverage < 1.0 with byte_fallback=True: rare glyphs fall
57
+ back to byte pieces instead of bloating the vocab with singletons.
58
+ - input_sentence_size + shuffle_input_sentence: without shuffling,
59
+ SentencePiece samples from the START of the concatenated corpus,
60
+ which silently biases vocab toward whichever domain file comes
61
+ first if you hand it multiple files back to back.
62
+ - split_digits: keeps numbers as individual digit tokens, which
63
+ generally helps arithmetic/math task tokenization consistency.
64
+ """
65
+ data_files = [str(Path(p)) for p in data_files]
66
+ if not data_files:
67
+ raise ValueError('data_files is empty')
68
+
69
+ missing = [f for f in data_files if not Path(f).exists()]
70
+ if missing:
71
+ raise FileNotFoundError(f'Missing input files: {missing}')
72
+
73
+ kwargs = dict(
74
+ input=','.join(data_files),
75
+ model_prefix=model_prefix,
76
+ vocab_size=int(vocab_size),
77
+ model_type=model_type,
78
+ character_coverage=character_coverage,
79
+ pad_id=pad_id,
80
+ unk_id=unk_id,
81
+ bos_id=bos_id,
82
+ eos_id=eos_id,
83
+ byte_fallback=byte_fallback,
84
+ hard_vocab_limit=False,
85
+ normalization_rule_name='nmt_nfkc',
86
+ add_dummy_prefix=add_dummy_prefix,
87
+ num_threads=num_threads,
88
+ input_sentence_size=input_sentence_size,
89
+ shuffle_input_sentence=shuffle_input_sentence,
90
+ max_sentence_length=max_sentence_length,
91
+ split_digits=split_digits,
92
+ allow_whitespace_only_pieces=allow_whitespace_only_pieces,
93
+ train_extremely_large_corpus=train_extremely_large_corpus,
94
+ )
95
+
96
+ logger.info(f"Training SentencePiece: vocab_size={vocab_size} model_type={model_type} "
97
+ f"files={len(data_files)}")
98
+ spm.SentencePieceTrainer.train(**kwargs)
99
+
100
+ model_path = f'{model_prefix}.model'
101
+ _validate_trained_model(
102
+ model_path, vocab_size,
103
+ expected_pad=pad_id, expected_unk=unk_id, expected_bos=bos_id, expected_eos=eos_id,
104
+ )
105
+ return model_path
106
+
107
+
108
+ def _validate_trained_model(
109
+ model_path: str,
110
+ expected_vocab_size: int,
111
+ expected_pad: int,
112
+ expected_unk: int,
113
+ expected_bos: int,
114
+ expected_eos: int,
115
+ ) -> None:
116
+ """
117
+ Self-critique validation pass — checks the things that actually broke
118
+ in the previous (tiktoken) tokenizer, not just "does it load".
119
+ """
120
+ sp = spm.SentencePieceProcessor(model_file=model_path)
121
+
122
+ # 1. Vocab size sanity
123
+ actual_vocab = sp.vocab_size()
124
+ if actual_vocab != expected_vocab_size:
125
+ logger.warning(f"Trained vocab_size={actual_vocab} differs from requested={expected_vocab_size} "
126
+ f"(hard_vocab_limit=False allows this if the corpus is small)")
127
+
128
+ # 2. Special token IDs must be EXACTLY what was requested — not just
129
+ # ">= 0". A previous bug class involved special-token ids silently
130
+ # drifting from what calling code assumed. Check explicitly, not
131
+ # loosely.
132
+ checks = [
133
+ ('pad', sp.pad_id(), expected_pad),
134
+ ('unk', sp.unk_id(), expected_unk),
135
+ ('bos', sp.bos_id(), expected_bos),
136
+ ('eos', sp.eos_id(), expected_eos),
137
+ ]
138
+ for name, actual, expected in checks:
139
+ if actual < 0:
140
+ raise ValueError(f'Trained model missing <{name}> special token')
141
+ if actual != expected:
142
+ raise ValueError(
143
+ f'<{name}> id drift: requested {expected}, SentencePiece '
144
+ f'assigned {actual}. This mismatch is exactly the class of '
145
+ f'bug that broke a previous tokenizer version — refusing '
146
+ f'to silently proceed.'
147
+ )
148
+
149
+ # 3. Basic round-trip: encode -> decode must reproduce recognizable text
150
+ probe = "The quick brown fox jumps over 42 lazy dogs. def foo(): return None"
151
+ ids = sp.encode(probe, out_type=int)
152
+ if not ids:
153
+ raise ValueError('Validation encode produced empty output')
154
+ decoded = sp.decode(ids)
155
+ if not decoded.strip():
156
+ raise ValueError('Validation round-trip produced empty decode')
157
+
158
+ # 4. SPECIFIC regression check for the actual reported failure mode:
159
+ # repetitive-token degenerate decode ("to to to", ",,,"). This won't
160
+ # catch a MODEL that's actually stuck in a repetition loop (that's a
161
+ # decoding-strategy issue, separate from the tokenizer), but it DOES
162
+ # catch a tokenizer that maps distinct ids to the same or corrupted
163
+ # text, which was the real bug here: encode the same repeated-word
164
+ # probe multiple times and confirm token ids are stable and decode
165
+ # is exact, not degenerating into duplicated/garbled pieces.
166
+ repeat_probe = "to to to , , , the the the"
167
+ repeat_ids = sp.encode(repeat_probe, out_type=int)
168
+ repeat_decoded = sp.decode(repeat_ids)
169
+ # Re-encoding the decoded output should reproduce the same ids
170
+ # (idempotency) — this is the real symptom check: a corrupted rank/id
171
+ # mapping breaks exactly this property even when a single encode/decode
172
+ # pass looks fine.
173
+ reencoded_ids = sp.encode(repeat_decoded, out_type=int)
174
+ if reencoded_ids != repeat_ids:
175
+ raise ValueError(
176
+ f'Round-trip idempotency FAILED on repeated-token probe: '
177
+ f'encode->decode->encode did not reproduce the same ids. '
178
+ f'original={repeat_ids} reencoded={reencoded_ids}. This is '
179
+ f'the specific failure signature of an id/rank mapping bug.'
180
+ )
181
+
182
+ # 5. Byte-fallback sanity: an unusual/rare unicode character must not
183
+ # crash and must not silently become <unk> if byte_fallback is on —
184
+ # it should decompose into byte pieces instead.
185
+ exotic_probe = "emoji test \U0001F600 and rare char \u0800"
186
+ exotic_ids = sp.encode(exotic_probe, out_type=int)
187
+ if not exotic_ids:
188
+ raise ValueError('Byte-fallback validation: exotic-character probe produced empty encode')
189
+ exotic_decoded = sp.decode(exotic_ids)
190
+ if not exotic_decoded.strip():
191
+ raise ValueError('Byte-fallback validation: exotic-character round-trip produced empty decode')
192
+
193
+ logger.info(f"✓ Validation OK: vocab={actual_vocab} probe_tokens={len(ids)} "
194
+ f"round-trip idempotency verified, byte-fallback verified")
195
+
196
+
197
+ class TokenizerWrapper:
198
+ def __init__(self, model_path: str):
199
+ model_path = str(Path(model_path))
200
+ if not Path(model_path).exists():
201
+ raise FileNotFoundError(model_path)
202
+ self.sp = spm.SentencePieceProcessor(model_file=model_path)
203
+ self.vocab_size = int(self.sp.vocab_size())
204
+ self.pad_id = self.sp.pad_id()
205
+ self.unk_id = self.sp.unk_id()
206
+ self.bos_id = self.sp.bos_id()
207
+ self.eos_id = self.sp.eos_id()
208
+ for name, val in [('pad', self.pad_id), ('unk', self.unk_id), ('bos', self.bos_id), ('eos', self.eos_id)]:
209
+ if val < 0:
210
+ raise ValueError(f'SentencePiece model missing <{name}>')
211
+ self._special_ids = {self.pad_id, self.bos_id, self.eos_id}
212
+
213
+ def encode(self, text: str, add_bos: bool = True, add_eos: bool = False) -> List[int]:
214
+ if text is None:
215
+ raise ValueError('encode() received None')
216
+ if text == '':
217
+ ids: List[int] = []
218
+ else:
219
+ ids = list(self.sp.encode(text, out_type=int))
220
+ if add_bos:
221
+ ids = [self.bos_id] + ids
222
+ if add_eos:
223
+ ids = ids + [self.eos_id]
224
+ return ids
225
+
226
+ def encode_large_text(
227
+ self,
228
+ text: str,
229
+ add_bos: bool = True,
230
+ add_eos: bool = False,
231
+ chunk_chars: int = 2_000_000,
232
+ ) -> List[int]:
233
+ """OOM-safe encode for large (multi-MB+) strings, e.g. a whole
234
+ training corpus file read in one go. `encode()` above calls
235
+ self.sp.encode() on the ENTIRE string in one call -- even though
236
+ SentencePiece's C++ core is fast (not a slow Python BPE loop),
237
+ the RETURN VALUE is still built as one giant Python list[int],
238
+ and Python's own string handling means a large UTF-8 file can
239
+ already be 2-4x its byte size once decoded into a `str` object.
240
+ For a real, confirmed example: a 300MB single-file corpus OOM'd
241
+ a ~12-13GB Colab/Kaggle box within ~30 seconds calling encode()
242
+ on the whole file at once -- the failure was in Python-level
243
+ memory (str + list[int] materialization), not inside
244
+ SentencePiece's own encoding step.
245
+
246
+ This method processes `text` in fixed-size character chunks,
247
+ converting each chunk's ids to the accumulator and discarding
248
+ the chunk's own Python list before moving to the next chunk, so
249
+ peak memory is bounded by `chunk_chars` instead of len(text).
250
+ Chunk boundaries are placed at whitespace (never mid-word/
251
+ mid-token), so the output is IDENTICAL to what encode(text,
252
+ add_bos, add_eos) would produce on the whole string at once --
253
+ this is not an approximation, it's the same encoding, just
254
+ computed incrementally. Verified by direct comparison across
255
+ many chunk sizes and edge cases (empty string, no-whitespace
256
+ text, exact chunk-boundary alignment) against the whole-string
257
+ encode() path.
258
+
259
+ Use this instead of encode() for anything that might be large
260
+ (a whole file's contents) -- keep encode() for genuinely small
261
+ strings (a single sentence/prompt at inference time) where the
262
+ chunking overhead isn't worth paying.
263
+ """
264
+ if text is None:
265
+ raise ValueError('encode_large_text() received None')
266
+
267
+ if text == '':
268
+ ids: List[int] = []
269
+ if add_bos:
270
+ ids = [self.bos_id] + ids
271
+ if add_eos:
272
+ ids = ids + [self.eos_id]
273
+ return ids
274
+
275
+ all_ids: List[int] = []
276
+ pos = 0
277
+ n = len(text)
278
+ while pos < n:
279
+ end = min(pos + chunk_chars, n)
280
+ is_eof = (end >= n)
281
+ window = text[pos:end]
282
+ if is_eof:
283
+ piece = window
284
+ pos = end
285
+ else:
286
+ # Cut at the LAST whitespace boundary inside this
287
+ # window, so no token is split across two encode()
288
+ # calls (which would silently tokenize the seam
289
+ # differently than encoding the whole string at once).
290
+ cut = max(window.rfind(' '), window.rfind('\n'))
291
+ if cut <= 0:
292
+ # cut == -1: no whitespace anywhere in this window
293
+ # (pathological -- one token/URL longer than
294
+ # chunk_chars). cut == 0: the window's FIRST
295
+ # character is whitespace, which would make
296
+ # `piece = window[:0]` empty and, worse, silently
297
+ # drop that whitespace character (it's before the
298
+ # cut, so it's never included in this piece OR the
299
+ # next one, since pos would advance past it without
300
+ # encoding it). Both cases: take the whole window
301
+ # as one piece and advance past it -- correctness
302
+ # at this one seam matters less than a silently
303
+ # dropped character or a stalled loop, and this
304
+ # only ever triggers when chunk_chars is smaller
305
+ # than a typical word (never true at the real
306
+ # ~2MB default; only reachable with a deliberately
307
+ # tiny chunk_chars in tests).
308
+ piece = window
309
+ pos = end
310
+ else:
311
+ piece = window[:cut]
312
+ pos += cut
313
+ if not piece:
314
+ # Should be unreachable now that cut<=0 is folded into
315
+ # the "take whole window" branch above, but keep this
316
+ # as a hard backstop against any future edge case that
317
+ # produces an empty piece with is_eof=False -- without
318
+ # it, pos would never advance and the loop would spin
319
+ # forever.
320
+ if not is_eof:
321
+ pos += 1
322
+ continue
323
+ piece_ids = self.sp.encode(piece, out_type=int)
324
+ all_ids.extend(piece_ids)
325
+ del piece_ids
326
+
327
+ if add_bos:
328
+ all_ids = [self.bos_id] + all_ids
329
+ if add_eos:
330
+ all_ids = all_ids + [self.eos_id]
331
+ return all_ids
332
+
333
+ def encode_batch(
334
+ self,
335
+ texts: Sequence[str],
336
+ add_bos: bool = True,
337
+ add_eos: bool = False,
338
+ skip_errors: bool = False,
339
+ ) -> List[List[int]]:
340
+ out: List[List[int]] = []
341
+ for i, t in enumerate(texts):
342
+ try:
343
+ out.append(self.encode(t, add_bos=add_bos, add_eos=add_eos))
344
+ except Exception as e:
345
+ if skip_errors:
346
+ logger.warning(f"encode_batch: skipping item {i} ({e})")
347
+ continue
348
+ raise
349
+ return out
350
+
351
+ def decode(self, ids: Sequence[int], skip_special_tokens: bool = True) -> str:
352
+ # Drop anything outside the valid piece-id range first. This is
353
+ # required, not cosmetic: PyTorch's ignore_index=-100 convention for
354
+ # masked label positions means `ids` is very commonly a raw labels
355
+ # tensor, and sp.decode() raises IndexError on any id < 0 or
356
+ # >= vocab_size instead of skipping it.
357
+ ids = [int(i) for i in ids if 0 <= int(i) < self.vocab_size]
358
+ if skip_special_tokens:
359
+ filtered = [i for i in ids if i not in self._special_ids]
360
+ else:
361
+ filtered = [i for i in ids if i != self.pad_id]
362
+ return self.sp.decode(filtered)
363
+
364
+ def decode_batch(self, batch_ids: Sequence[Sequence[int]], skip_special_tokens: bool = True) -> List[str]:
365
+ return [self.decode(ids, skip_special_tokens=skip_special_tokens) for ids in batch_ids]
366
+
367
+ def save_config(self, path: str) -> None:
368
+ Path(path).write_text(json.dumps({
369
+ 'vocab_size': self.vocab_size,
370
+ 'pad_id': self.pad_id,
371
+ 'unk_id': self.unk_id,
372
+ 'bos_id': self.bos_id,
373
+ 'eos_id': self.eos_id,
374
+ }, indent=2), encoding='utf-8')
375
+
376
+ @classmethod
377
+ def from_config(cls, model_path: str, config_path: Optional[str] = None) -> 'TokenizerWrapper':
378
+ """Load and, if a config is given, verify special-id consistency against it."""
379
+ tok = cls(model_path)
380
+ if config_path and Path(config_path).exists():
381
+ cfg = json.loads(Path(config_path).read_text(encoding='utf-8'))
382
+ mismatches = {
383
+ k: (cfg[k], getattr(tok, k))
384
+ for k in ('vocab_size', 'pad_id', 'unk_id', 'bos_id', 'eos_id')
385
+ if k in cfg and cfg[k] != getattr(tok, k)
386
+ }
387
+ if mismatches:
388
+ raise ValueError(f'Tokenizer/config mismatch: {mismatches}')
389
+ return tok