/
ump-team
/
ump-infra
Обзор
Документация
Войти
/
ump-team
/
ump-infra
Код
Запросы
1
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
dev/test
scripts/checkpoint.py
264 строки
10 KB
Dmitry Kochenov
v0.0.5: рефакторинг валидации фаз и улучшение PR-логики
19 июл 2026, 23:40
19 июл 2026, 23:40
a686387
Код
Авторство
О чём код?
"""checkpoint — LangGraph-style checkpoints для UMP v0.0.3. Реализация P0-5 из аудиторского отчёта v0.0.2: - save_checkpoint(step_id, phase, state) — сохранить snapshot после фазы - load_checkpoint(step_id, phase) — загрузить snapshot - list_checkpoints(step_id) — список доступных checkpoint'ов - latest_checkpoint(step_id) — последний сохранённый checkpoint - restore_from_checkpoint(step_id) — восстановить progress из последнего CP Хранение: `_meta/checkpoints/step-{step_id}-{phase}.json` (HMAC-signed). Каждый checkpoint — атомарная запись (temp + os.replace), подписан HMAC-SHA256 для защиты от подделки (как immutable-snapshot в protect_files). Интеграция с orchestrate_step.py: - set_phase(status='completed') → автоматически save_checkpoint - init_progress_file → restore_from_checkpoint (если есть checkpoint для этого step_id, предлагаем resume) См. `.agent/runtime/checkpoints-protocol.md` — дизайн-документ. """ from __future__ import annotations import contextlib import hashlib import hmac import json import os import sys import tempfile from datetime import datetime from pathlib import Path # Единая точка правды для BASE-резолвинга. try: from ump.paths import detect_base except ImportError: _here = Path(__file__).resolve().parent sys.path.insert(0, str(_here.parent)) from ump.paths import detect_base # type: ignore[no-redef] BASE = detect_base(__file__) CHECKPOINT_DIR = BASE / '_meta' / 'checkpoints' def _get_hmac_key() -> bytes: """HMAC-ключ — тот же, что в protect_files._get_hmac_key. Приоритет: 1. env var UMP_PROTECT_KEY (для CI) 2. host-specific файл ~/.config/ump/protect.key 3. Авто-генерация host-specific ключа """ env_key = os.environ.get('UMP_PROTECT_KEY') if env_key: return env_key.encode('utf-8') host_key_path = Path.home() / '.config' / 'ump' / 'protect.key' if host_key_path.exists(): return host_key_path.read_bytes().strip() new_key = os.urandom(32) host_key_path.parent.mkdir(parents=True, exist_ok=True) host_key_path.write_bytes(new_key) host_key_path.chmod(0o600) return new_key def _sign_payload(payload: dict) -> str: """Подписать payload HMAC-SHA256 (как в protect_files._sign_snapshot).""" content = json.dumps(payload, sort_keys=True, ensure_ascii=False) return hmac.new(_get_hmac_key(), content.encode('utf-8'), hashlib.sha256).hexdigest() def _verify_signature(checkpoint: dict) -> bool: """Проверить подпись checkpoint'а (защита от подделки).""" sig = checkpoint.get('_signature') if not sig: return False payload = {k: v for k, v in checkpoint.items() if k != '_signature'} expected = _sign_payload(payload) return hmac.compare_digest(sig, expected) def _checkpoint_path(step_id: str, phase: str) -> Path: """Путь к файлу checkpoint'а для пары (step_id, phase). Имя файла: step-{safe_step}-{phase}.json, где safe_step заменяет . и / на - (как progress_file_for в orchestrate_step.py). """ safe_step = step_id.replace('.', '-').replace('/', '-').replace('\\', '-') safe_phase = phase.replace('/', '-') return CHECKPOINT_DIR / f'step-{safe_step}-{safe_phase}.json' def save_checkpoint( step_id: str, phase: str, state: dict, *, extra: dict | None = None, ) -> Path: """Сохранить checkpoint после завершения фазы. Args: step_id: идентификатор шага (например, '1.6'). phase: имя фазы (например, 'PREFLIGHT', 'PHASE1'). state: произвольный JSON-сериализуемый dict — snapshot состояния. Обычно это progress dict из orchestrate_step.save_progress. extra: опциональные доп. поля (git_branch, commit_sha, ...). Returns: Path к сохранённому файлу checkpoint'а. Side effects: - Создаёт CHECKPOINT_DIR если не существует. - Атомарная запись: temp + os.replace (no partial writes). - Подпись HMAC-SHA256 для защиты от подделки. """ CHECKPOINT_DIR.mkdir(parents=True, exist_ok=True) checkpoint: dict = { 'step_id': step_id, 'phase': phase, 'saved_at': datetime.now().isoformat(), 'state': state, } if extra: checkpoint['extra'] = extra checkpoint['_signature'] = _sign_payload( {k: v for k, v in checkpoint.items() if k != '_signature'} ) target = _checkpoint_path(step_id, phase) fd, tmp_path = tempfile.mkstemp( dir=str(CHECKPOINT_DIR), prefix='.tmp-cp-', suffix='.json' ) try: with os.fdopen(fd, 'w', encoding='utf-8') as f: json.dump(checkpoint, f, ensure_ascii=False, indent=2) f.flush() os.fsync(f.fileno()) os.replace(tmp_path, target) except Exception: with contextlib.suppress(OSError): os.unlink(tmp_path) raise return target def load_checkpoint(step_id: str, phase: str) -> dict | None: """Загрузить checkpoint для пары (step_id, phase). Returns: dict с полями step_id, phase, saved_at, state, [extra], _signature или None, если checkpoint не существует. Raises: ValueError: если подпись checkpoint'а невалидна (возможна подделка). """ path = _checkpoint_path(step_id, phase) if not path.exists(): return None try: data = json.loads(path.read_text(encoding='utf-8')) except (OSError, json.JSONDecodeError) as e: raise ValueError(f'checkpoint corrupted: {path}: {e}') from e if not _verify_signature(data): raise ValueError( f'checkpoint signature invalid: {path}. ' f'Possible tampering — refuse to load.' ) return data def list_checkpoints(step_id: str) -> list[dict]: """Список всех checkpoint'ов для шага, отсортированных по phase-порядку. Возвращает list of dict'ов с полями step_id, phase, saved_at, path. Проверяет подпись каждого checkpoint'а; corrupted/tampered пропускаются с WARNING в stderr. """ if not CHECKPOINT_DIR.exists(): return [] safe_step = step_id.replace('.', '-').replace('/', '-').replace('\\', '-') pattern = f'step-{safe_step}-*.json' results: list[dict] = [] for path in sorted(CHECKPOINT_DIR.glob(pattern)): try: data = json.loads(path.read_text(encoding='utf-8')) if not _verify_signature(data): print( f'WARNING: checkpoint {path} — signature invalid, skipped', file=sys.stderr, ) continue results.append({ 'step_id': data['step_id'], 'phase': data['phase'], 'saved_at': data['saved_at'], 'path': str(path), }) except (OSError, json.JSONDecodeError) as e: print(f'WARNING: checkpoint {path} — corrupted: {e}', file=sys.stderr) continue return results def latest_checkpoint(step_id: str) -> dict | None: """Вернуть последний checkpoint для шага (по saved_at). Returns: dict с полями step_id, phase, saved_at, state, [extra], path или None, если checkpoint'ов нет. """ cps = list_checkpoints(step_id) if not cps: return None # list_checkpoints уже отсортирован по phase-порядку; берём последний. # Дополнительно проверяем по saved_at (на случай ручного переопределения). latest_meta = max(cps, key=lambda c: c['saved_at']) full = load_checkpoint(step_id, latest_meta['phase']) if full is None: return None full['path'] = latest_meta['path'] return full def restore_from_checkpoint(step_id: str) -> dict | None: """Восстановить progress dict из последнего checkpoint'а. Возвращает state dict (готовый к передаче в save_progress), или None если checkpoint'ов нет. Caller (init_progress_file в orchestrate_step.py) должен использовать это для resume вместо создания пустого progress. Returns: dict — состояние шага из последнего checkpoint'а, или None. """ latest = latest_checkpoint(step_id) if latest is None: return None return latest.get('state') def delete_checkpoint(step_id: str, phase: str) -> bool: """Удалить checkpoint для пары (step_id, phase). Используется при --reset-progress для полной очистки истории шага. Returns True если файл был удалён, False если не существовал. """ path = _checkpoint_path(step_id, phase) if not path.exists(): return False with contextlib.suppress(OSError): path.unlink() return True def delete_all_checkpoints(step_id: str) -> int: """Удалить все checkpoint'ы для шага. Возвращает количество удалённых.""" cps = list_checkpoints(step_id) count = 0 for cp in cps: if delete_checkpoint(step_id, cp['phase']): count += 1 return count