File size: 23,623 Bytes
58258b8
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
# Multi-Node SDG Pipeline Runbook

End-to-end recipe for running the SDG pipeline distributed across **4 worker
nodes Γ— 8 H100 GPUs = 32 shards**, with cross-node MongoDB caching and HuggingFace
output gathering. Tested on ByteDance Arnold (Kubernetes-backed Azure ND96isr_H100_v5).

- Target config: `sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml`
- Final dataset: published to `<your-hf-username>/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1`
- Wall time: **~40–60 hours** end-to-end
- HF username placeholder used below: `fzzhang` β€” replace with yours

If the cluster dies, follow this runbook top-to-bottom on a fresh allocation.

---

## 0. Cluster prerequisites

- 4 worker nodes, each with 8Γ— H100 (80 GB), shared no-filesystem (per-pod overlay).
- NVIDIA driver supporting at least CUDA 12.4 (driver 535+ is fine).
- HF account with write access to the target dataset namespace.
- Each worker has these Arnold-provided env vars: `ARNOLD_WORKER_0_HOST`,
  `ARNOLD_WORKER_0_PORT`, `MY_POD_NAME` (ends with `worker-N`).
- Cross-worker IPv6 traffic on the `ARNOLD_WORKER_*_PORT` ports is open by default β€”
  verify in Β§3 below.

---

## 1. Per-worker bootstrap (run on EVERY worker, in parallel)

### 1.1 Clone the repo

```bash
cd ~/lmc_muon  # or your preferred dir
# git clone or rsync this repo into ./self-verified-distillation
cd self-verified-distillation
```

### 1.2 Create the conda env and install pinned dependencies

The driver caps at CUDA 12.4, so vLLM 0.8.5 + torch 2.6.0+cu124 is the highest
stack that runs. Newer vLLM (β‰₯0.9) ships cu126/cu130 wheels which **will not
run on driver 535**.

```bash
conda create -n sdg python=3.12 -y
conda activate sdg

pip install -U uv
uv pip install -r sdg/requirements.txt

# Downgrade to cu12.4-compatible stack
uv pip uninstall vllm torch torchvision torchaudio flashinfer-python triton
uv pip install torch==2.6.0 torchvision==0.21.0 \
  --index-url https://download.pytorch.org/whl/cu124
uv pip install "vllm==0.8.5" "transformers>=4.51.3,<5.0"

# Sanity check β€” MUST print "OK 1.0"
python -c "import torch; x=torch.zeros(1).cuda(); print('OK', (x+1).item())"
```

### 1.3 HuggingFace login

```bash
huggingface-cli login   # paste token; or set HF_TOKEN env var
python -c "from huggingface_hub import HfApi; print(HfApi().whoami()['name'])"
# should print your HF username (e.g. fzzhang)
```

---

## 2. MongoDB on worker_0 (ONLY worker_0)

### 2.1 Install mongod

```bash
cd ~/lmc_muon/self-verified-distillation
curl -O https://fastdl.mongodb.org/linux/mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz
tar xzf mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz
export PATH="$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:$PATH"
which mongod   # must print a path

# Persist for new shells
echo "export PATH=\"$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:\$PATH\"" >> ~/.bashrc
```

### 2.2 Start mongod bound to all interfaces with IPv6

The `--ipv6` flag is mandatory; without it `--bind_ip_all` is IPv4-only and
cross-worker traffic (IPv6 fabric) cannot reach mongod.

```bash
mkdir -p /tmp/sdg_mongo
nohup mongod --port $ARNOLD_WORKER_0_PORT --dbpath /tmp/sdg_mongo \
  --bind_ip_all --ipv6 > mongo.log 2>&1 &

sleep 5
grep -i "Listening on" mongo.log
# Expected: BOTH 0.0.0.0:<port> AND [::]:<port>

# Local sanity
python3 -c "from pymongo import MongoClient; \
  MongoClient(f'mongodb://[::1]:$ARNOLD_WORKER_0_PORT', serverSelectionTimeoutMS=5000).admin.command('ping'); \
  print('mongo OK locally')"
```

---

## 3. Verify cross-worker MongoDB reachability

On workers 1, 2, 3:

```bash
python3 -c "
import os
from pymongo import MongoClient
host = os.environ['ARNOLD_WORKER_0_HOST']
port = os.environ['ARNOLD_WORKER_0_PORT']
uri = f'mongodb://[{host}]:{port}'
MongoClient(uri, serverSelectionTimeoutMS=5000).admin.command('ping')
print(f'mongo OK from {os.uname().nodename} β†’ {uri}')
"
```

