"""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())