codemaivanngu commited on
Commit
45ffa08
·
verified ·
1 Parent(s): 2fd9c15

Expose pinned 64-request per GPU migration helpers

Browse files
Files changed (2) hide show
  1. eval-queue-64-5f912bc.py +445 -0
  2. migrate-eval-64-5f912bc.py +181 -0
eval-queue-64-5f912bc.py ADDED
@@ -0,0 +1,445 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ #!/usr/bin/env python3
2
+ """Two independent GPU workers; tier barriers; bounded resumable internal eval."""
3
+ import argparse
4
+ import concurrent.futures as cf
5
+ from collections import deque
6
+ import contextlib
7
+ import fcntl
8
+ import json
9
+ import os
10
+ from pathlib import Path
11
+ import signal
12
+ import statistics
13
+ import subprocess
14
+ import sys
15
+ import tempfile
16
+ import time
17
+ import threading
18
+ import urllib.request
19
+ import contract_eval as E
20
+ import queue_data as D
21
+ from context_check import CONTEXT_LENGTH
22
+
23
+
24
+ class Deadline(Exception): pass
25
+
26
+
27
+ def atomic_json(path,value):
28
+ tmp=path.with_suffix(path.suffix+".tmp")
29
+ tmp.write_bytes(E.encoded(value)); os.replace(tmp,path)
30
+
31
+
32
+ @contextlib.contextmanager
33
+ def locked(path,blocking=True):
34
+ with path.open("a+") as f:
35
+ try: fcntl.flock(f,fcntl.LOCK_EX | (0 if blocking else fcntl.LOCK_NB))
36
+ except BlockingIOError:
37
+ yield None; return
38
+ try: yield f
39
+ finally: fcntl.flock(f,fcntl.LOCK_UN)
40
+
41
+
42
+ def journal(path):
43
+ rows={}
44
+ if not path.exists(): return rows
45
+ end=0
46
+ with path.open("rb") as f:
47
+ while line:=f.readline():
48
+ try: row=json.loads(line)
49
+ except ValueError:
50
+ if f.read(1): raise ValueError("corrupt interior journal line")
51
+ if line.endswith(b"\n"): raise ValueError("corrupt complete journal line")
52
+ # Preserve only an interrupted final write, never silently lose a full row.
53
+ tail=path.with_name(path.name+".interrupted-tail")
54
+ if tail.exists(): raise ValueError("inspect previous interrupted tail first")
55
+ tail.write_bytes(line)
56
+ with path.open("r+b") as repair: repair.truncate(end)
57
+ break
58
+ if row["id"] in rows: raise ValueError("duplicate journal ID")
59
+ rows[row["id"]]=row; end=f.tell()
60
+ if rows and path.stat().st_size:
61
+ with path.open("rb") as f:
62
+ f.seek(-1,2); missing_newline=f.read(1)!=b"\n"
63
+ if missing_newline:
64
+ with path.open("ab") as f: f.write(b"\n"); f.flush(); os.fsync(f.fileno())
65
+ return rows
66
+
67
+
68
+ def append(path,row):
69
+ with path.open("ab") as f:
70
+ f.write(E.encoded(row)+b"\n"); f.flush(); os.fsync(f.fileno())
71
+
72
+
73
+ def stop_group(process):
74
+ # Only the process group created by this queue; never pkill/ray stop.
75
+ try: os.killpg(process.pid,signal.SIGTERM)
76
+ except ProcessLookupError: return
77
+ try: process.wait(timeout=3)
78
+ except subprocess.TimeoutExpired: pass
79
+ try: os.killpg(process.pid,signal.SIGKILL)
80
+ except ProcessLookupError: pass
81
+ try: process.wait(timeout=3)
82
+ except subprocess.TimeoutExpired: pass
83
+
84
+
85
+ def score_job(python,job,deadline):
86
+ remaining=deadline-time.time()
87
+ if remaining<=0: raise Deadline()
88
+ timeout=min(120.,remaining)
89
+ env={"PATH":"/usr/bin:/bin","OPENBLAS_NUM_THREADS":"1","OMP_NUM_THREADS":"1",
90
+ "PYTHONDONTWRITEBYTECODE":"1","PYTHONUNBUFFERED":"1","LANG":"C.UTF-8"}
91
+ with tempfile.TemporaryDirectory(prefix="simct-score-") as work, tempfile.TemporaryFile() as out, tempfile.TemporaryFile() as err:
92
+ env["HOME"]=work; env["TMPDIR"]=work
93
+ p=subprocess.Popen([python,str(E.HERE/"internal_worker.py")],stdin=subprocess.PIPE,
94
+ stdout=out,stderr=err,cwd=work,env=env,start_new_session=True)
95
+ try:
96
+ try: p.communicate(E.encoded(job),timeout=timeout)
97
+ except subprocess.TimeoutExpired:
98
+ if timeout<120: raise Deadline()
99
+ return {"passed":False,"timeout":True}
100
+ out.seek(0); raw=out.read(1024**2)
101
+ if p.returncode:
102
+ raise RuntimeError("scorer runtime error: "+raw.decode(errors="replace")[:1000])
103
+ result=json.loads(raw)
104
+ if type(result.get("passed")) is not bool: raise ValueError("invalid scorer result")
105
+ return result
106
+ finally: stop_group(p)
107
+
108
+
109
+ def preflight(python):
110
+ cases=[("gsm8k",{"gold":"#### 2"},"#### 2","#### 3"),
111
+ ("math500",{"gold":"2"},r"\boxed{2}",r"\boxed{3}"),
112
+ ("mbpp",{"tests":["assert add(2,3)==5"],"setup":""},"def add(a,b): return a+b","def add(a,b): return a-b"),
113
+ ("live-code-bench-v6",{"tests":[{"input":"2 3\n","output":"5\n"}],"fn_name":None},"a,b=map(int,input().split());print(a+b)","print(0)"),
114
+ ("live-code-bench-v6",{"tests":[{"input":"2\n3","output":"5"}],"fn_name":"add"},"class Solution:\n def add(self,a,b): return a+b","class Solution:\n def add(self,a,b): return a-b")]
115
+ evidence=[]
116
+ for benchmark,item,good,bad in cases:
117
+ for correct,text in ((True,good),(False,bad)):
118
+ result=score_job(python,{"profile":D.PROFILE,"benchmark":benchmark,"item":item,"text":text},time.time()+150)
119
+ if result["passed"]!=correct or result.get("timeout"): raise RuntimeError("scorer qualification failed")
120
+ evidence.append({"benchmark":benchmark,"functional":bool(item.get("fn_name")),"correct":correct,**result})
121
+ return evidence
122
+
123
+
124
+ def state_path(root): return root/"state.json"
125
+
126
+
127
+ def read_state(root,plan_hash,plan):
128
+ path=state_path(root)
129
+ if path.exists():
130
+ state=E.read_json(path)
131
+ if state["plan_sha256"]!=plan_hash: raise ValueError("queue plan changed")
132
+ return state
133
+ start=time.time()
134
+ state={"plan_sha256":plan_hash,"started":start,"deadline":start+plan["hours"]*3600,
135
+ "admit_until":start+plan["admit_hours"]*3600,"jobs":{},"durations":[]}
136
+ atomic_json(path,state);return state
137
+
138
+
139
+ def next_jobs(plan,state,now):
140
+ pending=[j for j in plan["jobs"] if state["jobs"].get(j["id"],{}).get("status")!="completed"]
141
+ if not pending or now>=state["admit_until"]: return []
142
+ if any(x.get("status")=="failed" for x in state["jobs"].values()):
143
+ raise RuntimeError("queue has a failed job; inspect error and use retry after fixing it")
144
+ tier=min(j["tier"] for j in pending)
145
+ estimate=max(state["durations"],default=0.)
146
+ if now+estimate>state["deadline"]: return []
147
+ return [j for j in pending if j["tier"]==tier]
148
+
149
+
150
+ def cell_contract(plan_hash,job,data,benchmark,seed,server):
151
+ return {"plan_sha256":plan_hash,"checkpoint_sha256":job["checkpoint"]["sha256"],
152
+ "data_sha256":data[benchmark]["sha256"],"benchmark":benchmark,"seed":seed,
153
+ "profile":D.PROFILE,"server":server}
154
+
155
+
156
+ def generate_one(base,item,benchmark,seed,deadline):
157
+ if time.time()>=deadline: raise Deadline()
158
+ payload=E.generation_payload("eval-gemma",item,benchmark,seed)
159
+ request=urllib.request.Request(base+"/v1/chat/completions",data=E.encoded(payload),headers={"Content-Type":"application/json"})
160
+ opener=urllib.request.build_opener(urllib.request.ProxyHandler({}))
161
+ with opener.open(request,timeout=min(600.,max(.1,deadline-time.time()))) as response:
162
+ result=json.load(response)
163
+ E.validate_response(result)
164
+ return {"id":item["id"],"seed":seed,"request_sha256":E.digest(E.encoded(payload)),"response":result}
165
+
166
+
167
+ def run_cell(root,plan,plan_hash,job,benchmark,seed,base,server,args,deadline):
168
+ cell=root/"cells"/job["id"]/benchmark/str(seed);cell.mkdir(parents=True,exist_ok=True)
169
+ contract=cell_contract(plan_hash,job,plan["data"],benchmark,seed,server)
170
+ manifest=cell/"contract.json"
171
+ if manifest.exists():
172
+ if E.read_json(manifest)!=contract: raise ValueError("cell resume contract mismatch")
173
+ else: E.write_new(manifest,contract)
174
+ dataset=E.read_json(plan["data"][benchmark]["path"])
175
+ items={x["id"]:x for x in dataset["items"]}
176
+ if (cell/"metrics.json").exists():
177
+ complete=E.read_json(cell/"metrics.json")
178
+ for name in ("responses","scores"):
179
+ if E.file_hash(cell/(name+".jsonl"))!=complete[name+"_sha256"]:
180
+ raise ValueError("completed journal changed; no repair permitted")
181
+ responses=journal(cell/"responses.jsonl"); scores=journal(cell/"scores.jsonl")
182
+ if not set(scores)<=set(responses)<=set(items): raise ValueError("orphan/unknown results")
183
+ for key,row in responses.items():
184
+ expected=E.generation_payload("eval-gemma",items[key],benchmark,seed)
185
+ if row["seed"]!=seed or row["request_sha256"]!=E.digest(E.encoded(expected)): raise ValueError("response request mismatch")
186
+ E.validate_response(row["response"])
187
+ for key,row in scores.items():
188
+ if row["response_sha256"]!=E.digest(E.encoded(responses[key])) or type(row.get("passed")) is not bool:
189
+ raise ValueError("score does not match response")
190
+ if (cell/"metrics.json").exists():
191
+ m=E.read_json(cell/"metrics.json")
192
+ if len(scores)!=len(items) or m["responses_sha256"]!=E.file_hash(cell/"responses.jsonl") or m["scores_sha256"]!=E.file_hash(cell/"scores.jsonl"):
193
+ raise ValueError("completed cell integrity mismatch")
194
+ return m
195
+ started=time.time()
196
+ previous_seconds=E.read_json(cell/"timing.json").get("seconds",0.) if (cell/"timing.json").exists() else 0.
197
+ pending=iter([x for x in items.values() if x["id"] not in scores])
198
+ futures={}
199
+ def process_item(item):
200
+ # Data decoding remains in parent, one problem per scoring process.
201
+ row=responses.get(item["id"])
202
+ if row is None: row=generate_one(base,item,benchmark,seed,deadline)
203
+ return row
204
+ score_buffer=getattr(args,"score_buffer",64)
205
+ if score_buffer<args.concurrency: raise ValueError("score buffer must cover generation concurrency")
206
+ backlog=deque()
207
+ def grade_item(item):
208
+ # Decode potentially large/private tests only inside an active scorer slot.
209
+ row=responses[item["id"]]
210
+ payload={"profile":D.PROFILE,"benchmark":benchmark,"item":D.score_item(benchmark,item),
211
+ "text":row["response"]["choices"][0]["message"]["content"]}
212
+ return score_job(args.score_python,payload,deadline)
213
+ print("PIPELINE",json.dumps({"generation_concurrency":args.concurrency,
214
+ "score_workers":args.score_workers,"score_buffer":score_buffer}),flush=True)
215
+ try:
216
+ with cf.ThreadPoolExecutor(max_workers=args.concurrency) as generation, cf.ThreadPoolExecutor(max_workers=args.score_workers) as grading:
217
+ exhausted=False
218
+ while futures or backlog or not exhausted:
219
+ if time.time()>=deadline: raise Deadline()
220
+ scoring=sum(kind=="score" for kind,_ in futures.values())
221
+ while backlog and scoring<args.score_workers:
222
+ item=backlog.popleft()
223
+ futures[grading.submit(grade_item,item)]=( "score",item)
224
+ scoring+=1
225
+ generating=sum(kind=="generation" for kind,_ in futures.values())
226
+ # Reserve backlog capacity for in-flight generation; scorer backlog
227
+ # no longer consumes generation slots. Backpressure remains bounded.
228
+ while not exhausted and generating<args.concurrency and len(backlog)+generating<score_buffer:
229
+ try: item=next(pending)
230
+ except StopIteration: exhausted=True;break
231
+ futures[generation.submit(process_item,item)]=("generation",item)
232
+ generating+=1
233
+ if not futures: break
234
+ done,_=cf.wait(futures,timeout=min(1.,max(.1,deadline-time.time())),return_when=cf.FIRST_COMPLETED)
235
+ for future in done:
236
+ kind,item=futures.pop(future);value=future.result()
237
+ if kind=="generation":
238
+ if item["id"] not in responses:
239
+ append(cell/"responses.jsonl",value);responses[item["id"]]=value
240
+ backlog.append(item)
241
+ else:
242
+ value.update(id=item["id"],response_sha256=E.digest(E.encoded(responses[item["id"]])))
243
+ append(cell/"scores.jsonl",value);scores[item["id"]]=value
244
+ if len(scores)%25==0: print(f"PROGRESS {job['id']} {benchmark} seed={seed} {len(scores)}/{len(items)}",flush=True)
245
+ finally:
246
+ atomic_json(cell/"timing.json",{"seconds":previous_seconds+time.time()-started})
247
+ if len(scores)!=len(items): raise ValueError("incomplete cell")
248
+ metrics={"status":"completed","contract":contract,"count":len(items),
249
+ "score":sum(x["passed"] for x in scores.values())/len(items),
250
+ "truncated_fraction":sum(x["response"]["choices"][0]["finish_reason"]=="length" for x in responses.values())/len(items),
251
+ "seconds":previous_seconds+time.time()-started,
252
+ "responses_sha256":E.file_hash(cell/"responses.jsonl"),"scores_sha256":E.file_hash(cell/"scores.jsonl")}
253
+ E.write_new(cell/"metrics.json",metrics)
254
+ print("CELL_COMPLETE",job["id"],benchmark,seed,metrics["score"],flush=True)
255
+ return metrics
256
+
257
+
258
+ def run_checkpoint(root,plan,plan_hash,job,args,deadline):
259
+ actual=E.checkpoint_identity(job["checkpoint"]["path"])
260
+ if time.time()>=deadline: raise Deadline("wall budget during checkpoint verification")
261
+ if actual!=job["checkpoint"]: raise ValueError("checkpoint changed after plan")
262
+ port=31000+1000*args.gpu;base=f"http://127.0.0.1:{port}"
263
+ import socket
264
+ with socket.socket() as sock:
265
+ if sock.connect_ex(("127.0.0.1",port))==0: raise RuntimeError("server port already occupied")
266
+ env=dict(os.environ,CUDA_VISIBLE_DEVICES=str(args.gpu),HF_HUB_OFFLINE="1",TRANSFORMERS_OFFLINE="1",
267
+ HF_DATASETS_OFFLINE="1",OMP_NUM_THREADS="4",TOKENIZERS_PARALLELISM="false",PYTHONUNBUFFERED="1",
268
+ PYTHONPATH=str(E.HERE.parents[1]/"experiments/modal/vendor")+":"+str(E.HERE.parents[1]))
269
+ for name in ("HTTP_PROXY","HTTPS_PROXY","ALL_PROXY","http_proxy","https_proxy","all_proxy","WANDB_API_KEY","HF_TOKEN"):
270
+ env.pop(name,None)
271
+ env.pop("SGLANG_ALLOW_OVERWRITE_LONGER_CONTEXT_LEN",None)
272
+ subprocess.run(["bash",str(E.HERE.parents[1]/"experiments/runai/python-b200-host.sh"),
273
+ str(E.HERE/"context_check.py"),"--plan",str(args.plan.resolve()),
274
+ "--checkpoint",actual["path"]],env=dict(env,CUDA_VISIBLE_DEVICES=""),
275
+ check=True,timeout=max(.1,min(300.,deadline-time.time())))
276
+ if time.time()>=deadline: raise Deadline("wall budget during context qualification")
277
+ log=root/(job["id"]+f"-gpu{args.gpu}-server.log")
278
+ command=["bash",str(E.HERE.parents[1]/"experiments/runai/python-b200-host.sh"),"-m","sglang.launch_server",
279
+ "--model-path",actual["path"],"--served-model-name","eval-gemma","--host","127.0.0.1","--port",str(port),
280
+ "--tp-size","1","--mem-fraction-static","0.8","--context-length",str(CONTEXT_LENGTH),
281
+ "--attention-backend","triton","--disable-cuda-graph","--random-seed","42"]
282
+ with log.open("ab") as output:
283
+ process=subprocess.Popen(command,env=env,stdout=output,stderr=subprocess.STDOUT,start_new_session=True)
284
+ watchdog=threading.Timer(max(0.,deadline-time.time()),stop_group,args=(process,))
285
+ watchdog.daemon=True; watchdog.start()
286
+ try:
287
+ startup=min(deadline,time.time()+900)
288
+ while True:
289
+ if time.time()>=deadline: raise Deadline("wall budget")
290
+ if time.time()>=startup: raise RuntimeError("server readiness exceeded 900 seconds; see "+str(log))
291
+ if process.poll() is not None: raise RuntimeError("SGLang exited; see "+str(log))
292
+ try:
293
+ # Short health timeout; identity verification follows readiness.
294
+ opener=urllib.request.build_opener(urllib.request.ProxyHandler({}))
295
+ with opener.open(base+"/health",timeout=2): pass
296
+ break
297
+ except (OSError,ValueError): time.sleep(2)
298
+ server=E.verify_server(base,actual["path"],"eval-gemma")
299
+ for seed in plan["seeds"]:
300
+ for benchmark in E.CAPS:
301
+ run_cell(root,plan,plan_hash,job,benchmark,seed,base,server,args,deadline)
302
+ if E.verify_server(base,actual["path"],"eval-gemma")!=server: raise ValueError("server identity drift")
303
+ finally:
304
+ watchdog.cancel(); stop_group(process)
305
+
306
+
307
+ def worker(args):
308
+ root=args.plan.resolve().parent;plan=E.read_json(args.plan);plan_hash=E.file_hash(args.plan)
309
+ if plan["profile"]!=D.PROFILE or plan["source"]!=D.script_hashes(): raise ValueError("queue source/profile changed")
310
+ for data in plan["data"].values():
311
+ if E.file_hash(data["path"])!=data["sha256"]: raise ValueError("prepared data changed")
312
+ qualification=preflight(args.score_python)
313
+ with locked(root/"state.lock"):
314
+ state=read_state(root,plan_hash,plan)
315
+ if "qualification" in state and state["qualification"]!=qualification: raise ValueError("grader version drift")
316
+ state["qualification"]=qualification;atomic_json(state_path(root),state)
317
+ gpu_lock=Path(tempfile.gettempdir())/f"simct-eval-gpu{args.gpu}.lock"
318
+ with locked(gpu_lock,blocking=False) as own:
319
+ if own is None: raise RuntimeError("another eval worker owns this GPU")
320
+ while True:
321
+ with locked(root/"state.lock"):
322
+ state=read_state(root,plan_hash,plan)
323
+ candidates=next_jobs(plan,state,time.time())
324
+ if not candidates:
325
+ print("QUEUE_STOP admission/deadline/complete",flush=True);return
326
+ usage=subprocess.check_output(["nvidia-smi","-i",str(args.gpu),"--query-gpu=memory.used","--format=csv,noheader,nounits"],text=True).strip()
327
+ if int(usage)>=1024:
328
+ print("WAIT_GPU",args.gpu,usage,flush=True);time.sleep(15);continue
329
+ claimed=False
330
+ for job in candidates:
331
+ with locked(root/(job["id"]+".lock"),blocking=False) as lock:
332
+ if lock is None: continue
333
+ with locked(root/"state.lock"):
334
+ state=read_state(root,plan_hash,plan)
335
+ if job["id"] not in {j["id"] for j in next_jobs(plan,state,time.time())}: continue
336
+ state["jobs"][job["id"]]={"status":"running","gpu":args.gpu,"started":time.time()}
337
+ atomic_json(state_path(root),state)
338
+ started=time.time();claimed=True
339
+ try:
340
+ run_checkpoint(root,plan,plan_hash,job,args,state["deadline"])
341
+ result={"status":"completed","seconds":time.time()-started}
342
+ except Deadline as exc:
343
+ result={"status":"partial","reason":str(exc) or "wall budget"}
344
+ except Exception as exc:
345
+ result=({"status":"partial","reason":"wall budget"} if time.time()>=state["deadline"] else
346
+ {"status":"failed","error":type(exc).__name__+": "+str(exc)})
347
+ with locked(root/"state.lock"):
348
+ state=read_state(root,plan_hash,plan);state["jobs"][job["id"]]=result
349
+ if result["status"]=="completed": state["durations"].append(result["seconds"])
350
+ atomic_json(state_path(root),state)
351
+ print("JOB",job["id"],json.dumps(result),flush=True)
352
+ if result["status"]=="failed": raise RuntimeError(result["error"])
353
+ if result["status"]!="completed": return
354
+ break
355
+ if not claimed: time.sleep(5)
356
+
357
+
358
+ def summarize(args):
359
+ root=args.plan.resolve().parent;plan=E.read_json(args.plan);ph=E.file_hash(args.plan)
360
+ output={"profile":plan["profile"],"seeds":plan["seeds"],"checkpoints":{},"scope":plan["protocol"]["scope"]}
361
+ for job in plan["jobs"]:
362
+ benchmarks={}
363
+ for benchmark in E.CAPS:
364
+ paths=[root/"cells"/job["id"]/benchmark/str(seed)/"metrics.json" for seed in plan["seeds"]]
365
+ if not all(p.exists() for p in paths): continue
366
+ rows=[E.read_json(p) for p in paths]
367
+ for seed,row,path in zip(plan["seeds"],rows,paths):
368
+ c=row["contract"]
369
+ if c["data_sha256"]!=plan["data"][benchmark]["sha256"] or c["benchmark"]!=benchmark or c["profile"]!=D.PROFILE or row["count"]!=plan["data"][benchmark]["count"]:
370
+ raise ValueError("summary data contract mismatch")
371
+ for name in ("responses","scores"):
372
+ if E.file_hash(path.parent/(name+".jsonl"))!=row[name+"_sha256"]:
373
+ raise ValueError("summary journal integrity mismatch")
374
+ if c["plan_sha256"]!=ph or c["seed"]!=seed or c["checkpoint_sha256"]!=job["checkpoint"]["sha256"] or row["status"]!="completed":
375
+ raise ValueError("summary cell contract mismatch")
376
+ values=[x["score"] for x in rows]
377
+ benchmarks[benchmark]={"mean":statistics.mean(values),"sample_std":statistics.stdev(values),"scores":values}
378
+ result={"benchmarks":benchmarks,"complete":len(benchmarks)==4}
379
+ if result["complete"]: result["average"]=statistics.mean(x["mean"] for x in benchmarks.values())
380
+ output["checkpoints"][job["id"]]=result
381
+ atomic_json(root/"summary-3seeds.json",output);print(json.dumps(output,indent=2))
382
+
383
+
384
+ def recover_startup(args):
385
+ """Fork only a zero-result failed-startup queue; retain its original clock."""
386
+ old=args.from_plan.resolve();oldroot=old.parent
387
+ with locked(oldroot/"state.lock"), contextlib.ExitStack() as locks:
388
+ plan=E.read_json(old);state=E.read_json(oldroot/"state.json")
389
+ if state["plan_sha256"]!=E.file_hash(old) or plan["profile"]!=D.PROFILE:
390
+ raise ValueError("old plan identity mismatch")
391
+ if state["durations"] or not state["jobs"] or any(x["status"]!="failed" for x in state["jobs"].values()):
392
+ raise ValueError("recovery requires exclusively failed startup jobs")
393
+ if list(oldroot.glob("cells/**/*.json*")):
394
+ raise ValueError("results exist; use an audited migration instead")
395
+ for job in plan["jobs"]:
396
+ if locks.enter_context(locked(oldroot/(job["id"]+".lock"),blocking=False)) is None:
397
+ raise ValueError("old worker still owns a checkpoint")
398
+ for data in plan["data"].values():
399
+ if E.file_hash(data["path"])!=data["sha256"]: raise ValueError("old data changed")
400
+ root=args.out.resolve();root.mkdir(parents=True,exist_ok=False)
401
+ plan["source"]=D.script_hashes();plan["protocol"]["context_length"]=CONTEXT_LENGTH
402
+ plan["recovery"]={"from_plan":str(old),"sha256":E.file_hash(old),"reason":"Gemma2 native context startup fix; original clock retained"}
403
+ E.write_new(root/"plan.json",plan)
404
+ state.update(plan_sha256=E.file_hash(root/"plan.json"),jobs={},durations=[])
405
+ E.write_new(root/"state.json",state)
406
+ print("RECOVERED_PLAN="+str(root/"plan.json"))
407
+ print("REMAINING_HOURS="+str(max(0.,state["deadline"]-time.time())/3600))
408
+ print("Original deadline retained; no GPU started.")
409
+
410
+
411
+ def main():
412
+ p=argparse.ArgumentParser(description=__doc__);sub=p.add_subparsers(dest="cmd",required=True)
413
+ q=sub.add_parser("prepare")
414
+ q.add_argument("--out",type=Path,required=True)
415
+ q.add_argument("--shared",type=Path,default=Path("/workspace/storage-shared/nlp/tungks"))
416
+ q.add_argument("--company-data",type=Path,default=Path("/workspace/storage-shared/nlp/hadv17/lm-eval/benchmarks"))
417
+ q.add_argument("--proxy",default="http://10.30.154.118:80")
418
+ q.add_argument("--hours",type=float,default=20.)
419
+ q.set_defaults(func=D.prepare)
420
+ q=sub.add_parser("preflight");q.add_argument("--score-python",default="/usr/bin/python3.12");q.set_defaults(func=lambda a:print(json.dumps(preflight(a.score_python),indent=2)))
421
+ q=sub.add_parser("worker")
422
+ q.add_argument("--plan",type=Path,required=True);q.add_argument("--gpu",type=int,choices=range(8),required=True)
423
+ q.add_argument("--score-python",default="/usr/bin/python3.12")
424
+ q.add_argument("--concurrency",type=int,default=8);q.add_argument("--score-workers",type=int,default=2)
425
+ q.add_argument("--score-buffer",type=int,default=64,help="Maximum waiting plus in-flight generated items; tests decoded only by active scorers")
426
+ q.add_argument("--internal-code-execution",action="store_true",required=True,
427
+ help="Explicit company-internal profile; resource limits are not an OS sandbox")
428
+ q.set_defaults(func=worker)
429
+ q=sub.add_parser("summarize");q.add_argument("--plan",type=Path,required=True);q.set_defaults(func=summarize)
430
+ q=sub.add_parser("recover-startup");q.add_argument("--from-plan",type=Path,required=True);q.add_argument("--out",type=Path,required=True);q.set_defaults(func=recover_startup)
431
+ q=sub.add_parser("retry");q.add_argument("--plan",type=Path,required=True)
432
+ def retry(a):
433
+ root=a.plan.resolve().parent
434
+ with locked(root/"state.lock"):
435
+ s=E.read_json(state_path(root))
436
+ for key,value in list(s["jobs"].items()):
437
+ if value["status"]=="failed": del s["jobs"][key]
438
+ atomic_json(state_path(root),s)
439
+ q.set_defaults(func=retry)
440
+ a=p.parse_args()
441
+ if a.cmd=="prepare" and not 1<a.hours<=20: p.error("budget must be >1 and <=20 hours")
442
+ if a.cmd=="worker" and not (1<=a.concurrency<=64 and 1<=a.score_workers<=16 and a.concurrency<=a.score_buffer<=256): p.error("invalid concurrency")
443
+ a.func(a)
444
+
445
+ if __name__=="__main__": main()
migrate-eval-64-5f912bc.py ADDED
@@ -0,0 +1,181 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ #!/usr/bin/env python3
2
+ """Stop exact old workers, fork verified journals, add SimCT; never launch GPUs."""
3
+ import argparse, contextlib, copy, hashlib, json, os, shutil, signal, sys, time
4
+ from pathlib import Path
5
+
6
+
7
+ def proc(pid):
8
+ p=Path('/proc')/str(pid)
9
+ try:
10
+ raw=(p/'stat').read_text().rsplit(')',1)[1].split()
11
+ return dict(pid=int(pid),ppid=int(raw[1]),group=int(raw[2]),start=raw[19],state=raw[0],
12
+ argv=[x.decode() for x in (p/'cmdline').read_bytes().split(b'\0') if x])
13
+ except (OSError, UnicodeError): return None
14
+
15
+
16
+ def stop_workers(plan,queue,pids):
17
+ handles=[]; roots=[]
18
+ try:
19
+ for pid in pids:
20
+ info=proc(pid)
21
+ if info is None: continue
22
+ args=info['argv']
23
+ if str(queue) not in args or 'worker' not in args or '--plan' not in args or args[args.index('--plan')+1]!=str(plan):
24
+ raise ValueError('PID identity mismatch: '+str(pid))
25
+ fd=os.pidfd_open(pid); again=proc(pid)
26
+ if not again or again['start']!=info['start']: raise ValueError('PID changed')
27
+ handles.append(fd); roots.append(info)
28
+ # Freeze only validated workers so they cannot start a new job during teardown.
29
+ for fd in handles: signal.pidfd_send_signal(fd,signal.SIGSTOP)
30
+ table={int(p.name):proc(int(p.name)) for p in Path('/proc').iterdir() if p.name.isdigit()}
31
+ descendants={r['pid'] for r in roots}
32
+ while True:
33
+ more={pid for pid,x in table.items() if x and x['ppid'] in descendants}
34
+ if more<=descendants: break
35
+ descendants |= more
36
+ children=[]
37
+ for pid in descendants-{r['pid'] for r in roots}:
38
+ before=table[pid]
39
+ try: fd=os.pidfd_open(pid)
40
+ except ProcessLookupError: continue
41
+ after=proc(pid)
42
+ if not after or after['start']!=before['start']: os.close(fd);continue
43
+ children.append(fd)
44
+ try:
45
+ for fd in children+handles:
46
+ try: signal.pidfd_send_signal(fd,signal.SIGTERM)
47
+ except ProcessLookupError: pass
48
+ for fd in handles:
49
+ try: signal.pidfd_send_signal(fd,signal.SIGCONT)
50
+ except ProcessLookupError: pass
51
+ limit=time.monotonic()+30
52
+ while time.monotonic()<limit:
53
+ alive=[pid for pid in descendants if (x:=proc(pid)) and x['state']!='Z' and x['start']==table[pid]['start']]
54
+ if not alive: return
55
+ time.sleep(.5)
56
+ raise RuntimeError('Some owned processes remain; no migration performed: '+str(alive))
57
+ finally:
58
+ for fd in children: os.close(fd)
59
+ finally:
60
+ for fd in handles:
61
+ try: signal.pidfd_send_signal(fd,signal.SIGCONT)
62
+ except ProcessLookupError: pass
63
+ os.close(fd)
64
+
65
+
66
+ def main():
67
+ a=argparse.ArgumentParser();a.add_argument('--plan',type=Path,required=True);a.add_argument('--source',type=Path,required=True)
68
+ a.add_argument('--out',type=Path,required=True);a.add_argument('--worker-pids',nargs='+',type=int)
69
+ a.add_argument('--queue-update',type=Path)
70
+ a.add_argument('--simct',type=Path);args=a.parse_args()
71
+ old=args.plan.resolve(strict=True); source=args.source.resolve(strict=True); out=args.out.resolve()
72
+ if out.exists() or out.is_relative_to(old.parent): raise ValueError('Output must be a new sibling directory')
73
+ sys.path.insert(0,str(source/'scripts/evaluation'))
74
+ import eval_queue as Q
75
+ E,D=Q.E,Q.D
76
+ plan=E.read_json(old); oldhash=E.file_hash(old)
77
+ if plan['source']!=D.script_hashes(): raise ValueError('Source differs from running plan')
78
+ expected_count=25 if args.queue_update else 17
79
+ allowed=('sft','atomic','fixed','simct') if args.queue_update else ('sft','atomic','fixed')
80
+ if len(plan['jobs'])!=expected_count or any(j['mode'] not in allowed for j in plan['jobs']): raise ValueError('Unexpected old jobs')
81
+ update_bytes=None
82
+ if args.queue_update:
83
+ update_bytes=args.queue_update.read_bytes()
84
+ if plan['source']['eval_queue.py'] not in ('c64478394072d27fff38b432e6976e900a9b9a6a9065105590257ab34e8fa36d','616acf5249f5b7365b3f93ff7a17546c96b9d71ffdf5306b7981a6079bd0ce3e','af7cdaa0fab262f597babc64a17d5e1861ec2fbb35644612a096785962678a51'): raise ValueError('Unsupported original queue version')
85
+ if hashlib.sha256(update_bytes).hexdigest()!='c4c23d92039237e6a0717f9221f302d4965df98a1e02544e869a3e79987cde6f': raise ValueError('Unsupported queue update')
86
+ initial_state=E.read_json(old.parent/'state.json')
87
+ if time.time()>=initial_state['admit_until']: raise ValueError('Admission expired; workers left untouched')
88
+ for j in plan['jobs']:
89
+ if j['mode'] in D.RUNS and D.RUNS[j['mode']] not in Path(j['checkpoint']['path']).parts: raise ValueError('Wrong historical MP run')
90
+ extra=[]
91
+ if not args.queue_update:
92
+ if args.simct is None: raise ValueError('--simct required')
93
+ summary=E.read_json(args.simct/'run-summary.json')
94
+ if summary['kd_algorithm']!='span_ctkd' or summary['status']!='completed' or summary['optimizer_updates']!=312: raise ValueError('SimCT summary invalid')
95
+ if summary['student']!=next(j['checkpoint']['path'] for j in plan['jobs'] if j['mode']=='sft'): raise ValueError('Different SFT initialization')
96
+ # Hash before stopping to avoid wasting idle GPU time on checkpoint inventory.
97
+ for step in (312,156,80,240,40,200,120,280):
98
+ print('HASH_SIMCT',step,flush=True)
99
+ identity=E.checkpoint_identity(args.simct/f'step{step}')
100
+ tier=next(j['tier'] for j in plan['jobs'] if j['mode']=='atomic' and j['step']==step)
101
+ extra.append(dict(id=f'simct-{step}',mode='simct',step=step,tier=tier,checkpoint=identity))
102
+ for data in plan['data'].values():
103
+ if E.file_hash(data['path'])!=data['sha256']: raise ValueError('Data changed')
104
+ pids=args.worker_pids
105
+ if pids is None:
106
+ pids=[]
107
+ for entry in Path('/proc').iterdir():
108
+ if not entry.name.isdigit(): continue
109
+ info=proc(int(entry.name))
110
+ if info and 'worker' in info['argv'] and '--plan' in info['argv']:
111
+ argv=info['argv']
112
+ if argv[argv.index('--plan')+1]==str(old): pids.append(info['pid'])
113
+ print('MATCHED_WORKERS',pids,flush=True)
114
+ stop_workers(old,source/'scripts/evaluation/eval_queue.py',pids)
115
+ with contextlib.ExitStack() as stack:
116
+ for gpu in (0,1):
117
+ if stack.enter_context(Q.locked(Path(f'/tmp/simct-eval-gpu{gpu}.lock'),blocking=False)) is None: raise ValueError('GPU worker still active')
118
+ stack.enter_context(Q.locked(old.parent/'state.lock'))
119
+ for j in plan['jobs']:
120
+ if stack.enter_context(Q.locked(old.parent/(j['id']+'.lock'),blocking=False)) is None: raise ValueError('Job still active')
121
+ state=E.read_json(old.parent/'state.json')
122
+ if state['plan_sha256']!=oldhash or E.file_hash(old)!=oldhash: raise ValueError('Plan changed')
123
+ if time.time()>=state['admit_until']: raise ValueError('Original admission window expired; no implicit extension')
124
+ out.mkdir()
125
+ receipt={'from_plan':str(old),'old_plan_sha256':oldhash,'source':str(source),'concurrency':16,'original_files':{},'status':'incomplete'}
126
+ E.write_new(out/'migration.json',receipt)
127
+ new=copy.deepcopy(plan);new['jobs']+=extra;new['jobs'].sort(key=lambda j:j['tier'])
128
+ if update_bytes is not None:
129
+ target=out/'source'
130
+ shutil.copytree(source,target,ignore=shutil.ignore_patterns('.git','__pycache__','remote_artifacts'))
131
+ (target/'scripts/evaluation/eval_queue.py').write_bytes(update_bytes)
132
+ new['source']={name:E.file_hash(target/'scripts/evaluation'/name) for name in plan['source']}
133
+ if {k for k in new['source'] if new['source'][k]!=plan['source'][k]}!={'eval_queue.py'}: raise ValueError('Unexpected source changes')
134
+ receipt.update(source=str(target),previous_source=str(source),concurrency_by_gpu={"0":64,"1":64},score_workers=16,score_buffer=128)
135
+ receipt.pop('concurrency',None)
136
+ new['execution_transition']={'generation_concurrency_by_gpu':{'0':64,'1':64},'score_workers':16,'score_buffer':128,'old_source':plan['source']}
137
+
138
+ new['migration']={'from_plan':str(old),'sha256':oldhash,'reason':('Decouple generation/scoring; preserve journals and clock' if update_bytes is not None else 'Add eight SimCT checkpoints; preserve journals and original clock')}
139
+ E.write_new(out/'plan.json',new);newhash=E.file_hash(out/'plan.json')
140
+ if (old.parent/'cells').exists(): shutil.copytree(old.parent/'cells',out/'cells')
141
+ jobs={j['id']:j for j in plan['jobs']};datasets={b:{x['id']:x for x in E.read_json(d['path'])['items']} for b,d in plan['data'].items()}
142
+ totals={'responses':0,'scores':0,'metrics':0}
143
+ for f in (old.parent/'cells').rglob('*'):
144
+ if f.is_file(): receipt['original_files'][str(f.relative_to(old.parent))]=E.file_hash(f)
145
+ for cell in (out/'cells').glob('*/*/*'):
146
+ if not cell.is_dir(): continue
147
+ jid,b,seed=cell.relative_to(out/'cells').parts;seed=int(seed)
148
+ if jid not in jobs or b not in datasets or seed not in plan['seeds']: raise ValueError('Unknown cell')
149
+ contract=E.read_json(cell/'contract.json')
150
+ expected=Q.cell_contract(oldhash,jobs[jid],plan['data'],b,seed,contract['server'])
151
+ if contract!=expected: raise ValueError('Old cell contract mismatch')
152
+ metrics=E.read_json(cell/'metrics.json') if (cell/'metrics.json').exists() else None
153
+ if metrics:
154
+ if metrics['contract']!=contract or metrics['status']!='completed': raise ValueError('Old metrics mismatch')
155
+ for name in ('responses','scores'):
156
+ if E.file_hash(cell/(name+'.jsonl'))!=metrics[name+'_sha256']: raise ValueError('Completed journal changed')
157
+ responses=Q.journal(cell/'responses.jsonl');scores=Q.journal(cell/'scores.jsonl');items=datasets[b]
158
+ if not set(scores)<=set(responses)<=set(items): raise ValueError('Unknown result ID')
159
+ for key,row in responses.items():
160
+ payload=E.generation_payload('eval-gemma',items[key],b,seed)
161
+ if row['seed']!=seed or row['request_sha256']!=E.digest(E.encoded(payload)): raise ValueError('Request mismatch')
162
+ E.validate_response(row['response'])
163
+ for key,row in scores.items():
164
+ if type(row.get('passed')) is not bool or row['response_sha256']!=E.digest(E.encoded(responses[key])): raise ValueError('Score mismatch')
165
+ if metrics:
166
+ if len(scores)!=len(items) or metrics['count']!=len(items) or metrics['score']!=sum(x['passed'] for x in scores.values())/len(items): raise ValueError('Metric count/score mismatch')
167
+ metrics['contract']={**contract,'plan_sha256':newhash};Q.atomic_json(cell/'metrics.json',metrics);totals['metrics']+=1
168
+ Q.atomic_json(cell/'contract.json',{**contract,'plan_sha256':newhash})
169
+ totals['responses']+=len(responses);totals['scores']+=len(scores)
170
+ for rel,digest in receipt['original_files'].items():
171
+ if E.file_hash(old.parent/rel)!=digest: raise ValueError('Source journal changed during migration')
172
+ state['plan_sha256']=newhash
173
+ state['jobs']={k:v for k,v in state['jobs'].items() if v['status']=='completed'}
174
+ for jid in state['jobs']:
175
+ if not all((out/'cells'/jid/b/str(s)/'metrics.json').exists() for b in E.CAPS for s in plan['seeds']): raise ValueError('Completed job lacks cells')
176
+ E.write_new(out/'state.json',state)
177
+ receipt.update(status='verified',new_plan_sha256=newhash,totals=totals,deadline=state['deadline'])
178
+ Q.atomic_json(out/'migration.json',receipt)
179
+ print('MIGRATION_PASS',json.dumps(totals),flush=True);print('NEW_PLAN='+str(out/'plan.json'));print('NEW_SOURCE='+receipt['source']);print('No GPU launched. Original deadline retained.')
180
+
181
+ if __name__=='__main__': main()