You should see `mongo OK from ...` on each. If any worker times out, the
platform is blocking inter-worker traffic on that port β€” do not proceed.

---

## 4. Patch the config on each worker

Point `mongo_uri` at the IPv6 MongoDB and set `num_shards=32` to use all GPUs.

```bash
# Run on EACH worker
python3 <<'EOF'
import os, yaml
host = os.environ['ARNOLD_WORKER_0_HOST']
port = os.environ['ARNOLD_WORKER_0_PORT']
path = 'sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml'
with open(path) as f: cfg = yaml.safe_load(f)
cfg['mongo_uri'] = f'mongodb://[{host}]:{port}'
cfg['num_shards'] = 32
with open(path, 'w') as f: yaml.dump(cfg, f, sort_keys=False)
print('updated:', cfg['mongo_uri'])
EOF
```

---

## 5. Helper scripts

All orchestration scripts live in `scripts/` (committed to the repo, so each
worker has them after `git clone`). They are designed to be **runnable as-is**;
override behavior with env vars rather than editing the files:

| Script | Purpose | Run on |
|---|---|---|
| `scripts/launch_shards.sh` | Launch 8 shards on this worker, one per GPU. Honors `CONFIG`, `LIMIT`, `NUM_SHARDS`. | every worker |
| `scripts/upload_my_shards.py` | Poll for this worker's 8 stats.json, then upload to staging HF repo. Honors `HF_USER`, `RUN_NAME`. | every worker |
| `scripts/combine_and_publish.py` | Poll for 4 worker-done markers, gather, combine, publish to HF (full launcher.py stats schema). Honors `HF_USER`, `RUN_NAME`, `NUM_WORKERS`, `NUM_SHARDS`. | worker_0 only |
| `scripts/monitor.sh` | Per-worker live status (GPUs, current shard progress, completion counts, mongo cache size via cross-worker IPv6). | any worker |

Defaults target the math53K config + `HF_USER=fzzhang` + 4 workers Γ— 32 shards.
**Why scripts/, not inline heredocs:** during the original run, terminal
autocomplete corrupted multi-line heredoc pastes (Python files came out
unparseable). Committed files bypass paste entirely β€” you just `git pull` and
run them.

**Verify they're present and parse on each worker:**

```bash
ls scripts/
python -c "import ast; ast.parse(open('scripts/upload_my_shards.py').read()); ast.parse(open('scripts/combine_and_publish.py').read()); print('OK')"
bash -n scripts/launch_shards.sh && bash -n scripts/monitor.sh && echo "shell OK"
```

If you need a different HF account or a different config, export the env vars
before running (no file edits needed):

```bash
export HF_USER=youraccount
export RUN_NAME=your_run_name_matching_your_config
export CONFIG=sdg/configs/your_config.yaml
```

---

## 6. Per-worker uploader

`scripts/upload_my_shards.py` runs on **every worker** β€” sleep-polls every 2
minutes for that worker's 8 `stats.json` to appear, then uploads each
`shards/shard_NNN/` dir to `<HF_USER>/<RUN_NAME>-staging` and writes a
`workers/worker_N.done` marker so the combiner knows that worker is finished.

Consumes near-zero resources while waiting β€” safe to launch any time (before,
during, or after the shards run). No file edits needed; override `HF_USER` /
`RUN_NAME` via env if you're not using defaults (see Β§5).

---

## 7. Combiner (worker_0 only)

`scripts/combine_and_publish.py` runs on **worker_0 only** β€” sleep-polls for all
4 worker-done markers, downloads the staging repo, concatenates the 32
`output.jsonl` files, aggregates stats with the **full launcher.py schema**
(`summary`, `checks`, `failure_breakdown`, `failure_breakdown_examples`,
`num_shards`), enforces a row-count vs `total_passed` integrity check, and
publishes to `<HF_USER>/<RUN_NAME>` with a HuggingFace dataset-card README.

Set `HF_USER` / `RUN_NAME` / `NUM_WORKERS` / `NUM_SHARDS` via env if defaults
don't match.

---

## 8. Optional: smoke test before committing 40+ hours

Tiny end-to-end validation (~15 min on all 32 GPUs). Confirms multi-worker
plumbing without burning real time. Run **on each worker**:

