svd-code / RUNBOOK.md
fzzhang's picture
Upload folder using huggingface_hub
58258b8 verified
|
Raw History Blame Contribute Delete
23.6 kB
# 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.