Cion-lab commited on
Commit
f6f1007
·
verified ·
1 Parent(s): d5751b9

phase4_session: plan() tiles to the horizon, kill the process group, df_free_gb helper, plan before the pull, session-relative TRAIN timeout

Browse files
Files changed (1) hide show
  1. kernels/phase4_session.py +105 -57
kernels/phase4_session.py CHANGED
@@ -19,6 +19,7 @@ No credentials in this file: they are fetched at run time by ounce100m_credentia
19
 
20
  import json
21
  import os
 
22
  import subprocess
23
  import sys
24
  import time
@@ -63,20 +64,33 @@ def sh(argv, label, timeout=None, env=None):
63
  # child's traceback arrives in the same stream instead of a separate tail that hides it (E-023).
64
  import threading
65
  p = subprocess.Popen(argv, cwd=WORK, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
66
- text=True, bufsize=1, env=e)
67
  state = {"killed": False}
68
 
69
  def _kill():
 
 
 
70
  state["killed"] = True
71
  try:
 
 
72
  p.kill()
 
 
 
 
73
  except Exception:
74
- pass
 
75
 
76
  timer = threading.Timer(timeout, _kill) if timeout else None
77
  if timer:
78
  timer.daemon = True
79
  timer.start()
 
 
 
80
  lines = []
81
  for line in (p.stdout or []):
82
  lines.append(line.rstrip("\n"))
@@ -84,6 +98,7 @@ def sh(argv, label, timeout=None, env=None):
84
  rc = p.wait()
85
  if timer:
86
  timer.cancel()
 
87
  if state["killed"]:
88
  rc = -9