```bash
LIMIT=128 ./scripts/launch_shards.sh
```

Wait ~15 min, then check completion (on each worker):

```bash
grep -l "PIPELINE COMPLETE" logs/worker_*/shard_*.log | wc -l   # should be 8
```

If all four workers report 8, plumbing is verified. Clear the smoke outputs
before the real run:

```bash
# On each worker
rm -rf output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1
rm -rf logs/worker_*
# Optionally wipe the cache too (on worker_0)
pkill mongod && sleep 2
rm -rf /tmp/sdg_mongo/*
# Restart mongod (see Β§2.2)
```

---

## 9. Production launch

### 9.1 Launch shards on each worker

```bash
./scripts/launch_shards.sh
```

### 9.2 Start the per-worker uploader on each worker

It will sleep-poll until its 8 shards complete, then upload.

```bash
nohup python scripts/upload_my_shards.py > upload.log 2>&1 &
disown
```

### 9.3 Start the combiner on worker_0 ONLY

It will sleep-poll until all 4 worker-done markers exist on HF, then combine
and publish.

```bash
nohup python scripts/combine_and_publish.py > combine.log 2>&1 &
disown
```

After this, you can disconnect all terminals. Work continues under `nohup`.

### 9.4 (Recommended) Verify integrity once shards complete

Before the uploader kicks off uploads, you can sanity-check each shard locally.
For `selection_policy: first_valid` (the default), `output.jsonl` line count
must equal `total_passed` in `stats.json`. Run **on each worker**:

```bash
RUN=output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1/shards
for d in $RUN/shard_*; do
  lines=$(wc -l < $d/output.jsonl 2>/dev/null)
  passed=$(python3 -c "import json; print(json.load(open('$d/stats.json'))['summary']['total_passed'])" 2>/dev/null)
  echo "$(basename $d): $lines rows | total_passed=$passed"
done
```

Lines must match `total_passed` for every shard. The combiner re-enforces this
check globally and will `raise RuntimeError` if anything diverges β€” catching it
locally first is faster.

---

## 10. Monitoring

### 10.1 Per-worker live view: `scripts/monitor.sh`

```bash
watch -n 30 ./scripts/monitor.sh
# Or if `watch` isn't installed in the container:
while true; do clear; ./scripts/monitor.sh; sleep 30; done
```

Each refresh shows: per-GPU util, the most-recent tqdm state for each shard
(decoded from `\r`-overwritten progress bars), per-worker completion counts,
and the cluster-wide MongoDB cache size (queried over IPv6 so it works from any
worker, not just worker_0).

⚠ **Do not use `tail -1 | cut -cN`** to read shard progress β€” tqdm writes one
physical line with `\r` overwrites, so `cut -cN` shows the line's *beginning*
(which never changes) while progress advances at the *end*. `monitor.sh`
correctly does `tail -c 500 | tr '\r' '\n' | tail -1`.

Expected sequence in the per-shard line over the run:
1. `Loading model:` / `Capturing CUDA graphs` (~2 min vLLM warmup)
2. `Processed prompts: N/512 [...]` (Stage 3 generation, batch size 64 Γ— n=8 = 512)
3. `[Generation] N/1656 examples (...%)` (between-batch progress lines)
4. `Round 0: M pending rows` (Stage 4 UQ starts; tqdm batches now sized by pending rows)
5. `Cycle: ... Factual: ... Correctness: ...` (inside a UQ round)
6. `PIPELINE COMPLETE` (shard finished)

To definitively confirm a shard is in UQ rather than generation, grep the log:

```bash
grep -E "Round [0-9]+:|Cycle:|Factual:|Correctness:|\[Generation\]|PIPELINE" \
  logs/worker_0/shard_0.log | tail -20
```

You'll see `[Generation] 1656/1656 examples (100.0%)` followed by `Round 0:` if
generation is done.

### 10.2 Aggregate signal: MongoDB cache size

`monitor.sh` already prints this. Standalone:

```bash
python3 -c "
from pymongo import MongoClient
import os
host = os.environ['ARNOLD_WORKER_0_HOST']
port = os.environ['ARNOLD_WORKER_0_PORT']
print(MongoClient(f'mongodb://[{host}]:{port}').sdg_cache.inference_cache.count_documents({}))
"
```

