File size: 21,154 Bytes
1a0a7fb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
"""Prepare an immutable expanded-data pilot backup; never upload or load a model.

Usage: python scripts/prepare_expansion_backup.py --pilot /path/to/completed-run
    --output /path/to/new-export --fleet-plan /path/to/fleet/plan.json

The output is accepted by scripts/publish_hf_snapshot.py. Publication is separate.
Only trusted local project checkpoints may be supplied (PyTorch pickle format).
"""
import argparse
from datetime import datetime, timezone
import hashlib
import json
from pathlib import Path
import re
import shutil
import subprocess
import sys
import tempfile

PROJECT = Path(__file__).resolve().parents[1]
sys.path.insert(0, str(PROJECT))
from scripts.fleet_campaign import verify_candidate, verify_correctness
from scripts.publish_hf_snapshot import relative_name, verify_artifacts

SPLITS = ('train', 'validation', 'calibration', 'test', 'holdout')
COMPLETE = {'completed_step_target', 'completed_epochs', 'validation_early_stop'}


def digest(path):
    result = hashlib.sha256()
    with Path(path).open('rb') as handle:
        for chunk in iter(lambda: handle.read(1024 * 1024), b''):
            result.update(chunk)
    return result.hexdigest()


def read_json(path):
    return json.loads(Path(path).read_text())


def write_json(path, value):
    Path(path).write_text(json.dumps(value, indent=2, sort_keys=True) + '\n')


def verified_dataset(path):
    value = read_json(path / 'manifest.json')
    if set(value.get('split_sha256', {})) != set(SPLITS):
        raise ValueError('Dataset must contain exactly five registered splits')
    for split, checksum in value['split_sha256'].items():
        if digest(path / (split + '.jsonl')) != checksum:
            raise ValueError('Dataset split checksum mismatch')
    return value


def verify_source(root, manifest):
    revision = manifest['git_commit']
    if not re.fullmatch('[0-9a-f]{40}', revision):
        raise ValueError('Training source revision is not exact')
    if not manifest.get('source_sha256'):
        raise ValueError('Training source hash inventory is missing')
    for name, expected in manifest['source_sha256'].items():
        relative_name(name)
        content = subprocess.check_output(['git', 'show', revision + ':' + name], cwd=root,
                                          stderr=subprocess.DEVNULL, timeout=30)
        if hashlib.sha256(content).hexdigest() != expected:
            raise ValueError('Training source does not match committed revision')
    return revision


def finite_tensors(value, torch):
    if isinstance(value, torch.Tensor):
        return value.device.type == 'cpu' and bool(torch.isfinite(value).all())
    if isinstance(value, dict):
        return all(finite_tensors(item, torch) for item in value.values())
    if isinstance(value, (list, tuple)):
        return all(finite_tensors(item, torch) for item in value)
    return True


