phase4_prep: pin f6f10076; drive the leg check from the real planner, verify the manifest against itself and the mapped bytes, resume through load_cursor
Browse files- kernels/phase4_prep.py +80 -36
kernels/phase4_prep.py
CHANGED
|
@@ -13,18 +13,18 @@ import hashlib, json, os, subprocess, sys, time
|
|
| 13 |
|
| 14 |
os.chdir("/kaggle/working")
|
| 15 |
sys.path.insert(0, "/kaggle/working")
|
| 16 |
-
REV = os.environ.get("OUNCE100M_REV") or "
|
| 17 |
WANT = {
|
| 18 |
"ounce100m_credentials.py": ("ounce100m_credentials.py",
|
| 19 |
"6525f62f03f2d73650a1eb4f70fcb52d1194caad4ca88b2d8bd8fd54f88339b6"),
|
| 20 |
"shard_dataset.py": ("train/shard_dataset.py",
|
| 21 |
"adbcd96fbab9505a5e8dec4d2e83d00366a73f702786bffbde430195cf6f0471"),
|
| 22 |
"hubckpt.py": ("train/hubckpt.py",
|
| 23 |
-
"
|
| 24 |
"train_ounce100m.py": ("train/train_ounce100m.py",
|
| 25 |
-
"
|
| 26 |
"phase4_session.py": ("kernels/phase4_session.py",
|
| 27 |
-
"
|
| 28 |
}
|
| 29 |
BASE = "https://huggingface.co/Cion-lab/ounce100m-code/resolve/" + REV
|
| 30 |
for p, (rp, want) in sorted(WANT.items()):
|
|
@@ -45,11 +45,22 @@ def run(argv, label, timeout, env=None):
|
|
| 45 |
t0 = time.time()
|
| 46 |
e = dict(os.environ); e["PYTHONPATH"] = "/kaggle/working"
|
| 47 |
e.update(env or {})
|
| 48 |
-
|
| 49 |
-
env=e, bufsize=1)
|
| 50 |
import threading
|
|
|
|
|
|
|
| 51 |
killed = []
|
| 52 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 53 |
timer.daemon = True
|
| 54 |
timer.start()
|
| 55 |
out = []
|
|
@@ -71,12 +82,12 @@ def run(argv, label, timeout, env=None):
|
|
| 71 |
# Stage 1 -- the launcher itself, in rehearsal mode. On a CPU box torchrun is never reached: PREP_ONLY
|
| 72 |
# stops right after the pointer read, the mix download and the planner. Everything upstream of the first
|
| 73 |
# optimiser step is exercised with the real files and the real Hub.
|
|
|
|
|
|
|
|
|
|
|
|
|
| 74 |
rc1, _ = run([sys.executable, "phase4_session.py", REV], "PREP_LAUNCHER", 3000,
|
| 75 |
-
env=
|
| 76 |
-
"PLANNING_TOK_PER_S": os.environ.get("PLANNING_TOK_PER_S", "9696"),
|
| 77 |
-
"SESSION_GPU_HOURS": os.environ.get("SESSION_GPU_HOURS", "11.0"),
|
| 78 |
-
"QUOTA_LEFT_HOURS": os.environ.get("QUOTA_LEFT_HOURS", "27.5"),
|
| 79 |
-
"MIN_FREE_GB": "6"})
|
| 80 |
if rc1 != 0:
|
| 81 |
print("VERDICT PREP_STOP launcher rehearsal rc", rc1, flush=True)
|
| 82 |
raise SystemExit(2)
|
|
@@ -108,50 +119,83 @@ for b in bad[:10]:
|
|
| 108 |
print(" BAD", b, flush=True)
|
| 109 |
assert not bad, "the published mix does not match its own manifest"
|
| 110 |
|
|
|
|
| 111 |
import shard_dataset as SD
|
| 112 |
-
SEQ,
|
| 113 |
-
seqs_per_step =
|
| 114 |
tps = seqs_per_step * SEQ
|
| 115 |
-
|
|
|
|
| 116 |
store = SD.PackedTokenStore(ROOT)
|
| 117 |
n_samples = store.total_tokens // SEQ
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 118 |
print("READER", format(store.total_tokens, ","), "tokens in", len(store.shards), "shards ->",
|
| 119 |
format(n_samples, ","), "windows of", SEQ, "|", steps, "steps x", tps, "tok =",
|
| 120 |
-
format(steps * tps, ","), "tokens", flush=True)
|
| 121 |
assert n_samples >= steps * seqs_per_step, "not enough windows for the horizon"
|
| 122 |
|
| 123 |
-
#
|
| 124 |
-
|
| 125 |
-
|
| 126 |
-
|
| 127 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 128 |
avail = n_samples - start * seqs_per_step
|
| 129 |
-
|
| 130 |
-
|
| 131 |
-
|
| 132 |
-
print("LEG
|
| 133 |
-
start,
|
| 134 |
-
|
| 135 |
-
|
|
|
|
|
|
|
| 136 |
|
| 137 |
-
# Resume-equivalence
|
|
|
|
|
|
|
|
|
|
| 138 |
full = SD.make_dataset(store, SEQ, 20260919, start_sample=0)
|
| 139 |
-
for boundary in (
|
| 140 |
-
|
| 141 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 142 |
assert len(resumed) == len(full) - s0, (len(resumed), len(full), s0)
|
| 143 |
same = True
|
| 144 |
for j in (0, 1, len(resumed) // 2, len(resumed) - 1):
|
| 145 |
-
a = full[s0 + j][
|
| 146 |
-
|
| 147 |
-
|
| 148 |
print("RESUME_EQUIV step", boundary, "start_sample", s0, "len", len(resumed), "identical", same,
|
| 149 |
flush=True)
|
| 150 |
assert same, boundary
|
| 151 |
print("CHECK_OK windows", format(n_samples, ","), "shards", len(store.shards), flush=True)
|
| 152 |
"""
|
| 153 |
open("prepcheck.py", "w").write(CHECK)
|
| 154 |
-
rc2, _ = run([sys.executable, "prepcheck.py"], "CHECK_PUBLISHED_MIX", 2400)
|
| 155 |
|
| 156 |
# Stage 3 -- the reader's cursor format round-trips, and the trainer imports and prints --help on this
|
| 157 |
# image. Both are cheap and both have bitten before (E-028: a build script that could not import its
|
|
|
|
| 13 |
|
| 14 |
os.chdir("/kaggle/working")
|
| 15 |
sys.path.insert(0, "/kaggle/working")
|
| 16 |
+
REV = os.environ.get("OUNCE100M_REV") or "f6f100763eec8782851f1b418fe869df6a4dcb92"
|
| 17 |
WANT = {
|
| 18 |
"ounce100m_credentials.py": ("ounce100m_credentials.py",
|
| 19 |
"6525f62f03f2d73650a1eb4f70fcb52d1194caad4ca88b2d8bd8fd54f88339b6"),
|
| 20 |
"shard_dataset.py": ("train/shard_dataset.py",
|
| 21 |
"adbcd96fbab9505a5e8dec4d2e83d00366a73f702786bffbde430195cf6f0471"),
|
| 22 |
"hubckpt.py": ("train/hubckpt.py",
|
| 23 |
+
"568c31b300906cb8d78a59ef890baa9731b18e064e3d69b0eaea4b45c546f23e"),
|
| 24 |
"train_ounce100m.py": ("train/train_ounce100m.py",
|
| 25 |
+
"c4c178c0ac17c68de1521d85612a0dc604e6d4748dfdc8bd3e374512bc90ce96"),
|
| 26 |
"phase4_session.py": ("kernels/phase4_session.py",
|
| 27 |
+
"6e55edab4e9eedd78363c572ebdd409558e0c3f76e13f4ebfbe54e6c4dd4964d"),
|
| 28 |
}
|
| 29 |
BASE = "https://huggingface.co/Cion-lab/ounce100m-code/resolve/" + REV
|
| 30 |
for p, (rp, want) in sorted(WANT.items()):
|
|
|
|
| 45 |
t0 = time.time()
|
| 46 |
e = dict(os.environ); e["PYTHONPATH"] = "/kaggle/working"
|
| 47 |
e.update(env or {})
|
| 48 |
+
import signal as _sig
|
|
|
|
| 49 |
import threading
|
| 50 |
+
p = subprocess.Popen(argv, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True,
|
| 51 |
+
env=e, bufsize=1, start_new_session=True)
|
| 52 |
killed = []
|
| 53 |
+
|
| 54 |
+
def _kill():
|
| 55 |
+
# Kill the group, not the child: a staged `python -c snapshot_download` orphaned by a timeout
|
| 56 |
+
# keeps writing into mixroot while the next stage reads it.
|
| 57 |
+
killed.append(True)
|
| 58 |
+
try:
|
| 59 |
+
os.killpg(os.getpgid(p.pid), _sig.SIGTERM)
|
| 60 |
+
except Exception:
|
| 61 |
+
p.kill()
|
| 62 |
+
|
| 63 |
+
timer = threading.Timer(timeout, _kill)
|
| 64 |
timer.daemon = True
|
| 65 |
timer.start()
|
| 66 |
out = []
|
|
|
|
| 82 |
# Stage 1 -- the launcher itself, in rehearsal mode. On a CPU box torchrun is never reached: PREP_ONLY
|
| 83 |
# stops right after the pointer read, the mix download and the planner. Everything upstream of the first
|
| 84 |
# optimiser step is exercised with the real files and the real Hub.
|
| 85 |
+
PLANENV = {"PLANNING_TOK_PER_S": os.environ.get("PLANNING_TOK_PER_S", "9358"),
|
| 86 |
+
"SESSION_GPU_HOURS": os.environ.get("SESSION_GPU_HOURS", "6.9"),
|
| 87 |
+
"QUOTA_LEFT_HOURS": os.environ.get("QUOTA_LEFT_HOURS", "27.4"),
|
| 88 |
+
"MIN_FREE_GB": "6"}
|
| 89 |
rc1, _ = run([sys.executable, "phase4_session.py", REV], "PREP_LAUNCHER", 3000,
|
| 90 |
+
env=dict(PLANENV, PHASE4_PREP_ONLY="1"))
|
|
|
|
|
|
|
|
|
|
|
|
|
| 91 |
if rc1 != 0:
|
| 92 |
print("VERDICT PREP_STOP launcher rehearsal rc", rc1, flush=True)
|
| 93 |
raise SystemExit(2)
|
|
|
|
| 119 |
print(" BAD", b, flush=True)
|
| 120 |
assert not bad, "the published mix does not match its own manifest"
|
| 121 |
|
| 122 |
+
import phase4_session as S # constants and plan() only -- main() is behind __name__
|
| 123 |
import shard_dataset as SD
|
| 124 |
+
SEQ, steps = S.SEQ_LEN, S.HORIZON_STEPS
|
| 125 |
+
seqs_per_step = S.MICRO_BATCH * S.WORLD * S.ACCUM # the launcher's own geometry, not retyped
|
| 126 |
tps = seqs_per_step * SEQ
|
| 127 |
+
assert tps == S.TOKENS_PER_STEP, f"prep disagrees with the launcher about tokens/step: {tps}"
|
| 128 |
+
assert steps == int(S.TOKENS / tps), f"horizon disagrees: {steps}"
|
| 129 |
store = SD.PackedTokenStore(ROOT)
|
| 130 |
n_samples = store.total_tokens // SEQ
|
| 131 |
+
fingerprint = SD.mix_fingerprint(ROOT, man)
|
| 132 |
+
# The manifest is the corpus definition, and three different things have to agree about it: the record
|
| 133 |
+
# list, the counts it claims, and the bytes the reader actually mapped. A rebuild that truncated the list
|
| 134 |
+
# while leaving total_tokens alone would otherwise certify 1.11 B tokens of a smaller corpus (review).
|
| 135 |
+
if len(man["shards"]) != int(man.get("n_shards") or -1) or len(store.shards) != len(man["shards"]):
|
| 136 |
+
raise AssertionError(f"shard records {len(man['shards'])} vs n_shards {man.get('n_shards')} vs "
|
| 137 |
+
f"files opened {len(store.shards)}")
|
| 138 |
+
summed = sum(int(s["tokens"]) for s in man["shards"])
|
| 139 |
+
if summed != tok_s or store.total_tokens != summed:
|
| 140 |
+
raise AssertionError(f"manifest does not add up: records sum {summed}, total_tokens says {tok_s}, "
|
| 141 |
+
f"bytes the reader mapped say {store.total_tokens}")
|
| 142 |
print("READER", format(store.total_tokens, ","), "tokens in", len(store.shards), "shards ->",
|
| 143 |
format(n_samples, ","), "windows of", SEQ, "|", steps, "steps x", tps, "tok =",
|
| 144 |
+
format(steps * tps, ","), "tokens | fingerprint", fingerprint, flush=True)
|
| 145 |
assert n_samples >= steps * seqs_per_step, "not enough windows for the horizon"
|
| 146 |
|
| 147 |
+
# The schedule the run will actually live through: drive the real planner from step 0 until it reports the
|
| 148 |
+
# horizon, and require every leg to start where the previous one stopped. This is the check that catches a
|
| 149 |
+
# planner which stops one interval short of the end -- the last four steps of the run, unreachable, while
|
| 150 |
+
# every later session fails its own budget test for want of a whole interval.
|
| 151 |
+
legs, start = [], 0
|
| 152 |
+
while start < steps and len(legs) < 40:
|
| 153 |
+
pl = S.plan(start, 15.0)
|
| 154 |
+
if not pl["planned_steps"]:
|
| 155 |
+
raise AssertionError(f"the planner refuses to start at step {start} with "
|
| 156 |
+
f"{pl['usable_hours']} usable hours -- the run can never finish")
|
| 157 |
+
end = steps if pl["stop_after_steps"] == 0 else pl["stop_after_steps"]
|
| 158 |
+
assert end > start, (start, end)
|
| 159 |
+
need = (steps - start) * seqs_per_step # the trainer's own guard, at this start step
|
| 160 |
avail = n_samples - start * seqs_per_step
|
| 161 |
+
assert avail >= need, (start, need, avail)
|
| 162 |
+
store.tokens((end * seqs_per_step - 1) * SEQ, SEQ) # last window this leg reads must be addressable
|
| 163 |
+
legs.append((start, end))
|
| 164 |
+
print("LEG %5d -> %5d pct %5.1f -> %5.1f need %9d avail %9d" % (
|
| 165 |
+
start, end, 100.0 * start / steps, 100.0 * end / steps, need, avail), flush=True)
|
| 166 |
+
start = end
|
| 167 |
+
print("SESSIONS", len(legs), "ends_at", legs[-1][1], "horizon", steps, "contiguous",
|
| 168 |
+
[l[0] for l in legs[1:]] == [l[1] for l in legs[:-1]], flush=True)
|
| 169 |
+
assert legs[-1][1] == steps and [l[0] for l in legs[1:]] == [l[1] for l in legs[:-1]], legs
|
| 170 |
|
| 171 |
+
# Resume-equivalence through the cursor the trainer would actually write, at boundaries the run crosses:
|
| 172 |
+
# the first checkpoint, a mid-session one, and the last interval before the horizon. Building the resumed
|
| 173 |
+
# dataset from load_cursor() rather than a hand-computed offset is the point -- the derivation from
|
| 174 |
+
# step to samples_consumed is part of what is being tested.
|
| 175 |
full = SD.make_dataset(store, SEQ, 20260919, start_sample=0)
|
| 176 |
+
for boundary in sorted({S.CKPT_EVERY, legs[len(legs) // 2][0], steps - S.CKPT_EVERY}):
|
| 177 |
+
d = "/tmp/cur_%d" % boundary
|
| 178 |
+
os.makedirs(d, exist_ok=True)
|
| 179 |
+
SD.save_cursor(d, SD.cursor_dict(SEQ, boundary * seqs_per_step, 20260919, fingerprint,
|
| 180 |
+
step=boundary))
|
| 181 |
+
c = SD.load_cursor(d)
|
| 182 |
+
assert c["step"] * seqs_per_step == c["samples_consumed"], c
|
| 183 |
+
assert c["dataset_files_sha"] == fingerprint, c
|
| 184 |
+
s0 = c["samples_consumed"]
|
| 185 |
+
resumed = SD.make_dataset(store, SEQ, c["shuffle_seed"], start_sample=s0)
|
| 186 |
assert len(resumed) == len(full) - s0, (len(resumed), len(full), s0)
|
| 187 |
same = True
|
| 188 |
for j in (0, 1, len(resumed) // 2, len(resumed) - 1):
|
| 189 |
+
a, b = full[s0 + j], resumed[j]
|
| 190 |
+
same = same and a["input_ids"].tolist() == b["input_ids"].tolist() \
|
| 191 |
+
and a["labels"].tolist() == b["labels"].tolist()
|
| 192 |
print("RESUME_EQUIV step", boundary, "start_sample", s0, "len", len(resumed), "identical", same,
|
| 193 |
flush=True)
|
| 194 |
assert same, boundary
|
| 195 |
print("CHECK_OK windows", format(n_samples, ","), "shards", len(store.shards), flush=True)
|
| 196 |
"""
|
| 197 |
open("prepcheck.py", "w").write(CHECK)
|
| 198 |
+
rc2, _ = run([sys.executable, "prepcheck.py"], "CHECK_PUBLISHED_MIX", 2400, env=dict(PLANENV))
|
| 199 |
|
| 200 |
# Stage 3 -- the reader's cursor format round-trips, and the trainer imports and prints --help on this
|
| 201 |
# image. Both are cheap and both have bitten before (E-028: a build script that could not import its
|