Grows continuously. Expect hundreds of thousands by mid-run, well over a
million during UQ. **Use the IPv6 host, not `127.0.0.1`** β€” only worker_0 can
reach loopback mongod.

### 10.3 Uploader / combiner status

```bash
tail -f upload.log    # on any worker (file lives in working dir)
tail -f combine.log   # on worker_0
```

The combiner prints `DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME>`
when finished β€” that's the success signal for the entire pipeline.

---

## 11. Troubleshooting

### Shards exit immediately
- `Exit 127`: `mongod` (or `python`) not on PATH. `which mongod`, re-source `~/.bashrc`.
- `Qwen2Tokenizer has no attribute all_special_tokens_extended`: transformers
  too new for vLLM 0.8.5. `uv pip install "transformers>=4.51.3,<5.0"`.
- `CUDA driver version is insufficient`: torch+CUDA wheel doesn't match driver.
  Reinstall via Β§1.2 exactly.

### MongoDB connection refused from other workers
- mongod log shows only `0.0.0.0:PORT` (no `[::]:PORT`): missing `--ipv6` flag.
  `pkill mongod`, restart with `--bind_ip_all --ipv6`.

### `monitor.sh` shows shards stuck at `0%` for hours
- Almost always a display lie, not real stuckness. `tail -1 | cut -cN` would
  show the beginning of a long `\r`-overwritten line forever; the actual
  progress is at the *end*. Use `scripts/monitor.sh` (uses `tail -c 500 | tr '\r' '\n' | tail -1`).
- Confirm work is happening: GPU util β‰₯ 80%, log file `mtime` updates,
  MongoDB cache count growing.

### `monitor.sh` shows `mongo cache: <unreachable>` from a non-zero worker
- Workers other than worker_0 cannot reach `127.0.0.1:PORT`. The script uses
  `ARNOLD_WORKER_0_HOST` over IPv6, but if that env var isn't set in the
  current shell, the query fails. Re-source whatever sets Arnold env vars or
  run from a fresh login shell.

### Inline Python heredoc files won't parse
- Terminal autocomplete corrupts `cat > file <<'EOF' ... EOF`. Symptoms:
  `unterminated string literal`, `import os, re, time Path` (extra word
  appended), or scripts cut off mid-line.
- **Fix**: don't paste β€” use the committed files in `scripts/`. They survive
  `git pull`. If you must hand-write, use `nano <file>` (editor paste bypasses
  shell autocomplete) or base64-encode locally and `echo '...' | base64 -d > file`.

### Trial gets reclaimed mid-run
- Shard outputs on local disk are lost (no shared FS) and the MongoDB cache at
  `/tmp/sdg_mongo` dies with worker_0.
- **Mitigation**: start the uploaders in Β§9.2 *at launch time*, not at the
  end. They sleep-poll for completion and push completed work to HF
  incrementally. If a worker dies after some shards finish but before its 8
  are all done, you lose only the unfinished ones.
- Once a worker has written its `worker_N.done` marker on HF, that worker's
  data is durable. The combiner can run from any machine afterward, including
  a fresh trial.

### Some shards fail mid-run
- Each shard is independent; just relaunch the failed ones.
- Find failures: `grep -L "PIPELINE COMPLETE" logs/worker_*/shard_*.log`
- Relaunch one: `CUDA_VISIBLE_DEVICES=I nohup python -m sdg.generate --config <CONFIG> --shard-id N --num-shards 32 > logs/worker_W/shard_N.log 2>&1 &`
- MongoDB cache means most work is replayed for free on the rerun.

### Published `stats.json` is missing `failure_breakdown` fields
- Old `combine_and_publish.py` versions wrote a simplified stats schema.
  Current `scripts/combine_and_publish.py` includes the full schema.
- If you have a published dataset with the old stats, you can re-aggregate
  from `gathered/shards/*/stats.json` (still present on worker_0 after the
  combine) and re-upload `stats.json` β€” same logic as the combiner's stats
  section, just runs standalone.

---

## 12. After publish: verify and clean up

### 12.1 Verify the published dataset loads

```bash
python3 -c "
from datasets import load_dataset
ds = load_dataset('fzzhang/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1', split='train')
print('rows:', len(ds))
print('keys:', list(ds[0].keys()))
print('roles:', [m['role'] for m in ds[0]['messages']])
"
```

Should print the row count (~21k), `keys: ['messages']`, and the roles
(`['user', 'assistant']`). Confirms the HF dataset-card README correctly
auto-loads `output.jsonl`.

