opensysone / source /scripts /fleet_status.py
andyshu's picture
Back up verified OpenSysOne training snapshot and pinned source
2d5c26a verified
Raw History Blame Contribute Delete
10.3 kB
"""Read-only fleet status; stdlib only, no model loading or remote installation.
python scripts/fleet_status.py [--campaign /absolute/fleet/run] [--json]
"""
import argparse
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timezone
import json
from pathlib import Path
import re
import shlex
import subprocess
import sys
sys.dont_write_bytecode = True
ROOT = Path(__file__).resolve().parents[1]
if str(ROOT) not in sys.path:
sys.path.insert(0, str(ROOT))
from scripts.fleet_campaign import SSH
QUERY_TIMEOUT_SECONDS = 15
# Execute this self-contained stdlib reader with each candidate's existing
# interpreter. It never imports project/model code or writes remote files.
READER = r'''
import json
import math
from pathlib import Path
import sys
project, campaign = map(Path, sys.argv[1:3])
warnings = []
def read_json(path, required=False):
if not path.exists():
if required:
raise ValueError('Missing ' + str(path))
return {}
return json.loads(path.read_text())
def proc_matches(pid, expected):
if not pid or not expected:
return False
try:
actual = Path(f'/proc/{pid}/cmdline').read_bytes().split(b'\0')[:-1]
except FileNotFoundError:
return False
return actual == [part.encode() for part in expected]
def jsonl(path, last_only=False):
if not path.exists():
return []
lines = path.read_text().splitlines()
indices = range(len(lines)-1,-1,-1) if last_only else range(len(lines))
rows = []
for index in indices:
if not lines[index].strip():
continue
try:
row = json.loads(lines[index])
except json.JSONDecodeError:
if index == len(lines)-1:
warnings.append(path.name + ': trailing partial record ignored')
continue
raise ValueError(f'Malformed {path.name} at line {index+1}')
rows.append(row)
if last_only:
break
return rows
def code(path):
try:
return int(path.read_text().strip())
except FileNotFoundError:
return None
def finite(value):
return isinstance(value,(int,float)) and not isinstance(value,bool) and math.isfinite(value)
state = read_json(campaign/'state.json',required=True)
training = campaign/'training'
# Fleet plans may supply a training directory distinct from the conventional one.
if len(sys.argv) > 3:
training = Path(sys.argv[3])
expected_supervisor = [sys.executable,str(project/'scripts/launch_24h.py'),'--campaign',str(campaign)]
processes = {'supervisor':proc_matches(state.get('supervisor_pid'),expected_supervisor),
'trainer':proc_matches(state.get('child_pid'),state.get('child_command')),
'api':proc_matches(state.get('api_pid'),state.get('api_command'))}
latest = jsonl(training/'training.jsonl',last_only=True)
validation = jsonl(training/'validation.jsonl')
selected = read_json(training/'best_validation_selection.json')
metric = selected.get('metric')
score = selected.get('score')
accuracy = selected.get('accuracy')
selected_step = None
if finite(score):
for row in reversed(validation):
selection = row.get('selection',{})
candidate_score = selection.get('score')
if (row.get('improved') is True and selection.get('metric') == metric
and finite(candidate_score) and abs(candidate_score-score) <= 1e-8):
selected_step = row.get('step')
if not finite(accuracy):
accuracy = row.get('metrics',{}).get('overall',{}).get('accuracy')
break
exits = {'supervisor':code(campaign/'exit_code')}
for stage in ('training','evaluation','harness_check'):
recorded = code(campaign/(stage+'_exit_code'))
if recorded is not None:
exits[stage] = recorded
result = {'status':state.get('status'),'stage':state.get('stage'),'processes':processes,
'latest_logged_step':latest[0].get('step') if latest else None,
'selected_step':selected_step,'selection_metric':metric,
'selected_cv_nll':score if metric == 'crossfit_temperature_nll_v1' and finite(score) else None,
'selected_accuracy':accuracy if finite(accuracy) else None,
'recorded_exit_codes':exits,'heartbeat_utc':state.get('heartbeat_utc'),
'api_ready':state.get('api_ready',False)}
if warnings:
result['warnings'] = warnings
print(json.dumps(result,allow_nan=False))
'''
def process_matches(pid, expected):
if not pid or not expected:
return False
try:
actual = Path(f'/proc/{pid}/cmdline').read_bytes().split(b'\0')[:-1]
except FileNotFoundError:
return False
return actual == [part.encode() for part in expected]
def query_candidate(candidate):
result = {key:candidate.get(key) for key in ('name','host','campaign')}
host = candidate.get('host','local')
result['host'] = host
try:
if not re.fullmatch(r'[a-zA-Z0-9_.@-]+',host) or host.startswith('-'):
raise ValueError('Invalid SSH host')
for field in ('project','python','campaign','training'):
if not Path(candidate[field]).is_absolute():
raise ValueError(field + ' must be an absolute path')
arguments = [candidate['python'],'-c',READER,candidate['project'],candidate['campaign'],candidate['training']]
command = arguments if host == 'local' else [*SSH,host,shlex.join(arguments)]
completed = subprocess.run(command,capture_output=True,text=True,timeout=QUERY_TIMEOUT_SECONDS)
if completed.returncode:
# Python tracebacks contain filenames, not model state or environment.
detail = completed.stderr.strip().splitlines()
raise RuntimeError(f'reader exited {completed.returncode}: ' + (detail[-1] if detail else 'no stderr'))
result.update(json.loads(completed.stdout))
except (KeyError,TypeError,ValueError,OSError,RuntimeError,subprocess.TimeoutExpired) as error:
# A failed remote cannot hide status from the other candidates. Do not
# stringify TimeoutExpired: its command contains the entire reader code.
message = f'read timed out after {QUERY_TIMEOUT_SECONDS}s' if isinstance(error,subprocess.TimeoutExpired) else str(error)
result.update(status='unavailable',error=message)
return result
def collect_status(campaign):
campaign = Path(campaign).resolve()
plan = json.loads((campaign/'plan.json').read_text())
state = json.loads((campaign/'state.json').read_text())
candidates = plan.get('candidates',[])
with ThreadPoolExecutor(max_workers=max(1,min(8,len(candidates)))) as executor:
reports = list(executor.map(query_candidate,candidates))
recorded_exit = campaign/'exit_code'
fleet = {'campaign':str(campaign),'status':state.get('status'),'stage':state.get('stage'),
'processes':{key:process_matches(state.get(key+'_pid'),state.get(key+'_command'))
for key in ('supervisor','child','api')},
'heartbeat_utc':state.get('heartbeat_utc'),'api_ready':state.get('api_ready',False),
'last_recorded_exit_code':int(recorded_exit.read_text().strip()) if recorded_exit.exists() else None,
'selection_cutoff':plan.get('selection_cutoff'),'final_deadline':plan.get('final_deadline')}
if state.get('url'):
fleet['url'] = state['url']
return {'queried_utc':datetime.now(timezone.utc).isoformat(),'fleet':fleet,'candidates':reports}
def human_report(report):
fleet = report['fleet']
alive = ','.join(key for key,value in fleet['processes'].items() if value) or 'none'
print(f"Fleet: {fleet['status']} / {fleet['stage']} | live: {alive} | API ready: {fleet['api_ready']}")
print(f"Selection: {fleet['selection_cutoff']} | final deadline: {fleet['final_deadline']}")
headings = ['candidate','host','status','logged step','selected step','CV NLL','accuracy','live','exits']
rows = []
for candidate in report['candidates']:
processes = candidate.get('processes',{})
live = ','.join(key for key,value in processes.items() if value) or '-'
score,accuracy = candidate.get('selected_cv_nll'),candidate.get('selected_accuracy')
exits = ','.join(f'{key}:{value}' for key,value in candidate.get('recorded_exit_codes',{}).items() if value is not None) or '-'
status = candidate.get('stage') if candidate.get('status') == 'running' else candidate.get('status')
rows.append([candidate['name'],candidate['host'],status or candidate.get('status') or 'unknown',
str(candidate.get('latest_logged_step') if candidate.get('latest_logged_step') is not None else '-'),
str(candidate.get('selected_step') if candidate.get('selected_step') is not None else '-'),
f'{score:.6f}' if score is not None else '-',f'{100*accuracy:.2f}%' if accuracy is not None else '-',live,exits])
widths = [max(len(str(row[i])) for row in [headings,*rows]) for i in range(len(headings))]
for row in [headings,*rows]:
print(' '.join(str(value).ljust(width) for value,width in zip(row,widths)))
for candidate in report['candidates']:
if candidate.get('error'):
print(f"{candidate['name']}: {candidate['error']}")
for warning in candidate.get('warnings',[]):
print(f"{candidate['name']}: {warning}")
print('Selected step is omitted when JSON evidence does not identify it; logged updates may be newer than checkpoints.')
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('--campaign',help='Defaults to ~/ai/opensysone/runs/LAST_FLEET_CAMPAIGN')
parser.add_argument('--json',action='store_true',help='Emit machine-readable JSON')
args = parser.parse_args()
try:
campaign = args.campaign or (Path.home()/'ai/opensysone/runs/LAST_FLEET_CAMPAIGN').read_text().strip()
report = collect_status(campaign)
except (OSError,ValueError,KeyError) as error:
print(json.dumps({'error':str(error)}),file=sys.stderr)
return 1
if args.json:
print(json.dumps(report,indent=2,allow_nan=False))
else:
human_report(report)
return 0
if __name__ == '__main__':
raise SystemExit(main())