"""Back up one completed fleet deployment; run under nohup with --watch to wait. Uses implicit Hugging Face authentication. This process never controls training or serving, creates a repository, changes visibility, or selects a candidate. """ import argparse from datetime import datetime, timezone import hashlib import json import logging import math import os from pathlib import Path import shutil import signal import subprocess import sys import tempfile import time ROOT = Path(__file__).resolve().parents[1] EXPORT_ROOT = Path.home() / 'ai/opensysone/exports' POLL_SECONDS = 30 MAX_ATTEMPTS = 3 REQUIRED_CHECKS = ('branch_chunks_1', 'branch_chunks_2', 'branch_chunks_4', 'question_isolation', 'candidate_permutation', 'repeat') def utc(): return datetime.now(timezone.utc).isoformat() def read_json(path): return json.loads(Path(path).read_text()) def write_json(path, data): path = Path(path) temporary = path.with_suffix(path.suffix + '.tmp') temporary.write_text(json.dumps(data, indent=2, sort_keys=True) + '\n') temporary.replace(path) def digest(path, algorithm='sha256', git_blob=False): path = Path(path) value = hashlib.new(algorithm) if git_blob: value.update(f'blob {path.stat().st_size}\0'.encode()) with path.open('rb') as handle: for chunk in iter(lambda: handle.read(1024 * 1024), b''): value.update(chunk) return value.hexdigest() def ready_artifacts(campaign): """Return only the frozen, successfully evaluated deployment, else wait/fail.""" campaign = Path(campaign).resolve() state = read_json(campaign / 'state.json') if state.get('status') != 'complete' or state.get('api_ready') is not True: return None exit_path = campaign / 'exit_code' if not exit_path.exists(): return None # Completion state precedes the supervisor's final exit write. if exit_path.read_text().strip() != '0': # Resume leaves the previous exit file until the new supervisor exits. # A genuinely failed exit remains ineligible through the watch deadline. return None selection = read_json(campaign / 'selection.json') selected = selection['selected'] deployment = read_json(campaign / 'deployment.json') evaluation = campaign / 'evaluation' metrics = read_json(evaluation / 'metrics.json') manifest = read_json(evaluation / 'manifest.json') probe = read_json(campaign / 'api_probe.json') model = evaluation / 'model.pt' checkpoint = Path(selected['checkpoint']).resolve() if not checkpoint.is_relative_to(campaign) or not selected.get('eligible'): raise ValueError('Selected checkpoint is outside this eligible fleet snapshot') if metrics.get('status') != 'complete': raise ValueError('Final evaluation is incomplete') for split in ('test', 'holdout'): if not all(key in metrics.get(split, {}) for key in ('trained', 'calibrated', 'base', 'base_calibrated', 'calibrated_difference_95pct')): raise ValueError('Final evaluation is missing required completed metrics') counts = [] for kind in ('trained', 'calibrated', 'base', 'base_calibrated'): overall = metrics[split][kind].get('overall', {}) if not isinstance(overall.get('n'), int) or overall['n'] <= 0: raise ValueError('Final evaluation has no scored decisions') counts.append(overall['n']) for key in ('accuracy', 'nll', 'brier_multiclass_sum', 'ece_top_label_10_equal_width_bins'): value = overall.get(key) if not isinstance(value, (int, float)) or not math.isfinite(value) or value < 0: raise ValueError('Final evaluation has incomplete or nonfinite metrics') if len(set(counts)) != 1: raise ValueError('Final evaluation decision counts differ') if (Path(deployment['campaign']).resolve() != campaign or Path(deployment['model']).resolve() != model or Path(deployment['selection']).resolve() != campaign / 'selection.json' or Path(deployment['evaluation']).resolve() != evaluation / 'metrics.json' or Path(state['model']).resolve() != model): raise ValueError('Deployment does not match this completed fleet') model_sha = digest(model) if not all(value == model_sha for value in (metrics.get('model_sha256'), deployment.get('sha256'), probe.get('checkpoint_sha256'))): raise ValueError('Calibrated model/deployment/probe checksum mismatch') checkpoint_sha = digest(checkpoint) if not all(value == checkpoint_sha for value in (selected.get('checkpoint_sha256'), manifest.get('checkpoint_sha256'))): raise ValueError('Selected checkpoint checksum mismatch') if manifest.get('data_signature') != selected.get('data_signature'): raise ValueError('Evaluation data signature differs from the frozen selection') checks = read_json(evaluation / 'correctness.json') for name in REQUIRED_CHECKS: value = checks.get(name + '_probability_max_abs') if not isinstance(value, (int, float)) or not math.isfinite(value) or not 0 <= value <= 1e-4: raise ValueError('Final correctness evidence is missing or failed') return {'campaign': campaign, 'selected': selected, 'model': model, 'checkpoint': checkpoint, 'model_sha256': model_sha, 'checkpoint_sha256': checkpoint_sha, 'evaluation_manifest': manifest} def verify_export(directory): manifest = read_json(directory / 'backup_manifest.json') for name, expected in manifest['files'].items(): path = directory / name if not path.is_file() or path.stat().st_size != expected['size'] or digest(path) != expected['sha256']: raise ValueError('Immutable export checksum mismatch') return manifest def build_export(artifacts, export_root=EXPORT_ROOT, source_root=ROOT): campaign = artifacts['campaign'] export_root = Path(export_root) export_root.mkdir(parents=True, exist_ok=True) destination = export_root / (campaign.name + '-final') if destination.exists(): manifest = verify_export(destination) if (manifest['model_sha256'] != artifacts['model_sha256'] or manifest['checkpoint_sha256'] != artifacts['checkpoint_sha256']): raise ValueError('Existing export belongs to a different artifact') return destination revision = subprocess.check_output(['git', 'rev-parse', 'HEAD'], cwd=source_root, text=True).strip() # Archive committed source only, with an explicit guard against weight/key files. tracked = subprocess.check_output(['git', 'ls-tree', '-r', '--name-only', revision], cwd=source_root, text=True).splitlines() for name in tracked: path = Path(name) if (path.suffix.lower() in ('.pt', '.pth', '.safetensors', '.pem', '.key') or path.name in ('.env', 'token', 'id_rsa', 'id_ed25519') or path.name.startswith('.env.')): raise ValueError('Committed source archive contains a prohibited file type') for name in ('experiment.py', 'training_model.py'): committed = subprocess.check_output(['git', 'show', revision + ':' + name], cwd=source_root) if hashlib.sha256(committed).hexdigest() != artifacts['evaluation_manifest']['source_sha256'].get(name): raise ValueError('Committed core source differs from final evaluation provenance') temporary = Path(tempfile.mkdtemp(prefix=destination.name + '.tmp-', dir=export_root)) try: sources = {'model.pt': artifacts['model'], 'selected/best.pt': artifacts['checkpoint']} for name in ('selection.json', 'plan.json', 'state.json', 'deployment.json', 'api_probe.json', 'exit_code'): sources['fleet/' + name] = campaign / name for name in ('manifest.json', 'metrics.json', 'correctness.json', 'data_filter.json'): path = campaign / 'evaluation' / name if path.exists(): sources['evaluation/' + name] = path selected = artifacts['selected'] for field, name in (('prediction_path', 'validation_predictions.json'), ('correctness_path', 'training_correctness.json')): path = Path(selected[field]).resolve() if not path.is_relative_to(campaign): raise ValueError('Selection evidence escapes the fleet snapshot') sources['selected/' + name] = path training_manifest = artifacts['checkpoint'].parent / 'manifest.json' if training_manifest.exists(): sources['selected/training_manifest.json'] = training_manifest for name, source in sources.items(): target = temporary / name target.parent.mkdir(parents=True, exist_ok=True) shutil.copyfile(source, target) subprocess.run(['git', 'archive', '--format=tar.gz', '--output=' + str(temporary / 'source.tar.gz'), revision], cwd=source_root, check=True, capture_output=True, timeout=60) if (digest(temporary / 'model.pt') != artifacts['model_sha256'] or digest(temporary / 'selected/best.pt') != artifacts['checkpoint_sha256']): raise ValueError('Artifact changed during export') files = {str(path.relative_to(temporary)): {'size': path.stat().st_size, 'sha256': digest(path)} for path in sorted(temporary.rglob('*')) if path.is_file()} write_json(temporary / 'backup_manifest.json', { 'format': 1, 'created_utc': utc(), 'campaign': str(campaign), 'source_commit': revision, 'evaluation_source_commit': artifacts['evaluation_manifest'].get('source_commit'), 'model_sha256': artifacts['model_sha256'], 'checkpoint_sha256': artifacts['checkpoint_sha256'], 'scope': 'Calibrated decision-scoring adapter/head and frozen selection evidence; base weights excluded.', 'files': files}) for path in temporary.rglob('*'): if path.is_file(): path.chmod(0o444) temporary.rename(destination) except BaseException: shutil.rmtree(temporary, ignore_errors=True) raise return destination def publish_export(api, repo_id, directory, progress, operation_factory=None): """Payload commit -> remote integrity check -> atomic pointer/model-card commit.""" manifest = verify_export(directory) prefix = 'final/' + Path(manifest['campaign']).name initial = api.repo_info(repo_id, repo_type='model', timeout=30) progress('uploading_payload', repository_private=initial.private, export=str(directory)) uploaded = api.upload_folder(repo_id=repo_id, repo_type='model', folder_path=str(directory), path_in_repo=prefix, commit_message='Back up completed OpenSysOne fleet deployment') payload_revision = uploaded.oid progress('verifying_payload', payload_commit=payload_revision) files = {**manifest['files'], 'backup_manifest.json': { 'size': (directory / 'backup_manifest.json').stat().st_size, 'sha256': digest(directory / 'backup_manifest.json')}} remote = api.get_paths_info(repo_id, [prefix + '/' + name for name in files], revision=payload_revision, repo_type='model') by_path = {item.path: item for item in remote} for name, expected in files.items(): item = by_path.get(prefix + '/' + name) if item is None or item.size != expected['size']: raise ValueError('Remote payload size mismatch') lfs = getattr(item, 'lfs', None) if name.endswith('.pt') and lfs is None: raise ValueError('Remote model is missing LFS checksum metadata') actual = (lfs.get('sha256') if isinstance(lfs, dict) else lfs.sha256) if lfs else item.blob_id expected_hash = expected['sha256'] if lfs else digest(directory / name, 'sha1', git_blob=True) if actual != expected_hash: raise ValueError('Remote payload checksum mismatch') downloaded = api.hf_hub_download(repo_id, prefix + '/backup_manifest.json', revision=payload_revision, repo_type='model', etag_timeout=30) if Path(downloaded).read_bytes() != (directory / 'backup_manifest.json').read_bytes(): raise ValueError('Remote manifest differs from immutable export') # Preserve the existing card, including its unknown/explicit license metadata. current = api.repo_info(repo_id, repo_type='model', timeout=30) if current.private != initial.private: raise ValueError('Repository visibility changed during backup') card_files = api.get_paths_info(repo_id, ['README.md'], revision=current.sha, repo_type='model') card = Path(api.hf_hub_download(repo_id, 'README.md', revision=current.sha, repo_type='model', etag_timeout=30)).read_text() if card_files else '# OpenSysOne\n' start, end = '', '' section = (f'{start}\n\n## Completed fleet deployment\n\n' '**This completed release supersedes the earlier training-stage Current status.** ' 'Descriptions of unfinished finalization and the training snapshot above are historical. ' '`FINAL_MODEL.json` now identifies the verified calibrated release.\n\n' f'Final artifact: [`{prefix}/model.pt`]({prefix}/model.pt). ' f'Integrity and provenance: [`backup_manifest.json`]({prefix}/backup_manifest.json). ' f'Completed test/holdout results: [`metrics.json`]({prefix}/evaluation/metrics.json).\n\n' 'This contains trained adapter/head parameters and calibration metadata; pinned base weights ' 'are required separately. It scores explicit candidate decisions through the Jev-compatible ' 'harness. Public benchmark metrics do not establish general intelligence or universal calibration.\n\n' f'Source revision: `{manifest["source_commit"]}`. No license metadata has been changed.\n\n{end}') if start in card and end in card: before, rest = card.split(start, 1) card = before + section + rest.split(end, 1)[1] else: card = card.rstrip() + '\n\n' + section + '\n' pointer = {'repo_id': repo_id, 'path': prefix + '/model.pt', 'campaign': manifest['campaign'], 'model_sha256': manifest['model_sha256'], 'payload_commit': payload_revision, 'manifest_path': prefix + '/backup_manifest.json', 'manifest_sha256': files['backup_manifest.json']['sha256'], 'source_commit': manifest['source_commit'], 'published_utc': utc()} if operation_factory is None: from huggingface_hub import CommitOperationAdd operation_factory = CommitOperationAdd operations = [operation_factory(path_in_repo='FINAL_MODEL.json', path_or_fileobj=(json.dumps(pointer, indent=2) + '\n').encode()), operation_factory(path_in_repo='README.md', path_or_fileobj=card.encode())] progress('publishing_pointer', payload_commit=payload_revision) published = api.create_commit(repo_id, operations=operations, repo_type='model', parent_commit=current.sha, commit_message='Point to verified completed OpenSysOne model') for name, expected in (('FINAL_MODEL.json', operations[0].path_or_fileobj), ('README.md', card.encode())): path = api.hf_hub_download(repo_id, name, revision=published.oid, repo_type='model', etag_timeout=30) if Path(path).read_bytes() != expected: raise ValueError('Published pointer or model card verification failed') return {**pointer, 'pointer_commit': published.oid, 'repository_private': initial.private} def deadline_alarm(signum, frame): raise TimeoutError('Backup deadline exceeded') def stop_requested(signum, frame): raise KeyboardInterrupt('Backup stopped') def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('--campaign', required=True) parser.add_argument('--repo-id', required=True) parser.add_argument('--output', required=True, help='New directory for backup-only status and logs') parser.add_argument('--watch', action='store_true', help='Wait in this process; use nohup to detach') args = parser.parse_args() output = Path(args.output).resolve() output.mkdir(parents=True, exist_ok=False) campaign = Path(args.campaign).resolve() state = {'status': 'running', 'stage': 'starting', 'pid': os.getpid(), 'command': [sys.executable, *sys.argv], 'campaign': str(campaign), 'repo_id': args.repo_id, 'started_utc': utc()} def progress(stage, **values): state.update(stage=stage, heartbeat_utc=utc(), **values) write_json(output / 'state.json', state) code = 1 try: Path('/proc/self/oom_score_adj').write_text('0') plan = read_json(campaign / 'plan.json') deadline = datetime.fromisoformat(plan['final_deadline'].replace('Z', '+00:00')).timestamp() + 1800 progress('waiting_for_completion', watchdog_deadline_utc=datetime.fromtimestamp(deadline, timezone.utc).isoformat()) signal.signal(signal.SIGALRM, deadline_alarm) signal.signal(signal.SIGTERM, stop_requested) signal.alarm(max(1, math.ceil(deadline - time.time()))) while time.time() < deadline: artifacts = ready_artifacts(campaign) if artifacts is not None: break if not args.watch: raise ValueError('Fleet is not ready for final backup') progress('waiting_for_completion') time.sleep(min(POLL_SECONDS, max(0, deadline - time.time()))) else: raise TimeoutError('Fleet did not complete before backup deadline') progress('creating_immutable_export') directory = build_export(artifacts) os.environ.update(HF_HUB_DISABLE_PROGRESS_BARS='1', HF_HUB_DISABLE_XET='1', HF_HUB_DOWNLOAD_TIMEOUT='60', HF_HUB_ETAG_TIMEOUT='30') from huggingface_hub import HfApi, set_client_factory import httpx for name in ('huggingface_hub', 'httpx', 'httpcore'): logging.getLogger(name).setLevel(logging.CRITICAL) set_client_factory(lambda: httpx.Client(timeout=httpx.Timeout(60, connect=15), follow_redirects=True)) api = HfApi() for attempt in range(1, MAX_ATTEMPTS + 1): if time.time() >= deadline: raise TimeoutError('Backup deadline exceeded') progress('upload_attempt', attempt=attempt) signal.alarm(max(1, min(600, math.ceil(deadline - time.time())))) try: current_artifacts = ready_artifacts(campaign) if (current_artifacts is None or current_artifacts['model_sha256'] != artifacts['model_sha256'] or current_artifacts['checkpoint_sha256'] != artifacts['checkpoint_sha256']): raise ValueError('Completed fleet changed before backup') result = publish_export(api, args.repo_id, directory, progress) progress('complete', status='complete', finished_utc=utc(), **result) code = 0 break except Exception as error: # Raw HTTP errors may contain signed URLs or tokens. Persist only type. progress('attempt_failed', error_type=type(error).__name__) if attempt == MAX_ATTEMPTS or time.time() >= deadline: raise signal.alarm(max(1, math.ceil(deadline - time.time()))) time.sleep(min(POLL_SECONDS, max(0, deadline - time.time()))) except BaseException as error: progress('failed', status='failed', error_type=type(error).__name__, finished_utc=utc()) finally: signal.alarm(0) (output / 'exit_code').write_text(str(code) + '\n') return code if __name__ == '__main__': raise SystemExit(main())