### 12.2 Clean up local + intermediate state

```bash
# On each worker β€” free disk if you plan to reuse the trial
rm -rf output/ logs/ gathered/ combined_output.jsonl combined_stats.json upload.log

# On worker_0
pkill mongod
rm -rf /tmp/sdg_mongo combine.log

# Optionally delete the intermediate staging HF repo
HF_USER=fzzhang RUN_NAME=qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 python -c "
import os
from huggingface_hub import HfApi
repo = f\"{os.environ['HF_USER']}/{os.environ['RUN_NAME']}-staging\"
HfApi().delete_repo(repo_id=repo, repo_type='dataset')
print('deleted', repo)
"
```

The final HF dataset is independent of the staging repo and the MongoDB
cache β€” deleting them after the run is purely cosmetic.

---

## Step-by-step playbook

Every command in execution order. The "Where" column tells you which worker to
run the command on. Commands marked "every worker" should be run on workers 0,
1, 2, AND 3 (in parallel terminals to save wall time). Steps numbered with a
letter (e.g. 5a) are part of the same logical step but run on different
workers.

### Step 1 β€” get the code (every worker)

```bash
cd ~/lmc_muon                # or wherever you keep the repo
git clone <your-repo-url> self-verified-distillation
cd self-verified-distillation
```

### Step 2 β€” Python env + pinned deps (every worker)

```bash
conda create -n sdg python=3.12 -y
conda activate sdg
pip install -U uv
uv pip install -r sdg/requirements.txt
uv pip uninstall vllm torch torchvision torchaudio flashinfer-python triton
uv pip install torch==2.6.0 torchvision==0.21.0 \
  --index-url https://download.pytorch.org/whl/cu124
uv pip install "vllm==0.8.5" "transformers>=4.51.3,<5.0"
python -c "import torch; x=torch.zeros(1).cuda(); print('OK', (x+1).item())"
# Must print: OK 1.0
```

### Step 3 β€” HuggingFace login (every worker)

```bash
huggingface-cli login   # paste your token
python -c "from huggingface_hub import HfApi; print(HfApi().whoami()['name'])"
# Must print your HF username
```

### Step 4 β€” install mongod (worker_0 ONLY)

```bash
cd ~/lmc_muon/self-verified-distillation
curl -O https://fastdl.mongodb.org/linux/mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz
tar xzf mongodb-linux-x86_64-ubuntu2204-8.0.4.tgz
export PATH="$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:$PATH"
echo "export PATH=\"$PWD/mongodb-linux-x86_64-ubuntu2204-8.0.4/bin:\$PATH\"" >> ~/.bashrc
which mongod   # must print a path
```

### Step 5 β€” start mongod (worker_0 ONLY)

```bash
mkdir -p /tmp/sdg_mongo
nohup mongod --port $ARNOLD_WORKER_0_PORT --dbpath /tmp/sdg_mongo \
  --bind_ip_all --ipv6 > mongo.log 2>&1 &
sleep 5
grep -i "Listening on" mongo.log
# Must show BOTH 0.0.0.0:<port> AND [::]:<port>

python3 -c "from pymongo import MongoClient; \
  MongoClient(f'mongodb://[::1]:$ARNOLD_WORKER_0_PORT', serverSelectionTimeoutMS=5000).admin.command('ping'); \
  print('mongo OK locally')"
```

### Step 6 β€” verify cross-worker mongo reachability (workers 1, 2, 3)

```bash
python3 -c "
import os
from pymongo import MongoClient
host = os.environ['ARNOLD_WORKER_0_HOST']
port = os.environ['ARNOLD_WORKER_0_PORT']
uri = f'mongodb://[{host}]:{port}'
MongoClient(uri, serverSelectionTimeoutMS=5000).admin.command('ping')
print(f'mongo OK from {os.uname().nodename} β†’ {uri}')
"
```

All three must print `mongo OK from ...`. If any times out, the platform is
blocking inter-worker traffic β€” stop here.

### Step 7 β€” patch the config (every worker)

```bash
python3 <<'EOF'
import os, yaml
host = os.environ['ARNOLD_WORKER_0_HOST']
port = os.environ['ARNOLD_WORKER_0_PORT']
path = 'sdg/configs/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1.yaml'
with open(path) as f: cfg = yaml.safe_load(f)
cfg['mongo_uri'] = f'mongodb://[{host}]:{port}'
cfg['num_shards'] = 32
with open(path, 'w') as f: yaml.dump(cfg, f, sort_keys=False)
print('updated:', cfg['mongo_uri'])
EOF
```