def prepare(pilot, output, fleet_plan, repo_id='andyshu/opensysone', source_root=PROJECT):
    import torch
    torch.set_num_threads(4)
    if torch.cuda.is_initialized():
        raise ValueError('Backup must run without initialized CUDA')
    pilot, output, source_root = Path(pilot).resolve(), Path(output).resolve(), Path(source_root).resolve()
    if (pilot / 'artifacts').is_dir():
        pilot = pilot / 'artifacts'
    if output.exists() or output.is_relative_to(source_root) or output == source_root:
        raise ValueError('Output must be a new directory outside source checkout')
    if (pilot.parent / 'exit_code').read_text().strip() != '0':
        raise ValueError('Pilot must have completed with exit code zero')
    manifest, summary = read_json(pilot / 'manifest.json'), read_json(pilot / 'summary.json')
    if summary.get('status') not in COMPLETE or summary.get('final_correctness_status') != 'passed':
        raise ValueError('Pilot is not complete with final correctness gates passed')
    verify_correctness(pilot)
    revision = verify_source(source_root, manifest)
    selected_hash, checkpoint_hash = digest(pilot / 'best.pt'), digest(pilot / 'checkpoint.pt')
    if summary.get('best_sha256') != selected_hash or summary.get('checkpoint_sha256') != checkpoint_hash:
        raise ValueError('Pilot summary does not match durable checkpoints')
    best = torch.load(pilot / 'best.pt', map_location='cpu', weights_only=False)
    latest = torch.load(pilot / 'checkpoint.pt', map_location='cpu', weights_only=False)
    for saved in (best, latest):
        if saved.get('format') != 'opensysone-adapter-v1' or not saved.get('trainable_state') or 'temperature' in saved:
            raise ValueError('Expected raw trained adapter/head checkpoint')
        if saved.get('source_commit') != revision or saved.get('config') != manifest['config']:
            raise ValueError('Checkpoint source/config differs from pilot manifest')
        if saved.get('data_signature') != manifest['data_signature'] or saved.get('model_provenance') != manifest['model_provenance']:
            raise ValueError('Checkpoint data/model provenance differs from pilot manifest')
        if not finite_tensors(saved.get('trainable_state'), torch):
            raise ValueError('Checkpoint trainable tensors are nonfinite or non-CPU')
    if latest['step'] != summary['completed_steps'] or not 0 <= best['step'] <= latest['step']:
        raise ValueError('Checkpoint steps disagree with completed pilot')
    for key in ('best_validation_macro_nll', 'best_validation_selection_score', 'selection_metric'):
        if latest.get(key) != best.get(key) or summary.get(key) != best.get(key):
            raise ValueError('Resumable checkpoint does not preserve selected best metrics')
    if not {'optimizer', 'random_state', 'torch_rng', 'cuda_rng'}.issubset(latest):
        raise ValueError('Resumable checkpoint lacks optimizer or random states')
    optimizer = latest['optimizer']
    if not optimizer.get('state') or not finite_tensors(optimizer, torch):
        raise ValueError('Resumable optimizer state is missing or nonfinite')
    optimizer_steps = sorted({float(item['step']) for item in optimizer['state'].values()})
    if optimizer_steps != [float(latest['step'])] or latest['step'] <= 0:
        raise ValueError('Optimizer step does not match the completed pilot')

    dataset = Path(best['config']['dataset']).resolve()
    data = verified_dataset(dataset)
    base_root = Path(data['base_dataset']['path']).resolve()
    base = verified_dataset(base_root)
    if (data['base_dataset']['manifest_sha256'] != digest(base_root / 'manifest.json') or
            data['base_dataset']['split_sha256'] != base['split_sha256'] or
            any(data['split_sha256'][s] != base['split_sha256'][s] for s in SPLITS[1:])):
        raise ValueError('Expanded dataset parent/protected split lineage mismatch')
    def signature(value):
        return hashlib.sha256(json.dumps({'data': value['split_sha256'], 'model': best['model_provenance'],
            'implementation': manifest['source_sha256']['training_model.py'],
            'max_tokens': best['config']['max_tokens']}, sort_keys=True).encode()).hexdigest()
    expected_transition = {'kind': 'training_split_only', 'parent_dataset': str(base_root),
        'dataset': str(dataset), 'parent_data_signature': signature(base), 'data_signature': signature(data),
        'parent_manifest_sha256': digest(base_root / 'manifest.json'),
        'manifest_sha256': digest(dataset / 'manifest.json'),
        'protected_split_sha256': {s: base['split_sha256'][s] for s in SPLITS[1:]}}
    for saved in (best, latest):
        initial = saved.get('initialization', {})
        if (initial.get('kind') != 'warm_start' or initial.get('restores_optimizer') is not False or
                initial.get('restores_rng') is not False or initial.get('data_transition') != expected_transition or
                saved['data_signature'] != signature(data)):
            raise ValueError('Expanded pilot lacks verified fresh-optimizer data transition')
    if data['source_script_sha256'] != manifest['source_sha256']['scripts/prepare_expanded_data.py']:
        raise ValueError('Expanded dataset builder differs from committed training source')
    diagnostics = data['diagnostics']
    diagnostic_name = relative_name(diagnostics['path'])
    if diagnostics.get('selection_eligible') is not False or digest(dataset / diagnostic_name) != diagnostics['sha256']:
        raise ValueError('Diagnostic data checksum or selection exclusion mismatch')
    plan_hash = digest(fleet_plan)
    plan = read_json(fleet_plan)
    plan_captured_utc = datetime.now(timezone.utc).isoformat()
    reference_path = Path(plan['reference_predictions'])
    if digest(reference_path) != plan['reference_sha256']:
        raise ValueError('Fixed fleet validation reference checksum mismatch')
    reference = read_json(reference_path)
    parent_path = Path(best['initialization']['parent_checkpoint'])
    parent_hash = best['initialization']['parent_checkpoint_sha256']
    if digest(parent_path) != parent_hash:
        raise ValueError('Warm-start parent checkpoint checksum mismatch')
    parent = torch.load(parent_path, map_location='cpu', weights_only=False)
    parent_manifest = read_json(parent_path.parent / 'manifest.json')
    parent_revision = verify_source(source_root, parent_manifest)
    if (parent['source_commit'] != parent_revision or parent_revision != best['initialization']['parent_source_commit'] or
            parent['step'] != best['initialization']['parent_step'] or parent['data_signature'] != signature(base) or
            parent['model_provenance'] != best['model_provenance'] or not finite_tensors(parent['trainable_state'], torch)):
        raise ValueError('Warm-start parent provenance mismatch')
    verify_correctness(parent_path.parent)
    parent_predictions = parent_path.parent / f"validation_step_{parent['step']:06d}_predictions.json"
    parent_metrics = verify_candidate(parent, read_json(parent_predictions), reference)
    if read_json(parent_path.parent / 'best_validation_selection.json') != parent_metrics['selection']:
        raise ValueError('Parent selection evidence disagrees with durable checkpoint')
    predictions_path = None
    choices = [pilot / f"validation_step_{best['step']:06d}_predictions.json",
               pilot / 'best_validation_predictions.json']
    if best['step'] == 0:
        choices.append(pilot / 'initial_validation_predictions.json')
    for path in choices:
        if path.is_file():
            try:
                metrics = verify_candidate(best, read_json(path), reference)
            except ValueError:
                continue
            predictions_path = path
            break
    if predictions_path is None:
        raise ValueError('No matching fixed-validation evidence for durable best')
    if read_json(pilot / 'best_validation_selection.json') != metrics['selection']:
        raise ValueError('Pilot selection evidence disagrees with durable checkpoint')

    output.parent.mkdir(parents=True, exist_ok=True)
    staging = Path(tempfile.mkdtemp(prefix=output.name + '.staging-', dir=output.parent))
    files = {}
    def copy(source, name, expected=None):
        source = Path(source)
        name = relative_name(name)
        before = digest(source)
        if expected is not None and expected != before:
            raise ValueError('Source checksum differs from recorded provenance')
        target = staging / name
        target.parent.mkdir(parents=True, exist_ok=True)
        shutil.copyfile(source, target)
        if digest(source) != before or digest(target) != before:
            raise ValueError('Source changed while creating immutable backup')
        files[name] = {'bytes': target.stat().st_size, 'sha256': before, 'source_path': str(source)}
    artifact_directory = 'artifacts/expanded-gx10-4b-pilot'
    try:
        copy(fleet_plan, 'fleet/plan.json', plan_hash)
        copy(reference_path, 'fleet/reference-validation-predictions.json', plan['reference_sha256'])
        copy(Path(__file__).resolve(), 'backup-tools/prepare_expansion_backup.py')
        test_helper = PROJECT / 'tests/test_expansion_backup.py'
        if test_helper.is_file():
            copy(test_helper, 'backup-tools/test_expansion_backup.py')
        copy(pilot / 'best.pt', artifact_directory + '/best.pt', selected_hash)
        copy(pilot / 'checkpoint.pt', artifact_directory + '/checkpoint.pt', checkpoint_hash)
        prediction_name = artifact_directory + f"/validation_step_{best['step']:06d}_predictions.json"
        copy(predictions_path, prediction_name)
        for name in ('manifest.json', 'summary.json', 'correctness_initial.json', 'correctness_final.json',
                     'data_filter.json', 'best_validation_selection.json'):
            copy(pilot / name, artifact_directory + '/' + name)
        copy(pilot.parent / 'exit_code', artifact_directory + '/exit_code')
        parent_directory = 'artifacts/warm-start-parent'
        copy(parent_path, parent_directory + '/best.pt', parent_hash)
        parent_prediction_name = parent_directory + '/' + parent_predictions.name
        copy(parent_predictions, parent_prediction_name)
        for name in ('manifest.json', 'correctness_initial.json', 'data_filter.json', 'best_validation_selection.json',
                     'parent-snapshot.json', 'source-data-manifest.json'):
            copy(parent_path.parent / name, parent_directory + '/' + name)
        data_prefix = 'data/' + dataset.name
        copy(dataset / 'manifest.json', data_prefix + '/manifest.json', expected_transition['manifest_sha256'])
        copy(base_root / 'manifest.json', 'data/' + base_root.name + '-manifest.json', expected_transition['parent_manifest_sha256'])
        for split in SPLITS:
            copy(dataset / (split + '.jsonl'), data_prefix + '/' + split + '.jsonl', data['split_sha256'][split])
        copy(dataset / diagnostic_name, data_prefix + '/' + diagnostic_name, diagnostics['sha256'])
        proof_directory = dataset.with_name(dataset.name + '-proof')
        if proof_directory.is_dir():
            for name in ('tokenization-proof.json', 'data_filter.json', 'postbuild-audit.json',
                         'prepare_expanded_data.py', 'tokenize_only.py'):
                copy(proof_directory / name, data_prefix + '/proof/' + name)
        for family, info in data['sources'].items():
            if 'license_reference' in info:
                recorded = info['files'][info['license_reference']]
                relative = relative_name(recorded['local_path'])
                copy(dataset / relative, data_prefix + '/' + relative, recorded['sha256'])
        notices = ['# Dataset attribution and transformations', '',
            'This private backup retains the upstream licences recorded in the dataset manifest.',
            'No blanket licence is assigned to this mixed-source collection or to the model.',
            'The transformed five splits and diagnostics are copied exactly; raw source datasets and token caches are omitted.',
            'Source pins, raw-file hashes, URLs, normalization, sampling and grouping transformations are recorded in manifest.json.',
            'Calibration/test/holdout files are copied as opaque bytes; their predictions are never used by this helper.', '',
            '| Family | Upstream repository | Pinned revision | Recorded licence |', '| --- | --- | --- | --- |']
        notices += [f"| {family} | {info['repo']} | {info['revision']} | {info['license']} |"
                    for family, info in sorted(data['sources'].items())]
        notice_path = staging / data_prefix / 'ATTRIBUTION.md'
        notice_path.write_text('\n'.join(notices) + '\n')
        files[str(notice_path.relative_to(staging))] = {'bytes': notice_path.stat().st_size, 'sha256': digest(notice_path)}
        artifact = {'name': 'expanded-gx10-4b-pilot', 'source_training_directory': str(pilot),
            'artifact_directory': artifact_directory, 'source_commit': revision,
            'source_file_sha256': manifest['source_sha256'], 'base_model': best['model_provenance'],
            'data_signature': best['data_signature'], 'initialization': best['initialization'],
            'selected': {'path': artifact_directory + '/best.pt', 'step': best['step'], 'sha256': selected_hash,
                'predictions_path': prediction_name, 'raw_macro_nll': metrics['macro_nll'],
                'selection_metric': metrics['selection_metric'], 'selection_score': metrics['selection_score'],
                'accuracy': metrics['accuracy'], 'prediction_count': metrics['count']},
            'resumable': {'path': artifact_directory + '/checkpoint.pt', 'step': latest['step'],
                'sha256': checkpoint_hash, 'optimizer_steps': optimizer_steps,
                'paired_selected_best_step': best['step'], 'config': latest['config']}}
        parent_artifact = {'name': 'warm-start-parent', 'role': 'weights-only initialization; no optimizer carry-over',
            'source_training_directory': str(parent_path.parent), 'artifact_directory': parent_directory,
            'source_commit': parent_revision, 'source_file_sha256': parent_manifest['source_sha256'],
            'base_model': parent['model_provenance'], 'data_signature': parent['data_signature'],
            'selected': {'path': parent_directory + '/best.pt', 'step': parent['step'], 'sha256': parent_hash,
                'predictions_path': parent_prediction_name, 'raw_macro_nll': parent_metrics['macro_nll'],
                'selection_metric': parent_metrics['selection_metric'], 'selection_score': parent_metrics['selection_score'],
                'accuracy': parent_metrics['accuracy'], 'prediction_count': parent_metrics['count']}}
        backup = {'format': 'opensysone-backup-v1', 'created_utc': datetime.now(timezone.utc).isoformat(),
            'target_repository': repo_id, 'staging_directory': str(output), 'uploaded': False,
            'purpose': 'Completed expanded-training pilot with selected and resumable weights plus exact transformed dataset',
            'calibrated': False, 'reserved_predictions_accessed': False, 'gpu_initialized_during_backup': False,
            'dataset': {'path': data_prefix, 'payload_included': True, 'diagnostics_selection_eligible': False,
                        'manifest_sha256': expected_transition['manifest_sha256'], 'split_sha256': data['split_sha256']},
            'validation_reference': {'path': 'fleet/reference-validation-predictions.json',
                                     'sha256': plan['reference_sha256'], 'source_path': str(reference_path)},
            'fleet_plan': {'path': 'fleet/plan.json', 'sha256': plan_hash, 'captured_utc': plan_captured_utc},
            'artifacts': [artifact, parent_artifact], 'files': files, 'total_payload_bytes': sum(item['bytes'] for item in files.values()),
            'helper_sha256': digest(__file__),
            'restore_requirements': ['Restore the pinned base weights separately.',
                'Use the archived exact training source revision and recreate dataset paths or verify signatures after path relocation.',
                'Keep checkpoint.pt, best.pt and matching validation evidence together; checkpoint.pt includes Adam/RNG state.',
                'Each upstream dataset retains its recorded licence; no common project or dataset licence is asserted.']}
        write_json(staging / 'backup-manifest.json', backup)
        sums = {name: record['sha256'] for name, record in files.items()}
        sums['backup-manifest.json'] = digest(staging / 'backup-manifest.json')
        (staging / 'SHA256SUMS').write_text(''.join(f'{checksum}  {name}\n' for name, checksum in sorted(sums.items())))
        verify_artifacts(staging, repo_id)
        if torch.cuda.is_initialized():
            raise ValueError('Backup unexpectedly initialized CUDA')
        for path in staging.rglob('*'):
            if path.is_file():
                path.chmod(0o444)
        staging.rename(output)
        return backup
    except BaseException:
        shutil.rmtree(staging)
        raise


def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument('--pilot', required=True, help='Completed local pilot run or its artifacts directory')
    parser.add_argument('--output', required=True, help='Fresh immutable artifact directory outside the source checkout')
    parser.add_argument('--fleet-plan', required=True)
    parser.add_argument('--repo-id', default='andyshu/opensysone')
    args = parser.parse_args()
    Path('/proc/self/oom_score_adj').write_text('0')
    try:
        result = prepare(args.pilot, args.output, args.fleet_plan, args.repo_id)
    except Exception as error:
        print(json.dumps({'status': 'failed', 'error_type': type(error).__name__}), file=sys.stderr)
        return 1
    print(json.dumps({'status': 'prepared', 'output': args.output, 'payload_bytes': result['total_payload_bytes'],
                      'source_commit': result['artifacts'][0]['source_commit'], 'uploaded': False}))
    return 0


if __name__ == '__main__':
    raise SystemExit(main())