Add verified-journal SimCT queue migration helper d574482
Browse files- EXTEND-EVAL-SIMCT.md +29 -0
- extend-eval-simct-d574482.py +149 -0
EXTEND-EVAL-SIMCT.md
ADDED
|
@@ -0,0 +1,29 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
# Extend an active company eval with SimCT
|
| 2 |
+
|
| 3 |
+
Run extend-eval-simct.py using system Python 3.12. Pass --source pointing
|
| 4 |
+
to the unchanged checkout that generated --plan, --out a new sibling
|
| 5 |
+
directory, --simct the completed SimCT checkpoint root, and --worker-pids
|
| 6 |
+
the currently observed two Python worker PIDs.
|
| 7 |
+
|
| 8 |
+
The tool validates the 17-job plan, fixed new MP run names, source and data
|
| 9 |
+
hashes, SimCT completion/SFT initialization, and hashes eight SimCT checkpoints
|
| 10 |
+
before stopping the exact workers via Linux pidfds. It sends TERM to their
|
| 11 |
+
snapshotted descendants and workers, waits up to 30 seconds, then requires
|
| 12 |
+
GPU and job locks to be free. It never kills unrelated processes or launches GPUs.
|
| 13 |
+
An interrupted request may be regenerated. Original journals remain untouched.
|
| 14 |
+
|
| 15 |
+
A new 25-job plan copies journals; validates request, score and completed-metric
|
| 16 |
+
hashes; changes only their plan hash bindings; retains an audit receipt and
|
| 17 |
+
original deadline. Partial trailing writes can be repaired in the copy only.
|
| 18 |
+
Any validation failure leaves the output marked incomplete; do not launch it.
|
| 19 |
+
Only status=verified in migration.json permits launching worker --concurrency 16.
|
| 20 |
+
Use the ORIGINAL source checkout for new workers, because source hashes and
|
| 21 |
+
server identities must remain unchanged. Do not use recover-startup.
|
| 22 |
+
|
| 23 |
+
Tests: synthetic one-cell migration verifies 25 jobs, retained deadline, exact
|
| 24 |
+
response bytes and unchanged original cells; unrelated PID is rejected.
|
| 25 |
+
These tests do not simulate production SGLang teardown or establish numerical
|
| 26 |
+
identity across concurrency settings. Per-item seeds/decoding parameters stay
|
| 27 |
+
fixed, but batch scheduling can affect floating-point execution. Record the
|
| 28 |
+
concurrency switch in the migration receipt. No guarantee that all jobs fit
|
| 29 |
+
the remaining original deadline.
|
extend-eval-simct-d574482.py
ADDED
|
@@ -0,0 +1,149 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 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,required=True)
|
| 69 |
+
a.add_argument('--simct',type=Path,required=True);args=a.parse_args()
|
| 70 |
+
old=args.plan.resolve(strict=True); source=args.source.resolve(strict=True); out=args.out.resolve()
|
| 71 |
+
if out.exists() or out.is_relative_to(old.parent): raise ValueError('Output must be a new sibling directory')
|
| 72 |
+
sys.path.insert(0,str(source/'scripts/evaluation'))
|
| 73 |
+
import eval_queue as Q
|
| 74 |
+
E,D=Q.E,Q.D
|
| 75 |
+
plan=E.read_json(old); oldhash=E.file_hash(old)
|
| 76 |
+
if plan['source']!=D.script_hashes(): raise ValueError('Source differs from running plan')
|
| 77 |
+
if len(plan['jobs'])!=17 or any(j['mode'] not in ('sft','atomic','fixed') for j in plan['jobs']): raise ValueError('Unexpected old jobs')
|
| 78 |
+
for j in plan['jobs']:
|
| 79 |
+
if j['mode'] in D.RUNS and D.RUNS[j['mode']] not in Path(j['checkpoint']['path']).parts: raise ValueError('Wrong historical MP run')
|
| 80 |
+
summary=E.read_json(args.simct/'run-summary.json')
|
| 81 |
+
if summary['kd_algorithm']!='span_ctkd' or summary['status']!='completed' or summary['optimizer_updates']!=312: raise ValueError('SimCT summary invalid')
|
| 82 |
+
if summary['student']!=next(j['checkpoint']['path'] for j in plan['jobs'] if j['mode']=='sft'): raise ValueError('Different SFT initialization')
|
| 83 |
+
# Hash before stopping to avoid wasting idle GPU time on checkpoint inventory.
|
| 84 |
+
extra=[]
|
| 85 |
+
for step in (312,156,80,240,40,200,120,280):
|
| 86 |
+
print('HASH_SIMCT',step,flush=True)
|
| 87 |
+
identity=E.checkpoint_identity(args.simct/f'step{step}')
|
| 88 |
+
tier=next(j['tier'] for j in plan['jobs'] if j['mode']=='atomic' and j['step']==step)
|
| 89 |
+
extra.append(dict(id=f'simct-{step}',mode='simct',step=step,tier=tier,checkpoint=identity))
|
| 90 |
+
for data in plan['data'].values():
|
| 91 |
+
if E.file_hash(data['path'])!=data['sha256']: raise ValueError('Data changed')
|
| 92 |
+
stop_workers(old,source/'scripts/evaluation/eval_queue.py',args.worker_pids)
|
| 93 |
+
with contextlib.ExitStack() as stack:
|
| 94 |
+
for gpu in (0,1):
|
| 95 |
+
if stack.enter_context(Q.locked(Path(f'/tmp/simct-eval-gpu{gpu}.lock'),blocking=False)) is None: raise ValueError('GPU worker still active')
|
| 96 |
+
stack.enter_context(Q.locked(old.parent/'state.lock'))
|
| 97 |
+
for j in plan['jobs']:
|
| 98 |
+
if stack.enter_context(Q.locked(old.parent/(j['id']+'.lock'),blocking=False)) is None: raise ValueError('Job still active')
|
| 99 |
+
state=E.read_json(old.parent/'state.json')
|
| 100 |
+
if state['plan_sha256']!=oldhash or E.file_hash(old)!=oldhash: raise ValueError('Plan changed')
|
| 101 |
+
if time.time()>=state['admit_until']: raise ValueError('Original admission window expired; no implicit extension')
|
| 102 |
+
out.mkdir()
|
| 103 |
+
receipt={'from_plan':str(old),'old_plan_sha256':oldhash,'source':str(source),'concurrency':16,'original_files':{},'status':'incomplete'}
|
| 104 |
+
E.write_new(out/'migration.json',receipt)
|
| 105 |
+
new=copy.deepcopy(plan);new['jobs']+=extra;new['jobs'].sort(key=lambda j:j['tier'])
|
| 106 |
+
new['migration']={'from_plan':str(old),'sha256':oldhash,'reason':'Add eight SimCT checkpoints; preserve journals and original clock'}
|
| 107 |
+
E.write_new(out/'plan.json',new);newhash=E.file_hash(out/'plan.json')
|
| 108 |
+
if (old.parent/'cells').exists(): shutil.copytree(old.parent/'cells',out/'cells')
|
| 109 |
+
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()}
|
| 110 |
+
totals={'responses':0,'scores':0,'metrics':0}
|
| 111 |
+
for f in (old.parent/'cells').rglob('*'):
|
| 112 |
+
if f.is_file(): receipt['original_files'][str(f.relative_to(old.parent))]=E.file_hash(f)
|
| 113 |
+
for cell in (out/'cells').glob('*/*/*'):
|
| 114 |
+
if not cell.is_dir(): continue
|
| 115 |
+
jid,b,seed=cell.relative_to(out/'cells').parts;seed=int(seed)
|
| 116 |
+
if jid not in jobs or b not in datasets or seed not in plan['seeds']: raise ValueError('Unknown cell')
|
| 117 |
+
contract=E.read_json(cell/'contract.json')
|
| 118 |
+
expected=Q.cell_contract(oldhash,jobs[jid],plan['data'],b,seed,contract['server'])
|
| 119 |
+
if contract!=expected: raise ValueError('Old cell contract mismatch')
|
| 120 |
+
metrics=E.read_json(cell/'metrics.json') if (cell/'metrics.json').exists() else None
|
| 121 |
+
if metrics:
|
| 122 |
+
if metrics['contract']!=contract or metrics['status']!='completed': raise ValueError('Old metrics mismatch')
|
| 123 |
+
for name in ('responses','scores'):
|
| 124 |
+
if E.file_hash(cell/(name+'.jsonl'))!=metrics[name+'_sha256']: raise ValueError('Completed journal changed')
|
| 125 |
+
responses=Q.journal(cell/'responses.jsonl');scores=Q.journal(cell/'scores.jsonl');items=datasets[b]
|
| 126 |
+
if not set(scores)<=set(responses)<=set(items): raise ValueError('Unknown result ID')
|
| 127 |
+
for key,row in responses.items():
|
| 128 |
+
payload=E.generation_payload('eval-gemma',items[key],b,seed)
|
| 129 |
+
if row['seed']!=seed or row['request_sha256']!=E.digest(E.encoded(payload)): raise ValueError('Request mismatch')
|
| 130 |
+
E.validate_response(row['response'])
|
| 131 |
+
for key,row in scores.items():
|
| 132 |
+
if type(row.get('passed')) is not bool or row['response_sha256']!=E.digest(E.encoded(responses[key])): raise ValueError('Score mismatch')
|
| 133 |
+
if metrics:
|
| 134 |
+
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')
|
| 135 |
+
metrics['contract']={**contract,'plan_sha256':newhash};Q.atomic_json(cell/'metrics.json',metrics);totals['metrics']+=1
|
| 136 |
+
Q.atomic_json(cell/'contract.json',{**contract,'plan_sha256':newhash})
|
| 137 |
+
totals['responses']+=len(responses);totals['scores']+=len(scores)
|
| 138 |
+
for rel,digest in receipt['original_files'].items():
|
| 139 |
+
if E.file_hash(old.parent/rel)!=digest: raise ValueError('Source journal changed during migration')
|
| 140 |
+
state['plan_sha256']=newhash
|
| 141 |
+
state['jobs']={k:v for k,v in state['jobs'].items() if v['status']=='completed'}
|
| 142 |
+
for jid in state['jobs']:
|
| 143 |
+
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')
|
| 144 |
+
E.write_new(out/'state.json',state)
|
| 145 |
+
receipt.update(status='verified',new_plan_sha256=newhash,totals=totals,deadline=state['deadline'])
|
| 146 |
+
Q.atomic_json(out/'migration.json',receipt)
|
| 147 |
+
print('MIGRATION_PASS',json.dumps(totals),flush=True);print('NEW_PLAN='+str(out/'plan.json'));print('No GPU launched. Original deadline retained.')
|
| 148 |
+
|
| 149 |
+
if __name__=='__main__': main()
|