### Step 8 β€” verify scripts/ are intact (every worker)

```bash
python -c "import ast; ast.parse(open('scripts/upload_my_shards.py').read()); ast.parse(open('scripts/combine_and_publish.py').read()); print('OK')"
bash -n scripts/launch_shards.sh && bash -n scripts/monitor.sh && echo "shell OK"
```

### Step 9 β€” (optional) smoke test (every worker)

```bash
LIMIT=128 ./scripts/launch_shards.sh
# Wait ~15-20 min, then check:
grep -l "PIPELINE COMPLETE" logs/worker_*/shard_*.log | wc -l    # must be 8
```

If smoke passes, clear outputs before the real run:

```bash
# Every worker:
rm -rf output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1
rm -rf logs/worker_*
# Worker_0 only β€” wipe cache and restart mongod:
pkill mongod && sleep 2 && rm -rf /tmp/sdg_mongo/*
# Then redo Step 5 on worker_0
```

### Step 10 β€” launch the real run (every worker)

```bash
./scripts/launch_shards.sh
# Expected: 8 PIDs printed. ~5 min later, all 8 GPUs should be at ~95% util.
```

### Step 11 β€” start the per-worker uploader (every worker)

Safe to launch immediately β€” it sleep-polls until the worker's shards finish.

```bash
nohup python scripts/upload_my_shards.py > upload.log 2>&1 &
disown
```

### Step 12 β€” start the combiner (worker_0 ONLY)

```bash
nohup python scripts/combine_and_publish.py > combine.log 2>&1 &
disown
```

### Step 13 β€” monitor (any worker, anytime)

```bash
watch -n 30 ./scripts/monitor.sh
# Or, if `watch` isn't installed:
while true; do clear; ./scripts/monitor.sh; sleep 30; done
```

After this you can disconnect all terminals. The shards, uploaders, and
combiner all run under `nohup`.

### Step 14 β€” (recommended) per-shard integrity check after shards finish (every worker)

When `monitor.sh` shows `stats.json: 8/8` and `PIPELINE COMPLETE: 8/8` on a
worker, sanity-check its outputs:

```bash
RUN=output/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1/shards
for d in $RUN/shard_*; do
  lines=$(wc -l < $d/output.jsonl 2>/dev/null)
  passed=$(python3 -c "import json; print(json.load(open('$d/stats.json'))['summary']['total_passed'])" 2>/dev/null)
  echo "$(basename $d): $lines rows | total_passed=$passed"
done
```

`lines` must equal `total_passed` for every shard. If they diverge, the
combiner will refuse to publish β€” fix the broken shard before its uploader
starts pushing.

### Step 15 β€” wait for `DONE` (worker_0)

```bash
tail -f combine.log
# Watch for: DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME>
```

### Step 16 β€” verify the published dataset loads (any machine with HF access)

```bash
python3 -c "
from datasets import load_dataset
ds = load_dataset('fzzhang/qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1', split='train')
print('rows:', len(ds))
print('keys:', list(ds[0].keys()))
print('roles:', [m['role'] for m in ds[0]['messages']])
"
```

### Step 17 β€” (optional) cleanup

```bash
# Every worker
rm -rf output/ logs/ gathered/ combined_output.jsonl combined_stats.json upload.log

# Worker_0 only
pkill mongod
rm -rf /tmp/sdg_mongo combine.log

# Delete the intermediate staging HF repo (worker_0)
HF_USER=fzzhang \
RUN_NAME=qwen3_4b_openthoughts3_math53K_instill_n8_valredundancy5_round1 \
python -c "
import os
from huggingface_hub import HfApi
repo = f\"{os.environ['HF_USER']}/{os.environ['RUN_NAME']}-staging\"
HfApi().delete_repo(repo_id=repo, repo_type='dataset')
print('deleted', repo)
"
```

---

## Success signal

The whole pipeline succeeds when `combine.log` on worker_0 prints:
```
DONE: https://huggingface.co/datasets/<HF_USER>/<RUN_NAME>
```

Everything else (per-shard logs, mongo cache count, GPU util, staging repo,
upload markers) is plumbing that supports that one line.