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