opensysone / source /scripts /publish_hf_final.py
andyshu's picture
Back up verified OpenSysOne training snapshot and pinned source
2d5c26a verified
Raw History Blame Contribute Delete
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())