Download source/scripts/fleet_status.py from andyshu/opensysone: direct link, hf CLI and curl.
- Browser
- Download file 10.3 kB
-
https://huggingface.co/andyshu/opensysone/resolve/main/source/scripts/fleet_status.py
- Command line
-
hf download hf://andyshu/opensysone/source/scripts/fleet_status.py
-
curl -L -o fleet_status.py https://huggingface.co/andyshu/opensysone/resolve/main/source/scripts/fleet_status.py
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()) | |