89
  print(f"{label}_RC {rc}{' TIMEOUT' if state['killed'] else ''} seconds "
@@ -106,33 +121,57 @@ def fetch(rev, want):
106
 
107
 
108
  def plan(start_step, free_gb):
109
- """Largest whole number of checkpoint intervals that fits this session's remaining hours."""
110
- seconds_per_ckpt = CKPT_EVERY * TOKENS_PER_STEP / PLANNING_TOK_PER_S
 
 
 
 
 
 
 
111
  # The reserve comes off the total, not off one term: `min(session, quota - reserve)` let an 11 h
112
  # session plan 11 h of training and leave nothing for the code fetch, the 2.3 GB mix download, the
113
  # push/verify cycles and the end-of-run evaluation (review finding).
114
  usable = max(0.0, min(SESSION_GPU_HOURS, QUOTA_LEFT_HOURS) - RESERVE_HOURS)
115
- n = int(usable * 3600.0 // seconds_per_ckpt)
116
- remaining_intervals = max(0, (HORIZON_STEPS - start_step + CKPT_EVERY - 1) // CKPT_EVERY)
117
- n = min(n, remaining_intervals)
118
- target = start_step + n * CKPT_EVERY
119
- # The trainer refuses a --stop-after-steps that is not a whole multiple of its own interval, and
120
- # 3,814 is not one (3814 % 381 = 385). So the session that would cross the horizon is told to run
121
- # to it instead -- otherwise the final leg of the run could never be launched.
122
- stop = 0 if target >= HORIZON_STEPS else min((HORIZON_STEPS // CKPT_EVERY) * CKPT_EVERY, target)
123
  return {"start_step": start_step, "horizon_steps": HORIZON_STEPS, "ckpt_every": CKPT_EVERY,
124
  "tokens_per_step": TOKENS_PER_STEP, "planning_tok_per_s": PLANNING_TOK_PER_S,
125
- "hours_per_checkpoint_interval": round(seconds_per_ckpt / 3600.0, 3),
126
  "session_gpu_hours": SESSION_GPU_HOURS, "quota_left_hours": QUOTA_LEFT_HOURS,
127
  "reserve_hours": RESERVE_HOURS, "usable_hours": round(usable, 3),
128
- "intervals_this_session": n, "stop_after_steps": stop,
129
- "runs_to_horizon_this_session": stop == 0 and n > 0,
130
- "tokens_at_stop": (HORIZON_STEPS if stop == 0 else stop) * TOKENS_PER_STEP,
131
- "pct_of_run_at_stop": round(100.0 * (HORIZON_STEPS if stop == 0 else stop)
132
- * TOKENS_PER_STEP / TOKENS, 2),
133
  "free_gb_before_start": free_gb}
134
 
135
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
136
  def main():
137
  os.chdir(WORK)
138
  sys.path.insert(0, WORK)
@@ -143,38 +182,30 @@ def main():
143
  "shard_dataset.py": ("train/shard_dataset.py",
144
  "adbcd96fbab9505a5e8dec4d2e83d00366a73f702786bffbde430195cf6f0471"),
145
  "hubckpt.py": ("train/hubckpt.py",
146
- "c8c958417dc48db92f8be5c4b505581c746be9cc10110a358230ac47c480d5bf"),
147
  "train_ounce100m.py": ("train/train_ounce100m.py",
148
- "487edc1be54f03b5f1a312901580cc3e1504c412322ed074ceff869febbfb8bc"),
149
  })
150
 
151
  import ounce100m_credentials
152
  print("creds:", json.dumps(ounce100m_credentials.install(verify=True)), flush=True)
153
- rc, out = sh(["bash", "-c", "nvidia-smi --query-gpu=name,memory.used,memory.total "
154
- "--format=csv,noheader; df -k /kaggle/working | tail -1; free -g | head -2"],
155
- "env")
156
- # `df -k` not `-h`: the human-readable form prints "318G"/"500M"/"1.5T" and a float() on those raises,
157
- # which would take down the session before training even starts.
158
- free_gb = None
159
- for line in out.splitlines():
160
- if "overlay" in line or line.startswith("/dev/"):
161
- parts = line.split()
162
- if len(parts) >= 4:
163
- try:
164
- free_gb = float(parts[3]) / 1048576.0
165
- except ValueError:
166
- free_gb = None
167
  if free_gb is None:
168
- print("VERDICT REFUSED_TO_START: could not read the free space on /kaggle/working from "
169
- "df output -- refusing to guess whether the mix and two checkpoints will fit", flush=True)
170
  raise SystemExit(10)
171
- if free_gb < MIN_FREE_GB:
172
- print(f"VERDICT REFUSED_TO_START: {free_gb:.1f} GB free on /kaggle/working, this run needs at "
173
- f"least {MIN_FREE_GB} GB (2.3 GB mix + ~1.7 GB per checkpoint + the unpacked resume dir). "
174
- "A Kaggle GPU instance starts with ~20 GB, so a low number here means something is left "
175
- "over from a previous session on this box.", flush=True)
 
176
  raise SystemExit(10)
177
- print(f"disk: {free_gb:.1f} GB free", flush=True)
178
 
179
  # The Hub pointer is the only truth about where the run is (never local state -- §3.13).
180
  import hubckpt
@@ -192,20 +223,40 @@ def main():
192
  print("VERDICT RUN_ALREADY_COMPLETE at step", start_step, flush=True)
193
  return
194
 
 
 
 
 
 
 
 
 
 
 
 
195
  # Data: the published mix, fetched by anonymous read so this works on any instance. The checks run
196
  # here rather than inside the downloader so a quoting mistake cannot turn into a silent pass.
197
  rc, _ = sh([sys.executable, "-c",
198
  "import sys; sys.path.insert(0, '/kaggle/working');\n"
199
  "from huggingface_hub import snapshot_download\n"
200
  f"p = snapshot_download(repo_id='{MIX_REPO}', repo_type='dataset',\n"
201
- " local_dir='/kaggle/working/mixroot')\n"
202
  "print('downloaded to', p)\n"],
203
  "fetch_mix", 3600)
204
  if rc != 0:
205
  raise SystemExit(7)
 
 
 
 
 
206
  man = json.load(open(os.path.join(WORK, "mixroot", "manifest.json")))
207
  windows = int(man["total_tokens"]) // SEQ_LEN
208
  needed_windows = HORIZON_STEPS * MICRO_BATCH * ACCUM * WORLD
 
 
 
 
209
  print(f"mix: {man['total_tokens']:,} train tokens in {man['n_shards']} shards, "
210
  f">= {man['distinct_sources_per_shard_min']} sources per shard, "
211
  f"{windows:,} windows of {SEQ_LEN}; the horizon needs {needed_windows:,} "
@@ -219,21 +270,13 @@ def main():
219
  "implicitly by launching anyway.", flush=True)
220
  raise SystemExit(9)
221
 
222
- pl = plan(start_step, free_gb)
223
- print("PLAN_JSON", json.dumps(pl, indent=1), flush=True)
224
- if not pl["intervals_this_session"]:
225
- print(f"VERDICT REFUSED_TO_START: {pl['usable_hours']} usable hours will not reach one "
226
- f"{pl['hours_per_checkpoint_interval']} h checkpoint interval. Launching anyway would "
227
- "burn quota on work that cannot be resumed from. Raise SESSION_GPU_HOURS or wait for "
228
- "quota.", flush=True)
229
- raise SystemExit(8)
230
  if os.environ.get("PHASE4_PREP_ONLY") == "1":
231
  # Everything a session can get wrong before it bills GPU time has now been checked: the pinned
232
  # code hashes, the credentials resolve, the checkpoint pointer is readable, the published mix is
233
  # downloadable and large enough, and the plan lands on a checkpoint boundary. A CPU instance can
234
  # run all of that for free (E-031's arithmetic is exactly the kind of thing to find there).
235
  print("VERDICT PREP_ONLY_OK stop_after_steps", pl["stop_after_steps"],
236
- "intervals", pl["intervals_this_session"], flush=True)
237
  return
238
 
239
  argv = [sys.executable, "-u", "train_ounce100m.py",
@@ -245,15 +288,20 @@ def main():
245
  "--log-every", "20", "--val-tokens", "2000000",
246
  "--resume", "auto", "--stop-after-steps", str(pl["stop_after_steps"])]
247
  rc, out = sh(["torchrun", "--nproc_per_node=2"] + argv, "TRAIN",
248
- timeout=int((pl["intervals_this_session"] * pl["hours_per_checkpoint_interval"] + 1.5)
249
- * 3600))
 
 
 
250
  for line in out.splitlines():
251
  if line.startswith(("CKPT ", "RUN_JSON ", "resume from", "auto-resume", "precision:",
252
  "params:", "validation loss", "segment boundary")):
253
  print("KEY>", line[:400], flush=True)
254
- print("VERDICT PHASE4_SESSION rc", rc, "steps", start_step, "->", pl["stop_after_steps"],
255
- "tokens_at_stop", pl["tokens_at_stop"], flush=True)
256
- sh(["bash", "-c", "df -h /kaggle/working | tail -1"], "disk_after")
 
 
257
  # The session's own wall clock is the measurement that sharpens every later plan.
258
  raise SystemExit(rc)
259
 
 
19
 
20
  import json
21
  import os
22
+ import signal
23
  import subprocess
24
  import sys
25
  import time
 
64
  # child's traceback arrives in the same stream instead of a separate tail that hides it (E-023).
65
  import threading
66
  p = subprocess.Popen(argv, cwd=WORK, stdout=subprocess.PIPE, stderr=subprocess.STDOUT,
67
+ text=True, bufsize=1, env=e, start_new_session=True)
68
  state = {"killed": False}
69
 
70
  def _kill():
71
+ # torchrun is not the process tree. SIGKILLing only the child leaves both ranks alive on the two
72
+ # T4s, still stepping and still rolling latest.json, while the next session is launched into the
73
+ # same container -- so the whole group goes, politely first (review finding).
74
  state["killed"] = True
75
  try:
76
+ os.killpg(os.getpgid(p.pid), signal.SIGTERM)
77
+ except Exception:
78
  p.kill()
79
+
80
+ def _hardkill():
81
+ try:
82
+ os.killpg(os.getpgid(p.pid), signal.SIGKILL)
83
  except Exception:
84
+ p.kill()
85
+ p.kill()
86
 
87
  timer = threading.Timer(timeout, _kill) if timeout else None
88
  if timer:
89
  timer.daemon = True
90
  timer.start()
91
+ hard = threading.Timer(timeout + 120, _hardkill)
92
+ hard.daemon = True
93
+ hard.start()
94
  lines = []
95
  for line in (p.stdout or []):
96
  lines.append(line.rstrip("\n"))
 
98
  rc = p.wait()
99
  if timer:
100
  timer.cancel()
101
+ hard.cancel()
102
  if state["killed"]:
103
  rc = -9
104
  print(f"{label}_RC {rc}{' TIMEOUT' if state['killed'] else ''} seconds "
 
121
 
122
 
123
  def plan(start_step, free_gb):
124
+ """How far this session can train and still end on a checkpoint the Hub has verified.
125
+
126
+ The segment has to stop where a checkpoint is pushed: the 381-step grid, or the horizon. Running to the
127
+ horizon is preferred whenever it fits, because 3,814 is not a multiple of 381 -- stopping at the last
128
+ grid point leaves four steps orphaned and demands one more session whose only budget test is "can you
129
+ afford a whole 2.97 h interval", which the run could fail forever while its last seven minutes of work
130
+ stayed unstarted (review finding).
131
+ """
132
+ sec_per_step = TOKENS_PER_STEP / PLANNING_TOK_PER_S
133
  # The reserve comes off the total, not off one term: `min(session, quota - reserve)` let an 11 h
134
  # session plan 11 h of training and leave nothing for the code fetch, the 2.3 GB mix download, the
135
  # push/verify cycles and the end-of-run evaluation (review finding).
136
  usable = max(0.0, min(SESSION_GPU_HOURS, QUOTA_LEFT_HOURS) - RESERVE_HOURS)
137
+ steps_that_fit = int(usable * 3600.0 // sec_per_step)
138
+ remaining = HORIZON_STEPS - start_step
139
+ if steps_that_fit >= remaining:
140
+ stop, planned = 0, remaining # the trainer runs to its own horizon
141
+ else:
142
+ n = max(0, steps_that_fit // CKPT_EVERY)
143
+ stop, planned = start_step + n * CKPT_EVERY, n * CKPT_EVERY
144
+ end = HORIZON_STEPS if stop == 0 else stop
145
  return {"start_step": start_step, "horizon_steps": HORIZON_STEPS, "ckpt_every": CKPT_EVERY,
146
  "tokens_per_step": TOKENS_PER_STEP, "planning_tok_per_s": PLANNING_TOK_PER_S,
147
+ "hours_per_checkpoint_interval": round(CKPT_EVERY * sec_per_step / 3600.0, 3),
148
  "session_gpu_hours": SESSION_GPU_HOURS, "quota_left_hours": QUOTA_LEFT_HOURS,
149
  "reserve_hours": RESERVE_HOURS, "usable_hours": round(usable, 3),
150
+ "intervals_this_session": max(1, planned // CKPT_EVERY) if planned else 0,
151
+ "planned_steps": planned, "stop_after_steps": stop,
152
+ "runs_to_horizon_this_session": stop == 0 and planned > 0,
153
+ "tokens_at_stop": end * TOKENS_PER_STEP,
154
+ "pct_of_run_at_stop": round(100.0 * end * TOKENS_PER_STEP / TOKENS, 2),
155
  "free_gb_before_start": free_gb}
156
 
157
 
158
+ def df_free_gb(path=WORK):
159
+ """Free space on the working volume, in GB, or None if `df` could not be read.
160
+
161
+ `df -k`, not `-h`: the human-readable form prints "318G"/"500M"/"1.5T" and float() on those raises,
162
+ which would take the session down before training started.
163
+ """
164
+ rc, out = sh(["df", "-k", path], "df", timeout=60)
165
+ for line in out.splitlines():
166
+ parts = line.split()
167
+ if len(parts) >= 4 and ("overlay" in line or line.startswith("/dev/")):
168
+ try:
169
+ return float(parts[3]) / 1048576.0
170
+ except ValueError:
171
+ return None
172
+ return None
173
+
174
+
175
  def main():
176
  os.chdir(WORK)
177
  sys.path.insert(0, WORK)
 
182
  "shard_dataset.py": ("train/shard_dataset.py",
183
  "adbcd96fbab9505a5e8dec4d2e83d00366a73f702786bffbde430195cf6f0471"),
184
  "hubckpt.py": ("train/hubckpt.py",
185
+ "568c31b300906cb8d78a59ef890baa9731b18e064e3d69b0eaea4b45c546f23e"),
186
  "train_ounce100m.py": ("train/train_ounce100m.py",
187
+ "c4c178c0ac17c68de1521d85612a0dc604e6d4748dfdc8bd3e374512bc90ce96"),
188
  })
189
 
190
  import ounce100m_credentials
191
  print("creds:", json.dumps(ounce100m_credentials.install(verify=True)), flush=True)
192
+ sh(["bash", "-c", "nvidia-smi --query-gpu=name,memory.used,memory.total "
193
+ "--format=csv,noheader; free -g | head -2"], "env")
194
+ # The mix is 2.3 GB and it lands during this session, so the pre-download floor is the run's floor plus
195
+ # the pull. Measuring again after it, rather than trusting the arithmetic, is the point.
196
+ free_gb = df_free_gb()
 
 
 
 
 
 
 
 
 
197
  if free_gb is None:
198
+ print("VERDICT REFUSED_TO_START: could not read the free space on /kaggle/working from `df` -- "
199
+ "refusing to guess whether the mix and two checkpoints will fit", flush=True)
200
  raise SystemExit(10)
201
+ if free_gb < MIN_FREE_GB + 2.3:
202
+ print(f"VERDICT REFUSED_TO_START: {free_gb:.1f} GB free on /kaggle/working before the pull; this "
203
+ f"run needs {MIN_FREE_GB} GB free once the 2.3 GB mix has landed (~1.7 GB per checkpoint and "
204
+ "the two unpacked resume directories, on top of it). A Kaggle instance starts near 19.5 GB, "
205
+ "so a low number here means something is left over from a previous session on this box.",
206
+ flush=True)
207
  raise SystemExit(10)
208
+ print(f"disk: {free_gb:.1f} GB free before the mix lands", flush=True)
209
 
210
  # The Hub pointer is the only truth about where the run is (never local state -- §3.13).
211
  import hubckpt
 
223
  print("VERDICT RUN_ALREADY_COMPLETE at step", start_step, flush=True)
224
  return
225
 
226
+ # Plan before downloading anything: an unusable budget should refuse in seconds, not after pulling
227
+ # 2.3 GB onto a billed GPU session.
228
+ pl = plan(start_step, free_gb)
229
+ print("PLAN_JSON", json.dumps(pl, indent=1), flush=True)
230
+ if not pl["planned_steps"]:
231
+ print(f"VERDICT REFUSED_TO_START: {pl['usable_hours']} usable hours will not reach one "
232
+ f"{pl['hours_per_checkpoint_interval']} h checkpoint interval ({pl['start_step']} → "
233
+ f"{HORIZON_STEPS} left). Launching anyway would burn quota on work that cannot be resumed "
234
+ "from. Raise SESSION_GPU_HOURS or wait for quota.", flush=True)
235
+ raise SystemExit(8)
236
+
237
  # Data: the published mix, fetched by anonymous read so this works on any instance. The checks run
238
  # here rather than inside the downloader so a quoting mistake cannot turn into a silent pass.
239
  rc, _ = sh([sys.executable, "-c",
240
  "import sys; sys.path.insert(0, '/kaggle/working');\n"
241
  "from huggingface_hub import snapshot_download\n"
242
  f"p = snapshot_download(repo_id='{MIX_REPO}', repo_type='dataset',\n"
243
+ " local_dir='/kaggle/working/mixroot', max_workers=4)\n"
244
  "print('downloaded to', p)\n"],
245
  "fetch_mix", 3600)
246
  if rc != 0:
247
  raise SystemExit(7)
248
+ free_gb = df_free_gb()
249
+ if free_gb is not None and free_gb < MIN_FREE_GB:
250
+ print(f"VERDICT REFUSED_TO_START: only {free_gb:.1f} GB free after the mix landed, below the "
251
+ f"{MIN_FREE_GB} GB the checkpoints and the two resume directories need", flush=True)
252
+ raise SystemExit(10)
253
  man = json.load(open(os.path.join(WORK, "mixroot", "manifest.json")))
254
  windows = int(man["total_tokens"]) // SEQ_LEN
255
  needed_windows = HORIZON_STEPS * MICRO_BATCH * ACCUM * WORLD
256
+ if len(man.get("shards") or []) != int(man.get("n_shards") or -1):
257
+ print(f"VERDICT REFUSED_TO_START: the published manifest lists {len(man.get('shards') or [])} "
258
+ f"shard records but claims n_shards={man.get('n_shards')}", flush=True)
259
+ raise SystemExit(9)
260
  print(f"mix: {man['total_tokens']:,} train tokens in {man['n_shards']} shards, "
261
  f">= {man['distinct_sources_per_shard_min']} sources per shard, "
262
  f"{windows:,} windows of {SEQ_LEN}; the horizon needs {needed_windows:,} "
 
270
  "implicitly by launching anyway.", flush=True)
271
  raise SystemExit(9)
272
 
 
 
 
 
 
 
 
 
273
  if os.environ.get("PHASE4_PREP_ONLY") == "1":
274
  # Everything a session can get wrong before it bills GPU time has now been checked: the pinned
275
  # code hashes, the credentials resolve, the checkpoint pointer is readable, the published mix is
276
  # downloadable and large enough, and the plan lands on a checkpoint boundary. A CPU instance can
277
  # run all of that for free (E-031's arithmetic is exactly the kind of thing to find there).
278
  print("VERDICT PREP_ONLY_OK stop_after_steps", pl["stop_after_steps"],
279
+ "planned_steps", pl["planned_steps"], flush=True)
280
  return
281
 
282
  argv = [sys.executable, "-u", "train_ounce100m.py",
 
288
  "--log-every", "20", "--val-tokens", "2000000",
289
  "--resume", "auto", "--stop-after-steps", str(pl["stop_after_steps"])]
290
  rc, out = sh(["torchrun", "--nproc_per_node=2"] + argv, "TRAIN",
291
+ # Counted from the session's own start, not the child's: the ceiling is on the container,
292
+ # and by now the code fetch and the 2.3 GB pull have already spent part of the reserve.
293
+ # Killing at ceiling-minus-15-min takes the whole process group with it (see sh) instead
294
+ # of being caught mid-checkpoint by a platform kill that leaves both ranks running.
295
+ timeout=max(600, int(SESSION_GPU_HOURS * 3600 - (time.time() - T0) - 900)))
296
  for line in out.splitlines():
297
  if line.startswith(("CKPT ", "RUN_JSON ", "resume from", "auto-resume", "precision:",
298
  "params:", "validation loss", "segment boundary")):
299
  print("KEY>", line[:400], flush=True)
300
+ end_step = HORIZON_STEPS if pl["stop_after_steps"] == 0 else pl["stop_after_steps"]
301
+ print("VERDICT PHASE4_SESSION rc", rc, "steps", start_step, "->", end_step,
302
+ "tokens_at_stop", pl["tokens_at_stop"], "session_seconds",
303
+ round(time.time() - T0, 1), flush=True)
304
+ print(f"disk at exit: {df_free_gb()} GB free", flush=True)
305
  # The session's own wall clock is the measurement that sharpens every later plan.
306
  raise SystemExit(rc)
307